Skip to main content

aisimulate_core/replay/
evidence.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::collections::{BTreeMap, BTreeSet};
5
6use anyhow::{Result, ensure};
7use blake3::Hasher;
8use serde::Serialize;
9use uuid::Uuid;
10
11use crate::engine::{
12    PressureEvent as EnginePressureEvent, PressureKind as EnginePressureKind,
13    PressureState as SchedulerPressureState,
14};
15
16use crate::replay::{ReplayCaptureOptions, TraceCollector};
17
18#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize)]
19#[serde(rename_all = "snake_case")]
20pub enum WorkerPool {
21    Agg,
22    Prefill,
23    Decode,
24}
25
26impl WorkerPool {
27    const fn tag(self) -> u8 {
28        match self {
29            Self::Agg => 0,
30            Self::Prefill => 1,
31            Self::Decode => 2,
32        }
33    }
34
35    const fn as_str(self) -> &'static str {
36        match self {
37            Self::Agg => "agg",
38            Self::Prefill => "prefill",
39            Self::Decode => "decode",
40        }
41    }
42}
43
44#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
45#[serde(rename_all = "snake_case")]
46pub enum WorkerLifecycleTransitionKind {
47    WorkerStarting,
48    WorkerReady,
49    WorkerDraining,
50    WorkerRemoved,
51}
52
53#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
54pub struct WorkerLifecycleTransition {
55    pub worker_id: usize,
56    pub transition: WorkerLifecycleTransitionKind,
57    pub prior_state: Option<&'static str>,
58    pub state: &'static str,
59    #[serde(skip_serializing_if = "Option::is_none")]
60    pub reason: Option<&'static str>,
61    #[serde(skip_serializing_if = "Option::is_none")]
62    pub origin_operation_ordinal: Option<u64>,
63}
64
65#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)]
66pub struct WorkerPoolState {
67    pub active: Vec<usize>,
68    pub starting: Vec<usize>,
69    pub draining: Vec<usize>,
70}
71
72#[derive(Clone, Debug, PartialEq, Serialize)]
73pub struct LifecycleOperation {
74    pub operation_ordinal: u64,
75    pub at_ms: f64,
76    pub pool: WorkerPool,
77    pub cause: &'static str,
78    pub planner_tick_ordinal: Option<u64>,
79    pub origin_operation_ordinal: Option<u64>,
80    pub transitions: Vec<WorkerLifecycleTransition>,
81    pub state_after_batch: WorkerPoolState,
82    pub topology_released_request_uuids: Vec<String>,
83}
84
85#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
86#[serde(rename_all = "snake_case")]
87pub enum PressureKind {
88    VllmPreemption,
89    SglangRetraction,
90}
91
92#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
93pub struct EnginePressureState {
94    pub running_requests: usize,
95    pub waiting_requests: Option<usize>,
96    pub active_blocks: usize,
97}
98
99#[derive(Clone, Debug, PartialEq, Serialize)]
100pub struct PressureRecord {
101    pub pressure_ordinal: u64,
102    pub at_ms: f64,
103    pub pool: WorkerPool,
104    pub worker_id: u64,
105    pub dp_rank: u32,
106    pub kind: PressureKind,
107    pub request_uuid: String,
108    pub state_before: EnginePressureState,
109    pub state_after: EnginePressureState,
110    pub request_active_blocks_before: usize,
111    pub logical_available_blocks_before: Option<usize>,
112    pub required_blocks_before: Option<usize>,
113    pub readmitted_at_ms: Option<f64>,
114}
115
116#[derive(Clone, Debug, Default, PartialEq, Serialize)]
117pub struct PressureEvidence {
118    pub records: Vec<PressureRecord>,
119    pub vllm_preemptions_total: u64,
120    pub sglang_retractions_total: u64,
121}
122
123/// Replay boundary at which one batch of native KV observations was ingested.
124#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd, Serialize)]
125#[serde(rename_all = "snake_case")]
126pub enum KvIngestBoundary {
127    PassStart,
128    PassEnd,
129    SchedulerCommand,
130    OffloadTick,
131    WorkerLifecycle,
132}
133
134impl KvIngestBoundary {
135    const fn tag(self) -> u8 {
136        match self {
137            Self::PassStart => 0,
138            Self::PassEnd => 1,
139            Self::SchedulerCommand => 2,
140            Self::OffloadTick => 3,
141            Self::WorkerLifecycle => 4,
142        }
143    }
144
145    const fn as_str(self) -> &'static str {
146        match self {
147            Self::PassStart => "pass_start",
148            Self::PassEnd => "pass_end",
149            Self::SchedulerCommand => "scheduler_command",
150            Self::OffloadTick => "offload_tick",
151            Self::WorkerLifecycle => "worker_lifecycle",
152        }
153    }
154}
155
156#[derive(Clone, Debug, Default, PartialEq, Serialize)]
157pub struct KvIngestBoundaryStats {
158    pub batches: u64,
159    pub events: u64,
160    pub first_at_ms: f64,
161    pub last_at_ms: f64,
162}
163
164#[derive(Clone, Debug, Default, PartialEq, Serialize)]
165pub struct KvIngestEvidence {
166    pub encoding: String,
167    pub blake3_256: String,
168    pub batches: u64,
169    pub events: u64,
170    pub blocks: u64,
171    pub kind_counts: BTreeMap<String, u64>,
172    pub pool_counts: BTreeMap<String, u64>,
173    pub tier_counts: BTreeMap<String, u64>,
174    pub boundaries: BTreeMap<String, KvIngestBoundaryStats>,
175}
176
177/// Evidence owned and returned by one replay execution.
178#[derive(Clone, Debug, Default, PartialEq, Serialize)]
179pub struct OfflineRuntimeEvidence {
180    pub lifecycle_operations: Vec<LifecycleOperation>,
181    #[serde(skip_serializing_if = "Option::is_none")]
182    pub pressure: Option<PressureEvidence>,
183    #[serde(skip_serializing_if = "Option::is_none")]
184    pub kv_ingest: Option<KvIngestEvidence>,
185}
186
187/// Explicit execution-local evidence state. It deliberately is not stored in
188/// a thread-local: two replays on one thread cannot observe or overwrite each
189/// other's capture state.
190#[derive(Debug)]
191pub(crate) struct ReplayEvidenceCollector {
192    options: ReplayCaptureOptions,
193    lifecycle_operations: Vec<LifecycleOperation>,
194    pressure_records: Vec<PressureRecord>,
195    kv_ingest: Option<KvIngestAccumulator>,
196    outstanding_pressure: BTreeMap<(Uuid, WorkerPool), Vec<u64>>,
197    startup_origins: BTreeMap<(WorkerPool, usize), u64>,
198    drain_origins: BTreeMap<(WorkerPool, usize), u64>,
199}
200
201impl ReplayEvidenceCollector {
202    pub(crate) fn new(options: ReplayCaptureOptions) -> Self {
203        Self {
204            options,
205            lifecycle_operations: Vec::new(),
206            pressure_records: Vec::new(),
207            kv_ingest: options
208                .capture_canonical_evidence
209                .then(KvIngestAccumulator::new),
210            outstanding_pressure: BTreeMap::new(),
211            startup_origins: BTreeMap::new(),
212            drain_origins: BTreeMap::new(),
213        }
214    }
215
216    pub(crate) fn options(&self) -> ReplayCaptureOptions {
217        self.options
218    }
219
220    pub(crate) fn startup_origin(&self, pool: WorkerPool, worker_id: usize) -> Option<u64> {
221        self.startup_origins.get(&(pool, worker_id)).copied()
222    }
223
224    pub(crate) fn drain_origin(&self, pool: WorkerPool, worker_id: usize) -> Option<u64> {
225        self.drain_origins.get(&(pool, worker_id)).copied()
226    }
227
228    #[allow(clippy::too_many_arguments)]
229    pub(crate) fn record_lifecycle_operation(
230        &mut self,
231        at_ms: f64,
232        pool: WorkerPool,
233        cause: &'static str,
234        planner_tick_ordinal: Option<u64>,
235        origin_operation_ordinal: Option<u64>,
236        mut transitions: Vec<WorkerLifecycleTransition>,
237        state_after_batch: WorkerPoolState,
238        topology_released_request_uuids: Vec<Uuid>,
239    ) -> Option<u64> {
240        if !self.options.capture_lifecycle_evidence
241            || (transitions.is_empty() && topology_released_request_uuids.is_empty())
242        {
243            return None;
244        }
245
246        let operation_ordinal = u64::try_from(self.lifecycle_operations.len()).ok()?;
247        for transition in &mut transitions {
248            if transition.origin_operation_ordinal.is_none() {
249                transition.origin_operation_ordinal = Some(operation_ordinal);
250            }
251            let key = (pool, transition.worker_id);
252            match transition.transition {
253                WorkerLifecycleTransitionKind::WorkerStarting => {
254                    self.startup_origins.insert(key, operation_ordinal);
255                }
256                WorkerLifecycleTransitionKind::WorkerDraining => {
257                    self.drain_origins.insert(key, operation_ordinal);
258                }
259                WorkerLifecycleTransitionKind::WorkerReady => {
260                    self.startup_origins.remove(&key);
261                }
262                WorkerLifecycleTransitionKind::WorkerRemoved => {
263                    self.startup_origins.remove(&key);
264                    self.drain_origins.remove(&key);
265                }
266            }
267        }
268        let mut seen = BTreeSet::new();
269        self.lifecycle_operations.push(LifecycleOperation {
270            operation_ordinal,
271            at_ms,
272            pool,
273            cause,
274            planner_tick_ordinal,
275            origin_operation_ordinal,
276            transitions,
277            state_after_batch,
278            topology_released_request_uuids: topology_released_request_uuids
279                .into_iter()
280                .filter(|uuid| seen.insert(*uuid))
281                .map(|uuid| uuid.to_string())
282                .collect(),
283        });
284        Some(operation_ordinal)
285    }
286
287    /// Record one engine pressure transition and attach its stable ordinal to
288    /// the matching per-request report record.
289    #[allow(clippy::too_many_arguments)]
290    pub(crate) fn record_pressure(
291        &mut self,
292        collector: &mut TraceCollector,
293        at_ms: f64,
294        pool: WorkerPool,
295        worker_id: u64,
296        dp_rank: u32,
297        kind: PressureKind,
298        request_uuid: Uuid,
299        state_before: EnginePressureState,
300        state_after: EnginePressureState,
301        request_active_blocks_before: usize,
302        logical_available_blocks_before: Option<usize>,
303        required_blocks_before: Option<usize>,
304    ) -> Option<u64> {
305        if !self.options.capture_canonical_evidence {
306            return None;
307        }
308        let pressure_ordinal = u64::try_from(self.pressure_records.len()).ok()?;
309        self.pressure_records.push(PressureRecord {
310            pressure_ordinal,
311            at_ms,
312            pool,
313            worker_id,
314            dp_rank,
315            kind,
316            request_uuid: request_uuid.to_string(),
317            state_before,
318            state_after,
319            request_active_blocks_before,
320            logical_available_blocks_before,
321            required_blocks_before,
322            readmitted_at_ms: None,
323        });
324        self.outstanding_pressure
325            .entry((request_uuid, pool))
326            .or_default()
327            .push(pressure_ordinal);
328        collector.on_pressure_reference(request_uuid, pressure_ordinal);
329        Some(pressure_ordinal)
330    }
331
332    pub(crate) fn record_native_pressure(
333        &mut self,
334        collector: &mut TraceCollector,
335        pool: WorkerPool,
336        worker_id: u64,
337        dp_rank: u32,
338        event: EnginePressureEvent,
339    ) -> Option<u64> {
340        let kind = match event.kind {
341            EnginePressureKind::VllmPreemption => PressureKind::VllmPreemption,
342            EnginePressureKind::SglangRetraction => PressureKind::SglangRetraction,
343        };
344        self.record_pressure(
345            collector,
346            event.at_ms,
347            pool,
348            worker_id,
349            dp_rank,
350            kind,
351            event.request_id,
352            lower_pressure_state(event.state_before),
353            lower_pressure_state(event.state_after),
354            event.request_active_blocks_before,
355            event.logical_available_blocks_before,
356            event.required_blocks_before,
357        )
358    }
359
360    pub(crate) fn record_pressure_readmission(
361        &mut self,
362        request_uuid: Uuid,
363        pool: WorkerPool,
364        at_ms: f64,
365    ) {
366        if !self.options.capture_canonical_evidence {
367            return;
368        }
369        let key = (request_uuid, pool);
370        let Some(ordinals) = self.outstanding_pressure.get_mut(&key) else {
371            return;
372        };
373        let Some(pressure_ordinal) = ordinals.pop() else {
374            return;
375        };
376        let remove_key = ordinals.is_empty();
377        if remove_key {
378            self.outstanding_pressure.remove(&key);
379        }
380        if let Some(record) = self
381            .pressure_records
382            .get_mut(usize::try_from(pressure_ordinal).expect("pressure ordinal must fit usize"))
383        {
384            record.readmitted_at_ms = Some(at_ms);
385        }
386    }
387
388    pub(crate) fn record_kv_ingest(
389        &mut self,
390        pool: WorkerPool,
391        boundary: KvIngestBoundary,
392        at_ms: f64,
393        event_count: usize,
394        encode_events: impl FnOnce(&mut KvIngestEventEncoder<'_>) -> Result<()>,
395    ) -> Result<()> {
396        if !self.options.capture_canonical_evidence {
397            return Ok(());
398        }
399        ensure!(
400            at_ms.is_finite(),
401            "canonical KV ingestion rejects non-finite timestamp {at_ms}"
402        );
403        self.kv_ingest
404            .as_mut()
405            .expect("canonical KV accumulator was not initialized")
406            .record_batch(pool, boundary, at_ms, event_count, encode_events)
407    }
408
409    pub(crate) fn finish(self) -> OfflineRuntimeEvidence {
410        let Self {
411            options,
412            lifecycle_operations,
413            pressure_records,
414            kv_ingest,
415            ..
416        } = self;
417        let pressure = options.capture_canonical_evidence.then(|| {
418            let vllm_preemptions_total = pressure_records
419                .iter()
420                .filter(|record| record.kind == PressureKind::VllmPreemption)
421                .count() as u64;
422            let sglang_retractions_total = pressure_records
423                .iter()
424                .filter(|record| record.kind == PressureKind::SglangRetraction)
425                .count() as u64;
426            PressureEvidence {
427                records: pressure_records,
428                vllm_preemptions_total,
429                sglang_retractions_total,
430            }
431        });
432        let kv_ingest = kv_ingest.map(KvIngestAccumulator::finish);
433        OfflineRuntimeEvidence {
434            lifecycle_operations,
435            pressure,
436            kv_ingest,
437        }
438    }
439}
440
441#[derive(Debug)]
442struct KvIngestAccumulator {
443    hasher: Hasher,
444    evidence: KvIngestEvidence,
445}
446
447impl KvIngestAccumulator {
448    const ENCODING: &'static str = "dynamo.offline-kv-ingest.v1";
449
450    fn new() -> Self {
451        let mut hasher = Hasher::new();
452        put_bytes(&mut hasher, b"dynamo.offline-kv-ingest");
453        put_u32(&mut hasher, 1);
454        Self {
455            hasher,
456            evidence: KvIngestEvidence {
457                encoding: Self::ENCODING.to_string(),
458                ..KvIngestEvidence::default()
459            },
460        }
461    }
462
463    fn finish(mut self) -> KvIngestEvidence {
464        self.evidence.blake3_256 = self.hasher.finalize().to_hex().to_string();
465        self.evidence
466    }
467
468    fn record_batch(
469        &mut self,
470        pool: WorkerPool,
471        boundary: KvIngestBoundary,
472        at_ms: f64,
473        event_count: usize,
474        encode_events: impl FnOnce(&mut KvIngestEventEncoder<'_>) -> Result<()>,
475    ) -> Result<()> {
476        let batch_ordinal = self.evidence.batches;
477        self.evidence.batches = self
478            .evidence
479            .batches
480            .checked_add(1)
481            .expect("KV ingestion batch count overflow");
482        let event_count = to_u64(event_count, "KV event count")?;
483        put_u8(&mut self.hasher, pool.tag());
484        put_u8(&mut self.hasher, boundary.tag());
485        put_u64(&mut self.hasher, batch_ordinal);
486        put_f64(&mut self.hasher, at_ms);
487        put_u64(&mut self.hasher, event_count);
488
489        increment(&mut self.evidence.pool_counts, pool.as_str(), event_count);
490        let boundary_stats = self
491            .evidence
492            .boundaries
493            .entry(boundary.as_str().to_string())
494            .or_insert_with(|| KvIngestBoundaryStats {
495                first_at_ms: normalize_zero(at_ms),
496                ..KvIngestBoundaryStats::default()
497            });
498        boundary_stats.batches += 1;
499        boundary_stats.events += event_count;
500        boundary_stats.last_at_ms = normalize_zero(at_ms);
501
502        let mut encoder = KvIngestEventEncoder {
503            hasher: &mut self.hasher,
504            evidence: &mut self.evidence,
505        };
506        encode_events(&mut encoder)
507    }
508}
509
510fn increment(counts: &mut BTreeMap<String, u64>, key: &str, amount: u64) {
511    *counts.entry(key.to_string()).or_default() += amount;
512}
513
514fn normalize_zero(value: f64) -> f64 {
515    if value == 0.0 { 0.0 } else { value }
516}
517
518fn to_u64(value: usize, context: &str) -> Result<u64> {
519    u64::try_from(value).map_err(|_| anyhow::anyhow!("{context} exceeds u64"))
520}
521
522fn put_u8(hasher: &mut Hasher, value: u8) {
523    hasher.update(&[value]);
524}
525
526fn put_u32(hasher: &mut Hasher, value: u32) {
527    hasher.update(&value.to_be_bytes());
528}
529
530fn put_u64(hasher: &mut Hasher, value: u64) {
531    hasher.update(&value.to_be_bytes());
532}
533
534fn put_f64(hasher: &mut Hasher, value: f64) {
535    put_u64(hasher, normalize_zero(value).to_bits());
536}
537
538fn put_bytes(hasher: &mut Hasher, bytes: &[u8]) {
539    put_u64(
540        hasher,
541        u64::try_from(bytes.len()).expect("static KV digest domain length exceeds u64"),
542    );
543    hasher.update(bytes);
544}
545
546fn put_optional_u32(hasher: &mut Hasher, value: Option<u32>) {
547    match value {
548        Some(value) => {
549            put_u8(hasher, 1);
550            put_u32(hasher, value);
551        }
552        None => put_u8(hasher, 0),
553    }
554}
555
556fn put_optional_u64(hasher: &mut Hasher, value: Option<u64>) {
557    match value {
558        Some(value) => {
559            put_u8(hasher, 1);
560            put_u64(hasher, value);
561        }
562        None => put_u8(hasher, 0),
563    }
564}
565
566/// Neutral encoder used by observation adapters to contribute KV events to
567/// canonical replay evidence without exposing Router types to Replay.
568pub struct KvIngestEventEncoder<'a> {
569    hasher: &'a mut Hasher,
570    evidence: &'a mut KvIngestEvidence,
571}
572
573impl KvIngestEventEncoder<'_> {
574    pub fn begin_event(
575        &mut self,
576        worker_id: u64,
577        dp_rank: u32,
578        storage_tier_tag: u8,
579        storage_tier_name: &'static str,
580        event_id: u64,
581    ) {
582        self.evidence.events = self
583            .evidence
584            .events
585            .checked_add(1)
586            .expect("KV ingestion event count overflow");
587        put_u64(self.hasher, worker_id);
588        put_u32(self.hasher, dp_rank);
589        put_u8(self.hasher, storage_tier_tag);
590        put_u64(self.hasher, event_id);
591        increment(&mut self.evidence.tier_counts, storage_tier_name, 1);
592    }
593
594    pub fn begin_kind(&mut self, tag: u8, name: &'static str) {
595        put_u8(self.hasher, tag);
596        increment(&mut self.evidence.kind_counts, name, 1);
597    }
598
599    pub fn add_blocks(&mut self, count: usize, context: &str) -> Result<()> {
600        self.evidence.blocks = self
601            .evidence
602            .blocks
603            .checked_add(to_u64(count, context)?)
604            .expect("KV ingestion block count overflow");
605        Ok(())
606    }
607
608    pub fn put_len(&mut self, value: usize, context: &str) -> Result<()> {
609        put_u64(self.hasher, to_u64(value, context)?);
610        Ok(())
611    }
612
613    pub fn put_u8(&mut self, value: u8) {
614        put_u8(self.hasher, value);
615    }
616
617    pub fn put_u64(&mut self, value: u64) {
618        put_u64(self.hasher, value);
619    }
620
621    pub fn put_optional_u32(&mut self, value: Option<u32>) {
622        put_optional_u32(self.hasher, value);
623    }
624
625    pub fn put_optional_u64(&mut self, value: Option<u64>) {
626        put_optional_u64(self.hasher, value);
627    }
628}
629
630fn lower_pressure_state(state: SchedulerPressureState) -> EnginePressureState {
631    EnginePressureState {
632        running_requests: state.running_requests,
633        waiting_requests: state.waiting_requests,
634        active_blocks: state.active_blocks,
635    }
636}
637
638impl Default for ReplayEvidenceCollector {
639    fn default() -> Self {
640        Self::new(ReplayCaptureOptions::default())
641    }
642}
643
644pub(crate) fn common_origin(mut origins: impl Iterator<Item = u64>) -> Option<u64> {
645    let first = origins.next()?;
646    origins.all(|origin| origin == first).then_some(first)
647}
648
649#[cfg(test)]
650mod tests {
651    use crate::engine::{PressureEvent, PressureKind, PressureState};
652    use uuid::Uuid;
653
654    use crate::replay::{ReplayCaptureOptions, ReplayDeterminism, TraceCollector};
655
656    use super::{
657        KvIngestBoundary, ReplayEvidenceCollector, WorkerLifecycleTransition,
658        WorkerLifecycleTransitionKind, WorkerPool, WorkerPoolState,
659    };
660
661    #[test]
662    fn execution_local_pressure_ordinals_link_to_request_records() {
663        let uuid = Uuid::from_u128(42);
664        let mut trace = TraceCollector::default();
665        trace.set_capture_per_request(true);
666        trace.on_arrival(uuid, 0.0, 4, 2);
667        trace.on_admit(uuid, 0.0, 0);
668        let mut evidence = ReplayEvidenceCollector::new(ReplayCaptureOptions {
669            capture_per_request: true,
670            capture_canonical_evidence: true,
671            determinism: ReplayDeterminism::CanonicalV1,
672            ..Default::default()
673        });
674        evidence.record_native_pressure(
675            &mut trace,
676            WorkerPool::Agg,
677            7,
678            1,
679            PressureEvent {
680                at_ms: 2.0,
681                kind: PressureKind::VllmPreemption,
682                request_id: uuid,
683                state_before: PressureState {
684                    running_requests: 1,
685                    waiting_requests: Some(0),
686                    active_blocks: 4,
687                },
688                state_after: PressureState {
689                    running_requests: 0,
690                    waiting_requests: Some(1),
691                    active_blocks: 0,
692                },
693                request_active_blocks_before: 4,
694                logical_available_blocks_before: None,
695                required_blocks_before: None,
696            },
697        );
698        evidence.record_pressure_readmission(uuid, WorkerPool::Agg, 3.0);
699        trace.on_terminal(uuid, 4.0, crate::replay::ReplayTerminalStatus::Completed);
700        trace.set_runtime_evidence(evidence.finish());
701
702        let report = trace.finish();
703        assert_eq!(report.per_request[0].pressure_record_ordinals, vec![0]);
704        let pressure = report.runtime_evidence.pressure.unwrap();
705        assert_eq!(pressure.vllm_preemptions_total, 1);
706        assert_eq!(pressure.records[0].readmitted_at_ms, Some(3.0));
707    }
708
709    #[test]
710    fn lifecycle_origins_are_owned_by_one_collector() {
711        let mut evidence = ReplayEvidenceCollector::new(ReplayCaptureOptions {
712            capture_lifecycle_evidence: true,
713            ..Default::default()
714        });
715        evidence.record_lifecycle_operation(
716            1.0,
717            WorkerPool::Decode,
718            "planner_scale",
719            Some(0),
720            None,
721            vec![WorkerLifecycleTransition {
722                worker_id: 3,
723                transition: WorkerLifecycleTransitionKind::WorkerStarting,
724                prior_state: None,
725                state: "starting",
726                reason: None,
727                origin_operation_ordinal: None,
728            }],
729            WorkerPoolState {
730                starting: vec![3],
731                ..Default::default()
732            },
733            Vec::new(),
734        );
735        let origin = evidence.startup_origin(WorkerPool::Decode, 3);
736        evidence.record_lifecycle_operation(
737            2.0,
738            WorkerPool::Decode,
739            "worker_ready_event",
740            None,
741            origin,
742            vec![WorkerLifecycleTransition {
743                worker_id: 3,
744                transition: WorkerLifecycleTransitionKind::WorkerReady,
745                prior_state: Some("starting"),
746                state: "active",
747                reason: None,
748                origin_operation_ordinal: origin,
749            }],
750            WorkerPoolState {
751                active: vec![3],
752                ..Default::default()
753            },
754            Vec::new(),
755        );
756
757        let evidence = evidence.finish();
758        assert_eq!(evidence.lifecycle_operations.len(), 2);
759        assert_eq!(
760            evidence.lifecycle_operations[1].origin_operation_ordinal,
761            Some(0)
762        );
763    }
764
765    #[test]
766    fn canonical_kv_ingest_is_execution_local_and_byte_stable() {
767        let capture = || {
768            let mut evidence = ReplayEvidenceCollector::new(ReplayCaptureOptions {
769                capture_canonical_evidence: true,
770                ..Default::default()
771            });
772            evidence
773                .record_kv_ingest(
774                    WorkerPool::Agg,
775                    KvIngestBoundary::PassEnd,
776                    7.0,
777                    1,
778                    |encoder| {
779                        encoder.begin_event(3, 0, 0, "device", 9);
780                        encoder.begin_kind(0, "stored");
781                        encoder.add_blocks(2, "test blocks")?;
782                        encoder.put_len(2, "test blocks")?;
783                        encoder.put_u64(11);
784                        encoder.put_u64(12);
785                        Ok(())
786                    },
787                )
788                .unwrap();
789            evidence.finish().kv_ingest.unwrap()
790        };
791
792        let first = capture();
793        let second = capture();
794        assert_eq!(first, second);
795        assert_eq!(first.batches, 1);
796        assert_eq!(first.events, 1);
797        assert_eq!(first.blocks, 2);
798        assert_eq!(first.boundaries["pass_end"].last_at_ms, 7.0);
799        assert_eq!(first.blake3_256.len(), 64);
800    }
801}