use crate::agent::agent_tests::*;
use zeph_config::EnsembleConfig;
use zeph_llm::LlmError;
use zeph_llm::any::AnyProvider;
use zeph_llm::mock::MockProvider;
use zeph_orchestration::{
DagScheduler, GraphStatus, RuleBasedRouter, TaskEvent, TaskGraph, TaskNode, TaskOutcome,
TaskStatus,
};
fn running_task_graph(handle_id: &str) -> TaskGraph {
let mut graph = TaskGraph::new("ensemble verify test goal");
let mut node = TaskNode::new(0, "task-0", "produce output");
node.status = TaskStatus::Running;
node.assigned_agent = Some(handle_id.to_owned());
graph.tasks.push(node);
graph.status = GraphStatus::Running;
graph
}
fn base_orchestration_config() -> crate::config::OrchestrationConfig {
crate::config::OrchestrationConfig {
enabled: true,
verify_completeness: true,
..crate::config::OrchestrationConfig::default()
}
}
fn complete_json() -> String {
r#"{"complete": true, "gaps": [], "confidence": 0.9}"#.to_string()
}
#[cfg(feature = "scheduler")]
#[tokio::test]
async fn default_off_regression_never_touches_ensemble_metrics() {
let config = base_orchestration_config(); let graph = running_task_graph("handle-1");
let mut scheduler =
DagScheduler::resume_from(graph, &config, Box::new(RuleBasedRouter), vec![], None).unwrap();
scheduler
.event_sender()
.try_send(TaskEvent {
task_id: zeph_orchestration::TaskId(0),
agent_handle_id: "handle-1".to_owned(),
outcome: TaskOutcome::Completed {
output: "task output".into(),
artifacts: vec![],
tool_trace: None,
},
})
.unwrap();
let provider = mock_provider(vec![complete_json()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::no_tools();
let (metrics_tx, metrics_rx) = watch::channel(MetricsSnapshot::default());
let mut agent = crate::agent::Agent::new(provider, channel, registry, None, 5, executor)
.with_metrics(metrics_tx);
agent.services.orchestration.orchestration_config = config;
let token = tokio_util::sync::CancellationToken::new();
let status = agent
.run_scheduler_loop(&mut scheduler, 1, token)
.await
.unwrap();
assert_eq!(status, GraphStatus::Completed);
let snapshot = metrics_rx.borrow().clone();
assert_eq!(
snapshot.orchestration.ensemble_degraded_total, 0,
"default-off path must never increment ensemble_degraded_total"
);
assert!(
snapshot
.orchestration
.ensemble_last_agreement_ratio
.is_none(),
"default-off path must never populate ensemble_last_agreement_ratio"
);
assert!(snapshot.orchestration.ensemble_member_stats.is_empty());
}
#[cfg(feature = "scheduler")]
#[tokio::test]
async fn full_quorum_ensemble_path_updates_agreement_ratio() {
let config = crate::config::OrchestrationConfig {
ensemble: EnsembleConfig {
enabled: true,
verify: true,
members: vec!["m1".into(), "m2".into(), "m3".into()],
..EnsembleConfig::default()
},
..base_orchestration_config()
};
let graph = running_task_graph("handle-1");
let mut scheduler =
DagScheduler::resume_from(graph, &config, Box::new(RuleBasedRouter), vec![], None).unwrap();
scheduler
.event_sender()
.try_send(TaskEvent {
task_id: zeph_orchestration::TaskId(0),
agent_handle_id: "handle-1".to_owned(),
outcome: TaskOutcome::Completed {
output: "task output".into(),
artifacts: vec![],
tool_trace: None,
},
})
.unwrap();
let provider = mock_provider(vec![]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::no_tools();
let (metrics_tx, metrics_rx) = watch::channel(MetricsSnapshot::default());
let mut agent = crate::agent::Agent::new(provider, channel, registry, None, 5, executor)
.with_metrics(metrics_tx)
.with_ensemble_members(vec![
(
"m1".to_owned(),
AnyProvider::Mock(MockProvider::with_responses(vec![complete_json()])),
),
(
"m2".to_owned(),
AnyProvider::Mock(MockProvider::with_responses(vec![complete_json()])),
),
(
"m3".to_owned(),
AnyProvider::Mock(MockProvider::with_responses(vec![complete_json()])),
),
]);
agent.services.orchestration.orchestration_config = config;
let token = tokio_util::sync::CancellationToken::new();
let status = agent
.run_scheduler_loop(&mut scheduler, 1, token)
.await
.unwrap();
assert_eq!(status, GraphStatus::Completed);
let snapshot = metrics_rx.borrow().clone();
assert_eq!(
snapshot.orchestration.ensemble_degraded_total, 0,
"full-quorum merge must not be counted as degraded"
);
assert_eq!(
snapshot.orchestration.ensemble_last_agreement_ratio,
Some(1.0),
"unanimous 3-of-3 agreement must be recorded"
);
assert_eq!(snapshot.orchestration.ensemble_member_stats.len(), 3);
}
#[cfg(feature = "scheduler")]
#[tokio::test]
async fn quorum_fallback_increments_degraded_counter_and_uses_single_provider_result() {
let config = crate::config::OrchestrationConfig {
ensemble: EnsembleConfig {
enabled: true,
verify: true,
members: vec!["m1".into(), "m2".into(), "m3".into()],
..EnsembleConfig::default()
},
..base_orchestration_config()
};
let graph = running_task_graph("handle-1");
let mut scheduler =
DagScheduler::resume_from(graph, &config, Box::new(RuleBasedRouter), vec![], None).unwrap();
scheduler
.event_sender()
.try_send(TaskEvent {
task_id: zeph_orchestration::TaskId(0),
agent_handle_id: "handle-1".to_owned(),
outcome: TaskOutcome::Completed {
output: "task output".into(),
artifacts: vec![],
tool_trace: None,
},
})
.unwrap();
let provider = mock_provider(vec![complete_json()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::no_tools();
let (metrics_tx, metrics_rx) = watch::channel(MetricsSnapshot::default());
let mut agent = crate::agent::Agent::new(provider, channel, registry, None, 5, executor)
.with_metrics(metrics_tx)
.with_ensemble_members(vec![
(
"m1".to_owned(),
AnyProvider::Mock(MockProvider::with_responses(vec![complete_json()])),
),
(
"m2".to_owned(),
AnyProvider::Mock(MockProvider::default().with_errors(vec![LlmError::Unavailable])),
),
(
"m3".to_owned(),
AnyProvider::Mock(MockProvider::default().with_errors(vec![LlmError::Unavailable])),
),
]);
agent.services.orchestration.orchestration_config = config;
let token = tokio_util::sync::CancellationToken::new();
let status = agent
.run_scheduler_loop(&mut scheduler, 1, token)
.await
.unwrap();
assert_eq!(status, GraphStatus::Completed);
let snapshot = metrics_rx.borrow().clone();
assert_eq!(
snapshot.orchestration.ensemble_degraded_total, 1,
"below-quorum (1 of 3 responders) must increment the degraded counter exactly once"
);
assert!(
snapshot
.orchestration
.ensemble_last_agreement_ratio
.is_none(),
"a quorum-fallback round must not populate ensemble_last_agreement_ratio \
(that field is only set on a successful merge)"
);
}
#[cfg(feature = "scheduler")]
#[tokio::test]
async fn bootstrap_shrunk_ensemble_falls_back_before_verify_dispatch() {
use std::sync::atomic::Ordering;
let config = crate::config::OrchestrationConfig {
ensemble: EnsembleConfig {
enabled: true,
verify: true,
members: vec!["m1".into(), "m2".into(), "m3".into()],
..EnsembleConfig::default()
},
..base_orchestration_config()
};
let graph = running_task_graph("handle-1");
let mut scheduler =
DagScheduler::resume_from(graph, &config, Box::new(RuleBasedRouter), vec![], None).unwrap();
scheduler
.event_sender()
.try_send(TaskEvent {
task_id: zeph_orchestration::TaskId(0),
agent_handle_id: "handle-1".to_owned(),
outcome: TaskOutcome::Completed {
output: "task output".into(),
artifacts: vec![],
tool_trace: None,
},
})
.unwrap();
let provider = mock_provider(vec![complete_json()]);
let channel = MockChannel::new(vec![]);
let registry = create_test_registry();
let executor = MockToolExecutor::no_tools();
let (metrics_tx, metrics_rx) = watch::channel(MetricsSnapshot::default());
let (m1_provider, m1_peak) =
MockProvider::with_responses(vec![complete_json()]).with_concurrency_tracking();
let (m2_provider, m2_peak) =
MockProvider::with_responses(vec![complete_json()]).with_concurrency_tracking();
let mut agent = crate::agent::Agent::new(provider, channel, registry, None, 5, executor)
.with_metrics(metrics_tx)
.with_ensemble_members(vec![
("m1".to_owned(), AnyProvider::Mock(m1_provider)),
("m2".to_owned(), AnyProvider::Mock(m2_provider)),
]);
agent.services.orchestration.orchestration_config = config;
let token = tokio_util::sync::CancellationToken::new();
let status = agent
.run_scheduler_loop(&mut scheduler, 1, token)
.await
.unwrap();
assert_eq!(status, GraphStatus::Completed);
let snapshot = metrics_rx.borrow().clone();
assert_eq!(
snapshot.orchestration.ensemble_degraded_total, 1,
"an even-length (2) effective ensemble must trigger the S1 shrinkage-guard fallback"
);
assert!(
snapshot
.orchestration
.ensemble_last_agreement_ratio
.is_none()
);
assert_eq!(
m1_peak.load(Ordering::SeqCst),
0,
"member m1 must never be dispatched to — the S1 gate rejects the shrunk effective \
ensemble before EnsembleVerifier::verify() is ever called"
);
assert_eq!(
m2_peak.load(Ordering::SeqCst),
0,
"member m2 must never be dispatched to — the S1 gate rejects the shrunk effective \
ensemble before EnsembleVerifier::verify() is ever called"
);
}