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}