Skip to main content

canwu_sim/
ingress.rs

1use super::{
2    BoundaryReceipt, BoundaryRequest, CanwuError, CauseRef, CommandAdmission, CommandAttemptId,
3    CommandAttemptOutcome, CommandAttemptRecord, CommandAuthority, CommandContext, CommandEnvelope,
4    CommandId, CommandIngress, CommandOutcome, CommandReceipt, CommandRecord, CommandRejection,
5    CommandRequest, CommandRequestId, CommandTransactionCheckpoint, Deserialize, EntityRef,
6    ErrorCode, IngressId, IngressTransactionCheckpoint, InteractionPolicy, Issuer, PayloadSchema,
7    RejectionTransactionCheckpoint, Serialize, SimDuration, SimTime, Simulation, SystemCadence,
8    Value, canonical_hash, claim_counter, invalid_snapshot_error, is_expected_command_rejection,
9    resolve_command_authority, runtime_entity_exists, runtime_entity_identity_exists,
10    runtime_has_unqueued_command_history, validate_command_ingress_policy, validate_runtime_cause,
11};
12use std::cmp::Reverse;
13
14#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
15#[serde(rename_all = "snake_case")]
16pub enum IngressClass {
17    Command,
18    Communication,
19    Acknowledgement,
20    Information,
21    ScheduledSystem,
22}
23
24#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
25pub struct PluginIngressDescriptor {
26    pub name: String,
27    pub description: String,
28    pub class: IngressClass,
29    pub payload_schema: PayloadSchema,
30}
31
32#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
33pub struct PluginIngressRequest {
34    pub plugin: String,
35    pub packet_type: String,
36    pub due_at: SimTime,
37    pub priority: i32,
38    pub payload: Value,
39    pub affected_entities: Vec<EntityRef>,
40    #[serde(default, skip_serializing_if = "Option::is_none")]
41    pub cause: Option<CauseRef>,
42}
43
44impl PluginIngressRequest {
45    #[must_use]
46    pub fn new(
47        plugin: impl Into<String>,
48        packet_type: impl Into<String>,
49        due_at: SimTime,
50        payload: Value,
51    ) -> Self {
52        Self {
53            plugin: plugin.into(),
54            packet_type: packet_type.into(),
55            due_at,
56            priority: 0,
57            payload,
58            affected_entities: Vec::new(),
59            cause: None,
60        }
61    }
62
63    #[must_use]
64    pub const fn with_priority(mut self, priority: i32) -> Self {
65        self.priority = priority;
66        self
67    }
68
69    #[must_use]
70    pub fn with_entity(mut self, entity: EntityRef) -> Self {
71        self.affected_entities.push(entity);
72        self
73    }
74
75    #[must_use]
76    pub fn caused_by(mut self, cause: CauseRef) -> Self {
77        self.cause = Some(cause);
78        self
79    }
80}
81
82#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
83#[serde(tag = "type", rename_all = "snake_case")]
84pub enum IngressPayload {
85    Command {
86        request: Box<CommandRequest>,
87    },
88    Plugin {
89        plugin: String,
90        packet_type: String,
91        payload: Value,
92        affected_entities: Vec<EntityRef>,
93    },
94    Calendar {
95        cadences: Vec<SystemCadence>,
96    },
97}
98
99#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
100pub struct IngressRecord {
101    pub id: IngressId,
102    pub issued_at: SimTime,
103    #[serde(default, skip_serializing_if = "is_zero")]
104    pub eligible_boundary_count: u64,
105    pub due_at: SimTime,
106    pub class: IngressClass,
107    pub priority: i32,
108    pub payload: IngressPayload,
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub cause: Option<CauseRef>,
111}
112
113#[allow(clippy::trivially_copy_pass_by_ref)]
114const fn is_zero(value: &u64) -> bool {
115    *value == 0
116}
117
118#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
119pub struct IngressReceipt {
120    pub ingress_id: IngressId,
121    pub issued_at: SimTime,
122    pub due_at: SimTime,
123}
124
125#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
126pub(crate) struct IngressQueueKey {
127    pub due_at: SimTime,
128    pub class: IngressClass,
129    pub priority: Reverse<i32>,
130    pub issued_at: SimTime,
131    pub id: IngressId,
132}
133
134impl IngressQueueKey {
135    #[must_use]
136    pub(crate) const fn from_record(record: &IngressRecord) -> Self {
137        Self {
138            due_at: record.due_at,
139            class: record.class,
140            priority: Reverse(record.priority),
141            issued_at: record.issued_at,
142            id: record.id,
143        }
144    }
145}
146
147impl Simulation {
148    pub fn submit(&mut self, envelope: CommandEnvelope) -> Result<CommandReceipt, CanwuError> {
149        match self.admit_command(None, None, envelope, CommandIngress::LegacyDirect, false)? {
150            CommandOutcome::Accepted { receipt } => Ok(receipt),
151            CommandOutcome::Rejected { rejection } => Err(rejection.error),
152        }
153    }
154
155    pub fn enqueue_command(
156        &mut self,
157        due_at: SimTime,
158        priority: i32,
159        request: CommandRequest,
160    ) -> Result<IngressReceipt, CanwuError> {
161        self.ensure_runtime_ready()?;
162        self.ensure_canonical_ingress_can_start()?;
163        self.ensure_command_ingress_family(CommandIngress::LiveRequest)?;
164        if let Some(existing) = self
165            .state
166            .evidence
167            .archived_ingress_requests
168            .get(&request.request_id)
169        {
170            let input_hash = canonical_hash(
171                "canwu.archive.ingress.command.v1",
172                &(due_at, priority, &request),
173            )?;
174            if existing.input_hash == input_hash {
175                return Ok(existing.receipt.clone());
176            }
177            return Err(CanwuError::new(
178                ErrorCode::IdempotencyConflict,
179                format!(
180                    "command request {} is already queued with different ingress content",
181                    request.request_id
182                ),
183            ));
184        }
185        for record in &self.state.evidence.ingress {
186            let IngressPayload::Command { request: existing } = &record.payload else {
187                continue;
188            };
189            if existing.request_id != request.request_id {
190                continue;
191            }
192            if existing.as_ref() == &request
193                && record.due_at == due_at
194                && record.priority == priority
195            {
196                return Ok(IngressReceipt {
197                    ingress_id: record.id,
198                    issued_at: record.issued_at,
199                    due_at: record.due_at,
200                });
201            }
202            return Err(CanwuError::new(
203                ErrorCode::IdempotencyConflict,
204                format!(
205                    "command request {} is already queued with different ingress content",
206                    request.request_id
207                ),
208            ));
209        }
210        if self
211            .state
212            .evidence
213            .command_attempts
214            .iter()
215            .any(|attempt| attempt.request_id == Some(request.request_id))
216            || self
217                .state
218                .evidence
219                .archived_command_requests
220                .contains_key(&request.request_id)
221        {
222            return Err(CanwuError::new(
223                ErrorCode::IdempotencyConflict,
224                format!(
225                    "command request {} was already processed outside canonical ingress",
226                    request.request_id
227                ),
228            ));
229        }
230        if request
231            .envelope
232            .expected_time
233            .is_some_and(|expected| expected != due_at)
234        {
235            return Err(CanwuError::new(
236                ErrorCode::SimulationTimeConflict,
237                "queued command expected time must equal its due simulation time",
238            ));
239        }
240        self.append_ingress(
241            due_at,
242            IngressClass::Command,
243            priority,
244            IngressPayload::Command {
245                request: Box::new(request),
246            },
247            None,
248            false,
249        )
250    }
251
252    pub fn enqueue_plugin_ingress(
253        &mut self,
254        mut request: PluginIngressRequest,
255    ) -> Result<IngressReceipt, CanwuError> {
256        self.ensure_runtime_ready()?;
257        self.ensure_canonical_ingress_can_start()?;
258        if self
259            .state
260            .metadata
261            .run_configuration
262            .declared()
263            .is_some_and(|configuration| configuration.interaction == InteractionPolicy::ReadOnly)
264        {
265            return Err(CanwuError::new(
266                ErrorCode::InteractionReadOnly,
267                "the run interaction policy rejects newly authored plugin ingress",
268            ));
269        }
270        let key = (request.plugin.clone(), request.packet_type.clone());
271        let descriptor = self.plugins.ingress.get(&key).ok_or_else(|| {
272            CanwuError::new(
273                ErrorCode::InvalidPayload,
274                format!(
275                    "plugin ingress type {}.{} is not registered",
276                    request.plugin, request.packet_type
277                ),
278            )
279        })?;
280        descriptor.payload_schema.validate(&request.payload)?;
281        request.affected_entities.sort();
282        request.affected_entities.dedup();
283        if request
284            .affected_entities
285            .iter()
286            .any(|entity| !runtime_entity_identity_exists(&self.state, entity))
287        {
288            return Err(CanwuError::new(
289                ErrorCode::EntityNotFound,
290                "plugin ingress references an unknown entity identity",
291            ));
292        }
293        if let Some(cause) = &request.cause {
294            if matches!(
295                cause,
296                CauseRef::Boundary(_) | CauseRef::Command(_) | CauseRef::Event(_)
297            ) {
298                return Err(CanwuError::new(
299                    ErrorCode::InvalidPayload,
300                    "boundary, command, and event causes are reserved for plugin-generated ingress",
301                ));
302            }
303            validate_runtime_cause(&self.state, cause)?;
304        }
305        self.append_ingress(
306            request.due_at,
307            descriptor.class,
308            request.priority,
309            IngressPayload::Plugin {
310                plugin: request.plugin,
311                packet_type: request.packet_type,
312                payload: request.payload,
313                affected_entities: request.affected_entities,
314            },
315            request.cause,
316            false,
317        )
318    }
319
320    pub fn schedule_calendar_boundary(
321        &mut self,
322        due_at: SimTime,
323        mut cadences: Vec<SystemCadence>,
324    ) -> Result<IngressReceipt, CanwuError> {
325        self.ensure_runtime_ready()?;
326        self.ensure_canonical_ingress_can_start()?;
327        if cadences.contains(&SystemCadence::EventDriven) {
328            return Err(CanwuError::new(
329                ErrorCode::InvalidBoundary,
330                "calendar ingress cannot declare event-driven cadence",
331            ));
332        }
333        cadences.sort();
334        cadences.dedup();
335        if cadences.is_empty() {
336            return Err(CanwuError::new(
337                ErrorCode::InvalidBoundary,
338                "calendar ingress requires at least one scheduled cadence",
339            ));
340        }
341        self.append_ingress(
342            due_at,
343            IngressClass::ScheduledSystem,
344            0,
345            IngressPayload::Calendar { cadences },
346            Some(CauseRef::System("canwu.core.calendar".to_owned())),
347            false,
348        )
349    }
350
351    pub(super) fn append_ingress(
352        &mut self,
353        due_at: SimTime,
354        class: IngressClass,
355        priority: i32,
356        payload: IngressPayload,
357        cause: Option<CauseRef>,
358        after_current_boundary: bool,
359    ) -> Result<IngressReceipt, CanwuError> {
360        if due_at < self.state.scheduler.now {
361            return Err(CanwuError::new(
362                ErrorCode::LateIngress,
363                format!(
364                    "ingress due at {due_at} cannot be queued after committed time {}",
365                    self.state.scheduler.now
366                ),
367            ));
368        }
369        let transaction = IngressTransactionCheckpoint::capture(&self.state);
370        let (id, next_id) = claim_counter(self.state.counters.next_ingress_id, "ingress ID")?;
371        let boundary_count = self
372            .state
373            .evidence
374            .archived
375            .boundary_count
376            .checked_add(
377                u64::try_from(self.state.evidence.boundaries.len()).map_err(|_| {
378                    CanwuError::new(
379                        ErrorCode::IdentifierExhausted,
380                        "boundary count exceeds the ingress journal range",
381                    )
382                })?,
383            )
384            .ok_or_else(|| {
385                CanwuError::new(
386                    ErrorCode::IdentifierExhausted,
387                    "boundary count exceeds the ingress journal range",
388                )
389            })?;
390        let eligible_boundary_count = if after_current_boundary {
391            boundary_count.checked_add(1).ok_or_else(|| {
392                CanwuError::new(
393                    ErrorCode::IdentifierExhausted,
394                    "ingress boundary eligibility exceeds the journal range",
395                )
396            })?
397        } else {
398            boundary_count
399        };
400        let record = IngressRecord {
401            id: IngressId::new(id),
402            issued_at: self.state.scheduler.now,
403            eligible_boundary_count,
404            due_at,
405            class,
406            priority,
407            payload,
408            cause,
409        };
410        let queue_key = IngressQueueKey::from_record(&record);
411        self.state.counters.next_ingress_id = next_id;
412        self.state.scheduler.pending_ingress.insert(queue_key);
413        self.state.evidence.ingress.push(record.clone());
414        self.state.metadata.plugin_registration_closed = true;
415        if let Err(error) = self.refresh_checkpoint_hash() {
416            transaction.restore(&mut self.state, &queue_key);
417            return Err(error);
418        }
419        Ok(IngressReceipt {
420            ingress_id: record.id,
421            issued_at: record.issued_at,
422            due_at: record.due_at,
423        })
424    }
425
426    pub fn process_command(
427        &mut self,
428        request: CommandRequest,
429    ) -> Result<CommandOutcome, CanwuError> {
430        self.ensure_runtime_ready()?;
431        if self.state.evidence.archived.ingress_count != 0
432            || !self.state.evidence.ingress.is_empty()
433        {
434            return Err(CanwuError::new(
435                ErrorCode::MixedCommandIngress,
436                "direct command requests cannot bypass an active canonical ingress journal",
437            ));
438        }
439        self.admit_command(
440            Some(request.request_id),
441            Some(request.expected_revision),
442            request.envelope,
443            CommandIngress::LiveRequest,
444            true,
445        )
446    }
447
448    pub(super) fn admit_command(
449        &mut self,
450        request_id: Option<CommandRequestId>,
451        expected_revision: Option<u64>,
452        envelope: CommandEnvelope,
453        ingress: CommandIngress,
454        record_attempt: bool,
455    ) -> Result<CommandOutcome, CanwuError> {
456        self.ensure_runtime_ready()?;
457        self.ensure_command_ingress_family(ingress)?;
458        if let Some(cached) =
459            self.cached_command_outcome(request_id, expected_revision, &envelope)?
460        {
461            return Ok(cached);
462        }
463
464        let revision_before = self.revision();
465        let admission = CommandAdmission {
466            request_id,
467            expected_revision,
468            expected_time: envelope.expected_time,
469            revision_before,
470            ingress,
471        };
472        let attempt_id = if record_attempt {
473            let (value, _) = claim_counter(
474                self.state.counters.next_command_attempt_id,
475                "command attempt ID",
476            )?;
477            CommandAttemptId::new(value)
478        } else {
479            CommandAttemptId::default()
480        };
481        let authority = match resolve_command_authority(&envelope) {
482            Ok(authority) => authority,
483            Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
484                return self.record_command_rejection(attempt_id, admission, envelope, error);
485            }
486            Err(error) => return Err(error),
487        };
488        if let Err(error) = self.validate_command_ingress(&envelope.issuer, &authority, admission) {
489            if is_expected_command_rejection(&error.code) && record_attempt {
490                return self.record_command_rejection(attempt_id, admission, envelope, error);
491            }
492            return Err(error);
493        }
494        if let Some(expected_time) = envelope.expected_time
495            && expected_time != self.state.scheduler.now
496        {
497            let error = CanwuError::new(
498                ErrorCode::SimulationTimeConflict,
499                format!(
500                    "command expected time {expected_time}, but simulation is at {}",
501                    self.state.scheduler.now
502                ),
503            );
504            if record_attempt {
505                return self.record_command_rejection(attempt_id, admission, envelope, error);
506            }
507            return Err(error);
508        }
509
510        let (command_id_value, next_command_id) =
511            claim_counter(self.state.counters.next_command_id, "command ID")?;
512        let (correlation_id, next_correlation_id) =
513            claim_counter(self.state.counters.next_correlation_id, "correlation ID")?;
514        let command_id = CommandId::new(command_id_value);
515        let context = CommandContext {
516            issuer: envelope.issuer.clone(),
517            authority,
518            run_policy: self.state.metadata.run_configuration.command_policy(),
519            ingress: admission.ingress,
520            attempt_id: record_attempt.then_some(attempt_id),
521            command_id,
522            request_id: admission.request_id,
523            revision: admission.revision_before,
524            simulation_time: self.state.scheduler.now,
525            expected_revision: admission.expected_revision,
526            expected_time: envelope.expected_time,
527        };
528        let prepared = match self.prepare_command(&envelope, &context) {
529            Ok(prepared) => prepared,
530            Err(error) if is_expected_command_rejection(&error.code) && record_attempt => {
531                return self.record_command_rejection(attempt_id, admission, envelope, error);
532            }
533            Err(error) => return Err(error),
534        };
535        let next_attempt_id = if record_attempt {
536            let (claimed_id, next_attempt_id) = claim_counter(
537                self.state.counters.next_command_attempt_id,
538                "command attempt ID",
539            )?;
540            if claimed_id != attempt_id.get() {
541                return Err(CanwuError::new(
542                    ErrorCode::InvalidSnapshot,
543                    "command attempt allocation changed during application",
544                ));
545            }
546            Some(next_attempt_id)
547        } else {
548            None
549        };
550        let revision = self.next_state_revision()?;
551        let transaction = CommandTransactionCheckpoint::capture(&self.state);
552        let event_start = self.state.evidence.events.len();
553        self.state.counters.next_command_id = next_command_id;
554        self.state.counters.next_correlation_id = next_correlation_id;
555        self.invalidate_commitments(prepared.commitment_invalidation());
556
557        if let Err(error) = self.apply_prepared(prepared, command_id, correlation_id) {
558            transaction.restore(&mut self.state);
559            if is_expected_command_rejection(&error.code) && record_attempt {
560                return self.record_command_rejection(attempt_id, admission, envelope, error);
561            }
562            return Err(error);
563        }
564        let emitted_events: Vec<_> = self.state.evidence.events[event_start..]
565            .iter()
566            .map(|event| event.id)
567            .collect();
568        self.state.metadata.plugin_registration_closed = true;
569        self.state.evidence.commands.push(CommandRecord {
570            id: command_id,
571            attempt_id: record_attempt.then_some(attempt_id),
572            accepted_at: self.state.scheduler.now,
573            envelope: envelope.clone(),
574            emitted_events: if record_attempt {
575                emitted_events.clone()
576            } else {
577                Vec::new()
578            },
579        });
580        if let Some(next_attempt_id) = next_attempt_id {
581            self.state.counters.next_command_attempt_id = next_attempt_id;
582            self.state
583                .evidence
584                .command_attempts
585                .push(CommandAttemptRecord {
586                    id: attempt_id,
587                    at: self.state.scheduler.now,
588                    revision_before: admission.revision_before,
589                    ingress: admission.ingress,
590                    request_id: admission.request_id,
591                    expected_revision: admission.expected_revision,
592                    envelope,
593                    outcome: CommandAttemptOutcome::Accepted { command_id },
594                });
595        }
596        self.state.counters.state_revision = revision;
597        if let Err(error) = self.refresh_checkpoint_hash() {
598            transaction.restore(&mut self.state);
599            return Err(error);
600        }
601
602        Ok(CommandOutcome::Accepted {
603            receipt: CommandReceipt {
604                attempt_id: record_attempt.then_some(attempt_id),
605                command_id,
606                request_id: admission.request_id,
607                revision,
608                accepted_at: self.state.scheduler.now,
609                emitted_events,
610            },
611        })
612    }
613
614    fn ensure_command_ingress_family(&self, ingress: CommandIngress) -> Result<(), CanwuError> {
615        let has_legacy_commands = self.state.evidence.archived_legacy_commands
616            || self
617                .state
618                .evidence
619                .commands
620                .iter()
621                .any(|record| record.attempt_id.is_none());
622        let has_tracked_attempts = self.state.evidence.archived_tracked_attempts
623            || !self.state.evidence.command_attempts.is_empty()
624            || !self.state.evidence.ingress.is_empty();
625        if (ingress == CommandIngress::LegacyDirect && has_tracked_attempts)
626            || (ingress != CommandIngress::LegacyDirect && has_legacy_commands)
627        {
628            return Err(CanwuError::new(
629                ErrorCode::MixedCommandIngress,
630                "legacy-direct commands and tracked request/replay attempts cannot coexist in one run",
631            ));
632        }
633        Ok(())
634    }
635
636    pub(super) fn ensure_canonical_ingress_can_start(&self) -> Result<(), CanwuError> {
637        if runtime_has_unqueued_command_history(&self.state) {
638            return Err(CanwuError::new(
639                ErrorCode::MixedCommandIngress,
640                "canonical ingress cannot be added after direct command history",
641            ));
642        }
643        Ok(())
644    }
645
646    fn cached_command_outcome(
647        &self,
648        request_id: Option<CommandRequestId>,
649        expected_revision: Option<u64>,
650        envelope: &CommandEnvelope,
651    ) -> Result<Option<CommandOutcome>, CanwuError> {
652        let Some(request_id) = request_id else {
653            return Ok(None);
654        };
655        if let Some(cached) = self
656            .state
657            .evidence
658            .archived_command_requests
659            .get(&request_id)
660        {
661            let input_hash = canonical_hash(
662                "canwu.archive.command.request.v1",
663                &(expected_revision, envelope),
664            )?;
665            if cached.input_hash != input_hash {
666                return Ok(Some(CommandOutcome::Rejected {
667                    rejection: CommandRejection {
668                        attempt_id: None,
669                        request_id: Some(request_id),
670                        retained_revision: self.revision(),
671                        rejected_at: self.state.scheduler.now,
672                        error: CanwuError::new(
673                            ErrorCode::IdempotencyConflict,
674                            "this command request ID was already used for different input",
675                        ),
676                    },
677                }));
678            }
679            return Ok(Some(cached.outcome.clone()));
680        }
681        let Some(attempt) = self
682            .state
683            .evidence
684            .command_attempts
685            .iter()
686            .find(|attempt| attempt.request_id == Some(request_id))
687        else {
688            return Ok(None);
689        };
690        if attempt.expected_revision != expected_revision || &attempt.envelope != envelope {
691            return Ok(Some(CommandOutcome::Rejected {
692                rejection: CommandRejection {
693                    attempt_id: None,
694                    request_id: Some(request_id),
695                    retained_revision: self.revision(),
696                    rejected_at: self.state.scheduler.now,
697                    error: CanwuError::new(
698                        ErrorCode::IdempotencyConflict,
699                        "this command request ID was already used for different input",
700                    ),
701                },
702            }));
703        }
704        Ok(Some(self.command_outcome_from_attempt(attempt)?))
705    }
706
707    pub(super) fn command_outcome_from_attempt(
708        &self,
709        attempt: &CommandAttemptRecord,
710    ) -> Result<CommandOutcome, CanwuError> {
711        let request_id = attempt.request_id.ok_or_else(|| {
712            invalid_snapshot_error("tracked command attempt is missing its request ID")
713        })?;
714        let committed_revision = attempt.revision_before.checked_add(1).ok_or_else(|| {
715            invalid_snapshot_error("cached command attempt revision is exhausted")
716        })?;
717        match &attempt.outcome {
718            CommandAttemptOutcome::Accepted { command_id } => {
719                let retained_number = command_id
720                    .get()
721                    .checked_sub(self.state.evidence.archived.command_count)
722                    .and_then(|value| value.checked_sub(1))
723                    .ok_or_else(|| {
724                        invalid_snapshot_error(
725                            "accepted command attempt references archived command evidence",
726                        )
727                    })?;
728                let index = usize::try_from(retained_number).map_err(|_| {
729                    invalid_snapshot_error(
730                        "accepted command attempt exceeds the retained command index space",
731                    )
732                })?;
733                let record = self
734                    .state
735                    .evidence
736                    .commands
737                    .get(index)
738                    .filter(|record| record.id == *command_id)
739                    .ok_or_else(|| {
740                        invalid_snapshot_error(
741                            "accepted command attempt references a missing command",
742                        )
743                    })?;
744                Ok(CommandOutcome::Accepted {
745                    receipt: CommandReceipt {
746                        attempt_id: Some(attempt.id),
747                        command_id: *command_id,
748                        request_id: Some(request_id),
749                        revision: committed_revision,
750                        accepted_at: record.accepted_at,
751                        emitted_events: record.emitted_events.clone(),
752                    },
753                })
754            }
755            CommandAttemptOutcome::Rejected { error } => Ok(CommandOutcome::Rejected {
756                rejection: CommandRejection {
757                    attempt_id: Some(attempt.id),
758                    request_id: Some(request_id),
759                    retained_revision: committed_revision,
760                    rejected_at: attempt.at,
761                    error: error.clone(),
762                },
763            }),
764        }
765    }
766
767    fn record_command_rejection(
768        &mut self,
769        attempt_id: CommandAttemptId,
770        admission: CommandAdmission,
771        envelope: CommandEnvelope,
772        error: CanwuError,
773    ) -> Result<CommandOutcome, CanwuError> {
774        let (claimed_id, next_attempt_id) = claim_counter(
775            self.state.counters.next_command_attempt_id,
776            "command attempt ID",
777        )?;
778        if claimed_id != attempt_id.get() {
779            return Err(CanwuError::new(
780                ErrorCode::InvalidSnapshot,
781                "command attempt allocation changed during rejection",
782            ));
783        }
784        let revision = self.next_state_revision()?;
785        let attempt = CommandAttemptRecord {
786            id: attempt_id,
787            at: self.state.scheduler.now,
788            revision_before: admission.revision_before,
789            ingress: admission.ingress,
790            request_id: admission.request_id,
791            expected_revision: admission.expected_revision,
792            envelope,
793            outcome: CommandAttemptOutcome::Rejected {
794                error: error.clone(),
795            },
796        };
797        let transaction = RejectionTransactionCheckpoint::capture(&self.state);
798        self.state.counters.next_command_attempt_id = next_attempt_id;
799        self.state.counters.state_revision = revision;
800        self.state.metadata.plugin_registration_closed = true;
801        self.state.evidence.command_attempts.push(attempt);
802        if let Err(hash_error) = self.refresh_checkpoint_hash() {
803            transaction.restore(&mut self.state);
804            return Err(hash_error);
805        }
806        Ok(CommandOutcome::Rejected {
807            rejection: CommandRejection {
808                attempt_id: Some(attempt_id),
809                request_id: admission.request_id,
810                retained_revision: revision,
811                rejected_at: self.state.scheduler.now,
812                error,
813            },
814        })
815    }
816
817    fn validate_command_ingress(
818        &self,
819        issuer: &Issuer,
820        authority: &CommandAuthority,
821        admission: CommandAdmission,
822    ) -> Result<(), CanwuError> {
823        validate_command_ingress_policy(
824            &self.state.metadata.run_configuration,
825            issuer,
826            authority,
827            admission,
828            &|entity| runtime_entity_exists(&self.state, entity),
829        )
830    }
831
832    pub fn advance_canonical(
833        &mut self,
834        duration: SimDuration,
835    ) -> Result<Vec<BoundaryReceipt>, CanwuError> {
836        self.ensure_runtime_ready()?;
837        if duration.is_negative() {
838            return Err(CanwuError::new(
839                ErrorCode::InvalidDuration,
840                "canonical simulation time cannot advance by a negative duration",
841            ));
842        }
843        let target = self
844            .state
845            .scheduler
846            .now
847            .checked_add(duration)
848            .ok_or_else(|| {
849                CanwuError::new(
850                    ErrorCode::InvalidDuration,
851                    "canonical simulation target time exceeds the supported range",
852                )
853            })?;
854        let mut receipts = Vec::new();
855        while let Some(next_due) = self.next_canonical_due_time()
856            && next_due <= target
857        {
858            let at = next_due.max(self.state.scheduler.now);
859            receipts.push(self.settle_boundary(BoundaryRequest::at(at))?);
860        }
861        if self.state.scheduler.now < target {
862            self.advance_to(target)?;
863        }
864        Ok(receipts)
865    }
866
867    pub fn step_canonical(&mut self) -> Result<Option<BoundaryReceipt>, CanwuError> {
868        self.ensure_runtime_ready()?;
869        let Some(next_due) = self.next_canonical_due_time() else {
870            return Ok(None);
871        };
872        self.settle_boundary(BoundaryRequest::at(next_due.max(self.state.scheduler.now)))
873            .map(Some)
874    }
875
876    fn next_canonical_due_time(&self) -> Option<SimTime> {
877        let scheduled = self.state.scheduler.actions.keys().next().map(|key| key.at);
878        let ingress = self
879            .state
880            .scheduler
881            .pending_ingress
882            .first()
883            .map(|key| key.due_at);
884        match (scheduled, ingress) {
885            (Some(left), Some(right)) => Some(left.min(right)),
886            (Some(value), None) | (None, Some(value)) => Some(value),
887            (None, None) => None,
888        }
889    }
890
891    pub(super) fn take_due_ingress(&mut self, at: SimTime) -> Vec<IngressId> {
892        let mut admitted = Vec::new();
893        while self
894            .state
895            .scheduler
896            .pending_ingress
897            .first()
898            .is_some_and(|key| key.due_at <= at)
899        {
900            let key = self
901                .state
902                .scheduler
903                .pending_ingress
904                .pop_first()
905                .expect("pending ingress was checked as non-empty");
906            admitted.push(key.id);
907        }
908        admitted
909    }
910}