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