use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::time::interval;
#[cfg(feature = "self-modify")]
use crate::self_modify::{
memory::{AgentMemory, MemoryConfig, ModificationRecord, Outcome},
task_gen::{AnomalySeverity, DegradationSignal, TaskGenConfig, TaskGenerator},
};
use crate::self_tune::{
anomaly::{AnomalyDetector, DetectorConfig, Severity},
controller::{Controller, ControllerConfig},
telemetry_bus::{TelemetryBus, TelemetrySnapshot},
};
#[derive(Debug, Clone)]
pub struct OrchestratorConfig {
pub poll_interval: Duration,
pub recent_tasks_cap: usize,
pub task_gen_severity_threshold: Severity,
pub auto_adjust_params: bool,
pub detector: DetectorConfig,
pub controller: ControllerConfig,
#[cfg(feature = "self-modify")]
pub task_gen: TaskGenConfig,
#[cfg(feature = "self-modify")]
pub memory: MemoryConfig,
#[cfg(feature = "self-modify")]
pub gate_config: Option<crate::self_modify::gate::GateConfig>,
}
impl Default for OrchestratorConfig {
fn default() -> Self {
Self {
poll_interval: Duration::from_secs(5),
recent_tasks_cap: 100,
task_gen_severity_threshold: Severity::Warn,
auto_adjust_params: true,
detector: DetectorConfig::default(),
controller: ControllerConfig::default(),
#[cfg(feature = "self-modify")]
task_gen: TaskGenConfig::default(),
#[cfg(feature = "self-modify")]
memory: MemoryConfig::default(),
#[cfg(feature = "self-modify")]
gate_config: None,
}
}
}
#[derive(Debug, Clone, Default)]
pub struct OrchestratorStatus {
pub snapshots_processed: u64,
pub anomalies_detected: u64,
pub tasks_generated: u64,
pub param_adjustments: u64,
pub running: bool,
pub recent_task_names: Vec<String>,
}
pub struct SelfImprovementOrchestrator {
config: OrchestratorConfig,
bus: Arc<TelemetryBus>,
status: Arc<Mutex<OrchestratorStatus>>,
detector: AnomalyDetector,
controller: Controller,
#[cfg(feature = "self-modify")]
task_gen: TaskGenerator,
#[cfg(feature = "self-modify")]
memory: Arc<Mutex<AgentMemory>>,
#[cfg(feature = "self-modify")]
deployment_pipeline: Option<crate::self_modify::deployment::StagedDeploymentPipeline>,
}
impl SelfImprovementOrchestrator {
pub fn new(config: OrchestratorConfig, bus: Arc<TelemetryBus>) -> Self {
let detector = AnomalyDetector::new(config.detector.clone());
let controller = Controller::new(config.controller.clone());
#[cfg(feature = "self-modify")]
let task_gen = TaskGenerator::new(config.task_gen.clone());
#[cfg(feature = "self-modify")]
let memory = Arc::new(Mutex::new(AgentMemory::new(config.memory.clone())));
#[cfg(feature = "self-modify")]
let deployment_pipeline = config.gate_config.as_ref().map(|gc| {
crate::self_modify::deployment::StagedDeploymentPipeline::new(
crate::self_modify::gate::ValidationGate::new(gc.clone()),
)
});
Self {
config,
bus,
status: Arc::new(Mutex::new(OrchestratorStatus::default())),
detector,
controller,
#[cfg(feature = "self-modify")]
task_gen,
#[cfg(feature = "self-modify")]
memory,
#[cfg(feature = "self-modify")]
deployment_pipeline,
}
}
pub fn status_handle(&self) -> Arc<Mutex<OrchestratorStatus>> {
Arc::clone(&self.status)
}
#[cfg(feature = "self-modify")]
pub fn memory_handle(&self) -> Arc<Mutex<AgentMemory>> {
Arc::clone(&self.memory)
}
#[cfg(feature = "self-modify")]
pub fn add_deployment_target(
&mut self,
target: Box<dyn crate::self_modify::deployment::DeploymentTarget>,
) {
if let Some(ref mut pipeline) = self.deployment_pipeline {
pipeline.add_target(target);
}
}
pub async fn run(mut self) {
if let Ok(mut s) = self.status.lock() {
s.running = true;
}
let mut rx = self.bus.subscribe();
let mut poll = interval(self.config.poll_interval);
poll.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
let snap: TelemetrySnapshot = tokio::select! {
Ok(s) = rx.recv() => s,
_ = poll.tick() => {
self.bus.latest().await
}
};
self.process_snapshot(snap);
}
}
pub fn process_snapshot(&mut self, snap: TelemetrySnapshot) {
let anomalies = self.detector.observe(&snap);
if self.config.auto_adjust_params {
self.controller.observe(&snap);
}
#[cfg(feature = "self-modify")]
let mut new_tasks: Vec<String> = Vec::new();
#[cfg(not(feature = "self-modify"))]
let new_tasks: Vec<String> = Vec::new();
#[cfg(feature = "self-modify")]
for anomaly in &anomalies {
let passes_threshold = match self.config.task_gen_severity_threshold {
Severity::Info => true,
Severity::Warn => {
anomaly.severity == Severity::Warn || anomaly.severity == Severity::Critical
}
Severity::Critical => anomaly.severity == Severity::Critical,
};
if !passes_threshold {
continue;
}
let sev = match anomaly.severity {
Severity::Info => AnomalySeverity::Info,
Severity::Warn => AnomalySeverity::Warn,
Severity::Critical => AnomalySeverity::Critical,
};
let metric_name = anomaly
.message
.split(':')
.next()
.unwrap_or("unknown")
.trim()
.to_string();
let signal = DegradationSignal::Anomaly {
metric: metric_name.clone(),
severity: sev,
observed: anomaly.metric_value,
baseline: anomaly.score, };
let now_ms = snap.captured_at.elapsed().as_millis() as u64;
if let Some(task) = self
.task_gen
.generate_at(signal, std::time::Instant::now(), now_ms)
{
new_tasks.push(task.name.clone());
if let Ok(mut mem) = self.memory.lock() {
mem.record_modification(ModificationRecord {
id: task.id.clone(),
description: task.description.clone(),
affected_files: task.affected_files.clone(),
outcome: Outcome::Pending,
metric_deltas: std::collections::HashMap::new(),
notes: format!(
"Generated from anomaly: metric={} severity={:?}",
metric_name, anomaly.severity
),
timestamp_ms: now_ms,
});
}
}
}
#[cfg(feature = "self-modify")]
if let Some(ref pipeline) = self.deployment_pipeline {
if !new_tasks.is_empty() {
use crate::self_modify::deployment::{CargoCheckRunner, ParamChange};
let runner = CargoCheckRunner::new(".");
let changes: Vec<ParamChange> = new_tasks
.iter()
.map(|name| ParamChange {
param_name: name.clone(),
old_value: 0.0,
new_value: 1.0,
})
.collect();
let outcome =
pipeline.deploy(format!("auto-{}", new_tasks.len()), &runner, &changes);
match &outcome {
crate::self_modify::deployment::DeploymentOutcome::Deployed {
changes_applied,
..
} => {
tracing::info!(
target: "self_tune::orchestrator",
changes = changes_applied,
"Auto-deployment succeeded"
);
}
crate::self_modify::deployment::DeploymentOutcome::Rejected {
failed_checks,
} => {
tracing::warn!(
target: "self_tune::orchestrator",
checks = ?failed_checks,
"Auto-deployment rejected by validation gate"
);
}
crate::self_modify::deployment::DeploymentOutcome::AwaitingReview {
change_id,
} => {
tracing::info!(
target: "self_tune::orchestrator",
change_id = %change_id,
"Auto-deployment awaiting human review"
);
}
other => {
tracing::debug!(
target: "self_tune::orchestrator",
outcome = ?other,
"Auto-deployment outcome"
);
}
}
}
}
let adj_count = self.controller.audit_log().len() as u64;
if let Ok(mut s) = self.status.lock() {
s.snapshots_processed += 1;
s.anomalies_detected += anomalies.len() as u64;
s.tasks_generated += new_tasks.len() as u64;
s.param_adjustments = adj_count;
for name in new_tasks {
if s.recent_task_names.len() >= self.config.recent_tasks_cap {
s.recent_task_names.remove(0);
}
s.recent_task_names.push(name);
}
}
}
pub fn get_param(&self, param: crate::self_tune::controller::Param) -> f64 {
self.controller.get(param)
}
pub fn status_snapshot(&self) -> OrchestratorStatus {
self.status.lock().map(|s| s.clone()).unwrap_or_default()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::self_tune::telemetry_bus::{BusConfig, TelemetrySnapshot};
fn make_bus() -> Arc<TelemetryBus> {
Arc::new(TelemetryBus::new(BusConfig {
emit_interval: Duration::from_secs(60), queue_capacity: 100,
}))
}
fn make_orc() -> SelfImprovementOrchestrator {
SelfImprovementOrchestrator::new(OrchestratorConfig::default(), make_bus())
}
fn snap_with_latency(p95_us: u64) -> TelemetrySnapshot {
let mut s = TelemetrySnapshot::zero();
s.p95_1m_us = p95_us as f64;
s.avg_latency_us = (p95_us / 2) as f64;
s
}
#[test]
fn test_new_orchestrator_starts_with_zero_status() {
let orc = make_orc();
let status = orc.status_snapshot();
assert_eq!(status.snapshots_processed, 0);
assert_eq!(status.anomalies_detected, 0);
assert_eq!(status.tasks_generated, 0);
assert!(!status.running);
}
#[test]
fn test_status_handle_is_shared() {
let orc = make_orc();
let h1 = orc.status_handle();
let h2 = orc.status_handle();
assert!(Arc::ptr_eq(&h1, &h2));
}
#[cfg(feature = "self-modify")]
#[test]
fn test_memory_handle_is_shared() {
let orc = make_orc();
let h1 = orc.memory_handle();
let h2 = orc.memory_handle();
assert!(Arc::ptr_eq(&h1, &h2));
}
#[test]
fn test_process_snapshot_increments_count() {
let mut orc = make_orc();
orc.process_snapshot(TelemetrySnapshot::zero());
orc.process_snapshot(TelemetrySnapshot::zero());
let s = orc.status_snapshot();
assert_eq!(s.snapshots_processed, 2);
}
#[test]
fn test_process_snapshot_counts_snapshots() {
let mut orc = make_orc();
for _ in 0..20 {
orc.process_snapshot(snap_with_latency(5_000));
}
let s = orc.status_snapshot();
assert_eq!(s.snapshots_processed, 20);
}
#[test]
fn test_process_snapshot_spike_triggers_anomaly_detection() {
let mut orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
task_gen_severity_threshold: Severity::Info,
..OrchestratorConfig::default()
},
make_bus(),
);
for i in 0..30 {
let v = if i % 2 == 0 { 4_000 } else { 6_000 };
orc.process_snapshot(snap_with_latency(v));
}
for _ in 0..5 {
orc.process_snapshot(snap_with_latency(500_000)); }
let s = orc.status_snapshot();
assert!(
s.anomalies_detected > 0,
"spike should trigger anomaly detection"
);
}
#[test]
fn test_process_snapshot_auto_adjust_records_adjustments() {
let mut orc = make_orc();
for _ in 0..5 {
let mut snap = TelemetrySnapshot::zero();
snap.drop_rate = 0.5; orc.process_snapshot(snap);
}
let s = orc.status_snapshot();
assert!(s.param_adjustments > 0); }
#[test]
fn test_process_snapshot_auto_adjust_disabled_no_controller_effect() {
let mut orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
auto_adjust_params: false,
..OrchestratorConfig::default()
},
make_bus(),
);
for _ in 0..5 {
let mut snap = TelemetrySnapshot::zero();
snap.drop_rate = 0.9;
orc.process_snapshot(snap);
}
assert_eq!(orc.controller.audit_log().len(), 0);
}
#[test]
fn test_info_threshold_gates_task_gen_below_warn() {
let mut orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
task_gen_severity_threshold: Severity::Info,
..OrchestratorConfig::default()
},
make_bus(),
);
for i in 0..30 {
orc.process_snapshot(snap_with_latency(if i % 2 == 0 { 4_000 } else { 6_000 }));
}
for _ in 0..5 {
orc.process_snapshot(snap_with_latency(500_000));
}
let s = orc.status_snapshot();
assert!(s.tasks_generated <= s.anomalies_detected);
}
#[cfg(feature = "self-modify")]
#[test]
fn test_deployment_pipeline_wired_when_gate_config_set() {
use crate::self_modify::gate::GateConfig;
let config = OrchestratorConfig {
gate_config: Some(GateConfig::default()),
..OrchestratorConfig::default()
};
let orc = SelfImprovementOrchestrator::new(config, make_bus());
assert!(orc.deployment_pipeline.is_some());
}
#[cfg(feature = "self-modify")]
#[test]
fn test_deployment_pipeline_absent_when_gate_config_none() {
let orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
gate_config: None,
..OrchestratorConfig::default()
},
make_bus(),
);
assert!(orc.deployment_pipeline.is_none());
}
#[cfg(feature = "self-modify")]
#[test]
fn test_add_deployment_target_noop_when_pipeline_absent() {
use crate::self_modify::deployment::InMemoryParamTarget;
let mut orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
gate_config: None,
..OrchestratorConfig::default()
},
make_bus(),
);
orc.add_deployment_target(Box::new(InMemoryParamTarget::new("t")));
}
#[cfg(feature = "self-modify")]
#[test]
fn test_default_config_gate_config_is_none() {
assert!(OrchestratorConfig::default().gate_config.is_none());
}
#[cfg(feature = "self-modify")]
#[test]
fn test_generated_tasks_recorded_in_memory() {
let mut orc = SelfImprovementOrchestrator::new(
OrchestratorConfig {
task_gen_severity_threshold: Severity::Info,
..OrchestratorConfig::default()
},
make_bus(),
);
let mem_handle = orc.memory_handle();
for i in 0..30 {
orc.process_snapshot(snap_with_latency(if i % 2 == 0 { 4_000 } else { 6_000 }));
}
for _ in 0..5 {
orc.process_snapshot(snap_with_latency(500_000));
}
let s = orc.status_snapshot();
if s.tasks_generated > 0 {
let mem = mem_handle.lock().unwrap();
let pending: Vec<_> = mem
.modifications()
.filter(|r| r.outcome == crate::self_modify::memory::Outcome::Pending)
.collect();
assert!(
!pending.is_empty(),
"generated tasks should be in memory as Pending"
);
}
}
#[test]
fn test_get_param_returns_default_value() {
let orc = make_orc();
let v = orc.get_param(crate::self_tune::controller::Param::DedupChannelBuf);
assert!(v > 0.0, "default param value should be positive");
}
#[test]
fn test_default_config_auto_adjust_enabled() {
assert!(OrchestratorConfig::default().auto_adjust_params);
}
#[test]
fn test_default_config_threshold_is_warn() {
assert_eq!(
OrchestratorConfig::default().task_gen_severity_threshold,
Severity::Warn
);
}
#[test]
fn test_default_config_poll_interval_is_5s() {
assert_eq!(
OrchestratorConfig::default().poll_interval,
Duration::from_secs(5)
);
}
}