Skip to main content

canwu_sim/runtime/
mod.rs

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