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}