use std::os::unix::net as path_std_os_unix_net;
use std::sync::mpsc as path_std_sync_mpsc;
use std::{io as path_std_io, time as path_std_time};
use super::*;
use crate::estimated_cost::AgentCostSnapshot;
fn entry(id: &str, parent: Option<&str>, started_at: Option<u64>) -> SessionAgentListEntry {
SessionAgentListEntry {
agent_id: tau_proto::AgentId::parse(id).expect("valid test id"),
lifecycle: SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::Active,
},
persistence: SessionAgentPersistence::Durable,
facts: SessionAgentFacts::Available {
started_at: started_at.map(tau_proto::UnixMicros::new),
parent_agent: parent.map(|id| tau_proto::AgentId::parse(id).expect("valid parent id")),
role: "engineer".to_owned(),
display_name: None,
},
work_status: Some(tau_proto::SessionAgentWorkStatus::default()),
turn_activity: Some(tau_proto::AgentTurnActivity::Idle),
}
}
#[test]
fn default_filter_is_mode_not_suspended() {
let mut active_auto = entry("auto", None, Some(1));
active_auto.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::ActiveAuto,
};
let mut suspended = entry("suspended", None, Some(2));
suspended.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::Suspended,
};
let mut unavailable = entry("unavailable", None, Some(3));
unavailable.lifecycle = SessionAgentLifecycle::Unavailable;
let visible = visible_agents(
vec![suspended, unavailable, active_auto],
AgentListFilter::default(),
);
assert_eq!(
visible
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<Vec<_>>(),
vec!["auto"]
);
}
#[test]
fn active_picker_filters_navigation_and_runtime_state_independently() {
let cases = [
(
"active-running",
tau_proto::AgentNavigationMode::Active,
tau_proto::AgentRuntimeState::Running,
true,
),
(
"active-idle",
tau_proto::AgentNavigationMode::Active,
tau_proto::AgentRuntimeState::Idle,
true,
),
(
"auto-running",
tau_proto::AgentNavigationMode::ActiveAuto,
tau_proto::AgentRuntimeState::Running,
true,
),
(
"auto-idle",
tau_proto::AgentNavigationMode::ActiveAuto,
tau_proto::AgentRuntimeState::Idle,
false,
),
(
"suspended-running",
tau_proto::AgentNavigationMode::Suspended,
tau_proto::AgentRuntimeState::Running,
false,
),
(
"suspended-idle",
tau_proto::AgentNavigationMode::Suspended,
tau_proto::AgentRuntimeState::Idle,
false,
),
];
let rows = cases
.iter()
.enumerate()
.map(|(index, (id, mode, runtime, _))| {
let mut row = entry(id, None, Some(index as u64));
row.lifecycle = SessionAgentLifecycle::Live {
runtime_state: *runtime,
navigation_mode: *mode,
};
row
})
.collect();
let visible = picker_agents(rows, AgentPickerFilter::Active)
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<std::collections::HashSet<_>>();
for (id, _, _, expected) in cases {
assert_eq!(visible.contains(id), expected, "{id}");
}
}
#[test]
fn all_picker_includes_suspended_agents_and_preserves_runtime_column() {
let mut running = entry("suspended-running", None, Some(1));
running.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Running,
navigation_mode: tau_proto::AgentNavigationMode::Suspended,
};
let mut idle = entry("auto-idle", None, Some(2));
idle.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::ActiveAuto,
};
idle.facts = SessionAgentFacts::Missing;
let output = format_rows(&picker_agents(vec![running, idle], AgentPickerFilter::All));
assert!(output.contains("suspended-running\tlive\trunning\tidle\tsuspended\t"));
assert!(output.contains("auto-idle\tlive\tidle\tidle\tactive_auto\tdurable\tmissing\t"));
}
#[test]
fn picker_rows_append_canonical_cost_and_status() {
let zero = entry("zero", None, Some(1));
let mut nonzero = entry("nonzero", None, Some(2));
nonzero.work_status = Some(
tau_proto::SessionAgentWorkStatus::new(
tau_proto::AgentWorkStatusPhase::Working,
Some("verify \\ picker\u{202e} rows".to_owned()),
)
.expect("valid status"),
);
let unavailable = entry("unavailable", None, Some(3));
let output = format_picker_rows(&[zero, nonzero, unavailable], |agent_id| {
match agent_id.as_str() {
"zero" => Some(AgentCostSnapshot::new(
tau_proto::EstimatedApiCost::default(),
tau_proto::EstimatedApiCost::default(),
)),
"nonzero" => Some(AgentCostSnapshot::new(
tau_proto::EstimatedApiCost::from_picodollars(2_140_000_000_000),
tau_proto::EstimatedApiCost::from_picodollars(4_280_000_000_000),
)),
_ => None,
}
});
let extras = output
.lines()
.map(|row| row.split('\t').skip(11).collect::<Vec<_>>())
.collect::<Vec<_>>();
assert_eq!(
extras,
[
vec!["$.00/$.00", "❓", "-", "💤"],
vec!["$2.1/$4.3", "🚀", r"verify \\ picker\\u{202E} rows", "💤"],
vec!["-/-", "❓", "-", "💤"],
]
);
}
#[test]
fn picker_work_status_symbols_are_complete() {
use tau_proto::AgentWorkStatusPhase::{Blocked, Done, Unknown, Unreported, Waiting, Working};
assert_eq!(
[
None,
Some(Unreported),
Some(Working),
Some(Done),
Some(Blocked),
Some(Waiting),
Some(Unknown)
]
.map(work_status_symbol),
["❓", "❓", "🚀", "✅", "⛔️", "⏳", "❓"]
);
}
#[test]
fn cli_ui_skill_documents_status_indicator_symbols() {
const CLI_UI_SKILL: &str =
include_str!("../../../tau-skills/self-knowledge/tau-self-knowledge-cli-ui.md");
let work_status_cases = [
(
None,
"❓",
"no status is reported, or a previous status is no longer reliable",
),
(
Some(tau_proto::AgentWorkStatusPhase::Working),
"🚀",
"working",
),
(Some(tau_proto::AgentWorkStatusPhase::Done), "✅", "done"),
(
Some(tau_proto::AgentWorkStatusPhase::Blocked),
"⛔️",
"blocked pending external intervention",
),
(
Some(tau_proto::AgentWorkStatusPhase::Waiting),
"⏳",
"waiting for an expected self-resolving event",
),
];
for (phase, glyph, description) in work_status_cases {
assert_eq!(work_status_symbol(phase), glyph);
assert!(
CLI_UI_SKILL.contains(&format!("| `{glyph}` | {description} |")),
"skill must document {glyph} as {description}"
);
}
assert_eq!(
work_status_symbol(Some(tau_proto::AgentWorkStatusPhase::Unreported)),
"❓"
);
assert_eq!(
work_status_symbol(Some(tau_proto::AgentWorkStatusPhase::Unknown)),
"❓"
);
let turn_activity_cases = [
(
tau_proto::AgentTurnActivity::Responding,
"✨",
"the provider is generating a response",
),
(
tau_proto::AgentTurnActivity::Manipulating,
"🔨",
"an active tool is mutating state or has no more-specific category",
),
(
tau_proto::AgentTurnActivity::Fetching,
"🌐",
"an active tool is fetching data",
),
(
tau_proto::AgentTurnActivity::Waiting,
"⏳",
"an active tool is waiting",
),
(
tau_proto::AgentTurnActivity::TimerScheduled,
"🕔",
"no higher-priority work is active and a timer is scheduled",
),
(
tau_proto::AgentTurnActivity::Idle,
"💤",
"idle: no response, tool call, or ambient activity is active",
),
];
for (activity, glyph, description) in turn_activity_cases {
assert_eq!(turn_activity_symbol(activity), glyph);
assert!(
CLI_UI_SKILL.contains(&format!("| `{glyph}` | {description} |")),
"skill must document {glyph} as {description}"
);
}
assert!(CLI_UI_SKILL.contains("`⏳⏳` means"));
}
#[test]
fn pickers_keep_live_agents_without_available_creation_facts() {
let mut missing = entry("missing", None, Some(1));
missing.facts = SessionAgentFacts::Missing;
let mut invalid = entry("invalid", None, Some(2));
invalid.facts = SessionAgentFacts::Invalid;
let mut unreadable = entry("unreadable", None, Some(3));
unreadable.facts = SessionAgentFacts::Unreadable;
let mut unavailable = entry("unavailable", None, Some(4));
unavailable.lifecycle = SessionAgentLifecycle::Unavailable;
let mut unloaded = entry("unloaded", None, Some(5));
unloaded.lifecycle = SessionAgentLifecycle::Unloaded;
for filter in [AgentPickerFilter::Active, AgentPickerFilter::All] {
let visible = picker_agents(
vec![
missing.clone(),
invalid.clone(),
unreadable.clone(),
unavailable.clone(),
unloaded.clone(),
],
filter,
)
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<Vec<_>>();
assert_eq!(visible, vec!["invalid", "missing", "unreadable"]);
}
}
#[test]
fn picker_revalidation_preserves_active_or_all_category() {
let mut running = entry("auto", None, Some(1));
running.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Running,
navigation_mode: tau_proto::AgentNavigationMode::ActiveAuto,
};
let mut idle = running.clone();
idle.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::ActiveAuto,
};
let selected = &running.agent_id;
assert!(picker_selection_is_current(
&[running.clone()],
selected,
AgentPickerFilter::Active
));
assert!(!picker_selection_is_current(
&[idle.clone()],
selected,
AgentPickerFilter::Active
));
assert!(picker_selection_is_current(
&[idle],
selected,
AgentPickerFilter::All
));
}
#[test]
fn all_category_filter_is_additive() {
let mut suspended = entry("suspended", None, Some(1));
suspended.lifecycle = SessionAgentLifecycle::Live {
runtime_state: tau_proto::AgentRuntimeState::Idle,
navigation_mode: tau_proto::AgentNavigationMode::Suspended,
};
let mut unavailable = entry("unavailable", None, Some(2));
unavailable.lifecycle = SessionAgentLifecycle::Unavailable;
unavailable.facts = SessionAgentFacts::Unreadable;
let mut unloaded = entry("unloaded", None, Some(3));
unloaded.lifecycle = SessionAgentLifecycle::Unloaded;
let visible = visible_agents(
vec![unloaded, unavailable, suspended],
AgentListFilter {
include_suspended: true,
include_unavailable: true,
include_unloaded: true,
},
);
assert_eq!(visible.len(), 3);
}
#[test]
fn public_filter_flags_map_independently_and_all_enables_every_flag() {
let base = crate::cli::AgentListArgs {
session_id: "s1".parse().expect("valid session id"),
include_suspended: false,
include_unavailable: false,
include_unloaded: false,
all: false,
};
for (expected, args) in [
(
AgentListFilter {
include_suspended: true,
..AgentListFilter::default()
},
crate::cli::AgentListArgs {
include_suspended: true,
..base.clone()
},
),
(
AgentListFilter {
include_unavailable: true,
..AgentListFilter::default()
},
crate::cli::AgentListArgs {
include_unavailable: true,
..base.clone()
},
),
(
AgentListFilter {
include_unloaded: true,
..AgentListFilter::default()
},
crate::cli::AgentListArgs {
include_unloaded: true,
..base.clone()
},
),
] {
assert_eq!(AgentListFilter::from_args(&args), expected);
}
assert_eq!(
AgentListFilter::from_args(&crate::cli::AgentListArgs { all: true, ..base }),
AgentListFilter {
include_suspended: true,
include_unavailable: true,
include_unloaded: true,
}
);
}
#[test]
fn topological_order_places_parent_before_child() {
let rows = vec![
entry("child", Some("parent"), Some(1)),
entry("later-root", None, Some(3)),
entry("parent", None, Some(9)),
entry("early-root", None, Some(2)),
];
let ordered = topological_order(rows)
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<Vec<_>>();
assert_eq!(ordered, vec!["early-root", "later-root", "parent", "child"]);
}
#[test]
fn topological_order_preserves_cycle_rows() {
let rows = vec![
entry("b", Some("a"), Some(2)),
entry("a", Some("b"), Some(1)),
];
let ordered = topological_order(rows)
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<Vec<_>>();
assert_eq!(ordered, vec!["a", "b"]);
}
#[test]
fn topological_order_is_stable_for_legacy_and_orphan_rows() {
let rows = vec![
entry("unknown-b", None, None),
entry("orphan", Some("missing"), Some(2)),
entry("self-parent", Some("self-parent"), Some(1)),
entry("unknown-a", None, None),
];
let expected = vec!["self-parent", "orphan", "unknown-a", "unknown-b"];
for rows in [rows.clone(), rows.into_iter().rev().collect()] {
assert_eq!(
topological_order(rows)
.into_iter()
.map(|agent| agent.agent_id.to_string())
.collect::<Vec<_>>(),
expected
);
}
}
#[test]
fn format_rows_escapes_free_text() {
let mut row = entry("agent", None, Some(42));
row.facts = SessionAgentFacts::Available {
started_at: Some(tau_proto::UnixMicros::new(42)),
parent_agent: None,
role: "role\twith\\slash".to_owned(),
display_name: Some("line\nname".to_owned()),
};
let output = format_rows(&[row]);
let fields = output.trim_end().split('\t').collect::<Vec<_>>();
assert_eq!(fields.len(), 11);
assert_eq!(fields[0], "agent");
assert_eq!(fields[7], "role\\twith\\\\slash");
assert_eq!(fields[10], "line\\nname");
}
#[test]
fn format_rows_matches_exact_eleven_column_contract() {
let mut row = entry("agent", None, None);
row.lifecycle = SessionAgentLifecycle::Unavailable;
row.persistence = SessionAgentPersistence::Ephemeral;
row.facts = SessionAgentFacts::Invalid;
assert_eq!(
format_rows(&[row]),
"agent\tunavailable\t-\t-\t-\tephemeral\tinvalid\t-\t-\t-\t-\n"
);
}
#[test]
fn selected_agent_id_rejects_multiline_or_invalid_rows() {
assert_eq!(
selected_agent_id("agent-1\tlive")
.expect("valid selection")
.as_str(),
"agent-1"
);
assert!(selected_agent_id("bad/id\tlive").is_err());
assert!(selected_agent_id("agent-1\nagent-2").is_err());
}
#[test]
fn detailed_activity_mappings_cover_all_categories() {
use tau_proto::AgentTurnActivity::{
Fetching, Idle, Manipulating, Responding, TimerScheduled, Waiting,
};
let cases = [
(Responding, "responding", "✨"),
(Manipulating, "manipulating", "🔨"),
(Fetching, "fetching", "🌐"),
(Waiting, "waiting", "⏳"),
(TimerScheduled, "timer_scheduled", "🕔"),
(Idle, "idle", "💤"),
];
for (activity, name, emoji) in cases {
assert_eq!(turn_activity_name(activity), name);
assert_eq!(turn_activity_symbol(activity), emoji);
}
}
#[cfg(unix)]
#[test]
fn roster_request_times_out_on_silent_peer() {
assert_roster_request_times_out(RosterTimeoutPeerBehavior::Silent);
}
#[cfg(unix)]
#[test]
fn roster_request_deadline_survives_unrelated_frames() {
assert_roster_request_times_out(RosterTimeoutPeerBehavior::UnrelatedFrames);
}
#[cfg(unix)]
#[test]
fn roster_request_deadline_stops_partial_frame_trickle() {
assert_roster_request_times_out(RosterTimeoutPeerBehavior::PartialFrameTrickle);
}
#[cfg(unix)]
#[derive(Clone, Copy)]
enum RosterTimeoutPeerBehavior {
Silent,
UnrelatedFrames,
PartialFrameTrickle,
}
#[cfg(unix)]
fn assert_roster_request_times_out(behavior: RosterTimeoutPeerBehavior) {
let temp = tempfile::tempdir().expect("tempdir");
let socket_path = temp.path().join("harness.sock");
let listener = path_std_os_unix_net::UnixListener::bind(&socket_path).expect("bind listener");
listener
.set_nonblocking(true)
.expect("nonblocking listener");
let (stop_tx, stop_rx) = path_std_sync_mpsc::channel();
let server = std::thread::spawn(move || run_roster_timeout_peer(listener, stop_rx, behavior));
let started = path_std_time::Instant::now();
let result = request_at_socket_with_timeout_typed(
&socket_path,
&tau_proto::SessionId::parse("s1").expect("session id"),
SessionAgentListScope::Current,
Duration::from_millis(200),
);
let elapsed = started.elapsed();
let _ = stop_tx.send(());
let server_result = server.join().expect("server thread");
server_result.expect("peer must admit s1 and observe its current-scope roster request");
assert!(
matches!(
result,
Err(CliError::Participant(ref message))
if message == "agent roster request timed out"
),
"expected exact roster timeout, got {result:?}"
);
assert!(
elapsed < Duration::from_secs(1),
"roster timeout took {elapsed:?}"
);
}
#[cfg(unix)]
fn run_roster_timeout_peer(
listener: path_std_os_unix_net::UnixListener,
stop: path_std_sync_mpsc::Receiver<()>,
behavior: RosterTimeoutPeerBehavior,
) -> Result<(), String> {
use std::io::Write as _;
let lifetime_deadline = path_std_time::Instant::now() + Duration::from_millis(1_500);
let stream = loop {
match listener.accept() {
Ok((stream, _)) => break stream,
Err(error) if error.kind() == path_std_io::ErrorKind::WouldBlock => {
if stop.try_recv().is_ok() {
return Err("client stopped before connecting".to_owned());
}
if path_std_time::Instant::now() >= lifetime_deadline {
return Err("timed out accepting roster client".to_owned());
}
std::thread::sleep(Duration::from_millis(2));
}
Err(error) => return Err(format!("accept roster client: {error}")),
}
};
stream
.set_read_timeout(Some(Duration::from_secs(1)))
.map_err(|error| format!("set peer read timeout: {error}"))?;
stream
.set_write_timeout(Some(Duration::from_secs(1)))
.map_err(|error| format!("set peer write timeout: {error}"))?;
let read_stream = stream
.try_clone()
.map_err(|error| format!("clone peer stream: {error}"))?;
let mut reader = tau_proto::HarnessInputReader::new(path_std_io::BufReader::new(read_stream));
let hello = reader
.read_message()
.map_err(|error| format!("read Hello: {error}"))?
.ok_or_else(|| "client disconnected before Hello".to_owned())?;
let HarnessInputMessage::Hello(hello) = hello else {
return Err(format!("expected Hello, got {hello:?}"));
};
let session_id =
tau_proto::SessionId::parse("s1").map_err(|error| format!("parse session id: {error}"))?;
let client_name = tau_proto::ExtensionName::parse("tau-list-agents")
.map_err(|error| format!("parse client name: {error}"))?;
let expected_hello = crate::ui_client::hello_message(client_name, Some(&session_id));
if HarnessInputMessage::Hello(hello.clone()) != expected_hello {
return Err(format!("unexpected roster Hello: {hello:?}"));
}
let mut writer = tau_proto::HarnessOutputWriter::new(path_std_io::BufWriter::new(
stream
.try_clone()
.map_err(|error| format!("clone peer writer: {error}"))?,
));
writer
.write_message(&HarnessOutputMessage::SessionAccepted(
tau_proto::SessionAccepted {
session_id,
harness_protocol_version: None,
},
))
.map_err(|error| format!("write SessionAccepted: {error}"))?;
writer
.flush()
.map_err(|error| format!("flush SessionAccepted: {error}"))?;
let request = reader
.read_message()
.map_err(|error| format!("read GetSessionAgentList: {error}"))?
.ok_or_else(|| "client disconnected before GetSessionAgentList".to_owned())?;
let HarnessInputMessage::GetSessionAgentList(request) = request else {
return Err(format!("expected GetSessionAgentList, got {request:?}"));
};
if request.request_id.is_empty()
|| request.session_id.as_str() != "s1"
|| request.scope != SessionAgentListScope::Current
{
return Err(format!("unexpected roster request: {request:?}"));
}
match behavior {
RosterTimeoutPeerBehavior::Silent => wait_for_roster_peer_stop(&stop, lifetime_deadline),
RosterTimeoutPeerBehavior::UnrelatedFrames => {
let mut index = 0_u64;
while !roster_peer_should_stop(&stop, lifetime_deadline) {
if let Err(error) =
writer.write_message(&HarnessOutputMessage::PeerSessionProbeResult(
tau_proto::PeerSessionProbeResult {
request_id: format!("unrelated-{index}"),
available: false,
},
))
{
return roster_peer_finish_after_write_error(
&stop,
"write unrelated frame",
error,
);
}
if let Err(error) = writer.flush() {
return roster_peer_finish_after_write_error(
&stop,
"flush unrelated frame",
error,
);
}
index += 1;
std::thread::sleep(Duration::from_millis(5));
}
Ok(())
}
RosterTimeoutPeerBehavior::PartialFrameTrickle => {
let frame = tau_proto::encode_message_to_vec(
&HarnessOutputMessage::PeerSessionProbeResult(tau_proto::PeerSessionProbeResult {
request_id: "unrelated-".repeat(128),
available: false,
}),
)
.map_err(|error| format!("encode unrelated frame: {error}"))?;
let mut stream = stream;
for byte in frame {
if roster_peer_should_stop(&stop, lifetime_deadline) {
return Ok(());
}
if let Err(error) = stream.write_all(&[byte]) {
return roster_peer_finish_after_write_error(
&stop,
"write partial frame",
error,
);
}
std::thread::sleep(Duration::from_millis(10));
}
wait_for_roster_peer_stop(&stop, lifetime_deadline)
}
}
}
#[cfg(unix)]
fn wait_for_roster_peer_stop(
stop: &path_std_sync_mpsc::Receiver<()>,
lifetime_deadline: path_std_time::Instant,
) -> Result<(), String> {
let remaining = lifetime_deadline.saturating_duration_since(path_std_time::Instant::now());
match stop.recv_timeout(remaining) {
Ok(()) => Ok(()),
Err(path_std_sync_mpsc::RecvTimeoutError::Timeout) => Ok(()),
Err(path_std_sync_mpsc::RecvTimeoutError::Disconnected) => {
Err("roster test dropped cleanup signal".to_owned())
}
}
}
#[cfg(unix)]
fn roster_peer_should_stop(
stop: &path_std_sync_mpsc::Receiver<()>,
lifetime_deadline: path_std_time::Instant,
) -> bool {
stop.try_recv().is_ok() || path_std_time::Instant::now() >= lifetime_deadline
}
#[cfg(unix)]
fn roster_peer_finish_after_write_error<E: std::fmt::Display>(
stop: &path_std_sync_mpsc::Receiver<()>,
action: &str,
error: E,
) -> Result<(), String> {
match stop.recv_timeout(Duration::from_millis(20)) {
Ok(()) => Ok(()),
Err(_) => Err(format!("{action}: {error}")),
}
}
#[cfg(target_os = "linux")]
#[test]
fn roster_request_deadline_bounds_saturated_backlog_connect() {
let temp = tempfile::tempdir().expect("tempdir");
let socket_path = temp.path().join("harness.sock");
let listener = socket2::Socket::new(socket2::Domain::UNIX, socket2::Type::STREAM, None)
.expect("listener socket");
listener
.bind(&socket2::SockAddr::unix(&socket_path).expect("socket address"))
.expect("bind listener");
listener.listen(1).expect("listen");
let mut backlog = Vec::new();
for _ in 0..16 {
let socket = socket2::Socket::new(socket2::Domain::UNIX, socket2::Type::STREAM, None)
.expect("backlog socket");
match socket.connect_timeout(
&socket2::SockAddr::unix(&socket_path).expect("socket address"),
Duration::from_millis(10),
) {
Ok(()) => backlog.push(socket),
Err(_) => break,
}
}
assert!(!backlog.is_empty());
let started = path_std_time::Instant::now();
let result = request_at_socket_with_timeout_typed(
&socket_path,
&tau_proto::SessionId::parse("s1").expect("session id"),
SessionAgentListScope::Current,
Duration::from_millis(20),
);
assert!(result.is_err());
assert!(started.elapsed() < Duration::from_secs(1));
}