1use crate::coding::claude::parser::{
10 SessionTurnBuilder, SourceFidelityOracle, TraceRecord, TurnDraft, build_source_fidelity_oracle,
11};
12use crate::coding::tokenizer::{TokenizerFactory, TokenizerWorker, last_word_overlap_start};
13use anyhow::{Result, anyhow, bail};
14use crossbeam_channel::{Receiver, Sender, bounded, unbounded};
15use dynamo_data_gen::{sequence_hashes_for_tokens, write_empty_files};
16use rustc_hash::FxHashMap;
17use serde::Serialize;
18use serde_json::{Map, Value, json};
19use std::cmp::Reverse;
20use std::collections::{BTreeMap, BTreeSet, BinaryHeap, VecDeque};
21use std::fs::File;
22use std::io::{BufWriter, Write};
23use std::path::Path;
24use std::thread::{self, JoinHandle};
25
26#[derive(Debug, Clone, Copy)]
27pub struct ExportConfig {
28 pub block_size: usize,
29 pub delta_overlap_words: usize,
30 pub tokenizer_workers: usize,
31}
32
33#[derive(Debug, Clone, Default)]
34pub struct ExportStats {
35 pub row_count: usize,
36 pub tool_row_count: usize,
37 pub sidecar_count: usize,
38 pub max_heap_len: usize,
39 pub fidelity: FidelityReport,
40}
41
42#[derive(Debug, Clone, Default)]
43pub struct FidelityReport {
44 pub requests_verified: usize,
45 pub compactions_verified: usize,
46 pub usage_requests_verified: usize,
47 pub tools_verified: usize,
48 pub child_links_verified: usize,
49 pub background_tools: usize,
50 pub background_agents: usize,
51 pub background_completions_missing: usize,
52 pub background_titles_unreplayable: usize,
53 pub cache_prefix_blocks_verified: usize,
54 pub compaction_prefix_blocks_verified: usize,
55 pub post_compaction_prefix_blocks_verified: usize,
56 pub unmatched_tool_calls: usize,
57 pub unmatched_tool_results: usize,
58 pub unresolved_child_sessions: usize,
59}
60
61#[derive(Serialize)]
67struct ClaudeToolReplayMetadata {
68 source_request_id: String,
69 #[serde(skip_serializing_if = "Option::is_none")]
70 consumer_request_id: Option<String>,
71 #[serde(skip_serializing_if = "Option::is_none")]
72 child_session_id: Option<String>,
73 execution_mode: String,
74}
75
76impl FidelityReport {
77 pub fn render(&self) -> String {
78 let ordinary_requests = self
79 .requests_verified
80 .saturating_sub(self.compactions_verified);
81 format!(
82 "Fidelity: requests={0}/{0} compactions={1}/{1} usage={2}/{15} tools={3}/{3} child_links={4}/{4} cache_prefix_blocks={5} compaction_prefix_blocks={6} post_compaction_prefix_blocks={7}\nBackground: tools={8} agents={9} missing_completions={10} title_requests_unreplayable={11}\nLimitations: synthetic_kv_hashes={0} unmatched_tool_calls={12} unmatched_tool_results={13} unresolved_child_sessions={14}",
83 self.requests_verified,
84 self.compactions_verified,
85 self.usage_requests_verified,
86 self.tools_verified,
87 self.child_links_verified,
88 self.cache_prefix_blocks_verified,
89 self.compaction_prefix_blocks_verified,
90 self.post_compaction_prefix_blocks_verified,
91 self.background_tools,
92 self.background_agents,
93 self.background_completions_missing,
94 self.background_titles_unreplayable,
95 self.unmatched_tool_calls,
96 self.unmatched_tool_results,
97 self.unresolved_child_sessions,
98 ordinary_requests,
99 )
100 }
101}
102
103struct FidelityVerifier {
104 oracle: SourceFidelityOracle,
105 seen_requests: BTreeSet<(String, String)>,
106 seen_compactions: BTreeSet<(String, String)>,
107 tools_by_class: BTreeMap<String, usize>,
108 tool_count: usize,
109 tool_errors: usize,
110 child_links: usize,
111 background_tools: usize,
112 background_agents: usize,
113 usage_requests: usize,
114 cache_prefix_blocks_verified: usize,
115 compaction_prefix_blocks_verified: usize,
116 post_compaction_prefix_blocks_verified: usize,
117 next_turn_by_session: FxHashMap<String, usize>,
118 previous_hashes_by_session: FxHashMap<String, Vec<u64>>,
119 previous_input_length_by_session: FxHashMap<String, usize>,
120 previous_was_compaction_by_session: FxHashMap<String, bool>,
121 expected_next_cache_read_by_session: FxHashMap<String, usize>,
122 export_sessions: BTreeSet<String>,
123 causal_references: Vec<(String, usize, String)>,
124 child_session_references: Vec<String>,
125}
126
127#[derive(Debug, Clone, Eq, Ord, PartialEq, PartialOrd)]
128struct HeapEntry {
129 request_start_ms: i64,
130 turn_index: usize,
131 export_session_id: String,
132 session_id: String,
133}
134
135#[derive(Debug)]
136struct OverlapBase {
137 previous_text: String,
138 previous_tokens: Vec<u32>,
139}
140
141#[derive(Debug)]
142struct ReadyTurn {
143 current_text: String,
144 tokens: Vec<u32>,
145}
146
147#[derive(Debug)]
148struct HeadTurn {
149 turn: TurnDraft,
150 turn_key: u64,
151 scheduled: bool,
152 ready: Option<ReadyTurn>,
153}
154
155#[derive(Debug)]
156struct SessionState {
157 builder: SessionTurnBuilder,
158 head: Option<HeadTurn>,
159 overlap_base: Option<OverlapBase>,
160 replay_base: Option<Vec<u32>>,
161 next_turn_key: u64,
162}
163
164#[derive(Debug)]
165struct TokenizeJob {
166 session_id: String,
167 turn_key: u64,
168 current_text: String,
169 overlap_start: Option<usize>,
170 previous_overlap_text: Option<String>,
171 previous_tokens: Option<Vec<u32>>,
172 overlap_words: usize,
173}
174
175#[derive(Debug)]
176struct TokenizeResponse {
177 session_id: String,
178 turn_key: u64,
179 outcome: Result<ReadyTurn, String>,
180}
181
182impl FidelityVerifier {
183 fn new(oracle: SourceFidelityOracle) -> Self {
184 Self {
185 oracle,
186 seen_requests: BTreeSet::new(),
187 seen_compactions: BTreeSet::new(),
188 tools_by_class: BTreeMap::new(),
189 tool_count: 0,
190 tool_errors: 0,
191 child_links: 0,
192 background_tools: 0,
193 background_agents: 0,
194 usage_requests: 0,
195 cache_prefix_blocks_verified: 0,
196 compaction_prefix_blocks_verified: 0,
197 post_compaction_prefix_blocks_verified: 0,
198 next_turn_by_session: FxHashMap::default(),
199 previous_hashes_by_session: FxHashMap::default(),
200 previous_input_length_by_session: FxHashMap::default(),
201 previous_was_compaction_by_session: FxHashMap::default(),
202 expected_next_cache_read_by_session: FxHashMap::default(),
203 export_sessions: BTreeSet::new(),
204 causal_references: Vec::new(),
205 child_session_references: Vec::new(),
206 }
207 }
208
209 fn observe(
210 &mut self,
211 turn: &TurnDraft,
212 replay_tokens: &[u32],
213 input_sequence_hashes: &[u64],
214 block_size: usize,
215 ) -> Result<()> {
216 let key = (turn.session_id.clone(), turn.source_request_id.clone());
217 if let Some(compaction) = &turn.compaction {
218 let expected = self.oracle.compactions.get(&key).ok_or_else(|| {
219 anyhow!(
220 "fidelity verification found unexpected compaction {} in session {}",
221 turn.source_request_id,
222 turn.session_id
223 )
224 })?;
225 if compaction != expected || !self.seen_compactions.insert(key) {
226 bail!(
227 "fidelity verification compaction mismatch for session {} sequence {}",
228 turn.session_id,
229 compaction.sequence
230 );
231 }
232 let expected_turn = self
233 .next_turn_by_session
234 .get(&turn.session_id)
235 .copied()
236 .unwrap_or_default();
237 let expected_start = compaction
238 .ended_at_ms
239 .saturating_sub(compaction.duration_ms);
240 if turn.turn_index != expected_turn
241 || turn.request_start_ms != expected_start
242 || turn.assistant_end_ms != compaction.ended_at_ms
243 || turn.observed_input_length != Some(compaction.pre_tokens)
244 || turn.cache_read_input_tokens.is_some()
245 || replay_tokens.len() != compaction.pre_tokens
246 {
247 bail!(
248 "fidelity verification compaction timing/cache mismatch for session {} sequence {}",
249 turn.session_id,
250 compaction.sequence
251 );
252 }
253 } else {
254 let expected = self.oracle.requests.get(&key).ok_or_else(|| {
255 anyhow!(
256 "fidelity verification found unexpected request {} in session {}",
257 turn.source_request_id,
258 turn.session_id
259 )
260 })?;
261 if !self.seen_requests.insert(key) {
262 bail!(
263 "fidelity verification found duplicate request {} in session {}",
264 turn.source_request_id,
265 turn.session_id
266 );
267 }
268 let expected_turn = self
269 .next_turn_by_session
270 .entry(turn.session_id.clone())
271 .or_default();
272 if turn.turn_index != *expected_turn {
273 bail!(
274 "fidelity verification expected turn {} for session {}, got {}",
275 *expected_turn,
276 turn.session_id,
277 turn.turn_index
278 );
279 }
280 *expected_turn += 1;
281 if turn.request_start_ms != expected.request_start_ms
282 || turn.assistant_end_ms != expected.assistant_end_ms
283 {
284 bail!(
285 "fidelity verification timing mismatch for session {} turn {}: expected {}..{}, got {}..{}",
286 turn.session_id,
287 turn.turn_index,
288 expected.request_start_ms,
289 expected.assistant_end_ms,
290 turn.request_start_ms,
291 turn.assistant_end_ms
292 );
293 }
294 if let Some(output_length) = expected.output_length
295 && turn.output_length != output_length
296 {
297 bail!(
298 "fidelity verification output mismatch for session {} turn {}: expected {}, got {}",
299 turn.session_id,
300 turn.turn_index,
301 output_length,
302 turn.output_length
303 );
304 }
305 if let Some(input_length) = expected.input_length {
306 self.usage_requests += 1;
307 if replay_tokens.len() != input_length
308 || turn.cache_read_input_tokens != expected.cache_read_input_tokens
309 || turn.cache_creation_input_tokens != expected.cache_creation_input_tokens
310 {
311 bail!(
312 "fidelity verification input/cache mismatch for session {} turn {}",
313 turn.session_id,
314 turn.turn_index
315 );
316 }
317 }
318 }
319 self.export_sessions.insert(turn.export_session_id.clone());
320 if turn.request_start_ms > turn.assistant_start_ms
321 || turn.assistant_start_ms > turn.assistant_end_ms
322 {
323 bail!(
324 "invalid request timing for session {} turn {}",
325 turn.session_id,
326 turn.turn_index
327 );
328 }
329 let expected_hashes = replay_tokens.len().div_ceil(block_size);
330 if input_sequence_hashes.len() != expected_hashes {
331 bail!(
332 "fidelity verification expected {} hashes for session {} turn {}, got {}",
333 expected_hashes,
334 turn.session_id,
335 turn.turn_index,
336 input_sequence_hashes.len()
337 );
338 }
339 let previous_was_compaction = self
340 .previous_was_compaction_by_session
341 .get(&turn.session_id)
342 .copied()
343 .unwrap_or(false);
344 let previous_input_length = self
345 .previous_input_length_by_session
346 .get(&turn.session_id)
347 .copied();
348 let previous_hashes = self.previous_hashes_by_session.get(&turn.session_id);
349 if turn.compaction.is_some() && previous_hashes.is_none() {
350 bail!(
351 "fidelity verification cannot recover compaction prefix for session {}",
352 turn.session_id
353 );
354 }
355 if let (Some(previous_hashes), Some(previous_input_length)) =
356 (previous_hashes, previous_input_length)
357 {
358 let verifiable_blocks = if let Some(compaction) = &turn.compaction {
359 previous_input_length.min(compaction.pre_tokens.saturating_sub(1)) / block_size
360 } else {
361 let cached_blocks = turn.cache_read_input_tokens.unwrap_or(0) / block_size;
362 cached_blocks
363 .min(previous_input_length / block_size)
364 .min(previous_hashes.len())
365 .min(input_sequence_hashes.len())
366 };
367 if previous_hashes[..verifiable_blocks] != input_sequence_hashes[..verifiable_blocks] {
368 bail!(
369 "fidelity verification cached prefix mismatch for session {} turn {}",
370 turn.session_id,
371 turn.turn_index
372 );
373 }
374 self.cache_prefix_blocks_verified += verifiable_blocks;
375 if turn.compaction.is_some() {
376 if verifiable_blocks == 0 {
377 bail!(
378 "fidelity verification found no recoverable compaction prefix blocks for session {}",
379 turn.session_id
380 );
381 }
382 self.compaction_prefix_blocks_verified += verifiable_blocks;
383 } else if previous_was_compaction {
384 let cached_tokens = turn.cache_read_input_tokens.unwrap_or(0);
385 let cached_blocks = cached_tokens / block_size;
386 let cache_creation_tokens = turn.cache_creation_input_tokens.unwrap_or(0);
387 if cached_blocks == 0
388 || cache_creation_tokens == 0
389 || cached_tokens > previous_input_length
390 || verifiable_blocks != cached_blocks
391 {
392 bail!(
393 "fidelity verification found post-compaction cache miss for session {}",
394 turn.session_id
395 );
396 }
397 self.post_compaction_prefix_blocks_verified += verifiable_blocks;
398 self.expected_next_cache_read_by_session.insert(
399 turn.session_id.clone(),
400 cached_tokens.saturating_add(cache_creation_tokens),
401 );
402 }
403 }
404 if turn.compaction.is_none()
405 && !previous_was_compaction
406 && let Some(expected_cache_read) = self
407 .expected_next_cache_read_by_session
408 .remove(&turn.session_id)
409 && turn.cache_read_input_tokens != Some(expected_cache_read)
410 {
411 bail!(
412 "fidelity verification expected {} post-compaction cache-read tokens for session {}, got {:?}",
413 expected_cache_read,
414 turn.session_id,
415 turn.cache_read_input_tokens
416 );
417 }
418 self.previous_hashes_by_session
419 .insert(turn.session_id.clone(), input_sequence_hashes.to_vec());
420 self.previous_input_length_by_session
421 .insert(turn.session_id.clone(), replay_tokens.len());
422 self.previous_was_compaction_by_session
423 .insert(turn.session_id.clone(), turn.compaction.is_some());
424
425 for tool in &turn.tools {
426 if tool.started_at_ms > tool.ended_at_ms {
427 bail!(
428 "invalid tool timing for {} in session {}",
429 tool.tool_call_id,
430 turn.session_id
431 );
432 }
433 self.tool_count += 1;
434 *self
435 .tools_by_class
436 .entry(tool.tool_class.clone())
437 .or_insert(0) += 1;
438 self.tool_errors += usize::from(tool.is_error);
439 self.child_links += usize::from(tool.child_session_id.is_some());
440 if let Some(child_session_id) = &tool.child_session_id {
441 self.child_session_references.push(child_session_id.clone());
442 }
443 self.background_tools += usize::from(tool.execution_mode == "background");
444 self.background_agents +=
445 usize::from(tool.execution_mode == "background" && tool.child_session_id.is_some());
446 if !matches!(tool.execution_mode.as_str(), "blocking" | "background") {
447 bail!(
448 "fidelity verification found invalid execution mode {} for {}",
449 tool.execution_mode,
450 tool.tool_call_id
451 );
452 }
453 if tool.child_session_id.as_deref() == Some(turn.export_session_id.as_str()) {
454 bail!(
455 "fidelity verification found self-referential child session for {}",
456 tool.tool_call_id
457 );
458 }
459 if let Some(consumer_turn_index) = tool.consumer_turn_index {
460 if consumer_turn_index <= turn.turn_index {
461 bail!(
462 "fidelity verification found non-forward consumer for {}",
463 tool.tool_call_id
464 );
465 }
466 self.causal_references.push((
467 turn.session_id.clone(),
468 consumer_turn_index,
469 tool.tool_call_id.clone(),
470 ));
471 }
472 }
473 Ok(())
474 }
475
476 fn finish(
477 self,
478 request_rows: usize,
479 tool_rows: usize,
480 sidecar_rows: usize,
481 ) -> Result<FidelityReport> {
482 for (session_id, consumer_turn_index, tool_call_id) in &self.causal_references {
483 let turn_count = self
484 .next_turn_by_session
485 .get(session_id)
486 .copied()
487 .unwrap_or(0);
488 if *consumer_turn_index >= turn_count {
489 bail!(
490 "fidelity verification found missing consumer turn {} for {}",
491 consumer_turn_index,
492 tool_call_id
493 );
494 }
495 }
496 let unresolved_child_sessions = self
497 .child_session_references
498 .iter()
499 .filter(|session_id| !self.export_sessions.contains(*session_id))
500 .count();
501 let source_request_rows = self.oracle.requests.len() + self.oracle.compactions.len();
502 let seen_request_rows = self.seen_requests.len() + self.seen_compactions.len();
503 if request_rows != source_request_rows
504 || request_rows != seen_request_rows
505 || self.seen_compactions.len() != self.oracle.compactions.len()
506 || sidecar_rows != request_rows
507 {
508 bail!(
509 "fidelity verification request mismatch: source={} ({} compactions), emitted={}, sidecar={}",
510 source_request_rows,
511 self.oracle.compactions.len(),
512 request_rows,
513 sidecar_rows
514 );
515 }
516 if tool_rows != self.oracle.paired_tools
517 || self.tool_count != self.oracle.paired_tools
518 || self.tool_errors != self.oracle.tool_errors
519 || self.tools_by_class != self.oracle.tools_by_class
520 {
521 bail!(
522 "fidelity verification tool mismatch: count={}/{}, errors={}/{}, classes_equal={}",
523 self.oracle.paired_tools,
524 tool_rows,
525 self.oracle.tool_errors,
526 self.tool_errors,
527 self.tools_by_class == self.oracle.tools_by_class
528 );
529 }
530 if self.child_links != self.oracle.child_links
531 || self.background_tools != self.oracle.background_tools
532 || self.background_agents != self.oracle.background_agents
533 {
534 bail!(
535 "fidelity verification agent mismatch: child_links={}/{}, background_tools={}/{}, background_agents={}/{}",
536 self.oracle.child_links,
537 self.child_links,
538 self.oracle.background_tools,
539 self.background_tools,
540 self.oracle.background_agents,
541 self.background_agents
542 );
543 }
544 Ok(FidelityReport {
545 requests_verified: request_rows,
546 compactions_verified: self.seen_compactions.len(),
547 usage_requests_verified: self.usage_requests,
548 tools_verified: tool_rows,
549 child_links_verified: self.child_links,
550 background_tools: self.background_tools,
551 background_agents: self.background_agents,
552 background_completions_missing: self.oracle.background_completions_missing,
553 background_titles_unreplayable: self.oracle.background_titles,
554 cache_prefix_blocks_verified: self.cache_prefix_blocks_verified,
555 compaction_prefix_blocks_verified: self.compaction_prefix_blocks_verified,
556 post_compaction_prefix_blocks_verified: self.post_compaction_prefix_blocks_verified,
557 unmatched_tool_calls: self.oracle.unmatched_tool_calls,
558 unmatched_tool_results: self.oracle.unmatched_tool_results,
559 unresolved_child_sessions,
560 })
561 }
562}
563
564pub fn write_streamed_request_trace_rows<F>(
565 output_path: &Path,
566 sidecar_path: &Path,
567 sessions: FxHashMap<String, Vec<TraceRecord>>,
568 preserve_session_ids: bool,
569 tokenizer_factory: F,
570 config: ExportConfig,
571) -> Result<ExportStats>
572where
573 F: TokenizerFactory,
574{
575 if config.block_size == 0 {
576 bail!("block_size must be greater than 0");
577 }
578 if config.tokenizer_workers == 0 {
579 bail!("tokenizer_workers must be greater than 0");
580 }
581
582 let mut verifier = FidelityVerifier::new(build_source_fidelity_oracle(&sessions)?);
583 let mut parser_tokenizer = tokenizer_factory.create_worker()?;
584 let mut states = FxHashMap::default();
585 let mut heap = BinaryHeap::new();
586 let mut unscheduled_sessions = VecDeque::new();
587 let mut stats = ExportStats::default();
588
589 for (session_id, records) in sessions {
590 let mut builder =
591 SessionTurnBuilder::new(session_id.clone(), records, preserve_session_ids);
592 let Some(first_turn) = builder.next_turn(&mut parser_tokenizer)? else {
593 continue;
594 };
595
596 let head = HeadTurn {
597 turn: first_turn,
598 turn_key: 0,
599 scheduled: false,
600 ready: None,
601 };
602 states.insert(
603 session_id.clone(),
604 SessionState {
605 builder,
606 head: Some(head),
607 overlap_base: None,
608 replay_base: None,
609 next_turn_key: 1,
610 },
611 );
612 push_heap_entry(&mut heap, &session_id, states.get(&session_id).unwrap());
613 unscheduled_sessions.push_back(session_id);
614 }
615
616 if states.is_empty() {
617 write_empty_files(output_path, Some(sidecar_path))?;
618 stats.fidelity = verifier.finish(0, 0, 0)?;
619 return Ok(stats);
620 }
621
622 stats.max_heap_len = heap.len();
623 let trace_start_ms = states
624 .values()
625 .filter_map(|state| state.head.as_ref())
626 .map(|head| head.turn.request_start_ms)
627 .min()
628 .unwrap_or_default();
629 let mut output = create_writer(output_path)?;
630 let mut sidecar = create_writer(sidecar_path)?;
631
632 let (job_tx, job_rx) = bounded::<TokenizeJob>(config.tokenizer_workers);
633 let (result_tx, result_rx) = unbounded::<TokenizeResponse>();
634 let workers = spawn_tokenizer_workers(
635 tokenizer_factory,
636 config.tokenizer_workers,
637 job_rx,
638 result_tx,
639 );
640
641 let mut inflight_jobs = 0_usize;
642 while !heap.is_empty() {
643 schedule_pending_jobs(
644 &mut states,
645 &mut unscheduled_sessions,
646 &job_tx,
647 &mut inflight_jobs,
648 config.delta_overlap_words,
649 config.tokenizer_workers,
650 )?;
651
652 let Some(Reverse(entry)) = heap.peek() else {
653 break;
654 };
655 let head_ready = states
656 .get(&entry.session_id)
657 .and_then(|state| state.head.as_ref())
658 .and_then(|head| head.ready.as_ref())
659 .is_some();
660 if !head_ready {
661 let response = result_rx
662 .recv()
663 .map_err(|_| anyhow!("tokenizer worker channel closed unexpectedly"))?;
664 inflight_jobs = inflight_jobs.saturating_sub(1);
665 apply_tokenize_response(&mut states, response)?;
666 continue;
667 }
668
669 let Reverse(entry) = heap.pop().unwrap();
670 let session_id = entry.session_id.clone();
671 let (turn, ready_turn) = {
672 let state = states
673 .get_mut(&session_id)
674 .ok_or_else(|| anyhow!("missing session state for {}", session_id))?;
675 let mut head = state
676 .head
677 .take()
678 .ok_or_else(|| anyhow!("missing head for session {}", session_id))?;
679 let ready_turn = head
680 .ready
681 .take()
682 .ok_or_else(|| anyhow!("missing tokenized result for session {}", session_id))?;
683 (head.turn, ready_turn)
684 };
685
686 let next_turn = {
687 let state = states
688 .get_mut(&session_id)
689 .ok_or_else(|| anyhow!("missing session state for {}", session_id))?;
690 state.builder.next_turn(&mut parser_tokenizer)?
691 };
692 let replay_tokens = {
693 let state = states
694 .get(&session_id)
695 .ok_or_else(|| anyhow!("missing session state for {}", session_id))?;
696 materialize_replay_tokens(&turn, &ready_turn.tokens, state.replay_base.as_deref())
697 };
698 let input_sequence_hashes = sequence_hashes_for_tokens(&replay_tokens, config.block_size)?;
699 verifier.observe(
700 &turn,
701 &replay_tokens,
702 &input_sequence_hashes,
703 config.block_size,
704 )?;
705 let request_id = turn.compaction.as_ref().map_or_else(
706 || canonical_request_id(&turn.export_session_id, turn.turn_index),
707 |compaction| {
708 canonical_compaction_request_id(&turn.export_session_id, compaction.sequence)
709 },
710 );
711 let mut agent_context = Map::from_iter([(
712 "session_id".to_string(),
713 Value::String(turn.export_session_id.clone()),
714 )]);
715 if let Some(parent_session_id) = &turn.export_parent_session_id {
716 agent_context.insert(
717 "parent_session_id".to_string(),
718 Value::String(parent_session_id.clone()),
719 );
720 }
721 let mut request = Map::from_iter([
722 ("request_id".to_string(), json!(request_id)),
723 ("model".to_string(), json!(turn.model)),
724 ("input_tokens".to_string(), json!(replay_tokens.len())),
725 ("output_tokens".to_string(), json!(turn.output_length)),
726 (
727 "request_received_ms".to_string(),
728 json!(nonnegative_ms(turn.request_start_ms)),
729 ),
730 (
731 "total_time_ms".to_string(),
732 json!((turn.assistant_end_ms - turn.request_start_ms).max(0) as f64),
733 ),
734 (
735 "replay".to_string(),
736 json!({
737 "trace_block_size": config.block_size,
738 "input_length": replay_tokens.len(),
739 "input_sequence_hashes": input_sequence_hashes,
740 }),
741 ),
742 ]);
743 if let Some(cached_tokens) = turn.cache_read_input_tokens {
744 request.insert("cached_tokens".to_string(), json!(cached_tokens));
745 }
746 if let Some(compaction) = &turn.compaction {
747 request.insert(
748 "claude".to_string(),
749 json!({
750 "compaction": {
751 "trigger": compaction.trigger,
752 "pre_tokens": compaction.pre_tokens,
753 "post_tokens": compaction.post_tokens,
754 "duration_ms": compaction.duration_ms,
755 "cache_fidelity": "recoverable_cache_safe_prefix",
756 "output_fidelity": "tokenized_compact_summary",
757 }
758 }),
759 );
760 }
761 let event = json!({
762 "schema": "dynamo.request.trace.v1",
763 "event_type": "request_end",
764 "event_time_unix_ms": nonnegative_ms(turn.assistant_end_ms),
765 "event_source": "harness",
766 "agent_context": agent_context,
767 "request": request,
768 });
769 let row = json!({
770 "timestamp": nonnegative_ms(turn.assistant_end_ms - trace_start_ms),
771 "event": event,
772 });
773
774 write_json_line(&mut output, &row)?;
775 for tool in &turn.tools {
776 let event_type = if tool.is_error {
777 "tool_error"
778 } else {
779 "tool_end"
780 };
781 let claude = ClaudeToolReplayMetadata {
782 source_request_id: request_id.clone(),
783 consumer_request_id: tool
784 .consumer_turn_index
785 .map(|turn_index| canonical_request_id(&turn.export_session_id, turn_index)),
786 child_session_id: tool.child_session_id.clone(),
787 execution_mode: tool.execution_mode.clone(),
788 };
789 let tool_row = json!({
790 "timestamp": nonnegative_ms(tool.ended_at_ms - trace_start_ms),
791 "event": {
792 "schema": "dynamo.request.trace.v1",
793 "event_type": event_type,
794 "event_time_unix_ms": nonnegative_ms(tool.ended_at_ms),
795 "event_source": "harness",
796 "agent_context": agent_context,
797 "tool": {
798 "tool_call_id": tool.tool_call_id,
799 "tool_class": tool.tool_class,
800 "claude": claude,
801 "started_at_unix_ms": nonnegative_ms(tool.started_at_ms),
802 "ended_at_unix_ms": nonnegative_ms(tool.ended_at_ms),
803 "duration_ms": (tool.ended_at_ms - tool.started_at_ms).max(0) as f64,
804 "status": if tool.is_error { "error" } else { "succeeded" },
805 "output_bytes": tool.output_bytes,
806 "error_type": if tool.is_error { Some("claude_tool_error") } else { None },
807 }
808 }
809 });
810 write_json_line(&mut output, &tool_row)?;
811 stats.tool_row_count += 1;
812 }
813 write_json_line(&mut sidecar, &turn.sidecar)?;
814 stats.row_count += 1;
815 stats.sidecar_count += 1;
816
817 let state = states
818 .get_mut(&session_id)
819 .ok_or_else(|| anyhow!("missing session state for {}", session_id))?;
820 state.overlap_base = Some(OverlapBase {
821 previous_text: ready_turn.current_text,
822 previous_tokens: ready_turn.tokens,
823 });
824 state.replay_base = Some(replay_tokens);
825
826 if let Some(next_turn) = next_turn {
827 let turn_key = state.next_turn_key;
828 state.next_turn_key += 1;
829 state.head = Some(HeadTurn {
830 turn: next_turn,
831 turn_key,
832 scheduled: false,
833 ready: None,
834 });
835 push_heap_entry(&mut heap, &session_id, state);
836 unscheduled_sessions.push_back(session_id);
837 stats.max_heap_len = stats.max_heap_len.max(heap.len());
838 continue;
839 }
840
841 states.remove(&session_id);
842 }
843
844 drop(job_tx);
845 for worker in workers {
846 worker
847 .join()
848 .map_err(|_| anyhow!("tokenizer worker panicked"))?;
849 }
850 stats.fidelity = verifier.finish(stats.row_count, stats.tool_row_count, stats.sidecar_count)?;
851 output.flush()?;
852 sidecar.flush()?;
853 Ok(stats)
854}
855
856fn create_writer(path: &Path) -> Result<BufWriter<File>> {
857 if let Some(parent) = path.parent() {
858 std::fs::create_dir_all(parent)?;
859 }
860 Ok(BufWriter::new(File::create(path)?))
861}
862
863fn write_json_line(writer: &mut impl Write, value: &impl Serialize) -> Result<()> {
864 serde_json::to_writer(&mut *writer, value)?;
865 writer.write_all(b"\n")?;
866 Ok(())
867}
868
869fn nonnegative_ms(value: i64) -> u64 {
870 value.max(0) as u64
871}
872
873fn canonical_request_id(session_id: &str, turn_index: usize) -> String {
874 format!("claude:{session_id}:{turn_index}")
875}
876
877fn canonical_compaction_request_id(session_id: &str, sequence: usize) -> String {
878 format!("claude:{session_id}:compact:{sequence}")
879}
880
881fn materialize_replay_tokens(
882 turn: &TurnDraft,
883 rendered_tokens: &[u32],
884 previous_tokens: Option<&[u32]>,
885) -> Vec<u32> {
886 let Some(input_length) = turn.observed_input_length else {
887 return rendered_tokens.to_vec();
888 };
889
890 if turn.compaction.is_some() {
891 let shared_length = previous_tokens
892 .map(|tokens| tokens.len())
893 .unwrap_or_default()
894 .min(input_length.saturating_sub(1));
895 let mut tokens = Vec::with_capacity(input_length);
896 if let Some(previous_tokens) = previous_tokens {
897 tokens.extend_from_slice(&previous_tokens[..shared_length]);
898 }
899 while tokens.len() < input_length {
900 tokens.push(synthetic_token(
901 &turn.export_session_id,
902 turn.turn_index,
903 tokens.len(),
904 rendered_tokens,
905 ));
906 }
907 return tokens;
908 }
909
910 let cached_length = turn.cache_read_input_tokens.unwrap_or(0).min(input_length);
911 let mut tokens = Vec::with_capacity(input_length);
912 if let Some(previous_tokens) = previous_tokens {
913 tokens.extend_from_slice(&previous_tokens[..cached_length.min(previous_tokens.len())]);
914 }
915 while tokens.len() < cached_length {
916 tokens.push(synthetic_token(
917 &turn.export_session_id,
918 turn.turn_index.saturating_sub(1),
919 tokens.len(),
920 rendered_tokens,
921 ));
922 }
923 while tokens.len() < input_length {
924 tokens.push(synthetic_token(
925 &turn.export_session_id,
926 turn.turn_index,
927 tokens.len(),
928 rendered_tokens,
929 ));
930 }
931 tokens
932}
933
934fn synthetic_token(
935 session_id: &str,
936 turn_index: usize,
937 position: usize,
938 rendered_tokens: &[u32],
939) -> u32 {
940 let mut hash = 0x811c_9dc5_u32;
941 for byte in session_id.bytes() {
942 hash = (hash ^ u32::from(byte)).wrapping_mul(0x0100_0193);
943 }
944 hash = (hash ^ turn_index as u32).wrapping_mul(0x0100_0193);
945 hash = (hash ^ position as u32).wrapping_mul(0x0100_0193);
946 if rendered_tokens.is_empty() {
947 hash
948 } else {
949 hash ^ rendered_tokens[position % rendered_tokens.len()]
950 }
951}
952
953fn push_heap_entry(
954 heap: &mut BinaryHeap<Reverse<HeapEntry>>,
955 session_id: &str,
956 state: &SessionState,
957) {
958 if let Some(head) = state.head.as_ref() {
959 heap.push(Reverse(HeapEntry {
960 request_start_ms: head.turn.request_start_ms,
961 turn_index: head.turn.turn_index,
962 export_session_id: head.turn.export_session_id.clone(),
963 session_id: session_id.to_string(),
964 }));
965 }
966}
967
968fn schedule_pending_jobs(
969 states: &mut FxHashMap<String, SessionState>,
970 unscheduled_sessions: &mut VecDeque<String>,
971 job_tx: &Sender<TokenizeJob>,
972 inflight_jobs: &mut usize,
973 overlap_words: usize,
974 worker_limit: usize,
975) -> Result<()> {
976 while *inflight_jobs < worker_limit {
977 let Some(session_id) = unscheduled_sessions.pop_front() else {
978 return Ok(());
979 };
980 let Some(state) = states.get_mut(&session_id) else {
981 continue;
982 };
983 let Some(head) = state.head.as_mut() else {
984 continue;
985 };
986 if head.scheduled || head.ready.is_some() {
987 continue;
988 }
989
990 let overlap_base = state.overlap_base.take();
991 let current_text = std::mem::take(&mut head.turn.input_text);
992 let (overlap_start, previous_overlap_text, previous_tokens) =
993 prepare_overlap_inputs(overlap_base, ¤t_text, overlap_words);
994 let job = TokenizeJob {
995 session_id: session_id.clone(),
996 turn_key: head.turn_key,
997 current_text,
998 overlap_start,
999 previous_overlap_text,
1000 previous_tokens,
1001 overlap_words,
1002 };
1003 job_tx
1004 .send(job)
1005 .map_err(|_| anyhow!("failed to schedule tokenization job"))?;
1006 head.scheduled = true;
1007 *inflight_jobs += 1;
1008 }
1009 Ok(())
1010}
1011
1012fn apply_tokenize_response(
1013 states: &mut FxHashMap<String, SessionState>,
1014 response: TokenizeResponse,
1015) -> Result<()> {
1016 let Some(state) = states.get_mut(&response.session_id) else {
1017 return Ok(());
1018 };
1019 let Some(head) = state.head.as_mut() else {
1020 return Ok(());
1021 };
1022 if head.turn_key != response.turn_key {
1023 return Ok(());
1024 }
1025 head.scheduled = false;
1026 match response.outcome {
1027 Ok(ready) => {
1028 head.ready = Some(ready);
1029 Ok(())
1030 }
1031 Err(message) => bail!("{message}"),
1032 }
1033}
1034
1035fn prepare_overlap_inputs(
1036 overlap_base: Option<OverlapBase>,
1037 current_text: &str,
1038 overlap_words: usize,
1039) -> (Option<usize>, Option<String>, Option<Vec<u32>>) {
1040 if overlap_words == 0 {
1041 return (None, None, None);
1042 }
1043 let Some(overlap_base) = overlap_base else {
1044 return (None, None, None);
1045 };
1046 if !current_text.starts_with(&overlap_base.previous_text) {
1047 return (None, None, None);
1048 }
1049
1050 let overlap_start = last_word_overlap_start(&overlap_base.previous_text, overlap_words);
1051 (
1052 Some(overlap_start),
1053 Some(overlap_base.previous_text[overlap_start..].to_string()),
1054 Some(overlap_base.previous_tokens),
1055 )
1056}
1057
1058fn spawn_tokenizer_workers<F>(
1059 factory: F,
1060 worker_count: usize,
1061 job_rx: Receiver<TokenizeJob>,
1062 result_tx: Sender<TokenizeResponse>,
1063) -> Vec<JoinHandle<()>>
1064where
1065 F: TokenizerFactory,
1066{
1067 (0..worker_count)
1068 .map(|_| {
1069 let job_rx = job_rx.clone();
1070 let result_tx = result_tx.clone();
1071 let factory = factory.clone();
1072 thread::spawn(move || {
1073 let mut tokenizer = match factory.create_worker() {
1074 Ok(tokenizer) => tokenizer,
1075 Err(error) => {
1076 let _ = result_tx.send(TokenizeResponse {
1077 session_id: "__worker_init__".to_string(),
1078 turn_key: 0,
1079 outcome: Err(format!(
1080 "failed to initialize tokenizer worker: {error:#}"
1081 )),
1082 });
1083 return;
1084 }
1085 };
1086 while let Ok(job) = job_rx.recv() {
1087 let outcome = tokenize_job(&mut tokenizer, &job)
1088 .map(|tokens| ReadyTurn {
1089 current_text: job.current_text,
1090 tokens,
1091 })
1092 .map_err(|error| {
1093 format!("failed to tokenize session {}: {error:#}", job.session_id)
1094 });
1095 let _ = result_tx.send(TokenizeResponse {
1096 session_id: job.session_id,
1097 turn_key: job.turn_key,
1098 outcome,
1099 });
1100 }
1101 })
1102 })
1103 .collect()
1104}
1105
1106fn tokenize_job(tokenizer: &mut impl TokenizerWorker, job: &TokenizeJob) -> Result<Vec<u32>> {
1107 let Some(overlap_start) = job.overlap_start else {
1108 return tokenizer.encode(&job.current_text);
1109 };
1110 let Some(previous_overlap_text) = job.previous_overlap_text.as_deref() else {
1111 return tokenizer.encode(&job.current_text);
1112 };
1113 let Some(previous_tokens) = job.previous_tokens.as_deref() else {
1114 return tokenizer.encode(&job.current_text);
1115 };
1116 if job.overlap_words == 0 || !job.current_text.is_char_boundary(overlap_start) {
1117 return tokenizer.encode(&job.current_text);
1118 }
1119
1120 let previous_overlap_tokens = tokenizer.encode(previous_overlap_text)?;
1121 let prefix_token_count = previous_tokens
1122 .len()
1123 .saturating_sub(previous_overlap_tokens.len());
1124 let suffix_tokens = tokenizer.encode(&job.current_text[overlap_start..])?;
1125 let mut merged = Vec::with_capacity(prefix_token_count + suffix_tokens.len());
1126 merged.extend_from_slice(&previous_tokens[..prefix_token_count]);
1127 merged.extend(suffix_tokens);
1128 Ok(merged)
1129}
1130
1131#[cfg(test)]
1132mod tests {
1133 use super::{
1134 ExportConfig, HeadTurn, ReadyTurn, SessionState, TurnDraft, apply_tokenize_response,
1135 write_streamed_request_trace_rows,
1136 };
1137 use crate::coding::claude::parser::{SessionTurnBuilder, TraceRecord};
1138 use crate::coding::tokenizer::{TokenizerFactory, TokenizerWorker};
1139 use anyhow::Result;
1140 use rustc_hash::FxHashMap;
1141 use serde_json::{Value, json};
1142 use std::sync::{Arc, Mutex};
1143 use std::thread;
1144 use std::time::Duration;
1145 use tempfile::TempDir;
1146
1147 #[derive(Clone, Default)]
1148 struct StubFactory {
1149 calls: Arc<Mutex<Vec<String>>>,
1150 }
1151
1152 struct StubWorker {
1153 calls: Arc<Mutex<Vec<String>>>,
1154 }
1155
1156 impl TokenizerFactory for StubFactory {
1157 type Worker = StubWorker;
1158
1159 fn create_worker(&self) -> Result<Self::Worker> {
1160 Ok(StubWorker {
1161 calls: self.calls.clone(),
1162 })
1163 }
1164 }
1165
1166 impl TokenizerWorker for StubWorker {
1167 fn encode(&mut self, text: &str) -> Result<Vec<u32>> {
1168 if text.contains("slow") {
1169 thread::sleep(Duration::from_millis(20));
1170 }
1171 self.calls.lock().unwrap().push(text.to_string());
1172 Ok(text
1173 .split_whitespace()
1174 .map(|word| word.len() as u32)
1175 .collect())
1176 }
1177 }
1178
1179 fn make_record(
1180 session_id: &str,
1181 row_type: &str,
1182 timestamp_ms: i64,
1183 source_order: u64,
1184 raw: Value,
1185 ) -> TraceRecord {
1186 TraceRecord {
1187 session_id: session_id.to_string(),
1188 parent_session_id: None,
1189 row_type: row_type.to_string(),
1190 timestamp_ms,
1191 source_order,
1192 raw,
1193 }
1194 }
1195
1196 #[test]
1197 fn stale_result_is_dropped_by_turn_key() {
1198 let mut states = FxHashMap::default();
1199 states.insert(
1200 "session-a".to_string(),
1201 SessionState {
1202 builder: SessionTurnBuilder::new("session-a".to_string(), Vec::new(), true),
1203 head: Some(HeadTurn {
1204 turn: TurnDraft {
1205 session_id: "session-a".to_string(),
1206 source_request_id: "req-1".to_string(),
1207 export_session_id: "session-a".to_string(),
1208 export_parent_session_id: None,
1209 turn_index: 1,
1210 model: "test-model".to_string(),
1211 input_text: String::new(),
1212 output_length: 1,
1213 observed_input_length: None,
1214 cache_read_input_tokens: None,
1215 cache_creation_input_tokens: None,
1216 request_start_ms: 1,
1217 assistant_start_ms: 1,
1218 assistant_end_ms: 2,
1219 delay_ms: None,
1220 tools: Vec::new(),
1221 sidecar: json!({}),
1222 compaction: None,
1223 },
1224 turn_key: 9,
1225 scheduled: true,
1226 ready: None,
1227 }),
1228 overlap_base: None,
1229 replay_base: None,
1230 next_turn_key: 10,
1231 },
1232 );
1233
1234 apply_tokenize_response(
1235 &mut states,
1236 super::TokenizeResponse {
1237 session_id: "session-a".to_string(),
1238 turn_key: 7,
1239 outcome: Ok(ReadyTurn {
1240 current_text: "stale".to_string(),
1241 tokens: vec![1],
1242 }),
1243 },
1244 )
1245 .unwrap();
1246
1247 assert!(
1248 states
1249 .get("session-a")
1250 .unwrap()
1251 .head
1252 .as_ref()
1253 .unwrap()
1254 .ready
1255 .is_none()
1256 );
1257 }
1258
1259 #[test]
1260 fn streamed_writer_preserves_global_order_with_parallel_tokenization() {
1261 let temp = TempDir::new().unwrap();
1262 let output_path = temp.path().join("trace.jsonl");
1263 let sidecar_path = temp.path().join("trace.sidecar.jsonl");
1264 let mut sessions = FxHashMap::default();
1265 sessions.insert(
1266 "session-a".to_string(),
1267 vec![
1268 make_record(
1269 "session-a",
1270 "user",
1271 1_000,
1272 0,
1273 json!({"type":"user","message":{"role":"user","content":"slow first a"}}),
1274 ),
1275 make_record(
1276 "session-a",
1277 "assistant",
1278 2_000,
1279 1,
1280 json!({"type":"assistant","message":{"id":"a-1","content":[{"type":"text","text":"done a"}],"usage":{"input_tokens":4,"cache_read_input_tokens":0,"cache_creation_input_tokens":0,"output_tokens":3}}}),
1281 ),
1282 make_record(
1283 "session-a",
1284 "user",
1285 2_100,
1286 2,
1287 json!({"type":"user","message":{"role":"user","content":"follow a"}}),
1288 ),
1289 make_record(
1290 "session-a",
1291 "assistant",
1292 2_200,
1293 3,
1294 json!({"type":"assistant","message":{"id":"a-2","content":[{"type":"text","text":"done a 2"}],"usage":{"input_tokens":2,"cache_read_input_tokens":4,"cache_creation_input_tokens":0,"output_tokens":4}}}),
1295 ),
1296 ],
1297 );
1298 sessions.insert(
1299 "session-b".to_string(),
1300 vec![
1301 make_record(
1302 "session-b",
1303 "user",
1304 900,
1305 4,
1306 json!({"type":"user","message":{"role":"user","content":"first b"}}),
1307 ),
1308 make_record(
1309 "session-b",
1310 "assistant",
1311 1_100,
1312 5,
1313 json!({"type":"assistant","message":{"id":"b-1","content":[{"type":"text","text":"done b"}],"usage":{"output_tokens":2}}}),
1314 ),
1315 ],
1316 );
1317
1318 let stats = write_streamed_request_trace_rows(
1319 &output_path,
1320 &sidecar_path,
1321 sessions,
1322 true,
1323 StubFactory::default(),
1324 ExportConfig {
1325 block_size: 2,
1326 delta_overlap_words: 50,
1327 tokenizer_workers: 2,
1328 },
1329 )
1330 .unwrap();
1331
1332 let rows = std::fs::read_to_string(&output_path)
1333 .unwrap()
1334 .lines()
1335 .map(|line| serde_json::from_str::<Value>(line).unwrap())
1336 .collect::<Vec<_>>();
1337 let sidecar_rows = std::fs::read_to_string(&sidecar_path)
1338 .unwrap()
1339 .lines()
1340 .map(|line| serde_json::from_str::<Value>(line).unwrap())
1341 .collect::<Vec<_>>();
1342
1343 assert_eq!(stats.row_count, 3);
1344 assert_eq!(stats.sidecar_count, 3);
1345 assert!(stats.max_heap_len <= 2);
1346 assert_eq!(rows.len(), 3);
1347 assert_eq!(sidecar_rows.len(), 3);
1348 assert_eq!(rows[0]["event"]["agent_context"]["session_id"], "session-b");
1349 assert_eq!(rows[1]["event"]["agent_context"]["session_id"], "session-a");
1350 assert!(
1351 rows[1]["event"]["agent_context"]
1352 .get("session_final")
1353 .is_none()
1354 );
1355 assert_eq!(rows[2]["event"]["request"]["request_received_ms"], 2_100);
1356 assert_eq!(rows[1]["event"]["request"]["replay"]["input_length"], 4);
1357 assert_eq!(rows[2]["event"]["request"]["replay"]["input_length"], 6);
1358 let first_hashes = rows[1]["event"]["request"]["replay"]["input_sequence_hashes"]
1359 .as_array()
1360 .unwrap();
1361 let second_hashes = rows[2]["event"]["request"]["replay"]["input_sequence_hashes"]
1362 .as_array()
1363 .unwrap();
1364 assert_eq!(first_hashes.as_slice(), &second_hashes[..2]);
1365 }
1366
1367 #[test]
1368 fn streamed_writer_replays_cache_safe_compaction() {
1369 use dynamo_data_gen::request_trace::{
1370 agentic::lower_agentic_mooncake_rows, load::load_request_trace_records,
1371 };
1372
1373 let temp = TempDir::new().unwrap();
1374 let output_path = temp.path().join("trace.jsonl");
1375 let sidecar_path = temp.path().join("trace.sidecar.jsonl");
1376 let mut sessions = FxHashMap::default();
1377 sessions.insert(
1378 "session-a".to_string(),
1379 vec![
1380 make_record(
1381 "session-a",
1382 "user",
1383 1_000,
1384 0,
1385 json!({"type":"user","message":{"role":"user","content":"first prompt"}}),
1386 ),
1387 make_record(
1388 "session-a",
1389 "assistant",
1390 1_100,
1391 1,
1392 json!({"type":"assistant","requestId":"req-0","message":{"id":"a-0","model":"test-model","content":[{"type":"text","text":"first answer"}],"usage":{"input_tokens":8,"cache_read_input_tokens":0,"cache_creation_input_tokens":0,"output_tokens":2}}}),
1393 ),
1394 make_record(
1395 "session-a",
1396 "system",
1397 2_000,
1398 2,
1399 json!({"type":"system","subtype":"compact_boundary","compactMetadata":{"trigger":"manual","preTokens":10,"postTokens":3,"durationMs":500}}),
1400 ),
1401 make_record(
1402 "session-a",
1403 "user",
1404 2_000,
1405 3,
1406 json!({"type":"user","isCompactSummary":true,"message":{"role":"user","content":"compact summary"}}),
1407 ),
1408 make_record(
1409 "session-a",
1410 "assistant",
1411 2_100,
1412 4,
1413 json!({"type":"assistant","requestId":"req-1","message":{"id":"a-1","model":"test-model","content":[{"type":"text","text":"after compact"}],"usage":{"input_tokens":2,"cache_read_input_tokens":4,"cache_creation_input_tokens":6,"output_tokens":2}}}),
1414 ),
1415 make_record(
1416 "session-a",
1417 "user",
1418 2_200,
1419 5,
1420 json!({"type":"user","message":{"role":"user","content":"next prompt"}}),
1421 ),
1422 make_record(
1423 "session-a",
1424 "assistant",
1425 2_300,
1426 6,
1427 json!({"type":"assistant","requestId":"req-2","message":{"id":"a-2","model":"test-model","content":[{"type":"text","text":"next answer"}],"usage":{"input_tokens":2,"cache_read_input_tokens":10,"cache_creation_input_tokens":2,"output_tokens":2}}}),
1428 ),
1429 ],
1430 );
1431
1432 let no_prefix_error = write_streamed_request_trace_rows(
1433 &temp.path().join("no-prefix.jsonl"),
1434 &temp.path().join("no-prefix.sidecar.jsonl"),
1435 sessions.clone(),
1436 true,
1437 StubFactory::default(),
1438 ExportConfig {
1439 block_size: 16,
1440 delta_overlap_words: 50,
1441 tokenizer_workers: 1,
1442 },
1443 )
1444 .unwrap_err();
1445 assert!(
1446 no_prefix_error
1447 .to_string()
1448 .contains("no recoverable compaction prefix")
1449 );
1450
1451 let mut no_summary_write = sessions.clone();
1452 let first_post = no_summary_write
1453 .get_mut("session-a")
1454 .unwrap()
1455 .iter_mut()
1456 .find(|record| record.raw["requestId"] == "req-1")
1457 .unwrap();
1458 first_post.raw["message"]["usage"]["cache_creation_input_tokens"] = json!(0);
1459 let no_summary_write_error = write_streamed_request_trace_rows(
1460 &temp.path().join("no-summary-write.jsonl"),
1461 &temp.path().join("no-summary-write.sidecar.jsonl"),
1462 no_summary_write,
1463 true,
1464 StubFactory::default(),
1465 ExportConfig {
1466 block_size: 2,
1467 delta_overlap_words: 50,
1468 tokenizer_workers: 1,
1469 },
1470 )
1471 .unwrap_err();
1472 assert!(
1473 no_summary_write_error
1474 .to_string()
1475 .contains("post-compaction cache miss")
1476 );
1477
1478 let stats = write_streamed_request_trace_rows(
1479 &output_path,
1480 &sidecar_path,
1481 sessions,
1482 true,
1483 StubFactory::default(),
1484 ExportConfig {
1485 block_size: 2,
1486 delta_overlap_words: 50,
1487 tokenizer_workers: 1,
1488 },
1489 )
1490 .unwrap();
1491
1492 let rows = std::fs::read_to_string(&output_path)
1493 .unwrap()
1494 .lines()
1495 .map(|line| serde_json::from_str::<Value>(line).unwrap())
1496 .collect::<Vec<_>>();
1497 let sidecars = std::fs::read_to_string(&sidecar_path)
1498 .unwrap()
1499 .lines()
1500 .map(|line| serde_json::from_str::<Value>(line).unwrap())
1501 .collect::<Vec<_>>();
1502
1503 assert_eq!(stats.row_count, 4);
1504 assert_eq!(stats.sidecar_count, 4);
1505 assert_eq!(stats.fidelity.compactions_verified, 1);
1506 assert_eq!(stats.fidelity.compaction_prefix_blocks_verified, 4);
1507 assert_eq!(stats.fidelity.post_compaction_prefix_blocks_verified, 2);
1508 assert_eq!(rows.len(), 4);
1509 assert_eq!(sidecars.len(), 4);
1510 assert_eq!(
1511 rows[1]["event"]["request"]["request_id"],
1512 "claude:session-a:compact:0"
1513 );
1514 assert_eq!(rows[1]["event"]["request"]["request_received_ms"], 1_500);
1515 assert_eq!(rows[1]["event"]["event_time_unix_ms"], 2_000);
1516 assert_eq!(rows[1]["event"]["request"]["total_time_ms"], 500.0);
1517 assert!(rows[1]["event"]["request"].get("cached_tokens").is_none());
1518 assert_eq!(rows[1]["event"]["request"]["replay"]["input_length"], 10);
1519 assert_eq!(
1520 rows[1]["event"]["request"]["claude"]["compaction"]["pre_tokens"],
1521 10
1522 );
1523 assert_eq!(
1524 rows[1]["event"]["request"]["claude"]["compaction"]["post_tokens"],
1525 3
1526 );
1527
1528 let hashes = rows
1529 .iter()
1530 .map(|row| {
1531 row["event"]["request"]["replay"]["input_sequence_hashes"]
1532 .as_array()
1533 .unwrap()
1534 })
1535 .collect::<Vec<_>>();
1536 assert_eq!(hashes[0], &hashes[1][..4]);
1537 assert_eq!(&hashes[1][..2], &hashes[2][..2]);
1538 assert_ne!(hashes[1][2], hashes[2][2]);
1539 assert_eq!(&hashes[2][..5], &hashes[3][..5]);
1540 assert_ne!(hashes[2][5], hashes[3][5]);
1541
1542 let loaded = load_request_trace_records(&[output_path]).unwrap();
1543 assert_eq!(loaded.requests.len(), 4);
1544 let mut agentic_rows = Vec::new();
1545 lower_agentic_mooncake_rows(loaded, |_, row| {
1546 agentic_rows.push(row);
1547 Ok(())
1548 })
1549 .unwrap();
1550 assert_eq!(agentic_rows.len(), 4);
1551 assert_eq!(agentic_rows[1].request_id, "claude:session-a:compact:0");
1552 }
1553
1554 #[test]
1555 fn streamed_writer_emits_canonical_tool_terminal_events() {
1556 use dynamo_data_gen::request_trace::load::load_request_trace_records;
1557
1558 let temp = TempDir::new().unwrap();
1559 let output_path = temp.path().join("trace.jsonl");
1560 let sidecar_path = temp.path().join("trace.sidecar.jsonl");
1561 let mut sessions = FxHashMap::default();
1562 sessions.insert(
1563 "session-a".to_string(),
1564 vec![
1565 make_record(
1566 "session-a",
1567 "user",
1568 1_000,
1569 0,
1570 json!({"type":"user","message":{"role":"user","content":"run"}}),
1571 ),
1572 make_record(
1573 "session-a",
1574 "assistant",
1575 1_100,
1576 1,
1577 json!({"type":"assistant","requestId":"req-1","message":{"id":"a-1","content":[{"type":"tool_use","id":"raw-1","name":"Bash","input":{}}],"usage":{"input_tokens":2,"cache_read_input_tokens":0,"cache_creation_input_tokens":0,"output_tokens":3}}}),
1578 ),
1579 make_record(
1580 "session-a",
1581 "user",
1582 1_200,
1583 2,
1584 json!({"type":"user","message":{"role":"user","content":[{"type":"tool_result","tool_use_id":"raw-1","content":"bad","is_error":true}]}}),
1585 ),
1586 make_record(
1587 "session-a",
1588 "ai-title",
1589 0,
1590 3,
1591 json!({"type":"ai-title","aiTitle":"Background title"}),
1592 ),
1593 ],
1594 );
1595
1596 let stats = write_streamed_request_trace_rows(
1597 &output_path,
1598 &sidecar_path,
1599 sessions,
1600 true,
1601 StubFactory::default(),
1602 ExportConfig {
1603 block_size: 2,
1604 delta_overlap_words: 50,
1605 tokenizer_workers: 1,
1606 },
1607 )
1608 .unwrap();
1609
1610 assert_eq!(stats.row_count, 1);
1611 assert_eq!(stats.tool_row_count, 1);
1612 assert_eq!(stats.fidelity.requests_verified, 1);
1613 assert_eq!(stats.fidelity.tools_verified, 1);
1614 assert_eq!(stats.fidelity.background_titles_unreplayable, 1);
1615 let rows = std::fs::read_to_string(&output_path).unwrap();
1616 assert!(rows.lines().any(|line| {
1617 let row: Value = serde_json::from_str(line).unwrap();
1618 row["event"]["event_type"] == "tool_error"
1619 && row["event"]["tool"]["tool_class"] == "Bash"
1620 }));
1621 let loaded = load_request_trace_records(&[output_path]).unwrap();
1622 assert_eq!(loaded.tools.len(), 1);
1623 }
1624
1625 #[test]
1626 fn request_trace_preserves_child_identity_and_anonymized_causality() {
1627 use dynamo_data_gen::request_trace::{
1628 agentic::lower_agentic_mooncake_rows, load::load_request_trace_records,
1629 };
1630
1631 let temp = TempDir::new().unwrap();
1632 let output_path = temp.path().join("trace.jsonl");
1633 let sidecar_path = temp.path().join("trace.sidecar.jsonl");
1634 let mut sessions = FxHashMap::default();
1635 sessions.insert(
1636 "root-session".to_string(),
1637 vec![
1638 make_record(
1639 "root-session",
1640 "user",
1641 1_000,
1642 0,
1643 json!({"type":"user","message":{"role":"user","content":"spawn child"}}),
1644 ),
1645 make_record(
1646 "root-session",
1647 "assistant",
1648 1_100,
1649 1,
1650 json!({"type":"assistant","requestId":"root-1","message":{"id":"root-1","content":[{"type":"tool_use","id":"agent-call","name":"Agent","input":{"run_in_background":true}}],"usage":{"output_tokens":2}}}),
1651 ),
1652 make_record(
1653 "root-session",
1654 "user",
1655 1_150,
1656 2,
1657 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"}]}}),
1658 ),
1659 make_record(
1660 "root-session",
1661 "user",
1662 1_300,
1663 4,
1664 json!({"type":"user","message":{"role":"user","content":"continue parent work"}}),
1665 ),
1666 make_record(
1667 "root-session",
1668 "assistant",
1669 1_400,
1670 5,
1671 json!({"type":"assistant","requestId":"root-2","message":{"id":"root-2","content":[{"type":"text","text":"working"}],"usage":{"output_tokens":1}}}),
1672 ),
1673 make_record(
1674 "root-session",
1675 "queue-operation",
1676 1_800,
1677 6,
1678 json!({"type":"queue-operation","operation":"enqueue","content":"<tool-use-id>agent-call</tool-use-id><status>completed</status>done"}),
1679 ),
1680 make_record(
1681 "root-session",
1682 "user",
1683 1_850,
1684 7,
1685 json!({"type":"user","message":{"role":"user","content":"child done"}}),
1686 ),
1687 make_record(
1688 "root-session",
1689 "assistant",
1690 1_950,
1691 8,
1692 json!({"type":"assistant","requestId":"root-3","message":{"id":"root-3","content":[{"type":"text","text":"finished"}],"usage":{"output_tokens":1}}}),
1693 ),
1694 ],
1695 );
1696 sessions.insert(
1697 "child-agent".to_string(),
1698 vec![
1699 make_record(
1700 "root-session",
1701 "user",
1702 1_200,
1703 3,
1704 json!({"type":"user","isSidechain":true,"agentId":"child-agent","message":{"role":"user","content":"investigate"}}),
1705 ),
1706 make_record(
1707 "root-session",
1708 "assistant",
1709 1_700,
1710 9,
1711 json!({"type":"assistant","isSidechain":true,"agentId":"child-agent","message":{"id":"child-1","content":[{"type":"text","text":"result"}],"usage":{"output_tokens":1}}}),
1712 ),
1713 ],
1714 );
1715
1716 let config = ExportConfig {
1717 block_size: 2,
1718 delta_overlap_words: 50,
1719 tokenizer_workers: 2,
1720 };
1721 let stats = write_streamed_request_trace_rows(
1722 &output_path,
1723 &sidecar_path,
1724 sessions.clone(),
1725 true,
1726 StubFactory::default(),
1727 config,
1728 )
1729 .unwrap();
1730
1731 let anonymous_stats = write_streamed_request_trace_rows(
1732 &temp.path().join("anonymous.jsonl"),
1733 &temp.path().join("anonymous.sidecar.jsonl"),
1734 sessions,
1735 false,
1736 StubFactory::default(),
1737 config,
1738 )
1739 .unwrap();
1740 assert_eq!(anonymous_stats.fidelity.requests_verified, 4);
1741
1742 assert_eq!(stats.fidelity.requests_verified, 4);
1743 assert_eq!(stats.fidelity.tools_verified, 1);
1744 assert_eq!(stats.fidelity.child_links_verified, 1);
1745 assert_eq!(stats.fidelity.background_tools, 1);
1746 assert_eq!(stats.fidelity.background_agents, 1);
1747
1748 let rows = std::fs::read_to_string(&output_path)
1749 .unwrap()
1750 .lines()
1751 .map(|line| serde_json::from_str::<Value>(line).unwrap())
1752 .collect::<Vec<_>>();
1753 let child = rows
1754 .iter()
1755 .find(|row| row["event"]["agent_context"]["session_id"] == "child-agent")
1756 .unwrap();
1757 assert_eq!(
1758 child["event"]["agent_context"]["parent_session_id"],
1759 "root-session"
1760 );
1761 assert!(
1762 child["event"]["agent_context"]
1763 .get("session_final")
1764 .is_none()
1765 );
1766
1767 let loaded = load_request_trace_records(&[output_path]).unwrap();
1768 let mut agentic_rows = Vec::new();
1769 lower_agentic_mooncake_rows(loaded, |_, row| {
1770 agentic_rows.push(row);
1771 Ok(())
1772 })
1773 .unwrap();
1774 assert_eq!(agentic_rows.len(), 4);
1775 let by_id = agentic_rows
1776 .iter()
1777 .map(|row| (row.request_id.as_str(), row))
1778 .collect::<std::collections::HashMap<_, _>>();
1779 assert_eq!(
1780 by_id["claude:root-session:0"].branches,
1781 vec!["claude:child-agent:0"]
1782 );
1783 assert_eq!(
1784 by_id["claude:child-agent:0"].request_kind.as_deref(),
1785 Some("background_agent")
1786 );
1787 assert_eq!(
1788 by_id["claude:root-session:1"].wait_for,
1789 vec!["claude:root-session:0"]
1790 );
1791 assert!(
1792 by_id["claude:root-session:2"]
1793 .wait_for
1794 .contains(&"claude:child-agent:0".to_string())
1795 );
1796 assert_eq!(by_id["claude:root-session:2"].tool_wait_ms, Some(100.0));
1797 }
1798}