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, ®istered.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 ®istered.plugin,
705 ®istered.contract.writes,
706 &self.plugins.state_owners,
707 &self.plugins.record_schemas,
708 &directives,
709 )?;
710 self.apply_directives(
711 ®istered.plugin,
712 directives,
713 ®istered.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}