1use 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#[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#[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#[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 #[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
566pub 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}