use super::support::*;
use super::*;
#[test]
fn worker_handle_is_joined_after_run_finished() {
let temp = tempfile::TempDir::new().unwrap();
let paths = crate::config::McPaths::from_root_and_project_dir(
temp.path().join("mc"),
temp.path().join("project"),
);
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = MissionControlApp::new(
TuiSessionConfig {
config: EffectiveConfig {
provider: Some("openai".to_string()),
model: Some("test-model".to_string()),
no_color: false,
file_autocomplete_respects_gitignore: true,
custom_providers: std::collections::BTreeMap::new(),
thinking_level: crate::thinking::ThinkingLevel::Default,
auth: None,
paths,
},
settings: crate::config::Settings::default(),
appearance: crate::appearance::RuntimeAppearance::default(),
theme_cli_override: false,
initial_prompt: None,
instructions: Vec::new(),
discovered_skills: SkillDiscovery::default(),
skills: SkillDiscovery::default(),
commands: CommandRegistry::mvp(),
manager: SessionManager::new(temp.path().join("sessions")),
active_session: None,
cwd: temp.path().to_path_buf(),
herdr_reporter: None,
},
sender,
);
app.active_run = true;
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(|| {}),
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
let mut ui_state = state::MissionControlState::default();
app.join_completed_worker(&mut ui_state);
assert!(!app.active_run);
assert!(app.worker.is_none());
}
#[test]
fn panicked_tracked_workers_retire_through_shared_failure_fallback() {
for worker_kind in ["provider", "auth", "compaction"] {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = test_app(&temp, sender);
let outcome = Arc::new(WorkerOutcomeState::default());
let worker_id = outcome.worker_id();
let handle = thread::spawn(|| panic!("panic payload must not reach the UI"));
while !handle.is_finished() {
thread::yield_now();
}
app.active_run = true;
app.worker = Some(WorkerState {
handle,
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome: Arc::clone(&outcome),
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::default(),
accepts_steering: worker_kind != "compaction",
});
let mut ui_state = state::MissionControlState::default();
ui_state.start_running_prompt(format!("{worker_kind} prompt"));
ui_state.active_worker_id = Some(worker_id);
assert!(app.finish_worker_if_ready(&mut ui_state), "{worker_kind}");
assert!(app.worker.is_none());
assert!(!app.active_run);
assert!(ui_state.active_worker_id.is_none());
assert_eq!(
ui_state.status,
"worker exited unexpectedly before completion"
);
assert!(
ui_state
.transcript
.iter()
.any(|line| line == "error: worker exited unexpectedly before completion")
);
assert!(outcome.completion_delivery_failed());
assert!(matches!(
outcome.final_event(),
Some(WorkerFinalEvent::Error(message))
if message == "worker exited unexpectedly before completion"
));
}
}
#[test]
fn finished_worker_is_joined_without_run_finished_event() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = test_app(&temp, sender);
app.active_run = true;
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(|| {}),
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
let mut ui_state = state::MissionControlState::default();
assert!(app.finish_worker_if_ready(&mut ui_state));
assert!(!app.active_run);
assert!(app.worker.is_none());
}
#[test]
fn cancel_finished_worker_joins_and_clears_state() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = test_app(&temp, sender);
app.active_run = true;
let cancel = Arc::new(AtomicBool::new(false));
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(|| {}),
cancel: cancel.clone(),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
let mut ui_state = state::MissionControlState::default();
let status = app.cancel_active_worker(&mut ui_state);
assert!(status.is_none());
assert!(cancel.load(Ordering::SeqCst));
assert!(!app.active_run);
assert!(app.worker.is_none());
}
#[test]
fn cancel_running_worker_reconciles_on_subsequent_ready_check() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = test_app(&temp, sender);
app.active_run = true;
let cancel = Arc::new(AtomicBool::new(false));
let worker_cancel = cancel.clone();
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(move || {
while !worker_cancel.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
thread::sleep(Duration::from_millis(25));
}),
cancel: cancel.clone(),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
let mut ui_state = state::MissionControlState::default();
let status = app.cancel_active_worker(&mut ui_state);
assert!(status.is_some());
assert!(cancel.load(Ordering::SeqCst));
assert!(!app.active_run);
assert!(app.worker.is_some());
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
assert!(app.finish_worker_if_ready(&mut ui_state));
assert!(app.worker.is_none());
}
#[test]
fn cancel_blocked_worker_reports_blocking_call_limitation_without_joining() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(2);
let mut app = test_app(&temp, sender);
app.active_run = true;
let cancel = Arc::new(AtomicBool::new(false));
let worker_cancel = cancel.clone();
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(move || {
while !worker_cancel.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
thread::sleep(CANCEL_JOIN_GRACE_PERIOD + Duration::from_millis(25));
}),
cancel: cancel.clone(),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
let mut ui_state = state::MissionControlState::default();
let status = app.cancel_active_worker(&mut ui_state);
assert_eq!(
status.as_deref(),
Some(
"cancellation requested; provider worker is still shutting down its current blocking call"
)
);
assert!(cancel.load(Ordering::SeqCst));
assert!(!app.active_run);
assert!(app.worker.is_some());
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
assert!(app.finish_worker_if_ready(&mut ui_state));
}
#[test]
fn worker_outcome_reconciles_after_saturated_completion_channel() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(1);
sender.send(TuiEvent::Error("fill".to_string())).unwrap();
let outcome = Arc::new(WorkerOutcomeState::default());
let worker_outcome = Arc::clone(&outcome);
let worker_sender = sender.clone();
let handle = thread::spawn(move || {
send_completion(
&worker_sender,
&worker_outcome,
Some(WorkerFinalEvent::Done),
);
});
while !handle.is_finished() {
thread::yield_now();
}
assert!(outcome.completion_delivery_failed());
let mut app = test_app(&temp, sender);
app.active_run = true;
app.worker = Some(WorkerState {
handle,
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::new(),
accepts_steering: false,
});
let mut ui_state = state::MissionControlState::default();
ui_state.start_running_prompt("inspect".to_string());
ui_state.assistant_streaming = true;
assert!(app.finish_worker_if_ready(&mut ui_state));
assert_eq!(ui_state.status, "run complete");
assert!(ui_state.running_prompt.is_none());
assert!(!ui_state.assistant_streaming);
assert!(app.worker.is_none());
assert!(matches!(receiver.recv().unwrap(), TuiEvent::Error(message) if message == "fill"));
assert!(receiver.try_recv().is_err());
}
#[test]
fn completion_send_attempts_run_finished_after_saturated_final_status() {
let (sender, receiver) = bounded::<TuiEvent>(1);
sender.send(TuiEvent::Error("fill".to_string())).unwrap();
let outcome = WorkerOutcomeState::default();
send_completion(&sender, &outcome, Some(WorkerFinalEvent::Done));
assert!(matches!(receiver.recv().unwrap(), TuiEvent::Error(message) if message == "fill"));
assert!(receiver.try_recv().is_err());
}