use std::{
sync::mpsc::{self, Receiver, Sender},
time::{Duration, Instant},
};
use super::*;
const OPERATION_TIMEOUT: Duration = Duration::from_secs(1);
fn blocked_worker() -> (Sender<()>, Receiver<()>, std::thread::JoinHandle<()>) {
let (release_sender, release_receiver) = mpsc::channel();
let (done_sender, done_receiver) = mpsc::channel();
let worker = std::thread::spawn(move || {
release_receiver.recv().unwrap();
done_sender.send(()).unwrap();
});
(release_sender, done_receiver, worker)
}
fn assert_blocked_operation_completes<T>(
receiver: Receiver<T>,
operation: std::thread::JoinHandle<()>,
release_sender: Sender<()>,
worker_done: Receiver<()>,
) -> T {
let receive_result = receiver.recv_timeout(OPERATION_TIMEOUT);
let release_result = release_sender.send(());
let operation_result = operation.join();
let worker_result = worker_done.recv_timeout(OPERATION_TIMEOUT);
assert!(release_result.is_ok(), "blocked worker release failed");
assert!(operation_result.is_ok(), "operation helper panicked");
assert!(worker_result.is_ok(), "blocked worker did not exit");
receive_result.expect("operation deadlocked")
}
fn wait_for_finished_worker(worker: &std::thread::JoinHandle<()>) {
let deadline = Instant::now() + OPERATION_TIMEOUT;
while !worker.is_finished() && Instant::now() < deadline {
std::thread::yield_now();
}
assert!(
worker.is_finished(),
"title worker did not finish before the deadlock guard expired"
);
}
fn run_operation<T: Send + 'static>(
operation: impl FnOnce() -> T + Send + 'static,
) -> (Receiver<T>, std::thread::JoinHandle<()>) {
let (sender, receiver) = mpsc::channel();
let thread = std::thread::spawn(move || sender.send(operation()).unwrap());
(receiver, thread)
}
#[test]
fn print_mode_joins_already_finished_title_generation_before_returning() {
let worker = std::thread::spawn(|| {});
wait_for_finished_worker(&worker);
let mut handle = Some(worker);
finish_title_generation_for_mode(&mut handle, InvocationMode::Print);
assert!(handle.is_none());
}
#[test]
fn print_mode_detaches_active_title_generation_non_blocking() {
let (release_sender, worker_done, worker) = blocked_worker();
let (receiver, operation) = run_operation(move || {
let mut handle = Some(worker);
finish_title_generation_for_mode(&mut handle, InvocationMode::Print);
handle.is_none()
});
assert!(assert_blocked_operation_completes(
receiver,
operation,
release_sender,
worker_done,
));
}
#[test]
fn print_mode_title_guard_finish_detaches_blocked_thread_non_blocking() {
let (release_sender, worker_done, worker) = blocked_worker();
let (receiver, operation) = run_operation(move || {
let mut guard = TitleGenerationGuard::new(
Some(worker),
InvocationMode::Print,
AgentCancellation::default(),
false,
);
guard.finish();
true
});
assert!(assert_blocked_operation_completes(
receiver,
operation,
release_sender,
worker_done,
));
}
#[test]
fn print_mode_title_guard_drop_detaches_blocked_thread_non_blocking() {
let (release_sender, worker_done, worker) = blocked_worker();
let (receiver, operation) = run_operation(move || {
let _guard = TitleGenerationGuard::new(
Some(worker),
InvocationMode::Print,
AgentCancellation::default(),
false,
);
});
assert_blocked_operation_completes(receiver, operation, release_sender, worker_done);
}
#[test]
fn print_mode_title_guard_skips_join_when_canceled() {
let (release_sender, worker_done, worker) = blocked_worker();
let cancel = Arc::new(AtomicBool::new(true));
let (receiver, operation) = run_operation(move || {
let _guard = TitleGenerationGuard::new(
Some(worker),
InvocationMode::Print,
AgentCancellation::new(cancel),
false,
);
});
assert_blocked_operation_completes(receiver, operation, release_sender, worker_done);
}
#[test]
fn exclusive_writer_title_cleanup_waits_even_after_cancellation() {
for canceled in [false, true] {
let (release_sender, worker_done, worker) = blocked_worker();
let (entered_sender, entered_receiver) = mpsc::channel();
let (receiver, operation) = run_operation(move || {
let guard = TitleGenerationGuard::new(
Some(worker),
InvocationMode::Print,
AgentCancellation::new(Arc::new(AtomicBool::new(canceled))),
true,
);
entered_sender.send(()).unwrap();
drop(guard);
});
entered_receiver.recv_timeout(OPERATION_TIMEOUT).unwrap();
let before_cleanup = receiver.recv_timeout(Duration::from_millis(25));
release_sender.send(()).unwrap();
operation.join().unwrap();
worker_done.recv_timeout(OPERATION_TIMEOUT).unwrap();
assert!(matches!(
before_cleanup,
Err(mpsc::RecvTimeoutError::Timeout)
));
receiver.recv_timeout(OPERATION_TIMEOUT).unwrap();
}
}
#[test]
fn non_print_modes_leave_active_title_generation_non_blocking() {
for mode in [InvocationMode::MissionControl, InvocationMode::Subagent] {
let (release_sender, worker_done, worker) = blocked_worker();
let (receiver, operation) = run_operation(move || {
let mut handle = Some(worker);
finish_title_generation_for_mode(&mut handle, mode);
handle
});
let handle =
assert_blocked_operation_completes(receiver, operation, release_sender, worker_done)
.unwrap();
handle.join().unwrap();
}
}
#[test]
fn non_print_modes_join_already_finished_title_generation() {
for mode in [InvocationMode::MissionControl, InvocationMode::Subagent] {
let worker = std::thread::spawn(|| {});
wait_for_finished_worker(&worker);
let mut handle = Some(worker);
finish_title_generation_for_mode(&mut handle, mode);
assert!(handle.is_none());
}
}
#[test]
fn herdr_primary_lifecycle_reports_thinking_tool_done_in_order() {
let temp = tempfile::TempDir::new().unwrap();
std::fs::write(temp.path().join("file.txt"), "file text").unwrap();
let provider = ScriptedProvider::new(vec![read_done("herdr_read"), text_done("done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let reporter = herdr_owner.reporter();
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
tools: Some(&tools),
herdr_reporter: Some(reporter),
..run_request("read file", temp.path())
},
)
.unwrap();
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("working".to_string(), "running read".to_string()),
("idle".to_string(), "done".to_string()),
]
);
}
#[test]
fn herdr_subagent_invocation_reports_nothing_even_with_parent_reporter() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("child done")]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let reporter = herdr_owner.reporter();
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
tools: Some(&tools),
invocation_mode: InvocationMode::Subagent,
herdr_reporter: Some(reporter),
..run_request("child task", temp.path())
},
)
.unwrap();
assert!(lines.lock().unwrap().is_empty());
}
#[test]
fn herdr_primary_provider_failure_reports_blocked_not_incomplete() {
let temp = tempfile::TempDir::new().unwrap();
let provider = FailingProvider { events: Vec::new() };
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
herdr_reporter: Some(herdr_owner.reporter()),
..run_request("fail", temp.path())
},
)
.unwrap_err();
assert_eq!(error.to_string(), "provider broke");
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("blocked".to_string(), "needs attention".to_string()),
]
);
}
#[test]
fn herdr_primary_cancellation_before_provider_reports_cancelled() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![text_done("never")]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let cancel = Arc::new(AtomicBool::new(true));
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
cancellation: AgentCancellation::new(cancel),
herdr_reporter: Some(herdr_owner.reporter()),
..run_request("cancel", temp.path())
},
)
.unwrap_err();
assert!(is_run_canceled(&error));
assert!(provider.requests().is_empty());
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "cancelled".to_string()),
]
);
}
#[test]
fn herdr_primary_cancellation_after_provider_stream_reports_cancelled() {
let temp = tempfile::TempDir::new().unwrap();
let cancel = Arc::new(AtomicBool::new(false));
let provider = CancelingProvider::new(Arc::clone(&cancel));
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let error = agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
cancellation: AgentCancellation::new(cancel),
herdr_reporter: Some(herdr_owner.reporter()),
..run_request("cancel during stream", temp.path())
},
)
.unwrap_err();
assert!(is_run_canceled(&error));
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "cancelled".to_string()),
]
);
}
#[test]
fn herdr_primary_auto_continue_has_no_intermediate_terminal_report() {
let temp = tempfile::TempDir::new().unwrap();
let session = crate::sessions::SessionManager::new(temp.path().join("sessions"))
.create()
.unwrap();
let provider = AutoContinueProvider::new(vec![
AutoContinueStep::TextThenEligibleTimeout("partial "),
AutoContinueStep::TextThenDone("done"),
]);
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
session: Some(&session),
herdr_reporter: Some(herdr_owner.reporter()),
..run_request("recover", temp.path())
},
)
.unwrap();
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "done".to_string()),
]
);
}
#[test]
fn herdr_hook_blocked_tool_can_still_finish_done() {
let temp = tempfile::TempDir::new().unwrap();
let provider = ScriptedProvider::new(vec![
vec![write_call("blocked_write", "blocked.txt", "nope"), done()],
text_done("finished"),
]);
let tools = ToolRuntime::new(temp.path()).unwrap();
let hooks = HookRuntime::new(
temp.path(),
before_tool_hook(
"gate",
"exit 9",
Some(crate::config::HookFailurePolicy::Block),
),
true,
)
.unwrap();
let agent = AgentSession::new("model", &[], &SkillDiscovery::default());
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
agent
.run_print_with_tools_streaming_output_cancellable(
&provider,
AgentRunRequest {
tools: Some(&tools),
hooks: Some(&hooks),
herdr_reporter: Some(herdr_owner.reporter()),
..run_request("write", temp.path())
},
)
.unwrap();
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("working".to_string(), "running write".to_string()),
("idle".to_string(), "done".to_string()),
]
);
}
#[test]
fn herdr_runner_setup_failure_is_blocked_before_provider_request() {
let temp = tempfile::TempDir::new().unwrap();
let config = crate::config::EffectiveConfig {
provider: Some(crate::providers::OPENAI_CODEX_PROVIDER.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: Some(crate::config::ProviderCredential::OAuth {
access: "test-access".to_string(),
account_id: Some("test-account".to_string()),
}),
paths: crate::config::McPaths::from_root(temp.path().join("mc")),
};
let (herdr_owner, lines) = crate::herdr::HerdrOwner::new_for_test();
let error = crate::agent::runner::run_provider_once_streaming(
&config,
&[],
&SkillDiscovery::default(),
crate::agent::runner::ProviderRunOptions {
settings: Some(crate::config::Settings::default()),
prompt: "setup failure",
session: None,
cwd: temp.path(),
output_sink: None,
selected_primary_agent: None,
cancellation: None,
session_title_notifier: None,
herdr_reporter: Some(herdr_owner.reporter()),
invocation_mode: InvocationMode::Print,
disabled_tools: None,
disabled_subagent_profiles: None,
subagent_profile_discovery: None,
mcp: None,
},
)
.unwrap_err();
assert!(!error.to_string().is_empty());
assert_eq!(
herdr_statuses(&lines),
vec![
("working".to_string(), "thinking".to_string()),
("blocked".to_string(), "needs attention".to_string()),
]
);
}