1use super::{
2 ArmyId, BoundaryReceipt, BoundaryRequest, CanwuError, CauseRef, CommandAttemptId, CommandId,
3 CommandPolicyContext, CommandRequestId, CommandTransactionCheckpoint, Deserialize, EntityRef,
4 ErrorCode, EventId, IngressId, IngressTransactionCheckpoint, InteractionPolicy, LetterId,
5 PayloadSchema, PersonId, RejectionTransactionCheckpoint, Serialize, SimDuration, SimTime,
6 Simulation, SystemCadence, TerritoryId, Value, canonical_hash, claim_counter,
7 invalid_snapshot_error, is_expected_command_rejection, resolve_command_authority,
8 runtime_entity_exists, runtime_entity_identity_exists, runtime_has_unqueued_command_history,
9 validate_command_ingress_policy, validate_runtime_cause,
10};
11use std::cmp::Reverse;
12
13#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
14#[serde(tag = "type", content = "id", rename_all = "snake_case")]
15pub enum Issuer {
16 Actor(PersonId),
17 Human(String),
18 Ai(String),
19 Institution(String),
20 Replay(String),
21 Experiment(String),
22 Debug,
23 System(String),
24}
25
26#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
27#[serde(rename_all = "snake_case")]
28pub enum CommandIngress {
29 LegacyDirect,
30 LiveRequest,
31 FrozenReplay,
32}
33
34#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
35#[serde(tag = "type", rename_all = "snake_case")]
36pub enum DecisionOrigin {
37 Actor {
38 actor: PersonId,
39 },
40 Institution {
41 institution: EntityRef,
42 responsible_actor: Option<PersonId>,
43 },
44 Council {
45 council_id: String,
46 },
47 NoResponsibleActor {
48 reason: String,
49 },
50}
51
52#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
53pub struct CommandAuthority {
54 pub decision_origin: DecisionOrigin,
55 pub seat_id: Option<String>,
56 pub permission_profile_id: Option<String>,
57 pub command_subject: Option<EntityRef>,
58}
59
60impl CommandAuthority {
61 #[must_use]
62 pub const fn for_actor(actor: PersonId) -> Self {
63 Self {
64 decision_origin: DecisionOrigin::Actor { actor },
65 seat_id: None,
66 permission_profile_id: None,
67 command_subject: None,
68 }
69 }
70
71 #[must_use]
72 pub fn no_responsible_actor(reason: impl Into<String>) -> Self {
73 Self {
74 decision_origin: DecisionOrigin::NoResponsibleActor {
75 reason: reason.into(),
76 },
77 seat_id: None,
78 permission_profile_id: None,
79 command_subject: None,
80 }
81 }
82}
83
84#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
85pub struct CommandContext {
86 pub issuer: Issuer,
87 pub authority: CommandAuthority,
88 #[serde(default, skip_serializing_if = "Option::is_none")]
91 pub decision_controller_id: Option<String>,
92 pub run_policy: CommandPolicyContext,
93 pub ingress: CommandIngress,
94 pub attempt_id: Option<CommandAttemptId>,
95 pub command_id: CommandId,
96 pub request_id: Option<CommandRequestId>,
97 pub revision: u64,
98 pub simulation_time: SimTime,
99 pub expected_revision: Option<u64>,
100 pub expected_time: Option<SimTime>,
101}
102
103#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
104#[serde(tag = "type", rename_all = "snake_case")]
105pub enum Command {
106 OrderMovement {
107 subject: EntityRef,
108 destination: TerritoryId,
109 #[serde(default, skip_serializing_if = "Vec::is_empty")]
110 cargo: Vec<LetterId>,
111 },
112 DebugSetArmyMorale {
113 army: ArmyId,
114 morale: u16,
115 },
116 Plugin {
117 plugin: String,
118 command: String,
119 payload: Value,
120 },
121}
122
123#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
124pub struct CommandEnvelope {
125 pub issuer: Issuer,
126 #[serde(default, skip_serializing_if = "Option::is_none")]
127 pub authority: Option<CommandAuthority>,
128 pub command: Command,
129 pub expected_time: Option<SimTime>,
130}
131
132impl CommandEnvelope {
133 #[must_use]
134 pub const fn new(issuer: Issuer, command: Command) -> Self {
135 Self {
136 issuer,
137 authority: None,
138 command,
139 expected_time: None,
140 }
141 }
142
143 #[must_use]
144 pub const fn at_time(mut self, expected_time: SimTime) -> Self {
145 self.expected_time = Some(expected_time);
146 self
147 }
148
149 #[must_use]
150 pub fn with_authority(mut self, authority: CommandAuthority) -> Self {
151 self.authority = Some(authority);
152 self
153 }
154}
155
156#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
157pub struct CommandRequest {
158 pub request_id: CommandRequestId,
159 pub expected_revision: u64,
165 pub envelope: CommandEnvelope,
166}
167
168impl CommandRequest {
169 #[must_use]
170 pub const fn new(
171 request_id: CommandRequestId,
172 expected_revision: u64,
173 envelope: CommandEnvelope,
174 ) -> Self {
175 Self {
176 request_id,
177 expected_revision,
178 envelope,
179 }
180 }
181}
182
183#[derive(Clone, Copy)]
184pub(super) struct CommandAdmission {
185 pub(super) request_id: Option<CommandRequestId>,
186 pub(super) expected_revision: Option<u64>,
187 pub(super) expected_time: Option<SimTime>,
188 pub(super) revision_before: u64,
189 pub(super) ingress: CommandIngress,
190}
191
192#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
193pub struct CommandRecord {
194 pub id: CommandId,
195 #[serde(default, skip_serializing_if = "Option::is_none")]
196 pub attempt_id: Option<CommandAttemptId>,
197 pub accepted_at: SimTime,
198 pub envelope: CommandEnvelope,
199 #[serde(default, skip_serializing_if = "Vec::is_empty")]
200 pub emitted_events: Vec<EventId>,
201}
202
203#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
204pub struct CommandReceipt {
205 pub attempt_id: Option<CommandAttemptId>,
206 pub command_id: CommandId,
207 pub request_id: Option<CommandRequestId>,
208 pub revision: u64,
210 pub accepted_at: SimTime,
211 pub emitted_events: Vec<EventId>,
212}
213
214#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
215pub struct CommandRejection {
216 pub attempt_id: Option<CommandAttemptId>,
217 pub request_id: Option<CommandRequestId>,
218 pub retained_revision: u64,
221 pub rejected_at: SimTime,
222 pub error: CanwuError,
223}
224
225#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
226#[serde(tag = "decision", rename_all = "snake_case")]
227pub enum CommandOutcome {
228 Accepted { receipt: CommandReceipt },
229 Rejected { rejection: CommandRejection },
230}
231
232#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
233#[serde(tag = "decision", rename_all = "snake_case")]
234pub enum CommandAttemptOutcome {
235 Accepted { command_id: CommandId },
236 Rejected { error: CanwuError },
237}
238
239#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
240pub struct CommandAttemptRecord {
241 pub id: CommandAttemptId,
242 pub at: SimTime,
243 pub revision_before: u64,
245 pub ingress: CommandIngress,
246 pub request_id: Option<CommandRequestId>,
247 pub expected_revision: Option<u64>,
248 pub envelope: CommandEnvelope,
249 pub outcome: CommandAttemptOutcome,
250}
251
252#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
253#[serde(rename_all = "snake_case")]
254pub enum IngressClass {
255 Command,
256 Communication,
257 Acknowledgement,
258 Information,
259 Decision,
260 ScheduledSystem,
261}
262
263#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
264pub struct PluginIngressDescriptor {
265 pub name: String,
266 pub description: String,
267 pub class: IngressClass,
268 pub payload_schema: PayloadSchema,
269}
270
271#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
272pub struct PluginIngressRequest {
273 pub plugin: String,
274 pub packet_type: String,
275 pub due_at: SimTime,
276 pub priority: i32,
277 pub payload: Value,
278 pub affected_entities: Vec<EntityRef>,
279 #[serde(default, skip_serializing_if = "Option::is_none")]
280 pub cause: Option<CauseRef>,
281 #[serde(default, skip_serializing_if = "Vec::is_empty")]
282 pub archive_retention: Vec<PluginArchiveRetention>,
283}
284
285#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
286#[serde(deny_unknown_fields)]
287pub struct PluginArchiveRetention {
288 pub namespace: String,
289 pub object_id: String,
290}
291
292#[derive(Clone, Debug, Eq, PartialEq)]
296pub struct PluginIngressPermit {
297 pub(super) plugin: String,
298 pub(super) packet_type: String,
299 pub(super) semantic_hash: String,
300 pub(super) token: String,
301}
302
303#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
304#[serde(tag = "maintenance", rename_all = "snake_case")]
305pub enum MaintenanceIngressRequest {
306 DecisionArchive {
307 commit: super::VerifiedDecisionArchiveCommit,
308 },
309 OwnerAuthorized {
310 commit: super::VerifiedOwnerAuthorizedMaintenanceCommit,
311 },
312}
313
314#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
315#[serde(rename_all = "snake_case")]
316pub enum MaintenanceDisposition {
317 Applied,
318 RejectedStale,
319}
320
321#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
322#[serde(deny_unknown_fields)]
323pub struct MaintenanceRejectionReceipt {
324 pub token: String,
325 pub expected_source_root: String,
326 pub observed_source_root: String,
327 pub reason: String,
328}
329
330#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
331#[serde(deny_unknown_fields)]
332pub struct MaintenanceChangeRecord {
333 pub kind: String,
334 pub token: String,
335 pub disposition: MaintenanceDisposition,
336 pub source_root: String,
337 pub target_root: String,
338 #[serde(default, skip_serializing_if = "Option::is_none")]
339 pub rejection: Option<MaintenanceRejectionReceipt>,
340}
341
342impl PluginIngressRequest {
343 #[must_use]
344 pub fn new(
345 plugin: impl Into<String>,
346 packet_type: impl Into<String>,
347 due_at: SimTime,
348 payload: Value,
349 ) -> Self {
350 Self {
351 plugin: plugin.into(),
352 packet_type: packet_type.into(),
353 due_at,
354 priority: 0,
355 payload,
356 affected_entities: Vec::new(),
357 cause: None,
358 archive_retention: Vec::new(),
359 }
360 }
361
362 #[must_use]
363 pub const fn with_priority(mut self, priority: i32) -> Self {
364 self.priority = priority;
365 self
366 }
367
368 #[must_use]
369 pub fn with_entity(mut self, entity: EntityRef) -> Self {
370 self.affected_entities.push(entity);
371 self
372 }
373
374 #[must_use]
375 pub fn caused_by(mut self, cause: CauseRef) -> Self {
376 self.cause = Some(cause);
377 self
378 }
379
380 #[must_use]
381 pub fn with_archive_retention(
382 mut self,
383 retention: impl IntoIterator<Item = PluginArchiveRetention>,
384 ) -> Self {
385 self.archive_retention.extend(retention);
386 self
387 }
388}
389
390#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
391#[serde(tag = "type", rename_all = "snake_case")]
392pub enum IngressPayload {
393 Command {
394 request: Box<CommandRequest>,
395 },
396 Plugin {
397 plugin: String,
398 packet_type: String,
399 payload: Value,
400 affected_entities: Vec<EntityRef>,
401 #[serde(default, skip_serializing_if = "Vec::is_empty")]
402 archive_retention: Vec<PluginArchiveRetention>,
403 },
404 Calendar {
405 cadences: Vec<SystemCadence>,
406 },
407 Decision {
408 request: Box<super::DecisionIngressRequest>,
409 },
410 Maintenance {
411 request: Box<MaintenanceIngressRequest>,
412 },
413}
414
415#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
416pub struct IngressRecord {
417 pub id: IngressId,
418 pub issued_at: SimTime,
419 #[serde(default, skip_serializing_if = "is_zero")]
420 pub eligible_boundary_count: u64,
421 pub due_at: SimTime,
422 pub class: IngressClass,
423 pub priority: i32,
424 pub payload: IngressPayload,
425 #[serde(default, skip_serializing_if = "Option::is_none")]
426 pub cause: Option<CauseRef>,
427}
428
429#[allow(clippy::trivially_copy_pass_by_ref)]
430const fn is_zero(value: &u64) -> bool {
431 *value == 0
432}
433
434#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
435pub struct IngressReceipt {
436 pub ingress_id: IngressId,
437 pub issued_at: SimTime,
438 pub due_at: SimTime,
439}
440
441#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
442pub(crate) struct IngressQueueKey {
443 pub due_at: SimTime,
444 pub class: IngressClass,
445 pub priority: Reverse<i32>,
446 pub issued_at: SimTime,
447 pub id: IngressId,
448}
449
450impl IngressQueueKey {
451 #[must_use]
452 pub(crate) const fn from_record(record: &IngressRecord) -> Self {
453 Self {
454 due_at: record.due_at,
455 class: record.class,
456 priority: Reverse(record.priority),
457 issued_at: record.issued_at,
458 id: record.id,
459 }
460 }
461}
462
463impl Simulation {
464 pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
465 match self.admit_command(
466 None,
467 None,
468 envelope,
469 CommandIngress::LegacyDirect,
470 None,
471 false,
472 )? {
473 CommandOutcome::Accepted { receipt } => Ok(receipt),
474 CommandOutcome::Rejected { rejection } => Err(rejection.error),
475 }
476 }
477
478 pub fn enqueue_command(
479 &mut self,
480 due_at: SimTime,
481 priority: i32,
482 request: CommandRequest,
483 ) -> Result<IngressReceipt, CanwuError> {
484 self.ensure_runtime_ready()?;
485 self.ensure_canonical_ingress_can_start()?;
486 self.ensure_command_ingress_family(CommandIngress::LiveRequest)?;
487 if let Some(existing) = self
488 .state
489 .evidence
490 .archived_ingress_requests
491 .get(&request.request_id)
492 {
493 let input_hash = canonical_hash(
494 "canwu.archive.ingress.command.v1",
495 &(due_at, priority, &request),
496 )?;
497 if existing.input_hash == input_hash {
498 return Ok(existing.receipt.clone());
499 }
500 return Err(CanwuError::new(
501 ErrorCode::IdempotencyConflict,
502 format!(
503 "command request {} is already queued with different ingress content",
504 request.request_id
505 ),
506 ));
507 }
508 for record in &self.state.evidence.ingress {
509 let IngressPayload::Command { request: existing } = &record.payload else {
510 continue;
511 };
512 if existing.request_id != request.request_id {
513 continue;
514 }
515 if existing.as_ref() == &request
516 && record.due_at == due_at
517 && record.priority == priority
518 {
519 return Ok(IngressReceipt {
520 ingress_id: record.id,
521 issued_at: record.issued_at,
522 due_at: record.due_at,
523 });
524 }
525 return Err(CanwuError::new(
526 ErrorCode::IdempotencyConflict,
527 format!(
528 "command request {} is already queued with different ingress content",
529 request.request_id
530 ),
531 ));
532 }
533 if self.command_request_id_is_in_use(request.request_id) {
534 return Err(CanwuError::new(
535 ErrorCode::IdempotencyConflict,
536 format!(
537 "command request {} is already reserved or processed",
538 request.request_id
539 ),
540 ));
541 }
542 if request
543 .envelope
544 .expected_time
545 .is_some_and(|expected| expected != due_at)
546 {
547 return Err(CanwuError::new(
548 ErrorCode::SimulationTimeConflict,
549 "queued command expected time must equal its due simulation time",
550 ));
551 }
552 self.append_ingress(
553 due_at,
554 IngressClass::Command,
555 priority,
556 IngressPayload::Command {
557 request: Box::new(request),
558 },
559 None,
560 false,
561 )
562 }
563
564 pub fn enqueue_plugin_ingress(
565 &mut self,
566 request: PluginIngressRequest,
567 ) -> Result<IngressReceipt, CanwuError> {
568 self.enqueue_plugin_ingress_inner(request, None, false)
569 }
570
571 pub fn enqueue_permitted_plugin_ingress(
574 &mut self,
575 request: PluginIngressRequest,
576 permit: &PluginIngressPermit,
577 ) -> Result<IngressReceipt, CanwuError> {
578 self.enqueue_plugin_ingress_inner(request, Some(permit), false)
579 }
580
581 pub(super) fn replay_plugin_ingress(
582 &mut self,
583 request: PluginIngressRequest,
584 ) -> Result<IngressReceipt, CanwuError> {
585 self.enqueue_plugin_ingress_inner(request, None, true)
586 }
587
588 fn enqueue_plugin_ingress_inner(
589 &mut self,
590 mut request: PluginIngressRequest,
591 permit: Option<&PluginIngressPermit>,
592 replay: bool,
593 ) -> Result<IngressReceipt, CanwuError> {
594 self.ensure_runtime_ready()?;
595 self.ensure_canonical_ingress_can_start()?;
596 if !replay
597 && self
598 .state
599 .metadata
600 .run_configuration
601 .declared()
602 .is_some_and(|configuration| {
603 configuration.interaction == InteractionPolicy::ReadOnly
604 })
605 {
606 return Err(CanwuError::new(
607 ErrorCode::InteractionReadOnly,
608 "the run interaction policy rejects newly authored plugin ingress",
609 ));
610 }
611 let key = (request.plugin.clone(), request.packet_type.clone());
612 let descriptor = self.plugins.ingress.get(&key).ok_or_else(|| {
613 CanwuError::new(
614 ErrorCode::InvalidPayload,
615 format!(
616 "plugin ingress type {}.{} is not registered",
617 request.plugin, request.packet_type
618 ),
619 )
620 })?;
621 let internal = self.plugins.internal_ingress.contains(&key);
622 if internal && !replay {
623 let semantic_hash = self
624 .plugins
625 .descriptors
626 .get(&request.plugin)
627 .map(|descriptor| descriptor.semantic_hash.as_str())
628 .ok_or_else(|| {
629 CanwuError::new(
630 ErrorCode::PluginNotActive,
631 "internal plugin ingress owner is unavailable",
632 )
633 })?;
634 let expected_token = super::canonical_hash(
635 "canwu.plugin.internal-ingress-permit.v1",
636 &(&request.plugin, &request.packet_type, semantic_hash),
637 )?;
638 if permit.is_none_or(|permit| {
639 permit.plugin != request.plugin
640 || permit.packet_type != request.packet_type
641 || permit.semantic_hash != semantic_hash
642 || permit.token != expected_token
643 }) {
644 return Err(CanwuError::new(
645 ErrorCode::InvalidAuthority,
646 "plugin-owned internal ingress requires its opaque registration permit",
647 ));
648 }
649 }
650 if !request.archive_retention.is_empty() && !internal && !replay {
651 return Err(CanwuError::new(
652 ErrorCode::InvalidAuthority,
653 "archive retention may be attached only to plugin-owned internal ingress",
654 ));
655 }
656 request.archive_retention.sort();
657 request.archive_retention.dedup();
658 if request.archive_retention.len() > 32_768
659 || request.archive_retention.iter().any(|retention| {
660 retention.namespace.is_empty()
661 || retention.namespace.len() > 128
662 || retention.object_id.is_empty()
663 || retention.object_id.len() > 256
664 || !retention.namespace.bytes().all(|byte| {
665 byte.is_ascii_lowercase()
666 || byte.is_ascii_digit()
667 || matches!(byte, b'.' | b'-' | b'_')
668 })
669 || !retention.object_id.is_ascii()
670 })
671 {
672 return Err(CanwuError::new(
673 ErrorCode::InvalidPayload,
674 "plugin ingress archive retention is malformed or exceeds its hard limit",
675 ));
676 }
677 self.plugins.validate_archive_retention(
678 &request.plugin,
679 &request.packet_type,
680 &request.payload,
681 &request.archive_retention,
682 )?;
683 descriptor.payload_schema.validate(&request.payload)?;
684 request.affected_entities.sort();
685 request.affected_entities.dedup();
686 if request
687 .affected_entities
688 .iter()
689 .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
690 {
691 return Err(CanwuError::new(
692 ErrorCode::EntityNotFound,
693 "plugin ingress references an unknown entity identity",
694 ));
695 }
696 if let Some(cause) = &request.cause {
697 if matches!(
698 cause,
699 CauseRef::Boundary(_) | CauseRef::Command(_) | CauseRef::Event(_)
700 ) {
701 return Err(CanwuError::new(
702 ErrorCode::InvalidPayload,
703 "boundary, command, and event causes are reserved for plugin-generated ingress",
704 ));
705 }
706 validate_runtime_cause(&self.state, cause)?;
707 }
708 self.append_ingress(
709 request.due_at,
710 descriptor.class,
711 request.priority,
712 IngressPayload::Plugin {
713 plugin: request.plugin,
714 packet_type: request.packet_type,
715 payload: request.payload,
716 affected_entities: request.affected_entities,
717 archive_retention: request.archive_retention,
718 },
719 request.cause,
720 false,
721 )
722 }
723
724 pub fn schedule_calendar_boundary(
725 &mut self,
726 due_at: SimTime,
727 mut cadences: Vec<SystemCadence>,
728 ) -> Result<IngressReceipt, CanwuError> {
729 self.ensure_runtime_ready()?;
730 self.ensure_canonical_ingress_can_start()?;
731 if cadences.contains(&SystemCadence::EventDriven) {
732 return Err(CanwuError::new(
733 ErrorCode::InvalidBoundary,
734 "calendar ingress cannot declare event-driven cadence",
735 ));
736 }
737 cadences.sort();
738 cadences.dedup();
739 if cadences.is_empty() {
740 return Err(CanwuError::new(
741 ErrorCode::InvalidBoundary,
742 "calendar ingress requires at least one scheduled cadence",
743 ));
744 }
745 self.append_ingress(
746 due_at,
747 IngressClass::ScheduledSystem,
748 0,
749 IngressPayload::Calendar { cadences },
750 Some(CauseRef::System("canwu.core.calendar".to_owned())),
751 false,
752 )
753 }
754
755 pub(super) fn append_ingress(
756 &mut self,
757 due_at: SimTime,
758 class: IngressClass,
759 priority: i32,
760 payload: IngressPayload,
761 cause: Option<CauseRef>,
762 after_current_boundary: bool,
763 ) -> Result<IngressReceipt, CanwuError> {
764 if due_at < self.state.scheduler.now {
765 return Err(CanwuError::new(
766 ErrorCode::LateIngress,
767 format!(
768 "ingress due at {due_at} cannot be queued after committed time {}",
769 self.state.scheduler.now
770 ),
771 ));
772 }
773 let transaction = IngressTransactionCheckpoint::capture(&self.state);
774 let (id, next_id) = claim_counter(self.state.counters.next_ingress_id, "ingress ID")?;
775 let boundary_count = self
776 .state
777 .evidence
778 .archived
779 .boundary_count
780 .checked_add(
781 u64::try_from(self.state.evidence.boundaries.len()).map_err(|_| {
782 CanwuError::new(
783 ErrorCode::IdentifierExhausted,
784 "boundary count exceeds the ingress journal range",
785 )
786 })?,
787 )
788 .ok_or_else(|| {
789 CanwuError::new(
790 ErrorCode::IdentifierExhausted,
791 "boundary count exceeds the ingress journal range",
792 )
793 })?;
794 let eligible_boundary_count = if after_current_boundary {
795 boundary_count.checked_add(1).ok_or_else(|| {
796 CanwuError::new(
797 ErrorCode::IdentifierExhausted,
798 "ingress boundary eligibility exceeds the journal range",
799 )
800 })?
801 } else {
802 boundary_count
803 };
804 let record = IngressRecord {
805 id: IngressId::new(id),
806 issued_at: self.state.scheduler.now,
807 eligible_boundary_count,
808 due_at,
809 class,
810 priority,
811 payload,
812 cause,
813 };
814 let queue_key = IngressQueueKey::from_record(&record);
815 self.state.counters.next_ingress_id = next_id;
816 self.state.scheduler.pending_ingress.insert(queue_key);
817 self.state.evidence.ingress.push(record.clone());
818 self.state.metadata.plugin_registration_closed = true;
819 if let Err(error) = self.refresh_checkpoint_hash() {
820 transaction.restore(&mut self.state, &queue_key);
821 return Err(error);
822 }
823 Ok(IngressReceipt {
824 ingress_id: record.id,
825 issued_at: record.issued_at,
826 due_at: record.due_at,
827 })
828 }
829
830 pub fn process_command(
831 &mut self,
832 request: CommandRequest,
833 ) -> Result<CommandOutcome, CanwuError> {
834 self.ensure_runtime_ready()?;
835 if self.state.evidence.archived.ingress_count != 0
836 || !self.state.evidence.ingress.is_empty()
837 {
838 return Err(CanwuError::new(
839 ErrorCode::MixedCommandIngress,
840 "direct command requests cannot bypass an active canonical ingress journal",
841 ));
842 }
843 self.admit_command(
844 Some(request.request_id),
845 Some(request.expected_revision),
846 request.envelope,
847 CommandIngress::LiveRequest,
848 None,
849 true,
850 )
851 }
852
853 pub(super) fn admit_command(
854 &mut self,
855 request_id: Option<CommandRequestId>,
856 expected_revision: Option<u64>,
857 envelope: CommandEnvelope,
858 ingress: CommandIngress,
859 decision_controller_id: Option<String>,
860 record_attempt: bool,
861 ) -> Result<CommandOutcome, CanwuError> {
862 self.ensure_runtime_ready()?;
863 self.ensure_command_ingress_family(ingress)?;
864 if let Some(cached) =
865 self.cached_command_outcome(request_id, expected_revision, &envelope)?
866 {
867 return Ok(cached);
868 }
869
870 let revision_before = self.revision();
871 let admission = CommandAdmission {
872 request_id,
873 expected_revision,
874 expected_time: envelope.expected_time,
875 revision_before,
876 ingress,
877 };
878 let attempt_id = if record_attempt {
879 let (value, _) = claim_counter(
880 self.state.counters.next_command_attempt_id,
881 "command attempt ID",
882 )?;
883 CommandAttemptId::new(value)
884 } else {
885 CommandAttemptId::default()
886 };
887 let authority = match resolve_command_authority(&envelope) {
888 Ok(authority) => authority,
889 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
890 return self.record_command_rejection(attempt_id, admission, envelope, error);
891 }
892 Err(error) => return Err(error),
893 };
894 if let Err(error) = self.validate_command_ingress(&envelope.issuer, &authority, admission) {
895 if is_expected_command_rejection(&error.code) && record_attempt {
896 return self.record_command_rejection(attempt_id, admission, envelope, error);
897 }
898 return Err(error);
899 }
900 if let Some(expected_time) = envelope.expected_time
901 && expected_time != self.state.scheduler.now
902 {
903 let error = CanwuError::new(
904 ErrorCode::SimulationTimeConflict,
905 format!(
906 "command expected time {expected_time}, but simulation is at {}",
907 self.state.scheduler.now
908 ),
909 );
910 if record_attempt {
911 return self.record_command_rejection(attempt_id, admission, envelope, error);
912 }
913 return Err(error);
914 }
915
916 let (command_id_value, next_command_id) =
917 claim_counter(self.state.counters.next_command_id, "command ID")?;
918 let (correlation_id, next_correlation_id) =
919 claim_counter(self.state.counters.next_correlation_id, "correlation ID")?;
920 let command_id = CommandId::new(command_id_value);
921 let context = CommandContext {
922 issuer: envelope.issuer.clone(),
923 authority,
924 decision_controller_id,
925 run_policy: self.state.metadata.run_configuration.command_policy(),
926 ingress: admission.ingress,
927 attempt_id: record_attempt.then_some(attempt_id),
928 command_id,
929 request_id: admission.request_id,
930 revision: admission.revision_before,
931 simulation_time: self.state.scheduler.now,
932 expected_revision: admission.expected_revision,
933 expected_time: envelope.expected_time,
934 };
935 let prepared = match self.prepare_command(&envelope, &context) {
936 Ok(prepared) => prepared,
937 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
938 return self.record_command_rejection(attempt_id, admission, envelope, error);
939 }
940 Err(error) => return Err(error),
941 };
942 let next_attempt_id = if record_attempt {
943 let (claimed_id, next_attempt_id) = claim_counter(
944 self.state.counters.next_command_attempt_id,
945 "command attempt ID",
946 )?;
947 if claimed_id != attempt_id.get() {
948 return Err(CanwuError::new(
949 ErrorCode::InvalidSnapshot,
950 "command attempt allocation changed during application",
951 ));
952 }
953 Some(next_attempt_id)
954 } else {
955 None
956 };
957 let revision = self.next_state_revision()?;
958 let transaction = CommandTransactionCheckpoint::capture(&self.state);
959 let event_start = self.state.evidence.events.len();
960 self.state.counters.next_command_id = next_command_id;
961 self.state.counters.next_correlation_id = next_correlation_id;
962 self.invalidate_commitments(prepared.commitment_invalidation());
963
964 if let Err(error) = self.apply_prepared(prepared, command_id, correlation_id) {
965 transaction.restore(&mut self.state);
966 if is_expected_command_rejection(&error.code) && record_attempt {
967 return self.record_command_rejection(attempt_id, admission, envelope, error);
968 }
969 return Err(error);
970 }
971 let emitted_events: Vec<_> = self.state.evidence.events[event_start..]
972 .iter()
973 .map(|event| event.id)
974 .collect();
975 self.state.metadata.plugin_registration_closed = true;
976 self.state.evidence.commands.push(CommandRecord {
977 id: command_id,
978 attempt_id: record_attempt.then_some(attempt_id),
979 accepted_at: self.state.scheduler.now,
980 envelope: envelope.clone(),
981 emitted_events: if record_attempt {
982 emitted_events.clone()
983 } else {
984 Vec::new()
985 },
986 });
987 if let Some(next_attempt_id) = next_attempt_id {
988 self.state.counters.next_command_attempt_id = next_attempt_id;
989 self.state
990 .evidence
991 .command_attempts
992 .push(CommandAttemptRecord {
993 id: attempt_id,
994 at: self.state.scheduler.now,
995 revision_before: admission.revision_before,
996 ingress: admission.ingress,
997 request_id: admission.request_id,
998 expected_revision: admission.expected_revision,
999 envelope,
1000 outcome: CommandAttemptOutcome::Accepted { command_id },
1001 });
1002 }
1003 self.state.counters.state_revision = revision;
1004 if let Err(error) = self.refresh_checkpoint_hash() {
1005 transaction.restore(&mut self.state);
1006 return Err(error);
1007 }
1008
1009 Ok(CommandOutcome::Accepted {
1010 receipt: CommandReceipt {
1011 attempt_id: record_attempt.then_some(attempt_id),
1012 command_id,
1013 request_id: admission.request_id,
1014 revision,
1015 accepted_at: self.state.scheduler.now,
1016 emitted_events,
1017 },
1018 })
1019 }
1020
1021 fn ensure_command_ingress_family(&self, ingress: CommandIngress) -> Result<(), CanwuError> {
1022 let has_legacy_commands = self.state.evidence.archived_legacy_commands
1023 || self
1024 .state
1025 .evidence
1026 .commands
1027 .iter()
1028 .any(|record| record.attempt_id.is_none());
1029 let has_tracked_attempts = self.state.evidence.archived_tracked_attempts
1030 || !self.state.evidence.command_attempts.is_empty()
1031 || !self.state.evidence.ingress.is_empty();
1032 if (ingress == CommandIngress::LegacyDirect && has_tracked_attempts)
1033 || (ingress != CommandIngress::LegacyDirect && has_legacy_commands)
1034 {
1035 return Err(CanwuError::new(
1036 ErrorCode::MixedCommandIngress,
1037 "legacy-direct commands and tracked request/replay attempts cannot coexist in one run",
1038 ));
1039 }
1040 Ok(())
1041 }
1042
1043 pub(super) fn ensure_canonical_ingress_can_start(&self) -> Result<(), CanwuError> {
1044 if runtime_has_unqueued_command_history(&self.state) {
1045 return Err(CanwuError::new(
1046 ErrorCode::MixedCommandIngress,
1047 "canonical ingress cannot be added after direct command history",
1048 ));
1049 }
1050 Ok(())
1051 }
1052
1053 fn cached_command_outcome(
1054 &self,
1055 request_id: Option<CommandRequestId>,
1056 expected_revision: Option<u64>,
1057 envelope: &CommandEnvelope,
1058 ) -> Result<Option<CommandOutcome>, CanwuError> {
1059 let Some(request_id) = request_id else {
1060 return Ok(None);
1061 };
1062 if let Some(cached) = self
1063 .state
1064 .evidence
1065 .archived_command_requests
1066 .get(&request_id)
1067 {
1068 let input_hash = canonical_hash(
1069 "canwu.archive.command.request.v1",
1070 &(expected_revision, envelope),
1071 )?;
1072 if cached.input_hash != input_hash {
1073 return Ok(Some(CommandOutcome::Rejected {
1074 rejection: CommandRejection {
1075 attempt_id: None,
1076 request_id: Some(request_id),
1077 retained_revision: self.revision(),
1078 rejected_at: self.state.scheduler.now,
1079 error: CanwuError::new(
1080 ErrorCode::IdempotencyConflict,
1081 "this command request ID was already used for different input",
1082 ),
1083 },
1084 }));
1085 }
1086 return Ok(Some(cached.outcome.clone()));
1087 }
1088 let Some(attempt) = self
1089 .state
1090 .evidence
1091 .command_attempts
1092 .iter()
1093 .find(|attempt| attempt.request_id == Some(request_id))
1094 else {
1095 return Ok(None);
1096 };
1097 if attempt.expected_revision != expected_revision || &attempt.envelope != envelope {
1098 return Ok(Some(CommandOutcome::Rejected {
1099 rejection: CommandRejection {
1100 attempt_id: None,
1101 request_id: Some(request_id),
1102 retained_revision: self.revision(),
1103 rejected_at: self.state.scheduler.now,
1104 error: CanwuError::new(
1105 ErrorCode::IdempotencyConflict,
1106 "this command request ID was already used for different input",
1107 ),
1108 },
1109 }));
1110 }
1111 Ok(Some(self.command_outcome_from_attempt(attempt)?))
1112 }
1113
1114 pub(super) fn command_outcome_from_attempt(
1115 &self,
1116 attempt: &CommandAttemptRecord,
1117 ) -> Result<CommandOutcome, CanwuError> {
1118 let request_id = attempt.request_id.ok_or_else(|| {
1119 invalid_snapshot_error("tracked command attempt is missing its request ID")
1120 })?;
1121 let committed_revision = attempt.revision_before.checked_add(1).ok_or_else(|| {
1122 invalid_snapshot_error("cached command attempt revision is exhausted")
1123 })?;
1124 match &attempt.outcome {
1125 CommandAttemptOutcome::Accepted { command_id } => {
1126 let retained_number = command_id
1127 .get()
1128 .checked_sub(self.state.evidence.archived.command_count)
1129 .and_then(|value| value.checked_sub(1))
1130 .ok_or_else(|| {
1131 invalid_snapshot_error(
1132 "accepted command attempt references archived command evidence",
1133 )
1134 })?;
1135 let index = usize::try_from(retained_number).map_err(|_| {
1136 invalid_snapshot_error(
1137 "accepted command attempt exceeds the retained command index space",
1138 )
1139 })?;
1140 let record = self
1141 .state
1142 .evidence
1143 .commands
1144 .get(index)
1145 .filter(|record| record.id == *command_id)
1146 .ok_or_else(|| {
1147 invalid_snapshot_error(
1148 "accepted command attempt references a missing command",
1149 )
1150 })?;
1151 Ok(CommandOutcome::Accepted {
1152 receipt: CommandReceipt {
1153 attempt_id: Some(attempt.id),
1154 command_id: *command_id,
1155 request_id: Some(request_id),
1156 revision: committed_revision,
1157 accepted_at: record.accepted_at,
1158 emitted_events: record.emitted_events.clone(),
1159 },
1160 })
1161 }
1162 CommandAttemptOutcome::Rejected { error } => Ok(CommandOutcome::Rejected {
1163 rejection: CommandRejection {
1164 attempt_id: Some(attempt.id),
1165 request_id: Some(request_id),
1166 retained_revision: committed_revision,
1167 rejected_at: attempt.at,
1168 error: error.clone(),
1169 },
1170 }),
1171 }
1172 }
1173
1174 fn record_command_rejection(
1175 &mut self,
1176 attempt_id: CommandAttemptId,
1177 admission: CommandAdmission,
1178 envelope: CommandEnvelope,
1179 error: CanwuError,
1180 ) -> Result<CommandOutcome, CanwuError> {
1181 let (claimed_id, next_attempt_id) = claim_counter(
1182 self.state.counters.next_command_attempt_id,
1183 "command attempt ID",
1184 )?;
1185 if claimed_id != attempt_id.get() {
1186 return Err(CanwuError::new(
1187 ErrorCode::InvalidSnapshot,
1188 "command attempt allocation changed during rejection",
1189 ));
1190 }
1191 let revision = self.next_state_revision()?;
1192 let attempt = CommandAttemptRecord {
1193 id: attempt_id,
1194 at: self.state.scheduler.now,
1195 revision_before: admission.revision_before,
1196 ingress: admission.ingress,
1197 request_id: admission.request_id,
1198 expected_revision: admission.expected_revision,
1199 envelope,
1200 outcome: CommandAttemptOutcome::Rejected {
1201 error: error.clone(),
1202 },
1203 };
1204 let transaction = RejectionTransactionCheckpoint::capture(&self.state);
1205 self.state.counters.next_command_attempt_id = next_attempt_id;
1206 self.state.counters.state_revision = revision;
1207 self.state.metadata.plugin_registration_closed = true;
1208 self.state.evidence.command_attempts.push(attempt);
1209 if let Err(hash_error) = self.refresh_checkpoint_hash() {
1210 transaction.restore(&mut self.state);
1211 return Err(hash_error);
1212 }
1213 Ok(CommandOutcome::Rejected {
1214 rejection: CommandRejection {
1215 attempt_id: Some(attempt_id),
1216 request_id: admission.request_id,
1217 retained_revision: revision,
1218 rejected_at: self.state.scheduler.now,
1219 error,
1220 },
1221 })
1222 }
1223
1224 fn validate_command_ingress(
1225 &self,
1226 issuer: &Issuer,
1227 authority: &CommandAuthority,
1228 admission: CommandAdmission,
1229 ) -> Result<(), CanwuError> {
1230 validate_command_ingress_policy(
1231 &self.state.metadata.run_configuration,
1232 issuer,
1233 authority,
1234 admission,
1235 &|entity| runtime_entity_exists(&self.state, entity),
1236 )
1237 }
1238
1239 pub fn advance_canonical(
1240 &mut self,
1241 duration: SimDuration,
1242 ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
1243 self.ensure_runtime_ready()?;
1244 if duration.is_negative() {
1245 return Err(CanwuError::new(
1246 ErrorCode::InvalidDuration,
1247 "canonical simulation time cannot advance by a negative duration",
1248 ));
1249 }
1250 let target = self
1251 .state
1252 .scheduler
1253 .now
1254 .checked_add(duration)
1255 .ok_or_else(|| {
1256 CanwuError::new(
1257 ErrorCode::InvalidDuration,
1258 "canonical simulation target time exceeds the supported range",
1259 )
1260 })?;
1261 let mut receipts = Vec::new();
1262 while let Some(next_due) = self.next_canonical_due_time()
1263 && next_due <= target
1264 {
1265 let at = next_due.max(self.state.scheduler.now);
1266 receipts.push(self.settle_boundary(BoundaryRequest::at(at))?);
1267 }
1268 if self.state.scheduler.now < target {
1269 self.advance_to(target)?;
1270 }
1271 Ok(receipts)
1272 }
1273
1274 pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
1275 self.ensure_runtime_ready()?;
1276 let Some(next_due) = self.next_canonical_due_time() else {
1277 return Ok(None);
1278 };
1279 self.settle_boundary(BoundaryRequest::at(next_due.max(self.state.scheduler.now)))
1280 .map(Some)
1281 }
1282
1283 fn next_canonical_due_time(&self) -> Option<SimTime> {
1284 let scheduled = self.state.scheduler.actions.keys().next().map(|key| key.at);
1285 let ingress = self
1286 .state
1287 .scheduler
1288 .pending_ingress
1289 .first()
1290 .map(|key| key.due_at);
1291 match (scheduled, ingress) {
1292 (Some(left), Some(right)) => Some(left.min(right)),
1293 (Some(value), None) | (None, Some(value)) => Some(value),
1294 (None, None) => None,
1295 }
1296 }
1297
1298 pub(super) fn take_due_ingress(&mut self, at: SimTime) -> Vec<IngressId> {
1299 let mut admitted = Vec::new();
1300 while self
1301 .state
1302 .scheduler
1303 .pending_ingress
1304 .first()
1305 .is_some_and(|key| key.due_at <= at)
1306 {
1307 let key = self
1308 .state
1309 .scheduler
1310 .pending_ingress
1311 .pop_first()
1312 .expect("pending ingress was checked as non-empty");
1313 admitted.push(key.id);
1314 }
1315 admitted
1316 }
1317}