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