1use crate::error::MacpError;
2use crate::log_store::{EntryKind, LogEntry};
3use crate::mode_registry::ModeRegistry;
4use crate::pb::Envelope;
5use crate::policy::registry::PolicyRegistry;
6use crate::registry::PersistedSession;
7use crate::session::{
8 extract_ttl_ms, parse_session_start_payload, validate_canonical_session_start_payload_for_mode,
9 Session, SessionState,
10};
11
12pub fn replay_session(
19 session_id: &str,
20 log_entries: &[LogEntry],
21 registry: &ModeRegistry,
22 policy_registry: Option<&PolicyRegistry>,
23) -> Result<Session, MacpError> {
24 if let Some(session) =
26 try_replay_from_checkpoint(session_id, log_entries, registry, policy_registry)?
27 {
28 return Ok(session);
29 }
30
31 replay_from_start(session_id, log_entries, registry, policy_registry)
32}
33
34fn try_replay_from_checkpoint(
37 session_id: &str,
38 log_entries: &[LogEntry],
39 registry: &ModeRegistry,
40 _policy_registry: Option<&PolicyRegistry>,
41) -> Result<Option<Session>, MacpError> {
42 let checkpoint_idx = log_entries
43 .iter()
44 .rposition(|e| e.entry_kind == EntryKind::Checkpoint);
45
46 let idx = match checkpoint_idx {
47 Some(idx) => idx,
48 None => return Ok(None),
49 };
50
51 let checkpoint = &log_entries[idx];
52 let persisted: PersistedSession =
53 serde_json::from_slice(&checkpoint.raw_payload).map_err(|_| MacpError::InvalidPayload)?;
54 let mut session = Session::from(persisted);
55 session.session_id = session_id.into();
56
57 if !session.policy_version.is_empty() && session.policy_definition.is_none() {
64 tracing::warn!(
65 session_id,
66 policy_version = %session.policy_version,
67 "checkpoint missing policy_definition; falling back to full replay for deterministic policy resolution"
68 );
69 return Ok(None);
70 }
71
72 let mode = registry
73 .get_mode(&session.mode)
74 .ok_or(MacpError::UnknownMode)?;
75
76 for entry in &log_entries[idx + 1..] {
78 replay_entry(&mut session, session_id, entry, &mode)?;
79 }
80
81 Ok(Some(session))
82}
83
84fn replay_entry(
86 session: &mut Session,
87 session_id: &str,
88 entry: &LogEntry,
89 mode: &crate::mode_registry::ModeRef<'_>,
90) -> Result<(), MacpError> {
91 match entry.entry_kind {
92 EntryKind::Incoming => {
93 let replay_env = Envelope {
94 macp_version: if entry.macp_version.is_empty() {
95 "1.0".into()
96 } else {
97 entry.macp_version.clone()
98 },
99 mode: if entry.mode.is_empty() {
100 session.mode.clone()
101 } else {
102 entry.mode.clone()
103 },
104 message_type: entry.message_type.clone(),
105 message_id: entry.message_id.clone(),
106 session_id: session_id.into(),
107 sender: entry.sender.clone(),
108 timestamp_unix_ms: if entry.timestamp_unix_ms != 0 {
111 entry.timestamp_unix_ms
112 } else {
113 entry.received_at_ms
114 },
115 payload: entry.raw_payload.clone(),
116 };
117
118 if session.state != SessionState::Open {
119 if !replay_env.message_id.is_empty() {
120 session.seen_message_ids.insert(replay_env.message_id);
121 }
122 return Ok(());
123 }
124
125 mode.authorize_sender(session, &replay_env)?;
126 let ctx = macp_core::mode::MessageContext::new(if entry.received_at_ms != 0 {
129 entry.received_at_ms
130 } else {
131 entry.timestamp_unix_ms
132 });
133 let response = mode.on_message_at(session, &replay_env, &ctx)?;
134 session.apply_mode_response(response);
135 if !replay_env.message_id.is_empty() {
136 session.seen_message_ids.insert(replay_env.message_id);
137 }
138 }
139 EntryKind::Internal => match entry.message_type.as_str() {
140 "TtlExpired" => {
141 session.state = SessionState::Expired;
142 }
143 "SessionCancel" => {
146 let _ = session.cancel();
147 }
148 "SessionSuspend" => {
152 let at = if entry.received_at_ms != 0 {
153 entry.received_at_ms
154 } else {
155 entry.timestamp_unix_ms
156 };
157 let _ = session.suspend(at);
158 }
159 "SessionResume" => {
160 let at = if entry.received_at_ms != 0 {
161 entry.received_at_ms
162 } else {
163 entry.timestamp_unix_ms
164 };
165 let _ = session.resume(at);
166 }
167 _ => {}
168 },
169 EntryKind::Checkpoint => {
170 }
172 }
173 Ok(())
174}
175
176pub fn validate_replay_consistency(
192 session_id: &str,
193 replayed: &Session,
194 snapshot: &Session,
195) -> u32 {
196 let mut mismatches = 0u32;
197 if replayed.state != snapshot.state {
198 mismatches += 1;
199 tracing::warn!(
200 session_id,
201 replayed_state = ?replayed.state,
202 snapshot_state = ?snapshot.state,
203 "replay/snapshot state mismatch"
204 );
205 }
206 if replayed.seen_message_ids.len() != snapshot.seen_message_ids.len() {
207 mismatches += 1;
208 tracing::warn!(
209 session_id,
210 replayed_dedup = replayed.seen_message_ids.len(),
211 snapshot_dedup = snapshot.seen_message_ids.len(),
212 "replay/snapshot dedup count mismatch"
213 );
214 }
215 if replayed.participants != snapshot.participants {
216 mismatches += 1;
217 tracing::warn!(session_id, "replay/snapshot participants mismatch");
218 }
219 if replayed.mode_version != snapshot.mode_version
220 || replayed.configuration_version != snapshot.configuration_version
221 || replayed.policy_version != snapshot.policy_version
222 {
223 mismatches += 1;
224 tracing::warn!(
225 session_id,
226 "replay/snapshot bound-version mismatch (mode/configuration/policy)"
227 );
228 }
229 if replayed.mode_state != snapshot.mode_state {
233 mismatches += 1;
234 tracing::warn!(
235 session_id,
236 replayed_len = replayed.mode_state.len(),
237 snapshot_len = snapshot.mode_state.len(),
238 "replay/snapshot mode_state mismatch"
239 );
240 }
241 if replayed.accumulated_suspended_ms != snapshot.accumulated_suspended_ms {
245 mismatches += 1;
246 tracing::warn!(
247 session_id,
248 replayed_accumulated_suspended_ms = replayed.accumulated_suspended_ms,
249 snapshot_accumulated_suspended_ms = snapshot.accumulated_suspended_ms,
250 "replay/snapshot accumulated_suspended_ms mismatch"
251 );
252 }
253 if replayed.suspended_at_ms != snapshot.suspended_at_ms {
254 mismatches += 1;
255 tracing::warn!(
256 session_id,
257 replayed_suspended_at_ms = ?replayed.suspended_at_ms,
258 snapshot_suspended_at_ms = ?snapshot.suspended_at_ms,
259 "replay/snapshot suspended_at_ms mismatch"
260 );
261 }
262 if replayed.suspension_intervals != snapshot.suspension_intervals {
267 mismatches += 1;
268 tracing::warn!(
269 session_id,
270 replayed_suspension_cycles = replayed.suspension_intervals.len(),
271 snapshot_suspension_cycles = snapshot.suspension_intervals.len(),
272 "replay/snapshot suspension_intervals mismatch"
273 );
274 }
275 mismatches
276}
277
278fn replay_from_start(
280 session_id: &str,
281 log_entries: &[LogEntry],
282 registry: &ModeRegistry,
283 policy_registry: Option<&PolicyRegistry>,
284) -> Result<Session, MacpError> {
285 let start_entry = log_entries
287 .iter()
288 .find(|e| e.entry_kind == EntryKind::Incoming && e.message_type == "SessionStart")
289 .ok_or(MacpError::InvalidPayload)?;
290
291 let mode_name = if start_entry.mode.is_empty() {
293 return Err(MacpError::InvalidPayload);
296 } else {
297 &start_entry.mode
298 };
299
300 let mode = registry.get_mode(mode_name).ok_or(MacpError::UnknownMode)?;
301
302 let require_complete_start = registry.requires_strict_session_start(mode_name);
304 let start_payload = if start_entry.raw_payload.is_empty() && !require_complete_start {
305 crate::pb::SessionStartPayload::default()
306 } else {
307 parse_session_start_payload(&start_entry.raw_payload)?
308 };
309 if require_complete_start {
313 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
314 }
315
316 let ttl_ms = if !require_complete_start && start_payload.ttl_ms == 0 {
317 60_000i64
319 } else {
320 extract_ttl_ms(&start_payload)?
321 };
322
323 let started_at_unix_ms = start_entry.received_at_ms;
325 let ttl_expiry = started_at_unix_ms.saturating_add(ttl_ms);
326
327 let env = Envelope {
328 macp_version: if start_entry.macp_version.is_empty() {
329 "1.0".into()
330 } else {
331 start_entry.macp_version.clone()
332 },
333 mode: mode_name.to_string(),
334 message_type: "SessionStart".into(),
335 message_id: start_entry.message_id.clone(),
336 session_id: session_id.into(),
337 sender: start_entry.sender.clone(),
338 timestamp_unix_ms: if start_entry.timestamp_unix_ms != 0 {
339 start_entry.timestamp_unix_ms
340 } else {
341 start_entry.received_at_ms
342 },
343 payload: start_entry.raw_payload.clone(),
344 };
345
346 let mut session = Session::builder(session_id, mode_name, start_entry.sender.clone())
347 .semantics_rev(start_entry.semantics_rev)
350 .max_suspend_ms(start_entry.bound_max_suspend_ms.unwrap_or(0))
353 .ttl_expiry(ttl_expiry)
354 .ttl_ms(ttl_ms)
355 .started_at_unix_ms(started_at_unix_ms)
356 .participants(start_payload.participants.clone())
357 .intent(start_payload.intent.clone())
358 .mode_version(
364 start_entry
365 .bound_mode_version
366 .clone()
367 .unwrap_or_else(|| start_payload.mode_version.clone()),
368 )
369 .configuration_version(start_payload.configuration_version.clone())
370 .policy_version(start_payload.policy_version.clone())
371 .context_id(start_payload.context_id.clone())
372 .extensions(start_payload.extensions.clone())
373 .roots(start_payload.roots.clone())
374 .policy_definition(if !start_payload.policy_version.is_empty() {
375 policy_registry.and_then(|pr| pr.resolve(&start_payload.policy_version).ok())
376 } else {
377 None
378 })
379 .build();
380
381 let response = mode.on_session_start(&session, &env)?;
383 session.seen_message_ids.insert(env.message_id.clone());
384 session.apply_mode_response(response);
385
386 for entry in log_entries.iter().skip(1) {
388 replay_entry(&mut session, session_id, entry, &mode)?;
389 }
390
391 Ok(session)
392}
393
394#[cfg(test)]
395mod tests {
396 use super::*;
397 use crate::decision_pb::ProposalPayload;
398 use crate::decision_pb::VotePayload;
399 use crate::log_store::EntryKind;
400 use crate::pb::{CommitmentPayload, SessionStartPayload};
401 use prost::Message;
402
403 fn make_registry() -> ModeRegistry {
404 ModeRegistry::build_default(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator))
405 }
406
407 fn start_payload_bytes() -> Vec<u8> {
408 SessionStartPayload {
409 intent: "test".into(),
410 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
411 mode_version: "1.0.0".into(),
412 configuration_version: "cfg-1".into(),
413 policy_version: "policy-1".into(),
414 ttl_ms: 60_000,
415 context_id: String::new(),
416 extensions: std::collections::HashMap::new(),
417 roots: vec![],
418 max_suspend_ms: 0,
419 }
420 .encode_to_vec()
421 }
422
423 fn incoming_entry(
424 message_id: &str,
425 message_type: &str,
426 sender: &str,
427 payload: Vec<u8>,
428 received_at_ms: i64,
429 ) -> LogEntry {
430 LogEntry {
431 message_id: message_id.into(),
432 received_at_ms,
433 sender: sender.into(),
434 message_type: message_type.into(),
435 raw_payload: payload,
436 entry_kind: EntryKind::Incoming,
437 session_id: "s1".into(),
438 mode: "macp.mode.decision.v1".into(),
439 macp_version: "1.0".into(),
440 timestamp_unix_ms: received_at_ms,
441 bound_mode_version: None,
442 semantics_rev: 0,
443 bound_max_suspend_ms: None,
444 compacted_incoming_ordinals: 0,
445 }
446 }
447
448 fn internal_entry(message_type: &str, received_at_ms: i64) -> LogEntry {
449 LogEntry {
450 message_id: String::new(),
451 received_at_ms,
452 sender: "_runtime".into(),
453 message_type: message_type.into(),
454 raw_payload: vec![],
455 entry_kind: EntryKind::Internal,
456 session_id: "s1".into(),
457 mode: "macp.mode.decision.v1".into(),
458 macp_version: "1.0".into(),
459 timestamp_unix_ms: received_at_ms,
460 bound_mode_version: None,
461 semantics_rev: 0,
462 bound_max_suspend_ms: None,
463 compacted_incoming_ordinals: 0,
464 }
465 }
466
467 #[test]
468 fn replay_rebuilds_decision_session() {
469 let registry = make_registry();
470 let proposal = ProposalPayload {
471 proposal_id: "p1".into(),
472 option: "deploy".into(),
473 rationale: "ready".into(),
474 supporting_data: vec![],
475 }
476 .encode_to_vec();
477 let vote = VotePayload {
478 proposal_id: "p1".into(),
479 vote: "approve".into(),
480 reason: "lgtm".into(),
481 }
482 .encode_to_vec();
483 let commitment = CommitmentPayload {
484 commitment_id: "c1".into(),
485 action: "decision.selected".into(),
486 authority_scope: "payments".into(),
487 reason: "bound".into(),
488 mode_version: "1.0.0".into(),
489 policy_version: "policy-1".into(),
490 configuration_version: "cfg-1".into(),
491 outcome_positive: true,
492 supersedes: None,
493 }
494 .encode_to_vec();
495
496 let entries = vec![
497 incoming_entry(
498 "m1",
499 "SessionStart",
500 "agent://orchestrator",
501 start_payload_bytes(),
502 1000,
503 ),
504 incoming_entry("m2", "Proposal", "agent://orchestrator", proposal, 2000),
505 incoming_entry("m3", "Vote", "agent://fraud", vote, 3000),
506 incoming_entry("m4", "Commitment", "agent://orchestrator", commitment, 4000),
507 ];
508
509 let session = replay_session("s1", &entries, ®istry, None).unwrap();
510 assert_eq!(session.state, SessionState::Resolved);
511 assert_eq!(session.session_id, "s1");
512 assert!(session.seen_message_ids.contains("m1"));
513 assert!(session.seen_message_ids.contains("m2"));
514 assert!(session.seen_message_ids.contains("m3"));
515 assert!(session.seen_message_ids.contains("m4"));
516 assert!(session.resolution.is_some());
517 }
518
519 #[test]
520 fn replay_preserves_original_ttl() {
521 let registry = make_registry();
522 let original_time = 1_700_000_000_000i64;
523 let entries = vec![incoming_entry(
524 "m1",
525 "SessionStart",
526 "agent://orchestrator",
527 start_payload_bytes(),
528 original_time,
529 )];
530
531 let session = replay_session("s1", &entries, ®istry, None).unwrap();
532 assert_eq!(session.started_at_unix_ms, original_time);
533 assert_eq!(session.ttl_expiry, original_time + 60_000);
534 assert_eq!(session.ttl_ms, 60_000);
535 }
536
537 #[test]
538 fn replay_handles_ttl_expired() {
539 let registry = make_registry();
540 let entries = vec![
541 incoming_entry(
542 "m1",
543 "SessionStart",
544 "agent://orchestrator",
545 start_payload_bytes(),
546 1000,
547 ),
548 internal_entry("TtlExpired", 61001),
549 ];
550
551 let session = replay_session("s1", &entries, ®istry, None).unwrap();
552 assert_eq!(session.state, SessionState::Expired);
553 }
554
555 #[test]
556 fn replay_handles_session_cancel() {
557 let registry = make_registry();
558 let entries = vec![
559 incoming_entry(
560 "m1",
561 "SessionStart",
562 "agent://orchestrator",
563 start_payload_bytes(),
564 1000,
565 ),
566 internal_entry("SessionCancel", 5000),
567 ];
568
569 let session = replay_session("s1", &entries, ®istry, None).unwrap();
570 assert_eq!(session.state, SessionState::Cancelled);
572 }
573
574 #[test]
575 fn replay_fails_when_accepted_history_no_longer_applies() {
576 let registry = make_registry();
577 let vote = VotePayload {
578 proposal_id: "p1".into(),
579 vote: "approve".into(),
580 reason: String::new(),
581 }
582 .encode_to_vec();
583 let entries = vec![
584 incoming_entry(
585 "m1",
586 "SessionStart",
587 "agent://orchestrator",
588 start_payload_bytes(),
589 1000,
590 ),
591 incoming_entry("m2", "Vote", "agent://fraud", vote, 2000),
592 ];
593
594 let err = replay_session("s1", &entries, ®istry, None).unwrap_err();
595 let msg = err.to_string();
598 assert!(
599 msg == "InvalidTransition" || msg == "InvalidPayload" || msg == "Forbidden",
600 "unexpected error: {msg}"
601 );
602 }
603
604 #[test]
605 fn replay_empty_log_returns_error() {
606 let registry = make_registry();
607 let result = replay_session("s1", &[], ®istry, None);
608 assert!(result.is_err());
609 }
610
611 #[test]
612 fn backward_compat_old_log_entry_without_new_fields() {
613 let json = r#"{"message_id":"m1","received_at_ms":1000,"sender":"test","message_type":"Message","raw_payload":[],"entry_kind":"Incoming"}"#;
615 let entry: LogEntry = serde_json::from_str(json).unwrap();
616 assert_eq!(entry.session_id, "");
617 assert_eq!(entry.mode, "");
618 assert_eq!(entry.macp_version, "");
619 }
620
621 #[test]
632 fn replay_from_checkpoint_restores_state() {
633 use crate::registry::PersistedSession;
634
635 let registry = make_registry();
636 let start_payload = SessionStartPayload {
637 intent: "test".into(),
638 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
639 mode_version: "1.0.0".into(),
640 configuration_version: "cfg-1".into(),
641 policy_version: String::new(),
642 ttl_ms: 60_000,
643 context_id: String::new(),
644 extensions: std::collections::HashMap::new(),
645 roots: vec![],
646 max_suspend_ms: 0,
647 }
648 .encode_to_vec();
649
650 let proposal = ProposalPayload {
652 proposal_id: "p1".into(),
653 option: "deploy".into(),
654 rationale: "ready".into(),
655 supporting_data: vec![],
656 }
657 .encode_to_vec();
658
659 let full_entries = vec![
660 incoming_entry(
661 "m1",
662 "SessionStart",
663 "agent://orchestrator",
664 start_payload,
665 1000,
666 ),
667 incoming_entry(
668 "m2",
669 "Proposal",
670 "agent://orchestrator",
671 proposal.clone(),
672 2000,
673 ),
674 ];
675 let full_session = replay_session("s1", &full_entries, ®istry, None).unwrap();
676
677 let mut persisted = PersistedSession::from(&full_session);
679 persisted.intent = "restored-from-checkpoint".into();
683 let checkpoint_payload = serde_json::to_vec(&persisted).unwrap();
684 let checkpoint = LogEntry {
685 message_id: String::new(),
686 received_at_ms: 3000,
687 sender: "_runtime".into(),
688 message_type: "Checkpoint".into(),
689 raw_payload: checkpoint_payload,
690 entry_kind: EntryKind::Checkpoint,
691 session_id: "s1".into(),
692 mode: "macp.mode.decision.v1".into(),
693 macp_version: "1.0".into(),
694 timestamp_unix_ms: 3000,
695 bound_mode_version: None,
696 semantics_rev: 0,
697 bound_max_suspend_ms: None,
698 compacted_incoming_ordinals: 0,
699 };
700
701 let vote = VotePayload {
703 proposal_id: "p1".into(),
704 vote: "approve".into(),
705 reason: "lgtm".into(),
706 }
707 .encode_to_vec();
708
709 let entries_with_checkpoint = vec![
711 full_entries[0].clone(),
712 full_entries[1].clone(),
713 checkpoint,
714 incoming_entry("m3", "Vote", "agent://fraud", vote, 4000),
715 ];
716
717 let session = replay_session("s1", &entries_with_checkpoint, ®istry, None).unwrap();
718 assert_eq!(session.state, SessionState::Open);
719 assert_eq!(
720 session.intent, "restored-from-checkpoint",
721 "the checkpoint fast path must have been taken, else this test \
722 proves nothing about the checkpoint"
723 );
724 assert!(session.seen_message_ids.contains("m1"));
726 assert!(session.seen_message_ids.contains("m2"));
727 assert!(session.seen_message_ids.contains("m3"));
728 }
729
730 #[test]
743 fn replay_from_checkpoint_restores_suspension_intervals() {
744 use crate::registry::PersistedSession;
745
746 let registry = make_registry();
747 let start_payload = SessionStartPayload {
748 intent: "test".into(),
749 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
750 mode_version: "1.0.0".into(),
751 configuration_version: "cfg-1".into(),
752 policy_version: String::new(),
753 ttl_ms: 60_000,
754 context_id: String::new(),
755 extensions: std::collections::HashMap::new(),
756 roots: vec![],
757 max_suspend_ms: 0,
758 }
759 .encode_to_vec();
760
761 let prefix = vec![
762 incoming_entry(
763 "m1",
764 "SessionStart",
765 "agent://orchestrator",
766 start_payload,
767 1_000,
768 ),
769 internal_entry("SessionSuspend", 1_050),
770 internal_entry("SessionResume", 1_300),
771 ];
772 let before_checkpoint = replay_session("s1", &prefix, ®istry, None).unwrap();
773 assert_eq!(before_checkpoint.suspension_intervals, vec![(1_050, 1_300)]);
774
775 let mut persisted = PersistedSession::from(&before_checkpoint);
776 persisted.intent = "restored-from-checkpoint".into();
779 let checkpoint_payload = serde_json::to_vec(&persisted).unwrap();
782 let checkpoint = LogEntry {
783 message_id: String::new(),
784 received_at_ms: 1_400,
785 sender: "_runtime".into(),
786 message_type: "Checkpoint".into(),
787 raw_payload: checkpoint_payload,
788 entry_kind: EntryKind::Checkpoint,
789 session_id: "s1".into(),
790 mode: "macp.mode.decision.v1".into(),
791 macp_version: "1.0".into(),
792 timestamp_unix_ms: 1_400,
793 bound_mode_version: None,
794 semantics_rev: 0,
795 bound_max_suspend_ms: None,
796 compacted_incoming_ordinals: 0,
797 };
798
799 let entries = vec![
802 prefix[0].clone(),
803 prefix[1].clone(),
804 prefix[2].clone(),
805 checkpoint,
806 internal_entry("SessionSuspend", 1_500),
807 internal_entry("SessionResume", 1_600),
808 ];
809 let session = replay_session("s1", &entries, ®istry, None).unwrap();
810 assert_eq!(session.state, SessionState::Open);
811 assert_eq!(
812 session.intent, "restored-from-checkpoint",
813 "the checkpoint fast path must have been taken, else this test \
814 proves nothing about the snapshot round-trip"
815 );
816 assert_eq!(
817 session.suspension_intervals,
818 vec![(1_050, 1_300), (1_500, 1_600)],
819 "the pre-checkpoint pause must come from the snapshot and the \
820 post-checkpoint pause from the replayed tail"
821 );
822 }
823
824 #[test]
825 fn replay_without_checkpoint_still_works() {
826 let registry = make_registry();
828 let entries = vec![incoming_entry(
829 "m1",
830 "SessionStart",
831 "agent://orchestrator",
832 start_payload_bytes(),
833 1000,
834 )];
835 let session = replay_session("s1", &entries, ®istry, None).unwrap();
836 assert_eq!(session.state, SessionState::Open);
837 assert!(session.seen_message_ids.contains("m1"));
838 }
839
840 fn ext_registry_with_dyn_mode(version: &str) -> ModeRegistry {
841 let registry = make_registry();
842 registry
843 .register_extension(macp_pb::pb::ModeDescriptor {
844 mode: "ext.dyn.v1".into(),
845 mode_version: version.into(),
846 message_types: vec!["SessionStart".into(), "Commitment".into()],
847 terminal_message_types: vec!["Commitment".into()],
848 ..Default::default()
849 })
850 .unwrap();
851 registry
852 }
853
854 fn ext_start_entry(bound_mode_version: Option<String>) -> LogEntry {
855 let payload = SessionStartPayload {
857 participants: vec!["alice".into()],
858 configuration_version: "cfg-1".into(),
859 ttl_ms: 60_000,
860 ..Default::default()
861 }
862 .encode_to_vec();
863 LogEntry {
864 message_id: "m1".into(),
865 received_at_ms: 1000,
866 sender: "alice".into(),
867 message_type: "SessionStart".into(),
868 raw_payload: payload,
869 entry_kind: EntryKind::Incoming,
870 session_id: "s1".into(),
871 mode: "ext.dyn.v1".into(),
872 macp_version: "1.0".into(),
873 timestamp_unix_ms: 1000,
874 bound_mode_version,
875 semantics_rev: 0,
876 bound_max_suspend_ms: None,
877 compacted_incoming_ordinals: 0,
878 }
879 }
880
881 #[test]
885 fn replay_uses_recorded_mode_version_binding() {
886 let registry = ext_registry_with_dyn_mode("9.9.9");
887 let entries = vec![ext_start_entry(Some("2.5.0".into()))];
888 let session = replay_session("s1", &entries, ®istry, None).unwrap();
889 assert_eq!(session.mode_version, "2.5.0");
890 }
891
892 #[test]
896 fn replay_legacy_entry_without_binding_keeps_empty_version() {
897 let registry = ext_registry_with_dyn_mode("9.9.9");
898 let entries = vec![ext_start_entry(None)];
899 let session = replay_session("s1", &entries, ®istry, None).unwrap();
900 assert_eq!(session.mode_version, "");
901 }
902
903 #[test]
906 fn legacy_log_entry_json_without_binding_field_deserializes() {
907 let json = serde_json::json!({
908 "message_id": "m1",
909 "received_at_ms": 1000,
910 "sender": "alice",
911 "message_type": "SessionStart",
912 "raw_payload": [],
913 "entry_kind": "Incoming",
914 "session_id": "s1",
915 "mode": "ext.dyn.v1",
916 "macp_version": "1.0",
917 "timestamp_unix_ms": 1000
918 });
919 let entry: LogEntry = serde_json::from_value(json).unwrap();
920 assert_eq!(entry.bound_mode_version, None);
921 assert_eq!(entry.semantics_rev, 0);
924 assert_eq!(entry.bound_max_suspend_ms, None);
926 }
927
928 #[test]
931 fn replay_uses_recorded_max_suspend_cap() {
932 let registry = ext_registry_with_dyn_mode("1.0.0");
933 let mut entry = ext_start_entry(Some("1.0.0".into()));
934 entry.bound_max_suspend_ms = Some(1234);
935 let session = replay_session("s1", &[entry], ®istry, None).unwrap();
936 assert_eq!(session.max_suspend_ms, 1234);
937 assert_eq!(session.effective_max_suspend_ms(), 1234);
938 }
939
940 #[test]
943 fn replay_legacy_entry_keeps_default_cap_semantics() {
944 let registry = ext_registry_with_dyn_mode("1.0.0");
945 let entries = vec![ext_start_entry(None)];
946 let session = replay_session("s1", &entries, ®istry, None).unwrap();
947 assert_eq!(session.max_suspend_ms, 0);
948 assert_eq!(
949 session.effective_max_suspend_ms(),
950 macp_core::session::MAX_SUSPEND_MS
951 );
952 }
953
954 #[test]
959 fn replay_preserves_recorded_semantics_rev() {
960 let registry = ext_registry_with_dyn_mode("1.0.0");
961 let entries = vec![ext_start_entry(Some("1.0.0".into()))]; let session = replay_session("s1", &entries, ®istry, None).unwrap();
963 assert_eq!(session.semantics_rev, 0);
964 assert_ne!(
965 session.semantics_rev,
966 macp_core::session::CURRENT_SEMANTICS_REV
967 );
968 }
969
970 #[test]
971 fn replay_consistency_flags_state_and_dedup_divergence() {
972 let a = Session::builder("s1", "macp.mode.decision.v1", "agent://a")
973 .mode_version("1.0.0")
974 .configuration_version("cfg-1")
975 .build();
976 assert_eq!(validate_replay_consistency("s1", &a, &a.clone()), 0);
978
979 let mut b = a.clone();
981 b.state = SessionState::Resolved;
982 b.seen_message_ids.insert("m1".into());
983 assert_eq!(validate_replay_consistency("s1", &a, &b), 2);
984
985 let mut c = a.clone();
987 c.mode_state = vec![7, 7, 7];
988 assert_eq!(validate_replay_consistency("s1", &a, &c), 1);
989
990 let mut d = a.clone();
994 d.accumulated_suspended_ms = 5_000;
995 assert_eq!(validate_replay_consistency("s1", &a, &d), 1);
996 d.suspended_at_ms = Some(1_000);
997 assert_eq!(validate_replay_consistency("s1", &a, &d), 2);
998 d.suspension_intervals = vec![(1_000, 6_000)];
1000 assert_eq!(validate_replay_consistency("s1", &a, &d), 3);
1001
1002 let mut e = b.clone();
1005 e.mode_state = vec![7, 7, 7];
1006 e.accumulated_suspended_ms = 5_000;
1007 e.suspended_at_ms = Some(1_000);
1008 e.suspension_intervals = vec![(1_000, 6_000)];
1009 assert_eq!(validate_replay_consistency("s1", &a, &e), 6);
1010 }
1011
1012 const HANDOFF_TIMEOUT_MS: i64 = 100;
1026
1027 fn handoff_policy_registry() -> PolicyRegistry {
1028 let registry = PolicyRegistry::new();
1029 registry
1030 .register(macp_core::policy::PolicyDefinition {
1031 policy_id: "handoff-auto-accept".into(),
1032 mode: "macp.mode.handoff.v1".into(),
1033 description: "implicit accept after 100ms".into(),
1034 rules: serde_json::json!({
1035 "acceptance": { "implicit_accept_timeout_ms": HANDOFF_TIMEOUT_MS },
1036 "commitment": { "authority": "initiator_only" }
1037 }),
1038 schema_version: 1,
1039 })
1040 .unwrap();
1041 registry
1042 }
1043
1044 fn handoff_entry(
1045 message_id: &str,
1046 message_type: &str,
1047 payload: Vec<u8>,
1048 envelope_ms: i64,
1049 received_ms: i64,
1050 ) -> LogEntry {
1051 LogEntry {
1052 message_id: message_id.into(),
1053 received_at_ms: received_ms,
1054 sender: "alice".into(),
1055 message_type: message_type.into(),
1056 raw_payload: payload,
1057 entry_kind: EntryKind::Incoming,
1058 session_id: "s1".into(),
1059 mode: "macp.mode.handoff.v1".into(),
1060 macp_version: "1.0".into(),
1061 timestamp_unix_ms: envelope_ms,
1062 bound_mode_version: None,
1063 semantics_rev: 0,
1064 bound_max_suspend_ms: None,
1065 compacted_incoming_ordinals: 0,
1066 }
1067 }
1068
1069 fn implicit_accept_entry(deadline_ms: i64) -> LogEntry {
1080 let payload = crate::handoff_pb::HandoffAcceptPayload {
1081 handoff_id: "h1".into(),
1082 accepted_by: "bob".into(),
1083 reason: "implicit accept (timeout)".into(),
1084 implicit: true,
1085 }
1086 .encode_to_vec();
1087 let mut entry = handoff_entry(
1088 "implicit-accept:h1",
1089 "HandoffAccept",
1090 payload,
1091 deadline_ms,
1092 deadline_ms,
1093 );
1094 entry.sender = "bob".into();
1095 entry
1096 }
1097
1098 fn handoff_history(
1102 semantics_rev: u32,
1103 commit_envelope_ms: i64,
1104 commit_received_ms: i64,
1105 ) -> Vec<LogEntry> {
1106 let start_payload = SessionStartPayload {
1107 intent: "escalate".into(),
1108 participants: vec!["alice".into(), "bob".into()],
1109 mode_version: "1.0.0".into(),
1110 configuration_version: "cfg-1".into(),
1111 policy_version: "handoff-auto-accept".into(),
1112 ttl_ms: 60_000,
1113 context_id: String::new(),
1114 extensions: std::collections::HashMap::new(),
1115 roots: vec![],
1116 max_suspend_ms: 0,
1117 }
1118 .encode_to_vec();
1119 let offer = crate::handoff_pb::HandoffOfferPayload {
1120 handoff_id: "h1".into(),
1121 target_participant: "bob".into(),
1122 scope: "support".into(),
1123 reason: "escalate".into(),
1124 }
1125 .encode_to_vec();
1126 let commitment = CommitmentPayload {
1127 commitment_id: "c1".into(),
1128 action: "handoff.accepted".into(),
1129 authority_scope: "support".into(),
1130 reason: "bound".into(),
1131 mode_version: "1.0.0".into(),
1132 policy_version: "handoff-auto-accept".into(),
1133 configuration_version: "cfg-1".into(),
1134 outcome_positive: true,
1135 supersedes: None,
1136 }
1137 .encode_to_vec();
1138
1139 let mut start = handoff_entry("m1", "SessionStart", start_payload, 1_000, 1_000);
1140 start.semantics_rev = semantics_rev;
1141 vec![
1142 start,
1143 handoff_entry("m2", "HandoffOffer", offer, 1_000, 1_000),
1146 handoff_entry(
1147 "m3",
1148 "Commitment",
1149 commitment,
1150 commit_envelope_ms,
1151 commit_received_ms,
1152 ),
1153 ]
1154 }
1155
1156 fn assert_implicitly_accepted(session: &Session) {
1159 assert_eq!(session.state, SessionState::Resolved);
1160 let state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
1161 let offer = &state["offers"]["h1"];
1162 assert_eq!(offer["disposition"], "Accepted");
1163 assert_eq!(offer["accepted_by"], "bob");
1164 assert_eq!(offer["outcome_reason"], "implicit accept (timeout)");
1165 }
1166
1167 #[test]
1172 fn legacy_rev0_handoff_history_replays_under_envelope_clock() {
1173 let registry = make_registry();
1174 let policies = handoff_policy_registry();
1175 let entries = handoff_history(0, 1_300, 1_050);
1176
1177 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1178 assert_eq!(session.semantics_rev, 0);
1179 assert_implicitly_accepted(&session);
1180
1181 for rev in [1, macp_core::session::CURRENT_SEMANTICS_REV] {
1186 let mut newer = entries.clone();
1187 newer[0].semantics_rev = rev;
1188 assert!(
1189 replay_session("s1", &newer, ®istry, Some(&policies)).is_err(),
1190 "rev {rev} must not reproduce the rev-0 outcome"
1191 );
1192 }
1193 }
1194
1195 #[test]
1199 fn legacy_rev1_handoff_history_replays_under_acceptance_clock() {
1200 let registry = make_registry();
1201 let policies = handoff_policy_registry();
1202 let entries = handoff_history(1, 1_050, 1_300);
1203
1204 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1205 assert_eq!(session.semantics_rev, 1);
1206 assert_implicitly_accepted(&session);
1207
1208 let mut legacy = entries.clone();
1210 legacy[0].semantics_rev = 0;
1211 assert!(replay_session("s1", &legacy, ®istry, Some(&policies)).is_err());
1212 }
1213
1214 #[test]
1233 fn current_rev_handoff_history_replays_identically_to_rev1() {
1234 let registry = make_registry();
1235 let policies = handoff_policy_registry();
1236
1237 let rev1 = replay_session(
1238 "s1",
1239 &handoff_history(1, 1_050, 1_300),
1240 ®istry,
1241 Some(&policies),
1242 )
1243 .unwrap();
1244
1245 let mut current_entries =
1249 handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1250 let commitment = current_entries.pop().expect("commitment is last");
1251 current_entries.push(implicit_accept_entry(1_100));
1252 current_entries.push(commitment);
1253
1254 let current = replay_session("s1", ¤t_entries, ®istry, Some(&policies)).unwrap();
1255
1256 assert_implicitly_accepted(¤t);
1257 assert_eq!(current.state, rev1.state);
1258 assert_eq!(current.mode_state, rev1.mode_state);
1259 assert_eq!(current.resolution, rev1.resolution);
1260
1261 assert_implicitly_accepted(&rev1);
1264 assert!(!rev1.seen_message_ids.contains("implicit-accept:h1"));
1265 assert!(current.seen_message_ids.contains("implicit-accept:h1"));
1268 }
1269
1270 #[test]
1282 fn rev2_commitment_without_synthetic_entry_fails_replay() {
1283 let registry = make_registry();
1284 let policies = handoff_policy_registry();
1285
1286 let entries = handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1290 assert!(
1291 replay_session("s1", &entries, ®istry, Some(&policies)).is_err(),
1292 "rev 2 must not infer an accept the history does not record"
1293 );
1294
1295 let mut legacy = entries.clone();
1298 legacy[0].semantics_rev = 1;
1299 let session = replay_session("s1", &legacy, ®istry, Some(&policies))
1300 .expect("rev 1 keeps the interim in-Commitment implicit accept");
1301 assert_implicitly_accepted(&session);
1302
1303 let mut with_synthetic = entries.clone();
1306 let commitment = with_synthetic.pop().expect("commitment is last");
1307 with_synthetic.push(implicit_accept_entry(1_100));
1308 with_synthetic.push(commitment);
1309 let session = replay_session("s1", &with_synthetic, ®istry, Some(&policies))
1310 .expect("rev 2 resolves once the synthetic entry is in history");
1311 assert_implicitly_accepted(&session);
1312 }
1313
1314 fn handoff_history_with_suspension(
1319 semantics_rev: u32,
1320 suspend_at_ms: i64,
1321 resume_at_ms: i64,
1322 commit_ms: i64,
1323 ) -> Vec<LogEntry> {
1324 let mut entries = handoff_history(semantics_rev, commit_ms, commit_ms);
1327 let commit = entries.pop().expect("commitment is the last entry");
1328 entries.push(internal_entry("SessionSuspend", suspend_at_ms));
1329 entries.push(internal_entry("SessionResume", resume_at_ms));
1330 entries.push(commit);
1331 entries
1332 }
1333
1334 #[test]
1344 fn legacy_rev1_handoff_history_with_suspension_still_implicitly_accepts() {
1345 let registry = make_registry();
1346 let policies = handoff_policy_registry();
1347 let entries = handoff_history_with_suspension(1, 1_050, 1_300, 1_300);
1348
1349 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1350 assert_eq!(session.semantics_rev, 1);
1351 assert_eq!(session.accumulated_suspended_ms, 250);
1352 assert_implicitly_accepted(&session);
1353
1354 let mut rev2 = entries.clone();
1358 rev2[0].semantics_rev = macp_core::session::CURRENT_SEMANTICS_REV;
1359 assert!(
1360 replay_session("s1", &rev2, ®istry, Some(&policies)).is_err(),
1361 "rev 2 must not reproduce the rev-1 outcome"
1362 );
1363 }
1364
1365 #[test]
1380 fn rev2_handoff_history_implicitly_accepts_on_unsuspended_time() {
1381 let registry = make_registry();
1382 let policies = handoff_policy_registry();
1383 let mut entries = handoff_history_with_suspension(
1384 macp_core::session::CURRENT_SEMANTICS_REV,
1385 1_050,
1386 1_300,
1387 1_450,
1388 );
1389 let commitment = entries.pop().expect("commitment is last");
1390 entries.push(implicit_accept_entry(1_350));
1391 entries.push(commitment);
1392
1393 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1394 assert_eq!(session.accumulated_suspended_ms, 250);
1395 assert_implicitly_accepted(&session);
1396 assert!(session.seen_message_ids.contains("implicit-accept:h1"));
1397 }
1398
1399 fn handoff_history_with_two_suspensions(semantics_rev: u32, commit_ms: i64) -> Vec<LogEntry> {
1424 let mut entries = handoff_history(semantics_rev, commit_ms, commit_ms);
1425 let commit = entries.pop().expect("commitment is the last entry");
1426 entries.push(internal_entry("SessionSuspend", 1_050));
1427 entries.push(internal_entry("SessionResume", 1_300));
1428 entries.push(internal_entry("SessionSuspend", 1_330));
1429 entries.push(internal_entry("SessionResume", 1_500));
1430 entries.push(commit);
1431 entries
1432 }
1433
1434 #[test]
1446 fn rev2_handoff_history_subtracts_every_suspension_pair() {
1447 let registry = make_registry();
1448 let policies = handoff_policy_registry();
1449
1450 let rev1 = handoff_history_with_two_suspensions(1, 1_510);
1451 let session = replay_session("s1", &rev1, ®istry, Some(&policies)).unwrap();
1452 assert_eq!(session.semantics_rev, 1);
1453 assert_eq!(session.accumulated_suspended_ms, 420);
1454 assert_implicitly_accepted(&session);
1455
1456 let rev2 =
1457 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_510);
1458 assert!(
1459 replay_session("s1", &rev2, ®istry, Some(&policies)).is_err(),
1460 "rev 2 must subtract both pauses (90ms unsuspended < 100ms timeout)"
1461 );
1462 }
1463
1464 #[test]
1472 fn rev2_handoff_history_accepts_on_unsuspended_time_across_two_pauses() {
1473 let registry = make_registry();
1474 let policies = handoff_policy_registry();
1475 let mut entries =
1476 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_600);
1477 let commitment = entries.pop().expect("commitment is last");
1478 entries.push(implicit_accept_entry(1_520));
1479 entries.push(commitment);
1480
1481 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1482 assert_eq!(session.accumulated_suspended_ms, 420);
1483 assert_implicitly_accepted(&session);
1484 }
1485
1486 #[test]
1495 fn replay_rebuilds_suspension_intervals_from_the_log() {
1496 let registry = make_registry();
1497 let policies = handoff_policy_registry();
1498
1499 let rev1 = handoff_history_with_two_suspensions(1, 1_510);
1500 let session = replay_session("s1", &rev1, ®istry, Some(&policies)).unwrap();
1501 assert_eq!(
1502 session.suspension_intervals,
1503 vec![(1_050, 1_300), (1_330, 1_500)]
1504 );
1505
1506 let mut rev2 =
1507 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_600);
1508 let commitment = rev2.pop().expect("commitment is last");
1511 rev2.push(implicit_accept_entry(1_520));
1512 rev2.push(commitment);
1513 let session = replay_session("s1", &rev2, ®istry, Some(&policies)).unwrap();
1514 assert_eq!(
1515 session.suspension_intervals,
1516 vec![(1_050, 1_300), (1_330, 1_500)]
1517 );
1518 assert_eq!(session.unsuspended_deadline(1_000, 100), 1_520);
1523 }
1524
1525 fn handoff_history_with_reserved_message_ids(semantics_rev: u32) -> Vec<LogEntry> {
1547 let mut entries = handoff_history(semantics_rev, 1_050, 1_300);
1550 entries[0].message_id = "implicit-accept:squatted-at-start".into();
1551 entries[2].message_id = squatted_commitment_id(semantics_rev).into();
1552 if semantics_rev >= 2 {
1553 let commitment = entries.pop().expect("commitment is last");
1561 entries.push(implicit_accept_entry(1_100));
1562 entries.push(commitment);
1563 }
1564 entries
1565 }
1566
1567 fn squatted_commitment_id(semantics_rev: u32) -> &'static str {
1570 if semantics_rev >= 2 {
1571 "implicit-accept:squatted-at-commit"
1572 } else {
1573 "implicit-accept:h1"
1574 }
1575 }
1576
1577 #[test]
1589 fn reserved_prefix_entry_replays_at_every_rev() {
1590 let registry = make_registry();
1591 let policies = handoff_policy_registry();
1592
1593 for rev in [1, macp_core::session::CURRENT_SEMANTICS_REV] {
1594 let entries = handoff_history_with_reserved_message_ids(rev);
1595 let session = replay_session("s1", &entries, ®istry, Some(&policies))
1596 .unwrap_or_else(|e| panic!("rev {rev} must replay reserved ids, got {e}"));
1597 assert_eq!(session.semantics_rev, rev);
1598 assert_implicitly_accepted(&session);
1599 assert!(session
1602 .seen_message_ids
1603 .contains("implicit-accept:squatted-at-start"));
1604 assert!(session
1605 .seen_message_ids
1606 .contains(squatted_commitment_id(rev)));
1607 }
1608
1609 let mode = crate::mode::handoff::HandoffMode::new(std::sync::Arc::new(
1613 macp_policy::DefaultPolicyEvaluator,
1614 ));
1615 let entries =
1616 handoff_history_with_reserved_message_ids(macp_core::session::CURRENT_SEMANTICS_REV);
1617 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1618 for message_id in ["implicit-accept:squatted-at-start", "implicit-accept:h1"] {
1619 let env = Envelope {
1620 macp_version: "1.0".into(),
1621 mode: session.mode.clone(),
1622 message_type: "HandoffContext".into(),
1623 message_id: message_id.into(),
1624 session_id: "s1".into(),
1625 sender: "alice".into(),
1626 timestamp_unix_ms: 1_000,
1627 payload: vec![],
1628 };
1629 assert!(matches!(
1630 crate::mode::Mode::validate_client_envelope(&mode, &session, &env).unwrap_err(),
1631 MacpError::InvalidEnvelope
1632 ));
1633 }
1634 }
1635
1636 #[test]
1660 fn synthetic_shaped_entry_reaches_dispatch_not_the_client_boundary() {
1661 let registry = make_registry();
1662 let policies = handoff_policy_registry();
1663
1664 let mut entries = handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1665 let commitment = entries.pop().expect("commitment is the last entry");
1671 entries.push(implicit_accept_entry(1_100));
1672 entries.push(commitment);
1673
1674 let session = replay_session("s1", &entries, ®istry, Some(&policies))
1675 .expect("the synthetic entry must replay through dispatch at rev >= 2");
1676 assert_implicitly_accepted(&session);
1677 assert!(session.seen_message_ids.contains("implicit-accept:h1"));
1680 }
1681}