#[cfg(unix)]
use super::transport::{
ConnectDisposition, classify_connect_result, poll_timeout_ms, write_line_to_socket,
};
use super::*;
use std::sync::{
Barrier, Mutex,
atomic::{AtomicBool, Ordering as AtomicOrdering},
};
#[test]
fn sequence_range_reservation_rejects_partial_cleanup_at_u64_max() {
assert_eq!(
allocate_report_seq_range_at(u64::MAX - 1, 0, 2),
None,
"release cleanup must reserve both values or emit neither"
);
assert_eq!(
allocate_report_seq_range_at(u64::MAX - 2, 0, 2),
Some((u64::MAX - 1, u64::MAX))
);
assert_eq!(allocate_report_seq_range_at(u64::MAX, 0, 1), None);
assert_eq!(allocate_report_seq_range_at(0, 0, 0), None);
}
#[cfg(unix)]
#[test]
fn poll_timeout_rounds_up_to_one_millisecond_without_unit_loss() {
assert_eq!(poll_timeout_ms(Duration::from_millis(100)), 100);
assert_eq!(poll_timeout_ms(Duration::from_millis(1)), 1);
assert_eq!(poll_timeout_ms(Duration::from_nanos(999_999)), 1);
}
#[cfg(unix)]
#[test]
fn connect_errors_are_classified_as_pending_without_retrying_connect() {
assert_eq!(
classify_connect_result(0, None),
ConnectDisposition::Connected
);
for error in [
libc::EINTR,
libc::EINPROGRESS,
libc::EALREADY,
libc::EAGAIN,
libc::EWOULDBLOCK,
] {
assert_eq!(
classify_connect_result(-1, Some(error)),
ConnectDisposition::Pending
);
}
assert_eq!(
classify_connect_result(-1, Some(libc::ECONNREFUSED)),
ConnectDisposition::Failed
);
}
#[derive(Default)]
struct MapEnv(Vec<(&'static str, OsString)>);
impl MapEnv {
fn with(mut self, key: &'static str, value: impl Into<OsString>) -> Self {
self.0.push((key, value.into()));
self
}
}
impl HerdrEnv for MapEnv {
fn get_os(&self, key: &str) -> Option<OsString> {
self.0
.iter()
.find(|(candidate, _)| *candidate == key)
.map(|(_, value)| value.clone())
}
}
#[derive(Default)]
struct RecordingTransport {
lines: Mutex<Vec<String>>,
}
impl HerdrTransport for RecordingTransport {
fn write_line(&self, _socket_path: &Path, line: &str) -> io::Result<()> {
self.lines
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.push(line.to_string());
Ok(())
}
}
struct BlockingTransport {
lines: Mutex<Vec<String>>,
entered: Arc<Barrier>,
proceed: Arc<Barrier>,
block_first: AtomicBool,
}
impl HerdrTransport for BlockingTransport {
fn write_line(&self, _socket_path: &Path, line: &str) -> io::Result<()> {
self.lines
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.push(line.to_string());
if self.block_first.swap(false, AtomicOrdering::SeqCst) {
self.entered.wait();
self.proceed.wait();
}
Ok(())
}
}
struct FixedClock(u64);
impl HerdrClock for FixedClock {
fn epoch_ms(&self) -> u64 {
self.0
}
}
struct FailingTransport;
impl HerdrTransport for FailingTransport {
fn write_line(&self, _socket_path: &Path, _line: &str) -> io::Result<()> {
Err(io::Error::new(io::ErrorKind::NotFound, "herdr unavailable"))
}
}
fn reporter_with_transport(transport: Arc<dyn HerdrTransport>) -> HerdrReporter {
HerdrOwner::new_with_transport_and_clock(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport,
Arc::new(FixedClock(1_000_000_000)),
)
.reporter()
}
fn report_value(line: &str) -> serde_json::Value {
serde_json::from_str(line).unwrap()
}
fn report_params(line: &str) -> serde_json::Value {
report_value(line)["params"].clone()
}
#[test]
fn report_agent_uses_exact_semantic_shape_without_custom_status() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report(HerdrAgentState::Blocked, Some("needs attention"));
let lines = transport.lines.lock().unwrap();
let value = report_value(&lines[0]);
assert_eq!(value["method"], "pane.report_agent");
assert_eq!(value["params"]["pane_id"], "pane-1");
assert_eq!(value["params"]["source"], SOURCE);
assert_eq!(value["params"]["agent"], AGENT);
assert_eq!(value["params"]["state"], "blocked");
assert_eq!(value["params"]["message"], "needs attention");
assert!(value["params"]["seq"].is_u64());
assert!(value["params"].get("custom_status").is_none());
assert!(value["params"].get("agent_session_path").is_none());
assert!(value["params"].get("cwd").is_none());
assert!(value["params"].get("model").is_none());
assert!(value["params"].get("provider").is_none());
}
#[test]
fn lifecycle_helpers_use_bounded_safe_messages_and_states() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_ready();
reporter.report_thinking();
reporter.report_tool("read");
reporter.report_bash();
reporter.report_done();
reporter.report_cancelled();
reporter.report_blocked();
let lines = transport.lines.lock().unwrap();
let messages = lines
.iter()
.map(|line| {
let params = report_params(line);
(
params["state"].as_str().unwrap().to_string(),
params["message"].as_str().unwrap().to_string(),
)
})
.collect::<Vec<_>>();
assert_eq!(
messages,
vec![
("idle".to_string(), "ready".to_string()),
("working".to_string(), "thinking".to_string()),
("working".to_string(), "running read".to_string()),
("working".to_string(), "running bash".to_string()),
("idle".to_string(), "done".to_string()),
("idle".to_string(), "cancelled".to_string()),
("blocked".to_string(), "needs attention".to_string()),
]
);
assert!(messages.iter().all(|(_, message)| message.len() <= 80));
}
#[test]
fn tool_name_sanitizer_stays_bounded_for_internal_labels() {
assert_eq!(
sanitize_tool_name("bad/tool name with spaces and $args"),
"badtoolnamewithspacesandargs"
);
assert_eq!(sanitize_tool_name("!!!"), "tool");
assert_eq!(
sanitize_tool_name(&"x".repeat(MAX_TOOL_NAME_CHARS + 1)).len(),
MAX_TOOL_NAME_CHARS
);
}
#[test]
fn report_agent_session_is_independent_and_sends_only_session_id() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_agent_session("session-123");
reporter.report_agent_session_with_source("session-456", HerdrSessionStartSource::Resume);
let lines = transport.lines.lock().unwrap();
let value = report_value(&lines[0]);
assert_eq!(value["method"], "pane.report_agent_session");
assert_eq!(value["params"]["pane_id"], "pane-1");
assert_eq!(value["params"]["source"], SOURCE);
assert_eq!(value["params"]["agent"], AGENT);
assert_eq!(value["params"]["agent_session_id"], "session-123");
assert!(value["params"].get("session_start_source").is_none());
assert!(value["params"].get("agent_session_path").is_none());
assert!(value["params"].get("state").is_none());
assert!(value["params"].get("message").is_none());
assert!(value["params"]["seq"].is_u64());
let sourced = report_value(&lines[1]);
assert_eq!(sourced["params"]["agent_session_id"], "session-456");
assert_eq!(sourced["params"]["session_start_source"], "resume");
}
#[test]
fn metadata_uses_schema_guards_and_summary_token_patch() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_metadata(Some("Fix auth"), Some("working on auth"));
reporter.clear_metadata();
let lines = transport.lines.lock().unwrap();
let value = report_value(&lines[0]);
assert_eq!(value["method"], "pane.report_metadata");
assert_eq!(value["params"]["pane_id"], "pane-1");
assert_eq!(value["params"]["source"], METADATA_SOURCE);
assert_eq!(value["params"]["agent"], AGENT);
assert_eq!(value["params"]["applies_to_source"], SOURCE);
assert_eq!(value["params"]["title"], "Fix auth");
assert_eq!(value["params"]["display_agent"], AGENT);
assert_eq!(value["params"]["state_labels"]["idle"], "ready");
assert_eq!(value["params"]["state_labels"]["working"], "working");
assert_eq!(
value["params"]["state_labels"]["blocked"],
"needs attention"
);
assert_eq!(value["params"]["tokens"]["summary"], "working on auth");
assert_eq!(value["params"]["clear_title"], false);
assert_eq!(value["params"]["clear_display_agent"], false);
assert_eq!(value["params"]["clear_state_labels"], false);
}
#[test]
fn metadata_null_summary_clears_token_and_title_clear_uses_schema_flag() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_metadata(None, None);
reporter.report_summary(None);
reporter.report_title(None);
let lines = transport.lines.lock().unwrap();
let full_clear = report_value(&lines[0]);
assert_eq!(full_clear["method"], "pane.report_metadata");
assert_eq!(
full_clear["params"]["tokens"]["summary"],
serde_json::Value::Null
);
assert_eq!(full_clear["params"]["clear_title"], true);
assert_eq!(full_clear["params"]["clear_display_agent"], false);
assert_eq!(full_clear["params"]["clear_state_labels"], false);
let token_clear = report_value(&lines[1]);
assert_eq!(
token_clear["params"]["tokens"]["summary"],
serde_json::Value::Null
);
assert!(token_clear["params"].get("title").is_none());
assert!(token_clear["params"].get("state_labels").is_none());
let title_clear = report_value(&lines[2]);
assert_eq!(title_clear["params"]["clear_title"], true);
assert_eq!(title_clear["params"]["display_agent"], AGENT);
assert_eq!(title_clear["params"]["state_labels"]["working"], "working");
assert_eq!(title_clear["params"]["clear_display_agent"], false);
assert_eq!(title_clear["params"]["clear_state_labels"], false);
assert!(title_clear["params"].get("tokens").is_none());
}
#[test]
fn title_replacement_resends_complete_presentation_snapshot() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_title(Some("new title"));
reporter.report_title(None);
let lines = transport.lines.lock().unwrap();
let set = report_value(&lines[0]);
assert_eq!(set["params"]["title"], "new title");
assert_eq!(set["params"]["display_agent"], AGENT);
assert_eq!(set["params"]["state_labels"]["working"], "working");
assert_eq!(set["params"]["clear_display_agent"], false);
assert_eq!(set["params"]["clear_state_labels"], false);
let clear = report_value(&lines[1]);
assert_eq!(clear["params"]["clear_title"], true);
assert_eq!(clear["params"]["display_agent"], AGENT);
assert_eq!(
clear["params"]["state_labels"]["blocked"],
"needs attention"
);
assert_eq!(clear["params"]["clear_display_agent"], false);
assert_eq!(clear["params"]["clear_state_labels"], false);
}
#[test]
fn invalid_metadata_values_are_dropped_but_none_explicitly_clears() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_title(Some("\n"));
reporter.report_summary(Some("\0"));
reporter.report_metadata(Some("\n"), Some("valid summary"));
assert!(transport.lines.lock().unwrap().is_empty());
reporter.report_title(None);
let lines = transport.lines.lock().unwrap();
assert_eq!(lines.len(), 1);
assert_eq!(report_value(&lines[0])["params"]["clear_title"], true);
}
#[test]
fn session_start_sources_serialize_as_the_exact_typed_values() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
for (session_id, source, expected) in [
("startup-id", HerdrSessionStartSource::Startup, "startup"),
("resume-id", HerdrSessionStartSource::Resume, "resume"),
("new-id", HerdrSessionStartSource::New, "new"),
("select-id", HerdrSessionStartSource::Select, "select"),
] {
reporter.report_agent_session_with_source(session_id, source);
let lines = transport.lines.lock().unwrap();
assert_eq!(
report_value(lines.last().unwrap())["params"]["session_start_source"],
expected
);
}
}
#[test]
fn invalid_session_ids_are_rejected_without_rewriting_identity() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
reporter.report_agent_session("keep\nid");
reporter.report_agent_session(&"x".repeat(MAX_SESSION_ID_BYTES + 1));
assert!(transport.lines.lock().unwrap().is_empty());
let exact = " leading and trailing ";
reporter.report_agent_session(exact);
let lines = transport.lines.lock().unwrap();
assert_eq!(report_value(&lines[0])["params"]["agent_session_id"], exact);
}
#[test]
fn release_clears_owned_metadata_before_final_release() {
let transport = Arc::new(RecordingTransport::default());
let owner = HerdrOwner::new_with_transport_and_clock(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport.clone(),
Arc::new(FixedClock(1_000_000_000)),
);
let reporter = owner.reporter();
reporter.report_metadata(Some("Title"), Some("summary"));
owner.release();
let lines = transport.lines.lock().unwrap();
assert_eq!(lines.len(), 3);
let clear = report_value(&lines[1]);
assert_eq!(clear["method"], "pane.report_metadata");
assert_eq!(clear["params"]["source"], METADATA_SOURCE);
assert_eq!(clear["params"]["applies_to_source"], SOURCE);
assert_eq!(clear["params"]["clear_title"], true);
assert_eq!(clear["params"]["clear_display_agent"], true);
assert_eq!(clear["params"]["clear_state_labels"], true);
assert!(clear["params"].get("display_agent").is_none());
assert!(clear["params"].get("state_labels").is_none());
assert_eq!(
clear["params"]["tokens"]["summary"],
serde_json::Value::Null
);
let release = report_value(&lines[2]);
assert_eq!(release["method"], "pane.release_agent");
assert_eq!(release["params"]["pane_id"], "pane-1");
assert_eq!(release["params"]["source"], SOURCE);
assert_eq!(release["params"]["agent"], AGENT);
assert!(release["params"]["seq"].is_u64());
let clear_seq = clear["params"]["seq"].as_u64().unwrap();
let release_seq = release["params"]["seq"].as_u64().unwrap();
assert_eq!(release_seq, clear_seq + 1);
}
#[test]
fn release_is_idempotent_and_fences_all_stale_clones() {
let transport = Arc::new(RecordingTransport::default());
let owner = HerdrOwner::new_with_transport_and_clock(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport.clone(),
Arc::new(FixedClock(1_000_000_000)),
);
let reporter = owner.reporter();
let stale_clone = reporter.clone();
reporter.report_thinking();
owner.release();
stale_clone.report_done();
stale_clone.report_agent_session("late-session");
stale_clone.report_summary(Some("late summary"));
let lines = transport.lines.lock().unwrap();
assert_eq!(lines.len(), 3);
assert_eq!(report_value(&lines[2])["method"], "pane.release_agent");
}
#[test]
fn metadata_and_messages_are_locally_bounded_before_send() {
let transport = Arc::new(RecordingTransport::default());
let reporter = reporter_with_transport(transport.clone());
let long = "x".repeat(200);
reporter.report_metadata(Some(&long), Some(&long));
reporter.report(HerdrAgentState::Working, Some(&long));
let lines = transport.lines.lock().unwrap();
let metadata = report_value(&lines[0]);
let agent = report_value(&lines[1]);
assert_eq!(
metadata["params"]["title"]
.as_str()
.unwrap()
.chars()
.count(),
64
);
assert_eq!(
metadata["params"]["tokens"]["summary"]
.as_str()
.unwrap()
.chars()
.count(),
64
);
assert_eq!(
agent["params"]["message"].as_str().unwrap().chars().count(),
64
);
}
#[test]
fn sequence_is_timestamp_seeded_and_monotonic_across_constructions() {
let first_transport = Arc::new(RecordingTransport::default());
let first_reporter = reporter_with_transport(first_transport.clone());
first_reporter.report_thinking();
let second_transport = Arc::new(RecordingTransport::default());
let second_owner = HerdrOwner::new_with_transport_and_clock(
"pane-2".to_string(),
PathBuf::from("/tmp/herdr.sock"),
second_transport.clone(),
Arc::new(FixedClock(1_000_000_001)),
);
let second_reporter = second_owner.reporter();
second_reporter.report_done();
let first_seq = report_params(&first_transport.lines.lock().unwrap()[0])["seq"]
.as_u64()
.unwrap();
let second_seq = report_params(&second_transport.lines.lock().unwrap()[0])["seq"]
.as_u64()
.unwrap();
assert!(first_seq >= 1_000_000_000_000);
assert!(second_seq > first_seq);
}
#[test]
fn sequence_allocator_uses_max_of_previous_plus_one_and_clock_floor() {
let (first, _) = allocate_report_seq_range_at(0, 2_000_000_000, 1).unwrap();
let (second, _) = allocate_report_seq_range_at(first, 1, 1).unwrap();
assert!(first >= 2_000_000_000_000);
assert!(second > first);
assert_eq!(allocate_report_seq_range_at(u64::MAX, 1, 1), None);
assert_eq!(
allocate_report_seq_range_at(u64::MAX - 1, u64::MAX, 1),
Some((u64::MAX, u64::MAX))
);
assert_eq!(allocate_report_seq_range_at(u64::MAX, u64::MAX, 1), None);
}
#[test]
fn herdr_reporter_resolver_requires_enabled_env_and_pane() {
let enabled = HerdrSettings {
enabled: true,
..Default::default()
};
assert!(HerdrOwner::from_env_map(&HerdrSettings::default(), MapEnv::default()).is_none());
assert!(HerdrOwner::from_env_map(&enabled, MapEnv::default()).is_none());
assert!(HerdrOwner::from_env_map(&enabled, MapEnv::default().with("HERDR_ENV", "0")).is_none());
assert!(HerdrOwner::from_env_map(&enabled, MapEnv::default().with("HERDR_ENV", "1")).is_none());
let owner = HerdrOwner::from_env_map(
&enabled,
MapEnv::default()
.with("HERDR_ENV", "1")
.with("HERDR_PANE_ID", " pane-1 ")
.with("HERDR_SOCKET_PATH", "/tmp/herdr.sock"),
)
.unwrap();
let reporter = owner.reporter();
assert_eq!(reporter.shared.pane_id, "pane-1");
assert_eq!(
reporter.shared.socket_path,
PathBuf::from("/tmp/herdr.sock")
);
}
#[test]
fn turn_reporter_does_not_release_application_authority() {
let (owner, lines) = HerdrOwner::new_for_test();
let reporter = owner.reporter();
{
let _turn = HerdrTurnReporter::start(Some(reporter));
}
let lines = lines.lock().unwrap();
assert_eq!(lines.len(), 2);
assert!(
lines
.iter()
.all(|line| report_value(line)["method"] == "pane.report_agent")
);
}
#[test]
fn turn_reporter_terminal_outcomes_are_explicit_and_drop_is_blocked() {
let (owner, lines) = HerdrOwner::new_for_test();
let reporter = owner.reporter();
let mut done = HerdrTurnReporter::start(Some(reporter.clone()));
done.done();
done.done();
let mut cancelled = HerdrTurnReporter::start(Some(reporter.clone()));
cancelled.cancelled();
cancelled.cancelled();
let mut blocked = HerdrTurnReporter::start(Some(reporter));
blocked.blocked();
let lines = lines.lock().unwrap();
let states = lines
.iter()
.map(|line| {
let params = report_params(line);
(
params["state"].as_str().unwrap().to_string(),
params["message"].as_str().unwrap().to_string(),
)
})
.collect::<Vec<_>>();
assert_eq!(
states,
vec![
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "done".to_string()),
("working".to_string(), "thinking".to_string()),
("idle".to_string(), "cancelled".to_string()),
("working".to_string(), "thinking".to_string()),
("blocked".to_string(), "needs attention".to_string()),
]
);
}
#[test]
fn release_fence_waits_for_inflight_report_and_rejects_stale_handles() {
let entered = Arc::new(Barrier::new(2));
let proceed = Arc::new(Barrier::new(2));
let transport = Arc::new(BlockingTransport {
lines: Mutex::new(Vec::new()),
entered: Arc::clone(&entered),
proceed: Arc::clone(&proceed),
block_first: AtomicBool::new(true),
});
let owner = HerdrOwner::new_with_transport_and_clock(
"pane-1".to_string(),
PathBuf::from("/tmp/herdr.sock"),
transport.clone(),
Arc::new(FixedClock(1_000_000_000)),
);
let reporter = owner.reporter();
let stale = reporter.clone();
let report_thread = std::thread::spawn(move || reporter.report_thinking());
entered.wait();
assert_eq!(transport.lines.lock().unwrap().len(), 1);
let release_thread = std::thread::spawn(move || owner.release());
assert_eq!(transport.lines.lock().unwrap().len(), 1);
proceed.wait();
report_thread.join().unwrap();
release_thread.join().unwrap();
stale.report_done();
let lines = transport.lines.lock().unwrap();
assert_eq!(lines.len(), 3);
assert_eq!(report_value(&lines[0])["method"], "pane.report_agent");
assert_eq!(report_value(&lines[1])["method"], "pane.report_metadata");
assert_eq!(report_value(&lines[2])["method"], "pane.release_agent");
}
#[cfg(unix)]
#[test]
fn unix_socket_transport_writes_a_complete_line() {
use std::{io::Read, os::unix::net::UnixListener};
let directory = tempfile::tempdir().unwrap();
let socket_path = directory.path().join("herdr.sock");
let listener = UnixListener::bind(&socket_path).unwrap();
let reader = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let mut received = String::new();
stream.read_to_string(&mut received).unwrap();
received
});
write_line_to_socket(&socket_path, "bounded test line\n").unwrap();
assert_eq!(reader.join().unwrap(), "bounded test line\n");
}
#[test]
fn transport_failure_is_best_effort_even_for_release() {
let owner = HerdrOwner::new_with_transport(
"pane-1".to_string(),
PathBuf::from("/tmp/missing-herdr.sock"),
Arc::new(FailingTransport),
);
let reporter = owner.reporter();
reporter.report_thinking();
reporter.report_metadata(Some("title"), Some("summary"));
owner.release();
reporter.report_done();
}