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}
282
283impl PluginIngressRequest {
284 #[must_use]
285 pub fn new(
286 plugin: impl Into<String>,
287 packet_type: impl Into<String>,
288 due_at: SimTime,
289 payload: Value,
290 ) -> Self {
291 Self {
292 plugin: plugin.into(),
293 packet_type: packet_type.into(),
294 due_at,
295 priority: 0,
296 payload,
297 affected_entities: Vec::new(),
298 cause: None,
299 }
300 }
301
302 #[must_use]
303 pub const fn with_priority(mut self, priority: i32) -> Self {
304 self.priority = priority;
305 self
306 }
307
308 #[must_use]
309 pub fn with_entity(mut self, entity: EntityRef) -> Self {
310 self.affected_entities.push(entity);
311 self
312 }
313
314 #[must_use]
315 pub fn caused_by(mut self, cause: CauseRef) -> Self {
316 self.cause = Some(cause);
317 self
318 }
319}
320
321#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
322#[serde(tag = "type", rename_all = "snake_case")]
323pub enum IngressPayload {
324 Command {
325 request: Box<CommandRequest>,
326 },
327 Plugin {
328 plugin: String,
329 packet_type: String,
330 payload: Value,
331 affected_entities: Vec<EntityRef>,
332 },
333 Calendar {
334 cadences: Vec<SystemCadence>,
335 },
336 Decision {
337 request: Box<super::DecisionIngressRequest>,
338 },
339}
340
341#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
342pub struct IngressRecord {
343 pub id: IngressId,
344 pub issued_at: SimTime,
345 #[serde(default, skip_serializing_if = "is_zero")]
346 pub eligible_boundary_count: u64,
347 pub due_at: SimTime,
348 pub class: IngressClass,
349 pub priority: i32,
350 pub payload: IngressPayload,
351 #[serde(default, skip_serializing_if = "Option::is_none")]
352 pub cause: Option<CauseRef>,
353}
354
355#[allow(clippy::trivially_copy_pass_by_ref)]
356const fn is_zero(value: &u64) -> bool {
357 *value == 0
358}
359
360#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
361pub struct IngressReceipt {
362 pub ingress_id: IngressId,
363 pub issued_at: SimTime,
364 pub due_at: SimTime,
365}
366
367#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
368pub(crate) struct IngressQueueKey {
369 pub due_at: SimTime,
370 pub class: IngressClass,
371 pub priority: Reverse<i32>,
372 pub issued_at: SimTime,
373 pub id: IngressId,
374}
375
376impl IngressQueueKey {
377 #[must_use]
378 pub(crate) const fn from_record(record: &IngressRecord) -> Self {
379 Self {
380 due_at: record.due_at,
381 class: record.class,
382 priority: Reverse(record.priority),
383 issued_at: record.issued_at,
384 id: record.id,
385 }
386 }
387}
388
389impl Simulation {
390 pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
391 match self.admit_command(
392 None,
393 None,
394 envelope,
395 CommandIngress::LegacyDirect,
396 None,
397 false,
398 )? {
399 CommandOutcome::Accepted { receipt } => Ok(receipt),
400 CommandOutcome::Rejected { rejection } => Err(rejection.error),
401 }
402 }
403
404 pub fn enqueue_command(
405 &mut self,
406 due_at: SimTime,
407 priority: i32,
408 request: CommandRequest,
409 ) -> Result<IngressReceipt, CanwuError> {
410 self.ensure_runtime_ready()?;
411 self.ensure_canonical_ingress_can_start()?;
412 self.ensure_command_ingress_family(CommandIngress::LiveRequest)?;
413 if let Some(existing) = self
414 .state
415 .evidence
416 .archived_ingress_requests
417 .get(&request.request_id)
418 {
419 let input_hash = canonical_hash(
420 "canwu.archive.ingress.command.v1",
421 &(due_at, priority, &request),
422 )?;
423 if existing.input_hash == input_hash {
424 return Ok(existing.receipt.clone());
425 }
426 return Err(CanwuError::new(
427 ErrorCode::IdempotencyConflict,
428 format!(
429 "command request {} is already queued with different ingress content",
430 request.request_id
431 ),
432 ));
433 }
434 for record in &self.state.evidence.ingress {
435 let IngressPayload::Command { request: existing } = &record.payload else {
436 continue;
437 };
438 if existing.request_id != request.request_id {
439 continue;
440 }
441 if existing.as_ref() == &request
442 && record.due_at == due_at
443 && record.priority == priority
444 {
445 return Ok(IngressReceipt {
446 ingress_id: record.id,
447 issued_at: record.issued_at,
448 due_at: record.due_at,
449 });
450 }
451 return Err(CanwuError::new(
452 ErrorCode::IdempotencyConflict,
453 format!(
454 "command request {} is already queued with different ingress content",
455 request.request_id
456 ),
457 ));
458 }
459 if self.command_request_id_is_in_use(request.request_id) {
460 return Err(CanwuError::new(
461 ErrorCode::IdempotencyConflict,
462 format!(
463 "command request {} is already reserved or processed",
464 request.request_id
465 ),
466 ));
467 }
468 if request
469 .envelope
470 .expected_time
471 .is_some_and(|expected| expected != due_at)
472 {
473 return Err(CanwuError::new(
474 ErrorCode::SimulationTimeConflict,
475 "queued command expected time must equal its due simulation time",
476 ));
477 }
478 self.append_ingress(
479 due_at,
480 IngressClass::Command,
481 priority,
482 IngressPayload::Command {
483 request: Box::new(request),
484 },
485 None,
486 false,
487 )
488 }
489
490 pub fn enqueue_plugin_ingress(
491 &mut self,
492 mut request: PluginIngressRequest,
493 ) -> Result<IngressReceipt, CanwuError> {
494 self.ensure_runtime_ready()?;
495 self.ensure_canonical_ingress_can_start()?;
496 if self
497 .state
498 .metadata
499 .run_configuration
500 .declared()
501 .is_some_and(|configuration| configuration.interaction == InteractionPolicy::ReadOnly)
502 {
503 return Err(CanwuError::new(
504 ErrorCode::InteractionReadOnly,
505 "the run interaction policy rejects newly authored plugin ingress",
506 ));
507 }
508 let key = (request.plugin.clone(), request.packet_type.clone());
509 let descriptor = self.plugins.ingress.get(&key).ok_or_else(|| {
510 CanwuError::new(
511 ErrorCode::InvalidPayload,
512 format!(
513 "plugin ingress type {}.{} is not registered",
514 request.plugin, request.packet_type
515 ),
516 )
517 })?;
518 descriptor.payload_schema.validate(&request.payload)?;
519 request.affected_entities.sort();
520 request.affected_entities.dedup();
521 if request
522 .affected_entities
523 .iter()
524 .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
525 {
526 return Err(CanwuError::new(
527 ErrorCode::EntityNotFound,
528 "plugin ingress references an unknown entity identity",
529 ));
530 }
531 if let Some(cause) = &request.cause {
532 if matches!(
533 cause,
534 CauseRef::Boundary(_) | CauseRef::Command(_) | CauseRef::Event(_)
535 ) {
536 return Err(CanwuError::new(
537 ErrorCode::InvalidPayload,
538 "boundary, command, and event causes are reserved for plugin-generated ingress",
539 ));
540 }
541 validate_runtime_cause(&self.state, cause)?;
542 }
543 self.append_ingress(
544 request.due_at,
545 descriptor.class,
546 request.priority,
547 IngressPayload::Plugin {
548 plugin: request.plugin,
549 packet_type: request.packet_type,
550 payload: request.payload,
551 affected_entities: request.affected_entities,
552 },
553 request.cause,
554 false,
555 )
556 }
557
558 pub fn schedule_calendar_boundary(
559 &mut self,
560 due_at: SimTime,
561 mut cadences: Vec<SystemCadence>,
562 ) -> Result<IngressReceipt, CanwuError> {
563 self.ensure_runtime_ready()?;
564 self.ensure_canonical_ingress_can_start()?;
565 if cadences.contains(&SystemCadence::EventDriven) {
566 return Err(CanwuError::new(
567 ErrorCode::InvalidBoundary,
568 "calendar ingress cannot declare event-driven cadence",
569 ));
570 }
571 cadences.sort();
572 cadences.dedup();
573 if cadences.is_empty() {
574 return Err(CanwuError::new(
575 ErrorCode::InvalidBoundary,
576 "calendar ingress requires at least one scheduled cadence",
577 ));
578 }
579 self.append_ingress(
580 due_at,
581 IngressClass::ScheduledSystem,
582 0,
583 IngressPayload::Calendar { cadences },
584 Some(CauseRef::System("canwu.core.calendar".to_owned())),
585 false,
586 )
587 }
588
589 pub(super) fn append_ingress(
590 &mut self,
591 due_at: SimTime,
592 class: IngressClass,
593 priority: i32,
594 payload: IngressPayload,
595 cause: Option<CauseRef>,
596 after_current_boundary: bool,
597 ) -> Result<IngressReceipt, CanwuError> {
598 if due_at < self.state.scheduler.now {
599 return Err(CanwuError::new(
600 ErrorCode::LateIngress,
601 format!(
602 "ingress due at {due_at} cannot be queued after committed time {}",
603 self.state.scheduler.now
604 ),
605 ));
606 }
607 let transaction = IngressTransactionCheckpoint::capture(&self.state);
608 let (id, next_id) = claim_counter(self.state.counters.next_ingress_id, "ingress ID")?;
609 let boundary_count = self
610 .state
611 .evidence
612 .archived
613 .boundary_count
614 .checked_add(
615 u64::try_from(self.state.evidence.boundaries.len()).map_err(|_| {
616 CanwuError::new(
617 ErrorCode::IdentifierExhausted,
618 "boundary count exceeds the ingress journal range",
619 )
620 })?,
621 )
622 .ok_or_else(|| {
623 CanwuError::new(
624 ErrorCode::IdentifierExhausted,
625 "boundary count exceeds the ingress journal range",
626 )
627 })?;
628 let eligible_boundary_count = if after_current_boundary {
629 boundary_count.checked_add(1).ok_or_else(|| {
630 CanwuError::new(
631 ErrorCode::IdentifierExhausted,
632 "ingress boundary eligibility exceeds the journal range",
633 )
634 })?
635 } else {
636 boundary_count
637 };
638 let record = IngressRecord {
639 id: IngressId::new(id),
640 issued_at: self.state.scheduler.now,
641 eligible_boundary_count,
642 due_at,
643 class,
644 priority,
645 payload,
646 cause,
647 };
648 let queue_key = IngressQueueKey::from_record(&record);
649 self.state.counters.next_ingress_id = next_id;
650 self.state.scheduler.pending_ingress.insert(queue_key);
651 self.state.evidence.ingress.push(record.clone());
652 self.state.metadata.plugin_registration_closed = true;
653 if let Err(error) = self.refresh_checkpoint_hash() {
654 transaction.restore(&mut self.state, &queue_key);
655 return Err(error);
656 }
657 Ok(IngressReceipt {
658 ingress_id: record.id,
659 issued_at: record.issued_at,
660 due_at: record.due_at,
661 })
662 }
663
664 pub fn process_command(
665 &mut self,
666 request: CommandRequest,
667 ) -> Result<CommandOutcome, CanwuError> {
668 self.ensure_runtime_ready()?;
669 if self.state.evidence.archived.ingress_count != 0
670 || !self.state.evidence.ingress.is_empty()
671 {
672 return Err(CanwuError::new(
673 ErrorCode::MixedCommandIngress,
674 "direct command requests cannot bypass an active canonical ingress journal",
675 ));
676 }
677 self.admit_command(
678 Some(request.request_id),
679 Some(request.expected_revision),
680 request.envelope,
681 CommandIngress::LiveRequest,
682 None,
683 true,
684 )
685 }
686
687 pub(super) fn admit_command(
688 &mut self,
689 request_id: Option<CommandRequestId>,
690 expected_revision: Option<u64>,
691 envelope: CommandEnvelope,
692 ingress: CommandIngress,
693 decision_controller_id: Option<String>,
694 record_attempt: bool,
695 ) -> Result<CommandOutcome, CanwuError> {
696 self.ensure_runtime_ready()?;
697 self.ensure_command_ingress_family(ingress)?;
698 if let Some(cached) =
699 self.cached_command_outcome(request_id, expected_revision, &envelope)?
700 {
701 return Ok(cached);
702 }
703
704 let revision_before = self.revision();
705 let admission = CommandAdmission {
706 request_id,
707 expected_revision,
708 expected_time: envelope.expected_time,
709 revision_before,
710 ingress,
711 };
712 let attempt_id = if record_attempt {
713 let (value, _) = claim_counter(
714 self.state.counters.next_command_attempt_id,
715 "command attempt ID",
716 )?;
717 CommandAttemptId::new(value)
718 } else {
719 CommandAttemptId::default()
720 };
721 let authority = match resolve_command_authority(&envelope) {
722 Ok(authority) => authority,
723 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
724 return self.record_command_rejection(attempt_id, admission, envelope, error);
725 }
726 Err(error) => return Err(error),
727 };
728 if let Err(error) = self.validate_command_ingress(&envelope.issuer, &authority, admission) {
729 if is_expected_command_rejection(&error.code) && record_attempt {
730 return self.record_command_rejection(attempt_id, admission, envelope, error);
731 }
732 return Err(error);
733 }
734 if let Some(expected_time) = envelope.expected_time
735 && expected_time != self.state.scheduler.now
736 {
737 let error = CanwuError::new(
738 ErrorCode::SimulationTimeConflict,
739 format!(
740 "command expected time {expected_time}, but simulation is at {}",
741 self.state.scheduler.now
742 ),
743 );
744 if record_attempt {
745 return self.record_command_rejection(attempt_id, admission, envelope, error);
746 }
747 return Err(error);
748 }
749
750 let (command_id_value, next_command_id) =
751 claim_counter(self.state.counters.next_command_id, "command ID")?;
752 let (correlation_id, next_correlation_id) =
753 claim_counter(self.state.counters.next_correlation_id, "correlation ID")?;
754 let command_id = CommandId::new(command_id_value);
755 let context = CommandContext {
756 issuer: envelope.issuer.clone(),
757 authority,
758 decision_controller_id,
759 run_policy: self.state.metadata.run_configuration.command_policy(),
760 ingress: admission.ingress,
761 attempt_id: record_attempt.then_some(attempt_id),
762 command_id,
763 request_id: admission.request_id,
764 revision: admission.revision_before,
765 simulation_time: self.state.scheduler.now,
766 expected_revision: admission.expected_revision,
767 expected_time: envelope.expected_time,
768 };
769 let prepared = match self.prepare_command(&envelope, &context) {
770 Ok(prepared) => prepared,
771 Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
772 return self.record_command_rejection(attempt_id, admission, envelope, error);
773 }
774 Err(error) => return Err(error),
775 };
776 let next_attempt_id = if record_attempt {
777 let (claimed_id, next_attempt_id) = claim_counter(
778 self.state.counters.next_command_attempt_id,
779 "command attempt ID",
780 )?;
781 if claimed_id != attempt_id.get() {
782 return Err(CanwuError::new(
783 ErrorCode::InvalidSnapshot,
784 "command attempt allocation changed during application",
785 ));
786 }
787 Some(next_attempt_id)
788 } else {
789 None
790 };
791 let revision = self.next_state_revision()?;
792 let transaction = CommandTransactionCheckpoint::capture(&self.state);
793 let event_start = self.state.evidence.events.len();
794 self.state.counters.next_command_id = next_command_id;
795 self.state.counters.next_correlation_id = next_correlation_id;
796 self.invalidate_commitments(prepared.commitment_invalidation());
797
798 if let Err(error) = self.apply_prepared(prepared, command_id, correlation_id) {
799 transaction.restore(&mut self.state);
800 if is_expected_command_rejection(&error.code) && record_attempt {
801 return self.record_command_rejection(attempt_id, admission, envelope, error);
802 }
803 return Err(error);
804 }
805 let emitted_events: Vec<_> = self.state.evidence.events[event_start..]
806 .iter()
807 .map(|event| event.id)
808 .collect();
809 self.state.metadata.plugin_registration_closed = true;
810 self.state.evidence.commands.push(CommandRecord {
811 id: command_id,
812 attempt_id: record_attempt.then_some(attempt_id),
813 accepted_at: self.state.scheduler.now,
814 envelope: envelope.clone(),
815 emitted_events: if record_attempt {
816 emitted_events.clone()
817 } else {
818 Vec::new()
819 },
820 });
821 if let Some(next_attempt_id) = next_attempt_id {
822 self.state.counters.next_command_attempt_id = next_attempt_id;
823 self.state
824 .evidence
825 .command_attempts
826 .push(CommandAttemptRecord {
827 id: attempt_id,
828 at: self.state.scheduler.now,
829 revision_before: admission.revision_before,
830 ingress: admission.ingress,
831 request_id: admission.request_id,
832 expected_revision: admission.expected_revision,
833 envelope,
834 outcome: CommandAttemptOutcome::Accepted { command_id },
835 });
836 }
837 self.state.counters.state_revision = revision;
838 if let Err(error) = self.refresh_checkpoint_hash() {
839 transaction.restore(&mut self.state);
840 return Err(error);
841 }
842
843 Ok(CommandOutcome::Accepted {
844 receipt: CommandReceipt {
845 attempt_id: record_attempt.then_some(attempt_id),
846 command_id,
847 request_id: admission.request_id,
848 revision,
849 accepted_at: self.state.scheduler.now,
850 emitted_events,
851 },
852 })
853 }
854
855 fn ensure_command_ingress_family(&self, ingress: CommandIngress) -> Result<(), CanwuError> {
856 let has_legacy_commands = self.state.evidence.archived_legacy_commands
857 || self
858 .state
859 .evidence
860 .commands
861 .iter()
862 .any(|record| record.attempt_id.is_none());
863 let has_tracked_attempts = self.state.evidence.archived_tracked_attempts
864 || !self.state.evidence.command_attempts.is_empty()
865 || !self.state.evidence.ingress.is_empty();
866 if (ingress == CommandIngress::LegacyDirect && has_tracked_attempts)
867 || (ingress != CommandIngress::LegacyDirect && has_legacy_commands)
868 {
869 return Err(CanwuError::new(
870 ErrorCode::MixedCommandIngress,
871 "legacy-direct commands and tracked request/replay attempts cannot coexist in one run",
872 ));
873 }
874 Ok(())
875 }
876
877 pub(super) fn ensure_canonical_ingress_can_start(&self) -> Result<(), CanwuError> {
878 if runtime_has_unqueued_command_history(&self.state) {
879 return Err(CanwuError::new(
880 ErrorCode::MixedCommandIngress,
881 "canonical ingress cannot be added after direct command history",
882 ));
883 }
884 Ok(())
885 }
886
887 fn cached_command_outcome(
888 &self,
889 request_id: Option<CommandRequestId>,
890 expected_revision: Option<u64>,
891 envelope: &CommandEnvelope,
892 ) -> Result<Option<CommandOutcome>, CanwuError> {
893 let Some(request_id) = request_id else {
894 return Ok(None);
895 };
896 if let Some(cached) = self
897 .state
898 .evidence
899 .archived_command_requests
900 .get(&request_id)
901 {
902 let input_hash = canonical_hash(
903 "canwu.archive.command.request.v1",
904 &(expected_revision, envelope),
905 )?;
906 if cached.input_hash != input_hash {
907 return Ok(Some(CommandOutcome::Rejected {
908 rejection: CommandRejection {
909 attempt_id: None,
910 request_id: Some(request_id),
911 retained_revision: self.revision(),
912 rejected_at: self.state.scheduler.now,
913 error: CanwuError::new(
914 ErrorCode::IdempotencyConflict,
915 "this command request ID was already used for different input",
916 ),
917 },
918 }));
919 }
920 return Ok(Some(cached.outcome.clone()));
921 }
922 let Some(attempt) = self
923 .state
924 .evidence
925 .command_attempts
926 .iter()
927 .find(|attempt| attempt.request_id == Some(request_id))
928 else {
929 return Ok(None);
930 };
931 if attempt.expected_revision != expected_revision || &attempt.envelope != envelope {
932 return Ok(Some(CommandOutcome::Rejected {
933 rejection: CommandRejection {
934 attempt_id: None,
935 request_id: Some(request_id),
936 retained_revision: self.revision(),
937 rejected_at: self.state.scheduler.now,
938 error: CanwuError::new(
939 ErrorCode::IdempotencyConflict,
940 "this command request ID was already used for different input",
941 ),
942 },
943 }));
944 }
945 Ok(Some(self.command_outcome_from_attempt(attempt)?))
946 }
947
948 pub(super) fn command_outcome_from_attempt(
949 &self,
950 attempt: &CommandAttemptRecord,
951 ) -> Result<CommandOutcome, CanwuError> {
952 let request_id = attempt.request_id.ok_or_else(|| {
953 invalid_snapshot_error("tracked command attempt is missing its request ID")
954 })?;
955 let committed_revision = attempt.revision_before.checked_add(1).ok_or_else(|| {
956 invalid_snapshot_error("cached command attempt revision is exhausted")
957 })?;
958 match &attempt.outcome {
959 CommandAttemptOutcome::Accepted { command_id } => {
960 let retained_number = command_id
961 .get()
962 .checked_sub(self.state.evidence.archived.command_count)
963 .and_then(|value| value.checked_sub(1))
964 .ok_or_else(|| {
965 invalid_snapshot_error(
966 "accepted command attempt references archived command evidence",
967 )
968 })?;
969 let index = usize::try_from(retained_number).map_err(|_| {
970 invalid_snapshot_error(
971 "accepted command attempt exceeds the retained command index space",
972 )
973 })?;
974 let record = self
975 .state
976 .evidence
977 .commands
978 .get(index)
979 .filter(|record| record.id == *command_id)
980 .ok_or_else(|| {
981 invalid_snapshot_error(
982 "accepted command attempt references a missing command",
983 )
984 })?;
985 Ok(CommandOutcome::Accepted {
986 receipt: CommandReceipt {
987 attempt_id: Some(attempt.id),
988 command_id: *command_id,
989 request_id: Some(request_id),
990 revision: committed_revision,
991 accepted_at: record.accepted_at,
992 emitted_events: record.emitted_events.clone(),
993 },
994 })
995 }
996 CommandAttemptOutcome::Rejected { error } => Ok(CommandOutcome::Rejected {
997 rejection: CommandRejection {
998 attempt_id: Some(attempt.id),
999 request_id: Some(request_id),
1000 retained_revision: committed_revision,
1001 rejected_at: attempt.at,
1002 error: error.clone(),
1003 },
1004 }),
1005 }
1006 }
1007
1008 fn record_command_rejection(
1009 &mut self,
1010 attempt_id: CommandAttemptId,
1011 admission: CommandAdmission,
1012 envelope: CommandEnvelope,
1013 error: CanwuError,
1014 ) -> Result<CommandOutcome, CanwuError> {
1015 let (claimed_id, next_attempt_id) = claim_counter(
1016 self.state.counters.next_command_attempt_id,
1017 "command attempt ID",
1018 )?;
1019 if claimed_id != attempt_id.get() {
1020 return Err(CanwuError::new(
1021 ErrorCode::InvalidSnapshot,
1022 "command attempt allocation changed during rejection",
1023 ));
1024 }
1025 let revision = self.next_state_revision()?;
1026 let attempt = CommandAttemptRecord {
1027 id: attempt_id,
1028 at: self.state.scheduler.now,
1029 revision_before: admission.revision_before,
1030 ingress: admission.ingress,
1031 request_id: admission.request_id,
1032 expected_revision: admission.expected_revision,
1033 envelope,
1034 outcome: CommandAttemptOutcome::Rejected {
1035 error: error.clone(),
1036 },
1037 };
1038 let transaction = RejectionTransactionCheckpoint::capture(&self.state);
1039 self.state.counters.next_command_attempt_id = next_attempt_id;
1040 self.state.counters.state_revision = revision;
1041 self.state.metadata.plugin_registration_closed = true;
1042 self.state.evidence.command_attempts.push(attempt);
1043 if let Err(hash_error) = self.refresh_checkpoint_hash() {
1044 transaction.restore(&mut self.state);
1045 return Err(hash_error);
1046 }
1047 Ok(CommandOutcome::Rejected {
1048 rejection: CommandRejection {
1049 attempt_id: Some(attempt_id),
1050 request_id: admission.request_id,
1051 retained_revision: revision,
1052 rejected_at: self.state.scheduler.now,
1053 error,
1054 },
1055 })
1056 }
1057
1058 fn validate_command_ingress(
1059 &self,
1060 issuer: &Issuer,
1061 authority: &CommandAuthority,
1062 admission: CommandAdmission,
1063 ) -> Result<(), CanwuError> {
1064 validate_command_ingress_policy(
1065 &self.state.metadata.run_configuration,
1066 issuer,
1067 authority,
1068 admission,
1069 &|entity| runtime_entity_exists(&self.state, entity),
1070 )
1071 }
1072
1073 pub fn advance_canonical(
1074 &mut self,
1075 duration: SimDuration,
1076 ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
1077 self.ensure_runtime_ready()?;
1078 if duration.is_negative() {
1079 return Err(CanwuError::new(
1080 ErrorCode::InvalidDuration,
1081 "canonical simulation time cannot advance by a negative duration",
1082 ));
1083 }
1084 let target = self
1085 .state
1086 .scheduler
1087 .now
1088 .checked_add(duration)
1089 .ok_or_else(|| {
1090 CanwuError::new(
1091 ErrorCode::InvalidDuration,
1092 "canonical simulation target time exceeds the supported range",
1093 )
1094 })?;
1095 let mut receipts = Vec::new();
1096 while let Some(next_due) = self.next_canonical_due_time()
1097 && next_due <= target
1098 {
1099 let at = next_due.max(self.state.scheduler.now);
1100 receipts.push(self.settle_boundary(BoundaryRequest::at(at))?);
1101 }
1102 if self.state.scheduler.now < target {
1103 self.advance_to(target)?;
1104 }
1105 Ok(receipts)
1106 }
1107
1108 pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
1109 self.ensure_runtime_ready()?;
1110 let Some(next_due) = self.next_canonical_due_time() else {
1111 return Ok(None);
1112 };
1113 self.settle_boundary(BoundaryRequest::at(next_due.max(self.state.scheduler.now)))
1114 .map(Some)
1115 }
1116
1117 fn next_canonical_due_time(&self) -> Option<SimTime> {
1118 let scheduled = self.state.scheduler.actions.keys().next().map(|key| key.at);
1119 let ingress = self
1120 .state
1121 .scheduler
1122 .pending_ingress
1123 .first()
1124 .map(|key| key.due_at);
1125 match (scheduled, ingress) {
1126 (Some(left), Some(right)) => Some(left.min(right)),
1127 (Some(value), None) | (None, Some(value)) => Some(value),
1128 (None, None) => None,
1129 }
1130 }
1131
1132 pub(super) fn take_due_ingress(&mut self, at: SimTime) -> Vec<IngressId> {
1133 let mut admitted = Vec::new();
1134 while self
1135 .state
1136 .scheduler
1137 .pending_ingress
1138 .first()
1139 .is_some_and(|key| key.due_at <= at)
1140 {
1141 let key = self
1142 .state
1143 .scheduler
1144 .pending_ingress
1145 .pop_first()
1146 .expect("pending ingress was checked as non-empty");
1147 admitted.push(key.id);
1148 }
1149 admitted
1150 }
1151}