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}