Skip to main content

macp_runtime/
runtime.rs

1use chrono::Utc;
2use std::sync::Arc;
3
4use crate::error::MacpError;
5use crate::extensions::ExtensionProviderRegistry;
6use crate::log_store::{EntryKind, LogEntry, LogStore};
7use crate::metrics::RuntimeMetrics;
8use crate::mode_registry::ModeRegistry;
9use crate::pb::{Envelope, ModeDescriptor};
10use crate::policy::registry::PolicyRegistry;
11use crate::policy::PolicyDefinition;
12use crate::registry::SessionRegistry;
13use crate::session::{
14    extract_ttl_ms, parse_session_start_payload, validate_canonical_session_start_payload,
15    validate_session_id_for_acceptance, Session, SessionState,
16};
17use crate::storage::StorageBackend;
18use crate::stream_bus::SessionStreamBus;
19
20#[derive(Debug)]
21pub struct ProcessResult {
22    pub session_state: SessionState,
23    pub duplicate: bool,
24}
25
26#[derive(Clone, Debug)]
27pub enum SessionLifecycleEvent {
28    Created { session_id: String },
29    Resolved { session_id: String },
30    Expired { session_id: String },
31    Suspended { session_id: String },
32    Resumed { session_id: String },
33    Cancelled { session_id: String },
34}
35
36pub struct Runtime {
37    pub storage: Arc<dyn StorageBackend>,
38    pub registry: Arc<SessionRegistry>,
39    pub log_store: Arc<LogStore>,
40    stream_bus: Arc<SessionStreamBus>,
41    signal_bus: tokio::sync::broadcast::Sender<Envelope>,
42    session_lifecycle_bus: tokio::sync::broadcast::Sender<SessionLifecycleEvent>,
43    mode_registry: Arc<ModeRegistry>,
44    policy_registry: Arc<PolicyRegistry>,
45    #[allow(dead_code)] // plumbed for future session-extension providers; register API TBD
46    extensions: Arc<ExtensionProviderRegistry>,
47    metrics: Arc<RuntimeMetrics>,
48    checkpoint_interval: usize,
49}
50
51impl Runtime {
52    pub fn new(
53        storage: Arc<dyn StorageBackend>,
54        registry: Arc<SessionRegistry>,
55        log_store: Arc<LogStore>,
56    ) -> Self {
57        Self::with_mode_registry(
58            storage,
59            registry,
60            log_store,
61            Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
62                macp_policy::DefaultPolicyEvaluator,
63            ))),
64        )
65    }
66
67    pub fn with_mode_registry(
68        storage: Arc<dyn StorageBackend>,
69        registry: Arc<SessionRegistry>,
70        log_store: Arc<LogStore>,
71        mode_registry: Arc<ModeRegistry>,
72    ) -> Self {
73        Self::with_registries(
74            storage,
75            registry,
76            log_store,
77            mode_registry,
78            Arc::new(PolicyRegistry::new()),
79        )
80    }
81
82    pub fn with_registries(
83        storage: Arc<dyn StorageBackend>,
84        registry: Arc<SessionRegistry>,
85        log_store: Arc<LogStore>,
86        mode_registry: Arc<ModeRegistry>,
87        policy_registry: Arc<PolicyRegistry>,
88    ) -> Self {
89        let checkpoint_interval = std::env::var("MACP_CHECKPOINT_INTERVAL")
90            .ok()
91            .and_then(|v| v.parse().ok())
92            .unwrap_or(0); // 0 = disabled by default
93        let (signal_tx, _) = tokio::sync::broadcast::channel(256);
94        let (session_lifecycle_tx, _) = tokio::sync::broadcast::channel(64);
95        Self {
96            storage,
97            registry,
98            log_store,
99            stream_bus: Arc::new(SessionStreamBus::default()),
100            signal_bus: signal_tx,
101            session_lifecycle_bus: session_lifecycle_tx,
102            mode_registry,
103            policy_registry,
104            extensions: Arc::new(ExtensionProviderRegistry::new()),
105            metrics: Arc::new(RuntimeMetrics::new()),
106            checkpoint_interval,
107        }
108    }
109
110    /// Returns all mode names the runtime can handle (standards-track + extensions).
111    /// Used by Initialize and GetManifest to advertise full capability.
112    pub fn registered_mode_names(&self) -> Vec<String> {
113        self.mode_registry.all_mode_names()
114    }
115
116    /// Returns only standards-track mode descriptors for ListModes.
117    pub fn standard_mode_descriptors(&self) -> Vec<ModeDescriptor> {
118        self.mode_registry.standard_mode_descriptors()
119    }
120
121    /// Returns only extension mode descriptors for ListExtModes.
122    pub fn extension_mode_descriptors(&self) -> Vec<ModeDescriptor> {
123        self.mode_registry.extension_mode_descriptors()
124    }
125
126    pub fn register_extension(&self, descriptor: ModeDescriptor) -> Result<(), String> {
127        self.mode_registry.register_extension(descriptor)
128    }
129
130    pub fn unregister_extension(&self, mode: &str) -> Result<(), String> {
131        self.mode_registry.unregister_extension(mode)
132    }
133
134    pub fn promote_mode(&self, mode: &str, new_name: Option<&str>) -> Result<String, String> {
135        self.mode_registry.promote_mode(mode, new_name)
136    }
137
138    pub fn subscribe_mode_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
139        self.mode_registry.subscribe_changes()
140    }
141
142    pub fn mode_registry(&self) -> &Arc<ModeRegistry> {
143        &self.mode_registry
144    }
145
146    // ── Policy registry delegation ──────────────────────────────────
147
148    pub fn register_policy(&self, definition: PolicyDefinition) -> Result<(), String> {
149        self.policy_registry.register(definition)
150    }
151
152    pub fn unregister_policy(&self, policy_id: &str) -> Result<(), String> {
153        self.policy_registry.unregister(policy_id)
154    }
155
156    pub fn get_policy(&self, policy_id: &str) -> Option<PolicyDefinition> {
157        self.policy_registry.get(policy_id)
158    }
159
160    pub fn list_policies(&self, mode_filter: Option<&str>) -> Vec<PolicyDefinition> {
161        self.policy_registry.list(mode_filter)
162    }
163
164    pub fn subscribe_policy_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
165        self.policy_registry.subscribe_changes()
166    }
167
168    pub fn policy_registry(&self) -> &Arc<PolicyRegistry> {
169        &self.policy_registry
170    }
171
172    pub fn metrics(&self) -> &Arc<RuntimeMetrics> {
173        &self.metrics
174    }
175
176    pub fn subscribe_session_stream(
177        &self,
178        session_id: &str,
179    ) -> tokio::sync::broadcast::Receiver<Envelope> {
180        self.stream_bus.subscribe(session_id)
181    }
182
183    pub fn subscribe_signals(&self) -> tokio::sync::broadcast::Receiver<Envelope> {
184        self.signal_bus.subscribe()
185    }
186
187    pub fn subscribe_session_lifecycle(
188        &self,
189    ) -> tokio::sync::broadcast::Receiver<SessionLifecycleEvent> {
190        self.session_lifecycle_bus.subscribe()
191    }
192
193    /// RFC-MACP-0006-A1: Replay accepted envelopes from the session log for
194    /// passive subscribe. Returns `Incoming` log entries as reconstructed
195    /// `Envelope` values, starting from `after_sequence` (0-based log index).
196    pub async fn get_session_envelopes_after(
197        &self,
198        session_id: &str,
199        after_sequence: u64,
200    ) -> Vec<Envelope> {
201        self.log_store
202            .get_incoming_after(session_id, after_sequence)
203            .await
204            .into_iter()
205            .map(|(_idx, entry)| Envelope {
206                macp_version: if entry.macp_version.is_empty() {
207                    "1.0".into()
208                } else {
209                    entry.macp_version
210                },
211                mode: entry.mode,
212                message_type: entry.message_type,
213                message_id: entry.message_id,
214                session_id: entry.session_id,
215                sender: entry.sender,
216                timestamp_unix_ms: if entry.timestamp_unix_ms != 0 {
217                    entry.timestamp_unix_ms
218                } else {
219                    entry.received_at_ms
220                },
221                payload: entry.raw_payload,
222            })
223            .collect()
224    }
225
226    fn publish_accepted_envelope(&self, env: &Envelope) {
227        if !env.session_id.is_empty() {
228            self.stream_bus.publish(&env.session_id, env.clone());
229        }
230    }
231
232    fn make_incoming_entry(env: &Envelope) -> LogEntry {
233        LogEntry {
234            message_id: env.message_id.clone(),
235            received_at_ms: Utc::now().timestamp_millis(),
236            sender: env.sender.clone(),
237            message_type: env.message_type.clone(),
238            raw_payload: env.payload.clone(),
239            entry_kind: EntryKind::Incoming,
240            session_id: env.session_id.clone(),
241            mode: env.mode.clone(),
242            macp_version: env.macp_version.clone(),
243            timestamp_unix_ms: env.timestamp_unix_ms,
244        }
245    }
246
247    fn make_internal_entry(
248        message_type: &str,
249        payload: &[u8],
250        session_id: &str,
251        mode: &str,
252    ) -> LogEntry {
253        let now = Utc::now().timestamp_millis();
254        LogEntry {
255            message_id: String::new(),
256            received_at_ms: now,
257            sender: "_runtime".into(),
258            message_type: message_type.into(),
259            raw_payload: payload.to_vec(),
260            entry_kind: EntryKind::Internal,
261            session_id: session_id.into(),
262            mode: mode.into(),
263            macp_version: "1.0".into(),
264            timestamp_unix_ms: now,
265        }
266    }
267
268    async fn save_session_to_storage(&self, session: &Session) {
269        if let Err(err) = self.storage.save_session(session).await {
270            tracing::warn!(
271                session_id = %session.session_id,
272                error = %err,
273                "failed to persist session snapshot"
274            );
275        }
276    }
277
278    async fn maybe_expire_session(
279        &self,
280        session_id: &str,
281        session: &mut Session,
282    ) -> Result<bool, MacpError> {
283        let now = Utc::now().timestamp_millis();
284        // An Open session past its deadline, or a Suspended session that has
285        // exceeded the MAX_SUSPEND_MS cap (RFC-MACP-0001 §7.5), expires.
286        let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
287            || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
288        if expires {
289            let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
290            self.storage
291                .append_log_entry(session_id, &entry)
292                .await
293                .map_err(|_| MacpError::StorageFailed)?;
294            self.log_store.append(session_id, entry).await;
295            session.state = SessionState::Expired;
296            session.suspended_at_ms = None;
297            self.metrics.record_session_expired(&session.mode);
298            tracing::info!(session_id, "session expired via TTL");
299            let _ = self
300                .session_lifecycle_bus
301                .send(SessionLifecycleEvent::Expired {
302                    session_id: session_id.to_string(),
303                });
304            return Ok(true);
305        }
306        Ok(false)
307    }
308
309    pub async fn process(
310        &self,
311        env: &Envelope,
312        max_open_sessions: Option<usize>,
313    ) -> Result<ProcessResult, MacpError> {
314        match env.message_type.as_str() {
315            "SessionStart" => self.process_session_start(env, max_open_sessions).await,
316            "Signal" | "Progress" => self.process_signal(env).await,
317            _ => self.process_message(env).await,
318        }
319    }
320
321    async fn process_session_start(
322        &self,
323        env: &Envelope,
324        max_open_sessions: Option<usize>,
325    ) -> Result<ProcessResult, MacpError> {
326        if env.mode.trim().is_empty() {
327            return Err(MacpError::InvalidEnvelope);
328        }
329        validate_session_id_for_acceptance(&env.session_id)?;
330        let mode_name = env.mode.as_str();
331        let mode = self
332            .mode_registry
333            .get_mode(mode_name)
334            .ok_or(MacpError::UnknownMode)?;
335
336        let start_payload = parse_session_start_payload(&env.payload)?;
337        let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
338        if require_complete_start {
339            validate_canonical_session_start_payload(&start_payload)?;
340        }
341
342        // Validate mode_version matches the registered descriptor's version
343        if let Some(descriptor_version) = self.mode_registry.get_mode_version(mode_name) {
344            if !start_payload.mode_version.is_empty()
345                && start_payload.mode_version != descriptor_version
346            {
347                tracing::warn!(
348                    mode = mode_name,
349                    payload_version = %start_payload.mode_version,
350                    descriptor_version = %descriptor_version,
351                    "mode_version mismatch"
352                );
353                return Err(MacpError::InvalidEnvelope);
354            }
355        }
356
357        let ttl_ms = extract_ttl_ms(&start_payload)?;
358
359        let mut guard = self.registry.sessions.write().await;
360        if let Some(existing) = guard.get(&env.session_id) {
361            if existing.seen_message_ids.contains(&env.message_id) {
362                return Ok(ProcessResult {
363                    session_state: existing.state.clone(),
364                    duplicate: true,
365                });
366            }
367            return Err(MacpError::SessionAlreadyExists);
368        }
369
370        // Enforce max_open_sessions atomically under the write lock to
371        // prevent TOCTOU races where concurrent SessionStart requests
372        // both pass a read-lock count check before either is inserted.
373        if let Some(max_open) = max_open_sessions {
374            let now = Utc::now().timestamp_millis();
375            let count = guard
376                .values()
377                .filter(|s| {
378                    s.initiator_sender == env.sender
379                        && s.state == SessionState::Open
380                        && now <= s.ttl_expiry
381                })
382                .count();
383            if count >= max_open {
384                return Err(MacpError::RateLimited);
385            }
386        }
387
388        // Resolve the governance policy for this session.
389        // RFC-MACP-0012 §6.1: policy_version is resolved at SessionStart; empty
390        // resolves to "policy.default". The resolved PolicyDescriptor is stored
391        // immutably on the session for deterministic replay (RFC-MACP-0003 §3).
392        let effective_policy_version = if start_payload.policy_version.is_empty() {
393            crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
394        } else {
395            start_payload.policy_version.clone()
396        };
397        let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
398            Ok(policy) => {
399                // RFC 6.1: reject if policy mode doesn't match session mode
400                if policy.mode != "*" && policy.mode != mode_name {
401                    return Err(MacpError::InvalidPolicyDefinition);
402                }
403                Some(policy)
404            }
405            Err(_) => {
406                return Err(MacpError::UnknownPolicyVersion);
407            }
408        };
409
410        let accepted_at = Utc::now().timestamp_millis();
411        // RFC-MACP-0003 §2: TTL deadline is computed from the SessionStart
412        // envelope's timestamp_unix_ms, not wall-clock time. This ensures
413        // deterministic replay. Fall back to accepted_at if envelope has no timestamp.
414        let ttl_base = if env.timestamp_unix_ms > 0 {
415            env.timestamp_unix_ms
416        } else {
417            accepted_at
418        };
419        let ttl_expiry = ttl_base.saturating_add(ttl_ms);
420        let session = Session {
421            session_id: env.session_id.clone(),
422            state: SessionState::Open,
423            ttl_expiry,
424            ttl_ms,
425            started_at_unix_ms: accepted_at,
426            resolution: None,
427            mode: mode_name.to_string(),
428            mode_state: vec![],
429            participants: start_payload.participants.clone(),
430            seen_message_ids: std::collections::HashSet::new(),
431            intent: start_payload.intent.clone(),
432            mode_version: start_payload.mode_version.clone(),
433            configuration_version: start_payload.configuration_version.clone(),
434            policy_version: effective_policy_version,
435            context_id: start_payload.context_id.clone(),
436            extensions: start_payload.extensions.clone(),
437            roots: start_payload.roots.clone(),
438            initiator_sender: env.sender.clone(),
439            participant_message_counts: std::collections::HashMap::new(),
440            participant_last_seen: std::collections::HashMap::new(),
441            policy_definition,
442            suspended_at_ms: None,
443            accumulated_suspended_ms: 0,
444        };
445
446        let response = mode.on_session_start(&session, env)?;
447
448        // 1. Create storage directory and write log entry (COMMIT POINT)
449        self.storage
450            .create_session_storage(&env.session_id)
451            .await
452            .map_err(|_| MacpError::StorageFailed)?;
453        let incoming_entry = Self::make_incoming_entry(env);
454        self.storage
455            .append_log_entry(&env.session_id, &incoming_entry)
456            .await
457            .map_err(|_| MacpError::StorageFailed)?;
458
459        // 2. Update in-memory caches
460        self.log_store.create_session_log(&env.session_id).await;
461        self.log_store.append(&env.session_id, incoming_entry).await;
462
463        let mut session = session;
464        session.seen_message_ids.insert(env.message_id.clone());
465        session.apply_mode_response(response);
466
467        let result_state = session.state.clone();
468        // 3. Session save — fatal on SessionStart to ensure snapshot durability.
469        // For subsequent messages, the log entry (COMMIT POINT) is already persisted
470        // so snapshot failure is recoverable via replay.
471        if let Err(err) = self.storage.save_session(&session).await {
472            tracing::error!(
473                session_id = %session.session_id,
474                error = %err,
475                "failed to persist session snapshot at SessionStart"
476            );
477            return Err(MacpError::StorageFailed);
478        }
479        self.metrics.record_session_start(mode_name);
480        tracing::info!(
481            session_id = %env.session_id,
482            mode = mode_name,
483            sender = %env.sender,
484            "session started"
485        );
486        guard.insert(env.session_id.clone(), session);
487        self.publish_accepted_envelope(env);
488        let _ = self
489            .session_lifecycle_bus
490            .send(SessionLifecycleEvent::Created {
491                session_id: env.session_id.clone(),
492            });
493
494        Ok(ProcessResult {
495            session_state: result_state,
496            duplicate: false,
497        })
498    }
499
500    /// Process a session-scoped message following the RFC-MACP-0001 Section 7.3
501    /// terminal-state transition order:
502    /// 1. Check session OPEN
503    /// 2. Validate message (mode.authorize_sender + mode.on_message)
504    /// 3. Accept into history (log_store.append)
505    /// 4. Transition to RESOLVED (session.apply_mode_response)
506    /// 5. Reject subsequent messages (enforced by step 1 on next call)
507    async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
508        let mut guard = self.registry.sessions.write().await;
509        let session = guard
510            .get_mut(&env.session_id)
511            .ok_or(MacpError::UnknownSession)?;
512
513        // Per-message kernel invariants (dedup, mode-binding, TTL, the monotonic
514        // OPEN gate) live in `macp_modes::step` so any consumer of the
515        // coordination core runs the identical checks. The runtime is the first
516        // caller: it drives the phases here so it can interpose its append-only
517        // write between validation and commit (a failed write must not consume a
518        // dedup slot) — which a single all-in-one step could not preserve.
519        let now_ms = chrono::Utc::now().timestamp_millis();
520        match macp_modes::step::check_preconditions(session, env, now_ms)? {
521            macp_modes::step::Precheck::Duplicate => {
522                return Ok(ProcessResult {
523                    session_state: session.state.clone(),
524                    duplicate: true,
525                });
526            }
527            macp_modes::step::Precheck::Expired => {
528                // Durable expiry via the existing path: it appends the
529                // `TtlExpired` log entry, updates metrics/lifecycle, and marks
530                // the session Expired. `check_preconditions` and
531                // `maybe_expire_session` share the same strict `>`, OPEN-guarded
532                // rule, so this always expires.
533                let expired = self.maybe_expire_session(&env.session_id, session).await?;
534                debug_assert!(expired, "check_preconditions reported Expired");
535                self.save_session_to_storage(session).await;
536                return Err(MacpError::TtlExpired);
537            }
538            macp_modes::step::Precheck::Proceed => {}
539        }
540
541        let mode = self
542            .mode_registry
543            .get_mode(&session.mode)
544            .ok_or(MacpError::UnknownMode)?;
545        mode.authorize_sender(session, env)?;
546        let response = mode.on_message(session, env)?;
547
548        // 1. COMMIT POINT: write log entry to disk
549        let incoming_entry = Self::make_incoming_entry(env);
550        self.storage
551            .append_log_entry(&env.session_id, &incoming_entry)
552            .await
553            .map_err(|_| MacpError::StorageFailed)?;
554
555        // 2. Update in-memory state via the shared commit phase (consume dedup
556        //    slot, record participant activity, apply mode response) — the exact
557        //    sequence a library consumer runs through `macp_modes::step`.
558        self.log_store.append(&env.session_id, incoming_entry).await;
559        let result_state = macp_modes::step::commit(session, env, response, now_ms);
560
561        self.metrics.record_message_accepted(&session.mode);
562        if env.message_type == "Commitment" {
563            self.metrics.record_commitment_accepted(&session.mode);
564        }
565
566        tracing::debug!(
567            session_id = %env.session_id,
568            message_type = %env.message_type,
569            sender = %env.sender,
570            state = ?result_state,
571            "message accepted"
572        );
573
574        if result_state == SessionState::Resolved {
575            self.metrics.record_session_resolved(&session.mode);
576            tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
577            let _ = self
578                .session_lifecycle_bus
579                .send(SessionLifecycleEvent::Resolved {
580                    session_id: env.session_id.clone(),
581                });
582        }
583
584        // 3. Best-effort session save + checkpoint
585        self.save_session_to_storage(session).await;
586        if result_state == SessionState::Resolved {
587            if !self.maybe_compact_log(&env.session_id, session).await {
588                self.force_insert_checkpoint(&env.session_id, session).await;
589            }
590        } else {
591            self.maybe_insert_checkpoint(&env.session_id, session).await;
592        }
593        self.publish_accepted_envelope(env);
594
595        Ok(ProcessResult {
596            session_state: result_state,
597            duplicate: false,
598        })
599    }
600
601    /// Process a Signal or Progress envelope. Signals are informational out-of-band
602    /// notifications. Progress messages carry structured ProgressPayload.
603    /// Neither mutates session state — both are broadcast to subscribers.
604    async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
605        // RFC-MACP-0001 §4 / RFC-MACP-0010: validate SignalPayload structure.
606        // signal_type must be non-empty when a payload is present.
607        if env.message_type == "Signal" && !env.payload.is_empty() {
608            let signal: crate::pb::SignalPayload =
609                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
610            if signal.signal_type.trim().is_empty() {
611                return Err(MacpError::InvalidPayload);
612            }
613        }
614        // RFC-MACP-0001: validate ProgressPayload structure for Progress messages.
615        if env.message_type == "Progress" && !env.payload.is_empty() {
616            let _: crate::pb::ProgressPayload =
617                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
618        }
619        tracing::debug!(
620            sender = %env.sender,
621            message_id = %env.message_id,
622            message_type = %env.message_type,
623            "signal received"
624        );
625        let _ = self.signal_bus.send(env.clone());
626        Ok(ProcessResult {
627            session_state: SessionState::Open,
628            duplicate: false,
629        })
630    }
631
632    pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
633        let mut guard = self.registry.sessions.write().await;
634        let changed = if let Some(session) = guard.get_mut(session_id) {
635            self.maybe_expire_session(session_id, session)
636                .await
637                .unwrap_or(false)
638        } else {
639            return None;
640        };
641        if changed {
642            if let Some(session) = guard.get(session_id) {
643                self.save_session_to_storage(session).await;
644            }
645        }
646        guard.get(session_id).cloned()
647    }
648
649    /// Cancel a session. The `cancelled_by` parameter MUST be the authenticated
650    /// sender of the CancelSession RPC (RFC-MACP-0001 Section 7.3: CancelSession
651    /// is a Core control-plane message; mode authorization does not apply).
652    pub async fn cancel_session(
653        &self,
654        session_id: &str,
655        reason: &str,
656        cancelled_by: &str,
657    ) -> Result<ProcessResult, MacpError> {
658        let mut guard = self.registry.sessions.write().await;
659        let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
660
661        self.maybe_expire_session(session_id, session).await?;
662
663        // Already terminal (Resolved/Expired/Cancelled): nothing to do. An Open
664        // or Suspended session can still be cancelled (RFC-MACP-0001 §7.2/§7.3).
665        if session.state.is_terminal() {
666            let result_state = session.state.clone();
667            self.save_session_to_storage(session).await;
668            return Ok(ProcessResult {
669                session_state: result_state,
670                duplicate: false,
671            });
672        }
673
674        // RFC-MACP-0001: runtime encodes a proper SessionCancelPayload with
675        // `cancelled_by` set to the authenticated sender identity.
676        let cancel_payload = crate::pb::SessionCancelPayload {
677            reason: reason.to_string(),
678            cancelled_by: cancelled_by.to_string(),
679        };
680        let cancel_entry = Self::make_internal_entry(
681            "SessionCancel",
682            &prost::Message::encode_to_vec(&cancel_payload),
683            session_id,
684            &session.mode,
685        );
686        self.storage
687            .append_log_entry(session_id, &cancel_entry)
688            .await
689            .map_err(|_| MacpError::StorageFailed)?;
690        self.log_store.append(session_id, cancel_entry).await;
691        // RFC-MACP-0001 §7.3: cancellation terminates as CANCELLED (distinct
692        // from EXPIRED) — `cancel()` also clears any suspension marker.
693        let _ = session.cancel();
694        self.save_session_to_storage(session).await;
695        if !self.maybe_compact_log(session_id, session).await {
696            self.force_insert_checkpoint(session_id, session).await;
697        }
698        self.metrics.record_session_cancelled(&session.mode);
699        tracing::info!(session_id, reason, "session cancelled");
700        let _ = self
701            .session_lifecycle_bus
702            .send(SessionLifecycleEvent::Cancelled {
703                session_id: session_id.to_string(),
704            });
705
706        Ok(ProcessResult {
707            session_state: SessionState::Cancelled,
708            duplicate: false,
709        })
710    }
711
712    /// Suspend an `Open` session (RFC-MACP-0001 §7.5). Appends a `SessionSuspend`
713    /// annotation, transitions Open -> Suspended, and emits a lifecycle event.
714    /// The session's TTL is banked and restored on resume.
715    pub async fn suspend_session(
716        &self,
717        session_id: &str,
718        reason: &str,
719        suspended_by: &str,
720    ) -> Result<ProcessResult, MacpError> {
721        let mut guard = self.registry.sessions.write().await;
722        let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
723
724        self.maybe_expire_session(session_id, session).await?;
725        if session.state != SessionState::Open {
726            return Err(MacpError::SessionNotOpen);
727        }
728
729        let now_ms = chrono::Utc::now().timestamp_millis();
730        let payload = crate::pb::SessionSuspendPayload {
731            reason: reason.to_string(),
732            suspended_by: suspended_by.to_string(),
733        };
734        let entry = Self::make_internal_entry(
735            "SessionSuspend",
736            &prost::Message::encode_to_vec(&payload),
737            session_id,
738            &session.mode,
739        );
740        self.storage
741            .append_log_entry(session_id, &entry)
742            .await
743            .map_err(|_| MacpError::StorageFailed)?;
744        self.log_store.append(session_id, entry).await;
745        session.suspend(now_ms)?;
746        self.save_session_to_storage(session).await;
747        self.metrics.record_session_suspended(&session.mode);
748        tracing::info!(session_id, reason, "session suspended");
749        let _ = self
750            .session_lifecycle_bus
751            .send(SessionLifecycleEvent::Suspended {
752                session_id: session_id.to_string(),
753            });
754
755        Ok(ProcessResult {
756            session_state: SessionState::Suspended,
757            duplicate: false,
758        })
759    }
760
761    /// Resume a `Suspended` session (RFC-MACP-0001 §7.5), banking the suspended
762    /// duration into the TTL deadline. If the `MAX_SUSPEND_MS` cap is exceeded,
763    /// the session is force-expired instead.
764    pub async fn resume_session(
765        &self,
766        session_id: &str,
767        reason: &str,
768        resumed_by: &str,
769    ) -> Result<ProcessResult, MacpError> {
770        let mut guard = self.registry.sessions.write().await;
771        let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
772
773        if session.state != SessionState::Suspended {
774            return Err(MacpError::SessionNotOpen);
775        }
776
777        let now_ms = chrono::Utc::now().timestamp_millis();
778        let banked_before = session
779            .suspended_at_ms
780            .map(|at| (now_ms - at).max(0))
781            .unwrap_or(0);
782        let payload = crate::pb::SessionResumePayload {
783            reason: reason.to_string(),
784            resumed_by: resumed_by.to_string(),
785            banked_ms: banked_before,
786        };
787        let entry = Self::make_internal_entry(
788            "SessionResume",
789            &prost::Message::encode_to_vec(&payload),
790            session_id,
791            &session.mode,
792        );
793        self.storage
794            .append_log_entry(session_id, &entry)
795            .await
796            .map_err(|_| MacpError::StorageFailed)?;
797        self.log_store.append(session_id, entry).await;
798
799        // `resume` banks the TTL; if the suspend cap is exceeded it force-expires.
800        match session.resume(now_ms) {
801            Ok(()) => {
802                self.save_session_to_storage(session).await;
803                self.metrics.record_session_resumed(&session.mode);
804                tracing::info!(session_id, reason, "session resumed");
805                let _ = self
806                    .session_lifecycle_bus
807                    .send(SessionLifecycleEvent::Resumed {
808                        session_id: session_id.to_string(),
809                    });
810                Ok(ProcessResult {
811                    session_state: SessionState::Open,
812                    duplicate: false,
813                })
814            }
815            Err(_) => {
816                // MAX_SUSPEND_MS exceeded: the session is now Expired.
817                self.save_session_to_storage(session).await;
818                self.metrics.record_session_expired(&session.mode);
819                let _ = self
820                    .session_lifecycle_bus
821                    .send(SessionLifecycleEvent::Expired {
822                        session_id: session_id.to_string(),
823                    });
824                Err(MacpError::TtlExpired)
825            }
826        }
827    }
828
829    /// Best-effort log compaction for terminal sessions.
830    /// Returns `true` if compaction succeeded, `false` if skipped or failed.
831    async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
832        match crate::storage::compaction::compact_session_log(&*self.storage, session_id, session)
833            .await
834        {
835            Ok(()) => true,
836            Err(e) => {
837                tracing::debug!(
838                    session_id,
839                    error = %e,
840                    "log compaction skipped (backend may not support it)"
841                );
842                false
843            }
844        }
845    }
846
847    /// Force a checkpoint entry regardless of interval settings.
848    /// Used as a fallback when compaction fails on terminal sessions.
849    async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
850        let persisted = crate::registry::PersistedSession::from(session);
851        let raw_payload = match serde_json::to_vec(&persisted) {
852            Ok(bytes) => bytes,
853            Err(e) => {
854                tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
855                return;
856            }
857        };
858        let now = Utc::now().timestamp_millis();
859        let checkpoint = LogEntry {
860            message_id: String::new(),
861            received_at_ms: now,
862            sender: "_runtime".into(),
863            message_type: "Checkpoint".into(),
864            raw_payload,
865            entry_kind: EntryKind::Checkpoint,
866            session_id: session_id.into(),
867            mode: session.mode.clone(),
868            macp_version: String::new(),
869            timestamp_unix_ms: now,
870        };
871        if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
872            tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
873            return;
874        }
875        self.log_store.append(session_id, checkpoint).await;
876        tracing::debug!(
877            session_id,
878            "forced checkpoint inserted for terminal session"
879        );
880    }
881
882    /// Insert a checkpoint entry if the log has reached the configured interval.
883    async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
884        if self.checkpoint_interval == 0 {
885            return;
886        }
887        let log_len = self
888            .log_store
889            .get_log(session_id)
890            .await
891            .map(|l| l.len())
892            .unwrap_or(0);
893        // Only checkpoint at interval boundaries, and not on the first entry
894        if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
895            return;
896        }
897        self.force_insert_checkpoint(session_id, session).await;
898        tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
899    }
900
901    /// Expire all sessions that have exceeded their TTL.
902    /// Called by the background cleanup task to proactively transition
903    /// stale sessions without waiting for the next incoming message.
904    pub async fn cleanup_expired_sessions(&self) {
905        let now = Utc::now().timestamp_millis();
906        let mut guard = self.registry.sessions.write().await;
907        let expired_ids: Vec<String> = guard
908            .iter()
909            .filter(|(_, s)| s.state == SessionState::Open && now > s.ttl_expiry)
910            .map(|(id, _)| id.clone())
911            .collect();
912
913        for session_id in &expired_ids {
914            if let Some(session) = guard.get_mut(session_id) {
915                if session.state != SessionState::Open || now <= session.ttl_expiry {
916                    continue;
917                }
918                let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
919                if let Err(e) = self.storage.append_log_entry(session_id, &entry).await {
920                    tracing::warn!(
921                        session_id,
922                        error = %e,
923                        "failed to write TTL expiry during cleanup"
924                    );
925                    continue;
926                }
927                self.log_store.append(session_id, entry).await;
928                session.state = SessionState::Expired;
929                self.metrics.record_session_expired(&session.mode);
930                self.save_session_to_storage(session).await;
931                if !self.maybe_compact_log(session_id, session).await {
932                    self.force_insert_checkpoint(session_id, session).await;
933                }
934                tracing::info!(session_id, "session expired via background cleanup");
935                let _ = self
936                    .session_lifecycle_bus
937                    .send(SessionLifecycleEvent::Expired {
938                        session_id: session_id.clone(),
939                    });
940            }
941        }
942
943        if !expired_ids.is_empty() {
944            tracing::info!(
945                count = expired_ids.len(),
946                "background cleanup expired sessions"
947            );
948        }
949    }
950
951    /// Evict resolved/expired sessions older than `retention_secs` from memory.
952    /// Sessions remain queryable from durable storage even after eviction.
953    pub async fn evict_stale_sessions(&self, retention_secs: u64) {
954        let now = Utc::now().timestamp_millis();
955        let cutoff = now - (retention_secs as i64 * 1000);
956        let mut guard = self.registry.sessions.write().await;
957        let evict_ids: Vec<String> = guard
958            .iter()
959            .filter(|(_, s)| {
960                matches!(s.state, SessionState::Resolved | SessionState::Expired)
961                    && s.started_at_unix_ms < cutoff
962            })
963            .map(|(id, _)| id.clone())
964            .collect();
965
966        for id in &evict_ids {
967            guard.remove(id);
968        }
969
970        if !evict_ids.is_empty() {
971            tracing::info!(
972                count = evict_ids.len(),
973                "evicted stale sessions from memory"
974            );
975        }
976    }
977}
978
979#[cfg(test)]
980mod tests {
981    use super::*;
982    use crate::decision_pb::ProposalPayload;
983    use crate::pb::{CommitmentPayload, SessionStartPayload};
984    use prost::Message;
985
986    fn new_sid() -> String {
987        uuid::Uuid::new_v4().as_hyphenated().to_string()
988    }
989
990    fn make_runtime() -> Runtime {
991        let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
992        let registry = Arc::new(SessionRegistry::new());
993        let log_store = Arc::new(LogStore::new());
994        Runtime::new(storage, registry, log_store)
995    }
996
997    fn session_start(participants: Vec<String>) -> Vec<u8> {
998        SessionStartPayload {
999            intent: "intent".into(),
1000            participants,
1001            mode_version: "1.0.0".into(),
1002            configuration_version: "cfg-1".into(),
1003            policy_version: String::new(),
1004            ttl_ms: 1_000,
1005            context_id: String::new(),
1006            extensions: std::collections::HashMap::new(),
1007            roots: vec![],
1008        }
1009        .encode_to_vec()
1010    }
1011
1012    fn env(
1013        mode: &str,
1014        message_type: &str,
1015        message_id: &str,
1016        session_id: &str,
1017        sender: &str,
1018        payload: Vec<u8>,
1019    ) -> Envelope {
1020        Envelope {
1021            macp_version: "1.0".into(),
1022            mode: mode.into(),
1023            message_type: message_type.into(),
1024            message_id: message_id.into(),
1025            session_id: session_id.into(),
1026            sender: sender.into(),
1027            timestamp_unix_ms: Utc::now().timestamp_millis(),
1028            payload,
1029        }
1030    }
1031
1032    #[tokio::test]
1033    async fn standard_session_start_is_strict() {
1034        let rt = make_runtime();
1035        let sid = new_sid();
1036        let bad = SessionStartPayload {
1037            ttl_ms: 0,
1038            ..Default::default()
1039        }
1040        .encode_to_vec();
1041        let err = rt
1042            .process(
1043                &env(
1044                    "macp.mode.decision.v1",
1045                    "SessionStart",
1046                    "m1",
1047                    &sid,
1048                    "agent://orchestrator",
1049                    bad,
1050                ),
1051                None,
1052            )
1053            .await
1054            .unwrap_err();
1055        assert!(matches!(
1056            err,
1057            MacpError::InvalidPayload | MacpError::InvalidTtl
1058        ));
1059    }
1060
1061    #[tokio::test]
1062    async fn empty_mode_is_rejected() {
1063        let rt = make_runtime();
1064        let sid = new_sid();
1065        let err = rt
1066            .process(
1067                &env(
1068                    "",
1069                    "SessionStart",
1070                    "m1",
1071                    &sid,
1072                    "agent://orchestrator",
1073                    session_start(vec!["agent://fraud".into()]),
1074                ),
1075                None,
1076            )
1077            .await
1078            .unwrap_err();
1079        assert_eq!(err.to_string(), "InvalidEnvelope");
1080    }
1081
1082    #[tokio::test]
1083    async fn rejected_messages_do_not_enter_dedup_state() {
1084        let rt = make_runtime();
1085        let sid = new_sid();
1086        rt.process(
1087            &env(
1088                "macp.mode.decision.v1",
1089                "SessionStart",
1090                "m1",
1091                &sid,
1092                "agent://orchestrator",
1093                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1094            ),
1095            None,
1096        )
1097        .await
1098        .unwrap();
1099
1100        let bad = rt
1101            .process(
1102                &env(
1103                    "macp.mode.decision.v1",
1104                    "Proposal",
1105                    "m2",
1106                    &sid,
1107                    "agent://fraud",
1108                    b"not-protobuf".to_vec(),
1109                ),
1110                None,
1111            )
1112            .await
1113            .unwrap_err();
1114        assert_eq!(bad.to_string(), "InvalidPayload");
1115
1116        let good = ProposalPayload {
1117            proposal_id: "p1".into(),
1118            option: "step-up".into(),
1119            rationale: "risk".into(),
1120            supporting_data: vec![],
1121        }
1122        .encode_to_vec();
1123        let result = rt
1124            .process(
1125                &env(
1126                    "macp.mode.decision.v1",
1127                    "Proposal",
1128                    "m2",
1129                    &sid,
1130                    "agent://orchestrator",
1131                    good,
1132                ),
1133                None,
1134            )
1135            .await
1136            .unwrap();
1137        assert!(!result.duplicate);
1138    }
1139
1140    #[tokio::test]
1141    async fn get_session_transitions_expired_sessions() {
1142        let rt = make_runtime();
1143        let sid = new_sid();
1144        let payload = SessionStartPayload {
1145            intent: "intent".into(),
1146            participants: vec!["agent://fraud".into()],
1147            mode_version: "1.0.0".into(),
1148            configuration_version: "cfg-1".into(),
1149            policy_version: String::new(),
1150            ttl_ms: 1,
1151            context_id: String::new(),
1152            extensions: std::collections::HashMap::new(),
1153            roots: vec![],
1154        }
1155        .encode_to_vec();
1156        rt.process(
1157            &env(
1158                "macp.mode.decision.v1",
1159                "SessionStart",
1160                "m1",
1161                &sid,
1162                "agent://orchestrator",
1163                payload,
1164            ),
1165            None,
1166        )
1167        .await
1168        .unwrap();
1169        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1170        let session = rt.get_session_checked(&sid).await.unwrap();
1171        assert_eq!(session.state, SessionState::Expired);
1172    }
1173
1174    #[tokio::test]
1175    async fn multi_round_requires_standard_session_start() {
1176        let rt = make_runtime();
1177        let sid = new_sid();
1178        // multi-round is now standards-track: empty mode_version should fail
1179        let payload = SessionStartPayload {
1180            participants: vec!["creator".into(), "other".into()],
1181            ..Default::default()
1182        }
1183        .encode_to_vec();
1184        let err = rt
1185            .process(
1186                &env(
1187                    "ext.multi_round.v1",
1188                    "SessionStart",
1189                    "m1",
1190                    &sid,
1191                    "creator",
1192                    payload,
1193                ),
1194                None,
1195            )
1196            .await
1197            .unwrap_err();
1198        assert!(matches!(
1199            err,
1200            MacpError::InvalidPayload | MacpError::InvalidTtl
1201        ));
1202    }
1203
1204    #[tokio::test]
1205    async fn multi_round_valid_session_start() {
1206        let rt = make_runtime();
1207        let sid = new_sid();
1208        let payload = session_start(vec!["alice".into(), "bob".into()]);
1209        rt.process(
1210            &env(
1211                "ext.multi_round.v1",
1212                "SessionStart",
1213                "m1",
1214                &sid,
1215                "coordinator",
1216                payload,
1217            ),
1218            None,
1219        )
1220        .await
1221        .unwrap();
1222        let session = rt.get_session_checked(&sid).await.unwrap();
1223        assert_eq!(session.mode, "ext.multi_round.v1");
1224        assert_eq!(session.participants, vec!["alice", "bob"]);
1225    }
1226
1227    #[tokio::test]
1228    async fn duplicate_session_start_message_id_returns_duplicate() {
1229        let rt = make_runtime();
1230        let sid = new_sid();
1231        let payload = session_start(vec!["agent://fraud".into()]);
1232        rt.process(
1233            &env(
1234                "macp.mode.decision.v1",
1235                "SessionStart",
1236                "m1",
1237                &sid,
1238                "agent://orchestrator",
1239                payload.clone(),
1240            ),
1241            None,
1242        )
1243        .await
1244        .unwrap();
1245
1246        let result = rt
1247            .process(
1248                &env(
1249                    "macp.mode.decision.v1",
1250                    "SessionStart",
1251                    "m1",
1252                    &sid,
1253                    "agent://orchestrator",
1254                    payload,
1255                ),
1256                None,
1257            )
1258            .await
1259            .unwrap();
1260        assert!(result.duplicate);
1261    }
1262
1263    #[tokio::test]
1264    async fn non_start_mode_mismatch_rejected() {
1265        let rt = make_runtime();
1266        let sid = new_sid();
1267        rt.process(
1268            &env(
1269                "macp.mode.decision.v1",
1270                "SessionStart",
1271                "m1",
1272                &sid,
1273                "agent://orchestrator",
1274                session_start(vec!["agent://fraud".into()]),
1275            ),
1276            None,
1277        )
1278        .await
1279        .unwrap();
1280
1281        let proposal = ProposalPayload {
1282            proposal_id: "p1".into(),
1283            option: "step-up".into(),
1284            rationale: "risk".into(),
1285            supporting_data: vec![],
1286        }
1287        .encode_to_vec();
1288        let err = rt
1289            .process(
1290                &env(
1291                    "macp.mode.task.v1",
1292                    "Proposal",
1293                    "m2",
1294                    &sid,
1295                    "agent://orchestrator",
1296                    proposal,
1297                ),
1298                None,
1299            )
1300            .await
1301            .unwrap_err();
1302        assert_eq!(err.to_string(), "InvalidEnvelope");
1303    }
1304
1305    #[tokio::test]
1306    async fn cancel_idempotent_on_already_expired() {
1307        let rt = make_runtime();
1308        let sid = new_sid();
1309        let payload = SessionStartPayload {
1310            intent: "intent".into(),
1311            participants: vec!["agent://fraud".into()],
1312            mode_version: "1.0.0".into(),
1313            configuration_version: "cfg-1".into(),
1314            policy_version: String::new(),
1315            ttl_ms: 1,
1316            context_id: String::new(),
1317            extensions: std::collections::HashMap::new(),
1318            roots: vec![],
1319        }
1320        .encode_to_vec();
1321        rt.process(
1322            &env(
1323                "macp.mode.decision.v1",
1324                "SessionStart",
1325                "m1",
1326                &sid,
1327                "agent://orchestrator",
1328                payload,
1329            ),
1330            None,
1331        )
1332        .await
1333        .unwrap();
1334        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1335        let result = rt
1336            .cancel_session(&sid, "cleanup", "agent://orchestrator")
1337            .await
1338            .unwrap();
1339        assert_eq!(result.session_state, SessionState::Expired);
1340    }
1341
1342    #[tokio::test]
1343    async fn accepted_envelopes_are_published_in_order() {
1344        let rt = make_runtime();
1345        let sid = new_sid();
1346        let mut events = rt.subscribe_session_stream(&sid);
1347
1348        let start = env(
1349            "macp.mode.decision.v1",
1350            "SessionStart",
1351            "m1",
1352            &sid,
1353            "agent://orchestrator",
1354            session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1355        );
1356        rt.process(&start, None).await.unwrap();
1357        let first = events.recv().await.unwrap();
1358        assert_eq!(first.message_id, "m1");
1359        assert_eq!(first.message_type, "SessionStart");
1360
1361        let proposal = ProposalPayload {
1362            proposal_id: "p1".into(),
1363            option: "step-up".into(),
1364            rationale: "risk".into(),
1365            supporting_data: vec![],
1366        }
1367        .encode_to_vec();
1368        let proposal_env = env(
1369            "macp.mode.decision.v1",
1370            "Proposal",
1371            "m2",
1372            &sid,
1373            "agent://orchestrator",
1374            proposal,
1375        );
1376        rt.process(&proposal_env, None).await.unwrap();
1377        let second = events.recv().await.unwrap();
1378        assert_eq!(second.message_id, "m2");
1379        assert_eq!(second.message_type, "Proposal");
1380    }
1381
1382    #[tokio::test]
1383    async fn commitment_versions_are_carried_into_resolution() {
1384        let rt = make_runtime();
1385        let sid = new_sid();
1386        rt.process(
1387            &env(
1388                "macp.mode.proposal.v1",
1389                "SessionStart",
1390                "m1",
1391                &sid,
1392                "agent://buyer",
1393                session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
1394            ),
1395            None,
1396        )
1397        .await
1398        .unwrap();
1399
1400        let proposal = crate::proposal_pb::ProposalPayload {
1401            proposal_id: "p1".into(),
1402            title: "offer".into(),
1403            summary: "summary".into(),
1404            details: vec![],
1405            tags: vec![],
1406        }
1407        .encode_to_vec();
1408        rt.process(
1409            &env(
1410                "macp.mode.proposal.v1",
1411                "Proposal",
1412                "m2",
1413                &sid,
1414                "agent://seller",
1415                proposal,
1416            ),
1417            None,
1418        )
1419        .await
1420        .unwrap();
1421        let accept = crate::proposal_pb::AcceptPayload {
1422            proposal_id: "p1".into(),
1423            reason: String::new(),
1424        }
1425        .encode_to_vec();
1426        rt.process(
1427            &env(
1428                "macp.mode.proposal.v1",
1429                "Accept",
1430                "m3",
1431                &sid,
1432                "agent://seller",
1433                accept.clone(),
1434            ),
1435            None,
1436        )
1437        .await
1438        .unwrap();
1439        rt.process(
1440            &env(
1441                "macp.mode.proposal.v1",
1442                "Accept",
1443                "m4",
1444                &sid,
1445                "agent://buyer",
1446                accept,
1447            ),
1448            None,
1449        )
1450        .await
1451        .unwrap();
1452        let commitment = CommitmentPayload {
1453            commitment_id: "c1".into(),
1454            action: "proposal.accepted".into(),
1455            authority_scope: "commercial".into(),
1456            reason: "bound".into(),
1457            mode_version: "1.0.0".into(),
1458            policy_version: "policy.default".into(),
1459            configuration_version: "cfg-1".into(),
1460            outcome_positive: true,
1461            supersedes: None,
1462        }
1463        .encode_to_vec();
1464        let result = rt
1465            .process(
1466                &env(
1467                    "macp.mode.proposal.v1",
1468                    "Commitment",
1469                    "m5",
1470                    &sid,
1471                    "agent://buyer",
1472                    commitment,
1473                ),
1474                None,
1475            )
1476            .await
1477            .unwrap();
1478        assert_eq!(result.session_state, SessionState::Resolved);
1479    }
1480
1481    #[tokio::test]
1482    async fn max_open_sessions_enforced_under_write_lock() {
1483        let rt = make_runtime();
1484        let sid1 = new_sid();
1485        let sid2 = new_sid();
1486        let sid3 = new_sid();
1487        rt.process(
1488            &env(
1489                "macp.mode.decision.v1",
1490                "SessionStart",
1491                "m1",
1492                &sid1,
1493                "agent://orchestrator",
1494                session_start(vec!["agent://fraud".into()]),
1495            ),
1496            Some(1),
1497        )
1498        .await
1499        .unwrap();
1500
1501        let err = rt
1502            .process(
1503                &env(
1504                    "macp.mode.decision.v1",
1505                    "SessionStart",
1506                    "m2",
1507                    &sid2,
1508                    "agent://orchestrator",
1509                    session_start(vec!["agent://fraud".into()]),
1510                ),
1511                Some(1),
1512            )
1513            .await
1514            .unwrap_err();
1515        assert!(matches!(err, MacpError::RateLimited));
1516
1517        rt.process(
1518            &env(
1519                "macp.mode.decision.v1",
1520                "SessionStart",
1521                "m3",
1522                &sid3,
1523                "agent://other",
1524                session_start(vec!["agent://fraud".into()]),
1525            ),
1526            Some(1),
1527        )
1528        .await
1529        .unwrap();
1530    }
1531
1532    #[tokio::test]
1533    async fn weak_session_id_rejected() {
1534        let rt = make_runtime();
1535        let err = rt
1536            .process(
1537                &env(
1538                    "macp.mode.decision.v1",
1539                    "SessionStart",
1540                    "m1",
1541                    "s1",
1542                    "agent://orchestrator",
1543                    session_start(vec!["agent://fraud".into()]),
1544                ),
1545                None,
1546            )
1547            .await
1548            .unwrap_err();
1549        assert_eq!(err.to_string(), "InvalidSessionId");
1550    }
1551
1552    #[tokio::test]
1553    async fn log_append_failure_rejects_session_start() {
1554        use std::io;
1555        struct FailingBackend;
1556        #[async_trait::async_trait]
1557        impl StorageBackend for FailingBackend {
1558            async fn save_session(&self, _: &Session) -> io::Result<()> {
1559                Ok(())
1560            }
1561            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1562                Ok(None)
1563            }
1564            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1565                Ok(vec![])
1566            }
1567            async fn delete_session(&self, _: &str) -> io::Result<()> {
1568                Ok(())
1569            }
1570            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1571                Ok(vec![])
1572            }
1573            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1574                Err(io::Error::other("disk full"))
1575            }
1576            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1577                Ok(vec![])
1578            }
1579            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1580                Ok(())
1581            }
1582        }
1583
1584        let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
1585        let registry = Arc::new(SessionRegistry::new());
1586        let log_store = Arc::new(LogStore::new());
1587        let rt = Runtime::new(storage, registry, log_store);
1588        let sid = new_sid();
1589
1590        let err = rt
1591            .process(
1592                &env(
1593                    "macp.mode.decision.v1",
1594                    "SessionStart",
1595                    "m1",
1596                    &sid,
1597                    "agent://orchestrator",
1598                    session_start(vec!["agent://fraud".into()]),
1599                ),
1600                None,
1601            )
1602            .await
1603            .unwrap_err();
1604        assert_eq!(err.to_string(), "StorageFailed");
1605    }
1606
1607    #[tokio::test]
1608    async fn log_append_failure_rejects_in_session_message() {
1609        use std::io;
1610        use std::sync::atomic::{AtomicUsize, Ordering};
1611
1612        struct FailOnSecondAppend {
1613            count: AtomicUsize,
1614        }
1615        #[async_trait::async_trait]
1616        impl StorageBackend for FailOnSecondAppend {
1617            async fn save_session(&self, _: &Session) -> io::Result<()> {
1618                Ok(())
1619            }
1620            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1621                Ok(None)
1622            }
1623            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1624                Ok(vec![])
1625            }
1626            async fn delete_session(&self, _: &str) -> io::Result<()> {
1627                Ok(())
1628            }
1629            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1630                Ok(vec![])
1631            }
1632            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1633                let n = self.count.fetch_add(1, Ordering::SeqCst);
1634                if n >= 1 {
1635                    Err(io::Error::other("disk full"))
1636                } else {
1637                    Ok(())
1638                }
1639            }
1640            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1641                Ok(vec![])
1642            }
1643            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1644                Ok(())
1645            }
1646        }
1647
1648        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
1649            count: AtomicUsize::new(0),
1650        });
1651        let registry = Arc::new(SessionRegistry::new());
1652        let log_store = Arc::new(LogStore::new());
1653        let rt = Runtime::new(storage, registry, log_store);
1654        let sid = new_sid();
1655
1656        // SessionStart succeeds (first append)
1657        rt.process(
1658            &env(
1659                "macp.mode.decision.v1",
1660                "SessionStart",
1661                "m1",
1662                &sid,
1663                "agent://orchestrator",
1664                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1665            ),
1666            None,
1667        )
1668        .await
1669        .unwrap();
1670
1671        // Proposal fails (second append)
1672        let proposal = ProposalPayload {
1673            proposal_id: "p1".into(),
1674            option: "step-up".into(),
1675            rationale: "risk".into(),
1676            supporting_data: vec![],
1677        }
1678        .encode_to_vec();
1679        let err = rt
1680            .process(
1681                &env(
1682                    "macp.mode.decision.v1",
1683                    "Proposal",
1684                    "m2",
1685                    &sid,
1686                    "agent://orchestrator",
1687                    proposal,
1688                ),
1689                None,
1690            )
1691            .await
1692            .unwrap_err();
1693        assert_eq!(err.to_string(), "StorageFailed");
1694
1695        // Verify the message was not added to dedup state
1696        let session = rt.get_session_checked(&sid).await.unwrap();
1697        assert!(!session.seen_message_ids.contains("m2"));
1698    }
1699
1700    #[tokio::test]
1701    async fn cancel_session_fails_if_log_append_fails() {
1702        use std::io;
1703        use std::sync::atomic::{AtomicUsize, Ordering};
1704
1705        struct FailOnSecondAppend {
1706            count: AtomicUsize,
1707        }
1708        #[async_trait::async_trait]
1709        impl StorageBackend for FailOnSecondAppend {
1710            async fn save_session(&self, _: &Session) -> io::Result<()> {
1711                Ok(())
1712            }
1713            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1714                Ok(None)
1715            }
1716            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1717                Ok(vec![])
1718            }
1719            async fn delete_session(&self, _: &str) -> io::Result<()> {
1720                Ok(())
1721            }
1722            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1723                Ok(vec![])
1724            }
1725            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1726                let n = self.count.fetch_add(1, Ordering::SeqCst);
1727                if n >= 1 {
1728                    Err(io::Error::other("disk full"))
1729                } else {
1730                    Ok(())
1731                }
1732            }
1733            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1734                Ok(vec![])
1735            }
1736            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1737                Ok(())
1738            }
1739        }
1740
1741        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
1742            count: AtomicUsize::new(0),
1743        });
1744        let registry = Arc::new(SessionRegistry::new());
1745        let log_store = Arc::new(LogStore::new());
1746        let rt = Runtime::new(storage, registry, log_store);
1747        let sid = new_sid();
1748
1749        rt.process(
1750            &env(
1751                "macp.mode.decision.v1",
1752                "SessionStart",
1753                "m1",
1754                &sid,
1755                "agent://orchestrator",
1756                session_start(vec!["agent://fraud".into()]),
1757            ),
1758            None,
1759        )
1760        .await
1761        .unwrap();
1762
1763        let err = rt
1764            .cancel_session(&sid, "test cancel", "agent://orchestrator")
1765            .await
1766            .unwrap_err();
1767        assert_eq!(err.to_string(), "StorageFailed");
1768    }
1769
1770    #[tokio::test]
1771    async fn ttl_expiration_rejects_message() {
1772        let rt = make_runtime();
1773        let sid = new_sid();
1774        let payload = SessionStartPayload {
1775            intent: "intent".into(),
1776            participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
1777            mode_version: "1.0.0".into(),
1778            configuration_version: "cfg-1".into(),
1779            policy_version: String::new(),
1780            ttl_ms: 1,
1781            context_id: String::new(),
1782            extensions: std::collections::HashMap::new(),
1783            roots: vec![],
1784        }
1785        .encode_to_vec();
1786        rt.process(
1787            &env(
1788                "macp.mode.decision.v1",
1789                "SessionStart",
1790                "m1",
1791                &sid,
1792                "agent://orchestrator",
1793                payload,
1794            ),
1795            None,
1796        )
1797        .await
1798        .unwrap();
1799        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1800        let proposal = ProposalPayload {
1801            proposal_id: "p1".into(),
1802            option: "step-up".into(),
1803            rationale: "risk".into(),
1804            supporting_data: vec![],
1805        }
1806        .encode_to_vec();
1807        let err = rt
1808            .process(
1809                &env(
1810                    "macp.mode.decision.v1",
1811                    "Proposal",
1812                    "m2",
1813                    &sid,
1814                    "agent://orchestrator",
1815                    proposal,
1816                ),
1817                None,
1818            )
1819            .await
1820            .unwrap_err();
1821        assert_eq!(err.to_string(), "TtlExpired");
1822    }
1823
1824    #[tokio::test]
1825    async fn cleanup_expired_sessions_marks_expired() {
1826        let rt = make_runtime();
1827        let sid = new_sid();
1828        let payload = SessionStartPayload {
1829            intent: "intent".into(),
1830            participants: vec!["agent://fraud".into()],
1831            mode_version: "1.0.0".into(),
1832            configuration_version: "cfg-1".into(),
1833            policy_version: String::new(),
1834            ttl_ms: 1,
1835            context_id: String::new(),
1836            extensions: std::collections::HashMap::new(),
1837            roots: vec![],
1838        }
1839        .encode_to_vec();
1840        rt.process(
1841            &env(
1842                "macp.mode.decision.v1",
1843                "SessionStart",
1844                "m1",
1845                &sid,
1846                "agent://orchestrator",
1847                payload,
1848            ),
1849            None,
1850        )
1851        .await
1852        .unwrap();
1853        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1854        rt.cleanup_expired_sessions().await;
1855        let session = rt.get_session_checked(&sid).await.unwrap();
1856        assert_eq!(session.state, SessionState::Expired);
1857    }
1858
1859    #[tokio::test]
1860    async fn evict_stale_sessions_removes_resolved() {
1861        let rt = make_runtime();
1862        let sid = new_sid();
1863        // Start a decision session
1864        rt.process(
1865            &env(
1866                "macp.mode.decision.v1",
1867                "SessionStart",
1868                "m1",
1869                &sid,
1870                "agent://orchestrator",
1871                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1872            ),
1873            None,
1874        )
1875        .await
1876        .unwrap();
1877        // Send a Proposal
1878        let proposal = ProposalPayload {
1879            proposal_id: "p1".into(),
1880            option: "step-up".into(),
1881            rationale: "risk".into(),
1882            supporting_data: vec![],
1883        }
1884        .encode_to_vec();
1885        rt.process(
1886            &env(
1887                "macp.mode.decision.v1",
1888                "Proposal",
1889                "m2",
1890                &sid,
1891                "agent://orchestrator",
1892                proposal,
1893            ),
1894            None,
1895        )
1896        .await
1897        .unwrap();
1898        // Commit to resolve the session
1899        let commitment = CommitmentPayload {
1900            commitment_id: "c1".into(),
1901            action: "decision.selected".into(),
1902            authority_scope: "payments".into(),
1903            reason: "bound".into(),
1904            mode_version: "1.0.0".into(),
1905            policy_version: "policy.default".into(),
1906            configuration_version: "cfg-1".into(),
1907            outcome_positive: true,
1908            supersedes: None,
1909        }
1910        .encode_to_vec();
1911        let result = rt
1912            .process(
1913                &env(
1914                    "macp.mode.decision.v1",
1915                    "Commitment",
1916                    "m3",
1917                    &sid,
1918                    "agent://orchestrator",
1919                    commitment,
1920                ),
1921                None,
1922            )
1923            .await
1924            .unwrap();
1925        assert_eq!(result.session_state, SessionState::Resolved);
1926        // Wait a moment so the session's started_at_unix_ms is strictly in the past
1927        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1928        // Evict with retention = 0 (evict immediately)
1929        rt.evict_stale_sessions(0).await;
1930        // Session should no longer be in the in-memory registry
1931        assert!(rt.registry.get_session(&sid).await.is_none());
1932    }
1933
1934    #[tokio::test]
1935    async fn session_start_with_wrong_mode_version_rejected() {
1936        let rt = make_runtime();
1937        let sid = new_sid();
1938        let payload = SessionStartPayload {
1939            intent: "test".into(),
1940            participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
1941            mode_version: "99.0.0".into(), // wrong version
1942            configuration_version: "cfg-1".into(),
1943            policy_version: String::new(),
1944            ttl_ms: 60_000,
1945            context_id: String::new(),
1946            extensions: std::collections::HashMap::new(),
1947            roots: vec![],
1948        }
1949        .encode_to_vec();
1950
1951        let err = rt
1952            .process(
1953                &env(
1954                    "macp.mode.decision.v1",
1955                    "SessionStart",
1956                    "m1",
1957                    &sid,
1958                    "agent://orchestrator",
1959                    payload,
1960                ),
1961                None,
1962            )
1963            .await
1964            .unwrap_err();
1965        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
1966    }
1967
1968    #[tokio::test]
1969    async fn signal_empty_signal_type_rejected() {
1970        let rt = make_runtime();
1971        // Use non-default data so proto3 serializes a non-empty payload
1972        let signal_payload = crate::pb::SignalPayload {
1973            signal_type: String::new(),
1974            data: b"some data".to_vec(),
1975            confidence: 0.0,
1976            correlation_session_id: String::new(),
1977        }
1978        .encode_to_vec();
1979        let signal = Envelope {
1980            macp_version: "1.0".into(),
1981            mode: String::new(),
1982            message_type: "Signal".into(),
1983            message_id: "sig-1".into(),
1984            session_id: String::new(),
1985            sender: "agent://a".into(),
1986            timestamp_unix_ms: 0,
1987            payload: signal_payload,
1988        };
1989        let err = rt.process_signal(&signal).await.unwrap_err();
1990        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
1991    }
1992
1993    #[tokio::test]
1994    async fn signal_valid_payload_accepted() {
1995        let rt = make_runtime();
1996        let signal_payload = crate::pb::SignalPayload {
1997            signal_type: "heartbeat".into(),
1998            data: vec![],
1999            confidence: 0.8,
2000            correlation_session_id: String::new(),
2001        }
2002        .encode_to_vec();
2003        let signal = Envelope {
2004            macp_version: "1.0".into(),
2005            mode: String::new(),
2006            message_type: "Signal".into(),
2007            message_id: "sig-2".into(),
2008            session_id: String::new(),
2009            sender: "agent://a".into(),
2010            timestamp_unix_ms: 0,
2011            payload: signal_payload,
2012        };
2013        rt.process_signal(&signal).await.unwrap();
2014    }
2015
2016    #[tokio::test]
2017    async fn signal_empty_payload_accepted() {
2018        let rt = make_runtime();
2019        let signal = Envelope {
2020            macp_version: "1.0".into(),
2021            mode: String::new(),
2022            message_type: "Signal".into(),
2023            message_id: "sig-3".into(),
2024            session_id: String::new(),
2025            sender: "agent://a".into(),
2026            timestamp_unix_ms: 0,
2027            payload: vec![],
2028        };
2029        rt.process_signal(&signal).await.unwrap();
2030    }
2031}