Skip to main content

dynamo_bench/coding/claude/
parser.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use crate::coding::common::{
5    anonymized_session_id, canonical_json_string, content_blocks, flatten_block_content_text,
6    object_field, parse_utc_timestamp_ms,
7};
8use crate::coding::tokenizer::TokenizerWorker;
9use anyhow::Result;
10use rustc_hash::FxHashMap;
11use serde_json::{Map, Value, json};
12use std::collections::{BTreeMap, BTreeSet};
13use std::fs::File;
14use std::io::{BufRead, BufReader};
15use std::path::PathBuf;
16
17#[derive(Clone, Debug)]
18pub struct TraceRecord {
19    pub session_id: String,
20    pub parent_session_id: Option<String>,
21    pub row_type: String,
22    pub timestamp_ms: i64,
23    pub source_order: u64,
24    pub raw: Value,
25}
26
27#[derive(Clone, Debug)]
28struct ConversationEntry {
29    kind: String,
30    rendered: String,
31}
32
33#[derive(Clone, Debug)]
34struct ToolCallSummary {
35    raw_id: Option<String>,
36    name: String,
37    normalized_id: Option<String>,
38    arg_size_chars: usize,
39    started_at_ms: i64,
40}
41
42#[derive(Clone, Debug)]
43struct CachedProgressMetrics {
44    progress_event_count: usize,
45    agent_ids: BTreeSet<String>,
46    assistant_text_blocks: usize,
47    tool_counts: BTreeMap<String, usize>,
48    tool_result_count: usize,
49    tool_error_count: usize,
50    tool_intervals: Vec<(i64, i64)>,
51    first_ts: i64,
52    last_ts: i64,
53}
54
55#[derive(Debug)]
56struct AssistantGroupSummary {
57    entries_by_record: BTreeMap<usize, Vec<ConversationEntry>>,
58    model: String,
59    output_length: usize,
60    assistant_text_blocks: usize,
61    top_level_tool_calls: Vec<ToolCallSummary>,
62    raw_task_tool_ids: Vec<String>,
63    input_length: Option<usize>,
64    cache_read_input_tokens: Option<usize>,
65    cache_creation_input_tokens: Option<usize>,
66    start_ms: i64,
67    end_ms: i64,
68}
69
70#[derive(Clone, Debug)]
71pub struct ToolDraft {
72    pub tool_call_id: String,
73    pub tool_class: String,
74    pub started_at_ms: i64,
75    pub ended_at_ms: i64,
76    pub is_error: bool,
77    pub output_bytes: usize,
78    pub child_session_id: Option<String>,
79    pub consumer_turn_index: Option<usize>,
80    pub execution_mode: String,
81}
82
83#[derive(Clone, Debug)]
84pub struct TurnDraft {
85    pub session_id: String,
86    pub source_request_id: String,
87    pub export_session_id: String,
88    pub export_parent_session_id: Option<String>,
89    pub turn_index: usize,
90    pub model: String,
91    pub input_text: String,
92    pub output_length: usize,
93    pub observed_input_length: Option<usize>,
94    pub cache_read_input_tokens: Option<usize>,
95    pub cache_creation_input_tokens: Option<usize>,
96    pub request_start_ms: i64,
97    pub assistant_start_ms: i64,
98    pub assistant_end_ms: i64,
99    pub delay_ms: Option<i64>,
100    pub tools: Vec<ToolDraft>,
101    pub sidecar: Value,
102    pub compaction: Option<CompactionMetadata>,
103}
104
105#[derive(Clone, Debug, Eq, PartialEq)]
106pub struct CompactionMetadata {
107    pub sequence: usize,
108    pub trigger: String,
109    pub pre_tokens: usize,
110    pub post_tokens: usize,
111    pub duration_ms: i64,
112    pub ended_at_ms: i64,
113}
114
115#[derive(Clone, Debug, Default)]
116pub struct SourceRequestExpectation {
117    pub input_length: Option<usize>,
118    pub cache_read_input_tokens: Option<usize>,
119    pub cache_creation_input_tokens: Option<usize>,
120    pub output_length: Option<usize>,
121    pub request_start_ms: i64,
122    pub assistant_end_ms: i64,
123}
124
125#[derive(Clone, Debug, Default)]
126pub struct SourceFidelityOracle {
127    pub requests: BTreeMap<(String, String), SourceRequestExpectation>,
128    pub compactions: BTreeMap<(String, String), CompactionMetadata>,
129    pub tools_by_class: BTreeMap<String, usize>,
130    pub paired_tools: usize,
131    pub tool_errors: usize,
132    pub child_links: usize,
133    pub background_tools: usize,
134    pub background_agents: usize,
135    pub background_completions_missing: usize,
136    pub background_titles: usize,
137    pub unmatched_tool_calls: usize,
138    pub unmatched_tool_results: usize,
139}
140
141#[derive(Debug)]
142struct PendingCompaction {
143    metadata: CompactionMetadata,
144    prompt_text: String,
145}
146
147#[derive(Debug, Default)]
148struct ToolIdNormalizer {
149    raw_to_normalized: FxHashMap<String, String>,
150}
151
152impl ToolIdNormalizer {
153    fn normalize(&mut self, raw_id: Option<&str>) -> Option<String> {
154        let raw_id = raw_id?;
155        if let Some(existing) = self.raw_to_normalized.get(raw_id) {
156            return Some(existing.clone());
157        }
158        let next_id = self.raw_to_normalized.len() + 1;
159        let normalized = format!("tool_{next_id:04}");
160        self.raw_to_normalized
161            .insert(raw_id.to_string(), normalized.clone());
162        Some(normalized)
163    }
164}
165
166#[derive(Debug)]
167pub struct SessionTurnBuilder {
168    session_id: String,
169    export_session_id: String,
170    export_parent_session_id: Option<String>,
171    records: Vec<TraceRecord>,
172    top_level_indices: Vec<usize>,
173    progress_metrics_index: FxHashMap<String, CachedProgressMetrics>,
174    request_index_by_group_key: FxHashMap<String, usize>,
175    request_start_ms_by_group_key: FxHashMap<String, i64>,
176    top_level_cursor: usize,
177    normalizer: ToolIdNormalizer,
178    conversation_entries: Vec<ConversationEntry>,
179    prompt_text: String,
180    pending_request_start_ms: Option<i64>,
181    previous_assistant_end_ms: Option<i64>,
182    turn_index: usize,
183    pending_compaction: Option<PendingCompaction>,
184    previous_model: Option<String>,
185    compaction_sequence: usize,
186    preserve_session_ids: bool,
187}
188
189impl SessionTurnBuilder {
190    pub fn new(trace_id: String, records: Vec<TraceRecord>, preserve_session_ids: bool) -> Self {
191        let root_session_id = records
192            .first()
193            .map(|record| record.session_id.clone())
194            .unwrap_or_else(|| trace_id.clone());
195        let is_subagent = trace_id != root_session_id;
196        let parent_session_id = records
197            .iter()
198            .find_map(|record| record.parent_session_id.clone())
199            .unwrap_or_else(|| root_session_id.clone());
200        let export_session_id = if preserve_session_ids {
201            trace_id.clone()
202        } else {
203            anonymized_session_id(&trace_id)
204        };
205        let export_parent_session_id = is_subagent.then(|| {
206            if preserve_session_ids {
207                parent_session_id
208            } else {
209                anonymized_session_id(&parent_session_id)
210            }
211        });
212
213        let progress_index = build_progress_index(&records);
214        let progress_metrics_index = build_progress_metrics_index(&progress_index, &records);
215        let top_level_indices: Vec<usize> = records
216            .iter()
217            .enumerate()
218            .filter_map(|(index, record)| {
219                let is_top_level =
220                    matches!(record.row_type.as_str(), "user" | "assistant" | "system")
221                        && (is_subagent
222                            || !record
223                                .raw
224                                .get("isSidechain")
225                                .and_then(Value::as_bool)
226                                .unwrap_or(false));
227                is_top_level.then_some(index)
228            })
229            .collect();
230        let mut request_index_by_group_key = FxHashMap::default();
231        let mut request_start_ms_by_group_key = FxHashMap::default();
232        for index in &top_level_indices {
233            let record = &records[*index];
234            if record.row_type != "assistant" {
235                continue;
236            }
237            let group_key = assistant_group_key(record);
238            let next_index = request_index_by_group_key.len();
239            request_index_by_group_key
240                .entry(group_key.clone())
241                .or_insert(next_index);
242            request_start_ms_by_group_key
243                .entry(group_key)
244                .and_modify(|start: &mut i64| *start = (*start).min(record.timestamp_ms))
245                .or_insert(record.timestamp_ms);
246        }
247
248        Self {
249            session_id: trace_id,
250            export_session_id,
251            export_parent_session_id,
252            records,
253            top_level_indices,
254            progress_metrics_index,
255            request_index_by_group_key,
256            request_start_ms_by_group_key,
257            top_level_cursor: 0,
258            normalizer: ToolIdNormalizer::default(),
259            conversation_entries: Vec::new(),
260            prompt_text: String::new(),
261            pending_request_start_ms: None,
262            previous_assistant_end_ms: None,
263            turn_index: 0,
264            pending_compaction: None,
265            previous_model: None,
266            compaction_sequence: 0,
267            preserve_session_ids,
268        }
269    }
270
271    pub fn next_turn(&mut self, tokenizer: &mut impl TokenizerWorker) -> Result<Option<TurnDraft>> {
272        while self.top_level_cursor < self.top_level_indices.len() {
273            let record_index = self.top_level_indices[self.top_level_cursor];
274            let record = &self.records[record_index];
275
276            if record.row_type == "system" {
277                if is_compact_boundary(record) {
278                    let metadata = compaction_metadata(record, self.compaction_sequence)?;
279                    self.compaction_sequence += 1;
280                    self.pending_compaction = Some(PendingCompaction {
281                        metadata,
282                        prompt_text: self.prompt_text.clone(),
283                    });
284                }
285                self.top_level_cursor += 1;
286                continue;
287            }
288
289            if record.row_type == "user" {
290                if should_skip_user_record(record)? {
291                    self.top_level_cursor += 1;
292                    continue;
293                }
294
295                let request_start_ms = record.timestamp_ms;
296                let message = object_field(&record.raw, "message");
297                let rendered_entries = render_user_entries(message, &mut self.normalizer)?;
298                if is_compact_summary(record) {
299                    let summary_text = flatten_block_content_text(
300                        message
301                            .and_then(|message| message.get("content"))
302                            .unwrap_or(&Value::Null),
303                    )?;
304                    self.replace_conversation_entries(rendered_entries);
305                    self.pending_request_start_ms = Some(request_start_ms);
306                    self.top_level_cursor += 1;
307
308                    let Some(pending) = self.pending_compaction.take() else {
309                        continue;
310                    };
311
312                    let output_length = tokenizer.encode(&summary_text)?.len();
313                    let input_text = if pending.prompt_text.is_empty() {
314                        "[system] Compact the conversation.".to_string()
315                    } else {
316                        format!(
317                            "{}\n[system] Compact the conversation.",
318                            pending.prompt_text
319                        )
320                    };
321                    let source_request_id = format!("compact:{}", pending.metadata.sequence);
322                    let mut sidecar = Map::new();
323                    sidecar.insert(
324                        "session_id".to_string(),
325                        Value::String(self.export_session_id.clone()),
326                    );
327                    if let Some(parent_session_id) = &self.export_parent_session_id {
328                        sidecar.insert(
329                            "parent_session_id".to_string(),
330                            Value::String(parent_session_id.clone()),
331                        );
332                    }
333                    sidecar.insert("turn_index".to_string(), json!(self.turn_index));
334                    sidecar.insert(
335                        "source_request_id".to_string(),
336                        Value::String(source_request_id.clone()),
337                    );
338                    sidecar.insert("request_kind".to_string(), json!("compaction"));
339                    sidecar.insert(
340                        "input_fidelity".to_string(),
341                        json!("claude_cache_safe_fork"),
342                    );
343                    sidecar.insert(
344                        "replay_hash_fidelity".to_string(),
345                        json!("synthetic_usage_shaped"),
346                    );
347                    sidecar.insert(
348                        "compaction".to_string(),
349                        compaction_json(&pending.metadata, output_length),
350                    );
351                    let request_start_ms = pending
352                        .metadata
353                        .ended_at_ms
354                        .saturating_sub(pending.metadata.duration_ms);
355                    self.previous_assistant_end_ms = Some(pending.metadata.ended_at_ms);
356                    return Ok(Some(TurnDraft {
357                        session_id: self.session_id.clone(),
358                        source_request_id,
359                        export_session_id: self.export_session_id.clone(),
360                        export_parent_session_id: self.export_parent_session_id.clone(),
361                        turn_index: self.turn_index,
362                        model: self
363                            .previous_model
364                            .clone()
365                            .unwrap_or_else(|| "unknown".to_string()),
366                        input_text,
367                        output_length,
368                        observed_input_length: Some(pending.metadata.pre_tokens),
369                        cache_read_input_tokens: None,
370                        cache_creation_input_tokens: None,
371                        request_start_ms,
372                        assistant_start_ms: pending.metadata.ended_at_ms,
373                        assistant_end_ms: pending.metadata.ended_at_ms,
374                        delay_ms: None,
375                        tools: Vec::new(),
376                        sidecar: Value::Object(sidecar),
377                        compaction: Some(pending.metadata),
378                    }));
379                } else {
380                    self.pending_compaction = None;
381                    self.extend_conversation_entries(rendered_entries);
382                }
383                self.pending_request_start_ms = Some(request_start_ms);
384                self.top_level_cursor += 1;
385                continue;
386            }
387
388            self.pending_compaction = None;
389            let group_key = assistant_group_key(record);
390            let mut group_indices = vec![record_index];
391            let mut interleaved_user_indices = Vec::new();
392            self.top_level_cursor += 1;
393            while self.top_level_cursor < self.top_level_indices.len() {
394                let next_index = self.top_level_indices[self.top_level_cursor];
395                let next_record = &self.records[next_index];
396                if next_record.row_type == "system" && !is_compact_boundary(next_record) {
397                    self.top_level_cursor += 1;
398                    continue;
399                }
400                if next_record.row_type == "user" && is_tool_result_user_record(next_record) {
401                    interleaved_user_indices.push(next_index);
402                    self.top_level_cursor += 1;
403                    continue;
404                }
405                if next_record.row_type != "assistant"
406                    || assistant_group_key(next_record) != group_key
407                {
408                    break;
409                }
410                group_indices.push(next_index);
411                self.top_level_cursor += 1;
412            }
413
414            let mut group_summary = summarize_assistant_group(
415                &self.records,
416                &group_indices,
417                &mut self.normalizer,
418                tokenizer,
419            )?;
420            let input_text = self.prompt_text.clone();
421            let request_start_ms = self
422                .pending_request_start_ms
423                .take()
424                .unwrap_or(group_summary.start_ms);
425            let tools = pair_tool_results(
426                &group_summary.top_level_tool_calls,
427                &interleaved_user_indices,
428                &self.records,
429                &self.request_index_by_group_key,
430                &self.request_start_ms_by_group_key,
431                &group_key,
432                self.preserve_session_ids,
433            )?;
434
435            let top_level_tool_names = group_summary
436                .top_level_tool_calls
437                .iter()
438                .map(|tool_call| tool_call.name.clone())
439                .collect::<Vec<_>>();
440            let used_task_tool = top_level_tool_names.iter().any(|name| name == "Task");
441            let top_level_tool_calls = group_summary
442                .top_level_tool_calls
443                .iter()
444                .map(|tool_call| {
445                    json!({
446                        "name": tool_call.name,
447                        "tool_id": tool_call.normalized_id,
448                        "arg_size_chars": tool_call.arg_size_chars,
449                    })
450                })
451                .collect::<Vec<_>>();
452
453            let mut sidecar = Map::new();
454            sidecar.insert(
455                "session_id".to_string(),
456                Value::String(self.export_session_id.clone()),
457            );
458            if let Some(parent_session_id) = &self.export_parent_session_id {
459                sidecar.insert(
460                    "parent_session_id".to_string(),
461                    Value::String(parent_session_id.clone()),
462                );
463            }
464            sidecar.insert("turn_index".to_string(), json!(self.turn_index));
465            sidecar.insert(
466                "source_request_id".to_string(),
467                Value::String(group_key.clone()),
468            );
469            sidecar.insert(
470                "num_messages_in_context".to_string(),
471                json!(self.conversation_entries.len()),
472            );
473            sidecar.insert(
474                "context_shape".to_string(),
475                Value::Array(
476                    self.conversation_entries
477                        .iter()
478                        .map(|entry| Value::String(entry.kind.clone()))
479                        .collect(),
480                ),
481            );
482            sidecar.insert(
483                "tool_rounds_before_answer".to_string(),
484                json!(count_trailing_tool_results(&self.conversation_entries)),
485            );
486            sidecar.insert("used_task_tool".to_string(), Value::Bool(used_task_tool));
487            sidecar.insert(
488                "assistant_text_blocks".to_string(),
489                json!(group_summary.assistant_text_blocks),
490            );
491            sidecar.insert(
492                "top_level_tool_call_count".to_string(),
493                json!(group_summary.top_level_tool_calls.len()),
494            );
495            sidecar.insert(
496                "top_level_tool_names".to_string(),
497                Value::Array(
498                    top_level_tool_names
499                        .iter()
500                        .cloned()
501                        .map(Value::String)
502                        .collect(),
503                ),
504            );
505            sidecar.insert(
506                "top_level_tool_calls".to_string(),
507                Value::Array(top_level_tool_calls),
508            );
509            sidecar.insert(
510                "input_fidelity".to_string(),
511                Value::String(
512                    if group_summary.input_length.is_some() {
513                        "claude_usage_cache_prefix"
514                    } else {
515                        "rendered_transcript"
516                    }
517                    .to_string(),
518                ),
519            );
520            sidecar.insert(
521                "replay_hash_fidelity".to_string(),
522                Value::String(
523                    if group_summary.input_length.is_some() {
524                        "synthetic_usage_shaped"
525                    } else {
526                        "rendered_transcript"
527                    }
528                    .to_string(),
529                ),
530            );
531            if let Some(input_length) = group_summary.input_length {
532                sidecar.insert("observed_input_tokens".to_string(), json!(input_length));
533            }
534            if let Some(cache_read) = group_summary.cache_read_input_tokens {
535                sidecar.insert(
536                    "observed_cache_read_input_tokens".to_string(),
537                    json!(cache_read),
538                );
539            }
540            if let Some(cache_creation) = group_summary.cache_creation_input_tokens {
541                sidecar.insert(
542                    "observed_cache_creation_input_tokens".to_string(),
543                    json!(cache_creation),
544                );
545            }
546
547            let progress_metrics = aggregate_progress_metrics(
548                &group_summary.raw_task_tool_ids,
549                &self.progress_metrics_index,
550                &mut self.normalizer,
551            );
552            if let Some(progress_map) = progress_metrics.as_object() {
553                for (key, value) in progress_map {
554                    sidecar.insert(key.clone(), value.clone());
555                }
556            }
557
558            let model = group_summary.model;
559            self.previous_model = Some(model.clone());
560            let turn = TurnDraft {
561                session_id: self.session_id.clone(),
562                source_request_id: group_key,
563                export_session_id: self.export_session_id.clone(),
564                export_parent_session_id: self.export_parent_session_id.clone(),
565                turn_index: self.turn_index,
566                model,
567                input_text,
568                output_length: group_summary.output_length,
569                observed_input_length: group_summary.input_length,
570                cache_read_input_tokens: group_summary.cache_read_input_tokens,
571                cache_creation_input_tokens: group_summary.cache_creation_input_tokens,
572                request_start_ms,
573                assistant_start_ms: group_summary.start_ms,
574                assistant_end_ms: group_summary.end_ms,
575                delay_ms: self
576                    .previous_assistant_end_ms
577                    .map(|previous_end| (request_start_ms - previous_end).max(0)),
578                tools,
579                sidecar: Value::Object(sidecar),
580                compaction: None,
581            };
582
583            let mut ordered_indices = group_indices;
584            ordered_indices.extend(interleaved_user_indices.iter().copied());
585            ordered_indices.sort_unstable();
586            let mut ordered_entries = Vec::new();
587            for index in ordered_indices {
588                if let Some(entries) = group_summary.entries_by_record.remove(&index) {
589                    ordered_entries.extend(entries);
590                } else {
591                    ordered_entries.extend(render_user_entries(
592                        object_field(&self.records[index].raw, "message"),
593                        &mut self.normalizer,
594                    )?);
595                }
596            }
597            self.extend_conversation_entries(ordered_entries);
598            self.pending_request_start_ms = interleaved_user_indices
599                .last()
600                .map(|index| self.records[*index].timestamp_ms);
601            self.previous_assistant_end_ms = Some(group_summary.end_ms);
602            self.turn_index += 1;
603            return Ok(Some(turn));
604        }
605
606        Ok(None)
607    }
608
609    fn replace_conversation_entries(&mut self, entries: Vec<ConversationEntry>) {
610        self.prompt_text = render_entry_buffer(&entries);
611        self.conversation_entries = entries;
612    }
613
614    fn extend_conversation_entries(&mut self, entries: Vec<ConversationEntry>) {
615        append_rendered_entries(&mut self.prompt_text, &entries);
616        self.conversation_entries.extend(entries);
617    }
618}
619
620pub fn load_trace_records(trace_files: &[PathBuf]) -> Result<FxHashMap<String, Vec<TraceRecord>>> {
621    let mut sessions: FxHashMap<String, Vec<TraceRecord>> = FxHashMap::default();
622    let mut source_order = 0_u64;
623
624    for trace_file in trace_files {
625        let file = File::open(trace_file)?;
626        let reader = BufReader::new(file);
627        for (line_number, line) in reader.lines().enumerate() {
628            let line = line?;
629            if line.trim().is_empty() {
630                continue;
631            }
632            let payload: Value = serde_json::from_str(&line).map_err(|error| {
633                anyhow::anyhow!(
634                    "invalid JSON in {}:{}: {}",
635                    trace_file.display(),
636                    line_number + 1,
637                    error
638                )
639            })?;
640
641            let session_id = payload
642                .get("sessionId")
643                .and_then(Value::as_str)
644                .map(str::to_string);
645            let agent_id = payload
646                .get("agentId")
647                .and_then(Value::as_str)
648                .map(str::to_string);
649            let row_type = payload
650                .get("type")
651                .and_then(Value::as_str)
652                .map(str::to_string);
653            let (Some(session_id), Some(row_type)) = (session_id, row_type) else {
654                source_order += 1;
655                continue;
656            };
657            let timestamp_ms = match payload.get("timestamp").and_then(Value::as_str) {
658                Some(timestamp) => match parse_utc_timestamp_ms(timestamp) {
659                    Ok(timestamp_ms) => timestamp_ms,
660                    Err(_) => {
661                        source_order += 1;
662                        continue;
663                    }
664                },
665                None if row_type == "ai-title" => 0,
666                None => {
667                    source_order += 1;
668                    continue;
669                }
670            };
671
672            let trace_id = agent_id.clone().unwrap_or_else(|| session_id.clone());
673            sessions.entry(trace_id).or_default().push(TraceRecord {
674                session_id,
675                parent_session_id: None,
676                row_type,
677                timestamp_ms,
678                source_order,
679                raw: payload,
680            });
681            source_order += 1;
682        }
683    }
684
685    let mut parent_by_session = FxHashMap::default();
686    for (parent_session_id, records) in &sessions {
687        for record in records {
688            if let Some(child_session_id) = record
689                .raw
690                .get("toolUseResult")
691                .and_then(Value::as_object)
692                .and_then(|result| result.get("agentId"))
693                .and_then(Value::as_str)
694            {
695                parent_by_session
696                    .entry(child_session_id.to_string())
697                    .or_insert_with(|| parent_session_id.clone());
698            }
699        }
700    }
701
702    for (session_id, records) in &mut sessions {
703        let parent_session_id = parent_by_session.get(session_id).cloned();
704        for record in records.iter_mut() {
705            record.parent_session_id.clone_from(&parent_session_id);
706        }
707        records.sort_by_key(|record| record.source_order);
708    }
709
710    Ok(sessions)
711}
712
713pub fn build_source_fidelity_oracle(
714    sessions: &FxHashMap<String, Vec<TraceRecord>>,
715) -> Result<SourceFidelityOracle> {
716    let mut oracle = SourceFidelityOracle::default();
717    let mut tool_calls: FxHashMap<(String, String), String> = FxHashMap::default();
718    let mut background_titles = BTreeSet::new();
719    let mut background_tool_ids = BTreeSet::new();
720    let background_completions = sessions
721        .iter()
722        .flat_map(|(trace_id, records)| {
723            records.iter().filter_map(|record| {
724                let content = (record.row_type == "queue-operation"
725                    && record.raw.get("operation").and_then(Value::as_str) == Some("enqueue"))
726                .then(|| record.raw.get("content").and_then(Value::as_str))
727                .flatten()?;
728                Some((
729                    (trace_id.clone(), queued_tool_id(content)?.to_string()),
730                    !content.contains("<status>completed</status>"),
731                ))
732            })
733        })
734        .collect::<FxHashMap<_, _>>();
735
736    for (trace_id, records) in sessions {
737        let root_session_id = records
738            .first()
739            .map(|record| record.session_id.as_str())
740            .unwrap_or(trace_id);
741        let is_subagent = trace_id != root_session_id;
742        let mut pending_request_start_ms = None;
743        let mut compaction_sequence = 0;
744
745        for record in records {
746            if record.row_type == "ai-title" {
747                let title = record
748                    .raw
749                    .get("aiTitle")
750                    .and_then(Value::as_str)
751                    .unwrap_or_default();
752                background_titles.insert((record.session_id.clone(), title.to_string()));
753                continue;
754            }
755            let is_top_level = matches!(record.row_type.as_str(), "user" | "assistant" | "system")
756                && (is_subagent
757                    || !record
758                        .raw
759                        .get("isSidechain")
760                        .and_then(Value::as_bool)
761                        .unwrap_or(false));
762            if !is_top_level {
763                continue;
764            }
765            if is_compact_boundary(record) {
766                let metadata = compaction_metadata(record, compaction_sequence)?;
767                let source_request_id = format!("compact:{}", metadata.sequence);
768                oracle
769                    .compactions
770                    .insert((trace_id.clone(), source_request_id), metadata);
771                compaction_sequence += 1;
772                continue;
773            }
774            if record.row_type == "system" {
775                continue;
776            }
777
778            if record.row_type == "user" {
779                if should_skip_user_record(record)? {
780                    continue;
781                }
782                pending_request_start_ms = Some(record.timestamp_ms);
783                let message = object_field(&record.raw, "message");
784                for block in content_blocks(message.and_then(|message| message.get("content"))) {
785                    if block.get("type").and_then(Value::as_str) != Some("tool_result") {
786                        continue;
787                    }
788                    let Some(raw_id) = block.get("tool_use_id").and_then(Value::as_str) else {
789                        continue;
790                    };
791                    let key = (trace_id.clone(), raw_id.to_string());
792                    let Some(tool_class) = tool_calls.remove(&key) else {
793                        oracle.unmatched_tool_results += 1;
794                        continue;
795                    };
796                    oracle.paired_tools += 1;
797                    *oracle.tools_by_class.entry(tool_class).or_insert(0) += 1;
798                    let launch_error = block
799                        .get("is_error")
800                        .and_then(Value::as_bool)
801                        .unwrap_or(false);
802                    let mut is_async = false;
803                    if let Some(result) = record.raw.get("toolUseResult").and_then(Value::as_object)
804                    {
805                        is_async = result
806                            .get("isAsync")
807                            .and_then(Value::as_bool)
808                            .unwrap_or(false)
809                            || result.get("backgroundTaskId").is_some();
810                        if result.get("agentId").and_then(Value::as_str).is_some() {
811                            oracle.child_links += 1;
812                            oracle.background_agents += usize::from(is_async);
813                        }
814                        if is_async {
815                            oracle.background_tools += 1;
816                            background_tool_ids.insert(key);
817                        }
818                    }
819                    oracle.tool_errors += usize::from(if is_async {
820                        background_completions
821                            .get(&(trace_id.clone(), raw_id.to_string()))
822                            .copied()
823                            .unwrap_or(launch_error)
824                    } else {
825                        launch_error
826                    });
827                }
828                continue;
829            }
830
831            let group_key = assistant_group_key(record);
832            let request_key = (trace_id.clone(), group_key);
833            if !oracle.requests.contains_key(&request_key) {
834                oracle.requests.insert(
835                    request_key.clone(),
836                    SourceRequestExpectation {
837                        request_start_ms: pending_request_start_ms
838                            .take()
839                            .unwrap_or(record.timestamp_ms),
840                        assistant_end_ms: record.timestamp_ms,
841                        ..Default::default()
842                    },
843                );
844            }
845            let expectation = oracle
846                .requests
847                .get_mut(&request_key)
848                .expect("request expectation was inserted");
849            expectation.assistant_end_ms = expectation.assistant_end_ms.max(record.timestamp_ms);
850            let Some(message) = object_field(&record.raw, "message") else {
851                continue;
852            };
853            if let Some(usage) = object_field(&Value::Object(message.clone()), "usage") {
854                let input_length = [
855                    "input_tokens",
856                    "cache_creation_input_tokens",
857                    "cache_read_input_tokens",
858                ]
859                .into_iter()
860                .filter_map(|key| usage.get(key).and_then(Value::as_u64))
861                .fold(0_usize, |total, value| total.saturating_add(value as usize));
862                expectation.input_length = Some(
863                    expectation
864                        .input_length
865                        .unwrap_or_default()
866                        .max(input_length),
867                );
868                let cache_read = usage
869                    .get("cache_read_input_tokens")
870                    .and_then(Value::as_u64)
871                    .unwrap_or(0) as usize;
872                expectation.cache_read_input_tokens = Some(
873                    expectation
874                        .cache_read_input_tokens
875                        .unwrap_or_default()
876                        .max(cache_read),
877                );
878                let cache_creation = usage
879                    .get("cache_creation_input_tokens")
880                    .and_then(Value::as_u64)
881                    .unwrap_or(0) as usize;
882                expectation.cache_creation_input_tokens = Some(
883                    expectation
884                        .cache_creation_input_tokens
885                        .unwrap_or_default()
886                        .max(cache_creation),
887                );
888                if let Some(output_length) = usage.get("output_tokens").and_then(Value::as_u64) {
889                    expectation.output_length = Some(
890                        expectation
891                            .output_length
892                            .unwrap_or_default()
893                            .max(output_length as usize),
894                    );
895                }
896            }
897            for block in content_blocks(message.get("content")) {
898                if block.get("type").and_then(Value::as_str) != Some("tool_use") {
899                    continue;
900                }
901                let Some(raw_id) = block.get("id").and_then(Value::as_str) else {
902                    continue;
903                };
904                let tool_class = block
905                    .get("name")
906                    .and_then(Value::as_str)
907                    .unwrap_or("unknown")
908                    .to_string();
909                tool_calls.insert((trace_id.clone(), raw_id.to_string()), tool_class);
910            }
911        }
912    }
913
914    oracle.background_titles = background_titles.len();
915    oracle.background_completions_missing = background_tool_ids
916        .iter()
917        .filter(|tool_id| !background_completions.contains_key(*tool_id))
918        .count();
919    oracle.unmatched_tool_calls = tool_calls.len();
920    Ok(oracle)
921}
922
923fn queued_tool_id(content: &str) -> Option<&str> {
924    let start = content.find("<tool-use-id>")? + "<tool-use-id>".len();
925    let end = content[start..].find("</tool-use-id>")? + start;
926    Some(&content[start..end])
927}
928
929pub(crate) fn assistant_group_key(record: &TraceRecord) -> String {
930    if let Some(request_id) = record.raw.get("requestId").and_then(Value::as_str) {
931        return request_id.to_string();
932    }
933    if let Some(message_id) = object_field(&record.raw, "message")
934        .and_then(|message| message.get("id"))
935        .and_then(Value::as_str)
936    {
937        return message_id.to_string();
938    }
939    if let Some(uuid) = record.raw.get("uuid").and_then(Value::as_str) {
940        return uuid.to_string();
941    }
942    format!("row-{}", record.source_order)
943}
944
945fn is_compact_boundary(record: &TraceRecord) -> bool {
946    record.row_type == "system"
947        && record
948            .raw
949            .get("subtype")
950            .and_then(Value::as_str)
951            .map(|subtype| subtype == "compact_boundary")
952            .unwrap_or(false)
953}
954
955fn is_compact_summary(record: &TraceRecord) -> bool {
956    record.row_type == "user"
957        && record
958            .raw
959            .get("isCompactSummary")
960            .and_then(Value::as_bool)
961            .unwrap_or(false)
962}
963
964fn compaction_metadata(record: &TraceRecord, sequence: usize) -> Result<CompactionMetadata> {
965    let metadata = record
966        .raw
967        .get("compactMetadata")
968        .and_then(Value::as_object)
969        .ok_or_else(|| anyhow::anyhow!("compact_boundary is missing compactMetadata"))?;
970    let trigger = metadata
971        .get("trigger")
972        .and_then(Value::as_str)
973        .ok_or_else(|| anyhow::anyhow!("compactMetadata is missing trigger"))?
974        .to_string();
975    let pre_tokens = metadata
976        .get("preTokens")
977        .and_then(Value::as_u64)
978        .and_then(|value| usize::try_from(value).ok())
979        .ok_or_else(|| anyhow::anyhow!("compactMetadata has invalid preTokens"))?;
980    if pre_tokens == 0 {
981        anyhow::bail!("compactMetadata preTokens leaves no recoverable prefix");
982    }
983    let post_tokens = metadata
984        .get("postTokens")
985        .and_then(Value::as_u64)
986        .and_then(|value| usize::try_from(value).ok())
987        .ok_or_else(|| anyhow::anyhow!("compactMetadata has invalid postTokens"))?;
988    let duration_ms = metadata
989        .get("durationMs")
990        .and_then(Value::as_u64)
991        .and_then(|value| i64::try_from(value).ok())
992        .ok_or_else(|| anyhow::anyhow!("compactMetadata has invalid durationMs"))?;
993    Ok(CompactionMetadata {
994        sequence,
995        trigger,
996        pre_tokens,
997        post_tokens,
998        duration_ms,
999        ended_at_ms: record.timestamp_ms,
1000    })
1001}
1002
1003fn compaction_json(metadata: &CompactionMetadata, summary_output_tokens: usize) -> Value {
1004    json!({
1005        "trigger": metadata.trigger,
1006        "pre_tokens": metadata.pre_tokens,
1007        "post_tokens": metadata.post_tokens,
1008        "duration_ms": metadata.duration_ms,
1009        "summary_output_tokens": summary_output_tokens,
1010        "cache_fidelity": "recoverable_cache_safe_prefix",
1011        "output_fidelity": "tokenized_compact_summary",
1012    })
1013}
1014
1015fn is_local_command_wrapper_text(text: &str) -> bool {
1016    let stripped = text.trim();
1017    [
1018        "<command-name>",
1019        "<command-message>",
1020        "<command-args>",
1021        "<local-command-caveat>",
1022        "<local-command-stdout>",
1023        "<local-command-stderr>",
1024    ]
1025    .iter()
1026    .any(|prefix| stripped.starts_with(prefix))
1027}
1028
1029fn should_skip_user_record(record: &TraceRecord) -> Result<bool> {
1030    if record.row_type != "user" {
1031        return Ok(false);
1032    }
1033    if record
1034        .raw
1035        .get("isMeta")
1036        .and_then(Value::as_bool)
1037        .unwrap_or(false)
1038    {
1039        return Ok(true);
1040    }
1041
1042    let Some(message) = object_field(&record.raw, "message") else {
1043        return Ok(false);
1044    };
1045    let blocks = content_blocks(message.get("content"));
1046    if blocks.is_empty() {
1047        return Ok(false);
1048    }
1049    if blocks
1050        .iter()
1051        .any(|block| block.get("type").and_then(Value::as_str) != Some("text"))
1052    {
1053        return Ok(false);
1054    }
1055
1056    let texts = blocks
1057        .iter()
1058        .map(|block| {
1059            block
1060                .get("text")
1061                .and_then(Value::as_str)
1062                .unwrap_or_default()
1063                .to_string()
1064        })
1065        .collect::<Vec<_>>();
1066    Ok(!texts.is_empty() && texts.iter().all(|text| is_local_command_wrapper_text(text)))
1067}
1068
1069fn is_tool_result_user_record(record: &TraceRecord) -> bool {
1070    if record.row_type != "user" {
1071        return false;
1072    }
1073    let Some(message) = object_field(&record.raw, "message") else {
1074        return false;
1075    };
1076    let blocks = content_blocks(message.get("content"));
1077    !blocks.is_empty()
1078        && blocks
1079            .iter()
1080            .all(|block| block.get("type").and_then(Value::as_str) == Some("tool_result"))
1081}
1082
1083fn pair_tool_results(
1084    calls: &[ToolCallSummary],
1085    user_indices: &[usize],
1086    records: &[TraceRecord],
1087    request_index_by_group_key: &FxHashMap<String, usize>,
1088    request_start_ms_by_group_key: &FxHashMap<String, i64>,
1089    current_group_key: &str,
1090    preserve_session_ids: bool,
1091) -> Result<Vec<ToolDraft>> {
1092    let calls_by_id = calls
1093        .iter()
1094        .filter_map(|call| call.raw_id.as_deref().map(|id| (id, call)))
1095        .collect::<FxHashMap<_, _>>();
1096    let mut tools = Vec::new();
1097
1098    for index in user_indices {
1099        let record = &records[*index];
1100        let Some(message) = object_field(&record.raw, "message") else {
1101            continue;
1102        };
1103        for block in content_blocks(message.get("content")) {
1104            if block.get("type").and_then(Value::as_str) != Some("tool_result") {
1105                continue;
1106            }
1107            let Some(raw_id) = block.get("tool_use_id").and_then(Value::as_str) else {
1108                continue;
1109            };
1110            let Some(call) = calls_by_id.get(raw_id) else {
1111                continue;
1112            };
1113            let content = flatten_block_content_text(block.get("content").unwrap_or(&Value::Null))?;
1114            let tool_result = record.raw.get("toolUseResult").and_then(Value::as_object);
1115            let is_async = tool_result
1116                .and_then(|result| result.get("isAsync"))
1117                .and_then(Value::as_bool)
1118                .unwrap_or(false)
1119                || tool_result.is_some_and(|result| result.get("backgroundTaskId").is_some());
1120            let child_session_id = tool_result
1121                .and_then(|result| result.get("agentId"))
1122                .and_then(Value::as_str)
1123                .map(|id| {
1124                    if preserve_session_ids {
1125                        id.to_string()
1126                    } else {
1127                        anonymized_session_id(id)
1128                    }
1129                });
1130            let launch_error = block
1131                .get("is_error")
1132                .and_then(Value::as_bool)
1133                .unwrap_or(false);
1134            let (ended_at_ms, output_bytes, is_error, consumer_turn_index) = if is_async {
1135                async_tool_completion(
1136                    records,
1137                    raw_id,
1138                    record.timestamp_ms,
1139                    request_index_by_group_key,
1140                    request_start_ms_by_group_key,
1141                )
1142                .unwrap_or((
1143                    record.timestamp_ms,
1144                    content.len(),
1145                    launch_error,
1146                    None,
1147                ))
1148            } else {
1149                (
1150                    record.timestamp_ms,
1151                    content.len(),
1152                    launch_error,
1153                    next_consumer_turn(
1154                        records,
1155                        record.source_order,
1156                        current_group_key,
1157                        request_index_by_group_key,
1158                    ),
1159                )
1160            };
1161            tools.push(ToolDraft {
1162                tool_call_id: call
1163                    .normalized_id
1164                    .clone()
1165                    .unwrap_or_else(|| "tool_unknown".to_string()),
1166                tool_class: call.name.clone(),
1167                started_at_ms: call.started_at_ms,
1168                ended_at_ms,
1169                is_error,
1170                output_bytes,
1171                child_session_id,
1172                consumer_turn_index,
1173                execution_mode: if is_async {
1174                    "background".to_string()
1175                } else {
1176                    "blocking".to_string()
1177                },
1178            });
1179        }
1180    }
1181
1182    Ok(tools)
1183}
1184
1185fn next_consumer_turn(
1186    records: &[TraceRecord],
1187    after_source_order: u64,
1188    current_group_key: &str,
1189    request_index_by_group_key: &FxHashMap<String, usize>,
1190) -> Option<usize> {
1191    records
1192        .iter()
1193        .filter(|record| record.source_order > after_source_order && record.row_type == "assistant")
1194        .find_map(|record| {
1195            let group_key = assistant_group_key(record);
1196            (group_key != current_group_key)
1197                .then(|| request_index_by_group_key.get(&group_key).copied())
1198                .flatten()
1199        })
1200}
1201
1202fn async_tool_completion(
1203    records: &[TraceRecord],
1204    raw_tool_id: &str,
1205    after_timestamp_ms: i64,
1206    request_index_by_group_key: &FxHashMap<String, usize>,
1207    request_start_ms_by_group_key: &FxHashMap<String, i64>,
1208) -> Option<(i64, usize, bool, Option<usize>)> {
1209    let tool_marker = format!("<tool-use-id>{raw_tool_id}</tool-use-id>");
1210    let completion = records
1211        .iter()
1212        .filter(|record| {
1213            record.timestamp_ms >= after_timestamp_ms
1214                && record.row_type == "queue-operation"
1215                && record.raw.get("operation").and_then(Value::as_str) == Some("enqueue")
1216                && record
1217                    .raw
1218                    .get("content")
1219                    .and_then(Value::as_str)
1220                    .is_some_and(|content| content.contains(&tool_marker))
1221        })
1222        .min_by_key(|record| (record.timestamp_ms, record.source_order))?;
1223    let content = completion
1224        .raw
1225        .get("content")
1226        .and_then(Value::as_str)
1227        .unwrap_or_default();
1228    let consumer_turn_index = request_index_by_group_key
1229        .iter()
1230        .filter_map(|(group_key, turn_index)| {
1231            let start_ms = *request_start_ms_by_group_key.get(group_key)?;
1232            (start_ms > completion.timestamp_ms).then_some((start_ms, *turn_index))
1233        })
1234        .min()
1235        .map(|(_, turn_index)| turn_index);
1236    Some((
1237        completion.timestamp_ms,
1238        content.len(),
1239        !content.contains("<status>completed</status>"),
1240        consumer_turn_index,
1241    ))
1242}
1243
1244fn sanitize_structure(value: &Value, normalizer: &mut ToolIdNormalizer) -> Value {
1245    match value {
1246        Value::Object(map) => {
1247            let mut sanitized = Map::new();
1248            for (key, item) in map {
1249                if matches!(
1250                    key.as_str(),
1251                    "tool_use_id" | "toolUseID" | "parentToolUseID"
1252                ) && let Some(raw_id) = item.as_str()
1253                    && let Some(normalized) = normalizer.normalize(Some(raw_id))
1254                {
1255                    sanitized.insert(key.clone(), Value::String(normalized));
1256                    continue;
1257                }
1258                sanitized.insert(key.clone(), sanitize_structure(item, normalizer));
1259            }
1260            Value::Object(sanitized)
1261        }
1262        Value::Array(items) => Value::Array(
1263            items
1264                .iter()
1265                .map(|item| sanitize_structure(item, normalizer))
1266                .collect(),
1267        ),
1268        _ => value.clone(),
1269    }
1270}
1271
1272fn count_trailing_tool_results(entries: &[ConversationEntry]) -> usize {
1273    entries
1274        .iter()
1275        .rev()
1276        .take_while(|entry| entry.kind == "user_tool_result")
1277        .count()
1278}
1279
1280fn append_rendered_entries(buffer: &mut String, entries: &[ConversationEntry]) {
1281    for entry in entries {
1282        if !buffer.is_empty() {
1283            buffer.push('\n');
1284        }
1285        buffer.push_str(&entry.rendered);
1286    }
1287}
1288
1289fn render_entry_buffer(entries: &[ConversationEntry]) -> String {
1290    let mut buffer = String::new();
1291    append_rendered_entries(&mut buffer, entries);
1292    buffer
1293}
1294
1295fn render_user_entries(
1296    message: Option<&Map<String, Value>>,
1297    normalizer: &mut ToolIdNormalizer,
1298) -> Result<Vec<ConversationEntry>> {
1299    let mut rendered_entries = Vec::new();
1300    let Some(message) = message else {
1301        return Ok(rendered_entries);
1302    };
1303
1304    for block in content_blocks(message.get("content")) {
1305        let block_type = block
1306            .get("type")
1307            .and_then(Value::as_str)
1308            .unwrap_or("unknown");
1309        if matches!(block_type, "thinking" | "redacted_thinking") {
1310            continue;
1311        }
1312        if block_type == "text" {
1313            let text = block
1314                .get("text")
1315                .and_then(Value::as_str)
1316                .unwrap_or_default();
1317            if !text.is_empty() {
1318                rendered_entries.push(ConversationEntry {
1319                    kind: "user_text".to_string(),
1320                    rendered: format!("[user] {text}"),
1321                });
1322            }
1323            continue;
1324        }
1325        if block_type == "tool_result" {
1326            let normalized_id =
1327                normalizer.normalize(block.get("tool_use_id").and_then(Value::as_str));
1328            let content_text =
1329                flatten_block_content_text(block.get("content").unwrap_or(&Value::Null))?;
1330            let is_error = block
1331                .get("is_error")
1332                .and_then(Value::as_bool)
1333                .unwrap_or(false);
1334            let header = format!(
1335                "[user_tool_result id={} error={}]",
1336                normalized_id.unwrap_or_else(|| "tool_unknown".to_string()),
1337                if is_error { "true" } else { "false" }
1338            );
1339            let rendered = if content_text.is_empty() {
1340                header
1341            } else {
1342                format!("{header} {content_text}")
1343            };
1344            rendered_entries.push(ConversationEntry {
1345                kind: "user_tool_result".to_string(),
1346                rendered,
1347            });
1348            continue;
1349        }
1350
1351        let sanitized = sanitize_structure(&block, normalizer);
1352        rendered_entries.push(ConversationEntry {
1353            kind: "user_block".to_string(),
1354            rendered: format!(
1355                "[user_block type={block_type}] {}",
1356                canonical_json_string(&sanitized)?
1357            ),
1358        });
1359    }
1360
1361    Ok(rendered_entries)
1362}
1363
1364fn summarize_assistant_group(
1365    records: &[TraceRecord],
1366    group_indices: &[usize],
1367    normalizer: &mut ToolIdNormalizer,
1368    tokenizer: &mut impl TokenizerWorker,
1369) -> Result<AssistantGroupSummary> {
1370    let mut entries = Vec::new();
1371    let mut entries_by_record = BTreeMap::new();
1372    let mut tool_calls = Vec::new();
1373    let mut raw_task_tool_ids = Vec::new();
1374    let mut assistant_text_blocks = 0;
1375    let mut output_lengths = Vec::new();
1376    let mut input_lengths = Vec::new();
1377    let mut cache_read_lengths = Vec::new();
1378    let mut cache_creation_lengths = Vec::new();
1379    let mut model = None;
1380
1381    for index in group_indices {
1382        let record = &records[*index];
1383        let mut record_entries = Vec::new();
1384        let Some(message) = object_field(&record.raw, "message") else {
1385            continue;
1386        };
1387        if model.is_none() {
1388            model = message
1389                .get("model")
1390                .and_then(Value::as_str)
1391                .map(str::to_string);
1392        }
1393        if let Some(usage) = object_field(&Value::Object(message.clone()), "usage") {
1394            let input_tokens = usage
1395                .get("input_tokens")
1396                .and_then(Value::as_u64)
1397                .unwrap_or(0) as usize;
1398            let cache_read = usage
1399                .get("cache_read_input_tokens")
1400                .and_then(Value::as_u64)
1401                .unwrap_or(0) as usize;
1402            let cache_creation = usage
1403                .get("cache_creation_input_tokens")
1404                .and_then(Value::as_u64)
1405                .unwrap_or(0) as usize;
1406            input_lengths.push(
1407                input_tokens
1408                    .saturating_add(cache_read)
1409                    .saturating_add(cache_creation),
1410            );
1411            cache_read_lengths.push(cache_read);
1412            cache_creation_lengths.push(cache_creation);
1413            if let Some(output_tokens) = usage.get("output_tokens").and_then(Value::as_u64) {
1414                output_lengths.push(output_tokens as usize);
1415            }
1416        }
1417
1418        for block in content_blocks(message.get("content")) {
1419            let block_type = block
1420                .get("type")
1421                .and_then(Value::as_str)
1422                .unwrap_or("unknown");
1423            if matches!(block_type, "thinking" | "redacted_thinking") {
1424                continue;
1425            }
1426            if block_type == "text" {
1427                let text = block
1428                    .get("text")
1429                    .and_then(Value::as_str)
1430                    .unwrap_or_default();
1431                if !text.is_empty() {
1432                    assistant_text_blocks += 1;
1433                    record_entries.push(ConversationEntry {
1434                        kind: "assistant_text".to_string(),
1435                        rendered: format!("[assistant] {text}"),
1436                    });
1437                }
1438                continue;
1439            }
1440            if block_type == "tool_use" {
1441                let raw_id = block.get("id").and_then(Value::as_str);
1442                let normalized_id = normalizer.normalize(raw_id);
1443                let tool_name = block
1444                    .get("name")
1445                    .and_then(Value::as_str)
1446                    .unwrap_or("unknown")
1447                    .to_string();
1448                let args_json = canonical_json_string(&sanitize_structure(
1449                    block.get("input").unwrap_or(&Value::Null),
1450                    normalizer,
1451                ))?;
1452                record_entries.push(ConversationEntry {
1453                    kind: "assistant_tool_use".to_string(),
1454                    rendered: format!(
1455                        "[assistant_tool_use id={} name={} args={}]",
1456                        normalized_id
1457                            .clone()
1458                            .unwrap_or_else(|| "tool_unknown".to_string()),
1459                        tool_name,
1460                        args_json
1461                    ),
1462                });
1463                tool_calls.push(ToolCallSummary {
1464                    raw_id: raw_id.map(str::to_string),
1465                    name: tool_name.clone(),
1466                    normalized_id,
1467                    arg_size_chars: args_json.len(),
1468                    started_at_ms: record.timestamp_ms,
1469                });
1470                if matches!(tool_name.as_str(), "Agent" | "Task")
1471                    && let Some(raw_id) = raw_id
1472                {
1473                    raw_task_tool_ids.push(raw_id.to_string());
1474                }
1475                continue;
1476            }
1477
1478            let sanitized = sanitize_structure(&block, normalizer);
1479            record_entries.push(ConversationEntry {
1480                kind: "assistant_block".to_string(),
1481                rendered: format!(
1482                    "[assistant_block type={block_type}] {}",
1483                    canonical_json_string(&sanitized)?
1484                ),
1485            });
1486        }
1487        entries.extend(record_entries.iter().cloned());
1488        entries_by_record.insert(*index, record_entries);
1489    }
1490
1491    let output_length = if let Some(max_length) = output_lengths.into_iter().max() {
1492        max_length
1493    } else {
1494        let rendered_text = render_entry_buffer(&entries);
1495        tokenizer.encode(&rendered_text)?.len()
1496    };
1497
1498    let start_ms = group_indices
1499        .first()
1500        .map(|index| records[*index].timestamp_ms)
1501        .unwrap_or_default();
1502    let end_ms = group_indices
1503        .last()
1504        .map(|index| records[*index].timestamp_ms)
1505        .unwrap_or(start_ms);
1506
1507    Ok(AssistantGroupSummary {
1508        entries_by_record,
1509        model: model.unwrap_or_else(|| "unknown".to_string()),
1510        output_length,
1511        assistant_text_blocks,
1512        top_level_tool_calls: tool_calls,
1513        raw_task_tool_ids,
1514        input_length: input_lengths.into_iter().max(),
1515        cache_read_input_tokens: cache_read_lengths.into_iter().max(),
1516        cache_creation_input_tokens: cache_creation_lengths.into_iter().max(),
1517        start_ms,
1518        end_ms,
1519    })
1520}
1521
1522fn progress_timestamp_ms(record: &TraceRecord) -> i64 {
1523    record
1524        .raw
1525        .get("data")
1526        .and_then(Value::as_object)
1527        .and_then(|data| data.get("message"))
1528        .and_then(Value::as_object)
1529        .and_then(|message| message.get("timestamp"))
1530        .and_then(Value::as_str)
1531        .and_then(|timestamp| parse_utc_timestamp_ms(timestamp).ok())
1532        .unwrap_or(record.timestamp_ms)
1533}
1534
1535fn build_progress_index(records: &[TraceRecord]) -> FxHashMap<String, Vec<usize>> {
1536    let mut progress_index: FxHashMap<String, Vec<usize>> = FxHashMap::default();
1537    for (index, record) in records.iter().enumerate() {
1538        if !record.row_type.contains("progress") {
1539            continue;
1540        }
1541        let Some(parent_tool_use_id) = record.raw.get("parentToolUseID").and_then(Value::as_str)
1542        else {
1543            continue;
1544        };
1545        progress_index
1546            .entry(parent_tool_use_id.to_string())
1547            .or_default()
1548            .push(index);
1549    }
1550
1551    for indices in progress_index.values_mut() {
1552        indices.sort_by_key(|index| {
1553            let record = &records[*index];
1554            (progress_timestamp_ms(record), record.source_order)
1555        });
1556    }
1557    progress_index
1558}
1559
1560fn build_progress_metrics_index(
1561    progress_index: &FxHashMap<String, Vec<usize>>,
1562    records: &[TraceRecord],
1563) -> FxHashMap<String, CachedProgressMetrics> {
1564    progress_index
1565        .iter()
1566        .map(|(tool_id, indices)| {
1567            (
1568                tool_id.clone(),
1569                summarize_progress_indices(indices, records),
1570            )
1571        })
1572        .collect()
1573}
1574
1575fn aggregate_progress_metrics(
1576    task_tool_ids: &[String],
1577    progress_metrics_index: &FxHashMap<String, CachedProgressMetrics>,
1578    normalizer: &mut ToolIdNormalizer,
1579) -> Value {
1580    let mut selected_metrics = Vec::new();
1581    let mut seen_task_ids = BTreeSet::new();
1582    for task_tool_id in task_tool_ids {
1583        if !seen_task_ids.insert(task_tool_id.as_str()) {
1584            continue;
1585        }
1586        if let Some(metrics) = progress_metrics_index.get(task_tool_id) {
1587            selected_metrics.push(metrics);
1588        }
1589    }
1590
1591    if selected_metrics.is_empty() {
1592        return json!({
1593            "task_parent_tool_ids": task_tool_ids
1594                .iter()
1595                .filter_map(|tool_id| normalizer.normalize(Some(tool_id)))
1596                .collect::<Vec<_>>(),
1597            "nested_progress_event_count": 0,
1598            "nested_agent_count": 0,
1599            "nested_tool_call_count": 0,
1600            "nested_tool_result_count": 0,
1601            "nested_tool_error_count": 0,
1602            "nested_tool_counts": BTreeMap::<String, usize>::new(),
1603            "nested_tool_names": Vec::<String>::new(),
1604            "nested_tool_total_latency_ms": 0,
1605            "nested_tool_max_latency_ms": 0,
1606            "nested_tool_avg_latency_ms": 0,
1607            "nested_tool_max_parallelism": 0,
1608            "nested_assistant_text_blocks": 0,
1609            "task_duration_ms": 0,
1610        });
1611    }
1612
1613    let mut agent_ids = BTreeSet::new();
1614    let mut assistant_text_blocks = 0_usize;
1615    let mut tool_counts: BTreeMap<String, usize> = BTreeMap::new();
1616    let mut tool_intervals = Vec::new();
1617    let mut tool_result_count = 0_usize;
1618    let mut tool_error_count = 0_usize;
1619    let mut first_ts = i64::MAX;
1620    let mut last_ts = i64::MIN;
1621    let mut progress_event_count = 0_usize;
1622
1623    for metrics in selected_metrics {
1624        progress_event_count += metrics.progress_event_count;
1625        assistant_text_blocks += metrics.assistant_text_blocks;
1626        tool_result_count += metrics.tool_result_count;
1627        tool_error_count += metrics.tool_error_count;
1628        first_ts = first_ts.min(metrics.first_ts);
1629        last_ts = last_ts.max(metrics.last_ts);
1630        agent_ids.extend(metrics.agent_ids.iter().cloned());
1631        tool_intervals.extend(metrics.tool_intervals.iter().copied());
1632        for (tool_name, count) in &metrics.tool_counts {
1633            *tool_counts.entry(tool_name.clone()).or_insert(0) += count;
1634        }
1635    }
1636
1637    let total_latency: i64 = tool_intervals
1638        .iter()
1639        .map(|(start, end)| (end - start).max(0))
1640        .sum();
1641    let max_latency = tool_intervals
1642        .iter()
1643        .map(|(start, end)| (end - start).max(0))
1644        .max()
1645        .unwrap_or(0);
1646    let avg_latency = if tool_intervals.is_empty() {
1647        0
1648    } else {
1649        total_latency / tool_intervals.len() as i64
1650    };
1651
1652    let mut parallel_events = Vec::new();
1653    for (start_ms, end_ms) in &tool_intervals {
1654        parallel_events.push((*start_ms, 1_i32));
1655        parallel_events.push((*end_ms, -1_i32));
1656    }
1657    parallel_events.sort_by_key(|(timestamp, delta)| (*timestamp, -delta));
1658    let mut current_parallelism = 0_i32;
1659    let mut max_parallelism = 0_i32;
1660    for (_, delta) in parallel_events {
1661        current_parallelism += delta;
1662        max_parallelism = max_parallelism.max(current_parallelism);
1663    }
1664
1665    json!({
1666        "task_parent_tool_ids": task_tool_ids
1667            .iter()
1668            .filter_map(|tool_id| normalizer.normalize(Some(tool_id)))
1669            .collect::<Vec<_>>(),
1670        "nested_progress_event_count": progress_event_count,
1671        "nested_agent_count": agent_ids.len(),
1672        "nested_tool_call_count": tool_counts.values().sum::<usize>(),
1673        "nested_tool_result_count": tool_result_count,
1674        "nested_tool_error_count": tool_error_count,
1675        "nested_tool_counts": tool_counts,
1676        "nested_tool_names": tool_counts.keys().cloned().collect::<Vec<_>>(),
1677        "nested_tool_total_latency_ms": total_latency,
1678        "nested_tool_max_latency_ms": max_latency,
1679        "nested_tool_avg_latency_ms": avg_latency,
1680        "nested_tool_max_parallelism": max_parallelism,
1681        "nested_assistant_text_blocks": assistant_text_blocks,
1682        "task_duration_ms": (last_ts - first_ts).max(0),
1683    })
1684}
1685
1686fn summarize_progress_indices(indices: &[usize], records: &[TraceRecord]) -> CachedProgressMetrics {
1687    let mut agent_ids = BTreeSet::new();
1688    let mut assistant_text_blocks = 0_usize;
1689    let mut tool_counts: BTreeMap<String, usize> = BTreeMap::new();
1690    let mut tool_start_times: FxHashMap<String, i64> = FxHashMap::default();
1691    let mut tool_intervals = Vec::new();
1692    let mut tool_result_count = 0_usize;
1693    let mut tool_error_count = 0_usize;
1694    let first_ts = progress_timestamp_ms(&records[indices[0]]);
1695    let last_ts = progress_timestamp_ms(&records[*indices.last().unwrap()]);
1696
1697    for index in indices {
1698        let record = &records[*index];
1699        let timestamp_ms = progress_timestamp_ms(record);
1700        let Some(data) = record.raw.get("data").and_then(Value::as_object) else {
1701            continue;
1702        };
1703        if let Some(agent_id) = data.get("agentId").and_then(Value::as_str) {
1704            agent_ids.insert(agent_id.to_string());
1705        }
1706
1707        let Some(nested_message) = data.get("message").and_then(Value::as_object) else {
1708            continue;
1709        };
1710        let nested_type = nested_message
1711            .get("type")
1712            .and_then(Value::as_str)
1713            .unwrap_or_default();
1714        let nested_payload = nested_message.get("message").and_then(Value::as_object);
1715
1716        if nested_type == "assistant" {
1717            let Some(nested_payload) = nested_payload else {
1718                continue;
1719            };
1720            for block in content_blocks(nested_payload.get("content")) {
1721                let block_type = block
1722                    .get("type")
1723                    .and_then(Value::as_str)
1724                    .unwrap_or_default();
1725                if block_type == "text" {
1726                    assistant_text_blocks += 1;
1727                    continue;
1728                }
1729                if block_type != "tool_use" {
1730                    continue;
1731                }
1732                let Some(raw_id) = block.get("id").and_then(Value::as_str) else {
1733                    continue;
1734                };
1735                let tool_name = block
1736                    .get("name")
1737                    .and_then(Value::as_str)
1738                    .unwrap_or("unknown")
1739                    .to_string();
1740                *tool_counts.entry(tool_name).or_insert(0) += 1;
1741                tool_start_times.insert(raw_id.to_string(), timestamp_ms);
1742            }
1743            continue;
1744        }
1745
1746        if nested_type != "user" {
1747            continue;
1748        }
1749        let Some(nested_payload) = nested_payload else {
1750            continue;
1751        };
1752        for block in content_blocks(nested_payload.get("content")) {
1753            if block.get("type").and_then(Value::as_str) != Some("tool_result") {
1754                continue;
1755            }
1756            let Some(raw_tool_id) = block.get("tool_use_id").and_then(Value::as_str) else {
1757                continue;
1758            };
1759            if let Some(start_ms) = tool_start_times.remove(raw_tool_id) {
1760                tool_intervals.push((start_ms, timestamp_ms));
1761            }
1762            tool_result_count += 1;
1763            if block
1764                .get("is_error")
1765                .and_then(Value::as_bool)
1766                .unwrap_or(false)
1767            {
1768                tool_error_count += 1;
1769            }
1770        }
1771    }
1772
1773    CachedProgressMetrics {
1774        progress_event_count: indices.len(),
1775        agent_ids,
1776        assistant_text_blocks,
1777        tool_counts,
1778        tool_result_count,
1779        tool_error_count,
1780        tool_intervals,
1781        first_ts,
1782        last_ts,
1783    }
1784}
1785
1786#[cfg(test)]
1787mod tests {
1788    use super::{
1789        SessionTurnBuilder, TraceRecord, build_source_fidelity_oracle, load_trace_records,
1790    };
1791    use crate::coding::common::anonymized_session_id;
1792    use crate::coding::tokenizer::TokenizerWorker;
1793    use anyhow::Result;
1794    use rustc_hash::FxHashMap;
1795    use serde_json::{Value, json};
1796    use tempfile::TempDir;
1797
1798    struct StubTokenizer;
1799
1800    impl TokenizerWorker for StubTokenizer {
1801        fn encode(&mut self, text: &str) -> Result<Vec<u32>> {
1802            Ok(vec![text.len() as u32])
1803        }
1804    }
1805
1806    fn make_record(
1807        row_type: &str,
1808        timestamp_ms: i64,
1809        source_order: u64,
1810        raw: Value,
1811    ) -> TraceRecord {
1812        TraceRecord {
1813            session_id: "session-1".to_string(),
1814            parent_session_id: None,
1815            row_type: row_type.to_string(),
1816            timestamp_ms,
1817            source_order,
1818            raw,
1819        }
1820    }
1821
1822    #[test]
1823    fn groups_interleaved_fragments_and_pairs_tool_results() {
1824        let records = vec![
1825            make_record(
1826                "user",
1827                1_000,
1828                0,
1829                json!({"type":"user","message":{"role":"user","content":"start"}}),
1830            ),
1831            make_record(
1832                "assistant",
1833                1_100,
1834                1,
1835                json!({"type":"assistant","requestId":"req-1","message":{"id":"msg-1","content":[{"type":"text","text":"working"}],"usage":{"input_tokens":2,"cache_creation_input_tokens":3,"cache_read_input_tokens":5,"output_tokens":7}}}),
1836            ),
1837            make_record(
1838                "assistant",
1839                1_200,
1840                2,
1841                json!({"type":"assistant","requestId":"req-1","message":{"id":"msg-1","content":[{"type":"tool_use","id":"raw-1","name":"Read","input":{}}],"usage":{"input_tokens":2,"cache_creation_input_tokens":3,"cache_read_input_tokens":5,"output_tokens":7}}}),
1842            ),
1843            make_record(
1844                "system",
1845                1_250,
1846                3,
1847                json!({"type":"system","subtype":"turn_duration"}),
1848            ),
1849            make_record(
1850                "user",
1851                1_300,
1852                4,
1853                json!({"type":"user","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"raw-1","content":"ok"}]}}),
1854            ),
1855            make_record(
1856                "assistant",
1857                1_400,
1858                5,
1859                json!({"type":"assistant","requestId":"req-1","message":{"id":"msg-1","content":[{"type":"tool_use","id":"raw-2","name":"Bash","input":{}}],"usage":{"input_tokens":2,"cache_creation_input_tokens":3,"cache_read_input_tokens":5,"output_tokens":7}}}),
1860            ),
1861            make_record(
1862                "user",
1863                1_500,
1864                6,
1865                json!({"type":"user","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"raw-2","content":"failed","is_error":true}]}}),
1866            ),
1867            make_record(
1868                "assistant",
1869                1_600,
1870                7,
1871                json!({"type":"assistant","requestId":"req-2","message":{"id":"msg-2","content":[{"type":"text","text":"done"}],"usage":{"input_tokens":1,"cache_creation_input_tokens":0,"cache_read_input_tokens":10,"output_tokens":2}}}),
1872            ),
1873        ];
1874
1875        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
1876        let first = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
1877        let second = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
1878
1879        assert_eq!(first.request_start_ms, 1_000);
1880        assert_eq!(first.assistant_end_ms, 1_400);
1881        assert_eq!(first.observed_input_length, Some(10));
1882        assert_eq!(first.output_length, 7);
1883        assert_eq!(first.tools.len(), 2);
1884        assert_eq!(first.tools[0].tool_class, "Read");
1885        assert_eq!(first.tools[0].started_at_ms, 1_200);
1886        assert_eq!(first.tools[0].ended_at_ms, 1_300);
1887        assert!(first.tools[1].is_error);
1888        assert_eq!(second.request_start_ms, 1_500);
1889        assert!(builder.next_turn(&mut StubTokenizer).unwrap().is_none());
1890    }
1891
1892    #[test]
1893    fn loader_preserves_source_order_for_compaction_markers() {
1894        let temp = TempDir::new().unwrap();
1895        let trace = temp.path().join("session.jsonl");
1896        std::fs::write(
1897            &trace,
1898            concat!(
1899                "{\"type\":\"system\",\"subtype\":\"compact_boundary\",\"sessionId\":\"session-1\",\"timestamp\":\"2026-01-01T00:00:00.002Z\",\"compactMetadata\":{\"trigger\":\"manual\",\"preTokens\":10,\"postTokens\":3,\"durationMs\":1}}\n",
1900                "{\"type\":\"user\",\"isCompactSummary\":true,\"sessionId\":\"session-1\",\"timestamp\":\"2026-01-01T00:00:00.001Z\",\"message\":{\"content\":\"summary\"}}\n",
1901                "{\"type\":\"assistant\",\"sessionId\":\"session-1\",\"timestamp\":\"2026-01-01T00:00:00.003Z\",\"message\":{\"id\":\"msg-1\",\"content\":[{\"type\":\"text\",\"text\":\"done\"}],\"usage\":{\"output_tokens\":1}}}\n"
1902            ),
1903        )
1904        .unwrap();
1905
1906        let sessions = load_trace_records(&[trace]).unwrap();
1907        let records = sessions.get("session-1").unwrap();
1908        assert_eq!(records[0].row_type, "system");
1909        assert_eq!(records[1].row_type, "user");
1910
1911        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records.clone(), true);
1912        let compaction = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
1913        assert!(compaction.compaction.is_some());
1914        let turn = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
1915        assert_eq!(turn.input_text, "[user] summary");
1916    }
1917
1918    #[test]
1919    fn loader_infers_immediate_parent_from_agent_result() {
1920        let temp = TempDir::new().unwrap();
1921        let parent = temp.path().join("parent.jsonl");
1922        let child = temp.path().join("child.jsonl");
1923        std::fs::write(
1924            &parent,
1925            "{\"type\":\"user\",\"sessionId\":\"root\",\"agentId\":\"parent\",\"timestamp\":\"2026-01-01T00:00:00.001Z\",\"toolUseResult\":{\"agentId\":\"child\"},\"message\":{\"content\":[{\"type\":\"tool_result\",\"tool_use_id\":\"agent-call\",\"content\":\"done\"}]}}\n",
1926        )
1927        .unwrap();
1928        std::fs::write(
1929            &child,
1930            "{\"type\":\"assistant\",\"sessionId\":\"root\",\"agentId\":\"child\",\"timestamp\":\"2026-01-01T00:00:00.002Z\",\"message\":{\"id\":\"msg-1\",\"content\":[{\"type\":\"text\",\"text\":\"done\"}],\"usage\":{\"output_tokens\":1}}}\n",
1931        )
1932        .unwrap();
1933
1934        let sessions = load_trace_records(&[parent, child]).unwrap();
1935        assert_eq!(
1936            sessions["child"][0].parent_session_id.as_deref(),
1937            Some("parent")
1938        );
1939    }
1940
1941    #[test]
1942    fn fidelity_oracle_ignores_root_sidechain_compaction() {
1943        let mut sessions = FxHashMap::default();
1944        sessions.insert(
1945            "session-1".to_string(),
1946            vec![make_record(
1947                "system",
1948                1_000,
1949                0,
1950                json!({"type":"system","subtype":"compact_boundary","isSidechain":true,"compactMetadata":{"trigger":"manual","preTokens":10,"postTokens":3,"durationMs":500}}),
1951            )],
1952        );
1953
1954        assert!(
1955            build_source_fidelity_oracle(&sessions)
1956                .unwrap()
1957                .compactions
1958                .is_empty()
1959        );
1960    }
1961
1962    #[test]
1963    fn compact_boundary_restarts_transcript_from_summary() {
1964        let records = vec![
1965            make_record(
1966                "user",
1967                1_000,
1968                0,
1969                json!({"type":"user","message":{"role":"user","content":"before compact"}}),
1970            ),
1971            make_record(
1972                "assistant",
1973                2_000,
1974                1,
1975                json!({"type":"assistant","message":{"id":"assistant-1","content":[{"type":"text","text":"first answer"}],"usage":{"output_tokens":3}}}),
1976            ),
1977            make_record(
1978                "system",
1979                3_000,
1980                2,
1981                json!({"type":"system","subtype":"compact_boundary","compactMetadata":{"trigger":"manual","preTokens":10,"postTokens":3,"durationMs":500}}),
1982            ),
1983            make_record(
1984                "user",
1985                3_001,
1986                3,
1987                json!({"type":"user","isCompactSummary":true,"message":{"role":"user","content":"compacted summary"}}),
1988            ),
1989            make_record(
1990                "assistant",
1991                4_000,
1992                4,
1993                json!({"type":"assistant","message":{"id":"assistant-2","content":[{"type":"text","text":"second answer"}],"usage":{"output_tokens":5}}}),
1994            ),
1995        ];
1996
1997        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
1998        let mut tokenizer = StubTokenizer;
1999        let mut turns = Vec::new();
2000        while let Some(turn) = builder.next_turn(&mut tokenizer).unwrap() {
2001            turns.push(turn);
2002        }
2003
2004        assert_eq!(
2005            turns
2006                .iter()
2007                .filter(|turn| turn.compaction.is_none())
2008                .map(|turn| turn.input_text.as_str())
2009                .collect::<Vec<_>>(),
2010            vec!["[user] before compact", "[user] compacted summary"]
2011        );
2012        assert_eq!(
2013            turns
2014                .iter()
2015                .filter(|turn| turn.compaction.is_some())
2016                .count(),
2017            1
2018        );
2019    }
2020
2021    #[test]
2022    fn compact_boundary_survives_ignored_rows_around_summary() {
2023        let records = vec![
2024            make_record(
2025                "user",
2026                1_000,
2027                0,
2028                json!({"type":"user","message":{"role":"user","content":"before compact"}}),
2029            ),
2030            make_record(
2031                "assistant",
2032                2_000,
2033                1,
2034                json!({"type":"assistant","message":{"id":"assistant-1","content":[{"type":"text","text":"first answer"}],"usage":{"output_tokens":3}}}),
2035            ),
2036            make_record(
2037                "system",
2038                3_000,
2039                2,
2040                json!({"type":"system","subtype":"compact_boundary","compactMetadata":{"trigger":"manual","preTokens":10,"postTokens":3,"durationMs":500}}),
2041            ),
2042            make_record(
2043                "system",
2044                3_001,
2045                3,
2046                json!({"type":"system","subtype":"turn_duration"}),
2047            ),
2048            make_record(
2049                "user",
2050                3_002,
2051                4,
2052                json!({"type":"user","isMeta":true,"message":{"role":"user","content":"<local-command-caveat>ignore me</local-command-caveat>"}}),
2053            ),
2054            make_record(
2055                "user",
2056                3_003,
2057                5,
2058                json!({"type":"user","isCompactSummary":true,"message":{"role":"user","content":"compacted summary"}}),
2059            ),
2060            make_record(
2061                "user",
2062                3_004,
2063                6,
2064                json!({"type":"user","isMeta":true,"message":{"role":"user","content":"<local-command-caveat>ignore me</local-command-caveat>"}}),
2065            ),
2066            make_record(
2067                "user",
2068                3_005,
2069                7,
2070                json!({"type":"user","message":{"role":"user","content":"<command-name>/compact</command-name>\n<command-message>compact</command-message>"}}),
2071            ),
2072            make_record(
2073                "user",
2074                3_006,
2075                8,
2076                json!({"type":"user","message":{"role":"user","content":"<local-command-stdout>Compacted</local-command-stdout>"}}),
2077            ),
2078            make_record(
2079                "assistant",
2080                4_000,
2081                9,
2082                json!({"type":"assistant","message":{"id":"assistant-2","content":[{"type":"text","text":"second answer"}],"usage":{"output_tokens":5}}}),
2083            ),
2084        ];
2085
2086        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
2087        let mut tokenizer = StubTokenizer;
2088        let mut turns = Vec::new();
2089        while let Some(turn) = builder.next_turn(&mut tokenizer).unwrap() {
2090            turns.push(turn);
2091        }
2092
2093        assert_eq!(
2094            turns
2095                .iter()
2096                .filter(|turn| turn.compaction.is_none())
2097                .map(|turn| turn.input_text.as_str())
2098                .collect::<Vec<_>>(),
2099            vec!["[user] before compact", "[user] compacted summary"]
2100        );
2101    }
2102
2103    #[test]
2104    fn orphan_compact_summary_still_replaces_transcript() {
2105        let records = vec![
2106            make_record(
2107                "user",
2108                1_000,
2109                0,
2110                json!({"type":"user","message":{"role":"user","content":"old prompt"}}),
2111            ),
2112            make_record(
2113                "assistant",
2114                2_000,
2115                1,
2116                json!({"type":"assistant","message":{"id":"assistant-1","content":[{"type":"text","text":"old answer"}],"usage":{"output_tokens":2}}}),
2117            ),
2118            make_record(
2119                "user",
2120                3_000,
2121                2,
2122                json!({"type":"user","isCompactSummary":true,"message":{"role":"user","content":"summary only"}}),
2123            ),
2124            make_record(
2125                "assistant",
2126                4_000,
2127                3,
2128                json!({"type":"assistant","message":{"id":"assistant-2","content":[{"type":"text","text":"new answer"}],"usage":{"output_tokens":2}}}),
2129            ),
2130        ];
2131
2132        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
2133        let mut tokenizer = StubTokenizer;
2134        assert_eq!(
2135            builder
2136                .next_turn(&mut tokenizer)
2137                .unwrap()
2138                .unwrap()
2139                .input_text,
2140            "[user] old prompt"
2141        );
2142        assert_eq!(
2143            builder
2144                .next_turn(&mut tokenizer)
2145                .unwrap()
2146                .unwrap()
2147                .input_text,
2148            "[user] summary only"
2149        );
2150    }
2151
2152    #[test]
2153    fn background_agent_joins_only_after_completion_notification() {
2154        let records = vec![
2155            make_record(
2156                "user",
2157                1_000,
2158                0,
2159                json!({"type":"user","message":{"role":"user","content":"start agent"}}),
2160            ),
2161            make_record(
2162                "assistant",
2163                1_100,
2164                1,
2165                json!({"type":"assistant","requestId":"req-1","message":{"id":"msg-1","content":[{"type":"tool_use","id":"agent-call","name":"Agent","input":{"run_in_background":true}}],"usage":{"output_tokens":2}}}),
2166            ),
2167            make_record(
2168                "user",
2169                1_150,
2170                2,
2171                json!({"type":"user","toolUseResult":{"isAsync":true,"agentId":"child-agent","status":"async_launched"},"message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"agent-call","content":"launched"}]}}),
2172            ),
2173            make_record(
2174                "user",
2175                1_300,
2176                3,
2177                json!({"type":"user","message":{"role":"user","content":"keep working"}}),
2178            ),
2179            make_record(
2180                "assistant",
2181                1_400,
2182                4,
2183                json!({"type":"assistant","requestId":"req-2","message":{"id":"msg-2","content":[{"type":"text","text":"still working"}],"usage":{"output_tokens":2}}}),
2184            ),
2185            make_record(
2186                "queue-operation",
2187                1_800,
2188                5,
2189                json!({"type":"queue-operation","operation":"enqueue","content":"<tool-use-id>agent-call</tool-use-id><status>completed</status>done"}),
2190            ),
2191            make_record(
2192                "assistant",
2193                1_810,
2194                6,
2195                json!({"type":"assistant","requestId":"req-2","message":{"id":"msg-2","content":[{"type":"text","text":"late fragment"}],"usage":{"output_tokens":2}}}),
2196            ),
2197            make_record(
2198                "user",
2199                1_850,
2200                7,
2201                json!({"type":"user","message":{"role":"user","content":"use result"}}),
2202            ),
2203            make_record(
2204                "assistant",
2205                1_900,
2206                8,
2207                json!({"type":"assistant","requestId":"req-3","message":{"id":"msg-3","content":[{"type":"text","text":"finished"}],"usage":{"output_tokens":1}}}),
2208            ),
2209        ];
2210
2211        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
2212        let turns = std::iter::from_fn(|| builder.next_turn(&mut StubTokenizer).transpose())
2213            .collect::<Result<Vec<_>>>()
2214            .unwrap();
2215
2216        assert_eq!(turns.len(), 3);
2217        assert_eq!(turns[0].tools.len(), 1);
2218        let tool = &turns[0].tools[0];
2219        assert_eq!(tool.ended_at_ms, 1_800);
2220        assert_eq!(tool.consumer_turn_index, Some(2));
2221        assert_eq!(tool.child_session_id.as_deref(), Some("child-agent"));
2222        assert_eq!(tool.execution_mode, "background");
2223    }
2224
2225    #[test]
2226    fn background_bash_uses_terminal_notification_status() {
2227        let records = vec![
2228            make_record(
2229                "user",
2230                1_000,
2231                0,
2232                json!({"type":"user","message":{"role":"user","content":"start command"}}),
2233            ),
2234            make_record(
2235                "assistant",
2236                1_100,
2237                1,
2238                json!({"type":"assistant","requestId":"req-1","message":{"id":"msg-1","content":[{"type":"tool_use","id":"bash-call","name":"Bash","input":{"run_in_background":true}}],"usage":{"output_tokens":2}}}),
2239            ),
2240            make_record(
2241                "user",
2242                1_150,
2243                2,
2244                json!({"type":"user","toolUseResult":{"backgroundTaskId":"task-1"},"message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"bash-call","content":"launched","is_error":false}]}}),
2245            ),
2246            make_record(
2247                "assistant",
2248                1_300,
2249                3,
2250                json!({"type":"assistant","requestId":"req-2","message":{"id":"msg-2","content":[{"type":"text","text":"other work"}],"usage":{"output_tokens":1}}}),
2251            ),
2252            make_record(
2253                "queue-operation",
2254                1_500,
2255                4,
2256                json!({"type":"queue-operation","operation":"enqueue","content":"<task-id>task-1</task-id><tool-use-id>bash-call</tool-use-id><status>failed</status>"}),
2257            ),
2258            make_record(
2259                "user",
2260                1_510,
2261                5,
2262                json!({"type":"user","message":{"role":"user","content":"<tool-use-id>bash-call</tool-use-id><status>failed</status>"}}),
2263            ),
2264            make_record(
2265                "assistant",
2266                1_600,
2267                6,
2268                json!({"type":"assistant","requestId":"req-3","message":{"id":"msg-3","content":[{"type":"text","text":"handled"}],"usage":{"output_tokens":1}}}),
2269            ),
2270        ];
2271        let mut sessions = FxHashMap::default();
2272        sessions.insert("session-1".to_string(), records.clone());
2273        let oracle = build_source_fidelity_oracle(&sessions).unwrap();
2274        assert_eq!(oracle.background_tools, 1);
2275        assert_eq!(oracle.background_agents, 0);
2276        assert_eq!(oracle.background_completions_missing, 0);
2277        assert_eq!(oracle.tool_errors, 1);
2278
2279        let mut builder = SessionTurnBuilder::new("session-1".to_string(), records, true);
2280        let first = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
2281        let tool = &first.tools[0];
2282        assert_eq!(tool.ended_at_ms, 1_500);
2283        assert_eq!(tool.consumer_turn_index, Some(2));
2284        assert!(tool.is_error);
2285        assert!(tool.child_session_id.is_none());
2286        assert_eq!(tool.execution_mode, "background");
2287    }
2288
2289    #[test]
2290    fn child_identity_anonymizes_child_and_parent_consistently() {
2291        let records = vec![
2292            make_record(
2293                "user",
2294                1_000,
2295                0,
2296                json!({"type":"user","isSidechain":true,"message":{"role":"user","content":"task"}}),
2297            ),
2298            make_record(
2299                "assistant",
2300                2_000,
2301                1,
2302                json!({"type":"assistant","isSidechain":true,"message":{"id":"child-1","content":[{"type":"text","text":"done"}],"usage":{"output_tokens":1}}}),
2303            ),
2304        ];
2305        let mut builder = SessionTurnBuilder::new("child-agent".to_string(), records, false);
2306
2307        let turn = builder.next_turn(&mut StubTokenizer).unwrap().unwrap();
2308
2309        assert_eq!(turn.export_session_id, anonymized_session_id("child-agent"));
2310        assert_eq!(
2311            turn.export_parent_session_id.as_deref(),
2312            Some(anonymized_session_id("session-1").as_str())
2313        );
2314    }
2315}