Skip to main content

canwu_sim/runtime/
mod.rs

1mod boundary;
2mod decision;
3mod error;
4mod event_payloads;
5mod hashing;
6mod ingress;
7mod knowledge;
8mod legacy_world;
9mod maintenance;
10mod manifest;
11mod page_store;
12mod persistence;
13mod plugins;
14mod policy;
15mod random;
16mod records;
17mod replay;
18mod revision;
19mod scenario;
20mod scheduling;
21mod settlement;
22mod state;
23mod transactions;
24mod validation;
25mod view;
26
27pub use hashing::{canonical_byte_hash, canonical_hash};
28
29pub use boundary::{
30    BoundaryChange, BoundaryContext, BoundaryDirective, BoundaryEmission, BoundaryEmissionKind,
31    BoundaryIngressGeneration, BoundaryKnowledgeChange, BoundaryProposal, BoundaryReceipt,
32    BoundaryRecord, BoundaryRequest, BoundarySystemContract, BoundarySystemHandler,
33    KnowledgeWriteGrant, OutboxEntry, PluginIngressTarget, RandomDecisionResolution,
34    ReservationAllocation, ReservationDisposition, ReservationOffer, ReservationOfferRecord,
35    ReservationPoolKey, ReservationRef, ReservationRequest, ReservationRequestRecord,
36};
37pub use canwu_core::{
38    DomainRecordVersionRef, DomainRecordVersionSource, EvidenceRef, HolderKnowledgeRecordId,
39    KnowledgeHolderPolicy, KnowledgeHolderRef, KnowledgeRecordId, KnowledgeRecordKind,
40    KnowledgeSchemaId,
41};
42pub use canwu_decision::{
43    ControllerDecision, DECISION_ARCHIVE_BUCKET_PAGE_FORMAT_VERSION,
44    DECISION_ARCHIVE_FORMAT_VERSION, DecisionAction, DecisionArchiveBlob,
45    DecisionArchiveBucketPage, DecisionArchivePageKey, DecisionArchiveProvider,
46    DecisionArchiveReceipt, DecisionArchiveRecord, DecisionArchiveStore,
47    DecisionArchiveStoreOutcome, DecisionAttemptErrorCode, DecisionAttemptOutcome,
48    DecisionAttemptRecord, DecisionAuthority, DecisionContext, DecisionController,
49    DecisionControllerBinding, DecisionError, DecisionErrorCode, DecisionExternalEvidence,
50    DecisionFactorContribution, DecisionHistoryCursor, DecisionHistoryKey, DecisionHistoryLocation,
51    DecisionHistoryPage, DecisionHistoryQueryBudget, DecisionHotState, DecisionLocatorScaleMetrics,
52    DecisionMutation, DecisionOption, DecisionOptionEvaluation, DecisionOptionWeight,
53    DecisionOutcome, DecisionPolicy, DecisionPolicyIdentity, DecisionPolicyKind,
54    DecisionRandomEvidence, DecisionRule, DecisionState, DecisionTicket, DecisionTicketDraft,
55    DecisionTicketState, DecisionTrace, ExternalDecisionOption, ExternalDecisionRequest,
56    ExternalDecisionResponse, ExternalPolicy, HumanDecisionResponse, HumanPolicy, LlmModelIdentity,
57    LlmPolicy, MAX_DECISION_ARCHIVE_BATCH_ENTRIES, MAX_DECISION_HISTORY_PAGE_BYTES,
58    MAX_DECISION_HISTORY_PAGE_SIZE, OrderedRulePolicy, PersistentDecisionLog, PolicyDecision,
59    PreparedDecisionArchive, QueuedExternalPolicy, QueuedHumanPolicy, QueuedLlmPolicy, RuleChoice,
60    RulePolicy, TraceLocatorScaleMetrics, UtilityEvaluator, UtilityPolicy, UtilityProfile,
61    VerifiedDecisionArchiveCommit, WeightedUtilityEvaluator, WeightedUtilityPolicy,
62    format8_decision_locator_scale_probe, format8_trace_locator_scale_probe,
63};
64pub use decision::{
65    DECISION_REQUEST_COMMITMENT_DOMAIN, DecisionEvaluation, DecisionIngressRequest,
66    PreparedDecisionIngress,
67};
68pub use ingress::{
69    IngressClass, IngressPayload, IngressReceipt, IngressRecord, MaintenanceChangeRecord,
70    MaintenanceDisposition, MaintenanceIngressRequest, MaintenanceRejectionReceipt,
71    PluginArchiveRetention, PluginIngressDescriptor, PluginIngressPermit, PluginIngressRequest,
72};
73pub use knowledge::{
74    KnowledgeLimitsV1, KnowledgeSubjectSchema, KnowledgeSubjectTargetKind, PluginKnowledgeSchema,
75};
76pub use maintenance::{
77    MAX_OWNER_AUTHORIZED_MUTATIONS, MAX_OWNER_AUTHORIZED_PARTICIPANTS,
78    OWNER_AUTHORIZED_MAINTENANCE_FORMAT_VERSION, OwnerAuthorizedMaintenanceDraft,
79    OwnerAuthorizedMaintenanceRequest, OwnerAuthorizedMutation, OwnerAuthorizedParticipantDraft,
80    OwnerAuthorizedParticipantProposal, OwnerAuthorizedParticipantRole,
81    OwnerAuthorizedRecordExpectation, VerifiedOwnerAuthorizedMaintenanceCommit,
82};
83pub use manifest::{ArtifactManifest, RUN_MANIFEST_FORMAT_VERSION, RunManifest};
84pub use page_store::{
85    MAX_STATE_DELTA_PAGES, MAX_STATE_PAGE_BYTES, PreparedStateDelta, STATE_PAGE_CODEC,
86    STATE_PAGE_FORMAT_VERSION, STATE_PAGE_RETENTION_FORMAT_VERSION, StatePageBlob,
87    StatePageProvider, StatePageRetentionHandle, StatePageRetentionLedger, StatePageRetentionPhase,
88    StatePageStore, prepare_state_delta, state_page_id, verify_state_delta,
89};
90pub use persistence::{
91    ArchiveProvider, ArchiveReachabilityManifest, ArchiveStore, ArchiveStoreOutcome,
92    ArchivedEvidenceLocator, ArchivedEvidenceReceipt, ArchivedPluginIngressProvenance,
93    ArchivedSegmentHeader, CHECKPOINT_JOURNAL_FORMAT_VERSION, CheckpointJournal,
94    CompactedSimulation, EvidenceArchiveIndex, EvidenceCursor, EvidenceDependency,
95    EvidenceIndexEntry, EvidenceItemLocator, EvidenceJournalKind, EvidenceJournalRoots,
96    EvidenceJournalSegment, EvidenceNestedLocator, EvidenceRequirement, EvidenceSealToken,
97    IDENTITY_EVIDENCE_DEPENDENCIES_FIELD, IDENTITY_EVIDENCE_DEPENDENCIES_FORMAT_VERSION,
98    IdentityEvidenceDependenciesV1, PAGED_CHECKPOINT_FORMAT_VERSION,
99    PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FIELD,
100    PAYLOAD_REQUIRED_EVIDENCE_CONTINUATION_FORMAT_VERSION, PagedCheckpointScaleMetrics,
101    PagedSimulationCheckpoint, PayloadRequiredEvidenceContinuationV1,
102    PortablePagedSimulationCheckpoint, PreparedEvidenceSeal, PreparedPagedSimulationCheckpoint,
103    ReplayJournal, SimulationCheckpoint, SimulationSnapshot, format8_paged_checkpoint_scale_probe,
104    identity_evidence_dependencies_property_v1, payload_required_evidence_continuation_property_v1,
105};
106pub use plugins::{
107    MaintenanceDependencyResolverDescriptor, OwnerAuthorizedMaintenanceParticipant,
108    PluginArchiveObjectProvider, PluginArchiveReachabilityParticipant,
109};
110pub use policy::{
111    CommandPolicyContext, ControllerPolicy, InteractionPolicy, ObservationPolicy,
112    RUN_CONFIGURATION_FORMAT_VERSION, RunConfiguration, RunConfigurationSnapshot, RunPurpose,
113    SeatBinding, SeatPolicy, TracePolicy,
114};
115pub use random::{
116    KeyedDrawReservation, RandomAlgorithm, RandomDrawAddress, RandomDrawOutcome,
117    RandomDrawProducer, RandomDrawRecord, RandomOperationAddressV1, RandomOperationTarget,
118    RandomSample, RandomStreamKey, RandomStreamState,
119};
120pub use records::{
121    DomainRecord, DomainRecordChange, DomainRecordClass, DomainRecordCommitmentRoots,
122    DomainRecordDraft, DomainRecordLifecycle, DomainRecordMutation, DomainRecordMutationPolicy,
123    DomainRecordOperation, DomainRecordPageRoots, DomainRecordSchema, DomainReference,
124    DomainReferenceSchema, DomainReferenceTarget, DomainReferenceTargetKind, PatriciaStoreMetrics,
125    PersistentDomainRecordStore, format8_patricia_scale_probe,
126};
127
128use canwu_core::{
129    ArmyId, BoundaryId, CommandAttemptId, CommandId, CommandRequestId, DecisionRequestId,
130    DecisionTicketId, DecisionTraceId, DeterministicRng, DomainRecordKind, DomainRecordRef,
131    DomainRecordType, EntityRef, EventId, FieldSchema, GovernmentId, IngressId, LetterId, PersonId,
132    RandomDrawId, ResourceId, RouteId, SchemaRegistry, TerritoryId, TypeSchema,
133    TypedDomainRecordRef,
134};
135pub use canwu_event::{CauseRef, EventAudience, EventKind, SimEvent};
136pub use canwu_knowledge::{
137    ActorKnowledge, ArmyKnowledge, EstimateRange, KnowledgeCursor, KnowledgeHistoryView,
138    KnowledgeOrigin, KnowledgeQuery, KnowledgeReadCut, KnowledgeRecord, KnowledgeRecordDraft,
139    KnowledgeRecordView, KnowledgeSnapshot, KnowledgeSource, KnowledgeSubject,
140    KnowledgeSubjectTarget,
141};
142use canwu_time::{SimDuration, SimTime};
143pub use legacy_world::{
144    Army, Government, LetterCargo, LetterStatus, MapPoint, Person, PersonTransitState, Route,
145    Territory, TransitState, WorldSnapshot,
146};
147use serde::{Deserialize, Serialize};
148use serde_json::Value;
149use std::cell::RefCell;
150use std::collections::{BTreeMap, BTreeSet, HashSet};
151use std::panic::{AssertUnwindSafe, catch_unwind};
152use std::rc::Rc;
153
154use event_payloads::{
155    DebugFieldChanged, KNOWLEDGE_PUBLISHED, KnowledgePublished, MoveOrdered, PLUGIN,
156    PersonMoveOrdered, RuntimeEventPayload,
157};
158use hashing::{
159    ControlCommitmentMaterial, StateHashMaterial, authoritative_run_identity,
160    boundary_state_hash_for_commitments, checkpoint_hash_for_commitments,
161    checkpoint_hash_for_configuration, commitment_roots_are_canonical, compute_boundary_hash,
162    decision_commitment_root, identity_commitment_root, is_canonical_hash,
163    knowledge_commitment_root, plugin_component_commitment_root, random_stream_commitment_root,
164    runtime_commitment_roots, scheduler_commitment_root, snapshot_boundary_head_state_hash,
165    snapshot_checkpoint_hash, snapshot_commitment_roots, snapshot_is_at_boundary_head, state_hash,
166    world_commitment_root,
167};
168use ingress::IngressQueueKey;
169use revision::{
170    PersistedAdmissionCursors, authoritative_revision_count, boundaries_before_attempts,
171};
172use settlement::{PendingBoundaryRandomDraw, boundary_has_event_ingress, boundary_system_due};
173use state::{
174    CommitmentDomains, JournalCommitmentRoots, RuntimeCommitmentCache,
175    RuntimeCommitmentRootUpdates, RuntimeCounters, RuntimeCurrentState,
176    RuntimeDomainCommitmentRoots, RuntimeEvidence, RuntimeMetadata, RuntimeScheduler, RuntimeState,
177};
178use transactions::{
179    BoundaryTransactionCheckpoint, ClockTransactionCheckpoint, CommandTransactionCheckpoint,
180    IngressTransactionCheckpoint, RejectionTransactionCheckpoint,
181    ScheduledBatchTransactionCheckpoint,
182};
183use validation::{
184    RuntimeValidationContext, claim_counter, core_world_entity_exists,
185    has_unqueued_command_history, proposal_entity_exists, proposal_entity_identity_exists,
186    runtime_current_entity_exists, runtime_entity_exists,
187    runtime_entity_exists_with_record_overlay, runtime_entity_identity_exists,
188    runtime_has_unqueued_command_history, snapshot_entity_exists_in_history,
189    validate_directives_with_context, validate_domain_dependents_with_records,
190    validate_run_configuration_entities, validate_runtime_cause,
191    validate_runtime_domain_dependents, validate_snapshot,
192};
193
194pub const ENGINE_VERSION: &str = env!("CARGO_PKG_VERSION");
195/// Format 8 binds every decision attempt to its complete ingress request and
196/// activates content-addressed state-page persistence. Older
197/// snapshots, journals, and sub-contract versions are rejected before any
198/// mutable runtime state is constructed.
199pub const SNAPSHOT_FORMAT_VERSION: u32 = 8;
200/// Version of the authoritative revision commitment.
201pub const STATE_REVISION_FORMAT_VERSION: u32 = 3;
202/// Version of persisted monotonic boundary-admission cursors.
203pub const ADMISSION_CURSOR_FORMAT_VERSION: u32 = 3;
204/// Version of the domain-separated checkpoint commitment contract.
205pub const COMMITMENT_FORMAT_VERSION: u32 = 4;
206/// Maximum nested depth of the compatibility synchronous event-reactor path.
207///
208/// New plugin mechanics should use phased boundary systems instead of relying
209/// on recursively emitted immediate events.
210pub const MAX_SYNCHRONOUS_REACTION_DEPTH: usize = 32;
211const CORE_STATE_NAMESPACE: &str = "canwu.core";
212const MAX_DOMAIN_RECORD_QUERY_LIMIT: usize = 10_000;
213const GENESIS_BOUNDARY_HASH: &str =
214    "0000000000000000000000000000000000000000000000000000000000000000";
215
216fn fresh_authority_root_seed(root_seed: u64, run_manifest_hash: &str) -> Result<u64, CanwuError> {
217    let digest =
218        hashing::canonical_hash("canwu.authority-root.v1", &(root_seed, run_manifest_hash))?;
219    u64::from_str_radix(&digest[..16], 16)
220        .map_err(|error| {
221            CanwuError::new(
222                ErrorCode::InvalidAuthority,
223                format!("invalid authority root derivation: {error}"),
224            )
225        })
226        .and_then(|seed| {
227            (seed != 0).then_some(seed).ok_or_else(|| {
228                CanwuError::new(
229                    ErrorCode::InvalidAuthority,
230                    "the derived authority root cannot be zero",
231                )
232            })
233        })
234}
235
236fn reject_unknown_current_fields(
237    input: &Value,
238    encoded: &Value,
239    path: &str,
240) -> Result<(), CanwuError> {
241    match (input, encoded) {
242        (Value::Object(input), Value::Object(encoded)) => {
243            for (key, value) in input {
244                let field_path = if path.is_empty() {
245                    key.clone()
246                } else {
247                    format!("{path}.{key}")
248                };
249                let Some(expected) = encoded.get(key) else {
250                    return Err(invalid_snapshot_error(format!(
251                        "format 8 wire contains unknown field `{field_path}`"
252                    )));
253                };
254                reject_unknown_current_fields(value, expected, &field_path)?;
255            }
256        }
257        (Value::Array(input), Value::Array(encoded)) => {
258            if input.len() != encoded.len() {
259                return Err(invalid_snapshot_error(format!(
260                    "format 8 wire array `{path}` changed shape during decoding"
261                )));
262            }
263            for (index, (value, expected)) in input.iter().zip(encoded).enumerate() {
264                reject_unknown_current_fields(value, expected, &format!("{path}[{index}]"))?;
265            }
266        }
267        _ => {}
268    }
269    Ok(())
270}
271
272fn deserialize_current_json<T>(json: &str, label: &str) -> Result<T, CanwuError>
273where
274    T: for<'de> Deserialize<'de> + Serialize,
275{
276    let input: Value = serde_json::from_str(json).map_err(|error| {
277        invalid_snapshot_error(format!("could not parse format 8 {label}: {error}"))
278    })?;
279    deserialize_current_value(&input, label)
280}
281
282fn deserialize_current_value<T>(input: &Value, label: &str) -> Result<T, CanwuError>
283where
284    T: for<'de> Deserialize<'de> + Serialize,
285{
286    let decoded: T = serde_json::from_value(input.clone()).map_err(|error| {
287        invalid_snapshot_error(format!("could not deserialize format 8 {label}: {error}"))
288    })?;
289    let encoded = serde_json::to_value(&decoded).map_err(|error| {
290        invalid_snapshot_error(format!("could not re-encode format 8 {label}: {error}"))
291    })?;
292    reject_unknown_current_fields(input, &encoded, "")?;
293    Ok(decoded)
294}
295
296fn deserialize_current_snapshot_json(json: &str) -> Result<SimulationSnapshot, CanwuError> {
297    let input: Value = serde_json::from_str(json).map_err(|error| {
298        invalid_snapshot_error(format!("could not parse format 8 snapshot: {error}"))
299    })?;
300    let engine_version = input.get("engine_version").and_then(Value::as_str);
301    let snapshot_format_version = input.get("snapshot_format_version").and_then(Value::as_u64);
302    if engine_version != Some(ENGINE_VERSION)
303        || snapshot_format_version != Some(u64::from(SNAPSHOT_FORMAT_VERSION))
304    {
305        return Err(CanwuError::new(
306            ErrorCode::UnsupportedSnapshotVersion,
307            format!(
308                "the JSON snapshot loader accepts only engine {ENGINE_VERSION} format {SNAPSHOT_FORMAT_VERSION}; pre-8 formats are not supported"
309            ),
310        ));
311    }
312    deserialize_current_value(&input, "snapshot")
313}
314
315fn validate_current_snapshot_contract(snapshot: &SimulationSnapshot) -> Result<(), CanwuError> {
316    if snapshot.commitment_format_version != COMMITMENT_FORMAT_VERSION
317        || snapshot.revision_format_version != STATE_REVISION_FORMAT_VERSION
318        || snapshot.replay_revision_format_version != STATE_REVISION_FORMAT_VERSION
319        || snapshot.admission_cursor_format_version != ADMISSION_CURSOR_FORMAT_VERSION
320        || snapshot.authority_root_seed == 0
321        || snapshot.legacy_rng.is_some()
322    {
323        return Err(invalid_snapshot_error(
324            "format 8 snapshots must use the current commitment, revision, admission, and authority contracts",
325        ));
326    }
327    let Some(run_manifest @ RunManifest::Declared { .. }) = snapshot.run_manifest.as_ref() else {
328        return Err(invalid_snapshot_error(
329            "format 8 snapshots require a declared run manifest",
330        ));
331    };
332    let Some(initial_scenario) = snapshot.initial_scenario.as_ref() else {
333        return Err(invalid_snapshot_error(
334            "format 8 snapshots must retain their canonical initial scenario",
335        ));
336    };
337    if matches!(
338        snapshot.run_configuration,
339        Some(
340            RunConfigurationSnapshot::LegacyUnspecified | RunConfigurationSnapshot::ManifestOnlyV1
341        )
342    ) {
343        return Err(invalid_snapshot_error(
344            "format 8 snapshots cannot use legacy or manifest-only run configuration provenance",
345        ));
346    }
347    manifest::validate(run_manifest, Some(initial_scenario))?;
348    let expected_manifest_hash = manifest::hash(run_manifest)?;
349    if snapshot.run_manifest_hash != expected_manifest_hash {
350        return Err(invalid_snapshot_error(
351            "format 8 snapshot run manifest hash is inconsistent",
352        ));
353    }
354    let run_configuration = snapshot
355        .run_configuration
356        .as_ref()
357        .ok_or_else(|| invalid_snapshot_error("format 8 snapshots require run configuration"))?;
358    let (_, authority_manifest_hash) =
359        authoritative_run_identity(run_manifest, &expected_manifest_hash, run_configuration)?;
360    let expected_authority_root =
361        fresh_authority_root_seed(snapshot.root_seed, &authority_manifest_hash)?;
362    if snapshot.authority_root_seed != expected_authority_root {
363        return Err(invalid_snapshot_error(
364            "format 8 snapshot authority root is not bound to its run identity",
365        ));
366    }
367    Ok(())
368}
369
370/// One deterministic trusted-host page of records from an authoritative read cut.
371#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
372pub struct DomainRecordPage {
373    pub kind: DomainRecordKind,
374    pub revision: u64,
375    pub records: Vec<DomainRecord>,
376    pub next: Option<DomainRecordRef>,
377}
378
379fn validate_domain_record_page_request(
380    kind: &DomainRecordKind,
381    after: Option<&DomainRecordRef>,
382    limit: usize,
383) -> Result<(), CanwuError> {
384    if limit == 0 || limit > MAX_DOMAIN_RECORD_QUERY_LIMIT {
385        return Err(CanwuError::new(
386            ErrorCode::ValueOutOfRange,
387            format!(
388                "domain-record query limit must be between 1 and {MAX_DOMAIN_RECORD_QUERY_LIMIT}"
389            ),
390        ));
391    }
392    if after.is_some_and(|cursor| cursor.kind != *kind) {
393        return Err(CanwuError::new(
394            ErrorCode::InvalidPayload,
395            "domain-record page cursor has the wrong kind",
396        ));
397    }
398    Ok(())
399}
400
401fn domain_record_candidates(
402    records: &impl records::DomainRecordRead,
403    kind: &DomainRecordKind,
404    after: Option<&DomainRecordRef>,
405    limit: usize,
406) -> BTreeMap<DomainRecordRef, DomainRecord> {
407    let (lower, excluded) = after.map_or_else(
408        || {
409            (
410                DomainRecordRef {
411                    kind: kind.clone(),
412                    id: String::new(),
413                },
414                false,
415            )
416        },
417        |cursor| (cursor.clone(), true),
418    );
419    records
420        .range_from(lower, excluded)
421        .take_while(|(reference, _)| reference.kind == *kind)
422        .take(limit)
423        .map(|(reference, record)| (reference.clone(), record.clone()))
424        .collect()
425}
426
427fn retained_domain_record_version(
428    state: &RuntimeState,
429    reference: &DomainRecordVersionRef,
430) -> Option<DomainRecord> {
431    if reference.version == 0 {
432        return None;
433    }
434    let record = match reference.established_by {
435        DomainRecordVersionSource::InitialScenario => state
436            .metadata
437            .initial_scenario
438            .as_ref()?
439            .domain_records
440            .get(
441                *state
442                    .metadata
443                    .initial_domain_record_indexes
444                    .get(&reference.record)?,
445            )
446            .filter(|record| record.version == reference.version),
447        DomainRecordVersionSource::BoundaryChange {
448            boundary,
449            change_index,
450        } => state
451            .evidence
452            .retained_boundary(boundary)?
453            .record_changes
454            .get(usize::try_from(change_index).ok()?)
455            .map(|change| &change.current)
456            .filter(|record| {
457                record.reference == reference.record && record.version == reference.version
458            }),
459    }?;
460    Some(record.clone())
461}
462
463fn current_domain_record_version(
464    state: &RuntimeState,
465    reference: &DomainRecordRef,
466) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
467    let Some(record) = state.current.domain_records.get(reference) else {
468        return Ok(None);
469    };
470    let Some(version) = state.metadata.current_domain_record_versions.get(reference) else {
471        return Err(CanwuError::new(
472            ErrorCode::InvalidSnapshot,
473            "current domain-record provenance index is missing a live record",
474        ));
475    };
476    if version.version != record.version {
477        return Err(CanwuError::new(
478            ErrorCode::InvalidSnapshot,
479            "current domain-record provenance index disagrees with live state",
480        ));
481    }
482    Ok(Some(version.clone()))
483}
484
485fn build_current_domain_record_versions(
486    initial_scenario: Option<&Scenario>,
487    boundaries: &[BoundaryRecord],
488    current_records: &[DomainRecord],
489) -> Result<BTreeMap<DomainRecordRef, DomainRecordVersionRef>, CanwuError> {
490    let mut versions = initial_scenario
491        .into_iter()
492        .flat_map(|scenario| scenario.domain_records.iter())
493        .map(|record| {
494            (
495                record.reference.clone(),
496                DomainRecordVersionRef {
497                    record: record.reference.clone(),
498                    version: record.version,
499                    established_by: DomainRecordVersionSource::InitialScenario,
500                },
501            )
502        })
503        .collect::<BTreeMap<_, _>>();
504    for boundary in boundaries {
505        for (change_index, change) in boundary.record_changes.iter().enumerate() {
506            let change_index = u64::try_from(change_index).map_err(|_| {
507                CanwuError::new(
508                    ErrorCode::IdentifierExhausted,
509                    "domain-record change index exceeds the persistent identifier space",
510                )
511            })?;
512            versions.insert(
513                change.current.reference.clone(),
514                DomainRecordVersionRef {
515                    record: change.current.reference.clone(),
516                    version: change.current.version,
517                    established_by: DomainRecordVersionSource::BoundaryChange {
518                        boundary: boundary.id,
519                        change_index,
520                    },
521                },
522            );
523        }
524    }
525    let current_references = current_records
526        .iter()
527        .map(|record| record.reference.clone())
528        .collect::<BTreeSet<_>>();
529    for record in current_records {
530        let Some(version) = versions.get(&record.reference) else {
531            return Err(CanwuError::new(
532                ErrorCode::InvalidSnapshot,
533                "current domain-record state has no exact provenance index entry",
534            ));
535        };
536        if version.version != record.version {
537            return Err(CanwuError::new(
538                ErrorCode::InvalidSnapshot,
539                "current domain-record provenance index disagrees with snapshot state",
540            ));
541        }
542    }
543    versions.retain(|reference, _| current_references.contains(reference));
544    Ok(versions)
545}
546
547fn retained_evidence_time(state: &RuntimeState, reference: &EvidenceRef) -> Option<SimTime> {
548    match reference {
549        EvidenceRef::Command(id) => state
550            .evidence
551            .retained_command(*id)
552            .map(|record| record.accepted_at),
553        EvidenceRef::CommandAttempt(id) => state
554            .evidence
555            .retained_command_attempt(*id)
556            .map(|record| record.at),
557        EvidenceRef::Event(id) => state
558            .evidence
559            .retained_event(*id)
560            .map(|record| record.timestamp),
561        EvidenceRef::Ingress(id) => state
562            .evidence
563            .retained_ingress(*id)
564            .map(|record| record.issued_at),
565        EvidenceRef::Boundary(id) => state
566            .evidence
567            .retained_boundary(*id)
568            .map(|record| record.at),
569        EvidenceRef::RandomDraw(id) => state
570            .evidence
571            .retained_random_draw(*id)
572            .map(|record| record.at),
573        EvidenceRef::DomainRecordVersion(version) => match version.established_by {
574            DomainRecordVersionSource::InitialScenario => {
575                retained_domain_record_version(state, version).map(|_| state.scheduler.initial_time)
576            }
577            DomainRecordVersionSource::BoundaryChange { boundary, .. } => {
578                retained_domain_record_version(state, version).and_then(|_| {
579                    state
580                        .evidence
581                        .retained_boundary(boundary)
582                        .map(|record| record.at)
583                })
584            }
585        },
586    }
587}
588
589pub use error::{CanwuError, ErrorCode};
590pub use hashing::CommitmentRoots;
591
592use ingress::CommandAdmission;
593pub use ingress::{
594    Command, CommandAttemptOutcome, CommandAttemptRecord, CommandAuthority, CommandContext,
595    CommandEnvelope, CommandIngress, CommandOutcome, CommandReceipt, CommandRecord,
596    CommandRejection, CommandRequest, DecisionOrigin, Issuer,
597};
598
599#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
600#[repr(u8)]
601#[serde(rename_all = "snake_case")]
602pub enum BoundaryPhase {
603    EventIngress = 1,
604    BoundarySnapshot = 2,
605    DerivedFieldSolve = 3,
606    PerceptionAndAttentionRefresh = 4,
607    DecisionAndAcceptedEffectIntake = 5,
608    ReservationAndAllocation = 6,
609    DomainDeltaProposal = 7,
610    InvariantValidation = 8,
611    AtomicDomainCommit = 9,
612    HistoricalCandidateEvaluation = 10,
613    ConditionalTransitionCommit = 11,
614    StrategicAggregation = 12,
615    PerspectiveAndReportMaterialization = 13,
616    SaveReplayAndDiagnosticHashing = 14,
617}
618
619impl BoundaryPhase {
620    pub const ALL: [Self; 14] = [
621        Self::EventIngress,
622        Self::BoundarySnapshot,
623        Self::DerivedFieldSolve,
624        Self::PerceptionAndAttentionRefresh,
625        Self::DecisionAndAcceptedEffectIntake,
626        Self::ReservationAndAllocation,
627        Self::DomainDeltaProposal,
628        Self::InvariantValidation,
629        Self::AtomicDomainCommit,
630        Self::HistoricalCandidateEvaluation,
631        Self::ConditionalTransitionCommit,
632        Self::StrategicAggregation,
633        Self::PerspectiveAndReportMaterialization,
634        Self::SaveReplayAndDiagnosticHashing,
635    ];
636}
637
638#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
639#[serde(rename_all = "snake_case")]
640pub enum SystemCadence {
641    EventDriven,
642    SubDaily,
643    Daily,
644    Monthly,
645    Seasonal,
646    Annual,
647    EraScheduled,
648}
649
650#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
651#[serde(rename_all = "snake_case")]
652pub enum StateVisibility {
653    SameBoundary,
654    NextBoundary,
655}
656
657#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
658pub struct StateKey {
659    pub namespace: String,
660    pub name: String,
661}
662
663impl StateKey {
664    #[must_use]
665    pub fn new(namespace: impl Into<String>, name: impl Into<String>) -> Self {
666        Self {
667            namespace: namespace.into(),
668            name: name.into(),
669        }
670    }
671
672    #[must_use]
673    pub fn core_people() -> Self {
674        Self::new(CORE_STATE_NAMESPACE, "people")
675    }
676
677    #[must_use]
678    pub fn core_governments() -> Self {
679        Self::new(CORE_STATE_NAMESPACE, "governments")
680    }
681
682    #[must_use]
683    pub fn core_territories() -> Self {
684        Self::new(CORE_STATE_NAMESPACE, "territories")
685    }
686
687    #[must_use]
688    pub fn core_routes() -> Self {
689        Self::new(CORE_STATE_NAMESPACE, "routes")
690    }
691
692    #[must_use]
693    pub fn core_armies() -> Self {
694        Self::new(CORE_STATE_NAMESPACE, "armies")
695    }
696
697    #[must_use]
698    pub fn core_knowledge() -> Self {
699        Self::new(CORE_STATE_NAMESPACE, "knowledge")
700    }
701
702    #[must_use]
703    pub fn core_commands() -> Self {
704        Self::new(CORE_STATE_NAMESPACE, "commands")
705    }
706
707    #[must_use]
708    pub fn core_events() -> Self {
709        Self::new(CORE_STATE_NAMESPACE, "events")
710    }
711
712    #[must_use]
713    pub fn core_ingress() -> Self {
714        Self::new(CORE_STATE_NAMESPACE, "ingress")
715    }
716
717    /// Current decision controllers, tickets, and request outcomes.
718    #[must_use]
719    pub fn core_decisions() -> Self {
720        Self::new(CORE_STATE_NAMESPACE, "decisions")
721    }
722
723    /// Administrative read access to current plugin-owned domain records.
724    #[must_use]
725    pub fn core_domain_records() -> Self {
726        Self::new(CORE_STATE_NAMESPACE, "domain_records")
727    }
728
729    /// Retained or archived boundary and random-draw evidence identities.
730    #[must_use]
731    pub fn core_evidence() -> Self {
732        Self::new(CORE_STATE_NAMESPACE, "evidence")
733    }
734}
735
736#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
737pub struct SystemContract {
738    pub name: String,
739    pub phase: BoundaryPhase,
740    pub cadence: SystemCadence,
741    pub reads: Vec<StateKey>,
742    pub writes: Vec<StateKey>,
743    pub visibility: StateVisibility,
744}
745
746impl SystemContract {
747    #[must_use]
748    pub fn event_driven(name: impl Into<String>, phase: BoundaryPhase) -> Self {
749        Self {
750            name: name.into(),
751            phase,
752            cadence: SystemCadence::EventDriven,
753            reads: Vec::new(),
754            writes: Vec::new(),
755            visibility: StateVisibility::SameBoundary,
756        }
757    }
758}
759
760pub use scenario::{DemoIds, Scenario, demo_scenario};
761use scenario::{
762    base_schema, canonicalize_scenario, require_plugin_aware_initial_records, validate_scenario,
763    validate_scenario_state, validate_strict_id_order,
764};
765
766use plugins::PluginComponentKey;
767pub use plugins::{
768    PLUGIN_DESCRIPTOR_FORMAT_VERSION, PayloadProperty, PayloadSchema, PayloadValueType,
769    PluginActionDescriptor, PluginCommandHandler, PluginComponentRecord, PluginDescriptor,
770    PluginRegistrar, PluginRegistry, SimulationPlugin, SimulationSystemHandler, SystemDirective,
771};
772
773pub use view::SimulationView;
774use view::SimulationViewState;
775
776#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
777enum BoundaryWriteStage {
778    Ordinary,
779    Transition,
780    Aggregation,
781    Perspective,
782}
783
784#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
785enum DomainRecordCommitStage {
786    Maintenance,
787    Ordinary,
788    Transition,
789    Aggregation,
790    Perspective,
791    Deferred,
792}
793
794impl DomainRecordCommitStage {
795    const ALL: [Self; 6] = [
796        Self::Maintenance,
797        Self::Ordinary,
798        Self::Transition,
799        Self::Aggregation,
800        Self::Perspective,
801        Self::Deferred,
802    ];
803
804    const fn ordinal(self) -> u8 {
805        match self {
806            Self::Maintenance => 1,
807            Self::Ordinary => 2,
808            Self::Transition => 3,
809            Self::Aggregation => 4,
810            Self::Perspective => 5,
811            Self::Deferred => 6,
812        }
813    }
814}
815
816#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
817struct DomainHistoryCut {
818    boundary: usize,
819    stage: u8,
820}
821
822impl DomainHistoryCut {
823    const GENESIS: Self = Self {
824        boundary: 0,
825        stage: 0,
826    };
827
828    const fn after_boundaries(boundary: usize) -> Self {
829        Self { boundary, stage: 6 }
830    }
831
832    const fn after_stage(boundary: usize, stage: DomainRecordCommitStage) -> Self {
833        Self {
834            boundary,
835            stage: stage.ordinal(),
836        }
837    }
838}
839
840#[derive(Clone, Debug, Default)]
841struct BoundaryDomainEntityCuts {
842    changes: BTreeMap<DomainRecordRef, Vec<DomainEntityStageChange>>,
843}
844
845impl BoundaryDomainEntityCuts {
846    fn record(&mut self, stage: DomainRecordCommitStage, change: &DomainRecordChange) {
847        let previous_live = change
848            .previous
849            .as_ref()
850            .is_some_and(domain_record_is_live_entity);
851        let current_live = domain_record_is_live_entity(&change.current);
852        if previous_live != current_live {
853            self.changes
854                .entry(change.current.reference.clone())
855                .or_default()
856                .push(DomainEntityStageChange {
857                    stage,
858                    plugin: change.plugin.clone(),
859                    system: change.system.clone(),
860                    previous_live,
861                    current_live,
862                });
863        }
864    }
865
866    fn is_live(
867        &self,
868        final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
869        reference: &DomainRecordRef,
870        stage: Option<DomainRecordCommitStage>,
871    ) -> bool {
872        let mut live = final_records
873            .get(reference)
874            .is_some_and(domain_record_is_live_entity);
875        if let Some(changes) = self.changes.get(reference) {
876            for change in changes.iter().rev() {
877                if stage.is_some_and(|stage| change.stage <= stage) {
878                    break;
879                }
880                live = change.previous_live;
881            }
882        }
883        live
884    }
885
886    fn is_live_for_proposal(
887        &self,
888        final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
889        reference: &DomainRecordRef,
890        phase: BoundaryPhase,
891        commit_stage: DomainRecordCommitStage,
892        plugin: &str,
893        system: &str,
894    ) -> bool {
895        let visible_after = match phase {
896            BoundaryPhase::DomainDeltaProposal => None,
897            BoundaryPhase::HistoricalCandidateEvaluation => Some(DomainRecordCommitStage::Ordinary),
898            BoundaryPhase::StrategicAggregation => Some(DomainRecordCommitStage::Transition),
899            BoundaryPhase::PerspectiveAndReportMaterialization => {
900                Some(DomainRecordCommitStage::Aggregation)
901            }
902            BoundaryPhase::EventIngress
903            | BoundaryPhase::BoundarySnapshot
904            | BoundaryPhase::DerivedFieldSolve
905            | BoundaryPhase::PerceptionAndAttentionRefresh
906            | BoundaryPhase::DecisionAndAcceptedEffectIntake
907            | BoundaryPhase::ReservationAndAllocation
908            | BoundaryPhase::InvariantValidation
909            | BoundaryPhase::AtomicDomainCommit
910            | BoundaryPhase::ConditionalTransitionCommit
911            | BoundaryPhase::SaveReplayAndDiagnosticHashing => return false,
912        };
913        let before_proposal = self.is_live(final_records, reference, visible_after);
914        self.changes
915            .get(reference)
916            .and_then(|changes| {
917                changes.iter().find(|change| {
918                    change.stage == commit_stage
919                        && change.plugin == plugin
920                        && change.system == system
921                })
922            })
923            .map_or(before_proposal, |change| change.current_live)
924    }
925
926    fn identity_exists_for_proposal(
927        &self,
928        final_records: &BTreeMap<DomainRecordRef, DomainRecord>,
929        reference: &DomainRecordRef,
930        phase: BoundaryPhase,
931        commit_stage: DomainRecordCommitStage,
932        plugin: &str,
933        system: &str,
934    ) -> bool {
935        if !final_records.contains_key(reference) {
936            return false;
937        }
938        let visible_after = match phase {
939            BoundaryPhase::DomainDeltaProposal => None,
940            BoundaryPhase::HistoricalCandidateEvaluation => Some(DomainRecordCommitStage::Ordinary),
941            BoundaryPhase::StrategicAggregation => Some(DomainRecordCommitStage::Transition),
942            BoundaryPhase::PerspectiveAndReportMaterialization => {
943                Some(DomainRecordCommitStage::Aggregation)
944            }
945            BoundaryPhase::EventIngress
946            | BoundaryPhase::BoundarySnapshot
947            | BoundaryPhase::DerivedFieldSolve
948            | BoundaryPhase::PerceptionAndAttentionRefresh
949            | BoundaryPhase::DecisionAndAcceptedEffectIntake
950            | BoundaryPhase::ReservationAndAllocation
951            | BoundaryPhase::InvariantValidation
952            | BoundaryPhase::AtomicDomainCommit
953            | BoundaryPhase::ConditionalTransitionCommit
954            | BoundaryPhase::SaveReplayAndDiagnosticHashing => return false,
955        };
956        self.changes
957            .get(reference)
958            .and_then(|changes| {
959                changes
960                    .iter()
961                    .find(|change| !change.previous_live && change.current_live)
962            })
963            .is_none_or(|creation| {
964                visible_after.is_some_and(|stage| creation.stage <= stage)
965                    || (creation.stage == commit_stage
966                        && creation.plugin == plugin
967                        && creation.system == system)
968            })
969    }
970}
971
972#[derive(Clone, Debug)]
973struct DomainEntityStageChange {
974    stage: DomainRecordCommitStage,
975    plugin: String,
976    system: String,
977    previous_live: bool,
978    current_live: bool,
979}
980
981#[derive(Clone, Debug)]
982struct DomainRecordHistory {
983    lifetimes: BTreeMap<DomainRecordRef, DomainEntityLifetime>,
984}
985
986impl DomainRecordHistory {
987    fn from_initial_records(records: &BTreeMap<DomainRecordRef, DomainRecord>) -> Self {
988        let lifetimes = records
989            .values()
990            .filter(|record| record.class == DomainRecordClass::Entity)
991            .map(|record| {
992                (
993                    record.reference.clone(),
994                    DomainEntityLifetime {
995                        created_at: DomainHistoryCut::GENESIS,
996                        deleted_at: record.is_deleted().then_some(DomainHistoryCut::GENESIS),
997                    },
998                )
999            })
1000            .collect();
1001        Self { lifetimes }
1002    }
1003
1004    fn apply_boundary(
1005        &mut self,
1006        boundary: usize,
1007        cuts: &BoundaryDomainEntityCuts,
1008    ) -> Result<(), CanwuError> {
1009        for (reference, changes) in &cuts.changes {
1010            for change in changes {
1011                let cut = DomainHistoryCut::after_stage(boundary, change.stage);
1012                match (change.previous_live, change.current_live) {
1013                    (false, true) => {
1014                        if self
1015                            .lifetimes
1016                            .insert(
1017                                reference.clone(),
1018                                DomainEntityLifetime {
1019                                    created_at: cut,
1020                                    deleted_at: None,
1021                                },
1022                            )
1023                            .is_some()
1024                        {
1025                            return invalid_snapshot(
1026                                "domain entity history recreates an existing stable identity",
1027                            );
1028                        }
1029                    }
1030                    (true, false) => {
1031                        let Some(lifetime) = self.lifetimes.get_mut(reference) else {
1032                            return invalid_snapshot(
1033                                "domain entity history deletes an identity before creation",
1034                            );
1035                        };
1036                        if lifetime.deleted_at.replace(cut).is_some() {
1037                            return invalid_snapshot(
1038                                "domain entity history deletes the same identity more than once",
1039                            );
1040                        }
1041                    }
1042                    (false, false) | (true, true) => {}
1043                }
1044            }
1045        }
1046        Ok(())
1047    }
1048
1049    fn is_live(&self, reference: &DomainRecordRef, cut: DomainHistoryCut) -> bool {
1050        self.lifetimes.get(reference).is_some_and(|lifetime| {
1051            lifetime.created_at <= cut && lifetime.deleted_at.is_none_or(|deleted| cut < deleted)
1052        })
1053    }
1054
1055    fn exists(&self, reference: &DomainRecordRef, cut: DomainHistoryCut) -> bool {
1056        self.lifetimes
1057            .get(reference)
1058            .is_some_and(|lifetime| lifetime.created_at <= cut)
1059    }
1060
1061    fn before_time(snapshot: &SimulationSnapshot, at: SimTime) -> DomainHistoryCut {
1062        let count = snapshot
1063            .boundaries
1064            .partition_point(|boundary| boundary.at < at);
1065        DomainHistoryCut::after_boundaries(count)
1066    }
1067}
1068
1069#[derive(Clone, Copy, Debug)]
1070struct DomainEntityLifetime {
1071    created_at: DomainHistoryCut,
1072    deleted_at: Option<DomainHistoryCut>,
1073}
1074
1075fn domain_record_is_live_entity(record: &DomainRecord) -> bool {
1076    record.class == DomainRecordClass::Entity && !record.is_deleted()
1077}
1078
1079const fn boundary_write_stage(phase: BoundaryPhase) -> Option<BoundaryWriteStage> {
1080    match phase {
1081        BoundaryPhase::DomainDeltaProposal => Some(BoundaryWriteStage::Ordinary),
1082        BoundaryPhase::HistoricalCandidateEvaluation => Some(BoundaryWriteStage::Transition),
1083        BoundaryPhase::StrategicAggregation => Some(BoundaryWriteStage::Aggregation),
1084        BoundaryPhase::PerspectiveAndReportMaterialization => Some(BoundaryWriteStage::Perspective),
1085        BoundaryPhase::EventIngress
1086        | BoundaryPhase::BoundarySnapshot
1087        | BoundaryPhase::DerivedFieldSolve
1088        | BoundaryPhase::PerceptionAndAttentionRefresh
1089        | BoundaryPhase::DecisionAndAcceptedEffectIntake
1090        | BoundaryPhase::ReservationAndAllocation
1091        | BoundaryPhase::InvariantValidation
1092        | BoundaryPhase::AtomicDomainCommit
1093        | BoundaryPhase::ConditionalTransitionCommit
1094        | BoundaryPhase::SaveReplayAndDiagnosticHashing => None,
1095    }
1096}
1097
1098const fn domain_record_commit_stage(
1099    phase: BoundaryPhase,
1100    visibility: StateVisibility,
1101) -> Option<DomainRecordCommitStage> {
1102    let stage = match phase {
1103        BoundaryPhase::DomainDeltaProposal => DomainRecordCommitStage::Ordinary,
1104        BoundaryPhase::HistoricalCandidateEvaluation => DomainRecordCommitStage::Transition,
1105        BoundaryPhase::StrategicAggregation => DomainRecordCommitStage::Aggregation,
1106        BoundaryPhase::PerspectiveAndReportMaterialization => DomainRecordCommitStage::Perspective,
1107        BoundaryPhase::EventIngress
1108        | BoundaryPhase::BoundarySnapshot
1109        | BoundaryPhase::DerivedFieldSolve
1110        | BoundaryPhase::PerceptionAndAttentionRefresh
1111        | BoundaryPhase::DecisionAndAcceptedEffectIntake
1112        | BoundaryPhase::ReservationAndAllocation
1113        | BoundaryPhase::InvariantValidation
1114        | BoundaryPhase::AtomicDomainCommit
1115        | BoundaryPhase::ConditionalTransitionCommit
1116        | BoundaryPhase::SaveReplayAndDiagnosticHashing => return None,
1117    };
1118    Some(match visibility {
1119        StateVisibility::SameBoundary => stage,
1120        StateVisibility::NextBoundary => DomainRecordCommitStage::Deferred,
1121    })
1122}
1123
1124fn validate_type_schema(schema: &TypeSchema) -> Result<(), CanwuError> {
1125    if schema.type_name.trim().is_empty() || schema.type_name != schema.type_name.trim() {
1126        return Err(CanwuError::new(
1127            ErrorCode::InvalidPluginRegistration,
1128            "plugin schema type name must be non-empty and have no surrounding whitespace",
1129        ));
1130    }
1131    let mut field_names = BTreeSet::new();
1132    for field in &schema.fields {
1133        if field.name.trim().is_empty()
1134            || field.name != field.name.trim()
1135            || field.value_type.trim().is_empty()
1136            || field.value_type != field.value_type.trim()
1137            || field
1138                .reference_type
1139                .as_ref()
1140                .is_some_and(|value| value.trim().is_empty() || value != value.trim())
1141            || !field_names.insert(&field.name)
1142        {
1143            return Err(CanwuError::new(
1144                ErrorCode::InvalidPluginRegistration,
1145                format!("schema {} contains an invalid field", schema.type_name),
1146            ));
1147        }
1148    }
1149    Ok(())
1150}
1151
1152use scheduling::{ScheduleKey, ScheduledAction, ScheduledRecord};
1153
1154const fn one_u64() -> u64 {
1155    1
1156}
1157
1158#[allow(clippy::trivially_copy_pass_by_ref)]
1159const fn is_zero_u64(value: &u64) -> bool {
1160    *value == 0
1161}
1162
1163#[allow(clippy::trivially_copy_pass_by_ref)]
1164const fn is_zero_u32(value: &u32) -> bool {
1165    *value == 0
1166}
1167
1168#[allow(clippy::trivially_copy_pass_by_ref)]
1169const fn is_one_u64(value: &u64) -> bool {
1170    *value == 1
1171}
1172
1173fn command_attempt_slice_is_empty(value: &&[CommandAttemptRecord]) -> bool {
1174    value.is_empty()
1175}
1176
1177fn command_attempt_id_slice_is_empty(value: &&[CommandAttemptId]) -> bool {
1178    value.is_empty()
1179}
1180
1181fn domain_record_slice_is_empty(value: &&[DomainRecord]) -> bool {
1182    value.is_empty()
1183}
1184
1185fn domain_record_change_slice_is_empty(value: &&[DomainRecordChange]) -> bool {
1186    value.is_empty()
1187}
1188
1189fn ingress_record_slice_is_empty(value: &&[IngressRecord]) -> bool {
1190    value.is_empty()
1191}
1192
1193fn maintenance_change_slice_is_empty(value: &&[MaintenanceChangeRecord]) -> bool {
1194    value.is_empty()
1195}
1196
1197const BOUNDARY_STATE_HASH_V1_PREFIX: &str = "v1:";
1198
1199#[derive(Clone, Copy, Debug, Eq, PartialEq)]
1200enum BoundaryStateHashFormat {
1201    LegacyV0,
1202    CommitmentsV1,
1203}
1204
1205fn boundary_state_hash_format(value: Option<&str>) -> Result<BoundaryStateHashFormat, CanwuError> {
1206    match value {
1207        Some(value) if value.starts_with(BOUNDARY_STATE_HASH_V1_PREFIX) => {
1208            let hash = &value[BOUNDARY_STATE_HASH_V1_PREFIX.len()..];
1209            if !is_canonical_hash(hash) {
1210                return invalid_snapshot("boundary state commitment v1 is not canonical");
1211            }
1212            Ok(BoundaryStateHashFormat::CommitmentsV1)
1213        }
1214        Some(value) if is_canonical_hash(value) => Ok(BoundaryStateHashFormat::LegacyV0),
1215        Some(_) => invalid_snapshot("boundary state commitment format is unsupported"),
1216        None => Ok(BoundaryStateHashFormat::LegacyV0),
1217    }
1218}
1219
1220pub struct Simulation {
1221    state: RuntimeState,
1222    schema: SchemaRegistry,
1223    plugins: PluginRegistry,
1224    plugin_archive_provider: Rc<dyn PluginArchiveObjectProvider>,
1225    sync_reaction_depth: usize,
1226}
1227
1228impl Simulation {
1229    /// Creates a simulation after validating that scenario references are sound.
1230    pub fn new(seed: u64, scenario: Scenario) -> Result<Self, CanwuError> {
1231        require_plugin_aware_initial_records(&scenario)?;
1232        let run_manifest = RunManifest::for_scenario("canwu.inline", "scenario", "1", &scenario)?;
1233        Self::new_with_configuration_snapshot(
1234            seed,
1235            scenario,
1236            run_manifest,
1237            RunConfigurationSnapshot::CompatibilityV1,
1238        )
1239    }
1240
1241    /// Creates a simulation and activates the plugins required by initial
1242    /// application-defined records before returning a snapshot-capable runtime.
1243    pub fn new_with_plugins(
1244        seed: u64,
1245        scenario: Scenario,
1246        plugins: &[&dyn SimulationPlugin],
1247    ) -> Result<Self, CanwuError> {
1248        let run_manifest = RunManifest::for_scenario("canwu.inline", "scenario", "1", &scenario)?;
1249        Self::new_with_manifest_and_plugins(seed, scenario, run_manifest, plugins)
1250    }
1251
1252    /// Creates a simulation with an exact, persisted run environment identity.
1253    pub fn new_with_manifest(
1254        seed: u64,
1255        scenario: Scenario,
1256        run_manifest: RunManifest,
1257    ) -> Result<Self, CanwuError> {
1258        require_plugin_aware_initial_records(&scenario)?;
1259        Self::new_with_configuration_snapshot(
1260            seed,
1261            scenario,
1262            run_manifest,
1263            RunConfigurationSnapshot::CompatibilityV1,
1264        )
1265    }
1266
1267    /// Creates a manifested run and activates all initial domain packages before
1268    /// the runtime can be observed or snapshotted.
1269    pub fn new_with_manifest_and_plugins(
1270        seed: u64,
1271        scenario: Scenario,
1272        run_manifest: RunManifest,
1273        plugins: &[&dyn SimulationPlugin],
1274    ) -> Result<Self, CanwuError> {
1275        let simulation = Self::new_with_configuration_snapshot(
1276            seed,
1277            scenario,
1278            run_manifest,
1279            RunConfigurationSnapshot::CompatibilityV1,
1280        )?;
1281        Self::activate_initial_plugins(simulation, plugins)
1282    }
1283
1284    /// Creates a run whose six policy dimensions are persisted and bound to
1285    /// the run-configuration artifact in `run_manifest`.
1286    pub fn new_with_run_configuration(
1287        seed: u64,
1288        scenario: Scenario,
1289        run_manifest: RunManifest,
1290        mut run_configuration: RunConfiguration,
1291    ) -> Result<Self, CanwuError> {
1292        require_plugin_aware_initial_records(&scenario)?;
1293        run_configuration.canonicalize();
1294        Self::new_with_configuration_snapshot(
1295            seed,
1296            scenario,
1297            run_manifest,
1298            RunConfigurationSnapshot::Declared(run_configuration),
1299        )
1300    }
1301
1302    /// Creates a declared-policy run and activates all initial domain packages
1303    /// before the runtime can be observed or snapshotted.
1304    pub fn new_with_run_configuration_and_plugins(
1305        seed: u64,
1306        scenario: Scenario,
1307        run_manifest: RunManifest,
1308        mut run_configuration: RunConfiguration,
1309        plugins: &[&dyn SimulationPlugin],
1310    ) -> Result<Self, CanwuError> {
1311        run_configuration.canonicalize();
1312        let simulation = Self::new_with_configuration_snapshot(
1313            seed,
1314            scenario,
1315            run_manifest,
1316            RunConfigurationSnapshot::Declared(run_configuration),
1317        )?;
1318        Self::activate_initial_plugins(simulation, plugins)
1319    }
1320
1321    fn activate_initial_plugins(
1322        mut simulation: Self,
1323        plugins: &[&dyn SimulationPlugin],
1324    ) -> Result<Self, CanwuError> {
1325        for plugin in plugins {
1326            simulation.register_plugin(*plugin)?;
1327        }
1328        simulation.ensure_runtime_ready()?;
1329        Ok(simulation)
1330    }
1331
1332    fn new_with_configuration_snapshot(
1333        seed: u64,
1334        mut scenario: Scenario,
1335        mut run_manifest: RunManifest,
1336        run_configuration: RunConfigurationSnapshot,
1337    ) -> Result<Self, CanwuError> {
1338        canonicalize_scenario(&mut scenario);
1339        validate_scenario(&scenario)?;
1340        manifest::canonicalize(&mut run_manifest);
1341        manifest::validate(&run_manifest, Some(&scenario))?;
1342        manifest::validate_run_configuration(&run_manifest, &run_configuration)?;
1343        validate_run_configuration_entities(
1344            &run_configuration,
1345            &scenario.entities,
1346            &scenario.world,
1347            &scenario.domain_records,
1348        )?;
1349        let run_manifest_hash = manifest::hash(&run_manifest)?;
1350        if scenario
1351            .world
1352            .armies
1353            .iter()
1354            .any(|army| army.transit.is_some())
1355        {
1356            return Err(CanwuError::new(
1357                ErrorCode::InvalidSnapshot,
1358                "initial scenarios cannot contain transit without admitted command/event/queue evidence",
1359            ));
1360        }
1361        if scenario
1362            .world
1363            .people
1364            .iter()
1365            .any(|person| person.transit.is_some())
1366            || scenario
1367                .world
1368                .letters
1369                .iter()
1370                .any(|letter| letter.status == LetterStatus::InTransit)
1371        {
1372            return Err(CanwuError::new(
1373                ErrorCode::InvalidSnapshot,
1374                "initial scenarios cannot contain person or letter transit without admitted command/event/queue evidence",
1375            ));
1376        }
1377        let schema = base_schema();
1378        let plugins = PluginRegistry::default();
1379        let (_, authority_manifest_hash) =
1380            authoritative_run_identity(&run_manifest, &run_manifest_hash, &run_configuration)?;
1381        let authority_root_seed = fresh_authority_root_seed(seed, &authority_manifest_hash)?;
1382        let core_stream = RandomStreamState::initial(seed, random::core_report_delay_stream());
1383        let initial_scenario = Some(scenario.clone());
1384        let initial_domain_record_indexes = initial_scenario
1385            .as_ref()
1386            .map(|scenario| {
1387                scenario
1388                    .domain_records
1389                    .iter()
1390                    .enumerate()
1391                    .map(|(index, record)| (record.reference.clone(), index))
1392                    .collect()
1393            })
1394            .unwrap_or_default();
1395        let current_domain_record_versions = build_current_domain_record_versions(
1396            initial_scenario.as_ref(),
1397            &[],
1398            &scenario.domain_records,
1399        )?;
1400        let mut simulation = Self {
1401            state: RuntimeState {
1402                current: RuntimeCurrentState {
1403                    entities: scenario.entities.into_iter().collect(),
1404                    people: scenario
1405                        .world
1406                        .people
1407                        .into_iter()
1408                        .map(|value| (value.id, value))
1409                        .collect(),
1410                    letters: scenario
1411                        .world
1412                        .letters
1413                        .into_iter()
1414                        .map(|value| (value.id, value))
1415                        .collect(),
1416                    governments: scenario
1417                        .world
1418                        .governments
1419                        .into_iter()
1420                        .map(|value| (value.id, value))
1421                        .collect(),
1422                    territories: scenario
1423                        .world
1424                        .territories
1425                        .into_iter()
1426                        .map(|value| (value.id, value))
1427                        .collect(),
1428                    routes: scenario
1429                        .world
1430                        .routes
1431                        .into_iter()
1432                        .map(|value| (value.id, value))
1433                        .collect(),
1434                    armies: scenario
1435                        .world
1436                        .armies
1437                        .into_iter()
1438                        .map(|value| (value.id, value))
1439                        .collect(),
1440                    knowledge: scenario.knowledge,
1441                    plugin_components: BTreeMap::new(),
1442                    domain_records: PersistentDomainRecordStore::from_records(
1443                        scenario
1444                            .domain_records
1445                            .into_iter()
1446                            .map(|record| (record.reference.clone(), record))
1447                            .collect(),
1448                    )?,
1449                    decisions: DecisionState::default(),
1450                    root_seed: seed,
1451                    authority_root_seed,
1452                    random_streams: BTreeMap::from([(core_stream.key.clone(), core_stream)]),
1453                },
1454                scheduler: RuntimeScheduler {
1455                    initial_time: scenario.start_time,
1456                    now: scenario.start_time,
1457                    actions: BTreeMap::new(),
1458                    pending_ingress: BTreeSet::new(),
1459                },
1460                counters: RuntimeCounters {
1461                    next_event_id: 1,
1462                    next_command_id: 1,
1463                    next_command_attempt_id: 1,
1464                    next_ingress_id: 1,
1465                    next_boundary_id: 1,
1466                    next_random_draw_id: 1,
1467                    next_knowledge_record_id: 1,
1468                    next_schedule_sequence: 1,
1469                    next_correlation_id: 1,
1470                    next_decision_trace_id: 1,
1471                    state_revision: 0,
1472                    admitted_attempt_count: 0,
1473                    admitted_command_count: 0,
1474                    admitted_event_count: 0,
1475                },
1476                metadata: RuntimeMetadata {
1477                    initial_scenario,
1478                    initial_domain_record_indexes,
1479                    current_domain_record_versions,
1480                    run_manifest,
1481                    run_manifest_hash,
1482                    run_configuration,
1483                    checkpoint_hash: String::new(),
1484                    commitment_format_version: COMMITMENT_FORMAT_VERSION,
1485                    commitment_roots: None,
1486                    commitment_cache: None,
1487                    plugin_registration_closed: false,
1488                    replay_revision_format_version: STATE_REVISION_FORMAT_VERSION,
1489                },
1490                evidence: RuntimeEvidence {
1491                    archived: EvidenceCursor::default(),
1492                    archived_boundary_head: None,
1493                    archived_legacy_commands: false,
1494                    archived_tracked_attempts: false,
1495                    archived_unqueued_command_history: false,
1496                    archived_command_requests: BTreeMap::new(),
1497                    archived_ingress_requests: BTreeMap::new(),
1498                    archived_decision_requests: BTreeMap::new(),
1499                    archived_decision_command_requests: BTreeSet::new(),
1500                    events: Vec::new(),
1501                    commands: Vec::new(),
1502                    command_attempts: Vec::new(),
1503                    ingress: Vec::new(),
1504                    boundaries: Vec::new(),
1505                    random_draws: Vec::new(),
1506                    archived_segment_headers: Vec::new(),
1507                    archived_evidence_receipts: BTreeMap::new(),
1508                    keyed_draw_reservations: Vec::new(),
1509                },
1510            },
1511            schema,
1512            plugins,
1513            plugin_archive_provider: Rc::new(()),
1514            sync_reaction_depth: 0,
1515        };
1516        simulation.refresh_checkpoint_hash()?;
1517        Ok(simulation)
1518    }
1519
1520    pub fn demo(seed: u64) -> Result<(Self, DemoIds), CanwuError> {
1521        let (scenario, ids) = demo_scenario();
1522        Self::new(seed, scenario).map(|simulation| (simulation, ids))
1523    }
1524
1525    pub fn register_plugin<P: SimulationPlugin + ?Sized>(
1526        &mut self,
1527        plugin: &P,
1528    ) -> Result<(), CanwuError> {
1529        let plugin_name = plugin.name().trim();
1530        if plugin_name.is_empty() || plugin_name != plugin.name() {
1531            return Err(CanwuError::new(
1532                ErrorCode::InvalidPluginRegistration,
1533                "plugin name must be non-empty and have no surrounding whitespace",
1534            ));
1535        }
1536        let rehydrating = self.plugins.descriptors.contains_key(plugin_name)
1537            && !self.plugins.active_plugins.contains(plugin_name);
1538        if self.state.metadata.plugin_registration_closed && !rehydrating {
1539            return Err(CanwuError::new(
1540                ErrorCode::PluginRegistrationClosed,
1541                "new plugins must be registered before authoritative execution begins",
1542            ));
1543        }
1544        let state_start = self.state.clone();
1545        let schema_start = self.schema.clone();
1546        let plugins_start = self.plugins.clone();
1547        let result = (|| {
1548            self.plugins.register(plugin, &mut self.schema)?;
1549            self.invalidate_commitments(
1550                CommitmentDomains::RANDOM_STREAMS | CommitmentDomains::IDENTITY,
1551            );
1552            if !self.plugins.record_schemas.is_empty()
1553                && self.state.metadata.initial_scenario.is_none()
1554            {
1555                return Err(CanwuError::new(
1556                    ErrorCode::UnsupportedSnapshotVersion,
1557                    "this snapshot predates manifest-bound domain-record genesis and cannot activate record schemas",
1558                ));
1559            }
1560            records::validate_records_for_owner(
1561                &self.state.current.domain_records,
1562                &self.plugins.record_schemas,
1563                plugin_name,
1564                self.state.scheduler.now,
1565                &|entity| runtime_entity_exists(&self.state, entity),
1566            )?;
1567            let activation_records = self
1568                .state
1569                .current
1570                .domain_records
1571                .values()
1572                .filter(|record| record.owner == plugin_name)
1573                .cloned()
1574                .collect::<Vec<_>>();
1575            plugin.validate_activation(&activation_records)?;
1576            for stream in self.plugins.random_stream_owners.keys() {
1577                self.state
1578                    .current
1579                    .random_streams
1580                    .entry(stream.clone())
1581                    .or_insert_with(|| {
1582                        RandomStreamState::initial(self.state.current.root_seed, stream.clone())
1583                    });
1584            }
1585            self.refresh_checkpoint_hash()
1586        })();
1587        if let Err(error) = result {
1588            self.state = state_start;
1589            self.schema = schema_start;
1590            self.plugins = plugins_start;
1591            return Err(error);
1592        }
1593        Ok(())
1594    }
1595
1596    fn ensure_runtime_ready(&self) -> Result<(), CanwuError> {
1597        // Initial construction, snapshot restore, and plugin activation perform
1598        // the complete domain-record audit. Thereafter all public mutations go
1599        // through affected-closure validation on the persistent store. Repeating
1600        // the cold audit here would deserialize every untouched plugin payload
1601        // before every ingress and boundary, defeating Format-8 shard isolation.
1602        self.plugins.ensure_active()
1603    }
1604
1605    fn bound_initial_scenario(&self) -> Option<&Scenario> {
1606        self.state.metadata.initial_scenario.as_ref()
1607    }
1608
1609    #[must_use]
1610    pub const fn time(&self) -> SimTime {
1611        self.state.scheduler.now
1612    }
1613
1614    #[must_use]
1615    pub const fn run_manifest(&self) -> &RunManifest {
1616        &self.state.metadata.run_manifest
1617    }
1618
1619    #[must_use]
1620    pub const fn run_configuration(&self) -> &RunConfigurationSnapshot {
1621        &self.state.metadata.run_configuration
1622    }
1623
1624    #[must_use]
1625    /// Returns the persisted authoritative transaction revision.
1626    ///
1627    /// Accepted commands, persisted expected rejections, and completed
1628    /// settlement boundaries each advance it exactly once. Failed work, exact
1629    /// retries, bare clock movement, queued but unadmitted ingress, and plugin
1630    /// setup do not advance it; use the expected-time guard with external
1631    /// commands to detect clock and scheduled-work advancement.
1632    pub const fn revision(&self) -> u64 {
1633        self.state.counters.state_revision
1634    }
1635
1636    #[must_use]
1637    pub fn run_manifest_hash(&self) -> &str {
1638        &self.state.metadata.run_manifest_hash
1639    }
1640
1641    #[must_use]
1642    pub fn checkpoint_hash(&self) -> &str {
1643        &self.state.metadata.checkpoint_hash
1644    }
1645
1646    /// Hash of simulated state and causal evidence. Run-purpose, controller,
1647    /// seat, observation, interaction, and trace policy remain save identity
1648    /// but are deliberately excluded from this authoritative result identity.
1649    pub fn authoritative_state_hash(&self) -> Result<String, CanwuError> {
1650        self.compute_boundary_state_hash()
1651    }
1652
1653    pub fn entities(&self) -> impl Iterator<Item = &EntityRef> {
1654        self.state.current.entities.iter()
1655    }
1656
1657    #[must_use]
1658    pub fn entity_exists(&self, entity: &EntityRef) -> bool {
1659        runtime_entity_exists(&self.state, entity)
1660    }
1661
1662    #[must_use]
1663    pub fn world(&self) -> WorldSnapshot {
1664        WorldSnapshot {
1665            people: self.state.current.people.values().cloned().collect(),
1666            governments: self.state.current.governments.values().cloned().collect(),
1667            territories: self.state.current.territories.values().cloned().collect(),
1668            routes: self.state.current.routes.values().cloned().collect(),
1669            armies: self.state.current.armies.values().cloned().collect(),
1670            letters: self.state.current.letters.values().cloned().collect(),
1671        }
1672    }
1673
1674    #[must_use]
1675    pub fn knowledge(&self) -> &KnowledgeSnapshot {
1676        &self.state.current.knowledge
1677    }
1678
1679    #[must_use]
1680    pub fn events(&self) -> &[SimEvent] {
1681        &self.state.evidence.events
1682    }
1683
1684    #[must_use]
1685    pub fn command_log(&self) -> &[CommandRecord] {
1686        &self.state.evidence.commands
1687    }
1688
1689    #[must_use]
1690    pub fn command_attempts(&self) -> &[CommandAttemptRecord] {
1691        &self.state.evidence.command_attempts
1692    }
1693
1694    #[must_use]
1695    pub fn ingress_log(&self) -> &[IngressRecord] {
1696        &self.state.evidence.ingress
1697    }
1698
1699    #[must_use]
1700    pub fn domain_record(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
1701        self.state.current.domain_records.get(reference)
1702    }
1703
1704    /// Returns whether an exact domain-record version exists in current or retained evidence.
1705    #[must_use]
1706    pub fn domain_record_version_evidence_exists(
1707        &self,
1708        reference: &DomainRecordVersionRef,
1709    ) -> bool {
1710        !matches!(
1711            validation::resolve_evidence_reference(
1712                &validation::RuntimeValidationContext::new(&self.state),
1713                &EvidenceRef::DomainRecordVersion(reference.clone()),
1714            ),
1715            validation::EvidenceAvailability::Missing
1716        )
1717    }
1718
1719    /// Returns whether a generic evidence identity is retained or archived.
1720    #[must_use]
1721    pub fn evidence_exists(&self, reference: &EvidenceRef) -> bool {
1722        !matches!(
1723            validation::resolve_evidence_reference(
1724                &validation::RuntimeValidationContext::new(&self.state),
1725                reference,
1726            ),
1727            validation::EvidenceAvailability::Missing
1728        )
1729    }
1730
1731    /// Returns when retained evidence first became authoritative.
1732    ///
1733    /// `None` means the evidence is missing or only its compact archive
1734    /// receipt remains. Proposed same-boundary evidence is available through
1735    /// [`SimulationView::evidence_time`] while its boundary is being built.
1736    #[must_use]
1737    pub fn evidence_time(&self, reference: &EvidenceRef) -> Option<SimTime> {
1738        retained_evidence_time(&self.state, reference)
1739    }
1740
1741    /// Resolves the retained record body for one exact domain-record version.
1742    ///
1743    /// Returns `None` when the version is unavailable or only its compacted
1744    /// archive receipt remains.
1745    #[must_use]
1746    pub fn domain_record_version(
1747        &self,
1748        reference: &DomainRecordVersionRef,
1749    ) -> Option<DomainRecord> {
1750        retained_domain_record_version(&self.state, reference)
1751    }
1752
1753    /// Returns the exact evidence identity for the authoritative current
1754    /// version without manufacturing an initial-scenario source after
1755    /// compaction.
1756    pub fn current_domain_record_version(
1757        &self,
1758        reference: &DomainRecordRef,
1759    ) -> Result<Option<DomainRecordVersionRef>, CanwuError> {
1760        current_domain_record_version(&self.state, reference)
1761    }
1762
1763    #[must_use]
1764    pub fn typed_domain_record<T: DomainRecordType>(
1765        &self,
1766        reference: &TypedDomainRecordRef<T>,
1767    ) -> Option<&DomainRecord> {
1768        self.domain_record(reference.as_untyped())
1769    }
1770
1771    pub fn domain_records(&self) -> impl Iterator<Item = &DomainRecord> {
1772        self.state.current.domain_records.values()
1773    }
1774
1775    /// Returns one revision-bound page of authoritative domain records.
1776    ///
1777    /// Pass the revision returned by the first page on every subsequent page.
1778    /// A mutation between pages is rejected instead of mixing two read cuts.
1779    pub fn domain_record_page(
1780        &self,
1781        kind: &DomainRecordKind,
1782        after: Option<&DomainRecordRef>,
1783        limit: usize,
1784        expected_revision: Option<u64>,
1785    ) -> Result<DomainRecordPage, CanwuError> {
1786        validate_domain_record_page_request(kind, after, limit)?;
1787        let revision = self.revision();
1788        if expected_revision.is_some_and(|expected| expected != revision) {
1789            return Err(CanwuError::new(
1790                ErrorCode::SimulationRevisionConflict,
1791                format!(
1792                    "domain-record page expected revision {expected_revision:?}, current revision is {revision}"
1793                ),
1794            ));
1795        }
1796        let requested = limit.checked_add(1).unwrap_or(limit);
1797        let mut records =
1798            domain_record_candidates(&self.state.current.domain_records, kind, after, requested)
1799                .into_values()
1800                .collect::<Vec<_>>();
1801        let has_more = records.len() > limit;
1802        records.truncate(limit);
1803        let next = has_more
1804            .then(|| records.last().map(|record| record.reference.clone()))
1805            .flatten();
1806        Ok(DomainRecordPage {
1807            kind: kind.clone(),
1808            revision,
1809            records,
1810            next,
1811        })
1812    }
1813
1814    #[must_use]
1815    pub fn boundaries(&self) -> &[BoundaryRecord] {
1816        &self.state.evidence.boundaries
1817    }
1818
1819    #[must_use]
1820    pub fn random_draws(&self) -> &[RandomDrawRecord] {
1821        &self.state.evidence.random_draws
1822    }
1823
1824    #[must_use]
1825    pub fn boundary_head_hash(&self) -> Option<&str> {
1826        self.state.evidence.boundary_head_hash()
1827    }
1828
1829    #[must_use]
1830    pub const fn schema(&self) -> &SchemaRegistry {
1831        &self.schema
1832    }
1833
1834    pub fn plugin_descriptors(&self) -> impl Iterator<Item = &PluginDescriptor> {
1835        self.plugins.descriptors()
1836    }
1837
1838    /// Returns the persisted audience declaration for a plugin event.
1839    ///
1840    /// Built-in event visibility remains part of the public actor-relative
1841    /// projection. Unlisted plugin event types deliberately resolve to
1842    /// [`EventAudience::Private`].
1843    #[must_use]
1844    pub fn event_audience(&self, event: &SimEvent) -> EventAudience {
1845        match event.kind.event_type() {
1846            PLUGIN => event
1847                .kind
1848                .plugin_identity()
1849                .map_or(EventAudience::Private, |(plugin, event_type)| {
1850                    self.plugins.event_audience(plugin, event_type)
1851                }),
1852            KNOWLEDGE_PUBLISHED => KnowledgePublished::decode(&event.kind)
1853                .map_or(EventAudience::Private, |payload| {
1854                    EventAudience::KnowledgeHolder(payload.holder)
1855                }),
1856            _ => EventAudience::Private,
1857        }
1858    }
1859
1860    ///
1861    /// # Panics
1862    ///
1863    /// Panics only if a runtime object was constructed without its required
1864    /// Format 8 initial scenario, which is prevented by the public loaders.
1865    #[must_use]
1866    pub fn replay_journal(&self) -> ReplayJournal {
1867        ReplayJournal {
1868            engine_version: ENGINE_VERSION.to_owned(),
1869            snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
1870            root_seed: self.state.current.root_seed,
1871            initial_scenario: self
1872                .state
1873                .metadata
1874                .initial_scenario
1875                .clone()
1876                .expect("Format 8 runs always retain their initial scenario"),
1877            authority_root_seed: self.state.current.authority_root_seed,
1878            run_manifest: self.state.metadata.run_manifest.clone(),
1879            run_manifest_hash: self.state.metadata.run_manifest_hash.clone(),
1880            run_configuration: self.state.metadata.run_configuration.clone(),
1881            plugin_descriptors: self.plugins.descriptors().cloned().collect(),
1882            plugin_registration_closed: self.state.metadata.plugin_registration_closed,
1883            commands: self.state.evidence.commands.clone(),
1884            command_attempts: self.state.evidence.command_attempts.clone(),
1885            ingress: self.state.evidence.ingress.clone(),
1886            boundaries: self.state.evidence.boundaries.clone(),
1887            final_time: self.state.scheduler.now,
1888            checkpoint_hash: self.state.metadata.checkpoint_hash.clone(),
1889            commitment_format_version: self.state.metadata.commitment_format_version,
1890            revision_format_version: self.state.metadata.replay_revision_format_version,
1891            final_revision: self.state.counters.state_revision,
1892        }
1893    }
1894
1895    /// Returns the durable external-delivery outbox derived from committed
1896    /// boundary emissions. Entries are deterministic and never re-sent by
1897    /// exact replay; the host owns delivery retries and acknowledgement.
1898    pub fn outbox_entries(&self) -> Result<Vec<OutboxEntry>, CanwuError> {
1899        Self::outbox_entries_for_boundaries(
1900            &self.state.metadata.run_manifest_hash,
1901            &self.state.evidence.boundaries,
1902        )
1903    }
1904
1905    pub(crate) fn outbox_entries_for_boundaries(
1906        run_manifest_hash: &str,
1907        boundaries: &[BoundaryRecord],
1908    ) -> Result<Vec<OutboxEntry>, CanwuError> {
1909        let mut entries = Vec::new();
1910        for boundary in boundaries {
1911            for (index, emission) in boundary.emissions.iter().enumerate() {
1912                let emission_index = u64::try_from(index).map_err(|_| {
1913                    CanwuError::new(
1914                        ErrorCode::IdentifierExhausted,
1915                        "outbox emission index exceeds the persistent identifier space",
1916                    )
1917                })?;
1918                let delivery_id = canonical_hash(
1919                    "canwu.outbox.delivery.v1",
1920                    &(
1921                        run_manifest_hash,
1922                        boundary.id,
1923                        emission.event,
1924                        emission_index,
1925                    ),
1926                )?;
1927                entries.push(OutboxEntry {
1928                    delivery_id,
1929                    boundary: boundary.id,
1930                    event: emission.event,
1931                    emission_index,
1932                    plugin: emission.plugin.clone(),
1933                    system: emission.system.clone(),
1934                });
1935            }
1936        }
1937        Ok(entries)
1938    }
1939
1940    fn compute_boundary_state_hash_for(
1941        &mut self,
1942        format: BoundaryStateHashFormat,
1943    ) -> Result<String, CanwuError> {
1944        match format {
1945            BoundaryStateHashFormat::LegacyV0 => self.compute_boundary_state_hash(),
1946            BoundaryStateHashFormat::CommitmentsV1 => {
1947                let roots = self.refresh_runtime_commitment_roots()?;
1948                boundary_state_hash_for_commitments(&roots)
1949            }
1950        }
1951    }
1952
1953    fn compute_boundary_state_hash(&self) -> Result<String, CanwuError> {
1954        let world = self.world();
1955        let entities: Vec<_> = self.state.current.entities.iter().cloned().collect();
1956        let plugin_components: Vec<_> = self
1957            .state
1958            .current
1959            .plugin_components
1960            .values()
1961            .cloned()
1962            .collect();
1963        let domain_records: Vec<_> = self
1964            .state
1965            .current
1966            .domain_records
1967            .values()
1968            .cloned()
1969            .collect();
1970        let plugin_descriptors: Vec<_> = self.plugins.descriptors().cloned().collect();
1971        let scheduled: Vec<_> = self
1972            .state
1973            .scheduler
1974            .actions
1975            .iter()
1976            .map(|(key, action)| ScheduledRecord {
1977                key: key.clone(),
1978                action: action.clone(),
1979            })
1980            .collect();
1981        let random_streams: Vec<_> = self
1982            .state
1983            .current
1984            .random_streams
1985            .values()
1986            .cloned()
1987            .collect();
1988        let (authoritative_manifest, authoritative_manifest_hash) = authoritative_run_identity(
1989            &self.state.metadata.run_manifest,
1990            &self.state.metadata.run_manifest_hash,
1991            &self.state.metadata.run_configuration,
1992        )?;
1993        let initial_scenario = hashing::committed_initial_scenario(self.bound_initial_scenario());
1994        state_hash(&StateHashMaterial {
1995            engine_version: ENGINE_VERSION,
1996            snapshot_format_version: SNAPSHOT_FORMAT_VERSION,
1997            run_manifest: &authoritative_manifest,
1998            run_manifest_hash: &authoritative_manifest_hash,
1999            initial_time: self.state.scheduler.initial_time,
2000            initial_scenario: initial_scenario.as_ref(),
2001            now: self.state.scheduler.now,
2002            plugin_registration_closed: self.state.metadata.plugin_registration_closed,
2003            entities: hashing::committed_entities(&entities, &world),
2004            world: &world,
2005            knowledge: &self.state.current.knowledge,
2006            events: &self.state.evidence.events,
2007            commands: &self.state.evidence.commands,
2008            command_attempts: &self.state.evidence.command_attempts,
2009            ingress: &self.state.evidence.ingress,
2010            plugin_components: &plugin_components,
2011            domain_records: &domain_records,
2012            decisions: &self.state.current.decisions,
2013            plugin_descriptors: &plugin_descriptors,
2014            schema: &self.schema,
2015            scheduled: &scheduled,
2016            root_seed: self.state.current.root_seed,
2017            authority_root_seed: self.state.current.authority_root_seed,
2018            random_streams: &random_streams,
2019            random_draws: &self.state.evidence.random_draws,
2020            next_event_id: self.state.counters.next_event_id,
2021            next_command_id: self.state.counters.next_command_id,
2022            next_command_attempt_id: self.state.counters.next_command_attempt_id,
2023            next_ingress_id: self.state.counters.next_ingress_id,
2024            next_boundary_id: self.state.counters.next_boundary_id,
2025            next_random_draw_id: self.state.counters.next_random_draw_id,
2026            next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
2027            next_schedule_sequence: self.state.counters.next_schedule_sequence,
2028            next_correlation_id: self.state.counters.next_correlation_id,
2029            next_decision_trace_id: self.state.counters.next_decision_trace_id,
2030        })
2031    }
2032
2033    fn compute_commitment_root_updates(
2034        &self,
2035        needs: CommitmentDomains,
2036    ) -> Result<RuntimeCommitmentRootUpdates, CanwuError> {
2037        let world = needs
2038            .contains(CommitmentDomains::WORLD)
2039            .then(|| {
2040                let world = self.world();
2041                let entities: Vec<_> = self.state.current.entities.iter().cloned().collect();
2042                world_commitment_root(&world, &entities)
2043            })
2044            .transpose()?;
2045        let knowledge = needs
2046            .contains(CommitmentDomains::KNOWLEDGE)
2047            .then(|| knowledge_commitment_root(&self.state.current.knowledge))
2048            .transpose()?;
2049        let plugin_components = needs
2050            .contains(CommitmentDomains::PLUGIN_COMPONENTS)
2051            .then(|| {
2052                let values: Vec<_> = self
2053                    .state
2054                    .current
2055                    .plugin_components
2056                    .values()
2057                    .cloned()
2058                    .collect();
2059                plugin_component_commitment_root(&values)
2060            })
2061            .transpose()?;
2062        let domain_records = needs
2063            .contains(CommitmentDomains::DOMAIN_RECORDS)
2064            .then(|| self.state.current.domain_records.commitment_root())
2065            .transpose()?;
2066        let decisions = needs
2067            .contains(CommitmentDomains::DECISIONS)
2068            .then(|| decision_commitment_root(&self.state.current.decisions))
2069            .transpose()?;
2070        let scheduler = needs
2071            .contains(CommitmentDomains::SCHEDULER)
2072            .then(|| {
2073                let scheduled: Vec<_> = self
2074                    .state
2075                    .scheduler
2076                    .actions
2077                    .iter()
2078                    .map(|(key, action)| ScheduledRecord {
2079                        key: key.clone(),
2080                        action: action.clone(),
2081                    })
2082                    .collect();
2083                scheduler_commitment_root(self.state.scheduler.now, &scheduled)
2084            })
2085            .transpose()?;
2086        let random_streams = needs
2087            .contains(CommitmentDomains::RANDOM_STREAMS)
2088            .then(|| {
2089                let values: Vec<_> = self
2090                    .state
2091                    .current
2092                    .random_streams
2093                    .values()
2094                    .cloned()
2095                    .collect();
2096                random_stream_commitment_root(&values)
2097            })
2098            .transpose()?;
2099        let identity = if needs.contains(CommitmentDomains::IDENTITY) {
2100            let descriptors: Vec<_> = self.plugins.descriptors().cloned().collect();
2101            let (manifest, manifest_hash) = authoritative_run_identity(
2102                &self.state.metadata.run_manifest,
2103                &self.state.metadata.run_manifest_hash,
2104                &self.state.metadata.run_configuration,
2105            )?;
2106            let initial_scenario =
2107                hashing::committed_initial_scenario(self.bound_initial_scenario());
2108            Some(identity_commitment_root(
2109                ENGINE_VERSION,
2110                SNAPSHOT_FORMAT_VERSION,
2111                &manifest,
2112                &manifest_hash,
2113                self.state.scheduler.initial_time,
2114                initial_scenario.as_ref(),
2115                self.state.current.authority_root_seed,
2116                &descriptors,
2117                &self.schema,
2118            )?)
2119        } else {
2120            None
2121        };
2122        Ok(RuntimeCommitmentRootUpdates {
2123            world,
2124            knowledge,
2125            plugin_components,
2126            domain_records,
2127            decisions,
2128            scheduler,
2129            random_streams,
2130            identity,
2131        })
2132    }
2133
2134    fn invalidate_commitments(&mut self, domains: CommitmentDomains) {
2135        if let Some(cache) = self.state.metadata.commitment_cache.as_mut() {
2136            cache.invalidate(domains);
2137        }
2138    }
2139
2140    fn refresh_runtime_commitment_roots(&mut self) -> Result<CommitmentRoots, CanwuError> {
2141        if self.state.metadata.commitment_format_version != COMMITMENT_FORMAT_VERSION {
2142            return Err(CanwuError::new(
2143                ErrorCode::UnsupportedSnapshotVersion,
2144                format!(
2145                    "commitment format {} cannot produce boundary state commitment v1",
2146                    self.state.metadata.commitment_format_version
2147                ),
2148            ));
2149        }
2150        let needs = {
2151            if self.state.metadata.commitment_cache.is_none() {
2152                self.state.metadata.commitment_cache =
2153                    Some(RuntimeCommitmentCache::from_evidence(&self.state.evidence)?);
2154            }
2155            let cache = self
2156                .state
2157                .metadata
2158                .commitment_cache
2159                .as_mut()
2160                .ok_or_else(|| {
2161                    CanwuError::new(
2162                        ErrorCode::InvalidSnapshot,
2163                        "commitment cache is unavailable while refreshing runtime roots",
2164                    )
2165                })?;
2166            cache.sync(&self.state.evidence)?;
2167            cache.needs()
2168        };
2169        let updates = self.compute_commitment_root_updates(needs)?;
2170        let boundary_head = self.boundary_head_hash().map(str::to_owned);
2171        let control = ControlCommitmentMaterial {
2172            plugin_registration_closed: self.state.metadata.plugin_registration_closed,
2173            next_event_id: self.state.counters.next_event_id,
2174            next_command_id: self.state.counters.next_command_id,
2175            next_command_attempt_id: self.state.counters.next_command_attempt_id,
2176            next_ingress_id: self.state.counters.next_ingress_id,
2177            next_boundary_id: self.state.counters.next_boundary_id,
2178            next_random_draw_id: self.state.counters.next_random_draw_id,
2179            next_knowledge_record_id: self.state.counters.next_knowledge_record_id,
2180            next_schedule_sequence: self.state.counters.next_schedule_sequence,
2181            next_correlation_id: self.state.counters.next_correlation_id,
2182            next_decision_trace_id: self.state.counters.next_decision_trace_id,
2183        };
2184        let (domain_roots, journal_roots) = {
2185            let cache = self
2186                .state
2187                .metadata
2188                .commitment_cache
2189                .as_mut()
2190                .ok_or_else(|| {
2191                    CanwuError::new(
2192                        ErrorCode::InvalidSnapshot,
2193                        "commitment cache is unavailable while applying root updates",
2194                    )
2195                })?;
2196            cache.apply(updates);
2197            (cache.domain_roots()?, cache.roots())
2198        };
2199        runtime_commitment_roots(
2200            &domain_roots,
2201            &journal_roots,
2202            self.state.current.root_seed,
2203            boundary_head.as_deref(),
2204            &control,
2205        )
2206    }
2207
2208    fn refresh_checkpoint_hash(&mut self) -> Result<(), CanwuError> {
2209        if self.state.metadata.commitment_format_version == COMMITMENT_FORMAT_VERSION {
2210            let roots = self.refresh_runtime_commitment_roots()?;
2211            self.state.metadata.checkpoint_hash = checkpoint_hash_for_commitments(
2212                &roots,
2213                &self.state.metadata.run_manifest_hash,
2214                self.state.metadata.commitment_format_version,
2215                STATE_REVISION_FORMAT_VERSION,
2216                self.state.counters.state_revision,
2217                self.state.metadata.replay_revision_format_version,
2218            )?;
2219            self.state.metadata.commitment_roots = Some(roots);
2220        } else if self.state.metadata.commitment_format_version == 0 {
2221            let state_hash = self.compute_boundary_state_hash()?;
2222            self.state.metadata.checkpoint_hash = checkpoint_hash_for_configuration(
2223                &state_hash,
2224                self.boundary_head_hash(),
2225                &self.state.metadata.run_manifest_hash,
2226                &self.state.metadata.run_configuration,
2227                STATE_REVISION_FORMAT_VERSION,
2228                self.state.counters.state_revision,
2229                self.state.metadata.replay_revision_format_version,
2230            )?;
2231            self.state.metadata.commitment_roots = None;
2232            self.state.metadata.commitment_cache = None;
2233        } else {
2234            return Err(CanwuError::new(
2235                ErrorCode::UnsupportedSnapshotVersion,
2236                format!(
2237                    "commitment format {} is unsupported; this engine writes format {COMMITMENT_FORMAT_VERSION}",
2238                    self.state.metadata.commitment_format_version
2239                ),
2240            ));
2241        }
2242        Ok(())
2243    }
2244
2245    fn next_state_revision(&self) -> Result<u64, CanwuError> {
2246        self.state
2247            .counters
2248            .state_revision
2249            .checked_add(1)
2250            .ok_or_else(|| {
2251                CanwuError::new(
2252                    ErrorCode::IdentifierExhausted,
2253                    "authoritative state revision space is exhausted",
2254                )
2255            })
2256    }
2257
2258    fn advance_state_revision(&mut self) -> Result<u64, CanwuError> {
2259        let next = self.next_state_revision()?;
2260        self.state.counters.state_revision = next;
2261        Ok(next)
2262    }
2263
2264    #[must_use]
2265    pub fn snapshot(&self) -> SimulationSnapshot {
2266        let mut snapshot = self.checkpoint_state();
2267        snapshot.events.clone_from(&self.state.evidence.events);
2268        snapshot.commands.clone_from(&self.state.evidence.commands);
2269        snapshot
2270            .command_attempts
2271            .clone_from(&self.state.evidence.command_attempts);
2272        snapshot.ingress.clone_from(&self.state.evidence.ingress);
2273        snapshot
2274            .boundaries
2275            .clone_from(&self.state.evidence.boundaries);
2276        snapshot
2277            .random_draws
2278            .clone_from(&self.state.evidence.random_draws);
2279        snapshot
2280    }
2281
2282    pub fn snapshot_json(&self) -> Result<String, CanwuError> {
2283        serde_json::to_string_pretty(&self.snapshot()).map_err(|error| {
2284            CanwuError::new(
2285                ErrorCode::InvalidSnapshot,
2286                format!("could not serialize snapshot: {error}"),
2287            )
2288        })
2289    }
2290
2291    pub fn from_snapshot(snapshot: SimulationSnapshot) -> Result<Self, CanwuError> {
2292        if snapshot.snapshot_format_version != SNAPSHOT_FORMAT_VERSION
2293            || snapshot.engine_version != ENGINE_VERSION
2294        {
2295            return Err(CanwuError::new(
2296                ErrorCode::UnsupportedSnapshotVersion,
2297                format!(
2298                    "the typed snapshot loader accepts only engine {ENGINE_VERSION} format {SNAPSHOT_FORMAT_VERSION}; pre-8 formats are not supported"
2299                ),
2300            ));
2301        }
2302        validate_current_snapshot_contract(&snapshot)?;
2303        validate_scenario_state(&Scenario {
2304            start_time: snapshot.now,
2305            entities: snapshot.entities.clone(),
2306            world: snapshot.world.clone(),
2307            knowledge: snapshot.knowledge.clone(),
2308            domain_records: snapshot.domain_records.clone(),
2309        })?;
2310        let plugins = PluginRegistry::from_descriptors(snapshot.plugin_descriptors.clone())?;
2311        validate_snapshot(&snapshot, &plugins)?;
2312        let admitted_ingress: BTreeSet<_> = snapshot
2313            .boundaries
2314            .iter()
2315            .flat_map(|boundary| boundary.admitted_ingress.iter().copied())
2316            .collect();
2317        let pending_ingress = snapshot
2318            .ingress
2319            .iter()
2320            .filter(|record| !admitted_ingress.contains(&record.id))
2321            .map(IngressQueueKey::from_record)
2322            .collect();
2323        let initial_scenario = Some(snapshot.initial_scenario.clone().ok_or_else(|| {
2324            invalid_snapshot_error("format 8 validation requires an initial scenario")
2325        })?);
2326        let initial_domain_record_indexes = initial_scenario
2327            .as_ref()
2328            .map(|scenario| {
2329                scenario
2330                    .domain_records
2331                    .iter()
2332                    .enumerate()
2333                    .map(|(index, record)| (record.reference.clone(), index))
2334                    .collect()
2335            })
2336            .unwrap_or_default();
2337        let current_domain_record_versions = build_current_domain_record_versions(
2338            initial_scenario.as_ref(),
2339            &snapshot.boundaries,
2340            &snapshot.domain_records,
2341        )?;
2342        let mut simulation = Self {
2343            state: RuntimeState {
2344                current: RuntimeCurrentState {
2345                    entities: snapshot.entities.into_iter().collect(),
2346                    people: snapshot
2347                        .world
2348                        .people
2349                        .into_iter()
2350                        .map(|value| (value.id, value))
2351                        .collect(),
2352                    letters: snapshot
2353                        .world
2354                        .letters
2355                        .into_iter()
2356                        .map(|value| (value.id, value))
2357                        .collect(),
2358                    governments: snapshot
2359                        .world
2360                        .governments
2361                        .into_iter()
2362                        .map(|value| (value.id, value))
2363                        .collect(),
2364                    territories: snapshot
2365                        .world
2366                        .territories
2367                        .into_iter()
2368                        .map(|value| (value.id, value))
2369                        .collect(),
2370                    routes: snapshot
2371                        .world
2372                        .routes
2373                        .into_iter()
2374                        .map(|value| (value.id, value))
2375                        .collect(),
2376                    armies: snapshot
2377                        .world
2378                        .armies
2379                        .into_iter()
2380                        .map(|value| (value.id, value))
2381                        .collect(),
2382                    knowledge: snapshot.knowledge,
2383                    plugin_components: snapshot
2384                        .plugin_components
2385                        .into_iter()
2386                        .map(|record| {
2387                            (
2388                                component_key(
2389                                    &record.plugin,
2390                                    &record.state,
2391                                    &record.entity,
2392                                    &record.component,
2393                                ),
2394                                record,
2395                            )
2396                        })
2397                        .collect(),
2398                    domain_records: PersistentDomainRecordStore::from_records(
2399                        snapshot
2400                            .domain_records
2401                            .into_iter()
2402                            .map(|record| (record.reference.clone(), record))
2403                            .collect(),
2404                    )?,
2405                    decisions: snapshot.decisions,
2406                    root_seed: snapshot.root_seed,
2407                    authority_root_seed: snapshot.authority_root_seed,
2408                    random_streams: snapshot
2409                        .random_streams
2410                        .into_iter()
2411                        .map(|state| (state.key.clone(), state))
2412                        .collect(),
2413                },
2414                scheduler: RuntimeScheduler {
2415                    initial_time: snapshot.initial_time,
2416                    now: snapshot.now,
2417                    actions: snapshot
2418                        .scheduled
2419                        .into_iter()
2420                        .map(|record| (record.key, record.action))
2421                        .collect(),
2422                    pending_ingress,
2423                },
2424                counters: RuntimeCounters {
2425                    next_event_id: snapshot.next_event_id,
2426                    next_command_id: snapshot.next_command_id,
2427                    next_command_attempt_id: snapshot.next_command_attempt_id,
2428                    next_ingress_id: snapshot.next_ingress_id,
2429                    next_boundary_id: snapshot.next_boundary_id,
2430                    next_random_draw_id: snapshot.next_random_draw_id,
2431                    next_knowledge_record_id: snapshot.next_knowledge_record_id,
2432                    next_schedule_sequence: snapshot.next_schedule_sequence,
2433                    next_correlation_id: snapshot.next_correlation_id,
2434                    next_decision_trace_id: snapshot.next_decision_trace_id,
2435                    state_revision: snapshot.state_revision,
2436                    admitted_attempt_count: snapshot.admitted_attempt_count,
2437                    admitted_command_count: snapshot.admitted_command_count,
2438                    admitted_event_count: snapshot.admitted_event_count,
2439                },
2440                metadata: RuntimeMetadata {
2441                    initial_scenario,
2442                    initial_domain_record_indexes,
2443                    current_domain_record_versions,
2444                    run_manifest: snapshot.run_manifest.clone().ok_or_else(|| {
2445                        invalid_snapshot_error("snapshot is missing its run manifest")
2446                    })?,
2447                    run_manifest_hash: snapshot.run_manifest_hash.clone(),
2448                    run_configuration: snapshot.run_configuration.clone().ok_or_else(|| {
2449                        invalid_snapshot_error("snapshot is missing its run configuration")
2450                    })?,
2451                    checkpoint_hash: snapshot.checkpoint_hash.clone(),
2452                    commitment_format_version: snapshot.commitment_format_version,
2453                    commitment_roots: snapshot.commitment_roots.clone(),
2454                    commitment_cache: None,
2455                    plugin_registration_closed: snapshot.plugin_registration_closed,
2456                    replay_revision_format_version: snapshot.replay_revision_format_version,
2457                },
2458                evidence: RuntimeEvidence {
2459                    archived: EvidenceCursor::default(),
2460                    archived_boundary_head: None,
2461                    archived_legacy_commands: false,
2462                    archived_tracked_attempts: false,
2463                    archived_unqueued_command_history: false,
2464                    archived_command_requests: BTreeMap::new(),
2465                    archived_ingress_requests: BTreeMap::new(),
2466                    archived_decision_requests: BTreeMap::new(),
2467                    archived_decision_command_requests: BTreeSet::new(),
2468                    events: snapshot.events,
2469                    commands: snapshot.commands,
2470                    command_attempts: snapshot.command_attempts,
2471                    ingress: snapshot.ingress,
2472                    boundaries: snapshot.boundaries,
2473                    random_draws: snapshot.random_draws,
2474                    archived_segment_headers: Vec::new(),
2475                    archived_evidence_receipts: BTreeMap::new(),
2476                    keyed_draw_reservations: Vec::new(),
2477                },
2478            },
2479            schema: snapshot.schema,
2480            plugins,
2481            plugin_archive_provider: Rc::new(()),
2482            sync_reaction_depth: 0,
2483        };
2484        simulation.refresh_checkpoint_hash()?;
2485        Ok(simulation)
2486    }
2487
2488    pub fn from_snapshot_json(json: &str) -> Result<Self, CanwuError> {
2489        let snapshot = deserialize_current_snapshot_json(json)?;
2490        Self::from_snapshot(snapshot)
2491    }
2492
2493    pub fn from_snapshot_with_plugins(
2494        snapshot: SimulationSnapshot,
2495        plugins: &[&dyn SimulationPlugin],
2496    ) -> Result<Self, CanwuError> {
2497        let mut simulation = Self::from_snapshot(snapshot)?;
2498        for plugin in plugins {
2499            simulation.register_plugin(*plugin)?;
2500        }
2501        simulation.ensure_runtime_ready()?;
2502        Ok(simulation)
2503    }
2504
2505    pub fn from_snapshot_json_with_plugins(
2506        json: &str,
2507        plugins: &[&dyn SimulationPlugin],
2508    ) -> Result<Self, CanwuError> {
2509        let snapshot = deserialize_current_snapshot_json(json)?;
2510        Self::from_snapshot_with_plugins(snapshot, plugins)
2511    }
2512
2513    #[must_use]
2514    pub fn fork(&self) -> Self {
2515        Self {
2516            state: self.state.clone(),
2517            schema: self.schema.clone(),
2518            plugins: self.plugins.clone(),
2519            plugin_archive_provider: Rc::clone(&self.plugin_archive_provider),
2520            sync_reaction_depth: 0,
2521        }
2522    }
2523
2524    /// Attaches the caller-owned provider used to resolve package cold
2525    /// archives during normal command admission, boundary settlement, and
2526    /// detached queries. The provider is executable host context and is never
2527    /// serialized into a snapshot.
2528    pub fn set_plugin_archive_object_provider(
2529        &mut self,
2530        provider: Rc<dyn PluginArchiveObjectProvider>,
2531    ) {
2532        self.plugin_archive_provider = provider;
2533    }
2534
2535    /// Loads one package-owned archive object from the currently attached host
2536    /// provider. Callers must authenticate returned bytes against a committed
2537    /// package archive root before treating them as authoritative.
2538    pub fn plugin_archive_object(
2539        &self,
2540        namespace: &str,
2541        object_id: &str,
2542    ) -> Result<Option<Vec<u8>>, CanwuError> {
2543        self.plugin_archive_provider
2544            .load_plugin_archive_object(namespace, object_id)
2545    }
2546
2547    fn prepare_command(
2548        &self,
2549        envelope: &CommandEnvelope,
2550        context: &CommandContext,
2551    ) -> Result<PreparedCommand, CanwuError> {
2552        match &envelope.command {
2553            Command::OrderMovement {
2554                subject,
2555                destination,
2556                cargo,
2557            } => {
2558                let Some(actor) = decision_actor(&context.authority) else {
2559                    return Err(CanwuError::new(
2560                        ErrorCode::InvalidAuthority,
2561                        "movement commands require an accountable actor origin",
2562                    ));
2563                };
2564                let person = self.state.current.people.get(&actor).ok_or_else(|| {
2565                    CanwuError::new(
2566                        ErrorCode::ActorNotFound,
2567                        format!("actor {actor} was not found"),
2568                    )
2569                    .with_entity(EntityRef::Person(actor))
2570                })?;
2571                if context
2572                    .authority
2573                    .command_subject
2574                    .as_ref()
2575                    .is_some_and(|bound| bound != subject)
2576                {
2577                    return Err(CanwuError::new(
2578                        ErrorCode::InvalidAuthority,
2579                        "command subject does not match the movement subject",
2580                    )
2581                    .with_entity(subject.clone()));
2582                }
2583                if !self.state.current.territories.contains_key(destination) {
2584                    return Err(CanwuError::new(
2585                        ErrorCode::DestinationNotFound,
2586                        format!("destination {destination} was not found"),
2587                    )
2588                    .with_entity(EntityRef::Territory(*destination)));
2589                }
2590                if cargo.windows(2).any(|pair| pair[0] >= pair[1]) {
2591                    return Err(CanwuError::new(
2592                        ErrorCode::InvalidPayload,
2593                        "movement cargo IDs must be sorted and unique",
2594                    ));
2595                }
2596                match subject {
2597                    EntityRef::Army(army) => {
2598                        if !cargo.is_empty() {
2599                            return Err(CanwuError::new(
2600                                ErrorCode::InvalidPayload,
2601                                "army movement does not accept letter cargo yet",
2602                            ));
2603                        }
2604                        let army_state = self.state.current.armies.get(army).ok_or_else(|| {
2605                            CanwuError::new(
2606                                ErrorCode::ArmyNotFound,
2607                                format!("army {army} was not found"),
2608                            )
2609                            .with_entity(EntityRef::Army(*army))
2610                        })?;
2611                        if army_state.commander != person.id {
2612                            return Err(CanwuError::new(
2613                                ErrorCode::InvalidAuthority,
2614                                format!("{} does not command {}", person.name, army_state.name),
2615                            )
2616                            .with_entity(EntityRef::Person(person.id))
2617                            .with_entity(EntityRef::Army(*army)));
2618                        }
2619                        if army_state.transit.is_some() {
2620                            return Err(CanwuError::new(
2621                                ErrorCode::InvalidAuthority,
2622                                format!("{} is already moving", army_state.name),
2623                            )
2624                            .with_entity(EntityRef::Army(*army)));
2625                        }
2626                        let arrival_at =
2627                            self.movement_arrival_time(army_state.location, *destination)?;
2628                        Ok(PreparedCommand::ArmyMovement {
2629                            army: *army,
2630                            actor,
2631                            from: army_state.location,
2632                            destination: *destination,
2633                            arrival_at,
2634                        })
2635                    }
2636                    EntityRef::Person(person_id) => {
2637                        if *person_id != actor
2638                            || context
2639                                .authority
2640                                .command_subject
2641                                .as_ref()
2642                                .is_some_and(|subject| subject != &EntityRef::Person(*person_id))
2643                        {
2644                            return Err(CanwuError::new(
2645                                ErrorCode::InvalidAuthority,
2646                                "self-directed movement must bind the actor to the person subject",
2647                            )
2648                            .with_entity(EntityRef::Person(*person_id)));
2649                        }
2650                        let person_state =
2651                            self.state.current.people.get(person_id).ok_or_else(|| {
2652                                CanwuError::new(
2653                                    ErrorCode::EntityNotFound,
2654                                    format!("person {person_id} was not found"),
2655                                )
2656                                .with_entity(EntityRef::Person(*person_id))
2657                            })?;
2658                        if person_state.transit.is_some() {
2659                            return Err(CanwuError::new(
2660                                ErrorCode::InvalidAuthority,
2661                                format!("person {person_id} is already moving"),
2662                            )
2663                            .with_entity(EntityRef::Person(*person_id)));
2664                        }
2665                        for letter_id in cargo {
2666                            let letter =
2667                                self.state.current.letters.get(letter_id).ok_or_else(|| {
2668                                    CanwuError::new(
2669                                        ErrorCode::EntityNotFound,
2670                                        format!("letter {letter_id} was not found"),
2671                                    )
2672                                    .with_entity(
2673                                        EntityRef::Resource(ResourceId::new(letter_id.get())),
2674                                    )
2675                                })?;
2676                            if letter.status != LetterStatus::HeldByPerson
2677                                || letter.carrier != Some(*person_id)
2678                                || !self.state.current.people.contains_key(&letter.sender)
2679                                || !self.state.current.people.contains_key(&letter.recipient)
2680                            {
2681                                return Err(CanwuError::new(
2682                                    ErrorCode::InvalidAuthority,
2683                                    format!("letter {letter_id} is not held by the moving person"),
2684                                )
2685                                .with_entity(EntityRef::Resource(ResourceId::new(
2686                                    letter_id.get(),
2687                                ))));
2688                            }
2689                        }
2690                        let arrival_at = self
2691                            .movement_arrival_time(person_state.current_location, *destination)?;
2692                        Ok(PreparedCommand::MovePerson {
2693                            person: *person_id,
2694                            from: person_state.current_location,
2695                            destination: *destination,
2696                            cargo: cargo.clone(),
2697                            arrival_at,
2698                        })
2699                    }
2700                    _ => Err(CanwuError::new(
2701                        ErrorCode::InvalidAuthority,
2702                        "only army and person subjects support built-in movement",
2703                    )
2704                    .with_entity(subject.clone())),
2705                }
2706            }
2707            Command::DebugSetArmyMorale { army, morale } => {
2708                if envelope.issuer != Issuer::Debug {
2709                    return Err(CanwuError::new(
2710                        ErrorCode::InvalidAuthority,
2711                        "debug state edits require the explicit debug issuer",
2712                    ));
2713                }
2714                if *morale > 100 {
2715                    return Err(CanwuError::new(
2716                        ErrorCode::ValueOutOfRange,
2717                        "army morale must be between 0 and 100",
2718                    ));
2719                }
2720                let old_morale = self.state.current.armies.get(army).map_or_else(
2721                    || {
2722                        Err(CanwuError::new(
2723                            ErrorCode::ArmyNotFound,
2724                            format!("army {army} was not found"),
2725                        ))
2726                    },
2727                    |army_state| Ok(army_state.morale),
2728                )?;
2729                Ok(PreparedCommand::DebugMorale {
2730                    army: *army,
2731                    old_morale,
2732                    new_morale: *morale,
2733                })
2734            }
2735            Command::Plugin {
2736                plugin,
2737                command,
2738                payload,
2739            } => {
2740                let registered = self
2741                    .plugins
2742                    .commands
2743                    .get(&(plugin.clone(), command.clone()))
2744                    .ok_or_else(|| {
2745                        CanwuError::new(
2746                            ErrorCode::PluginCommandNotFound,
2747                            format!("plugin command {plugin}.{command} is not registered"),
2748                        )
2749                    })?;
2750                let handler = registered.handler;
2751                let descriptor = registered.descriptor.clone();
2752                descriptor.payload_schema.validate(payload)?;
2753                let reader = format!("{plugin}.{command}");
2754                let directives = catch_unwind(AssertUnwindSafe(|| {
2755                    handler(
2756                        &self.plugin_view(&reader, &descriptor.reads),
2757                        context,
2758                        payload,
2759                    )
2760                }))
2761                .map_err(|_| {
2762                    CanwuError::new(
2763                        ErrorCode::PluginPanicked,
2764                        format!("plugin command {plugin}.{command} panicked"),
2765                    )
2766                })??;
2767                validate_directives_with_context(
2768                    &RuntimeValidationContext::new(&self.state),
2769                    plugin,
2770                    &descriptor.writes,
2771                    &self.plugins.state_owners,
2772                    &self.plugins.record_schemas,
2773                    &directives,
2774                )?;
2775                Ok(PreparedCommand::Plugin {
2776                    plugin: plugin.clone(),
2777                    directives,
2778                    allowed_writes: descriptor.writes,
2779                })
2780            }
2781        }
2782    }
2783
2784    fn movement_arrival_time(
2785        &self,
2786        from: TerritoryId,
2787        to: TerritoryId,
2788    ) -> Result<SimTime, CanwuError> {
2789        let travel_minutes = if from == to {
2790            1
2791        } else {
2792            self.state
2793                .current
2794                .routes
2795                .values()
2796                .find(|route| route.connects(from, to))
2797                .ok_or_else(|| {
2798                    CanwuError::new(
2799                        ErrorCode::NoRoute,
2800                        format!("no direct route connects territory {from} to {to}"),
2801                    )
2802                })?
2803                .travel_minutes
2804        };
2805        if travel_minutes <= 0 {
2806            return Err(CanwuError::new(
2807                ErrorCode::InvalidDuration,
2808                "movement route duration must be positive",
2809            ));
2810        }
2811        self.state
2812            .scheduler
2813            .now
2814            .checked_add(SimDuration::minutes(travel_minutes))
2815            .ok_or_else(|| {
2816                CanwuError::new(
2817                    ErrorCode::InvalidDuration,
2818                    "movement arrival time exceeds the supported range",
2819                )
2820            })
2821    }
2822
2823    fn apply_prepared(
2824        &mut self,
2825        prepared: PreparedCommand,
2826        command_id: CommandId,
2827        correlation_id: u64,
2828    ) -> Result<(), CanwuError> {
2829        match prepared {
2830            PreparedCommand::ArmyMovement {
2831                army,
2832                actor,
2833                from,
2834                destination,
2835                arrival_at,
2836            } => {
2837                let army_state = self.state.current.armies.get_mut(&army).ok_or_else(|| {
2838                    CanwuError::new(ErrorCode::ArmyNotFound, "validated army disappeared")
2839                })?;
2840                army_state.transit = Some(TransitState {
2841                    from,
2842                    to: destination,
2843                    departed_at: self.state.scheduler.now,
2844                    arrives_at: arrival_at,
2845                });
2846                let event = self.emit(
2847                    MoveOrdered {
2848                        army,
2849                        from,
2850                        to: destination,
2851                        arrival_at,
2852                    }
2853                    .into_kind(),
2854                    vec![
2855                        EntityRef::Army(army),
2856                        EntityRef::Person(actor),
2857                        EntityRef::Territory(from),
2858                        EntityRef::Territory(destination),
2859                    ],
2860                    format!("Army {army} was ordered from {from} to {destination}"),
2861                    Some(CauseRef::Command(command_id)),
2862                    correlation_id,
2863                )?;
2864                self.schedule_at(
2865                    arrival_at,
2866                    ScheduledAction::ArmyArrival {
2867                        army,
2868                        destination,
2869                        order_event: event,
2870                        correlation_id,
2871                    },
2872                )?;
2873            }
2874            PreparedCommand::MovePerson {
2875                person,
2876                from,
2877                destination,
2878                cargo,
2879                arrival_at,
2880            } => {
2881                self.invalidate_commitments(CommitmentDomains::WORLD);
2882                let person_state = self.state.current.people.get_mut(&person).ok_or_else(|| {
2883                    CanwuError::new(ErrorCode::EntityNotFound, "validated person disappeared")
2884                })?;
2885                person_state.transit = Some(PersonTransitState {
2886                    from,
2887                    to: destination,
2888                    departed_at: self.state.scheduler.now,
2889                    arrives_at: arrival_at,
2890                });
2891                for letter_id in &cargo {
2892                    let letter =
2893                        self.state
2894                            .current
2895                            .letters
2896                            .get_mut(letter_id)
2897                            .ok_or_else(|| {
2898                                CanwuError::new(
2899                                    ErrorCode::EntityNotFound,
2900                                    "validated letter disappeared",
2901                                )
2902                            })?;
2903                    letter.status = LetterStatus::InTransit;
2904                    letter.carrier = Some(person);
2905                    letter.location = None;
2906                }
2907                let event = self.emit(
2908                    PersonMoveOrdered {
2909                        person,
2910                        from,
2911                        to: destination,
2912                        arrival_at,
2913                    }
2914                    .into_kind(),
2915                    std::iter::once(EntityRef::Person(person))
2916                        .chain(
2917                            cargo
2918                                .iter()
2919                                .copied()
2920                                .map(|id| EntityRef::Resource(ResourceId::new(id.get()))),
2921                        )
2922                        .chain([
2923                            EntityRef::Territory(from),
2924                            EntityRef::Territory(destination),
2925                        ])
2926                        .collect(),
2927                    format!("Person {person} was ordered from {from} to {destination}"),
2928                    Some(CauseRef::Command(command_id)),
2929                    correlation_id,
2930                )?;
2931                self.schedule_at(
2932                    arrival_at,
2933                    ScheduledAction::PersonArrival {
2934                        person,
2935                        destination,
2936                        order_event: event,
2937                        cargo,
2938                        correlation_id,
2939                    },
2940                )?;
2941            }
2942            PreparedCommand::DebugMorale {
2943                army,
2944                old_morale,
2945                new_morale,
2946            } => {
2947                self.state
2948                    .current
2949                    .armies
2950                    .get_mut(&army)
2951                    .ok_or_else(|| {
2952                        CanwuError::new(ErrorCode::ArmyNotFound, "validated army disappeared")
2953                    })?
2954                    .morale = new_morale;
2955                self.emit(
2956                    DebugFieldChanged {
2957                        entity: EntityRef::Army(army),
2958                        field: "morale".to_owned(),
2959                        old_value: old_morale.to_string(),
2960                        new_value: new_morale.to_string(),
2961                    }
2962                    .into_kind(),
2963                    vec![EntityRef::Army(army)],
2964                    format!(
2965                        "Debug command changed army {army} morale {old_morale} -> {new_morale}"
2966                    ),
2967                    Some(CauseRef::Command(command_id)),
2968                    correlation_id,
2969                )?;
2970            }
2971            PreparedCommand::Plugin {
2972                plugin,
2973                directives,
2974                allowed_writes,
2975            } => {
2976                self.apply_directives(
2977                    &plugin,
2978                    directives,
2979                    &allowed_writes,
2980                    &CauseRef::Command(command_id),
2981                    correlation_id,
2982                )?;
2983            }
2984        }
2985        Ok(())
2986    }
2987}
2988
2989enum PreparedCommand {
2990    ArmyMovement {
2991        army: ArmyId,
2992        actor: PersonId,
2993        from: TerritoryId,
2994        destination: TerritoryId,
2995        arrival_at: SimTime,
2996    },
2997    MovePerson {
2998        person: PersonId,
2999        from: TerritoryId,
3000        destination: TerritoryId,
3001        cargo: Vec<LetterId>,
3002        arrival_at: SimTime,
3003    },
3004    DebugMorale {
3005        army: ArmyId,
3006        old_morale: u16,
3007        new_morale: u16,
3008    },
3009    Plugin {
3010        plugin: String,
3011        directives: Vec<SystemDirective>,
3012        allowed_writes: Vec<StateKey>,
3013    },
3014}
3015
3016impl PreparedCommand {
3017    fn commitment_invalidation(&self) -> CommitmentDomains {
3018        match self {
3019            Self::ArmyMovement { .. } => {
3020                CommitmentDomains::WORLD
3021                    | CommitmentDomains::KNOWLEDGE
3022                    | CommitmentDomains::PLUGIN_COMPONENTS
3023                    | CommitmentDomains::SCHEDULER
3024            }
3025            Self::MovePerson { .. } => CommitmentDomains::WORLD | CommitmentDomains::SCHEDULER,
3026            Self::DebugMorale { .. } => {
3027                CommitmentDomains::WORLD
3028                    | CommitmentDomains::PLUGIN_COMPONENTS
3029                    | CommitmentDomains::SCHEDULER
3030            }
3031            Self::Plugin { .. } => {
3032                CommitmentDomains::PLUGIN_COMPONENTS | CommitmentDomains::SCHEDULER
3033            }
3034        }
3035    }
3036}
3037
3038fn validate_directives(
3039    plugin: &str,
3040    allowed_writes: &[StateKey],
3041    state_owners: &BTreeMap<StateKey, String>,
3042    record_schemas: &records::DomainRecordSchemas,
3043    entity_exists: &dyn Fn(&EntityRef) -> bool,
3044    directives: &[SystemDirective],
3045) -> Result<(), CanwuError> {
3046    for directive in directives {
3047        match directive {
3048            SystemDirective::SetComponent {
3049                state,
3050                entity,
3051                component,
3052                ..
3053            } => {
3054                if component.trim().is_empty() || component != component.trim() {
3055                    return Err(CanwuError::new(
3056                        ErrorCode::InvalidPayload,
3057                        "plugin component name must be non-empty and canonical",
3058                    ));
3059                }
3060                if !allowed_writes.contains(state) {
3061                    return Err(CanwuError::new(
3062                        ErrorCode::UndeclaredStateWrite,
3063                        format!(
3064                            "plugin {plugin} did not declare write access to {}.{}",
3065                            state.namespace, state.name
3066                        ),
3067                    ));
3068                }
3069                if state_owners.get(state).is_none_or(|owner| owner != plugin) {
3070                    return Err(CanwuError::new(
3071                        ErrorCode::UndeclaredStateWrite,
3072                        format!(
3073                            "plugin {plugin} does not own state {}.{}",
3074                            state.namespace, state.name
3075                        ),
3076                    ));
3077                }
3078                if is_domain_record_state(record_schemas, state) {
3079                    return Err(CanwuError::new(
3080                        ErrorCode::UndeclaredStateWrite,
3081                        "domain record state cannot be written as an immediate component",
3082                    ));
3083                }
3084                if !entity_exists(entity) {
3085                    return Err(CanwuError::new(
3086                        ErrorCode::EntityNotFound,
3087                        format!("plugin {plugin} targeted missing entity {entity}"),
3088                    )
3089                    .with_entity(entity.clone()));
3090                }
3091            }
3092            SystemDirective::Emit { event_type, .. }
3093                if event_type.trim().is_empty() || event_type != event_type.trim() =>
3094            {
3095                return Err(CanwuError::new(
3096                    ErrorCode::InvalidPayload,
3097                    "plugin event type must be non-empty and canonical",
3098                ));
3099            }
3100            SystemDirective::Emit { affected, .. }
3101                if affected.iter().any(|entity| !entity_exists(entity)) =>
3102            {
3103                return Err(CanwuError::new(
3104                    ErrorCode::EntityNotFound,
3105                    format!("plugin {plugin} emitted an event for a missing entity"),
3106                ));
3107            }
3108            SystemDirective::Schedule { after, directive } => {
3109                if *after <= SimDuration::ZERO {
3110                    return Err(CanwuError::new(
3111                        ErrorCode::InvalidDuration,
3112                        "plugin systems must schedule work strictly in the future",
3113                    ));
3114                }
3115                validate_directives(
3116                    plugin,
3117                    allowed_writes,
3118                    state_owners,
3119                    record_schemas,
3120                    entity_exists,
3121                    std::slice::from_ref(directive),
3122                )?;
3123            }
3124            SystemDirective::EnqueuePluginIngress {
3125                after,
3126                packet_type,
3127                affected,
3128                ..
3129            } => {
3130                if packet_type.trim().is_empty() || packet_type != packet_type.trim() {
3131                    return Err(CanwuError::new(
3132                        ErrorCode::InvalidPayload,
3133                        "plugin ingress type must be non-empty and canonical",
3134                    ));
3135                }
3136                if *after < SimDuration::ZERO {
3137                    return Err(CanwuError::new(
3138                        ErrorCode::InvalidDuration,
3139                        "plugin command ingress delay cannot be negative",
3140                    ));
3141                }
3142                if affected.iter().any(|entity| !entity_exists(entity)) {
3143                    return Err(CanwuError::new(
3144                        ErrorCode::EntityNotFound,
3145                        format!("plugin {plugin} queued ingress for a missing entity"),
3146                    ));
3147                }
3148            }
3149            SystemDirective::Emit { .. } => {}
3150        }
3151    }
3152    Ok(())
3153}
3154
3155fn resolve_command_authority(envelope: &CommandEnvelope) -> Result<CommandAuthority, CanwuError> {
3156    if let Some(authority) = &envelope.authority {
3157        return Ok(authority.clone());
3158    }
3159    match &envelope.issuer {
3160        Issuer::Actor(actor) => Ok(CommandAuthority::for_actor(*actor)),
3161        Issuer::Debug => Ok(CommandAuthority::no_responsible_actor("debug-command")),
3162        Issuer::System(system) => Ok(CommandAuthority::no_responsible_actor(format!(
3163            "system:{system}"
3164        ))),
3165        Issuer::Human(_)
3166        | Issuer::Ai(_)
3167        | Issuer::Institution(_)
3168        | Issuer::Replay(_)
3169        | Issuer::Experiment(_) => Err(CanwuError::new(
3170            ErrorCode::InvalidAuthority,
3171            "typed command origins require an explicit authority context",
3172        )),
3173    }
3174}
3175
3176fn validate_command_ingress_policy(
3177    run_configuration: &RunConfigurationSnapshot,
3178    issuer: &Issuer,
3179    authority: &CommandAuthority,
3180    admission: CommandAdmission,
3181    entity_exists: &dyn Fn(&EntityRef) -> bool,
3182) -> Result<(), CanwuError> {
3183    let CommandAdmission {
3184        request_id,
3185        expected_revision,
3186        expected_time,
3187        revision_before: current_revision,
3188        ingress,
3189    } = admission;
3190    if request_id.is_some_and(|id| id.get() == 0) {
3191        return Err(CanwuError::new(
3192            ErrorCode::InvalidPayload,
3193            "command request IDs must be nonzero",
3194        ));
3195    }
3196    if let Some(expected) = expected_revision
3197        && expected != current_revision
3198    {
3199        return Err(CanwuError::new(
3200            ErrorCode::SimulationRevisionConflict,
3201            format!(
3202                "command expected revision {expected}, but simulation is at revision {current_revision}"
3203            ),
3204        ));
3205    }
3206    validate_command_authority(authority, entity_exists)?;
3207    if matches!(issuer, Issuer::Replay(_)) != (ingress == CommandIngress::FrozenReplay) {
3208        return Err(CanwuError::new(
3209            ErrorCode::InvalidAuthority,
3210            "replay command origins are valid only for frozen replay ingress",
3211        ));
3212    }
3213
3214    let RunConfigurationSnapshot::Declared(configuration) = run_configuration else {
3215        return Ok(());
3216    };
3217    if ingress == CommandIngress::LegacyDirect {
3218        return Err(CanwuError::new(
3219            ErrorCode::InvalidAuthority,
3220            "declared runs require tracked request or frozen replay ingress",
3221        ));
3222    }
3223    let external = !matches!(issuer, Issuer::System(_));
3224    if configuration.require_idempotency_keys && external && request_id.is_none() {
3225        return Err(CanwuError::new(
3226            ErrorCode::MissingIdempotencyKey,
3227            "this run requires a stable command request ID",
3228        ));
3229    }
3230    if configuration.require_idempotency_keys && external && expected_revision.is_none() {
3231        return Err(CanwuError::new(
3232            ErrorCode::SimulationRevisionConflict,
3233            "this run requires an expected command revision",
3234        ));
3235    }
3236    if configuration.interaction == InteractionPolicy::ReadOnly
3237        && !matches!(issuer, Issuer::Replay(_) | Issuer::System(_))
3238    {
3239        return Err(CanwuError::new(
3240            ErrorCode::InteractionReadOnly,
3241            "the run interaction policy rejects newly authored authoritative commands",
3242        ));
3243    }
3244    if external && expected_time.is_none() {
3245        return Err(CanwuError::new(
3246            ErrorCode::SimulationTimeConflict,
3247            "declared external commands require an expected simulation time",
3248        ));
3249    }
3250
3251    match issuer {
3252        Issuer::Actor(_) => Err(CanwuError::new(
3253            ErrorCode::InvalidAuthority,
3254            "declared runs require a typed human, AI, institution, replay, experiment, debug, or system origin",
3255        )),
3256        Issuer::Human(controller) => {
3257            let Some(binding) = &configuration.seat_binding else {
3258                return Err(CanwuError::new(
3259                    ErrorCode::InvalidAuthority,
3260                    "human commands require the run's exact seat binding",
3261                ));
3262            };
3263            if configuration.controller != ControllerPolicy::HumanRoleBound
3264                || controller != &binding.controller_id
3265                || authority.seat_id.as_deref() != Some(binding.seat_id.as_str())
3266                || authority.permission_profile_id.as_deref()
3267                    != Some(binding.permission_profile_id.as_str())
3268                || !authority_matches_seat_binding(configuration.seat, binding, authority)
3269            {
3270                return Err(CanwuError::new(
3271                    ErrorCode::InvalidAuthority,
3272                    "human command origin does not match the active controller, seat binding, and permission profile",
3273                ));
3274            }
3275            Ok(())
3276        }
3277        Issuer::Ai(controller) | Issuer::Institution(controller) => {
3278            if !canonical_text(controller)
3279                || matches!(
3280                    authority.decision_origin,
3281                    DecisionOrigin::NoResponsibleActor { .. }
3282                )
3283            {
3284                return Err(CanwuError::new(
3285                    ErrorCode::InvalidAuthority,
3286                    "AI and institutional commands require a canonical controller and responsible decision origin",
3287                ));
3288            }
3289            Ok(())
3290        }
3291        Issuer::Replay(source) => {
3292            if !canonical_text(source)
3293                || ingress != CommandIngress::FrozenReplay
3294                || configuration.purpose != RunPurpose::Replay
3295                || configuration.controller != ControllerPolicy::ReplayController
3296                || configuration.interaction != InteractionPolicy::ReadOnly
3297            {
3298                return Err(CanwuError::new(
3299                    ErrorCode::InvalidAuthority,
3300                    "replay command sources require a replay-purpose, replay-controller, read-only run",
3301                ));
3302            }
3303            if let Some(binding) = &configuration.seat_binding
3304                && (source != &binding.controller_id
3305                    || authority.seat_id.as_deref() != Some(binding.seat_id.as_str())
3306                    || authority.permission_profile_id.as_deref()
3307                        != Some(binding.permission_profile_id.as_str())
3308                    || !authority_matches_seat_binding(configuration.seat, binding, authority))
3309            {
3310                return Err(CanwuError::new(
3311                    ErrorCode::InvalidAuthority,
3312                    "frozen replay input does not match its recorded controller and seat binding",
3313                ));
3314            }
3315            Ok(())
3316        }
3317        Issuer::Experiment(intervention) => {
3318            if configuration.interaction != InteractionPolicy::VersionedExperiment
3319                || !configuration.declared_interventions.contains(intervention)
3320            {
3321                return Err(CanwuError::new(
3322                    ErrorCode::InvalidAuthority,
3323                    "experiment commands must name an intervention declared by the run",
3324                ));
3325            }
3326            Ok(())
3327        }
3328        Issuer::Debug => {
3329            if !configuration.diagnostic_commands_enabled {
3330                return Err(CanwuError::new(
3331                    ErrorCode::InvalidAuthority,
3332                    "debug command authority is disabled by the run configuration",
3333                ));
3334            }
3335            Ok(())
3336        }
3337        Issuer::System(system) => {
3338            if !canonical_text(system)
3339                || !matches!(
3340                    authority.decision_origin,
3341                    DecisionOrigin::NoResponsibleActor { .. }
3342                )
3343            {
3344                return Err(CanwuError::new(
3345                    ErrorCode::InvalidAuthority,
3346                    "system commands require a canonical system ID and typed no-responsible-actor origin",
3347                ));
3348            }
3349            Ok(())
3350        }
3351    }
3352}
3353
3354fn validate_command_authority(
3355    authority: &CommandAuthority,
3356    entity_exists: &dyn Fn(&EntityRef) -> bool,
3357) -> Result<(), CanwuError> {
3358    if authority
3359        .seat_id
3360        .as_ref()
3361        .is_some_and(|value| !canonical_text(value))
3362        || authority
3363            .permission_profile_id
3364            .as_ref()
3365            .is_some_and(|value| !canonical_text(value))
3366        || authority.seat_id.is_some() != authority.permission_profile_id.is_some()
3367        || authority
3368            .command_subject
3369            .as_ref()
3370            .is_some_and(|entity| !entity_exists(entity))
3371    {
3372        return Err(CanwuError::new(
3373            ErrorCode::InvalidAuthority,
3374            "command authority contains an invalid seat, permission profile, or subject",
3375        ));
3376    }
3377    match &authority.decision_origin {
3378        DecisionOrigin::Actor { actor } => {
3379            if !entity_exists(&EntityRef::Person(*actor)) {
3380                return Err(CanwuError::new(
3381                    ErrorCode::InvalidAuthority,
3382                    "command decision origin references a missing actor",
3383                ));
3384            }
3385        }
3386        DecisionOrigin::Institution {
3387            institution,
3388            responsible_actor,
3389        } => {
3390            if !entity_exists(institution)
3391                || responsible_actor.is_some_and(|actor| !entity_exists(&EntityRef::Person(actor)))
3392            {
3393                return Err(CanwuError::new(
3394                    ErrorCode::InvalidAuthority,
3395                    "command decision origin references a missing institution or actor",
3396                ));
3397            }
3398        }
3399        DecisionOrigin::Council { council_id } if !canonical_text(council_id) => {
3400            return Err(CanwuError::new(
3401                ErrorCode::InvalidAuthority,
3402                "command council origin requires a canonical ID",
3403            ));
3404        }
3405        DecisionOrigin::NoResponsibleActor { reason } if !canonical_text(reason) => {
3406            return Err(CanwuError::new(
3407                ErrorCode::InvalidAuthority,
3408                "no-responsible-actor origins require a canonical reason",
3409            ));
3410        }
3411        DecisionOrigin::Council { .. } | DecisionOrigin::NoResponsibleActor { .. } => {}
3412    }
3413    Ok(())
3414}
3415
3416fn authority_matches_seat_binding(
3417    seat: SeatPolicy,
3418    binding: &SeatBinding,
3419    authority: &CommandAuthority,
3420) -> bool {
3421    match (seat, &authority.decision_origin) {
3422        (SeatPolicy::CharacterBound, DecisionOrigin::Actor { actor }) => {
3423            binding.actor == Some(*actor) && binding.institution.is_none()
3424        }
3425        (
3426            SeatPolicy::InstitutionBound,
3427            DecisionOrigin::Institution {
3428                institution,
3429                responsible_actor,
3430            },
3431        ) => {
3432            binding.institution.as_ref() == Some(institution)
3433                && binding
3434                    .actor
3435                    .is_none_or(|actor| Some(actor) == *responsible_actor)
3436        }
3437        (SeatPolicy::ObserverSeat | SeatPolicy::AdvisorSeat, origin) => {
3438            let actor_matches = binding.actor.is_none_or(
3439                |expected| matches!(origin, DecisionOrigin::Actor { actor } if *actor == expected),
3440            );
3441            let institution_matches = binding.institution.as_ref().is_none_or(|expected| {
3442                matches!(
3443                    origin,
3444                    DecisionOrigin::Institution { institution, .. } if institution == expected
3445                )
3446            });
3447            actor_matches && institution_matches
3448        }
3449        _ => false,
3450    }
3451}
3452
3453const fn decision_actor(authority: &CommandAuthority) -> Option<PersonId> {
3454    match &authority.decision_origin {
3455        DecisionOrigin::Actor { actor } => Some(*actor),
3456        DecisionOrigin::Institution {
3457            responsible_actor, ..
3458        } => *responsible_actor,
3459        DecisionOrigin::Council { .. } | DecisionOrigin::NoResponsibleActor { .. } => None,
3460    }
3461}
3462
3463const fn is_expected_command_rejection(code: &ErrorCode) -> bool {
3464    matches!(
3465        code,
3466        ErrorCode::ActorNotFound
3467            | ErrorCode::ArmyNotFound
3468            | ErrorCode::DestinationNotFound
3469            | ErrorCode::EntityNotFound
3470            | ErrorCode::IdempotencyConflict
3471            | ErrorCode::InteractionReadOnly
3472            | ErrorCode::InvalidAuthority
3473            | ErrorCode::InvalidDuration
3474            | ErrorCode::InvalidPayload
3475            | ErrorCode::MissingIdempotencyKey
3476            | ErrorCode::MixedCommandIngress
3477            | ErrorCode::NoRoute
3478            | ErrorCode::PluginCommandNotFound
3479            | ErrorCode::SimulationRevisionConflict
3480            | ErrorCode::SimulationTimeConflict
3481            | ErrorCode::ValueOutOfRange
3482    )
3483}
3484
3485fn canonical_text(value: &str) -> bool {
3486    !value.is_empty() && value == value.trim()
3487}
3488
3489fn component_key(
3490    plugin: &str,
3491    state: &StateKey,
3492    entity: &EntityRef,
3493    component: &str,
3494) -> PluginComponentKey {
3495    PluginComponentKey {
3496        plugin: plugin.to_owned(),
3497        state: state.clone(),
3498        entity: entity.clone(),
3499        component: component.to_owned(),
3500    }
3501}
3502
3503fn record_change_affected_entities(change: &DomainRecordChange) -> Vec<EntityRef> {
3504    (change.current.class == DomainRecordClass::Entity)
3505        .then(|| EntityRef::Domain(change.current.reference.clone()))
3506        .into_iter()
3507        .collect()
3508}
3509
3510fn is_domain_record_state(schemas: &records::DomainRecordSchemas, state: &StateKey) -> bool {
3511    schemas.contains_key(&DomainRecordKind::new(&state.namespace, &state.name))
3512}
3513
3514fn snapshot_command_attempt_preflight_error(
3515    snapshot: &SimulationSnapshot,
3516    attempt: &CommandAttemptRecord,
3517    history: &DomainRecordHistory,
3518    cut: DomainHistoryCut,
3519) -> Option<CanwuError> {
3520    let authority = match resolve_command_authority(&attempt.envelope) {
3521        Ok(authority) => authority,
3522        Err(error) => return Some(error),
3523    };
3524    let Some(run_configuration) = snapshot.run_configuration.as_ref() else {
3525        return Some(invalid_snapshot_error(
3526            "snapshot run configuration is required before command attempts",
3527        ));
3528    };
3529    if let Err(error) = validate_command_ingress_policy(
3530        run_configuration,
3531        &attempt.envelope.issuer,
3532        &authority,
3533        CommandAdmission {
3534            request_id: attempt.request_id,
3535            expected_revision: attempt.expected_revision,
3536            expected_time: attempt.envelope.expected_time,
3537            revision_before: attempt.revision_before,
3538            ingress: attempt.ingress,
3539        },
3540        &|entity| snapshot_entity_exists_in_history(snapshot, history, cut, entity),
3541    ) {
3542        return Some(error);
3543    }
3544    attempt.envelope.expected_time.and_then(|expected_time| {
3545        (expected_time != attempt.at).then(|| {
3546            CanwuError::new(
3547                ErrorCode::SimulationTimeConflict,
3548                format!(
3549                    "command expected time {expected_time}, but simulation is at {}",
3550                    attempt.at
3551                ),
3552            )
3553        })
3554    })
3555}
3556
3557fn invalid_snapshot_error(message: impl Into<String>) -> CanwuError {
3558    CanwuError::new(ErrorCode::InvalidSnapshot, message)
3559}
3560
3561fn invalid_snapshot<T>(message: impl Into<String>) -> Result<T, CanwuError> {
3562    Err(invalid_snapshot_error(message))
3563}
3564
3565#[cfg(test)]
3566mod tests;