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