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