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