From 6001e667588d3e2eba23f9c5ac5d9a4777b7a341 Mon Sep 17 00:00:00 2001 From: Stephen Belanger Date: Fri, 11 Sep 2026 22:57:34 +0800 Subject: [PATCH] fix: ignore unlinked Claude child transcripts --- bt-daemon/src/transcript_import/claude.rs | 145 +++++++++++++++------- bt-daemon/src/transcript_import/mod.rs | 143 +++++++++++++++++++++ bt-daemon/tests/replay.rs | 61 +++++++++ 3 files changed, 307 insertions(+), 42 deletions(-) diff --git a/bt-daemon/src/transcript_import/claude.rs b/bt-daemon/src/transcript_import/claude.rs index ddcc6dd..df9605c 100644 --- a/bt-daemon/src/transcript_import/claude.rs +++ b/bt-daemon/src/transcript_import/claude.rs @@ -1,11 +1,11 @@ use super::{ - envelope, file_session_id, find_jsonl_files, read_jsonl_records, string_at, timestamp_bounds, - timestamp_ms, validate_session_id, + envelope, file_session_id, find_jsonl_files, read_complete_jsonl_records, read_jsonl_records, + string_at, timestamp_bounds, timestamp_ms, validate_session_id, }; use crate::wire::Envelope; use anyhow::bail; use serde_json::{json, Value}; -use std::collections::HashMap; +use std::collections::{HashMap, HashSet}; use std::io::BufRead; use std::path::{Path, PathBuf}; @@ -279,7 +279,7 @@ fn subagents(path: &Path, records: &[Value]) -> anyhow::Result> { paths.sort(); let mut result_records = HashMap::)>::new(); - let mut calls = HashMap::)>::new(); + let mut calls = HashMap::, Option)>::new(); for (index, record) in records.iter().enumerate() { if let Some(agent_id) = string_at(record, "/toolUseResult/agentId") { let call_id = record @@ -306,56 +306,117 @@ fn subagents(path: &Path, records: &[Value]) -> anyhow::Result> { let Some(call_id) = string_at(block, "/id") else { continue; }; - calls.insert(call_id, (index, string_at(block, "/input/subagent_type"))); + calls.insert( + call_id, + ( + index, + string_at(block, "/input/subagent_type"), + string_at(block, "/input/prompt"), + ), + ); } } - let fallback_index = records.len().saturating_sub(1); - paths + let completed_calls = result_records + .values() + .filter_map(|(_, call_id)| call_id.as_deref()) + .collect::>(); + let mut unmatched_calls = calls + .iter() + .filter(|(call_id, _)| !completed_calls.contains(call_id.as_str())) + .map(|(call_id, (index, agent_type, prompt))| { + ( + call_id.clone(), + *index, + timestamp_ms(&records[*index]).unwrap_or_default(), + agent_type.clone(), + prompt.clone(), + ) + }) + .collect::>(); + unmatched_calls.sort_by_key(|(_, index, _, _, _)| *index); + + let mut children = Vec::new(); + for path in paths { + let Some(agent_id) = path + .file_stem() + .and_then(|name| name.to_str()) + .and_then(|name| name.strip_prefix("agent-")) + .map(str::to_owned) + else { + continue; + }; + let records = if result_records.contains_key(&agent_id) { + read_jsonl_records(&path)? + } else { + read_complete_jsonl_records(&path)? + }; + let start_ms = timestamp_bounds(&records).0; + children.push((path, agent_id, records, start_ms)); + } + children.sort_by(|left, right| left.3.cmp(&right.3).then_with(|| left.0.cmp(&right.0))); + + let subagents = children .into_iter() - .filter_map(|subagent_path| { - let name = subagent_path.file_stem()?.to_str()?; - let agent_id = name.strip_prefix("agent-")?.to_owned(); - let (record_index, call_id) = result_records - .get(&agent_id) - .cloned() - .unwrap_or((fallback_index, None)); + .filter_map(|(subagent_path, agent_id, subagent_records, child_start)| { + // A completed tool result is the strongest parent/child link. If + // the session was interrupted first, Claude still leaves both the + // Agent call and a sidechain whose initial prompt matches that call. + // Files with neither anchor are auxiliary or stale and must not + // become synthetic `subagent: agent` outputs on the final turn. + let (record_index, call_id) = if let Some(link) = result_records.get(&agent_id) { + link.clone() + } else { + let child_prompt = subagent_records.iter().find_map(|record| { + (record.get("isSidechain").and_then(Value::as_bool) == Some(true) + && string_at(record, "/agentId").as_deref() == Some(&agent_id)) + .then(|| string_at(record, "/message/content")) + .flatten() + })?; + let child_type = subagent_records + .iter() + .find_map(|record| string_at(record, "/attributionAgent")); + let call_index = unmatched_calls + .iter() + .enumerate() + .filter(|(_, (_, _, call_start, agent_type, prompt))| { + prompt.as_deref() == Some(child_prompt.as_str()) + && *call_start <= child_start + && match (agent_type.as_deref(), child_type.as_deref()) { + (Some(expected), Some(actual)) => expected == actual, + _ => true, + } + }) + .min_by_key(|(_, (_, _, call_start, _, _))| child_start - *call_start) + .map(|(index, _)| index)?; + let (call_id, index, _, _, _) = unmatched_calls.remove(call_index); + (index, Some(call_id)) + }; let (start_index, agent_type) = call_id .as_ref() .and_then(|call_id| calls.get(call_id)) - .cloned() + .map(|(index, agent_type, _)| (*index, agent_type.clone())) .unwrap_or((record_index, None)); - Some(( - subagent_path, + let (_, child_end) = timestamp_bounds(&subagent_records); + let start_ms = timestamp_ms(&records[start_index]) + .unwrap_or(child_start) + .min(child_start); + let end_ms = timestamp_ms(&records[record_index]) + .unwrap_or(child_end) + .max(child_end); + Some(Subagent { agent_id, + agent_type, + path: subagent_path.to_string_lossy().into_owned(), start_index, record_index, - agent_type, - )) + start_ms, + end_ms, + last_assistant_message: last_assistant_text(&subagent_records), + }) }) - .map( - |(subagent_path, agent_id, start_index, record_index, agent_type)| { - let subagent_records = read_jsonl_records(&subagent_path)?; - let (child_start, child_end) = timestamp_bounds(&subagent_records); - let start_ms = timestamp_ms(&records[start_index]) - .unwrap_or(child_start) - .min(child_start); - let end_ms = timestamp_ms(&records[record_index]) - .unwrap_or(child_end) - .max(child_end); - Ok(Subagent { - agent_id, - agent_type, - path: subagent_path.to_string_lossy().into_owned(), - start_index, - record_index, - start_ms, - end_ms, - last_assistant_message: last_assistant_text(&subagent_records), - }) - }, - ) - .collect() + .collect::>(); + Ok(subagents) } fn import_envelope( diff --git a/bt-daemon/src/transcript_import/mod.rs b/bt-daemon/src/transcript_import/mod.rs index 5d045b6..ecbc5dc 100644 --- a/bt-daemon/src/transcript_import/mod.rs +++ b/bt-daemon/src/transcript_import/mod.rs @@ -428,6 +428,28 @@ fn read_jsonl_records(path: &Path) -> anyhow::Result> { .collect() } +fn read_complete_jsonl_records(path: &Path) -> anyhow::Result> { + let contents = + std::fs::read(path).with_context(|| format!("read transcript {}", path.display()))?; + let lines = contents.split(|byte| *byte == b'\n').collect::>(); + let mut records = Vec::new(); + for (index, line) in lines.iter().enumerate() { + if line.iter().all(u8::is_ascii_whitespace) { + continue; + } + match serde_json::from_slice(line) { + Ok(record) => records.push(record), + Err(_) if index + 1 == lines.len() => break, + Err(error) => { + return Err(error).with_context(|| { + format!("parse transcript {} line {}", path.display(), index + 1) + }); + } + } + } + Ok(records) +} + fn envelope( source: &str, source_version: Option, @@ -682,6 +704,127 @@ mod tests { ); } + #[test] + fn claude_import_rejects_partial_completed_subagent_transcript() { + let temp = tempfile::tempdir().unwrap(); + let transcript = temp.path().join("session-a.jsonl"); + let subagent = temp.path().join("session-a/subagents/agent-child-a.jsonl"); + std::fs::create_dir_all(subagent.parent().unwrap()).unwrap(); + let records = [ + json!({"type":"user","sessionId":"session-a","timestamp":"2026-01-01T00:00:01Z","message":{"content":"delegate"}}), + json!({"type":"assistant","sessionId":"session-a","timestamp":"2026-01-01T00:00:02Z","message":{"content":[{"type":"tool_use","id":"call-a","name":"Agent","input":{"prompt":"review this change"}}]}}), + json!({"type":"user","sessionId":"session-a","timestamp":"2026-01-01T00:00:05Z","toolUseResult":{"agentId":"child-a"},"message":{"content":[{"type":"tool_result","tool_use_id":"call-a","content":"done"}]}}), + ]; + std::fs::write( + &transcript, + records + .iter() + .map(Value::to_string) + .collect::>() + .join("\n"), + ) + .unwrap(); + std::fs::write( + &subagent, + format!( + "{}\n{{\"type\":", + json!({"type":"assistant","sessionId":"session-a","timestamp":"2026-01-01T00:00:04Z","message":{"content":[{"type":"text","text":"reviewed"}]}}) + ), + ) + .unwrap(); + + let error = transcript_envelopes(&transcript, ImportSource::Claude).unwrap_err(); + assert!(format!("{error:#}").contains("line 2")); + } + + #[test] + fn claude_import_keeps_interrupted_subagent_with_matching_agent_call() { + let temp = tempfile::tempdir().unwrap(); + let transcript = temp.path().join("session-a.jsonl"); + let subagent = temp.path().join("session-a/subagents/agent-child-a.jsonl"); + std::fs::create_dir_all(subagent.parent().unwrap()).unwrap(); + let records = [ + json!({"type":"user","sessionId":"session-a","timestamp":"2026-01-01T00:00:01Z","message":{"content":"delegate"}}), + json!({"type":"assistant","sessionId":"session-a","timestamp":"2026-01-01T00:00:02Z","message":{"content":[{"type":"tool_use","id":"call-a","name":"Agent","input":{"subagent_type":"reviewer","prompt":"review this change"}}]}}), + ]; + std::fs::write( + &transcript, + records + .iter() + .map(Value::to_string) + .collect::>() + .join("\n"), + ) + .unwrap(); + std::fs::write( + &subagent, + format!( + "{}\n{{\"type\":", + json!({"type":"user","isSidechain":true,"agentId":"child-a","sessionId":"session-a","timestamp":"2026-01-01T00:00:03Z","message":{"content":"review this change"}}) + ), + ) + .unwrap(); + + let events = transcript_envelopes(&transcript, ImportSource::Claude).unwrap(); + assert!(events.iter().any(|event| { + event.event == "SubagentStart" + && event.payload["agent_id"] == json!("child-a") + && event.payload["agent_type"] == json!("reviewer") + })); + assert!(events.iter().any(|event| { + event.event == "SubagentStop" + && event.payload["agent_transcript_path"] + .as_str() + .is_some_and(|path| Path::new(path) == subagent) + })); + } + + #[test] + fn claude_import_matches_interrupted_subagents_by_native_time() { + let temp = tempfile::tempdir().unwrap(); + let transcript = temp.path().join("session-a.jsonl"); + let child_dir = temp.path().join("session-a/subagents"); + std::fs::create_dir_all(&child_dir).unwrap(); + let records = [ + json!({"type":"user","sessionId":"session-a","timestamp":"2026-01-01T00:00:01Z","message":{"content":"delegate"}}), + json!({"type":"assistant","sessionId":"session-a","timestamp":"2026-01-01T00:00:02Z","message":{"content":[{"type":"tool_use","id":"call-early","name":"Agent","input":{"subagent_type":"reviewer","prompt":"same prompt"}}]}}), + json!({"type":"assistant","sessionId":"session-a","timestamp":"2026-01-01T00:00:10Z","message":{"content":[{"type":"tool_use","id":"call-late","name":"Agent","input":{"subagent_type":"reviewer","prompt":"same prompt"}}]}}), + ]; + std::fs::write( + &transcript, + records + .iter() + .map(Value::to_string) + .collect::>() + .join("\n"), + ) + .unwrap(); + for (agent_id, timestamp) in [ + ("z-early", "2026-01-01T00:00:03Z"), + ("a-late", "2026-01-01T00:00:11Z"), + ] { + std::fs::write( + child_dir.join(format!("agent-{agent_id}.jsonl")), + json!({"type":"user","isSidechain":true,"agentId":agent_id,"sessionId":"session-a","timestamp":timestamp,"message":{"content":"same prompt"}}).to_string(), + ) + .unwrap(); + } + + let events = transcript_envelopes(&transcript, ImportSource::Claude).unwrap(); + let starts = events + .iter() + .filter(|event| event.event == "SubagentStart") + .map(|event| (event.payload["agent_id"].as_str().unwrap(), event.ts_ms)) + .collect::>(); + assert_eq!( + starts, + vec![ + ("z-early", 1_767_225_602_000), + ("a-late", 1_767_225_610_000) + ] + ); + } + #[test] fn codex_import_discovers_spawned_rollout_and_emits_lifecycle() { let temp = tempfile::tempdir().unwrap(); diff --git a/bt-daemon/tests/replay.rs b/bt-daemon/tests/replay.rs index 86156e5..6549bcd 100644 --- a/bt-daemon/tests/replay.rs +++ b/bt-daemon/tests/replay.rs @@ -552,6 +552,67 @@ async fn imports_native_claude_transcript_with_multiple_turns_and_tools() { })); } +#[tokio::test] +async fn imports_claude_followups_as_sequential_turns_after_a_compact_summary() { + let tmp = tempfile::tempdir().unwrap(); + let transcript = tmp.path().join("claude.jsonl"); + write_jsonl( + &transcript, + &[ + json!({"type":"user","timestamp":"2026-01-01T00:00:01Z","sessionId":"claude-followup","cwd":"/tmp/demo","message":{"content":"first request"}}), + json!({"type":"assistant","timestamp":"2026-01-01T00:00:02Z","sessionId":"claude-followup","requestId":"first-response","message":{"id":"first-response","model":"claude-test","content":[{"type":"text","text":"first response"}],"usage":{"input_tokens":1,"output_tokens":1}}}), + json!({"type":"user","timestamp":"2026-01-01T00:00:03Z","sessionId":"claude-followup","isCompactSummary":true,"message":{"content":"internal summary"}}), + json!({"type":"user","timestamp":"2026-01-01T00:00:04Z","sessionId":"claude-followup","message":{"content":"follow up"}}), + json!({"type":"assistant","timestamp":"2026-01-01T00:00:05Z","sessionId":"claude-followup","requestId":"second-response","message":{"id":"second-response","model":"claude-test","content":[{"type":"text","text":"second response"}],"usage":{"input_tokens":1,"output_tokens":1}}}), + ], + ); + let subagents = tmp.path().join("claude/subagents"); + std::fs::create_dir_all(&subagents).unwrap(); + write_jsonl( + &subagents.join("agent-orphan-followup.jsonl"), + &[ + json!({"type":"user","isSidechain":true,"agentId":"orphan-followup","timestamp":"2026-01-01T00:00:04Z","sessionId":"claude-followup","message":{"content":"follow up"}}), + json!({"type":"assistant","isSidechain":true,"agentId":"orphan-followup","timestamp":"2026-01-01T00:00:04Z","sessionId":"claude-followup","message":{"content":[{"type":"text","text":"follow up"}]}}), + ], + ); + write_jsonl( + &subagents.join("agent-orphan-summary.jsonl"), + &[ + json!({"type":"user","isSidechain":true,"agentId":"orphan-summary","timestamp":"2026-01-01T00:00:05Z","sessionId":"claude-followup","message":{"content":"summarize the session"}}), + json!({"type":"assistant","isSidechain":true,"agentId":"orphan-summary","timestamp":"2026-01-01T00:00:05Z","sessionId":"claude-followup","message":{"content":[{"type":"text","text":"internal summary"}]}}), + ], + ); + + let output = tmp.path().join("spans"); + import_transcript( + &transcript, + ImportSource::Claude, + options(&output), + None, + false, + ) + .await + .unwrap(); + + let rows = rows(&output.join("claude-followup.ndjson")); + let task_inserts = rows + .iter() + .filter_map(|op| op.get("Insert")) + .filter(|row| row.get("span_type").and_then(Value::as_str) == Some("task")) + .collect::>(); + assert_eq!( + task_inserts.len(), + 3, + "unlinked child transcripts must not become subagents" + ); + assert!(task_inserts.iter().all(|row| !row["name"] + .as_str() + .unwrap_or_default() + .starts_with("subagent:"))); + assert!(task_inserts.iter().any(|row| row["name"] == "Turn 1")); + assert!(task_inserts.iter().any(|row| row["name"] == "Turn 2")); +} + #[tokio::test] async fn imports_non_monotonic_claude_records_into_their_native_turns() { let tmp = tempfile::tempdir().unwrap();