diff --git a/bt-daemon/src/translate/antigravity.rs b/bt-daemon/src/translate/antigravity.rs index 3c80cb1..da71e8b 100644 --- a/bt-daemon/src/translate/antigravity.rs +++ b/bt-daemon/src/translate/antigravity.rs @@ -63,6 +63,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, @@ -84,6 +85,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, @@ -112,6 +114,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 @@ -152,7 +155,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), @@ -486,6 +489,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})), @@ -503,6 +507,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 123b507..4dacd87 100644 --- a/bt-daemon/src/translate/claude.rs +++ b/bt-daemon/src/translate/claude.rs @@ -150,6 +150,7 @@ struct TranscriptCursor { struct Subagent { span_id: String, + parent_span_id: String, transcript_path: Option, } @@ -248,6 +249,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, @@ -280,6 +282,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, @@ -325,6 +328,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 = hook.cwd.clone().unwrap_or_default(); let workspace = basename(&cwd); let mut metadata = ctx @@ -360,7 +364,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), @@ -422,6 +426,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() })); @@ -485,6 +490,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() })); @@ -529,6 +535,7 @@ impl ClaudeTranslator { hook.agent_id.clone(), Subagent { span_id: span_id.clone(), + parent_span_id: parent_id, transcript_path: None, }, ); @@ -594,6 +601,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), @@ -638,6 +646,11 @@ impl ClaudeTranslator { fn stop_subagent(&mut self, event: &Envelope, hook: SubagentHook, ops: &mut Vec) { let agent_id = hook.agent_id.clone(); let parent = self.ensure_subagent(&hook, 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 = hook.agent_transcript_path; if let Some(agent) = self.subagents.get_mut(&agent_id) { agent.transcript_path = path.clone(); @@ -667,6 +680,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() @@ -792,9 +806,11 @@ impl ClaudeTranslator { } self.emitted_tools.insert(tool.call_id.clone()); if let Some(pending) = self.pending_tools.remove(&tool.call_id) { - ops.push(SpanOp::Merge( - tool.into_update(pending.span_id, self.root_span_id.clone()), - )); + ops.push(SpanOp::Merge(tool.into_update( + pending.span_id, + self.root_span_id.clone(), + pending.parent_id, + ))); } else { let span_key = format!("tool:{}", tool.call_id); ops.push(SpanOp::Insert(tool.into_row( @@ -845,6 +861,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 @@ -882,6 +899,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), metadata: Some(Value::Object(metadata)), ..Default::default() @@ -909,6 +927,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() })); @@ -919,6 +938,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() })); @@ -954,6 +974,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() })); @@ -1035,6 +1056,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() @@ -1044,6 +1066,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() @@ -1053,6 +1076,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() @@ -1063,6 +1087,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() })); @@ -1652,10 +1677,11 @@ impl TranscriptTool { } } - fn into_update(self, span_id: String, root_span_id: String) -> SpanRow { + fn into_update(self, span_id: String, root_span_id: String, parent_id: String) -> SpanRow { SpanRow { span_id, root_span_id, + parent_span_ids: vec![parent_id], end_ms: Some(self.end_ms), output: self.output.clone(), metadata: Some(self.metadata()), diff --git a/bt-daemon/src/translate/codex.rs b/bt-daemon/src/translate/codex.rs index ca62e39..f2c7197 100644 --- a/bt-daemon/src/translate/codex.rs +++ b/bt-daemon/src/translate/codex.rs @@ -300,7 +300,7 @@ impl AgentTranslator for CodexTranslator { } "PreCompact" | "PostCompact" => { if let Some(hook) = decode::(payload) { - self.record_compaction_trigger(hook, &mut ops); + self.record_compaction_trigger(event, hook, &mut ops); } } _ => {} @@ -411,6 +411,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(), @@ -456,6 +470,7 @@ impl CodexTranslator { ops.push(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_span_id.clone().into_iter().collect(), input: Some(json!({ "model": model, "cwd": self.root_cwd, @@ -469,7 +484,12 @@ impl CodexTranslator { })); } - fn record_compaction_trigger(&mut self, hook: CompactHook, ops: &mut Vec) { + fn record_compaction_trigger( + &mut self, + event: &Envelope, + hook: CompactHook, + ops: &mut Vec, + ) { let turn_id = hook.turn_id; let trigger = hook.trigger.unwrap_or_else(|| "manual".to_string()); self.compaction_trigger_by_turn @@ -478,9 +498,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 = effective_transcript_path(event) + .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.effective_root_span_id.clone(), + parent_span_ids, metadata: Some(json!({ "compaction": { "trigger": trigger } })), ..Default::default() })); @@ -651,6 +676,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: scope.turn_parent_span_id.clone(), root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: self.scope_root_parent_span_ids(scope), input: Some(input), metadata: Some(metadata), ..Default::default() @@ -662,6 +688,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], metadata: Some(json!({ "model": m })), ..Default::default() })); @@ -754,7 +781,7 @@ impl CodexTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: self.root_span_id.clone(), root_span_id: self.effective_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), @@ -842,6 +869,7 @@ impl CodexTranslator { source: TurnInputSource, ops: &mut Vec, ) { + let turn_parent_span_id = scope.turn_parent_span_id.clone(); if let Some(turn) = scope.open_turns.last_mut() { match source { TurnInputSource::Authoritative => { @@ -870,6 +898,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.effective_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() @@ -1095,13 +1124,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))); } @@ -1113,6 +1142,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.effective_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))), @@ -1163,6 +1193,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: self.turn_parent_span_ids(&llm.turn_id), end_ms: Some(end), output, metadata: usage_metadata, @@ -1188,6 +1219,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 @@ -1203,6 +1235,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![turn_span_id], end_ms: Some(llm.last_output_ms), output, metadata: Some(json!({ @@ -1220,6 +1253,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id, root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: vec![scope.turn_parent_span_id.clone()], end_ms: Some(ts), output, ..Default::default() @@ -1233,6 +1267,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() @@ -1244,6 +1279,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.effective_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()), @@ -1285,6 +1321,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn_span.clone(), root_span_id: self.effective_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": { @@ -1346,6 +1383,7 @@ impl CodexTranslator { ops.push(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_span_id.clone().into_iter().collect(), end_ms: Some(end_ms), late_merge_key: Some(format!("session:stop:{end_ms}")), ..Default::default() @@ -1376,6 +1414,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: scope.turn_parent_span_id.clone(), root_span_id: self.effective_root_span_id.clone(), + parent_span_ids: self.scope_root_parent_span_ids(&scope), end_ms: Some(end), ..Default::default() })); @@ -1397,6 +1436,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.effective_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!({ @@ -1405,15 +1445,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.effective_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), @@ -1424,6 +1467,7 @@ impl CodexTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id, root_span_id: self.effective_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/grok.rs b/bt-daemon/src/translate/grok.rs index ba699f8..0e5314e 100644 --- a/bt-daemon/src/translate/grok.rs +++ b/bt-daemon/src/translate/grok.rs @@ -143,6 +143,7 @@ struct TurnCompletion { struct OpenLlm { span_id: String, + parent_span_id: String, prompt_id: Option, stream_start_ms: Option, output: BoundedOutput, @@ -153,11 +154,13 @@ struct OpenLlm { struct OpenTool { span_id: String, + parent_span_id: String, start_ms: i64, } struct CompletedTool { span_id: String, + parent_span_id: String, end_ms: i64, terminal_update_seen: bool, } @@ -209,6 +212,7 @@ struct GrokTranslator { session_id: String, session_span_id: String, root_span_id: String, + session_parent_span_ids: Vec, root_open: bool, root_closed: bool, root_generation: u32, @@ -224,6 +228,7 @@ struct GrokTranslator { events: TranscriptCursor, system_prompt: Option, first_llm_span_id: Option, + first_llm_parent_span_id: Option, first_llm_user_input: Option, pending: Option, cwd: Option, @@ -238,6 +243,7 @@ impl GrokTranslator { 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_closed: false, root_generation: 0, @@ -253,6 +259,7 @@ impl GrokTranslator { events: TranscriptCursor::default(), system_prompt: None, first_llm_span_id: None, + first_llm_parent_span_id: None, first_llm_user_input: None, pending: None, cwd: None, @@ -271,6 +278,7 @@ impl GrokTranslator { self.emitted_turns.clear(); self.emitted_chunks.clear(); self.first_llm_span_id = None; + self.first_llm_parent_span_id = None; self.first_llm_user_input = None; self.last_ts_ms = 0; } @@ -292,6 +300,7 @@ impl GrokTranslator { 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), metadata: Some(json!({ "resumed": true, @@ -308,6 +317,7 @@ impl GrokTranslator { .map(|config| config.attached_span_ids()) .unwrap_or_default(); self.root_span_id = attached.1.unwrap_or_else(|| self.session_span_id.clone()); + self.session_parent_span_ids = attached.0.clone().into_iter().collect(); self.root_open = true; let mut metadata = ctx @@ -336,7 +346,7 @@ impl GrokTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: self.session_span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: attached.0.into_iter().collect(), + parent_span_ids: self.session_parent_span_ids.clone(), name: "Grok".into(), span_type: SpanType::Task, start_ms: Some(ts_ms), @@ -346,12 +356,15 @@ impl GrokTranslator { } fn sync_user_input(&mut self, ops: &mut Vec) { + let session_span_id = self.session_span_id.clone(); let Some(turn) = self.current_turn.as_mut() else { return; }; if !turn.user_input_dirty { return; } + let turn_span_id = turn.span_id.clone(); + let first_llm_parent = self.first_llm_parent_span_id.clone(); let input = turn.user_input.message_content(); let mut metadata = Map::new(); turn.user_input @@ -360,6 +373,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: turn.span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![session_span_id], input: input.clone(), metadata: (!metadata.is_empty()).then_some(Value::Object(metadata)), ..Default::default() @@ -377,6 +391,9 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![first_llm_parent + .clone() + .unwrap_or_else(|| turn_span_id.clone())], input: first_llm_input( self.system_prompt.as_deref(), self.first_llm_user_input.as_ref(), @@ -528,7 +545,7 @@ impl GrokTranslator { ); let turn = self.current_turn.as_ref().expect("turn checked above"); let sequence = turn.llm_span_ids.len() + 1; - let turn_span_id = turn.span_id.clone(); + let llm_parent_span_id = turn.span_id.clone(); let turn_key = turn.key.clone(); let model = turn.model.clone().unwrap_or_else(|| "Grok".into()); let identity_prompt = prompt_id.as_deref().unwrap_or("unknown"); @@ -544,6 +561,7 @@ impl GrokTranslator { let is_first_llm = self.first_llm_span_id.is_none(); let input = if is_first_llm { self.first_llm_span_id = Some(span_id.clone()); + self.first_llm_parent_span_id = Some(llm_parent_span_id.clone()); first_llm_input( self.system_prompt.as_deref(), self.first_llm_user_input.as_ref(), @@ -574,7 +592,7 @@ impl GrokTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: vec![turn_span_id], + parent_span_ids: vec![llm_parent_span_id.clone()], name: format!("{model} call {sequence}"), span_type: SpanType::Llm, start_ms: Some(boundary_ms), @@ -589,6 +607,7 @@ impl GrokTranslator { .push(span_id.clone()); self.open_llm = Some(OpenLlm { span_id, + parent_span_id: llm_parent_span_id, prompt_id: prompt_id.clone(), stream_start_ms, output: BoundedOutput::default(), @@ -628,7 +647,7 @@ impl GrokTranslator { } "tool_call" => { self.close_llm(ts_ms, "tool_call", Some("llm:tool_call".into()), ops); - let Some(turn_span_id) = + let Some(tool_parent_span_id) = self.current_turn.as_ref().map(|turn| turn.span_id.clone()) else { return; @@ -654,7 +673,7 @@ impl GrokTranslator { ops.push(SpanOp::Insert(SpanRow { span_id: span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: vec![turn_span_id], + parent_span_ids: vec![tool_parent_span_id.clone()], name, span_type: SpanType::Tool, start_ms: Some(ts_ms), @@ -666,6 +685,7 @@ impl GrokTranslator { call_key, OpenTool { span_id, + parent_span_id: tool_parent_span_id, start_ms: ts_ms, }, ); @@ -690,6 +710,7 @@ impl GrokTranslator { ( OpenTool { span_id: completed.span_id, + parent_span_id: completed.parent_span_id, start_ms: completed.end_ms, }, true, @@ -699,10 +720,11 @@ impl GrokTranslator { return; }; let span_id = ids::span_id(&self.session_id, &format!("tool:{call_id}")); + let parent_span_id = turn.span_id.clone(); ops.push(SpanOp::Insert(SpanRow { span_id: span_id.clone(), root_span_id: self.root_span_id.clone(), - parent_span_ids: vec![turn.span_id.clone()], + parent_span_ids: vec![parent_span_id.clone()], name: update .get("title") .and_then(Value::as_str) @@ -720,6 +742,7 @@ impl GrokTranslator { ( OpenTool { span_id, + parent_span_id, start_ms: ts_ms, }, false, @@ -730,6 +753,7 @@ impl GrokTranslator { call_key.clone(), CompletedTool { span_id: tool.span_id.clone(), + parent_span_id: tool.parent_span_id.clone(), end_ms, terminal_update_seen: true, }, @@ -753,6 +777,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: tool.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![tool.parent_span_id], end_ms: Some(end_ms), output: update .get("rawOutput") @@ -779,6 +804,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm_span_id.clone(), root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![turn.span_id.clone()], metrics: Some(metrics.clone()), metadata: Some(json!({ "usage_scope": "turn", @@ -853,6 +879,7 @@ impl GrokTranslator { key, CompletedTool { span_id: completed.span_id, + parent_span_id: completed.parent_span_id.clone(), end_ms, terminal_update_seen: completed.terminal_update_seen, }, @@ -873,6 +900,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![completed.parent_span_id], end_ms: Some(end_ms), metadata: Some(Value::Object(metadata)), error: failed.then(|| format!("Grok tool outcome: {}", outcome.unwrap_or("error"))), @@ -916,6 +944,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: llm.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![llm.parent_span_id], end_ms: Some(end_ms.max(llm.last_ms).max(llm.start_ms)), output: Some(Value::Array(vec![output])), metadata: Some(Value::Object(metadata)), @@ -939,6 +968,7 @@ impl GrokTranslator { call_id.clone(), CompletedTool { span_id: tool.span_id.clone(), + parent_span_id: tool.parent_span_id.clone(), end_ms, terminal_update_seen: false, }, @@ -946,6 +976,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: tool.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![tool.parent_span_id], end_ms: Some(end_ms), metadata: Some(json!({ "incomplete": true, @@ -963,6 +994,7 @@ impl GrokTranslator { call_id.clone(), CompletedTool { span_id: tool.span_id.clone(), + parent_span_id: tool.parent_span_id.clone(), end_ms, terminal_update_seen: false, }, @@ -970,6 +1002,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: tool.span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: vec![tool.parent_span_id], end_ms: Some(end_ms), metadata: Some(json!({"incomplete": true, "close_reason": reason})), ..Default::default() @@ -1005,6 +1038,7 @@ impl GrokTranslator { 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(end_ms), output: Some(turn.assistant_output.into_value()), metadata: Some(Value::Object(metadata)), @@ -1072,6 +1106,7 @@ impl GrokTranslator { ops.push(SpanOp::Merge(SpanRow { span_id, root_span_id: self.root_span_id.clone(), + parent_span_ids: self.first_llm_parent_span_id.clone().into_iter().collect(), input: Some(input), metadata: Some(Value::Object(metadata)), ..Default::default() @@ -1180,6 +1215,7 @@ impl GrokTranslator { 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(terminal_ms), late_merge_key: Some(format!("session:terminal:{}", self.root_generation)), ..Default::default() @@ -1238,6 +1274,7 @@ impl AgentTranslator for GrokTranslator { 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(self.last_ts_ms), late_merge_key: Some(format!("session:finalize:{}", self.root_generation)), ..Default::default() diff --git a/bt-daemon/src/translate/opencode.rs b/bt-daemon/src/translate/opencode.rs index 0ba7b0a..24b1193 100644 --- a/bt-daemon/src/translate/opencode.rs +++ b/bt-daemon/src/translate/opencode.rs @@ -347,6 +347,7 @@ fn properties(payload: &Value, field: &str) -> Option { 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, @@ -358,6 +359,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, @@ -637,6 +639,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() }, @@ -666,6 +669,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() @@ -988,6 +992,7 @@ impl OpenCodeTranslator { return vec![]; } s.tool_starts.insert(call.clone(), ts_ms); + s.tool_parent_span_ids.insert(call.clone(), turn.clone()); if let Some(args) = event.output.and_then(|output| output.args) { s.tool_args.insert(call.clone(), args); } @@ -1020,6 +1025,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 { @@ -1065,10 +1071,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(ts_ms)), @@ -1177,6 +1184,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), @@ -1193,7 +1201,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(ts_ms)), @@ -1238,9 +1246,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!({ @@ -1257,6 +1272,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(), @@ -1267,6 +1283,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}), @@ -1281,6 +1298,7 @@ impl OpenCodeTranslator { ops.push(SpanOp::Merge(SpanRow { span_id: s.root_span_id.clone(), root_span_id: s.effective_root_span_id.clone(), + parent_span_ids: s.parent_span_ids.clone(), 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 072d2d8..f41d241 100644 --- a/bt-daemon/src/translate/pi.rs +++ b/bt-daemon/src/translate/pi.rs @@ -765,11 +765,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), @@ -846,6 +842,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() @@ -859,6 +856,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 91ed465..743ef55 100644 --- a/bt-daemon/tests/antigravity_translator.rs +++ b/bt-daemon/tests/antigravity_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; 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 1636e96..205615c 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}; @@ -92,6 +92,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; @@ -398,18 +414,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() }; @@ -422,7 +440,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 0fc3a21..0034953 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}; @@ -111,6 +116,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 { @@ -259,10 +265,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"})), tags: vec!["ci".into(), "docs".into()], ..SessionRoute::default() @@ -276,34 +286,33 @@ 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.tags, Some(vec!["ci".into(), "docs".into()])); + 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 5f9ae6a..f07a8e1 100644 --- a/bt-daemon/tests/codex_translator.rs +++ b/bt-daemon/tests/codex_translator.rs @@ -2,10 +2,14 @@ //! hook triggers into a session → turn → {llm, tool} span tree. Mirrors the //! happy-path shape of the TS `event-processor` tests. +#[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; @@ -180,6 +184,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). @@ -425,6 +430,54 @@ fn codex_native_prompt_ignores_surrounding_injected_user_rows() { ); } +#[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, + tags: Vec::new(), + }), + }; + 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. @@ -824,6 +877,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"); @@ -1216,6 +1270,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"); @@ -1481,6 +1536,7 @@ fn codex_subagent_with_malformed_optional_type_nests_under_spawning_turn() { .unwrap(), ); + assert_merges_preserve_insert_identity(&ops); let rows = reduce(ops); let root = find(&rows, SpanType::Task, "codex: app"); @@ -1580,3 +1636,88 @@ fn codex_rollout_routing_ignores_future_fields() { assert_eq!(tool.input, Some(json!("{\"command\":\"pwd\"}"))); assert_eq!(tool.output, Some(json!("/work/app"))); } + +/// `end_main_root` closes the root handle on every main-scope Stop, so a later +/// `SessionStart` (resume/compact) re-merges the root through the sink's +/// stateless path. That merge must still carry the attached parent. +#[test] +fn codex_root_source_merge_after_stop_keeps_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, + tags: Vec::new(), + }), + }; + 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(), + ); + // Resume in the same live session: the root handle is already closed. + ops.extend( + translator + .handle( + &envelope( + "attached-session", + "SessionStart", + path, + json!({ "source": "resume" }), + ), + &ctx, + ) + .unwrap(), + ); + ops.extend(translator.flush(&ctx).unwrap()); + + let root_span_id = ops + .iter() + .find_map(|op| match op { + SpanOp::Insert(row) if row.name.starts_with("codex:") => Some(row.span_id.clone()), + _ => None, + }) + .expect("root insert"); + let source_merge = ops + .iter() + .filter_map(|op| match op { + SpanOp::Merge(row) if row.span_id == root_span_id => Some(row), + _ => None, + }) + .find(|row| { + row.metadata + .as_ref() + .and_then(|m| m.get("session_source")) + .is_some() + }) + .expect("post-Stop root source merge"); + assert_eq!( + source_merge.parent_span_ids, + ["external-parent"], + "post-Stop root merge dropped the external parent" + ); + assert_merges_preserve_insert_identity(&ops); +} diff --git a/bt-daemon/tests/grok_translator.rs b/bt-daemon/tests/grok_translator.rs index d2c01c8..dee3ad8 100644 --- a/bt-daemon/tests/grok_translator.rs +++ b/bt-daemon/tests/grok_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::{AgentTranslator, Registry, SessionCtx, SpanOp, SpanType}; use serde_json::json; +use span_identity::{assert_merges_preserve_insert_identity, IdentityLedger}; use std::path::{Path, PathBuf}; fn fixture(name: &str) -> PathBuf { @@ -64,9 +69,30 @@ fn point_at(event: &mut Envelope, updates: &Path, events: &Path) { json!(std::fs::metadata(events).unwrap().len()); } +thread_local! { + /// One ledger per test (each `#[test]` runs on its own thread), so a merge + /// emitted by a later hook is still checked against the batch that inserted + /// the span. + static LEDGER: std::cell::RefCell = + std::cell::RefCell::new(IdentityLedger::default()); +} + +fn check_identity(ops: &[SpanOp]) { + LEDGER.with(|ledger| ledger.borrow_mut().check(ops)); +} + +/// Handle one event and assert every merge it emits still carries the parent +/// identity its insert used. +fn handle(translator: &mut dyn AgentTranslator, event: &Envelope) -> Vec { + let ops = translator.handle(event, &ctx()).unwrap(); + check_identity(&ops); + ops +} + fn drain(translator: &mut dyn AgentTranslator) -> Vec { let mut ops = Vec::new(); while let Some(batch) = translator.drain_pending(&ctx()).unwrap() { + check_identity(&batch); ops.extend(batch); } ops @@ -271,7 +297,7 @@ fn grok_enriches_emitted_spans_from_the_captured_working_directory() { event.payload["cwd"] = json!(repo); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let inserts: Vec<_> = ops .iter() @@ -351,7 +377,7 @@ fn grok_single_call_usage_and_missing_tool_start_are_recovered() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let llm = ops .iter() .find_map(|op| match op { @@ -379,7 +405,7 @@ fn grok_edge_fixture_preserves_turns_skips_malformed_records_and_marks_mismatche point_at(&mut event, &updates, &events); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let turns: Vec<_> = ops .iter() @@ -478,7 +504,7 @@ fn grok_partial_records_are_retried_and_complete_malformed_lines_do_not_stall() let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let first = translator.handle(&first_event, &ctx()).unwrap(); + let first = handle(&mut *translator, &first_event); assert!(!first .iter() .any(|op| matches!(op, SpanOp::Insert(row) if row.span_type == SpanType::Llm))); @@ -490,7 +516,7 @@ fn grok_partial_records_are_retried_and_complete_malformed_lines_do_not_stall() std::fs::write(&updates_path, completed).unwrap(); let mut second_event = envelope(0, 0); point_at(&mut second_event, &updates_path, &events_path); - let second = translator.handle(&second_event, &ctx()).unwrap(); + let second = handle(&mut *translator, &second_event); assert!(second .iter() .any(|op| matches!(op, SpanOp::Insert(row) if row.span_type == SpanType::Llm))); @@ -515,7 +541,7 @@ fn grok_processes_a_complete_terminal_record_without_a_trailing_newline() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); assert!(ops.iter().any(|op| { matches!(op, SpanOp::Merge(row) @@ -546,7 +572,7 @@ fn grok_events_failure_does_not_block_updates_and_recovers_enrichment() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let primary = translator.handle(&event, &ctx()).unwrap(); + let primary = handle(&mut *translator, &event); assert_eq!( primary .iter() @@ -562,7 +588,7 @@ fn grok_events_failure_does_not_block_updates_and_recovers_enrichment() { })); std::fs::write(&events_path, event_body).unwrap(); - let enrichment = translator.handle(&event, &ctx()).unwrap(); + let enrichment = handle(&mut *translator, &event); assert!(!enrichment.iter().any(|op| matches!(op, SpanOp::Insert(_)))); let enriched_tools: Vec<_> = enrichment .iter() @@ -578,7 +604,7 @@ fn grok_events_failure_does_not_block_updates_and_recovers_enrichment() { }) .collect(); assert_eq!(enriched_tools.len(), 1); - assert!(translator.handle(&event, &ctx()).unwrap().is_empty()); + assert!(handle(&mut *translator, &event).is_empty()); } #[test] @@ -599,7 +625,7 @@ fn grok_buffers_tool_events_until_the_matching_update_arrives() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let first = translator.handle(&first_event, &ctx()).unwrap(); + let first = handle(&mut *translator, &first_event); assert!(!first.iter().any(|op| { matches!(op, SpanOp::Merge(row) if row.metadata.as_ref().is_some_and(|metadata| metadata.get("outcome").is_some())) @@ -612,7 +638,7 @@ fn grok_buffers_tool_events_until_the_matching_update_arrives() { .unwrap(); let mut second_event = envelope(0, 0); point_at(&mut second_event, &updates_path, &events_path); - let second = translator.handle(&second_event, &ctx()).unwrap(); + let second = handle(&mut *translator, &second_event); assert!(second.iter().any(|op| { matches!(op, SpanOp::Merge(row) if row.output == Some(json!("done")) @@ -644,7 +670,7 @@ fn grok_drains_every_captured_event_batch_before_finishing_the_hook() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let mut ops = translator.handle(&envelope, &ctx()).unwrap(); + let mut ops = handle(&mut *translator, &envelope); ops.extend(drain(translator.as_mut())); assert_eq!( @@ -681,7 +707,7 @@ fn grok_updates_failure_keeps_primary_state_retryable() { assert!(translator.handle(&event, &ctx()).is_err()); std::fs::write(&updates_path, update_body).unwrap(); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); assert_eq!( ops.iter() .filter(|op| matches!(op, SpanOp::Insert(row) if row.name == "Turn 1")) @@ -711,7 +737,7 @@ fn grok_resets_semantic_state_when_the_transcript_is_replaced() { point_at(&mut first_event, &updates_path, &events_path); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - translator.handle(&first_event, &ctx()).unwrap(); + handle(&mut *translator, &first_event); let replacement = concat!( "{\"params\":{\"update\":{\"sessionUpdate\":\"user_message_chunk\",\"content\":{\"text\":\"new\"}},\"_meta\":{\"promptIndex\":0,\"agentTimestampMs\":2000}}}\n", @@ -720,7 +746,7 @@ fn grok_resets_semantic_state_when_the_transcript_is_replaced() { std::fs::write(&updates_path, replacement).unwrap(); let mut replacement_event = envelope(0, 0); point_at(&mut replacement_event, &updates_path, &events_path); - let ops = translator.handle(&replacement_event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &replacement_event); assert!(ops.iter().any(|op| { matches!(op, SpanOp::Insert(row) @@ -757,7 +783,7 @@ fn grok_supports_native_and_documented_terminal_events() { let mut translator = registry.create("grok", "grok-session"); let mut event = base.clone(); event.event = event_name.into(); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let root_id = ops .iter() .find_map(|op| match op { @@ -779,7 +805,7 @@ fn grok_supports_native_and_documented_terminal_events() { let mut translator = registry.create("grok", "grok-session"); let mut event = base.clone(); event.event = event_name.into(); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let turn_id = ops .iter() .find_map(|op| match op { @@ -823,7 +849,7 @@ fn grok_resume_extends_and_recloses_the_existing_session_root() { let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let first = translator.handle(&first_end, &ctx()).unwrap(); + let first = handle(&mut *translator, &first_end); let root_id = first .iter() .find_map(|op| match op { @@ -845,7 +871,7 @@ fn grok_resume_extends_and_recloses_the_existing_session_root() { let mut resumed = envelope(0, 0); resumed.ts_ms = 3_100; point_at(&mut resumed, &updates_path, &events_path); - let second = translator.handle(&resumed, &ctx()).unwrap(); + let second = handle(&mut *translator, &resumed); assert!(second.iter().any(|op| { matches!(op, SpanOp::Merge(row) if row.span_id == root_id @@ -856,7 +882,7 @@ fn grok_resume_extends_and_recloses_the_existing_session_root() { let mut second_end = resumed; second_end.event = "session_end".into(); second_end.ts_ms = 3_200; - let third = translator.handle(&second_end, &ctx()).unwrap(); + let third = handle(&mut *translator, &second_end); assert!(third.iter().any(|op| { matches!(op, SpanOp::Merge(row) if row.span_id == root_id @@ -881,7 +907,7 @@ fn grok_catch_up_and_open_state_bounds_drain_without_losing_later_records() { point_at(&mut event, &updates_path, &events_path); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let mut ops = translator.handle(&event, &ctx()).unwrap(); + let mut ops = handle(&mut *translator, &event); ops.extend(drain(translator.as_mut())); assert_eq!( ops.iter() @@ -910,7 +936,7 @@ fn grok_finalize_closes_open_spans() { point_at(&mut event, &updates_path, &events_path); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let handled = translator.handle(&event, &ctx()).unwrap(); + let handled = handle(&mut *translator, &event); let root_id = handled .iter() .find_map(|op| match op { @@ -976,7 +1002,7 @@ fn grok_assembles_user_and_identical_assistant_chunks() { point_at(&mut event, &updates_path, &events_path); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let ops = translator.handle(&event, &ctx()).unwrap(); + let ops = handle(&mut *translator, &event); let turn_id = ops .iter() .find_map(|op| match op { @@ -1033,7 +1059,7 @@ fn grok_late_update_completes_an_evicted_tool_without_reinserting_it() { point_at(&mut event, &updates_path, &events_path); let registry = Registry::default_agents(); let mut translator = registry.create("grok", "grok-session"); - let mut ops = translator.handle(&event, &ctx()).unwrap(); + let mut ops = handle(&mut *translator, &event); ops.extend(drain(translator.as_mut())); let tool_id = ops .iter() @@ -1063,3 +1089,83 @@ fn grok_late_update_completes_an_evicted_tool_without_reinserting_it() { && row.metadata.as_ref().is_some_and(|metadata| metadata["incomplete"] == false)) })); } + +/// An attached Grok session parents its root at the caller's span. Every +/// session-row merge — resume, terminal, finalize — must repeat that parent, +/// because each carries a `late_merge_key` and so is expected to land after +/// the terminal row, with no open handle to borrow identity from. +#[test] +fn grok_attached_session_merges_keep_the_external_parent() { + let mut components = SpanComponents::new(SpanObjectType::ProjectLogs); + components.span_id = Some("external-parent".into()); + components.root_span_id = Some("external-root".into()); + let attached = SessionCtx { + session_id: "grok-session".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 temp = tempfile::tempdir().unwrap(); + let updates_path = temp.path().join("updates.jsonl"); + let events_path = temp.path().join("events.jsonl"); + std::fs::write(&events_path, "").unwrap(); + std::fs::write( + &updates_path, + concat!( + "{\"params\":{\"update\":{\"sessionUpdate\":\"user_message_chunk\",\"content\":{\"text\":\"go\"}},\"_meta\":{\"promptIndex\":0,\"agentTimestampMs\":1000}}}\n", + "{\"params\":{\"update\":{\"sessionUpdate\":\"turn_completed\",\"usage\":{\"modelCalls\":0}},\"_meta\":{\"agentTimestampMs\":1100}}}\n" + ), + ) + .unwrap(); + let mut end = envelope(0, 0); + end.event = "SessionEnd".into(); + end.ts_ms = 1_200; + point_at(&mut end, &updates_path, &events_path); + + let registry = Registry::default_agents(); + let mut translator = registry.create("grok", "grok-session"); + let mut ops = translator.handle(&end, &attached).unwrap(); + ops.extend(translator.finalize(&attached).unwrap()); + + let root = ops + .iter() + .find_map(|op| match op { + SpanOp::Insert(row) if row.name == "Grok" => Some(row.clone()), + _ => None, + }) + .expect("root insert"); + assert_eq!(root.root_span_id, "external-root"); + assert_eq!(root.parent_span_ids, ["external-parent"]); + + let session_merges: Vec<_> = ops + .iter() + .filter_map(|op| match op { + SpanOp::Merge(row) if row.span_id == root.span_id => Some(row), + _ => None, + }) + .collect(); + assert!( + !session_merges.is_empty(), + "expected at least one session-row merge" + ); + for merge in session_merges { + assert_eq!( + merge.parent_span_ids, + ["external-parent"], + "session merge {:?} dropped the external parent", + merge.late_merge_key + ); + } + assert_merges_preserve_insert_identity(&ops); +} diff --git a/bt-daemon/tests/opencode_translator.rs b/bt-daemon/tests/opencode_translator.rs index 6179ef6..7ff2656 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 { @@ -336,6 +342,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 @@ -350,10 +368,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"})), tags: vec!["ci".into(), "docs".into()], ..SessionRoute::default() @@ -367,22 +389,50 @@ 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" + ); + assert!(inserted_root + .metadata + .as_ref() + .unwrap() + .get("username") + .is_some()); + 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!(root.metadata.as_ref().unwrap().get("username").is_some()); + assert_eq!(root.root_span_id, "external-root"); + assert_eq!(root.parent_span_ids, ["external-parent"]); assert_eq!(root.tags, Some(vec!["ci".into(), "docs".into()])); } @@ -434,6 +484,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( @@ -441,13 +492,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() ))); } @@ -726,3 +779,84 @@ fn opencode_accepts_flattened_tool_hook_payloads() { assert_eq!(tool.output.as_ref().unwrap(), "contents"); assert_eq!(tool.end_ms, Some(4)); } + +/// `session.idle` fires after every turn and merges the root summary while the +/// session stays resumable. That merge must keep the attached parent — the +/// sibling close-root arm already does. +#[test] +fn opencode_idle_root_merge_keeps_external_parent() { + 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 }), + ..SessionRoute::default() + } + .with_auth(BackendAuth { + token: "test".into(), + api_url: None, + app_url: None, + org_name: None, + org_id: None, + }), + ), + }; + let mut ops = translator + .handle( + &event( + "session.created", + 1, + json!({"properties":{"info":{"id":"native"}}}), + ), + &ctx, + ) + .unwrap(); + ops.extend( + translator + .handle( + &event("chat.message", 2, json!({"input":{"sessionID":"native"}})), + &ctx, + ) + .unwrap(), + ); + // Idle, not deleted: the session stays resumable for the next turn. + ops.extend( + translator + .handle( + &event( + "session.idle", + 3, + json!({"properties":{"sessionID":"native"}}), + ), + &ctx, + ) + .unwrap(), + ); + + let root_span_id = ops + .iter() + .find_map(|op| match op { + SpanOp::Insert(row) if row.name == "OpenCode" => Some(row.span_id.clone()), + _ => None, + }) + .expect("root insert"); + let idle_merge = ops + .iter() + .filter_map(|op| match op { + SpanOp::Merge(row) if row.span_id == root_span_id => Some(row), + _ => None, + }) + .next_back() + .expect("idle root merge"); + assert_eq!( + idle_merge.parent_span_ids, + ["external-parent"], + "idle root merge dropped the external parent" + ); + assert_merges_preserve_insert_identity(&ops); +} diff --git a/bt-daemon/tests/pi_translator.rs b/bt-daemon/tests/pi_translator.rs index 64bf0ac..fb1dd73 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 { @@ -230,10 +236,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"})), tags: vec!["ci".into(), "docs".into()], ..SessionRoute::default() @@ -247,19 +257,40 @@ 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"); + assert!(inserted_root + .metadata + .as_ref() + .unwrap() + .get("username") + .is_some()); + 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.metadata.as_ref().unwrap().get("username").is_some()); - 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"); assert_eq!(root.tags, Some(vec!["ci".into(), "docs".into()])); } diff --git a/bt-daemon/tests/support/span_identity.rs b/bt-daemon/tests/support/span_identity.rs new file mode 100644 index 0000000..763f547 --- /dev/null +++ b/bt-daemon/tests/support/span_identity.rs @@ -0,0 +1,48 @@ +use bt_daemon::SpanOp; +use std::collections::HashMap; + +/// Remembers the identity each span was inserted with, so merges emitted in a +/// later batch than their insert can still be checked. +#[derive(Default)] +pub(crate) struct IdentityLedger { + identities: HashMap)>, +} + +impl IdentityLedger { + /// Record this batch's inserts, then assert its merges repeat the identity + /// their span was inserted with. + pub(crate) fn check(&mut self, ops: &[SpanOp]) { + for op in ops { + match op { + SpanOp::Insert(row) => { + self.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) = self + .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 + ); + } + } + } + } +} + +/// Stateless merges must carry the same hierarchy as their original inserts. +#[allow(dead_code)] +pub(crate) fn assert_merges_preserve_insert_identity(ops: &[SpanOp]) { + IdentityLedger::default().check(ops); +}