From 096860fb718fef6d21c56d6dc295d9832bae3bba Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 21:57:45 -0400 Subject: [PATCH 01/16] Add ds_agent_ipc: shared socket path and line framing Co-Authored-By: Claude Haiku 4.5 --- Cargo.toml | 2 +- ds_agent_ipc/Cargo.toml | 10 +++++ ds_agent_ipc/src/lib.rs | 89 +++++++++++++++++++++++++++++++++++++++++ 3 files changed, 100 insertions(+), 1 deletion(-) create mode 100644 ds_agent_ipc/Cargo.toml create mode 100644 ds_agent_ipc/src/lib.rs diff --git a/Cargo.toml b/Cargo.toml index 1839e83..1a0d9bf 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,3 @@ [workspace] resolver = "3" -members = ["robocol", "fake_rc", "ds_cli", "capture/decode"] +members = ["robocol", "fake_rc", "ds_cli", "capture/decode", "ds_agent_ipc"] diff --git a/ds_agent_ipc/Cargo.toml b/ds_agent_ipc/Cargo.toml new file mode 100644 index 0000000..dd445e8 --- /dev/null +++ b/ds_agent_ipc/Cargo.toml @@ -0,0 +1,10 @@ +[package] +name = "ds_agent_ipc" +version = "0.1.0" +edition = "2024" +license = "MIT" +description = "Shared socket-path and line-framing logic for ds_agentd/ds_agent" + +[dependencies] +serde_json = "1" +libc = "0.2" diff --git a/ds_agent_ipc/src/lib.rs b/ds_agent_ipc/src/lib.rs new file mode 100644 index 0000000..4f20daf --- /dev/null +++ b/ds_agent_ipc/src/lib.rs @@ -0,0 +1,89 @@ +//! Socket path and line-framing shared by `ds_agentd` and `ds_agent`. Both +//! binaries must resolve the same path and speak the same framing, so this +//! logic lives once instead of being duplicated on both sides of the pipe. + +use std::io::{self, BufRead, Write}; +use std::path::PathBuf; + +/// Where `ds_agentd` binds its Unix domain socket and `ds_agent` connects. +/// Prefers `$XDG_RUNTIME_DIR` (cleaned up by the OS on logout); falls back to +/// a per-user path under `/tmp` so this also works in a plain SSH session. +pub fn socket_path() -> PathBuf { + if let Some(dir) = std::env::var_os("XDG_RUNTIME_DIR") { + return PathBuf::from(dir).join("ds_agentd.sock"); + } + let uid = unsafe { libc::getuid() }; + PathBuf::from(format!("/tmp/ds_agentd-{uid}.sock")) +} + +/// Writes one JSON value as a single line: the whole `\n`-terminated buffer +/// goes out in one `write_all`, so a reader never sees a partial object. +pub fn write_message(stream: &mut impl Write, value: &serde_json::Value) -> io::Result<()> { + let mut line = serde_json::to_vec(value).map_err(io::Error::other)?; + line.push(b'\n'); + stream.write_all(&line) +} + +/// Reads one `\n`-terminated line, with the newline stripped. `Ok(None)` +/// means clean EOF (the peer closed the connection); any other read error +/// (including an unterminated final line) still surfaces as `Err`. +pub fn read_message(reader: &mut impl BufRead) -> io::Result> { + let mut line = String::new(); + let n = reader.read_line(&mut line)?; + if n == 0 { + return Ok(None); + } + if line.ends_with('\n') { + line.pop(); + if line.ends_with('\r') { + line.pop(); + } + } + Ok(Some(line)) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Cursor; + + #[test] + fn round_trips_one_message() { + let mut buf = Vec::new(); + write_message(&mut buf, &serde_json::json!({"ok": true, "n": 1})).unwrap(); + let mut reader = std::io::BufReader::new(Cursor::new(buf)); + let line = read_message(&mut reader).unwrap().unwrap(); + let value: serde_json::Value = serde_json::from_str(&line).unwrap(); + assert_eq!(value["ok"], true); + assert_eq!(value["n"], 1); + } + + #[test] + fn reads_multiple_messages_in_sequence() { + let mut buf = Vec::new(); + write_message(&mut buf, &serde_json::json!({"seq": 1})).unwrap(); + write_message(&mut buf, &serde_json::json!({"seq": 2})).unwrap(); + let mut reader = std::io::BufReader::new(Cursor::new(buf)); + let first = read_message(&mut reader).unwrap().unwrap(); + let second = read_message(&mut reader).unwrap().unwrap(); + assert!(first.contains("\"seq\":1")); + assert!(second.contains("\"seq\":2")); + } + + #[test] + fn clean_eof_yields_none() { + let mut reader = std::io::BufReader::new(Cursor::new(Vec::new())); + assert!(read_message(&mut reader).unwrap().is_none()); + } + + #[test] + fn socket_path_prefers_xdg_runtime_dir() { + // SAFETY: test runs single-threaded within this process's env mutation. + unsafe { std::env::set_var("XDG_RUNTIME_DIR", "/tmp/ds-agent-test-runtime") }; + assert_eq!( + socket_path(), + PathBuf::from("/tmp/ds-agent-test-runtime/ds_agentd.sock") + ); + unsafe { std::env::remove_var("XDG_RUNTIME_DIR") }; + } +} From 4434f732c3bcfd98bea5f56992e2ff0516a6c559 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:03:04 -0400 Subject: [PATCH 02/16] Add ds_agentd DaemonState: event-driven cache of robot status Implements the DaemonState struct and all required methods (new, apply_event, find_config, find_opmode, to_json) with comprehensive unit tests covering: - Initial disconnected state with empty caches - Connection and disconnection state management - OpMode list population and lookup - Configuration list parsing and lookup - Telemetry caching by tag with replacement - JSON serialization of all state All six tests pass cleanly. Verified Telemetry struct has the expected fields (tag: String, strings: Vec, numbers: BTreeMap, Default impl). --- Cargo.lock | 17 ++++ Cargo.toml | 2 +- ds_agentd/Cargo.toml | 15 +++ ds_agentd/src/main.rs | 3 + ds_agentd/src/state.rs | 217 +++++++++++++++++++++++++++++++++++++++++ 5 files changed, 253 insertions(+), 1 deletion(-) create mode 100644 ds_agentd/Cargo.toml create mode 100644 ds_agentd/src/main.rs create mode 100644 ds_agentd/src/state.rs diff --git a/Cargo.lock b/Cargo.lock index 509a79c..ca50bba 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -29,6 +29,23 @@ dependencies = [ "error-code", ] +[[package]] +name = "ds_agent_ipc" +version = "0.1.0" +dependencies = [ + "libc", + "serde_json", +] + +[[package]] +name = "ds_agentd" +version = "0.1.0" +dependencies = [ + "ds_agent_ipc", + "robocol", + "serde_json", +] + [[package]] name = "ds_cli" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 1a0d9bf..82e6526 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,3 @@ [workspace] resolver = "3" -members = ["robocol", "fake_rc", "ds_cli", "capture/decode", "ds_agent_ipc"] +members = ["robocol", "fake_rc", "ds_cli", "capture/decode", "ds_agent_ipc", "ds_agentd"] diff --git a/ds_agentd/Cargo.toml b/ds_agentd/Cargo.toml new file mode 100644 index 0000000..d1a1091 --- /dev/null +++ b/ds_agentd/Cargo.toml @@ -0,0 +1,15 @@ +[package] +name = "ds_agentd" +version = "0.1.0" +edition = "2024" +license = "MIT" +description = "Headless daemon holding one persistent robocol connection for ds_agent" + +[[bin]] +name = "ds_agentd" +path = "src/main.rs" + +[dependencies] +robocol = { path = "../robocol" } +ds_agent_ipc = { path = "../ds_agent_ipc" } +serde_json = "1" diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs new file mode 100644 index 0000000..b6cfe31 --- /dev/null +++ b/ds_agentd/src/main.rs @@ -0,0 +1,3 @@ +mod state; + +fn main() {} diff --git a/ds_agentd/src/state.rs b/ds_agentd/src/state.rs new file mode 100644 index 0000000..4388958 --- /dev/null +++ b/ds_agentd/src/state.rs @@ -0,0 +1,217 @@ +//! In-memory cache of everything the robot controller has told us, built +//! from the `robocol::Event` stream. Pure data — no sockets, no threads — +//! so `status` can answer instantly and name-lookups (opmode/config) don't +//! need to round-trip the robot. + +use std::collections::HashMap; + +use robocol::client::Event; +use robocol::cmd::{ConfigMeta, OpModeMeta, parse_config_list}; +use robocol::packets::Telemetry; +use robocol::types::RobotState; + +#[allow(dead_code)] +pub struct DaemonState { + pub connected: bool, + pub peer: Option, + pub robot_state: Option, + pub opmodes: Vec, + /// Name most recently confirmed by `Event::OpModeInited`, so `run` with + /// no argument can fall back to it the way `ds_cli`'s REPL session does. + pub last_inited: Option, + pub active_config: Option, + pub configs: Vec, + pub device_types: Option, + pub last_scan: Option, + /// Most recent `CMD_DISCOVER_LYNX_MODULES_RESP` payload. Not keyed by + /// serial: `Event::LynxModules` only carries the RC's raw response, not + /// the serial it was requested for, so per-serial caching isn't possible + /// from this event alone — the `lynx-modules` command itself still + /// returns the right answer, via a fresh wait on this event, per call. + pub last_lynx_modules: Option, + pub telemetry: HashMap, +} + +#[allow(dead_code)] +impl DaemonState { + pub fn new() -> Self { + DaemonState { + connected: false, + peer: None, + robot_state: None, + opmodes: Vec::new(), + last_inited: None, + active_config: None, + configs: Vec::new(), + device_types: None, + last_scan: None, + last_lynx_modules: None, + telemetry: HashMap::new(), + } + } + + pub fn apply_event(&mut self, event: &Event) { + match event { + Event::Connected { peer } => { + self.connected = true; + self.peer = Some(peer.to_string()); + } + Event::Disconnected => { + self.connected = false; + self.peer = None; + self.robot_state = None; + } + Event::RobotState(state) => self.robot_state = Some(*state), + Event::OpModeList(list) => self.opmodes = list.clone(), + Event::OpModeInited(name) => self.last_inited = Some(name.clone()), + Event::ActiveConfiguration(extra) => self.active_config = Some(extra.clone()), + Event::ConfigurationList(extra) => self.configs = parse_config_list(extra), + Event::UserDeviceList(extra) => self.device_types = Some(extra.clone()), + Event::ScanResult(extra) => self.last_scan = Some(extra.clone()), + Event::LynxModules(extra) => self.last_lynx_modules = Some(extra.clone()), + Event::Telemetry(t) => { + self.telemetry.insert(t.tag.clone(), t.clone()); + } + _ => {} + } + } + + pub fn find_config(&self, name: &str) -> Option { + self.configs.iter().find(|c| c.name == name).cloned() + } + + pub fn find_opmode(&self, name: &str) -> bool { + self.opmodes.iter().any(|m| m.name == name) + } + + pub fn to_json(&self) -> serde_json::Value { + let opmodes: Vec = self + .opmodes + .iter() + .map(|m| serde_json::json!({"name": m.name, "flavor": m.flavor, "group": m.group})) + .collect(); + let telemetry: serde_json::Map = self + .telemetry + .iter() + .map(|(tag, t)| { + let strings: serde_json::Map = t + .strings + .iter() + .map(|(k, v)| (k.clone(), serde_json::Value::String(v.clone()))) + .collect(); + let numbers: serde_json::Map = t + .numbers + .iter() + .map(|(k, v)| (k.clone(), serde_json::json!(v))) + .collect(); + ( + tag.clone(), + serde_json::json!({"strings": strings, "numbers": numbers}), + ) + }) + .collect(); + serde_json::json!({ + "ok": true, + "connected": self.connected, + "peer": self.peer, + "robot_state": self.robot_state.map(|s| format!("{s:?}")), + "opmodes": opmodes, + "last_inited": self.last_inited, + "active_config": self.active_config, + "configs": self.configs, + "device_types": self.device_types, + "last_scan": self.last_scan, + "last_lynx_modules": self.last_lynx_modules, + "telemetry": telemetry, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::net::SocketAddr; + + #[test] + fn starts_disconnected_with_empty_caches() { + let state = DaemonState::new(); + assert!(!state.connected); + assert!(state.opmodes.is_empty()); + assert!(state.find_config("anything").is_none()); + assert!(!state.find_opmode("anything")); + } + + #[test] + fn connected_then_disconnected_clears_robot_state() { + let mut state = DaemonState::new(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state.apply_event(&Event::Connected { peer }); + state.apply_event(&Event::RobotState(RobotState::Running)); + assert!(state.connected); + assert_eq!(state.robot_state, Some(RobotState::Running)); + + state.apply_event(&Event::Disconnected); + assert!(!state.connected); + assert_eq!(state.robot_state, None); + assert_eq!(state.peer, None); + } + + #[test] + fn opmode_list_populates_find_opmode() { + let mut state = DaemonState::new(); + state.apply_event(&Event::OpModeList(vec![OpModeMeta { + name: "Duo (TeleOp)".into(), + flavor: "TELEOP".into(), + group: "drive".into(), + }])); + assert!(state.find_opmode("Duo (TeleOp)")); + assert!(!state.find_opmode("Nonexistent")); + } + + #[test] + fn opmode_inited_sets_last_inited() { + let mut state = DaemonState::new(); + state.apply_event(&Event::OpModeInited("Duo (TeleOp)".into())); + assert_eq!(state.last_inited.as_deref(), Some("Duo (TeleOp)")); + } + + #[test] + fn configuration_list_populates_find_config() { + let mut state = DaemonState::new(); + let json = r#"[{"isDirty":false,"location":"Local","name":"my_config","resourceId":1}]"#; + state.apply_event(&Event::ConfigurationList(json.to_string())); + let meta = state.find_config("my_config").expect("config found"); + assert_eq!(meta.location, "Local"); + assert_eq!(meta.resource_id, 1); + } + + #[test] + fn telemetry_keeps_latest_per_tag() { + let mut state = DaemonState::new(); + let first = Telemetry { + tag: "auto".into(), + strings: vec![("phase".into(), "start".into())], + ..Default::default() + }; + state.apply_event(&Event::Telemetry(first)); + + let second = Telemetry { + tag: "auto".into(), + strings: vec![("phase".into(), "end".into())], + ..Default::default() + }; + state.apply_event(&Event::Telemetry(second)); + + assert_eq!(state.telemetry.len(), 1); + assert_eq!(state.telemetry["auto"].strings[0].1, "end"); + } + + #[test] + fn to_json_reports_ok_true_and_connection_state() { + let state = DaemonState::new(); + let json = state.to_json(); + assert_eq!(json["ok"], true); + assert_eq!(json["connected"], false); + assert!(json["opmodes"].as_array().unwrap().is_empty()); + } +} From b5a25287795cbb5a31ad511ce4b1093e4bb7af4b Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:08:08 -0400 Subject: [PATCH 03/16] Add ds_agentd WaiterRegistry and Subscribers for event fan-out Co-Authored-By: Claude Haiku 4.5 --- ds_agentd/src/main.rs | 1 + ds_agentd/src/waiters.rs | 172 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 173 insertions(+) create mode 100644 ds_agentd/src/waiters.rs diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs index b6cfe31..51e347b 100644 --- a/ds_agentd/src/main.rs +++ b/ds_agentd/src/main.rs @@ -1,3 +1,4 @@ mod state; +mod waiters; fn main() {} diff --git a/ds_agentd/src/waiters.rs b/ds_agentd/src/waiters.rs new file mode 100644 index 0000000..4f13169 --- /dev/null +++ b/ds_agentd/src/waiters.rs @@ -0,0 +1,172 @@ +//! Two independent fan-out mechanisms sitting on top of the same +//! `robocol::Event` stream: `WaiterRegistry` resolves exactly one +//! request/response call per registration and then forgets it; +//! `Subscribers` keeps broadcasting to every `watch` connection until it +//! disconnects. + +use std::sync::mpsc::{self, Receiver, Sender}; + +use robocol::client::Event; + +#[allow(dead_code)] +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum WaiterKind { + OpModeList, + OpModeInited, + OpModeRunning, + ActiveConfiguration, + ConfigurationList, + Configuration, + UserDeviceList, + ScanResult, + LynxModules, +} + +#[allow(dead_code)] +pub fn waiter_kind_of(event: &Event) -> Option { + match event { + Event::OpModeList(_) => Some(WaiterKind::OpModeList), + Event::OpModeInited(_) => Some(WaiterKind::OpModeInited), + Event::OpModeRunning(_) => Some(WaiterKind::OpModeRunning), + Event::ActiveConfiguration(_) => Some(WaiterKind::ActiveConfiguration), + Event::ConfigurationList(_) => Some(WaiterKind::ConfigurationList), + Event::Configuration(_) => Some(WaiterKind::Configuration), + Event::UserDeviceList(_) => Some(WaiterKind::UserDeviceList), + Event::ScanResult(_) => Some(WaiterKind::ScanResult), + Event::LynxModules(_) => Some(WaiterKind::LynxModules), + _ => None, + } +} + +#[allow(dead_code)] +pub struct WaiterRegistry { + waiters: Vec<(WaiterKind, Sender)>, +} + +#[allow(dead_code)] +impl WaiterRegistry { + pub fn new() -> Self { + WaiterRegistry { + waiters: Vec::new(), + } + } + + pub fn register(&mut self, kind: WaiterKind) -> Receiver { + let (tx, rx) = mpsc::channel(); + self.waiters.push((kind, tx)); + rx + } + + /// Delivers `event` to every waiter registered for its kind, then drops + /// them — each registration is resolved (or forgotten, on send failure) + /// exactly once. Waiters of other kinds are left untouched. + pub fn resolve(&mut self, event: &Event) { + let Some(kind) = waiter_kind_of(event) else { + return; + }; + self.waiters.retain(|(k, tx)| { + if *k == kind { + let _ = tx.send(event.clone()); + false + } else { + true + } + }); + } +} + +#[allow(dead_code)] +pub struct Subscribers { + list: Vec>, +} + +#[allow(dead_code)] +impl Subscribers { + pub fn new() -> Self { + Subscribers { list: Vec::new() } + } + + pub fn subscribe(&mut self) -> Receiver { + let (tx, rx) = mpsc::channel(); + self.list.push(tx); + rx + } + + /// Sends `event` to every live subscriber, dropping any whose receiver + /// has gone away (their `watch` connection closed). + pub fn broadcast(&mut self, event: &Event) { + self.list.retain(|tx| tx.send(event.clone()).is_ok()); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn resolve_delivers_to_matching_kind_only() { + let mut registry = WaiterRegistry::new(); + let opmode_rx = registry.register(WaiterKind::OpModeList); + let config_rx = registry.register(WaiterKind::ActiveConfiguration); + + registry.resolve(&Event::OpModeList(vec![])); + + assert!(opmode_rx.try_recv().is_ok()); + assert!(config_rx.try_recv().is_err()); + } + + #[test] + fn resolve_forgets_waiter_after_one_delivery() { + let mut registry = WaiterRegistry::new(); + let rx = registry.register(WaiterKind::ScanResult); + + registry.resolve(&Event::ScanResult("first".into())); + registry.resolve(&Event::ScanResult("second".into())); + + assert_eq!(rx.try_recv().unwrap(), Event::ScanResult("first".into())); + assert!(rx.try_recv().is_err()); + } + + #[test] + fn resolve_ignores_events_with_no_waiter_kind() { + let mut registry = WaiterRegistry::new(); + let rx = registry.register(WaiterKind::OpModeList); + registry.resolve(&Event::Disconnected); + assert!(rx.try_recv().is_err()); + } + + #[test] + fn dropped_receiver_is_pruned_on_next_matching_event() { + let mut registry = WaiterRegistry::new(); + drop(registry.register(WaiterKind::ScanResult)); + // Resolving with a dropped receiver must not panic, and the dead + // entry must be gone afterward (checked indirectly: a second + // resolve with a fresh waiter only sees its own event). + registry.resolve(&Event::ScanResult("ignored".into())); + let rx = registry.register(WaiterKind::ScanResult); + registry.resolve(&Event::ScanResult("seen".into())); + assert_eq!(rx.try_recv().unwrap(), Event::ScanResult("seen".into())); + } + + #[test] + fn broadcast_reaches_every_subscriber() { + let mut subs = Subscribers::new(); + let a = subs.subscribe(); + let b = subs.subscribe(); + subs.broadcast(&Event::Disconnected); + assert_eq!(a.try_recv().unwrap(), Event::Disconnected); + assert_eq!(b.try_recv().unwrap(), Event::Disconnected); + } + + #[test] + fn broadcast_drops_disconnected_subscribers() { + let mut subs = Subscribers::new(); + { + let _rx = subs.subscribe(); // dropped immediately + } + let live = subs.subscribe(); + subs.broadcast(&Event::Disconnected); + assert_eq!(live.try_recv().unwrap(), Event::Disconnected); + assert_eq!(subs.list.len(), 1); + } +} From f5c5b25f85a44ad51501f6cf0893f7132cda0fc3 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:14:08 -0400 Subject: [PATCH 04/16] Add ds_agentd request dispatch for the full ds_cli command surface Implements the Request enum (serde-tagged on `cmd`) and handle_request, which runs each command against a live RobocolClient + DaemonState + WaiterRegistry and produces the JSON reply. Watch/Shutdown are left unreachable here since Task 5's connection layer intercepts them before calling handle_request. Co-Authored-By: Claude Sonnet 5 --- Cargo.lock | 1 + ds_agentd/Cargo.toml | 1 + ds_agentd/src/main.rs | 1 + ds_agentd/src/request.rs | 699 +++++++++++++++++++++++++++++++++++++++ 4 files changed, 702 insertions(+) create mode 100644 ds_agentd/src/request.rs diff --git a/Cargo.lock b/Cargo.lock index ca50bba..8015159 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -43,6 +43,7 @@ version = "0.1.0" dependencies = [ "ds_agent_ipc", "robocol", + "serde", "serde_json", ] diff --git a/ds_agentd/Cargo.toml b/ds_agentd/Cargo.toml index d1a1091..912192a 100644 --- a/ds_agentd/Cargo.toml +++ b/ds_agentd/Cargo.toml @@ -12,4 +12,5 @@ path = "src/main.rs" [dependencies] robocol = { path = "../robocol" } ds_agent_ipc = { path = "../ds_agent_ipc" } +serde = { version = "1", features = ["derive"] } serde_json = "1" diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs index 51e347b..9bf78d8 100644 --- a/ds_agentd/src/main.rs +++ b/ds_agentd/src/main.rs @@ -1,3 +1,4 @@ +mod request; mod state; mod waiters; diff --git a/ds_agentd/src/request.rs b/ds_agentd/src/request.rs new file mode 100644 index 0000000..b331707 --- /dev/null +++ b/ds_agentd/src/request.rs @@ -0,0 +1,699 @@ +//! The daemon's JSON protocol: parsing one request off the wire and turning +//! it into a `robocol` client call plus a reply, using `DaemonState` for +//! instant answers (`status`) and name lookups, and `WaiterRegistry` to +//! correlate a request with the `Event` that answers it. + +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use robocol::client::{Event, RobocolClient}; +use robocol::cmd::parse_config_list; +use robocol::packets::Gamepad; +use serde::Deserialize; + +use crate::state::DaemonState; +use crate::waiters::{WaiterKind, WaiterRegistry}; + +// Not yet constructed outside tests: Task 5's connection handling is what +// will call `serde_json::from_str::` on incoming socket lines. +#[allow(dead_code)] +#[derive(Debug, Deserialize)] +#[serde(tag = "cmd", rename_all = "kebab-case")] +pub enum Request { + Ping, + Status, + List, + Init { + name: String, + }, + Run { + name: Option, + }, + Stop, + Restart, + ActiveConfig, + Configs, + Config { + name: String, + }, + SaveConfig { + json: String, + }, + ActivateConfig { + name: String, + }, + DeleteConfig { + name: String, + }, + DeviceTypes, + Scan, + LynxModules { + serial: String, + }, + Gamepad { + #[serde(default)] + left_stick_x: Option, + #[serde(default)] + left_stick_y: Option, + #[serde(default)] + right_stick_x: Option, + #[serde(default)] + right_stick_y: Option, + #[serde(default)] + left_trigger: Option, + #[serde(default)] + right_trigger: Option, + #[serde(default)] + dpad_up: Option, + #[serde(default)] + dpad_down: Option, + #[serde(default)] + dpad_left: Option, + #[serde(default)] + dpad_right: Option, + #[serde(default)] + a: Option, + #[serde(default)] + b: Option, + #[serde(default)] + x: Option, + #[serde(default)] + y: Option, + #[serde(default)] + start: Option, + #[serde(default)] + back: Option, + #[serde(default)] + left_bumper: Option, + #[serde(default)] + right_bumper: Option, + #[serde(default)] + left_stick_button: Option, + #[serde(default)] + right_stick_button: Option, + #[serde(default)] + duration_ms: Option, + }, + Watch { + #[serde(default)] + types: Option>, + }, + Shutdown, +} + +// The functions below are only reachable via `handle_request`, which is +// itself only wired up to a real socket by Task 5 — allow dead_code until +// then, as Tasks 2/3 did for `state`/`waiters`. +#[allow(dead_code)] +fn err(code: &str, message: impl Into) -> serde_json::Value { + serde_json::json!({"ok": false, "error": code, "message": message.into()}) +} + +#[allow(dead_code)] +fn not_connected() -> serde_json::Value { + err("not_connected", "robot controller not connected yet") +} + +#[allow(dead_code)] +fn check_known_opmode(state: &Arc>, name: &str) -> Option { + let state = state.lock().expect("state mutex poisoned"); + if state.opmodes.is_empty() || state.find_opmode(name) { + None + } else { + Some(err( + "unknown_opmode", + format!("no opmode named {name:?}; run `list` first"), + )) + } +} + +/// Registers a waiter, fires `send` (which must issue exactly one +/// `robocol` request), and blocks up to `timeout` for the matching event. +/// Fails fast with `not_connected` rather than waiting out the timeout when +/// the daemon has never seen `Event::Connected`. +#[allow(dead_code)] +fn wait_for_event( + waiters: &Arc>, + state: &Arc>, + kind: WaiterKind, + timeout: Duration, + send: impl FnOnce(), +) -> Result { + if !state.lock().expect("state mutex poisoned").connected { + return Err(not_connected()); + } + let rx = waiters + .lock() + .expect("waiters mutex poisoned") + .register(kind); + send(); + rx.recv_timeout(timeout) + .map_err(|_| err("timeout", "robot controller did not respond in time")) +} + +#[allow(dead_code, clippy::too_many_arguments)] +fn build_gamepad( + left_stick_x: Option, + left_stick_y: Option, + right_stick_x: Option, + right_stick_y: Option, + left_trigger: Option, + right_trigger: Option, + dpad_up: Option, + dpad_down: Option, + dpad_left: Option, + dpad_right: Option, + a: Option, + b: Option, + x: Option, + y: Option, + start: Option, + back: Option, + left_bumper: Option, + right_bumper: Option, + left_stick_button: Option, + right_stick_button: Option, +) -> Gamepad { + Gamepad { + left_stick_x: left_stick_x.unwrap_or(0.0), + left_stick_y: left_stick_y.unwrap_or(0.0), + right_stick_x: right_stick_x.unwrap_or(0.0), + right_stick_y: right_stick_y.unwrap_or(0.0), + left_trigger: left_trigger.unwrap_or(0.0), + right_trigger: right_trigger.unwrap_or(0.0), + dpad_up: dpad_up.unwrap_or(false), + dpad_down: dpad_down.unwrap_or(false), + dpad_left: dpad_left.unwrap_or(false), + dpad_right: dpad_right.unwrap_or(false), + a: a.unwrap_or(false), + b: b.unwrap_or(false), + x: x.unwrap_or(false), + y: y.unwrap_or(false), + start: start.unwrap_or(false), + back: back.unwrap_or(false), + left_bumper: left_bumper.unwrap_or(false), + right_bumper: right_bumper.unwrap_or(false), + left_stick_button: left_stick_button.unwrap_or(false), + right_stick_button: right_stick_button.unwrap_or(false), + ..Gamepad::default() + } +} + +#[allow(dead_code)] +pub fn handle_request( + req: Request, + client: &Arc>, + state: &Arc>, + waiters: &Arc>, + timeout: Duration, +) -> serde_json::Value { + match req { + Request::Ping => serde_json::json!({"ok": true}), + Request::Status => state.lock().expect("state mutex poisoned").to_json(), + Request::List => { + let client = client.clone(); + match wait_for_event(waiters, state, WaiterKind::OpModeList, timeout, move || { + client + .lock() + .expect("client mutex poisoned") + .request_opmode_list(); + }) { + Ok(Event::OpModeList(list)) => serde_json::json!({ + "ok": true, + "opmodes": list.iter().map(|m| serde_json::json!({ + "name": m.name, "flavor": m.flavor, "group": m.group, + })).collect::>(), + }), + Ok(_) => { + unreachable!("WaiterKind::OpModeList only resolves with Event::OpModeList") + } + Err(e) => e, + } + } + Request::Init { name } => { + if let Some(e) = check_known_opmode(state, &name) { + return e; + } + let client = client.clone(); + let send_name = name.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::OpModeInited, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .init_opmode(&send_name); + }, + ) { + Ok(Event::OpModeInited(n)) => serde_json::json!({"ok": true, "name": n}), + Ok(_) => { + unreachable!("WaiterKind::OpModeInited only resolves with Event::OpModeInited") + } + Err(e) => e, + } + } + Request::Run { name } => { + let name = match name.or_else(|| { + state + .lock() + .expect("state mutex poisoned") + .last_inited + .clone() + }) { + Some(n) => n, + None => { + return err( + "bad_args", + "no opmode inited; call `init ` first or pass a name", + ); + } + }; + if let Some(e) = check_known_opmode(state, &name) { + return e; + } + let client = client.clone(); + let send_name = name.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::OpModeRunning, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .run_opmode(&send_name); + }, + ) { + Ok(Event::OpModeRunning(n)) => serde_json::json!({"ok": true, "name": n}), + Ok(_) => unreachable!( + "WaiterKind::OpModeRunning only resolves with Event::OpModeRunning" + ), + Err(e) => e, + } + } + Request::Stop => { + client.lock().expect("client mutex poisoned").stop_opmode(); + serde_json::json!({"ok": true}) + } + Request::Restart => { + client + .lock() + .expect("client mutex poisoned") + .restart_robot(); + serde_json::json!({"ok": true}) + } + Request::ActiveConfig => { + let client = client.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::ActiveConfiguration, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .request_active_config(); + }, + ) { + Ok(Event::ActiveConfiguration(extra)) => { + serde_json::json!({"ok": true, "config": extra}) + } + Ok(_) => unreachable!( + "WaiterKind::ActiveConfiguration only resolves with Event::ActiveConfiguration" + ), + Err(e) => e, + } + } + Request::Configs => { + let client = client.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::ConfigurationList, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .request_configurations(); + }, + ) { + Ok(Event::ConfigurationList(extra)) => { + serde_json::json!({"ok": true, "configs": parse_config_list(&extra)}) + } + Ok(_) => unreachable!( + "WaiterKind::ConfigurationList only resolves with Event::ConfigurationList" + ), + Err(e) => e, + } + } + Request::Config { name } => { + let Some(meta) = state + .lock() + .expect("state mutex poisoned") + .find_config(&name) + else { + return err( + "unknown_config", + format!("unknown config {name:?}; run `configs` first"), + ); + }; + let client = client.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::Configuration, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .request_particular_configuration(&meta); + }, + ) { + Ok(Event::Configuration(extra)) => serde_json::json!({"ok": true, "config": extra}), + Ok(_) => unreachable!( + "WaiterKind::Configuration only resolves with Event::Configuration" + ), + Err(e) => e, + } + } + Request::SaveConfig { json } => { + client + .lock() + .expect("client mutex poisoned") + .save_configuration(&json); + serde_json::json!({"ok": true}) + } + Request::ActivateConfig { name } => { + let Some(meta) = state + .lock() + .expect("state mutex poisoned") + .find_config(&name) + else { + return err( + "unknown_config", + format!("unknown config {name:?}; run `configs` first"), + ); + }; + client + .lock() + .expect("client mutex poisoned") + .activate_configuration(&meta); + serde_json::json!({"ok": true}) + } + Request::DeleteConfig { name } => { + let Some(meta) = state + .lock() + .expect("state mutex poisoned") + .find_config(&name) + else { + return err( + "unknown_config", + format!("unknown config {name:?}; run `configs` first"), + ); + }; + client + .lock() + .expect("client mutex poisoned") + .delete_configuration(&meta); + serde_json::json!({"ok": true}) + } + Request::DeviceTypes => { + let client = client.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::UserDeviceList, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .request_user_device_types(); + }, + ) { + Ok(Event::UserDeviceList(extra)) => { + serde_json::json!({"ok": true, "device_types": extra}) + } + Ok(_) => unreachable!( + "WaiterKind::UserDeviceList only resolves with Event::UserDeviceList" + ), + Err(e) => e, + } + } + Request::Scan => { + let client = client.clone(); + match wait_for_event(waiters, state, WaiterKind::ScanResult, timeout, move || { + client.lock().expect("client mutex poisoned").scan(); + }) { + Ok(Event::ScanResult(extra)) => serde_json::json!({"ok": true, "scan": extra}), + Ok(_) => { + unreachable!("WaiterKind::ScanResult only resolves with Event::ScanResult") + } + Err(e) => e, + } + } + Request::LynxModules { serial } => { + let client = client.clone(); + let send_serial = serial.clone(); + match wait_for_event( + waiters, + state, + WaiterKind::LynxModules, + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .discover_lynx_modules(&send_serial); + }, + ) { + Ok(Event::LynxModules(extra)) => { + serde_json::json!({"ok": true, "serial": serial, "modules": extra}) + } + Ok(_) => { + unreachable!("WaiterKind::LynxModules only resolves with Event::LynxModules") + } + Err(e) => e, + } + } + Request::Gamepad { + left_stick_x, + left_stick_y, + right_stick_x, + right_stick_y, + left_trigger, + right_trigger, + dpad_up, + dpad_down, + dpad_left, + dpad_right, + a, + b, + x, + y, + start, + back, + left_bumper, + right_bumper, + left_stick_button, + right_stick_button, + duration_ms, + } => { + let gamepad = build_gamepad( + left_stick_x, + left_stick_y, + right_stick_x, + right_stick_y, + left_trigger, + right_trigger, + dpad_up, + dpad_down, + dpad_left, + dpad_right, + a, + b, + x, + y, + start, + back, + left_bumper, + right_bumper, + left_stick_button, + right_stick_button, + ); + client + .lock() + .expect("client mutex poisoned") + .send_gamepad(gamepad); + if let Some(ms) = duration_ms { + std::thread::sleep(Duration::from_millis(ms)); + client + .lock() + .expect("client mutex poisoned") + .send_gamepad(Gamepad::default()); + } + serde_json::json!({"ok": true}) + } + // Watch and Shutdown are handled by the connection layer (Task 5), + // which needs the raw stream to push a live stream / exit the + // process; they never reach this dispatcher. + Request::Watch { .. } | Request::Shutdown => { + unreachable!("Watch and Shutdown are intercepted before handle_request is called") + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::net::SocketAddr; + use std::sync::mpsc; + + fn deps() -> (Arc>, Arc>) { + ( + Arc::new(Mutex::new(DaemonState::new())), + Arc::new(Mutex::new(WaiterRegistry::new())), + ) + } + + fn unreachable_client() -> Arc> { + // Binds a real socket but points at an address nothing answers, so + // it never emits Event::Connected — exactly the "not connected" case. + let config = robocol::ClientConfig { + bind_port: 0, + peer_addrs: vec!["127.0.0.1".parse().unwrap()], + peer_port: 1, // nothing listens on port 1 + ..Default::default() + }; + let (client, _events) = RobocolClient::start(config).expect("bind UDP socket"); + Arc::new(Mutex::new(client)) + } + + #[test] + fn status_answers_without_touching_the_client_or_waiters() { + let (state, waiters) = deps(); + let client = unreachable_client(); + let reply = handle_request( + Request::Status, + &client, + &state, + &waiters, + Duration::from_millis(50), + ); + assert_eq!(reply["ok"], true); + assert_eq!(reply["connected"], false); + } + + #[test] + fn list_fails_fast_with_not_connected() { + let (state, waiters) = deps(); + let client = unreachable_client(); + let reply = handle_request( + Request::List, + &client, + &state, + &waiters, + Duration::from_secs(5), + ); + assert_eq!(reply["ok"], false); + assert_eq!(reply["error"], "not_connected"); + } + + #[test] + fn init_rejects_unknown_opmode_without_waiting() { + let (state, waiters) = deps(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state + .lock() + .unwrap() + .apply_event(&Event::Connected { peer }); + state + .lock() + .unwrap() + .apply_event(&Event::OpModeList(vec![robocol::cmd::OpModeMeta { + name: "Duo (TeleOp)".into(), + flavor: "TELEOP".into(), + group: "drive".into(), + }])); + let client = unreachable_client(); + let reply = handle_request( + Request::Init { + name: "Nonexistent".into(), + }, + &client, + &state, + &waiters, + Duration::from_millis(50), + ); + assert_eq!(reply["ok"], false); + assert_eq!(reply["error"], "unknown_opmode"); + } + + #[test] + fn list_resolves_once_the_matching_event_arrives() { + let (state, waiters) = deps(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state + .lock() + .unwrap() + .apply_event(&Event::Connected { peer }); + let client = unreachable_client(); + + // Simulate the robot's reply landing on the waiter registry shortly + // after the request goes out, the way the real event-pump thread + // would (Task 6 wires that thread up for real). + let waiters_bg = waiters.clone(); + let (ready_tx, ready_rx) = mpsc::channel::<()>(); + std::thread::spawn(move || { + ready_rx.recv().unwrap(); + std::thread::sleep(Duration::from_millis(20)); + waiters_bg.lock().unwrap().resolve(&Event::OpModeList(vec![ + robocol::cmd::OpModeMeta { + name: "Solo (TeleOp)".into(), + flavor: "TELEOP".into(), + group: "drive".into(), + }, + ])); + }); + + // handle_request's `send` closure runs before recv_timeout blocks, + // so signal the background thread right as we call in. + ready_tx.send(()).unwrap(); + let reply = handle_request( + Request::List, + &client, + &state, + &waiters, + Duration::from_secs(1), + ); + assert_eq!(reply["ok"], true); + assert_eq!(reply["opmodes"][0]["name"], "Solo (TeleOp)"); + } + + #[test] + fn list_times_out_when_nothing_resolves_it() { + let (state, waiters) = deps(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state + .lock() + .unwrap() + .apply_event(&Event::Connected { peer }); + let client = unreachable_client(); + let reply = handle_request( + Request::List, + &client, + &state, + &waiters, + Duration::from_millis(50), + ); + assert_eq!(reply["ok"], false); + assert_eq!(reply["error"], "timeout"); + } +} From 625ec49c8711a83406f68a34c4ab4b5ad0049766 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:33:18 -0400 Subject: [PATCH 05/16] Add ds_agentd connection handling: watch streaming and shutdown Co-Authored-By: Claude Sonnet 5 --- ds_agentd/src/connection.rs | 241 ++++++++++++++++++++++++++++++++++++ ds_agentd/src/main.rs | 1 + 2 files changed, 242 insertions(+) create mode 100644 ds_agentd/src/connection.rs diff --git a/ds_agentd/src/connection.rs b/ds_agentd/src/connection.rs new file mode 100644 index 0000000..4399b5c --- /dev/null +++ b/ds_agentd/src/connection.rs @@ -0,0 +1,241 @@ +//! Per-connection handling: reads one request line, and either streams +//! `watch` events, exits the process for `shutdown`, or dispatches +//! everything else through `request::handle_request` for a single reply. + +use std::io::BufReader; +use std::os::unix::net::UnixStream; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use robocol::client::{Event, RobocolClient}; + +use crate::request::{Request, handle_request}; +use crate::state::DaemonState; +use crate::waiters::{Subscribers, WaiterRegistry}; + +// Not yet wired up outside tests: Task 6's accept loop is what will spawn a +// thread per connection calling `handle_connection`, as Tasks 2-4 noted for +// `state`/`waiters`/`request`. +#[allow(dead_code)] +fn event_type_name(event: &Event) -> Option<&'static str> { + match event { + Event::Connected { .. } => Some("connected"), + Event::Disconnected => Some("disconnected"), + Event::RobotState(_) => Some("state"), + Event::Telemetry(_) => Some("telemetry"), + Event::Stacktrace(_) => Some("stacktrace"), + Event::CommandDropped { .. } => Some("command_dropped"), + Event::ProtocolError(_) => Some("protocol_error"), + Event::WebcamAvailable(_) => Some("webcam_available"), + _ => None, + } +} + +#[allow(dead_code)] +fn event_to_watch_json(event: &Event, types: Option<&[String]>) -> Option { + let kind = event_type_name(event)?; + if let Some(types) = types + && !types.iter().any(|t| t == kind) + { + return None; + } + Some(match event { + Event::Connected { peer } => serde_json::json!({"event": kind, "peer": peer.to_string()}), + Event::Disconnected => serde_json::json!({"event": kind}), + Event::RobotState(s) => serde_json::json!({"event": kind, "state": format!("{s:?}")}), + Event::Telemetry(t) => { + let strings: serde_json::Map = t + .strings + .iter() + .map(|(k, v)| (k.clone(), serde_json::Value::String(v.clone()))) + .collect(); + let numbers: serde_json::Map = t + .numbers + .iter() + .map(|(k, v)| (k.clone(), serde_json::json!(v))) + .collect(); + serde_json::json!({"event": kind, "tag": t.tag, "strings": strings, "numbers": numbers}) + } + Event::Stacktrace(trace) => serde_json::json!({"event": kind, "trace": trace}), + Event::CommandDropped { name } => serde_json::json!({"event": kind, "name": name}), + Event::ProtocolError(message) => serde_json::json!({"event": kind, "message": message}), + Event::WebcamAvailable(available) => { + serde_json::json!({"event": kind, "available": available}) + } + _ => return None, + }) +} + +#[allow(dead_code)] +fn run_watch( + stream: &mut UnixStream, + subscribers: &Arc>, + types: Option>, +) { + let rx = subscribers + .lock() + .expect("subscribers mutex poisoned") + .subscribe(); + for event in rx { + if let Some(json) = event_to_watch_json(&event, types.as_deref()) + && ds_agent_ipc::write_message(stream, &json).is_err() + { + break; + } + } +} + +#[allow(dead_code)] +pub fn handle_connection( + mut stream: UnixStream, + client: Arc>, + state: Arc>, + waiters: Arc>, + subscribers: Arc>, + timeout: Duration, +) { + let reader_stream = match stream.try_clone() { + Ok(s) => s, + Err(_) => return, + }; + let mut reader = BufReader::new(reader_stream); + let Ok(Some(line)) = ds_agent_ipc::read_message(&mut reader) else { + return; + }; + let req: Request = match serde_json::from_str(&line) { + Ok(r) => r, + Err(e) => { + let _ = ds_agent_ipc::write_message( + &mut stream, + &serde_json::json!({"ok": false, "error": "bad_args", "message": e.to_string()}), + ); + return; + } + }; + match req { + Request::Watch { types } => run_watch(&mut stream, &subscribers, types), + Request::Shutdown => { + let _ = ds_agent_ipc::write_message(&mut stream, &serde_json::json!({"ok": true})); + std::process::exit(0); + } + other => { + let reply = handle_request(other, &client, &state, &waiters, timeout); + let _ = ds_agent_ipc::write_message(&mut stream, &reply); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::io::BufRead; + + #[allow(clippy::type_complexity)] + fn deps() -> ( + Arc>, + Arc>, + Arc>, + Arc>, + ) { + let config = robocol::ClientConfig { + bind_port: 0, + peer_port: 1, + ..Default::default() + }; + let (client, _events) = RobocolClient::start(config).expect("bind UDP socket"); + ( + Arc::new(Mutex::new(client)), + Arc::new(Mutex::new(DaemonState::new())), + Arc::new(Mutex::new(WaiterRegistry::new())), + Arc::new(Mutex::new(Subscribers::new())), + ) + } + + #[test] + fn ping_round_trips_over_a_real_socket_pair() { + let (client, state, waiters, subscribers) = deps(); + let (mut cli_side, daemon_side) = UnixStream::pair().unwrap(); + ds_agent_ipc::write_message(&mut cli_side, &serde_json::json!({"cmd": "ping"})).unwrap(); + handle_connection( + daemon_side, + client, + state, + waiters, + subscribers, + Duration::from_millis(50), + ); + + let mut reader = BufReader::new(cli_side); + let mut line = String::new(); + reader.read_line(&mut line).unwrap(); + let reply: serde_json::Value = serde_json::from_str(line.trim_end()).unwrap(); + assert_eq!(reply["ok"], true); + } + + #[test] + fn bad_json_yields_bad_args() { + let (client, state, waiters, subscribers) = deps(); + let (mut cli_side, daemon_side) = UnixStream::pair().unwrap(); + cli_side.set_nonblocking(false).unwrap(); + use std::io::Write; + cli_side.write_all(b"not json\n").unwrap(); + handle_connection( + daemon_side, + client, + state, + waiters, + subscribers, + Duration::from_millis(50), + ); + + let mut reader = BufReader::new(cli_side); + let mut line = String::new(); + reader.read_line(&mut line).unwrap(); + let reply: serde_json::Value = serde_json::from_str(line.trim_end()).unwrap(); + assert_eq!(reply["ok"], false); + assert_eq!(reply["error"], "bad_args"); + } + + #[test] + fn watch_streams_only_requested_types() { + let (client, state, waiters, subscribers) = deps(); + let (mut cli_side, daemon_side) = UnixStream::pair().unwrap(); + ds_agent_ipc::write_message( + &mut cli_side, + &serde_json::json!({"cmd": "watch", "types": ["disconnected"]}), + ) + .unwrap(); + + let subs = subscribers.clone(); + let handle = std::thread::spawn(move || { + handle_connection( + daemon_side, + client, + state, + waiters, + subs, + Duration::from_secs(1), + ); + }); + + // Give handle_connection time to subscribe before broadcasting. + std::thread::sleep(Duration::from_millis(50)); + subscribers + .lock() + .unwrap() + .broadcast(&Event::RobotState(robocol::types::RobotState::Running)); + subscribers.lock().unwrap().broadcast(&Event::Disconnected); + + let mut reader = BufReader::new(cli_side); + let mut line = String::new(); + reader.read_line(&mut line).unwrap(); + let event: serde_json::Value = serde_json::from_str(line.trim_end()).unwrap(); + assert_eq!(event["event"], "disconnected"); + + drop(reader); // closes the CLI side + // run_watch only notices a closed peer when its next write fails, so + // push one more event through to wake it and let it exit. + subscribers.lock().unwrap().broadcast(&Event::Disconnected); + handle.join().unwrap(); + } +} diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs index 9bf78d8..dda5e0e 100644 --- a/ds_agentd/src/main.rs +++ b/ds_agentd/src/main.rs @@ -1,3 +1,4 @@ +mod connection; mod request; mod state; mod waiters; From 3f176da341fd722a5ebc3cfb9d34c1a49d4dcf71 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:39:30 -0400 Subject: [PATCH 06/16] Wire ds_agentd into a runnable daemon: accept loop and process lifecycle Replaces the main.rs stub with argument parsing (--peer/--peer-port/ --bind-port), a RobocolClient event-pump thread that feeds DaemonState, WaiterRegistry, and Subscribers in order, and a Unix-socket accept loop that spawns one thread per connection. bind_socket handles the stale- socket (crashed prior run) and already-running (live daemon) cases. Also removes the #[allow(dead_code)] annotations Tasks 2-5 left on state/waiters/request/connection now that main.rs calls everything; clippy --all-targets -D warnings is clean with none of them restored. Co-Authored-By: Claude Sonnet 5 --- ds_agentd/src/connection.rs | 4 -- ds_agentd/src/main.rs | 136 +++++++++++++++++++++++++++++++++++- ds_agentd/src/request.rs | 6 -- ds_agentd/src/state.rs | 2 - ds_agentd/src/waiters.rs | 6 -- 5 files changed, 135 insertions(+), 19 deletions(-) diff --git a/ds_agentd/src/connection.rs b/ds_agentd/src/connection.rs index 4399b5c..0f40e71 100644 --- a/ds_agentd/src/connection.rs +++ b/ds_agentd/src/connection.rs @@ -16,7 +16,6 @@ use crate::waiters::{Subscribers, WaiterRegistry}; // Not yet wired up outside tests: Task 6's accept loop is what will spawn a // thread per connection calling `handle_connection`, as Tasks 2-4 noted for // `state`/`waiters`/`request`. -#[allow(dead_code)] fn event_type_name(event: &Event) -> Option<&'static str> { match event { Event::Connected { .. } => Some("connected"), @@ -31,7 +30,6 @@ fn event_type_name(event: &Event) -> Option<&'static str> { } } -#[allow(dead_code)] fn event_to_watch_json(event: &Event, types: Option<&[String]>) -> Option { let kind = event_type_name(event)?; if let Some(types) = types @@ -66,7 +64,6 @@ fn event_to_watch_json(event: &Event, types: Option<&[String]>) -> Option>, @@ -85,7 +82,6 @@ fn run_watch( } } -#[allow(dead_code)] pub fn handle_connection( mut stream: UnixStream, client: Arc>, diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs index dda5e0e..c7eb904 100644 --- a/ds_agentd/src/main.rs +++ b/ds_agentd/src/main.rs @@ -1,6 +1,140 @@ +//! Headless daemon: owns one persistent `robocol::RobocolClient` connection +//! and answers `ds_agent` over a Unix domain socket. See +//! `docs/superpowers/specs/2026-09-10-headless-driver-station-agent-design.md` +//! in the deck-station repo for the wire protocol this implements. + mod connection; mod request; mod state; mod waiters; -fn main() {} +use std::net::IpAddr; +use std::os::unix::net::{UnixListener, UnixStream}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use robocol::{ClientConfig, RobocolClient}; + +use connection::handle_connection; +use state::DaemonState; +use waiters::{Subscribers, WaiterRegistry}; + +const DEFAULT_TIMEOUT: Duration = Duration::from_secs(5); + +struct Args { + peer: Option, + peer_port: Option, + bind_port: Option, +} + +fn parse_args() -> Args { + let mut args = Args { + peer: None, + peer_port: None, + bind_port: None, + }; + let raw: Vec = std::env::args().skip(1).collect(); + let mut i = 0; + while i < raw.len() { + match raw[i].as_str() { + "--peer" => { + args.peer = raw.get(i + 1).and_then(|s| s.parse().ok()); + i += 2; + } + "--peer-port" => { + args.peer_port = raw.get(i + 1).and_then(|s| s.parse().ok()); + i += 2; + } + "--bind-port" => { + args.bind_port = raw.get(i + 1).and_then(|s| s.parse().ok()); + i += 2; + } + _ => i += 1, + } + } + args +} + +/// Timeout for request/response commands. `DS_AGENT_TIMEOUT_MS` lets tests +/// use something far shorter than the 5s production default. +fn resolve_timeout() -> Duration { + std::env::var("DS_AGENT_TIMEOUT_MS") + .ok() + .and_then(|s| s.parse().ok()) + .map(Duration::from_millis) + .unwrap_or(DEFAULT_TIMEOUT) +} + +/// Binds the daemon's Unix socket, clearing a stale one left behind by a +/// crashed previous instance. If a live daemon already owns the path, exits +/// the process immediately rather than running two daemons against one +/// robot connection. +fn bind_socket() -> UnixListener { + let path = ds_agent_ipc::socket_path(); + match UnixListener::bind(&path) { + Ok(listener) => listener, + Err(_) => { + if UnixStream::connect(&path).is_ok() { + eprintln!("ds_agentd already running at {}", path.display()); + std::process::exit(0); + } + let _ = std::fs::remove_file(&path); + UnixListener::bind(&path).expect("bind ds_agentd socket after clearing stale file") + } + } +} + +fn main() { + let args = parse_args(); + let mut config = ClientConfig::default(); + if let Some(peer) = args.peer { + config.peer_addrs = vec![peer]; + } + if let Some(port) = args.peer_port { + config.peer_port = port; + } + if let Some(port) = args.bind_port { + config.bind_port = port; + } + + let (client, events) = RobocolClient::start(config).expect("bind UDP socket"); + let client = Arc::new(Mutex::new(client)); + let state = Arc::new(Mutex::new(DaemonState::new())); + let waiters = Arc::new(Mutex::new(WaiterRegistry::new())); + let subscribers = Arc::new(Mutex::new(Subscribers::new())); + + { + let state = state.clone(); + let waiters = waiters.clone(); + let subscribers = subscribers.clone(); + std::thread::spawn(move || { + for event in events { + state + .lock() + .expect("state mutex poisoned") + .apply_event(&event); + waiters + .lock() + .expect("waiters mutex poisoned") + .resolve(&event); + subscribers + .lock() + .expect("subscribers mutex poisoned") + .broadcast(&event); + } + }); + } + + let timeout = resolve_timeout(); + let listener = bind_socket(); + for incoming in listener.incoming() { + let Ok(stream) = incoming else { continue }; + let client = client.clone(); + let state = state.clone(); + let waiters = waiters.clone(); + let subscribers = subscribers.clone(); + std::thread::spawn(move || { + handle_connection(stream, client, state, waiters, subscribers, timeout); + }); + } +} diff --git a/ds_agentd/src/request.rs b/ds_agentd/src/request.rs index b331707..6b21947 100644 --- a/ds_agentd/src/request.rs +++ b/ds_agentd/src/request.rs @@ -16,7 +16,6 @@ use crate::waiters::{WaiterKind, WaiterRegistry}; // Not yet constructed outside tests: Task 5's connection handling is what // will call `serde_json::from_str::` on incoming socket lines. -#[allow(dead_code)] #[derive(Debug, Deserialize)] #[serde(tag = "cmd", rename_all = "kebab-case")] pub enum Request { @@ -104,17 +103,14 @@ pub enum Request { // The functions below are only reachable via `handle_request`, which is // itself only wired up to a real socket by Task 5 — allow dead_code until // then, as Tasks 2/3 did for `state`/`waiters`. -#[allow(dead_code)] fn err(code: &str, message: impl Into) -> serde_json::Value { serde_json::json!({"ok": false, "error": code, "message": message.into()}) } -#[allow(dead_code)] fn not_connected() -> serde_json::Value { err("not_connected", "robot controller not connected yet") } -#[allow(dead_code)] fn check_known_opmode(state: &Arc>, name: &str) -> Option { let state = state.lock().expect("state mutex poisoned"); if state.opmodes.is_empty() || state.find_opmode(name) { @@ -131,7 +127,6 @@ fn check_known_opmode(state: &Arc>, name: &str) -> Option>, state: &Arc>, @@ -199,7 +194,6 @@ fn build_gamepad( } } -#[allow(dead_code)] pub fn handle_request( req: Request, client: &Arc>, diff --git a/ds_agentd/src/state.rs b/ds_agentd/src/state.rs index 4388958..09b4853 100644 --- a/ds_agentd/src/state.rs +++ b/ds_agentd/src/state.rs @@ -10,7 +10,6 @@ use robocol::cmd::{ConfigMeta, OpModeMeta, parse_config_list}; use robocol::packets::Telemetry; use robocol::types::RobotState; -#[allow(dead_code)] pub struct DaemonState { pub connected: bool, pub peer: Option, @@ -32,7 +31,6 @@ pub struct DaemonState { pub telemetry: HashMap, } -#[allow(dead_code)] impl DaemonState { pub fn new() -> Self { DaemonState { diff --git a/ds_agentd/src/waiters.rs b/ds_agentd/src/waiters.rs index 4f13169..525136f 100644 --- a/ds_agentd/src/waiters.rs +++ b/ds_agentd/src/waiters.rs @@ -8,7 +8,6 @@ use std::sync::mpsc::{self, Receiver, Sender}; use robocol::client::Event; -#[allow(dead_code)] #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum WaiterKind { OpModeList, @@ -22,7 +21,6 @@ pub enum WaiterKind { LynxModules, } -#[allow(dead_code)] pub fn waiter_kind_of(event: &Event) -> Option { match event { Event::OpModeList(_) => Some(WaiterKind::OpModeList), @@ -38,12 +36,10 @@ pub fn waiter_kind_of(event: &Event) -> Option { } } -#[allow(dead_code)] pub struct WaiterRegistry { waiters: Vec<(WaiterKind, Sender)>, } -#[allow(dead_code)] impl WaiterRegistry { pub fn new() -> Self { WaiterRegistry { @@ -75,12 +71,10 @@ impl WaiterRegistry { } } -#[allow(dead_code)] pub struct Subscribers { list: Vec>, } -#[allow(dead_code)] impl Subscribers { pub fn new() -> Self { Subscribers { list: Vec::new() } From 544d7469ee940d2fb0002006e7331eadf45c2c6b Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:45:58 -0400 Subject: [PATCH 07/16] Add ds_agent: non-interactive JSON CLI for ds_agentd Co-Authored-By: Claude Sonnet 5 --- Cargo.lock | 8 + Cargo.toml | 2 +- ds_agent/Cargo.toml | 14 ++ ds_agent/src/main.rs | 363 +++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 386 insertions(+), 1 deletion(-) create mode 100644 ds_agent/Cargo.toml create mode 100644 ds_agent/src/main.rs diff --git a/Cargo.lock b/Cargo.lock index 8015159..4bcc067 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -29,6 +29,14 @@ dependencies = [ "error-code", ] +[[package]] +name = "ds_agent" +version = "0.1.0" +dependencies = [ + "ds_agent_ipc", + "serde_json", +] + [[package]] name = "ds_agent_ipc" version = "0.1.0" diff --git a/Cargo.toml b/Cargo.toml index 82e6526..6b5a258 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,3 @@ [workspace] resolver = "3" -members = ["robocol", "fake_rc", "ds_cli", "capture/decode", "ds_agent_ipc", "ds_agentd"] +members = ["robocol", "fake_rc", "ds_cli", "capture/decode", "ds_agent_ipc", "ds_agentd", "ds_agent"] diff --git a/ds_agent/Cargo.toml b/ds_agent/Cargo.toml new file mode 100644 index 0000000..89e3fda --- /dev/null +++ b/ds_agent/Cargo.toml @@ -0,0 +1,14 @@ +[package] +name = "ds_agent" +version = "0.1.0" +edition = "2024" +license = "MIT" +description = "Non-interactive JSON CLI for ds_agentd, for scripting or an AI agent" + +[[bin]] +name = "ds_agent" +path = "src/main.rs" + +[dependencies] +ds_agent_ipc = { path = "../ds_agent_ipc" } +serde_json = "1" diff --git a/ds_agent/src/main.rs b/ds_agent/src/main.rs new file mode 100644 index 0000000..144c19c --- /dev/null +++ b/ds_agent/src/main.rs @@ -0,0 +1,363 @@ +//! Non-interactive CLI for `ds_agentd`. One invocation, one JSON request, +//! one reply on stdout — see `ds_cli` for the interactive REPL this mirrors. +//! +//! ```sh +//! ds_agent list +//! ds_agent init "Duo (TeleOp)" +//! ds_agent run +//! ds_agent gamepad --left-stick-y -0.5 --duration 300ms +//! ds_agent watch --types telemetry +//! ds_agent daemon status +//! ``` + +use std::io::BufReader; +use std::os::unix::net::UnixStream; +use std::time::{Duration, Instant}; + +fn build_request(args: &[String]) -> Result { + let Some(cmd) = args.first() else { + return Err("usage: ds_agent [args]".into()); + }; + match cmd.as_str() { + "list" => Ok(serde_json::json!({"cmd": "list"})), + "init" => { + let name = args.get(1).ok_or("usage: ds_agent init ")?; + Ok(serde_json::json!({"cmd": "init", "name": name})) + } + "run" => Ok(serde_json::json!({"cmd": "run", "name": args.get(1)})), + "stop" => Ok(serde_json::json!({"cmd": "stop"})), + "restart" => Ok(serde_json::json!({"cmd": "restart"})), + "active-config" => Ok(serde_json::json!({"cmd": "active-config"})), + "configs" => Ok(serde_json::json!({"cmd": "configs"})), + "config" => { + let name = args.get(1).ok_or("usage: ds_agent config ")?; + Ok(serde_json::json!({"cmd": "config", "name": name})) + } + "save-config" => { + let json = args.get(1).ok_or("usage: ds_agent save-config ")?; + Ok(serde_json::json!({"cmd": "save-config", "json": json})) + } + "activate-config" => { + let name = args + .get(1) + .ok_or("usage: ds_agent activate-config ")?; + Ok(serde_json::json!({"cmd": "activate-config", "name": name})) + } + "delete-config" => { + let name = args.get(1).ok_or("usage: ds_agent delete-config ")?; + Ok(serde_json::json!({"cmd": "delete-config", "name": name})) + } + "device-types" => Ok(serde_json::json!({"cmd": "device-types"})), + "scan" => Ok(serde_json::json!({"cmd": "scan"})), + "lynx-modules" => { + let serial = args.get(1).ok_or("usage: ds_agent lynx-modules ")?; + Ok(serde_json::json!({"cmd": "lynx-modules", "serial": serial})) + } + "gamepad" => build_gamepad_request(&args[1..]), + "status" => Ok(serde_json::json!({"cmd": "status"})), + "watch" => build_watch_request(&args[1..]), + other => Err(format!("unknown command: {other}")), + } +} + +const GAMEPAD_FLOAT_FLAGS: &[&str] = &[ + "left_stick_x", + "left_stick_y", + "right_stick_x", + "right_stick_y", + "left_trigger", + "right_trigger", +]; +const GAMEPAD_BOOL_FLAGS: &[&str] = &[ + "dpad_up", + "dpad_down", + "dpad_left", + "dpad_right", + "a", + "b", + "x", + "y", + "start", + "back", + "left_bumper", + "right_bumper", + "left_stick_button", + "right_stick_button", +]; + +fn parse_duration_ms(value: &str) -> Result { + if let Some(ms) = value.strip_suffix("ms") { + ms.parse().map_err(|_| format!("invalid duration: {value}")) + } else if let Some(s) = value.strip_suffix('s') { + s.parse::() + .map(|secs| (secs * 1000.0) as u64) + .map_err(|_| format!("invalid duration: {value}")) + } else { + value + .parse() + .map_err(|_| format!("invalid duration: {value}")) + } +} + +fn build_gamepad_request(args: &[String]) -> Result { + let mut map = serde_json::Map::new(); + map.insert("cmd".into(), serde_json::json!("gamepad")); + let mut i = 0; + while i < args.len() { + let flag = args[i] + .strip_prefix("--") + .ok_or_else(|| format!("expected --flag, got {}", args[i]))?; + let key = flag.replace('-', "_"); + if key == "duration" { + let value = args.get(i + 1).ok_or("missing value for --duration")?; + map.insert( + "duration_ms".into(), + serde_json::json!(parse_duration_ms(value)?), + ); + i += 2; + continue; + } + let value = args + .get(i + 1) + .ok_or_else(|| format!("missing value for --{flag}"))?; + if GAMEPAD_FLOAT_FLAGS.contains(&key.as_str()) { + let parsed: f32 = value + .parse() + .map_err(|_| format!("--{flag} expects a number"))?; + map.insert(key, serde_json::json!(parsed)); + } else if GAMEPAD_BOOL_FLAGS.contains(&key.as_str()) { + let parsed: bool = value + .parse() + .map_err(|_| format!("--{flag} expects true/false"))?; + map.insert(key, serde_json::json!(parsed)); + } else if key == "duration_ms" { + let parsed: u64 = value + .parse() + .map_err(|_| format!("--{flag} expects milliseconds"))?; + map.insert(key, serde_json::json!(parsed)); + } else { + return Err(format!("unknown gamepad flag: --{flag}")); + } + i += 2; + } + Ok(serde_json::Value::Object(map)) +} + +fn build_watch_request(args: &[String]) -> Result { + let mut types: Option> = None; + let mut i = 0; + while i < args.len() { + match args[i].as_str() { + "--types" => { + let list = args + .get(i + 1) + .ok_or("usage: ds_agent watch [--types a,b,c]")?; + types = Some(list.split(',').map(str::to_string).collect()); + i += 2; + } + other => return Err(format!("unknown watch flag: {other}")), + } + } + Ok(serde_json::json!({"cmd": "watch", "types": types})) +} + +fn connect_or_spawn() -> Result { + let path = ds_agent_ipc::socket_path(); + if let Ok(stream) = UnixStream::connect(&path) { + return Ok(stream); + } + let exe = std::env::current_exe().map_err(|e| e.to_string())?; + let daemon_path = exe.with_file_name("ds_agentd"); + std::process::Command::new(&daemon_path) + .stdin(std::process::Stdio::null()) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .spawn() + .map_err(|e| format!("failed to start {}: {e}", daemon_path.display()))?; + + let deadline = Instant::now() + Duration::from_secs(3); + loop { + if let Ok(stream) = UnixStream::connect(&path) { + return Ok(stream); + } + if Instant::now() >= deadline { + return Err("ds_agentd did not become ready in time".into()); + } + std::thread::sleep(Duration::from_millis(100)); + } +} + +fn run_daemon_subcommand(args: &[String]) -> i32 { + let path = ds_agent_ipc::socket_path(); + match args.first().map(String::as_str) { + Some("status") => match UnixStream::connect(&path) { + Ok(mut stream) => { + let _ = + ds_agent_ipc::write_message(&mut stream, &serde_json::json!({"cmd": "status"})); + let mut reader = BufReader::new(stream); + match ds_agent_ipc::read_message(&mut reader) { + Ok(Some(line)) => { + println!("{line}"); + 0 + } + _ => { + println!(r#"{{"ok":false,"error":"daemon_unreachable"}}"#); + 1 + } + } + } + Err(_) => { + println!(r#"{{"ok":true,"running":false}}"#); + 0 + } + }, + Some("stop") => match UnixStream::connect(&path) { + Ok(mut stream) => { + let _ = ds_agent_ipc::write_message( + &mut stream, + &serde_json::json!({"cmd": "shutdown"}), + ); + println!(r#"{{"ok":true}}"#); + 0 + } + Err(_) => { + println!(r#"{{"ok":true,"running":false}}"#); + 0 + } + }, + _ => { + eprintln!("usage: ds_agent daemon "); + 2 + } + } +} + +fn main() { + let raw_args: Vec = std::env::args().skip(1).collect(); + if raw_args.first().map(String::as_str) == Some("daemon") { + std::process::exit(run_daemon_subcommand(&raw_args[1..])); + } + + let request = match build_request(&raw_args) { + Ok(r) => r, + Err(msg) => { + eprintln!("{msg}"); + std::process::exit(2); + } + }; + let is_watch = request["cmd"] == "watch"; + + let mut stream = match connect_or_spawn() { + Ok(s) => s, + Err(msg) => { + println!(r#"{{"ok":false,"error":"daemon_unreachable","message":{msg:?}}}"#); + std::process::exit(1); + } + }; + if ds_agent_ipc::write_message(&mut stream, &request).is_err() { + eprintln!("failed to write to ds_agentd"); + std::process::exit(1); + } + + let mut reader = BufReader::new(stream.try_clone().expect("clone stream")); + if is_watch { + while let Ok(Some(line)) = ds_agent_ipc::read_message(&mut reader) { + println!("{line}"); + } + return; + } + match ds_agent_ipc::read_message(&mut reader) { + Ok(Some(line)) => { + println!("{line}"); + let ok = serde_json::from_str::(&line) + .map(|v| v["ok"] == true) + .unwrap_or(false); + std::process::exit(if ok { 0 } else { 1 }); + } + _ => { + eprintln!("no response from ds_agentd"); + std::process::exit(1); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn list_needs_no_arguments() { + let req = build_request(&["list".to_string()]).unwrap(); + assert_eq!(req, serde_json::json!({"cmd": "list"})); + } + + #[test] + fn init_requires_a_name() { + assert!(build_request(&["init".to_string()]).is_err()); + let req = build_request(&["init".to_string(), "Duo (TeleOp)".to_string()]).unwrap(); + assert_eq!( + req, + serde_json::json!({"cmd": "init", "name": "Duo (TeleOp)"}) + ); + } + + #[test] + fn run_without_a_name_sends_null() { + let req = build_request(&["run".to_string()]).unwrap(); + assert_eq!(req, serde_json::json!({"cmd": "run", "name": null})); + } + + #[test] + fn unknown_command_is_an_error() { + assert!(build_request(&["not-a-command".to_string()]).is_err()); + } + + #[test] + fn gamepad_parses_floats_bools_and_duration() { + let args: Vec = [ + "gamepad", + "--left-stick-y", + "-0.5", + "--a", + "true", + "--duration", + "300ms", + ] + .iter() + .map(|s| s.to_string()) + .collect(); + let req = build_request(&args).unwrap(); + assert_eq!(req["cmd"], "gamepad"); + assert_eq!(req["left_stick_y"], -0.5); + assert_eq!(req["a"], true); + assert_eq!(req["duration_ms"], 300); + } + + #[test] + fn gamepad_rejects_unknown_flags() { + let args: Vec = ["gamepad", "--not-a-flag", "1"] + .iter() + .map(|s| s.to_string()) + .collect(); + assert!(build_request(&args).is_err()); + } + + #[test] + fn watch_parses_comma_separated_types() { + let args: Vec = ["watch", "--types", "telemetry,disconnected"] + .iter() + .map(|s| s.to_string()) + .collect(); + let req = build_request(&args).unwrap(); + assert_eq!( + req["types"], + serde_json::json!(["telemetry", "disconnected"]) + ); + } + + #[test] + fn watch_with_no_flags_sends_null_types() { + let req = build_request(&["watch".to_string()]).unwrap(); + assert_eq!(req["types"], serde_json::Value::Null); + } +} From 95f801fd448c08782fa737f111db8872858e274f Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:51:33 -0400 Subject: [PATCH 08/16] ds_agent: log spawned daemon output and mirror daemon_unreachable to stderr Spawned ds_agentd's stdout/stderr now go to ds_agentd.log next to the socket instead of /dev/null, and the timeout/spawn-failure messages name that log path, so a startup crash (e.g. port already bound) is diagnosable instead of just "did not become ready in time". The daemon_unreachable JSON error is now written to both stdout and stderr in main() and in `daemon status`, per spec. Co-Authored-By: Claude Sonnet 5 --- ds_agent/src/main.rs | 40 ++++++++++++++++++++++++++++++++++------ 1 file changed, 34 insertions(+), 6 deletions(-) diff --git a/ds_agent/src/main.rs b/ds_agent/src/main.rs index 144c19c..1d140f4 100644 --- a/ds_agent/src/main.rs +++ b/ds_agent/src/main.rs @@ -161,6 +161,9 @@ fn build_watch_request(args: &[String]) -> Result { Ok(serde_json::json!({"cmd": "watch", "types": types})) } +/// Connects to the daemon, spawning it if needed. A spawned daemon's stdout +/// and stderr are redirected to `ds_agentd.log` next to the socket, so a +/// crash or panic during startup can be diagnosed after the fact. fn connect_or_spawn() -> Result { let path = ds_agent_ipc::socket_path(); if let Ok(stream) = UnixStream::connect(&path) { @@ -168,12 +171,27 @@ fn connect_or_spawn() -> Result { } let exe = std::env::current_exe().map_err(|e| e.to_string())?; let daemon_path = exe.with_file_name("ds_agentd"); + let log_path = path.with_file_name("ds_agentd.log"); + let log = std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(&log_path) + .map_err(|e| format!("failed to open {}: {e}", log_path.display()))?; + let log_err = log + .try_clone() + .map_err(|e| format!("failed to open {}: {e}", log_path.display()))?; std::process::Command::new(&daemon_path) .stdin(std::process::Stdio::null()) - .stdout(std::process::Stdio::null()) - .stderr(std::process::Stdio::null()) + .stdout(std::process::Stdio::from(log)) + .stderr(std::process::Stdio::from(log_err)) .spawn() - .map_err(|e| format!("failed to start {}: {e}", daemon_path.display()))?; + .map_err(|e| { + format!( + "failed to start {}: {e}; see {}", + daemon_path.display(), + log_path.display() + ) + })?; let deadline = Instant::now() + Duration::from_secs(3); loop { @@ -181,7 +199,10 @@ fn connect_or_spawn() -> Result { return Ok(stream); } if Instant::now() >= deadline { - return Err("ds_agentd did not become ready in time".into()); + return Err(format!( + "ds_agentd did not become ready in time; see {}", + log_path.display() + )); } std::thread::sleep(Duration::from_millis(100)); } @@ -201,7 +222,10 @@ fn run_daemon_subcommand(args: &[String]) -> i32 { 0 } _ => { - println!(r#"{{"ok":false,"error":"daemon_unreachable"}}"#); + let line = serde_json::json!({"ok": false, "error": "daemon_unreachable"}) + .to_string(); + println!("{line}"); + eprintln!("{line}"); 1 } } @@ -250,7 +274,11 @@ fn main() { let mut stream = match connect_or_spawn() { Ok(s) => s, Err(msg) => { - println!(r#"{{"ok":false,"error":"daemon_unreachable","message":{msg:?}}}"#); + let line = + serde_json::json!({"ok": false, "error": "daemon_unreachable", "message": msg}) + .to_string(); + println!("{line}"); + eprintln!("{line}"); std::process::exit(1); } }; From 3ba506250d1c7c95139a7beade512af06d8d4e6e Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 22:56:15 -0400 Subject: [PATCH 09/16] Add end-to-end test driving ds_agent/ds_agentd against fake_rc Co-Authored-By: Claude Sonnet 5 --- ds_agentd/tests/end_to_end.rs | 166 ++++++++++++++++++++++++++++++++++ 1 file changed, 166 insertions(+) create mode 100644 ds_agentd/tests/end_to_end.rs diff --git a/ds_agentd/tests/end_to_end.rs b/ds_agentd/tests/end_to_end.rs new file mode 100644 index 0000000..4055b65 --- /dev/null +++ b/ds_agentd/tests/end_to_end.rs @@ -0,0 +1,166 @@ +//! Drives the full `ds_agent` -> `ds_agentd` -> `fake_rc` stack the way a +//! human would drive `ds_cli` -> a real Control Hub, but scripted. + +use std::io::{BufRead, BufReader}; +use std::path::PathBuf; +use std::process::{Child, Command, Stdio}; +use std::time::{Duration, Instant}; + +struct ChildGuard(Child); +impl Drop for ChildGuard { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +fn workspace_root() -> PathBuf { + PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .parent() + .expect("workspace root") + .to_path_buf() +} + +/// `cargo test -p ds_agentd` only builds this crate's own binaries, so build +/// the two siblings the test drives before touching them. +fn build_sibling_bins() { + let status = Command::new(env!("CARGO")) + .args(["build", "-p", "ds_agent", "-p", "fake_rc"]) + .current_dir(workspace_root()) + .status() + .expect("run cargo build for ds_agent/fake_rc"); + assert!(status.success(), "building ds_agent/fake_rc failed"); +} + +fn bin_path(name: &str) -> PathBuf { + // CARGO_BIN_EXE_ds_agentd is set for this crate's own binary; sibling + // binaries built by build_sibling_bins() share the same target/debug + // dir, one level up from ds_agentd's own exe path. + let mut path = PathBuf::from(env!("CARGO_BIN_EXE_ds_agentd")); + path.pop(); + path.push(name); + path +} + +fn run_agent(socket_dir: &std::path::Path, args: &[&str]) -> serde_json::Value { + let output = Command::new(bin_path("ds_agent")) + .args(args) + .env("XDG_RUNTIME_DIR", socket_dir) + .output() + .expect("run ds_agent"); + let stdout = String::from_utf8_lossy(&output.stdout); + serde_json::from_str(stdout.trim()).unwrap_or_else(|e| { + panic!( + "ds_agent {args:?} did not print JSON: {e}\nstdout: {stdout}\nstderr: {}", + String::from_utf8_lossy(&output.stderr) + ) + }) +} + +#[test] +fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { + build_sibling_bins(); + + let tmp = tempdir(); + let fake_rc_port: u16 = 20950; + + let mut fake_rc = ChildGuard( + Command::new(bin_path("fake_rc")) + .arg(fake_rc_port.to_string()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("start fake_rc"), + ); + + let mut daemon = ChildGuard( + Command::new(bin_path("ds_agentd")) + .args([ + "--peer", + "127.0.0.1", + "--peer-port", + &fake_rc_port.to_string(), + "--bind-port", + "0", + ]) + .env("XDG_RUNTIME_DIR", &tmp) + .env("DS_AGENT_TIMEOUT_MS", "2000") + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("start ds_agentd"), + ); + + // Wait for the daemon's socket to exist before sending anything. + let socket_path = tmp.join("ds_agentd.sock"); + let deadline = Instant::now() + Duration::from_secs(3); + while !socket_path.exists() { + assert!(Instant::now() < deadline, "ds_agentd socket never appeared"); + std::thread::sleep(Duration::from_millis(50)); + } + + // Wait for the daemon to actually connect to fake_rc (discovery + a + // couple of heartbeats) before issuing commands that depend on it. + let deadline = Instant::now() + Duration::from_secs(3); + loop { + let status = run_agent(&tmp, &["status"]); + if status["connected"] == true { + break; + } + assert!( + Instant::now() < deadline, + "ds_agentd never connected to fake_rc" + ); + std::thread::sleep(Duration::from_millis(100)); + } + + let list = run_agent(&tmp, &["list"]); + assert_eq!(list["ok"], true); + let opmodes = list["opmodes"].as_array().expect("opmodes array"); + assert!(opmodes.iter().any(|m| m["name"] == "Duo (TeleOp)")); + + let init = run_agent(&tmp, &["init", "Duo (TeleOp)"]); + assert_eq!(init["ok"], true); + assert_eq!(init["name"], "Duo (TeleOp)"); + + let run = run_agent(&tmp, &["run"]); // no name: falls back to last_inited + assert_eq!(run["ok"], true); + assert_eq!(run["name"], "Duo (TeleOp)"); + + let stop = run_agent(&tmp, &["stop"]); + assert_eq!(stop["ok"], true); + + // Watch: fake_rc streams telemetry at 10 Hz, so a short-lived watch + // process should see at least one line before it's killed. + let mut watch = Command::new(bin_path("ds_agent")) + .args(["watch", "--types", "telemetry"]) + .env("XDG_RUNTIME_DIR", &tmp) + .stdout(Stdio::piped()) + .spawn() + .expect("start ds_agent watch"); + let stdout = watch.stdout.take().expect("watch stdout"); + let mut reader = BufReader::new(stdout); + let mut line = String::new(); + reader.read_line(&mut line).expect("read a telemetry line"); + let event: serde_json::Value = + serde_json::from_str(line.trim_end()).expect("watch line is JSON"); + assert_eq!(event["event"], "telemetry"); + let _ = watch.kill(); + let _ = watch.wait(); + + // Explicit teardown via `daemon stop`, ahead of ChildGuard's Drop kill, + // proves the shutdown command itself works. + let stopped = run_agent(&tmp, &["daemon", "stop"]); + assert_eq!(stopped["ok"], true); + std::thread::sleep(Duration::from_millis(200)); + let _ = daemon.0.try_wait(); + let _ = fake_rc.0.kill(); + + std::fs::remove_dir_all(&tmp).expect("clean up runtime dir"); +} + +fn tempdir() -> PathBuf { + let dir = PathBuf::from(format!("/tmp/ds_agent_e2e_{}", std::process::id())); + std::fs::create_dir_all(&dir).expect("create temp runtime dir"); + dir +} From ae455bbc8b26d9007eb6bda1cdb0988a5847c109 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:00:41 -0400 Subject: [PATCH 10/16] e2e: guard the watch child so a failing run cannot orphan it Co-Authored-By: Claude Sonnet 5 --- ds_agentd/tests/end_to_end.rs | 37 ++++++++++++++++++++--------------- 1 file changed, 21 insertions(+), 16 deletions(-) diff --git a/ds_agentd/tests/end_to_end.rs b/ds_agentd/tests/end_to_end.rs index 4055b65..88c4fa9 100644 --- a/ds_agentd/tests/end_to_end.rs +++ b/ds_agentd/tests/end_to_end.rs @@ -131,22 +131,27 @@ fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { assert_eq!(stop["ok"], true); // Watch: fake_rc streams telemetry at 10 Hz, so a short-lived watch - // process should see at least one line before it's killed. - let mut watch = Command::new(bin_path("ds_agent")) - .args(["watch", "--types", "telemetry"]) - .env("XDG_RUNTIME_DIR", &tmp) - .stdout(Stdio::piped()) - .spawn() - .expect("start ds_agent watch"); - let stdout = watch.stdout.take().expect("watch stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - reader.read_line(&mut line).expect("read a telemetry line"); - let event: serde_json::Value = - serde_json::from_str(line.trim_end()).expect("watch line is JSON"); - assert_eq!(event["event"], "telemetry"); - let _ = watch.kill(); - let _ = watch.wait(); + // process should see at least one line before it's killed. Scoped in a + // block so the ChildGuard drops (killing the watch) before `daemon + // stop` below, and so a panic on any of the asserts still kills it + // rather than leaking an orphan `ds_agent watch` process. + { + let mut watch = ChildGuard( + Command::new(bin_path("ds_agent")) + .args(["watch", "--types", "telemetry"]) + .env("XDG_RUNTIME_DIR", &tmp) + .stdout(Stdio::piped()) + .spawn() + .expect("start ds_agent watch"), + ); + let stdout = watch.0.stdout.take().expect("watch stdout"); + let mut reader = BufReader::new(stdout); + let mut line = String::new(); + reader.read_line(&mut line).expect("read a telemetry line"); + let event: serde_json::Value = + serde_json::from_str(line.trim_end()).expect("watch line is JSON"); + assert_eq!(event["event"], "telemetry"); + } // Explicit teardown via `daemon stop`, ahead of ChildGuard's Drop kill, // proves the shutdown command itself works. From f31feafa660353e51bca88793821b8ee27efcd18 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:15:57 -0400 Subject: [PATCH 11/16] ds_agentd: match waiters on payload and gate every command on connected MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit WaiterRegistry::register now takes a predicate so an unsolicited or stale same-kind event (the RC's NOTIFY_INIT_OP_MODE $Stop$Robot$ after a stop) can't resolve an unrelated request. Init/Run match on the opmode name; everything else accepts any event of its kind (LynxModules can't be payload-matched — its response doesn't carry the serial). handle_request now fails fast with not_connected for every command except ping/status, via one ensure_connected gate ahead of dispatch, instead of letting stop/restart/config writes/gamepad go out with no peer. Also drops the stale 'Task N' comments and the dead_code allow. Co-Authored-By: Claude Opus 5 --- ds_agentd/src/request.rs | 191 +++++++++++++++++++++++++++++++-------- ds_agentd/src/waiters.rs | 107 ++++++++++++++++++---- 2 files changed, 240 insertions(+), 58 deletions(-) diff --git a/ds_agentd/src/request.rs b/ds_agentd/src/request.rs index 6b21947..01e803a 100644 --- a/ds_agentd/src/request.rs +++ b/ds_agentd/src/request.rs @@ -12,10 +12,8 @@ use robocol::packets::Gamepad; use serde::Deserialize; use crate::state::DaemonState; -use crate::waiters::{WaiterKind, WaiterRegistry}; +use crate::waiters::{WaiterKind, WaiterPredicate, WaiterRegistry}; -// Not yet constructed outside tests: Task 5's connection handling is what -// will call `serde_json::from_str::` on incoming socket lines. #[derive(Debug, Deserialize)] #[serde(tag = "cmd", rename_all = "kebab-case")] pub enum Request { @@ -100,15 +98,24 @@ pub enum Request { Shutdown, } -// The functions below are only reachable via `handle_request`, which is -// itself only wired up to a real socket by Task 5 — allow dead_code until -// then, as Tasks 2/3 did for `state`/`waiters`. fn err(code: &str, message: impl Into) -> serde_json::Value { serde_json::json!({"ok": false, "error": code, "message": message.into()}) } -fn not_connected() -> serde_json::Value { - err("not_connected", "robot controller not connected yet") +/// Every command except `ping`/`status` fails fast with `not_connected` +/// when the daemon has never seen `Event::Connected`, rather than waiting +/// out a timeout or silently dropping a packet that has no peer to go to. +fn ensure_connected(state: &Arc>) -> Result<(), serde_json::Value> { + if state.lock().expect("state mutex poisoned").connected { + Ok(()) + } else { + Err(err("not_connected", "robot controller not connected yet")) + } +} + +/// Predicate for commands whose response carries nothing to match on. +fn any_event() -> WaiterPredicate { + Box::new(|_| true) } fn check_known_opmode(state: &Arc>, name: &str) -> Option { @@ -123,30 +130,27 @@ fn check_known_opmode(state: &Arc>, name: &str) -> Option>, - state: &Arc>, kind: WaiterKind, + pred: WaiterPredicate, timeout: Duration, send: impl FnOnce(), ) -> Result { - if !state.lock().expect("state mutex poisoned").connected { - return Err(not_connected()); - } let rx = waiters .lock() .expect("waiters mutex poisoned") - .register(kind); + .register(kind, pred); send(); rx.recv_timeout(timeout) .map_err(|_| err("timeout", "robot controller did not respond in time")) } -#[allow(dead_code, clippy::too_many_arguments)] +#[allow(clippy::too_many_arguments)] fn build_gamepad( left_stick_x: Option, left_stick_y: Option, @@ -201,17 +205,31 @@ pub fn handle_request( waiters: &Arc>, timeout: Duration, ) -> serde_json::Value { + match &req { + Request::Ping | Request::Status => {} + _ => { + if let Err(e) = ensure_connected(state) { + return e; + } + } + } match req { Request::Ping => serde_json::json!({"ok": true}), Request::Status => state.lock().expect("state mutex poisoned").to_json(), Request::List => { let client = client.clone(); - match wait_for_event(waiters, state, WaiterKind::OpModeList, timeout, move || { - client - .lock() - .expect("client mutex poisoned") - .request_opmode_list(); - }) { + match wait_for_event( + waiters, + WaiterKind::OpModeList, + any_event(), + timeout, + move || { + client + .lock() + .expect("client mutex poisoned") + .request_opmode_list(); + }, + ) { Ok(Event::OpModeList(list)) => serde_json::json!({ "ok": true, "opmodes": list.iter().map(|m| serde_json::json!({ @@ -230,10 +248,11 @@ pub fn handle_request( } let client = client.clone(); let send_name = name.clone(); + let want = name.clone(); match wait_for_event( waiters, - state, WaiterKind::OpModeInited, + Box::new(move |e| matches!(e, Event::OpModeInited(n) if *n == want)), timeout, move || { client @@ -270,10 +289,11 @@ pub fn handle_request( } let client = client.clone(); let send_name = name.clone(); + let want = name.clone(); match wait_for_event( waiters, - state, WaiterKind::OpModeRunning, + Box::new(move |e| matches!(e, Event::OpModeRunning(n) if *n == want)), timeout, move || { client @@ -304,8 +324,8 @@ pub fn handle_request( let client = client.clone(); match wait_for_event( waiters, - state, WaiterKind::ActiveConfiguration, + any_event(), timeout, move || { client @@ -327,8 +347,8 @@ pub fn handle_request( let client = client.clone(); match wait_for_event( waiters, - state, WaiterKind::ConfigurationList, + any_event(), timeout, move || { client @@ -360,8 +380,8 @@ pub fn handle_request( let client = client.clone(); match wait_for_event( waiters, - state, WaiterKind::Configuration, + any_event(), timeout, move || { client @@ -422,8 +442,8 @@ pub fn handle_request( let client = client.clone(); match wait_for_event( waiters, - state, WaiterKind::UserDeviceList, + any_event(), timeout, move || { client @@ -443,9 +463,15 @@ pub fn handle_request( } Request::Scan => { let client = client.clone(); - match wait_for_event(waiters, state, WaiterKind::ScanResult, timeout, move || { - client.lock().expect("client mutex poisoned").scan(); - }) { + match wait_for_event( + waiters, + WaiterKind::ScanResult, + any_event(), + timeout, + move || { + client.lock().expect("client mutex poisoned").scan(); + }, + ) { Ok(Event::ScanResult(extra)) => serde_json::json!({"ok": true, "scan": extra}), Ok(_) => { unreachable!("WaiterKind::ScanResult only resolves with Event::ScanResult") @@ -456,10 +482,11 @@ pub fn handle_request( Request::LynxModules { serial } => { let client = client.clone(); let send_serial = serial.clone(); + // First-match only: see the `WaiterKind::LynxModules` doc. match wait_for_event( waiters, - state, WaiterKind::LynxModules, + any_event(), timeout, move || { client @@ -535,9 +562,9 @@ pub fn handle_request( } serde_json::json!({"ok": true}) } - // Watch and Shutdown are handled by the connection layer (Task 5), - // which needs the raw stream to push a live stream / exit the - // process; they never reach this dispatcher. + // Watch and Shutdown are handled by the connection layer, which + // needs the raw stream to push a live stream / exit the process; + // they never reach this dispatcher. Request::Watch { .. } | Request::Shutdown => { unreachable!("Watch and Shutdown are intercepted before handle_request is called") } @@ -641,8 +668,8 @@ mod tests { let client = unreachable_client(); // Simulate the robot's reply landing on the waiter registry shortly - // after the request goes out, the way the real event-pump thread - // would (Task 6 wires that thread up for real). + // after the request goes out, the way the real event-pump thread in + // `main.rs` does. let waiters_bg = waiters.clone(); let (ready_tx, ready_rx) = mpsc::channel::<()>(); std::thread::spawn(move || { @@ -671,6 +698,92 @@ mod tests { assert_eq!(reply["opmodes"][0]["name"], "Solo (TeleOp)"); } + #[test] + fn gamepad_fails_fast_with_not_connected() { + let (state, waiters) = deps(); + let client = unreachable_client(); + let reply = handle_request( + Request::Gamepad { + left_stick_x: None, + left_stick_y: None, + right_stick_x: None, + right_stick_y: None, + left_trigger: None, + right_trigger: None, + dpad_up: None, + dpad_down: None, + dpad_left: None, + dpad_right: None, + a: None, + b: None, + x: None, + y: None, + start: None, + back: None, + left_bumper: None, + right_bumper: None, + left_stick_button: None, + right_stick_button: None, + duration_ms: None, + }, + &client, + &state, + &waiters, + Duration::from_millis(50), + ); + assert_eq!(reply["ok"], false); + assert_eq!(reply["error"], "not_connected"); + } + + #[test] + fn init_ignores_a_stale_stop_robot_inited_event() { + // Mirrors `ds_agent stop; ds_agent init Foo`: the RC's late + // `NOTIFY_INIT_OP_MODE $Stop$Robot$` lands while the `init Foo` + // waiter is registered and must not resolve it. + let (state, waiters) = deps(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state + .lock() + .unwrap() + .apply_event(&Event::Connected { peer }); + state + .lock() + .unwrap() + .apply_event(&Event::OpModeList(vec![robocol::cmd::OpModeMeta { + name: "Foo".into(), + flavor: "TELEOP".into(), + group: "drive".into(), + }])); + let client = unreachable_client(); + + let waiters_bg = waiters.clone(); + let (ready_tx, ready_rx) = mpsc::channel::<()>(); + std::thread::spawn(move || { + ready_rx.recv().unwrap(); + std::thread::sleep(Duration::from_millis(20)); + waiters_bg + .lock() + .unwrap() + .resolve(&Event::OpModeInited(robocol::cmd::DEFAULT_OP_MODE.into())); + std::thread::sleep(Duration::from_millis(50)); + waiters_bg + .lock() + .unwrap() + .resolve(&Event::OpModeInited("Foo".into())); + }); + + ready_tx.send(()).unwrap(); + let reply = handle_request( + Request::Init { name: "Foo".into() }, + &client, + &state, + &waiters, + Duration::from_secs(1), + ); + assert_eq!(reply["ok"], true); + assert_eq!(reply["name"], "Foo"); + } + #[test] fn list_times_out_when_nothing_resolves_it() { let (state, waiters) = deps(); diff --git a/ds_agentd/src/waiters.rs b/ds_agentd/src/waiters.rs index 525136f..a5c0220 100644 --- a/ds_agentd/src/waiters.rs +++ b/ds_agentd/src/waiters.rs @@ -4,7 +4,7 @@ //! `Subscribers` keeps broadcasting to every `watch` connection until it //! disconnects. -use std::sync::mpsc::{self, Receiver, Sender}; +use std::sync::mpsc::{self, Receiver, Sender, SyncSender, TrySendError}; use robocol::client::Event; @@ -18,6 +18,10 @@ pub enum WaiterKind { Configuration, UserDeviceList, ScanResult, + /// Cannot be payload-matched: `Event::LynxModules` carries only the RC's + /// raw response, not the serial it was requested for, so a `LynxModules` + /// waiter is resolved by the first matching-kind event (register it with + /// a `|_| true` predicate). LynxModules, } @@ -36,8 +40,13 @@ pub fn waiter_kind_of(event: &Event) -> Option { } } +/// Decides whether a same-kind event is the one a particular waiter asked +/// for (e.g. `OpModeInited(n)` with `n == requested_name`), so an unsolicited +/// or stale event of the right kind can't resolve the wrong request. +pub type WaiterPredicate = Box bool + Send>; + pub struct WaiterRegistry { - waiters: Vec<(WaiterKind, Sender)>, + waiters: Vec<(WaiterKind, WaiterPredicate, Sender)>, } impl WaiterRegistry { @@ -47,21 +56,22 @@ impl WaiterRegistry { } } - pub fn register(&mut self, kind: WaiterKind) -> Receiver { + pub fn register(&mut self, kind: WaiterKind, pred: WaiterPredicate) -> Receiver { let (tx, rx) = mpsc::channel(); - self.waiters.push((kind, tx)); + self.waiters.push((kind, pred, tx)); rx } - /// Delivers `event` to every waiter registered for its kind, then drops - /// them — each registration is resolved (or forgotten, on send failure) - /// exactly once. Waiters of other kinds are left untouched. + /// Delivers `event` to every waiter registered for its kind whose + /// predicate accepts it, then drops them — each registration is resolved + /// (or forgotten, on send failure) exactly once. Waiters of other kinds, + /// or whose predicate rejects the payload, are left untouched. pub fn resolve(&mut self, event: &Event) { let Some(kind) = waiter_kind_of(event) else { return; }; - self.waiters.retain(|(k, tx)| { - if *k == kind { + self.waiters.retain(|(k, pred, tx)| { + if *k == kind && pred(event) { let _ = tx.send(event.clone()); false } else { @@ -71,8 +81,13 @@ impl WaiterRegistry { } } +/// Per-subscriber queue depth. A stalled `ds_agent watch` reader must not +/// grow the daemon's memory without bound, so events past this are dropped +/// for that subscriber only. +const SUBSCRIBER_QUEUE: usize = 64; + pub struct Subscribers { - list: Vec>, + list: Vec>, } impl Subscribers { @@ -81,15 +96,19 @@ impl Subscribers { } pub fn subscribe(&mut self) -> Receiver { - let (tx, rx) = mpsc::channel(); + let (tx, rx) = mpsc::sync_channel(SUBSCRIBER_QUEUE); self.list.push(tx); rx } /// Sends `event` to every live subscriber, dropping any whose receiver - /// has gone away (their `watch` connection closed). + /// has gone away (their `watch` connection closed). A subscriber whose + /// queue is full just misses this event; it is not pruned. pub fn broadcast(&mut self, event: &Event) { - self.list.retain(|tx| tx.send(event.clone()).is_ok()); + self.list.retain(|tx| match tx.try_send(event.clone()) { + Ok(()) | Err(TrySendError::Full(_)) => true, + Err(TrySendError::Disconnected(_)) => false, + }); } } @@ -97,11 +116,15 @@ impl Subscribers { mod tests { use super::*; + fn any() -> WaiterPredicate { + Box::new(|_| true) + } + #[test] fn resolve_delivers_to_matching_kind_only() { let mut registry = WaiterRegistry::new(); - let opmode_rx = registry.register(WaiterKind::OpModeList); - let config_rx = registry.register(WaiterKind::ActiveConfiguration); + let opmode_rx = registry.register(WaiterKind::OpModeList, any()); + let config_rx = registry.register(WaiterKind::ActiveConfiguration, any()); registry.resolve(&Event::OpModeList(vec![])); @@ -112,7 +135,7 @@ mod tests { #[test] fn resolve_forgets_waiter_after_one_delivery() { let mut registry = WaiterRegistry::new(); - let rx = registry.register(WaiterKind::ScanResult); + let rx = registry.register(WaiterKind::ScanResult, any()); registry.resolve(&Event::ScanResult("first".into())); registry.resolve(&Event::ScanResult("second".into())); @@ -124,7 +147,7 @@ mod tests { #[test] fn resolve_ignores_events_with_no_waiter_kind() { let mut registry = WaiterRegistry::new(); - let rx = registry.register(WaiterKind::OpModeList); + let rx = registry.register(WaiterKind::OpModeList, any()); registry.resolve(&Event::Disconnected); assert!(rx.try_recv().is_err()); } @@ -132,16 +155,34 @@ mod tests { #[test] fn dropped_receiver_is_pruned_on_next_matching_event() { let mut registry = WaiterRegistry::new(); - drop(registry.register(WaiterKind::ScanResult)); + drop(registry.register(WaiterKind::ScanResult, any())); // Resolving with a dropped receiver must not panic, and the dead // entry must be gone afterward (checked indirectly: a second // resolve with a fresh waiter only sees its own event). registry.resolve(&Event::ScanResult("ignored".into())); - let rx = registry.register(WaiterKind::ScanResult); + let rx = registry.register(WaiterKind::ScanResult, any()); registry.resolve(&Event::ScanResult("seen".into())); assert_eq!(rx.try_recv().unwrap(), Event::ScanResult("seen".into())); } + #[test] + fn predicate_rejects_same_kind_event_with_other_payload() { + let mut registry = WaiterRegistry::new(); + let rx = registry.register( + WaiterKind::OpModeInited, + Box::new(|e| matches!(e, Event::OpModeInited(n) if n == "Foo")), + ); + + // The RC's unsolicited "$Stop$Robot$" must not resolve an `init Foo`. + registry.resolve(&Event::OpModeInited("$Stop$Robot$".into())); + assert!(rx.try_recv().is_err()); + assert_eq!(registry.waiters.len(), 1); + + registry.resolve(&Event::OpModeInited("Foo".into())); + assert_eq!(rx.try_recv().unwrap(), Event::OpModeInited("Foo".into())); + assert!(registry.waiters.is_empty()); + } + #[test] fn broadcast_reaches_every_subscriber() { let mut subs = Subscribers::new(); @@ -163,4 +204,32 @@ mod tests { assert_eq!(live.try_recv().unwrap(), Event::Disconnected); assert_eq!(subs.list.len(), 1); } + + #[test] + fn full_queue_drops_events_but_keeps_subscriber_while_disconnected_is_pruned() { + let mut subs = Subscribers::new(); + let slow = subs.subscribe(); // never drained + { + let _gone = subs.subscribe(); // dropped: must be pruned + } + assert_eq!(subs.list.len(), 2); + + // Overfill the slow subscriber's queue; broadcast must neither block + // nor prune it. + for _ in 0..(SUBSCRIBER_QUEUE + 10) { + subs.broadcast(&Event::Disconnected); + } + assert_eq!(subs.list.len(), 1, "full queue must not prune; dead one must"); + + // Exactly the queue depth got through; the excess was dropped. + let mut received = 0; + while slow.try_recv().is_ok() { + received += 1; + } + assert_eq!(received, SUBSCRIBER_QUEUE); + + // Once drained, the subscriber receives again. + subs.broadcast(&Event::Disconnected); + assert_eq!(slow.try_recv().unwrap(), Event::Disconnected); + } } From 458c822cbbaf3eb15d96bb90b8855970687d00ab Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:16:04 -0400 Subject: [PATCH 12/16] ds_agentd: bind the Unix socket before UDP, exit when robocol dies, keep $Stop$Robot$ out of last_inited bind_socket() now runs first so the 'already running -> exit 0' path is reachable with default ports, and a UDP bind failure (Deck Station or ds_cli holding 20884) prints a diagnosable message, removes the socket file so it isn't left stale, and exits 1 instead of panicking. The event pump exits the process when the robocol worker's event stream ends, so a daemon can't linger answering connected:true from a frozen cache. DaemonState no longer records $Stop$Robot$ as last_inited, so a bare 'run' after 'stop' still targets the user's last real opmode. Also drops connection.rs's stale 'Task 6' comment. Co-Authored-By: Claude Opus 5 --- ds_agentd/src/connection.rs | 3 --- ds_agentd/src/main.rs | 25 +++++++++++++++++++++++-- ds_agentd/src/state.rs | 24 +++++++++++++++++++++++- 3 files changed, 46 insertions(+), 6 deletions(-) diff --git a/ds_agentd/src/connection.rs b/ds_agentd/src/connection.rs index 0f40e71..8a5a3be 100644 --- a/ds_agentd/src/connection.rs +++ b/ds_agentd/src/connection.rs @@ -13,9 +13,6 @@ use crate::request::{Request, handle_request}; use crate::state::DaemonState; use crate::waiters::{Subscribers, WaiterRegistry}; -// Not yet wired up outside tests: Task 6's accept loop is what will spawn a -// thread per connection calling `handle_connection`, as Tasks 2-4 noted for -// `state`/`waiters`/`request`. fn event_type_name(event: &Event) -> Option<&'static str> { match event { Event::Connected { .. } => Some("connected"), diff --git a/ds_agentd/src/main.rs b/ds_agentd/src/main.rs index c7eb904..9181eff 100644 --- a/ds_agentd/src/main.rs +++ b/ds_agentd/src/main.rs @@ -97,7 +97,24 @@ fn main() { config.bind_port = port; } - let (client, events) = RobocolClient::start(config).expect("bind UDP socket"); + // Claim the Unix socket first: if another daemon already owns it this + // exits 0 before touching UDP, so the "already running" path is + // reachable even with the default bind port. + let listener = bind_socket(); + + let bind_port = config.bind_port; + let (client, events) = match RobocolClient::start(config) { + Ok(started) => started, + Err(e) => { + eprintln!( + "ds_agentd: cannot bind UDP port {bind_port}: {e} (is Deck Station or ds_cli already running?)" + ); + // Don't leave a socket file that ds_agent would mistake for a + // live daemon. + let _ = std::fs::remove_file(ds_agent_ipc::socket_path()); + std::process::exit(1); + } + }; let client = Arc::new(Mutex::new(client)); let state = Arc::new(Mutex::new(DaemonState::new())); let waiters = Arc::new(Mutex::new(WaiterRegistry::new())); @@ -122,11 +139,15 @@ fn main() { .expect("subscribers mutex poisoned") .broadcast(&event); } + // The robocol worker only drops its sender when it dies. A + // daemon with no robot connection must not linger answering + // `connected: true` from a frozen cache. + eprintln!("ds_agentd: robocol client stopped; exiting"); + std::process::exit(1); }); } let timeout = resolve_timeout(); - let listener = bind_socket(); for incoming in listener.incoming() { let Ok(stream) = incoming else { continue }; let client = client.clone(); diff --git a/ds_agentd/src/state.rs b/ds_agentd/src/state.rs index 09b4853..e6cb2c5 100644 --- a/ds_agentd/src/state.rs +++ b/ds_agentd/src/state.rs @@ -61,7 +61,12 @@ impl DaemonState { } Event::RobotState(state) => self.robot_state = Some(*state), Event::OpModeList(list) => self.opmodes = list.clone(), - Event::OpModeInited(name) => self.last_inited = Some(name.clone()), + // The RC reports `$Stop$Robot$` inited whenever an opmode ends; + // that is not a user choice, so a bare `run` after `stop` still + // targets the last real opmode. + Event::OpModeInited(name) if name != robocol::cmd::DEFAULT_OP_MODE => { + self.last_inited = Some(name.clone()) + } Event::ActiveConfiguration(extra) => self.active_config = Some(extra.clone()), Event::ConfigurationList(extra) => self.configs = parse_config_list(extra), Event::UserDeviceList(extra) => self.device_types = Some(extra.clone()), @@ -173,6 +178,23 @@ mod tests { assert_eq!(state.last_inited.as_deref(), Some("Duo (TeleOp)")); } + #[test] + fn stop_robot_inited_does_not_overwrite_last_inited() { + let mut state = DaemonState::new(); + state.apply_event(&Event::OpModeInited("Duo (TeleOp)".into())); + state.apply_event(&Event::OpModeInited( + robocol::cmd::DEFAULT_OP_MODE.to_string(), + )); + assert_eq!(state.last_inited.as_deref(), Some("Duo (TeleOp)")); + + // And it never becomes the fallback on its own either. + let mut fresh = DaemonState::new(); + fresh.apply_event(&Event::OpModeInited( + robocol::cmd::DEFAULT_OP_MODE.to_string(), + )); + assert_eq!(fresh.last_inited, None); + } + #[test] fn configuration_list_populates_find_config() { let mut state = DaemonState::new(); From f0524868702bc636e224c88e37e37530f730f13f Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:17:01 -0400 Subject: [PATCH 13/16] ds_agent: detach the spawned daemon, read the stop reply, pass DS_AGENTD_ARGS through The auto-spawned ds_agentd now gets its own process group so it outlives the one-shot CLI and any shell job control around it, and DS_AGENTD_ARGS (split on whitespace) is forwarded to it so a caller can pick the peer or bind port without running the daemon by hand. 'daemon stop' now waits for the daemon's {"ok":true} instead of fire-and-forget; 'daemon status' always carries a 'running' field. All unreachable-daemon exits go through one print_unreachable helper that emits the same JSON shape on stdout and stderr, including the post-connect write failure that used to print bare text. The spawn log lives at socket_path().with_extension("log") so the /tmp fallback is per-uid like the socket, and '-1s'/'nans' durations are rejected instead of wrapping. Co-Authored-By: Claude Opus 5 --- ds_agent/src/main.rs | 107 ++++++++++++++++++++++++++----------------- 1 file changed, 65 insertions(+), 42 deletions(-) diff --git a/ds_agent/src/main.rs b/ds_agent/src/main.rs index 1d140f4..c62134c 100644 --- a/ds_agent/src/main.rs +++ b/ds_agent/src/main.rs @@ -8,10 +8,12 @@ //! ds_agent gamepad --left-stick-y -0.5 --duration 300ms //! ds_agent watch --types telemetry //! ds_agent daemon status +//! DS_AGENTD_ARGS="--peer 192.168.43.1 --bind-port 0" ds_agent status # extra ds_agentd flags on auto-spawn //! ``` use std::io::BufReader; use std::os::unix::net::UnixStream; +use std::os::unix::process::CommandExt; use std::time::{Duration, Instant}; fn build_request(args: &[String]) -> Result { @@ -89,9 +91,10 @@ fn parse_duration_ms(value: &str) -> Result { if let Some(ms) = value.strip_suffix("ms") { ms.parse().map_err(|_| format!("invalid duration: {value}")) } else if let Some(s) = value.strip_suffix('s') { - s.parse::() - .map(|secs| (secs * 1000.0) as u64) - .map_err(|_| format!("invalid duration: {value}")) + match s.parse::() { + Ok(secs) if secs.is_finite() && secs >= 0.0 => Ok((secs * 1000.0) as u64), + _ => Err(format!("invalid duration: {value}")), + } } else { value .parse() @@ -161,9 +164,11 @@ fn build_watch_request(args: &[String]) -> Result { Ok(serde_json::json!({"cmd": "watch", "types": types})) } -/// Connects to the daemon, spawning it if needed. A spawned daemon's stdout -/// and stderr are redirected to `ds_agentd.log` next to the socket, so a -/// crash or panic during startup can be diagnosed after the fact. +/// Connects to the daemon, spawning it detached (its own process group) if +/// needed. A spawned daemon's stdout and stderr are redirected to a `.log` +/// file beside the socket, so a crash or panic during startup can be +/// diagnosed after the fact. If `DS_AGENTD_ARGS` is set, it is split on +/// whitespace and passed through as extra `ds_agentd` arguments. fn connect_or_spawn() -> Result { let path = ds_agent_ipc::socket_path(); if let Ok(stream) = UnixStream::connect(&path) { @@ -171,7 +176,7 @@ fn connect_or_spawn() -> Result { } let exe = std::env::current_exe().map_err(|e| e.to_string())?; let daemon_path = exe.with_file_name("ds_agentd"); - let log_path = path.with_file_name("ds_agentd.log"); + let log_path = path.with_extension("log"); let log = std::fs::OpenOptions::new() .create(true) .append(true) @@ -180,10 +185,15 @@ fn connect_or_spawn() -> Result { let log_err = log .try_clone() .map_err(|e| format!("failed to open {}: {e}", log_path.display()))?; - std::process::Command::new(&daemon_path) + let mut command = std::process::Command::new(&daemon_path); + if let Ok(extra) = std::env::var("DS_AGENTD_ARGS") { + command.args(extra.split_whitespace()); + } + command .stdin(std::process::Stdio::null()) .stdout(std::process::Stdio::from(log)) .stderr(std::process::Stdio::from(log_err)) + .process_group(0) .spawn() .map_err(|e| { format!( @@ -208,42 +218,54 @@ fn connect_or_spawn() -> Result { } } +/// Prints the `daemon_unreachable` error to stdout (for the caller's JSON +/// parser) and stderr (for a human), then exits 1. +fn print_unreachable(message: &str) -> ! { + let line = serde_json::json!({"ok": false, "error": "daemon_unreachable", "message": message}) + .to_string(); + println!("{line}"); + eprintln!("{line}"); + std::process::exit(1); +} + +/// Sends one request over an already-connected stream and reads one reply +/// line; `None` when the daemon closed without answering. +fn round_trip(mut stream: UnixStream, request: &serde_json::Value) -> Option { + ds_agent_ipc::write_message(&mut stream, request).ok()?; + let mut reader = BufReader::new(stream); + let line = ds_agent_ipc::read_message(&mut reader).ok().flatten()?; + serde_json::from_str(&line).ok() +} + fn run_daemon_subcommand(args: &[String]) -> i32 { let path = ds_agent_ipc::socket_path(); match args.first().map(String::as_str) { Some("status") => match UnixStream::connect(&path) { - Ok(mut stream) => { - let _ = - ds_agent_ipc::write_message(&mut stream, &serde_json::json!({"cmd": "status"})); - let mut reader = BufReader::new(stream); - match ds_agent_ipc::read_message(&mut reader) { - Ok(Some(line)) => { - println!("{line}"); - 0 - } - _ => { - let line = serde_json::json!({"ok": false, "error": "daemon_unreachable"}) - .to_string(); - println!("{line}"); - eprintln!("{line}"); - 1 + Ok(stream) => match round_trip(stream, &serde_json::json!({"cmd": "status"})) { + Some(mut status) => { + // `running` is always present: false when there is no + // socket to connect to, true when the daemon answered. + if let Some(obj) = status.as_object_mut() { + obj.insert("running".into(), serde_json::Value::Bool(true)); } + println!("{status}"); + 0 } - } + None => print_unreachable("ds_agentd accepted the connection but did not answer"), + }, Err(_) => { println!(r#"{{"ok":true,"running":false}}"#); 0 } }, Some("stop") => match UnixStream::connect(&path) { - Ok(mut stream) => { - let _ = ds_agent_ipc::write_message( - &mut stream, - &serde_json::json!({"cmd": "shutdown"}), - ); - println!(r#"{{"ok":true}}"#); - 0 - } + Ok(stream) => match round_trip(stream, &serde_json::json!({"cmd": "shutdown"})) { + Some(reply) if reply["ok"] == true => { + println!("{reply}"); + 0 + } + _ => print_unreachable("ds_agentd did not acknowledge shutdown"), + }, Err(_) => { println!(r#"{{"ok":true,"running":false}}"#); 0 @@ -273,18 +295,10 @@ fn main() { let mut stream = match connect_or_spawn() { Ok(s) => s, - Err(msg) => { - let line = - serde_json::json!({"ok": false, "error": "daemon_unreachable", "message": msg}) - .to_string(); - println!("{line}"); - eprintln!("{line}"); - std::process::exit(1); - } + Err(msg) => print_unreachable(&msg), }; if ds_agent_ipc::write_message(&mut stream, &request).is_err() { - eprintln!("failed to write to ds_agentd"); - std::process::exit(1); + print_unreachable("failed to write to ds_agentd"); } let mut reader = BufReader::new(stream.try_clone().expect("clone stream")); @@ -361,6 +375,15 @@ mod tests { assert_eq!(req["duration_ms"], 300); } + #[test] + fn duration_seconds_rejects_negative_and_non_finite() { + assert_eq!(parse_duration_ms("1.5s"), Ok(1500)); + assert_eq!(parse_duration_ms("0s"), Ok(0)); + assert!(parse_duration_ms("-1s").is_err()); + assert!(parse_duration_ms("nans").is_err()); + assert!(parse_duration_ms("infs").is_err()); + } + #[test] fn gamepad_rejects_unknown_flags() { let args: Vec = ["gamepad", "--not-a-flag", "1"] From 99f45fc07fb47f7a6b938eceabe11a4e784ba844 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:18:16 -0400 Subject: [PATCH 14/16] e2e: prove daemon stop exits the daemon, bound the watch read, cover auto-spawn The fake_rc test now polls try_wait for up to 2 s after 'daemon stop' and fails if the daemon is still alive, and reads the first watch line on a helper thread with a 5 s recv_timeout so a silent daemon can't hang the suite (the kill-on-drop EOFs the helper). A second test drives the auto-spawn path end to end: DS_AGENTD_ARGS points the spawned daemon at a dead port, 'status' answers ok/connected:false, the daemon's log file exists beside the socket, 'daemon status' reports running:true, and after 'daemon stop' the socket stops accepting. A Drop guard re-issues 'daemon stop' and removes the runtime dir so a panic mid-test can't orphan the detached daemon. Co-Authored-By: Claude Opus 5 --- ds_agentd/tests/end_to_end.rs | 181 ++++++++++++++++++++++++++-------- 1 file changed, 141 insertions(+), 40 deletions(-) diff --git a/ds_agentd/tests/end_to_end.rs b/ds_agentd/tests/end_to_end.rs index 88c4fa9..3096279 100644 --- a/ds_agentd/tests/end_to_end.rs +++ b/ds_agentd/tests/end_to_end.rs @@ -1,9 +1,14 @@ //! Drives the full `ds_agent` -> `ds_agentd` -> `fake_rc` stack the way a //! human would drive `ds_cli` -> a real Control Hub, but scripted. +//! +//! The two tests here use different runtime dirs and UDP ports, so cargo may +//! run them in parallel. use std::io::{BufRead, BufReader}; -use std::path::PathBuf; +use std::os::unix::net::UnixStream; +use std::path::{Path, PathBuf}; use std::process::{Child, Command, Stdio}; +use std::sync::mpsc; use std::time::{Duration, Instant}; struct ChildGuard(Child); @@ -22,30 +27,39 @@ fn workspace_root() -> PathBuf { } /// `cargo test -p ds_agentd` only builds this crate's own binaries, so build -/// the two siblings the test drives before touching them. -fn build_sibling_bins() { - let status = Command::new(env!("CARGO")) - .args(["build", "-p", "ds_agent", "-p", "fake_rc"]) +/// the sibling crates a test drives before touching them. +fn build_bins(packages: &[&str]) { + let mut cmd = Command::new(env!("CARGO")); + cmd.arg("build"); + for p in packages { + cmd.args(["-p", p]); + } + let status = cmd .current_dir(workspace_root()) .status() - .expect("run cargo build for ds_agent/fake_rc"); - assert!(status.success(), "building ds_agent/fake_rc failed"); + .unwrap_or_else(|e| panic!("run cargo build for {packages:?}: {e}")); + assert!(status.success(), "building {packages:?} failed"); } fn bin_path(name: &str) -> PathBuf { // CARGO_BIN_EXE_ds_agentd is set for this crate's own binary; sibling - // binaries built by build_sibling_bins() share the same target/debug - // dir, one level up from ds_agentd's own exe path. + // binaries built by build_bins() share the same target/debug dir, one + // level up from ds_agentd's own exe path. let mut path = PathBuf::from(env!("CARGO_BIN_EXE_ds_agentd")); path.pop(); path.push(name); path } -fn run_agent(socket_dir: &std::path::Path, args: &[&str]) -> serde_json::Value { +fn run_agent_with_env( + socket_dir: &Path, + args: &[&str], + env: &[(&str, &str)], +) -> serde_json::Value { let output = Command::new(bin_path("ds_agent")) .args(args) .env("XDG_RUNTIME_DIR", socket_dir) + .envs(env.iter().copied()) .output() .expect("run ds_agent"); let stdout = String::from_utf8_lossy(&output.stdout); @@ -57,11 +71,35 @@ fn run_agent(socket_dir: &std::path::Path, args: &[&str]) -> serde_json::Value { }) } +fn run_agent(socket_dir: &Path, args: &[&str]) -> serde_json::Value { + run_agent_with_env(socket_dir, args, &[]) +} + +fn tempdir(tag: &str) -> PathBuf { + let dir = PathBuf::from(format!("/tmp/ds_agent_e2e_{tag}{}", std::process::id())); + std::fs::create_dir_all(&dir).expect("create temp runtime dir"); + dir +} + +/// Polls `cond` every 50 ms until it holds or `limit` elapses. +fn wait_until(limit: Duration, mut cond: impl FnMut() -> bool) -> bool { + let deadline = Instant::now() + limit; + loop { + if cond() { + return true; + } + if Instant::now() >= deadline { + return false; + } + std::thread::sleep(Duration::from_millis(50)); + } +} + #[test] fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { - build_sibling_bins(); + build_bins(&["ds_agent", "fake_rc"]); - let tmp = tempdir(); + let tmp = tempdir(""); let fake_rc_port: u16 = 20950; let mut fake_rc = ChildGuard( @@ -93,26 +131,19 @@ fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { // Wait for the daemon's socket to exist before sending anything. let socket_path = tmp.join("ds_agentd.sock"); - let deadline = Instant::now() + Duration::from_secs(3); - while !socket_path.exists() { - assert!(Instant::now() < deadline, "ds_agentd socket never appeared"); - std::thread::sleep(Duration::from_millis(50)); - } + assert!( + wait_until(Duration::from_secs(3), || socket_path.exists()), + "ds_agentd socket never appeared" + ); // Wait for the daemon to actually connect to fake_rc (discovery + a // couple of heartbeats) before issuing commands that depend on it. - let deadline = Instant::now() + Duration::from_secs(3); - loop { - let status = run_agent(&tmp, &["status"]); - if status["connected"] == true { - break; - } - assert!( - Instant::now() < deadline, - "ds_agentd never connected to fake_rc" - ); - std::thread::sleep(Duration::from_millis(100)); - } + assert!( + wait_until(Duration::from_secs(3), || run_agent(&tmp, &["status"]) + ["connected"] + == true), + "ds_agentd never connected to fake_rc" + ); let list = run_agent(&tmp, &["list"]); assert_eq!(list["ok"], true); @@ -134,7 +165,9 @@ fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { // process should see at least one line before it's killed. Scoped in a // block so the ChildGuard drops (killing the watch) before `daemon // stop` below, and so a panic on any of the asserts still kills it - // rather than leaking an orphan `ds_agent watch` process. + // rather than leaking an orphan `ds_agent watch` process. The blocking + // read lives on a helper thread so a silent daemon fails the test + // within 5 s instead of hanging it; the kill-on-drop EOFs that thread. { let mut watch = ChildGuard( Command::new(bin_path("ds_agent")) @@ -145,27 +178,95 @@ fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { .expect("start ds_agent watch"), ); let stdout = watch.0.stdout.take().expect("watch stdout"); - let mut reader = BufReader::new(stdout); - let mut line = String::new(); - reader.read_line(&mut line).expect("read a telemetry line"); + let (tx, rx) = mpsc::channel::>(); + std::thread::spawn(move || { + let mut reader = BufReader::new(stdout); + let mut line = String::new(); + let _ = tx.send(reader.read_line(&mut line).map(|_| line)); + }); + let line = rx + .recv_timeout(Duration::from_secs(5)) + .expect("no telemetry line within 5s") + .expect("read a telemetry line"); let event: serde_json::Value = serde_json::from_str(line.trim_end()).expect("watch line is JSON"); assert_eq!(event["event"], "telemetry"); } // Explicit teardown via `daemon stop`, ahead of ChildGuard's Drop kill, - // proves the shutdown command itself works. + // proves the shutdown command itself works — and that the daemon + // actually exits on it, not just acknowledges. let stopped = run_agent(&tmp, &["daemon", "stop"]); assert_eq!(stopped["ok"], true); - std::thread::sleep(Duration::from_millis(200)); - let _ = daemon.0.try_wait(); + assert!( + wait_until(Duration::from_secs(2), || matches!( + daemon.0.try_wait(), + Ok(Some(_)) + )), + "daemon did not exit after `daemon stop`" + ); let _ = fake_rc.0.kill(); std::fs::remove_dir_all(&tmp).expect("clean up runtime dir"); } -fn tempdir() -> PathBuf { - let dir = PathBuf::from(format!("/tmp/ds_agent_e2e_{}", std::process::id())); - std::fs::create_dir_all(&dir).expect("create temp runtime dir"); - dir +/// Tears down whatever `ds_agent` auto-spawned even if the test panics +/// midway: the daemon is detached in its own process group, so nothing else +/// would reap it. +struct SpawnedDaemonGuard { + dir: PathBuf, + env: Vec<(&'static str, &'static str)>, +} + +impl Drop for SpawnedDaemonGuard { + fn drop(&mut self) { + let _ = Command::new(bin_path("ds_agent")) + .args(["daemon", "stop"]) + .env("XDG_RUNTIME_DIR", &self.dir) + .envs(self.env.iter().copied()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status(); + let _ = std::fs::remove_dir_all(&self.dir); + } +} + +#[test] +fn auto_spawns_daemon_logs_it_and_stops_it() { + build_bins(&["ds_agent"]); + + let tmp = tempdir("spawn_"); + // Port 1: nothing answers, so no fake_rc is needed and nothing clashes + // with the other test (or a running Deck Station on 20884). + let env: Vec<(&str, &str)> = vec![( + "DS_AGENTD_ARGS", + "--peer 127.0.0.1 --peer-port 1 --bind-port 0", + )]; + let guard = SpawnedDaemonGuard { + dir: tmp.clone(), + env: env.clone(), + }; + + // No daemon is running for this runtime dir, so this spawns one. + let status = run_agent_with_env(&tmp, &["status"], &env); + assert_eq!(status["ok"], true, "status reply: {status}"); + assert_eq!(status["connected"], false); + + let log = tmp.join("ds_agentd.log"); + assert!(log.exists(), "spawned daemon's log {} missing", log.display()); + + let daemon_status = run_agent_with_env(&tmp, &["daemon", "status"], &env); + assert_eq!(daemon_status["running"], true); + + let stopped = run_agent_with_env(&tmp, &["daemon", "stop"], &env); + assert_eq!(stopped["ok"], true, "daemon stop reply: {stopped}"); + + let socket = tmp.join("ds_agentd.sock"); + assert!( + wait_until(Duration::from_secs(2), || UnixStream::connect(&socket).is_err()), + "ds_agentd still accepting connections after `daemon stop`" + ); + + drop(guard); // best-effort second stop is a no-op; removes the dir + assert!(!tmp.exists(), "runtime dir not removed"); } From f7e006a50fb5c68a79a7c344b48fd2c6bd15a4f9 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:18:33 -0400 Subject: [PATCH 15/16] ds_agentd: rustfmt Co-Authored-By: Claude Opus 5 --- ds_agentd/src/waiters.rs | 6 +++++- ds_agentd/tests/end_to_end.rs | 22 ++++++++++++---------- 2 files changed, 17 insertions(+), 11 deletions(-) diff --git a/ds_agentd/src/waiters.rs b/ds_agentd/src/waiters.rs index a5c0220..052e181 100644 --- a/ds_agentd/src/waiters.rs +++ b/ds_agentd/src/waiters.rs @@ -219,7 +219,11 @@ mod tests { for _ in 0..(SUBSCRIBER_QUEUE + 10) { subs.broadcast(&Event::Disconnected); } - assert_eq!(subs.list.len(), 1, "full queue must not prune; dead one must"); + assert_eq!( + subs.list.len(), + 1, + "full queue must not prune; dead one must" + ); // Exactly the queue depth got through; the excess was dropped. let mut received = 0; diff --git a/ds_agentd/tests/end_to_end.rs b/ds_agentd/tests/end_to_end.rs index 3096279..0dfcce7 100644 --- a/ds_agentd/tests/end_to_end.rs +++ b/ds_agentd/tests/end_to_end.rs @@ -51,11 +51,7 @@ fn bin_path(name: &str) -> PathBuf { path } -fn run_agent_with_env( - socket_dir: &Path, - args: &[&str], - env: &[(&str, &str)], -) -> serde_json::Value { +fn run_agent_with_env(socket_dir: &Path, args: &[&str], env: &[(&str, &str)]) -> serde_json::Value { let output = Command::new(bin_path("ds_agent")) .args(args) .env("XDG_RUNTIME_DIR", socket_dir) @@ -139,9 +135,10 @@ fn drives_an_opmode_and_watches_telemetry_against_fake_rc() { // Wait for the daemon to actually connect to fake_rc (discovery + a // couple of heartbeats) before issuing commands that depend on it. assert!( - wait_until(Duration::from_secs(3), || run_agent(&tmp, &["status"]) - ["connected"] - == true), + wait_until( + Duration::from_secs(3), + || run_agent(&tmp, &["status"])["connected"] == true + ), "ds_agentd never connected to fake_rc" ); @@ -253,7 +250,11 @@ fn auto_spawns_daemon_logs_it_and_stops_it() { assert_eq!(status["connected"], false); let log = tmp.join("ds_agentd.log"); - assert!(log.exists(), "spawned daemon's log {} missing", log.display()); + assert!( + log.exists(), + "spawned daemon's log {} missing", + log.display() + ); let daemon_status = run_agent_with_env(&tmp, &["daemon", "status"], &env); assert_eq!(daemon_status["running"], true); @@ -263,7 +264,8 @@ fn auto_spawns_daemon_logs_it_and_stops_it() { let socket = tmp.join("ds_agentd.sock"); assert!( - wait_until(Duration::from_secs(2), || UnixStream::connect(&socket).is_err()), + wait_until(Duration::from_secs(2), || UnixStream::connect(&socket) + .is_err()), "ds_agentd still accepting connections after `daemon stop`" ); From 74f26fcf45e1e489d46d1083d0bf3012e5e49230 Mon Sep 17 00:00:00 2001 From: Vlad Date: Thu, 10 Sep 2026 23:27:50 -0400 Subject: [PATCH 16/16] ds_agentd: unregister waiters that time out A payload-matched waiter (init Foo) is only removed by a later OpModeInited("Foo"); if that never comes, each timed-out attempt left a permanent entry in WaiterRegistry. register now returns a WaiterId and wait_for_event unregisters it on timeout. Co-Authored-By: Claude Opus 5 --- ds_agentd/src/request.rs | 37 +++++++++++++++-- ds_agentd/src/waiters.rs | 86 ++++++++++++++++++++++++++++++++-------- 2 files changed, 104 insertions(+), 19 deletions(-) diff --git a/ds_agentd/src/request.rs b/ds_agentd/src/request.rs index 01e803a..dacd384 100644 --- a/ds_agentd/src/request.rs +++ b/ds_agentd/src/request.rs @@ -141,13 +141,20 @@ fn wait_for_event( timeout: Duration, send: impl FnOnce(), ) -> Result { - let rx = waiters + let (id, rx) = waiters .lock() .expect("waiters mutex poisoned") .register(kind, pred); send(); - rx.recv_timeout(timeout) - .map_err(|_| err("timeout", "robot controller did not respond in time")) + rx.recv_timeout(timeout).map_err(|_| { + // A payload-matched waiter (`init Foo`) would otherwise linger until + // a matching event that may never arrive. + waiters + .lock() + .expect("waiters mutex poisoned") + .unregister(id); + err("timeout", "robot controller did not respond in time") + }) } #[allow(clippy::too_many_arguments)] @@ -803,4 +810,28 @@ mod tests { assert_eq!(reply["ok"], false); assert_eq!(reply["error"], "timeout"); } + + #[test] + fn timed_out_init_waiter_is_unregistered() { + let (state, waiters) = deps(); + let peer: SocketAddr = "127.0.0.1:20884".parse().unwrap(); + state + .lock() + .unwrap() + .apply_event(&Event::Connected { peer }); + let client = unreachable_client(); + let reply = handle_request( + Request::Init { + name: "Never".into(), + }, + &client, + &state, + &waiters, + Duration::from_millis(50), + ); + assert_eq!(reply["error"], "timeout"); + // A name-matched waiter would otherwise linger until an + // `OpModeInited("Never")` that may never arrive. + assert_eq!(waiters.lock().unwrap().len(), 0); + } } diff --git a/ds_agentd/src/waiters.rs b/ds_agentd/src/waiters.rs index 052e181..ef79938 100644 --- a/ds_agentd/src/waiters.rs +++ b/ds_agentd/src/waiters.rs @@ -45,21 +45,54 @@ pub fn waiter_kind_of(event: &Event) -> Option { /// or stale event of the right kind can't resolve the wrong request. pub type WaiterPredicate = Box bool + Send>; +/// Handle returned by `register`, used to `unregister` a waiter that gave +/// up (timed out) before its event arrived. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct WaiterId(u64); + +struct Waiter { + id: WaiterId, + kind: WaiterKind, + pred: WaiterPredicate, + tx: Sender, +} + pub struct WaiterRegistry { - waiters: Vec<(WaiterKind, WaiterPredicate, Sender)>, + waiters: Vec, + next_id: u64, } impl WaiterRegistry { pub fn new() -> Self { WaiterRegistry { waiters: Vec::new(), + next_id: 0, } } - pub fn register(&mut self, kind: WaiterKind, pred: WaiterPredicate) -> Receiver { + pub fn register( + &mut self, + kind: WaiterKind, + pred: WaiterPredicate, + ) -> (WaiterId, Receiver) { let (tx, rx) = mpsc::channel(); - self.waiters.push((kind, pred, tx)); - rx + let id = WaiterId(self.next_id); + self.next_id += 1; + self.waiters.push(Waiter { id, kind, pred, tx }); + (id, rx) + } + + /// Forgets a waiter whose caller stopped waiting. A payload-matched + /// waiter (`init Foo`) is otherwise only removed by a later + /// `OpModeInited("Foo")`, which may never come, so timed-out callers + /// must unregister or the registry grows without bound. + pub fn unregister(&mut self, id: WaiterId) { + self.waiters.retain(|w| w.id != id); + } + + #[cfg(test)] + pub(crate) fn len(&self) -> usize { + self.waiters.len() } /// Delivers `event` to every waiter registered for its kind whose @@ -70,9 +103,9 @@ impl WaiterRegistry { let Some(kind) = waiter_kind_of(event) else { return; }; - self.waiters.retain(|(k, pred, tx)| { - if *k == kind && pred(event) { - let _ = tx.send(event.clone()); + self.waiters.retain(|w| { + if w.kind == kind && (w.pred)(event) { + let _ = w.tx.send(event.clone()); false } else { true @@ -123,8 +156,8 @@ mod tests { #[test] fn resolve_delivers_to_matching_kind_only() { let mut registry = WaiterRegistry::new(); - let opmode_rx = registry.register(WaiterKind::OpModeList, any()); - let config_rx = registry.register(WaiterKind::ActiveConfiguration, any()); + let (_, opmode_rx) = registry.register(WaiterKind::OpModeList, any()); + let (_, config_rx) = registry.register(WaiterKind::ActiveConfiguration, any()); registry.resolve(&Event::OpModeList(vec![])); @@ -135,7 +168,7 @@ mod tests { #[test] fn resolve_forgets_waiter_after_one_delivery() { let mut registry = WaiterRegistry::new(); - let rx = registry.register(WaiterKind::ScanResult, any()); + let (_, rx) = registry.register(WaiterKind::ScanResult, any()); registry.resolve(&Event::ScanResult("first".into())); registry.resolve(&Event::ScanResult("second".into())); @@ -147,7 +180,7 @@ mod tests { #[test] fn resolve_ignores_events_with_no_waiter_kind() { let mut registry = WaiterRegistry::new(); - let rx = registry.register(WaiterKind::OpModeList, any()); + let (_, rx) = registry.register(WaiterKind::OpModeList, any()); registry.resolve(&Event::Disconnected); assert!(rx.try_recv().is_err()); } @@ -155,12 +188,12 @@ mod tests { #[test] fn dropped_receiver_is_pruned_on_next_matching_event() { let mut registry = WaiterRegistry::new(); - drop(registry.register(WaiterKind::ScanResult, any())); + drop(registry.register(WaiterKind::ScanResult, any()).1); // Resolving with a dropped receiver must not panic, and the dead // entry must be gone afterward (checked indirectly: a second // resolve with a fresh waiter only sees its own event). registry.resolve(&Event::ScanResult("ignored".into())); - let rx = registry.register(WaiterKind::ScanResult, any()); + let (_, rx) = registry.register(WaiterKind::ScanResult, any()); registry.resolve(&Event::ScanResult("seen".into())); assert_eq!(rx.try_recv().unwrap(), Event::ScanResult("seen".into())); } @@ -168,7 +201,7 @@ mod tests { #[test] fn predicate_rejects_same_kind_event_with_other_payload() { let mut registry = WaiterRegistry::new(); - let rx = registry.register( + let (_, rx) = registry.register( WaiterKind::OpModeInited, Box::new(|e| matches!(e, Event::OpModeInited(n) if n == "Foo")), ); @@ -176,11 +209,32 @@ mod tests { // The RC's unsolicited "$Stop$Robot$" must not resolve an `init Foo`. registry.resolve(&Event::OpModeInited("$Stop$Robot$".into())); assert!(rx.try_recv().is_err()); - assert_eq!(registry.waiters.len(), 1); + assert_eq!(registry.len(), 1); registry.resolve(&Event::OpModeInited("Foo".into())); assert_eq!(rx.try_recv().unwrap(), Event::OpModeInited("Foo".into())); - assert!(registry.waiters.is_empty()); + assert_eq!(registry.len(), 0); + } + + #[test] + fn unregister_forgets_only_the_given_waiter() { + let mut registry = WaiterRegistry::new(); + let (timed_out, _rx1) = registry.register( + WaiterKind::OpModeInited, + Box::new(|e| matches!(e, Event::OpModeInited(n) if n == "Never")), + ); + let (_, rx2) = registry.register(WaiterKind::OpModeInited, any()); + assert_eq!(registry.len(), 2); + + registry.unregister(timed_out); + assert_eq!(registry.len(), 1); + + // Unregistering twice (or an unknown id) is a no-op; the survivor + // still resolves. + registry.unregister(timed_out); + registry.resolve(&Event::OpModeInited("Other".into())); + assert_eq!(rx2.try_recv().unwrap(), Event::OpModeInited("Other".into())); + assert_eq!(registry.len(), 0); } #[test]