Skip to main content

canwu_sim/runtime/
decision.rs

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}