Skip to main content

canwu_sim/runtime/
ingress.rs

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    /// Present only when the engine admitted this command as the selected
90    /// action of a validated `DecisionTicket` controller.
91    #[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    /// Must equal the persisted authoritative revision at command admission.
161    ///
162    /// Accepted commands, persisted expected rejections, and completed
163    /// settlement boundaries advance the revision. Bare clock movement does
164    /// not, so declared external commands also carry `envelope.expected_time`.
165    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    /// Authoritative revision after the accepted command commits.
210    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    /// Authoritative revision after persisted rejection evidence commits.
220    /// Non-persisted conflicts retain the already committed current revision.
221    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    /// Authoritative revision immediately before this attempt transaction.
245    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/// Opaque capability issued only while a plugin registers a kernel-internal
294/// ingress type. Hosts can pass the capability back but cannot construct or
295/// alter it.
296#[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    /// Terminal withdrawal of one still-pending plugin ingress item, recorded
415    /// strictly before that item's due time.
416    ///
417    /// The record is never queued or admitted itself: it takes effect when it
418    /// is appended, its `due_at` equals its `issued_at`, and the withdrawn
419    /// item never reaches a boundary. Nothing is rolled back; the withdrawn
420    /// record stays in the journal as evidence.
421    PluginCancellation {
422        /// The withdrawn plugin ingress item.
423        cancelled: IngressId,
424        authority: IngressCancellationAuthority,
425        reason: String,
426    },
427}
428
429/// Authority that withdrew a queued plugin ingress item.
430///
431/// Only the issuer of an item may withdraw it: the trusted host for items it
432/// enqueued through the ordinary host API, the owning plugin's opaque
433/// registration permit for its internal packet types, or a boundary system of
434/// the plugin that scheduled the item inside the engine.
435#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
436#[serde(rename_all = "snake_case")]
437pub enum IngressCancellationAuthority {
438    /// Host withdrawal of a host-enqueued item whose packet type any host may
439    /// enqueue without a permit.
440    Host,
441    /// Withdrawal through the owning plugin's [`PluginIngressPermit`] for the
442    /// item's exact internal packet type. Covers host-enqueued items of that
443    /// type and items the same plugin scheduled inside the engine.
444    PluginPermit,
445    /// Withdrawal by a boundary system of the plugin that scheduled the item
446    /// inside the engine, through [`super::BoundaryDirective::CancelPluginIngress`].
447    BoundarySystem,
448}
449
450/// Maximum UTF-8 byte length of a plugin ingress cancellation reason.
451pub 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/// Evidence presented with one cancellation request.
460#[derive(Clone, Copy)]
461pub(super) struct PluginIngressCancellationProof<'a> {
462    /// Opaque permit presented by a live host call.
463    pub(super) permit: Option<&'a PluginIngressPermit>,
464    /// Exact replay reproduces a recorded permit withdrawal without the token.
465    pub(super) replay: bool,
466    /// Plugin whose boundary system proposed the cancellation directive.
467    pub(super) boundary_plugin: Option<&'a str>,
468    /// Generation evidence staged earlier in the current boundary.
469    pub(super) current_generations: &'a [BoundaryIngressGeneration],
470}
471
472/// Who issued a queued plugin ingress item, as far as cancellation authority
473/// is concerned.
474#[derive(Clone, Copy, Debug, Eq, PartialEq)]
475pub(super) enum PluginIngressIssuer<'a> {
476    /// Enqueued by the trusted host (with or without a permit).
477    Host,
478    /// Scheduled inside the engine by this plugin's boundary system or command.
479    Plugin(&'a str),
480}
481
482/// Checks the issuer-only authority rule shared by live cancellation and
483/// snapshot validation. The caller has already proved any permit token.
484pub(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    /// Queues a plugin-owned internal ingress through an opaque capability
663    /// returned by [`super::plugins::PluginRegistrar::register_internal_ingress`].
664    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    /// Withdraws a pending plugin ingress item that the host enqueued through
680    /// [`Self::enqueue_plugin_ingress`].
681    ///
682    /// The item must still be queued and strictly before its due time;
683    /// otherwise (including archived IDs) the call fails with
684    /// [`ErrorCode::LateIngress`]. Items of internal packet types require
685    /// [`Self::cancel_permitted_plugin_ingress`], and items scheduled inside
686    /// the engine can be withdrawn only by their issuing plugin; both fail
687    /// here with [`ErrorCode::InvalidAuthority`]. Unknown IDs fail with
688    /// [`ErrorCode::EvidenceUnavailable`]; non-plugin targets and reasons
689    /// that are empty, untrimmed, or longer than
690    /// [`MAX_INGRESS_CANCELLATION_REASON_BYTES`] fail with
691    /// [`ErrorCode::InvalidPayload`]; declared read-only runs fail with
692    /// [`ErrorCode::InteractionReadOnly`].
693    ///
694    /// The withdrawal is appended to the ingress journal as a terminal
695    /// [`IngressPayload::PluginCancellation`] record; the returned receipt
696    /// names that record. The withdrawn item is never admitted and never
697    /// settles, and nothing is rolled back.
698    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    /// Withdraws a pending plugin ingress item of an internal packet type
713    /// through the owning plugin's opaque registration permit.
714    ///
715    /// The permit must match the item's exact plugin and packet type. It
716    /// covers host-enqueued items of that type and items the same plugin
717    /// scheduled inside the engine, but not items another plugin scheduled
718    /// into this packet type. Timing and journal rules match
719    /// [`Self::cancel_plugin_ingress`].
720    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    /// Resolves and authorizes one cancellation target, returning its exact
790    /// pending queue entry.
791    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    /// Returns who issued a retained plugin ingress record.
876    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    /// Appends one terminal cancellation record and removes its target from
936    /// the canonical queue in the same transaction.
937    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    /// Claims the next ingress identifier and its journal eligibility cut.
1182    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}