Skip to main content

canwu_sim/runtime/
mod.rs

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;
152/// Version of the independently migrated authoritative revision commitment.
153pub const STATE_REVISION_FORMAT_VERSION: u32 = 1;
154/// Version of persisted monotonic boundary-admission cursors.
155pub const ADMISSION_CURSOR_FORMAT_VERSION: u32 = 1;
156/// Version of the domain-separated checkpoint commitment contract.
157pub const COMMITMENT_FORMAT_VERSION: u32 = 1;
158/// Maximum nested depth of the compatibility synchronous event-reactor path.
159///
160/// New plugin mechanics should use phased boundary systems instead of relying
161/// on recursively emitted immediate events.
162pub 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    /// Creates a simulation after validating that scenario references are sound.
782    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    /// Creates a simulation and activates the plugins required by initial
794    /// application-defined records before returning a snapshot-capable runtime.
795    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    /// Creates a simulation with an exact, persisted run environment identity.
805    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    /// Creates a manifested run and activates all initial domain packages before
820    /// the runtime can be observed or snapshotted.
821    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    /// Creates a run whose six policy dimensions are persisted and bound to
837    /// the run-configuration artifact in `run_manifest`.
838    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    /// Creates a declared-policy run and activates all initial domain packages
855    /// before the runtime can be observed or snapshotted.
856    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    /// Returns the persisted authoritative transaction revision.
1158    ///
1159    /// Accepted commands, persisted expected rejections, and completed
1160    /// settlement boundaries each advance it exactly once. Failed work, exact
1161    /// retries, bare clock movement, queued but unadmitted ingress, and plugin
1162    /// setup do not advance it; use the expected-time guard with external
1163    /// commands to detect clock and scheduled-work advancement.
1164    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    /// Hash of simulated state and causal evidence. Run-purpose, controller,
1179    /// seat, observation, interaction, and trace policy remain save identity
1180    /// but are deliberately excluded from this authoritative result identity.
1181    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    /// Returns the persisted audience declaration for a plugin event.
1264    ///
1265    /// Built-in event visibility remains part of the public actor-relative
1266    /// projection. Unlisted plugin event types deliberately resolve to
1267    /// [`EventAudience::Private`].
1268    #[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;