Skip to main content

canwu_sim/
scheduling.rs

1use super::{
2    ActorKnowledge, ArmyId, ArmyKnowledge, AssertUnwindSafe, BTreeMap, BoundaryId, CanwuError,
3    CauseRef, ClockTransactionCheckpoint, CommitmentDomains, DeterministicRng, EntityRef,
4    ErrorCode, EstimateRange, EventId, EventKind, KnowledgeSource, PendingBoundaryRandomDraw,
5    PersonId, RandomDrawId, RandomDrawOutcome, RandomDrawProducer, RandomDrawRecord,
6    RandomStreamKey, RuntimeValidationContext, ScheduleKey, ScheduledAction,
7    ScheduledBatchTransactionCheckpoint, SimDuration, SimEvent, SimTime, Simulation,
8    SimulationView, SimulationViewState, StateKey, TerritoryId, catch_unwind, claim_counter,
9    random, validate_directives_with_context,
10};
11
12impl Simulation {
13    pub fn advance(&mut self, duration: SimDuration) -> Result<Vec<SimEvent>, CanwuError> {
14        self.ensure_runtime_ready()?;
15        if duration.is_negative() {
16            return Err(CanwuError::new(
17                ErrorCode::InvalidDuration,
18                "simulation time cannot advance by a negative duration",
19            ));
20        }
21        let target = self
22            .state
23            .scheduler
24            .now
25            .checked_add(duration)
26            .ok_or_else(|| {
27                CanwuError::new(
28                    ErrorCode::InvalidDuration,
29                    "simulation target time exceeds the supported range",
30                )
31            })?;
32        self.ensure_legacy_advance_does_not_cross_ingress(target)?;
33        self.advance_to(target)
34    }
35
36    pub fn step(&mut self) -> Result<Vec<SimEvent>, CanwuError> {
37        self.ensure_runtime_ready()?;
38        if self.state.scheduler.pending_ingress.first().is_some() {
39            return Err(CanwuError::new(
40                ErrorCode::InvalidBoundary,
41                "pending canonical ingress requires step_canonical",
42            ));
43        }
44        let Some(next_time) = self.state.scheduler.actions.keys().next().map(|key| key.at) else {
45            return Ok(Vec::new());
46        };
47        self.advance_to(next_time)
48    }
49
50    pub fn advance_until<F>(
51        &mut self,
52        maximum: SimDuration,
53        mut condition: F,
54    ) -> Result<Vec<SimEvent>, CanwuError>
55    where
56        F: FnMut(&Self) -> bool,
57    {
58        self.ensure_runtime_ready()?;
59        if maximum.is_negative() {
60            return Err(CanwuError::new(
61                ErrorCode::InvalidDuration,
62                "advance_until maximum cannot be negative",
63            ));
64        }
65        let target = self
66            .state
67            .scheduler
68            .now
69            .checked_add(maximum)
70            .ok_or_else(|| {
71                CanwuError::new(
72                    ErrorCode::InvalidDuration,
73                    "advance_until target time exceeds the supported range",
74                )
75            })?;
76        self.ensure_legacy_advance_does_not_cross_ingress(target)?;
77        let start = self.state.evidence.events.len();
78        while self.state.scheduler.now < target && !condition(self) {
79            let next_time = self
80                .state
81                .scheduler
82                .actions
83                .keys()
84                .next()
85                .map_or(target, |key| key.at.min(target));
86            self.advance_to(next_time)?;
87            if next_time == target {
88                break;
89            }
90        }
91        Ok(self.state.evidence.events[start..].to_vec())
92    }
93
94    pub(super) fn ensure_legacy_advance_does_not_cross_ingress(
95        &self,
96        target: SimTime,
97    ) -> Result<(), CanwuError> {
98        if self
99            .state
100            .scheduler
101            .pending_ingress
102            .first()
103            .is_some_and(|key| key.due_at <= target)
104        {
105            return Err(CanwuError::new(
106                ErrorCode::InvalidBoundary,
107                "legacy time advancement cannot cross pending canonical ingress; use advance_canonical",
108            ));
109        }
110        Ok(())
111    }
112
113    pub(super) fn advance_to(&mut self, target: SimTime) -> Result<Vec<SimEvent>, CanwuError> {
114        let start = self.state.evidence.events.len();
115        while let Some(boundary_time) = self.state.scheduler.actions.keys().next().map(|key| key.at)
116            && boundary_time <= target
117        {
118            let transaction = ScheduledBatchTransactionCheckpoint::capture(&self.state);
119            let result = (|| {
120                self.invalidate_commitments(CommitmentDomains::SCHEDULER);
121                self.state.scheduler.now = boundary_time;
122                while self
123                    .state
124                    .scheduler
125                    .actions
126                    .first_key_value()
127                    .is_some_and(|(key, _)| key.at == boundary_time)
128                {
129                    let (_, action) = self
130                        .state
131                        .scheduler
132                        .actions
133                        .pop_first()
134                        .expect("scheduler was checked as non-empty");
135                    self.execute_scheduled(action)?;
136                }
137                self.state.metadata.plugin_registration_closed = true;
138                self.refresh_checkpoint_hash()
139            })();
140            if let Err(error) = result {
141                transaction.restore(&mut self.state);
142                return Err(error);
143            }
144        }
145        let transaction = ClockTransactionCheckpoint::capture(&self.state);
146        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
147        self.state.scheduler.now = target;
148        self.state.metadata.plugin_registration_closed = true;
149        if let Err(error) = self.refresh_checkpoint_hash() {
150            transaction.restore(&mut self.state);
151            return Err(error);
152        }
153        Ok(self.state.evidence.events[start..].to_vec())
154    }
155
156    pub(super) fn advance_to_before_boundary(&mut self, target: SimTime) -> Result<(), CanwuError> {
157        while let Some(next) = self.state.scheduler.actions.keys().next().map(|key| key.at)
158            && next < target
159        {
160            self.advance_to(next)?;
161        }
162        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
163        self.state.scheduler.now = target;
164        self.state.metadata.plugin_registration_closed = true;
165        self.refresh_checkpoint_hash()
166    }
167
168    pub(super) fn execute_scheduled_at(&mut self, at: SimTime) -> Result<(), CanwuError> {
169        if self
170            .state
171            .scheduler
172            .actions
173            .first_key_value()
174            .is_some_and(|(key, _)| key.at == at)
175        {
176            self.invalidate_commitments(CommitmentDomains::SCHEDULER);
177        }
178        while self
179            .state
180            .scheduler
181            .actions
182            .first_key_value()
183            .is_some_and(|(key, _)| key.at == at)
184        {
185            let (_, action) = self
186                .state
187                .scheduler
188                .actions
189                .pop_first()
190                .expect("scheduler was checked as non-empty");
191            self.execute_scheduled(action)?;
192        }
193        Ok(())
194    }
195
196    fn execute_scheduled(&mut self, action: ScheduledAction) -> Result<(), CanwuError> {
197        match action {
198            ScheduledAction::ArmyArrival {
199                army,
200                destination,
201                order_event,
202                correlation_id,
203            } => self.execute_arrival(army, destination, order_event, correlation_id),
204            ScheduledAction::KnowledgeReport {
205                recipient,
206                army,
207                location,
208                observed_at,
209                dispatch_event,
210                correlation_id,
211            } => {
212                self.update_army_knowledge(
213                    recipient,
214                    army,
215                    location,
216                    observed_at,
217                    KnowledgeSource::Report {
218                        source_event: dispatch_event,
219                    },
220                    850,
221                );
222                self.emit(
223                    EventKind::KnowledgeUpdated {
224                        recipient,
225                        army,
226                        known_location: location,
227                    },
228                    vec![EntityRef::Person(recipient), EntityRef::Army(army)],
229                    format!(
230                        "Person {recipient} received a report locating army {army} at {location}"
231                    ),
232                    Some(CauseRef::Event(dispatch_event)),
233                    correlation_id,
234                )?;
235                Ok(())
236            }
237            ScheduledAction::PluginDirective {
238                plugin,
239                directive,
240                allowed_writes,
241                cause,
242                correlation_id,
243            } => {
244                let directives = vec![*directive];
245                validate_directives_with_context(
246                    &RuntimeValidationContext::new(&self.state),
247                    &plugin,
248                    &allowed_writes,
249                    &self.plugins.state_owners,
250                    &self.plugins.record_schemas,
251                    &directives,
252                )?;
253                self.apply_directives(&plugin, directives, &allowed_writes, &cause, correlation_id)
254            }
255        }
256    }
257
258    fn execute_arrival(
259        &mut self,
260        army: ArmyId,
261        destination: TerritoryId,
262        order_event: EventId,
263        correlation_id: u64,
264    ) -> Result<(), CanwuError> {
265        self.invalidate_commitments(CommitmentDomains::WORLD);
266        let commander = {
267            let army_state = self.state.current.armies.get_mut(&army).ok_or_else(|| {
268                CanwuError::new(ErrorCode::ArmyNotFound, "scheduled army no longer exists")
269            })?;
270            army_state.location = destination;
271            army_state.transit = None;
272            army_state.commander
273        };
274        let arrival_event = self.emit(
275            EventKind::ArmyArrived {
276                army,
277                territory: destination,
278            },
279            vec![EntityRef::Army(army), EntityRef::Territory(destination)],
280            format!("Army {army} arrived in territory {destination}"),
281            Some(CauseRef::Event(order_event)),
282            correlation_id,
283        )?;
284
285        self.update_army_knowledge(
286            commander,
287            army,
288            destination,
289            self.state.scheduler.now,
290            KnowledgeSource::CommandResponsibility,
291            1000,
292        );
293        self.emit(
294            EventKind::KnowledgeUpdated {
295                recipient: commander,
296                army,
297                known_location: destination,
298            },
299            vec![EntityRef::Person(commander), EntityRef::Army(army)],
300            format!("Commander {commander} learned that army {army} arrived at {destination}"),
301            Some(CauseRef::Event(arrival_event)),
302            correlation_id,
303        )?;
304
305        let recipients: Vec<_> = self
306            .state
307            .current
308            .people
309            .keys()
310            .copied()
311            .filter(|person| *person != commander)
312            .collect();
313        for recipient in recipients {
314            let (draw_id, jitter) = self.draw_random(
315                &random::core_report_delay_stream(),
316                12 * 60,
317                "knowledge report delivery jitter",
318                RandomDrawProducer::CoreSystem {
319                    system: "canwu.core.knowledge-report-delay".to_owned(),
320                },
321                CauseRef::Event(arrival_event),
322                correlation_id,
323            )?;
324            let jitter_minutes =
325                i64::try_from(jitter).expect("report jitter is bounded to a small integer");
326            let arrives_at = self
327                .state
328                .scheduler
329                .now
330                .checked_add(SimDuration::hours(36))
331                .and_then(|time| time.checked_add(SimDuration::minutes(jitter_minutes)))
332                .ok_or_else(|| {
333                    CanwuError::new(
334                        ErrorCode::InvalidDuration,
335                        "knowledge report arrival time exceeds the supported range",
336                    )
337                })?;
338            let dispatch_event = self.emit(
339                EventKind::ReportDispatched {
340                    recipient,
341                    army,
342                    arrives_at,
343                },
344                vec![EntityRef::Person(recipient), EntityRef::Army(army)],
345                format!("A report about army {army} was dispatched to person {recipient}"),
346                Some(CauseRef::Event(arrival_event)),
347                correlation_id,
348            )?;
349            self.record_random_outcome(
350                draw_id,
351                RandomDrawOutcome::KnowledgeReportDelivery {
352                    recipient,
353                    army,
354                    dispatch_event,
355                    arrives_at,
356                },
357            )?;
358            self.schedule_at(
359                arrives_at,
360                ScheduledAction::KnowledgeReport {
361                    recipient,
362                    army,
363                    location: destination,
364                    observed_at: self.state.scheduler.now,
365                    dispatch_event,
366                    correlation_id,
367                },
368            )?;
369        }
370        Ok(())
371    }
372
373    fn update_army_knowledge(
374        &mut self,
375        recipient: PersonId,
376        army: ArmyId,
377        location: TerritoryId,
378        observed_at: SimTime,
379        source: KnowledgeSource,
380        confidence_per_mille: u16,
381    ) {
382        self.invalidate_commitments(CommitmentDomains::KNOWLEDGE);
383        let (strength, known_name) = self.state.current.armies.get(&army).map_or_else(
384            || (0, None),
385            |value| (value.strength, Some(value.name.clone())),
386        );
387        let actor = self
388            .state
389            .current
390            .knowledge
391            .actors
392            .entry(recipient)
393            .or_insert_with(|| ActorKnowledge {
394                actor: recipient,
395                armies: BTreeMap::new(),
396            });
397        actor.armies.insert(
398            army,
399            ArmyKnowledge {
400                army,
401                known_name,
402                known_location: Some(location),
403                estimated_strength: EstimateRange {
404                    minimum: strength.saturating_mul(9) / 10,
405                    maximum: strength.saturating_mul(11) / 10,
406                },
407                observed_at,
408                learned_at: self.state.scheduler.now,
409                confidence_per_mille,
410                source,
411            },
412        );
413    }
414
415    pub(super) fn emit(
416        &mut self,
417        kind: EventKind,
418        affected_entities: Vec<EntityRef>,
419        summary: String,
420        cause: Option<CauseRef>,
421        correlation_id: u64,
422    ) -> Result<EventId, CanwuError> {
423        let previous_depth = self.sync_reaction_depth;
424        if previous_depth >= super::MAX_SYNCHRONOUS_REACTION_DEPTH {
425            return Err(CanwuError::new(
426                ErrorCode::SynchronousReactionLimit,
427                format!(
428                    "synchronous event reactors exceeded the maximum nested depth of {}",
429                    super::MAX_SYNCHRONOUS_REACTION_DEPTH
430                ),
431            ));
432        }
433        self.sync_reaction_depth = previous_depth + 1;
434        let result = self.emit_immediate(kind, affected_entities, summary, cause, correlation_id);
435        self.sync_reaction_depth = previous_depth;
436        result
437    }
438
439    fn emit_immediate(
440        &mut self,
441        kind: EventKind,
442        affected_entities: Vec<EntityRef>,
443        summary: String,
444        cause: Option<CauseRef>,
445        correlation_id: u64,
446    ) -> Result<EventId, CanwuError> {
447        let event = self.append_event(kind, affected_entities, summary, cause, correlation_id)?;
448        let id = event.id;
449
450        let systems = self.plugins.systems.clone();
451        for registered in systems {
452            let reader = format!("{}.{}", registered.plugin, registered.contract.name);
453            let directives = catch_unwind(AssertUnwindSafe(|| {
454                (registered.handler)(
455                    &self.plugin_view(&reader, &registered.contract.reads),
456                    &event,
457                )
458            }))
459            .map_err(|_| {
460                CanwuError::new(
461                    ErrorCode::PluginPanicked,
462                    format!(
463                        "plugin system {}.{} panicked",
464                        registered.plugin, registered.contract.name
465                    ),
466                )
467            })??;
468            validate_directives_with_context(
469                &RuntimeValidationContext::new(&self.state),
470                &registered.plugin,
471                &registered.contract.writes,
472                &self.plugins.state_owners,
473                &self.plugins.record_schemas,
474                &directives,
475            )?;
476            self.apply_directives(
477                &registered.plugin,
478                directives,
479                &registered.contract.writes,
480                &CauseRef::Event(id),
481                correlation_id,
482            )?;
483        }
484        Ok(id)
485    }
486
487    pub(super) fn append_event(
488        &mut self,
489        kind: EventKind,
490        affected_entities: Vec<EntityRef>,
491        summary: String,
492        cause: Option<CauseRef>,
493        correlation_id: u64,
494    ) -> Result<SimEvent, CanwuError> {
495        let (event_id, next_event_id) =
496            claim_counter(self.state.counters.next_event_id, "event ID")?;
497        let id = EventId::new(event_id);
498        self.state.counters.next_event_id = next_event_id;
499        let event = SimEvent {
500            id,
501            timestamp: self.state.scheduler.now,
502            kind,
503            affected_entities,
504            summary,
505            cause,
506            correlation_id,
507        };
508        self.state.evidence.events.push(event.clone());
509        Ok(event)
510    }
511
512    fn draw_random(
513        &mut self,
514        stream: &RandomStreamKey,
515        upper_exclusive: u64,
516        purpose: &str,
517        producer: RandomDrawProducer,
518        cause: CauseRef,
519        correlation_id: u64,
520    ) -> Result<(RandomDrawId, u64), CanwuError> {
521        if upper_exclusive == 0
522            || purpose.trim().is_empty()
523            || purpose != purpose.trim()
524            || correlation_id == 0
525        {
526            return Err(CanwuError::new(
527                ErrorCode::InvalidRandomDraw,
528                "random draws require a positive bound, canonical purpose, and correlation",
529            ));
530        }
531        let (draw_id, next_random_draw_id) =
532            claim_counter(self.state.counters.next_random_draw_id, "random draw ID")?;
533        self.invalidate_commitments(CommitmentDomains::RANDOM_STREAMS);
534        let state = self
535            .state
536            .current
537            .random_streams
538            .get_mut(stream)
539            .ok_or_else(|| {
540                CanwuError::new(
541                    ErrorCode::InvalidRandomStream,
542                    format!(
543                        "random stream {}.{}@{} is not initialized",
544                        stream.namespace, stream.name, stream.version
545                    ),
546                )
547            })?;
548        let next_position = state.position.checked_add(1).ok_or_else(|| {
549            CanwuError::new(
550                ErrorCode::IdentifierExhausted,
551                "random stream position is exhausted",
552            )
553        })?;
554        let position = state.position;
555        let mut generator = DeterministicRng::from_seed(state.generator_state);
556        let value = generator.range(upper_exclusive);
557        state.position = next_position;
558        state.generator_state = generator.state();
559        self.state.counters.next_random_draw_id = next_random_draw_id;
560        let id = RandomDrawId::new(draw_id);
561        self.state.evidence.random_draws.push(RandomDrawRecord {
562            id,
563            at: self.state.scheduler.now,
564            stream: stream.clone(),
565            position,
566            upper_exclusive,
567            value,
568            purpose: purpose.to_owned(),
569            producer,
570            outcome: None,
571            cause,
572            correlation_id,
573        });
574        Ok((id, value))
575    }
576
577    fn record_random_outcome(
578        &mut self,
579        id: RandomDrawId,
580        outcome: RandomDrawOutcome,
581    ) -> Result<(), CanwuError> {
582        let Some(draw) = self
583            .state
584            .evidence
585            .random_draws
586            .last_mut()
587            .filter(|draw| draw.id == id)
588        else {
589            return Err(CanwuError::new(
590                ErrorCode::InvalidRandomDraw,
591                "random draw outcome does not match the latest pending draw",
592            ));
593        };
594        if draw.outcome.replace(outcome).is_some() {
595            return Err(CanwuError::new(
596                ErrorCode::InvalidRandomDraw,
597                "random draw outcome was already recorded",
598            ));
599        }
600        Ok(())
601    }
602
603    pub(super) fn append_boundary_random_draws(
604        &mut self,
605        boundary: BoundaryId,
606        correlation_id: u64,
607        draws: Vec<PendingBoundaryRandomDraw>,
608    ) -> Result<Vec<RandomDrawId>, CanwuError> {
609        let mut ids = Vec::with_capacity(draws.len());
610        for pending in draws {
611            let (draw_id, next_random_draw_id) =
612                claim_counter(self.state.counters.next_random_draw_id, "random draw ID")?;
613            let id = RandomDrawId::new(draw_id);
614            self.state.counters.next_random_draw_id = next_random_draw_id;
615            self.state.evidence.random_draws.push(RandomDrawRecord {
616                id,
617                at: self.state.scheduler.now,
618                stream: pending.draw.stream,
619                position: pending.draw.position,
620                upper_exclusive: pending.draw.upper_exclusive,
621                value: pending.draw.value,
622                purpose: pending.draw.purpose,
623                producer: RandomDrawProducer::BoundarySystem {
624                    boundary,
625                    plugin: pending.plugin,
626                    system: pending.system,
627                },
628                outcome: Some(RandomDrawOutcome::BoundarySystemDecision),
629                cause: CauseRef::Boundary(boundary),
630                correlation_id,
631            });
632            ids.push(id);
633        }
634        Ok(ids)
635    }
636
637    pub(super) fn schedule_at(
638        &mut self,
639        at: SimTime,
640        action: ScheduledAction,
641    ) -> Result<(), CanwuError> {
642        if at <= self.state.scheduler.now {
643            return Err(CanwuError::new(
644                ErrorCode::InvalidDuration,
645                "scheduled work must target a strictly future simulation time",
646            ));
647        }
648        let (sequence, next_sequence) = claim_counter(
649            self.state.counters.next_schedule_sequence,
650            "schedule sequence",
651        )?;
652        let key = ScheduleKey { at, sequence };
653        self.state.counters.next_schedule_sequence = next_sequence;
654        self.invalidate_commitments(CommitmentDomains::SCHEDULER);
655        if self.state.scheduler.actions.insert(key, action).is_some() {
656            return Err(CanwuError::new(
657                ErrorCode::InvalidSnapshot,
658                "the runtime attempted to reuse a schedule key",
659            ));
660        }
661        Ok(())
662    }
663
664    pub(super) fn plugin_view<'a>(
665        &'a self,
666        reader: &'a str,
667        reads: &'a [StateKey],
668    ) -> SimulationView<'a> {
669        SimulationView {
670            state: SimulationViewState::Runtime(&self.state),
671            state_owners: &self.plugins.state_owners,
672            reader: Some(reader),
673            allowed_reads: Some(reads),
674            allowed_ingress: None,
675            ingress_plugin: None,
676            component_overlay: None,
677            proposed_components: None,
678            record_overlay: None,
679            proposed_records: None,
680            allocations: None,
681            allowed_reservations: None,
682            random_session: None,
683        }
684    }
685}