1mod boundary;
2mod decision;
3mod error;
4mod hashing;
5mod ingress;
6mod knowledge;
7mod legacy_v4;
8mod manifest;
9mod migration;
10mod persistence;
11mod plugins;
12mod policy;
13mod random;
14mod records;
15mod replay;
16mod scenario;
17mod scheduling;
18mod settlement;
19mod state;
20mod transactions;
21mod validation;
22mod view;
23
24pub use hashing::{canonical_byte_hash, canonical_hash};
25
26pub use boundary::{
27 BoundaryChange, BoundaryContext, BoundaryDirective, BoundaryEmission, BoundaryEmissionKind,
28 BoundaryIngressGeneration, BoundaryKnowledgeChange, BoundaryProposal, BoundaryReceipt,
29 BoundaryRecord, BoundaryRequest, BoundarySystemContract, BoundarySystemHandler,
30 KnowledgeWriteGrant, PluginIngressTarget, ReservationAllocation, ReservationDisposition,
31 ReservationOffer, ReservationOfferRecord, ReservationPoolKey, ReservationRef,
32 ReservationRequest, ReservationRequestRecord,
33};
34pub use canwu_core::{
35 DomainRecordVersionRef, DomainRecordVersionSource, EvidenceRef, HolderKnowledgeRecordId,
36 KnowledgeHolderPolicy, KnowledgeHolderRef, KnowledgeRecordId, KnowledgeRecordKind,
37 KnowledgeSchemaId,
38};
39pub use canwu_decision::{
40 ControllerDecision, DecisionAction, DecisionAttemptErrorCode, DecisionAttemptOutcome,
41 DecisionAttemptRecord, DecisionAuthority, DecisionContext, DecisionController,
42 DecisionControllerBinding, DecisionError, DecisionErrorCode, DecisionExternalEvidence,
43 DecisionFactorContribution, DecisionMutation, DecisionOption, DecisionOptionEvaluation,
44 DecisionOutcome, DecisionPolicy, DecisionPolicyIdentity, DecisionPolicyKind, DecisionRule,
45 DecisionState, DecisionTicket, DecisionTicketDraft, DecisionTicketState, DecisionTrace,
46 ExternalDecisionOption, ExternalDecisionRequest, ExternalDecisionResponse, ExternalPolicy,
47 HumanDecisionResponse, HumanPolicy, LlmModelIdentity, LlmPolicy, OrderedRulePolicy,
48 PolicyDecision, QueuedExternalPolicy, QueuedHumanPolicy, QueuedLlmPolicy, RuleChoice,
49 RulePolicy, UtilityEvaluator, UtilityPolicy, UtilityProfile, WeightedUtilityEvaluator,
50 WeightedUtilityPolicy,
51};
52pub use decision::{DecisionEvaluation, DecisionIngressRequest, PreparedDecisionIngress};
53pub use ingress::{
54 IngressClass, IngressPayload, IngressReceipt, IngressRecord, PluginIngressDescriptor,
55 PluginIngressRequest,
56};
57pub use knowledge::{
58 KnowledgeLimitsV1, KnowledgeSubjectSchema, KnowledgeSubjectTargetKind, PluginKnowledgeSchema,
59};
60pub use manifest::{ArtifactManifest, RUN_MANIFEST_FORMAT_VERSION, RunManifest};
61pub use persistence::{
62 ArchiveProvider, ArchiveStore, ArchiveStoreOutcome, ArchivedEvidenceLocator,
63 ArchivedEvidenceReceipt, ArchivedSegmentHeader, CHECKPOINT_JOURNAL_FORMAT_VERSION,
64 CheckpointJournal, CompactedSimulation, EvidenceArchiveIndex, EvidenceCursor,
65 EvidenceDependency, EvidenceIndexEntry, EvidenceItemLocator, EvidenceJournalKind,
66 EvidenceJournalRoots, EvidenceJournalSegment, EvidenceNestedLocator, EvidenceRequirement,
67 EvidenceSealToken, PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD,
68 PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION, PayloadRequiredEvidenceContinuationV1,
69 PreparedEvidenceSeal, ReplayJournal, SimulationCheckpoint, SimulationSnapshot,
70 payload_required_evidence_continuation_property_v1,
71};
72pub use policy::{
73 CommandPolicyContext, ControllerPolicy, InteractionPolicy, ObservationPolicy,
74 RUN_CONFIGURATION_FORMAT_VERSION, RunConfiguration, RunConfigurationSnapshot, RunPurpose,
75 SeatBinding, SeatPolicy, TracePolicy,
76};
77pub use random::{
78 KeyedDrawReservation, RandomAlgorithm, RandomDrawAddress, RandomDrawOutcome,
79 RandomDrawProducer, RandomDrawRecord, RandomOperationAddressV1, RandomOperationTarget,
80 RandomStreamKey, RandomStreamState,
81};
82pub use records::{
83 DomainRecord, DomainRecordChange, DomainRecordClass, DomainRecordDraft, DomainRecordLifecycle,
84 DomainRecordMutation, DomainRecordMutationPolicy, DomainRecordOperation, DomainRecordSchema,
85 DomainReference, DomainReferenceSchema, DomainReferenceTarget, DomainReferenceTargetKind,
86};
87
88use canwu_core::{
89 ArmyId, BoundaryId, CommandAttemptId, CommandId, CommandRequestId, DecisionRequestId,
90 DecisionTicketId, DecisionTraceId, DeterministicRng, DomainRecordKind, DomainRecordRef,
91 DomainRecordType, EntityRef, EventId, FieldSchema, GovernmentId, IngressId, LetterId, PersonId,
92 RandomDrawId, ResourceId, RouteId, SchemaRegistry, TerritoryId, TypeSchema,
93 TypedDomainRecordRef,
94};
95pub use canwu_event::{CauseRef, EventAudience, EventKind, SimEvent};
96pub use canwu_knowledge::{
97 ActorKnowledge, ArmyKnowledge, EstimateRange, KnowledgeCursor, KnowledgeHistoryView,
98 KnowledgeOrigin, KnowledgeQuery, KnowledgeReadCut, KnowledgeRecord, KnowledgeRecordDraft,
99 KnowledgeRecordView, KnowledgeSnapshot, KnowledgeSource, KnowledgeSubject,
100 KnowledgeSubjectTarget,
101};
102use canwu_time::{SimDuration, SimTime};
103use canwu_world::{
104 Army, Government, LetterCargo, LetterStatus, MapPoint, Person, PersonTransitState, Route,
105 Territory, TransitState, WorldSnapshot,
106};
107use serde::{Deserialize, Serialize};
108use serde_json::Value;
109use std::cell::RefCell;
110use std::collections::{BTreeMap, BTreeSet, HashSet};
111use std::panic::{AssertUnwindSafe, catch_unwind};
112
113use hashing::{
114 ControlCommitmentMaterial, StateHashMaterial, authoritative_run_identity,
115 boundary_state_hash_for_commitments, checkpoint_hash_for_commitments,
116 checkpoint_hash_for_configuration, commitment_roots_are_canonical, compute_boundary_hash,
117 decision_commitment_root, domain_record_commitment_root, identity_commitment_root,
118 is_canonical_hash, knowledge_commitment_root, plugin_component_commitment_root,
119 random_stream_commitment_root, runtime_commitment_roots, scheduler_commitment_root,
120 snapshot_boundary_head_state_hash, snapshot_checkpoint_hash, snapshot_commitment_roots,
121 snapshot_is_at_boundary_head, snapshot_state_hash, state_hash, world_commitment_root,
122};
123use ingress::IngressQueueKey;
124use migration::{
125 PersistedAdmissionCursors, authoritative_revision_count, boundaries_before_attempts,
126 inferred_run_configuration, migrate_snapshot,
127};
128use settlement::{PendingBoundaryRandomDraw, boundary_has_event_ingress, boundary_system_due};
129use state::{
130 CommitmentDomains, JournalCommitmentRoots, RuntimeCommitmentCache,
131 RuntimeCommitmentRootUpdates, RuntimeCounters, RuntimeCurrentState,
132 RuntimeDomainCommitmentRoots, RuntimeEvidence, RuntimeMetadata, RuntimeScheduler, RuntimeState,
133};
134use transactions::{
135 BoundaryTransactionCheckpoint, ClockTransactionCheckpoint, CommandTransactionCheckpoint,
136 IngressTransactionCheckpoint, RejectionTransactionCheckpoint,
137 ScheduledBatchTransactionCheckpoint,
138};
139use validation::{
140 RuntimeValidationContext, claim_counter, core_world_entity_exists,
141 has_unqueued_command_history, proposal_entity_exists, proposal_entity_identity_exists,
142 runtime_current_entity_exists, runtime_entity_exists,
143 runtime_entity_exists_with_record_overlay, runtime_entity_identity_exists,
144 runtime_has_unqueued_command_history, snapshot_entity_exists_in_history,
145 validate_directives_with_context, validate_domain_dependents_with_records,
146 validate_run_configuration_entities, validate_runtime_cause,
147 validate_runtime_domain_dependents, validate_snapshot,
148};
149
150pub const ENGINE_VERSION: &str = env!("CARGO_PKG_VERSION");
151pub const SNAPSHOT_FORMAT_VERSION: u32 = 5;
152pub const STATE_REVISION_FORMAT_VERSION: u32 = 1;
154pub const ADMISSION_CURSOR_FORMAT_VERSION: u32 = 1;
156pub const COMMITMENT_FORMAT_VERSION: u32 = 1;
158pub const MAX_SYNCHRONOUS_REACTION_DEPTH: usize = 32;
163const CORE_STATE_NAMESPACE: &str = "canwu.core";
164const GENESIS_BOUNDARY_HASH: &str =
165 "0000000000000000000000000000000000000000000000000000000000000000";
166
167pub use error::{CanwuError, ErrorCode};
168pub use hashing::CommitmentRoots;
169
170use ingress::CommandAdmission;
171pub use ingress::{
172 Command, CommandAttemptOutcome, CommandAttemptRecord, CommandAuthority, CommandContext,
173 CommandEnvelope, CommandIngress, CommandOutcome, CommandReceipt, CommandRecord,
174 CommandRejection, CommandRequest, DecisionOrigin, Issuer,
175};
176
177#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
178#[repr(u8)]
179#[serde(rename_all = "snake_case")]
180pub enum BoundaryPhase {
181 EventIngress = 1,
182 BoundarySnapshot = 2,
183 DerivedFieldSolve = 3,
184 PerceptionAndAttentionRefresh = 4,
185 DecisionAndAcceptedEffectIntake = 5,
186 ReservationAndAllocation = 6,
187 DomainDeltaProposal = 7,
188 InvariantValidation = 8,
189 AtomicDomainCommit = 9,
190 HistoricalCandidateEvaluation = 10,
191 ConditionalTransitionCommit = 11,
192 StrategicAggregation = 12,
193 PerspectiveAndReportMaterialization = 13,
194 SaveReplayAndDiagnosticHashing = 14,
195}
196
197impl BoundaryPhase {
198 pub const ALL: [Self; 14] = [
199 Self::EventIngress,
200 Self::BoundarySnapshot,
201 Self::DerivedFieldSolve,
202 Self::PerceptionAndAttentionRefresh,
203 Self::DecisionAndAcceptedEffectIntake,
204 Self::ReservationAndAllocation,
205 Self::DomainDeltaProposal,
206 Self::InvariantValidation,
207 Self::AtomicDomainCommit,
208 Self::HistoricalCandidateEvaluation,
209 Self::ConditionalTransitionCommit,
210 Self::StrategicAggregation,
211 Self::PerspectiveAndReportMaterialization,
212 Self::SaveReplayAndDiagnosticHashing,
213 ];
214}
215
216#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
217#[serde(rename_all = "snake_case")]
218pub enum SystemCadence {
219 EventDriven,
220 SubDaily,
221 Daily,
222 Monthly,
223 Seasonal,
224 Annual,
225 EraScheduled,
226}
227
228#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
229#[serde(rename_all = "snake_case")]
230pub enum StateVisibility {
231 SameBoundary,
232 NextBoundary,
233}
234
235#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
236pub struct StateKey {
237 pub namespace: String,
238 pub name: String,
239}
240
241impl StateKey {
242 #[must_use]
243 pub fn new(namespace: impl Into<String>, name: impl Into<String>) -> Self {
244 Self {
245 namespace: namespace.into(),
246 name: name.into(),
247 }
248 }
249
250 #[must_use]
251 pub fn core_people() -> Self {
252 Self::new(CORE_STATE_NAMESPACE, "people")
253 }
254
255 #[must_use]
256 pub fn core_governments() -> Self {
257 Self::new(CORE_STATE_NAMESPACE, "governments")
258 }
259
260 #[must_use]
261 pub fn core_territories() -> Self {
262 Self::new(CORE_STATE_NAMESPACE, "territories")
263 }
264
265 #[must_use]
266 pub fn core_routes() -> Self {
267 Self::new(CORE_STATE_NAMESPACE, "routes")
268 }
269
270 #[must_use]
271 pub fn core_armies() -> Self {
272 Self::new(CORE_STATE_NAMESPACE, "armies")
273 }
274
275 #[must_use]
276 pub fn core_knowledge() -> Self {
277 Self::new(CORE_STATE_NAMESPACE, "knowledge")
278 }
279
280 #[must_use]
281 pub fn core_commands() -> Self {
282 Self::new(CORE_STATE_NAMESPACE, "commands")
283 }
284
285 #[must_use]
286 pub fn core_events() -> Self {
287 Self::new(CORE_STATE_NAMESPACE, "events")
288 }
289
290 #[must_use]
291 pub fn core_ingress() -> Self {
292 Self::new(CORE_STATE_NAMESPACE, "ingress")
293 }
294}
295
296#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
297pub struct SystemContract {
298 pub name: String,
299 pub phase: BoundaryPhase,
300 pub cadence: SystemCadence,
301 pub reads: Vec<StateKey>,
302 pub writes: Vec<StateKey>,
303 pub visibility: StateVisibility,
304}
305
306impl SystemContract {
307 #[must_use]
308 pub fn event_driven(name: impl Into<String>, phase: BoundaryPhase) -> Self {
309 Self {
310 name: name.into(),
311 phase,
312 cadence: SystemCadence::EventDriven,
313 reads: Vec::new(),
314 writes: Vec::new(),
315 visibility: StateVisibility::SameBoundary,
316 }
317 }
318}
319
320pub use scenario::{DemoIds, Scenario, demo_scenario};
321use scenario::{
322 base_schema, canonicalize_scenario, require_plugin_aware_initial_records, validate_scenario,
323 validate_scenario_state, validate_strict_id_order,
324};
325
326use plugins::PluginComponentKey;
327pub use plugins::{
328 PayloadProperty, PayloadSchema, PayloadValueType, PluginActionDescriptor, PluginCommandHandler,
329 PluginComponentRecord, PluginDescriptor, PluginRegistrar, PluginRegistry, SimulationPlugin,
330 SimulationSystemHandler, SystemDirective,
331};
332
333pub use view::SimulationView;
334use view::SimulationViewState;
335
336#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
337enum BoundaryWriteStage {
338 Ordinary,
339 Transition,
340 Aggregation,
341 Perspective,
342}
343
344#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
345enum DomainRecordCommitStage {
346 Ordinary,
347 Transition,
348 Aggregation,
349 Perspective,
350 Deferred,
351}
352
353impl DomainRecordCommitStage {
354 const ALL: [Self; 5] = [
355 Self::Ordinary,
356 Self::Transition,
357 Self::Aggregation,
358 Self::Perspective,
359 Self::Deferred,
360 ];
361
362 const fn ordinal(self) -> u8 {
363 match self {
364 Self::Ordinary => 1,
365 Self::Transition => 2,
366 Self::Aggregation => 3,
367 Self::Perspective => 4,
368 Self::Deferred => 5,
369 }
370 }
371}
372
373#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
374struct DomainHistoryCut {
375 boundary: usize,
376 stage: u8,
377}
378
379impl DomainHistoryCut {
380 const GENESIS: Self = Self {
381 boundary: 0,
382 stage: 0,
383 };
384
385 const fn after_boundaries(boundary: usize) -> Self {
386 Self { boundary, stage: 5 }
387 }
388
389 const fn after_stage(boundary: usize, stage: DomainRecordCommitStage) -> Self {
390 Self {
391 boundary,
392 stage: stage.ordinal(),
393 }
394 }
395}
396
397#[derive(Clone, Debug, Default)]
398struct BoundaryDomainEntityCuts {
399 changes: BTreeMap<DomainRecordRef, Vec<DomainEntityStageChange>>,
400}
401
402impl BoundaryDomainEntityCuts {
403 fn record(&mut self, stage: DomainRecordCommitStage, change: &DomainRecordChange) {
404 let previous_live = change
405 .previous
406 .as_ref()
407 .is_some_and(domain_record_is_live_entity);
408 let current_live = domain_record_is_live_entity(&change.current);
409 if previous_live != current_live {
410 self.changes
411 .entry(change.current.reference.clone())
412 .or_default()
413 .push(DomainEntityStageChange {
414 stage,
415 plugin: change.plugin.clone(),
416 system: change.system.clone(),
417 previous_live,
418 current_live,
419 });
420 }
421 }
422
423 fn is_live(
424 &self,
425 final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
426 reference: &DomainRecordRef,
427 stage: Option<DomainRecordCommitStage>,
428 ) -> bool {
429 let mut live = final_records
430 .get(reference)
431 .is_some_and(domain_record_is_live_entity);
432 if let Some(changes) = self.changes.get(reference) {
433 for change in changes.iter().rev() {
434 if stage.is_some_and(|stage| change.stage <= stage) {
435 break;
436 }
437 live = change.previous_live;
438 }
439 }
440 live
441 }
442
443 fn is_live_for_proposal(
444 &self,
445 final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
446 reference: &DomainRecordRef,
447 phase: BoundaryPhase,
448 commit_stage: DomainRecordCommitStage,
449 plugin: &str,
450 system: &str,
451 ) -> bool {
452 let visible_after = match phase {
453 BoundaryPhase::DomainDeltaProposal => None,
454 BoundaryPhase::HistoricalCandidateEvaluation => Some(DomainRecordCommitStage::Ordinary),
455 BoundaryPhase::StrategicAggregation => Some(DomainRecordCommitStage::Transition),
456 BoundaryPhase::PerspectiveAndReportMaterialization => {
457 Some(DomainRecordCommitStage::Aggregation)
458 }
459 BoundaryPhase::EventIngress
460 | BoundaryPhase::BoundarySnapshot
461 | BoundaryPhase::DerivedFieldSolve
462 | BoundaryPhase::PerceptionAndAttentionRefresh
463 | BoundaryPhase::DecisionAndAcceptedEffectIntake
464 | BoundaryPhase::ReservationAndAllocation
465 | BoundaryPhase::InvariantValidation
466 | BoundaryPhase::AtomicDomainCommit
467 | BoundaryPhase::ConditionalTransitionCommit
468 | BoundaryPhase::SaveReplayAndDiagnosticHashing => return false,
469 };
470 let before_proposal = self.is_live(final_records, reference, visible_after);
471 self.changes
472 .get(reference)
473 .and_then(|changes| {
474 changes.iter().find(|change| {
475 change.stage == commit_stage
476 && change.plugin == plugin
477 && change.system == system
478 })
479 })
480 .map_or(before_proposal, |change| change.current_live)
481 }
482
483 fn identity_exists_for_proposal(
484 &self,
485 final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
486 reference: &DomainRecordRef,
487 phase: BoundaryPhase,
488 commit_stage: DomainRecordCommitStage,
489 plugin: &str,
490 system: &str,
491 ) -> bool {
492 if !final_records.contains_key(reference) {
493 return false;
494 }
495 let visible_after = match phase {
496 BoundaryPhase::DomainDeltaProposal => None,
497 BoundaryPhase::HistoricalCandidateEvaluation => Some(DomainRecordCommitStage::Ordinary),
498 BoundaryPhase::StrategicAggregation => Some(DomainRecordCommitStage::Transition),
499 BoundaryPhase::PerspectiveAndReportMaterialization => {
500 Some(DomainRecordCommitStage::Aggregation)
501 }
502 BoundaryPhase::EventIngress
503 | BoundaryPhase::BoundarySnapshot
504 | BoundaryPhase::DerivedFieldSolve
505 | BoundaryPhase::PerceptionAndAttentionRefresh
506 | BoundaryPhase::DecisionAndAcceptedEffectIntake
507 | BoundaryPhase::ReservationAndAllocation
508 | BoundaryPhase::InvariantValidation
509 | BoundaryPhase::AtomicDomainCommit
510 | BoundaryPhase::ConditionalTransitionCommit
511 | BoundaryPhase::SaveReplayAndDiagnosticHashing => return false,
512 };
513 self.changes
514 .get(reference)
515 .and_then(|changes| {
516 changes
517 .iter()
518 .find(|change| !change.previous_live && change.current_live)
519 })
520 .is_none_or(|creation| {
521 visible_after.is_some_and(|stage| creation.stage <= stage)
522 || (creation.stage == commit_stage
523 && creation.plugin == plugin
524 && creation.system == system)
525 })
526 }
527}
528
529#[derive(Clone, Debug)]
530struct DomainEntityStageChange {
531 stage: DomainRecordCommitStage,
532 plugin: String,
533 system: String,
534 previous_live: bool,
535 current_live: bool,
536}
537
538#[derive(Clone, Debug)]
539struct DomainRecordHistory {
540 lifetimes: BTreeMap<DomainRecordRef, DomainEntityLifetime>,
541}
542
543impl DomainRecordHistory {
544 fn from_initial_records(records: &BTreeMap<DomainRecordRef, DomainRecord>) -> Self {
545 let lifetimes = records
546 .values()
547 .filter(|record| record.class == DomainRecordClass::Entity)
548 .map(|record| {
549 (
550 record.reference.clone(),
551 DomainEntityLifetime {
552 created_at: DomainHistoryCut::GENESIS,
553 deleted_at: record.is_deleted().then_some(DomainHistoryCut::GENESIS),
554 },
555 )
556 })
557 .collect();
558 Self { lifetimes }
559 }
560
561 fn apply_boundary(
562 &mut self,
563 boundary: usize,
564 cuts: &BoundaryDomainEntityCuts,
565 ) -> Result<(), CanwuError> {
566 for (reference, changes) in &cuts.changes {
567 for change in changes {
568 let cut = DomainHistoryCut::after_stage(boundary, change.stage);
569 match (change.previous_live, change.current_live) {
570 (false, true) => {
571 if self
572 .lifetimes
573 .insert(
574 reference.clone(),
575 DomainEntityLifetime {
576 created_at: cut,
577 deleted_at: None,
578 },
579 )
580 .is_some()
581 {
582 return invalid_snapshot(
583 "domain entity history recreates an existing stable identity",
584 );
585 }
586 }
587 (true, false) => {
588 let Some(lifetime) = self.lifetimes.get_mut(reference) else {
589 return invalid_snapshot(
590 "domain entity history deletes an identity before creation",
591 );
592 };
593 if lifetime.deleted_at.replace(cut).is_some() {
594 return invalid_snapshot(
595 "domain entity history deletes the same identity more than once",
596 );
597 }
598 }
599 (false, false) | (true, true) => {}
600 }
601 }
602 }
603 Ok(())
604 }
605
606 fn is_live(&self, reference: &DomainRecordRef, cut: DomainHistoryCut) -> bool {
607 self.lifetimes.get(reference).is_some_and(|lifetime| {
608 lifetime.created_at <= cut && lifetime.deleted_at.is_none_or(|deleted| cut < deleted)
609 })
610 }
611
612 fn exists(&self, reference: &DomainRecordRef, cut: DomainHistoryCut) -> bool {
613 self.lifetimes
614 .get(reference)
615 .is_some_and(|lifetime| lifetime.created_at <= cut)
616 }
617
618 fn before_time(snapshot: &SimulationSnapshot, at: SimTime) -> DomainHistoryCut {
619 let count = snapshot
620 .boundaries
621 .partition_point(|boundary| boundary.at < at);
622 DomainHistoryCut::after_boundaries(count)
623 }
624}
625
626#[derive(Clone, Copy, Debug)]
627struct DomainEntityLifetime {
628 created_at: DomainHistoryCut,
629 deleted_at: Option<DomainHistoryCut>,
630}
631
632fn domain_record_is_live_entity(record: &DomainRecord) -> bool {
633 record.class == DomainRecordClass::Entity && !record.is_deleted()
634}
635
636const fn boundary_write_stage(phase: BoundaryPhase) -> Option<BoundaryWriteStage> {
637 match phase {
638 BoundaryPhase::DomainDeltaProposal => Some(BoundaryWriteStage::Ordinary),
639 BoundaryPhase::HistoricalCandidateEvaluation => Some(BoundaryWriteStage::Transition),
640 BoundaryPhase::StrategicAggregation => Some(BoundaryWriteStage::Aggregation),
641 BoundaryPhase::PerspectiveAndReportMaterialization => Some(BoundaryWriteStage::Perspective),
642 BoundaryPhase::EventIngress
643 | BoundaryPhase::BoundarySnapshot
644 | BoundaryPhase::DerivedFieldSolve
645 | BoundaryPhase::PerceptionAndAttentionRefresh
646 | BoundaryPhase::DecisionAndAcceptedEffectIntake
647 | BoundaryPhase::ReservationAndAllocation
648 | BoundaryPhase::InvariantValidation
649 | BoundaryPhase::AtomicDomainCommit
650 | BoundaryPhase::ConditionalTransitionCommit
651 | BoundaryPhase::SaveReplayAndDiagnosticHashing => None,
652 }
653}
654
655const fn domain_record_commit_stage(
656 phase: BoundaryPhase,
657 visibility: StateVisibility,
658) -> Option<DomainRecordCommitStage> {
659 let stage = match phase {
660 BoundaryPhase::DomainDeltaProposal => DomainRecordCommitStage::Ordinary,
661 BoundaryPhase::HistoricalCandidateEvaluation => DomainRecordCommitStage::Transition,
662 BoundaryPhase::StrategicAggregation => DomainRecordCommitStage::Aggregation,
663 BoundaryPhase::PerspectiveAndReportMaterialization => DomainRecordCommitStage::Perspective,
664 BoundaryPhase::EventIngress
665 | BoundaryPhase::BoundarySnapshot
666 | BoundaryPhase::DerivedFieldSolve
667 | BoundaryPhase::PerceptionAndAttentionRefresh
668 | BoundaryPhase::DecisionAndAcceptedEffectIntake
669 | BoundaryPhase::ReservationAndAllocation
670 | BoundaryPhase::InvariantValidation
671 | BoundaryPhase::AtomicDomainCommit
672 | BoundaryPhase::ConditionalTransitionCommit
673 | BoundaryPhase::SaveReplayAndDiagnosticHashing => return None,
674 };
675 Some(match visibility {
676 StateVisibility::SameBoundary => stage,
677 StateVisibility::NextBoundary => DomainRecordCommitStage::Deferred,
678 })
679}
680
681fn validate_type_schema(schema: &TypeSchema) -> Result<(), CanwuError> {
682 if schema.type_name.trim().is_empty() || schema.type_name != schema.type_name.trim() {
683 return Err(CanwuError::new(
684 ErrorCode::InvalidPluginRegistration,
685 "plugin schema type name must be non-empty and have no surrounding whitespace",
686 ));
687 }
688 let mut field_names = BTreeSet::new();
689 for field in &schema.fields {
690 if field.name.trim().is_empty()
691 || field.name != field.name.trim()
692 || field.value_type.trim().is_empty()
693 || field.value_type != field.value_type.trim()
694 || field
695 .reference_type
696 .as_ref()
697 .is_some_and(|value| value.trim().is_empty() || value != value.trim())
698 || !field_names.insert(&field.name)
699 {
700 return Err(CanwuError::new(
701 ErrorCode::InvalidPluginRegistration,
702 format!("schema {} contains an invalid field", schema.type_name),
703 ));
704 }
705 }
706 Ok(())
707}
708
709use scheduling::{ScheduleKey, ScheduledAction, ScheduledRecord};
710
711const fn one_u64() -> u64 {
712 1
713}
714
715#[allow(clippy::trivially_copy_pass_by_ref)]
716const fn is_zero_u64(value: &u64) -> bool {
717 *value == 0
718}
719
720#[allow(clippy::trivially_copy_pass_by_ref)]
721const fn is_zero_u32(value: &u32) -> bool {
722 *value == 0
723}
724
725#[allow(clippy::trivially_copy_pass_by_ref)]
726const fn is_one_u64(value: &u64) -> bool {
727 *value == 1
728}
729
730fn command_attempt_slice_is_empty(value: &&[CommandAttemptRecord]) -> bool {
731 value.is_empty()
732}
733
734fn command_attempt_id_slice_is_empty(value: &&[CommandAttemptId]) -> bool {
735 value.is_empty()
736}
737
738fn domain_record_slice_is_empty(value: &&[DomainRecord]) -> bool {
739 value.is_empty()
740}
741
742fn domain_record_change_slice_is_empty(value: &&[DomainRecordChange]) -> bool {
743 value.is_empty()
744}
745
746fn ingress_record_slice_is_empty(value: &&[IngressRecord]) -> bool {
747 value.is_empty()
748}
749
750const BOUNDARY_STATE_HASH_V1_PREFIX: &str = "v1:";
751
752#[derive(Clone, Copy, Debug, Eq, PartialEq)]
753enum BoundaryStateHashFormat {
754 LegacyV0,
755 CommitmentsV1,
756}
757
758fn boundary_state_hash_format(value: Option<&str>) -> Result<BoundaryStateHashFormat, CanwuError> {
759 match value {
760 Some(value) if value.starts_with(BOUNDARY_STATE_HASH_V1_PREFIX) => {
761 let hash = &value[BOUNDARY_STATE_HASH_V1_PREFIX.len()..];
762 if !is_canonical_hash(hash) {
763 return invalid_snapshot("boundary state commitment v1 is not canonical");
764 }
765 Ok(BoundaryStateHashFormat::CommitmentsV1)
766 }
767 Some(value) if is_canonical_hash(value) => Ok(BoundaryStateHashFormat::LegacyV0),
768 Some(_) => invalid_snapshot("boundary state commitment format is unsupported"),
769 None => Ok(BoundaryStateHashFormat::LegacyV0),
770 }
771}
772
773pub struct Simulation {
774 state: RuntimeState,
775 schema: SchemaRegistry,
776 plugins: PluginRegistry,
777 sync_reaction_depth: usize,
778}
779
780impl Simulation {
781 pub fn new(seed: u64, scenario: Scenario) -> Result<Self, CanwuError> {
783 require_plugin_aware_initial_records(&scenario)?;
784 let run_manifest = RunManifest::for_scenario("canwu.inline", "scenario", "1", &scenario)?;
785 Self::new_with_configuration_snapshot(
786 seed,
787 scenario,
788 run_manifest,
789 RunConfigurationSnapshot::CompatibilityV1,
790 )
791 }
792
793 pub fn new_with_plugins(
796 seed: u64,
797 scenario: Scenario,
798 plugins: &[&dyn SimulationPlugin],
799 ) -> Result<Self, CanwuError> {
800 let run_manifest = RunManifest::for_scenario("canwu.inline", "scenario", "1", &scenario)?;
801 Self::new_with_manifest_and_plugins(seed, scenario, run_manifest, plugins)
802 }
803
804 pub fn new_with_manifest(
806 seed: u64,
807 scenario: Scenario,
808 run_manifest: RunManifest,
809 ) -> Result<Self, CanwuError> {
810 require_plugin_aware_initial_records(&scenario)?;
811 Self::new_with_configuration_snapshot(
812 seed,
813 scenario,
814 run_manifest,
815 RunConfigurationSnapshot::CompatibilityV1,
816 )
817 }
818
819 pub fn new_with_manifest_and_plugins(
822 seed: u64,
823 scenario: Scenario,
824 run_manifest: RunManifest,
825 plugins: &[&dyn SimulationPlugin],
826 ) -> Result<Self, CanwuError> {
827 let simulation = Self::new_with_configuration_snapshot(
828 seed,
829 scenario,
830 run_manifest,
831 RunConfigurationSnapshot::CompatibilityV1,
832 )?;
833 Self::activate_initial_plugins(simulation, plugins)
834 }
835
836 pub fn new_with_run_configuration(
839 seed: u64,
840 scenario: Scenario,
841 run_manifest: RunManifest,
842 mut run_configuration: RunConfiguration,
843 ) -> Result<Self, CanwuError> {
844 require_plugin_aware_initial_records(&scenario)?;
845 run_configuration.canonicalize();
846 Self::new_with_configuration_snapshot(
847 seed,
848 scenario,
849 run_manifest,
850 RunConfigurationSnapshot::Declared(run_configuration),
851 )
852 }
853
854 pub fn new_with_run_configuration_and_plugins(
857 seed: u64,
858 scenario: Scenario,
859 run_manifest: RunManifest,
860 mut run_configuration: RunConfiguration,
861 plugins: &[&dyn SimulationPlugin],
862 ) -> Result<Self, CanwuError> {
863 run_configuration.canonicalize();
864 let simulation = Self::new_with_configuration_snapshot(
865 seed,
866 scenario,
867 run_manifest,
868 RunConfigurationSnapshot::Declared(run_configuration),
869 )?;
870 Self::activate_initial_plugins(simulation, plugins)
871 }
872
873 fn activate_initial_plugins(
874 mut simulation: Self,
875 plugins: &[&dyn SimulationPlugin],
876 ) -> Result<Self, CanwuError> {
877 for plugin in plugins {
878 simulation.register_plugin(*plugin)?;
879 }
880 simulation.ensure_runtime_ready()?;
881 Ok(simulation)
882 }
883
884 fn new_with_configuration_snapshot(
885 seed: u64,
886 mut scenario: Scenario,
887 mut run_manifest: RunManifest,
888 run_configuration: RunConfigurationSnapshot,
889 ) -> Result<Self, CanwuError> {
890 canonicalize_scenario(&mut scenario);
891 validate_scenario(&scenario)?;
892 manifest::canonicalize(&mut run_manifest);
893 manifest::validate(&run_manifest, Some(&scenario), false)?;
894 manifest::validate_run_configuration(&run_manifest, &run_configuration)?;
895 validate_run_configuration_entities(
896 &run_configuration,
897 &scenario.world,
898 &scenario.domain_records,
899 )?;
900 let run_manifest_hash = manifest::hash(&run_manifest)?;
901 if scenario
902 .world
903 .armies
904 .iter()
905 .any(|army| army.transit.is_some())
906 {
907 return Err(CanwuError::new(
908 ErrorCode::InvalidSnapshot,
909 "initial scenarios cannot contain transit without admitted command/event/queue evidence",
910 ));
911 }
912 if scenario
913 .world
914 .people
915 .iter()
916 .any(|person| person.transit.is_some())
917 || scenario
918 .world
919 .letters
920 .iter()
921 .any(|letter| letter.status == LetterStatus::InTransit)
922 {
923 return Err(CanwuError::new(
924 ErrorCode::InvalidSnapshot,
925 "initial scenarios cannot contain person or letter transit without admitted command/event/queue evidence",
926 ));
927 }
928 let schema = base_schema();
929 let plugins = PluginRegistry::default();
930 let core_stream = RandomStreamState::initial(seed, random::core_report_delay_stream());
931 let initial_scenario = Some(scenario.clone());
932 let mut simulation = Self {
933 state: RuntimeState {
934 current: RuntimeCurrentState {
935 people: scenario
936 .world
937 .people
938 .into_iter()
939 .map(|value| (value.id, value))
940 .collect(),
941 letters: scenario
942 .world
943 .letters
944 .into_iter()
945 .map(|value| (value.id, value))
946 .collect(),
947 governments: scenario
948 .world
949 .governments
950 .into_iter()
951 .map(|value| (value.id, value))
952 .collect(),
953 territories: scenario
954 .world
955 .territories
956 .into_iter()
957 .map(|value| (value.id, value))
958 .collect(),
959 routes: scenario
960 .world
961 .routes
962 .into_iter()
963 .map(|value| (value.id, value))
964 .collect(),
965 armies: scenario
966 .world
967 .armies
968 .into_iter()
969 .map(|value| (value.id, value))
970 .collect(),
971 knowledge: scenario.knowledge,
972 plugin_components: BTreeMap::new(),
973 domain_records: scenario
974 .domain_records
975 .into_iter()
976 .map(|record| (record.reference.clone(), record))
977 .collect(),
978 decisions: DecisionState::default(),
979 root_seed: seed,
980 random_streams: BTreeMap::from([(core_stream.key.clone(), core_stream)]),
981 },
982 scheduler: RuntimeScheduler {
983 initial_time: scenario.start_time,
984 now: scenario.start_time,
985 actions: BTreeMap::new(),
986 pending_ingress: BTreeSet::new(),
987 },
988 counters: RuntimeCounters {
989 next_event_id: 1,
990 next_command_id: 1,
991 next_command_attempt_id: 1,
992 next_ingress_id: 1,
993 next_boundary_id: 1,
994 next_random_draw_id: 1,
995 next_knowledge_record_id: 1,
996 next_schedule_sequence: 1,
997 next_correlation_id: 1,
998 next_decision_trace_id: 1,
999 state_revision: 0,
1000 admitted_attempt_count: 0,
1001 admitted_command_count: 0,
1002 admitted_event_count: 0,
1003 },
1004 metadata: RuntimeMetadata {
1005 initial_scenario,
1006 run_manifest,
1007 run_manifest_hash,
1008 run_configuration,
1009 checkpoint_hash: String::new(),
1010 commitment_format_version: COMMITMENT_FORMAT_VERSION,
1011 commitment_roots: None,
1012 commitment_cache: None,
1013 plugin_registration_closed: false,
1014 replay_revision_format_version: STATE_REVISION_FORMAT_VERSION,
1015 },
1016 evidence: RuntimeEvidence {
1017 archived: EvidenceCursor::default(),
1018 archived_boundary_head: None,
1019 archived_legacy_commands: false,
1020 archived_tracked_attempts: false,
1021 archived_unqueued_command_history: false,
1022 archived_command_requests: BTreeMap::new(),
1023 archived_ingress_requests: BTreeMap::new(),
1024 archived_decision_requests: BTreeMap::new(),
1025 archived_decision_command_requests: BTreeSet::new(),
1026 events: Vec::new(),
1027 commands: Vec::new(),
1028 command_attempts: Vec::new(),
1029 ingress: Vec::new(),
1030 boundaries: Vec::new(),
1031 random_draws: Vec::new(),
1032 archived_segment_headers: Vec::new(),
1033 archived_evidence_receipts: BTreeMap::new(),
1034 keyed_draw_reservations: Vec::new(),
1035 },
1036 },
1037 schema,
1038 plugins,
1039 sync_reaction_depth: 0,
1040 };
1041 simulation.refresh_checkpoint_hash()?;
1042 Ok(simulation)
1043 }
1044
1045 pub fn demo(seed: u64) -> Result<(Self, DemoIds), CanwuError> {
1046 let (scenario, ids) = demo_scenario();
1047 Self::new(seed, scenario).map(|simulation| (simulation, ids))
1048 }
1049
1050 pub fn register_plugin<P: SimulationPlugin + ?Sized>(
1051 &mut self,
1052 plugin: &P,
1053 ) -> Result<(), CanwuError> {
1054 let plugin_name = plugin.name().trim();
1055 if plugin_name.is_empty() || plugin_name != plugin.name() {
1056 return Err(CanwuError::new(
1057 ErrorCode::InvalidPluginRegistration,
1058 "plugin name must be non-empty and have no surrounding whitespace",
1059 ));
1060 }
1061 let rehydrating = self.plugins.descriptors.contains_key(plugin_name)
1062 && !self.plugins.active_plugins.contains(plugin_name);
1063 if self.state.metadata.plugin_registration_closed && !rehydrating {
1064 return Err(CanwuError::new(
1065 ErrorCode::PluginRegistrationClosed,
1066 "new plugins must be registered before authoritative execution begins",
1067 ));
1068 }
1069 let state_start = self.state.clone();
1070 let schema_start = self.schema.clone();
1071 let plugins_start = self.plugins.clone();
1072 let result = (|| {
1073 self.plugins.register(plugin, &mut self.schema)?;
1074 self.invalidate_commitments(
1075 CommitmentDomains::RANDOM_STREAMS | CommitmentDomains::IDENTITY,
1076 );
1077 if !self.plugins.record_schemas.is_empty()
1078 && self.state.metadata.initial_scenario.is_none()
1079 {
1080 return Err(CanwuError::new(
1081 ErrorCode::UnsupportedSnapshotVersion,
1082 "this snapshot predates manifest-bound domain-record genesis and cannot activate record schemas",
1083 ));
1084 }
1085 records::validate_records_for_owner(
1086 &self.state.current.domain_records,
1087 &self.plugins.record_schemas,
1088 plugin_name,
1089 self.state.scheduler.now,
1090 &|entity| runtime_entity_exists(&self.state, entity),
1091 )?;
1092 for stream in self.plugins.random_stream_owners.keys() {
1093 self.state
1094 .current
1095 .random_streams
1096 .entry(stream.clone())
1097 .or_insert_with(|| {
1098 RandomStreamState::initial(self.state.current.root_seed, stream.clone())
1099 });
1100 }
1101 self.refresh_checkpoint_hash()
1102 })();
1103 if let Err(error) = result {
1104 self.state = state_start;
1105 self.schema = schema_start;
1106 self.plugins = plugins_start;
1107 return Err(error);
1108 }
1109 Ok(())
1110 }
1111
1112 fn ensure_runtime_ready(&self) -> Result<(), CanwuError> {
1113 self.plugins.ensure_active()?;
1114 records::validate_record_store(
1115 &self.state.current.domain_records,
1116 &self.plugins.record_schemas,
1117 self.state.scheduler.now,
1118 &|entity| runtime_entity_exists(&self.state, entity),
1119 )
1120 }
1121
1122 fn domain_record_feature_enabled(&self) -> bool {
1123 !self.plugins.record_schemas.is_empty()
1124 || !self.state.current.domain_records.is_empty()
1125 || self
1126 .state
1127 .evidence
1128 .boundaries
1129 .iter()
1130 .any(|boundary| !boundary.record_changes.is_empty())
1131 }
1132
1133 fn bound_initial_scenario(&self) -> Option<&Scenario> {
1134 if self.domain_record_feature_enabled() {
1135 self.state.metadata.initial_scenario.as_ref()
1136 } else {
1137 None
1138 }
1139 }
1140
1141 #[must_use]
1142 pub const fn time(&self) -> SimTime {
1143 self.state.scheduler.now
1144 }
1145
1146 #[must_use]
1147 pub const fn run_manifest(&self) -> &RunManifest {
1148 &self.state.metadata.run_manifest
1149 }
1150
1151 #[must_use]
1152 pub const fn run_configuration(&self) -> &RunConfigurationSnapshot {
1153 &self.state.metadata.run_configuration
1154 }
1155
1156 #[must_use]
1157 pub const fn revision(&self) -> u64 {
1165 self.state.counters.state_revision
1166 }
1167
1168 #[must_use]
1169 pub fn run_manifest_hash(&self) -> &str {
1170 &self.state.metadata.run_manifest_hash
1171 }
1172
1173 #[must_use]
1174 pub fn checkpoint_hash(&self) -> &str {
1175 &self.state.metadata.checkpoint_hash
1176 }
1177
1178 pub fn authoritative_state_hash(&self) -> Result<String, CanwuError> {
1182 self.compute_boundary_state_hash()
1183 }
1184
1185 #[must_use]
1186 pub fn world(&self) -> WorldSnapshot {
1187 WorldSnapshot {
1188 people: self.state.current.people.values().cloned().collect(),
1189 governments: self.state.current.governments.values().cloned().collect(),
1190 territories: self.state.current.territories.values().cloned().collect(),
1191 routes: self.state.current.routes.values().cloned().collect(),
1192 armies: self.state.current.armies.values().cloned().collect(),
1193 letters: self.state.current.letters.values().cloned().collect(),
1194 }
1195 }
1196
1197 #[must_use]
1198 pub fn knowledge(&self) -> &KnowledgeSnapshot {
1199 &self.state.current.knowledge
1200 }
1201
1202 #[must_use]
1203 pub fn events(&self) -> &[SimEvent] {
1204 &self.state.evidence.events
1205 }
1206
1207 #[must_use]
1208 pub fn command_log(&self) -> &[CommandRecord] {
1209 &self.state.evidence.commands
1210 }
1211
1212 #[must_use]
1213 pub fn command_attempts(&self) -> &[CommandAttemptRecord] {
1214 &self.state.evidence.command_attempts
1215 }
1216
1217 #[must_use]
1218 pub fn ingress_log(&self) -> &[IngressRecord] {
1219 &self.state.evidence.ingress
1220 }
1221
1222 #[must_use]
1223 pub fn domain_record(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
1224 self.state.current.domain_records.get(reference)
1225 }
1226
1227 #[must_use]
1228 pub fn typed_domain_record<T: DomainRecordType>(
1229 &self,
1230 reference: &TypedDomainRecordRef<T>,
1231 ) -> Option<&DomainRecord> {
1232 self.domain_record(reference.as_untyped())
1233 }
1234
1235 pub fn domain_records(&self) -> impl Iterator<Item = &DomainRecord> {
1236 self.state.current.domain_records.values()
1237 }
1238
1239 #[must_use]
1240 pub fn boundaries(&self) -> &[BoundaryRecord] {
1241 &self.state.evidence.boundaries
1242 }
1243
1244 #[must_use]
1245 pub fn random_draws(&self) -> &[RandomDrawRecord] {
1246 &self.state.evidence.random_draws
1247 }
1248
1249 #[must_use]
1250 pub fn boundary_head_hash(&self) -> Option<&str> {
1251 self.state.evidence.boundary_head_hash()
1252 }
1253
1254 #[must_use]
1255 pub const fn schema(&self) -> &SchemaRegistry {
1256 &self.schema
1257 }
1258
1259 pub fn plugin_descriptors(&self) -> impl Iterator<Item = &PluginDescriptor> {
1260 self.plugins.descriptors()
1261 }
1262
1263 #[must_use]
1269 pub fn event_audience(&self, event: &SimEvent) -> EventAudience {
1270 match &event.kind {
1271 EventKind::Plugin { plugin, event_type } => {
1272 self.plugins.event_audience(plugin, event_type)
1273 }
1274 EventKind::KnowledgePublished { holder, .. } => {
1275 EventAudience::KnowledgeHolder(holder.clone())
1276 }
1277 EventKind::MoveOrdered { .. }
1278 | EventKind::PersonMoveOrdered { .. }
1279 | EventKind::ArmyArrived { .. }
1280 | EventKind::PersonArrived { .. }
1281 | EventKind::LetterDelivered { .. }
1282 | EventKind::ReportDispatched { .. }
1283 | EventKind::KnowledgeUpdated { .. }
1284 | EventKind::DebugFieldChanged { .. } => EventAudience::Private,
1285 }
1286 }
1287
1288 #[must_use]
1289 pub fn replay_journal(&self) -> ReplayJournal {
1290 ReplayJournal {
1291 engine_version: ENGINE_VERSION.to_owned(),
1292 snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
1293 root_seed: self.state.current.root_seed,
1294 run_manifest: self.state.metadata.run_manifest.clone(),
1295 run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
1296 run_configuration: self.state.metadata.run_configuration.clone(),
1297 plugin_descriptors: self.plugins.descriptors().cloned().collect(),
1298 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
1299 commands: self.state.evidence.commands.clone(),
1300 command_attempts: self.state.evidence.command_attempts.clone(),
1301 ingress: self.state.evidence.ingress.clone(),
1302 boundaries: self.state.evidence.boundaries.clone(),
1303 final_time: self.state.scheduler.now,
1304 checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
1305 commitment_format_version: self.state.metadata.commitment_format_version,
1306 revision_format_version: self.state.metadata.replay_revision_format_version,
1307 final_revision: self.state.counters.state_revision,
1308 }
1309 }
1310
1311 fn compute_boundary_state_hash_for(
1312 &mut self,
1313 format: BoundaryStateHashFormat,
1314 ) -> Result<String, CanwuError> {
1315 match format {
1316 BoundaryStateHashFormat::LegacyV0 => self.compute_boundary_state_hash(),
1317 BoundaryStateHashFormat::CommitmentsV1 => {
1318 let roots = self.refresh_runtime_commitment_roots()?;
1319 boundary_state_hash_for_commitments(&roots)
1320 }
1321 }
1322 }
1323
1324 fn compute_boundary_state_hash(&self) -> Result<String, CanwuError> {
1325 let world = self.world();
1326 let plugin_components: Vec<_> = self
1327 .state
1328 .current
1329 .plugin_components
1330 .values()
1331 .cloned()
1332 .collect();
1333 let domain_records: Vec<_> = self
1334 .state
1335 .current
1336 .domain_records
1337 .values()
1338 .cloned()
1339 .collect();
1340 let plugin_descriptors: Vec<_> = self.plugins.descriptors().cloned().collect();
1341 let scheduled: Vec<_> = self
1342 .state
1343 .scheduler
1344 .actions
1345 .iter()
1346 .map(|(key, action)| ScheduledRecord {
1347 key: key.clone(),
1348 action: action.clone(),
1349 })
1350 .collect();
1351 let random_streams: Vec<_> = self
1352 .state
1353 .current
1354 .random_streams
1355 .values()
1356 .cloned()
1357 .collect();
1358 let (authoritative_manifest, authoritative_manifest_hash) = authoritative_run_identity(
1359 &self.state.metadata.run_manifest,
1360 &self.state.metadata.run_manifest_hash,
1361 &self.state.metadata.run_configuration,
1362 )?;
1363 let initial_scenario = self.bound_initial_scenario();
1364 state_hash(&StateHashMaterial {
1365 engine_version: ENGINE_VERSION,
1366 snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
1367 run_manifest: &authoritative_manifest,
1368 run_manifest_hash: &authoritative_manifest_hash,
1369 initial_time: self.state.scheduler.initial_time,
1370 initial_scenario,
1371 now: self.state.scheduler.now,
1372 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
1373 world: &world,
1374 knowledge: &self.state.current.knowledge,
1375 events: &self.state.evidence.events,
1376 commands: &self.state.evidence.commands,
1377 command_attempts: &self.state.evidence.command_attempts,
1378 ingress: &self.state.evidence.ingress,
1379 plugin_components: &plugin_components,
1380 domain_records: &domain_records,
1381 decisions: &self.state.current.decisions,
1382 plugin_descriptors: &plugin_descriptors,
1383 schema: &self.schema,
1384 scheduled: &scheduled,
1385 root_seed: self.state.current.root_seed,
1386 random_streams: &random_streams,
1387 random_draws: &self.state.evidence.random_draws,
1388 next_event_id: self.state.counters.next_event_id,
1389 next_command_id: self.state.counters.next_command_id,
1390 next_command_attempt_id: self.state.counters.next_command_attempt_id,
1391 next_ingress_id: self.state.counters.next_ingress_id,
1392 next_boundary_id: self.state.counters.next_boundary_id,
1393 next_random_draw_id: self.state.counters.next_random_draw_id,
1394 next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
1395 next_schedule_sequence: self.state.counters.next_schedule_sequence,
1396 next_correlation_id: self.state.counters.next_correlation_id,
1397 next_decision_trace_id: self.state.counters.next_decision_trace_id,
1398 })
1399 }
1400
1401 fn compute_commitment_root_updates(
1402 &self,
1403 needs: CommitmentDomains,
1404 ) -> Result<RuntimeCommitmentRootUpdates, CanwuError> {
1405 let world = needs
1406 .contains(CommitmentDomains::WORLD)
1407 .then(|| world_commitment_root(&self.world()))
1408 .transpose()?;
1409 let knowledge = needs
1410 .contains(CommitmentDomains::KNOWLEDGE)
1411 .then(|| knowledge_commitment_root(&self.state.current.knowledge))
1412 .transpose()?;
1413 let plugin_components = needs
1414 .contains(CommitmentDomains::PLUGIN_COMPONENTS)
1415 .then(|| {
1416 let values: Vec<_> = self
1417 .state
1418 .current
1419 .plugin_components
1420 .values()
1421 .cloned()
1422 .collect();
1423 plugin_component_commitment_root(&values)
1424 })
1425 .transpose()?;
1426 let domain_records = needs
1427 .contains(CommitmentDomains::DOMAIN_RECORDS)
1428 .then(|| {
1429 let values: Vec<_> = self
1430 .state
1431 .current
1432 .domain_records
1433 .values()
1434 .cloned()
1435 .collect();
1436 domain_record_commitment_root(&values)
1437 })
1438 .transpose()?;
1439 let decisions = needs
1440 .contains(CommitmentDomains::DECISIONS)
1441 .then(|| decision_commitment_root(&self.state.current.decisions))
1442 .transpose()?;
1443 let scheduler = needs
1444 .contains(CommitmentDomains::SCHEDULER)
1445 .then(|| {
1446 let scheduled: Vec<_> = self
1447 .state
1448 .scheduler
1449 .actions
1450 .iter()
1451 .map(|(key, action)| ScheduledRecord {
1452 key: key.clone(),
1453 action: action.clone(),
1454 })
1455 .collect();
1456 scheduler_commitment_root(self.state.scheduler.now, &scheduled)
1457 })
1458 .transpose()?;
1459 let random_streams = needs
1460 .contains(CommitmentDomains::RANDOM_STREAMS)
1461 .then(|| {
1462 let values: Vec<_> = self
1463 .state
1464 .current
1465 .random_streams
1466 .values()
1467 .cloned()
1468 .collect();
1469 random_stream_commitment_root(&values)
1470 })
1471 .transpose()?;
1472 let identity = if needs.contains(CommitmentDomains::IDENTITY) {
1473 let descriptors: Vec<_> = self.plugins.descriptors().cloned().collect();
1474 let (manifest, manifest_hash) = authoritative_run_identity(
1475 &self.state.metadata.run_manifest,
1476 &self.state.metadata.run_manifest_hash,
1477 &self.state.metadata.run_configuration,
1478 )?;
1479 Some(identity_commitment_root(
1480 ENGINE_VERSION,
1481 SNAPSHOT_FORMAT_VERSION,
1482 &manifest,
1483 &manifest_hash,
1484 self.state.scheduler.initial_time,
1485 self.bound_initial_scenario(),
1486 &descriptors,
1487 &self.schema,
1488 )?)
1489 } else {
1490 None
1491 };
1492 Ok(RuntimeCommitmentRootUpdates {
1493 world,
1494 knowledge,
1495 plugin_components,
1496 domain_records,
1497 decisions,
1498 scheduler,
1499 random_streams,
1500 identity,
1501 })
1502 }
1503
1504 fn invalidate_commitments(&mut self, domains: CommitmentDomains) {
1505 if let Some(cache) = self.state.metadata.commitment_cache.as_mut() {
1506 cache.invalidate(domains);
1507 }
1508 }
1509
1510 fn refresh_runtime_commitment_roots(&mut self) -> Result<CommitmentRoots, CanwuError> {
1511 if self.state.metadata.commitment_format_version != COMMITMENT_FORMAT_VERSION {
1512 return Err(CanwuError::new(
1513 ErrorCode::UnsupportedSnapshotVersion,
1514 format!(
1515 "commitment format {} cannot produce boundary state commitment v1",
1516 self.state.metadata.commitment_format_version
1517 ),
1518 ));
1519 }
1520 let needs = {
1521 let cache = if let Some(cache) = self.state.metadata.commitment_cache.as_mut() {
1522 cache
1523 } else {
1524 self.state.metadata.commitment_cache =
1525 Some(RuntimeCommitmentCache::from_evidence(&self.state.evidence)?);
1526 self.state
1527 .metadata
1528 .commitment_cache
1529 .as_mut()
1530 .expect("the commitment cache was initialized")
1531 };
1532 cache.sync(&self.state.evidence)?;
1533 cache.needs()
1534 };
1535 let updates = self.compute_commitment_root_updates(needs)?;
1536 let boundary_head = self.boundary_head_hash().map(str::to_owned);
1537 let control = ControlCommitmentMaterial {
1538 plugin_registration_closed: self.state.metadata.plugin_registration_closed,
1539 next_event_id: self.state.counters.next_event_id,
1540 next_command_id: self.state.counters.next_command_id,
1541 next_command_attempt_id: self.state.counters.next_command_attempt_id,
1542 next_ingress_id: self.state.counters.next_ingress_id,
1543 next_boundary_id: self.state.counters.next_boundary_id,
1544 next_random_draw_id: self.state.counters.next_random_draw_id,
1545 next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
1546 next_schedule_sequence: self.state.counters.next_schedule_sequence,
1547 next_correlation_id: self.state.counters.next_correlation_id,
1548 next_decision_trace_id: self.state.counters.next_decision_trace_id,
1549 };
1550 let (domain_roots, journal_roots) = {
1551 let cache = self
1552 .state
1553 .metadata
1554 .commitment_cache
1555 .as_mut()
1556 .expect("the commitment cache was initialized");
1557 cache.apply(updates);
1558 (cache.domain_roots()?, cache.roots())
1559 };
1560 runtime_commitment_roots(
1561 &domain_roots,
1562 &journal_roots,
1563 self.state.current.root_seed,
1564 boundary_head.as_deref(),
1565 &control,
1566 )
1567 }
1568
1569 fn refresh_checkpoint_hash(&mut self) -> Result<(), CanwuError> {
1570 if self.state.metadata.commitment_format_version == COMMITMENT_FORMAT_VERSION {
1571 let roots = self.refresh_runtime_commitment_roots()?;
1572 self.state.metadata.checkpoint_hash = checkpoint_hash_for_commitments(
1573 &roots,
1574 &self.state.metadata.run_manifest_hash,
1575 self.state.metadata.commitment_format_version,
1576 STATE_REVISION_FORMAT_VERSION,
1577 self.state.counters.state_revision,
1578 self.state.metadata.replay_revision_format_version,
1579 )?;
1580 self.state.metadata.commitment_roots = Some(roots);
1581 } else if self.state.metadata.commitment_format_version == 0 {
1582 let state_hash = self.compute_boundary_state_hash()?;
1583 self.state.metadata.checkpoint_hash = checkpoint_hash_for_configuration(
1584 &state_hash,
1585 self.boundary_head_hash(),
1586 &self.state.metadata.run_manifest_hash,
1587 &self.state.metadata.run_configuration,
1588 STATE_REVISION_FORMAT_VERSION,
1589 self.state.counters.state_revision,
1590 self.state.metadata.replay_revision_format_version,
1591 )?;
1592 self.state.metadata.commitment_roots = None;
1593 self.state.metadata.commitment_cache = None;
1594 } else {
1595 return Err(CanwuError::new(
1596 ErrorCode::UnsupportedSnapshotVersion,
1597 format!(
1598 "commitment format {} is unsupported; this engine writes format {COMMITMENT_FORMAT_VERSION}",
1599 self.state.metadata.commitment_format_version
1600 ),
1601 ));
1602 }
1603 Ok(())
1604 }
1605
1606 fn next_state_revision(&self) -> Result<u64, CanwuError> {
1607 self.state
1608 .counters
1609 .state_revision
1610 .checked_add(1)
1611 .ok_or_else(|| {
1612 CanwuError::new(
1613 ErrorCode::IdentifierExhausted,
1614 "authoritative state revision space is exhausted",
1615 )
1616 })
1617 }
1618
1619 fn advance_state_revision(&mut self) -> Result<u64, CanwuError> {
1620 let next = self.next_state_revision()?;
1621 self.state.counters.state_revision = next;
1622 Ok(next)
1623 }
1624
1625 #[must_use]
1626 pub fn snapshot(&self) -> SimulationSnapshot {
1627 let mut snapshot = self.checkpoint_state();
1628 snapshot.events.clone_from(&self.state.evidence.events);
1629 snapshot.commands.clone_from(&self.state.evidence.commands);
1630 snapshot
1631 .command_attempts
1632 .clone_from(&self.state.evidence.command_attempts);
1633 snapshot.ingress.clone_from(&self.state.evidence.ingress);
1634 snapshot
1635 .boundaries
1636 .clone_from(&self.state.evidence.boundaries);
1637 snapshot
1638 .random_draws
1639 .clone_from(&self.state.evidence.random_draws);
1640 snapshot
1641 }
1642
1643 pub fn snapshot_json(&self) -> Result<String, CanwuError> {
1644 serde_json::to_string_pretty(&self.snapshot()).map_err(|error| {
1645 CanwuError::new(
1646 ErrorCode::InvalidSnapshot,
1647 format!("could not serialize snapshot: {error}"),
1648 )
1649 })
1650 }
1651
1652 pub fn from_snapshot(snapshot: SimulationSnapshot) -> Result<Self, CanwuError> {
1653 if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION
1654 || snapshot.engine_version != ENGINE_VERSION
1655 {
1656 return Err(CanwuError::new(
1657 ErrorCode::UnsupportedSnapshotVersion,
1658 format!(
1659 "the typed snapshot loader accepts only engine {ENGINE_VERSION} format {SNAPSHOT_FORMAT_VERSION}; legacy format-4 JSON must use the strict JSON loader"
1660 ),
1661 ));
1662 }
1663 let snapshot = migrate_snapshot(snapshot)?;
1664 validate_scenario_state(&Scenario {
1665 start_time: snapshot.now,
1666 world: snapshot.world.clone(),
1667 knowledge: snapshot.knowledge.clone(),
1668 domain_records: snapshot.domain_records.clone(),
1669 })?;
1670 let plugins = PluginRegistry::from_descriptors(snapshot.plugin_descriptors.clone())?;
1671 validate_snapshot(&snapshot, &plugins)?;
1672 let admitted_ingress: BTreeSet<_> = snapshot
1673 .boundaries
1674 .iter()
1675 .flat_map(|boundary| boundary.admitted_ingress.iter().copied())
1676 .collect();
1677 let pending_ingress = snapshot
1678 .ingress
1679 .iter()
1680 .filter(|record| !admitted_ingress.contains(&record.id))
1681 .map(IngressQueueKey::from_record)
1682 .collect();
1683 let initial_scenario = match snapshot.initial_scenario.clone() {
1684 Some(initial_scenario) => Some(initial_scenario),
1685 None if !snapshot.plugin_registration_closed => match snapshot.run_manifest.as_ref() {
1686 Some(run_manifest @ RunManifest::Declared { .. }) => {
1687 let initial_scenario = Scenario {
1688 start_time: snapshot.initial_time,
1689 world: snapshot.world.clone(),
1690 knowledge: snapshot.knowledge.clone(),
1691 domain_records: snapshot.domain_records.clone(),
1692 };
1693 manifest::validate(run_manifest, Some(&initial_scenario), true)?;
1694 Some(initial_scenario)
1695 }
1696 _ => None,
1697 },
1698 None => None,
1699 };
1700 let mut simulation = Self {
1701 state: RuntimeState {
1702 current: RuntimeCurrentState {
1703 people: snapshot
1704 .world
1705 .people
1706 .into_iter()
1707 .map(|value| (value.id, value))
1708 .collect(),
1709 letters: snapshot
1710 .world
1711 .letters
1712 .into_iter()
1713 .map(|value| (value.id, value))
1714 .collect(),
1715 governments: snapshot
1716 .world
1717 .governments
1718 .into_iter()
1719 .map(|value| (value.id, value))
1720 .collect(),
1721 territories: snapshot
1722 .world
1723 .territories
1724 .into_iter()
1725 .map(|value| (value.id, value))
1726 .collect(),
1727 routes: snapshot
1728 .world
1729 .routes
1730 .into_iter()
1731 .map(|value| (value.id, value))
1732 .collect(),
1733 armies: snapshot
1734 .world
1735 .armies
1736 .into_iter()
1737 .map(|value| (value.id, value))
1738 .collect(),
1739 knowledge: snapshot.knowledge,
1740 plugin_components: snapshot
1741 .plugin_components
1742 .into_iter()
1743 .map(|record| {
1744 (
1745 component_key(
1746 &record.plugin,
1747 &record.state,
1748 &record.entity,
1749 &record.component,
1750 ),
1751 record,
1752 )
1753 })
1754 .collect(),
1755 domain_records: snapshot
1756 .domain_records
1757 .into_iter()
1758 .map(|record| (record.reference.clone(), record))
1759 .collect(),
1760 decisions: snapshot.decisions,
1761 root_seed: snapshot.root_seed,
1762 random_streams: snapshot
1763 .random_streams
1764 .into_iter()
1765 .map(|state| (state.key.clone(), state))
1766 .collect(),
1767 },
1768 scheduler: RuntimeScheduler {
1769 initial_time: snapshot.initial_time,
1770 now: snapshot.now,
1771 actions: snapshot
1772 .scheduled
1773 .into_iter()
1774 .map(|record| (record.key, record.action))
1775 .collect(),
1776 pending_ingress,
1777 },
1778 counters: RuntimeCounters {
1779 next_event_id: snapshot.next_event_id,
1780 next_command_id: snapshot.next_command_id,
1781 next_command_attempt_id: snapshot.next_command_attempt_id,
1782 next_ingress_id: snapshot.next_ingress_id,
1783 next_boundary_id: snapshot.next_boundary_id,
1784 next_random_draw_id: snapshot.next_random_draw_id,
1785 next_knowledge_record_id: snapshot.next_knowledge_record_id,
1786 next_schedule_sequence: snapshot.next_schedule_sequence,
1787 next_correlation_id: snapshot.next_correlation_id,
1788 next_decision_trace_id: snapshot.next_decision_trace_id,
1789 state_revision: snapshot.state_revision,
1790 admitted_attempt_count: snapshot.admitted_attempt_count,
1791 admitted_command_count: snapshot.admitted_command_count,
1792 admitted_event_count: snapshot.admitted_event_count,
1793 },
1794 metadata: RuntimeMetadata {
1795 initial_scenario,
1796 run_manifest: snapshot.run_manifest.clone().ok_or_else(|| {
1797 invalid_snapshot_error("snapshot is missing its run manifest")
1798 })?,
1799 run_manifest_hash: snapshot.run_manifest_hash.clone(),
1800 run_configuration: snapshot.run_configuration.clone().ok_or_else(|| {
1801 invalid_snapshot_error("snapshot is missing its run configuration")
1802 })?,
1803 checkpoint_hash: snapshot.checkpoint_hash.clone(),
1804 commitment_format_version: snapshot.commitment_format_version,
1805 commitment_roots: snapshot.commitment_roots.clone(),
1806 commitment_cache: None,
1807 plugin_registration_closed: snapshot.plugin_registration_closed,
1808 replay_revision_format_version: snapshot.replay_revision_format_version,
1809 },
1810 evidence: RuntimeEvidence {
1811 archived: EvidenceCursor::default(),
1812 archived_boundary_head: None,
1813 archived_legacy_commands: false,
1814 archived_tracked_attempts: false,
1815 archived_unqueued_command_history: false,
1816 archived_command_requests: BTreeMap::new(),
1817 archived_ingress_requests: BTreeMap::new(),
1818 archived_decision_requests: BTreeMap::new(),
1819 archived_decision_command_requests: BTreeSet::new(),
1820 events: snapshot.events,
1821 commands: snapshot.commands,
1822 command_attempts: snapshot.command_attempts,
1823 ingress: snapshot.ingress,
1824 boundaries: snapshot.boundaries,
1825 random_draws: snapshot.random_draws,
1826 archived_segment_headers: Vec::new(),
1827 archived_evidence_receipts: BTreeMap::new(),
1828 keyed_draw_reservations: Vec::new(),
1829 },
1830 },
1831 schema: snapshot.schema,
1832 plugins,
1833 sync_reaction_depth: 0,
1834 };
1835 simulation.refresh_checkpoint_hash()?;
1836 Ok(simulation)
1837 }
1838
1839 pub fn from_snapshot_json(json: &str) -> Result<Self, CanwuError> {
1840 let snapshot = legacy_v4::deserialize_snapshot_json(json)?;
1841 Self::from_snapshot(snapshot)
1842 }
1843
1844 pub fn from_snapshot_with_plugins(
1845 snapshot: SimulationSnapshot,
1846 plugins: &[&dyn SimulationPlugin],
1847 ) -> Result<Self, CanwuError> {
1848 let mut simulation = Self::from_snapshot(snapshot)?;
1849 for plugin in plugins {
1850 simulation.register_plugin(*plugin)?;
1851 }
1852 simulation.ensure_runtime_ready()?;
1853 Ok(simulation)
1854 }
1855
1856 pub fn from_snapshot_json_with_plugins(
1857 json: &str,
1858 plugins: &[&dyn SimulationPlugin],
1859 ) -> Result<Self, CanwuError> {
1860 let snapshot = legacy_v4::deserialize_snapshot_json(json)?;
1861 Self::from_snapshot_with_plugins(snapshot, plugins)
1862 }
1863
1864 #[must_use]
1865 pub fn fork(&self) -> Self {
1866 Self {
1867 state: self.state.clone(),
1868 schema: self.schema.clone(),
1869 plugins: self.plugins.clone(),
1870 sync_reaction_depth: 0,
1871 }
1872 }
1873
1874 fn prepare_command(
1875 &self,
1876 envelope: &CommandEnvelope,
1877 context: &CommandContext,
1878 ) -> Result<PreparedCommand, CanwuError> {
1879 match &envelope.command {
1880 Command::OrderMovement {
1881 subject,
1882 destination,
1883 cargo,
1884 } => {
1885 let Some(actor) = decision_actor(&context.authority) else {
1886 return Err(CanwuError::new(
1887 ErrorCode::InvalidAuthority,
1888 "movement commands require an accountable actor origin",
1889 ));
1890 };
1891 let person = self.state.current.people.get(&actor).ok_or_else(|| {
1892 CanwuError::new(
1893 ErrorCode::ActorNotFound,
1894 format!("actor {actor} was not found"),
1895 )
1896 .with_entity(EntityRef::Person(actor))
1897 })?;
1898 if context
1899 .authority
1900 .command_subject
1901 .as_ref()
1902 .is_some_and(|bound| bound != subject)
1903 {
1904 return Err(CanwuError::new(
1905 ErrorCode::InvalidAuthority,
1906 "command subject does not match the movement subject",
1907 )
1908 .with_entity(subject.clone()));
1909 }
1910 if !self.state.current.territories.contains_key(destination) {
1911 return Err(CanwuError::new(
1912 ErrorCode::DestinationNotFound,
1913 format!("destination {destination} was not found"),
1914 )
1915 .with_entity(EntityRef::Territory(*destination)));
1916 }
1917 if cargo.windows(2).any(|pair| pair[0] >= pair[1]) {
1918 return Err(CanwuError::new(
1919 ErrorCode::InvalidPayload,
1920 "movement cargo IDs must be sorted and unique",
1921 ));
1922 }
1923 match subject {
1924 EntityRef::Army(army) => {
1925 if !cargo.is_empty() {
1926 return Err(CanwuError::new(
1927 ErrorCode::InvalidPayload,
1928 "army movement does not accept letter cargo yet",
1929 ));
1930 }
1931 let army_state = self.state.current.armies.get(army).ok_or_else(|| {
1932 CanwuError::new(
1933 ErrorCode::ArmyNotFound,
1934 format!("army {army} was not found"),
1935 )
1936 .with_entity(EntityRef::Army(*army))
1937 })?;
1938 if army_state.commander != person.id {
1939 return Err(CanwuError::new(
1940 ErrorCode::InvalidAuthority,
1941 format!("{} does not command {}", person.name, army_state.name),
1942 )
1943 .with_entity(EntityRef::Person(person.id))
1944 .with_entity(EntityRef::Army(*army)));
1945 }
1946 if army_state.transit.is_some() {
1947 return Err(CanwuError::new(
1948 ErrorCode::InvalidAuthority,
1949 format!("{} is already moving", army_state.name),
1950 )
1951 .with_entity(EntityRef::Army(*army)));
1952 }
1953 let arrival_at =
1954 self.movement_arrival_time(army_state.location, *destination)?;
1955 Ok(PreparedCommand::ArmyMovement {
1956 army: *army,
1957 actor,
1958 from: army_state.location,
1959 destination: *destination,
1960 arrival_at,
1961 })
1962 }
1963 EntityRef::Person(person_id) => {
1964 if *person_id != actor
1965 || context
1966 .authority
1967 .command_subject
1968 .as_ref()
1969 .is_some_and(|subject| subject != &EntityRef::Person(*person_id))
1970 {
1971 return Err(CanwuError::new(
1972 ErrorCode::InvalidAuthority,
1973 "self-directed movement must bind the actor to the person subject",
1974 )
1975 .with_entity(EntityRef::Person(*person_id)));
1976 }
1977 let person_state =
1978 self.state.current.people.get(person_id).ok_or_else(|| {
1979 CanwuError::new(
1980 ErrorCode::EntityNotFound,
1981 format!("person {person_id} was not found"),
1982 )
1983 .with_entity(EntityRef::Person(*person_id))
1984 })?;
1985 if person_state.transit.is_some() {
1986 return Err(CanwuError::new(
1987 ErrorCode::InvalidAuthority,
1988 format!("person {person_id} is already moving"),
1989 )
1990 .with_entity(EntityRef::Person(*person_id)));
1991 }
1992 for letter_id in cargo {
1993 let letter =
1994 self.state.current.letters.get(letter_id).ok_or_else(|| {
1995 CanwuError::new(
1996 ErrorCode::EntityNotFound,
1997 format!("letter {letter_id} was not found"),
1998 )
1999 .with_entity(
2000 EntityRef::Resource(ResourceId::new(letter_id.get())),
2001 )
2002 })?;
2003 if letter.status != LetterStatus::HeldByPerson
2004 || letter.carrier != Some(*person_id)
2005 || !self.state.current.people.contains_key(&letter.sender)
2006 || !self.state.current.people.contains_key(&letter.recipient)
2007 {
2008 return Err(CanwuError::new(
2009 ErrorCode::InvalidAuthority,
2010 format!("letter {letter_id} is not held by the moving person"),
2011 )
2012 .with_entity(EntityRef::Resource(ResourceId::new(
2013 letter_id.get(),
2014 ))));
2015 }
2016 }
2017 let arrival_at = self
2018 .movement_arrival_time(person_state.current_location, *destination)?;
2019 Ok(PreparedCommand::MovePerson {
2020 person: *person_id,
2021 from: person_state.current_location,
2022 destination: *destination,
2023 cargo: cargo.clone(),
2024 arrival_at,
2025 })
2026 }
2027 _ => Err(CanwuError::new(
2028 ErrorCode::InvalidAuthority,
2029 "only army and person subjects support built-in movement",
2030 )
2031 .with_entity(subject.clone())),
2032 }
2033 }
2034 Command::DebugSetArmyMorale { army, morale } => {
2035 if envelope.issuer != Issuer::Debug {
2036 return Err(CanwuError::new(
2037 ErrorCode::InvalidAuthority,
2038 "debug state edits require the explicit debug issuer",
2039 ));
2040 }
2041 if *morale > 100 {
2042 return Err(CanwuError::new(
2043 ErrorCode::ValueOutOfRange,
2044 "army morale must be between 0 and 100",
2045 ));
2046 }
2047 let old_morale = self.state.current.armies.get(army).map_or_else(
2048 || {
2049 Err(CanwuError::new(
2050 ErrorCode::ArmyNotFound,
2051 format!("army {army} was not found"),
2052 ))
2053 },
2054 |army_state| Ok(army_state.morale),
2055 )?;
2056 Ok(PreparedCommand::DebugMorale {
2057 army: *army,
2058 old_morale,
2059 new_morale: *morale,
2060 })
2061 }
2062 Command::Plugin {
2063 plugin,
2064 command,
2065 payload,
2066 } => {
2067 let registered = self
2068 .plugins
2069 .commands
2070 .get(&(plugin.clone(), command.clone()))
2071 .ok_or_else(|| {
2072 CanwuError::new(
2073 ErrorCode::PluginCommandNotFound,
2074 format!("plugin command {plugin}.{command} is not registered"),
2075 )
2076 })?;
2077 let handler = registered.handler;
2078 let descriptor = registered.descriptor.clone();
2079 descriptor.payload_schema.validate(payload)?;
2080 let reader = format!("{plugin}.{command}");
2081 let directives = catch_unwind(AssertUnwindSafe(|| {
2082 handler(
2083 &self.plugin_view(&reader, &descriptor.reads),
2084 context,
2085 payload,
2086 )
2087 }))
2088 .map_err(|_| {
2089 CanwuError::new(
2090 ErrorCode::PluginPanicked,
2091 format!("plugin command {plugin}.{command} panicked"),
2092 )
2093 })??;
2094 validate_directives_with_context(
2095 &RuntimeValidationContext::new(&self.state),
2096 plugin,
2097 &descriptor.writes,
2098 &self.plugins.state_owners,
2099 &self.plugins.record_schemas,
2100 &directives,
2101 )?;
2102 Ok(PreparedCommand::Plugin {
2103 plugin: plugin.clone(),
2104 directives,
2105 allowed_writes: descriptor.writes,
2106 })
2107 }
2108 }
2109 }
2110
2111 fn movement_arrival_time(
2112 &self,
2113 from: TerritoryId,
2114 to: TerritoryId,
2115 ) -> Result<SimTime, CanwuError> {
2116 let travel_minutes = if from == to {
2117 1
2118 } else {
2119 self.state
2120 .current
2121 .routes
2122 .values()
2123 .find(|route| route.connects(from, to))
2124 .ok_or_else(|| {
2125 CanwuError::new(
2126 ErrorCode::NoRoute,
2127 format!("no direct route connects territory {from} to {to}"),
2128 )
2129 })?
2130 .travel_minutes
2131 };
2132 if travel_minutes <= 0 {
2133 return Err(CanwuError::new(
2134 ErrorCode::InvalidDuration,
2135 "movement route duration must be positive",
2136 ));
2137 }
2138 self.state
2139 .scheduler
2140 .now
2141 .checked_add(SimDuration::minutes(travel_minutes))
2142 .ok_or_else(|| {
2143 CanwuError::new(
2144 ErrorCode::InvalidDuration,
2145 "movement arrival time exceeds the supported range",
2146 )
2147 })
2148 }
2149
2150 fn apply_prepared(
2151 &mut self,
2152 prepared: PreparedCommand,
2153 command_id: CommandId,
2154 correlation_id: u64,
2155 ) -> Result<(), CanwuError> {
2156 match prepared {
2157 PreparedCommand::ArmyMovement {
2158 army,
2159 actor,
2160 from,
2161 destination,
2162 arrival_at,
2163 } => {
2164 let army_state = self.state.current.armies.get_mut(&army).ok_or_else(|| {
2165 CanwuError::new(ErrorCode::ArmyNotFound, "validated army disappeared")
2166 })?;
2167 army_state.transit = Some(TransitState {
2168 from,
2169 to: destination,
2170 departed_at: self.state.scheduler.now,
2171 arrives_at: arrival_at,
2172 });
2173 let event = self.emit(
2174 EventKind::MoveOrdered {
2175 army,
2176 from,
2177 to: destination,
2178 arrival_at,
2179 },
2180 vec![
2181 EntityRef::Army(army),
2182 EntityRef::Person(actor),
2183 EntityRef::Territory(from),
2184 EntityRef::Territory(destination),
2185 ],
2186 format!("Army {army} was ordered from {from} to {destination}"),
2187 Some(CauseRef::Command(command_id)),
2188 correlation_id,
2189 )?;
2190 self.schedule_at(
2191 arrival_at,
2192 ScheduledAction::ArmyArrival {
2193 army,
2194 destination,
2195 order_event: event,
2196 correlation_id,
2197 },
2198 )?;
2199 }
2200 PreparedCommand::MovePerson {
2201 person,
2202 from,
2203 destination,
2204 cargo,
2205 arrival_at,
2206 } => {
2207 self.invalidate_commitments(CommitmentDomains::WORLD);
2208 let person_state = self.state.current.people.get_mut(&person).ok_or_else(|| {
2209 CanwuError::new(ErrorCode::EntityNotFound, "validated person disappeared")
2210 })?;
2211 person_state.transit = Some(PersonTransitState {
2212 from,
2213 to: destination,
2214 departed_at: self.state.scheduler.now,
2215 arrives_at: arrival_at,
2216 });
2217 for letter_id in &cargo {
2218 let letter =
2219 self.state
2220 .current
2221 .letters
2222 .get_mut(letter_id)
2223 .ok_or_else(|| {
2224 CanwuError::new(
2225 ErrorCode::EntityNotFound,
2226 "validated letter disappeared",
2227 )
2228 })?;
2229 letter.status = LetterStatus::InTransit;
2230 letter.carrier = Some(person);
2231 letter.location = None;
2232 }
2233 let event = self.emit(
2234 EventKind::PersonMoveOrdered {
2235 person,
2236 from,
2237 to: destination,
2238 arrival_at,
2239 },
2240 std::iter::once(EntityRef::Person(person))
2241 .chain(
2242 cargo
2243 .iter()
2244 .copied()
2245 .map(|id| EntityRef::Resource(ResourceId::new(id.get()))),
2246 )
2247 .chain([
2248 EntityRef::Territory(from),
2249 EntityRef::Territory(destination),
2250 ])
2251 .collect(),
2252 format!("Person {person} was ordered from {from} to {destination}"),
2253 Some(CauseRef::Command(command_id)),
2254 correlation_id,
2255 )?;
2256 self.schedule_at(
2257 arrival_at,
2258 ScheduledAction::PersonArrival {
2259 person,
2260 destination,
2261 order_event: event,
2262 cargo,
2263 correlation_id,
2264 },
2265 )?;
2266 }
2267 PreparedCommand::DebugMorale {
2268 army,
2269 old_morale,
2270 new_morale,
2271 } => {
2272 self.state
2273 .current
2274 .armies
2275 .get_mut(&army)
2276 .ok_or_else(|| {
2277 CanwuError::new(ErrorCode::ArmyNotFound, "validated army disappeared")
2278 })?
2279 .morale = new_morale;
2280 self.emit(
2281 EventKind::DebugFieldChanged {
2282 entity: EntityRef::Army(army),
2283 field: "morale".to_owned(),
2284 old_value: old_morale.to_string(),
2285 new_value: new_morale.to_string(),
2286 },
2287 vec![EntityRef::Army(army)],
2288 format!(
2289 "Debug command changed army {army} morale {old_morale} -> {new_morale}"
2290 ),
2291 Some(CauseRef::Command(command_id)),
2292 correlation_id,
2293 )?;
2294 }
2295 PreparedCommand::Plugin {
2296 plugin,
2297 directives,
2298 allowed_writes,
2299 } => {
2300 self.apply_directives(
2301 &plugin,
2302 directives,
2303 &allowed_writes,
2304 &CauseRef::Command(command_id),
2305 correlation_id,
2306 )?;
2307 }
2308 }
2309 Ok(())
2310 }
2311}
2312
2313enum PreparedCommand {
2314 ArmyMovement {
2315 army: ArmyId,
2316 actor: PersonId,
2317 from: TerritoryId,
2318 destination: TerritoryId,
2319 arrival_at: SimTime,
2320 },
2321 MovePerson {
2322 person: PersonId,
2323 from: TerritoryId,
2324 destination: TerritoryId,
2325 cargo: Vec<LetterId>,
2326 arrival_at: SimTime,
2327 },
2328 DebugMorale {
2329 army: ArmyId,
2330 old_morale: u16,
2331 new_morale: u16,
2332 },
2333 Plugin {
2334 plugin: String,
2335 directives: Vec<SystemDirective>,
2336 allowed_writes: Vec<StateKey>,
2337 },
2338}
2339
2340impl PreparedCommand {
2341 fn commitment_invalidation(&self) -> CommitmentDomains {
2342 match self {
2343 Self::ArmyMovement { .. } => {
2344 CommitmentDomains::WORLD
2345 | CommitmentDomains::KNOWLEDGE
2346 | CommitmentDomains::PLUGIN_COMPONENTS
2347 | CommitmentDomains::SCHEDULER
2348 }
2349 Self::MovePerson { .. } => CommitmentDomains::WORLD | CommitmentDomains::SCHEDULER,
2350 Self::DebugMorale { .. } => {
2351 CommitmentDomains::WORLD
2352 | CommitmentDomains::PLUGIN_COMPONENTS
2353 | CommitmentDomains::SCHEDULER
2354 }
2355 Self::Plugin { .. } => {
2356 CommitmentDomains::PLUGIN_COMPONENTS | CommitmentDomains::SCHEDULER
2357 }
2358 }
2359 }
2360}
2361
2362fn validate_directives(
2363 plugin: &str,
2364 allowed_writes: &[StateKey],
2365 state_owners: &BTreeMap<StateKey, String>,
2366 record_schemas: &records::DomainRecordSchemas,
2367 entity_exists: &dyn Fn(&EntityRef) -> bool,
2368 directives: &[SystemDirective],
2369) -> Result<(), CanwuError> {
2370 for directive in directives {
2371 match directive {
2372 SystemDirective::SetComponent {
2373 state,
2374 entity,
2375 component,
2376 ..
2377 } => {
2378 if component.trim().is_empty() || component != component.trim() {
2379 return Err(CanwuError::new(
2380 ErrorCode::InvalidPayload,
2381 "plugin component name must be non-empty and canonical",
2382 ));
2383 }
2384 if !allowed_writes.contains(state) {
2385 return Err(CanwuError::new(
2386 ErrorCode::UndeclaredStateWrite,
2387 format!(
2388 "plugin {plugin} did not declare write access to {}.{}",
2389 state.namespace, state.name
2390 ),
2391 ));
2392 }
2393 if state_owners.get(state).is_none_or(|owner| owner != plugin) {
2394 return Err(CanwuError::new(
2395 ErrorCode::UndeclaredStateWrite,
2396 format!(
2397 "plugin {plugin} does not own state {}.{}",
2398 state.namespace, state.name
2399 ),
2400 ));
2401 }
2402 if is_domain_record_state(record_schemas, state) {
2403 return Err(CanwuError::new(
2404 ErrorCode::UndeclaredStateWrite,
2405 "domain record state cannot be written as an immediate component",
2406 ));
2407 }
2408 if !entity_exists(entity) {
2409 return Err(CanwuError::new(
2410 ErrorCode::EntityNotFound,
2411 format!("plugin {plugin} targeted missing entity {entity}"),
2412 )
2413 .with_entity(entity.clone()));
2414 }
2415 }
2416 SystemDirective::Emit { event_type, .. }
2417 if event_type.trim().is_empty() || event_type != event_type.trim() =>
2418 {
2419 return Err(CanwuError::new(
2420 ErrorCode::InvalidPayload,
2421 "plugin event type must be non-empty and canonical",
2422 ));
2423 }
2424 SystemDirective::Emit { affected, .. }
2425 if affected.iter().any(|entity| !entity_exists(entity)) =>
2426 {
2427 return Err(CanwuError::new(
2428 ErrorCode::EntityNotFound,
2429 format!("plugin {plugin} emitted an event for a missing entity"),
2430 ));
2431 }
2432 SystemDirective::Schedule { after, directive } => {
2433 if *after <= SimDuration::ZERO {
2434 return Err(CanwuError::new(
2435 ErrorCode::InvalidDuration,
2436 "plugin systems must schedule work strictly in the future",
2437 ));
2438 }
2439 validate_directives(
2440 plugin,
2441 allowed_writes,
2442 state_owners,
2443 record_schemas,
2444 entity_exists,
2445 std::slice::from_ref(directive),
2446 )?;
2447 }
2448 SystemDirective::EnqueuePluginIngress {
2449 after,
2450 packet_type,
2451 affected,
2452 ..
2453 } => {
2454 if packet_type.trim().is_empty() || packet_type != packet_type.trim() {
2455 return Err(CanwuError::new(
2456 ErrorCode::InvalidPayload,
2457 "plugin ingress type must be non-empty and canonical",
2458 ));
2459 }
2460 if *after < SimDuration::ZERO {
2461 return Err(CanwuError::new(
2462 ErrorCode::InvalidDuration,
2463 "plugin command ingress delay cannot be negative",
2464 ));
2465 }
2466 if affected.iter().any(|entity| !entity_exists(entity)) {
2467 return Err(CanwuError::new(
2468 ErrorCode::EntityNotFound,
2469 format!("plugin {plugin} queued ingress for a missing entity"),
2470 ));
2471 }
2472 }
2473 SystemDirective::Emit { .. } => {}
2474 }
2475 }
2476 Ok(())
2477}
2478
2479fn resolve_command_authority(envelope: &CommandEnvelope) -> Result<CommandAuthority, CanwuError> {
2480 if let Some(authority) = &envelope.authority {
2481 return Ok(authority.clone());
2482 }
2483 match &envelope.issuer {
2484 Issuer::Actor(actor) => Ok(CommandAuthority::for_actor(*actor)),
2485 Issuer::Debug => Ok(CommandAuthority::no_responsible_actor("debug-command")),
2486 Issuer::System(system) => Ok(CommandAuthority::no_responsible_actor(format!(
2487 "system:{system}"
2488 ))),
2489 Issuer::Human(_)
2490 | Issuer::Ai(_)
2491 | Issuer::Institution(_)
2492 | Issuer::Replay(_)
2493 | Issuer::Experiment(_) => Err(CanwuError::new(
2494 ErrorCode::InvalidAuthority,
2495 "typed command origins require an explicit authority context",
2496 )),
2497 }
2498}
2499
2500fn validate_command_ingress_policy(
2501 run_configuration: &RunConfigurationSnapshot,
2502 issuer: &Issuer,
2503 authority: &CommandAuthority,
2504 admission: CommandAdmission,
2505 entity_exists: &dyn Fn(&EntityRef) -> bool,
2506) -> Result<(), CanwuError> {
2507 let CommandAdmission {
2508 request_id,
2509 expected_revision,
2510 expected_time,
2511 revision_before: current_revision,
2512 ingress,
2513 } = admission;
2514 if request_id.is_some_and(|id| id.get() == 0) {
2515 return Err(CanwuError::new(
2516 ErrorCode::InvalidPayload,
2517 "command request IDs must be nonzero",
2518 ));
2519 }
2520 if let Some(expected) = expected_revision
2521 && expected != current_revision
2522 {
2523 return Err(CanwuError::new(
2524 ErrorCode::SimulationRevisionConflict,
2525 format!(
2526 "command expected revision {expected}, but simulation is at revision {current_revision}"
2527 ),
2528 ));
2529 }
2530 validate_command_authority(authority, entity_exists)?;
2531 if matches!(issuer, Issuer::Replay(_)) != (ingress == CommandIngress::FrozenReplay) {
2532 return Err(CanwuError::new(
2533 ErrorCode::InvalidAuthority,
2534 "replay command origins are valid only for frozen replay ingress",
2535 ));
2536 }
2537
2538 let RunConfigurationSnapshot::Declared(configuration) = run_configuration else {
2539 return Ok(());
2540 };
2541 if ingress == CommandIngress::LegacyDirect {
2542 return Err(CanwuError::new(
2543 ErrorCode::InvalidAuthority,
2544 "declared runs require tracked request or frozen replay ingress",
2545 ));
2546 }
2547 let external = !matches!(issuer, Issuer::System(_));
2548 if configuration.require_idempotency_keys && external && request_id.is_none() {
2549 return Err(CanwuError::new(
2550 ErrorCode::MissingIdempotencyKey,
2551 "this run requires a stable command request ID",
2552 ));
2553 }
2554 if configuration.require_idempotency_keys && external && expected_revision.is_none() {
2555 return Err(CanwuError::new(
2556 ErrorCode::SimulationRevisionConflict,
2557 "this run requires an expected command revision",
2558 ));
2559 }
2560 if configuration.interaction == InteractionPolicy::ReadOnly
2561 && !matches!(issuer, Issuer::Replay(_) | Issuer::System(_))
2562 {
2563 return Err(CanwuError::new(
2564 ErrorCode::InteractionReadOnly,
2565 "the run interaction policy rejects newly authored authoritative commands",
2566 ));
2567 }
2568 if external && expected_time.is_none() {
2569 return Err(CanwuError::new(
2570 ErrorCode::SimulationTimeConflict,
2571 "declared external commands require an expected simulation time",
2572 ));
2573 }
2574
2575 match issuer {
2576 Issuer::Actor(_) => Err(CanwuError::new(
2577 ErrorCode::InvalidAuthority,
2578 "declared runs require a typed human, AI, institution, replay, experiment, debug, or system origin",
2579 )),
2580 Issuer::Human(controller) => {
2581 let Some(binding) = &configuration.seat_binding else {
2582 return Err(CanwuError::new(
2583 ErrorCode::InvalidAuthority,
2584 "human commands require the run's exact seat binding",
2585 ));
2586 };
2587 if configuration.controller != ControllerPolicy::HumanRoleBound
2588 || controller != &binding.controller_id
2589 || authority.seat_id.as_deref() != Some(binding.seat_id.as_str())
2590 || authority.permission_profile_id.as_deref()
2591 != Some(binding.permission_profile_id.as_str())
2592 || !authority_matches_seat_binding(configuration.seat, binding, authority)
2593 {
2594 return Err(CanwuError::new(
2595 ErrorCode::InvalidAuthority,
2596 "human command origin does not match the active controller, seat binding, and permission profile",
2597 ));
2598 }
2599 Ok(())
2600 }
2601 Issuer::Ai(controller) | Issuer::Institution(controller) => {
2602 if !canonical_text(controller)
2603 || matches!(
2604 authority.decision_origin,
2605 DecisionOrigin::NoResponsibleActor { .. }
2606 )
2607 {
2608 return Err(CanwuError::new(
2609 ErrorCode::InvalidAuthority,
2610 "AI and institutional commands require a canonical controller and responsible decision origin",
2611 ));
2612 }
2613 Ok(())
2614 }
2615 Issuer::Replay(source) => {
2616 if !canonical_text(source)
2617 || ingress != CommandIngress::FrozenReplay
2618 || configuration.purpose != RunPurpose::Replay
2619 || configuration.controller != ControllerPolicy::ReplayController
2620 || configuration.interaction != InteractionPolicy::ReadOnly
2621 {
2622 return Err(CanwuError::new(
2623 ErrorCode::InvalidAuthority,
2624 "replay command sources require a replay-purpose, replay-controller, read-only run",
2625 ));
2626 }
2627 if let Some(binding) = &configuration.seat_binding
2628 && (source != &binding.controller_id
2629 || authority.seat_id.as_deref() != Some(binding.seat_id.as_str())
2630 || authority.permission_profile_id.as_deref()
2631 != Some(binding.permission_profile_id.as_str())
2632 || !authority_matches_seat_binding(configuration.seat, binding, authority))
2633 {
2634 return Err(CanwuError::new(
2635 ErrorCode::InvalidAuthority,
2636 "frozen replay input does not match its recorded controller and seat binding",
2637 ));
2638 }
2639 Ok(())
2640 }
2641 Issuer::Experiment(intervention) => {
2642 if configuration.interaction != InteractionPolicy::VersionedExperiment
2643 || !configuration.declared_interventions.contains(intervention)
2644 {
2645 return Err(CanwuError::new(
2646 ErrorCode::InvalidAuthority,
2647 "experiment commands must name an intervention declared by the run",
2648 ));
2649 }
2650 Ok(())
2651 }
2652 Issuer::Debug => {
2653 if !configuration.diagnostic_commands_enabled {
2654 return Err(CanwuError::new(
2655 ErrorCode::InvalidAuthority,
2656 "debug command authority is disabled by the run configuration",
2657 ));
2658 }
2659 Ok(())
2660 }
2661 Issuer::System(system) => {
2662 if !canonical_text(system)
2663 || !matches!(
2664 authority.decision_origin,
2665 DecisionOrigin::NoResponsibleActor { .. }
2666 )
2667 {
2668 return Err(CanwuError::new(
2669 ErrorCode::InvalidAuthority,
2670 "system commands require a canonical system ID and typed no-responsible-actor origin",
2671 ));
2672 }
2673 Ok(())
2674 }
2675 }
2676}
2677
2678fn validate_command_authority(
2679 authority: &CommandAuthority,
2680 entity_exists: &dyn Fn(&EntityRef) -> bool,
2681) -> Result<(), CanwuError> {
2682 if authority
2683 .seat_id
2684 .as_ref()
2685 .is_some_and(|value| !canonical_text(value))
2686 || authority
2687 .permission_profile_id
2688 .as_ref()
2689 .is_some_and(|value| !canonical_text(value))
2690 || authority.seat_id.is_some() != authority.permission_profile_id.is_some()
2691 || authority
2692 .command_subject
2693 .as_ref()
2694 .is_some_and(|entity| !entity_exists(entity))
2695 {
2696 return Err(CanwuError::new(
2697 ErrorCode::InvalidAuthority,
2698 "command authority contains an invalid seat, permission profile, or subject",
2699 ));
2700 }
2701 match &authority.decision_origin {
2702 DecisionOrigin::Actor { actor } => {
2703 if !entity_exists(&EntityRef::Person(*actor)) {
2704 return Err(CanwuError::new(
2705 ErrorCode::InvalidAuthority,
2706 "command decision origin references a missing actor",
2707 ));
2708 }
2709 }
2710 DecisionOrigin::Institution {
2711 institution,
2712 responsible_actor,
2713 } => {
2714 if !entity_exists(institution)
2715 || responsible_actor.is_some_and(|actor| !entity_exists(&EntityRef::Person(actor)))
2716 {
2717 return Err(CanwuError::new(
2718 ErrorCode::InvalidAuthority,
2719 "command decision origin references a missing institution or actor",
2720 ));
2721 }
2722 }
2723 DecisionOrigin::Council { council_id } if !canonical_text(council_id) => {
2724 return Err(CanwuError::new(
2725 ErrorCode::InvalidAuthority,
2726 "command council origin requires a canonical ID",
2727 ));
2728 }
2729 DecisionOrigin::NoResponsibleActor { reason } if !canonical_text(reason) => {
2730 return Err(CanwuError::new(
2731 ErrorCode::InvalidAuthority,
2732 "no-responsible-actor origins require a canonical reason",
2733 ));
2734 }
2735 DecisionOrigin::Council { .. } | DecisionOrigin::NoResponsibleActor { .. } => {}
2736 }
2737 Ok(())
2738}
2739
2740fn authority_matches_seat_binding(
2741 seat: SeatPolicy,
2742 binding: &SeatBinding,
2743 authority: &CommandAuthority,
2744) -> bool {
2745 match (seat, &authority.decision_origin) {
2746 (SeatPolicy::CharacterBound, DecisionOrigin::Actor { actor }) => {
2747 binding.actor == Some(*actor) && binding.institution.is_none()
2748 }
2749 (
2750 SeatPolicy::InstitutionBound,
2751 DecisionOrigin::Institution {
2752 institution,
2753 responsible_actor,
2754 },
2755 ) => {
2756 binding.institution.as_ref() == Some(institution)
2757 && binding
2758 .actor
2759 .is_none_or(|actor| Some(actor) == *responsible_actor)
2760 }
2761 (SeatPolicy::ObserverSeat | SeatPolicy::AdvisorSeat, origin) => {
2762 let actor_matches = binding.actor.is_none_or(
2763 |expected| matches!(origin, DecisionOrigin::Actor { actor } if *actor == expected),
2764 );
2765 let institution_matches = binding.institution.as_ref().is_none_or(|expected| {
2766 matches!(
2767 origin,
2768 DecisionOrigin::Institution { institution, .. } if institution == expected
2769 )
2770 });
2771 actor_matches && institution_matches
2772 }
2773 _ => false,
2774 }
2775}
2776
2777const fn decision_actor(authority: &CommandAuthority) -> Option<PersonId> {
2778 match &authority.decision_origin {
2779 DecisionOrigin::Actor { actor } => Some(*actor),
2780 DecisionOrigin::Institution {
2781 responsible_actor, ..
2782 } => *responsible_actor,
2783 DecisionOrigin::Council { .. } | DecisionOrigin::NoResponsibleActor { .. } => None,
2784 }
2785}
2786
2787const fn is_expected_command_rejection(code: &ErrorCode) -> bool {
2788 matches!(
2789 code,
2790 ErrorCode::ActorNotFound
2791 | ErrorCode::ArmyNotFound
2792 | ErrorCode::DestinationNotFound
2793 | ErrorCode::EntityNotFound
2794 | ErrorCode::IdempotencyConflict
2795 | ErrorCode::InteractionReadOnly
2796 | ErrorCode::InvalidAuthority
2797 | ErrorCode::InvalidDuration
2798 | ErrorCode::InvalidPayload
2799 | ErrorCode::MissingIdempotencyKey
2800 | ErrorCode::MixedCommandIngress
2801 | ErrorCode::NoRoute
2802 | ErrorCode::PluginCommandNotFound
2803 | ErrorCode::SimulationRevisionConflict
2804 | ErrorCode::SimulationTimeConflict
2805 | ErrorCode::ValueOutOfRange
2806 )
2807}
2808
2809fn canonical_text(value: &str) -> bool {
2810 !value.is_empty() && value == value.trim()
2811}
2812
2813fn component_key(
2814 plugin: &str,
2815 state: &StateKey,
2816 entity: &EntityRef,
2817 component: &str,
2818) -> PluginComponentKey {
2819 PluginComponentKey {
2820 plugin: plugin.to_owned(),
2821 state: state.clone(),
2822 entity: entity.clone(),
2823 component: component.to_owned(),
2824 }
2825}
2826
2827fn record_change_affected_entities(change: &DomainRecordChange) -> Vec<EntityRef> {
2828 (change.current.class == DomainRecordClass::Entity)
2829 .then(|| EntityRef::Domain(change.current.reference.clone()))
2830 .into_iter()
2831 .collect()
2832}
2833
2834fn is_domain_record_state(schemas: &records::DomainRecordSchemas, state: &StateKey) -> bool {
2835 schemas.contains_key(&DomainRecordKind::new(&state.namespace, &state.name))
2836}
2837
2838fn snapshot_command_attempt_preflight_error(
2839 snapshot: &SimulationSnapshot,
2840 attempt: &CommandAttemptRecord,
2841 history: &DomainRecordHistory,
2842 cut: DomainHistoryCut,
2843) -> Option<CanwuError> {
2844 let authority = match resolve_command_authority(&attempt.envelope) {
2845 Ok(authority) => authority,
2846 Err(error) => return Some(error),
2847 };
2848 if let Err(error) = validate_command_ingress_policy(
2849 snapshot
2850 .run_configuration
2851 .as_ref()
2852 .expect("snapshot run configuration is validated before command attempts"),
2853 &attempt.envelope.issuer,
2854 &authority,
2855 CommandAdmission {
2856 request_id: attempt.request_id,
2857 expected_revision: attempt.expected_revision,
2858 expected_time: attempt.envelope.expected_time,
2859 revision_before: attempt.revision_before,
2860 ingress: attempt.ingress,
2861 },
2862 &|entity| snapshot_entity_exists_in_history(snapshot, history, cut, entity),
2863 ) {
2864 return Some(error);
2865 }
2866 attempt.envelope.expected_time.and_then(|expected_time| {
2867 (expected_time != attempt.at).then(|| {
2868 CanwuError::new(
2869 ErrorCode::SimulationTimeConflict,
2870 format!(
2871 "command expected time {expected_time}, but simulation is at {}",
2872 attempt.at
2873 ),
2874 )
2875 })
2876 })
2877}
2878
2879fn invalid_snapshot_error(message: impl Into<String>) -> CanwuError {
2880 CanwuError::new(ErrorCode::InvalidSnapshot, message)
2881}
2882
2883fn invalid_snapshot<T>(message: impl Into<String>) -> Result<T, CanwuError> {
2884 Err(invalid_snapshot_error(message))
2885}
2886
2887#[cfg(test)]
2888mod tests;