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