1use super::{
2 ArmyId, BoundaryIngressGeneration, BoundaryReceipt, BoundaryRequest, CanwuError, CauseRef,
3 CommandAttemptId, CommandId, CommandPolicyContext, CommandRequestId,
4 CommandTransactionCheckpoint, Deserialize, EntityRef, ErrorCode, EventId, IngressId,
5 IngressTransactionCheckpoint, InteractionPolicy, LetterId, PayloadSchema, PersonId,
6 RejectionTransactionCheckpoint, Serialize, SimDuration, SimTime, Simulation, SystemCadence,
7 TerritoryId, Value, canonical_hash, claim_counter, invalid_snapshot_error,
8 is_expected_command_rejection, resolve_command_authority, runtime_entity_exists,
9 runtime_entity_identity_exists, runtime_has_unqueued_command_history,
10 validate_command_ingress_policy, validate_runtime_cause,
11};
12use std::cmp::Reverse;
13
14#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
15#[serde(tag = "type", content = "id", rename_all = "snake_case")]
16pub enum Issuer {
17 Actor(PersonId),
18 Human(String),
19 Ai(String),
20 Institution(String),
21 Replay(String),
22 Experiment(String),
23 Debug,
24 System(String),
25}
26
27#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
28#[serde(rename_all = "snake_case")]
29pub enum CommandIngress {
30 LegacyDirect,
31 LiveRequest,
32 FrozenReplay,
33}
34
35#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
36#[serde(tag = "type", rename_all = "snake_case")]
37pub enum DecisionOrigin {
38 Actor {
39 actor: PersonId,
40 },
41 Institution {
42 institution: EntityRef,
43 responsible_actor: Option<PersonId>,
44 },
45 Council {
46 council_id: String,
47 },
48 NoResponsibleActor {
49 reason: String,
50 },
51}
52
53#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
54pub struct CommandAuthority {
55 pub decision_origin: DecisionOrigin,
56 pub seat_id: Option<String>,
57 pub permission_profile_id: Option<String>,
58 pub command_subject: Option<EntityRef>,
59}
60
61impl CommandAuthority {
62 #[must_use]
63 pub const fn for_actor(actor: PersonId) -> Self {
64 Self {
65 decision_origin: DecisionOrigin::Actor { actor },
66 seat_id: None,
67 permission_profile_id: None,
68 command_subject: None,
69 }
70 }
71
72 #[must_use]
73 pub fn no_responsible_actor(reason: impl Into<String>) -> Self {
74 Self {
75 decision_origin: DecisionOrigin::NoResponsibleActor {
76 reason: reason.into(),
77 },
78 seat_id: None,
79 permission_profile_id: None,
80 command_subject: None,
81 }
82 }
83}
84
85#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
86pub struct CommandContext {
87 pub issuer: Issuer,
88 pub authority: CommandAuthority,
89 #[serde(default, skip_serializing_if = "Option::is_none")]
92 pub decision_controller_id: Option<String>,
93 pub run_policy: CommandPolicyContext,
94 pub ingress: CommandIngress,
95 pub attempt_id: Option<CommandAttemptId>,
96 pub command_id: CommandId,
97 pub request_id: Option<CommandRequestId>,
98 pub revision: u64,
99 pub simulation_time: SimTime,
100 pub expected_revision: Option<u64>,
101 pub expected_time: Option<SimTime>,
102}
103
104#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
105#[serde(tag = "type", rename_all = "snake_case")]
106pub enum Command {
107 OrderMovement {
108 subject: EntityRef,
109 destination: TerritoryId,
110 #[serde(default, skip_serializing_if = "Vec::is_empty")]
111 cargo: Vec<LetterId>,
112 },
113 DebugSetArmyMorale {
114 army: ArmyId,
115 morale: u16,
116 },
117 Plugin {
118 plugin: String,
119 command: String,
120 payload: Value,
121 },
122}
123
124#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
125pub struct CommandEnvelope {
126 pub issuer: Issuer,
127 #[serde(default, skip_serializing_if = "Option::is_none")]
128 pub authority: Option<CommandAuthority>,
129 pub command: Command,
130 pub expected_time: Option<SimTime>,
131}
132
133impl CommandEnvelope {
134 #[must_use]
135 pub const fn new(issuer: Issuer, command: Command) -> Self {
136 Self {
137 issuer,
138 authority: None,
139 command,
140 expected_time: None,
141 }
142 }
143
144 #[must_use]
145 pub const fn at_time(mut self, expected_time: SimTime) -> Self {
146 self.expected_time = Some(expected_time);
147 self
148 }
149
150 #[must_use]
151 pub fn with_authority(mut self, authority: CommandAuthority) -> Self {
152 self.authority = Some(authority);
153 self
154 }
155}
156
157#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
158pub struct CommandRequest {
159 pub request_id: CommandRequestId,
160 pub expected_revision: u64,
166 pub envelope: CommandEnvelope,
167}
168
169impl CommandRequest {
170 #[must_use]
171 pub const fn new(
172 request_id: CommandRequestId,
173 expected_revision: u64,
174 envelope: CommandEnvelope,
175 ) -> Self {
176 Self {
177 request_id,
178 expected_revision,
179 envelope,
180 }
181 }
182}
183
184#[derive(Clone, Copy)]
185pub(super) struct CommandAdmission {
186 pub(super) request_id: Option<CommandRequestId>,
187 pub(super) expected_revision: Option<u64>,
188 pub(super) expected_time: Option<SimTime>,
189 pub(super) revision_before: u64,
190 pub(super) ingress: CommandIngress,
191}
192
193#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
194pub struct CommandRecord {
195 pub id: CommandId,
196 #[serde(default, skip_serializing_if = "Option::is_none")]
197 pub attempt_id: Option<CommandAttemptId>,
198 pub accepted_at: SimTime,
199 pub envelope: CommandEnvelope,
200 #[serde(default, skip_serializing_if = "Vec::is_empty")]
201 pub emitted_events: Vec<EventId>,
202}
203
204#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
205pub struct CommandReceipt {
206 pub attempt_id: Option<CommandAttemptId>,
207 pub command_id: CommandId,
208 pub request_id: Option<CommandRequestId>,
209 pub revision: u64,
211 pub accepted_at: SimTime,
212 pub emitted_events: Vec<EventId>,
213}
214
215#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
216pub struct CommandRejection {
217 pub attempt_id: Option<CommandAttemptId>,
218 pub request_id: Option<CommandRequestId>,
219 pub retained_revision: u64,
222 pub rejected_at: SimTime,
223 pub error: CanwuError,
224}
225
226#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
227#[serde(tag = "decision", rename_all = "snake_case")]
228pub enum CommandOutcome {
229 Accepted { receipt: CommandReceipt },
230 Rejected { rejection: CommandRejection },
231}
232
233#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
234#[serde(tag = "decision", rename_all = "snake_case")]
235pub enum CommandAttemptOutcome {
236 Accepted { command_id: CommandId },
237 Rejected { error: CanwuError },
238}
239
240#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
241pub struct CommandAttemptRecord {
242 pub id: CommandAttemptId,
243 pub at: SimTime,
244 pub revision_before: u64,
246 pub ingress: CommandIngress,
247 pub request_id: Option<CommandRequestId>,
248 pub expected_revision: Option<u64>,
249 pub envelope: CommandEnvelope,
250 pub outcome: CommandAttemptOutcome,
251}
252
253#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
254#[serde(rename_all = "snake_case")]
255pub enum IngressClass {
256 Command,
257 Communication,
258 Acknowledgement,
259 Information,
260 Decision,
261 ScheduledSystem,
262}
263
264#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
265pub struct PluginIngressDescriptor {
266 pub name: String,
267 pub description: String,
268 pub class: IngressClass,
269 pub payload_schema: PayloadSchema,
270}
271
272#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
273pub struct PluginIngressRequest {
274 pub plugin: String,
275 pub packet_type: String,
276 pub due_at: SimTime,
277 pub priority: i32,
278 pub payload: Value,
279 pub affected_entities: Vec<EntityRef>,
280 #[serde(default, skip_serializing_if = "Option::is_none")]
281 pub cause: Option<CauseRef>,
282 #[serde(default, skip_serializing_if = "Vec::is_empty")]
283 pub archive_retention: Vec<PluginArchiveRetention>,
284}
285
286#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
287#[serde(deny_unknown_fields)]
288pub struct PluginArchiveRetention {
289 pub namespace: String,
290 pub object_id: String,
291}
292
293#[derive(Clone, Debug, Eq, PartialEq)]
297pub struct PluginIngressPermit {
298 pub(super) plugin: String,
299 pub(super) packet_type: String,
300 pub(super) semantic_hash: String,
301 pub(super) token: String,
302}
303
304#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
305#[serde(tag = "maintenance", rename_all = "snake_case")]
306pub enum MaintenanceIngressRequest {
307 DecisionArchive {
308 commit: super::VerifiedDecisionArchiveCommit,
309 },
310 OwnerAuthorized {
311 commit: super::VerifiedOwnerAuthorizedMaintenanceCommit,
312 },
313}
314
315#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
316#[serde(rename_all = "snake_case")]
317pub enum MaintenanceDisposition {
318 Applied,
319 RejectedStale,
320}
321
322#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
323#[serde(deny_unknown_fields)]
324pub struct MaintenanceRejectionReceipt {
325 pub token: String,
326 pub expected_source_root: String,
327 pub observed_source_root: String,
328 pub reason: String,
329}
330
331#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
332#[serde(deny_unknown_fields)]
333pub struct MaintenanceChangeRecord {
334 pub kind: String,
335 pub token: String,
336 pub disposition: MaintenanceDisposition,
337 pub source_root: String,
338 pub target_root: String,
339 #[serde(default, skip_serializing_if = "Option::is_none")]
340 pub rejection: Option<MaintenanceRejectionReceipt>,
341}
342
343impl PluginIngressRequest {
344 #[must_use]
345 pub fn new(
346 plugin: impl Into<String>,
347 packet_type: impl Into<String>,
348 due_at: SimTime,
349 payload: Value,
350 ) -> Self {
351 Self {
352 plugin: plugin.into(),
353 packet_type: packet_type.into(),
354 due_at,
355 priority: 0,
356 payload,
357 affected_entities: Vec::new(),
358 cause: None,
359 archive_retention: Vec::new(),
360 }
361 }
362
363 #[must_use]
364 pub const fn with_priority(mut self, priority: i32) -> Self {
365 self.priority = priority;
366 self
367 }
368
369 #[must_use]
370 pub fn with_entity(mut self, entity: EntityRef) -> Self {
371 self.affected_entities.push(entity);
372 self
373 }
374
375 #[must_use]
376 pub fn caused_by(mut self, cause: CauseRef) -> Self {
377 self.cause = Some(cause);
378 self
379 }
380
381 #[must_use]
382 pub fn with_archive_retention(
383 mut self,
384 retention: impl IntoIterator<Item = PluginArchiveRetention>,
385 ) -> Self {
386 self.archive_retention.extend(retention);
387 self
388 }
389}
390
391#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
392#[serde(tag = "type", rename_all = "snake_case")]
393pub enum IngressPayload {
394 Command {
395 request: Box<CommandRequest>,
396 },
397 Plugin {
398 plugin: String,
399 packet_type: String,
400 payload: Value,
401 affected_entities: Vec<EntityRef>,
402 #[serde(default, skip_serializing_if = "Vec::is_empty")]
403 archive_retention: Vec<PluginArchiveRetention>,
404 },
405 Calendar {
406 cadences: Vec<SystemCadence>,
407 },
408 Decision {
409 request: Box<super::DecisionIngressRequest>,
410 },
411 Maintenance {
412 request: Box<MaintenanceIngressRequest>,
413 },
414 PluginCancellation {
422 cancelled: IngressId,
424 authority: IngressCancellationAuthority,
425 reason: String,
426 },
427}
428
429#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
436#[serde(rename_all = "snake_case")]
437pub enum IngressCancellationAuthority {
438 Host,
441 PluginPermit,
445 BoundarySystem,
448}
449
450pub const MAX_INGRESS_CANCELLATION_REASON_BYTES: usize = 1_024;
452
453pub(super) fn valid_ingress_cancellation_reason(reason: &str) -> bool {
454 !reason.is_empty()
455 && reason == reason.trim()
456 && reason.len() <= MAX_INGRESS_CANCELLATION_REASON_BYTES
457}
458
459#[derive(Clone, Copy)]
461pub(super) struct PluginIngressCancellationProof<'a> {
462 pub(super) permit: Option<&'a PluginIngressPermit>,
464 pub(super) replay: bool,
466 pub(super) boundary_plugin: Option<&'a str>,
468 pub(super) current_generations: &'a [BoundaryIngressGeneration],
470}
471
472#[derive(Clone, Copy, Debug, Eq, PartialEq)]
475pub(super) enum PluginIngressIssuer<'a> {
476 Host,
478 Plugin(&'a str),
480}
481
482pub(super) fn plugin_ingress_cancellation_authorized(
485 authority: IngressCancellationAuthority,
486 issuer: PluginIngressIssuer<'_>,
487 target_plugin: &str,
488 internal: bool,
489 boundary_plugin: Option<&str>,
490) -> bool {
491 match authority {
492 IngressCancellationAuthority::Host => issuer == PluginIngressIssuer::Host && !internal,
493 IngressCancellationAuthority::PluginPermit => {
494 internal
495 && match issuer {
496 PluginIngressIssuer::Host => true,
497 PluginIngressIssuer::Plugin(plugin) => plugin == target_plugin,
498 }
499 }
500 IngressCancellationAuthority::BoundarySystem => {
501 boundary_plugin.is_some_and(|plugin| issuer == PluginIngressIssuer::Plugin(plugin))
502 }
503 }
504}
505
506#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
507pub struct IngressRecord {
508 pub id: IngressId,
509 pub issued_at: SimTime,
510 #[serde(default, skip_serializing_if = "is_zero")]
511 pub eligible_boundary_count: u64,
512 pub due_at: SimTime,
513 pub class: IngressClass,
514 pub priority: i32,
515 pub payload: IngressPayload,
516 #[serde(default, skip_serializing_if = "Option::is_none")]
517 pub cause: Option<CauseRef>,
518}
519
520#[allow(clippy::trivially_copy_pass_by_ref)]
521const fn is_zero(value: &u64) -> bool {
522 *value == 0
523}
524
525#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
526pub struct IngressReceipt {
527 pub ingress_id: IngressId,
528 pub issued_at: SimTime,
529 pub due_at: SimTime,
530}
531
532#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
533pub(crate) struct IngressQueueKey {
534 pub due_at: SimTime,
535 pub class: IngressClass,
536 pub priority: Reverse<i32>,
537 pub issued_at: SimTime,
538 pub id: IngressId,
539}
540
541impl IngressQueueKey {
542 #[must_use]
543 pub(crate) const fn from_record(record: &IngressRecord) -> Self {
544 Self {
545 due_at: record.due_at,
546 class: record.class,
547 priority: Reverse(record.priority),
548 issued_at: record.issued_at,
549 id: record.id,
550 }
551 }
552}
553
554impl Simulation {
555 pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
556 match self.admit_command(
557 None,
558 None,
559 envelope,
560 CommandIngress::LegacyDirect,
561 None,
562 false,
563 )? {
564 CommandOutcome::Accepted { receipt } => Ok(receipt),
565 CommandOutcome::Rejected { rejection } => Err(rejection.error),
566 }
567 }
568
569 pub fn enqueue_command(
570 &mut self,
571 due_at: SimTime,
572 priority: i32,
573 request: CommandRequest,
574 ) -> Result<IngressReceipt, CanwuError> {
575 self.ensure_runtime_ready()?;
576 self.ensure_canonical_ingress_can_start()?;
577 self.ensure_command_ingress_family(CommandIngress::LiveRequest)?;
578 if let Some(existing) = self
579 .state
580 .evidence
581 .archived_ingress_requests
582 .get(&request.request_id)
583 {
584 let input_hash = canonical_hash(
585 "canwu.archive.ingress.command.v1",
586 &(due_at, priority, &request),
587 )?;
588 if existing.input_hash == input_hash {
589 return Ok(existing.receipt.clone());
590 }
591 return Err(CanwuError::new(
592 ErrorCode::IdempotencyConflict,
593 format!(
594 "command request {} is already queued with different ingress content",
595 request.request_id
596 ),
597 ));
598 }
599 for record in &self.state.evidence.ingress {
600 let IngressPayload::Command { request: existing } = &record.payload else {
601 continue;
602 };
603 if existing.request_id != request.request_id {
604 continue;
605 }
606 if existing.as_ref() == &request
607 && record.due_at == due_at
608 && record.priority == priority
609 {
610 return Ok(IngressReceipt {
611 ingress_id: record.id,
612 issued_at: record.issued_at,
613 due_at: record.due_at,
614 });
615 }
616 return Err(CanwuError::new(
617 ErrorCode::IdempotencyConflict,
618 format!(
619 "command request {} is already queued with different ingress content",
620 request.request_id
621 ),
622 ));
623 }
624 if self.command_request_id_is_in_use(request.request_id) {
625 return Err(CanwuError::new(
626 ErrorCode::IdempotencyConflict,
627 format!(
628 "command request {} is already reserved or processed",
629 request.request_id
630 ),
631 ));
632 }
633 if request
634 .envelope
635 .expected_time
636 .is_some_and(|expected| expected != due_at)
637 {
638 return Err(CanwuError::new(
639 ErrorCode::SimulationTimeConflict,
640 "queued command expected time must equal its due simulation time",
641 ));
642 }
643 self.append_ingress(
644 due_at,
645 IngressClass::Command,
646 priority,
647 IngressPayload::Command {
648 request: Box::new(request),
649 },
650 None,
651 false,
652 )
653 }
654
655 pub fn enqueue_plugin_ingress(
656 &mut self,
657 request: PluginIngressRequest,
658 ) -> Result<IngressReceipt, CanwuError> {
659 self.enqueue_plugin_ingress_inner(request, None, false)
660 }
661
662 pub fn enqueue_permitted_plugin_ingress(
665 &mut self,
666 request: PluginIngressRequest,
667 permit: &PluginIngressPermit,
668 ) -> Result<IngressReceipt, CanwuError> {
669 self.enqueue_plugin_ingress_inner(request, Some(permit), false)
670 }
671
672 pub(super) fn replay_plugin_ingress(
673 &mut self,
674 request: PluginIngressRequest,
675 ) -> Result<IngressReceipt, CanwuError> {
676 self.enqueue_plugin_ingress_inner(request, None, true)
677 }
678
679 pub fn cancel_plugin_ingress(
699 &mut self,
700 ingress_id: IngressId,
701 reason: impl Into<String>,
702 ) -> Result<IngressReceipt, CanwuError> {
703 self.cancel_plugin_ingress_inner(
704 ingress_id,
705 IngressCancellationAuthority::Host,
706 None,
707 reason.into(),
708 false,
709 )
710 }
711
712 pub fn cancel_permitted_plugin_ingress(
721 &mut self,
722 ingress_id: IngressId,
723 permit: &PluginIngressPermit,
724 reason: impl Into<String>,
725 ) -> Result<IngressReceipt, CanwuError> {
726 self.cancel_plugin_ingress_inner(
727 ingress_id,
728 IngressCancellationAuthority::PluginPermit,
729 Some(permit),
730 reason.into(),
731 false,
732 )
733 }
734
735 pub(super) fn replay_plugin_ingress_cancellation(
736 &mut self,
737 ingress_id: IngressId,
738 authority: IngressCancellationAuthority,
739 reason: String,
740 ) -> Result<IngressReceipt, CanwuError> {
741 if authority == IngressCancellationAuthority::BoundarySystem {
742 return Err(CanwuError::new(
743 ErrorCode::ReplayMismatch,
744 "boundary-system ingress cancellation must be reproduced by its boundary",
745 ));
746 }
747 self.cancel_plugin_ingress_inner(ingress_id, authority, None, reason, true)
748 }
749
750 fn cancel_plugin_ingress_inner(
751 &mut self,
752 ingress_id: IngressId,
753 authority: IngressCancellationAuthority,
754 permit: Option<&PluginIngressPermit>,
755 reason: String,
756 replay: bool,
757 ) -> Result<IngressReceipt, CanwuError> {
758 self.ensure_runtime_ready()?;
759 self.ensure_canonical_ingress_can_start()?;
760 if !replay
761 && self
762 .state
763 .metadata
764 .run_configuration
765 .declared()
766 .is_some_and(|configuration| {
767 configuration.interaction == InteractionPolicy::ReadOnly
768 })
769 {
770 return Err(CanwuError::new(
771 ErrorCode::InteractionReadOnly,
772 "the run interaction policy rejects newly authored plugin ingress cancellation",
773 ));
774 }
775 let target = self.plugin_ingress_cancellation_target(
776 ingress_id,
777 authority,
778 PluginIngressCancellationProof {
779 permit,
780 replay,
781 boundary_plugin: None,
782 current_generations: &[],
783 },
784 &reason,
785 )?;
786 self.append_plugin_ingress_cancellation(target, authority, reason, None, false)
787 }
788
789 pub(super) fn plugin_ingress_cancellation_target(
792 &self,
793 ingress_id: IngressId,
794 authority: IngressCancellationAuthority,
795 proof: PluginIngressCancellationProof<'_>,
796 reason: &str,
797 ) -> Result<IngressQueueKey, CanwuError> {
798 if !valid_ingress_cancellation_reason(reason) {
799 return Err(CanwuError::new(
800 ErrorCode::InvalidPayload,
801 format!(
802 "plugin ingress cancellation reason must be nonempty trimmed text of at most {MAX_INGRESS_CANCELLATION_REASON_BYTES} bytes"
803 ),
804 ));
805 }
806 let Some(record) = self.state.evidence.retained_ingress(ingress_id) else {
807 if ingress_id.get() != 0
808 && ingress_id.get() <= self.state.evidence.archived.ingress_count
809 {
810 return Err(CanwuError::new(
811 ErrorCode::LateIngress,
812 format!("ingress {ingress_id} is archived and no longer pending"),
813 ));
814 }
815 return Err(CanwuError::new(
816 ErrorCode::EvidenceUnavailable,
817 format!("ingress {ingress_id} does not exist"),
818 ));
819 };
820 let IngressPayload::Plugin {
821 plugin,
822 packet_type,
823 ..
824 } = &record.payload
825 else {
826 return Err(CanwuError::new(
827 ErrorCode::InvalidPayload,
828 format!("ingress {ingress_id} is not plugin ingress and cannot be cancelled"),
829 ));
830 };
831 let internal = self
832 .plugins
833 .internal_ingress
834 .contains(&(plugin.clone(), packet_type.clone()));
835 if authority == IngressCancellationAuthority::PluginPermit && !proof.replay {
836 let permitted = match proof.permit {
837 Some(permit) => self.plugin_ingress_permit_matches(plugin, packet_type, permit)?,
838 None => false,
839 };
840 if !permitted {
841 return Err(CanwuError::new(
842 ErrorCode::InvalidAuthority,
843 "plugin ingress cancellation requires the item's exact opaque registration permit",
844 ));
845 }
846 }
847 let issuer = self.plugin_ingress_issuer(record, plugin, proof.current_generations)?;
848 if !plugin_ingress_cancellation_authorized(
849 authority,
850 issuer,
851 plugin,
852 internal,
853 proof.boundary_plugin,
854 ) {
855 return Err(CanwuError::new(
856 ErrorCode::InvalidAuthority,
857 format!("only the issuer of ingress {ingress_id} may cancel it"),
858 ));
859 }
860 let key = IngressQueueKey::from_record(record);
861 if !self.state.scheduler.pending_ingress.contains(&key)
862 || record.due_at <= self.state.scheduler.now
863 {
864 return Err(CanwuError::new(
865 ErrorCode::LateIngress,
866 format!(
867 "ingress {ingress_id} is already due, admitted, or cancelled at {}",
868 self.state.scheduler.now
869 ),
870 ));
871 }
872 Ok(key)
873 }
874
875 pub(super) fn plugin_ingress_issuer<'a>(
877 &'a self,
878 record: &IngressRecord,
879 plugin: &'a str,
880 current_generations: &'a [BoundaryIngressGeneration],
881 ) -> Result<PluginIngressIssuer<'a>, CanwuError> {
882 match &record.cause {
883 None | Some(CauseRef::System(_)) => Ok(PluginIngressIssuer::Host),
884 Some(CauseRef::Command(_)) => Ok(PluginIngressIssuer::Plugin(plugin)),
885 Some(CauseRef::Boundary(boundary)) => self
886 .state
887 .evidence
888 .retained_boundary(*boundary)
889 .map_or(current_generations, |boundary| {
890 boundary.generated_ingress.as_slice()
891 })
892 .iter()
893 .find(|generation| generation.ingress == record.id)
894 .map(|generation| PluginIngressIssuer::Plugin(generation.plugin.as_str()))
895 .ok_or_else(|| {
896 CanwuError::new(
897 ErrorCode::InvalidSnapshot,
898 "boundary-generated ingress lacks its generation evidence",
899 )
900 }),
901 Some(CauseRef::Event(_)) => Err(CanwuError::new(
902 ErrorCode::InvalidAuthority,
903 "event-caused ingress has no cancellable issuer",
904 )),
905 }
906 }
907
908 fn plugin_ingress_permit_matches(
909 &self,
910 plugin: &str,
911 packet_type: &str,
912 permit: &PluginIngressPermit,
913 ) -> Result<bool, CanwuError> {
914 let semantic_hash = self
915 .plugins
916 .descriptors
917 .get(plugin)
918 .map(|descriptor| descriptor.semantic_hash.as_str())
919 .ok_or_else(|| {
920 CanwuError::new(
921 ErrorCode::PluginNotActive,
922 "internal plugin ingress owner is unavailable",
923 )
924 })?;
925 let expected_token = super::canonical_hash(
926 "canwu.plugin.internal-ingress-permit.v1",
927 &(plugin, packet_type, semantic_hash),
928 )?;
929 Ok(permit.plugin == plugin
930 && permit.packet_type == packet_type
931 && permit.semantic_hash == semantic_hash
932 && permit.token == expected_token)
933 }
934
935 pub(super) fn append_plugin_ingress_cancellation(
938 &mut self,
939 target: IngressQueueKey,
940 authority: IngressCancellationAuthority,
941 reason: String,
942 cause: Option<CauseRef>,
943 after_current_boundary: bool,
944 ) -> Result<IngressReceipt, CanwuError> {
945 let transaction = IngressTransactionCheckpoint::capture(&self.state);
946 let (id, next_id, eligible_boundary_count) =
947 self.next_ingress_identity(after_current_boundary)?;
948 let now = self.state.scheduler.now;
949 let record = IngressRecord {
950 id: IngressId::new(id),
951 issued_at: now,
952 eligible_boundary_count,
953 due_at: now,
954 class: target.class,
955 priority: 0,
956 payload: IngressPayload::PluginCancellation {
957 cancelled: target.id,
958 authority,
959 reason,
960 },
961 cause,
962 };
963 self.state.counters.next_ingress_id = next_id;
964 self.state.scheduler.pending_ingress.remove(&target);
965 self.state.scheduler.cancelled_ingress.insert(target.id);
966 self.state.evidence.ingress.push(record);
967 self.state.metadata.plugin_registration_closed = true;
968 if let Err(error) = self.refresh_checkpoint_hash() {
969 transaction.restore_cancellation(&mut self.state, target);
970 return Err(error);
971 }
972 Ok(IngressReceipt {
973 ingress_id: IngressId::new(id),
974 issued_at: now,
975 due_at: now,
976 })
977 }
978
979 fn enqueue_plugin_ingress_inner(
980 &mut self,
981 mut request: PluginIngressRequest,
982 permit: Option<&PluginIngressPermit>,
983 replay: bool,
984 ) -> Result<IngressReceipt, CanwuError> {
985 self.ensure_runtime_ready()?;
986 self.ensure_canonical_ingress_can_start()?;
987 if !replay
988 && self
989 .state
990 .metadata
991 .run_configuration
992 .declared()
993 .is_some_and(|configuration| {
994 configuration.interaction == InteractionPolicy::ReadOnly
995 })
996 {
997 return Err(CanwuError::new(
998 ErrorCode::InteractionReadOnly,
999 "the run interaction policy rejects newly authored plugin ingress",
1000 ));
1001 }
1002 let key = (request.plugin.clone(), request.packet_type.clone());
1003 let descriptor = self.plugins.ingress.get(&key).ok_or_else(|| {
1004 CanwuError::new(
1005 ErrorCode::InvalidPayload,
1006 format!(
1007 "plugin ingress type {}.{} is not registered",
1008 request.plugin, request.packet_type
1009 ),
1010 )
1011 })?;
1012 let internal = self.plugins.internal_ingress.contains(&key);
1013 if internal && !replay {
1014 let permitted = match permit {
1015 Some(permit) => self.plugin_ingress_permit_matches(
1016 &request.plugin,
1017 &request.packet_type,
1018 permit,
1019 )?,
1020 None => false,
1021 };
1022 if !permitted {
1023 return Err(CanwuError::new(
1024 ErrorCode::InvalidAuthority,
1025 "plugin-owned internal ingress requires its opaque registration permit",
1026 ));
1027 }
1028 }
1029 if !request.archive_retention.is_empty() && !internal && !replay {
1030 return Err(CanwuError::new(
1031 ErrorCode::InvalidAuthority,
1032 "archive retention may be attached only to plugin-owned internal ingress",
1033 ));
1034 }
1035 request.archive_retention.sort();
1036 request.archive_retention.dedup();
1037 if request.archive_retention.len() > 32_768
1038 || request.archive_retention.iter().any(|retention| {
1039 retention.namespace.is_empty()
1040 || retention.namespace.len() > 128
1041 || retention.object_id.is_empty()
1042 || retention.object_id.len() > 256
1043 || !retention.namespace.bytes().all(|byte| {
1044 byte.is_ascii_lowercase()
1045 || byte.is_ascii_digit()
1046 || matches!(byte, b'.' | b'-' | b'_')
1047 })
1048 || !retention.object_id.is_ascii()
1049 })
1050 {
1051 return Err(CanwuError::new(
1052 ErrorCode::InvalidPayload,
1053 "plugin ingress archive retention is malformed or exceeds its hard limit",
1054 ));
1055 }
1056 self.plugins.validate_archive_retention(
1057 &request.plugin,
1058 &request.packet_type,
1059 &request.payload,
1060 &request.archive_retention,
1061 )?;
1062 descriptor.payload_schema.validate(&request.payload)?;
1063 request.affected_entities.sort();
1064 request.affected_entities.dedup();
1065 if request
1066 .affected_entities
1067 .iter()
1068 .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
1069 {
1070 return Err(CanwuError::new(
1071 ErrorCode::EntityNotFound,
1072 "plugin ingress references an unknown entity identity",
1073 ));
1074 }
1075 if let Some(cause) = &request.cause {
1076 if matches!(
1077 cause,
1078 CauseRef::Boundary(_) | CauseRef::Command(_) | CauseRef::Event(_)
1079 ) {
1080 return Err(CanwuError::new(
1081 ErrorCode::InvalidPayload,
1082 "boundary, command, and event causes are reserved for plugin-generated ingress",
1083 ));
1084 }
1085 validate_runtime_cause(&self.state, cause)?;
1086 }
1087 self.append_ingress(
1088 request.due_at,
1089 descriptor.class,
1090 request.priority,
1091 IngressPayload::Plugin {
1092 plugin: request.plugin,
1093 packet_type: request.packet_type,
1094 payload: request.payload,
1095 affected_entities: request.affected_entities,
1096 archive_retention: request.archive_retention,
1097 },
1098 request.cause,
1099 false,
1100 )
1101 }
1102
1103 pub fn schedule_calendar_boundary(
1104 &mut self,
1105 due_at: SimTime,
1106 mut cadences: Vec<SystemCadence>,
1107 ) -> Result<IngressReceipt, CanwuError> {
1108 self.ensure_runtime_ready()?;
1109 self.ensure_canonical_ingress_can_start()?;
1110 if cadences.contains(&SystemCadence::EventDriven) {
1111 return Err(CanwuError::new(
1112 ErrorCode::InvalidBoundary,
1113 "calendar ingress cannot declare event-driven cadence",
1114 ));
1115 }
1116 cadences.sort();
1117 cadences.dedup();
1118 if cadences.is_empty() {
1119 return Err(CanwuError::new(
1120 ErrorCode::InvalidBoundary,
1121 "calendar ingress requires at least one scheduled cadence",
1122 ));
1123 }
1124 self.append_ingress(
1125 due_at,
1126 IngressClass::ScheduledSystem,
1127 0,
1128 IngressPayload::Calendar { cadences },
1129 Some(CauseRef::System("canwu.core.calendar".to_owned())),
1130 false,
1131 )
1132 }
1133
1134 pub(super) fn append_ingress(
1135 &mut self,
1136 due_at: SimTime,
1137 class: IngressClass,
1138 priority: i32,
1139 payload: IngressPayload,
1140 cause: Option<CauseRef>,
1141 after_current_boundary: bool,
1142 ) -> Result<IngressReceipt, CanwuError> {
1143 if due_at < self.state.scheduler.now {
1144 return Err(CanwuError::new(
1145 ErrorCode::LateIngress,
1146 format!(
1147 "ingress due at {due_at} cannot be queued after committed time {}",
1148 self.state.scheduler.now
1149 ),
1150 ));
1151 }
1152 let transaction = IngressTransactionCheckpoint::capture(&self.state);
1153 let (id, next_id, eligible_boundary_count) =
1154 self.next_ingress_identity(after_current_boundary)?;
1155 let record = IngressRecord {
1156 id: IngressId::new(id),
1157 issued_at: self.state.scheduler.now,
1158 eligible_boundary_count,
1159 due_at,
1160 class,
1161 priority,
1162 payload,
1163 cause,
1164 };
1165 let queue_key = IngressQueueKey::from_record(&record);
1166 self.state.counters.next_ingress_id = next_id;
1167 self.state.scheduler.pending_ingress.insert(queue_key);
1168 self.state.evidence.ingress.push(record.clone());
1169 self.state.metadata.plugin_registration_closed = true;
1170 if let Err(error) = self.refresh_checkpoint_hash() {
1171 transaction.restore(&mut self.state, &queue_key);
1172 return Err(error);
1173 }
1174 Ok(IngressReceipt {
1175 ingress_id: record.id,
1176 issued_at: record.issued_at,
1177 due_at: record.due_at,
1178 })
1179 }
1180
1181 fn next_ingress_identity(
1183 &self,
1184 after_current_boundary: bool,
1185 ) -> Result<(u64, u64, u64), CanwuError> {
1186 let (id, next_id) = claim_counter(self.state.counters.next_ingress_id, "ingress ID")?;
1187 let boundary_count = self
1188 .state
1189 .evidence
1190 .archived
1191 .boundary_count
1192 .checked_add(
1193 u64::try_from(self.state.evidence.boundaries.len()).map_err(|_| {
1194 CanwuError::new(
1195 ErrorCode::IdentifierExhausted,
1196 "boundary count exceeds the ingress journal range",
1197 )
1198 })?,
1199 )
1200 .ok_or_else(|| {
1201 CanwuError::new(
1202 ErrorCode::IdentifierExhausted,
1203 "boundary count exceeds the ingress journal range",
1204 )
1205 })?;
1206 let eligible_boundary_count = if after_current_boundary {
1207 boundary_count.checked_add(1).ok_or_else(|| {
1208 CanwuError::new(
1209 ErrorCode::IdentifierExhausted,
1210 "ingress boundary eligibility exceeds the journal range",
1211 )
1212 })?
1213 } else {
1214 boundary_count
1215 };
1216 Ok((id, next_id, eligible_boundary_count))
1217 }
1218
1219 pub fn process_command(
1220 &mut self,
1221 request: CommandRequest,
1222 ) -> Result<CommandOutcome, CanwuError> {
1223 self.ensure_runtime_ready()?;
1224 if self.state.evidence.archived.ingress_count != 0
1225 || !self.state.evidence.ingress.is_empty()
1226 {
1227 return Err(CanwuError::new(
1228 ErrorCode::MixedCommandIngress,
1229 "direct command requests cannot bypass an active canonical ingress journal",
1230 ));
1231 }
1232 self.admit_command(
1233 Some(request.request_id),
1234 Some(request.expected_revision),
1235 request.envelope,
1236 CommandIngress::LiveRequest,
1237 None,
1238 true,
1239 )
1240 }
1241
1242 pub(super) fn admit_command(
1243 &mut self,
1244 request_id: Option<CommandRequestId>,
1245 expected_revision: Option<u64>,
1246 envelope: CommandEnvelope,
1247 ingress: CommandIngress,
1248 decision_controller_id: Option<String>,
1249 record_attempt: bool,
1250 ) -> Result<CommandOutcome, CanwuError> {
1251 self.ensure_runtime_ready()?;
1252 self.ensure_command_ingress_family(ingress)?;
1253 if let Some(cached) =
1254 self.cached_command_outcome(request_id, expected_revision, &envelope)?
1255 {
1256 return Ok(cached);
1257 }
1258
1259 let revision_before = self.revision();
1260 let admission = CommandAdmission {
1261 request_id,
1262 expected_revision,
1263 expected_time: envelope.expected_time,
1264 revision_before,
1265 ingress,
1266 };
1267 let attempt_id = if record_attempt {
1268 let (value, _) = claim_counter(
1269 self.state.counters.next_command_attempt_id,
1270 "command attempt ID",
1271 )?;
1272 CommandAttemptId::new(value)
1273 } else {
1274 CommandAttemptId::default()
1275 };
1276 let authority = match resolve_command_authority(&envelope) {
1277 Ok(authority) => authority,
1278 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
1279 return self.record_command_rejection(attempt_id, admission, envelope, error);
1280 }
1281 Err(error) => return Err(error),
1282 };
1283 if let Err(error) = self.validate_command_ingress(&envelope.issuer, &authority, admission) {
1284 if is_expected_command_rejection(&error.code) && record_attempt {
1285 return self.record_command_rejection(attempt_id, admission, envelope, error);
1286 }
1287 return Err(error);
1288 }
1289 if let Some(expected_time) = envelope.expected_time
1290 && expected_time != self.state.scheduler.now
1291 {
1292 let error = CanwuError::new(
1293 ErrorCode::SimulationTimeConflict,
1294 format!(
1295 "command expected time {expected_time}, but simulation is at {}",
1296 self.state.scheduler.now
1297 ),
1298 );
1299 if record_attempt {
1300 return self.record_command_rejection(attempt_id, admission, envelope, error);
1301 }
1302 return Err(error);
1303 }
1304 if let Err(error) = self.validate_command_issuer(&envelope.issuer, &authority) {
1305 if record_attempt {
1306 return self.record_command_rejection(attempt_id, admission, envelope, error);
1307 }
1308 return Err(error);
1309 }
1310
1311 let (command_id_value, next_command_id) =
1312 claim_counter(self.state.counters.next_command_id, "command ID")?;
1313 let (correlation_id, next_correlation_id) =
1314 claim_counter(self.state.counters.next_correlation_id, "correlation ID")?;
1315 let command_id = CommandId::new(command_id_value);
1316 let context = CommandContext {
1317 issuer: envelope.issuer.clone(),
1318 authority,
1319 decision_controller_id,
1320 run_policy: self.state.metadata.run_configuration.command_policy(),
1321 ingress: admission.ingress,
1322 attempt_id: record_attempt.then_some(attempt_id),
1323 command_id,
1324 request_id: admission.request_id,
1325 revision: admission.revision_before,
1326 simulation_time: self.state.scheduler.now,
1327 expected_revision: admission.expected_revision,
1328 expected_time: envelope.expected_time,
1329 };
1330 let prepared = match self.prepare_command(&envelope, &context) {
1331 Ok(prepared) => prepared,
1332 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
1333 return self.record_command_rejection(attempt_id, admission, envelope, error);
1334 }
1335 Err(error) => return Err(error),
1336 };
1337 let next_attempt_id = if record_attempt {
1338 let (claimed_id, next_attempt_id) = claim_counter(
1339 self.state.counters.next_command_attempt_id,
1340 "command attempt ID",
1341 )?;
1342 if claimed_id != attempt_id.get() {
1343 return Err(CanwuError::new(
1344 ErrorCode::InvalidSnapshot,
1345 "command attempt allocation changed during application",
1346 ));
1347 }
1348 Some(next_attempt_id)
1349 } else {
1350 None
1351 };
1352 let revision = self.next_state_revision()?;
1353 let transaction = CommandTransactionCheckpoint::capture(&self.state);
1354 let event_start = self.state.evidence.events.len();
1355 self.state.counters.next_command_id = next_command_id;
1356 self.state.counters.next_correlation_id = next_correlation_id;
1357 self.invalidate_commitments(prepared.commitment_invalidation());
1358
1359 if let Err(error) = self.apply_prepared(prepared, command_id, correlation_id) {
1360 transaction.restore(&mut self.state);
1361 if is_expected_command_rejection(&error.code) && record_attempt {
1362 return self.record_command_rejection(attempt_id, admission, envelope, error);
1363 }
1364 return Err(error);
1365 }
1366 let emitted_events: Vec<_> = self.state.evidence.events[event_start..]
1367 .iter()
1368 .map(|event| event.id)
1369 .collect();
1370 self.state.metadata.plugin_registration_closed = true;
1371 self.state.evidence.commands.push(CommandRecord {
1372 id: command_id,
1373 attempt_id: record_attempt.then_some(attempt_id),
1374 accepted_at: self.state.scheduler.now,
1375 envelope: envelope.clone(),
1376 emitted_events: if record_attempt {
1377 emitted_events.clone()
1378 } else {
1379 Vec::new()
1380 },
1381 });
1382 if let Some(next_attempt_id) = next_attempt_id {
1383 self.state.counters.next_command_attempt_id = next_attempt_id;
1384 self.state
1385 .evidence
1386 .command_attempts
1387 .push(CommandAttemptRecord {
1388 id: attempt_id,
1389 at: self.state.scheduler.now,
1390 revision_before: admission.revision_before,
1391 ingress: admission.ingress,
1392 request_id: admission.request_id,
1393 expected_revision: admission.expected_revision,
1394 envelope,
1395 outcome: CommandAttemptOutcome::Accepted { command_id },
1396 });
1397 }
1398 self.state.counters.state_revision = revision;
1399 if let Err(error) = self.refresh_checkpoint_hash() {
1400 transaction.restore(&mut self.state);
1401 return Err(error);
1402 }
1403
1404 Ok(CommandOutcome::Accepted {
1405 receipt: CommandReceipt {
1406 attempt_id: record_attempt.then_some(attempt_id),
1407 command_id,
1408 request_id: admission.request_id,
1409 revision,
1410 accepted_at: self.state.scheduler.now,
1411 emitted_events,
1412 },
1413 })
1414 }
1415
1416 fn ensure_command_ingress_family(&self, ingress: CommandIngress) -> Result<(), CanwuError> {
1417 let has_legacy_commands = self.state.evidence.archived_legacy_commands
1418 || self
1419 .state
1420 .evidence
1421 .commands
1422 .iter()
1423 .any(|record| record.attempt_id.is_none());
1424 let has_tracked_attempts = self.state.evidence.archived_tracked_attempts
1425 || !self.state.evidence.command_attempts.is_empty()
1426 || !self.state.evidence.ingress.is_empty();
1427 if (ingress == CommandIngress::LegacyDirect && has_tracked_attempts)
1428 || (ingress != CommandIngress::LegacyDirect && has_legacy_commands)
1429 {
1430 return Err(CanwuError::new(
1431 ErrorCode::MixedCommandIngress,
1432 "legacy-direct commands and tracked request/replay attempts cannot coexist in one run",
1433 ));
1434 }
1435 Ok(())
1436 }
1437
1438 pub(super) fn ensure_canonical_ingress_can_start(&self) -> Result<(), CanwuError> {
1439 if runtime_has_unqueued_command_history(&self.state) {
1440 return Err(CanwuError::new(
1441 ErrorCode::MixedCommandIngress,
1442 "canonical ingress cannot be added after direct command history",
1443 ));
1444 }
1445 Ok(())
1446 }
1447
1448 fn cached_command_outcome(
1449 &self,
1450 request_id: Option<CommandRequestId>,
1451 expected_revision: Option<u64>,
1452 envelope: &CommandEnvelope,
1453 ) -> Result<Option<CommandOutcome>, CanwuError> {
1454 let Some(request_id) = request_id else {
1455 return Ok(None);
1456 };
1457 if let Some(cached) = self
1458 .state
1459 .evidence
1460 .archived_command_requests
1461 .get(&request_id)
1462 {
1463 let input_hash = canonical_hash(
1464 "canwu.archive.command.request.v1",
1465 &(expected_revision, envelope),
1466 )?;
1467 if cached.input_hash != input_hash {
1468 return Ok(Some(CommandOutcome::Rejected {
1469 rejection: CommandRejection {
1470 attempt_id: None,
1471 request_id: Some(request_id),
1472 retained_revision: self.revision(),
1473 rejected_at: self.state.scheduler.now,
1474 error: CanwuError::new(
1475 ErrorCode::IdempotencyConflict,
1476 "this command request ID was already used for different input",
1477 ),
1478 },
1479 }));
1480 }
1481 return Ok(Some(cached.outcome.clone()));
1482 }
1483 let Some(attempt) = self
1484 .state
1485 .evidence
1486 .command_attempts
1487 .iter()
1488 .find(|attempt| attempt.request_id == Some(request_id))
1489 else {
1490 return Ok(None);
1491 };
1492 if attempt.expected_revision != expected_revision || &attempt.envelope != envelope {
1493 return Ok(Some(CommandOutcome::Rejected {
1494 rejection: CommandRejection {
1495 attempt_id: None,
1496 request_id: Some(request_id),
1497 retained_revision: self.revision(),
1498 rejected_at: self.state.scheduler.now,
1499 error: CanwuError::new(
1500 ErrorCode::IdempotencyConflict,
1501 "this command request ID was already used for different input",
1502 ),
1503 },
1504 }));
1505 }
1506 Ok(Some(self.command_outcome_from_attempt(attempt)?))
1507 }
1508
1509 pub(super) fn command_outcome_from_attempt(
1510 &self,
1511 attempt: &CommandAttemptRecord,
1512 ) -> Result<CommandOutcome, CanwuError> {
1513 let request_id = attempt.request_id.ok_or_else(|| {
1514 invalid_snapshot_error("tracked command attempt is missing its request ID")
1515 })?;
1516 let committed_revision = attempt.revision_before.checked_add(1).ok_or_else(|| {
1517 invalid_snapshot_error("cached command attempt revision is exhausted")
1518 })?;
1519 match &attempt.outcome {
1520 CommandAttemptOutcome::Accepted { command_id } => {
1521 let retained_number = command_id
1522 .get()
1523 .checked_sub(self.state.evidence.archived.command_count)
1524 .and_then(|value| value.checked_sub(1))
1525 .ok_or_else(|| {
1526 invalid_snapshot_error(
1527 "accepted command attempt references archived command evidence",
1528 )
1529 })?;
1530 let index = usize::try_from(retained_number).map_err(|_| {
1531 invalid_snapshot_error(
1532 "accepted command attempt exceeds the retained command index space",
1533 )
1534 })?;
1535 let record = self
1536 .state
1537 .evidence
1538 .commands
1539 .get(index)
1540 .filter(|record| record.id == *command_id)
1541 .ok_or_else(|| {
1542 invalid_snapshot_error(
1543 "accepted command attempt references a missing command",
1544 )
1545 })?;
1546 Ok(CommandOutcome::Accepted {
1547 receipt: CommandReceipt {
1548 attempt_id: Some(attempt.id),
1549 command_id: *command_id,
1550 request_id: Some(request_id),
1551 revision: committed_revision,
1552 accepted_at: record.accepted_at,
1553 emitted_events: record.emitted_events.clone(),
1554 },
1555 })
1556 }
1557 CommandAttemptOutcome::Rejected { error } => Ok(CommandOutcome::Rejected {
1558 rejection: CommandRejection {
1559 attempt_id: Some(attempt.id),
1560 request_id: Some(request_id),
1561 retained_revision: committed_revision,
1562 rejected_at: attempt.at,
1563 error: error.clone(),
1564 },
1565 }),
1566 }
1567 }
1568
1569 fn record_command_rejection(
1570 &mut self,
1571 attempt_id: CommandAttemptId,
1572 admission: CommandAdmission,
1573 envelope: CommandEnvelope,
1574 error: CanwuError,
1575 ) -> Result<CommandOutcome, CanwuError> {
1576 let (claimed_id, next_attempt_id) = claim_counter(
1577 self.state.counters.next_command_attempt_id,
1578 "command attempt ID",
1579 )?;
1580 if claimed_id != attempt_id.get() {
1581 return Err(CanwuError::new(
1582 ErrorCode::InvalidSnapshot,
1583 "command attempt allocation changed during rejection",
1584 ));
1585 }
1586 let revision = self.next_state_revision()?;
1587 let attempt = CommandAttemptRecord {
1588 id: attempt_id,
1589 at: self.state.scheduler.now,
1590 revision_before: admission.revision_before,
1591 ingress: admission.ingress,
1592 request_id: admission.request_id,
1593 expected_revision: admission.expected_revision,
1594 envelope,
1595 outcome: CommandAttemptOutcome::Rejected {
1596 error: error.clone(),
1597 },
1598 };
1599 let transaction = RejectionTransactionCheckpoint::capture(&self.state);
1600 self.state.counters.next_command_attempt_id = next_attempt_id;
1601 self.state.counters.state_revision = revision;
1602 self.state.metadata.plugin_registration_closed = true;
1603 self.state.evidence.command_attempts.push(attempt);
1604 if let Err(hash_error) = self.refresh_checkpoint_hash() {
1605 transaction.restore(&mut self.state);
1606 return Err(hash_error);
1607 }
1608 Ok(CommandOutcome::Rejected {
1609 rejection: CommandRejection {
1610 attempt_id: Some(attempt_id),
1611 request_id: admission.request_id,
1612 retained_revision: revision,
1613 rejected_at: self.state.scheduler.now,
1614 error,
1615 },
1616 })
1617 }
1618
1619 fn validate_command_ingress(
1620 &self,
1621 issuer: &Issuer,
1622 authority: &CommandAuthority,
1623 admission: CommandAdmission,
1624 ) -> Result<(), CanwuError> {
1625 validate_command_ingress_policy(
1626 &self.state.metadata.run_configuration,
1627 issuer,
1628 authority,
1629 admission,
1630 &|entity| runtime_entity_exists(&self.state, entity),
1631 )
1632 }
1633
1634 pub fn advance_canonical(
1635 &mut self,
1636 duration: SimDuration,
1637 ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
1638 self.ensure_runtime_ready()?;
1639 if duration.is_negative() {
1640 return Err(CanwuError::new(
1641 ErrorCode::InvalidDuration,
1642 "canonical simulation time cannot advance by a negative duration",
1643 ));
1644 }
1645 let target = self
1646 .state
1647 .scheduler
1648 .now
1649 .checked_add(duration)
1650 .ok_or_else(|| {
1651 CanwuError::new(
1652 ErrorCode::InvalidDuration,
1653 "canonical simulation target time exceeds the supported range",
1654 )
1655 })?;
1656 let mut receipts = Vec::new();
1657 while let Some(next_due) = self.next_canonical_due_time()
1658 && next_due <= target
1659 {
1660 let at = next_due.max(self.state.scheduler.now);
1661 receipts.push(self.settle_boundary(BoundaryRequest::at(at))?);
1662 }
1663 if self.state.scheduler.now < target {
1664 self.advance_to(target)?;
1665 }
1666 Ok(receipts)
1667 }
1668
1669 pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
1670 self.ensure_runtime_ready()?;
1671 let Some(next_due) = self.next_canonical_due_time() else {
1672 return Ok(None);
1673 };
1674 self.settle_boundary(BoundaryRequest::at(next_due.max(self.state.scheduler.now)))
1675 .map(Some)
1676 }
1677
1678 fn next_canonical_due_time(&self) -> Option<SimTime> {
1679 let scheduled = self.state.scheduler.actions.keys().next().map(|key| key.at);
1680 let ingress = self
1681 .state
1682 .scheduler
1683 .pending_ingress
1684 .first()
1685 .map(|key| key.due_at);
1686 match (scheduled, ingress) {
1687 (Some(left), Some(right)) => Some(left.min(right)),
1688 (Some(value), None) | (None, Some(value)) => Some(value),
1689 (None, None) => None,
1690 }
1691 }
1692
1693 pub(super) fn take_due_ingress(&mut self, at: SimTime) -> Vec<IngressId> {
1694 let mut admitted = Vec::new();
1695 while self
1696 .state
1697 .scheduler
1698 .pending_ingress
1699 .first()
1700 .is_some_and(|key| key.due_at <= at)
1701 {
1702 let key = self
1703 .state
1704 .scheduler
1705 .pending_ingress
1706 .pop_first()
1707 .expect("pending ingress was checked as non-empty");
1708 admitted.push(key.id);
1709 }
1710 admitted
1711 }
1712}