use super::support::*;
use super::*;
#[test]
fn submit_active_prompt_queues_steering_without_provider_submission() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let running = Arc::new(AtomicBool::new(true));
let thread_running = Arc::clone(&running);
let steering = crate::agent::steering::AgentSteering::default();
app.active_run = true;
app.worker = Some(WorkerState {
handle: thread::spawn(move || {
while thread_running.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
}),
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome: Arc::new(WorkerOutcomeState::default()),
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: steering.clone(),
accepts_steering: true,
});
app.steering = steering.clone();
let mut ui_state = state::MissionControlState {
prompt_editor: state::PromptEditor::new("adjust course", "adjust course".len()),
..Default::default()
};
let submitted = app.submit(
"adjust course".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
assert!(!submitted);
assert_eq!(ui_state.status, "Steering queued");
assert_eq!(ui_state.pending_steering_count, 1);
assert!(ui_state.prompt_plain_text().is_empty());
assert_eq!(ui_state.prompt_cursor(), 0);
assert_eq!(steering.pending_count(), 1);
assert_eq!(
steering.drain_collapsed(),
Some(
"Steering update from user while current run was active:\n\nadjust course".to_string()
)
);
assert!(receiver.try_recv().is_err());
running.store(false, Ordering::SeqCst);
let worker = app.worker.take().unwrap();
worker.handle.join().unwrap();
}
#[test]
fn bash_mode_active_submit_preserves_input_without_queueing_steering() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let running = Arc::new(AtomicBool::new(true));
let thread_running = Arc::clone(&running);
let steering = crate::agent::steering::AgentSteering::default();
app.active_run = true;
app.worker = Some(WorkerState {
handle: thread::spawn(move || {
while thread_running.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
}),
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome: Arc::new(WorkerOutcomeState::default()),
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: steering.clone(),
accepts_steering: false,
});
app.steering = steering.clone();
let mut ui_state = state::MissionControlState {
prompt_editor: state::PromptEditor::new("keep this draft", "keep this draft".len()),
..Default::default()
};
let submitted = app.submit(
"keep this draft".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
assert!(!submitted);
assert_eq!(
ui_state.status,
"active run cannot accept steering; wait for it to finish"
);
assert_eq!(ui_state.prompt_plain_text(), "keep this draft");
assert_eq!(ui_state.prompt_cursor(), "keep this draft".len());
assert_eq!(ui_state.pending_steering_count, 0);
assert_eq!(steering.pending_count(), 0);
assert!(receiver.try_recv().is_err());
running.store(false, Ordering::SeqCst);
let worker = app.worker.take().unwrap();
worker.handle.join().unwrap();
}
#[test]
fn cancel_active_worker_clears_pending_steering_for_finished_worker() {
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 steering = crate::agent::steering::AgentSteering::default();
steering.try_enqueue("stale steering".to_string()).unwrap();
app.steering = steering.clone();
app.last_steering_pending_count = 1;
let outcome = Arc::new(WorkerOutcomeState::default());
outcome.mark_completion_delivery_failed();
app.worker = Some(WorkerState {
handle: thread::spawn(|| {}),
cancel,
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: steering.clone(),
accepts_steering: true,
});
while !app.worker.as_ref().unwrap().handle.is_finished() {
thread::yield_now();
}
let mut ui_state = state::MissionControlState {
pending_steering_count: 1,
..Default::default()
};
let status = app.cancel_active_worker(&mut ui_state);
assert!(status.is_none());
assert_eq!(steering.pending_count(), 0);
assert_eq!(ui_state.pending_steering_count, 0);
assert_eq!(app.last_steering_pending_count, 0);
assert!(app.worker.is_none());
}
#[test]
fn cancel_active_worker_returns_before_join_grace_while_worker_shuts_down() {
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 steering = crate::agent::steering::AgentSteering::default();
steering.try_enqueue("stale steering".to_string()).unwrap();
app.steering = steering.clone();
app.last_steering_pending_count = 1;
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,
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: steering.clone(),
accepts_steering: true,
});
let mut ui_state = state::MissionControlState {
pending_steering_count: 1,
..Default::default()
};
let started = Instant::now();
let status = app.cancel_active_worker(&mut ui_state);
assert!(started.elapsed() < CANCEL_JOIN_GRACE_PERIOD);
assert!(status.is_some());
assert_eq!(steering.pending_count(), 0);
assert_eq!(ui_state.pending_steering_count, 0);
assert_eq!(app.last_steering_pending_count, 0);
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_eq!(steering.pending_count(), 0);
assert_eq!(ui_state.pending_steering_count, 0);
}
#[test]
fn cancel_running_prompt_sets_visible_status_before_worker_reconciliation() {
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,
login_manual: None,
outcome,
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: crate::agent::steering::AgentSteering::default(),
accepts_steering: true,
});
let mut ui_state = state::MissionControlState::default();
ui_state.start_running_prompt("inspect repo".to_string());
let started = Instant::now();
app.cancel_running_prompt(&mut ui_state);
assert!(started.elapsed() < CANCEL_JOIN_GRACE_PERIOD);
assert_eq!(ui_state.status, "canceling: inspect repo");
assert!(matches!(
ui_state.running_prompt.as_ref().map(|prompt| prompt.status),
Some(state::RunningPromptStatus::Canceling)
));
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 submit_active_prompt_reports_full_steering_queue_and_restores_prompt() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let running = Arc::new(AtomicBool::new(true));
let thread_running = Arc::clone(&running);
let steering = crate::agent::steering::AgentSteering::default();
for index in 0..16 {
steering.try_enqueue(format!("message {index}")).unwrap();
}
app.active_run = true;
app.worker = Some(WorkerState {
handle: thread::spawn(move || {
while thread_running.load(Ordering::SeqCst) {
thread::sleep(Duration::from_millis(1));
}
}),
cancel: Arc::new(AtomicBool::new(false)),
login_manual: None,
outcome: Arc::new(WorkerOutcomeState::default()),
outcome_reconciled: false,
shutdown_policy: WorkerShutdownPolicy::Cancel,
steering: steering.clone(),
accepts_steering: true,
});
app.steering = steering.clone();
let mut ui_state = state::MissionControlState::default();
let submitted = app.submit(
"overflow".to_string(),
&mut ui_state,
&receiver,
test_area(),
);
assert!(!submitted);
assert_eq!(ui_state.status, "Steering queue full (16 max)");
assert_eq!(ui_state.prompt_plain_text(), "overflow");
assert_eq!(ui_state.prompt_cursor(), "overflow".len());
assert_eq!(steering.pending_count(), 16);
assert!(receiver.try_recv().is_err());
running.store(false, Ordering::SeqCst);
let worker = app.worker.take().unwrap();
worker.handle.join().unwrap();
}
#[test]
fn steering_pending_drop_shows_injected_feedback_and_clears_count() {
let temp = tempfile::TempDir::new().unwrap();
let (sender, _receiver) = bounded::<TuiEvent>(4);
let mut app = test_app(&temp, sender);
let steering = crate::agent::steering::AgentSteering::default();
steering.try_enqueue("adjust".to_string()).unwrap();
app.steering = steering.clone();
app.last_steering_pending_count = 1;
let mut ui_state = state::MissionControlState {
pending_steering_count: 1,
..Default::default()
};
assert!(steering.drain_collapsed().is_some());
let changed = app.sync_steering_feedback(&mut ui_state);
assert!(changed);
assert_eq!(ui_state.pending_steering_count, 0);
assert_eq!(ui_state.status, "Steering injected");
assert!(matches!(
ui_state.toast.as_ref().map(|toast| toast.message.as_str()),
Some("Steering injected")
));
}
#[test]
fn provider_worker_gets_steering_handle_without_queuing_submit_input() {
let temp = tempfile::TempDir::new().unwrap();
let base_url = start_sse_server("tui response");
let (sender, receiver) = bounded::<TuiEvent>(8);
let mut app = test_app(&temp, sender);
let provider_id = "local-ai".to_string();
app.config.provider = Some(provider_id.clone());
app.config.model = Some("test-model".to_string());
app.config.custom_providers = std::collections::BTreeMap::from([(
provider_id.clone(),
crate::config::make_custom_provider_config("Local", &base_url, "").unwrap(),
)]);
app.config.auth = Some(crate::config::ProviderCredential::NoAuth);
app.state.config = Some(app.config.clone());
app.state.auth_state = crate::config::AuthState::Ready {
provider: provider_id,
credential: crate::config::ProviderCredential::NoAuth,
};
let mut ui_state = state::MissionControlState::default();
let submitted = app.submit("hello".to_string(), &mut ui_state, &receiver, test_area());
assert!(submitted);
let steering = app.worker.as_ref().unwrap().steering.clone();
assert_eq!(steering.pending_count(), 0);
let worker = app.worker.take().unwrap();
worker.handle.join().unwrap();
}