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