1use 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}