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