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, ®istered.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 ®istered.plugin,
694 ®istered.contract.writes,
695 &self.plugins.state_owners,
696 &self.plugins.record_schemas,
697 &directives,
698 )?;
699 self.apply_directives(
700 ®istered.plugin,
701 directives,
702 ®istered.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}