diff --git a/bt-daemon/src/translate/antigravity.rs b/bt-daemon/src/translate/antigravity.rs index 2b330c4..747f455 100644 --- a/bt-daemon/src/translate/antigravity.rs +++ b/bt-daemon/src/translate/antigravity.rs @@ -60,6 +60,7 @@ struct AntigravityTranslator { session_id: String, session_span_id: String, root_span_id: String, + session_parent_span_ids: Vec, root_open: bool, root_ended: bool, turn: Option, @@ -81,6 +82,7 @@ impl AntigravityTranslator { session_id: session_id.to_string(), session_span_id: root.clone(), root_span_id: root, + session_parent_span_ids: Vec::new(), root_open: false, root_ended: false, turn: None, @@ -109,6 +111,7 @@ impl AntigravityTranslator { if let Some(external_root) = external_root_span_id { self.root_span_id = external_root; } + self.session_parent_span_ids = parent_span_id.into_iter().collect(); let workspace = event .payload @@ -147,7 +150,7 @@ impl AntigravityTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: parent_span_id.into_iter().collect(), + parent_span_ids: self.session_parent_span_ids.clone(), name: format!("Antigravity: {label}"), span_type: SpanType::Task, start_ms: Some(event.ts_ms), @@ -467,6 +470,7 @@ impl AntigravityTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], end_ms: Some(ts_ms), output: turn.last_output, metadata: Some(json!({"turn_number":turn.number})), @@ -484,6 +488,7 @@ impl AntigravityTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.session_parent_span_ids.clone(), end_ms: Some(ts_ms), error, ..Default::default() diff --git a/bt-daemon/src/translate/claude.rs b/bt-daemon/src/translate/claude.rs index ad5e58f..8b98362 100644 --- a/bt-daemon/src/translate/claude.rs +++ b/bt-daemon/src/translate/claude.rs @@ -52,6 +52,7 @@ struct TranscriptCursor { struct Subagent { span_id: String, + parent_span_id: String, transcript_path: Option, } @@ -78,6 +79,7 @@ struct ClaudeTranslator { session_id: String, session_span_id: String, root_span_id: String, + session_parent_span_ids: Vec, root_open: bool, root_ended: bool, turn: Option, @@ -108,6 +110,7 @@ impl ClaudeTranslator { session_id: session_id.to_string(), session_span_id: root.clone(), root_span_id: root, + session_parent_span_ids: Vec::new(), root_open: false, root_ended: false, turn: None, @@ -145,6 +148,7 @@ impl ClaudeTranslator { if let Some(root) = root_span_id { self.root_span_id = root.clone(); } + self.session_parent_span_ids = parent_span_id.into_iter().collect(); let cwd = string_field(&event.payload, "cwd").unwrap_or_default(); let workspace = basename(&cwd); let mut metadata = ctx @@ -180,7 +184,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: parent_span_id.into_iter().collect(), + parent_span_ids: self.session_parent_span_ids.clone(), name: format!("Claude Code: {workspace}"), span_type: SpanType::Task, start_ms: Some(event.ts_ms), @@ -216,6 +220,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: old.id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], end_ms: Some(event.ts_ms), ..Default::default() })); @@ -279,6 +284,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], metadata: explicit_skill_metadata(&self.pending_skills), ..Default::default() })); @@ -324,6 +330,7 @@ impl ClaudeTranslator { agent_id.to_string(), Subagent { span_id: span_id.clone(), + parent_span_id: parent_id, transcript_path: None, }, ); @@ -378,6 +385,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: pending.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![pending.parent_id], end_ms: Some(event.ts_ms), output, metadata: Some(metadata), @@ -424,6 +432,11 @@ impl ClaudeTranslator { return; }; let parent = self.ensure_subagent(&agent_id, event, ops); + let parent_span_ids = self + .subagents + .get(&agent_id) + .map(|agent| vec![agent.parent_span_id.clone()]) + .unwrap_or_else(|| vec![self.session_span_id.clone()]); let path = string_field(&event.payload, "agent_transcript_path"); if let Some(agent) = self.subagents.get_mut(&agent_id) { agent.transcript_path = path.clone(); @@ -453,6 +466,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: parent.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids, end_ms: Some(event.ts_ms), output: event.payload.get("last_assistant_message").cloned(), ..Default::default() @@ -618,6 +632,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], end_ms: Some(event.ts_ms), output: event .payload @@ -652,6 +667,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: tool.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![tool.parent_id], end_ms: Some(end_ms), error: Some(error.to_string()), ..Default::default() @@ -679,6 +695,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], end_ms: Some(event.ts_ms), ..Default::default() })); @@ -689,6 +706,7 @@ impl ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.session_parent_span_ids.clone(), end_ms: Some(event.ts_ms), ..Default::default() })); @@ -720,6 +738,7 @@ impl AgentTranslator for ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.session_parent_span_ids.clone(), metadata: Some(json!({ "claude_code_version": version })), ..Default::default() })); @@ -789,6 +808,7 @@ impl AgentTranslator for ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: tool.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![tool.parent_id], end_ms: Some(end_ms), error: Some("Session ended before tool completion".into()), ..Default::default() @@ -798,6 +818,7 @@ impl AgentTranslator for ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![self.session_span_id.clone()], end_ms: Some(end_ms), error: Some("Session ended before turn completion".into()), ..Default::default() @@ -807,6 +828,7 @@ impl AgentTranslator for ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: subagent.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![subagent.parent_span_id], end_ms: Some(end_ms), error: Some("Session ended before subagent completion".into()), ..Default::default() @@ -817,6 +839,7 @@ impl AgentTranslator for ClaudeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.session_parent_span_ids.clone(), end_ms: Some(end_ms), ..Default::default() })); diff --git a/bt-daemon/src/translate/codex.rs b/bt-daemon/src/translate/codex.rs index f6439fa..cc1873c 100644 --- a/bt-daemon/src/translate/codex.rs +++ b/bt-daemon/src/translate/codex.rs @@ -303,6 +303,20 @@ impl AgentTranslator for CodexTranslator { } impl CodexTranslator { + fn turn_parent_span_ids(&self, turn_id: &str) -> Vec { + vec![ids::span_id(&self.session_id, &format!("turn:{turn_id}"))] + } + + fn scope_root_parent_span_ids(&self, scope: &Scope) -> Vec { + match scope.kind { + ScopeKind::Main => self.external_parent_span_id.clone().into_iter().collect(), + ScopeKind::Subagent => vec![scope + .spawning_turn_span_id + .clone() + .unwrap_or_else(|| self.root_span_id.clone())], + } + } + fn start_catch_up(&mut self, ctx: &SessionCtx, finalize: bool) -> anyhow::Result> { anyhow::ensure!( self.pending.is_none(), @@ -344,9 +358,14 @@ impl CodexTranslator { if self.compaction_spans.remove(&turn_id) { self.compaction_trigger_by_turn.remove(&turn_id); let span_id = ids::span_id(&self.session_id, &format!("turn:{turn_id}")); + let parent_span_ids = str_field(payload, "transcript_path") + .and_then(|path| self.scopes.get(&path)) + .map(|scope| vec![scope.turn_parent_span_id.clone()]) + .unwrap_or_else(|| vec![self.root_span_id.clone()]); ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids, metadata: Some(json!({ "compaction": { "trigger": trigger } })), ..Default::default() })); @@ -506,6 +525,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: scope.turn_parent_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.scope_root_parent_span_ids(scope), input: Some(input), metadata: Some(json!({ "model": m })), ..Default::default() @@ -517,6 +537,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], metadata: Some(json!({ "model": m })), ..Default::default() })); @@ -604,7 +625,7 @@ impl CodexTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: self.root_span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: self.external_parent_span_id.clone().into_iter().collect(), + parent_span_ids: self.scope_root_parent_span_ids(scope), name, span_type: SpanType::Task, start_ms: Some(ts), @@ -681,6 +702,7 @@ impl CodexTranslator { let Some(text) = text else { return }; // Explicit skill mentions in the prompt (e.g. "$skill", "/skills name"). let names = explicit_skill_names(&text); + let turn_parent_span_id = scope.turn_parent_span_id.clone(); if let Some(turn) = scope.open_turns.last_mut() { for n in names { if !turn.explicit_skill_names.contains(&n) { @@ -690,6 +712,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![turn_parent_span_id], input: Some(json!(text)), metadata: explicit_skill_metadata(&turn.explicit_skill_names), ..Default::default() @@ -761,6 +784,7 @@ impl CodexTranslator { } } else if role == "user" { let names = explicit_skill_names(&text); + let turn_parent_span_id = scope.turn_parent_span_id.clone(); if let Some(turn) = scope.open_turns.last_mut() { for name in names { if !turn.explicit_skill_names.contains(&name) { @@ -771,6 +795,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![turn_parent_span_id], metadata: Some(metadata), ..Default::default() })); @@ -928,13 +953,13 @@ impl CodexTranslator { return; }; push_tool_result(scope, Some(&call_id), payload); - let Some((span_id, _turn_id)) = scope.open_tools.remove(&call_id) else { + let Some((span_id, turn_id)) = scope.open_tools.remove(&call_id) else { return; }; if let Some(turn) = scope .open_turns .iter_mut() - .find(|turn| turn.turn_id == _turn_id) + .find(|turn| turn.turn_id == turn_id) { turn.last_child_end_ms = Some(turn.last_child_end_ms.map_or(ts, |p| p.max(ts))); } @@ -946,6 +971,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: self.turn_parent_span_ids(&turn_id), end_ms: Some(ts), output, metadata: Some(tool_approval_metadata(Some(ToolApproval::Approved))), @@ -996,6 +1022,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: self.turn_parent_span_ids(&llm.turn_id), end_ms: Some(end), output, metadata: usage_metadata, @@ -1021,6 +1048,7 @@ impl CodexTranslator { return; }; let turn_id = scope.open_turns[turn_index].turn_id.clone(); + let turn_span_id = scope.open_turns[turn_index].span_id.clone(); if scope .open_llm @@ -1036,6 +1064,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![turn_span_id], end_ms: Some(llm.last_output_ms), output, metadata: Some(json!({ @@ -1053,6 +1082,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], end_ms: Some(ts), output, ..Default::default() @@ -1066,6 +1096,7 @@ impl CodexTranslator { end_ms: Option, ops: &mut Vec, ) { + let parent_span_ids = self.turn_parent_span_ids(turn_id); let call_ids: Vec = scope .open_tools .iter() @@ -1077,6 +1108,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: parent_span_ids.clone(), end_ms, metadata: Some(tool_approval_metadata(Some(ToolApproval::Approved))), error: Some(MISSING_TOOL_OUTPUT_ERROR.to_string()), @@ -1125,6 +1157,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn_span.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], name: "compaction".to_string(), span_type: SpanType::Task, metadata: Some(json!({ "compaction": { @@ -1187,6 +1220,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: self.root_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.external_parent_span_id.clone().into_iter().collect(), end_ms: Some(end_ms), ..Default::default() })); @@ -1216,6 +1250,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: scope.turn_parent_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: self.scope_root_parent_span_ids(&scope), end_ms: Some(end), ..Default::default() })); @@ -1237,6 +1272,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: self.turn_parent_span_ids(&llm.turn_id), end_ms: end_ms.or(Some(llm.last_output_ms)), output, metadata: Some(json!({ @@ -1245,15 +1281,18 @@ impl CodexTranslator { ..Default::default() })); } - let tools: Vec<(String, String)> = scope + let tools: Vec<(String, String, String)> = scope .open_tools .drain() - .map(|(_, (span_id, _))| (span_id, MISSING_TOOL_OUTPUT_ERROR.to_string())) + .map(|(_, (span_id, turn_id))| { + (span_id, turn_id, MISSING_TOOL_OUTPUT_ERROR.to_string()) + }) .collect(); - for (sid, error) in tools { + for (sid, turn_id, error) in tools { ops.push(SpanOp::Merge(SpanRow { span_id: sid, root_span_id: self.root_span_id.clone(), + parent_span_ids: self.turn_parent_span_ids(&turn_id), end_ms, metadata: Some(tool_approval_metadata(Some(ToolApproval::Approved))), error: Some(error), @@ -1264,6 +1303,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], end_ms, ..Default::default() })); diff --git a/bt-daemon/src/translate/opencode.rs b/bt-daemon/src/translate/opencode.rs index 3bb751c..9d21b4c 100644 --- a/bt-daemon/src/translate/opencode.rs +++ b/bt-daemon/src/translate/opencode.rs @@ -41,6 +41,7 @@ impl TranslatorFactory for OpenCodeTranslatorFactory { struct NativeSession { root_span_id: String, effective_root_span_id: String, + parent_span_ids: Vec, parent_session_id: Option, current_turn_span_id: Option, turn_number: u32, @@ -52,6 +53,7 @@ struct NativeSession { reasoning_parts: HashMap, tool_calls: HashMap>, tool_starts: HashMap, + tool_parent_span_ids: HashMap, tool_names: HashMap, tool_args: HashMap, tool_outputs: HashMap, @@ -206,6 +208,7 @@ impl OpenCodeTranslator { NativeSession { root_span_id: root_span_id.clone(), effective_root_span_id: effective_root_span_id.clone(), + parent_span_ids: parent_span_ids.clone(), parent_session_id: parent_id.map(str::to_owned), ..Default::default() }, @@ -234,6 +237,7 @@ impl OpenCodeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn, root_span_id: state.effective_root_span_id.clone(), + parent_span_ids: vec![state.root_span_id.clone()], end_ms: Some(event.ts_ms), output: state.current_output.take().map(Value::String), ..Default::default() @@ -472,6 +476,7 @@ impl OpenCodeTranslator { return vec![]; } s.tool_starts.insert(call.into(), event.ts_ms); + s.tool_parent_span_ids.insert(call.into(), turn.clone()); if let Some(a) = event.payload.pointer("/output/args") { s.tool_args.insert(call.into(), a.clone()); } @@ -522,6 +527,7 @@ impl OpenCodeTranslator { s.tool_outputs.remove(call); s.tool_errors.remove(call); s.tool_starts.remove(call); + s.tool_parent_span_ids.remove(call); return vec![]; } let Some(turn) = s.current_turn_span_id.clone() else { @@ -562,10 +568,11 @@ impl OpenCodeTranslator { .unwrap_or(Value::Null) } let had_start = s.tool_starts.contains_key(call); + let parent_span_id = s.tool_parent_span_ids.remove(call).unwrap_or(turn); 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: (!had_start).then_some(turn).into_iter().collect(), + parent_span_ids: vec![parent_span_id], name, span_type: SpanType::Tool, start_ms: (!had_start).then(|| s.tool_starts.remove(call).unwrap_or(event.ts_ms)), @@ -664,6 +671,7 @@ impl OpenCodeTranslator { }; s.denied_tools.insert(call.clone()); let had_start = s.tool_starts.contains_key(&call); + let parent_span_id = s.tool_parent_span_ids.get(&call).cloned().unwrap_or(turn); let mut metadata = with_tool_approval( json!({"tool_name":tool,"call_id":call}), Some(ToolApproval::Denied), @@ -680,7 +688,7 @@ impl OpenCodeTranslator { 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: (!had_start).then_some(turn).into_iter().collect(), + parent_span_ids: vec![parent_span_id], name: title.unwrap_or(tool), span_type: SpanType::Tool, start_ms: (!had_start).then(|| s.tool_starts.remove(&call).unwrap_or(event.ts_ms)), @@ -725,9 +733,16 @@ impl OpenCodeTranslator { continue; } let tool_name = s.tool_names.remove(&call).unwrap_or_else(|| "tool".into()); + let parent_span_ids = s + .tool_parent_span_ids + .remove(&call) + .or_else(|| s.current_turn_span_id.clone()) + .into_iter() + .collect(); ops.push(SpanOp::Merge(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, end_ms: Some(ts), metadata: Some(with_tool_approval( json!({ @@ -744,6 +759,7 @@ impl OpenCodeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn, root_span_id: s.effective_root_span_id.clone(), + parent_span_ids: vec![s.root_span_id.clone()], end_ms: Some(ts), output: s.current_output.take().map(Value::String), error: error.clone(), @@ -754,6 +770,7 @@ impl OpenCodeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: s.root_span_id, root_span_id: s.effective_root_span_id, + parent_span_ids: s.parent_span_ids, end_ms: Some(ts), metadata: Some( json!({"total_turns":s.turn_number,"total_tool_calls":s.tool_call_count}), diff --git a/bt-daemon/src/translate/pi.rs b/bt-daemon/src/translate/pi.rs index 994ea30..884114b 100644 --- a/bt-daemon/src/translate/pi.rs +++ b/bt-daemon/src/translate/pi.rs @@ -473,11 +473,7 @@ impl PiTranslator { 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: pending - .is_none() - .then(|| turn.clone()) - .into_iter() - .collect(), + parent_span_ids: vec![turn.clone()], name, span_type: SpanType::Tool, start_ms: pending.is_none().then_some(tracked.start_ms), @@ -554,6 +550,7 @@ impl PiTranslator { vec![SpanOp::Merge(SpanRow { span_id: id, root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![self.root_span_id.clone()], end_ms: Some(ts), error, ..Default::default() @@ -566,6 +563,7 @@ impl PiTranslator { SpanOp::Merge(SpanRow { span_id: self.root_span_id.clone(), root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: self.external_parent.clone().into_iter().collect(), end_ms: Some(ts), metadata: Some( json!({"total_turns":self.turn_seq,"total_tool_calls":self.total_tools}), diff --git a/bt-daemon/tests/antigravity_translator.rs b/bt-daemon/tests/antigravity_translator.rs index a0dbe8e..e284f12 100644 --- a/bt-daemon/tests/antigravity_translator.rs +++ b/bt-daemon/tests/antigravity_translator.rs @@ -1,6 +1,11 @@ -use bt_daemon::wire::Envelope; +#[path = "support/span_identity.rs"] +mod span_identity; + +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute, TraceDestination}; use bt_daemon::{Registry, SessionCtx, SpanOp, SpanRow, SpanType}; use serde_json::{json, Value}; +use span_identity::assert_merges_preserve_insert_identity; use std::collections::HashMap; fn jsonl(records: &[Value]) -> (String, Vec) { @@ -59,6 +64,7 @@ fn event( } fn reduce(ops: Vec) -> HashMap { + assert_merges_preserve_insert_identity(&ops); let mut rows: HashMap = HashMap::new(); for op in ops { match op { @@ -109,6 +115,51 @@ fn reduce(ops: Vec) -> HashMap { rows } +#[test] +fn attached_antigravity_root_merge_preserves_external_identity() { + let registry = Registry::default_agents(); + let mut translator = registry.create("antigravity", "conversation-1"); + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); + let ctx = SessionCtx { + session_id: "conversation-1".into(), + config: Some( + SessionRoute { + destination: Some(TraceDestination::ParentSpan { components }), + ..SessionRoute::default() + } + .with_auth(BackendAuth { + token: "test".into(), + api_url: None, + app_url: None, + org_name: None, + org_id: None, + }), + ), + }; + let ops = translator + .handle( + &event( + "Stop", + 1, + "/tmp/test-antigravity/transcript.jsonl", + "", + 0, + json!({"fullyIdle":true}), + ), + &ctx, + ) + .unwrap(); + let rows = reduce(ops); + let root = rows + .values() + .find(|row| row.name.starts_with("Antigravity:")) + .unwrap(); + assert_eq!(root.root_span_id, "external-root"); + assert_eq!(root.parent_span_ids, ["external-parent"]); +} + #[test] fn hooks_and_full_transcript_build_model_and_tool_spans() { let records = vec![ diff --git a/bt-daemon/tests/braintrust_sink.rs b/bt-daemon/tests/braintrust_sink.rs index b1d054d..2858ac5 100644 --- a/bt-daemon/tests/braintrust_sink.rs +++ b/bt-daemon/tests/braintrust_sink.rs @@ -7,7 +7,7 @@ use bt_daemon::wire::{BackendAuth, FlushMode, SessionConfig, TraceDestination}; use bt_daemon::{ BraintrustSinkConfig, BraintrustSinkFactory, SinkFactory, SpanOp, SpanRow, SpanType, }; -use serde_json::json; +use serde_json::{json, Value}; use wiremock::matchers::{method, path}; use wiremock::{Mock, MockServer, ResponseTemplate}; @@ -90,6 +90,22 @@ async fn logs3_bodies(server: &MockServer) -> String { .join("\n") } +async fn logs3_rows(server: &MockServer) -> Vec { + server + .received_requests() + .await + .unwrap() + .iter() + .filter(|request| request.url.path() == "/logs3") + .flat_map(|request| { + serde_json::from_slice::(&request.body).unwrap()["rows"] + .as_array() + .cloned() + .unwrap_or_default() + }) + .collect() +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn explicit_project_id_does_not_register_a_project_name() { let server = MockServer::start().await; @@ -253,18 +269,20 @@ async fn late_merge_updates_a_completed_span_without_an_open_handle() { sink.emit(&[SpanOp::Insert(row( "finished", - "finished", - &[], + "trace-root", + &["turn-parent"], "original name", - SpanType::Task, + SpanType::Tool, 1, Some(2), ))]) .await .unwrap(); + sink.flush().await.unwrap(); let mut late = SpanRow { span_id: "finished".into(), - root_span_id: "finished".into(), + root_span_id: "trace-root".into(), + parent_span_ids: vec!["turn-parent".into()], output: Some(json!({"status":"late"})), ..Default::default() }; @@ -277,7 +295,18 @@ async fn late_merge_updates_a_completed_span_without_an_open_handle() { bodies.contains("original name"), "initial row absent: {bodies}" ); - assert!(bodies.contains("late"), "late merge absent: {bodies}"); + let rows = logs3_rows(&server).await; + let late = rows + .iter() + .find(|row| row.pointer("/output/status") == Some(&json!("late"))) + .unwrap_or_else(|| panic!("late merge absent: {bodies}")); + assert_eq!(late["_is_merge"], true); + assert_eq!(late["root_span_id"], "trace-root"); + assert_eq!( + late["span_parents"], + json!(["turn-parent"]), + "stateless merge did not repeat the child parent identity: {bodies}" + ); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] diff --git a/bt-daemon/tests/claude_translator.rs b/bt-daemon/tests/claude_translator.rs index 1b39eee..ecb8ee7 100644 --- a/bt-daemon/tests/claude_translator.rs +++ b/bt-daemon/tests/claude_translator.rs @@ -1,6 +1,11 @@ -use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute}; +#[path = "support/span_identity.rs"] +mod span_identity; + +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute, TraceDestination}; use bt_daemon::{Registry, SessionCtx, SpanOp, SpanRow, SpanType}; use serde_json::{json, Value}; +use span_identity::assert_merges_preserve_insert_identity; use std::collections::HashMap; use std::path::{Path, PathBuf}; @@ -95,6 +100,7 @@ fn replay_from(name: &str, source: Source) -> Vec { } fn reduce(ops: Vec) -> HashMap { + assert_merges_preserve_insert_identity(&ops); let mut rows = HashMap::::new(); for op in ops { match op { @@ -243,10 +249,14 @@ fn claude_real_fixture_matches_session_turn_tool_and_token_contract() { fn claude_additional_metadata_reaches_roots_without_overriding_session_fields() { let registry = Registry::default_agents(); let mut translator = registry.create("claude-code", "session"); + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); let ctx = SessionCtx { session_id: "session".into(), config: Some( SessionRoute { + destination: Some(TraceDestination::ParentSpan { components }), additional_metadata: Some(json!({"team": "platform", "source": "custom"})), ..SessionRoute::default() } @@ -259,33 +269,32 @@ fn claude_additional_metadata_reaches_roots_without_overriding_session_fields() }), ), }; - let ops = translator - .handle( - &Envelope { - source: "claude-code".into(), - source_version: None, - plugin_version: None, - session_id: "session".into(), - event: "UserPromptSubmit".into(), - ts_ms: 1, - managed_run_id: None, - payload: json!({"session_id":"session","cwd":"/workspace","prompt":"go"}), - route: None, - config: None, - capture: None, - }, - &ctx, - ) + let event = |name: &str, ts_ms: i64| Envelope { + source: "claude-code".into(), + source_version: None, + plugin_version: None, + session_id: "session".into(), + event: name.into(), + ts_ms, + managed_run_id: None, + payload: json!({"session_id":"session","cwd":"/workspace","prompt":"go"}), + route: None, + config: None, + capture: None, + }; + let mut ops = translator + .handle(&event("UserPromptSubmit", 1), &ctx) .unwrap(); - let root = ops - .into_iter() - .find_map(|op| match op { - SpanOp::Insert(row) if row.name.starts_with("Claude Code:") => Some(row), - _ => None, - }) + ops.extend(translator.handle(&event("SessionEnd", 2), &ctx).unwrap()); + let rows = reduce(ops); + let root = rows + .values() + .find(|row| row.name.starts_with("Claude Code:")) .unwrap(); assert_eq!(root.metadata.as_ref().unwrap()["team"], "platform"); assert_eq!(root.metadata.as_ref().unwrap()["source"], "claude-code"); + assert_eq!(root.root_span_id, "external-root"); + assert_eq!(root.parent_span_ids, ["external-parent"]); } #[test] diff --git a/bt-daemon/tests/codex_translator.rs b/bt-daemon/tests/codex_translator.rs index 896cd88..bd4ebe2 100644 --- a/bt-daemon/tests/codex_translator.rs +++ b/bt-daemon/tests/codex_translator.rs @@ -2,9 +2,14 @@ //! hook triggers into a session → turn → {llm, tool} span tree. Mirrors the //! happy-path shape of the TS `event-processor` tests. -use bt_daemon::wire::{BackendAuth, Envelope, FlushMode, SessionConfig}; +#[path = "support/span_identity.rs"] +mod span_identity; + +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use bt_daemon::wire::{BackendAuth, Envelope, FlushMode, SessionConfig, TraceDestination}; use bt_daemon::{Registry, SessionCtx, SpanOp, SpanRow, SpanType}; use serde_json::{json, Value}; +use span_identity::assert_merges_preserve_insert_identity; use std::collections::HashMap; use std::io::Write; @@ -175,6 +180,7 @@ fn codex_happy_path_builds_session_turn_llm_tool_tree() { ); ops.extend(tr.flush(&ctx).unwrap()); + assert_merges_preserve_insert_identity(&ops); let rows = reduce(ops); // Root (session). @@ -251,6 +257,53 @@ fn codex_happy_path_builds_session_turn_llm_tool_tree() { ); } +#[test] +fn attached_codex_root_merge_preserves_external_parent() { + let tmp = tempfile::tempdir().unwrap(); + let transcript = tmp.path().join("rollout.jsonl"); + write_transcript(&transcript); + let path = transcript.to_str().unwrap(); + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); + let ctx = SessionCtx { + session_id: "attached-session".into(), + config: Some(SessionConfig { + auth: BackendAuth { + token: "test".into(), + api_url: None, + app_url: None, + org_name: None, + org_id: None, + }, + destination: Some(TraceDestination::ParentSpan { components }), + flush_mode: FlushMode::FireAndForget, + additional_metadata: None, + }), + }; + let registry = Registry::default_agents(); + let mut translator = registry.create("codex", "attached-session"); + let mut ops = translator + .handle( + &envelope("attached-session", "SessionStart", path, json!({})), + &ctx, + ) + .unwrap(); + ops.extend( + translator + .handle(&envelope("attached-session", "Stop", path, json!({})), &ctx) + .unwrap(), + ); + ops.extend(translator.flush(&ctx).unwrap()); + assert_merges_preserve_insert_identity(&ops); + let rows = reduce(ops); + let root = rows + .values() + .find(|row| row.name.starts_with("codex:")) + .unwrap(); + assert_eq!(root.parent_span_ids, ["external-parent"]); +} + #[test] fn codex_incremental_reads_advance_offset() { // Two reads: the second only sees records appended after the first. @@ -566,6 +619,7 @@ fn late_task_complete_is_correlated_by_turn_id() { .unwrap(), ); + assert_merges_preserve_insert_identity(&ops); let rows = reduce(ops); let t1 = find(&rows, SpanType::Task, "turn: t1"); let t2 = find(&rows, SpanType::Task, "turn: t2"); @@ -860,6 +914,7 @@ fn codex_compaction_relabels_turn_and_adds_compaction_llm() { ) .unwrap(), ); + assert_merges_preserve_insert_identity(&ops); let rows = reduce(ops); let compaction = find(&rows, SpanType::Task, "compaction"); @@ -1039,6 +1094,7 @@ fn codex_subagent_nests_under_spawning_turn() { .unwrap(), ); + assert_merges_preserve_insert_identity(&ops); let rows = reduce(ops); let root = find(&rows, SpanType::Task, "codex: app"); diff --git a/bt-daemon/tests/opencode_translator.rs b/bt-daemon/tests/opencode_translator.rs index a0e3803..b3d8bae 100644 --- a/bt-daemon/tests/opencode_translator.rs +++ b/bt-daemon/tests/opencode_translator.rs @@ -1,6 +1,11 @@ -use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute}; +#[path = "support/span_identity.rs"] +mod span_identity; + +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute, TraceDestination}; use bt_daemon::{Registry, SessionCtx, SpanOp, SpanRow, SpanType}; use serde_json::json; +use span_identity::assert_merges_preserve_insert_identity; use std::collections::HashMap; fn event(name: &str, ts_ms: i64, payload: serde_json::Value) -> Envelope { @@ -20,6 +25,7 @@ fn event(name: &str, ts_ms: i64, payload: serde_json::Value) -> Envelope { } fn reduce(ops: Vec) -> HashMap { + assert_merges_preserve_insert_identity(&ops); let mut rows = HashMap::new(); for op in ops { match op { @@ -140,6 +146,18 @@ fn opencode_child_sessions_share_the_parent_trace_root() { .unwrap(); ops.extend(translator.handle(&event("chat.message", 2, json!({"input":{"sessionID":"parent"},"output":{"parts":[{"type":"text","text":"delegate"}]}})), &ctx).unwrap()); ops.extend(translator.handle(&event("session.created", 3, json!({"properties":{"info":{"id":"child","parentID":"parent","title":"find docs (@research subagent)"}}})), &ctx).unwrap()); + ops.extend( + translator + .handle( + &event( + "session.deleted", + 4, + json!({"properties":{"sessionID":"child"}}), + ), + &ctx, + ) + .unwrap(), + ); let rows = reduce(ops); let parent = rows.values().find(|r| r.name == "OpenCode").unwrap(); let child = rows @@ -154,10 +172,14 @@ fn opencode_child_sessions_share_the_parent_trace_root() { fn opencode_additional_metadata_reaches_roots_without_overriding_session_fields() { let registry = Registry::default_agents(); let mut translator = registry.create("opencode", "root-session"); + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); let ctx = SessionCtx { session_id: "root-session".into(), config: Some( SessionRoute { + destination: Some(TraceDestination::ParentSpan { components }), additional_metadata: Some(json!({"team": "platform", "source": "custom"})), ..SessionRoute::default() } @@ -170,21 +192,44 @@ fn opencode_additional_metadata_reaches_roots_without_overriding_session_fields( }), ), }; - let rows = reduce( + let mut ops = translator + .handle( + &event( + "session.created", + 1, + json!({"properties":{"info":{"id":"native"}}}), + ), + &ctx, + ) + .unwrap(); + ops.extend( translator .handle( &event( - "session.created", - 1, - json!({"properties":{"info":{"id":"native"}}}), + "session.deleted", + 2, + json!({"properties":{"sessionID":"native"}}), ), &ctx, ) .unwrap(), ); + let inserted_root = ops + .iter() + .find_map(|op| match op { + SpanOp::Insert(row) if row.name == "OpenCode" => Some(row), + _ => None, + }) + .unwrap(); + assert_eq!(inserted_root.metadata.as_ref().unwrap()["team"], "platform"); + assert_eq!( + inserted_root.metadata.as_ref().unwrap()["source"], + "opencode" + ); + let rows = reduce(ops); let root = rows.values().next().unwrap(); - assert_eq!(root.metadata.as_ref().unwrap()["team"], "platform"); - assert_eq!(root.metadata.as_ref().unwrap()["source"], "opencode"); + assert_eq!(root.root_span_id, "external-root"); + assert_eq!(root.parent_span_ids, ["external-parent"]); } #[test] @@ -235,6 +280,7 @@ fn opencode_finalization_closes_a_missing_tool_completion() { session_id: "root-session".into(), config: None, }; + let mut ops = Vec::new(); for envelope in [ event("chat.message", 1, json!({"input":{"sessionID":"native"}})), event( @@ -242,13 +288,15 @@ fn opencode_finalization_closes_a_missing_tool_completion() { 2, json!({"input":{"sessionID":"native","callID":"call","tool":"read"},"output":{"args":{"path":"x"}}}), ), + event("chat.message", 3, json!({"input":{"sessionID":"native"}})), ] { - translator.handle(&envelope, &ctx).unwrap(); + ops.extend(translator.handle(&envelope, &ctx).unwrap()); } - let ops = translator.finalize(&ctx).unwrap(); + ops.extend(translator.finalize(&ctx).unwrap()); + assert_merges_preserve_insert_identity(&ops); assert!(ops.iter().any(|op| matches!( op, - SpanOp::Merge(row) if row.end_ms == Some(2) && row.error.is_some() + SpanOp::Merge(row) if row.end_ms == Some(3) && row.error.is_some() ))); } diff --git a/bt-daemon/tests/pi_translator.rs b/bt-daemon/tests/pi_translator.rs index 8df7c87..780636f 100644 --- a/bt-daemon/tests/pi_translator.rs +++ b/bt-daemon/tests/pi_translator.rs @@ -1,6 +1,11 @@ -use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute}; +#[path = "support/span_identity.rs"] +mod span_identity; + +use braintrust_sdk_rust::{SpanComponents, SpanObjectType}; +use bt_daemon::wire::{BackendAuth, Envelope, SessionRoute, TraceDestination}; use bt_daemon::{Registry, SessionCtx, SpanOp, SpanRow, SpanType}; use serde_json::json; +use span_identity::assert_merges_preserve_insert_identity; use std::collections::HashMap; fn event(name: &str, ts_ms: i64, native: serde_json::Value) -> Envelope { @@ -20,6 +25,7 @@ fn event(name: &str, ts_ms: i64, native: serde_json::Value) -> Envelope { } fn reduce(ops: Vec) -> HashMap { + assert_merges_preserve_insert_identity(&ops); let mut rows = HashMap::new(); for op in ops { match op { @@ -146,10 +152,14 @@ fn pi_builds_turn_llm_tool_compaction_and_shutdown_spans() { fn pi_additional_metadata_reaches_roots_without_overriding_session_fields() { let registry = Registry::default_agents(); let mut translator = registry.create("pi", "pi-session"); + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); let ctx = SessionCtx { session_id: "pi-session".into(), config: Some( SessionRoute { + destination: Some(TraceDestination::ParentSpan { components }), additional_metadata: Some(json!({"team": "platform", "source": "custom"})), ..SessionRoute::default() } @@ -162,18 +172,34 @@ fn pi_additional_metadata_reaches_roots_without_overriding_session_fields() { }), ), }; - let mut event = event("session_start", 1, json!({"reason":"new"})); - event.payload["trace_settings"] = json!({ + let mut start_event = event("session_start", 1, json!({"reason":"new"})); + start_event.payload["trace_settings"] = json!({ "additional_metadata": {"team": "payload"}, "parent_span_id": "payload-parent", "root_span_id": "payload-root", }); - let rows = reduce(translator.handle(&event, &ctx).unwrap()); + let mut ops = translator.handle(&start_event, &ctx).unwrap(); + ops.extend( + translator + .handle( + &event("session_shutdown", 2, json!({"reason":"quit"})), + &ctx, + ) + .unwrap(), + ); + let inserted_root = ops + .iter() + .find_map(|op| match op { + SpanOp::Insert(row) if row.name == "Pi" => Some(row), + _ => None, + }) + .unwrap(); + assert_eq!(inserted_root.metadata.as_ref().unwrap()["team"], "platform"); + assert_eq!(inserted_root.metadata.as_ref().unwrap()["source"], "pi"); + let rows = reduce(ops); let root = rows.values().next().unwrap(); - assert_eq!(root.metadata.as_ref().unwrap()["team"], "platform"); - assert_eq!(root.metadata.as_ref().unwrap()["source"], "pi"); - assert!(root.parent_span_ids.is_empty()); - assert_eq!(root.root_span_id, root.span_id); + assert_eq!(root.parent_span_ids, ["external-parent"]); + assert_eq!(root.root_span_id, "external-root"); } #[test] diff --git a/bt-daemon/tests/support/span_identity.rs b/bt-daemon/tests/support/span_identity.rs new file mode 100644 index 0000000..ee0c0d6 --- /dev/null +++ b/bt-daemon/tests/support/span_identity.rs @@ -0,0 +1,32 @@ +use bt_daemon::SpanOp; +use std::collections::HashMap; + +/// Stateless merges must carry the same hierarchy as their original inserts. +pub(crate) fn assert_merges_preserve_insert_identity(ops: &[SpanOp]) { + let mut identities: HashMap)> = HashMap::new(); + for op in ops { + match op { + SpanOp::Insert(row) => { + identities.insert( + row.span_id.clone(), + (row.root_span_id.clone(), row.parent_span_ids.clone()), + ); + } + SpanOp::Merge(row) => { + let (root_span_id, parent_span_ids) = identities + .get(&row.span_id) + .unwrap_or_else(|| panic!("merge missing insert for span {}", row.span_id)); + assert_eq!( + &row.root_span_id, root_span_id, + "merge changed root identity for span {}", + row.span_id + ); + assert_eq!( + &row.parent_span_ids, parent_span_ids, + "merge dropped parent identity for span {}", + row.span_id + ); + } + } + } +}