Skip to main content

canwu_sim/runtime/
scheduling.rs

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