Skip to main content

canwu_sim/runtime/
scheduling.rs

1use super::event_payloads::canonicalize_event_kind;
2use super::event_payloads::{
3    ArmyArrived, KnowledgeUpdated, LetterDelivered, PersonArrived, ReportDispatched,
4    RuntimeEventPayload,
5};
6use super::settlement::CommittedBoundaryRandomDraw;
7use super::{
8    ActorKnowledge, ArmyId, ArmyKnowledge, AssertUnwindSafe, BTreeMap, BoundaryId, CanwuError,
9    CauseRef, ClockTransactionCheckpoint, CommitmentDomains, DeterministicRng, EntityRef,
10    ErrorCode, EstimateRange, EventId, EventKind, KnowledgeSource, LetterId, LetterStatus,
11    PendingBoundaryRandomDraw, PersonId, RandomAlgorithm, RandomDrawAddress, RandomDrawId,
12    RandomDrawOutcome, RandomDrawProducer, RandomDrawRecord, RandomStreamKey,
13    RuntimeValidationContext, ScheduledBatchTransactionCheckpoint, SimDuration, SimEvent, SimTime,
14    Simulation, SimulationView, SimulationViewState, StateKey, SystemDirective, TerritoryId,
15    catch_unwind, claim_counter, random, validate_directives_with_context,
16};
17use serde::{Deserialize, Serialize};
18
19#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
20pub(super) struct ScheduleKey {
21    pub(super) at: SimTime,
22    pub(super) sequence: u64,
23}
24
25#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
26#[serde(tag = "type", rename_all = "snake_case")]
27pub(super) enum ScheduledAction {
28    ArmyArrival {
29        army: ArmyId,
30        destination: TerritoryId,
31        order_event: EventId,
32        correlation_id: u64,
33    },
34    PersonArrival {
35        person: PersonId,
36        destination: TerritoryId,
37        order_event: EventId,
38        cargo: Vec<LetterId>,
39        correlation_id: u64,
40    },
41    KnowledgeReport {
42        recipient: PersonId,
43        army: ArmyId,
44        location: TerritoryId,
45        observed_at: SimTime,
46        dispatch_event: EventId,
47        correlation_id: u64,
48    },
49    PluginDirective {
50        plugin: String,
51        directive: Box<SystemDirective>,
52        allowed_writes: Vec<StateKey>,
53        cause: CauseRef,
54        correlation_id: u64,
55    },
56}
57
58#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
59pub(super) struct ScheduledRecord {
60    pub(super) key: ScheduleKey,
61    pub(super) action: ScheduledAction,
62}
63
64impl Simulation {
65    pub fn advance(&mut self, duration: SimDuration) -> Result<Vec<SimEvent>, CanwuError> {
66        self.ensure_runtime_ready()?;
67        if duration.is_negative() {
68            return Err(CanwuError::new(
69                ErrorCode::InvalidDuration,
70                "simulation time cannot advance by a negative duration",
71            ));
72        }
73        let target = self
74            .state
75            .scheduler
76            .now
77            .checked_add(duration)
78            .ok_or_else(|| {
79                CanwuError::new(
80                    ErrorCode::InvalidDuration,
81                    "simulation target time exceeds the supported range",
82                )
83            })?;
84        self.ensure_legacy_advance_does_not_cross_ingress(target)?;
85        self.advance_to(target)
86    }
87
88    pub fn step(&mut self) -> Result<Vec<SimEvent>, CanwuError> {
89        self.ensure_runtime_ready()?;
90        if self.state.scheduler.pending_ingress.first().is_some() {
91            return Err(CanwuError::new(
92                ErrorCode::InvalidBoundary,
93                "pending canonical ingress requires step_canonical",
94            ));
95        }
96        let Some(next_time) = self.state.scheduler.actions.keys().next().map(|key| key.at) else {
97            return Ok(Vec::new());
98        };
99        self.advance_to(next_time)
100    }
101
102    pub fn advance_until<F>(
103        &mut self,
104        maximum: SimDuration,
105        mut condition: F,
106    ) -> Result<Vec<SimEvent>, CanwuError>
107    where
108        F: FnMut(&Self) -> bool,
109    {
110        self.ensure_runtime_ready()?;
111        if maximum.is_negative() {
112            return Err(CanwuError::new(
113                ErrorCode::InvalidDuration,
114                "advance_until maximum cannot be negative",
115            ));
116        }
117        let target = self
118            .state
119            .scheduler
120            .now
121            .checked_add(maximum)
122            .ok_or_else(|| {
123                CanwuError::new(
124                    ErrorCode::InvalidDuration,
125                    "advance_until target time exceeds the supported range",
126                )
127            })?;
128        self.ensure_legacy_advance_does_not_cross_ingress(target)?;
129        let start = self.state.evidence.events.len();
130        while self.state.scheduler.now < target && !condition(self) {
131            let next_time = self
132                .state
133                .scheduler
134                .actions
135                .keys()
136                .next()
137                .map_or(target, |key| key.at.min(target));
138            self.advance_to(next_time)?;
139            if next_time == target {
140                break;
141            }
142        }
143        Ok(self.state.evidence.events[start..].to_vec())
144    }
145
146    pub(super) fn ensure_legacy_advance_does_not_cross_ingress(
147        &self,
148        target: SimTime,
149    ) -> Result<(), CanwuError> {
150        if self
151            .state
152            .scheduler
153            .pending_ingress
154            .first()
155            .is_some_and(|key| key.due_at <= target)
156        {
157            return Err(CanwuError::new(
158                ErrorCode::InvalidBoundary,
159                "legacy time advancement cannot cross pending canonical ingress; use advance_canonical",
160            ));
161        }
162        Ok(())
163    }
164
165    pub(super) fn advance_to(&mut self, target: SimTime) -> Result<Vec<SimEvent>, CanwuError> {
166        let start = self.state.evidence.events.len();
167        while let Some(boundary_time) = self.state.scheduler.actions.keys().next().map(|key| key.at)
168            && boundary_time <= target
169        {
170            let transaction = ScheduledBatchTransactionCheckpoint::capture(&self.state);
171            let result = (|| {
172                self.invalidate_commitments(CommitmentDomains::SCHEDULER);
173                self.state.scheduler.now = boundary_time;
174                while self
175                    .state
176                    .scheduler
177                    .actions
178                    .first_key_value()
179                    .is_some_and(|(key, _)| key.at == boundary_time)
180                {
181                    let (_, action) = self
182                        .state
183                        .scheduler
184                        .actions
185                        .pop_first()
186                        .expect("scheduler was checked as non-empty");
187                    self.execute_scheduled(action)?;
188                }
189                self.state.metadata.plugin_registration_closed = true;
190                self.refresh_checkpoint_hash()
191            })();
192            if let Err(error) = result {
193                transaction.restore(&mut self.state);
194                return Err(error);
195            }
196        }
197        let transaction = ClockTransactionCheckpoint::capture(&self.state);
198        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
199        self.state.scheduler.now = target;
200        self.state.metadata.plugin_registration_closed = true;
201        if let Err(error) = self.refresh_checkpoint_hash() {
202            transaction.restore(&mut self.state);
203            return Err(error);
204        }
205        Ok(self.state.evidence.events[start..].to_vec())
206    }
207
208    pub(super) fn advance_to_before_boundary(&mut self, target: SimTime) -> Result<(), CanwuError> {
209        while let Some(next) = self.state.scheduler.actions.keys().next().map(|key| key.at)
210            && next < target
211        {
212            self.advance_to(next)?;
213        }
214        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
215        self.state.scheduler.now = target;
216        self.state.metadata.plugin_registration_closed = true;
217        self.refresh_checkpoint_hash()
218    }
219
220    pub(super) fn execute_scheduled_at(&mut self, at: SimTime) -> Result<(), CanwuError> {
221        if self
222            .state
223            .scheduler
224            .actions
225            .first_key_value()
226            .is_some_and(|(key, _)| key.at == at)
227        {
228            self.invalidate_commitments(CommitmentDomains::SCHEDULER);
229        }
230        while self
231            .state
232            .scheduler
233            .actions
234            .first_key_value()
235            .is_some_and(|(key, _)| key.at == at)
236        {
237            let (_, action) = self
238                .state
239                .scheduler
240                .actions
241                .pop_first()
242                .expect("scheduler was checked as non-empty");
243            self.execute_scheduled(action)?;
244        }
245        Ok(())
246    }
247
248    fn execute_scheduled(&mut self, action: ScheduledAction) -> Result<(), CanwuError> {
249        match action {
250            ScheduledAction::ArmyArrival {
251                army,
252                destination,
253                order_event,
254                correlation_id,
255            } => self.execute_arrival(army, destination, order_event, correlation_id),
256            ScheduledAction::PersonArrival {
257                person,
258                destination,
259                order_event,
260                cargo,
261                correlation_id,
262            } => self.execute_person_arrival(
263                person,
264                destination,
265                order_event,
266                &cargo,
267                correlation_id,
268            ),
269            ScheduledAction::KnowledgeReport {
270                recipient,
271                army,
272                location,
273                observed_at,
274                dispatch_event,
275                correlation_id,
276            } => {
277                self.update_army_knowledge(
278                    recipient,
279                    army,
280                    location,
281                    observed_at,
282                    KnowledgeSource::Report {
283                        source_event: dispatch_event,
284                    },
285                    850,
286                );
287                self.emit(
288                    KnowledgeUpdated {
289                        recipient,
290                        army,
291                        known_location: location,
292                    }
293                    .into_kind(),
294                    vec![EntityRef::Person(recipient), EntityRef::Army(army)],
295                    format!(
296                        "Person {recipient} received a report locating army {army} at {location}"
297                    ),
298                    Some(CauseRef::Event(dispatch_event)),
299                    correlation_id,
300                )?;
301                Ok(())
302            }
303            ScheduledAction::PluginDirective {
304                plugin,
305                directive,
306                allowed_writes,
307                cause,
308                correlation_id,
309            } => {
310                let directives = vec![*directive];
311                validate_directives_with_context(
312                    &RuntimeValidationContext::new(&self.state),
313                    &plugin,
314                    &allowed_writes,
315                    &self.plugins.state_owners,
316                    &self.plugins.record_schemas,
317                    &directives,
318                )?;
319                self.apply_directives(&plugin, directives, &allowed_writes, &cause, correlation_id)
320            }
321        }
322    }
323
324    fn execute_arrival(
325        &mut self,
326        army: ArmyId,
327        destination: TerritoryId,
328        order_event: EventId,
329        correlation_id: u64,
330    ) -> Result<(), CanwuError> {
331        self.invalidate_commitments(CommitmentDomains::WORLD);
332        let commander = {
333            let army_state = self.state.current.armies.get_mut(&army).ok_or_else(|| {
334                CanwuError::new(ErrorCode::ArmyNotFound, "scheduled army no longer exists")
335            })?;
336            army_state.location = destination;
337            army_state.transit = None;
338            army_state.commander
339        };
340        let arrival_event = self.emit(
341            ArmyArrived {
342                army,
343                territory: destination,
344            }
345            .into_kind(),
346            vec![EntityRef::Army(army), EntityRef::Territory(destination)],
347            format!("Army {army} arrived in territory {destination}"),
348            Some(CauseRef::Event(order_event)),
349            correlation_id,
350        )?;
351
352        self.update_army_knowledge(
353            commander,
354            army,
355            destination,
356            self.state.scheduler.now,
357            KnowledgeSource::CommandResponsibility,
358            1000,
359        );
360        self.emit(
361            KnowledgeUpdated {
362                recipient: commander,
363                army,
364                known_location: destination,
365            }
366            .into_kind(),
367            vec![EntityRef::Person(commander), EntityRef::Army(army)],
368            format!("Commander {commander} learned that army {army} arrived at {destination}"),
369            Some(CauseRef::Event(arrival_event)),
370            correlation_id,
371        )?;
372
373        let recipients: Vec<_> = self
374            .state
375            .current
376            .people
377            .keys()
378            .copied()
379            .filter(|person| *person != commander)
380            .collect();
381        for recipient in recipients {
382            let (draw_id, jitter) = self.draw_random(
383                &random::core_report_delay_stream(),
384                12 * 60,
385                "knowledge report delivery jitter",
386                RandomDrawProducer::CoreSystem {
387                    system: "canwu.core.knowledge-report-delay".to_owned(),
388                },
389                CauseRef::Event(arrival_event),
390                correlation_id,
391            )?;
392            let jitter_minutes =
393                i64::try_from(jitter).expect("report jitter is bounded to a small integer");
394            let arrives_at = self
395                .state
396                .scheduler
397                .now
398                .checked_add(SimDuration::hours(36))
399                .and_then(|time| time.checked_add(SimDuration::minutes(jitter_minutes)))
400                .ok_or_else(|| {
401                    CanwuError::new(
402                        ErrorCode::InvalidDuration,
403                        "knowledge report arrival time exceeds the supported range",
404                    )
405                })?;
406            let dispatch_event = self.emit(
407                ReportDispatched {
408                    recipient,
409                    army,
410                    arrives_at,
411                }
412                .into_kind(),
413                vec![EntityRef::Person(recipient), EntityRef::Army(army)],
414                format!("A report about army {army} was dispatched to person {recipient}"),
415                Some(CauseRef::Event(arrival_event)),
416                correlation_id,
417            )?;
418            self.record_random_outcome(
419                draw_id,
420                RandomDrawOutcome::KnowledgeReportDelivery {
421                    recipient,
422                    army,
423                    dispatch_event,
424                    arrives_at,
425                },
426            )?;
427            self.schedule_at(
428                arrives_at,
429                ScheduledAction::KnowledgeReport {
430                    recipient,
431                    army,
432                    location: destination,
433                    observed_at: self.state.scheduler.now,
434                    dispatch_event,
435                    correlation_id,
436                },
437            )?;
438        }
439        Ok(())
440    }
441
442    fn execute_person_arrival(
443        &mut self,
444        person: PersonId,
445        destination: TerritoryId,
446        order_event: EventId,
447        cargo: &[super::LetterId],
448        correlation_id: u64,
449    ) -> Result<(), CanwuError> {
450        self.invalidate_commitments(CommitmentDomains::WORLD);
451        {
452            let person_state = self.state.current.people.get_mut(&person).ok_or_else(|| {
453                CanwuError::new(
454                    ErrorCode::EntityNotFound,
455                    "scheduled person no longer exists",
456                )
457            })?;
458            let transit = person_state.transit.take().ok_or_else(|| {
459                CanwuError::new(
460                    ErrorCode::InvalidSnapshot,
461                    "person arrival has no matching transit",
462                )
463            })?;
464            if transit.to != destination || transit.arrives_at != self.state.scheduler.now {
465                return Err(CanwuError::new(
466                    ErrorCode::InvalidSnapshot,
467                    "person arrival disagrees with its transit state",
468                ));
469            }
470            person_state.current_location = destination;
471        }
472        let arrival_event = self.emit(
473            PersonArrived {
474                person,
475                territory: destination,
476            }
477            .into_kind(),
478            vec![EntityRef::Person(person), EntityRef::Territory(destination)],
479            format!("Person {person} arrived in territory {destination}"),
480            Some(CauseRef::Event(order_event)),
481            correlation_id,
482        )?;
483
484        for letter_id in cargo {
485            self.settle_letter_at_arrival(
486                *letter_id,
487                person,
488                destination,
489                arrival_event,
490                correlation_id,
491            )?;
492        }
493        let waiting_letters: Vec<_> = self
494            .state
495            .current
496            .letters
497            .values()
498            .filter(|letter| {
499                letter.status == LetterStatus::HeldAtLocation
500                    && letter.location == Some(destination)
501                    && letter.recipient == person
502            })
503            .map(|letter| letter.id)
504            .collect();
505        for letter_id in waiting_letters {
506            self.deliver_letter(
507                letter_id,
508                person,
509                destination,
510                arrival_event,
511                correlation_id,
512            )?;
513        }
514        Ok(())
515    }
516
517    fn settle_letter_at_arrival(
518        &mut self,
519        letter_id: super::LetterId,
520        carrier: PersonId,
521        destination: TerritoryId,
522        arrival_event: EventId,
523        correlation_id: u64,
524    ) -> Result<(), CanwuError> {
525        let recipient = self
526            .state
527            .current
528            .letters
529            .get(&letter_id)
530            .ok_or_else(|| CanwuError::new(ErrorCode::EntityNotFound, "arrival cargo disappeared"))?
531            .recipient;
532        let recipient_is_present =
533            self.state
534                .current
535                .people
536                .get(&recipient)
537                .is_some_and(|person| {
538                    person.current_location == destination && person.transit.is_none()
539                });
540        if recipient_is_present {
541            self.deliver_letter(
542                letter_id,
543                carrier,
544                destination,
545                arrival_event,
546                correlation_id,
547            )
548        } else {
549            let letter = self
550                .state
551                .current
552                .letters
553                .get_mut(&letter_id)
554                .ok_or_else(|| {
555                    CanwuError::new(ErrorCode::EntityNotFound, "arrival cargo disappeared")
556                })?;
557            letter.status = LetterStatus::HeldAtLocation;
558            letter.carrier = None;
559            letter.location = Some(destination);
560            Ok(())
561        }
562    }
563
564    fn deliver_letter(
565        &mut self,
566        letter_id: super::LetterId,
567        carrier: PersonId,
568        territory: TerritoryId,
569        arrival_event: EventId,
570        correlation_id: u64,
571    ) -> Result<(), CanwuError> {
572        let letter = self
573            .state
574            .current
575            .letters
576            .get_mut(&letter_id)
577            .ok_or_else(|| {
578                CanwuError::new(ErrorCode::EntityNotFound, "delivery letter disappeared")
579            })?;
580        let sender = letter.sender;
581        let recipient = letter.recipient;
582        letter.status = LetterStatus::Delivered;
583        letter.carrier = None;
584        letter.location = Some(territory);
585        letter.delivered_at = Some(self.state.scheduler.now);
586        self.emit(
587            LetterDelivered {
588                letter: letter_id,
589                carrier,
590                recipient,
591                territory,
592            }
593            .into_kind(),
594            vec![
595                EntityRef::Resource(super::ResourceId::new(letter_id.get())),
596                EntityRef::Person(sender),
597                EntityRef::Person(carrier),
598                EntityRef::Person(recipient),
599                EntityRef::Territory(territory),
600            ],
601            format!("Letter {letter_id} was delivered to person {recipient}"),
602            Some(CauseRef::Event(arrival_event)),
603            correlation_id,
604        )?;
605        Ok(())
606    }
607
608    fn update_army_knowledge(
609        &mut self,
610        recipient: PersonId,
611        army: ArmyId,
612        location: TerritoryId,
613        observed_at: SimTime,
614        source: KnowledgeSource,
615        confidence_per_mille: u16,
616    ) {
617        self.invalidate_commitments(CommitmentDomains::KNOWLEDGE);
618        let (strength, known_name) = self.state.current.armies.get(&army).map_or_else(
619            || (0, None),
620            |value| (value.strength, Some(value.name.clone())),
621        );
622        let actor = self
623            .state
624            .current
625            .knowledge
626            .actors
627            .entry(recipient)
628            .or_insert_with(|| ActorKnowledge {
629                actor: recipient,
630                armies: BTreeMap::new(),
631            });
632        actor.armies.insert(
633            army,
634            ArmyKnowledge {
635                army,
636                known_name,
637                known_location: Some(location),
638                estimated_strength: EstimateRange {
639                    minimum: strength.saturating_mul(9) / 10,
640                    maximum: strength.saturating_mul(11) / 10,
641                },
642                observed_at,
643                learned_at: self.state.scheduler.now,
644                confidence_per_mille,
645                source,
646            },
647        );
648    }
649
650    pub(super) fn emit(
651        &mut self,
652        kind: EventKind,
653        affected_entities: Vec<EntityRef>,
654        summary: String,
655        cause: Option<CauseRef>,
656        correlation_id: u64,
657    ) -> Result<EventId, CanwuError> {
658        let previous_depth = self.sync_reaction_depth;
659        if previous_depth >= super::MAX_SYNCHRONOUS_REACTION_DEPTH {
660            return Err(CanwuError::new(
661                ErrorCode::SynchronousReactionLimit,
662                format!(
663                    "synchronous event reactors exceeded the maximum nested depth of {}",
664                    super::MAX_SYNCHRONOUS_REACTION_DEPTH
665                ),
666            ));
667        }
668        self.sync_reaction_depth = previous_depth + 1;
669        let result = self.emit_immediate(kind, affected_entities, summary, cause, correlation_id);
670        self.sync_reaction_depth = previous_depth;
671        result
672    }
673
674    fn emit_immediate(
675        &mut self,
676        kind: EventKind,
677        affected_entities: Vec<EntityRef>,
678        summary: String,
679        cause: Option<CauseRef>,
680        correlation_id: u64,
681    ) -> Result<EventId, CanwuError> {
682        let event = self.append_event(kind, affected_entities, summary, cause, correlation_id)?;
683        let id = event.id;
684
685        let systems = self.plugins.systems.clone();
686        for registered in systems {
687            let reader = format!("{}.{}", registered.plugin, registered.contract.name);
688            let directives = catch_unwind(AssertUnwindSafe(|| {
689                (registered.handler)(
690                    &self.plugin_view(&reader, &registered.contract.reads),
691                    &event,
692                )
693            }))
694            .map_err(|_| {
695                CanwuError::new(
696                    ErrorCode::PluginPanicked,
697                    format!(
698                        "plugin system {}.{} panicked",
699                        registered.plugin, registered.contract.name
700                    ),
701                )
702            })??;
703            validate_directives_with_context(
704                &RuntimeValidationContext::new(&self.state),
705                &registered.plugin,
706                &registered.contract.writes,
707                &self.plugins.state_owners,
708                &self.plugins.record_schemas,
709                &directives,
710            )?;
711            self.apply_directives(
712                &registered.plugin,
713                directives,
714                &registered.contract.writes,
715                &CauseRef::Event(id),
716                correlation_id,
717            )?;
718        }
719        Ok(id)
720    }
721
722    pub(super) fn append_event(
723        &mut self,
724        mut kind: EventKind,
725        affected_entities: Vec<EntityRef>,
726        summary: String,
727        cause: Option<CauseRef>,
728        correlation_id: u64,
729    ) -> Result<SimEvent, CanwuError> {
730        canonicalize_event_kind(&mut kind).map_err(|error| {
731            CanwuError::new(
732                ErrorCode::InvalidPayload,
733                format!("event payload is not canonical: {error}"),
734            )
735        })?;
736        let (event_id, next_event_id) =
737            claim_counter(self.state.counters.next_event_id, "event ID")?;
738        let id = EventId::new(event_id);
739        self.state.counters.next_event_id = next_event_id;
740        let event = SimEvent {
741            id,
742            timestamp: self.state.scheduler.now,
743            kind,
744            affected_entities,
745            summary,
746            cause,
747            correlation_id,
748        };
749        self.state.evidence.events.push(event.clone());
750        Ok(event)
751    }
752
753    fn draw_random(
754        &mut self,
755        stream: &RandomStreamKey,
756        upper_exclusive: u64,
757        purpose: &str,
758        producer: RandomDrawProducer,
759        cause: CauseRef,
760        correlation_id: u64,
761    ) -> Result<(RandomDrawId, u64), CanwuError> {
762        if upper_exclusive == 0
763            || purpose.trim().is_empty()
764            || purpose != purpose.trim()
765            || correlation_id == 0
766        {
767            return Err(CanwuError::new(
768                ErrorCode::InvalidRandomDraw,
769                "random draws require a positive bound, canonical purpose, and correlation",
770            ));
771        }
772        let (draw_id, next_random_draw_id) =
773            claim_counter(self.state.counters.next_random_draw_id, "random draw ID")?;
774        self.invalidate_commitments(CommitmentDomains::RANDOM_STREAMS);
775        let state = self
776            .state
777            .current
778            .random_streams
779            .get_mut(stream)
780            .ok_or_else(|| {
781                CanwuError::new(
782                    ErrorCode::InvalidRandomStream,
783                    format!(
784                        "random stream {}.{}@{} is not initialized",
785                        stream.namespace, stream.name, stream.version
786                    ),
787                )
788            })?;
789        let next_position = state.position.checked_add(1).ok_or_else(|| {
790            CanwuError::new(
791                ErrorCode::IdentifierExhausted,
792                "random stream position is exhausted",
793            )
794        })?;
795        let position = state.position;
796        let mut generator = DeterministicRng::from_seed(state.generator_state);
797        let value = match state.algorithm {
798            RandomAlgorithm::SplitMix64V1 => generator.range_modulo(upper_exclusive),
799            RandomAlgorithm::SplitMix64V2 => generator.range(upper_exclusive),
800        };
801        state.position = next_position;
802        state.generator_state = generator.state();
803        self.state.counters.next_random_draw_id = next_random_draw_id;
804        let id = RandomDrawId::new(draw_id);
805        self.state.evidence.random_draws.push(RandomDrawRecord {
806            id,
807            at: self.state.scheduler.now,
808            stream: stream.clone(),
809            address: RandomDrawAddress::Sequential { position },
810            operation_evidence: None,
811            upper_exclusive,
812            value,
813            purpose: purpose.to_owned(),
814            producer,
815            outcome: None,
816            cause,
817            correlation_id,
818        });
819        Ok((id, value))
820    }
821
822    fn record_random_outcome(
823        &mut self,
824        id: RandomDrawId,
825        outcome: RandomDrawOutcome,
826    ) -> Result<(), CanwuError> {
827        let Some(draw) = self
828            .state
829            .evidence
830            .random_draws
831            .last_mut()
832            .filter(|draw| draw.id == id)
833        else {
834            return Err(CanwuError::new(
835                ErrorCode::InvalidRandomDraw,
836                "random draw outcome does not match the latest pending draw",
837            ));
838        };
839        if draw.outcome.replace(outcome).is_some() {
840            return Err(CanwuError::new(
841                ErrorCode::InvalidRandomDraw,
842                "random draw outcome was already recorded",
843            ));
844        }
845        Ok(())
846    }
847
848    pub(super) fn append_boundary_random_draws(
849        &mut self,
850        boundary: BoundaryId,
851        correlation_id: u64,
852        draws: Vec<PendingBoundaryRandomDraw>,
853        mut outcomes: BTreeMap<(RandomStreamKey, RandomDrawAddress), RandomDrawOutcome>,
854    ) -> Result<Vec<CommittedBoundaryRandomDraw>, CanwuError> {
855        let mut ids = Vec::with_capacity(draws.len());
856        for pending in draws {
857            let (draw_id, next_random_draw_id) =
858                claim_counter(self.state.counters.next_random_draw_id, "random draw ID")?;
859            let id = RandomDrawId::new(draw_id);
860            self.state.counters.next_random_draw_id = next_random_draw_id;
861            let stream = pending.draw.stream;
862            let address = pending.draw.address;
863            let outcome = outcomes
864                .remove(&(stream.clone(), address.clone()))
865                .unwrap_or(RandomDrawOutcome::BoundarySystemDecision);
866            self.state.evidence.random_draws.push(RandomDrawRecord {
867                id,
868                at: self.state.scheduler.now,
869                stream: stream.clone(),
870                address: address.clone(),
871                operation_evidence: pending.draw.operation_evidence,
872                upper_exclusive: pending.draw.upper_exclusive,
873                value: pending.draw.value,
874                purpose: pending.draw.purpose,
875                producer: RandomDrawProducer::BoundarySystem {
876                    boundary,
877                    plugin: pending.plugin,
878                    system: pending.system,
879                },
880                outcome: Some(outcome),
881                cause: CauseRef::Boundary(boundary),
882                correlation_id,
883            });
884            ids.push(CommittedBoundaryRandomDraw {
885                id,
886                stream,
887                address,
888            });
889        }
890        if !outcomes.is_empty() {
891            return Err(CanwuError::new(
892                ErrorCode::InvalidRandomDraw,
893                "random decision references a draw that was not committed by its boundary",
894            ));
895        }
896        Ok(ids)
897    }
898
899    pub(super) fn schedule_at(
900        &mut self,
901        at: SimTime,
902        action: ScheduledAction,
903    ) -> Result<(), CanwuError> {
904        if at <= self.state.scheduler.now {
905            return Err(CanwuError::new(
906                ErrorCode::InvalidDuration,
907                "scheduled work must target a strictly future simulation time",
908            ));
909        }
910        let (sequence, next_sequence) = claim_counter(
911            self.state.counters.next_schedule_sequence,
912            "schedule sequence",
913        )?;
914        let key = ScheduleKey { at, sequence };
915        self.state.counters.next_schedule_sequence = next_sequence;
916        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
917        if self.state.scheduler.actions.insert(key, action).is_some() {
918            return Err(CanwuError::new(
919                ErrorCode::InvalidSnapshot,
920                "the runtime attempted to reuse a schedule key",
921            ));
922        }
923        Ok(())
924    }
925
926    pub(super) fn plugin_view<'a>(
927        &'a self,
928        reader: &'a str,
929        reads: &'a [StateKey],
930    ) -> SimulationView<'a> {
931        SimulationView {
932            state: SimulationViewState::Runtime(&self.state),
933            state_owners: &self.plugins.state_owners,
934            reader: Some(reader),
935            allowed_reads: Some(reads),
936            allowed_ingress: None,
937            ingress_plugin: None,
938            component_overlay: None,
939            proposed_components: None,
940            record_overlay: None,
941            proposed_records: None,
942            boundary_id: None,
943            proposal_evidence: None,
944            knowledge_overlay: None,
945            allocations: None,
946            allowed_reservations: None,
947            random_session: None,
948            plugin_archive_provider: self.plugin_archive_provider.as_ref(),
949            transitions: None,
950        }
951    }
952}