Skip to main content

canwu_sim/runtime/
ingress.rs

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    /// Present only when the engine admitted this command as the selected
89    /// action of a validated `DecisionTicket` controller.
90    #[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    /// Must equal the persisted authoritative revision at command admission.
160    ///
161    /// Accepted commands, persisted expected rejections, and completed
162    /// settlement boundaries advance the revision. Bare clock movement does
163    /// not, so declared external commands also carry `envelope.expected_time`.
164    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    /// Authoritative revision after the accepted command commits.
209    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    /// Authoritative revision after persisted rejection evidence commits.
219    /// Non-persisted conflicts retain the already committed current revision.
220    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    /// Authoritative revision immediately before this attempt transaction.
244    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}