use crate::analyzer::{AnalysisOutcome, AnalysisProvenance, AnalysisResult};
use crate::config::OrchestratorConfig;
use crate::openspec::{Change, ProposalMetadata};
use crate::parallel::{
MergeResult, MergeResultOrigin, MergeTaskOutcome, ParallelEvent, ParallelExecutor,
};
use crate::tui::queue::DynamicQueue;
use std::future::Future;
use std::pin::Pin;
use std::process::Command;
use std::sync::atomic::Ordering;
use std::sync::Arc;
use std::time::Duration;
use tempfile::TempDir;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
const EVENT_WAIT: Duration = Duration::from_secs(5);
const NO_STOP_WINDOW: Duration = Duration::from_millis(150);
const PENDING_MERGE_NOTICE: &str = "Waiting for 1 pending background merge/base-lane task(s)";
type AnalysisFuture<'a> = Pin<Box<dyn Future<Output = AnalysisOutcome> + Send + 'a>>;
fn test_config(workspace_base: &std::path::Path) -> OrchestratorConfig {
OrchestratorConfig {
apply_command: Some("echo apply {change_id}".to_string()),
archive_command: Some("echo archive {change_id}".to_string()),
analyze_command: Some("echo analyze".to_string()),
acceptance_command: Some("echo acceptance".to_string()),
resolve_command: Some("echo resolve".to_string()),
workspace_base_dir: Some(workspace_base.to_string_lossy().to_string()),
..Default::default()
}
}
fn test_change(id: &str) -> Change {
Change {
id: id.to_string(),
completed_tasks: 0,
total_tasks: 1,
last_modified: String::new(),
dependencies: Vec::new(),
metadata: ProposalMetadata::default(),
}
}
fn init_minimal_git_repo(repo_root: &std::path::Path) {
for args in [
vec!["init", "-b", "main"],
vec!["config", "user.email", "test@example.com"],
vec!["config", "user.name", "Test User"],
] {
let output = Command::new("git")
.args(args)
.current_dir(repo_root)
.output()
.expect("run git setup command");
assert!(output.status.success(), "git setup command failed");
}
std::fs::write(repo_root.join("README.md"), "base\n").expect("write base file");
for args in [vec!["add", "-A"], vec!["commit", "-m", "Base"]] {
let output = Command::new("git")
.args(args)
.current_dir(repo_root)
.output()
.expect("run git commit command");
assert!(output.status.success(), "git commit command failed");
}
}
fn cancelling_analyzer(
cancel_token: CancellationToken,
) -> impl for<'a> Fn(&'a [Change], &'a [String], u32) -> AnalysisFuture<'a> + Send + Sync {
move |_changes: &[Change], _in_flight: &[String], _iteration: u32| -> AnalysisFuture<'_> {
cancel_token.cancel();
Box::pin(async move {
AnalysisOutcome::new(
AnalysisResult {
order: Vec::new(),
dependencies: std::collections::HashMap::new(),
groups: None,
},
AnalysisProvenance::HealthyLlm,
)
})
}
}
async fn recv_until_stopped(events: &mut mpsc::Receiver<ParallelEvent>, logs: &mut Vec<String>) {
while let Some(event) = events.recv().await {
match event {
ParallelEvent::Log(entry) => logs.push(entry.message),
ParallelEvent::Stopped => return,
ParallelEvent::AllCompleted => {
panic!("a cancelled run must not report AllCompleted; logs: {logs:?}")
}
_ => {}
}
}
panic!("scheduler event stream closed before ParallelEvent::Stopped; logs: {logs:?}");
}
async fn recv_until_pending_merge_notice(
events: &mut mpsc::Receiver<ParallelEvent>,
logs: &mut Vec<String>,
) {
if logs
.iter()
.any(|message| message.contains(PENDING_MERGE_NOTICE))
{
return;
}
while let Some(event) = events.recv().await {
match event {
ParallelEvent::Log(entry) => {
let is_notice = entry.message.contains(PENDING_MERGE_NOTICE);
logs.push(entry.message);
if is_notice {
return;
}
}
ParallelEvent::Stopped => panic!(
"terminal stop was established while the background merge was still pending, \
without ever telling the operator why stopping waits; logs: {logs:?}"
),
ParallelEvent::AllCompleted => {
panic!("a cancelled run must not report AllCompleted; logs: {logs:?}")
}
_ => {}
}
}
panic!(
"scheduler event stream closed before the pending background merge wait was announced; \
logs: {logs:?}"
);
}
#[tokio::test]
async fn idle_parallel_stop_cancelled_scheduler_releases_registered_execution_handles() {
let repo_dir = TempDir::new().unwrap();
init_minimal_git_repo(repo_dir.path());
let workspace_base = TempDir::new().expect("create workspace base");
let (event_tx, mut events) = mpsc::channel(256);
let mut executor = ParallelExecutor::new(
repo_dir.path().to_path_buf(),
test_config(workspace_base.path()),
Some(event_tx),
);
let cancel_token = CancellationToken::new();
executor.set_cancel_token(cancel_token.clone());
let queue = Arc::new(DynamicQueue::new());
executor.set_dynamic_queue(queue.clone());
queue
.register_kill_token("in-flight-change".to_string(), CancellationToken::new())
.await;
let done = queue
.request_cancellation("in-flight-change")
.await
.expect("the in-flight change must own a registered execution handle");
assert_eq!(
queue.registered_execution_count().await,
1,
"test setup must leave one registered execution handle"
);
let analyzer = cancelling_analyzer(cancel_token);
let scheduler = tokio::spawn(async move {
executor
.execute_with_order_based_reanalysis(vec![test_change("queued-a")], analyzer)
.await
});
let mut logs = Vec::new();
tokio::time::timeout(EVENT_WAIT, recv_until_stopped(&mut events, &mut logs))
.await
.expect("cancelled scheduler must reach terminal stop");
assert_eq!(
queue.registered_execution_count().await,
0,
"registered execution handles must be released before terminal stop, otherwise a later \
idle stop in the same session is misclassified as a force stop; logs: {logs:?}"
);
assert!(
done.is_cancelled(),
"releasing a handle must fire its done handshake so waiters learn the task ended"
);
scheduler
.await
.expect("scheduler task must not panic")
.expect("cancelled scheduler loop must not fail");
}
#[tokio::test]
async fn idle_parallel_stop_cancelled_scheduler_handles_pending_merge_before_stopped() {
let repo_dir = TempDir::new().unwrap();
init_minimal_git_repo(repo_dir.path());
let workspace_base = TempDir::new().expect("create workspace base");
let (event_tx, mut events) = mpsc::channel(256);
let mut executor = ParallelExecutor::new(
repo_dir.path().to_path_buf(),
test_config(workspace_base.path()),
Some(event_tx),
);
let cancel_token = CancellationToken::new();
executor.set_cancel_token(cancel_token.clone());
let (merge_result_tx, merge_result_rx) = mpsc::channel(8);
let merge_task = merge_result_tx.clone();
executor.merge_result_channel_override = Some((merge_result_tx, merge_result_rx));
let pending_merge_count = executor.pending_merge_count.clone();
pending_merge_count.store(1, Ordering::SeqCst);
let analyzer = cancelling_analyzer(cancel_token);
let mut scheduler = tokio::spawn(async move {
executor
.execute_with_order_based_reanalysis(vec![test_change("queued-a")], analyzer)
.await
});
let mut logs = Vec::new();
let premature =
tokio::time::timeout(NO_STOP_WINDOW, recv_until_stopped(&mut events, &mut logs)).await;
assert!(
premature.is_err(),
"the cancelled scheduler must not emit Stopped while a background merge is pending; \
logs: {logs:?}"
);
assert_eq!(
pending_merge_count.load(Ordering::SeqCst),
1,
"the pending merge must still be outstanding while the scheduler waits"
);
assert!(
tokio::time::timeout(Duration::from_millis(1), &mut scheduler)
.await
.is_err(),
"the scheduler loop must still be inside its cleanup barrier"
);
tokio::time::timeout(
EVENT_WAIT,
recv_until_pending_merge_notice(&mut events, &mut logs),
)
.await
.expect("the operator must see why stopping waits before terminal stop is established");
merge_task
.send(MergeResult {
change_id: "merging-change".to_string(),
workspace_name: "merging-change".to_string(),
origin: MergeResultOrigin::PostArchiveMerge,
outcome: MergeTaskOutcome::deferred("base lane busy", false),
})
.await
.expect("scheduler must still hold the merge-result receiver");
tokio::time::timeout(EVENT_WAIT, recv_until_stopped(&mut events, &mut logs))
.await
.expect("the scheduler must stop once the pending merge reported");
assert_eq!(
pending_merge_count.load(Ordering::SeqCst),
0,
"the pending merge result must be handled before terminal stop; logs: {logs:?}"
);
scheduler
.await
.expect("scheduler task must not panic")
.expect("cancelled scheduler loop must not fail");
}