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 macp_core::MACP_VERSION.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() {
145 "TtlExpired" => {
146 session.state = SessionState::Expired;
147 }
148 "SessionCancel" => {
151 let _ = session.cancel();
152 }
153 "SessionSuspend" => {
157 let at = if entry.received_at_ms != 0 {
158 entry.received_at_ms
159 } else {
160 entry.timestamp_unix_ms
161 };
162 let _ = session.suspend(at);
163 }
164 "SessionResume" => {
165 let at = if entry.received_at_ms != 0 {
166 entry.received_at_ms
167 } else {
168 entry.timestamp_unix_ms
169 };
170 let _ = session.resume(at);
171 }
172 other => {
184 tracing::warn!(
185 session_id,
186 message_type = other,
187 received_at_ms = entry.received_at_ms,
188 "replay: unrecognized internal log entry type; skipping"
189 );
190 }
191 },
192 EntryKind::Checkpoint => {
193 }
195 }
196 Ok(())
197}
198
199pub fn validate_replay_consistency(
215 session_id: &str,
216 replayed: &Session,
217 snapshot: &Session,
218) -> u32 {
219 let mut mismatches = 0u32;
220 if replayed.state != snapshot.state {
221 mismatches += 1;
222 tracing::warn!(
223 session_id,
224 replayed_state = ?replayed.state,
225 snapshot_state = ?snapshot.state,
226 "replay/snapshot state mismatch"
227 );
228 }
229 if replayed.seen_message_ids.len() != snapshot.seen_message_ids.len() {
230 mismatches += 1;
231 tracing::warn!(
232 session_id,
233 replayed_dedup = replayed.seen_message_ids.len(),
234 snapshot_dedup = snapshot.seen_message_ids.len(),
235 "replay/snapshot dedup count mismatch"
236 );
237 }
238 if replayed.participants != snapshot.participants {
239 mismatches += 1;
240 tracing::warn!(session_id, "replay/snapshot participants mismatch");
241 }
242 if replayed.mode_version != snapshot.mode_version
243 || replayed.configuration_version != snapshot.configuration_version
244 || replayed.policy_version != snapshot.policy_version
245 {
246 mismatches += 1;
247 tracing::warn!(
248 session_id,
249 "replay/snapshot bound-version mismatch (mode/configuration/policy)"
250 );
251 }
252 if replayed.mode_state != snapshot.mode_state {
256 mismatches += 1;
257 tracing::warn!(
258 session_id,
259 replayed_len = replayed.mode_state.len(),
260 snapshot_len = snapshot.mode_state.len(),
261 "replay/snapshot mode_state mismatch"
262 );
263 }
264 if replayed.accumulated_suspended_ms != snapshot.accumulated_suspended_ms {
268 mismatches += 1;
269 tracing::warn!(
270 session_id,
271 replayed_accumulated_suspended_ms = replayed.accumulated_suspended_ms,
272 snapshot_accumulated_suspended_ms = snapshot.accumulated_suspended_ms,
273 "replay/snapshot accumulated_suspended_ms mismatch"
274 );
275 }
276 if replayed.suspended_at_ms != snapshot.suspended_at_ms {
277 mismatches += 1;
278 tracing::warn!(
279 session_id,
280 replayed_suspended_at_ms = ?replayed.suspended_at_ms,
281 snapshot_suspended_at_ms = ?snapshot.suspended_at_ms,
282 "replay/snapshot suspended_at_ms mismatch"
283 );
284 }
285 if replayed.suspension_intervals != snapshot.suspension_intervals {
290 mismatches += 1;
291 tracing::warn!(
292 session_id,
293 replayed_suspension_cycles = replayed.suspension_intervals.len(),
294 snapshot_suspension_cycles = snapshot.suspension_intervals.len(),
295 "replay/snapshot suspension_intervals mismatch"
296 );
297 }
298 mismatches
299}
300
301fn replay_from_start(
303 session_id: &str,
304 log_entries: &[LogEntry],
305 registry: &ModeRegistry,
306 policy_registry: Option<&PolicyRegistry>,
307) -> Result<Session, MacpError> {
308 let start_entry = log_entries
310 .iter()
311 .find(|e| e.entry_kind == EntryKind::Incoming && e.message_type == "SessionStart")
312 .ok_or(MacpError::InvalidPayload)?;
313
314 let mode_name = if start_entry.mode.is_empty() {
316 return Err(MacpError::InvalidPayload);
319 } else {
320 &start_entry.mode
321 };
322
323 let mode = registry.get_mode(mode_name).ok_or(MacpError::UnknownMode)?;
324
325 let require_complete_start = registry.requires_strict_session_start(mode_name);
327 let start_payload = if start_entry.raw_payload.is_empty() && !require_complete_start {
328 crate::pb::SessionStartPayload::default()
329 } else {
330 parse_session_start_payload(&start_entry.raw_payload)?
331 };
332 if require_complete_start {
336 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
337 }
338
339 let ttl_ms = if !require_complete_start && start_payload.ttl_ms == 0 {
340 60_000i64
342 } else {
343 extract_ttl_ms(&start_payload)?
344 };
345
346 let started_at_unix_ms = start_entry.received_at_ms;
348 let ttl_expiry = started_at_unix_ms.saturating_add(ttl_ms);
349
350 let env = Envelope {
351 macp_version: if start_entry.macp_version.is_empty() {
352 macp_core::MACP_VERSION.into()
353 } else {
354 start_entry.macp_version.clone()
355 },
356 mode: mode_name.to_string(),
357 message_type: "SessionStart".into(),
358 message_id: start_entry.message_id.clone(),
359 session_id: session_id.into(),
360 sender: start_entry.sender.clone(),
361 timestamp_unix_ms: if start_entry.timestamp_unix_ms != 0 {
362 start_entry.timestamp_unix_ms
363 } else {
364 start_entry.received_at_ms
365 },
366 payload: start_entry.raw_payload.clone(),
367 };
368
369 let mut session = Session::builder(session_id, mode_name, start_entry.sender.clone())
370 .semantics_rev(start_entry.semantics_rev)
373 .max_suspend_ms(start_entry.bound_max_suspend_ms.unwrap_or(0))
376 .ttl_expiry(ttl_expiry)
377 .ttl_ms(ttl_ms)
378 .started_at_unix_ms(started_at_unix_ms)
379 .participants(start_payload.participants.clone())
380 .intent(start_payload.intent.clone())
381 .mode_version(
387 start_entry
388 .bound_mode_version
389 .clone()
390 .unwrap_or_else(|| start_payload.mode_version.clone()),
391 )
392 .configuration_version(start_payload.configuration_version.clone())
393 .policy_version(start_payload.policy_version.clone())
394 .context_id(start_payload.context_id.clone())
395 .extensions(start_payload.extensions.clone())
396 .roots(start_payload.roots.clone())
397 .policy_definition(if !start_payload.policy_version.is_empty() {
398 policy_registry.and_then(|pr| pr.resolve(&start_payload.policy_version).ok())
399 } else {
400 None
401 })
402 .build();
403
404 let response = mode.on_session_start(&session, &env)?;
406 session.seen_message_ids.insert(env.message_id.clone());
407 session.apply_mode_response(response);
408
409 for entry in log_entries.iter().skip(1) {
411 replay_entry(&mut session, session_id, entry, &mode)?;
412 }
413
414 Ok(session)
415}
416
417#[cfg(test)]
418mod tests {
419 use super::*;
420 use crate::decision_pb::ProposalPayload;
421 use crate::decision_pb::VotePayload;
422 use crate::log_store::EntryKind;
423 use crate::pb::{CommitmentPayload, SessionResumePayload, SessionStartPayload};
424 use prost::Message;
425
426 fn make_registry() -> ModeRegistry {
427 ModeRegistry::build_default(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator))
428 }
429
430 fn start_payload_bytes() -> Vec<u8> {
431 SessionStartPayload {
432 intent: "test".into(),
433 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
434 mode_version: "1.0.0".into(),
435 configuration_version: "cfg-1".into(),
436 policy_version: "policy-1".into(),
437 ttl_ms: 60_000,
438 context_id: String::new(),
439 extensions: std::collections::HashMap::new(),
440 roots: vec![],
441 max_suspend_ms: 0,
442 }
443 .encode_to_vec()
444 }
445
446 fn incoming_entry(
447 message_id: &str,
448 message_type: &str,
449 sender: &str,
450 payload: Vec<u8>,
451 received_at_ms: i64,
452 ) -> LogEntry {
453 LogEntry {
454 message_id: message_id.into(),
455 received_at_ms,
456 sender: sender.into(),
457 message_type: message_type.into(),
458 raw_payload: payload,
459 entry_kind: EntryKind::Incoming,
460 session_id: "s1".into(),
461 mode: "macp.mode.decision.v1".into(),
462 macp_version: "1.0".into(),
463 timestamp_unix_ms: received_at_ms,
464 bound_mode_version: None,
465 semantics_rev: 0,
466 bound_max_suspend_ms: None,
467 compacted_incoming_ordinals: 0,
468 }
469 }
470
471 fn internal_entry(message_type: &str, received_at_ms: i64) -> LogEntry {
472 LogEntry {
473 message_id: String::new(),
474 received_at_ms,
475 sender: "_runtime".into(),
476 message_type: message_type.into(),
477 raw_payload: vec![],
478 entry_kind: EntryKind::Internal,
479 session_id: "s1".into(),
480 mode: "macp.mode.decision.v1".into(),
481 macp_version: "1.0".into(),
482 timestamp_unix_ms: received_at_ms,
483 bound_mode_version: None,
484 semantics_rev: 0,
485 bound_max_suspend_ms: None,
486 compacted_incoming_ordinals: 0,
487 }
488 }
489
490 #[test]
491 fn replay_rebuilds_decision_session() {
492 let registry = make_registry();
493 let proposal = ProposalPayload {
494 proposal_id: "p1".into(),
495 option: "deploy".into(),
496 rationale: "ready".into(),
497 supporting_data: vec![],
498 }
499 .encode_to_vec();
500 let vote = VotePayload {
501 proposal_id: "p1".into(),
502 vote: "approve".into(),
503 reason: "lgtm".into(),
504 }
505 .encode_to_vec();
506 let commitment = CommitmentPayload {
507 commitment_id: "c1".into(),
508 action: "decision.selected".into(),
509 authority_scope: "payments".into(),
510 reason: "bound".into(),
511 mode_version: "1.0.0".into(),
512 policy_version: "policy-1".into(),
513 configuration_version: "cfg-1".into(),
514 outcome_positive: true,
515 supersedes: None,
516 }
517 .encode_to_vec();
518
519 let entries = vec![
520 incoming_entry(
521 "m1",
522 "SessionStart",
523 "agent://orchestrator",
524 start_payload_bytes(),
525 1000,
526 ),
527 incoming_entry("m2", "Proposal", "agent://orchestrator", proposal, 2000),
528 incoming_entry("m3", "Vote", "agent://fraud", vote, 3000),
529 incoming_entry("m4", "Commitment", "agent://orchestrator", commitment, 4000),
530 ];
531
532 let session = replay_session("s1", &entries, ®istry, None).unwrap();
533 assert_eq!(session.state, SessionState::Resolved);
534 assert_eq!(session.session_id, "s1");
535 assert!(session.seen_message_ids.contains("m1"));
536 assert!(session.seen_message_ids.contains("m2"));
537 assert!(session.seen_message_ids.contains("m3"));
538 assert!(session.seen_message_ids.contains("m4"));
539 assert!(session.resolution.is_some());
540 }
541
542 #[test]
543 fn replay_preserves_original_ttl() {
544 let registry = make_registry();
545 let original_time = 1_700_000_000_000i64;
546 let entries = vec![incoming_entry(
547 "m1",
548 "SessionStart",
549 "agent://orchestrator",
550 start_payload_bytes(),
551 original_time,
552 )];
553
554 let session = replay_session("s1", &entries, ®istry, None).unwrap();
555 assert_eq!(session.started_at_unix_ms, original_time);
556 assert_eq!(session.ttl_expiry, original_time + 60_000);
557 assert_eq!(session.ttl_ms, 60_000);
558 }
559
560 #[test]
561 fn replay_handles_ttl_expired() {
562 let registry = make_registry();
563 let entries = vec![
564 incoming_entry(
565 "m1",
566 "SessionStart",
567 "agent://orchestrator",
568 start_payload_bytes(),
569 1000,
570 ),
571 internal_entry("TtlExpired", 61001),
572 ];
573
574 let session = replay_session("s1", &entries, ®istry, None).unwrap();
575 assert_eq!(session.state, SessionState::Expired);
576 }
577
578 #[test]
579 fn replay_handles_session_cancel() {
580 let registry = make_registry();
581 let entries = vec![
582 incoming_entry(
583 "m1",
584 "SessionStart",
585 "agent://orchestrator",
586 start_payload_bytes(),
587 1000,
588 ),
589 internal_entry("SessionCancel", 5000),
590 ];
591
592 let session = replay_session("s1", &entries, ®istry, None).unwrap();
593 assert_eq!(session.state, SessionState::Cancelled);
595 }
596
597 #[test]
605 fn replay_warns_but_continues_on_unrecognized_internal_entry() {
606 let registry = make_registry();
607 let proposal = ProposalPayload {
608 proposal_id: "p1".into(),
609 option: "deploy".into(),
610 rationale: "after unrecognized entry".into(),
611 supporting_data: vec![],
612 }
613 .encode_to_vec();
614 let entries = vec![
615 incoming_entry(
616 "m1",
617 "SessionStart",
618 "agent://orchestrator",
619 start_payload_bytes(),
620 1000,
621 ),
622 internal_entry("SomeFutureRuntimeAnnotation", 2000),
623 incoming_entry("m2", "Proposal", "agent://orchestrator", proposal, 3000),
624 ];
625
626 let session = replay_session("s1", &entries, ®istry, None)
627 .expect("an unrecognized Internal entry must not fail replay");
628 assert_eq!(session.state, SessionState::Open);
629 assert!(
630 session.seen_message_ids.contains("m1") && session.seen_message_ids.contains("m2"),
631 "replay must continue past the unrecognized entry and apply the \
632 message that follows it, not merely avoid erroring"
633 );
634 }
635
636 #[test]
637 fn replay_fails_when_accepted_history_no_longer_applies() {
638 let registry = make_registry();
639 let vote = VotePayload {
640 proposal_id: "p1".into(),
641 vote: "approve".into(),
642 reason: String::new(),
643 }
644 .encode_to_vec();
645 let entries = vec![
646 incoming_entry(
647 "m1",
648 "SessionStart",
649 "agent://orchestrator",
650 start_payload_bytes(),
651 1000,
652 ),
653 incoming_entry("m2", "Vote", "agent://fraud", vote, 2000),
654 ];
655
656 let err = replay_session("s1", &entries, ®istry, None).unwrap_err();
657 let msg = err.to_string();
660 assert!(
661 msg == "InvalidTransition" || msg == "InvalidPayload" || msg == "Forbidden",
662 "unexpected error: {msg}"
663 );
664 }
665
666 #[test]
667 fn replay_empty_log_returns_error() {
668 let registry = make_registry();
669 let result = replay_session("s1", &[], ®istry, None);
670 assert!(result.is_err());
671 }
672
673 #[test]
674 fn backward_compat_old_log_entry_without_new_fields() {
675 let json = r#"{"message_id":"m1","received_at_ms":1000,"sender":"test","message_type":"Message","raw_payload":[],"entry_kind":"Incoming"}"#;
677 let entry: LogEntry = serde_json::from_str(json).unwrap();
678 assert_eq!(entry.session_id, "");
679 assert_eq!(entry.mode, "");
680 assert_eq!(entry.macp_version, "");
681 }
682
683 #[test]
694 fn replay_from_checkpoint_restores_state() {
695 use crate::registry::PersistedSession;
696
697 let registry = make_registry();
698 let start_payload = SessionStartPayload {
699 intent: "test".into(),
700 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
701 mode_version: "1.0.0".into(),
702 configuration_version: "cfg-1".into(),
703 policy_version: String::new(),
704 ttl_ms: 60_000,
705 context_id: String::new(),
706 extensions: std::collections::HashMap::new(),
707 roots: vec![],
708 max_suspend_ms: 0,
709 }
710 .encode_to_vec();
711
712 let proposal = ProposalPayload {
714 proposal_id: "p1".into(),
715 option: "deploy".into(),
716 rationale: "ready".into(),
717 supporting_data: vec![],
718 }
719 .encode_to_vec();
720
721 let full_entries = vec![
722 incoming_entry(
723 "m1",
724 "SessionStart",
725 "agent://orchestrator",
726 start_payload,
727 1000,
728 ),
729 incoming_entry(
730 "m2",
731 "Proposal",
732 "agent://orchestrator",
733 proposal.clone(),
734 2000,
735 ),
736 ];
737 let full_session = replay_session("s1", &full_entries, ®istry, None).unwrap();
738
739 let mut persisted = PersistedSession::from(&full_session);
741 persisted.intent = "restored-from-checkpoint".into();
745 let checkpoint_payload = serde_json::to_vec(&persisted).unwrap();
746 let checkpoint = LogEntry {
747 message_id: String::new(),
748 received_at_ms: 3000,
749 sender: "_runtime".into(),
750 message_type: "Checkpoint".into(),
751 raw_payload: checkpoint_payload,
752 entry_kind: EntryKind::Checkpoint,
753 session_id: "s1".into(),
754 mode: "macp.mode.decision.v1".into(),
755 macp_version: "1.0".into(),
756 timestamp_unix_ms: 3000,
757 bound_mode_version: None,
758 semantics_rev: 0,
759 bound_max_suspend_ms: None,
760 compacted_incoming_ordinals: 0,
761 };
762
763 let vote = VotePayload {
765 proposal_id: "p1".into(),
766 vote: "approve".into(),
767 reason: "lgtm".into(),
768 }
769 .encode_to_vec();
770
771 let entries_with_checkpoint = vec![
773 full_entries[0].clone(),
774 full_entries[1].clone(),
775 checkpoint,
776 incoming_entry("m3", "Vote", "agent://fraud", vote, 4000),
777 ];
778
779 let session = replay_session("s1", &entries_with_checkpoint, ®istry, None).unwrap();
780 assert_eq!(session.state, SessionState::Open);
781 assert_eq!(
782 session.intent, "restored-from-checkpoint",
783 "the checkpoint fast path must have been taken, else this test \
784 proves nothing about the checkpoint"
785 );
786 assert!(session.seen_message_ids.contains("m1"));
788 assert!(session.seen_message_ids.contains("m2"));
789 assert!(session.seen_message_ids.contains("m3"));
790 }
791
792 #[test]
805 fn replay_from_checkpoint_restores_suspension_intervals() {
806 use crate::registry::PersistedSession;
807
808 let registry = make_registry();
809 let start_payload = SessionStartPayload {
810 intent: "test".into(),
811 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
812 mode_version: "1.0.0".into(),
813 configuration_version: "cfg-1".into(),
814 policy_version: String::new(),
815 ttl_ms: 60_000,
816 context_id: String::new(),
817 extensions: std::collections::HashMap::new(),
818 roots: vec![],
819 max_suspend_ms: 0,
820 }
821 .encode_to_vec();
822
823 let prefix = vec![
824 incoming_entry(
825 "m1",
826 "SessionStart",
827 "agent://orchestrator",
828 start_payload,
829 1_000,
830 ),
831 internal_entry("SessionSuspend", 1_050),
832 internal_entry("SessionResume", 1_300),
833 ];
834 let before_checkpoint = replay_session("s1", &prefix, ®istry, None).unwrap();
835 assert_eq!(before_checkpoint.suspension_intervals, vec![(1_050, 1_300)]);
836
837 let mut persisted = PersistedSession::from(&before_checkpoint);
838 persisted.intent = "restored-from-checkpoint".into();
841 let checkpoint_payload = serde_json::to_vec(&persisted).unwrap();
844 let checkpoint = LogEntry {
845 message_id: String::new(),
846 received_at_ms: 1_400,
847 sender: "_runtime".into(),
848 message_type: "Checkpoint".into(),
849 raw_payload: checkpoint_payload,
850 entry_kind: EntryKind::Checkpoint,
851 session_id: "s1".into(),
852 mode: "macp.mode.decision.v1".into(),
853 macp_version: "1.0".into(),
854 timestamp_unix_ms: 1_400,
855 bound_mode_version: None,
856 semantics_rev: 0,
857 bound_max_suspend_ms: None,
858 compacted_incoming_ordinals: 0,
859 };
860
861 let entries = vec![
864 prefix[0].clone(),
865 prefix[1].clone(),
866 prefix[2].clone(),
867 checkpoint,
868 internal_entry("SessionSuspend", 1_500),
869 internal_entry("SessionResume", 1_600),
870 ];
871 let session = replay_session("s1", &entries, ®istry, None).unwrap();
872 assert_eq!(session.state, SessionState::Open);
873 assert_eq!(
874 session.intent, "restored-from-checkpoint",
875 "the checkpoint fast path must have been taken, else this test \
876 proves nothing about the snapshot round-trip"
877 );
878 assert_eq!(
879 session.suspension_intervals,
880 vec![(1_050, 1_300), (1_500, 1_600)],
881 "the pre-checkpoint pause must come from the snapshot and the \
882 post-checkpoint pause from the replayed tail"
883 );
884 }
885
886 #[test]
896 fn replay_ignores_a_disagreeing_banked_ms_payload() {
897 let registry = make_registry();
898 let start_payload = SessionStartPayload {
899 intent: "test".into(),
900 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
901 mode_version: "1.0.0".into(),
902 configuration_version: "cfg-1".into(),
903 policy_version: String::new(),
904 ttl_ms: 60_000,
905 context_id: String::new(),
906 extensions: std::collections::HashMap::new(),
907 roots: vec![],
908 max_suspend_ms: 0,
909 }
910 .encode_to_vec();
911
912 let start_at = 1_000;
913 let suspend_at = 1_050;
914 let resume_at = 1_300;
915 let resume_payload = SessionResumePayload {
919 reason: "test".into(),
920 resumed_by: "agent://orchestrator".into(),
921 banked_ms: 10_000,
922 }
923 .encode_to_vec();
924
925 let entries = vec![
926 incoming_entry(
927 "m1",
928 "SessionStart",
929 "agent://orchestrator",
930 start_payload,
931 start_at,
932 ),
933 internal_entry("SessionSuspend", suspend_at),
934 LogEntry {
935 message_id: String::new(),
936 received_at_ms: resume_at,
937 sender: "_runtime".into(),
938 message_type: "SessionResume".into(),
939 raw_payload: resume_payload,
940 entry_kind: EntryKind::Internal,
941 session_id: "s1".into(),
942 mode: "macp.mode.decision.v1".into(),
943 macp_version: "1.0".into(),
944 timestamp_unix_ms: resume_at,
945 bound_mode_version: None,
946 semantics_rev: 0,
947 bound_max_suspend_ms: None,
948 compacted_incoming_ordinals: 0,
949 },
950 ];
951
952 let session = replay_session("s1", &entries, ®istry, None).unwrap();
953 assert_eq!(session.state, SessionState::Open);
954 assert_eq!(
955 session.ttl_expiry,
956 start_at + 60_000 + (resume_at - suspend_at),
957 "replay must derive the banked duration from the recorded \
958 suspend/resume timestamps (250ms here), never from the \
959 payload's banked_ms field (10000ms here)"
960 );
961 }
962
963 #[test]
964 fn replay_without_checkpoint_still_works() {
965 let registry = make_registry();
967 let entries = vec![incoming_entry(
968 "m1",
969 "SessionStart",
970 "agent://orchestrator",
971 start_payload_bytes(),
972 1000,
973 )];
974 let session = replay_session("s1", &entries, ®istry, None).unwrap();
975 assert_eq!(session.state, SessionState::Open);
976 assert!(session.seen_message_ids.contains("m1"));
977 }
978
979 fn ext_registry_with_dyn_mode(version: &str) -> ModeRegistry {
980 let registry = make_registry();
981 registry
982 .register_extension(macp_pb::pb::ModeDescriptor {
983 mode: "ext.dyn.v1".into(),
984 mode_version: version.into(),
985 message_types: vec!["SessionStart".into(), "Commitment".into()],
986 terminal_message_types: vec!["Commitment".into()],
987 ..Default::default()
988 })
989 .unwrap();
990 registry
991 }
992
993 fn ext_start_entry(bound_mode_version: Option<String>) -> LogEntry {
994 let payload = SessionStartPayload {
996 participants: vec!["alice".into()],
997 configuration_version: "cfg-1".into(),
998 ttl_ms: 60_000,
999 ..Default::default()
1000 }
1001 .encode_to_vec();
1002 LogEntry {
1003 message_id: "m1".into(),
1004 received_at_ms: 1000,
1005 sender: "alice".into(),
1006 message_type: "SessionStart".into(),
1007 raw_payload: payload,
1008 entry_kind: EntryKind::Incoming,
1009 session_id: "s1".into(),
1010 mode: "ext.dyn.v1".into(),
1011 macp_version: "1.0".into(),
1012 timestamp_unix_ms: 1000,
1013 bound_mode_version,
1014 semantics_rev: 0,
1015 bound_max_suspend_ms: None,
1016 compacted_incoming_ordinals: 0,
1017 }
1018 }
1019
1020 #[test]
1024 fn replay_uses_recorded_mode_version_binding() {
1025 let registry = ext_registry_with_dyn_mode("9.9.9");
1026 let entries = vec![ext_start_entry(Some("2.5.0".into()))];
1027 let session = replay_session("s1", &entries, ®istry, None).unwrap();
1028 assert_eq!(session.mode_version, "2.5.0");
1029 }
1030
1031 #[test]
1035 fn replay_legacy_entry_without_binding_keeps_empty_version() {
1036 let registry = ext_registry_with_dyn_mode("9.9.9");
1037 let entries = vec![ext_start_entry(None)];
1038 let session = replay_session("s1", &entries, ®istry, None).unwrap();
1039 assert_eq!(session.mode_version, "");
1040 }
1041
1042 #[test]
1045 fn legacy_log_entry_json_without_binding_field_deserializes() {
1046 let json = serde_json::json!({
1047 "message_id": "m1",
1048 "received_at_ms": 1000,
1049 "sender": "alice",
1050 "message_type": "SessionStart",
1051 "raw_payload": [],
1052 "entry_kind": "Incoming",
1053 "session_id": "s1",
1054 "mode": "ext.dyn.v1",
1055 "macp_version": "1.0",
1056 "timestamp_unix_ms": 1000
1057 });
1058 let entry: LogEntry = serde_json::from_value(json).unwrap();
1059 assert_eq!(entry.bound_mode_version, None);
1060 assert_eq!(entry.semantics_rev, 0);
1063 assert_eq!(entry.bound_max_suspend_ms, None);
1065 }
1066
1067 #[test]
1070 fn replay_uses_recorded_max_suspend_cap() {
1071 let registry = ext_registry_with_dyn_mode("1.0.0");
1072 let mut entry = ext_start_entry(Some("1.0.0".into()));
1073 entry.bound_max_suspend_ms = Some(1234);
1074 let session = replay_session("s1", &[entry], ®istry, None).unwrap();
1075 assert_eq!(session.max_suspend_ms, 1234);
1076 assert_eq!(session.effective_max_suspend_ms(), 1234);
1077 }
1078
1079 #[test]
1082 fn replay_legacy_entry_keeps_default_cap_semantics() {
1083 let registry = ext_registry_with_dyn_mode("1.0.0");
1084 let entries = vec![ext_start_entry(None)];
1085 let session = replay_session("s1", &entries, ®istry, None).unwrap();
1086 assert_eq!(session.max_suspend_ms, 0);
1087 assert_eq!(
1088 session.effective_max_suspend_ms(),
1089 macp_core::session::MAX_SUSPEND_MS
1090 );
1091 }
1092
1093 #[test]
1098 fn replay_preserves_recorded_semantics_rev() {
1099 let registry = ext_registry_with_dyn_mode("1.0.0");
1100 let entries = vec![ext_start_entry(Some("1.0.0".into()))]; let session = replay_session("s1", &entries, ®istry, None).unwrap();
1102 assert_eq!(session.semantics_rev, 0);
1103 assert_ne!(
1104 session.semantics_rev,
1105 macp_core::session::CURRENT_SEMANTICS_REV
1106 );
1107 }
1108
1109 #[test]
1110 fn replay_consistency_flags_state_and_dedup_divergence() {
1111 let a = Session::builder("s1", "macp.mode.decision.v1", "agent://a")
1112 .mode_version("1.0.0")
1113 .configuration_version("cfg-1")
1114 .build();
1115 assert_eq!(validate_replay_consistency("s1", &a, &a.clone()), 0);
1117
1118 let mut b = a.clone();
1120 b.state = SessionState::Resolved;
1121 b.seen_message_ids.insert("m1".into());
1122 assert_eq!(validate_replay_consistency("s1", &a, &b), 2);
1123
1124 let mut c = a.clone();
1126 c.mode_state = vec![7, 7, 7];
1127 assert_eq!(validate_replay_consistency("s1", &a, &c), 1);
1128
1129 let mut d = a.clone();
1133 d.accumulated_suspended_ms = 5_000;
1134 assert_eq!(validate_replay_consistency("s1", &a, &d), 1);
1135 d.suspended_at_ms = Some(1_000);
1136 assert_eq!(validate_replay_consistency("s1", &a, &d), 2);
1137 d.suspension_intervals = vec![(1_000, 6_000)];
1139 assert_eq!(validate_replay_consistency("s1", &a, &d), 3);
1140
1141 let mut e = b.clone();
1144 e.mode_state = vec![7, 7, 7];
1145 e.accumulated_suspended_ms = 5_000;
1146 e.suspended_at_ms = Some(1_000);
1147 e.suspension_intervals = vec![(1_000, 6_000)];
1148 assert_eq!(validate_replay_consistency("s1", &a, &e), 6);
1149 }
1150
1151 const HANDOFF_TIMEOUT_MS: i64 = 100;
1165
1166 fn handoff_policy_registry() -> PolicyRegistry {
1167 let registry = PolicyRegistry::new();
1168 registry
1169 .register(macp_core::policy::PolicyDefinition {
1170 policy_id: "handoff-auto-accept".into(),
1171 mode: "macp.mode.handoff.v1".into(),
1172 description: "implicit accept after 100ms".into(),
1173 rules: serde_json::json!({
1174 "acceptance": { "implicit_accept_timeout_ms": HANDOFF_TIMEOUT_MS },
1175 "commitment": { "authority": "initiator_only" }
1176 }),
1177 schema_version: 1,
1178 })
1179 .unwrap();
1180 registry
1181 }
1182
1183 fn handoff_entry(
1184 message_id: &str,
1185 message_type: &str,
1186 payload: Vec<u8>,
1187 envelope_ms: i64,
1188 received_ms: i64,
1189 ) -> LogEntry {
1190 LogEntry {
1191 message_id: message_id.into(),
1192 received_at_ms: received_ms,
1193 sender: "alice".into(),
1194 message_type: message_type.into(),
1195 raw_payload: payload,
1196 entry_kind: EntryKind::Incoming,
1197 session_id: "s1".into(),
1198 mode: "macp.mode.handoff.v1".into(),
1199 macp_version: "1.0".into(),
1200 timestamp_unix_ms: envelope_ms,
1201 bound_mode_version: None,
1202 semantics_rev: 0,
1203 bound_max_suspend_ms: None,
1204 compacted_incoming_ordinals: 0,
1205 }
1206 }
1207
1208 fn implicit_accept_entry(deadline_ms: i64) -> LogEntry {
1219 let payload = crate::handoff_pb::HandoffAcceptPayload {
1220 handoff_id: "h1".into(),
1221 accepted_by: "bob".into(),
1222 reason: "implicit accept (timeout)".into(),
1223 implicit: true,
1224 }
1225 .encode_to_vec();
1226 let mut entry = handoff_entry(
1227 "implicit-accept:h1",
1228 "HandoffAccept",
1229 payload,
1230 deadline_ms,
1231 deadline_ms,
1232 );
1233 entry.sender = "bob".into();
1234 entry
1235 }
1236
1237 fn handoff_history(
1241 semantics_rev: u32,
1242 commit_envelope_ms: i64,
1243 commit_received_ms: i64,
1244 ) -> Vec<LogEntry> {
1245 let start_payload = SessionStartPayload {
1246 intent: "escalate".into(),
1247 participants: vec!["alice".into(), "bob".into()],
1248 mode_version: "1.0.0".into(),
1249 configuration_version: "cfg-1".into(),
1250 policy_version: "handoff-auto-accept".into(),
1251 ttl_ms: 60_000,
1252 context_id: String::new(),
1253 extensions: std::collections::HashMap::new(),
1254 roots: vec![],
1255 max_suspend_ms: 0,
1256 }
1257 .encode_to_vec();
1258 let offer = crate::handoff_pb::HandoffOfferPayload {
1259 handoff_id: "h1".into(),
1260 target_participant: "bob".into(),
1261 scope: "support".into(),
1262 reason: "escalate".into(),
1263 }
1264 .encode_to_vec();
1265 let commitment = CommitmentPayload {
1266 commitment_id: "c1".into(),
1267 action: "handoff.accepted".into(),
1268 authority_scope: "support".into(),
1269 reason: "bound".into(),
1270 mode_version: "1.0.0".into(),
1271 policy_version: "handoff-auto-accept".into(),
1272 configuration_version: "cfg-1".into(),
1273 outcome_positive: true,
1274 supersedes: None,
1275 }
1276 .encode_to_vec();
1277
1278 let mut start = handoff_entry("m1", "SessionStart", start_payload, 1_000, 1_000);
1279 start.semantics_rev = semantics_rev;
1280 vec![
1281 start,
1282 handoff_entry("m2", "HandoffOffer", offer, 1_000, 1_000),
1285 handoff_entry(
1286 "m3",
1287 "Commitment",
1288 commitment,
1289 commit_envelope_ms,
1290 commit_received_ms,
1291 ),
1292 ]
1293 }
1294
1295 fn assert_implicitly_accepted(session: &Session) {
1298 assert_eq!(session.state, SessionState::Resolved);
1299 let state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
1300 let offer = &state["offers"]["h1"];
1301 assert_eq!(offer["disposition"], "Accepted");
1302 assert_eq!(offer["accepted_by"], "bob");
1303 assert_eq!(offer["outcome_reason"], "implicit accept (timeout)");
1304 }
1305
1306 #[test]
1311 fn legacy_rev0_handoff_history_replays_under_envelope_clock() {
1312 let registry = make_registry();
1313 let policies = handoff_policy_registry();
1314 let entries = handoff_history(0, 1_300, 1_050);
1315
1316 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1317 assert_eq!(session.semantics_rev, 0);
1318 assert_implicitly_accepted(&session);
1319
1320 for rev in [1, macp_core::session::CURRENT_SEMANTICS_REV] {
1325 let mut newer = entries.clone();
1326 newer[0].semantics_rev = rev;
1327 assert!(
1328 replay_session("s1", &newer, ®istry, Some(&policies)).is_err(),
1329 "rev {rev} must not reproduce the rev-0 outcome"
1330 );
1331 }
1332 }
1333
1334 #[test]
1338 fn legacy_rev1_handoff_history_replays_under_acceptance_clock() {
1339 let registry = make_registry();
1340 let policies = handoff_policy_registry();
1341 let entries = handoff_history(1, 1_050, 1_300);
1342
1343 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1344 assert_eq!(session.semantics_rev, 1);
1345 assert_implicitly_accepted(&session);
1346
1347 let mut legacy = entries.clone();
1349 legacy[0].semantics_rev = 0;
1350 assert!(replay_session("s1", &legacy, ®istry, Some(&policies)).is_err());
1351 }
1352
1353 #[test]
1372 fn current_rev_handoff_history_replays_identically_to_rev1() {
1373 let registry = make_registry();
1374 let policies = handoff_policy_registry();
1375
1376 let rev1 = replay_session(
1377 "s1",
1378 &handoff_history(1, 1_050, 1_300),
1379 ®istry,
1380 Some(&policies),
1381 )
1382 .unwrap();
1383
1384 let mut current_entries =
1388 handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1389 let commitment = current_entries.pop().expect("commitment is last");
1390 current_entries.push(implicit_accept_entry(1_100));
1391 current_entries.push(commitment);
1392
1393 let current = replay_session("s1", ¤t_entries, ®istry, Some(&policies)).unwrap();
1394
1395 assert_implicitly_accepted(¤t);
1396 assert_eq!(current.state, rev1.state);
1397 assert_eq!(current.mode_state, rev1.mode_state);
1398 assert_eq!(current.resolution, rev1.resolution);
1399
1400 assert_implicitly_accepted(&rev1);
1403 assert!(!rev1.seen_message_ids.contains("implicit-accept:h1"));
1404 assert!(current.seen_message_ids.contains("implicit-accept:h1"));
1407 }
1408
1409 #[test]
1421 fn rev2_commitment_without_synthetic_entry_fails_replay() {
1422 let registry = make_registry();
1423 let policies = handoff_policy_registry();
1424
1425 let entries = handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1429 assert!(
1430 replay_session("s1", &entries, ®istry, Some(&policies)).is_err(),
1431 "rev 2 must not infer an accept the history does not record"
1432 );
1433
1434 let mut legacy = entries.clone();
1437 legacy[0].semantics_rev = 1;
1438 let session = replay_session("s1", &legacy, ®istry, Some(&policies))
1439 .expect("rev 1 keeps the interim in-Commitment implicit accept");
1440 assert_implicitly_accepted(&session);
1441
1442 let mut with_synthetic = entries.clone();
1445 let commitment = with_synthetic.pop().expect("commitment is last");
1446 with_synthetic.push(implicit_accept_entry(1_100));
1447 with_synthetic.push(commitment);
1448 let session = replay_session("s1", &with_synthetic, ®istry, Some(&policies))
1449 .expect("rev 2 resolves once the synthetic entry is in history");
1450 assert_implicitly_accepted(&session);
1451 }
1452
1453 fn handoff_history_with_suspension(
1458 semantics_rev: u32,
1459 suspend_at_ms: i64,
1460 resume_at_ms: i64,
1461 commit_ms: i64,
1462 ) -> Vec<LogEntry> {
1463 let mut entries = handoff_history(semantics_rev, commit_ms, commit_ms);
1466 let commit = entries.pop().expect("commitment is the last entry");
1467 entries.push(internal_entry("SessionSuspend", suspend_at_ms));
1468 entries.push(internal_entry("SessionResume", resume_at_ms));
1469 entries.push(commit);
1470 entries
1471 }
1472
1473 #[test]
1483 fn legacy_rev1_handoff_history_with_suspension_still_implicitly_accepts() {
1484 let registry = make_registry();
1485 let policies = handoff_policy_registry();
1486 let entries = handoff_history_with_suspension(1, 1_050, 1_300, 1_300);
1487
1488 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1489 assert_eq!(session.semantics_rev, 1);
1490 assert_eq!(session.accumulated_suspended_ms, 250);
1491 assert_implicitly_accepted(&session);
1492
1493 let mut rev2 = entries.clone();
1497 rev2[0].semantics_rev = macp_core::session::CURRENT_SEMANTICS_REV;
1498 assert!(
1499 replay_session("s1", &rev2, ®istry, Some(&policies)).is_err(),
1500 "rev 2 must not reproduce the rev-1 outcome"
1501 );
1502 }
1503
1504 #[test]
1519 fn rev2_handoff_history_implicitly_accepts_on_unsuspended_time() {
1520 let registry = make_registry();
1521 let policies = handoff_policy_registry();
1522 let mut entries = handoff_history_with_suspension(
1523 macp_core::session::CURRENT_SEMANTICS_REV,
1524 1_050,
1525 1_300,
1526 1_450,
1527 );
1528 let commitment = entries.pop().expect("commitment is last");
1529 entries.push(implicit_accept_entry(1_350));
1530 entries.push(commitment);
1531
1532 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1533 assert_eq!(session.accumulated_suspended_ms, 250);
1534 assert_implicitly_accepted(&session);
1535 assert!(session.seen_message_ids.contains("implicit-accept:h1"));
1536 }
1537
1538 fn handoff_history_with_two_suspensions(semantics_rev: u32, commit_ms: i64) -> Vec<LogEntry> {
1563 let mut entries = handoff_history(semantics_rev, commit_ms, commit_ms);
1564 let commit = entries.pop().expect("commitment is the last entry");
1565 entries.push(internal_entry("SessionSuspend", 1_050));
1566 entries.push(internal_entry("SessionResume", 1_300));
1567 entries.push(internal_entry("SessionSuspend", 1_330));
1568 entries.push(internal_entry("SessionResume", 1_500));
1569 entries.push(commit);
1570 entries
1571 }
1572
1573 #[test]
1585 fn rev2_handoff_history_subtracts_every_suspension_pair() {
1586 let registry = make_registry();
1587 let policies = handoff_policy_registry();
1588
1589 let rev1 = handoff_history_with_two_suspensions(1, 1_510);
1590 let session = replay_session("s1", &rev1, ®istry, Some(&policies)).unwrap();
1591 assert_eq!(session.semantics_rev, 1);
1592 assert_eq!(session.accumulated_suspended_ms, 420);
1593 assert_implicitly_accepted(&session);
1594
1595 let rev2 =
1596 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_510);
1597 assert!(
1598 replay_session("s1", &rev2, ®istry, Some(&policies)).is_err(),
1599 "rev 2 must subtract both pauses (90ms unsuspended < 100ms timeout)"
1600 );
1601 }
1602
1603 #[test]
1611 fn rev2_handoff_history_accepts_on_unsuspended_time_across_two_pauses() {
1612 let registry = make_registry();
1613 let policies = handoff_policy_registry();
1614 let mut entries =
1615 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_600);
1616 let commitment = entries.pop().expect("commitment is last");
1617 entries.push(implicit_accept_entry(1_520));
1618 entries.push(commitment);
1619
1620 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1621 assert_eq!(session.accumulated_suspended_ms, 420);
1622 assert_implicitly_accepted(&session);
1623 }
1624
1625 #[test]
1634 fn replay_rebuilds_suspension_intervals_from_the_log() {
1635 let registry = make_registry();
1636 let policies = handoff_policy_registry();
1637
1638 let rev1 = handoff_history_with_two_suspensions(1, 1_510);
1639 let session = replay_session("s1", &rev1, ®istry, Some(&policies)).unwrap();
1640 assert_eq!(
1641 session.suspension_intervals,
1642 vec![(1_050, 1_300), (1_330, 1_500)]
1643 );
1644
1645 let mut rev2 =
1646 handoff_history_with_two_suspensions(macp_core::session::CURRENT_SEMANTICS_REV, 1_600);
1647 let commitment = rev2.pop().expect("commitment is last");
1650 rev2.push(implicit_accept_entry(1_520));
1651 rev2.push(commitment);
1652 let session = replay_session("s1", &rev2, ®istry, Some(&policies)).unwrap();
1653 assert_eq!(
1654 session.suspension_intervals,
1655 vec![(1_050, 1_300), (1_330, 1_500)]
1656 );
1657 assert_eq!(session.unsuspended_deadline(1_000, 100), 1_520);
1662 }
1663
1664 fn handoff_history_with_reserved_message_ids(semantics_rev: u32) -> Vec<LogEntry> {
1686 let mut entries = handoff_history(semantics_rev, 1_050, 1_300);
1689 entries[0].message_id = "implicit-accept:squatted-at-start".into();
1690 entries[2].message_id = squatted_commitment_id(semantics_rev).into();
1691 if semantics_rev >= 2 {
1692 let commitment = entries.pop().expect("commitment is last");
1700 entries.push(implicit_accept_entry(1_100));
1701 entries.push(commitment);
1702 }
1703 entries
1704 }
1705
1706 fn squatted_commitment_id(semantics_rev: u32) -> &'static str {
1709 if semantics_rev >= 2 {
1710 "implicit-accept:squatted-at-commit"
1711 } else {
1712 "implicit-accept:h1"
1713 }
1714 }
1715
1716 #[test]
1728 fn reserved_prefix_entry_replays_at_every_rev() {
1729 let registry = make_registry();
1730 let policies = handoff_policy_registry();
1731
1732 for rev in [1, macp_core::session::CURRENT_SEMANTICS_REV] {
1733 let entries = handoff_history_with_reserved_message_ids(rev);
1734 let session = replay_session("s1", &entries, ®istry, Some(&policies))
1735 .unwrap_or_else(|e| panic!("rev {rev} must replay reserved ids, got {e}"));
1736 assert_eq!(session.semantics_rev, rev);
1737 assert_implicitly_accepted(&session);
1738 assert!(session
1741 .seen_message_ids
1742 .contains("implicit-accept:squatted-at-start"));
1743 assert!(session
1744 .seen_message_ids
1745 .contains(squatted_commitment_id(rev)));
1746 }
1747
1748 let mode = crate::mode::handoff::HandoffMode::new(std::sync::Arc::new(
1752 macp_policy::DefaultPolicyEvaluator,
1753 ));
1754 let entries =
1755 handoff_history_with_reserved_message_ids(macp_core::session::CURRENT_SEMANTICS_REV);
1756 let session = replay_session("s1", &entries, ®istry, Some(&policies)).unwrap();
1757 for message_id in ["implicit-accept:squatted-at-start", "implicit-accept:h1"] {
1758 let env = Envelope {
1759 macp_version: "1.0".into(),
1760 mode: session.mode.clone(),
1761 message_type: "HandoffContext".into(),
1762 message_id: message_id.into(),
1763 session_id: "s1".into(),
1764 sender: "alice".into(),
1765 timestamp_unix_ms: 1_000,
1766 payload: vec![],
1767 };
1768 assert!(matches!(
1769 crate::mode::Mode::validate_client_envelope(&mode, &session, &env).unwrap_err(),
1770 MacpError::InvalidEnvelope
1771 ));
1772 }
1773 }
1774
1775 #[test]
1799 fn synthetic_shaped_entry_reaches_dispatch_not_the_client_boundary() {
1800 let registry = make_registry();
1801 let policies = handoff_policy_registry();
1802
1803 let mut entries = handoff_history(macp_core::session::CURRENT_SEMANTICS_REV, 1_050, 1_300);
1804 let commitment = entries.pop().expect("commitment is the last entry");
1810 entries.push(implicit_accept_entry(1_100));
1811 entries.push(commitment);
1812
1813 let session = replay_session("s1", &entries, ®istry, Some(&policies))
1814 .expect("the synthetic entry must replay through dispatch at rev >= 2");
1815 assert_implicitly_accepted(&session);
1816 assert!(session.seen_message_ids.contains("implicit-accept:h1"));
1819 }
1820
1821 fn multi_round_entry(
1824 message_id: &str,
1825 message_type: &str,
1826 sender: &str,
1827 payload: Vec<u8>,
1828 received_ms: i64,
1829 ) -> LogEntry {
1830 LogEntry {
1831 message_id: message_id.into(),
1832 received_at_ms: received_ms,
1833 sender: sender.into(),
1834 message_type: message_type.into(),
1835 raw_payload: payload,
1836 entry_kind: EntryKind::Incoming,
1837 session_id: "s1".into(),
1838 mode: "ext.multi_round.v1".into(),
1839 macp_version: "1.0".into(),
1840 timestamp_unix_ms: received_ms,
1841 bound_mode_version: None,
1842 semantics_rev: 0,
1843 bound_max_suspend_ms: None,
1844 compacted_incoming_ordinals: 0,
1845 }
1846 }
1847
1848 fn multi_round_collision_history(semantics_rev: u32) -> Vec<LogEntry> {
1856 let start_payload = SessionStartPayload {
1857 intent: "converge".into(),
1858 participants: vec!["alice".into()],
1859 mode_version: "1.0.0".into(),
1860 configuration_version: "cfg-1".into(),
1861 policy_version: String::new(),
1862 ttl_ms: 60_000,
1863 context_id: String::new(),
1864 extensions: std::collections::HashMap::new(),
1865 roots: vec![],
1866 max_suspend_ms: 0,
1867 }
1868 .encode_to_vec();
1869
1870 let contribute = macp_pb::multi_round_pb::ContributePayload {
1871 value: r#"{"value":"x"}"#.into(),
1872 }
1873 .encode_to_vec();
1874
1875 let commitment = CommitmentPayload {
1876 commitment_id: "c1".into(),
1877 action: "multi_round.converged".into(),
1878 authority_scope: "test".into(),
1879 reason: "converged".into(),
1880 mode_version: "1.0.0".into(),
1881 policy_version: String::new(),
1882 configuration_version: "cfg-1".into(),
1883 outcome_positive: true,
1884 supersedes: None,
1885 }
1886 .encode_to_vec();
1887
1888 let mut start =
1889 multi_round_entry("m1", "SessionStart", "coordinator", start_payload, 1_000);
1890 start.semantics_rev = semantics_rev;
1891 vec![
1892 start,
1893 multi_round_entry("m2", "Contribute", "alice", contribute, 2_000),
1894 multi_round_entry("m3", "Commitment", "coordinator", commitment, 3_000),
1895 ]
1896 }
1897
1898 fn converged_value(session: &Session) -> String {
1899 let resolution: serde_json::Value = serde_json::from_slice(
1900 session
1901 .resolution
1902 .as_deref()
1903 .expect("session must have resolved"),
1904 )
1905 .unwrap();
1906 resolution["converged_value"]
1907 .as_str()
1908 .expect("converged_value must be a string")
1909 .to_string()
1910 }
1911
1912 #[test]
1920 fn legacy_rev2_multi_round_history_replays_to_the_original_collision() {
1921 let registry = make_registry();
1922 let entries = multi_round_collision_history(2);
1923
1924 let session = replay_session("s1", &entries, ®istry, None).unwrap();
1925 assert_eq!(session.semantics_rev, 2);
1926 assert_eq!(session.state, SessionState::Resolved);
1927 assert_eq!(converged_value(&session), "x");
1928
1929 let mut current = entries.clone();
1933 current[0].semantics_rev = macp_core::session::CURRENT_SEMANTICS_REV;
1934 let current_session = replay_session("s1", ¤t, ®istry, None).unwrap();
1935 assert_ne!(converged_value(¤t_session), "x");
1936 }
1937
1938 #[test]
1942 fn rev3_multi_round_history_replays_to_the_corrected_value() {
1943 let registry = make_registry();
1944 let mut entries = multi_round_collision_history(macp_core::session::CURRENT_SEMANTICS_REV);
1945
1946 let session = replay_session("s1", &entries, ®istry, None).unwrap();
1947 assert_eq!(session.state, SessionState::Resolved);
1948 assert_eq!(converged_value(&session), r#"{"value":"x"}"#);
1949
1950 entries[0].semantics_rev = 2;
1953 let legacy_session = replay_session("s1", &entries, ®istry, None).unwrap();
1954 assert_eq!(converged_value(&legacy_session), "x");
1955 }
1956}