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 offers: Vec::new(),
633 }
634 }
635
636 fn receipt(event_ids: Vec<EventId>) -> OperationalReceipt {
637 OperationalReceipt {
638 receipt_id: ReceiptId::derive(&event_ids, "trip.rebooking_sent"),
639 event_ids,
640 severity: ReceiptSeverity::Success,
641 title: LocalizedText::new("Sent"),
642 body: LocalizedText::new("The rebooking was sent."),
643 status_code: "trip.rebooking_sent".to_owned(),
644 artifact_refs: Vec::new(),
645 }
646 }
647
648 #[test]
649 fn two_identical_executions_agree() {
650 let execution = TurnExecution::new(committed_record(), turn("fatto"));
651 assert_eq!(execution.same_as(&execution.clone()), Ok(()));
652 }
653
654 #[test]
655 fn a_different_plan_is_reported_before_anything_else() {
656 let text = "Cambia il nome in Lisbona";
657 let plan = crate::providers::UnderstandingBuilder::of(text)
658 .apply(
659 "trip.set_name",
660 "tok_1",
661 serde_json::Value::Null,
662 "Cambia il nome",
663 )
664 .build()
665 .unwrap();
666 let other = crate::providers::UnderstandingBuilder::of(text)
667 .ask("Cambia")
668 .build()
669 .unwrap();
670
671 let mut left = TurnExecution::new(committed_record(), turn("fatto"));
672 let mut right = left.clone();
673 left.record.understanding = Some(plan);
674 right.record.understanding = Some(other);
675 right.record.command_outcomes.clear();
677
678 let divergence = same_turn(&left, &right).unwrap_err();
679 assert_eq!(
680 divergence,
681 ReplayDivergence::PlanDiffers {
682 detail: "act count 1 then 0".to_owned(),
683 }
684 );
685
686 right.record.understanding = None;
687 assert_eq!(
688 same_plan(&left.record, &right.record).unwrap_err(),
689 ReplayDivergence::PlanPresenceDiffers {
690 left: true,
691 right: false,
692 }
693 );
694 }
695
696 #[test]
697 fn a_different_command_outcome_is_named_by_position() {
698 let left = TurnExecution::new(committed_record(), turn("fatto"));
699 let mut right = left.clone();
700 right.record.command_outcomes[0].outcome = CommandOutcome::IdempotentReplay;
701
702 assert_eq!(
703 same_turn(&left, &right).unwrap_err(),
704 ReplayDivergence::CommandDiffers {
705 index: 0,
706 detail: "outcome committed then idempotent_replay".to_owned(),
707 }
708 );
709
710 let mut fewer = left.clone();
711 fewer.record.command_outcomes.clear();
712 assert_eq!(
713 same_commands(&left.record, &fewer.record).unwrap_err(),
714 ReplayDivergence::CommandCountDiffers { left: 1, right: 0 }
715 );
716 }
717
718 #[test]
719 fn different_events_and_blocks_are_reported_in_order() {
720 let left = TurnExecution::new(committed_record(), turn("fatto"));
721 let mut right = left.clone();
722 right.record.event_ids = vec![event(2)];
723 right.record.command_outcomes[0].outcome = CommandOutcome::Committed {
724 new_revision: CaseRevision(2),
725 event_ids: vec![event(2)],
726 };
727 assert!(
728 matches!(
729 same_turn(&left, &right).unwrap_err(),
730 ReplayDivergence::CommandDiffers { .. }
731 ),
732 "the command outcome differs before the event list does"
733 );
734
735 let mut only_events = left.clone();
736 only_events.record.event_ids = vec![event(2)];
737 assert_eq!(
738 same_events(&left.record, &only_events.record).unwrap_err(),
739 ReplayDivergence::EventsDiffer {
740 index: 0,
741 left: event(1),
742 right: event(2),
743 }
744 );
745
746 let other_answer = TurnExecution::new(committed_record(), turn("done"));
747 assert!(matches!(
748 same_turn(&left, &other_answer).unwrap_err(),
749 ReplayDivergence::ResponseDiffers(AssertionFailure::BlocksDiffer { index: 0 })
750 ));
751 }
752
753 #[test]
754 fn a_record_with_a_decision_an_origin_and_its_events_explains_itself() {
755 let record = committed_record();
756 ReplayEvidence::new(&record)
757 .with_origin(CommandId::nil(), direct_origin())
758 .with_receipts(&[receipt(vec![event(1)])])
759 .explains_its_turn()
760 .expect("the record accounts for its turn");
761 }
762
763 #[test]
764 fn a_command_without_a_policy_decision_is_a_gap() {
765 let mut record = committed_record();
766 record.policy_decisions.clear();
767 let gaps = ReplayEvidence::new(&record)
768 .with_origin(CommandId::nil(), direct_origin())
769 .explains_its_turn()
770 .unwrap_err();
771 assert_eq!(
772 gaps.gaps,
773 vec![ReplayGap::CommandWithoutPolicyDecision {
774 command_ref: command_ref(),
775 }]
776 );
777 assert!(gaps.to_string().contains("has no policy decision"));
778 }
779
780 #[test]
781 fn a_command_without_an_origin_is_a_gap() {
782 let record = committed_record();
783 let gaps = ReplayEvidence::new(&record)
784 .explains_its_turn()
785 .unwrap_err();
786 assert_eq!(
787 gaps.gaps,
788 vec![ReplayGap::CommandWithoutOrigin {
789 command_ref: command_ref(),
790 }]
791 );
792 }
793
794 #[test]
795 fn an_origin_that_does_not_satisfy_the_recorded_policy_is_a_gap() {
796 let mut record = committed_record();
797 let policy = CommandPolicy {
798 risk: RiskClass::ExternalRegulated,
799 confirmation: ConfirmationPolicy::ExplicitClick,
800 ..CommandPolicy::conservative()
801 };
802 record.policy_decisions =
803 vec![PolicySnapshot::conservative().decide(command_ref(), &policy, &direct_origin())];
804
805 let gaps = ReplayEvidence::new(&record)
806 .with_origin(CommandId::nil(), direct_origin())
807 .explains_its_turn()
808 .unwrap_err();
809
810 assert!(gaps.gaps.contains(&ReplayGap::OriginDoesNotSatisfyPolicy {
811 command_ref: command_ref(),
812 risk: RiskClass::ExternalRegulated,
813 confirmation: ConfirmationPolicy::ExplicitClick,
814 }));
815 assert!(gaps.gaps.contains(&ReplayGap::RefusedCommandCommitted {
816 command_ref: command_ref(),
817 reason_key: reason::CONFIRMATION_REQUIRED.to_owned(),
818 }));
819 }
820
821 #[test]
822 fn an_unaccounted_decision_and_an_invented_event_are_gaps() {
823 let mut record = committed_record();
824 record.command_outcomes.clear();
825 record.event_ids.clear();
826
827 let gaps = ReplayEvidence::new(&record)
828 .with_receipts(&[receipt(vec![event(9)]), receipt(Vec::new())])
829 .explains_its_turn()
830 .unwrap_err();
831
832 assert!(gaps.gaps.contains(&ReplayGap::DecisionWithoutOutcome {
833 command_ref: command_ref(),
834 }));
835 assert!(
836 gaps.gaps
837 .iter()
838 .any(|gap| matches!(gap, ReplayGap::ReceiptCitesUnrecordedEvent { .. }))
839 );
840 assert!(
841 gaps.gaps
842 .iter()
843 .any(|gap| matches!(gap, ReplayGap::ReceiptWithoutEvents { .. }))
844 );
845 }
846
847 #[test]
848 fn a_command_claiming_an_unlisted_event_is_a_gap() {
849 let mut record = committed_record();
850 record.event_ids.clear();
851 let gaps = ReplayEvidence::new(&record)
852 .with_origin(CommandId::nil(), direct_origin())
853 .explains_its_turn()
854 .unwrap_err();
855 assert_eq!(
856 gaps.gaps,
857 vec![ReplayGap::CommittedEventNotRecorded {
858 command_ref: command_ref(),
859 event_id: event(1),
860 }]
861 );
862 assert_eq!(ReplayEvidence::new(&record).record().turn_id, TurnId::nil());
863 }
864}