Skip to main content

aisimulate_core/replay/
replayer.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Public replay facade over the mechanically moved topology runtimes.
5
6use std::collections::VecDeque;
7use std::time::Instant;
8
9use anyhow::Result as AnyResult;
10use uuid::Uuid;
11
12use crate::replay::OfflineDisaggReplayConfig;
13use crate::replay::agg::AggRuntimeImpl;
14use crate::replay::artifact::{
15    ReplayArtifactKvEventVisibility, ReplayArtifactSink, ReplayArtifacts,
16};
17use crate::replay::components::{
18    AdmissionQueue, NoReplayMetadata, ReplayAdmissionMetadata, ReplayEngineObservation, ReplayMode,
19};
20use crate::replay::core::round_robin::{AggregatedRoundRobinPlacement, PoolRoundRobinPlacement};
21use crate::replay::core::{NoEngineEvents, PlacementPolicy, WorkerTopology};
22use crate::replay::disagg::DisaggRuntimeImpl;
23use crate::replay::engine::{ReplayEngineConfig, ReplayEngineFactory};
24use crate::replay::error::{
25    placement_boundary, runtime_error, scaling_boundary, telemetry_boundary,
26};
27use crate::replay::loadgen::ReplayRequestPayload;
28use crate::replay::loadgen::WorkloadDriver;
29use crate::replay::protocol::{DirectRequest, ReplayPromptTokenSource, ReplayRequestContext};
30use crate::replay::scaling::ReplayScalingPolicy;
31use crate::replay::telemetry::{ReplayTelemetryObserver, ReplayTelemetrySnapshot};
32use crate::replay::{
33    ReplayCaptureOptions, ReplayDeterminism, ReplayError, ReplayReport, ReplayResult, ReplaySpec,
34    ReplayTopology, SlaThresholds, WorkerStage,
35};
36
37/// Runtime composition supplied by the built-in engine stack or a Dynamo
38/// adapter. The adapter owns concrete Router/Planner construction; Replay only
39/// sees the already-neutral placement and scaling contracts.
40pub trait ReplayComposition {
41    type Metadata: ReplayAdmissionMetadata;
42    type Observation: ReplayEngineObservation;
43    type AggregatedPlacement: PlacementPolicy<
44            ReplayRequestPayload,
45            Metadata = Self::Metadata,
46            Observation = <Self::Observation as ReplayEngineObservation>::Batch,
47        >;
48    type DisaggregatedPlacement: PlacementPolicy<
49            ReplayRequestPayload,
50            Metadata = Self::Metadata,
51            Observation = <Self::Observation as ReplayEngineObservation>::Batch,
52        >;
53
54    fn validate_spec(&self, _spec: &ReplaySpec) -> ReplayResult<()> {
55        Ok(())
56    }
57
58    fn create_aggregated_placement(
59        &mut self,
60        dp_size: u32,
61        topology: Vec<WorkerTopology>,
62    ) -> AnyResult<Self::AggregatedPlacement>;
63
64    fn create_disaggregated_placements(
65        &mut self,
66        prefill_dp_size: u32,
67        prefill_topology: Vec<WorkerTopology>,
68        decode_dp_size: u32,
69        decode_topology: Vec<WorkerTopology>,
70    ) -> AnyResult<(Self::DisaggregatedPlacement, Self::DisaggregatedPlacement)>;
71
72    /// Return the run-owned scaling policy, if this composition has one.
73    fn take_scaling_policy(&mut self) -> AnyResult<Option<Box<dyn ReplayScalingPolicy>>> {
74        Ok(None)
75    }
76
77    /// Inform policy construction about explicitly requested deterministic
78    /// selection. Implementations should use
79    /// [`ReplayDeterminism::selector_seed`] and leave normal runs unseeded.
80    fn set_determinism(&mut self, _determinism: ReplayDeterminism) -> ReplayResult<()> {
81        Ok(())
82    }
83}
84
85/// Classifies every fallible callback from a placement policy at the boundary
86/// where the policy enters the otherwise policy-neutral runtime.
87struct PlacementPolicyBoundary<P>(P);
88
89impl<Request, P> PlacementPolicy<Request> for PlacementPolicyBoundary<P>
90where
91    P: PlacementPolicy<Request>,
92{
93    type Metadata = P::Metadata;
94    type Observation = P::Observation;
95
96    fn place(
97        &mut self,
98        request: &Request,
99        metadata: Self::Metadata,
100        session_id: Option<String>,
101        now_ms: f64,
102    ) -> AnyResult<crate::replay::core::PlacementEffects> {
103        self.0
104            .place(request, metadata, session_id, now_ms)
105            .map_err(placement_boundary)
106    }
107
108    fn observe(
109        &mut self,
110        observation: Self::Observation,
111        now_ms: f64,
112    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
113        self.0
114            .observe(observation, now_ms)
115            .map_err(placement_boundary)
116    }
117
118    fn cancel_pending(&mut self, request_id: Uuid) -> bool {
119        self.0.cancel_pending(request_id)
120    }
121
122    fn request_terminal(
123        &mut self,
124        request_id: Uuid,
125        now_ms: f64,
126    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
127        self.0
128            .request_terminal(request_id, now_ms)
129            .map_err(placement_boundary)
130    }
131
132    fn prefill_completed(
133        &mut self,
134        request_id: Uuid,
135        now_ms: f64,
136    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
137        self.0
138            .prefill_completed(request_id, now_ms)
139            .map_err(placement_boundary)
140    }
141
142    fn pending_count(&self) -> usize {
143        self.0.pending_count()
144    }
145
146    fn worker_ready(
147        &mut self,
148        worker: WorkerTopology,
149        now_ms: f64,
150    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
151        self.0
152            .worker_ready(worker, now_ms)
153            .map_err(placement_boundary)
154    }
155
156    fn worker_draining(
157        &mut self,
158        worker: WorkerTopology,
159        now_ms: f64,
160    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
161        self.0
162            .worker_draining(worker, now_ms)
163            .map_err(placement_boundary)
164    }
165
166    fn worker_removed(
167        &mut self,
168        worker: WorkerTopology,
169        now_ms: f64,
170    ) -> AnyResult<Vec<crate::replay::core::Placement>> {
171        self.0
172            .worker_removed(worker, now_ms)
173            .map_err(placement_boundary)
174    }
175
176    fn topology_settled(&mut self, now_ms: f64) -> AnyResult<Vec<crate::replay::core::Placement>> {
177        self.0.topology_settled(now_ms).map_err(placement_boundary)
178    }
179}
180
181/// Classifies Planner/scaling callbacks without exposing policy-specific types
182/// to the aggregated or disaggregated runtime.
183struct ScalingPolicyBoundary(Box<dyn ReplayScalingPolicy>);
184
185impl ReplayScalingPolicy for ScalingPolicyBoundary {
186    fn capture_lifecycle_evidence(&self) -> bool {
187        self.0.capture_lifecycle_evidence()
188    }
189
190    fn initial_tick_ms(&mut self) -> AnyResult<f64> {
191        self.0.initial_tick_ms().map_err(scaling_boundary)
192    }
193
194    fn on_tick(
195        &mut self,
196        snapshot: crate::replay::scaling::ReplayScalingSnapshot,
197    ) -> AnyResult<crate::replay::scaling::ReplayScalingDecision> {
198        self.0.on_tick(snapshot).map_err(scaling_boundary)
199    }
200}
201
202/// Classifies observer failures without exposing adapter-specific types to the
203/// topology runtimes.
204struct TelemetryObserverBoundary(Box<dyn ReplayTelemetryObserver>);
205
206impl ReplayTelemetryObserver for TelemetryObserverBoundary {
207    fn on_sample(&mut self, snapshot: ReplayTelemetrySnapshot) -> AnyResult<()> {
208        self.0.on_sample(snapshot).map_err(telemetry_boundary)
209    }
210}
211
212/// Replay-owned runtime input used by compatibility runners that already
213/// lowered a trace into the shared workload driver.
214///
215/// Serializable callers should keep using [`ReplaySpec::requests`]. Dynamo's
216/// legacy entrypoints use this seam to preserve multi-turn, concurrency, and
217/// agentic scheduling without recompiling Replay sources in the Dynamo crate.
218#[doc(hidden)]
219#[allow(clippy::large_enum_variant)] // Preserve the inline workload through runtime construction.
220pub enum ReplayRuntimeInput {
221    Requests(VecDeque<DirectRequest>),
222    Workload(WorkloadDriver),
223}
224
225/// Built-in engine-only composition: Round-robin placement and fixed capacity.
226#[derive(Debug, Default, Clone, Copy)]
227pub struct RoundRobinComposition;
228
229impl ReplayComposition for RoundRobinComposition {
230    type Metadata = NoReplayMetadata;
231    type Observation = NoEngineEvents;
232    type AggregatedPlacement = AggregatedRoundRobinPlacement<()>;
233    type DisaggregatedPlacement = PoolRoundRobinPlacement<()>;
234
235    fn validate_spec(&self, spec: &ReplaySpec) -> ReplayResult<()> {
236        if spec.adapters.placement.provider != "round_robin" {
237            return Err(ReplayError::InvalidSpec(format!(
238                "engine composition requires round_robin placement, got {:?}",
239                spec.adapters.placement.provider
240            )));
241        }
242        if spec.adapters.scaling.provider != "none" {
243            return Err(ReplayError::InvalidSpec(format!(
244                "engine composition does not provide scaling, got {:?}",
245                spec.adapters.scaling.provider
246            )));
247        }
248        Ok(())
249    }
250
251    fn create_aggregated_placement(
252        &mut self,
253        dp_size: u32,
254        topology: Vec<WorkerTopology>,
255    ) -> AnyResult<Self::AggregatedPlacement> {
256        Ok(AggregatedRoundRobinPlacement::new(dp_size, topology))
257    }
258
259    fn create_disaggregated_placements(
260        &mut self,
261        _prefill_dp_size: u32,
262        prefill_topology: Vec<WorkerTopology>,
263        _decode_dp_size: u32,
264        decode_topology: Vec<WorkerTopology>,
265    ) -> AnyResult<(Self::DisaggregatedPlacement, Self::DisaggregatedPlacement)> {
266        Ok((
267            PoolRoundRobinPlacement::new(prefill_topology),
268            PoolRoundRobinPlacement::new(decode_topology),
269        ))
270    }
271}
272
273/// Owns one replay execution: canonical spec, engine construction, and
274/// the selected placement/scaling composition.
275pub struct Replayer<C = RoundRobinComposition> {
276    spec: ReplaySpec,
277    factory: ReplayEngineFactory,
278    composition: C,
279    runtime_input: Option<ReplayRuntimeInput>,
280    capture: ReplayCaptureOptions,
281    telemetry: Option<(f64, Box<dyn ReplayTelemetryObserver>)>,
282}
283
284impl Replayer<RoundRobinComposition> {
285    pub fn new(spec: ReplaySpec, factory: ReplayEngineFactory) -> ReplayResult<Self> {
286        Self::with_composition(spec, factory, RoundRobinComposition)
287    }
288
289    /// Run a fixed, aggregated single-worker replay and retain detailed
290    /// request/output/native-KV observations from the same Replayer-owned
291    /// aggregated runtime that produces the normal report.
292    ///
293    /// This contract intentionally targets one worker artifact. Multi-worker,
294    /// scaling, and disaggregated runs should consume the normal report and
295    /// placement/scaling observation contracts instead of creating a second
296    /// scheduler loop solely for artifact generation.
297    pub fn run_with_artifacts(
298        self,
299        visibility: ReplayArtifactKvEventVisibility,
300    ) -> ReplayResult<(ReplayReport, ReplayArtifacts)> {
301        let sink = ReplayArtifactSink::new(visibility);
302        let report = self.run_inner(Some(sink.clone()))?;
303        Ok((report, sink.take()?))
304    }
305}
306
307impl<C: ReplayComposition> Replayer<C> {
308    pub fn with_composition(
309        spec: ReplaySpec,
310        factory: ReplayEngineFactory,
311        composition: C,
312    ) -> ReplayResult<Self> {
313        spec.validate()?;
314        composition.validate_spec(&spec)?;
315        Ok(Self {
316            spec,
317            factory,
318            composition,
319            runtime_input: None,
320            capture: ReplayCaptureOptions::default(),
321            telemetry: None,
322        })
323    }
324
325    /// Override the serializable request list with an already-lowered,
326    /// Replay-owned runtime input.
327    #[doc(hidden)]
328    pub fn with_runtime_input(mut self, input: ReplayRuntimeInput) -> Self {
329        self.runtime_input = Some(input);
330        self
331    }
332
333    /// Configure detailed capture and canonical determinism for this run.
334    pub fn with_capture_options(mut self, options: ReplayCaptureOptions) -> Self {
335        self.capture = options;
336        self
337    }
338
339    /// Attach a policy-neutral observer sampled at a fixed virtual-time
340    /// interval. Telemetry remains disabled unless this method is called.
341    pub fn with_telemetry_observer(
342        mut self,
343        sample_interval_ms: f64,
344        observer: Box<dyn ReplayTelemetryObserver>,
345    ) -> ReplayResult<Self> {
346        if !sample_interval_ms.is_finite() || sample_interval_ms <= 0.0 {
347            return Err(ReplayError::InvalidSpec(format!(
348                "telemetry sample interval must be finite and positive, got {sample_interval_ms}"
349            )));
350        }
351        self.telemetry = Some((sample_interval_ms, observer));
352        Ok(self)
353    }
354
355    pub fn run(self) -> ReplayResult<ReplayReport> {
356        self.run_inner(None)
357    }
358
359    fn run_inner(
360        mut self,
361        artifact_sink: Option<ReplayArtifactSink>,
362    ) -> ReplayResult<ReplayReport> {
363        let wall_start = Instant::now();
364        self.composition.set_determinism(self.capture.determinism)?;
365        let engine_config = ReplayEngineConfig::parse(&self.spec.engine)?;
366        engine_config.validate_topology(&self.spec.topology)?;
367        let runtime_input = match self.runtime_input.take() {
368            Some(mut input) => {
369                apply_runtime_determinism(&mut input, self.capture.determinism);
370                input
371            }
372            None => {
373                ReplayRuntimeInput::Requests(lower_requests(&self.spec, self.capture.determinism)?)
374            }
375        };
376        let mode = self
377            .spec
378            .max_in_flight
379            .map_or(ReplayMode::Trace, |max_in_flight| ReplayMode::Concurrency {
380                max_in_flight,
381            });
382        let scaling = self
383            .composition
384            .take_scaling_policy()
385            .map_err(|error| ReplayError::Scaling(format!("{error:#}")))?;
386        let telemetry = self.telemetry.take();
387
388        let collector = match &self.spec.topology {
389            ReplayTopology::Aggregated { workers } => {
390                let role_factory = self.factory.role_factory(
391                    &engine_config,
392                    WorkerStage::Aggregated,
393                    C::Observation::capture_engine_kv_events(WorkerStage::Aggregated)
394                        || artifact_sink.is_some(),
395                )?;
396                let startup_time_ms = positive_delay(workers.startup_delay_ms);
397                if artifact_sink.is_some()
398                    && (workers.initial_workers != 1
399                        || role_factory.dp_size() != 1
400                        || scaling.is_some())
401                {
402                    return Err(ReplayError::InvalidSpec(
403                        "detailed replay artifacts require fixed aggregated topology with one logical DP1 worker"
404                            .to_string(),
405                    ));
406                }
407
408                let mut runtime = AggRuntimeImpl::<
409                    PlacementPolicyBoundary<C::AggregatedPlacement>,
410                    C::Observation,
411                    C::Metadata,
412                >::new_composed(
413                    role_factory,
414                    admission_queue(runtime_input, mode),
415                    workers.initial_workers,
416                    startup_time_ms,
417                    |dp_size, topology| {
418                        self.composition
419                            .create_aggregated_placement(dp_size, topology)
420                            .map(PlacementPolicyBoundary)
421                            .map_err(placement_boundary)
422                    },
423                )
424                .map_err(runtime_error)?
425                .with_capture_options(self.capture)
426                .with_per_request_records(
427                    self.spec.record_per_request || self.capture.effective_per_request(),
428                )
429                .with_max_sim_time_ms(self.spec.max_sim_time_ms);
430                if let Some(sink) = artifact_sink {
431                    runtime = runtime.with_artifact_sink(sink);
432                }
433                if let Some(policy) = scaling {
434                    runtime = runtime.with_scaling_policy(Box::new(ScalingPolicyBoundary(policy)));
435                }
436                if let Some((sample_interval_ms, observer)) = telemetry {
437                    runtime = runtime.with_telemetry_observer(
438                        sample_interval_ms,
439                        Box::new(TelemetryObserverBoundary(observer)),
440                    );
441                }
442                runtime.run().map_err(runtime_error)?.0
443            }
444            ReplayTopology::Disaggregated {
445                prefill,
446                decode,
447                handoff_latency_ms,
448            } => {
449                if artifact_sink.is_some() {
450                    return Err(ReplayError::InvalidSpec(
451                        "detailed replay artifacts require aggregated topology".to_string(),
452                    ));
453                }
454                let prefill_factory = self.factory.role_factory(
455                    &engine_config,
456                    WorkerStage::Prefill,
457                    C::Observation::capture_engine_kv_events(WorkerStage::Prefill),
458                )?;
459                let decode_factory = self.factory.role_factory(
460                    &engine_config,
461                    WorkerStage::Decode,
462                    C::Observation::capture_engine_kv_events(WorkerStage::Decode),
463                )?;
464                let config = OfflineDisaggReplayConfig {
465                    prefill_factory,
466                    decode_factory,
467                    prefill_startup_time_ms: positive_delay(prefill.startup_delay_ms),
468                    decode_startup_time_ms: positive_delay(decode.startup_delay_ms),
469                    num_prefill_workers: prefill.initial_workers,
470                    num_decode_workers: decode.initial_workers,
471                    handoff_latency_ms: *handoff_latency_ms,
472                };
473                let mut runtime = DisaggRuntimeImpl::<
474                    PlacementPolicyBoundary<C::DisaggregatedPlacement>,
475                    C::Observation,
476                    C::Metadata,
477                >::new_composed(
478                    &config,
479                    admission_queue(runtime_input, mode),
480                    false,
481                    |prefill_dp, prefill_topology, decode_dp, decode_topology| {
482                        self.composition
483                            .create_disaggregated_placements(
484                                prefill_dp,
485                                prefill_topology,
486                                decode_dp,
487                                decode_topology,
488                            )
489                            .map(|(prefill, decode)| {
490                                (
491                                    PlacementPolicyBoundary(prefill),
492                                    PlacementPolicyBoundary(decode),
493                                )
494                            })
495                            .map_err(placement_boundary)
496                    },
497                )
498                .map_err(runtime_error)?
499                .with_capture_options(self.capture)
500                .with_per_request_records(
501                    self.spec.record_per_request || self.capture.effective_per_request(),
502                )
503                .with_max_sim_time_ms(self.spec.max_sim_time_ms);
504                if let Some(policy) = scaling {
505                    runtime = runtime.with_scaling_policy(Box::new(ScalingPolicyBoundary(policy)));
506                }
507                if let Some((sample_interval_ms, observer)) = telemetry {
508                    runtime = runtime.with_telemetry_observer(
509                        sample_interval_ms,
510                        Box::new(TelemetryObserverBoundary(observer)),
511                    );
512                }
513                runtime.run().map_err(runtime_error)?.0
514            }
515        };
516
517        Ok(finish_report(collector, self.spec.sla)
518            .with_wall_time_ms(wall_start.elapsed().as_secs_f64() * 1_000.0))
519    }
520}
521
522fn admission_queue<Metadata: ReplayAdmissionMetadata>(
523    input: ReplayRuntimeInput,
524    mode: ReplayMode,
525) -> AdmissionQueue<Metadata> {
526    match input {
527        ReplayRuntimeInput::Requests(requests) => AdmissionQueue::new_requests(requests, mode),
528        ReplayRuntimeInput::Workload(driver) => AdmissionQueue::new_workload(driver, mode),
529    }
530}
531
532fn positive_delay(delay_ms: f64) -> Option<f64> {
533    (delay_ms > 0.0).then_some(delay_ms)
534}
535
536fn lower_requests(
537    spec: &ReplaySpec,
538    determinism: ReplayDeterminism,
539) -> ReplayResult<VecDeque<DirectRequest>> {
540    let mut pending = spec
541        .requests
542        .iter()
543        .enumerate()
544        .map(|(index, request)| -> ReplayResult<_> {
545            let request_id = match determinism {
546                ReplayDeterminism::Random => Uuid::new_v4(),
547                ReplayDeterminism::CanonicalV1 => Uuid::from_u128(
548                    u128::try_from(index)
549                        .expect("usize always fits u128")
550                        .checked_add(1)
551                        .expect("replay request index overflow"),
552                ),
553            };
554            let (tokens, prompt_token_source) = match &request.input_token_ids {
555                Some(tokens) => (tokens.clone(), ReplayPromptTokenSource::Materialized),
556                None => {
557                    let seed = u32::try_from(index)
558                        .unwrap_or(u32::MAX)
559                        .wrapping_mul(1_000_003);
560                    (
561                        (0..request.input_tokens)
562                            .map(|offset| {
563                                seed.wrapping_add(u32::try_from(offset).unwrap_or(u32::MAX))
564                            })
565                            .collect(),
566                        ReplayPromptTokenSource::LengthOnlySynthetic,
567                    )
568                }
569            };
570            let routing = request.routing_metadata()?;
571            Ok(DirectRequest {
572                tokens,
573                max_output_tokens: request.output_tokens,
574                output_token_ids: request.output_token_ids.clone(),
575                uuid: Some(request_id),
576                dp_rank: 0,
577                preferred_dp_rank: request.dp_rank,
578                arrival_timestamp_ms: Some(request.arrival_time_ms),
579                priority: routing.priority,
580                strict_priority: routing.strict_priority,
581                policy_class: routing.policy_class,
582                replay_context: Some(ReplayRequestContext {
583                    authored_id: request.id.clone(),
584                    session_id: request.session_id.clone(),
585                    turn_index: request.turn_index,
586                    metadata: request.metadata.clone(),
587                    prompt_token_source,
588                }),
589            })
590        })
591        .collect::<ReplayResult<Vec<_>>>()?;
592    pending.sort_by(|left, right| {
593        left.arrival_timestamp_ms
594            .expect("ReplaySpec request always has an arrival")
595            .total_cmp(
596                &right
597                    .arrival_timestamp_ms
598                    .expect("ReplaySpec request always has an arrival"),
599            )
600    });
601    Ok(pending.into())
602}
603
604fn apply_runtime_determinism(input: &mut ReplayRuntimeInput, determinism: ReplayDeterminism) {
605    if determinism != ReplayDeterminism::CanonicalV1 {
606        return;
607    }
608    match input {
609        ReplayRuntimeInput::Requests(requests) => {
610            for (index, request) in requests.iter_mut().enumerate() {
611                request.uuid = Some(Uuid::from_u128(
612                    u128::try_from(index)
613                        .expect("usize always fits u128")
614                        .checked_add(1)
615                        .expect("replay request index overflow"),
616                ));
617            }
618        }
619        ReplayRuntimeInput::Workload(driver) => {
620            driver.set_deterministic_request_ids(1);
621        }
622    }
623}
624
625fn finish_report(mut collector: crate::replay::TraceCollector, sla: SlaThresholds) -> ReplayReport {
626    collector.set_sla_thresholds(sla);
627    collector.finish()
628}
629
630#[cfg(test)]
631mod tests {
632    use super::*;
633    use crate::replay::{
634        ProviderSpec, ReplayAdapters, ReplayRequest, ReplayTopology, WorkerPoolSpec,
635    };
636
637    #[test]
638    fn replay_spec_lowering_preserves_correlation_routing_and_prompt_provenance() {
639        let spec = ReplaySpec {
640            version: 1,
641            topology: ReplayTopology::Aggregated {
642                workers: WorkerPoolSpec::default(),
643            },
644            engine: serde_json::Value::Null,
645            adapters: ReplayAdapters {
646                placement: ProviderSpec::round_robin(),
647                scaling: ProviderSpec::no_scaling(),
648            },
649            max_sim_time_ms: None,
650            max_in_flight: None,
651            record_per_request: true,
652            sla: Default::default(),
653            requests: vec![
654                ReplayRequest {
655                    id: "length-only".into(),
656                    arrival_time_ms: 0.0,
657                    input_tokens: 3,
658                    input_token_ids: None,
659                    output_tokens: 2,
660                    output_token_ids: None,
661                    dp_rank: Some(2),
662                    session_id: Some("session-a".into()),
663                    turn_index: Some(4),
664                    metadata: serde_json::json!({
665                        "priority": -7,
666                        "strict_priority": 9,
667                        "policy_class": "latency",
668                        "caller_tag": "preserved"
669                    }),
670                },
671                ReplayRequest {
672                    id: "materialized".into(),
673                    arrival_time_ms: 1.0,
674                    input_tokens: 2,
675                    input_token_ids: Some(vec![41, 42]),
676                    output_tokens: 1,
677                    output_token_ids: None,
678                    dp_rank: None,
679                    session_id: None,
680                    turn_index: None,
681                    metadata: serde_json::Value::Null,
682                },
683            ],
684        };
685
686        let lowered = lower_requests(&spec, ReplayDeterminism::CanonicalV1)
687            .unwrap()
688            .into_iter()
689            .collect::<Vec<_>>();
690        let first = &lowered[0];
691        assert_eq!(first.uuid, Some(Uuid::from_u128(1)));
692        assert_eq!(first.priority, -7);
693        assert_eq!(first.strict_priority, 9);
694        assert_eq!(first.policy_class.as_deref(), Some("latency"));
695        assert_eq!(first.preferred_dp_rank, Some(2));
696        assert!(!first.prompt_tokens_are_placement_safe());
697        let context = first.replay_context.as_ref().unwrap();
698        assert_eq!(context.authored_id, "length-only");
699        assert_eq!(context.session_id.as_deref(), Some("session-a"));
700        assert_eq!(context.turn_index, Some(4));
701        assert_eq!(context.metadata["caller_tag"], "preserved");
702
703        assert_eq!(lowered[1].tokens, vec![41, 42]);
704        assert!(lowered[1].prompt_tokens_are_placement_safe());
705    }
706
707    #[test]
708    fn random_lowering_does_not_use_ordinal_request_uuids() {
709        let spec = ReplaySpec {
710            version: 1,
711            topology: ReplayTopology::aggregated(1),
712            engine: serde_json::Value::Null,
713            adapters: ReplayAdapters::default(),
714            max_sim_time_ms: None,
715            max_in_flight: None,
716            record_per_request: false,
717            sla: Default::default(),
718            requests: vec![ReplayRequest {
719                id: "random".into(),
720                arrival_time_ms: 0.0,
721                input_tokens: 1,
722                input_token_ids: Some(vec![1]),
723                output_tokens: 1,
724                output_token_ids: None,
725                dp_rank: None,
726                session_id: None,
727                turn_index: None,
728                metadata: serde_json::Value::Null,
729            }],
730        };
731
732        let first = lower_requests(&spec, ReplayDeterminism::Random)
733            .unwrap()
734            .pop_front()
735            .unwrap();
736        assert_ne!(first.uuid, Some(Uuid::from_u128(1)));
737    }
738}