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