Skip to main content

canwu_sim/runtime/
mod.rs

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