Skip to main content

canwu_sim/runtime/
decision.rs

1use super::{
2    CanwuError, Command, CommandAuthority, CommandEnvelope, CommandIngress, CommandOutcome,
3    CommandRequest, CommandRequestId, ControllerDecision, DecisionAction, DecisionAttemptErrorCode,
4    DecisionAttemptOutcome, DecisionAttemptRecord, DecisionAuthority, DecisionController,
5    DecisionError, DecisionMutation, DecisionPolicy, DecisionPolicyKind, DecisionRequestId,
6    DecisionTicket, DecisionTicketId, DecisionTrace, DecisionTraceId, EntityRef, ErrorCode,
7    IngressClass, IngressPayload, IngressReceipt, Issuer, SimTime, Simulation, canonical_hash,
8    claim_counter, invalid_snapshot_error, runtime_entity_identity_exists,
9};
10use serde::{Deserialize, Serialize};
11
12#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
13pub struct DecisionIngressRequest {
14    pub request_id: DecisionRequestId,
15    pub expected_revision: u64,
16    pub mutation: DecisionMutation,
17    #[serde(default, skip_serializing_if = "Option::is_none")]
18    pub command: Option<Box<CommandRequest>>,
19}
20
21impl DecisionIngressRequest {
22    #[must_use]
23    pub const fn new(
24        request_id: DecisionRequestId,
25        expected_revision: u64,
26        mutation: DecisionMutation,
27    ) -> Self {
28        Self {
29            request_id,
30            expected_revision,
31            mutation,
32            command: None,
33        }
34    }
35
36    #[must_use]
37    pub fn with_command(mut self, command: CommandRequest) -> Self {
38        self.command = Some(Box::new(command));
39        self
40    }
41}
42
43#[derive(Clone, Debug, PartialEq)]
44pub enum DecisionEvaluation {
45    Pending(super::PolicyDecision),
46    Prepared(PreparedDecisionIngress),
47}
48
49#[derive(Clone, Debug, PartialEq)]
50pub struct PreparedDecisionIngress {
51    pub request: DecisionIngressRequest,
52    pub selected_action: Option<DecisionAction>,
53}
54
55impl Simulation {
56    #[must_use]
57    pub const fn decision_state(&self) -> &super::DecisionState {
58        &self.state.current.decisions
59    }
60
61    #[must_use]
62    pub fn decision_ticket(&self, id: DecisionTicketId) -> Option<&DecisionTicket> {
63        self.state.current.decisions.ticket(id)
64    }
65
66    #[must_use]
67    pub fn decision_traces(&self) -> &[DecisionTrace] {
68        &self.state.current.decisions.traces
69    }
70
71    #[must_use]
72    pub fn decision_attempts(&self) -> &[DecisionAttemptRecord] {
73        &self.state.current.decisions.attempts
74    }
75
76    pub fn prepare_decision(
77        &self,
78        decision_request_id: DecisionRequestId,
79        command_request_id: Option<CommandRequestId>,
80        ticket_id: DecisionTicketId,
81        policy: &dyn DecisionPolicy,
82    ) -> Result<DecisionEvaluation, CanwuError> {
83        self.prepare_decision_at(
84            self.state.scheduler.now,
85            decision_request_id,
86            command_request_id,
87            ticket_id,
88            policy,
89        )
90    }
91
92    pub fn prepare_decision_at(
93        &self,
94        due_at: SimTime,
95        decision_request_id: DecisionRequestId,
96        command_request_id: Option<CommandRequestId>,
97        ticket_id: DecisionTicketId,
98        policy: &dyn DecisionPolicy,
99    ) -> Result<DecisionEvaluation, CanwuError> {
100        self.ensure_runtime_ready()?;
101        if due_at < self.state.scheduler.now {
102            return Err(CanwuError::new(
103                ErrorCode::SimulationTimeConflict,
104                "a decision cannot be prepared behind committed simulation time",
105            ));
106        }
107        let ticket = self.decision_ticket(ticket_id).ok_or_else(|| {
108            CanwuError::new(
109                ErrorCode::InvalidDecision,
110                format!("decision ticket {ticket_id} was not found"),
111            )
112        })?;
113        if ticket.deadline.is_some_and(|deadline| deadline < due_at) {
114            return Err(CanwuError::new(
115                ErrorCode::InvalidDecision,
116                format!("decision ticket {ticket_id} has expired"),
117            ));
118        }
119        let controller = self
120            .state
121            .current
122            .decisions
123            .controller(&ticket.assigned_controller)
124            .ok_or_else(|| {
125                CanwuError::new(
126                    ErrorCode::InvalidDecision,
127                    "decision ticket names an unknown controller",
128                )
129            })?;
130        match DecisionController::evaluate(ticket, controller, policy).map_err(decision_error)? {
131            ControllerDecision::Pending(decision) => Ok(DecisionEvaluation::Pending(decision)),
132            ControllerDecision::Authoritative { decision, action } => {
133                let command = match &action {
134                    Some(DecisionAction::Command { command }) => {
135                        let request_id = command_request_id.ok_or_else(|| {
136                            CanwuError::new(
137                                ErrorCode::InvalidDecision,
138                                "a selected command option requires a command request ID",
139                            )
140                        })?;
141                        let command: Command =
142                            serde_json::from_value(command.clone()).map_err(|error| {
143                                CanwuError::new(
144                                    ErrorCode::InvalidDecision,
145                                    format!("decision option contains an invalid command: {error}"),
146                                )
147                            })?;
148                        Some(CommandRequest::new(
149                            request_id,
150                            self.revision(),
151                            CommandEnvelope::new(controller_issuer(controller), command)
152                                .with_authority(controller_authority(controller))
153                                .at_time(due_at),
154                        ))
155                    }
156                    Some(DecisionAction::None) | None => {
157                        if command_request_id.is_some() {
158                            return Err(CanwuError::new(
159                                ErrorCode::InvalidDecision,
160                                "a non-command option cannot reserve a command request ID",
161                            ));
162                        }
163                        None
164                    }
165                };
166                let mutation = DecisionMutation::Resolve {
167                    ticket_id,
168                    expected_version: ticket.version,
169                    controller_id: controller.id.clone(),
170                    policy: controller.policy.clone(),
171                    decision,
172                    command_request_id,
173                };
174                let request = DecisionIngressRequest {
175                    request_id: decision_request_id,
176                    expected_revision: self.revision(),
177                    mutation,
178                    command: command.map(Box::new),
179                };
180                Ok(DecisionEvaluation::Prepared(PreparedDecisionIngress {
181                    request,
182                    selected_action: action,
183                }))
184            }
185        }
186    }
187
188    pub fn enqueue_decision(
189        &mut self,
190        due_at: SimTime,
191        priority: i32,
192        request: DecisionIngressRequest,
193    ) -> Result<IngressReceipt, CanwuError> {
194        self.ensure_runtime_ready()?;
195        self.ensure_canonical_ingress_can_start()?;
196        if self
197            .state
198            .metadata
199            .run_configuration
200            .declared()
201            .is_some_and(|configuration| {
202                configuration.interaction == super::InteractionPolicy::ReadOnly
203            })
204        {
205            return Err(CanwuError::new(
206                ErrorCode::InteractionReadOnly,
207                "the run interaction policy rejects newly authored decision ingress",
208            ));
209        }
210        if request.request_id.get() == 0 {
211            return Err(CanwuError::new(
212                ErrorCode::InvalidDecision,
213                "decision request IDs must be nonzero",
214            ));
215        }
216        let input_hash = canonical_hash(
217            "canwu.ingress.decision-request.v1",
218            &(due_at, priority, &request),
219        )?;
220        if let Some(existing) = self
221            .state
222            .evidence
223            .archived_decision_requests
224            .get(&request.request_id)
225        {
226            if existing.input_hash == input_hash {
227                return Ok(existing.receipt.clone());
228            }
229            return Err(CanwuError::new(
230                ErrorCode::IdempotencyConflict,
231                format!(
232                    "decision request {} is already queued with different content",
233                    request.request_id
234                ),
235            ));
236        }
237        for record in &self.state.evidence.ingress {
238            let IngressPayload::Decision { request: existing } = &record.payload else {
239                continue;
240            };
241            if existing.request_id != request.request_id {
242                continue;
243            }
244            if existing.as_ref() == &request
245                && record.due_at == due_at
246                && record.priority == priority
247            {
248                return Ok(IngressReceipt {
249                    ingress_id: record.id,
250                    issued_at: record.issued_at,
251                    due_at: record.due_at,
252                });
253            }
254            return Err(CanwuError::new(
255                ErrorCode::IdempotencyConflict,
256                format!(
257                    "decision request {} is already queued with different content",
258                    request.request_id
259                ),
260            ));
261        }
262        if request.expected_revision != self.revision() {
263            return Err(CanwuError::new(
264                ErrorCode::SimulationRevisionConflict,
265                format!(
266                    "decision request {} expected revision {}, current revision is {}",
267                    request.request_id,
268                    request.expected_revision,
269                    self.revision()
270                ),
271            ));
272        }
273        if let Some(command) = &request.command {
274            if command.request_id.get() == 0 {
275                return Err(CanwuError::new(
276                    ErrorCode::InvalidDecision,
277                    "nested decision command request IDs must be nonzero",
278                ));
279            }
280            if command.expected_revision != request.expected_revision
281                || command.envelope.expected_time != Some(due_at)
282            {
283                return Err(CanwuError::new(
284                    ErrorCode::InvalidDecision,
285                    "nested decision command must use the decision request revision and due-time guards",
286                ));
287            }
288            if self.command_request_id_is_in_use(command.request_id) {
289                return Err(CanwuError::new(
290                    ErrorCode::IdempotencyConflict,
291                    format!(
292                        "nested decision command request {} is already reserved or processed",
293                        command.request_id
294                    ),
295                ));
296            }
297        }
298        self.append_ingress(
299            due_at,
300            IngressClass::Decision,
301            priority,
302            IngressPayload::Decision {
303                request: Box::new(request),
304            },
305            None,
306            false,
307        )
308    }
309
310    pub fn drive_decision(
311        &mut self,
312        due_at: SimTime,
313        priority: i32,
314        decision_request_id: DecisionRequestId,
315        command_request_id: Option<CommandRequestId>,
316        ticket_id: DecisionTicketId,
317        policy: &dyn DecisionPolicy,
318    ) -> Result<DecisionEvaluation, CanwuError> {
319        let evaluation = self.prepare_decision_at(
320            due_at,
321            decision_request_id,
322            command_request_id,
323            ticket_id,
324            policy,
325        )?;
326        if let DecisionEvaluation::Prepared(prepared) = &evaluation {
327            self.enqueue_decision(due_at, priority, prepared.request.clone())?;
328        }
329        Ok(evaluation)
330    }
331
332    pub(super) fn apply_decision_request(
333        &mut self,
334        request: DecisionIngressRequest,
335    ) -> Result<Option<CommandOutcome>, CanwuError> {
336        let decision_request_id = request.request_id;
337        let decision_expected_revision = request.expected_revision;
338        let revision_before = self.revision();
339        if request.expected_revision != self.revision() {
340            return Ok(self.record_decision_rejection(
341                request.request_id,
342                request.expected_revision,
343                DecisionAttemptErrorCode::SimulationRevisionConflict,
344                format!(
345                    "decision request {} expected revision {}, current revision is {}",
346                    request.request_id,
347                    request.expected_revision,
348                    self.revision()
349                ),
350            ));
351        }
352        if let Some(command) = &request.command
353            && !self.command_request_id_is_unique_for_admitted_decision(command.request_id)
354        {
355            return Ok(self.record_decision_rejection(
356                request.request_id,
357                request.expected_revision,
358                DecisionAttemptErrorCode::CommandRequestConflict,
359                format!(
360                    "nested decision command request {} is not unique at admission",
361                    command.request_id
362                ),
363            ));
364        }
365        if let Err(error) = self.validate_decision_mutation_entities(&request.mutation) {
366            return Ok(self.record_decision_rejection(
367                request.request_id,
368                request.expected_revision,
369                DecisionAttemptErrorCode::EntityUnavailable,
370                error.message,
371            ));
372        }
373        let trace_claim = if matches!(request.mutation, DecisionMutation::Resolve { .. }) {
374            let (id, next_id) = claim_counter(
375                self.state.counters.next_decision_trace_id,
376                "decision trace ID",
377            )?;
378            Some((DecisionTraceId::new(id), next_id))
379        } else {
380            None
381        };
382        let mut decisions = self.state.current.decisions.clone();
383        let prepared = match decisions.apply(
384            request.mutation,
385            self.state.scheduler.now,
386            trace_claim.map(|(id, _)| id),
387        ) {
388            Ok(prepared) => prepared,
389            Err(error) => {
390                return Ok(self.record_decision_rejection(
391                    request.request_id,
392                    request.expected_revision,
393                    error.code.into(),
394                    error.message,
395                ));
396            }
397        };
398        let controller = prepared
399            .trace
400            .as_ref()
401            .and_then(|trace| decisions.controller(&trace.controller_id));
402        let decision_controller_id = controller.map(|controller| controller.id.clone());
403        match (&prepared.action, &request.command) {
404            (Some(DecisionAction::Command { command }), Some(request)) => {
405                let expected: Command = match serde_json::from_value(command.clone()) {
406                    Ok(command) => command,
407                    Err(error) => {
408                        return Ok(self.record_decision_rejection(
409                            decision_request_id,
410                            decision_expected_revision,
411                            DecisionAttemptErrorCode::InvalidDecision,
412                            format!("decision option contains an invalid command: {error}"),
413                        ));
414                    }
415                };
416                if request.envelope.command != expected
417                    || request.expected_revision != self.revision()
418                    || prepared
419                        .trace
420                        .as_ref()
421                        .and_then(|trace| trace.command_request_id)
422                        != Some(request.request_id)
423                {
424                    return Ok(self.record_decision_rejection(
425                        decision_request_id,
426                        decision_expected_revision,
427                        DecisionAttemptErrorCode::InvalidDecision,
428                        "nested command does not match the selected decision option".to_owned(),
429                    ));
430                }
431                let controller = controller.ok_or_else(|| {
432                    invalid_snapshot_error("decision trace does not resolve its controller binding")
433                })?;
434                if request.envelope.issuer != controller_issuer(controller)
435                    || request.envelope.authority.as_ref()
436                        != Some(&controller_authority(controller))
437                    || request.envelope.expected_time != Some(self.state.scheduler.now)
438                {
439                    return Ok(self.record_decision_rejection(
440                        decision_request_id,
441                        decision_expected_revision,
442                        DecisionAttemptErrorCode::InvalidDecision,
443                        "nested command issuer, authority, or time guard was not derived from the decision controller".to_owned(),
444                    ));
445                }
446            }
447            (Some(DecisionAction::None) | None, None) => {}
448            _ => {
449                return Ok(self.record_decision_rejection(
450                    decision_request_id,
451                    decision_expected_revision,
452                    DecisionAttemptErrorCode::InvalidDecision,
453                    "decision action and nested command disagree".to_owned(),
454                ));
455            }
456        }
457        let trace_id = prepared.trace.as_ref().map(|trace| trace.id);
458        let command_request_id = request.command.as_ref().map(|request| request.request_id);
459        decisions.attempts.push(DecisionAttemptRecord {
460            request_id: decision_request_id,
461            at: self.state.scheduler.now,
462            revision_before,
463            expected_revision: decision_expected_revision,
464            outcome: DecisionAttemptOutcome::Accepted {
465                trace_id,
466                command_request_id,
467            },
468        });
469        if let Some((_, next_id)) = trace_claim {
470            self.state.counters.next_decision_trace_id = next_id;
471        }
472        self.state.current.decisions = decisions;
473        self.invalidate_commitments(super::CommitmentDomains::DECISIONS);
474        let Some(command) = request.command else {
475            return Ok(None);
476        };
477        let CommandRequest {
478            request_id,
479            expected_revision,
480            envelope,
481        } = *command;
482        self.admit_command(
483            Some(request_id),
484            Some(expected_revision),
485            envelope,
486            CommandIngress::LiveRequest,
487            decision_controller_id,
488            true,
489        )
490        .map(Some)
491    }
492
493    fn record_decision_rejection(
494        &mut self,
495        request_id: DecisionRequestId,
496        expected_revision: u64,
497        code: DecisionAttemptErrorCode,
498        message: String,
499    ) -> Option<CommandOutcome> {
500        self.state
501            .current
502            .decisions
503            .attempts
504            .push(DecisionAttemptRecord {
505                request_id,
506                at: self.state.scheduler.now,
507                revision_before: self.revision(),
508                expected_revision,
509                outcome: DecisionAttemptOutcome::Rejected { code, message },
510            });
511        self.invalidate_commitments(super::CommitmentDomains::DECISIONS);
512        None
513    }
514
515    pub(super) fn command_request_id_is_in_use(&self, request_id: CommandRequestId) -> bool {
516        self.state
517            .evidence
518            .archived_command_requests
519            .contains_key(&request_id)
520            || self
521                .state
522                .evidence
523                .archived_ingress_requests
524                .contains_key(&request_id)
525            || self
526                .state
527                .evidence
528                .archived_decision_command_requests
529                .contains(&request_id)
530            || self
531                .state
532                .evidence
533                .command_attempts
534                .iter()
535                .any(|attempt| attempt.request_id == Some(request_id))
536            || self
537                .state
538                .evidence
539                .ingress
540                .iter()
541                .any(|record| ingress_command_request_id(record) == Some(request_id))
542    }
543
544    fn command_request_id_is_unique_for_admitted_decision(
545        &self,
546        request_id: CommandRequestId,
547    ) -> bool {
548        !self
549            .state
550            .evidence
551            .archived_command_requests
552            .contains_key(&request_id)
553            && !self
554                .state
555                .evidence
556                .archived_ingress_requests
557                .contains_key(&request_id)
558            && !self
559                .state
560                .evidence
561                .archived_decision_command_requests
562                .contains(&request_id)
563            && !self
564                .state
565                .evidence
566                .command_attempts
567                .iter()
568                .any(|attempt| attempt.request_id == Some(request_id))
569            && self
570                .state
571                .evidence
572                .ingress
573                .iter()
574                .filter(|record| ingress_command_request_id(record) == Some(request_id))
575                .count()
576                == 1
577    }
578
579    fn validate_decision_mutation_entities(
580        &self,
581        mutation: &DecisionMutation,
582    ) -> Result<(), CanwuError> {
583        let entity_exists =
584            |entity: &EntityRef| runtime_entity_identity_exists(&self.state, entity);
585        let validate_authority = |authority: &DecisionAuthority| {
586            let valid = match authority {
587                DecisionAuthority::Actor { actor } => entity_exists(&EntityRef::Person(*actor)),
588                DecisionAuthority::Institution {
589                    institution,
590                    responsible_actor,
591                } => {
592                    entity_exists(institution)
593                        && responsible_actor
594                            .is_none_or(|actor| entity_exists(&EntityRef::Person(actor)))
595                }
596                DecisionAuthority::Council { .. }
597                | DecisionAuthority::NoResponsibleActor { .. } => true,
598            };
599            valid.then_some(()).ok_or_else(|| {
600                CanwuError::new(
601                    ErrorCode::InvalidDecision,
602                    "decision controller authority references an unknown entity",
603                )
604            })
605        };
606        match mutation {
607            DecisionMutation::RegisterController { controller } => {
608                validate_authority(&controller.authority)?;
609                if controller
610                    .command_subject
611                    .as_ref()
612                    .is_some_and(|entity| !entity_exists(entity))
613                {
614                    return Err(CanwuError::new(
615                        ErrorCode::InvalidDecision,
616                        "decision controller command subject references an unknown entity",
617                    ));
618                }
619            }
620            DecisionMutation::Open { ticket } if !entity_exists(&ticket.decision_maker) => {
621                return Err(CanwuError::new(
622                    ErrorCode::InvalidDecision,
623                    "decision maker references an unknown entity",
624                ));
625            }
626            DecisionMutation::Open { .. }
627            | DecisionMutation::ReplaceOptions { .. }
628            | DecisionMutation::Resolve { .. }
629            | DecisionMutation::Cancel { .. } => {}
630        }
631        Ok(())
632    }
633}
634
635fn ingress_command_request_id(record: &super::IngressRecord) -> Option<CommandRequestId> {
636    match &record.payload {
637        IngressPayload::Command { request } => Some(request.request_id),
638        IngressPayload::Decision { request } => {
639            request.command.as_ref().map(|request| request.request_id)
640        }
641        IngressPayload::Plugin { .. } | IngressPayload::Calendar { .. } => None,
642    }
643}
644
645pub(super) fn controller_issuer(controller: &super::DecisionControllerBinding) -> Issuer {
646    match controller.policy.kind {
647        DecisionPolicyKind::Human => Issuer::Human(controller.id.clone()),
648        DecisionPolicyKind::Utility
649        | DecisionPolicyKind::Rule
650        | DecisionPolicyKind::External
651        | DecisionPolicyKind::Llm => Issuer::Ai(controller.id.clone()),
652    }
653}
654
655pub(super) fn controller_authority(
656    controller: &super::DecisionControllerBinding,
657) -> CommandAuthority {
658    let decision_origin = match &controller.authority {
659        DecisionAuthority::Actor { actor } => super::DecisionOrigin::Actor { actor: *actor },
660        DecisionAuthority::Institution {
661            institution,
662            responsible_actor,
663        } => super::DecisionOrigin::Institution {
664            institution: institution.clone(),
665            responsible_actor: *responsible_actor,
666        },
667        DecisionAuthority::Council { council_id } => super::DecisionOrigin::Council {
668            council_id: council_id.clone(),
669        },
670        DecisionAuthority::NoResponsibleActor { reason } => {
671            super::DecisionOrigin::NoResponsibleActor {
672                reason: reason.clone(),
673            }
674        }
675    };
676    CommandAuthority {
677        decision_origin,
678        seat_id: controller.seat_id.clone(),
679        permission_profile_id: controller.permission_profile_id.clone(),
680        command_subject: controller.command_subject.clone(),
681    }
682}
683
684#[allow(clippy::needless_pass_by_value)]
685pub(super) fn decision_error(error: DecisionError) -> CanwuError {
686    CanwuError::new(ErrorCode::InvalidDecision, error.to_string())
687}