1use std::collections::BTreeMap;
50use std::fmt;
51
52use turnframe_core::command::{
53 CommandBatch, CommandEnvelope, CommandOrigin, ConfirmationPolicy, RiskClass, origin_satisfies,
54};
55use turnframe_core::event::OperationalReceipt;
56use turnframe_core::hash::canonical_value;
57use turnframe_core::ids::{CommandId, EventId, ReceiptId};
58use turnframe_core::reduce::CommandRef;
59use turnframe_core::replay::{CommandOutcome, ReplayRecord};
60use turnframe_core::response::AssistantTurn;
61use turnframe_core::understanding::Understanding;
62
63use crate::assertions::{AssertionFailure, identical_blocks};
64
65#[derive(Debug, Clone, PartialEq)]
71pub struct TurnExecution {
72 pub record: ReplayRecord,
74 pub response: AssistantTurn,
76}
77
78impl TurnExecution {
79 #[must_use]
81 pub fn new(record: ReplayRecord, response: AssistantTurn) -> Self {
82 Self { record, response }
83 }
84
85 pub fn same_as(&self, other: &Self) -> Result<(), ReplayDivergence> {
91 same_turn(self, other)
92 }
93}
94
95#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
101#[non_exhaustive]
102pub enum ReplayDivergence {
103 #[error("one execution recorded a normalized plan ({left}) and the other did not ({right})")]
105 PlanPresenceDiffers {
106 left: bool,
108 right: bool,
110 },
111 #[error("the normalized plans differ: {detail}")]
113 PlanDiffers {
114 detail: String,
116 },
117 #[error("the executions produced {left} and {right} command(s)")]
119 CommandCountDiffers {
120 left: usize,
122 right: usize,
124 },
125 #[error("command {index} differs: {detail}")]
127 CommandDiffers {
128 index: usize,
130 detail: String,
132 },
133 #[error("the executions committed {left} and {right} event(s)")]
135 EventCountDiffers {
136 left: usize,
138 right: usize,
140 },
141 #[error("event {index} differs: {left} then {right}")]
143 EventsDiffer {
144 index: usize,
146 left: EventId,
148 right: EventId,
150 },
151 #[error("the answers differ: {0}")]
153 ResponseDiffers(AssertionFailure),
154}
155
156pub fn same_plan(left: &ReplayRecord, right: &ReplayRecord) -> Result<(), ReplayDivergence> {
163 match (&left.understanding, &right.understanding) {
164 (None, None) => Ok(()),
165 (left, right) if left.is_some() != right.is_some() => {
166 Err(ReplayDivergence::PlanPresenceDiffers {
167 left: left.is_some(),
168 right: right.is_some(),
169 })
170 }
171 (Some(left), Some(right)) => {
172 if canonical_value(left).ok() == canonical_value(right).ok() {
173 Ok(())
174 } else {
175 Err(ReplayDivergence::PlanDiffers {
176 detail: first_plan_difference(left, right),
177 })
178 }
179 }
180 _ => Ok(()),
181 }
182}
183
184fn first_plan_difference(left: &Understanding, right: &Understanding) -> String {
186 if left.acts.len() != right.acts.len() {
187 return format!("act count {} then {}", left.acts.len(), right.acts.len());
188 }
189 for (a, b) in left.acts.iter().zip(&right.acts) {
190 if a != b {
191 return format!("act {} then {}", a.id, b.id);
192 }
193 }
194 if left.questions != right.questions {
195 return format!(
196 "questions ({} then {})",
197 left.questions.len(),
198 right.questions.len()
199 );
200 }
201 if left.constraints != right.constraints {
202 return format!(
203 "constraints ({} then {})",
204 left.constraints.len(),
205 right.constraints.len()
206 );
207 }
208 "other content".to_owned()
209}
210
211pub fn same_commands(left: &ReplayRecord, right: &ReplayRecord) -> Result<(), ReplayDivergence> {
218 if left.command_outcomes.len() != right.command_outcomes.len() {
219 return Err(ReplayDivergence::CommandCountDiffers {
220 left: left.command_outcomes.len(),
221 right: right.command_outcomes.len(),
222 });
223 }
224 for (index, (a, b)) in left
225 .command_outcomes
226 .iter()
227 .zip(&right.command_outcomes)
228 .enumerate()
229 {
230 let detail = if a.command_ref != b.command_ref {
231 Some(format!("{} then {}", a.command_ref, b.command_ref))
232 } else if a.idempotency_key != b.idempotency_key {
233 Some("idempotency key".to_owned())
234 } else if a.case_ref.key() != b.case_ref.key() {
235 Some("case".to_owned())
236 } else if a.case_ref.expected_revision != b.case_ref.expected_revision {
237 Some(format!(
238 "expected revision {} then {}",
239 a.case_ref.expected_revision.value(),
240 b.case_ref.expected_revision.value()
241 ))
242 } else if a.outcome != b.outcome {
243 Some(format!(
244 "outcome {} then {}",
245 outcome_label(&a.outcome),
246 outcome_label(&b.outcome)
247 ))
248 } else {
249 None
250 };
251 if let Some(detail) = detail {
252 return Err(ReplayDivergence::CommandDiffers { index, detail });
253 }
254 }
255 Ok(())
256}
257
258fn outcome_label(outcome: &CommandOutcome) -> &'static str {
260 match outcome {
261 CommandOutcome::Committed { .. } => "committed",
262 CommandOutcome::IdempotentReplay => "idempotent_replay",
263 CommandOutcome::RevisionConflict { .. } => "revision_conflict",
264 CommandOutcome::Rejected { .. } => "rejected",
265 CommandOutcome::Failed { .. } => "failed",
266 CommandOutcome::OutcomeUnknown { .. } => "outcome_unknown",
267 CommandOutcome::AwaitingConfirmation { .. } => "awaiting_confirmation",
268 _ => "other",
269 }
270}
271
272pub fn same_events(left: &ReplayRecord, right: &ReplayRecord) -> Result<(), ReplayDivergence> {
279 if left.event_ids.len() != right.event_ids.len() {
280 return Err(ReplayDivergence::EventCountDiffers {
281 left: left.event_ids.len(),
282 right: right.event_ids.len(),
283 });
284 }
285 for (index, (a, b)) in left.event_ids.iter().zip(&right.event_ids).enumerate() {
286 if a != b {
287 return Err(ReplayDivergence::EventsDiffer {
288 index,
289 left: *a,
290 right: *b,
291 });
292 }
293 }
294 Ok(())
295}
296
297pub fn same_blocks(left: &AssistantTurn, right: &AssistantTurn) -> Result<(), ReplayDivergence> {
303 identical_blocks(left, right).map_err(ReplayDivergence::ResponseDiffers)
304}
305
306pub fn same_turn(left: &TurnExecution, right: &TurnExecution) -> Result<(), ReplayDivergence> {
313 same_plan(&left.record, &right.record)?;
314 same_commands(&left.record, &right.record)?;
315 same_events(&left.record, &right.record)?;
316 same_blocks(&left.response, &right.response)
317}
318
319#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
321#[non_exhaustive]
322pub enum ReplayGap {
323 #[error("command {command_ref} has no policy decision")]
325 CommandWithoutPolicyDecision {
326 command_ref: CommandRef,
328 },
329 #[error("command {command_ref} has no origin in the evidence")]
331 CommandWithoutOrigin {
332 command_ref: CommandRef,
334 },
335 #[error(
337 "command {command_ref} carries an origin that does not satisfy its policy \
338 (risk {risk:?}, confirmation {confirmation:?})"
339 )]
340 OriginDoesNotSatisfyPolicy {
341 command_ref: CommandRef,
343 risk: RiskClass,
345 confirmation: ConfirmationPolicy,
347 },
348 #[error("command {command_ref} committed although policy refused it ({reason_key})")]
350 RefusedCommandCommitted {
351 command_ref: CommandRef,
353 reason_key: String,
355 },
356 #[error("policy decision for {command_ref} has no command outcome")]
358 DecisionWithoutOutcome {
359 command_ref: CommandRef,
361 },
362 #[error("command {command_ref} committed event {event_id}, which the record does not list")]
364 CommittedEventNotRecorded {
365 command_ref: CommandRef,
367 event_id: EventId,
369 },
370 #[error("receipt {receipt_id} ({status_code}) cites no event")]
372 ReceiptWithoutEvents {
373 receipt_id: ReceiptId,
375 status_code: String,
377 },
378 #[error(
380 "receipt {receipt_id} ({status_code}) cites event {event_id}, which the record does not list"
381 )]
382 ReceiptCitesUnrecordedEvent {
383 receipt_id: ReceiptId,
385 status_code: String,
387 event_id: EventId,
389 },
390}
391
392#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
394pub struct ReplayGaps {
395 pub gaps: Vec<ReplayGap>,
397}
398
399impl fmt::Display for ReplayGaps {
400 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
401 write!(f, "the replay record does not explain its turn:")?;
402 for gap in &self.gaps {
403 write!(f, "\n- {gap}")?;
404 }
405 Ok(())
406 }
407}
408
409#[derive(Debug, Clone)]
412pub struct ReplayEvidence<'a> {
413 record: &'a ReplayRecord,
414 origins: BTreeMap<CommandId, CommandOrigin>,
415 receipts: Vec<OperationalReceipt>,
416}
417
418impl<'a> ReplayEvidence<'a> {
419 #[must_use]
421 pub fn new(record: &'a ReplayRecord) -> Self {
422 Self {
423 record,
424 origins: BTreeMap::new(),
425 receipts: Vec::new(),
426 }
427 }
428
429 #[must_use]
431 pub fn with_batch<C>(mut self, batch: &CommandBatch<C>) -> Self {
432 for envelope in &batch.envelopes {
433 self.origins
434 .insert(envelope.command_id, envelope.origin.clone());
435 }
436 self
437 }
438
439 #[must_use]
441 pub fn with_envelope<C>(mut self, envelope: &CommandEnvelope<C>) -> Self {
442 self.origins
443 .insert(envelope.command_id, envelope.origin.clone());
444 self
445 }
446
447 #[must_use]
449 pub fn with_origin(mut self, command_id: CommandId, origin: CommandOrigin) -> Self {
450 self.origins.insert(command_id, origin);
451 self
452 }
453
454 #[must_use]
456 pub fn with_receipts(mut self, receipts: &[OperationalReceipt]) -> Self {
457 self.receipts.extend_from_slice(receipts);
458 self
459 }
460
461 #[must_use]
463 pub fn record(&self) -> &ReplayRecord {
464 self.record
465 }
466
467 pub fn explains_its_turn(&self) -> Result<(), ReplayGaps> {
479 let mut gaps = Vec::new();
480 for outcome in &self.record.command_outcomes {
481 let command_ref = outcome.command_ref;
482 match self
483 .record
484 .policy_decisions
485 .iter()
486 .find(|decision| decision.command_ref == command_ref)
487 {
488 None => gaps.push(ReplayGap::CommandWithoutPolicyDecision { command_ref }),
489 Some(decision) => {
490 match self.origins.get(&command_ref.command_id) {
491 None => gaps.push(ReplayGap::CommandWithoutOrigin { command_ref }),
492 Some(origin) => {
493 if !origin_satisfies(origin, &decision.policy) {
494 gaps.push(ReplayGap::OriginDoesNotSatisfyPolicy {
495 command_ref,
496 risk: decision.policy.risk,
497 confirmation: decision.policy.confirmation,
498 });
499 }
500 }
501 }
502 if !decision.allowed
503 && matches!(outcome.outcome, CommandOutcome::Committed { .. })
504 {
505 gaps.push(ReplayGap::RefusedCommandCommitted {
506 command_ref,
507 reason_key: decision.reason_key.clone(),
508 });
509 }
510 }
511 }
512 if let CommandOutcome::Committed { event_ids, .. } = &outcome.outcome {
513 for event_id in event_ids {
514 if !self.record.event_ids.contains(event_id) {
515 gaps.push(ReplayGap::CommittedEventNotRecorded {
516 command_ref,
517 event_id: *event_id,
518 });
519 }
520 }
521 }
522 }
523 for decision in &self.record.policy_decisions {
524 if !self
525 .record
526 .command_outcomes
527 .iter()
528 .any(|outcome| outcome.command_ref == decision.command_ref)
529 {
530 gaps.push(ReplayGap::DecisionWithoutOutcome {
531 command_ref: decision.command_ref,
532 });
533 }
534 }
535 for receipt in &self.receipts {
536 if receipt.event_ids.is_empty() {
537 gaps.push(ReplayGap::ReceiptWithoutEvents {
538 receipt_id: receipt.receipt_id,
539 status_code: receipt.status_code.clone(),
540 });
541 }
542 for event_id in &receipt.event_ids {
543 if !self.record.event_ids.contains(event_id) {
544 gaps.push(ReplayGap::ReceiptCitesUnrecordedEvent {
545 receipt_id: receipt.receipt_id,
546 status_code: receipt.status_code.clone(),
547 event_id: *event_id,
548 });
549 }
550 }
551 }
552 if gaps.is_empty() {
553 Ok(())
554 } else {
555 Err(ReplayGaps { gaps })
556 }
557 }
558}
559
560#[cfg(test)]
561mod tests {
562 use super::*;
563 use turnframe_core::case::CaseRef;
564 use turnframe_core::command::{CommandPolicy, IdempotencyKey};
565 use turnframe_core::event::ReceiptSeverity;
566 use turnframe_core::ids::{AccountId, BatchId, BlockId, CaseRevision, ConversationId, TurnId};
567 use turnframe_core::locale::LocalizedText;
568 use turnframe_core::policy::{PolicySnapshot, reason};
569 use turnframe_core::replay::CommandOutcomeRecord;
570 use turnframe_core::response::{GeneratedTransition, ReplayToken, ResponseBlock};
571
572 fn command_ref() -> CommandRef {
573 CommandRef {
574 batch_id: BatchId::nil(),
575 command_id: CommandId::nil(),
576 }
577 }
578
579 fn event(byte: u8) -> EventId {
580 EventId::from(uuid::Uuid::from_bytes([byte; 16]))
581 }
582
583 fn record() -> ReplayRecord {
584 ReplayRecord::received(
585 TurnId::nil(),
586 ConversationId::nil(),
587 AccountId::from("aurora"),
588 chrono::DateTime::from_timestamp(1_700_000_000, 0).unwrap_or_default(),
589 )
590 }
591
592 fn direct_origin() -> CommandOrigin {
593 CommandOrigin::DirectSafeUserAct {
594 evidence_digest: turnframe_core::hash::Digest::of_bytes(b"evidence"),
595 }
596 }
597
598 fn committed_record() -> ReplayRecord {
600 let mut record = record();
601 let policy = CommandPolicy::low_risk();
602 record
603 .policy_decisions
604 .push(PolicySnapshot::conservative().decide(command_ref(), &policy, &direct_origin()));
605 record.command_outcomes.push(CommandOutcomeRecord {
606 command_ref: command_ref(),
607 idempotency_key: IdempotencyKey("idem-1".to_owned()),
608 case_ref: CaseRef::new("trip", "trip-1", CaseRevision(1)),
609 origin: Some(direct_origin()),
610 outcome: CommandOutcome::Committed {
611 new_revision: CaseRevision(2),
612 event_ids: vec![event(1)],
613 },
614 });
615 record.event_ids.push(event(1));
616 record
617 }
618
619 fn turn(text: &str) -> AssistantTurn {
620 AssistantTurn {
621 turn_id: TurnId::nil(),
622 conversation_id: ConversationId::nil(),
623 blocks: vec![ResponseBlock::Transition(GeneratedTransition {
624 block_id: BlockId::from("t1"),
625 text: text.to_owned(),
626 facts_used: Vec::new(),
627 })],
628 replay_token: ReplayToken::from("rt"),
629 subjects: Vec::new(),
630 expectations: Vec::new(),
631 done: Vec::new(),
632 }
633 }
634
635 fn receipt(event_ids: Vec<EventId>) -> OperationalReceipt {
636 OperationalReceipt {
637 receipt_id: ReceiptId::derive(&event_ids, "trip.rebooking_sent"),
638 event_ids,
639 severity: ReceiptSeverity::Success,
640 title: LocalizedText::new("Sent"),
641 body: LocalizedText::new("The rebooking was sent."),
642 status_code: "trip.rebooking_sent".to_owned(),
643 artifact_refs: Vec::new(),
644 }
645 }
646
647 #[test]
648 fn two_identical_executions_agree() {
649 let execution = TurnExecution::new(committed_record(), turn("fatto"));
650 assert_eq!(execution.same_as(&execution.clone()), Ok(()));
651 }
652
653 #[test]
654 fn a_different_plan_is_reported_before_anything_else() {
655 let text = "Cambia il nome in Lisbona";
656 let plan = crate::providers::UnderstandingBuilder::of(text)
657 .apply(
658 "trip.set_name",
659 "tok_1",
660 serde_json::Value::Null,
661 "Cambia il nome",
662 )
663 .build()
664 .unwrap();
665 let other = crate::providers::UnderstandingBuilder::of(text)
666 .ask("Cambia")
667 .build()
668 .unwrap();
669
670 let mut left = TurnExecution::new(committed_record(), turn("fatto"));
671 let mut right = left.clone();
672 left.record.understanding = Some(plan);
673 right.record.understanding = Some(other);
674 right.record.command_outcomes.clear();
676
677 let divergence = same_turn(&left, &right).unwrap_err();
678 assert_eq!(
679 divergence,
680 ReplayDivergence::PlanDiffers {
681 detail: "act count 1 then 0".to_owned(),
682 }
683 );
684
685 right.record.understanding = None;
686 assert_eq!(
687 same_plan(&left.record, &right.record).unwrap_err(),
688 ReplayDivergence::PlanPresenceDiffers {
689 left: true,
690 right: false,
691 }
692 );
693 }
694
695 #[test]
696 fn a_different_command_outcome_is_named_by_position() {
697 let left = TurnExecution::new(committed_record(), turn("fatto"));
698 let mut right = left.clone();
699 right.record.command_outcomes[0].outcome = CommandOutcome::IdempotentReplay;
700
701 assert_eq!(
702 same_turn(&left, &right).unwrap_err(),
703 ReplayDivergence::CommandDiffers {
704 index: 0,
705 detail: "outcome committed then idempotent_replay".to_owned(),
706 }
707 );
708
709 let mut fewer = left.clone();
710 fewer.record.command_outcomes.clear();
711 assert_eq!(
712 same_commands(&left.record, &fewer.record).unwrap_err(),
713 ReplayDivergence::CommandCountDiffers { left: 1, right: 0 }
714 );
715 }
716
717 #[test]
718 fn different_events_and_blocks_are_reported_in_order() {
719 let left = TurnExecution::new(committed_record(), turn("fatto"));
720 let mut right = left.clone();
721 right.record.event_ids = vec![event(2)];
722 right.record.command_outcomes[0].outcome = CommandOutcome::Committed {
723 new_revision: CaseRevision(2),
724 event_ids: vec![event(2)],
725 };
726 assert!(
727 matches!(
728 same_turn(&left, &right).unwrap_err(),
729 ReplayDivergence::CommandDiffers { .. }
730 ),
731 "the command outcome differs before the event list does"
732 );
733
734 let mut only_events = left.clone();
735 only_events.record.event_ids = vec![event(2)];
736 assert_eq!(
737 same_events(&left.record, &only_events.record).unwrap_err(),
738 ReplayDivergence::EventsDiffer {
739 index: 0,
740 left: event(1),
741 right: event(2),
742 }
743 );
744
745 let other_answer = TurnExecution::new(committed_record(), turn("done"));
746 assert!(matches!(
747 same_turn(&left, &other_answer).unwrap_err(),
748 ReplayDivergence::ResponseDiffers(AssertionFailure::BlocksDiffer { index: 0 })
749 ));
750 }
751
752 #[test]
753 fn a_record_with_a_decision_an_origin_and_its_events_explains_itself() {
754 let record = committed_record();
755 ReplayEvidence::new(&record)
756 .with_origin(CommandId::nil(), direct_origin())
757 .with_receipts(&[receipt(vec![event(1)])])
758 .explains_its_turn()
759 .expect("the record accounts for its turn");
760 }
761
762 #[test]
763 fn a_command_without_a_policy_decision_is_a_gap() {
764 let mut record = committed_record();
765 record.policy_decisions.clear();
766 let gaps = ReplayEvidence::new(&record)
767 .with_origin(CommandId::nil(), direct_origin())
768 .explains_its_turn()
769 .unwrap_err();
770 assert_eq!(
771 gaps.gaps,
772 vec![ReplayGap::CommandWithoutPolicyDecision {
773 command_ref: command_ref(),
774 }]
775 );
776 assert!(gaps.to_string().contains("has no policy decision"));
777 }
778
779 #[test]
780 fn a_command_without_an_origin_is_a_gap() {
781 let record = committed_record();
782 let gaps = ReplayEvidence::new(&record)
783 .explains_its_turn()
784 .unwrap_err();
785 assert_eq!(
786 gaps.gaps,
787 vec![ReplayGap::CommandWithoutOrigin {
788 command_ref: command_ref(),
789 }]
790 );
791 }
792
793 #[test]
794 fn an_origin_that_does_not_satisfy_the_recorded_policy_is_a_gap() {
795 let mut record = committed_record();
796 let policy = CommandPolicy {
797 risk: RiskClass::ExternalRegulated,
798 confirmation: ConfirmationPolicy::ExplicitClick,
799 ..CommandPolicy::conservative()
800 };
801 record.policy_decisions =
802 vec![PolicySnapshot::conservative().decide(command_ref(), &policy, &direct_origin())];
803
804 let gaps = ReplayEvidence::new(&record)
805 .with_origin(CommandId::nil(), direct_origin())
806 .explains_its_turn()
807 .unwrap_err();
808
809 assert!(gaps.gaps.contains(&ReplayGap::OriginDoesNotSatisfyPolicy {
810 command_ref: command_ref(),
811 risk: RiskClass::ExternalRegulated,
812 confirmation: ConfirmationPolicy::ExplicitClick,
813 }));
814 assert!(gaps.gaps.contains(&ReplayGap::RefusedCommandCommitted {
815 command_ref: command_ref(),
816 reason_key: reason::CONFIRMATION_REQUIRED.to_owned(),
817 }));
818 }
819
820 #[test]
821 fn an_unaccounted_decision_and_an_invented_event_are_gaps() {
822 let mut record = committed_record();
823 record.command_outcomes.clear();
824 record.event_ids.clear();
825
826 let gaps = ReplayEvidence::new(&record)
827 .with_receipts(&[receipt(vec![event(9)]), receipt(Vec::new())])
828 .explains_its_turn()
829 .unwrap_err();
830
831 assert!(gaps.gaps.contains(&ReplayGap::DecisionWithoutOutcome {
832 command_ref: command_ref(),
833 }));
834 assert!(
835 gaps.gaps
836 .iter()
837 .any(|gap| matches!(gap, ReplayGap::ReceiptCitesUnrecordedEvent { .. }))
838 );
839 assert!(
840 gaps.gaps
841 .iter()
842 .any(|gap| matches!(gap, ReplayGap::ReceiptWithoutEvents { .. }))
843 );
844 }
845
846 #[test]
847 fn a_command_claiming_an_unlisted_event_is_a_gap() {
848 let mut record = committed_record();
849 record.event_ids.clear();
850 let gaps = ReplayEvidence::new(&record)
851 .with_origin(CommandId::nil(), direct_origin())
852 .explains_its_turn()
853 .unwrap_err();
854 assert_eq!(
855 gaps.gaps,
856 vec![ReplayGap::CommittedEventNotRecorded {
857 command_ref: command_ref(),
858 event_id: event(1),
859 }]
860 );
861 assert_eq!(ReplayEvidence::new(&record).record().turn_id, TurnId::nil());
862 }
863}