1use crate::model::require_text;
2use crate::{
3 DecisionError, DecisionErrorCode, DecisionExternalEvidence, DecisionFactorContribution,
4 DecisionOption, DecisionOptionEvaluation, DecisionOptionWeight, DecisionOutcome,
5 DecisionPolicyIdentity, DecisionPolicyKind, DecisionStage, DecisionTicket, PolicyDecision,
6};
7use serde::{Deserialize, Serialize};
8use std::cmp::Ordering;
9use std::collections::BTreeMap;
10
11pub trait DecisionPolicy {
12 fn identity(&self) -> &DecisionPolicyIdentity;
13 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError>;
14}
15
16pub trait UtilityEvaluator {
17 fn evaluate(
18 &self,
19 ticket: &DecisionTicket,
20 option: &DecisionOption,
21 ) -> Result<DecisionOptionEvaluation, DecisionError>;
22}
23
24pub trait UtilityPolicy: DecisionPolicy + UtilityEvaluator {}
25
26impl<T: DecisionPolicy + UtilityEvaluator> UtilityPolicy for T {}
27
28#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
29pub struct UtilityProfile {
30 pub weights: BTreeMap<String, i64>,
31}
32
33#[derive(Clone, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
34pub struct WeightedUtilityEvaluator {
35 pub profile: UtilityProfile,
36}
37
38impl WeightedUtilityEvaluator {
39 #[must_use]
40 pub const fn new(profile: UtilityProfile) -> Self {
41 Self { profile }
42 }
43}
44
45impl UtilityEvaluator for WeightedUtilityEvaluator {
46 fn evaluate(
47 &self,
48 _ticket: &DecisionTicket,
49 option: &DecisionOption,
50 ) -> Result<DecisionOptionEvaluation, DecisionError> {
51 if !option.is_available() {
52 return Ok(DecisionOptionEvaluation {
53 option_id: option.id.clone(),
54 available: false,
55 score: None,
56 factors: Vec::new(),
57 blockers: option.blockers.clone(),
58 });
59 }
60 let mut score = 0_i64;
61 let mut factors = Vec::new();
62 for (factor, value) in &option.utility_inputs {
63 let weight = self
64 .profile
65 .weights
66 .get(factor)
67 .copied()
68 .unwrap_or_default();
69 let contribution = value.checked_mul(weight).ok_or_else(|| {
70 DecisionError::new(
71 DecisionErrorCode::InvalidDecision,
72 format!("utility contribution for factor {factor} exceeds the i64 range"),
73 )
74 })?;
75 score = score.checked_add(contribution).ok_or_else(|| {
76 DecisionError::new(
77 DecisionErrorCode::InvalidDecision,
78 "utility score exceeds the i64 range",
79 )
80 })?;
81 factors.push(DecisionFactorContribution {
82 factor: factor.clone(),
83 value: *value,
84 weight,
85 contribution,
86 });
87 }
88 Ok(DecisionOptionEvaluation {
89 option_id: option.id.clone(),
90 available: true,
91 score: Some(score),
92 factors,
93 blockers: Vec::new(),
94 })
95 }
96}
97
98#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
99pub struct WeightedUtilityPolicy {
100 pub identity: DecisionPolicyIdentity,
101 pub evaluator: WeightedUtilityEvaluator,
102}
103
104impl WeightedUtilityPolicy {
105 #[must_use]
106 pub fn new(id: impl Into<String>, version: impl Into<String>, profile: UtilityProfile) -> Self {
107 Self {
108 identity: DecisionPolicyIdentity::new(DecisionPolicyKind::Utility, id, version),
109 evaluator: WeightedUtilityEvaluator::new(profile),
110 }
111 }
112}
113
114impl UtilityEvaluator for WeightedUtilityPolicy {
115 fn evaluate(
116 &self,
117 ticket: &DecisionTicket,
118 option: &DecisionOption,
119 ) -> Result<DecisionOptionEvaluation, DecisionError> {
120 self.evaluator.evaluate(ticket, option)
121 }
122}
123
124impl DecisionPolicy for WeightedUtilityPolicy {
125 fn identity(&self) -> &DecisionPolicyIdentity {
126 &self.identity
127 }
128
129 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
130 let mut evaluations = ticket
131 .options
132 .iter()
133 .map(|option| self.evaluate(ticket, option))
134 .collect::<Result<Vec<_>, _>>()?;
135 evaluations.sort_by(|left, right| left.option_id.cmp(&right.option_id));
136 let selected = evaluations
137 .iter()
138 .filter_map(|evaluation| evaluation.score.map(|score| (score, &evaluation.option_id)))
139 .max_by(|left, right| left.0.cmp(&right.0).then_with(|| right.1.cmp(left.1)))
140 .map(|(_, option_id)| option_id.clone());
141 let Some(option_id) = selected else {
142 return Ok(PolicyDecision {
143 outcome: DecisionOutcome::Deferred {
144 reason: "no available option".to_owned(),
145 },
146 summary: "utility policy deferred because every option was blocked".to_owned(),
147 evaluations,
148 external: None,
149 random: None,
150 stage: None,
151 fired_guards: Vec::new(),
152 });
153 };
154 Ok(PolicyDecision {
155 outcome: DecisionOutcome::Selected {
156 option_id: option_id.clone(),
157 },
158 summary: format!("utility policy selected {option_id}"),
159 evaluations,
160 external: None,
161 random: None,
162 stage: None,
163 fired_guards: Vec::new(),
164 })
165 }
166}
167
168#[derive(Clone, Debug, Eq, PartialEq)]
169pub enum RuleChoice {
170 Select(String),
172 Defer(String),
174 Exclude {
177 option_id: String,
178 reason: String,
179 },
180 NoMatch,
181}
182
183pub trait DecisionRule {
184 fn id(&self) -> &str;
185 fn evaluate(&self, ticket: &DecisionTicket) -> Result<RuleChoice, DecisionError>;
186}
187
188pub trait RulePolicy: DecisionPolicy {
189 fn rules(&self) -> &[Box<dyn DecisionRule>];
190}
191
192pub struct OrderedRulePolicy {
193 identity: DecisionPolicyIdentity,
194 rules: Vec<Box<dyn DecisionRule>>,
195}
196
197impl OrderedRulePolicy {
198 #[must_use]
199 pub fn new(
200 id: impl Into<String>,
201 version: impl Into<String>,
202 rules: Vec<Box<dyn DecisionRule>>,
203 ) -> Self {
204 Self {
205 identity: DecisionPolicyIdentity::new(DecisionPolicyKind::Rule, id, version),
206 rules,
207 }
208 }
209}
210
211impl RulePolicy for OrderedRulePolicy {
212 fn rules(&self) -> &[Box<dyn DecisionRule>] {
213 &self.rules
214 }
215}
216
217impl DecisionPolicy for OrderedRulePolicy {
218 fn identity(&self) -> &DecisionPolicyIdentity {
219 &self.identity
220 }
221
222 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
223 let run = run_ordered_rules(&self.rules, ticket)?;
224 let evaluations = run.exclusion_evaluations();
225 let (outcome, summary) = match run.terminal {
226 Some(RuleTerminal::Select { rule, option_id }) => (
227 DecisionOutcome::Selected { option_id },
228 format!("rule {rule} selected an option"),
229 ),
230 Some(RuleTerminal::Defer { rule, reason }) => (
231 DecisionOutcome::Deferred {
232 reason: reason.clone(),
233 },
234 format!("rule {rule} deferred: {reason}"),
235 ),
236 None => (
237 DecisionOutcome::Deferred {
238 reason: "no rule matched".to_owned(),
239 },
240 "ordered rule policy exhausted its rules".to_owned(),
241 ),
242 };
243 Ok(PolicyDecision {
244 outcome,
245 summary,
246 evaluations,
247 external: None,
248 random: None,
249 stage: None,
250 fired_guards: Vec::new(),
251 })
252 }
253}
254
255enum RuleTerminal {
256 Select { rule: String, option_id: String },
257 Defer { rule: String, reason: String },
258}
259
260#[derive(Default)]
261struct OrderedRuleRun {
262 fired: Vec<String>,
263 excluded: BTreeMap<String, String>,
265 terminal: Option<RuleTerminal>,
266}
267
268impl OrderedRuleRun {
269 fn exclusion_evaluations(&self) -> Vec<DecisionOptionEvaluation> {
270 self.excluded
271 .iter()
272 .map(|(option_id, blocker)| excluded_evaluation(option_id, blocker))
273 .collect()
274 }
275}
276
277fn excluded_evaluation(option_id: &str, blocker: &str) -> DecisionOptionEvaluation {
278 DecisionOptionEvaluation {
279 option_id: option_id.to_owned(),
280 available: false,
281 score: None,
282 factors: Vec::new(),
283 blockers: vec![blocker.to_owned()],
284 }
285}
286
287fn run_ordered_rules(
291 rules: &[Box<dyn DecisionRule>],
292 ticket: &DecisionTicket,
293) -> Result<OrderedRuleRun, DecisionError> {
294 let mut run = OrderedRuleRun::default();
295 for rule in rules {
296 match rule.evaluate(ticket)? {
297 RuleChoice::NoMatch => {}
298 RuleChoice::Select(option_id) => {
299 if run.excluded.contains_key(&option_id) {
300 return Err(DecisionError::new(
301 DecisionErrorCode::InvalidDecision,
302 format!(
303 "rule {} selected option {option_id} after an earlier rule excluded it",
304 rule.id()
305 ),
306 ));
307 }
308 run.fired.push(rule.id().to_owned());
309 run.terminal = Some(RuleTerminal::Select {
310 rule: rule.id().to_owned(),
311 option_id,
312 });
313 break;
314 }
315 RuleChoice::Defer(reason) => {
316 run.fired.push(rule.id().to_owned());
317 run.terminal = Some(RuleTerminal::Defer {
318 rule: rule.id().to_owned(),
319 reason,
320 });
321 break;
322 }
323 RuleChoice::Exclude { option_id, reason } => {
324 if ticket.option(&option_id).is_none() {
325 return Err(DecisionError::new(
326 DecisionErrorCode::InvalidOption,
327 format!("rule {} excluded unknown option {option_id}", rule.id()),
328 ));
329 }
330 require_text(&reason, "rule exclusion reason")?;
331 run.fired.push(rule.id().to_owned());
332 run.excluded
333 .entry(option_id)
334 .or_insert_with(|| format!("excluded by {}: {reason}", rule.id()));
335 }
336 }
337 }
338 Ok(run)
339}
340
341pub struct GuardedUtilityPolicy {
374 identity: DecisionPolicyIdentity,
375 guards: OrderedRulePolicy,
376 utility: WeightedUtilityEvaluator,
377 near_equivalence_margin: u64,
378 random_tie_break: bool,
379}
380
381impl GuardedUtilityPolicy {
382 #[must_use]
383 pub fn new(
384 id: impl Into<String>,
385 version: impl Into<String>,
386 guards: OrderedRulePolicy,
387 utility: WeightedUtilityEvaluator,
388 near_equivalence_margin: u64,
389 random_tie_break: bool,
390 ) -> Self {
391 let mut identity = DecisionPolicyIdentity::new(DecisionPolicyKind::Utility, id, version);
392 identity.semantic_hash = Some(guarded_utility_semantic_hash(
393 &guards,
394 &utility,
395 near_equivalence_margin,
396 random_tie_break,
397 ));
398 Self {
399 identity,
400 guards,
401 utility,
402 near_equivalence_margin,
403 random_tie_break,
404 }
405 }
406
407 #[must_use]
408 pub const fn guards(&self) -> &OrderedRulePolicy {
409 &self.guards
410 }
411
412 #[must_use]
413 pub const fn utility(&self) -> &WeightedUtilityEvaluator {
414 &self.utility
415 }
416
417 #[must_use]
418 pub const fn near_equivalence_margin(&self) -> u64 {
419 self.near_equivalence_margin
420 }
421
422 #[must_use]
423 pub const fn random_tie_break(&self) -> bool {
424 self.random_tie_break
425 }
426}
427
428fn guarded_utility_semantic_hash(
429 guards: &OrderedRulePolicy,
430 utility: &WeightedUtilityEvaluator,
431 near_equivalence_margin: u64,
432 random_tie_break: bool,
433) -> String {
434 fn text(hasher: &mut blake3::Hasher, value: &str) {
435 hasher.update(&(value.len() as u64).to_be_bytes());
436 hasher.update(value.as_bytes());
437 }
438 let mut hasher = blake3::Hasher::new();
439 text(&mut hasher, "canwu.decision.guarded-utility-policy.v1");
440 let guard_identity = guards.identity();
441 text(&mut hasher, policy_kind_name(guard_identity.kind));
442 text(&mut hasher, &guard_identity.id);
443 text(&mut hasher, &guard_identity.version);
444 hasher.update(&(guards.rules().len() as u64).to_be_bytes());
445 for rule in guards.rules() {
446 text(&mut hasher, rule.id());
447 }
448 hasher.update(&(utility.profile.weights.len() as u64).to_be_bytes());
449 for (factor, weight) in &utility.profile.weights {
450 text(&mut hasher, factor);
451 hasher.update(&weight.to_be_bytes());
452 }
453 hasher.update(&near_equivalence_margin.to_be_bytes());
454 hasher.update(&[u8::from(random_tie_break)]);
455 hasher.finalize().to_hex().to_string()
456}
457
458const fn policy_kind_name(kind: DecisionPolicyKind) -> &'static str {
459 match kind {
460 DecisionPolicyKind::Utility => "utility",
461 DecisionPolicyKind::Rule => "rule",
462 DecisionPolicyKind::Random => "random",
463 DecisionPolicyKind::Human => "human",
464 DecisionPolicyKind::External => "external",
465 DecisionPolicyKind::Llm => "llm",
466 }
467}
468
469impl UtilityEvaluator for GuardedUtilityPolicy {
470 fn evaluate(
471 &self,
472 ticket: &DecisionTicket,
473 option: &DecisionOption,
474 ) -> Result<DecisionOptionEvaluation, DecisionError> {
475 self.utility.evaluate(ticket, option)
476 }
477}
478
479impl DecisionPolicy for GuardedUtilityPolicy {
480 fn identity(&self) -> &DecisionPolicyIdentity {
481 &self.identity
482 }
483
484 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
485 let mut run = run_ordered_rules(self.guards.rules(), ticket)?;
486 let fired_guards = std::mem::take(&mut run.fired);
487 if let Some(terminal) = run.terminal.take() {
488 let (outcome, summary) = match terminal {
489 RuleTerminal::Select { rule, option_id } => (
490 DecisionOutcome::Selected {
491 option_id: option_id.clone(),
492 },
493 format!("guard {rule} selected {option_id}"),
494 ),
495 RuleTerminal::Defer { rule, reason } => (
496 DecisionOutcome::Deferred {
497 reason: reason.clone(),
498 },
499 format!("guard {rule} deferred: {reason}"),
500 ),
501 };
502 return Ok(PolicyDecision {
503 outcome,
504 summary,
505 evaluations: run.exclusion_evaluations(),
506 external: None,
507 random: None,
508 stage: Some(DecisionStage::Guard),
509 fired_guards,
510 });
511 }
512 let mut evaluations = ticket
513 .options
514 .iter()
515 .map(|option| match run.excluded.get(&option.id) {
516 Some(blocker) => {
517 let mut evaluation = excluded_evaluation(&option.id, blocker);
518 evaluation
519 .blockers
520 .splice(0..0, option.blockers.iter().cloned());
521 Ok(evaluation)
522 }
523 None => self.utility.evaluate(ticket, option),
524 })
525 .collect::<Result<Vec<_>, _>>()?;
526 evaluations.sort_by(|left, right| left.option_id.cmp(&right.option_id));
529 let scored = evaluations
530 .iter()
531 .filter_map(|evaluation| {
532 evaluation
533 .score
534 .map(|score| (score, evaluation.option_id.as_str()))
535 })
536 .collect::<Vec<_>>();
537 let Some(&(best, winner)) = scored
538 .iter()
539 .max_by(|left, right| left.0.cmp(&right.0).then_with(|| right.1.cmp(left.1)))
540 else {
541 return Ok(PolicyDecision {
542 outcome: DecisionOutcome::Deferred {
543 reason: "no available option".to_owned(),
544 },
545 summary: "guarded utility policy deferred because no option remained after guards"
546 .to_owned(),
547 evaluations,
548 external: None,
549 random: None,
550 stage: Some(DecisionStage::Utility),
551 fired_guards,
552 });
553 };
554 let margin = i128::from(self.near_equivalence_margin);
555 let candidates = scored
556 .iter()
557 .filter(|(score, _)| i128::from(best) - i128::from(*score) <= margin)
558 .map(|(_, option_id)| DecisionOptionWeight::new(*option_id, 1))
559 .collect::<Vec<_>>();
560 if self.random_tie_break && candidates.len() > 1 {
561 let summary = format!(
562 "guarded utility policy left {} near-equivalent options for a random tie-break",
563 candidates.len()
564 );
565 return Ok(PolicyDecision {
566 outcome: DecisionOutcome::PendingRandom { candidates },
567 summary,
568 evaluations,
569 external: None,
570 random: None,
571 stage: Some(DecisionStage::Random),
572 fired_guards,
573 });
574 }
575 let option_id = winner.to_owned();
576 Ok(PolicyDecision {
577 outcome: DecisionOutcome::Selected {
578 option_id: option_id.clone(),
579 },
580 summary: format!("guarded utility policy selected {option_id}"),
581 evaluations,
582 external: None,
583 random: None,
584 stage: Some(DecisionStage::Utility),
585 fired_guards,
586 })
587 }
588}
589
590#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
591pub struct HumanDecisionResponse {
592 pub ticket_version: u64,
593 pub option_id: String,
594 pub operator_id: String,
595}
596
597pub trait HumanPolicy: DecisionPolicy {
598 fn submitted_response(&self, ticket: &DecisionTicket) -> Option<HumanDecisionResponse>;
599}
600
601#[derive(Clone, Debug)]
602pub struct QueuedHumanPolicy {
603 identity: DecisionPolicyIdentity,
604 responses: BTreeMap<canwu_core::DecisionTicketId, HumanDecisionResponse>,
605}
606
607impl QueuedHumanPolicy {
608 #[must_use]
609 pub fn new(id: impl Into<String>, version: impl Into<String>) -> Self {
610 Self {
611 identity: DecisionPolicyIdentity::new(DecisionPolicyKind::Human, id, version),
612 responses: BTreeMap::new(),
613 }
614 }
615
616 pub fn submit(
617 &mut self,
618 ticket_id: canwu_core::DecisionTicketId,
619 response: HumanDecisionResponse,
620 ) -> Result<(), DecisionError> {
621 if let Some(existing) = self.responses.get(&ticket_id) {
622 return match existing.ticket_version.cmp(&response.ticket_version) {
623 Ordering::Less => {
624 self.responses.insert(ticket_id, response);
625 Ok(())
626 }
627 Ordering::Equal => Err(DecisionError::new(
628 DecisionErrorCode::DuplicateResponse,
629 "a response has already been submitted for this decision ticket version",
630 )),
631 Ordering::Greater => Err(DecisionError::new(
632 DecisionErrorCode::VersionConflict,
633 "a stale decision response cannot replace a newer queued response",
634 )),
635 };
636 }
637 self.responses.insert(ticket_id, response);
638 Ok(())
639 }
640}
641
642impl HumanPolicy for QueuedHumanPolicy {
643 fn submitted_response(&self, ticket: &DecisionTicket) -> Option<HumanDecisionResponse> {
644 self.responses.get(&ticket.id).cloned()
645 }
646}
647
648impl DecisionPolicy for QueuedHumanPolicy {
649 fn identity(&self) -> &DecisionPolicyIdentity {
650 &self.identity
651 }
652
653 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
654 let Some(response) = self.submitted_response(ticket) else {
655 return Ok(PolicyDecision::pending("awaiting human selection"));
656 };
657 if response.ticket_version != ticket.version {
658 return Err(DecisionError::new(
659 DecisionErrorCode::VersionConflict,
660 "human response targets a stale decision ticket version",
661 ));
662 }
663 Ok(PolicyDecision {
664 outcome: DecisionOutcome::Selected {
665 option_id: response.option_id,
666 },
667 summary: format!("human operator {} selected an option", response.operator_id),
668 evaluations: Vec::new(),
669 external: Some(DecisionExternalEvidence {
670 provider: "human".to_owned(),
671 model: None,
672 prompt_contract: None,
673 request_id: Some(response.operator_id),
674 metadata: BTreeMap::new(),
675 }),
676 random: None,
677 stage: None,
678 fired_guards: Vec::new(),
679 })
680 }
681}
682
683#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
684pub struct ExternalDecisionOption {
685 pub id: String,
686 pub label: String,
687 pub description: String,
688 pub metadata: serde_json::Value,
689}
690
691#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
692pub struct ExternalDecisionRequest {
693 pub ticket_id: canwu_core::DecisionTicketId,
694 pub ticket_version: u64,
695 pub definition: String,
696 pub summary: String,
697 pub context: crate::DecisionContext,
698 pub options: Vec<ExternalDecisionOption>,
699}
700
701impl From<&DecisionTicket> for ExternalDecisionRequest {
702 fn from(ticket: &DecisionTicket) -> Self {
703 Self {
704 ticket_id: ticket.id,
705 ticket_version: ticket.version,
706 definition: ticket.definition.clone(),
707 summary: ticket.summary.clone(),
708 context: ticket.context.clone(),
709 options: ticket
710 .options
711 .iter()
712 .filter(|option| option.is_available())
713 .map(|option| ExternalDecisionOption {
714 id: option.id.clone(),
715 label: option.label.clone(),
716 description: option.description.clone(),
717 metadata: option.metadata.clone(),
718 })
719 .collect(),
720 }
721 }
722}
723
724#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
725pub struct ExternalDecisionResponse {
726 pub ticket_version: u64,
727 pub option_id: String,
728 pub provider: String,
729 pub request_id: String,
730 #[serde(default)]
731 pub metadata: BTreeMap<String, String>,
732}
733
734pub trait ExternalPolicy: DecisionPolicy {
735 fn external_request(&self, ticket: &DecisionTicket) -> ExternalDecisionRequest {
736 ticket.into()
737 }
738
739 fn submitted_response(&self, ticket: &DecisionTicket) -> Option<ExternalDecisionResponse>;
740}
741
742#[derive(Clone, Debug)]
743pub struct QueuedExternalPolicy {
744 identity: DecisionPolicyIdentity,
745 responses: BTreeMap<canwu_core::DecisionTicketId, ExternalDecisionResponse>,
746}
747
748impl QueuedExternalPolicy {
749 #[must_use]
750 pub fn new(id: impl Into<String>, version: impl Into<String>) -> Self {
751 Self {
752 identity: DecisionPolicyIdentity::new(DecisionPolicyKind::External, id, version),
753 responses: BTreeMap::new(),
754 }
755 }
756
757 pub fn submit(
758 &mut self,
759 ticket_id: canwu_core::DecisionTicketId,
760 response: ExternalDecisionResponse,
761 ) -> Result<(), DecisionError> {
762 if let Some(existing) = self.responses.get(&ticket_id) {
763 return match existing.ticket_version.cmp(&response.ticket_version) {
764 Ordering::Less => {
765 self.responses.insert(ticket_id, response);
766 Ok(())
767 }
768 Ordering::Equal => Err(DecisionError::new(
769 DecisionErrorCode::DuplicateResponse,
770 "a response has already been submitted for this decision ticket version",
771 )),
772 Ordering::Greater => Err(DecisionError::new(
773 DecisionErrorCode::VersionConflict,
774 "a stale decision response cannot replace a newer queued response",
775 )),
776 };
777 }
778 self.responses.insert(ticket_id, response);
779 Ok(())
780 }
781}
782
783impl ExternalPolicy for QueuedExternalPolicy {
784 fn submitted_response(&self, ticket: &DecisionTicket) -> Option<ExternalDecisionResponse> {
785 self.responses.get(&ticket.id).cloned()
786 }
787}
788
789impl DecisionPolicy for QueuedExternalPolicy {
790 fn identity(&self) -> &DecisionPolicyIdentity {
791 &self.identity
792 }
793
794 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
795 let Some(response) = self.submitted_response(ticket) else {
796 return Ok(PolicyDecision::pending("awaiting external policy response"));
797 };
798 if response.ticket_version != ticket.version {
799 return Err(DecisionError::new(
800 DecisionErrorCode::VersionConflict,
801 "external response targets a stale decision ticket version",
802 ));
803 }
804 Ok(PolicyDecision {
805 outcome: DecisionOutcome::Selected {
806 option_id: response.option_id,
807 },
808 summary: format!("external provider {} selected an option", response.provider),
809 evaluations: Vec::new(),
810 external: Some(DecisionExternalEvidence {
811 provider: response.provider,
812 model: None,
813 prompt_contract: None,
814 request_id: Some(response.request_id),
815 metadata: response.metadata,
816 }),
817 random: None,
818 stage: None,
819 fired_guards: Vec::new(),
820 })
821 }
822}
823
824#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
825pub struct LlmModelIdentity {
826 pub provider: String,
827 pub model: String,
828 pub prompt_contract: String,
829}
830
831pub trait LlmPolicy: ExternalPolicy {
832 fn model_identity(&self) -> &LlmModelIdentity;
833}
834
835#[derive(Clone, Debug)]
836pub struct QueuedLlmPolicy {
837 identity: DecisionPolicyIdentity,
838 model: LlmModelIdentity,
839 responses: BTreeMap<canwu_core::DecisionTicketId, ExternalDecisionResponse>,
840}
841
842impl QueuedLlmPolicy {
843 #[must_use]
844 pub fn new(id: impl Into<String>, version: impl Into<String>, model: LlmModelIdentity) -> Self {
845 Self {
846 identity: DecisionPolicyIdentity::new(DecisionPolicyKind::Llm, id, version),
847 model,
848 responses: BTreeMap::new(),
849 }
850 }
851
852 pub fn submit(
853 &mut self,
854 ticket_id: canwu_core::DecisionTicketId,
855 response: ExternalDecisionResponse,
856 ) -> Result<(), DecisionError> {
857 if let Some(existing) = self.responses.get(&ticket_id) {
858 return match existing.ticket_version.cmp(&response.ticket_version) {
859 Ordering::Less => {
860 self.responses.insert(ticket_id, response);
861 Ok(())
862 }
863 Ordering::Equal => Err(DecisionError::new(
864 DecisionErrorCode::DuplicateResponse,
865 "a response has already been submitted for this decision ticket version",
866 )),
867 Ordering::Greater => Err(DecisionError::new(
868 DecisionErrorCode::VersionConflict,
869 "a stale decision response cannot replace a newer queued response",
870 )),
871 };
872 }
873 self.responses.insert(ticket_id, response);
874 Ok(())
875 }
876}
877
878impl ExternalPolicy for QueuedLlmPolicy {
879 fn submitted_response(&self, ticket: &DecisionTicket) -> Option<ExternalDecisionResponse> {
880 self.responses.get(&ticket.id).cloned()
881 }
882}
883
884impl LlmPolicy for QueuedLlmPolicy {
885 fn model_identity(&self) -> &LlmModelIdentity {
886 &self.model
887 }
888}
889
890impl DecisionPolicy for QueuedLlmPolicy {
891 fn identity(&self) -> &DecisionPolicyIdentity {
892 &self.identity
893 }
894
895 fn decide(&self, ticket: &DecisionTicket) -> Result<PolicyDecision, DecisionError> {
896 let Some(response) = self.submitted_response(ticket) else {
897 return Ok(PolicyDecision::pending(
898 "awaiting constrained LLM option selection",
899 ));
900 };
901 if response.ticket_version != ticket.version {
902 return Err(DecisionError::new(
903 DecisionErrorCode::VersionConflict,
904 "LLM response targets a stale decision ticket version",
905 ));
906 }
907 if response.provider != self.model.provider {
908 return Err(DecisionError::new(
909 DecisionErrorCode::PolicyMismatch,
910 "LLM response provider does not match the configured model identity",
911 ));
912 }
913 Ok(PolicyDecision {
914 outcome: DecisionOutcome::Selected {
915 option_id: response.option_id,
916 },
917 summary: format!(
918 "LLM {}:{} selected an existing option",
919 self.model.provider, self.model.model
920 ),
921 evaluations: Vec::new(),
922 external: Some(DecisionExternalEvidence {
923 provider: response.provider,
924 model: Some(self.model.model.clone()),
925 prompt_contract: Some(self.model.prompt_contract.clone()),
926 request_id: Some(response.request_id),
927 metadata: response.metadata,
928 }),
929 random: None,
930 stage: None,
931 fired_guards: Vec::new(),
932 })
933 }
934}