Skip to main content

dynamo_bench/coding/claude/
export.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Claude-specific request-trace export orchestration.
5//!
6//! Handles session scheduling, parallel tokenization with text-overlap reuse,
7//! and global ordering across sessions.
8
9use 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/// Claude-only evidence used to reconstruct tool scheduling after export.
62///
63/// Live tool events cannot know their future consumer request. Claude's saved
64/// session can, so the exporter stores that post-hoc evidence under
65/// `tool.claude` without extending the live request-trace tool API.
66#[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, &current_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}