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