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