From a2370730ef2afc5b0cc54294f8aa999eefdf1da0 Mon Sep 17 00:00:00 2001 From: Stephen Belanger Date: Tue, 25 Aug 2026 23:01:45 +0800 Subject: [PATCH] Add automatic distributed tracing across agents --- bt-daemon/Cargo.lock | 106 +- bt-daemon/Cargo.toml | 1 + bt-daemon/docs/protocol.md | 24 +- bt-daemon/src/correlation.rs | 510 ++++++ bt-daemon/src/dispatch.rs | 173 +- bt-daemon/src/journal.rs | 1 + bt-daemon/src/lib.rs | 3 + bt-daemon/src/process.rs | 206 +++ bt-daemon/src/server.rs | 525 +++++- bt-daemon/src/transcript_import/mod.rs | 1 + bt-daemon/src/translate/opencode.rs | 40 +- bt-daemon/src/translate/pi.rs | 64 +- bt-daemon/src/wire/envelope.rs | 54 + bt-daemon/src/wire/mod.rs | 4 +- bt-daemon/tests/claude_translator.rs | 6 + bt-daemon/tests/codex_translator.rs | 1 + bt-daemon/tests/distributed_tracing.rs | 2220 ++++++++++++++++++++++++ bt-daemon/tests/opencode_translator.rs | 1 + bt-daemon/tests/pi_translator.rs | 7 + bt-daemon/tests/pipeline.rs | 1 + bt-daemon/tests/support/distributed.rs | 489 ++++++ bt-daemon/tests/support/mod.rs | 1 + 22 files changed, 4353 insertions(+), 85 deletions(-) create mode 100644 bt-daemon/src/correlation.rs create mode 100644 bt-daemon/src/process.rs create mode 100644 bt-daemon/tests/distributed_tracing.rs create mode 100644 bt-daemon/tests/support/distributed.rs diff --git a/bt-daemon/Cargo.lock b/bt-daemon/Cargo.lock index dfc514e..225091e 100644 --- a/bt-daemon/Cargo.lock +++ b/bt-daemon/Cargo.lock @@ -274,6 +274,7 @@ dependencies = [ "serde", "serde_json", "sha2", + "sysinfo", "tempfile", "thiserror 2.0.19", "tokio", @@ -887,7 +888,7 @@ dependencies = [ "js-sys", "log", "wasm-bindgen", - "windows-core", + "windows-core 0.62.2", ] [[package]] @@ -1150,6 +1151,15 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "ntapi" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c3b335231dfd352ffb0f8017f3b6027a4917f7df785ea2143d8af2adc66980ae" +dependencies = [ + "winapi", +] + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -1774,6 +1784,19 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "sysinfo" +version = "0.33.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4fc858248ea01b66f19d8e8a6d55f41deaf91e9d495246fd01368d99935c6c01" +dependencies = [ + "core-foundation-sys", + "libc", + "memchr", + "ntapi", + "windows", +] + [[package]] name = "tempfile" version = "3.27.0" @@ -2202,19 +2225,74 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "winapi" +version = "0.3.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c839a674fcd7a98952e593242ea400abe93992746761e38641405d28b00f419" +dependencies = [ + "winapi-i686-pc-windows-gnu", + "winapi-x86_64-pc-windows-gnu", +] + +[[package]] +name = "winapi-i686-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac3b87c63620426dd9b991e5ce0329eff545bccbbb34f3be09ff6fb6ab51b7b6" + +[[package]] +name = "winapi-x86_64-pc-windows-gnu" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" + +[[package]] +name = "windows" +version = "0.57.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "12342cb4d8e3b046f3d80effd474a7a02447231330ef77d71daa6fbc40681143" +dependencies = [ + "windows-core 0.57.0", + "windows-targets", +] + +[[package]] +name = "windows-core" +version = "0.57.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d2ed2439a290666cd67ecce2b0ffaad89c2a56b976b736e6ece670297897832d" +dependencies = [ + "windows-implement 0.57.0", + "windows-interface 0.57.0", + "windows-result 0.1.2", + "windows-targets", +] + [[package]] name = "windows-core" version = "0.62.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b8e83a14d34d0623b51dce9581199302a221863196a1dde71a7663a4c2be9deb" dependencies = [ - "windows-implement", - "windows-interface", + "windows-implement 0.60.2", + "windows-interface 0.59.3", "windows-link", - "windows-result", + "windows-result 0.4.1", "windows-strings", ] +[[package]] +name = "windows-implement" +version = "0.57.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9107ddc059d5b6fbfbffdfa7a7fe3e22a226def0b2608f72e9d552763d3e1ad7" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "windows-implement" version = "0.60.2" @@ -2226,6 +2304,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "windows-interface" +version = "0.57.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29bee4b38ea3cde66011baa44dba677c432a78593e202392d1e9070cf2a7fca7" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "windows-interface" version = "0.59.3" @@ -2243,6 +2332,15 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f0805222e57f7521d6a62e36fa9163bc891acd422f971defe97d64e70d0a4fe5" +[[package]] +name = "windows-result" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5e383302e8ec8515204254685643de10811af0ed97ea37210dc26fb0032647f8" +dependencies = [ + "windows-targets", +] + [[package]] name = "windows-result" version = "0.4.1" diff --git a/bt-daemon/Cargo.toml b/bt-daemon/Cargo.toml index d404440..65784c7 100644 --- a/bt-daemon/Cargo.toml +++ b/bt-daemon/Cargo.toml @@ -27,6 +27,7 @@ regex = "1" serde = { version = "1", features = ["derive"] } serde_json = "1" sha2 = "0.10" +sysinfo = { version = "0.33.1", default-features = false, features = ["system"] } thiserror = "2" tempfile = "3" tokio = { version = "1", features = ["rt-multi-thread", "macros", "net", "io-util", "sync", "time", "process", "signal", "fs"] } diff --git a/bt-daemon/docs/protocol.md b/bt-daemon/docs/protocol.md index 6bc832f..57acb25 100644 --- a/bt-daemon/docs/protocol.md +++ b/bt-daemon/docs/protocol.md @@ -105,7 +105,13 @@ The hot path. Params are the **Envelope** (see below). Request result: ```json { "accepted": true } ``` -`accepted: true` means enqueued to the session's ordered queue and journaled. +`accepted: true` means durably recorded: normally journaled and enqueued to the +session's ordered queue, or held in the daemon's private correlation journal +while multiple parent calls remain indistinguishable. +For an event that opens a tool call, it also means the daemon has made that +active-tool marker visible to local child-session correlation. A child hook +that runs immediately after its parent's blocking pre-tool hook can therefore +attach without an intervening flush. The daemon never fails the caller's turn for a downstream (Braintrust) error; those are handled asynchronously and surfaced via `status.get`. The queue is bounded, so a session whose sink has stalled applies backpressure here instead @@ -223,6 +229,22 @@ Field notes: - **`managed_run_id`** is present only for events inherited from a `bt trace run` process tree. It groups native sessions for the final invocation flush and is not trace metadata. +- **`capture`** is optional daemon-owned process evidence. After + `initialize`, the daemon snapshots the connecting client's PID and bounded + ancestry as PID/start-time pairs and adds it before journaling. Hooks do not + construct this field and it contains no command line, environment, or + working-directory data. The daemon uses it as a local side channel to attach + an instrumented child session to an active parent tool span; agent-native + input and output payloads are reduced to hashes only when more than one + active call is a candidate. For Codex, whose hook payload references its + rollout rather than embedding the prompt, matching also considers at most + the final 256 KiB of native JSONL records. Active parent-tool state is + atomically snapshotted under `/correlation/parents` before a + blocking tool lifecycle hook is acknowledged. The compact snapshot contains + only process identities, span attachment components, non-secret routing, and + hashed matching fingerprints. It contains no raw prompt, tool input/output, + command line, environment, or resolved credential, and lets a child attach + even if the daemon restarts between the parent spawn and child session start. - **`route`** carries non-secret auth selection and trace settings. `profile` is optional and resolves through `bt`'s default profile when absent; `org_name` optionally constrains organization selection. The daemon resolves diff --git a/bt-daemon/src/correlation.rs b/bt-daemon/src/correlation.rs new file mode 100644 index 0000000..8e4bfe7 --- /dev/null +++ b/bt-daemon/src/correlation.rs @@ -0,0 +1,510 @@ +//! Daemon-wide local correlation between active tool calls and child sessions. +//! +//! Agent translators remain source-specific, but their emitted tool rows are a +//! common contract. This registry observes those rows, indexes the active tools +//! by the process ancestry that produced their session, and resolves a child's +//! requested route to the exact spawning tool span. + +use crate::translate::{SpanOp, SpanRow, SpanType}; +use crate::wire::{CaptureContext, ProcessIdentity, SessionConfig, SessionRoute, TraceDestination}; +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use serde::{Deserialize, Serialize}; +use serde_json::{Map, Value}; +use sha2::{Digest, Sha256}; +use std::collections::{HashMap, HashSet}; +use std::sync::Mutex; + +#[derive(Debug, Clone)] +pub(crate) struct ParentLink { + pub route: SessionRoute, +} + +#[derive(Debug, Clone)] +pub(crate) enum Resolution { + Standalone, + Parent(Box), + Ambiguous(Vec), +} + +#[derive(Default)] +pub(crate) struct CorrelationRegistry { + state: Mutex, +} + +#[derive(Default)] +struct State { + process_sessions: HashMap>, + session_processes: HashMap>, + active_tools: HashMap>, + live_sessions: HashSet, +} + +#[derive(Clone)] +struct ActiveTool { + components: SpanComponents, + route: SessionRoute, + fingerprints: HashSet<[u8; 32]>, + active: bool, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub(crate) struct ActiveParentSnapshot { + version: u32, + #[serde(default)] + dirty: bool, + correlation_key: String, + processes: Vec, + tools: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct ActiveToolSnapshot { + components: SpanComponents, + route: SessionRoute, + fingerprints: Vec<[u8; 32]>, +} + +impl CorrelationRegistry { + pub(crate) fn observe_session(&self, key: &str, capture: Option<&CaptureContext>) { + let mut state = self.state.lock().unwrap(); + state.live_sessions.insert(key.to_string()); + let Some(capture) = capture else { return }; + for process in capture + .process_chain + .iter() + .filter(|process| process.start_time_secs != 0) + { + state + .session_processes + .entry(key.to_string()) + .or_default() + .insert(process.clone()); + state + .process_sessions + .entry(process.clone()) + .or_default() + .insert(key.to_string()); + } + } + + pub(crate) fn active_parent_snapshot(&self, key: &str) -> Option { + let state = self.state.lock().unwrap(); + let tools = state.active_tools.get(key)?; + let tools: Vec<_> = tools + .values() + .filter(|tool| tool.active) + .map(|tool| ActiveToolSnapshot { + components: tool.components.clone(), + route: tool.route.clone(), + fingerprints: tool.fingerprints.iter().copied().collect(), + }) + .collect(); + if tools.is_empty() { + return None; + } + Some(ActiveParentSnapshot { + version: 1, + dirty: false, + correlation_key: key.to_string(), + processes: state + .session_processes + .get(key) + .map(|processes| processes.iter().cloned().collect()) + .unwrap_or_default(), + tools, + }) + } + + pub(crate) fn dirty_active_parent_snapshot(&self, key: &str) -> Option { + self.active_parent_snapshot(key).map(|mut snapshot| { + snapshot.dirty = true; + snapshot + }) + } + + pub(crate) fn restore_active_parent(&self, snapshot: ActiveParentSnapshot) -> bool { + let restored_tools: Vec<_> = snapshot + .tools + .into_iter() + .filter(|tool| tool.components.span_id.is_some()) + .collect(); + if snapshot.version != 1 + || snapshot.dirty + || snapshot.processes.is_empty() + || restored_tools.is_empty() + { + return false; + } + let mut state = self.state.lock().unwrap(); + let key = snapshot.correlation_key; + for process in snapshot + .processes + .into_iter() + .filter(|process| process.start_time_secs != 0) + { + state + .session_processes + .entry(key.clone()) + .or_default() + .insert(process.clone()); + state + .process_sessions + .entry(process) + .or_default() + .insert(key.clone()); + } + if !state.session_processes.contains_key(&key) { + return false; + } + let tools = state.active_tools.entry(key).or_default(); + for snapshot in restored_tools { + let span_id = snapshot.components.span_id.clone().unwrap(); + tools.insert( + span_id, + ActiveTool { + components: snapshot.components, + route: snapshot.route, + fingerprints: snapshot.fingerprints.into_iter().collect(), + active: true, + }, + ); + } + !tools.is_empty() + } + + pub(crate) fn observe_ops( + &self, + key: &str, + route: &SessionRoute, + config: &SessionConfig, + ops: &[SpanOp], + ) -> bool { + let mut state = self.state.lock().unwrap(); + let mut changed = false; + for op in ops { + let row = match op { + SpanOp::Insert(row) | SpanOp::Merge(row) => row, + }; + if row.span_id.is_empty() { + continue; + } + if row.end_ms.is_some() { + if let Some(tools) = state.active_tools.get_mut(key) { + if let Some(tool) = tools.get_mut(&row.span_id) { + if let Some(output) = &row.output { + tool.fingerprints.extend(fingerprints(output)); + } + tool.active = false; + changed = true; + } + } + continue; + } + if row.span_type != SpanType::Tool { + continue; + } + if !matches!(op, SpanOp::Insert(_)) { + continue; + } + let components = span_components(config, row); + let fingerprints = row.input.as_ref().map(fingerprints).unwrap_or_default(); + state + .active_tools + .entry(key.to_string()) + .or_default() + .insert( + row.span_id.clone(), + ActiveTool { + components, + route: route.clone(), + fingerprints, + active: true, + }, + ); + changed = true; + } + changed + } + + pub(crate) fn resolve( + &self, + child_source: &str, + child_session_key: Option<&str>, + capture: Option<&CaptureContext>, + evidence: &Value, + ) -> Resolution { + let Some(capture) = capture else { + return Resolution::Standalone; + }; + let state = self.state.lock().unwrap(); + for process in capture + .process_chain + .iter() + .skip(minimum_ancestor_depth(child_source)) + .filter(|process| process.start_time_secs != 0) + { + let Some(sessions) = state.process_sessions.get(process) else { + continue; + }; + let mut candidates = Vec::new(); + for session in sessions { + if child_session_key.is_some_and(|child| child == session) { + continue; + } + if let Some(tools) = state.active_tools.get(session) { + candidates.extend(tools.values().filter(|tool| tool.active).cloned()); + } + } + if candidates.is_empty() { + continue; + } + if candidates.len() == 1 { + return Resolution::Parent(Box::new(to_link(candidates.pop().unwrap()))); + } + + let candidate_span_ids = candidates + .iter() + .filter_map(|candidate| candidate.components.span_id.clone()) + .collect(); + let evidence = fingerprints(evidence); + let mut scored: Vec<(usize, ActiveTool)> = candidates + .into_iter() + .map(|candidate| { + let score = candidate.fingerprints.intersection(&evidence).count(); + (score, candidate) + }) + .filter(|(score, _)| *score > 0) + .collect(); + scored.sort_by_key(|(score, _)| std::cmp::Reverse(*score)); + if let Some((best_score, best)) = scored.first().cloned() { + let tied = scored + .get(1) + .is_some_and(|(second, _)| *second == best_score); + if !tied { + return Resolution::Parent(Box::new(to_link(best))); + } + } + return Resolution::Ambiguous(candidate_span_ids); + } + Resolution::Standalone + } + + pub(crate) fn resolve_pending( + &self, + child_source: &str, + capture: Option<&CaptureContext>, + evidence: &Value, + candidate_span_ids: &[String], + ) -> Resolution { + let Some(capture) = capture else { + return Resolution::Standalone; + }; + let wanted: HashSet<&str> = candidate_span_ids.iter().map(String::as_str).collect(); + let evidence = fingerprints(evidence); + let state = self.state.lock().unwrap(); + for process in capture + .process_chain + .iter() + .skip(minimum_ancestor_depth(child_source)) + .filter(|process| process.start_time_secs != 0) + { + let Some(sessions) = state.process_sessions.get(process) else { + continue; + }; + let mut scored = Vec::new(); + let mut found = 0usize; + let mut any_active = false; + for session in sessions { + let Some(tools) = state.active_tools.get(session) else { + continue; + }; + for (span_id, tool) in tools { + if wanted.contains(span_id.as_str()) { + found += 1; + any_active |= tool.active; + let score = tool.fingerprints.intersection(&evidence).count(); + if score > 0 { + scored.push((score, tool.clone())); + } + } + } + } + scored.sort_by_key(|(score, _)| std::cmp::Reverse(*score)); + if let Some((best_score, best)) = scored.first().cloned() { + let tied = scored + .get(1) + .is_some_and(|(second, _)| *second == best_score); + if !tied { + return Resolution::Parent(Box::new(to_link(best))); + } + } + if found == wanted.len() && !any_active { + return Resolution::Standalone; + } + return Resolution::Ambiguous(candidate_span_ids.to_vec()); + } + Resolution::Standalone + } + + pub(crate) fn remove_session(&self, key: &str) { + let mut state = self.state.lock().unwrap(); + state.live_sessions.remove(key); + state.active_tools.remove(key); + if let Some(processes) = state.session_processes.remove(key) { + for process in processes { + if let Some(sessions) = state.process_sessions.get_mut(&process) { + sessions.remove(key); + if sessions.is_empty() { + state.process_sessions.remove(&process); + } + } + } + } + } + + pub(crate) fn has_active_tools(&self, key: &str) -> bool { + self.state + .lock() + .unwrap() + .active_tools + .get(key) + .is_some_and(|tools| tools.values().any(|tool| tool.active)) + } + + pub(crate) fn has_any_active_tools(&self) -> bool { + let state = self.state.lock().unwrap(); + state.live_sessions.iter().any(|key| { + state + .active_tools + .get(key) + .is_some_and(|tools| tools.values().any(|tool| tool.active)) + }) + } +} + +fn minimum_ancestor_depth(source: &str) -> usize { + // Claude and Codex connect from a short-lived `bt trace hook` process, so + // index 1 is still the current agent process. Pi and OpenCode connect + // in-process, making index 0 current. A real child must be beyond those + // current-session processes in either capture shape. + match source { + "claude-code" | "codex" => 2, + _ => 1, + } +} + +fn to_link(tool: ActiveTool) -> ParentLink { + let mut route = tool.route; + route.destination = Some(TraceDestination::ParentSpan { + components: tool.components.clone(), + }); + ParentLink { route } +} + +fn span_components(config: &SessionConfig, row: &SpanRow) -> SpanComponents { + let mut object_type = SpanObjectType::ProjectLogs; + let mut object_id = None; + let mut compute_object_metadata_args = None; + let mut propagated_event = None; + let mut effective_root = row.root_span_id.clone(); + match config.destination.as_ref() { + Some(TraceDestination::ProjectLogs { + project_id, + project_name, + }) => { + object_id = project_id.clone(); + let mut args = Map::new(); + if let Some(project_id) = project_id { + args.insert("project_id".into(), Value::String(project_id.clone())); + } + if let Some(project_name) = project_name { + args.insert("project_name".into(), Value::String(project_name.clone())); + } + compute_object_metadata_args = (!args.is_empty()).then_some(args); + } + Some(TraceDestination::Experiment { experiment_id }) => { + object_type = SpanObjectType::Experiment; + object_id = Some(experiment_id.clone()); + } + Some(TraceDestination::ParentSpan { components }) => { + object_type = components.object_type; + object_id = components.object_id.clone(); + compute_object_metadata_args = components.compute_object_metadata_args.clone(); + propagated_event = components.propagated_event.clone(); + if let Some(root) = &components.root_span_id { + effective_root = root.clone(); + } + } + None => {} + } + SpanComponents { + object_type, + object_id, + compute_object_metadata_args, + row_id: Some(row.span_id.clone()), + span_id: Some(row.span_id.clone()), + root_span_id: Some(effective_root), + span_parents: (!row.parent_span_ids.is_empty()).then(|| row.parent_span_ids.clone()), + propagated_event, + } +} + +fn fingerprints(value: &Value) -> HashSet<[u8; 32]> { + let mut strings = Vec::new(); + collect_strings(value, &mut strings); + let mut result = HashSet::new(); + for value in strings { + let normalized = normalize(&value); + if normalized.len() >= 8 { + result.insert(hash(normalized.as_bytes())); + } + let tokens: Vec<&str> = normalized.split_whitespace().collect(); + for width in 3..=tokens.len().min(8) { + for window in tokens.windows(width) { + result.insert(hash(window.join(" ").as_bytes())); + } + } + } + result +} + +fn collect_strings(value: &Value, strings: &mut Vec) { + match value { + Value::String(value) => strings.push(value.clone()), + Value::Array(values) => values + .iter() + .for_each(|value| collect_strings(value, strings)), + Value::Object(values) => values + .values() + .for_each(|value| collect_strings(value, strings)), + _ => {} + } +} + +fn normalize(value: &str) -> String { + value + .split_whitespace() + .collect::>() + .join(" ") + .to_lowercase() +} + +fn hash(value: &[u8]) -> [u8; 32] { + Sha256::digest(value).into() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn fingerprints_match_prompt_subsets_without_parsing_shell() { + let parent = fingerprints(&serde_json::json!({ + "command": "agent --prompt 'inspect the distributed tracing linkage carefully please'" + })); + let child = fingerprints(&serde_json::json!({ + "prompt": "inspect the distributed tracing linkage carefully please" + })); + assert!(parent.intersection(&child).count() > 0); + } +} diff --git a/bt-daemon/src/dispatch.rs b/bt-daemon/src/dispatch.rs index d065898..a0cb73a 100644 --- a/bt-daemon/src/dispatch.rs +++ b/bt-daemon/src/dispatch.rs @@ -3,10 +3,11 @@ //! are processed strictly in arrival order. Different sessions run //! concurrently. //! -//! Ack semantics: `event.log` is acked once the event is journaled and handed -//! to the session's queue (see [`Session::enqueue`]). Delivery to -//! Braintrust happens later in the actor; a downstream error never fails the -//! caller's turn. +//! Ack semantics: `event.log` is acked once the event is journaled and its +//! first translation batch has updated local correlation state. Tool-start +//! events drain bounded translator continuations before ack so the spawning +//! marker is guaranteed visible. A downstream error never fails the caller's +//! turn. use crate::sink::SinkFactory; use crate::translate::{Registry, SessionCtx}; @@ -30,7 +31,7 @@ pub struct Counters { } enum SessionMsg { - Event(Box), + Event(Box, oneshot::Sender<()>), Configure(Box, oneshot::Sender<()>), Flush(oneshot::Sender), Shutdown(oneshot::Sender<()>), @@ -46,6 +47,18 @@ pub struct ReplayPlan { pub through: u64, } +pub(crate) struct SessionOptions { + pub session_id: String, + pub source: String, + pub plugin_version: Option, + pub replay: Option, + pub config: crate::wire::SessionConfig, + pub correlation_key: String, + pub route: SessionRoute, + pub correlation: Arc, + pub data_dir: PathBuf, +} + /// Handle to one live session: its queue plus observable counters/state. pub struct Session { pub source: String, @@ -59,14 +72,21 @@ pub struct Session { impl Session { /// Spawn a session's actor task and return its handle. pub fn spawn( - session_id: String, - source: String, - plugin_version: Option, - replay: Option, - config: crate::wire::SessionConfig, + options: SessionOptions, translators: Arc, sink_factory: Arc, ) -> Arc { + let SessionOptions { + session_id, + source, + plugin_version, + replay, + config, + correlation_key, + route, + correlation, + data_dir, + } = options; let (tx, rx) = mpsc::channel(QUEUE_CAPACITY); let counters = Arc::new(Counters::default()); let last_error = Arc::new(Mutex::new(None)); @@ -83,6 +103,10 @@ impl Session { permalink: permalink.clone(), replay, config, + correlation_key, + route, + correlation, + data_dir, }; tokio::spawn(actor.run(rx)); @@ -100,11 +124,14 @@ impl Session { pub async fn enqueue(&self, env: Envelope) -> anyhow::Result<()> { self.touch(); self.counters.queued.fetch_add(1, Ordering::Relaxed); + let (reply_tx, reply_rx) = oneshot::channel(); self.tx - .send(SessionMsg::Event(Box::new(env))) + .send(SessionMsg::Event(Box::new(env), reply_tx)) .await .map_err(|_| anyhow::anyhow!("session actor is gone"))?; - Ok(()) + reply_rx + .await + .map_err(|_| anyhow::anyhow!("session actor dropped event acknowledgement")) } fn touch(&self) { @@ -212,6 +239,31 @@ struct SessionActor { permalink: Arc>>, replay: Option, config: crate::wire::SessionConfig, + correlation_key: String, + route: SessionRoute, + correlation: Arc, + data_dir: PathBuf, +} + +#[derive(Clone, Copy)] +enum BatchMode { + Live, + Replay, + Flush, +} + +impl BatchMode { + fn errors(self) -> (&'static str, &'static str) { + match self { + Self::Live => ("translate failed", "sink emit failed"), + Self::Replay => ("journal replay failed", "sink replay emit failed"), + Self::Flush => ("translate flush failed", "sink emit (flush) failed"), + } + } + + fn observes_correlation(self) -> bool { + !matches!(self, Self::Flush) + } } impl SessionActor { @@ -229,8 +281,9 @@ impl SessionActor { // callers waiting on flush don't hang. while let Some(msg) = rx.recv().await { match msg { - SessionMsg::Event(_) => { + SessionMsg::Event(_, reply) => { self.counters.queued.fetch_sub(1, Ordering::Relaxed); + let _ = reply.send(()); } SessionMsg::Configure(_, r) => { let _ = r.send(()); @@ -260,22 +313,41 @@ impl SessionActor { while let Some(msg) = rx.recv().await { match msg { - SessionMsg::Event(env) => { + SessionMsg::Event(env, reply) => { + let correlation_barrier = is_tool_lifecycle_event(&env.event); + let mut reply = Some(reply); if let Some(cfg) = &env.config { sink.configure(cfg); ctx.config = Some(cfg.clone()); self.refresh_permalink(sink.as_ref()); } let translated = translator.handle(&env, &ctx); - self.emit_translator_batches( - &mut translator, - &mut sink, - &ctx, - translated, - "translate failed", - "sink emit failed", - ) - .await; + if !correlation_barrier { + let _ = reply.take().expect("event reply").send(()); + } + let correlation_changed = self + .emit_translator_batches( + &mut translator, + &mut sink, + &ctx, + translated, + BatchMode::Live, + ) + .await; + if correlation_changed || correlation_barrier { + if let Err(error) = crate::server::persist_active_parent_snapshot( + &self.data_dir, + &self.correlation_key, + &self.correlation, + ) + .await + { + self.set_error(error); + } + } + if let Some(reply) = reply { + let _ = reply.send(()); + } self.counters.queued.fetch_sub(1, Ordering::Relaxed); } SessionMsg::Configure(config, reply) => { @@ -306,18 +378,27 @@ impl SessionActor { sink: &mut Box, ctx: &SessionCtx, first: anyhow::Result>, - translate_error: &str, - emit_error: &str, - ) { + mode: BatchMode, + ) -> bool { + let (translate_error, emit_error) = mode.errors(); let mut next = match first { Ok(ops) => Some(ops), Err(e) => { self.set_error(format!("{translate_error}: {e}")); - return; + return false; } }; + let mut correlation_changed = false; while let Some(ops) = next { if !ops.is_empty() { + if mode.observes_correlation() { + correlation_changed |= self.correlation.observe_ops( + &self.correlation_key, + &self.route, + ctx.config.as_ref().expect("session config"), + &ops, + ); + } match sink.emit(&ops).await { Ok(n) => { self.counters.spans_emitted.fetch_add(n, Ordering::Relaxed); @@ -333,6 +414,7 @@ impl SessionActor { } }; } + correlation_changed } /// Stream the journal through the translator, emitting each entry's spans @@ -376,15 +458,9 @@ impl SessionActor { } let env = crate::journal::envelope_from_redacted(entry); let translated = translator.handle(&env, ctx); - self.emit_translator_batches( - translator, - sink, - ctx, - translated, - "journal replay failed", - "sink replay emit failed", - ) - .await; + let _ = self + .emit_translator_batches(translator, sink, ctx, translated, BatchMode::Replay) + .await; } } @@ -395,15 +471,9 @@ impl SessionActor { ctx: &SessionCtx, ) { let translated = translator.flush(ctx); - self.emit_translator_batches( - translator, - sink, - ctx, - translated, - "translate flush failed", - "sink emit (flush) failed", - ) - .await; + let _ = self + .emit_translator_batches(translator, sink, ctx, translated, BatchMode::Flush) + .await; if let Err(e) = sink.flush().await { self.set_error(format!("sink flush failed: {e}")); } @@ -421,3 +491,16 @@ impl SessionActor { *self.last_error.lock().unwrap() = Some(msg); } } + +pub(crate) fn is_tool_lifecycle_event(event: &str) -> bool { + matches!( + event, + "PreToolUse" + | "PostToolUse" + | "PostToolUseFailure" + | "tool_execution_start" + | "tool_execution_end" + | "tool.execute.before" + | "tool.execute.after" + ) +} diff --git a/bt-daemon/src/journal.rs b/bt-daemon/src/journal.rs index 17c71e7..91e3089 100644 --- a/bt-daemon/src/journal.rs +++ b/bt-daemon/src/journal.rs @@ -262,6 +262,7 @@ pub fn envelope_from_redacted(r: RedactedEnvelope) -> Envelope { event: r.event, ts_ms: r.ts_ms, managed_run_id: r.managed_run_id, + capture: r.capture, payload: r.payload, route, config, diff --git a/bt-daemon/src/lib.rs b/bt-daemon/src/lib.rs index 0ef36ae..9c0605b 100644 --- a/bt-daemon/src/lib.rs +++ b/bt-daemon/src/lib.rs @@ -14,9 +14,11 @@ pub mod paths; mod client; mod command_output; +mod correlation; mod dispatch; mod ids; mod journal; +pub(crate) mod process; mod server; mod settings; mod setup; @@ -262,6 +264,7 @@ pub async fn run_hook( managed_run_id: std::env::var(MANAGED_RUN_ID_ENV) .ok() .filter(|value| !value.is_empty()), + capture: None, payload, route: Some(route), config: None, diff --git a/bt-daemon/src/process.rs b/bt-daemon/src/process.rs new file mode 100644 index 0000000..29e7aab --- /dev/null +++ b/bt-daemon/src/process.rs @@ -0,0 +1,206 @@ +//! Best-effort local process identity capture for cross-agent correlation. +//! +//! The daemon snapshots a connecting client's ancestry while that process is +//! still alive. Collection is deliberately metadata-only: no command lines, +//! environment variables, or working directories are read. + +use crate::wire::{CaptureContext, ProcessIdentity}; +use std::collections::HashSet; +use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, System}; + +const MAX_PROCESS_CHAIN_DEPTH: usize = 64; + +/// Capture `pid` and its ancestors, nearest process first. +/// +/// Process inspection can race with process exit and can be restricted by the +/// host. In either case this returns the useful prefix collected so far rather +/// than failing event capture. +pub(crate) fn capture_process_context(pid: u32) -> CaptureContext { + let mut system = System::new(); + let process_chain = build_process_chain(pid, |pid| { + let sysinfo_pid = Pid::from_u32(pid); + system.refresh_processes_specifics( + ProcessesToUpdate::Some(&[sysinfo_pid]), + true, + ProcessRefreshKind::nothing(), + ); + let process = system.process(sysinfo_pid)?; + Some(ProcessSnapshot { + identity: ProcessIdentity { + pid, + start_time_secs: process.start_time(), + }, + parent_pid: process.parent().map(Pid::as_u32), + }) + }); + CaptureContext { process_chain } +} + +#[derive(Debug, Clone)] +struct ProcessSnapshot { + identity: ProcessIdentity, + parent_pid: Option, +} + +fn build_process_chain( + start_pid: u32, + mut inspect: impl FnMut(u32) -> Option, +) -> Vec { + if start_pid == 0 { + return Vec::new(); + } + + let mut chain = Vec::new(); + let mut seen = HashSet::new(); + let mut current = start_pid; + + while chain.len() < MAX_PROCESS_CHAIN_DEPTH && seen.insert(current) { + let Some(snapshot) = inspect(current) else { + break; + }; + chain.push(snapshot.identity); + let Some(parent) = snapshot.parent_pid.filter(|parent| *parent != 0) else { + break; + }; + current = parent; + } + + chain +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::HashMap; + use std::process::Command; + use std::time::Duration; + + fn snapshot(pid: u32, parent_pid: Option) -> ProcessSnapshot { + ProcessSnapshot { + identity: ProcessIdentity { + pid, + start_time_secs: u64::from(pid) * 10, + }, + parent_pid, + } + } + + #[test] + fn builds_nearest_first_chain() { + let processes = HashMap::from([ + (30, snapshot(30, Some(20))), + (20, snapshot(20, Some(10))), + (10, snapshot(10, None)), + ]); + + let chain = build_process_chain(30, |pid| processes.get(&pid).cloned()); + + assert_eq!( + chain + .iter() + .map(|identity| identity.pid) + .collect::>(), + vec![30, 20, 10] + ); + } + + #[test] + fn returns_available_prefix_when_an_ancestor_disappears() { + let processes = HashMap::from([(30, snapshot(30, Some(20)))]); + + let chain = build_process_chain(30, |pid| processes.get(&pid).cloned()); + + assert_eq!( + chain + .iter() + .map(|identity| identity.pid) + .collect::>(), + vec![30] + ); + } + + #[test] + fn stops_at_cycles_and_depth_limit() { + let cycle = HashMap::from([(30, snapshot(30, Some(20))), (20, snapshot(20, Some(30)))]); + let chain = build_process_chain(30, |pid| cycle.get(&pid).cloned()); + assert_eq!( + chain + .iter() + .map(|identity| identity.pid) + .collect::>(), + vec![30, 20] + ); + + let chain = build_process_chain(1, |pid| snapshot(pid, Some(pid + 1)).into()); + assert_eq!(chain.len(), MAX_PROCESS_CHAIN_DEPTH); + } + + #[test] + fn captures_current_process_with_stable_identity() { + let first = capture_process_context(std::process::id()); + let second = capture_process_context(std::process::id()); + + assert_eq!( + first.process_chain.first().map(|process| process.pid), + Some(std::process::id()) + ); + assert_eq!( + first + .process_chain + .first() + .map(|process| process.start_time_secs), + second + .process_chain + .first() + .map(|process| process.start_time_secs) + ); + } + + #[test] + fn captures_spawned_child_ancestry() { + const CHILD_ENV: &str = "_BT_PROCESS_CAPTURE_TEST_CHILD"; + if std::env::var_os(CHILD_ENV).is_some() { + std::thread::sleep(Duration::from_secs(5)); + return; + } + + let mut child = Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "process::tests::captures_spawned_child_ancestry", + "--nocapture", + ]) + .env(CHILD_ENV, "1") + .spawn() + .unwrap(); + let child_pid = child.id(); + let mut captured = CaptureContext::default(); + for _ in 0..50 { + captured = capture_process_context(child_pid); + if captured + .process_chain + .iter() + .any(|process| process.pid == std::process::id()) + { + break; + } + std::thread::sleep(Duration::from_millis(10)); + } + let _ = child.kill(); + let _ = child.wait(); + + assert_eq!( + captured.process_chain.first().map(|process| process.pid), + Some(child_pid) + ); + assert!(captured + .process_chain + .iter() + .any(|process| process.pid == std::process::id())); + } + + #[test] + fn zero_pid_has_no_process_context() { + assert!(capture_process_context(0).process_chain.is_empty()); + } +} diff --git a/bt-daemon/src/server.rs b/bt-daemon/src/server.rs index 7192ae7..0359648 100644 --- a/bt-daemon/src/server.rs +++ b/bt-daemon/src/server.rs @@ -2,7 +2,7 @@ //! serves JSON-RPC connections, and shuts down gracefully (idle timeout, //! `daemon.shutdown`, or SIGINT/SIGTERM). -use crate::dispatch::{hydrate_transcript_reference, ReplayPlan, Session}; +use crate::dispatch::{hydrate_transcript_reference, ReplayPlan, Session, SessionOptions}; use crate::journal::{self, JournalWriter}; use crate::sink::SinkFactory; use crate::translate::Registry; @@ -15,12 +15,15 @@ use crate::wire::{ use crate::wire::{AuthSelection, BackendAuth, SessionRoute}; use crate::{paths, ServeArgs}; use async_trait::async_trait; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use sha2::{Digest, Sha256}; use std::collections::{HashMap, HashSet}; use std::path::PathBuf; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; -use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncSeekExt, AsyncWriteExt, BufReader}; use tokio::sync::Notify; /// Injected dependencies for `serve`, so `bt` / tests can supply a sink @@ -71,6 +74,18 @@ struct SessionAuthState { lease: AuthLease, } +#[derive(Default, Serialize, Deserialize)] +struct PendingSession { + #[serde(default, skip_serializing_if = "Option::is_none")] + linked_route: Option, + #[serde(default)] + events: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + candidate_span_ids: Vec, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + evidence: Vec, +} + /// One independent delivery pipeline for a source session and the exact route /// carried by its hook or import envelope. #[derive(Clone, Debug, PartialEq, Eq, Hash)] @@ -86,6 +101,10 @@ impl DeliveryKey { route: serde_json::to_string(route)?, }) } + + fn correlation_key(&self) -> String { + format!("{}\u{1f}{}", self.session_id, self.route) + } } pub struct Daemon { @@ -100,6 +119,10 @@ pub struct Daemon { managed_run_sessions: Mutex>>, auth_errors: Mutex>, sessions: Mutex>>, + correlation: Arc, + automatic_links: Mutex>, + pending_sessions: Mutex>, + correlation_locks: Mutex>>>, started: Instant, last_activity: Mutex, shutting_down: AtomicBool, @@ -120,6 +143,10 @@ impl Daemon { managed_run_sessions: Mutex::new(HashMap::new()), auth_errors: Mutex::new(HashMap::new()), sessions: Mutex::new(HashMap::new()), + correlation: Arc::new(crate::correlation::CorrelationRegistry::default()), + automatic_links: Mutex::new(HashMap::new()), + pending_sessions: Mutex::new(HashMap::new()), + correlation_locks: Mutex::new(HashMap::new()), started: Instant::now(), last_activity: Mutex::new(Instant::now()), shutting_down: AtomicBool::new(false), @@ -234,6 +261,14 @@ impl Daemon { *self.last_activity.lock().unwrap() = Instant::now(); } + fn correlation_lock(&self, key: &str) -> Arc> { + let mut locks = self.correlation_locks.lock().unwrap(); + locks + .entry(key.to_string()) + .or_insert_with(|| Arc::new(tokio::sync::Mutex::new(()))) + .clone() + } + async fn session_for(&self, env: &Envelope, key: &DeliveryKey) -> anyhow::Result> { { let map = self.sessions.lock().unwrap(); @@ -263,11 +298,17 @@ impl Daemon { return Ok(s.clone()); } let session = Session::spawn( - env.session_id.clone(), - env.source.clone(), - env.plugin_version.clone(), - Some(replay), - config, + SessionOptions { + session_id: env.session_id.clone(), + source: env.source.clone(), + plugin_version: env.plugin_version.clone(), + replay: Some(replay), + config, + correlation_key: key.correlation_key(), + route: route.clone(), + correlation: self.correlation.clone(), + data_dir: self.data_dir.clone(), + }, self.translators.clone(), self.sink_factory.clone(), ); @@ -288,6 +329,7 @@ impl Daemon { return; }; session.shutdown().await; + self.correlation.remove_session(&key.correlation_key()); self.session_auth.lock().await.remove(key); self.auth_errors.lock().unwrap().remove(key); @@ -323,6 +365,7 @@ impl Daemon { session.idle_for() >= idle_timeout && session.counters.queued.load(Ordering::Relaxed) == 0 }) + .filter(|(key, _)| !self.correlation.has_active_tools(&key.correlation_key())) .map(|(key, _)| key.clone()) .collect() } @@ -490,6 +533,7 @@ pub async fn run(args: ServeArgs, opts: ServeOptions) -> anyhow::Result<()> { let daemon = Daemon::new(opts, data_dir); collect_garbage(&daemon.data_dir).await; + restore_active_parent_snapshots(&daemon.data_dir, &daemon.correlation).await; let idle_timeout = Duration::from_secs(args.idle_timeout_secs); spawn_idle_watchdog(daemon.clone(), idle_timeout); spawn_session_reaper( @@ -547,6 +591,7 @@ async fn accept_loop(daemon: Arc, mut listener: Listener) -> anyhow::Res async fn serve_connection(daemon: Arc, stream: ServerStream) -> anyhow::Result<()> { let (read_half, mut write_half) = tokio::io::split(stream); let mut lines = BufReader::new(read_half).lines(); + let mut client = None; while let Some(line) = lines.next_line().await? { if line.trim().is_empty() { @@ -561,7 +606,7 @@ async fn serve_connection(daemon: Arc, stream: ServerStream) -> anyhow:: method, "request received" ); - let response = handle_request(&daemon, req).await; + let response = handle_request(&daemon, req, &mut client).await; if let Some(error) = &response.error { tracing::warn!( request_id = ?request_id, @@ -585,7 +630,8 @@ async fn serve_connection(daemon: Arc, stream: ServerStream) -> anyhow:: if note.method == method::EVENT_LOG { if let Some(params) = note.params { match serde_json::from_value::(params) { - Ok(env) => { + Ok(mut env) => { + attach_process_capture(&mut env, client.as_ref()); let _ = accept_event(&daemon, env).await; } Err(error) => tracing::warn!( @@ -618,14 +664,406 @@ async fn serve_connection(daemon: Arc, stream: ServerStream) -> anyhow:: Ok(()) } +fn attach_process_capture(env: &mut Envelope, client: Option<&crate::wire::ClientInfo>) { + if env.capture.is_some() { + return; + } + let Some(pid) = client.and_then(|client| client.pid) else { + return; + }; + let capture = crate::process::capture_process_context(pid); + if !capture.process_chain.is_empty() { + env.capture = Some(capture); + } +} + async fn accept_event(daemon: &Arc, mut env: Envelope) -> Result<(), String> { + daemon.touch(); + let requested_link_key = automatic_link_key(&env); + let correlation_lock = daemon.correlation_lock(&requested_link_key); + let _correlation_guard = correlation_lock.lock().await; + let mut state = daemon + .pending_sessions + .lock() + .unwrap() + .remove(&requested_link_key) + .or_else(|| { + daemon + .automatic_links + .lock() + .unwrap() + .get(&requested_link_key) + .cloned() + .map(|route| PendingSession { + linked_route: Some(route), + events: Vec::new(), + candidate_span_ids: Vec::new(), + evidence: Vec::new(), + }) + }); + if state.is_none() { + state = read_correlation_state(&daemon.data_dir, &requested_link_key).await; + } + + if let Some(mut state) = state { + if let Some(route) = state.linked_route.clone() { + for mut pending in std::mem::take(&mut state.events) { + pending.route = Some(route.clone()); + accept_resolved_event(daemon, pending).await?; + } + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + daemon + .automatic_links + .lock() + .unwrap() + .insert(requested_link_key, route.clone()); + env.route = Some(route); + drop(_correlation_guard); + return accept_resolved_and_retry_pending(daemon, env).await; + } + + let capture = env.capture.as_ref().or_else(|| { + state + .events + .first() + .and_then(|event| event.capture.as_ref()) + }); + state.evidence.push(correlation_evidence(&env).await); + let evidence = Value::Array(state.evidence.clone()); + match daemon.correlation.resolve_pending( + &env.source, + capture, + &evidence, + &state.candidate_span_ids, + ) { + crate::correlation::Resolution::Parent(parent) => { + state.linked_route = Some(parent.route.clone()); + state.events.push(env); + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + for mut event in std::mem::take(&mut state.events) { + event.route = Some(parent.route.clone()); + accept_resolved_event(daemon, event).await?; + } + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + daemon + .automatic_links + .lock() + .unwrap() + .insert(requested_link_key, parent.route); + return Ok(()); + } + crate::correlation::Resolution::Ambiguous(_) + | crate::correlation::Resolution::Standalone => { + state.events.push(env); + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + daemon + .pending_sessions + .lock() + .unwrap() + .insert(requested_link_key, state); + return Ok(()); + } + } + } + + if is_session_start(&env.event) { + let evidence = correlation_evidence(&env).await; + match daemon + .correlation + .resolve(&env.source, None, env.capture.as_ref(), &evidence) + { + crate::correlation::Resolution::Parent(parent) => { + let state = PendingSession { + linked_route: Some(parent.route.clone()), + events: Vec::new(), + candidate_span_ids: Vec::new(), + evidence: Vec::new(), + }; + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + daemon + .automatic_links + .lock() + .unwrap() + .insert(requested_link_key, parent.route.clone()); + env.route = Some(parent.route); + } + crate::correlation::Resolution::Ambiguous(candidate_span_ids) => { + let evidence = correlation_evidence(&env).await; + let state = PendingSession { + linked_route: None, + events: vec![env], + candidate_span_ids, + evidence: vec![evidence], + }; + write_correlation_state(&daemon.data_dir, &requested_link_key, &state).await?; + daemon + .pending_sessions + .lock() + .unwrap() + .insert(requested_link_key, state); + return Ok(()); + } + crate::correlation::Resolution::Standalone => {} + } + } + + drop(_correlation_guard); + accept_resolved_and_retry_pending(daemon, env).await +} + +async fn accept_resolved_and_retry_pending( + daemon: &Arc, + env: Envelope, +) -> Result<(), String> { + accept_resolved_event(daemon, env).await?; + retry_pending_sessions(daemon).await +} + +async fn retry_pending_sessions(daemon: &Arc) -> Result<(), String> { + let keys: Vec = daemon + .pending_sessions + .lock() + .unwrap() + .keys() + .cloned() + .collect(); + for key in keys { + let lock = daemon.correlation_lock(&key); + let _guard = lock.lock().await; + let Some(mut state) = daemon.pending_sessions.lock().unwrap().remove(&key) else { + continue; + }; + if state.linked_route.is_some() || state.events.is_empty() { + daemon.pending_sessions.lock().unwrap().insert(key, state); + continue; + } + let capture = state.events.iter().find_map(|event| event.capture.as_ref()); + let evidence = Value::Array(state.evidence.clone()); + match daemon.correlation.resolve_pending( + state + .events + .first() + .map(|event| event.source.as_str()) + .unwrap_or(""), + capture, + &evidence, + &state.candidate_span_ids, + ) { + crate::correlation::Resolution::Parent(parent) => { + state.linked_route = Some(parent.route.clone()); + write_correlation_state(&daemon.data_dir, &key, &state).await?; + for mut event in std::mem::take(&mut state.events) { + event.route = Some(parent.route.clone()); + accept_resolved_event(daemon, event).await?; + } + write_correlation_state(&daemon.data_dir, &key, &state).await?; + daemon + .automatic_links + .lock() + .unwrap() + .insert(key, parent.route); + } + crate::correlation::Resolution::Ambiguous(_) => { + daemon.pending_sessions.lock().unwrap().insert(key, state); + } + crate::correlation::Resolution::Standalone => { + for event in std::mem::take(&mut state.events) { + accept_resolved_event(daemon, event).await?; + } + remove_correlation_state(&daemon.data_dir, &key).await; + } + } + } + Ok(()) +} + +/// Build opaque matching evidence without teaching the daemon agent or shell +/// syntax. Codex hook payloads reference a JSONL rollout rather than carrying +/// the prompt/output directly, so include a bounded tail of those native JSON +/// values when available. Other agents are fully represented by their payload. +async fn correlation_evidence(env: &Envelope) -> Value { + const MAX_TRANSCRIPT_EVIDENCE_BYTES: u64 = 256 * 1024; + + let mut evidence = vec![env.payload.clone()]; + if env.source != "codex" { + return Value::Array(evidence); + } + let mirror = env.payload.get("_bt_transcript_mirror"); + let path = mirror + .and_then(|value| value.get("mirror")) + .and_then(Value::as_str) + .or_else(|| env.payload.get("transcript_path").and_then(Value::as_str)); + let Some(path) = path else { + return Value::Array(evidence); + }; + let Ok(mut file) = tokio::fs::File::open(path).await else { + return Value::Array(evidence); + }; + let Ok(metadata) = file.metadata().await else { + return Value::Array(evidence); + }; + let through = mirror + .and_then(|value| value.get("through")) + .and_then(Value::as_u64) + .unwrap_or(metadata.len()) + .min(metadata.len()); + let start = through.saturating_sub(MAX_TRANSCRIPT_EVIDENCE_BYTES); + if file.seek(std::io::SeekFrom::Start(start)).await.is_err() { + return Value::Array(evidence); + } + let mut bytes = Vec::with_capacity((through - start) as usize); + if file + .take(through - start) + .read_to_end(&mut bytes) + .await + .is_err() + { + return Value::Array(evidence); + } + let text = String::from_utf8_lossy(&bytes); + let text = if start == 0 { + text.as_ref() + } else { + text.split_once('\n').map(|(_, tail)| tail).unwrap_or("") + }; + evidence.extend( + text.lines() + .filter_map(|line| serde_json::from_str::(line).ok()), + ); + Value::Array(evidence) +} + +fn correlation_state_path(data_dir: &std::path::Path, key: &str) -> PathBuf { + let digest = Sha256::digest(key.as_bytes()); + data_dir + .join("correlation") + .join(format!("{digest:x}.json")) +} + +async fn read_correlation_state(data_dir: &std::path::Path, key: &str) -> Option { + let bytes = tokio::fs::read(correlation_state_path(data_dir, key)) + .await + .ok()?; + serde_json::from_slice(&bytes).ok() +} + +async fn write_correlation_state( + data_dir: &std::path::Path, + key: &str, + state: &PendingSession, +) -> Result<(), String> { + let dir = data_dir.join("correlation"); + tokio::fs::create_dir_all(&dir) + .await + .map_err(|error| format!("correlation journal directory failed: {error}"))?; + let path = correlation_state_path(data_dir, key); + let temp = path.with_extension(format!("{}.tmp", uuid::Uuid::new_v4())); + let bytes = serde_json::to_vec(state) + .map_err(|error| format!("correlation journal encoding failed: {error}"))?; + tokio::fs::write(&temp, bytes) + .await + .map_err(|error| format!("correlation journal write failed: {error}"))?; + if let Err(first_error) = tokio::fs::rename(&temp, &path).await { + // Windows cannot replace an existing destination with rename. + let _ = tokio::fs::remove_file(&path).await; + tokio::fs::rename(&temp, &path).await.map_err(|error| { + format!("correlation journal replace failed: {first_error}; retry failed: {error}") + })?; + } + Ok(()) +} + +async fn remove_correlation_state(data_dir: &std::path::Path, key: &str) { + let _ = tokio::fs::remove_file(correlation_state_path(data_dir, key)).await; +} + +fn active_parent_snapshot_path(data_dir: &std::path::Path, key: &str) -> PathBuf { + let digest = Sha256::digest(key.as_bytes()); + data_dir + .join("correlation") + .join("parents") + .join(format!("{digest:x}.json")) +} + +pub(crate) async fn persist_active_parent_snapshot( + data_dir: &std::path::Path, + key: &str, + correlation: &crate::correlation::CorrelationRegistry, +) -> Result<(), String> { + let path = active_parent_snapshot_path(data_dir, key); + let Some(snapshot) = correlation.active_parent_snapshot(key) else { + let _ = tokio::fs::remove_file(path).await; + return Ok(()); + }; + let bytes = serde_json::to_vec(&snapshot) + .map_err(|error| format!("active parent encoding failed: {error}"))?; + write_active_parent_snapshot(data_dir, path, bytes).await +} + +async fn mark_active_parent_snapshot_dirty( + data_dir: &std::path::Path, + key: &str, + correlation: &crate::correlation::CorrelationRegistry, +) -> Result<(), String> { + let Some(snapshot) = correlation.dirty_active_parent_snapshot(key) else { + return Ok(()); + }; + let bytes = serde_json::to_vec(&snapshot) + .map_err(|error| format!("active parent encoding failed: {error}"))?; + write_active_parent_snapshot(data_dir, active_parent_snapshot_path(data_dir, key), bytes).await +} + +async fn write_active_parent_snapshot( + data_dir: &std::path::Path, + path: PathBuf, + bytes: Vec, +) -> Result<(), String> { + let dir = data_dir.join("correlation").join("parents"); + tokio::fs::create_dir_all(&dir) + .await + .map_err(|error| format!("active parent directory failed: {error}"))?; + let temp = path.with_extension(format!("{}.tmp", uuid::Uuid::new_v4())); + tokio::fs::write(&temp, bytes) + .await + .map_err(|error| format!("active parent write failed: {error}"))?; + if let Err(first_error) = tokio::fs::rename(&temp, &path).await { + let _ = tokio::fs::remove_file(&path).await; + tokio::fs::rename(&temp, &path).await.map_err(|error| { + format!("active parent replace failed: {first_error}; retry failed: {error}") + })?; + } + Ok(()) +} + +async fn restore_active_parent_snapshots( + data_dir: &std::path::Path, + correlation: &crate::correlation::CorrelationRegistry, +) { + let Ok(mut entries) = tokio::fs::read_dir(data_dir.join("correlation").join("parents")).await + else { + return; + }; + while let Ok(Some(entry)) = entries.next_entry().await { + let Ok(bytes) = tokio::fs::read(entry.path()).await else { + continue; + }; + let Ok(snapshot) = + serde_json::from_slice::(&bytes) + else { + tracing::debug!(path = %entry.path().display(), "invalid active parent snapshot ignored"); + continue; + }; + correlation.restore_active_parent(snapshot); + } +} + +async fn accept_resolved_event(daemon: &Arc, mut env: Envelope) -> Result<(), String> { let source = env.source.clone(); let event = env.event.clone(); let session_id = env.session_id.clone(); let managed_run_id = env.managed_run_id.clone(); let route = env.route.clone(); tracing::info!(source, event, session_id, "event received"); - daemon.touch(); let session_lock = daemon.session_lock(&session_id); let _session_guard = session_lock.lock().await; @@ -634,10 +1072,24 @@ async fn accept_event(daemon: &Arc, mut env: Envelope) -> Result<(), Str .configure_event(&mut env) .await .map_err(|error| format!("session auth failed: {error}"))?; + if crate::dispatch::is_tool_lifecycle_event(&env.event) { + if let Err(error) = mark_active_parent_snapshot_dirty( + &daemon.data_dir, + &delivery_key.correlation_key(), + &daemon.correlation, + ) + .await + { + tracing::warn!(session_id = %env.session_id, %error, "active parent snapshot could not be marked dirty"); + } + } let session = daemon .session_for(&env, &delivery_key) .await .map_err(|error| format!("session init failed: {error}"))?; + daemon + .correlation + .observe_session(&delivery_key.correlation_key(), env.capture.as_ref()); daemon .append_to_journal(&mut env) .await @@ -664,7 +1116,24 @@ async fn accept_event(daemon: &Arc, mut env: Envelope) -> Result<(), Str result.map(|_| ()) } -async fn handle_request(daemon: &Arc, req: Request) -> Response { +fn is_session_start(event: &str) -> bool { + matches!(event, "SessionStart" | "session_start" | "session.created") +} + +fn automatic_link_key(env: &Envelope) -> String { + let route = env + .route + .as_ref() + .and_then(|route| serde_json::to_string(route).ok()) + .unwrap_or_default(); + format!("{}\u{1f}{}\u{1f}{route}", env.source, env.session_id) +} + +async fn handle_request( + daemon: &Arc, + req: Request, + client: &mut Option, +) -> Response { let id = req.id.clone(); let params = req.params.unwrap_or(serde_json::Value::Null); @@ -697,6 +1166,7 @@ async fn handle_request(daemon: &Arc, req: Request) -> Response { ), ); } + *client = Some(p.client); let result = InitializeResult { protocol_version: PROTOCOL_VERSION, daemon_version: daemon.version.clone(), @@ -707,7 +1177,8 @@ async fn handle_request(daemon: &Arc, req: Request) -> Response { Response::ok(id, serde_json::to_value(result).unwrap()) } method::EVENT_LOG => { - let env = parse!(Envelope); + let mut env = parse!(Envelope); + attach_process_capture(&mut env, client.as_ref()); match accept_event(daemon, env).await { Ok(()) => Response::ok( id, @@ -841,6 +1312,31 @@ async fn collect_garbage(data_dir: &std::path::Path) { journal::gc_old_journals(data_dir, RETENTION).await; journal::gc_old_managed_runs(data_dir, RETENTION).await; crate::transcript_mirror::gc_old_mirrors(data_dir, RETENTION).await; + gc_old_correlation_states(data_dir, RETENTION).await; +} + +async fn gc_old_correlation_states(data_dir: &std::path::Path, max_age: Duration) { + gc_old_correlation_dir(&data_dir.join("correlation"), max_age).await; + gc_old_correlation_dir(&data_dir.join("correlation").join("parents"), max_age).await; +} + +async fn gc_old_correlation_dir(dir: &std::path::Path, max_age: Duration) { + let Ok(mut entries) = tokio::fs::read_dir(dir).await else { + return; + }; + let now = SystemTime::now(); + while let Ok(Some(entry)) = entries.next_entry().await { + let old = entry + .metadata() + .await + .ok() + .and_then(|metadata| metadata.modified().ok()) + .and_then(|modified| now.duration_since(modified).ok()) + .is_some_and(|age| age > max_age); + if old && entry.file_type().await.is_ok_and(|kind| kind.is_file()) { + let _ = tokio::fs::remove_file(entry.path()).await; + } + } } /// Retire delivery pipelines that have gone quiet. Without this, every session @@ -889,7 +1385,10 @@ fn spawn_idle_watchdog(daemon: Arc, idle_timeout: Duration) { _ = tokio::time::sleep(tick) => {} } let idle_for = daemon.last_activity.lock().unwrap().elapsed(); - if idle_for >= idle_timeout && daemon.total_queued() == 0 { + if idle_for >= idle_timeout + && daemon.total_queued() == 0 + && !daemon.correlation.has_any_active_tools() + { tracing::info!("idle for {:?}; shutting down", idle_for); daemon.trigger_shutdown(); return; diff --git a/bt-daemon/src/transcript_import/mod.rs b/bt-daemon/src/transcript_import/mod.rs index d546fee..b4f539e 100644 --- a/bt-daemon/src/transcript_import/mod.rs +++ b/bt-daemon/src/transcript_import/mod.rs @@ -304,6 +304,7 @@ fn envelope( event: event.into(), ts_ms, managed_run_id: None, + capture: None, payload, route: None, config: None, diff --git a/bt-daemon/src/translate/opencode.rs b/bt-daemon/src/translate/opencode.rs index 69077a6..e8c9ad5 100644 --- a/bt-daemon/src/translate/opencode.rs +++ b/bt-daemon/src/translate/opencode.rs @@ -449,10 +449,33 @@ impl OpenCodeTranslator { return vec![]; }; if let Some(s) = self.sessions.get_mut(&sid) { + let Some(turn) = s.current_turn_span_id.clone() else { + return vec![]; + }; + if s.tool_starts.contains_key(call) { + return vec![]; + } s.tool_starts.insert(call.into(), event.ts_ms); if let Some(a) = event.payload.pointer("/output/args") { s.tool_args.insert(call.into(), a.clone()); } + let tool = event + .payload + .pointer("/input/tool") + .or_else(|| event.payload.get("tool")) + .and_then(Value::as_str) + .unwrap_or("tool"); + return vec![SpanOp::Insert(SpanRow { + span_id: ids::span_id(&self.daemon_session_id, &format!("tool:{sid}:{call}")), + root_span_id: s.effective_root_span_id.clone(), + parent_span_ids: vec![turn], + name: tool.into(), + span_type: SpanType::Tool, + start_ms: Some(event.ts_ms), + input: s.tool_args.get(call).cloned(), + metadata: Some(json!({"tool_name":tool,"call_id":call})), + ..Default::default() + })]; } vec![] } @@ -513,20 +536,27 @@ impl OpenCodeTranslator { .cloned() .unwrap_or(Value::Null) } - vec![SpanOp::Insert(SpanRow { + let had_start = s.tool_starts.contains_key(call); + let row = SpanRow { span_id: ids::span_id(&self.daemon_session_id, &format!("tool:{sid}:{call}")), root_span_id: s.effective_root_span_id.clone(), - parent_span_ids: vec![turn], + parent_span_ids: (!had_start).then_some(turn).into_iter().collect(), name, span_type: SpanType::Tool, - start_ms: s.tool_starts.remove(call).or(Some(event.ts_ms)), + start_ms: (!had_start).then(|| s.tool_starts.remove(call).unwrap_or(event.ts_ms)), end_ms: Some(event.ts_ms), - input: args, + input: (!had_start).then_some(args).flatten(), output, error, metadata: Some(metadata), ..Default::default() - })] + }; + s.tool_starts.remove(call); + vec![if had_start { + SpanOp::Merge(row) + } else { + SpanOp::Insert(row) + }] } fn permission(&mut self, event: &Envelope) -> Vec { diff --git a/bt-daemon/src/translate/pi.rs b/bt-daemon/src/translate/pi.rs index 46e5b1b..8360e39 100644 --- a/bt-daemon/src/translate/pi.rs +++ b/bt-daemon/src/translate/pi.rs @@ -49,6 +49,7 @@ struct PendingLlm { first_token_ms: Option, provider: Option, } +#[derive(Clone)] struct ToolStart { start_ms: i64, name: String, @@ -90,7 +91,7 @@ impl AgentTranslator for PiTranslator { .map(str::to_owned) } "message_end" => ops.extend(self.message_end(event, envelope.ts_ms)), - "tool_execution_start" => self.tool_start(event, envelope.ts_ms), + "tool_execution_start" => ops.extend(self.tool_start(event, envelope.ts_ms)), "tool_execution_end" => ops.extend(self.tool_end(event, envelope.ts_ms)), "agent_end" if event.get("willRetry").and_then(Value::as_bool) != Some(true) => { ops.extend(self.close_turn(envelope.ts_ms, None)); @@ -348,22 +349,45 @@ impl PiTranslator { ..Default::default() })] } - fn tool_start(&mut self, event: &Value, ts: i64) { + fn tool_start(&mut self, event: &Value, ts: i64) -> Vec { let Some(id) = event.get("toolCallId").and_then(Value::as_str) else { - return; + return vec![]; + }; + let Some((turn, _)) = &self.turn else { + return vec![]; }; + let name = event + .get("toolName") + .and_then(Value::as_str) + .unwrap_or("tool") + .to_string(); + let args = event.get("args").cloned().unwrap_or(Value::Null); + if self.tools.contains_key(id) { + return vec![]; + } self.tools.insert( id.into(), ToolStart { start_ms: ts, - name: event - .get("toolName") - .and_then(Value::as_str) - .unwrap_or("tool") - .into(), - args: event.get("args").cloned().unwrap_or(Value::Null), + name: name.clone(), + args: args.clone(), }, ); + vec![SpanOp::Insert(SpanRow { + span_id: ids::span_id(&self.session_id, &format!("tool:{}:{id}", self.turn_seq)), + root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![turn.clone()], + name, + span_type: SpanType::Tool, + start_ms: Some(ts), + input: Some(args), + metadata: Some(json!({ + "tool_name": event.get("toolName").and_then(Value::as_str).unwrap_or("tool"), + "tool_call_id": id, + "tool_approval": "approved", + })), + ..Default::default() + })] } fn tool_end(&mut self, event: &Value, ts: i64) -> Vec { let Some((turn, _)) = &self.turn else { @@ -373,7 +397,8 @@ impl PiTranslator { .get("toolCallId") .and_then(Value::as_str) .unwrap_or(""); - let tracked = self.tools.remove(call).unwrap_or(ToolStart { + let pending = self.tools.remove(call); + let tracked = pending.clone().unwrap_or(ToolStart { start_ms: ts, name: event .get("toolName") @@ -389,15 +414,19 @@ impl PiTranslator { .as_ref() .map(|s| format!("skill: {s}")) .unwrap_or_else(|| tracked.name.clone()); - vec![SpanOp::Insert(SpanRow { + let row = SpanRow { span_id: ids::span_id(&self.session_id, &format!("tool:{}:{call}", self.turn_seq)), root_span_id: self.effective_root_span_id.clone(), - parent_span_ids: vec![turn.clone()], + parent_span_ids: pending + .is_none() + .then(|| turn.clone()) + .into_iter() + .collect(), name, span_type: SpanType::Tool, - start_ms: Some(tracked.start_ms), + start_ms: pending.is_none().then_some(tracked.start_ms), end_ms: Some(ts), - input: Some(tracked.args), + input: pending.is_none().then_some(tracked.args), output: event.get("result").cloned(), metadata: Some(json!({ "tool_name": if skill.is_some() { "skill" } else { &tracked.name }, @@ -409,7 +438,12 @@ impl PiTranslator { })), error: failed.then(|| format_error(event.get("result"))), ..Default::default() - })] + }; + vec![if pending.is_some() { + SpanOp::Merge(row) + } else { + SpanOp::Insert(row) + }] } fn finish_special( &mut self, diff --git a/bt-daemon/src/wire/envelope.rs b/bt-daemon/src/wire/envelope.rs index 9f7ef73..7969dc1 100644 --- a/bt-daemon/src/wire/envelope.rs +++ b/bt-daemon/src/wire/envelope.rs @@ -4,6 +4,30 @@ use braintrust_sdk_rust::SpanComponents; use serde::{Deserialize, Serialize}; +/// One operating-system process observed while capturing an event. +/// +/// A PID alone is not a stable identity because operating systems reuse it. +/// `start_time_secs`, when available, distinguishes different occupants of the +/// same PID. +#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)] +pub struct ProcessIdentity { + pub pid: u32, + /// Epoch seconds reported by the operating system, or zero if unavailable. + #[serde(default)] + pub start_time_secs: u64, +} + +/// Process evidence captured at the hook or in-process adapter boundary. +/// +/// `process_chain` is ordered from the process that connected to the daemon +/// toward the operating-system root. It intentionally excludes command lines, +/// environment variables, and working directories. +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct CaptureContext { + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub process_chain: Vec, +} + /// One captured hook event, forwarded from a shim to the daemon. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct Envelope { @@ -27,6 +51,9 @@ pub struct Envelope { /// it only to flush the sessions created by one managed child process tree. #[serde(default, skip_serializing_if = "Option::is_none")] pub managed_run_id: Option, + /// Daemon-captured process ancestry for local cross-agent correlation. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub capture: Option, /// The raw agent-native hook payload; opaque except to the translator. pub payload: serde_json::Value, /// Non-secret, immutable routing intent for this session. New clients use @@ -251,6 +278,7 @@ impl Envelope { event: self.event.clone(), ts_ms: self.ts_ms, managed_run_id: self.managed_run_id.clone(), + capture: self.capture.clone(), payload: self.payload.clone(), route: self.route.clone(), } @@ -271,6 +299,8 @@ pub struct RedactedEnvelope { pub ts_ms: i64, #[serde(default, skip_serializing_if = "Option::is_none")] pub managed_run_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub capture: Option, pub payload: serde_json::Value, #[serde(default, skip_serializing_if = "Option::is_none")] pub route: Option, @@ -289,6 +319,12 @@ mod tests { event: "PostToolUse".into(), ts_ms: 1_753_639_552_123, managed_run_id: Some("run-1".into()), + capture: Some(CaptureContext { + process_chain: vec![ProcessIdentity { + pid: 42, + start_time_secs: 1_753_639_500, + }], + }), payload: serde_json::json!({ "session_id": "sess-1", "tool_name": "shell" }), route: Some(SessionRoute { auth: AuthSelection { @@ -324,6 +360,13 @@ mod tests { let back: Envelope = serde_json::from_str(&s).unwrap(); assert_eq!(back.session_id, "sess-1"); assert_eq!(back.managed_run_id.as_deref(), Some("run-1")); + assert_eq!( + back.capture.as_ref().unwrap().process_chain[0], + ProcessIdentity { + pid: 42, + start_time_secs: 1_753_639_500, + } + ); assert_eq!(back.route.unwrap().auth.profile.as_deref(), Some("work")); assert!(back.config.is_none()); assert!(!s.contains("sk-super-secret")); @@ -340,6 +383,17 @@ mod tests { ); assert_eq!(r.route.unwrap().auth.profile.as_deref(), Some("work")); assert_eq!(r.managed_run_id.as_deref(), Some("run-1")); + assert_eq!(r.capture.unwrap().process_chain[0].pid, 42); + } + + #[test] + fn capture_context_is_backward_compatible() { + let mut value = serde_json::to_value(sample()).unwrap(); + value.as_object_mut().unwrap().remove("capture"); + + let envelope: Envelope = serde_json::from_value(value).unwrap(); + + assert!(envelope.capture.is_none()); } #[test] diff --git a/bt-daemon/src/wire/mod.rs b/bt-daemon/src/wire/mod.rs index 7e112ad..c8a2349 100644 --- a/bt-daemon/src/wire/mod.rs +++ b/bt-daemon/src/wire/mod.rs @@ -10,8 +10,8 @@ mod methods; mod rpc; pub use envelope::{ - AuthSelection, AuthSource, BackendAuth, Envelope, FlushMode, RedactedEnvelope, SessionConfig, - SessionRoute, TraceDestination, + AuthSelection, AuthSource, BackendAuth, CaptureContext, Envelope, FlushMode, ProcessIdentity, + RedactedEnvelope, SessionConfig, SessionRoute, TraceDestination, }; pub use methods::{ method, Capabilities, ClientInfo, EventLogResult, FlushParams, FlushResult, InitializeParams, diff --git a/bt-daemon/tests/claude_translator.rs b/bt-daemon/tests/claude_translator.rs index e8d7e2f..eb2dfa3 100644 --- a/bt-daemon/tests/claude_translator.rs +++ b/bt-daemon/tests/claude_translator.rs @@ -80,6 +80,7 @@ fn replay_from(name: &str, source: Source) -> Vec { payload, route: None, config: None, + capture: None, }; ops.extend(translator.handle(&env, &ctx).unwrap()); while let Some(batch) = translator.drain_pending(&ctx).unwrap() { @@ -271,6 +272,7 @@ fn claude_additional_metadata_reaches_roots_without_overriding_session_fields() payload: json!({"session_id":"session","cwd":"/workspace","prompt":"go"}), route: None, config: None, + capture: None, }, &ctx, ) @@ -370,6 +372,7 @@ fn claude_permission_denied_and_failed_tools_are_first_class_spans() { payload, route: None, config: None, + capture: None, }; let mut ops = translator .handle( @@ -459,6 +462,7 @@ fn claude_pairs_tool_lifecycle_and_marks_explicit_skills_and_stop_failures() { payload, route: None, config: None, + capture: None, }; let mut ops = Vec::new(); for envelope in [ @@ -578,6 +582,7 @@ fn claude_groups_streamed_rows_and_reads_late_final_output_at_session_end() { payload, route: None, config: None, + capture: None, }; let mut ops = translator .handle( @@ -714,6 +719,7 @@ fn claude_large_catch_up_emits_one_historical_snapshot_per_batch() { payload, route: None, config: None, + capture: None, }; translator .handle( diff --git a/bt-daemon/tests/codex_translator.rs b/bt-daemon/tests/codex_translator.rs index 1bc16ed..1281456 100644 --- a/bt-daemon/tests/codex_translator.rs +++ b/bt-daemon/tests/codex_translator.rs @@ -65,6 +65,7 @@ fn envelope(session: &str, event: &str, transcript_path: &str, extra: Value) -> payload, route: None, config: None, + capture: None, } } diff --git a/bt-daemon/tests/distributed_tracing.rs b/bt-daemon/tests/distributed_tracing.rs new file mode 100644 index 0000000..c7470e1 --- /dev/null +++ b/bt-daemon/tests/distributed_tracing.rs @@ -0,0 +1,2220 @@ +//! Deterministic cross-agent linkage through the complete daemon pipeline. +//! +//! These tests intentionally use native agent envelopes rather than invoking +//! installed coding-agent binaries. That keeps the complete 4x4 compatibility +//! matrix in the ordinary test suite while still exercising IPC, journaling, +//! translation, automatic routing, and sink configuration. + +mod support; + +use async_trait::async_trait; +use bt_daemon::wire::{AuthSelection, BackendAuth, Envelope, SessionConfig, TraceDestination}; +use bt_daemon::{ + flush_session, forward_envelope, run_serve, run_status, shutdown_daemon, AuthLease, + AuthProvider, AuthResolveReason, HostInfo, Registry, ServeArgs, ServeOptions, Sink, + SinkFactory, SpanOp, SpanRow, SpanType, StatusArgs, +}; +use std::collections::HashMap; +use std::ffi::OsString; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; +use support::distributed::{AgentKind, DistributedFixtures, ProcessTree}; + +#[derive(Default)] +struct RecordedSession { + source: String, + configs: Mutex>, + ops: Mutex>, +} + +#[derive(Default)] +struct RecordingSinkFactory { + sessions: Mutex>>, +} + +impl RecordingSinkFactory { + fn session(&self, session_id: &str) -> Arc { + self.sessions + .lock() + .unwrap() + .get(session_id) + .unwrap_or_else(|| panic!("no sink created for session {session_id}")) + .clone() + } +} + +impl SinkFactory for RecordingSinkFactory { + fn create( + &self, + session_id: &str, + source: &str, + _plugin_version: Option<&str>, + ) -> anyhow::Result> { + let record = Arc::new(RecordedSession { + source: source.into(), + ..RecordedSession::default() + }); + self.sessions + .lock() + .unwrap() + .insert(session_id.into(), record.clone()); + Ok(Box::new(RecordingSink { record })) + } +} + +struct RecordingSink { + record: Arc, +} + +#[async_trait] +impl Sink for RecordingSink { + fn configure(&mut self, config: &SessionConfig) { + self.record.configs.lock().unwrap().push(config.clone()); + } + + async fn emit(&mut self, ops: &[SpanOp]) -> anyhow::Result { + self.record.ops.lock().unwrap().extend_from_slice(ops); + Ok(ops.len() as u64) + } + + async fn flush(&mut self) -> anyhow::Result<()> { + Ok(()) + } +} + +struct TestAuthProvider; + +#[async_trait] +impl AuthProvider for TestAuthProvider { + async fn resolve( + &self, + selection: &AuthSelection, + _reason: AuthResolveReason, + ) -> anyhow::Result { + Ok(AuthLease { + selection: selection.clone().canonicalized()?, + auth: BackendAuth { + token: "test-token".into(), + api_url: Some("http://127.0.0.1.invalid".into()), + app_url: Some("http://127.0.0.1.invalid".into()), + org_name: Some("test".into()), + org_id: Some("test-org".into()), + }, + expires_at_ms: None, + }) + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn every_instrumented_parent_child_agent_pair_shares_one_trace() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let mut pair_index = 0_u32; + + for parent_kind in AgentKind::ALL { + for child_kind in AgentKind::ALL { + pair_index += 1; + let label = format!("{}-to-{}", parent_kind.label(), child_kind.label()); + let parent_session = format!("parent-{label}"); + let child_session = format!("child-{label}"); + let call_id = format!("call-{label}"); + let delegated_prompt = format!("LINK_PROMPT_{label}"); + let child_output = format!("LINK_OUTPUT_{label}"); + let base_pid = 1_000 + pair_index * 10; + let parent_tree = ProcessTree::root(base_pid); + let child_tree = ProcessTree::child(base_pid + 2, base_pid + 1, &parent_tree); + let base_ts = 1_700_000_000_000 + i64::from(pair_index) * 1_000; + + forward_all( + &mut fixtures.start_turn( + parent_kind, + &parent_session, + &parent_tree, + "Delegate the marked task.", + base_ts, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + parent_kind, + &parent_session, + &parent_tree, + &call_id, + &delegated_prompt, + base_ts + 10, + ), + &socket, + &host, + ) + .await; + + // The child starts immediately after the blocking start-tool hook. + // No parent flush is inserted here: forwarding the start event must + // not return until its correlation marker is visible. + forward_all( + &mut fixtures.start_turn( + child_kind, + &child_session, + &child_tree, + &delegated_prompt, + base_ts + 20, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + child_kind, + &child_session, + &child_tree, + &child_output, + base_ts + 30, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + parent_kind, + &parent_session, + &parent_tree, + &call_id, + &delegated_prompt, + &child_output, + base_ts + 40, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + parent_kind, + &parent_session, + &parent_tree, + "Parent complete", + base_ts + 50, + ), + &socket, + &host, + ) + .await; + + flush(&child_session, &socket).await; + flush(&parent_session, &socket).await; + assert_pair_linked(recording.as_ref(), &parent_session, &child_session, &label); + } + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn every_agent_pair_links_when_the_daemon_restarts_between_spawn_and_child_start() { + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let mut pair_index = 0_u32; + + for parent_kind in AgentKind::ALL { + for child_kind in AgentKind::ALL { + pair_index += 1; + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let label = format!("restart-{}-to-{}", parent_kind.label(), child_kind.label()); + let parent_session = format!("parent-{label}"); + let child_session = format!("child-{label}"); + let call_id = format!("call-{label}"); + let delegated_prompt = format!("DURABLE_LINK_PROMPT_{label}"); + let base_pid = 30_000 + pair_index * 10; + let parent_tree = ProcessTree::root(base_pid); + let child_tree = ProcessTree::child(base_pid + 2, base_pid + 1, &parent_tree); + let base_ts = 1_702_000_000_000 + i64::from(pair_index) * 1_000; + + let first = spawn_daemon(socket.clone(), data_dir.clone(), recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + parent_kind, + &parent_session, + &parent_tree, + "Delegate across a daemon restart.", + base_ts, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + parent_kind, + &parent_session, + &parent_tree, + &call_id, + &delegated_prompt, + base_ts + 10, + ), + &socket, + &host, + ) + .await; + + let parent_snapshot_dir = data_dir.join("correlation").join("parents"); + let mut snapshots = tokio::fs::read_dir(&parent_snapshot_dir).await.unwrap(); + let snapshot = snapshots + .next_entry() + .await + .unwrap() + .expect("active snapshot"); + let bytes = tokio::fs::read(snapshot.path()).await.unwrap(); + let serialized = String::from_utf8(bytes).unwrap(); + assert!(!serialized.contains(&delegated_prompt)); + assert!(!serialized.contains("test-token")); + + shutdown_daemon(&socket).await.unwrap(); + first.await.unwrap(); + + let second = spawn_daemon(socket.clone(), data_dir, recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + child_kind, + &child_session, + &child_tree, + &delegated_prompt, + base_ts + 20, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + child_kind, + &child_session, + &child_tree, + "Child completed after daemon restart.", + base_ts + 30, + ), + &socket, + &host, + ) + .await; + flush(&child_session, &socket).await; + assert_pair_linked(recording.as_ref(), &parent_session, &child_session, &label); + + shutdown_daemon(&socket).await.unwrap(); + second.await.unwrap(); + } + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn mixed_recursive_hierarchy_survives_a_restart_at_every_spawn_boundary() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let kinds = AgentKind::ALL; + let sessions = [ + "durable-recursive-claude", + "durable-recursive-codex", + "durable-recursive-pi", + "durable-recursive-opencode", + ]; + let prompts = [ + "durable recursive level one", + "durable recursive level two", + "durable recursive level three", + ]; + let trees = [ + ProcessTree::root(35_000), + ProcessTree::child(35_002, 35_001, &ProcessTree::root(35_000)), + ProcessTree::child( + 35_004, + 35_003, + &ProcessTree::child(35_002, 35_001, &ProcessTree::root(35_000)), + ), + ProcessTree::child( + 35_006, + 35_005, + &ProcessTree::child( + 35_004, + 35_003, + &ProcessTree::child(35_002, 35_001, &ProcessTree::root(35_000)), + ), + ), + ]; + let mut fixtures = DistributedFixtures::new(); + + for level in 0..kinds.len() { + let daemon = spawn_daemon(socket.clone(), data_dir.clone(), recording.clone()).await; + let prompt = if level == 0 { + "begin durable recursive hierarchy" + } else { + prompts[level - 1] + }; + forward_all( + &mut fixtures.start_turn( + kinds[level], + sessions[level], + &trees[level], + prompt, + 1_704_000_000_000 + level as i64 * 100, + ), + &socket, + &host, + ) + .await; + if level < prompts.len() { + forward( + fixtures.open_tool( + kinds[level], + sessions[level], + &trees[level], + &format!("durable-recursive-call-{level}"), + prompts[level], + 1_704_000_000_010 + level as i64 * 100, + ), + &socket, + &host, + ) + .await; + } else { + forward( + fixtures.close_session( + kinds[level], + sessions[level], + &trees[level], + "deepest child complete", + 1_704_000_000_010 + level as i64 * 100, + ), + &socket, + &host, + ) + .await; + flush(sessions[level], &socket).await; + } + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); + } + + let root = recording.session(sessions[0]); + let root_id = session_root(&inserted_rows(&root), sessions[0], "durable recursive root") + .root_span_id + .clone(); + for level in 0..prompts.len() { + let parent = recording.session(sessions[level]); + let tool = inserted_rows(&parent) + .into_iter() + .find(|row| { + row.span_type == SpanType::Tool + && row + .input + .as_ref() + .is_some_and(|input| input.to_string().contains(prompts[level])) + }) + .unwrap_or_else(|| panic!("durable recursive level {level}: missing tool")); + let child = recording.session(sessions[level + 1]); + let child_root = session_root( + &inserted_rows(&child), + sessions[level + 1], + "durable recursive child", + ) + .clone(); + assert_eq!(child_root.parent_span_ids, vec![tool.span_id.clone()]); + let configs = child.configs.lock().unwrap(); + let components = configs + .iter() + .find_map(|config| match &config.destination { + Some(TraceDestination::ParentSpan { components }) => Some(components), + _ => None, + }) + .unwrap_or_else(|| panic!("durable recursive level {level}: child not attached")); + assert_eq!(components.span_id.as_deref(), Some(tool.span_id.as_str())); + assert_eq!(components.root_span_id.as_deref(), Some(root_id.as_str())); + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn concurrent_tools_are_disambiguated_by_prompt_fingerprints() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(9_000); + let child_tree = ProcessTree::child(9_002, 9_001, &parent_tree); + let parent_session = "concurrent-parent"; + let child_session = "concurrent-child"; + let first_prompt = "inspect the unrelated alpha candidate carefully"; + let selected_prompt = "inspect the selected beta candidate carefully"; + + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + parent_session, + &parent_tree, + "Run two delegated tasks.", + 1_700_100_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + parent_session, + &parent_tree, + "call-alpha", + first_prompt, + 1_700_100_000_010, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + parent_session, + &parent_tree, + "call-beta", + selected_prompt, + 1_700_100_000_011, + ), + &socket, + &host, + ) + .await; + + // SessionStart alone is ambiguous and is held. UserPromptSubmit supplies + // the opaque-content fingerprints that select call-beta. + forward_all( + &mut fixtures.start_turn( + AgentKind::Pi, + child_session, + &child_tree, + selected_prompt, + 1_700_100_000_020, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + AgentKind::Pi, + child_session, + &child_tree, + "selected beta complete", + 1_700_100_000_030, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + AgentKind::Claude, + parent_session, + &parent_tree, + "call-alpha", + first_prompt, + "alpha complete", + 1_700_100_000_040, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + AgentKind::Claude, + parent_session, + &parent_tree, + "call-beta", + selected_prompt, + "beta complete", + 1_700_100_000_041, + ), + &socket, + &host, + ) + .await; + flush(child_session, &socket).await; + flush(parent_session, &socket).await; + + let parent = recording.session(parent_session); + let child = recording.session(child_session); + let selected_tool = inserted_rows(&parent) + .into_iter() + .find(|row| { + row.span_type == SpanType::Tool + && row + .input + .as_ref() + .is_some_and(|input| input.to_string().contains(selected_prompt)) + }) + .expect("selected parent tool span"); + let child_root = session_root( + &inserted_rows(&child), + child_session, + "concurrent fingerprint match", + ) + .clone(); + assert_eq!(child_root.parent_span_ids, vec![selected_tool.span_id]); + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn mixed_agents_link_recursively_through_every_generation() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let kinds = [ + AgentKind::Claude, + AgentKind::Codex, + AgentKind::Pi, + AgentKind::OpenCode, + ]; + let sessions = [ + "recursive-claude", + "recursive-codex", + "recursive-pi", + "recursive-opencode", + ]; + let prompts = [ + "delegate recursive level one", + "delegate recursive level two", + "delegate recursive level three", + ]; + let trees = [ + ProcessTree::root(11_000), + ProcessTree::child(11_002, 11_001, &ProcessTree::root(11_000)), + ProcessTree::child( + 11_004, + 11_003, + &ProcessTree::child(11_002, 11_001, &ProcessTree::root(11_000)), + ), + ProcessTree::child( + 11_006, + 11_005, + &ProcessTree::child( + 11_004, + 11_003, + &ProcessTree::child(11_002, 11_001, &ProcessTree::root(11_000)), + ), + ), + ]; + + forward_all( + &mut fixtures.start_turn( + kinds[0], + sessions[0], + &trees[0], + "start recursive work", + 1_700_200_000_000, + ), + &socket, + &host, + ) + .await; + for level in 0..3 { + forward( + fixtures.open_tool( + kinds[level], + sessions[level], + &trees[level], + &format!("recursive-call-{level}"), + prompts[level], + 1_700_200_000_010 + level as i64 * 20, + ), + &socket, + &host, + ) + .await; + forward_all( + &mut fixtures.start_turn( + kinds[level + 1], + sessions[level + 1], + &trees[level + 1], + prompts[level], + 1_700_200_000_020 + level as i64 * 20, + ), + &socket, + &host, + ) + .await; + } + + forward( + fixtures.close_session( + kinds[3], + sessions[3], + &trees[3], + "leaf complete", + 1_700_200_000_100, + ), + &socket, + &host, + ) + .await; + for level in (0..3).rev() { + forward( + fixtures.close_tool( + kinds[level], + sessions[level], + &trees[level], + &format!("recursive-call-{level}"), + prompts[level], + "child complete", + 1_700_200_000_110 + (2 - level) as i64 * 20, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + kinds[level], + sessions[level], + &trees[level], + "level complete", + 1_700_200_000_120 + (2 - level) as i64 * 20, + ), + &socket, + &host, + ) + .await; + } + for session in sessions { + flush(session, &socket).await; + } + + let root_record = recording.session(sessions[0]); + let root_id = session_root(&inserted_rows(&root_record), sessions[0], "recursive root") + .root_span_id + .clone(); + for level in 0..3 { + let parent = recording.session(sessions[level]); + let child = recording.session(sessions[level + 1]); + let tool = inserted_rows(&parent) + .into_iter() + .find(|row| row.span_type == SpanType::Tool) + .unwrap_or_else(|| panic!("recursive level {level}: missing tool")); + let child_root = session_root( + &inserted_rows(&child), + sessions[level + 1], + "recursive child", + ) + .clone(); + assert_eq!(child_root.parent_span_ids, vec![tool.span_id.clone()]); + let configs = child.configs.lock().unwrap(); + let components = configs + .iter() + .find_map(|config| match &config.destination { + Some(TraceDestination::ParentSpan { components }) => Some(components), + _ => None, + }) + .unwrap_or_else(|| panic!("recursive level {level}: child not attached")); + assert_eq!(components.span_id.as_deref(), Some(tool.span_id.as_str())); + assert_eq!(components.root_span_id.as_deref(), Some(root_id.as_str())); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn indistinguishable_concurrent_tools_fail_safe_without_a_parent() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(12_000); + let child_tree = ProcessTree::child(12_002, 12_001, &parent_tree); + let prompt = "the exact same delegated task for both candidates"; + + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "ambiguous-parent", + &parent_tree, + "delegate twice", + 1_700_300_000_000, + ), + &socket, + &host, + ) + .await; + for call in ["ambiguous-a", "ambiguous-b"] { + forward( + fixtures.open_tool( + AgentKind::Claude, + "ambiguous-parent", + &parent_tree, + call, + prompt, + 1_700_300_000_010, + ), + &socket, + &host, + ) + .await; + } + forward_all( + &mut fixtures.start_turn( + AgentKind::Pi, + "ambiguous-child", + &child_tree, + prompt, + 1_700_300_000_020, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + AgentKind::Pi, + "ambiguous-child", + &child_tree, + "done", + 1_700_300_000_030, + ), + &socket, + &host, + ) + .await; + for call in ["ambiguous-a", "ambiguous-b"] { + forward( + fixtures.close_tool( + AgentKind::Claude, + "ambiguous-parent", + &parent_tree, + call, + prompt, + "same indistinguishable output", + 1_700_300_000_040, + ), + &socket, + &host, + ) + .await; + } + flush("ambiguous-child", &socket).await; + + let child = recording.session("ambiguous-child"); + assert!( + child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + )), + "an ambiguous child must never be attached by guessing" + ); + let child_root = session_root( + &inserted_rows(&child), + "ambiguous-child", + "ambiguous standalone", + ) + .clone(); + assert!(child_root.parent_span_ids.is_empty()); + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn resolved_child_link_survives_daemon_restart() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let first_recording = Arc::new(RecordingSinkFactory::default()); + let first_daemon = + spawn_daemon(socket.clone(), data_dir.clone(), first_recording.clone()).await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(13_000); + let child_tree = ProcessTree::child(13_002, 13_001, &parent_tree); + let prompt = "persist this exact distributed child link"; + + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "restart-parent", + &parent_tree, + "delegate", + 1_700_400_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "restart-parent", + &parent_tree, + "restart-call", + prompt, + 1_700_400_000_010, + ), + &socket, + &host, + ) + .await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Pi, + "restart-child", + &child_tree, + prompt, + 1_700_400_000_020, + ), + &socket, + &host, + ) + .await; + flush("restart-child", &socket).await; + shutdown_daemon(&socket).await.unwrap(); + first_daemon.await.unwrap(); + + let second_recording = Arc::new(RecordingSinkFactory::default()); + let second_daemon = spawn_daemon(socket.clone(), data_dir, second_recording.clone()).await; + forward( + fixtures.close_session( + AgentKind::Pi, + "restart-child", + &child_tree, + "complete after restart", + 1_700_400_000_030, + ), + &socket, + &host, + ) + .await; + flush("restart-child", &socket).await; + + let child = second_recording.session("restart-child"); + assert!(child.configs.lock().unwrap().iter().any(|config| matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + shutdown_daemon(&socket).await.unwrap(); + second_daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn ambiguous_child_evidence_survives_daemon_restart() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let parent_tree = ProcessTree::root(14_000); + let child_tree = ProcessTree::child(14_002, 14_001, &parent_tree); + let selected = "select durable beta evidence after restart"; + let mut fixtures = DistributedFixtures::new(); + + let first_recording = Arc::new(RecordingSinkFactory::default()); + let first_daemon = + spawn_daemon(socket.clone(), data_dir.clone(), first_recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "pending-parent", + &parent_tree, + "delegate twice", + 1_700_500_000_000, + ), + &socket, + &host, + ) + .await; + for (call, prompt) in [ + ("pending-alpha", "unrelated durable alpha evidence"), + ("pending-beta", selected), + ] { + forward( + fixtures.open_tool( + AgentKind::Claude, + "pending-parent", + &parent_tree, + call, + prompt, + 1_700_500_000_010, + ), + &socket, + &host, + ) + .await; + } + let mut child_start = fixtures.start_turn( + AgentKind::Pi, + "pending-child", + &child_tree, + selected, + 1_700_500_000_020, + ); + forward(child_start.remove(0), &socket, &host).await; + assert!(tokio::fs::read_dir(data_dir.join("correlation")) + .await + .unwrap() + .next_entry() + .await + .unwrap() + .is_some()); + shutdown_daemon(&socket).await.unwrap(); + first_daemon.await.unwrap(); + + let second_recording = Arc::new(RecordingSinkFactory::default()); + let second_daemon = spawn_daemon(socket.clone(), data_dir, second_recording.clone()).await; + // No parent event is sent after restart. Both candidates and their hashed + // evidence must come from the compact active-parent snapshot. + forward(child_start.remove(0), &socket, &host).await; + forward( + fixtures.close_session( + AgentKind::Pi, + "pending-child", + &child_tree, + "resolved", + 1_700_500_000_040, + ), + &socket, + &host, + ) + .await; + flush("pending-child", &socket).await; + + let parent = first_recording.session("pending-parent"); + let beta = inserted_rows(&parent) + .into_iter() + .find(|row| { + row.span_type == SpanType::Tool + && row + .input + .as_ref() + .is_some_and(|input| input.to_string().contains(selected)) + }) + .expect("replayed beta tool"); + let child = second_recording.session("pending-child"); + let child_root = session_root( + &inserted_rows(&child), + "pending-child", + "durable pending child", + ) + .clone(); + assert_eq!(child_root.parent_span_ids, vec![beta.span_id]); + + shutdown_daemon(&socket).await.unwrap(); + second_daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn completed_tools_are_not_resurrected_after_restart() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(34_000); + let prompt = "completed snapshot must be removed"; + + let first = spawn_daemon(socket.clone(), data_dir.clone(), recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "closed-restart-parent", + &parent_tree, + "delegate", + 1_703_000_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "closed-restart-parent", + &parent_tree, + "closed-restart-call", + prompt, + 1_703_000_000_010, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + AgentKind::Claude, + "closed-restart-parent", + &parent_tree, + "closed-restart-call", + prompt, + "done", + 1_703_000_000_020, + ), + &socket, + &host, + ) + .await; + shutdown_daemon(&socket).await.unwrap(); + first.await.unwrap(); + + let second = spawn_daemon(socket.clone(), data_dir, recording.clone()).await; + for (index, reused_pid) in [false, true].into_iter().enumerate() { + let child_session = format!("closed-restart-child-{index}"); + let child_tree = ProcessTree::child(34_010 + index as u32, 34_001, &parent_tree); + let mut events = fixtures.start_turn( + AgentKind::Pi, + &child_session, + &child_tree, + prompt, + 1_703_000_000_030 + index as i64, + ); + if reused_pid { + for event in &mut events { + for process in &mut event.capture.as_mut().unwrap().process_chain { + if process.pid == parent_tree.agent.pid { + process.start_time_secs += 1; + } + } + } + } + forward_all(&mut events, &socket, &host).await; + forward( + fixtures.close_session( + AgentKind::Pi, + &child_session, + &child_tree, + "standalone", + 1_703_000_000_040 + index as i64, + ), + &socket, + &host, + ) + .await; + flush(&child_session, &socket).await; + let child = recording.session(&child_session); + assert!(child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + } + + shutdown_daemon(&socket).await.unwrap(); + second.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn interrupted_tool_transition_snapshot_fails_safe_after_restart() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(34_100); + let child_tree = ProcessTree::child(34_102, 34_101, &parent_tree); + let prompt = "dirty transition must never resurrect a parent"; + + let first = spawn_daemon(socket.clone(), data_dir.clone(), recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "dirty-restart-parent", + &parent_tree, + "delegate", + 1_703_100_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "dirty-restart-parent", + &parent_tree, + "dirty-restart-call", + prompt, + 1_703_100_000_010, + ), + &socket, + &host, + ) + .await; + shutdown_daemon(&socket).await.unwrap(); + first.await.unwrap(); + + let parent_dir = data_dir.join("correlation").join("parents"); + let mut entries = tokio::fs::read_dir(parent_dir).await.unwrap(); + let path = entries.next_entry().await.unwrap().unwrap().path(); + let mut snapshot: serde_json::Value = + serde_json::from_slice(&tokio::fs::read(&path).await.unwrap()).unwrap(); + snapshot["dirty"] = serde_json::Value::Bool(true); + tokio::fs::write(&path, serde_json::to_vec(&snapshot).unwrap()) + .await + .unwrap(); + + let second = spawn_daemon(socket.clone(), data_dir, recording.clone()).await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Pi, + "dirty-restart-child", + &child_tree, + prompt, + 1_703_100_000_020, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + AgentKind::Pi, + "dirty-restart-child", + &child_tree, + "standalone", + 1_703_100_000_030, + ), + &socket, + &host, + ) + .await; + flush("dirty-restart-child", &socket).await; + let child = recording.session("dirty-restart-child"); + assert!(child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + + shutdown_daemon(&socket).await.unwrap(); + second.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn corrupt_active_parent_snapshot_is_ignored() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let parent_dir = data_dir.join("correlation").join("parents"); + tokio::fs::create_dir_all(&parent_dir).await.unwrap(); + tokio::fs::write(parent_dir.join("corrupt.json"), b"{not-json") + .await + .unwrap(); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let daemon = spawn_daemon(socket.clone(), data_dir, recording).await; + + assert!(run_status(StatusArgs { + socket: Some(socket.clone()), + session_id: None, + }) + .await + .unwrap() + .is_some()); + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn restored_parent_snapshot_does_not_pin_an_otherwise_idle_daemon() { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(36_000); + + let first = spawn_daemon(socket.clone(), data_dir.clone(), recording).await; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "idle-restored-parent", + &parent_tree, + "delegate before restart", + 1_705_000_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "idle-restored-parent", + &parent_tree, + "idle-restored-call", + "child may arrive after restart", + 1_705_000_000_010, + ), + &socket, + &host, + ) + .await; + shutdown_daemon(&socket).await.unwrap(); + first.await.unwrap(); + + let second = spawn_daemon_with_timeouts( + socket, + data_dir, + Arc::new(RecordingSinkFactory::default()), + 1, + 1, + ) + .await; + tokio::time::timeout(Duration::from_secs(4), second) + .await + .expect("restored snapshot pinned idle daemon") + .unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn ipc_initialize_pid_is_captured_without_hook_metadata() { + let (socket, daemon, _recording, tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let mut event = fixtures + .start_turn( + AgentKind::Claude, + "automatic-process-capture", + &ProcessTree::root(15_000), + "capture automatically", + 1_700_600_000_000, + ) + .remove(0); + event.capture = None; + forward(event, &socket, &host).await; + flush("automatic-process-capture", &socket).await; + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); + + let journal = tokio::fs::read_to_string( + tmp.path() + .join("data/journal/automatic-process-capture.ndjson"), + ) + .await + .unwrap(); + let row: serde_json::Value = serde_json::from_str(journal.lines().next().unwrap()).unwrap(); + let chain = row + .pointer("/capture/process_chain") + .and_then(serde_json::Value::as_array) + .expect("daemon-added process ancestry in journal"); + assert_eq!( + chain.first().and_then(|process| process.get("pid")), + Some(&serde_json::json!(std::process::id())) + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn completed_tools_are_not_candidates_for_later_children() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + + for (index, kind) in AgentKind::ALL.into_iter().enumerate() { + let pid = 16_000 + index as u32 * 10; + let parent_tree = ProcessTree::root(pid); + let child_tree = ProcessTree::child(pid + 2, pid + 1, &parent_tree); + let parent_session = format!("closed-parent-{}", kind.label()); + let child_session = format!("after-close-child-{}", kind.label()); + let prompt = format!("already completed prompt for {}", kind.label()); + let ts = 1_700_700_000_000 + index as i64 * 100; + forward_all( + &mut fixtures.start_turn(kind, &parent_session, &parent_tree, "run once", ts), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + kind, + &parent_session, + &parent_tree, + "closed-call", + &prompt, + ts + 10, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + kind, + &parent_session, + &parent_tree, + "closed-call", + &prompt, + "done", + ts + 20, + ), + &socket, + &host, + ) + .await; + forward_all( + &mut fixtures.start_turn(AgentKind::Pi, &child_session, &child_tree, &prompt, ts + 30), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + AgentKind::Pi, + &child_session, + &child_tree, + "standalone", + ts + 40, + ), + &socket, + &host, + ) + .await; + flush(&child_session, &socket).await; + let child = recording.session(&child_session); + assert!(child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn child_output_can_disambiguate_when_inputs_do_not_overlap() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(17_000); + let child_tree = ProcessTree::child(17_002, 17_001, &parent_tree); + let selected_output = "OUTPUT_FINGERPRINT_BETA_928374"; + + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "output-parent", + &parent_tree, + "run opaque file tasks", + 1_700_800_000_000, + ), + &socket, + &host, + ) + .await; + for (call, opaque_input) in [ + ("output-alpha", "agent --prompt-file /tmp/opaque-alpha"), + ("output-beta", "agent --prompt-file /tmp/opaque-beta"), + ] { + forward( + fixtures.open_tool( + AgentKind::Claude, + "output-parent", + &parent_tree, + call, + opaque_input, + 1_700_800_000_010, + ), + &socket, + &host, + ) + .await; + } + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "output-child", + &child_tree, + "perform the task loaded from the file", + 1_700_800_000_020, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_session( + AgentKind::Claude, + "output-child", + &child_tree, + selected_output, + 1_700_800_000_030, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + AgentKind::Claude, + "output-parent", + &parent_tree, + "output-alpha", + "agent --prompt-file /tmp/opaque-alpha", + "OUTPUT_FINGERPRINT_ALPHA_193746", + 1_700_800_000_040, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + AgentKind::Claude, + "output-parent", + &parent_tree, + "output-beta", + "agent --prompt-file /tmp/opaque-beta", + selected_output, + 1_700_800_000_050, + ), + &socket, + &host, + ) + .await; + flush("output-child", &socket).await; + + let parent = recording.session("output-parent"); + let beta = inserted_rows(&parent) + .into_iter() + .find(|row| { + row.span_type == SpanType::Tool + && row + .input + .as_ref() + .is_some_and(|input| input.to_string().contains("opaque-beta")) + }) + .expect("beta tool"); + let child = recording.session("output-child"); + let child_root = session_root( + &inserted_rows(&child), + "output-child", + "output fingerprint child", + ) + .clone(); + assert_eq!(child_root.parent_span_ids, vec![beta.span_id]); + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn concurrent_agent_sessions_in_one_process_choose_their_own_tools() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let shared_parent_tree = ProcessTree::root(18_000); + let mut cases = Vec::new(); + + for (index, parent_kind) in AgentKind::ALL.into_iter().enumerate() { + let child_kind = AgentKind::ALL[(index + 1) % AgentKind::ALL.len()]; + let parent_session = format!("multiplex-parent-{}", parent_kind.label()); + let child_session = format!("multiplex-child-{}", child_kind.label()); + let call_id = format!("multiplex-call-{index}"); + let prompt = format!( + "unique multiplexed prompt number {index} for {}", + parent_kind.label() + ); + let ts = 1_700_900_000_000 + index as i64 * 100; + forward_all( + &mut fixtures.start_turn( + parent_kind, + &parent_session, + &shared_parent_tree, + "serve one multiplexed child", + ts, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + parent_kind, + &parent_session, + &shared_parent_tree, + &call_id, + &prompt, + ts + 10, + ), + &socket, + &host, + ) + .await; + let child_tree = ProcessTree::child( + 18_100 + index as u32 * 2, + 18_101 + index as u32 * 2, + &shared_parent_tree, + ); + let events = fixtures.start_turn(child_kind, &child_session, &child_tree, &prompt, ts + 20); + cases.push(( + parent_kind, + child_kind, + parent_session, + child_session, + call_id, + prompt, + child_tree, + ts, + events, + )); + } + + for (_, _, parent_session, _, _, _, _, _, _) in &cases { + let parent = recording.session(parent_session); + assert!( + parent + .configs + .lock() + .unwrap() + .iter() + .all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + )), + "top-level session {parent_session} was mistaken for a child" + ); + } + + let mut starts = tokio::task::JoinSet::new(); + for (_, _, _, _, _, _, _, _, events) in &cases { + let events = events.clone(); + let socket = socket.clone(); + let host = host.clone(); + starts.spawn(async move { + for event in events { + forward(event, &socket, &host).await; + } + }); + } + while let Some(result) = starts.join_next().await { + result.unwrap(); + } + + for ( + parent_kind, + child_kind, + parent_session, + child_session, + call_id, + prompt, + child_tree, + ts, + _, + ) in &cases + { + forward( + fixtures.close_session( + *child_kind, + child_session, + child_tree, + "multiplexed child complete", + *ts + 30, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.close_tool( + *parent_kind, + parent_session, + &shared_parent_tree, + call_id, + prompt, + "multiplexed tool complete", + *ts + 40, + ), + &socket, + &host, + ) + .await; + flush(child_session, &socket).await; + assert_pair_linked( + recording.as_ref(), + parent_session, + child_session, + &format!("concurrent shared-process parent {}", parent_kind.label()), + ); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn pid_reuse_and_unavailable_start_times_fail_safe() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(19_000); + let prompt = "PID_REUSE_MUST_NOT_LINK_THIS_CHILD"; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "pid-parent", + &parent_tree, + "delegate safely", + 1_701_000_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "pid-parent", + &parent_tree, + "pid-call", + prompt, + 1_701_000_000_010, + ), + &socket, + &host, + ) + .await; + + for (child_session, mode) in [("pid-reused-child", 0), ("unknown-start-child", 1)] { + let child_tree = ProcessTree::child(19_002 + mode, 19_001, &parent_tree); + let mut events = fixtures.start_turn( + AgentKind::Pi, + child_session, + &child_tree, + prompt, + 1_701_000_000_020 + i64::from(mode), + ); + for event in &mut events { + let capture = event.capture.as_mut().unwrap(); + if mode == 0 { + for process in &mut capture.process_chain { + if process.pid == parent_tree.agent.pid { + process.start_time_secs += 1; + } + } + } else { + for process in &mut capture.process_chain { + process.start_time_secs = 0; + } + } + } + forward_all(&mut events, &socket, &host).await; + forward( + fixtures.close_session( + AgentKind::Pi, + child_session, + &child_tree, + "standalone", + 1_701_000_000_030 + i64::from(mode), + ), + &socket, + &host, + ) + .await; + flush(child_session, &socket).await; + let child = recording.session(child_session); + assert!(child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn open_tool_prevents_idle_retirement_during_a_long_child_spawn() { + let tmp = tempfile::tempdir().unwrap(); + let socket = test_endpoint(tmp.path()); + let data_dir = tmp.path().join("data"); + let recording = Arc::new(RecordingSinkFactory::default()); + let daemon = + spawn_daemon_with_timeouts(socket.clone(), data_dir, recording.clone(), 1, 1).await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(20_000); + let child_tree = ProcessTree::child(20_002, 20_001, &parent_tree); + let prompt = "LONG_RUNNING_CHILD_AFTER_IDLE_WINDOW"; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "idle-parent", + &parent_tree, + "start long child", + 1_701_100_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + "idle-parent", + &parent_tree, + "idle-call", + prompt, + 1_701_100_000_010, + ), + &socket, + &host, + ) + .await; + + tokio::time::sleep(Duration::from_millis(1_300)).await; + assert!(run_status(StatusArgs { + socket: Some(socket.clone()), + session_id: None, + }) + .await + .unwrap() + .is_some()); + forward_all( + &mut fixtures.start_turn( + AgentKind::Pi, + "idle-child", + &child_tree, + prompt, + 1_701_100_001_500, + ), + &socket, + &host, + ) + .await; + flush("idle-child", &socket).await; + assert_pair_linked( + recording.as_ref(), + "idle-parent", + "idle-child", + "active tool idle protection", + ); + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn concurrent_ambiguous_children_do_not_deadlock_or_cross_link() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(21_000); + let prompt = "identical concurrent ambiguous child request"; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + "deadlock-parent", + &parent_tree, + "spawn ambiguous children", + 1_701_200_000_000, + ), + &socket, + &host, + ) + .await; + for call in ["deadlock-a", "deadlock-b"] { + forward( + fixtures.open_tool( + AgentKind::Claude, + "deadlock-parent", + &parent_tree, + call, + prompt, + 1_701_200_000_010, + ), + &socket, + &host, + ) + .await; + } + + let child_cases = [ + ( + AgentKind::Pi, + "deadlock-child-pi", + ProcessTree::child(21_002, 21_001, &parent_tree), + ), + ( + AgentKind::OpenCode, + "deadlock-child-opencode", + ProcessTree::child(21_004, 21_003, &parent_tree), + ), + ]; + let mut starts = tokio::task::JoinSet::new(); + for (kind, session, tree) in &child_cases { + let events = fixtures.start_turn(*kind, session, tree, prompt, 1_701_200_000_020); + let socket = socket.clone(); + let host = host.clone(); + starts.spawn(async move { + for event in events { + forward(event, &socket, &host).await; + } + }); + } + tokio::time::timeout(Duration::from_secs(3), async { + while let Some(result) = starts.join_next().await { + result.unwrap(); + } + }) + .await + .expect("concurrent ambiguous child starts deadlocked"); + + for (kind, session, tree) in &child_cases { + forward( + fixtures.close_session(*kind, session, tree, "same child output", 1_701_200_000_030), + &socket, + &host, + ) + .await; + } + for call in ["deadlock-a", "deadlock-b"] { + forward( + fixtures.close_tool( + AgentKind::Claude, + "deadlock-parent", + &parent_tree, + call, + prompt, + "same parent output", + 1_701_200_000_040, + ), + &socket, + &host, + ) + .await; + } + for (_, session, _) in &child_cases { + flush(session, &socket).await; + let child = recording.session(session); + assert!(child.configs.lock().unwrap().iter().all(|config| !matches!( + config.destination, + Some(TraceDestination::ParentSpan { .. }) + ))); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn one_tool_can_fan_out_to_concurrent_children_of_every_agent() { + let (socket, daemon, recording, _tmp) = start_daemon().await; + let host = HostInfo { + serve_argv: vec![OsString::from("unused")], + version: "test".into(), + }; + let mut fixtures = DistributedFixtures::new(); + let parent_tree = ProcessTree::root(22_000); + let parent_session = "fanout-parent"; + let prompt = "fan out one parent operation to every child agent"; + forward_all( + &mut fixtures.start_turn( + AgentKind::Claude, + parent_session, + &parent_tree, + "fan out", + 1_701_300_000_000, + ), + &socket, + &host, + ) + .await; + forward( + fixtures.open_tool( + AgentKind::Claude, + parent_session, + &parent_tree, + "fanout-call", + prompt, + 1_701_300_000_010, + ), + &socket, + &host, + ) + .await; + + let mut children = Vec::new(); + let mut starts = tokio::task::JoinSet::new(); + for (index, kind) in AgentKind::ALL.into_iter().enumerate() { + let session = format!("fanout-child-{}", kind.label()); + let tree = ProcessTree::child( + 22_100 + index as u32 * 2, + 22_101 + index as u32 * 2, + &parent_tree, + ); + let events = fixtures.start_turn(kind, &session, &tree, prompt, 1_701_300_000_020); + let socket = socket.clone(); + let host = host.clone(); + starts.spawn(async move { + for event in events { + forward(event, &socket, &host).await; + } + }); + children.push((kind, session, tree)); + } + while let Some(result) = starts.join_next().await { + result.unwrap(); + } + + for (kind, session, tree) in &children { + forward( + fixtures.close_session(*kind, session, tree, "fanout complete", 1_701_300_000_030), + &socket, + &host, + ) + .await; + flush(session, &socket).await; + assert_pair_linked( + recording.as_ref(), + parent_session, + session, + &format!("fanout child {}", kind.label()), + ); + } + + shutdown_daemon(&socket).await.unwrap(); + daemon.await.unwrap(); +} + +fn assert_pair_linked( + recording: &RecordingSinkFactory, + parent_session: &str, + child_session: &str, + label: &str, +) { + let parent = recording.session(parent_session); + let child = recording.session(child_session); + assert!(!parent.source.is_empty(), "{label}: parent source absent"); + assert!(!child.source.is_empty(), "{label}: child source absent"); + + let parent_rows = inserted_rows(&parent); + let child_rows = inserted_rows(&child); + let parent_root = session_root(&parent_rows, parent_session, label); + let parent_tool = parent_rows + .iter() + .find(|row| row.span_type == SpanType::Tool) + .unwrap_or_else(|| panic!("{label}: parent emitted no tool span")); + let child_root = session_root(&child_rows, child_session, label); + + assert_eq!( + child_root.parent_span_ids, + vec![parent_tool.span_id.clone()], + "{label}: child root was not parented to the spawning tool" + ); + + let configs = child.configs.lock().unwrap(); + let attached = configs.iter().find_map(|config| match &config.destination { + Some(TraceDestination::ParentSpan { components }) => Some(components), + _ => None, + }); + let attached = attached.unwrap_or_else(|| panic!("{label}: child sink was never attached")); + assert_eq!( + attached.span_id.as_deref(), + Some(parent_tool.span_id.as_str()), + "{label}: sink attachment chose a different parent" + ); + assert_eq!( + attached.root_span_id.as_deref(), + Some(parent_root.root_span_id.as_str()), + "{label}: sink attachment chose a different trace root" + ); +} + +fn inserted_rows(record: &RecordedSession) -> Vec { + record + .ops + .lock() + .unwrap() + .iter() + .filter_map(|op| match op { + SpanOp::Insert(row) => Some(row.clone()), + SpanOp::Merge(_) => None, + }) + .collect() +} + +fn session_root<'a>(rows: &'a [SpanRow], session: &str, label: &str) -> &'a SpanRow { + rows.iter() + .find(|row| { + row.span_type == SpanType::Task + && row + .metadata + .as_ref() + .and_then(|metadata| metadata.get("session_id")) + .and_then(serde_json::Value::as_str) + == Some(session) + }) + .unwrap_or_else(|| panic!("{label}: no root span found for {session}")) +} + +async fn forward_all(events: &mut [Envelope], socket: &Path, host: &HostInfo) { + for event in events.iter().cloned() { + forward(event, socket, host).await; + } +} + +async fn forward(event: Envelope, socket: &Path, host: &HostInfo) { + forward_envelope(&event, socket, host, false) + .await + .unwrap_or_else(|error| { + panic!( + "forward {} {} event {}: {error}", + event.source, event.session_id, event.event + ) + }); +} + +async fn flush(session: &str, socket: &Path) { + let result = flush_session(session, socket, 5_000) + .await + .unwrap_or_else(|error| panic!("flush {session}: {error}")); + assert!( + result.flushed && result.pending == 0, + "session {session} did not fully flush: {result:?}" + ); +} + +async fn start_daemon() -> ( + PathBuf, + tokio::task::JoinHandle<()>, + Arc, + tempfile::TempDir, +) { + let tmp = tempfile::tempdir().unwrap(); + let data_dir = tmp.path().join("data"); + let socket = test_endpoint(tmp.path()); + let recording = Arc::new(RecordingSinkFactory::default()); + let daemon = spawn_daemon(socket.clone(), data_dir, recording.clone()).await; + (socket, daemon, recording, tmp) +} + +async fn spawn_daemon( + socket: PathBuf, + data_dir: PathBuf, + recording: Arc, +) -> tokio::task::JoinHandle<()> { + spawn_daemon_with_timeouts(socket, data_dir, recording, 0, 0).await +} + +async fn spawn_daemon_with_timeouts( + socket: PathBuf, + data_dir: PathBuf, + recording: Arc, + idle_timeout_secs: u64, + session_idle_timeout_secs: u64, +) -> tokio::task::JoinHandle<()> { + let opts = ServeOptions { + version: "test".into(), + translators: Arc::new(Registry::default_agents()), + sink_factory: recording.clone(), + auth_provider: Some(Arc::new(TestAuthProvider)), + }; + let args = ServeArgs { + socket: Some(socket.clone()), + data_dir: Some(data_dir), + idle_timeout_secs, + session_idle_timeout_secs, + }; + let daemon = tokio::spawn(async move { + let _ = run_serve(args, opts).await; + }); + for _ in 0..200 { + if matches!( + run_status(StatusArgs { + socket: Some(socket.clone()), + session_id: None, + }) + .await, + Ok(Some(_)) + ) { + return daemon; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + panic!("distributed tracing test daemon did not start"); +} + +fn test_endpoint(root: &Path) -> PathBuf { + #[cfg(unix)] + { + root.join("distributed.sock") + } + #[cfg(windows)] + { + let _ = root; + PathBuf::from(format!( + r"\\.\pipe\bt-daemon-distributed-test-{}", + uuid::Uuid::new_v4() + )) + } +} diff --git a/bt-daemon/tests/opencode_translator.rs b/bt-daemon/tests/opencode_translator.rs index 0d4bc69..4960075 100644 --- a/bt-daemon/tests/opencode_translator.rs +++ b/bt-daemon/tests/opencode_translator.rs @@ -15,6 +15,7 @@ fn event(name: &str, ts_ms: i64, payload: serde_json::Value) -> Envelope { payload, route: None, config: None, + capture: None, } } diff --git a/bt-daemon/tests/pi_translator.rs b/bt-daemon/tests/pi_translator.rs index 1d089f8..46705e4 100644 --- a/bt-daemon/tests/pi_translator.rs +++ b/bt-daemon/tests/pi_translator.rs @@ -15,6 +15,7 @@ fn event(name: &str, ts_ms: i64, native: serde_json::Value) -> Envelope { payload: json!({"event":native,"extension_version":"1.0.0","cwd":"."}), route: None, config: None, + capture: None, } } @@ -27,9 +28,15 @@ fn reduce(ops: Vec) -> HashMap { } SpanOp::Merge(update) => { let row: &mut SpanRow = rows.entry(update.span_id.clone()).or_default(); + if !update.name.is_empty() { + row.name = update.name; + } if update.end_ms.is_some() { row.end_ms = update.end_ms; } + if update.output.is_some() { + row.output = update.output; + } if update.metadata.is_some() { row.metadata = update.metadata; } diff --git a/bt-daemon/tests/pipeline.rs b/bt-daemon/tests/pipeline.rs index 4298dc3..16a1063 100644 --- a/bt-daemon/tests/pipeline.rs +++ b/bt-daemon/tests/pipeline.rs @@ -143,6 +143,7 @@ fn envelope(session_id: &str, event: &str, ts_ms: i64) -> Envelope { ..SessionRoute::default() }), config: None, + capture: None, } } diff --git a/bt-daemon/tests/support/distributed.rs b/bt-daemon/tests/support/distributed.rs new file mode 100644 index 0000000..23045df --- /dev/null +++ b/bt-daemon/tests/support/distributed.rs @@ -0,0 +1,489 @@ +use bt_daemon::wire::{CaptureContext, Envelope, ProcessIdentity, SessionRoute, TraceDestination}; +use serde_json::{json, Value}; +use std::collections::HashMap; +use std::fs::OpenOptions; +use std::io::Write; +use std::path::PathBuf; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum AgentKind { + Claude, + Codex, + Pi, + OpenCode, +} + +impl AgentKind { + pub const ALL: [Self; 4] = [Self::Claude, Self::Codex, Self::Pi, Self::OpenCode]; + + pub fn label(self) -> &'static str { + match self { + Self::Claude => "claude", + Self::Codex => "codex", + Self::Pi => "pi", + Self::OpenCode => "opencode", + } + } + + fn source(self) -> &'static str { + match self { + Self::Claude => "claude-code", + Self::Codex => "codex", + Self::Pi => "pi", + Self::OpenCode => "opencode", + } + } + + fn uses_command_hook(self) -> bool { + matches!(self, Self::Claude | Self::Codex) + } +} + +#[derive(Clone, Debug)] +pub struct ProcessTree { + pub agent: ProcessIdentity, + pub ancestors: Vec, +} + +impl ProcessTree { + pub fn root(pid: u32) -> Self { + Self { + agent: process(pid), + ancestors: Vec::new(), + } + } + + pub fn child(pid: u32, shell_pid: u32, parent: &Self) -> Self { + let mut ancestors = vec![process(shell_pid), parent.agent.clone()]; + ancestors.extend(parent.ancestors.iter().cloned()); + Self { + agent: process(pid), + ancestors, + } + } + + fn capture(&self, kind: AgentKind, hook_pid: u32) -> CaptureContext { + let mut process_chain = Vec::new(); + if kind.uses_command_hook() { + process_chain.push(process(hook_pid)); + } + process_chain.push(self.agent.clone()); + process_chain.extend(self.ancestors.iter().cloned()); + CaptureContext { process_chain } + } +} + +fn process(pid: u32) -> ProcessIdentity { + ProcessIdentity { + pid, + // Test PIDs are synthetic, so use a stable nonzero boot-relative value + // to exercise the same PID-reuse-safe identity used by live capture. + start_time_secs: 1_700_000_000 + u64::from(pid), + } +} + +/// Produces compact but agent-native lifecycle streams. Codex is deliberately +/// backed by a growing rollout mirror because its live translator learns tool +/// lifecycle from transcript records rather than hook payloads alone. +pub struct DistributedFixtures { + root: tempfile::TempDir, + codex_transcripts: HashMap, + next_hook_pid: u32, +} + +struct CodexTranscript { + path: PathBuf, + bytes: u64, +} + +impl DistributedFixtures { + pub fn new() -> Self { + Self { + root: tempfile::tempdir().expect("create distributed fixture directory"), + codex_transcripts: HashMap::new(), + next_hook_pid: 50_000, + } + } + + pub fn start_turn( + &mut self, + kind: AgentKind, + session: &str, + tree: &ProcessTree, + prompt: &str, + ts_ms: i64, + ) -> Vec { + match kind { + AgentKind::Claude => vec![ + self.envelope( + kind, + session, + tree, + "SessionStart", + ts_ms, + json!({ + "session_id": session, + "hook_event_name": "SessionStart", + "cwd": "/tmp/distributed", + "source": "startup" + }), + ), + self.envelope( + kind, + session, + tree, + "UserPromptSubmit", + ts_ms + 1, + json!({ + "session_id": session, + "hook_event_name": "UserPromptSubmit", + "cwd": "/tmp/distributed", + "prompt": prompt + }), + ), + ], + AgentKind::Codex => { + let records = [ + json!({"timestamp":iso(ts_ms),"type":"session_meta","payload":{"id":session,"cwd":"/tmp/distributed","cli_version":"test"}}), + json!({"timestamp":iso(ts_ms + 1),"type":"turn_context","payload":{"model":"mock-model"}}), + json!({"timestamp":iso(ts_ms + 2),"type":"event_msg","payload":{"type":"task_started","turn_id":"turn-1"}}), + json!({"timestamp":iso(ts_ms + 3),"type":"event_msg","payload":{"type":"user_message","message":prompt}}), + ]; + self.append_codex(session, &records); + let payload = self.codex_payload(session, "SessionStart", Value::Null); + vec![self.envelope(kind, session, tree, "SessionStart", ts_ms + 3, payload)] + } + AgentKind::Pi => vec![ + self.envelope( + kind, + session, + tree, + "session_start", + ts_ms, + pi_payload(session, json!({"reason":"new"})), + ), + self.envelope( + kind, + session, + tree, + "before_agent_start", + ts_ms + 1, + pi_payload(session, json!({"prompt":prompt})), + ), + ], + AgentKind::OpenCode => vec![ + self.envelope( + kind, + session, + tree, + "session.created", + ts_ms, + json!({"properties":{"info":{"id":session,"title":"Distributed test"}}}), + ), + self.envelope( + kind, + session, + tree, + "chat.message", + ts_ms + 1, + json!({ + "input":{"sessionID":session,"model":{"modelID":"mock-model"}}, + "output":{"parts":[{"type":"text","text":prompt}]} + }), + ), + ], + } + } + + pub fn open_tool( + &mut self, + kind: AgentKind, + session: &str, + tree: &ProcessTree, + call_id: &str, + delegated_prompt: &str, + ts_ms: i64, + ) -> Envelope { + let command = format!("agent --prompt {delegated_prompt}"); + match kind { + AgentKind::Claude => self.envelope( + kind, + session, + tree, + "PreToolUse", + ts_ms, + json!({ + "session_id":session, + "hook_event_name":"PreToolUse", + "cwd":"/tmp/distributed", + "tool_name":"Bash", + "tool_use_id":call_id, + "tool_input":{"command":command} + }), + ), + AgentKind::Codex => { + self.append_codex( + session, + &[json!({ + "timestamp":iso(ts_ms), + "type":"response_item", + "payload":{"type":"function_call","call_id":call_id,"name":"exec_command","arguments":json!({"cmd":command}).to_string()} + })], + ); + let payload = self.codex_payload( + session, + "PreToolUse", + json!({ + "tool_name":"exec_command", + "tool_use_id":call_id, + "tool_input":{"cmd":command} + }), + ); + self.envelope(kind, session, tree, "PreToolUse", ts_ms, payload) + } + AgentKind::Pi => self.envelope( + kind, + session, + tree, + "tool_execution_start", + ts_ms, + pi_payload( + session, + json!({"toolCallId":call_id,"toolName":"bash","args":{"command":command}}), + ), + ), + AgentKind::OpenCode => self.envelope( + kind, + session, + tree, + "tool.execute.before", + ts_ms, + json!({ + "input":{"sessionID":session,"callID":call_id,"tool":"bash"}, + "output":{"args":{"command":command}} + }), + ), + } + } + + #[allow(clippy::too_many_arguments)] + pub fn close_tool( + &mut self, + kind: AgentKind, + session: &str, + tree: &ProcessTree, + call_id: &str, + delegated_prompt: &str, + output: &str, + ts_ms: i64, + ) -> Envelope { + let command = format!("agent --prompt {delegated_prompt}"); + match kind { + AgentKind::Claude => self.envelope( + kind, + session, + tree, + "PostToolUse", + ts_ms, + json!({ + "session_id":session, + "hook_event_name":"PostToolUse", + "cwd":"/tmp/distributed", + "tool_name":"Bash", + "tool_use_id":call_id, + "tool_input":{"command":command}, + "tool_response":{"stdout":output,"stderr":"","interrupted":false} + }), + ), + AgentKind::Codex => { + self.append_codex( + session, + &[json!({ + "timestamp":iso(ts_ms), + "type":"response_item", + "payload":{"type":"function_call_output","call_id":call_id,"output":output} + })], + ); + let payload = self.codex_payload( + session, + "PostToolUse", + json!({ + "tool_name":"exec_command", + "tool_use_id":call_id, + "tool_input":{"cmd":command}, + "tool_response":{"output":output} + }), + ); + self.envelope(kind, session, tree, "PostToolUse", ts_ms, payload) + } + AgentKind::Pi => self.envelope( + kind, + session, + tree, + "tool_execution_end", + ts_ms, + pi_payload( + session, + json!({"toolCallId":call_id,"toolName":"bash","result":output,"isError":false}), + ), + ), + AgentKind::OpenCode => self.envelope( + kind, + session, + tree, + "tool.execute.after", + ts_ms, + json!({ + "input":{"sessionID":session,"callID":call_id,"tool":"bash"}, + "result":{"title":"Bash","output":output} + }), + ), + } + } + + pub fn close_session( + &mut self, + kind: AgentKind, + session: &str, + tree: &ProcessTree, + output: &str, + ts_ms: i64, + ) -> Envelope { + match kind { + AgentKind::Claude => self.envelope( + kind, + session, + tree, + "Stop", + ts_ms, + json!({ + "session_id":session, + "hook_event_name":"Stop", + "cwd":"/tmp/distributed", + "last_assistant_message":output + }), + ), + AgentKind::Codex => { + self.append_codex( + session, + &[json!({ + "timestamp":iso(ts_ms), + "type":"event_msg", + "payload":{"type":"task_complete","last_agent_message":output} + })], + ); + let payload = self.codex_payload(session, "Stop", Value::Null); + self.envelope(kind, session, tree, "Stop", ts_ms, payload) + } + AgentKind::Pi => self.envelope( + kind, + session, + tree, + "agent_end", + ts_ms, + pi_payload(session, json!({"messages":[]})), + ), + AgentKind::OpenCode => self.envelope( + kind, + session, + tree, + "session.deleted", + ts_ms, + json!({"properties":{"sessionID":session}}), + ), + } + } + + fn envelope( + &mut self, + kind: AgentKind, + session: &str, + tree: &ProcessTree, + event: &str, + ts_ms: i64, + payload: Value, + ) -> Envelope { + let hook_pid = self.next_hook_pid; + self.next_hook_pid += 1; + Envelope { + source: kind.source().into(), + source_version: Some("integration-test".into()), + plugin_version: Some("integration-test".into()), + session_id: session.into(), + event: event.into(), + ts_ms, + managed_run_id: None, + capture: Some(tree.capture(kind, hook_pid)), + payload, + route: Some(test_route()), + config: None, + } + } + + fn append_codex(&mut self, session: &str, records: &[Value]) { + let transcript = self + .codex_transcripts + .entry(session.into()) + .or_insert_with(|| CodexTranscript { + path: self.root.path().join(format!("{session}.jsonl")), + bytes: 0, + }); + let mut file = OpenOptions::new() + .create(true) + .append(true) + .open(&transcript.path) + .expect("open Codex fixture transcript"); + for record in records { + let line = serde_json::to_string(record).expect("serialize Codex fixture record"); + writeln!(file, "{line}").expect("append Codex fixture record"); + transcript.bytes += line.len() as u64 + 1; + } + } + + fn codex_payload(&self, session: &str, event: &str, extra: Value) -> Value { + let transcript = self + .codex_transcripts + .get(session) + .expect("Codex fixture transcript exists"); + let path = transcript.path.to_string_lossy(); + let mut payload = json!({ + "session_id":session, + "hook_event_name":event, + "transcript_path":path, + "_bt_transcript_mirror":{ + "path":path, + "mirror":path, + "through":transcript.bytes + } + }); + if let (Value::Object(payload), Value::Object(extra)) = (&mut payload, extra) { + payload.extend(extra); + } + payload + } +} + +fn pi_payload(session: &str, event: Value) -> Value { + json!({ + "event":event, + "extension_version":"integration-test", + "native_session_id":session, + "cwd":"/tmp/distributed" + }) +} + +fn test_route() -> SessionRoute { + SessionRoute { + destination: Some(TraceDestination::ProjectLogs { + project_id: None, + project_name: Some("distributed-tracing-test".into()), + }), + ..SessionRoute::default() + } +} + +fn iso(ts_ms: i64) -> String { + chrono::DateTime::from_timestamp_millis(ts_ms) + .expect("fixture timestamp") + .to_rfc3339() +} diff --git a/bt-daemon/tests/support/mod.rs b/bt-daemon/tests/support/mod.rs index 12d793e..e154973 100644 --- a/bt-daemon/tests/support/mod.rs +++ b/bt-daemon/tests/support/mod.rs @@ -2,6 +2,7 @@ pub mod agent_process; pub mod agents; +pub mod distributed; pub mod inference; pub mod ingest; pub mod server;