Skip to main content

canwu_decision/
engine.rs

1use crate::model::{canonicalize_options, require_text};
2use crate::{
3    DecisionAction, DecisionAttemptOutcome, DecisionAttemptRecord, DecisionControllerBinding,
4    DecisionError, DecisionErrorCode, DecisionMutation, DecisionOutcome, DecisionPolicy,
5    DecisionPolicyIdentity, DecisionTicket, DecisionTicketState, DecisionTrace, PolicyDecision,
6};
7use canwu_core::{CommandRequestId, DecisionTicketId, DecisionTraceId};
8use canwu_time::SimTime;
9use serde::{Deserialize, Serialize};
10use std::collections::BTreeMap;
11
12#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
13pub struct DecisionState {
14    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
15    pub controllers: BTreeMap<String, DecisionControllerBinding>,
16    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
17    pub tickets: BTreeMap<DecisionTicketId, DecisionTicket>,
18    #[serde(default, skip_serializing_if = "Vec::is_empty")]
19    pub traces: Vec<DecisionTrace>,
20    #[serde(default, skip_serializing_if = "Vec::is_empty")]
21    pub attempts: Vec<DecisionAttemptRecord>,
22}
23
24impl DecisionState {
25    #[must_use]
26    pub fn is_empty(&self) -> bool {
27        self.controllers.is_empty()
28            && self.tickets.is_empty()
29            && self.traces.is_empty()
30            && self.attempts.is_empty()
31    }
32
33    #[must_use]
34    pub fn controller(&self, id: &str) -> Option<&DecisionControllerBinding> {
35        self.controllers.get(id)
36    }
37
38    #[must_use]
39    pub fn ticket(&self, id: DecisionTicketId) -> Option<&DecisionTicket> {
40        self.tickets.get(&id)
41    }
42
43    pub fn open_tickets(&self) -> impl Iterator<Item = &DecisionTicket> {
44        self.tickets.values().filter(|ticket| ticket.is_open())
45    }
46
47    pub fn validate(&self) -> Result<(), DecisionError> {
48        for (id, controller) in &self.controllers {
49            if id != &controller.id {
50                return Err(DecisionError::new(
51                    DecisionErrorCode::InvalidController,
52                    "controller map key does not match its persisted identity",
53                ));
54            }
55            controller.validate()?;
56        }
57        for (id, ticket) in &self.tickets {
58            if id != &ticket.id || !self.controllers.contains_key(&ticket.assigned_controller) {
59                return Err(DecisionError::new(
60                    DecisionErrorCode::InvalidDecision,
61                    "ticket identity or assigned controller is invalid",
62                ));
63            }
64            ticket.validate()?;
65        }
66        let mut expected_trace_id = 1_u64;
67        for trace in &self.traces {
68            if trace.id.get() != expected_trace_id {
69                return Err(DecisionError::new(
70                    DecisionErrorCode::InvalidDecision,
71                    "decision traces must use contiguous IDs in journal order",
72                ));
73            }
74            expected_trace_id = expected_trace_id.checked_add(1).ok_or_else(|| {
75                DecisionError::new(
76                    DecisionErrorCode::InvalidDecision,
77                    "decision trace ID range is exhausted",
78                )
79            })?;
80            let ticket = self.tickets.get(&trace.ticket_id).ok_or_else(|| {
81                DecisionError::new(
82                    DecisionErrorCode::InvalidDecision,
83                    "decision trace references an unknown ticket",
84                )
85            })?;
86            if trace.ticket_version == 0
87                || trace.ticket_version > ticket.version
88                || !self.controllers.contains_key(&trace.controller_id)
89            {
90                return Err(DecisionError::new(
91                    DecisionErrorCode::InvalidDecision,
92                    "decision trace version or controller is invalid",
93                ));
94            }
95        }
96        for ticket in self.tickets.values() {
97            if let DecisionTicketState::Resolved { trace_id, .. } = ticket.state {
98                let trace_index =
99                    usize::try_from(trace_id.get().saturating_sub(1)).map_err(|_| {
100                        DecisionError::new(
101                            DecisionErrorCode::InvalidDecision,
102                            "decision trace ID exceeds platform range",
103                        )
104                    })?;
105                if self
106                    .traces
107                    .get(trace_index)
108                    .is_none_or(|trace| trace.id != trace_id || trace.ticket_id != ticket.id)
109                {
110                    return Err(DecisionError::new(
111                        DecisionErrorCode::InvalidDecision,
112                        "resolved ticket does not reference its persisted trace",
113                    ));
114                }
115            }
116        }
117        let mut request_ids = std::collections::BTreeSet::new();
118        for attempt in &self.attempts {
119            if attempt.request_id.get() == 0 || !request_ids.insert(attempt.request_id) {
120                return Err(DecisionError::new(
121                    DecisionErrorCode::InvalidDecision,
122                    "decision attempts must use unique nonzero request IDs",
123                ));
124            }
125            match &attempt.outcome {
126                DecisionAttemptOutcome::Accepted {
127                    trace_id,
128                    command_request_id,
129                } => {
130                    if command_request_id.is_some() && trace_id.is_none() {
131                        return Err(DecisionError::new(
132                            DecisionErrorCode::InvalidDecision,
133                            "accepted decision commands require a decision trace",
134                        ));
135                    }
136                }
137                DecisionAttemptOutcome::Rejected { message, .. } => {
138                    require_text(message, "decision rejection message")?;
139                }
140            }
141        }
142        Ok(())
143    }
144
145    pub fn apply(
146        &mut self,
147        mutation: DecisionMutation,
148        at: SimTime,
149        trace_id: Option<DecisionTraceId>,
150    ) -> Result<PreparedDecision, DecisionError> {
151        let prepared = match mutation {
152            DecisionMutation::RegisterController { controller } => {
153                controller.validate()?;
154                if self.controllers.contains_key(&controller.id) {
155                    return Err(DecisionError::new(
156                        DecisionErrorCode::DuplicateController,
157                        format!(
158                            "decision controller {} is already registered",
159                            controller.id
160                        ),
161                    ));
162                }
163                self.controllers.insert(controller.id.clone(), controller);
164                PreparedDecision::default()
165            }
166            DecisionMutation::Open { mut ticket } => {
167                ticket.validate()?;
168                if self.tickets.contains_key(&ticket.id) {
169                    return Err(DecisionError::new(
170                        DecisionErrorCode::DuplicateTicket,
171                        format!("decision ticket {} is already present", ticket.id),
172                    ));
173                }
174                if !self.controllers.contains_key(&ticket.assigned_controller) {
175                    return Err(DecisionError::new(
176                        DecisionErrorCode::InvalidController,
177                        format!(
178                            "decision ticket {} names unknown controller {}",
179                            ticket.id, ticket.assigned_controller
180                        ),
181                    ));
182                }
183                if ticket.deadline.is_some_and(|deadline| deadline < at) {
184                    return Err(DecisionError::new(
185                        DecisionErrorCode::InvalidDecision,
186                        "decision deadline precedes its admission time",
187                    ));
188                }
189                let persisted = DecisionTicket {
190                    id: ticket.id,
191                    definition: ticket.definition,
192                    decision_maker: ticket.decision_maker,
193                    assigned_controller: ticket.assigned_controller,
194                    summary: ticket.summary,
195                    context: ticket.context,
196                    options: std::mem::take(&mut ticket.options),
197                    opened_at: at,
198                    updated_at: at,
199                    deadline: ticket.deadline,
200                    version: 1,
201                    state: DecisionTicketState::Open,
202                };
203                self.tickets.insert(persisted.id, persisted);
204                PreparedDecision::default()
205            }
206            DecisionMutation::ReplaceOptions {
207                ticket_id,
208                expected_version,
209                context,
210                mut options,
211            } => {
212                context.validate()?;
213                canonicalize_options(&mut options)?;
214                let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
215                ticket.context = context;
216                ticket.options = options;
217                ticket.updated_at = at;
218                ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
219                    DecisionError::new(
220                        DecisionErrorCode::InvalidDecision,
221                        "decision ticket version is exhausted",
222                    )
223                })?;
224                PreparedDecision::default()
225            }
226            DecisionMutation::Resolve {
227                ticket_id,
228                expected_version,
229                controller_id,
230                policy,
231                decision,
232                command_request_id,
233            } => self.resolve(
234                ticket_id,
235                expected_version,
236                &controller_id,
237                policy,
238                decision,
239                command_request_id,
240                at,
241                trace_id.ok_or_else(|| {
242                    DecisionError::new(
243                        DecisionErrorCode::InvalidDecision,
244                        "decision resolution requires a claimed trace ID",
245                    )
246                })?,
247            )?,
248            DecisionMutation::Cancel {
249                ticket_id,
250                expected_version,
251                reason,
252            } => {
253                require_text(&reason, "decision cancellation reason")?;
254                let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
255                ticket.updated_at = at;
256                ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
257                    DecisionError::new(
258                        DecisionErrorCode::InvalidDecision,
259                        "decision ticket version is exhausted",
260                    )
261                })?;
262                ticket.state = DecisionTicketState::Cancelled { reason };
263                PreparedDecision::default()
264            }
265        };
266        self.advance_time(at)?;
267        self.validate()?;
268        Ok(prepared)
269    }
270
271    fn open_ticket_mut(
272        &mut self,
273        ticket_id: DecisionTicketId,
274        expected_version: u64,
275        at: SimTime,
276    ) -> Result<&mut DecisionTicket, DecisionError> {
277        let ticket = self.tickets.get_mut(&ticket_id).ok_or_else(|| {
278            DecisionError::new(
279                DecisionErrorCode::TicketNotFound,
280                format!("decision ticket {ticket_id} was not found"),
281            )
282        })?;
283        if !ticket.is_open() || ticket.deadline.is_some_and(|deadline| deadline < at) {
284            return Err(DecisionError::new(
285                DecisionErrorCode::ClosedTicket,
286                format!("decision ticket {ticket_id} is not open"),
287            ));
288        }
289        if ticket.version != expected_version {
290            return Err(DecisionError::new(
291                DecisionErrorCode::VersionConflict,
292                format!(
293                    "decision ticket {ticket_id} is at version {}, expected {expected_version}",
294                    ticket.version
295                ),
296            ));
297        }
298        Ok(ticket)
299    }
300
301    #[allow(clippy::too_many_arguments)]
302    fn resolve(
303        &mut self,
304        ticket_id: DecisionTicketId,
305        expected_version: u64,
306        controller_id: &str,
307        policy: DecisionPolicyIdentity,
308        decision: PolicyDecision,
309        command_request_id: Option<CommandRequestId>,
310        at: SimTime,
311        trace_id: DecisionTraceId,
312    ) -> Result<PreparedDecision, DecisionError> {
313        let controller = self.controllers.get(controller_id).ok_or_else(|| {
314            DecisionError::new(
315                DecisionErrorCode::InvalidController,
316                format!("decision controller {controller_id} was not found"),
317            )
318        })?;
319        if controller.policy != policy {
320            return Err(DecisionError::new(
321                DecisionErrorCode::PolicyMismatch,
322                "decision resolution policy does not match the persisted controller binding",
323            ));
324        }
325        let ticket = self.open_ticket_mut(ticket_id, expected_version, at)?;
326        if ticket.assigned_controller != controller_id {
327            return Err(DecisionError::new(
328                DecisionErrorCode::InvalidController,
329                "decision resolution came from a controller not assigned to the ticket",
330            ));
331        }
332        decision.validate(ticket)?;
333        let action = match &decision.outcome {
334            DecisionOutcome::Selected { option_id } => {
335                ticket.option(option_id).map(|option| option.action.clone())
336            }
337            DecisionOutcome::Deferred { .. } => None,
338            DecisionOutcome::Pending { .. } => {
339                return Err(DecisionError::new(
340                    DecisionErrorCode::InvalidDecision,
341                    "pending policy outcomes are not authoritative decision mutations",
342                ));
343            }
344        };
345        if matches!(action, Some(DecisionAction::Command { .. })) != command_request_id.is_some() {
346            return Err(DecisionError::new(
347                DecisionErrorCode::InvalidDecision,
348                "command actions require exactly one command request ID",
349            ));
350        }
351        let trace = DecisionTrace {
352            id: trace_id,
353            ticket_id,
354            ticket_version: ticket.version,
355            controller_id: controller_id.to_owned(),
356            policy,
357            decided_at: at,
358            outcome: decision.outcome.clone(),
359            summary: decision.summary,
360            evaluations: decision.evaluations,
361            external: decision.external,
362            command_request_id,
363        };
364        ticket.updated_at = at;
365        ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
366            DecisionError::new(
367                DecisionErrorCode::InvalidDecision,
368                "decision ticket version is exhausted",
369            )
370        })?;
371        if let DecisionOutcome::Selected { option_id } = &trace.outcome {
372            ticket.state = DecisionTicketState::Resolved {
373                option_id: option_id.clone(),
374                trace_id,
375            };
376        }
377        self.traces.push(trace.clone());
378        Ok(PreparedDecision {
379            trace: Some(trace),
380            action,
381        })
382    }
383
384    pub fn advance_time(&mut self, at: SimTime) -> Result<(), DecisionError> {
385        for ticket in self.tickets.values_mut() {
386            if ticket.is_open() && ticket.deadline.is_some_and(|deadline| deadline < at) {
387                ticket.updated_at = at;
388                ticket.version = ticket.version.checked_add(1).ok_or_else(|| {
389                    DecisionError::new(
390                        DecisionErrorCode::InvalidDecision,
391                        "decision ticket version is exhausted",
392                    )
393                })?;
394                ticket.state = DecisionTicketState::Expired;
395            }
396        }
397        Ok(())
398    }
399}
400
401#[derive(Clone, Debug, Default, Eq, PartialEq)]
402pub struct PreparedDecision {
403    pub trace: Option<DecisionTrace>,
404    pub action: Option<DecisionAction>,
405}
406
407#[derive(Clone, Debug, Eq, PartialEq)]
408pub enum ControllerDecision {
409    Authoritative {
410        decision: PolicyDecision,
411        action: Option<DecisionAction>,
412    },
413    Pending(PolicyDecision),
414}
415
416pub struct DecisionController;
417
418impl DecisionController {
419    pub fn evaluate(
420        ticket: &DecisionTicket,
421        controller: &DecisionControllerBinding,
422        policy: &dyn DecisionPolicy,
423    ) -> Result<ControllerDecision, DecisionError> {
424        if !ticket.is_open() {
425            return Err(DecisionError::new(
426                DecisionErrorCode::ClosedTicket,
427                "only open tickets can be evaluated",
428            ));
429        }
430        if ticket.assigned_controller != controller.id || policy.identity() != controller.policy {
431            return Err(DecisionError::new(
432                DecisionErrorCode::PolicyMismatch,
433                "runtime policy identity does not match the ticket controller binding",
434            ));
435        }
436        let decision = policy.decide(ticket)?;
437        decision.validate(ticket)?;
438        if matches!(decision.outcome, DecisionOutcome::Pending { .. }) {
439            return Ok(ControllerDecision::Pending(decision));
440        }
441        let action = match &decision.outcome {
442            DecisionOutcome::Selected { option_id } => {
443                ticket.option(option_id).map(|option| option.action.clone())
444            }
445            DecisionOutcome::Deferred { .. } | DecisionOutcome::Pending { .. } => None,
446        };
447        Ok(ControllerDecision::Authoritative { action, decision })
448    }
449}