1use super::{
2 BoundaryId, CanwuError, CauseRef, Command, CommandAuthority, CommandEnvelope, CommandIngress,
3 CommandOutcome, CommandRequest, CommandRequestId, ControllerDecision, DecisionAction,
4 DecisionAttemptErrorCode, DecisionAttemptOutcome, DecisionAttemptRecord, DecisionAuthority,
5 DecisionController, DecisionError, DecisionMutation, DecisionPolicy, DecisionPolicyKind,
6 DecisionRequestId, DecisionTicket, DecisionTicketId, DecisionTrace, DecisionTraceId, EntityRef,
7 ErrorCode, IngressClass, IngressPayload, IngressReceipt, Issuer, MaintenanceChangeRecord,
8 MaintenanceDisposition, MaintenanceIngressRequest, MaintenanceRejectionReceipt, SimTime,
9 Simulation, VerifiedDecisionArchiveCommit, canonical_hash, claim_counter,
10 invalid_snapshot_error, runtime_entity_identity_exists,
11};
12use serde::{Deserialize, Serialize};
13
14pub const DECISION_REQUEST_COMMITMENT_DOMAIN: &str = "canwu.decision.ingress-request.v1";
15
16#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
17pub struct DecisionIngressRequest {
18 pub request_id: DecisionRequestId,
19 pub expected_revision: u64,
20 pub mutation: DecisionMutation,
21 #[serde(default, skip_serializing_if = "Option::is_none")]
22 pub command: Option<Box<CommandRequest>>,
23}
24
25impl DecisionIngressRequest {
26 #[must_use]
27 pub const fn new(
28 request_id: DecisionRequestId,
29 expected_revision: u64,
30 mutation: DecisionMutation,
31 ) -> Self {
32 Self {
33 request_id,
34 expected_revision,
35 mutation,
36 command: None,
37 }
38 }
39
40 #[must_use]
41 pub fn with_command(mut self, command: CommandRequest) -> Self {
42 self.command = Some(Box::new(command));
43 self
44 }
45
46 #[must_use]
50 pub fn carries_random_evidence(&self) -> bool {
51 matches!(
52 &self.mutation,
53 DecisionMutation::Resolve { decision, .. } if decision.random.is_some()
54 )
55 }
56}
57
58#[derive(Clone, Debug, PartialEq)]
59pub enum DecisionEvaluation {
60 Pending(super::PolicyDecision),
61 Prepared(PreparedDecisionIngress),
62}
63
64#[derive(Clone, Debug, PartialEq)]
65pub struct PreparedDecisionIngress {
66 pub request: DecisionIngressRequest,
67 pub selected_action: Option<DecisionAction>,
68}
69
70impl Simulation {
71 #[must_use]
72 pub fn decision_ticket(&self, id: DecisionTicketId) -> Option<&DecisionTicket> {
73 self.state.current.decisions.ticket(id)
74 }
75
76 #[must_use]
77 pub fn decision_controller(&self, id: &str) -> Option<&super::DecisionControllerBinding> {
78 self.state.current.decisions.controller(id)
79 }
80
81 #[must_use]
82 pub fn decision_trace(&self, id: DecisionTraceId) -> Option<&DecisionTrace> {
83 self.state.current.decisions.trace(id)
84 }
85
86 #[must_use]
87 pub fn decision_attempt(&self, id: DecisionRequestId) -> Option<&DecisionAttemptRecord> {
88 self.state.current.decisions.attempt(id)
89 }
90
91 #[must_use]
92 pub fn decision_hot_state(&self) -> super::DecisionHotState {
93 self.state.current.decisions.decision_hot_state()
94 }
95
96 #[must_use]
97 pub fn decision_history_location(
98 &self,
99 key: &super::DecisionHistoryKey,
100 ) -> super::DecisionHistoryLocation {
101 self.state.current.decisions.decision_locator(key)
102 }
103
104 pub fn decision_history_location_with_provider(
105 &self,
106 key: &super::DecisionHistoryKey,
107 provider: &dyn super::DecisionArchiveProvider,
108 ) -> Result<super::DecisionHistoryLocation, CanwuError> {
109 self.state
110 .current
111 .decisions
112 .decision_locator_with_provider(key, provider)
113 .map_err(decision_error)
114 }
115
116 pub fn prepare_decision(
117 &self,
118 decision_request_id: DecisionRequestId,
119 command_request_id: Option<CommandRequestId>,
120 ticket_id: DecisionTicketId,
121 policy: &dyn DecisionPolicy,
122 ) -> Result<DecisionEvaluation, CanwuError> {
123 self.prepare_decision_at(
124 self.state.scheduler.now,
125 decision_request_id,
126 command_request_id,
127 ticket_id,
128 policy,
129 )
130 }
131
132 pub fn prepare_decision_at(
133 &self,
134 due_at: SimTime,
135 decision_request_id: DecisionRequestId,
136 command_request_id: Option<CommandRequestId>,
137 ticket_id: DecisionTicketId,
138 policy: &dyn DecisionPolicy,
139 ) -> Result<DecisionEvaluation, CanwuError> {
140 self.ensure_runtime_ready()?;
141 if due_at < self.state.scheduler.now {
142 return Err(CanwuError::new(
143 ErrorCode::SimulationTimeConflict,
144 "a decision cannot be prepared behind committed simulation time",
145 ));
146 }
147 let ticket = self.decision_ticket(ticket_id).ok_or_else(|| {
148 CanwuError::new(
149 ErrorCode::InvalidDecision,
150 format!("decision ticket {ticket_id} was not found"),
151 )
152 })?;
153 if ticket.deadline.is_some_and(|deadline| deadline < due_at) {
154 return Err(CanwuError::new(
155 ErrorCode::InvalidDecision,
156 format!("decision ticket {ticket_id} has expired"),
157 ));
158 }
159 let controller = self
160 .state
161 .current
162 .decisions
163 .controller(&ticket.assigned_controller)
164 .ok_or_else(|| {
165 CanwuError::new(
166 ErrorCode::InvalidDecision,
167 "decision ticket names an unknown controller",
168 )
169 })?;
170 super::persons::validate_decision_preparation(
171 &self.state.current.person_availability,
172 ticket,
173 controller,
174 )?;
175 match DecisionController::evaluate(ticket, controller, policy).map_err(decision_error)? {
176 ControllerDecision::Pending(decision) => Ok(DecisionEvaluation::Pending(decision)),
177 ControllerDecision::Authoritative { decision, action } => {
178 let command = match &action {
179 Some(DecisionAction::Command { command }) => {
180 let request_id = command_request_id.ok_or_else(|| {
181 CanwuError::new(
182 ErrorCode::InvalidDecision,
183 "a selected command option requires a command request ID",
184 )
185 })?;
186 let command: Command =
187 serde_json::from_value(command.clone()).map_err(|error| {
188 CanwuError::new(
189 ErrorCode::InvalidDecision,
190 format!("decision option contains an invalid command: {error}"),
191 )
192 })?;
193 Some(CommandRequest::new(
194 request_id,
195 self.revision(),
196 CommandEnvelope::new(controller_issuer(controller), command)
197 .with_authority(controller_authority(controller))
198 .at_time(due_at),
199 ))
200 }
201 Some(DecisionAction::None) | None => {
202 if command_request_id.is_some() {
203 return Err(CanwuError::new(
204 ErrorCode::InvalidDecision,
205 "a non-command option cannot reserve a command request ID",
206 ));
207 }
208 None
209 }
210 };
211 let mutation = DecisionMutation::Resolve {
212 ticket_id,
213 expected_version: ticket.version,
214 controller_id: controller.id.clone(),
215 policy: controller.policy.clone(),
216 decision,
217 command_request_id,
218 };
219 let request = DecisionIngressRequest {
220 request_id: decision_request_id,
221 expected_revision: self.revision(),
222 mutation,
223 command: command.map(Box::new),
224 };
225 Ok(DecisionEvaluation::Prepared(PreparedDecisionIngress {
226 request,
227 selected_action: action,
228 }))
229 }
230 }
231 }
232
233 pub fn enqueue_decision(
234 &mut self,
235 due_at: SimTime,
236 priority: i32,
237 request: DecisionIngressRequest,
238 ) -> Result<IngressReceipt, CanwuError> {
239 self.ensure_runtime_ready()?;
240 self.ensure_canonical_ingress_can_start()?;
241 if self
242 .state
243 .metadata
244 .run_configuration
245 .declared()
246 .is_some_and(|configuration| {
247 configuration.interaction == super::InteractionPolicy::ReadOnly
248 })
249 {
250 return Err(CanwuError::new(
251 ErrorCode::InteractionReadOnly,
252 "the run interaction policy rejects newly authored decision ingress",
253 ));
254 }
255 if request.request_id.get() == 0 {
256 return Err(CanwuError::new(
257 ErrorCode::InvalidDecision,
258 "decision request IDs must be nonzero",
259 ));
260 }
261 if request.carries_random_evidence() {
262 return Err(CanwuError::new(
263 ErrorCode::InvalidDecision,
264 "random draw evidence can only be generated by a boundary random decision resolution",
265 ));
266 }
267 let input_hash = canonical_hash(
268 "canwu.ingress.decision-request.v1",
269 &(due_at, priority, &request),
270 )?;
271 if let Some(existing) = self
272 .state
273 .evidence
274 .archived_decision_requests
275 .get(&request.request_id)
276 {
277 if existing.input_hash == input_hash {
278 return Ok(existing.receipt.clone());
279 }
280 return Err(CanwuError::new(
281 ErrorCode::IdempotencyConflict,
282 format!(
283 "decision request {} is already queued with different content",
284 request.request_id
285 ),
286 ));
287 }
288 for record in &self.state.evidence.ingress {
289 let IngressPayload::Decision { request: existing } = &record.payload else {
290 continue;
291 };
292 if existing.request_id != request.request_id {
293 continue;
294 }
295 if existing.as_ref() == &request
296 && record.due_at == due_at
297 && record.priority == priority
298 {
299 return Ok(IngressReceipt {
300 ingress_id: record.id,
301 issued_at: record.issued_at,
302 due_at: record.due_at,
303 });
304 }
305 return Err(CanwuError::new(
306 ErrorCode::IdempotencyConflict,
307 format!(
308 "decision request {} is already queued with different content",
309 request.request_id
310 ),
311 ));
312 }
313 if request.expected_revision != self.revision() {
314 return Err(CanwuError::new(
315 ErrorCode::SimulationRevisionConflict,
316 format!(
317 "decision request {} expected revision {}, current revision is {}",
318 request.request_id,
319 request.expected_revision,
320 self.revision()
321 ),
322 ));
323 }
324 if let Some(command) = &request.command {
325 if command.request_id.get() == 0 {
326 return Err(CanwuError::new(
327 ErrorCode::InvalidDecision,
328 "nested decision command request IDs must be nonzero",
329 ));
330 }
331 if command.expected_revision != request.expected_revision
332 || command.envelope.expected_time != Some(due_at)
333 {
334 return Err(CanwuError::new(
335 ErrorCode::InvalidDecision,
336 "nested decision command must use the decision request revision and due-time guards",
337 ));
338 }
339 if self.command_request_id_is_in_use(command.request_id) {
340 return Err(CanwuError::new(
341 ErrorCode::IdempotencyConflict,
342 format!(
343 "nested decision command request {} is already reserved or processed",
344 command.request_id
345 ),
346 ));
347 }
348 }
349 self.append_ingress(
350 due_at,
351 IngressClass::Decision,
352 priority,
353 IngressPayload::Decision {
354 request: Box::new(request),
355 },
356 None,
357 false,
358 )
359 }
360
361 pub(super) fn append_boundary_decision_ingress(
362 &mut self,
363 boundary_id: BoundaryId,
364 due_at: SimTime,
365 priority: i32,
366 request: DecisionIngressRequest,
367 ) -> Result<IngressReceipt, CanwuError> {
368 self.ensure_canonical_ingress_can_start()?;
369 let expected_revision = self.revision().checked_add(1).ok_or_else(|| {
370 CanwuError::new(
371 ErrorCode::IdentifierExhausted,
372 "boundary-generated decision revision is exhausted",
373 )
374 })?;
375 if request.request_id.get() == 0 || request.expected_revision != expected_revision {
376 return Err(CanwuError::new(
377 ErrorCode::InvalidDecision,
378 "boundary-generated decision requires a nonzero ID and the post-boundary revision",
379 ));
380 }
381 if self.decision_attempt(request.request_id).is_some()
382 || self
383 .state
384 .evidence
385 .archived_decision_requests
386 .contains_key(&request.request_id)
387 || self.state.evidence.ingress.iter().any(|record| {
388 matches!(
389 &record.payload,
390 IngressPayload::Decision { request: existing }
391 if existing.request_id == request.request_id
392 )
393 })
394 {
395 return Err(CanwuError::new(
396 ErrorCode::IdempotencyConflict,
397 format!(
398 "boundary-generated decision request {} is already reserved or processed",
399 request.request_id
400 ),
401 ));
402 }
403 if let Some(command) = &request.command
404 && (command.request_id.get() == 0
405 || command.expected_revision != request.expected_revision
406 || command.envelope.expected_time != Some(due_at)
407 || self.command_request_id_is_in_use(command.request_id))
408 {
409 return Err(CanwuError::new(
410 ErrorCode::InvalidDecision,
411 "boundary-generated decision command identity or guards are invalid",
412 ));
413 }
414 self.append_ingress(
415 due_at,
416 IngressClass::Decision,
417 priority,
418 IngressPayload::Decision {
419 request: Box::new(request),
420 },
421 Some(CauseRef::Boundary(boundary_id)),
422 true,
423 )
424 }
425
426 pub(super) fn enqueue_decision_archive_commit(
427 &mut self,
428 due_at: SimTime,
429 priority: i32,
430 commit: VerifiedDecisionArchiveCommit,
431 ) -> Result<IngressReceipt, CanwuError> {
432 self.ensure_runtime_ready()?;
433 self.ensure_canonical_ingress_can_start()?;
434 self.state
435 .current
436 .decisions
437 .commit_verified_decision_archive(&commit)
438 .map_err(decision_error)?;
439 for record in &self.state.evidence.ingress {
440 let IngressPayload::Maintenance { request } = &record.payload else {
441 continue;
442 };
443 let MaintenanceIngressRequest::DecisionArchive { commit: existing } = request.as_ref()
444 else {
445 continue;
446 };
447 if existing.token() == commit.token() {
448 if existing == &commit && record.due_at == due_at && record.priority == priority {
449 return Ok(IngressReceipt {
450 ingress_id: record.id,
451 issued_at: record.issued_at,
452 due_at: record.due_at,
453 });
454 }
455 return Err(CanwuError::new(
456 ErrorCode::IdempotencyConflict,
457 "decision archive token is already queued with different content",
458 ));
459 }
460 }
461 self.append_ingress(
462 due_at,
463 IngressClass::ScheduledSystem,
464 priority,
465 IngressPayload::Maintenance {
466 request: Box::new(MaintenanceIngressRequest::DecisionArchive { commit }),
467 },
468 Some(super::CauseRef::System(
469 "canwu.core.decision-archive".to_owned(),
470 )),
471 false,
472 )
473 }
474
475 pub(super) fn apply_maintenance_request(
476 &mut self,
477 request: MaintenanceIngressRequest,
478 ) -> Result<(MaintenanceChangeRecord, Vec<super::DomainRecordChange>), CanwuError> {
479 match request {
480 MaintenanceIngressRequest::DecisionArchive { commit } => {
481 let observed_source_root = self
482 .state
483 .current
484 .decisions
485 .hot_history_commitment()
486 .map_err(decision_error)?;
487 if observed_source_root != commit.source_root() {
488 return Ok((
489 MaintenanceChangeRecord {
490 kind: "decision_archive".to_owned(),
491 token: commit.token().to_owned(),
492 disposition: MaintenanceDisposition::RejectedStale,
493 source_root: observed_source_root.clone(),
494 target_root: observed_source_root.clone(),
495 rejection: Some(MaintenanceRejectionReceipt {
496 token: commit.token().to_owned(),
497 expected_source_root: commit.source_root().to_owned(),
498 observed_source_root,
499 reason:
500 "decision archive source root changed after durable admission"
501 .to_owned(),
502 }),
503 },
504 Vec::new(),
505 ));
506 }
507 self.state.current.decisions = self
508 .state
509 .current
510 .decisions
511 .commit_verified_decision_archive(&commit)
512 .map_err(decision_error)?;
513 self.invalidate_commitments(super::CommitmentDomains::DECISIONS);
514 let target_root = self
515 .state
516 .current
517 .decisions
518 .hot_history_commitment()
519 .map_err(decision_error)?;
520 Ok((
521 MaintenanceChangeRecord {
522 kind: "decision_archive".to_owned(),
523 token: commit.token().to_owned(),
524 disposition: MaintenanceDisposition::Applied,
525 source_root: observed_source_root,
526 target_root,
527 rejection: None,
528 },
529 Vec::new(),
530 ))
531 }
532 MaintenanceIngressRequest::OwnerAuthorized { commit } => {
533 super::maintenance::validate_verified_commit_authorization_structure(
534 &commit,
535 &self.plugins,
536 )?;
537 let observed_source_root = canonical_hash(
538 "canwu.owner-authorized.source-domain-root.v1",
539 self.state.current.domain_records.roots(),
540 )?;
541 if observed_source_root != commit.source_root() {
542 return Ok((
543 MaintenanceChangeRecord {
544 kind: "owner_authorized".to_owned(),
545 token: commit.token().to_owned(),
546 disposition: MaintenanceDisposition::RejectedStale,
547 source_root: observed_source_root.clone(),
548 target_root: observed_source_root.clone(),
549 rejection: Some(MaintenanceRejectionReceipt {
550 token: commit.token().to_owned(),
551 expected_source_root: commit.source_root().to_owned(),
552 observed_source_root,
553 reason:
554 "owner-authorized source root changed after durable admission"
555 .to_owned(),
556 }),
557 },
558 Vec::new(),
559 ));
560 }
561 let record_changes = self.apply_owner_authorized_maintenance(&commit)?;
562 let target_root = canonical_hash(
563 "canwu.owner-authorized.source-domain-root.v1",
564 self.state.current.domain_records.roots(),
565 )?;
566 Ok((
567 MaintenanceChangeRecord {
568 kind: "owner_authorized".to_owned(),
569 token: commit.token().to_owned(),
570 disposition: MaintenanceDisposition::Applied,
571 source_root: observed_source_root,
572 target_root,
573 rejection: None,
574 },
575 record_changes,
576 ))
577 }
578 }
579 }
580
581 pub fn drive_decision(
582 &mut self,
583 due_at: SimTime,
584 priority: i32,
585 decision_request_id: DecisionRequestId,
586 command_request_id: Option<CommandRequestId>,
587 ticket_id: DecisionTicketId,
588 policy: &dyn DecisionPolicy,
589 ) -> Result<DecisionEvaluation, CanwuError> {
590 let evaluation = self.prepare_decision_at(
591 due_at,
592 decision_request_id,
593 command_request_id,
594 ticket_id,
595 policy,
596 )?;
597 if let DecisionEvaluation::Prepared(prepared) = &evaluation {
598 self.enqueue_decision(due_at, priority, prepared.request.clone())?;
599 }
600 Ok(evaluation)
601 }
602
603 pub(super) fn apply_decision_request(
604 &mut self,
605 request: DecisionIngressRequest,
606 ) -> Result<Option<CommandOutcome>, CanwuError> {
607 let request_commitment = canonical_hash(DECISION_REQUEST_COMMITMENT_DOMAIN, &request)?;
608 let decision_request_id = request.request_id;
609 let decision_expected_revision = request.expected_revision;
610 let revision_before = self.revision();
611 if request.expected_revision != self.revision() {
612 return self.record_decision_rejection(
613 request.request_id,
614 request.expected_revision,
615 request_commitment.clone(),
616 DecisionAttemptErrorCode::SimulationRevisionConflict,
617 format!(
618 "decision request {} expected revision {}, current revision is {}",
619 request.request_id,
620 request.expected_revision,
621 self.revision()
622 ),
623 );
624 }
625 if let Some(command) = &request.command
626 && !self.command_request_id_is_unique_for_admitted_decision(command.request_id)
627 {
628 return self.record_decision_rejection(
629 request.request_id,
630 request.expected_revision,
631 request_commitment.clone(),
632 DecisionAttemptErrorCode::CommandRequestConflict,
633 format!(
634 "nested decision command request {} is not unique at admission",
635 command.request_id
636 ),
637 );
638 }
639 if let Err(error) = self.validate_decision_mutation_entities(&request.mutation) {
640 return self.record_decision_rejection(
641 request.request_id,
642 request.expected_revision,
643 request_commitment.clone(),
644 DecisionAttemptErrorCode::EntityUnavailable,
645 error.message,
646 );
647 }
648 if let Some((code, message)) = super::persons::decision_mutation_availability_error(
649 &request.mutation,
650 &self.state.current.decisions,
651 &self.state.current.person_availability,
652 ) {
653 return self.record_decision_rejection(
654 request.request_id,
655 request.expected_revision,
656 request_commitment.clone(),
657 code,
658 message,
659 );
660 }
661 let trace_claim = if matches!(request.mutation, DecisionMutation::Resolve { .. }) {
662 let (id, next_id) = claim_counter(
663 self.state.counters.next_decision_trace_id,
664 "decision trace ID",
665 )?;
666 Some((DecisionTraceId::new(id), next_id))
667 } else {
668 None
669 };
670 let mut decisions = self.state.current.decisions.clone();
671 let prepared = match decisions.apply(
672 request.mutation,
673 self.state.scheduler.now,
674 trace_claim.map(|(id, _)| id),
675 ) {
676 Ok(prepared) => prepared,
677 Err(error) => {
678 return self.record_decision_rejection(
679 request.request_id,
680 request.expected_revision,
681 request_commitment.clone(),
682 error.code.into(),
683 error.message,
684 );
685 }
686 };
687 let controller = prepared
688 .trace
689 .as_ref()
690 .and_then(|trace| decisions.controller(&trace.controller_id));
691 let decision_controller_id = controller.map(|controller| controller.id.clone());
692 match (&prepared.action, &request.command) {
693 (Some(DecisionAction::Command { command }), Some(request)) => {
694 let expected: Command = match serde_json::from_value(command.clone()) {
695 Ok(command) => command,
696 Err(error) => {
697 return self.record_decision_rejection(
698 decision_request_id,
699 decision_expected_revision,
700 request_commitment.clone(),
701 DecisionAttemptErrorCode::InvalidDecision,
702 format!("decision option contains an invalid command: {error}"),
703 );
704 }
705 };
706 if request.envelope.command != expected
707 || request.expected_revision != self.revision()
708 || prepared
709 .trace
710 .as_ref()
711 .and_then(|trace| trace.command_request_id)
712 != Some(request.request_id)
713 {
714 return self.record_decision_rejection(
715 decision_request_id,
716 decision_expected_revision,
717 request_commitment.clone(),
718 DecisionAttemptErrorCode::InvalidDecision,
719 "nested command does not match the selected decision option".to_owned(),
720 );
721 }
722 let controller = controller.ok_or_else(|| {
723 invalid_snapshot_error("decision trace does not resolve its controller binding")
724 })?;
725 if request.envelope.issuer != controller_issuer(controller)
726 || request.envelope.authority.as_ref()
727 != Some(&controller_authority(controller))
728 || request.envelope.expected_time != Some(self.state.scheduler.now)
729 {
730 return self.record_decision_rejection(
731 decision_request_id,
732 decision_expected_revision,
733 request_commitment.clone(),
734 DecisionAttemptErrorCode::InvalidDecision,
735 "nested command issuer, authority, or time guard was not derived from the decision controller".to_owned(),
736 );
737 }
738 }
739 (Some(DecisionAction::None) | None, None) => {}
740 _ => {
741 return self.record_decision_rejection(
742 decision_request_id,
743 decision_expected_revision,
744 request_commitment.clone(),
745 DecisionAttemptErrorCode::InvalidDecision,
746 "decision action and nested command disagree".to_owned(),
747 );
748 }
749 }
750 let trace_id = prepared.trace.as_ref().map(|trace| trace.id);
751 let command_request_id = request.command.as_ref().map(|request| request.request_id);
752 decisions
753 .append_attempt(DecisionAttemptRecord {
754 request_id: decision_request_id,
755 request_commitment,
756 at: self.state.scheduler.now,
757 revision_before,
758 expected_revision: decision_expected_revision,
759 outcome: DecisionAttemptOutcome::Accepted {
760 trace_id,
761 command_request_id,
762 },
763 })
764 .map_err(decision_error)?;
765 if let Some((_, next_id)) = trace_claim {
766 self.state.counters.next_decision_trace_id = next_id;
767 }
768 self.state.current.decisions = decisions;
769 self.invalidate_commitments(super::CommitmentDomains::DECISIONS);
770 let Some(command) = request.command else {
771 return Ok(None);
772 };
773 let CommandRequest {
774 request_id,
775 expected_revision,
776 envelope,
777 } = *command;
778 self.admit_command(
779 Some(request_id),
780 Some(expected_revision),
781 envelope,
782 CommandIngress::LiveRequest,
783 decision_controller_id,
784 true,
785 )
786 .map(Some)
787 }
788
789 fn record_decision_rejection(
790 &mut self,
791 request_id: DecisionRequestId,
792 expected_revision: u64,
793 request_commitment: String,
794 code: DecisionAttemptErrorCode,
795 message: String,
796 ) -> Result<Option<CommandOutcome>, CanwuError> {
797 self.state
798 .current
799 .decisions
800 .append_attempt(DecisionAttemptRecord {
801 request_id,
802 request_commitment,
803 at: self.state.scheduler.now,
804 revision_before: self.revision(),
805 expected_revision,
806 outcome: DecisionAttemptOutcome::Rejected { code, message },
807 })
808 .map_err(decision_error)?;
809 self.invalidate_commitments(super::CommitmentDomains::DECISIONS);
810 Ok(None)
811 }
812
813 pub(super) fn command_request_id_is_in_use(&self, request_id: CommandRequestId) -> bool {
814 self.state
815 .evidence
816 .archived_command_requests
817 .contains_key(&request_id)
818 || self
819 .state
820 .evidence
821 .archived_ingress_requests
822 .contains_key(&request_id)
823 || self
824 .state
825 .evidence
826 .archived_decision_command_requests
827 .contains(&request_id)
828 || self
829 .state
830 .evidence
831 .command_attempts
832 .iter()
833 .any(|attempt| attempt.request_id == Some(request_id))
834 || self
835 .state
836 .evidence
837 .ingress
838 .iter()
839 .any(|record| ingress_command_request_id(record) == Some(request_id))
840 }
841
842 fn command_request_id_is_unique_for_admitted_decision(
843 &self,
844 request_id: CommandRequestId,
845 ) -> bool {
846 !self
847 .state
848 .evidence
849 .archived_command_requests
850 .contains_key(&request_id)
851 && !self
852 .state
853 .evidence
854 .archived_ingress_requests
855 .contains_key(&request_id)
856 && !self
857 .state
858 .evidence
859 .archived_decision_command_requests
860 .contains(&request_id)
861 && !self
862 .state
863 .evidence
864 .command_attempts
865 .iter()
866 .any(|attempt| attempt.request_id == Some(request_id))
867 && self
868 .state
869 .evidence
870 .ingress
871 .iter()
872 .filter(|record| ingress_command_request_id(record) == Some(request_id))
873 .count()
874 == 1
875 }
876
877 fn validate_decision_mutation_entities(
878 &self,
879 mutation: &DecisionMutation,
880 ) -> Result<(), CanwuError> {
881 let entity_exists =
882 |entity: &EntityRef| runtime_entity_identity_exists(&self.state, entity);
883 let validate_authority = |authority: &DecisionAuthority| {
884 let valid = match authority {
885 DecisionAuthority::Actor { actor } => entity_exists(&EntityRef::Person(*actor)),
886 DecisionAuthority::Institution {
887 institution,
888 responsible_actor,
889 } => {
890 entity_exists(institution)
891 && responsible_actor
892 .is_none_or(|actor| entity_exists(&EntityRef::Person(actor)))
893 }
894 DecisionAuthority::Council { .. }
895 | DecisionAuthority::NoResponsibleActor { .. } => true,
896 };
897 valid.then_some(()).ok_or_else(|| {
898 CanwuError::new(
899 ErrorCode::InvalidDecision,
900 "decision controller authority references an unknown entity",
901 )
902 })
903 };
904 match mutation {
905 DecisionMutation::RegisterController { controller } => {
906 validate_authority(&controller.authority)?;
907 if controller
908 .command_subject
909 .as_ref()
910 .is_some_and(|entity| !entity_exists(entity))
911 {
912 return Err(CanwuError::new(
913 ErrorCode::InvalidDecision,
914 "decision controller command subject references an unknown entity",
915 ));
916 }
917 }
918 DecisionMutation::Open { ticket } if !entity_exists(&ticket.decision_maker) => {
919 return Err(CanwuError::new(
920 ErrorCode::InvalidDecision,
921 "decision maker references an unknown entity",
922 ));
923 }
924 DecisionMutation::Open { .. }
925 | DecisionMutation::ReplaceOptions { .. }
926 | DecisionMutation::Resolve { .. }
927 | DecisionMutation::Cancel { .. } => {}
928 }
929 Ok(())
930 }
931}
932
933fn ingress_command_request_id(record: &super::IngressRecord) -> Option<CommandRequestId> {
934 match &record.payload {
935 IngressPayload::Command { request } => Some(request.request_id),
936 IngressPayload::Decision { request } => {
937 request.command.as_ref().map(|request| request.request_id)
938 }
939 IngressPayload::Plugin { .. }
940 | IngressPayload::Calendar { .. }
941 | IngressPayload::Maintenance { .. }
942 | IngressPayload::PluginCancellation { .. } => None,
943 }
944}
945
946pub(super) fn controller_issuer(controller: &super::DecisionControllerBinding) -> Issuer {
947 match controller.policy.kind {
948 DecisionPolicyKind::Human => Issuer::Human(controller.id.clone()),
949 DecisionPolicyKind::Utility
950 | DecisionPolicyKind::Rule
951 | DecisionPolicyKind::Random
952 | DecisionPolicyKind::External
953 | DecisionPolicyKind::Llm => Issuer::Ai(controller.id.clone()),
954 }
955}
956
957pub(super) fn controller_authority(
958 controller: &super::DecisionControllerBinding,
959) -> CommandAuthority {
960 let decision_origin = match &controller.authority {
961 DecisionAuthority::Actor { actor } => super::DecisionOrigin::Actor { actor: *actor },
962 DecisionAuthority::Institution {
963 institution,
964 responsible_actor,
965 } => super::DecisionOrigin::Institution {
966 institution: institution.clone(),
967 responsible_actor: *responsible_actor,
968 },
969 DecisionAuthority::Council { council_id } => super::DecisionOrigin::Council {
970 council_id: council_id.clone(),
971 },
972 DecisionAuthority::NoResponsibleActor { reason } => {
973 super::DecisionOrigin::NoResponsibleActor {
974 reason: reason.clone(),
975 }
976 }
977 };
978 CommandAuthority {
979 decision_origin,
980 seat_id: controller.seat_id.clone(),
981 permission_profile_id: controller.permission_profile_id.clone(),
982 command_subject: controller.command_subject.clone(),
983 }
984}
985
986#[allow(clippy::needless_pass_by_value)]
987pub(super) fn decision_error(error: DecisionError) -> CanwuError {
988 CanwuError::new(ErrorCode::InvalidDecision, error.to_string())
989}