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(
184 session_id: &str,
185 replayed: &Session,
186 snapshot: &Session,
187) -> u32 {
188 let mut mismatches = 0u32;
189 if replayed.state != snapshot.state {
190 mismatches += 1;
191 tracing::warn!(
192 session_id,
193 replayed_state = ?replayed.state,
194 snapshot_state = ?snapshot.state,
195 "replay/snapshot state mismatch"
196 );
197 }
198 if replayed.seen_message_ids.len() != snapshot.seen_message_ids.len() {
199 mismatches += 1;
200 tracing::warn!(
201 session_id,
202 replayed_dedup = replayed.seen_message_ids.len(),
203 snapshot_dedup = snapshot.seen_message_ids.len(),
204 "replay/snapshot dedup count mismatch"
205 );
206 }
207 if replayed.participants != snapshot.participants {
208 mismatches += 1;
209 tracing::warn!(session_id, "replay/snapshot participants mismatch");
210 }
211 if replayed.mode_version != snapshot.mode_version
212 || replayed.configuration_version != snapshot.configuration_version
213 || replayed.policy_version != snapshot.policy_version
214 {
215 mismatches += 1;
216 tracing::warn!(
217 session_id,
218 "replay/snapshot bound-version mismatch (mode/configuration/policy)"
219 );
220 }
221 mismatches
222}
223
224fn replay_from_start(
226 session_id: &str,
227 log_entries: &[LogEntry],
228 registry: &ModeRegistry,
229 policy_registry: Option<&PolicyRegistry>,
230) -> Result<Session, MacpError> {
231 let start_entry = log_entries
233 .iter()
234 .find(|e| e.entry_kind == EntryKind::Incoming && e.message_type == "SessionStart")
235 .ok_or(MacpError::InvalidPayload)?;
236
237 let mode_name = if start_entry.mode.is_empty() {
239 return Err(MacpError::InvalidPayload);
242 } else {
243 &start_entry.mode
244 };
245
246 let mode = registry.get_mode(mode_name).ok_or(MacpError::UnknownMode)?;
247
248 let require_complete_start = registry.requires_strict_session_start(mode_name);
250 let start_payload = if start_entry.raw_payload.is_empty() && !require_complete_start {
251 crate::pb::SessionStartPayload::default()
252 } else {
253 parse_session_start_payload(&start_entry.raw_payload)?
254 };
255 if require_complete_start {
259 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
260 }
261
262 let ttl_ms = if !require_complete_start && start_payload.ttl_ms == 0 {
263 60_000i64
265 } else {
266 extract_ttl_ms(&start_payload)?
267 };
268
269 let started_at_unix_ms = start_entry.received_at_ms;
271 let ttl_expiry = started_at_unix_ms.saturating_add(ttl_ms);
272
273 let env = Envelope {
274 macp_version: if start_entry.macp_version.is_empty() {
275 "1.0".into()
276 } else {
277 start_entry.macp_version.clone()
278 },
279 mode: mode_name.to_string(),
280 message_type: "SessionStart".into(),
281 message_id: start_entry.message_id.clone(),
282 session_id: session_id.into(),
283 sender: start_entry.sender.clone(),
284 timestamp_unix_ms: if start_entry.timestamp_unix_ms != 0 {
285 start_entry.timestamp_unix_ms
286 } else {
287 start_entry.received_at_ms
288 },
289 payload: start_entry.raw_payload.clone(),
290 };
291
292 let mut session = Session::builder(session_id, mode_name, start_entry.sender.clone())
293 .semantics_rev(start_entry.semantics_rev)
296 .max_suspend_ms(start_entry.bound_max_suspend_ms.unwrap_or(0))
299 .ttl_expiry(ttl_expiry)
300 .ttl_ms(ttl_ms)
301 .started_at_unix_ms(started_at_unix_ms)
302 .participants(start_payload.participants.clone())
303 .intent(start_payload.intent.clone())
304 .mode_version(
310 start_entry
311 .bound_mode_version
312 .clone()
313 .unwrap_or_else(|| start_payload.mode_version.clone()),
314 )
315 .configuration_version(start_payload.configuration_version.clone())
316 .policy_version(start_payload.policy_version.clone())
317 .context_id(start_payload.context_id.clone())
318 .extensions(start_payload.extensions.clone())
319 .roots(start_payload.roots.clone())
320 .policy_definition(if !start_payload.policy_version.is_empty() {
321 policy_registry.and_then(|pr| pr.resolve(&start_payload.policy_version).ok())
322 } else {
323 None
324 })
325 .build();
326
327 let response = mode.on_session_start(&session, &env)?;
329 session.seen_message_ids.insert(env.message_id.clone());
330 session.apply_mode_response(response);
331
332 for entry in log_entries.iter().skip(1) {
334 replay_entry(&mut session, session_id, entry, &mode)?;
335 }
336
337 Ok(session)
338}
339
340#[cfg(test)]
341mod tests {
342 use super::*;
343 use crate::decision_pb::ProposalPayload;
344 use crate::decision_pb::VotePayload;
345 use crate::log_store::EntryKind;
346 use crate::pb::{CommitmentPayload, SessionStartPayload};
347 use prost::Message;
348
349 fn make_registry() -> ModeRegistry {
350 ModeRegistry::build_default(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator))
351 }
352
353 fn start_payload_bytes() -> Vec<u8> {
354 SessionStartPayload {
355 intent: "test".into(),
356 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
357 mode_version: "1.0.0".into(),
358 configuration_version: "cfg-1".into(),
359 policy_version: "policy-1".into(),
360 ttl_ms: 60_000,
361 context_id: String::new(),
362 extensions: std::collections::HashMap::new(),
363 roots: vec![],
364 max_suspend_ms: 0,
365 }
366 .encode_to_vec()
367 }
368
369 fn incoming_entry(
370 message_id: &str,
371 message_type: &str,
372 sender: &str,
373 payload: Vec<u8>,
374 received_at_ms: i64,
375 ) -> LogEntry {
376 LogEntry {
377 message_id: message_id.into(),
378 received_at_ms,
379 sender: sender.into(),
380 message_type: message_type.into(),
381 raw_payload: payload,
382 entry_kind: EntryKind::Incoming,
383 session_id: "s1".into(),
384 mode: "macp.mode.decision.v1".into(),
385 macp_version: "1.0".into(),
386 timestamp_unix_ms: received_at_ms,
387 bound_mode_version: None,
388 semantics_rev: 0,
389 bound_max_suspend_ms: None,
390 compacted_incoming_ordinals: 0,
391 }
392 }
393
394 fn internal_entry(message_type: &str, received_at_ms: i64) -> LogEntry {
395 LogEntry {
396 message_id: String::new(),
397 received_at_ms,
398 sender: "_runtime".into(),
399 message_type: message_type.into(),
400 raw_payload: vec![],
401 entry_kind: EntryKind::Internal,
402 session_id: "s1".into(),
403 mode: "macp.mode.decision.v1".into(),
404 macp_version: "1.0".into(),
405 timestamp_unix_ms: received_at_ms,
406 bound_mode_version: None,
407 semantics_rev: 0,
408 bound_max_suspend_ms: None,
409 compacted_incoming_ordinals: 0,
410 }
411 }
412
413 #[test]
414 fn replay_rebuilds_decision_session() {
415 let registry = make_registry();
416 let proposal = ProposalPayload {
417 proposal_id: "p1".into(),
418 option: "deploy".into(),
419 rationale: "ready".into(),
420 supporting_data: vec![],
421 }
422 .encode_to_vec();
423 let vote = VotePayload {
424 proposal_id: "p1".into(),
425 vote: "approve".into(),
426 reason: "lgtm".into(),
427 }
428 .encode_to_vec();
429 let commitment = CommitmentPayload {
430 commitment_id: "c1".into(),
431 action: "decision.selected".into(),
432 authority_scope: "payments".into(),
433 reason: "bound".into(),
434 mode_version: "1.0.0".into(),
435 policy_version: "policy-1".into(),
436 configuration_version: "cfg-1".into(),
437 outcome_positive: true,
438 supersedes: None,
439 }
440 .encode_to_vec();
441
442 let entries = vec![
443 incoming_entry(
444 "m1",
445 "SessionStart",
446 "agent://orchestrator",
447 start_payload_bytes(),
448 1000,
449 ),
450 incoming_entry("m2", "Proposal", "agent://orchestrator", proposal, 2000),
451 incoming_entry("m3", "Vote", "agent://fraud", vote, 3000),
452 incoming_entry("m4", "Commitment", "agent://orchestrator", commitment, 4000),
453 ];
454
455 let session = replay_session("s1", &entries, ®istry, None).unwrap();
456 assert_eq!(session.state, SessionState::Resolved);
457 assert_eq!(session.session_id, "s1");
458 assert!(session.seen_message_ids.contains("m1"));
459 assert!(session.seen_message_ids.contains("m2"));
460 assert!(session.seen_message_ids.contains("m3"));
461 assert!(session.seen_message_ids.contains("m4"));
462 assert!(session.resolution.is_some());
463 }
464
465 #[test]
466 fn replay_preserves_original_ttl() {
467 let registry = make_registry();
468 let original_time = 1_700_000_000_000i64;
469 let entries = vec![incoming_entry(
470 "m1",
471 "SessionStart",
472 "agent://orchestrator",
473 start_payload_bytes(),
474 original_time,
475 )];
476
477 let session = replay_session("s1", &entries, ®istry, None).unwrap();
478 assert_eq!(session.started_at_unix_ms, original_time);
479 assert_eq!(session.ttl_expiry, original_time + 60_000);
480 assert_eq!(session.ttl_ms, 60_000);
481 }
482
483 #[test]
484 fn replay_handles_ttl_expired() {
485 let registry = make_registry();
486 let entries = vec![
487 incoming_entry(
488 "m1",
489 "SessionStart",
490 "agent://orchestrator",
491 start_payload_bytes(),
492 1000,
493 ),
494 internal_entry("TtlExpired", 61001),
495 ];
496
497 let session = replay_session("s1", &entries, ®istry, None).unwrap();
498 assert_eq!(session.state, SessionState::Expired);
499 }
500
501 #[test]
502 fn replay_handles_session_cancel() {
503 let registry = make_registry();
504 let entries = vec![
505 incoming_entry(
506 "m1",
507 "SessionStart",
508 "agent://orchestrator",
509 start_payload_bytes(),
510 1000,
511 ),
512 internal_entry("SessionCancel", 5000),
513 ];
514
515 let session = replay_session("s1", &entries, ®istry, None).unwrap();
516 assert_eq!(session.state, SessionState::Cancelled);
518 }
519
520 #[test]
521 fn replay_fails_when_accepted_history_no_longer_applies() {
522 let registry = make_registry();
523 let vote = VotePayload {
524 proposal_id: "p1".into(),
525 vote: "approve".into(),
526 reason: String::new(),
527 }
528 .encode_to_vec();
529 let entries = vec![
530 incoming_entry(
531 "m1",
532 "SessionStart",
533 "agent://orchestrator",
534 start_payload_bytes(),
535 1000,
536 ),
537 incoming_entry("m2", "Vote", "agent://fraud", vote, 2000),
538 ];
539
540 let err = replay_session("s1", &entries, ®istry, None).unwrap_err();
541 let msg = err.to_string();
544 assert!(
545 msg == "InvalidTransition" || msg == "InvalidPayload" || msg == "Forbidden",
546 "unexpected error: {msg}"
547 );
548 }
549
550 #[test]
551 fn replay_empty_log_returns_error() {
552 let registry = make_registry();
553 let result = replay_session("s1", &[], ®istry, None);
554 assert!(result.is_err());
555 }
556
557 #[test]
558 fn backward_compat_old_log_entry_without_new_fields() {
559 let json = r#"{"message_id":"m1","received_at_ms":1000,"sender":"test","message_type":"Message","raw_payload":[],"entry_kind":"Incoming"}"#;
561 let entry: LogEntry = serde_json::from_str(json).unwrap();
562 assert_eq!(entry.session_id, "");
563 assert_eq!(entry.mode, "");
564 assert_eq!(entry.macp_version, "");
565 }
566
567 #[test]
568 fn replay_from_checkpoint_restores_state() {
569 use crate::registry::PersistedSession;
570
571 let registry = make_registry();
572
573 let proposal = ProposalPayload {
575 proposal_id: "p1".into(),
576 option: "deploy".into(),
577 rationale: "ready".into(),
578 supporting_data: vec![],
579 }
580 .encode_to_vec();
581
582 let full_entries = vec![
583 incoming_entry(
584 "m1",
585 "SessionStart",
586 "agent://orchestrator",
587 start_payload_bytes(),
588 1000,
589 ),
590 incoming_entry(
591 "m2",
592 "Proposal",
593 "agent://orchestrator",
594 proposal.clone(),
595 2000,
596 ),
597 ];
598 let full_session = replay_session("s1", &full_entries, ®istry, None).unwrap();
599
600 let persisted = PersistedSession::from(&full_session);
602 let checkpoint_payload = serde_json::to_vec(&persisted).unwrap();
603 let checkpoint = LogEntry {
604 message_id: String::new(),
605 received_at_ms: 3000,
606 sender: "_runtime".into(),
607 message_type: "Checkpoint".into(),
608 raw_payload: checkpoint_payload,
609 entry_kind: EntryKind::Checkpoint,
610 session_id: "s1".into(),
611 mode: "macp.mode.decision.v1".into(),
612 macp_version: "1.0".into(),
613 timestamp_unix_ms: 3000,
614 bound_mode_version: None,
615 semantics_rev: 0,
616 bound_max_suspend_ms: None,
617 compacted_incoming_ordinals: 0,
618 };
619
620 let vote = VotePayload {
622 proposal_id: "p1".into(),
623 vote: "approve".into(),
624 reason: "lgtm".into(),
625 }
626 .encode_to_vec();
627
628 let entries_with_checkpoint = vec![
630 full_entries[0].clone(),
631 full_entries[1].clone(),
632 checkpoint,
633 incoming_entry("m3", "Vote", "agent://fraud", vote, 4000),
634 ];
635
636 let session = replay_session("s1", &entries_with_checkpoint, ®istry, None).unwrap();
637 assert_eq!(session.state, SessionState::Open);
638 assert!(session.seen_message_ids.contains("m1"));
640 assert!(session.seen_message_ids.contains("m2"));
641 assert!(session.seen_message_ids.contains("m3"));
642 }
643
644 #[test]
645 fn replay_without_checkpoint_still_works() {
646 let registry = make_registry();
648 let entries = vec![incoming_entry(
649 "m1",
650 "SessionStart",
651 "agent://orchestrator",
652 start_payload_bytes(),
653 1000,
654 )];
655 let session = replay_session("s1", &entries, ®istry, None).unwrap();
656 assert_eq!(session.state, SessionState::Open);
657 assert!(session.seen_message_ids.contains("m1"));
658 }
659
660 fn ext_registry_with_dyn_mode(version: &str) -> ModeRegistry {
661 let registry = make_registry();
662 registry
663 .register_extension(macp_pb::pb::ModeDescriptor {
664 mode: "ext.dyn.v1".into(),
665 mode_version: version.into(),
666 message_types: vec!["SessionStart".into(), "Commitment".into()],
667 terminal_message_types: vec!["Commitment".into()],
668 ..Default::default()
669 })
670 .unwrap();
671 registry
672 }
673
674 fn ext_start_entry(bound_mode_version: Option<String>) -> LogEntry {
675 let payload = SessionStartPayload {
677 participants: vec!["alice".into()],
678 configuration_version: "cfg-1".into(),
679 ttl_ms: 60_000,
680 ..Default::default()
681 }
682 .encode_to_vec();
683 LogEntry {
684 message_id: "m1".into(),
685 received_at_ms: 1000,
686 sender: "alice".into(),
687 message_type: "SessionStart".into(),
688 raw_payload: payload,
689 entry_kind: EntryKind::Incoming,
690 session_id: "s1".into(),
691 mode: "ext.dyn.v1".into(),
692 macp_version: "1.0".into(),
693 timestamp_unix_ms: 1000,
694 bound_mode_version,
695 semantics_rev: 0,
696 bound_max_suspend_ms: None,
697 compacted_incoming_ordinals: 0,
698 }
699 }
700
701 #[test]
705 fn replay_uses_recorded_mode_version_binding() {
706 let registry = ext_registry_with_dyn_mode("9.9.9");
707 let entries = vec![ext_start_entry(Some("2.5.0".into()))];
708 let session = replay_session("s1", &entries, ®istry, None).unwrap();
709 assert_eq!(session.mode_version, "2.5.0");
710 }
711
712 #[test]
716 fn replay_legacy_entry_without_binding_keeps_empty_version() {
717 let registry = ext_registry_with_dyn_mode("9.9.9");
718 let entries = vec![ext_start_entry(None)];
719 let session = replay_session("s1", &entries, ®istry, None).unwrap();
720 assert_eq!(session.mode_version, "");
721 }
722
723 #[test]
726 fn legacy_log_entry_json_without_binding_field_deserializes() {
727 let json = serde_json::json!({
728 "message_id": "m1",
729 "received_at_ms": 1000,
730 "sender": "alice",
731 "message_type": "SessionStart",
732 "raw_payload": [],
733 "entry_kind": "Incoming",
734 "session_id": "s1",
735 "mode": "ext.dyn.v1",
736 "macp_version": "1.0",
737 "timestamp_unix_ms": 1000
738 });
739 let entry: LogEntry = serde_json::from_value(json).unwrap();
740 assert_eq!(entry.bound_mode_version, None);
741 assert_eq!(entry.semantics_rev, 0);
744 assert_eq!(entry.bound_max_suspend_ms, None);
746 }
747
748 #[test]
751 fn replay_uses_recorded_max_suspend_cap() {
752 let registry = ext_registry_with_dyn_mode("1.0.0");
753 let mut entry = ext_start_entry(Some("1.0.0".into()));
754 entry.bound_max_suspend_ms = Some(1234);
755 let session = replay_session("s1", &[entry], ®istry, None).unwrap();
756 assert_eq!(session.max_suspend_ms, 1234);
757 assert_eq!(session.effective_max_suspend_ms(), 1234);
758 }
759
760 #[test]
763 fn replay_legacy_entry_keeps_default_cap_semantics() {
764 let registry = ext_registry_with_dyn_mode("1.0.0");
765 let entries = vec![ext_start_entry(None)];
766 let session = replay_session("s1", &entries, ®istry, None).unwrap();
767 assert_eq!(session.max_suspend_ms, 0);
768 assert_eq!(
769 session.effective_max_suspend_ms(),
770 macp_core::session::MAX_SUSPEND_MS
771 );
772 }
773
774 #[test]
779 fn replay_preserves_recorded_semantics_rev() {
780 let registry = ext_registry_with_dyn_mode("1.0.0");
781 let entries = vec![ext_start_entry(Some("1.0.0".into()))]; let session = replay_session("s1", &entries, ®istry, None).unwrap();
783 assert_eq!(session.semantics_rev, 0);
784 assert_ne!(
785 session.semantics_rev,
786 macp_core::session::CURRENT_SEMANTICS_REV
787 );
788 }
789
790 #[test]
791 fn replay_consistency_flags_state_and_dedup_divergence() {
792 let a = Session::builder("s1", "macp.mode.decision.v1", "agent://a")
793 .mode_version("1.0.0")
794 .configuration_version("cfg-1")
795 .build();
796 assert_eq!(validate_replay_consistency("s1", &a, &a.clone()), 0);
798
799 let mut b = a.clone();
801 b.state = SessionState::Resolved;
802 b.seen_message_ids.insert("m1".into());
803 assert_eq!(validate_replay_consistency("s1", &a, &b), 2);
804 }
805}