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_for_mode,
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 §3.2: Replay accepted envelopes from the session log for
194    /// passive subscribe, strictly after `after_sequence` (1-based accepted
195    /// ordinal, exclusive; 0 = from the start). `Err(base)` when the
196    /// requested range was discarded by log compaction — the caller must
197    /// surface an explicit error, not silently skip missing history.
198    pub async fn get_session_envelopes_after(
199        &self,
200        session_id: &str,
201        after_sequence: u64,
202    ) -> Result<Vec<Envelope>, u64> {
203        Ok(self
204            .log_store
205            .get_incoming_after(session_id, after_sequence)
206            .await?
207            .into_iter()
208            .map(|(_idx, entry)| Envelope {
209                macp_version: if entry.macp_version.is_empty() {
210                    "1.0".into()
211                } else {
212                    entry.macp_version
213                },
214                mode: entry.mode,
215                message_type: entry.message_type,
216                message_id: entry.message_id,
217                session_id: entry.session_id,
218                sender: entry.sender,
219                timestamp_unix_ms: if entry.timestamp_unix_ms != 0 {
220                    entry.timestamp_unix_ms
221                } else {
222                    entry.received_at_ms
223                },
224                payload: entry.raw_payload,
225            })
226            .collect())
227    }
228
229    fn publish_accepted_envelope(&self, env: &Envelope) {
230        if !env.session_id.is_empty() {
231            self.stream_bus.publish(&env.session_id, env.clone());
232        }
233    }
234
235    /// Whether the session's bound policy requests info-level per-message
236    /// audit lines (`rules.audit.level == "info"`).
237    fn audit_verbose(session: &Session) -> bool {
238        session
239            .policy_definition
240            .as_ref()
241            .and_then(|p| p.rules.get("audit"))
242            .and_then(|a| a.get("level"))
243            .and_then(|l| l.as_str())
244            == Some("info")
245    }
246
247    fn make_incoming_entry(env: &Envelope, received_at_ms: i64) -> LogEntry {
248        LogEntry {
249            message_id: env.message_id.clone(),
250            received_at_ms,
251            sender: env.sender.clone(),
252            message_type: env.message_type.clone(),
253            raw_payload: env.payload.clone(),
254            entry_kind: EntryKind::Incoming,
255            session_id: env.session_id.clone(),
256            mode: env.mode.clone(),
257            macp_version: env.macp_version.clone(),
258            timestamp_unix_ms: env.timestamp_unix_ms,
259            bound_mode_version: None,
260            semantics_rev: 0,
261            bound_max_suspend_ms: None,
262            compacted_incoming_ordinals: 0,
263        }
264    }
265
266    fn make_internal_entry(
267        message_type: &str,
268        payload: &[u8],
269        session_id: &str,
270        mode: &str,
271    ) -> LogEntry {
272        let now = Utc::now().timestamp_millis();
273        LogEntry {
274            message_id: String::new(),
275            received_at_ms: now,
276            sender: "_runtime".into(),
277            message_type: message_type.into(),
278            raw_payload: payload.to_vec(),
279            entry_kind: EntryKind::Internal,
280            session_id: session_id.into(),
281            mode: mode.into(),
282            macp_version: "1.0".into(),
283            timestamp_unix_ms: now,
284            bound_mode_version: None,
285            semantics_rev: 0,
286            bound_max_suspend_ms: None,
287            compacted_incoming_ordinals: 0,
288        }
289    }
290
291    async fn save_session_to_storage(&self, session: &Session) {
292        if let Err(err) = self.storage.save_session(session).await {
293            tracing::warn!(
294                session_id = %session.session_id,
295                error = %err,
296                "failed to persist session snapshot"
297            );
298        }
299    }
300
301    async fn maybe_expire_session(
302        &self,
303        session_id: &str,
304        session: &mut Session,
305    ) -> Result<bool, MacpError> {
306        let now = Utc::now().timestamp_millis();
307        // An Open session past its deadline, or a Suspended session that has
308        // exceeded the MAX_SUSPEND_MS cap (RFC-MACP-0001 §7.5), expires.
309        let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
310            || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
311        if expires {
312            let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
313            self.storage
314                .append_log_entry(session_id, &entry)
315                .await
316                .map_err(|_| MacpError::StorageFailed)?;
317            self.log_store.append(session_id, entry).await;
318            session.state = SessionState::Expired;
319            session.suspended_at_ms = None;
320            self.metrics.record_session_expired(&session.mode);
321            tracing::info!(session_id, "session expired via TTL");
322            let _ = self
323                .session_lifecycle_bus
324                .send(SessionLifecycleEvent::Expired {
325                    session_id: session_id.to_string(),
326                });
327            return Ok(true);
328        }
329        Ok(false)
330    }
331
332    pub async fn process(
333        &self,
334        env: &Envelope,
335        max_open_sessions: Option<usize>,
336    ) -> Result<ProcessResult, MacpError> {
337        match env.message_type.as_str() {
338            "SessionStart" => self.process_session_start(env, max_open_sessions).await,
339            "Signal" | "Progress" => self.process_signal(env).await,
340            _ => self.process_message(env).await,
341        }
342    }
343
344    async fn process_session_start(
345        &self,
346        env: &Envelope,
347        max_open_sessions: Option<usize>,
348    ) -> Result<ProcessResult, MacpError> {
349        if env.mode.trim().is_empty() {
350            return Err(MacpError::InvalidEnvelope);
351        }
352        validate_session_id_for_acceptance(&env.session_id)?;
353        let mode_name = env.mode.as_str();
354        let mode = self
355            .mode_registry
356            .get_mode(mode_name)
357            .ok_or(MacpError::UnknownMode)?;
358
359        let start_payload = parse_session_start_payload(&env.payload)?;
360        // `requires_strict_session_start` stays the source of *whether* the
361        // canonical contract applies: it reads the registry's per-entry
362        // `strict_session_start` flag, which `promote_mode` sets for modes the
363        // core's static name list has never heard of. The `_for_mode` validator
364        // decides only *which* roster rule applies within that contract, and
365        // defaults to the strict one for any mode it does not recognise — so a
366        // promoted mode keeps full canonical validation.
367        let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
368        if require_complete_start {
369            validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
370        }
371
372        // Validate mode_version matches the registered descriptor's version.
373        // When the payload omits mode_version (only possible for non-strict
374        // extension modes), bind the descriptor's version instead of leaving the
375        // session bound to "" — an empty binding makes the Commitment version
376        // check vacuous (any commitment with mode_version "" would match).
377        // The bound value is recorded on the SessionStart log entry so replay
378        // uses the recorded binding, never the live registry.
379        let descriptor_version = self.mode_registry.get_mode_version(mode_name);
380        if let Some(descriptor_version) = &descriptor_version {
381            if !start_payload.mode_version.is_empty()
382                && &start_payload.mode_version != descriptor_version
383            {
384                tracing::warn!(
385                    mode = mode_name,
386                    payload_version = %start_payload.mode_version,
387                    descriptor_version = %descriptor_version,
388                    "mode_version mismatch"
389                );
390                return Err(MacpError::InvalidEnvelope);
391            }
392        }
393        let bound_mode_version: Option<String> = if start_payload.mode_version.is_empty() {
394            descriptor_version
395        } else {
396            None
397        };
398        let effective_mode_version = bound_mode_version
399            .clone()
400            .unwrap_or_else(|| start_payload.mode_version.clone());
401
402        let ttl_ms = extract_ttl_ms(&start_payload)?;
403
404        // Existing-session path: duplicate SessionStart handling. Take the
405        // shared handle under a brief map read, then check dedup under the
406        // session's own mutex (never await a session mutex while holding the
407        // map lock).
408        if let Some(existing) = self.registry.get_shared(&env.session_id).await {
409            let existing = existing.lock().await;
410            if existing.seen_message_ids.contains(&env.message_id) {
411                return Ok(ProcessResult {
412                    session_state: existing.state.clone(),
413                    duplicate: true,
414                });
415            }
416            return Err(MacpError::SessionAlreadyExists);
417        }
418
419        // Resolve the governance policy for this session.
420        // RFC-MACP-0012 §6.1: policy_version is resolved at SessionStart; empty
421        // resolves to "policy.default". The resolved PolicyDescriptor is stored
422        // immutably on the session for deterministic replay (RFC-MACP-0003 §3).
423        let effective_policy_version = if start_payload.policy_version.is_empty() {
424            crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
425        } else {
426            start_payload.policy_version.clone()
427        };
428        let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
429            Ok(policy) => {
430                // RFC 6.1: reject if policy mode doesn't match session mode
431                if policy.mode != "*" && policy.mode != mode_name {
432                    return Err(MacpError::InvalidPolicyDefinition);
433                }
434                Some(policy)
435            }
436            Err(_) => {
437                return Err(MacpError::UnknownPolicyVersion);
438            }
439        };
440
441        let accepted_at = Utc::now().timestamp_millis();
442        // RFC-MACP-0003 §2: TTL deadline is computed from the SessionStart
443        // envelope's timestamp_unix_ms, not wall-clock time. This ensures
444        // deterministic replay. Fall back to accepted_at if envelope has no timestamp.
445        let ttl_base = if env.timestamp_unix_ms > 0 {
446            env.timestamp_unix_ms
447        } else {
448            accepted_at
449        };
450        let ttl_expiry = ttl_base.saturating_add(ttl_ms);
451        // Resolve the suspension cap (RFC-MACP-0001 §7.5): the payload's
452        // positive value, else the runtime default. The RESOLVED value is
453        // bound on the session and recorded on the SessionStart log entry so
454        // replay uses it — never live configuration (RFC-MACP-0003 §2).
455        let bound_max_suspend_ms = if start_payload.max_suspend_ms > 0 {
456            start_payload.max_suspend_ms
457        } else {
458            macp_core::session::MAX_SUSPEND_MS
459        };
460        let session = Session::builder(env.session_id.clone(), mode_name, env.sender.clone())
461            .ttl_expiry(ttl_expiry)
462            .ttl_ms(ttl_ms)
463            .max_suspend_ms(bound_max_suspend_ms)
464            .started_at_unix_ms(accepted_at)
465            .participants(start_payload.participants.clone())
466            .intent(start_payload.intent.clone())
467            .mode_version(effective_mode_version)
468            .configuration_version(start_payload.configuration_version.clone())
469            .policy_version(effective_policy_version)
470            .context_id(start_payload.context_id.clone())
471            .extensions(start_payload.extensions.clone())
472            .roots(start_payload.roots.clone())
473            .policy_definition(policy_definition)
474            .build();
475
476        let response = mode.on_session_start(&session, env)?;
477        let semantics_rev = session.semantics_rev;
478
479        // Reserve the session id atomically (dedup + max_open TOCTOU safety),
480        // then do the storage I/O with the map lock RELEASED and only this
481        // session's mutex held — a slow fsync on one SessionStart no longer
482        // stalls every other session.
483        let shared = std::sync::Arc::new(tokio::sync::Mutex::new(session));
484        // Lock our own reservation BEFORE publishing it, so any concurrent
485        // access to this session id blocks until start completes or rolls back.
486        let mut session_guard = shared
487            .clone()
488            .try_lock_owned()
489            .expect("freshly created mutex is uncontended");
490        {
491            let mut map = self.registry.sessions.write().await;
492            if map.contains_key(&env.session_id) {
493                // Lost a same-id race after the earlier existence check.
494                return Err(MacpError::SessionAlreadyExists);
495            }
496            if let Some(max_open) = max_open_sessions {
497                let now = Utc::now().timestamp_millis();
498                let mut count = 0usize;
499                for arc in map.values() {
500                    // Never await a session mutex under the map lock: a
501                    // locked entry is in-flight and therefore Open —
502                    // counting it is the conservative direction for a
503                    // rate limit.
504                    let counts = match arc.try_lock() {
505                        Ok(s) => {
506                            s.initiator_sender == env.sender
507                                && s.state == SessionState::Open
508                                && now <= s.ttl_expiry
509                        }
510                        Err(_) => true,
511                    };
512                    if counts {
513                        count += 1;
514                    }
515                }
516                if count >= max_open {
517                    return Err(MacpError::RateLimited);
518                }
519            }
520            map.insert(env.session_id.clone(), std::sync::Arc::clone(&shared));
521        }
522
523        // Roll back the reservation on any storage failure: poison the
524        // placeholder (non-Open) BEFORE removing it so a waiter that already
525        // cloned the Arc fails the OPEN gate instead of processing a message
526        // for a session whose SessionStart never committed.
527        let rollback = |runtime: &Self, session_guard: &mut Session| {
528            session_guard.state = SessionState::Expired;
529            let registry = std::sync::Arc::clone(&runtime.registry);
530            let sid = env.session_id.clone();
531            async move {
532                let mut map = registry.sessions.write().await;
533                map.remove(&sid);
534            }
535        };
536
537        // 1. Create storage directory and write log entry (COMMIT POINT)
538        if self
539            .storage
540            .create_session_storage(&env.session_id)
541            .await
542            .is_err()
543        {
544            rollback(self, &mut session_guard).await;
545            return Err(MacpError::StorageFailed);
546        }
547        let mut incoming_entry = Self::make_incoming_entry(env, accepted_at);
548        incoming_entry.bound_mode_version = bound_mode_version;
549        incoming_entry.semantics_rev = semantics_rev;
550        incoming_entry.bound_max_suspend_ms = Some(bound_max_suspend_ms);
551        if self
552            .storage
553            .append_log_entry(&env.session_id, &incoming_entry)
554            .await
555            .is_err()
556        {
557            rollback(self, &mut session_guard).await;
558            return Err(MacpError::StorageFailed);
559        }
560
561        // 2. Update in-memory caches
562        self.log_store.create_session_log(&env.session_id).await;
563        self.log_store.append(&env.session_id, incoming_entry).await;
564
565        session_guard
566            .seen_message_ids
567            .insert(env.message_id.clone());
568        session_guard.apply_mode_response(response);
569
570        let result_state = session_guard.state.clone();
571        // 3. Session snapshot — best-effort AFTER the durable append. The log
572        // entry above is the COMMIT POINT: once it is durable, the session
573        // exists and replay reconstructs it, so a snapshot failure must NOT
574        // fail (or roll back) the start. The previous fatal+rollback here was
575        // incoherent past the commit point — it could not un-append the
576        // durable SessionStart, so the "failed" session resurrected on
577        // restart, and a same-id client retry appended a SECOND SessionStart
578        // that made the log unreplayable.
579        if let Err(err) = self.storage.save_session(&session_guard).await {
580            tracing::warn!(
581                session_id = %session_guard.session_id,
582                error = %err,
583                "failed to persist session snapshot at SessionStart (recoverable via replay)"
584            );
585        }
586        self.metrics.record_session_start(mode_name);
587        tracing::info!(
588            session_id = %env.session_id,
589            mode = mode_name,
590            sender = %env.sender,
591            "session started"
592        );
593        // Publish while still holding the session mutex — publish order must
594        // equal acceptance order (process_message publishes under the mutex
595        // too). Publishing after the drop let a subscriber observe a later
596        // message's broadcast BEFORE this SessionStart's, breaking the FIFO
597        // premise the subscribe-window dedupe relies on.
598        self.publish_accepted_envelope(env);
599        drop(session_guard);
600        let _ = self
601            .session_lifecycle_bus
602            .send(SessionLifecycleEvent::Created {
603                session_id: env.session_id.clone(),
604            });
605
606        Ok(ProcessResult {
607            session_state: result_state,
608            duplicate: false,
609        })
610    }
611
612    /// Process a session-scoped message following the RFC-MACP-0001 Section 7.3
613    /// terminal-state transition order:
614    /// 1. Check session OPEN
615    /// 2. Validate message (mode.authorize_sender + mode.on_message)
616    /// 3. Accept into history (log_store.append)
617    /// 4. Transition to RESOLVED (session.apply_mode_response)
618    /// 5. Reject subsequent messages (enforced by step 1 on next call)
619    async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
620        // Per-session serialization (RFC-0001 §8.1): clone the shared handle
621        // under a brief map read, then hold ONLY this session's mutex across
622        // validate + append (fsync) + commit. Different sessions' appends
623        // proceed in parallel; the same session's appends stay strictly
624        // ordered (which also keeps RocksDB's per-session next_seq
625        // read-modify-write safe).
626        let shared = self
627            .registry
628            .get_shared(&env.session_id)
629            .await
630            .ok_or(MacpError::UnknownSession)?;
631        let mut session_guard = shared.lock().await;
632        let session = &mut *session_guard;
633
634        // Per-message kernel invariants (dedup, mode-binding, TTL, the monotonic
635        // OPEN gate) live in `macp_modes::step` so any consumer of the
636        // coordination core runs the identical checks. The runtime is the first
637        // caller: it drives the phases here so it can interpose its append-only
638        // write between validation and commit (a failed write must not consume a
639        // dedup slot) — which a single all-in-one step could not preserve.
640        let now_ms = chrono::Utc::now().timestamp_millis();
641        match macp_modes::step::check_preconditions(session, env, now_ms)? {
642            macp_modes::step::Precheck::Duplicate => {
643                return Ok(ProcessResult {
644                    session_state: session.state.clone(),
645                    duplicate: true,
646                });
647            }
648            macp_modes::step::Precheck::Expired => {
649                // Durable expiry via the existing path: it appends the
650                // `TtlExpired` log entry, updates metrics/lifecycle, and marks
651                // the session Expired. `check_preconditions` and
652                // `maybe_expire_session` share the same strict `>`, OPEN-guarded
653                // rule, so this always expires.
654                let expired = self.maybe_expire_session(&env.session_id, session).await?;
655                debug_assert!(expired, "check_preconditions reported Expired");
656                self.save_session_to_storage(session).await;
657                return Err(MacpError::TtlExpired);
658            }
659            macp_modes::step::Precheck::Proceed => {}
660        }
661
662        let mode = self
663            .mode_registry
664            .get_mode(&session.mode)
665            .ok_or(MacpError::UnknownMode)?;
666        mode.authorize_sender(session, env)?;
667        // One acceptance clock for both the mode call and the log entry, so
668        // replay (which re-reads received_at_ms) observes the identical time.
669        let accepted_at_ms = Utc::now().timestamp_millis();
670        let response = mode.on_message_at(
671            session,
672            env,
673            &macp_core::mode::MessageContext::new(accepted_at_ms),
674        )?;
675
676        // 1. COMMIT POINT: write log entry to disk
677        let incoming_entry = Self::make_incoming_entry(env, accepted_at_ms);
678        self.storage
679            .append_log_entry(&env.session_id, &incoming_entry)
680            .await
681            .map_err(|_| MacpError::StorageFailed)?;
682
683        // 2. Update in-memory state via the shared commit phase (consume dedup
684        //    slot, record participant activity, apply mode response) — the exact
685        //    sequence a library consumer runs through `macp_modes::step`.
686        self.log_store.append(&env.session_id, incoming_entry).await;
687        let result_state = macp_modes::step::commit(session, env, response, now_ms);
688
689        self.metrics.record_message_accepted(&session.mode);
690        if env.message_type == "Commitment" {
691            self.metrics.record_commitment_accepted(&session.mode);
692        }
693
694        // Policy-driven audit verbosity (E3b): a bound policy may request
695        // per-message audit lines at info level via an `audit.level` rules
696        // block ("info"); default stays debug. Mode rule schemas ignore
697        // unknown blocks, so `audit` composes with any mode's rules.
698        if Self::audit_verbose(session) {
699            tracing::info!(
700                session_id = %env.session_id,
701                message_type = %env.message_type,
702                sender = %env.sender,
703                state = ?result_state,
704                "message accepted (audit)"
705            );
706        } else {
707            tracing::debug!(
708                session_id = %env.session_id,
709                message_type = %env.message_type,
710                sender = %env.sender,
711                state = ?result_state,
712                "message accepted"
713            );
714        }
715
716        if result_state == SessionState::Resolved {
717            self.metrics.record_session_resolved(&session.mode);
718            tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
719            let _ = self
720                .session_lifecycle_bus
721                .send(SessionLifecycleEvent::Resolved {
722                    session_id: env.session_id.clone(),
723                });
724        }
725
726        // 3. Best-effort session save + checkpoint
727        self.save_session_to_storage(session).await;
728        if result_state == SessionState::Resolved {
729            if !self.maybe_compact_log(&env.session_id, session).await {
730                self.force_insert_checkpoint(&env.session_id, session).await;
731            }
732        } else {
733            self.maybe_insert_checkpoint(&env.session_id, session).await;
734        }
735        self.publish_accepted_envelope(env);
736
737        Ok(ProcessResult {
738            session_state: result_state,
739            duplicate: false,
740        })
741    }
742
743    /// Process a Signal or Progress envelope. Signals are informational out-of-band
744    /// notifications. Progress messages carry structured ProgressPayload.
745    /// Neither mutates session state — both are broadcast to subscribers.
746    async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
747        // RFC-MACP-0001 §4 / RFC-MACP-0010: validate SignalPayload structure.
748        // signal_type must be non-empty when a payload is present.
749        if env.message_type == "Signal" && !env.payload.is_empty() {
750            let signal: crate::pb::SignalPayload =
751                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
752            if signal.signal_type.trim().is_empty() {
753                return Err(MacpError::InvalidPayload);
754            }
755        }
756        // RFC-MACP-0001: validate ProgressPayload structure for Progress messages.
757        if env.message_type == "Progress" && !env.payload.is_empty() {
758            let _: crate::pb::ProgressPayload =
759                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
760        }
761        tracing::debug!(
762            sender = %env.sender,
763            message_id = %env.message_id,
764            message_type = %env.message_type,
765            "signal received"
766        );
767        let _ = self.signal_bus.send(env.clone());
768        Ok(ProcessResult {
769            session_state: SessionState::Open,
770            duplicate: false,
771        })
772    }
773
774    pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
775        let shared = self.registry.get_shared(session_id).await?;
776        let mut session = shared.lock().await;
777        let changed = self
778            .maybe_expire_session(session_id, &mut session)
779            .await
780            .unwrap_or(false);
781        if changed {
782            self.save_session_to_storage(&session).await;
783        }
784        Some(session.clone())
785    }
786
787    /// Cancel a session. The `cancelled_by` parameter MUST be the authenticated
788    /// sender of the CancelSession RPC (RFC-MACP-0001 Section 7.3: CancelSession
789    /// is a Core control-plane message; mode authorization does not apply).
790    pub async fn cancel_session(
791        &self,
792        session_id: &str,
793        reason: &str,
794        cancelled_by: &str,
795    ) -> Result<ProcessResult, MacpError> {
796        let shared = self
797            .registry
798            .get_shared(session_id)
799            .await
800            .ok_or(MacpError::UnknownSession)?;
801        let mut session_guard = shared.lock().await;
802        let session = &mut *session_guard;
803
804        self.maybe_expire_session(session_id, session).await?;
805
806        // Already terminal (Resolved/Expired/Cancelled): nothing to do. An Open
807        // or Suspended session can still be cancelled (RFC-MACP-0001 §7.2/§7.3).
808        if session.state.is_terminal() {
809            let result_state = session.state.clone();
810            self.save_session_to_storage(session).await;
811            return Ok(ProcessResult {
812                session_state: result_state,
813                duplicate: false,
814            });
815        }
816
817        // RFC-MACP-0001: runtime encodes a proper SessionCancelPayload with
818        // `cancelled_by` set to the authenticated sender identity.
819        let cancel_payload = crate::pb::SessionCancelPayload {
820            reason: reason.to_string(),
821            cancelled_by: cancelled_by.to_string(),
822        };
823        let cancel_entry = Self::make_internal_entry(
824            "SessionCancel",
825            &prost::Message::encode_to_vec(&cancel_payload),
826            session_id,
827            &session.mode,
828        );
829        self.storage
830            .append_log_entry(session_id, &cancel_entry)
831            .await
832            .map_err(|_| MacpError::StorageFailed)?;
833        self.log_store.append(session_id, cancel_entry).await;
834        // RFC-MACP-0001 §7.3: cancellation terminates as CANCELLED (distinct
835        // from EXPIRED) — `cancel()` also clears any suspension marker.
836        let _ = session.cancel();
837        self.save_session_to_storage(session).await;
838        if !self.maybe_compact_log(session_id, session).await {
839            self.force_insert_checkpoint(session_id, session).await;
840        }
841        self.metrics.record_session_cancelled(&session.mode);
842        tracing::info!(session_id, reason, "session cancelled");
843        let _ = self
844            .session_lifecycle_bus
845            .send(SessionLifecycleEvent::Cancelled {
846                session_id: session_id.to_string(),
847            });
848
849        Ok(ProcessResult {
850            session_state: SessionState::Cancelled,
851            duplicate: false,
852        })
853    }
854
855    /// Suspend an `Open` session (RFC-MACP-0001 §7.5). Appends a `SessionSuspend`
856    /// annotation, transitions Open -> Suspended, and emits a lifecycle event.
857    /// The session's TTL is banked and restored on resume.
858    pub async fn suspend_session(
859        &self,
860        session_id: &str,
861        reason: &str,
862        suspended_by: &str,
863    ) -> Result<ProcessResult, MacpError> {
864        let shared = self
865            .registry
866            .get_shared(session_id)
867            .await
868            .ok_or(MacpError::UnknownSession)?;
869        let mut session_guard = shared.lock().await;
870        let session = &mut *session_guard;
871
872        self.maybe_expire_session(session_id, session).await?;
873        if session.state != SessionState::Open {
874            return Err(MacpError::SessionNotOpen);
875        }
876
877        let now_ms = chrono::Utc::now().timestamp_millis();
878        let payload = crate::pb::SessionSuspendPayload {
879            reason: reason.to_string(),
880            suspended_by: suspended_by.to_string(),
881        };
882        let entry = Self::make_internal_entry(
883            "SessionSuspend",
884            &prost::Message::encode_to_vec(&payload),
885            session_id,
886            &session.mode,
887        );
888        self.storage
889            .append_log_entry(session_id, &entry)
890            .await
891            .map_err(|_| MacpError::StorageFailed)?;
892        self.log_store.append(session_id, entry).await;
893        session.suspend(now_ms)?;
894        self.save_session_to_storage(session).await;
895        self.metrics.record_session_suspended(&session.mode);
896        tracing::info!(session_id, reason, "session suspended");
897        let _ = self
898            .session_lifecycle_bus
899            .send(SessionLifecycleEvent::Suspended {
900                session_id: session_id.to_string(),
901            });
902
903        Ok(ProcessResult {
904            session_state: SessionState::Suspended,
905            duplicate: false,
906        })
907    }
908
909    /// Resume a `Suspended` session (RFC-MACP-0001 §7.5), banking the suspended
910    /// duration into the TTL deadline. If the `MAX_SUSPEND_MS` cap is exceeded,
911    /// the session is force-expired instead.
912    pub async fn resume_session(
913        &self,
914        session_id: &str,
915        reason: &str,
916        resumed_by: &str,
917    ) -> Result<ProcessResult, MacpError> {
918        let shared = self
919            .registry
920            .get_shared(session_id)
921            .await
922            .ok_or(MacpError::UnknownSession)?;
923        let mut session_guard = shared.lock().await;
924        let session = &mut *session_guard;
925
926        if session.state != SessionState::Suspended {
927            return Err(MacpError::SessionNotOpen);
928        }
929
930        let now_ms = chrono::Utc::now().timestamp_millis();
931        let banked_before = session
932            .suspended_at_ms
933            .map(|at| (now_ms - at).max(0))
934            .unwrap_or(0);
935        let payload = crate::pb::SessionResumePayload {
936            reason: reason.to_string(),
937            resumed_by: resumed_by.to_string(),
938            banked_ms: banked_before,
939        };
940        let entry = Self::make_internal_entry(
941            "SessionResume",
942            &prost::Message::encode_to_vec(&payload),
943            session_id,
944            &session.mode,
945        );
946        self.storage
947            .append_log_entry(session_id, &entry)
948            .await
949            .map_err(|_| MacpError::StorageFailed)?;
950        self.log_store.append(session_id, entry).await;
951
952        // `resume` banks the TTL; if the suspend cap is exceeded it force-expires.
953        match session.resume(now_ms) {
954            Ok(()) => {
955                self.save_session_to_storage(session).await;
956                self.metrics.record_session_resumed(&session.mode);
957                tracing::info!(session_id, reason, "session resumed");
958                let _ = self
959                    .session_lifecycle_bus
960                    .send(SessionLifecycleEvent::Resumed {
961                        session_id: session_id.to_string(),
962                    });
963                Ok(ProcessResult {
964                    session_state: SessionState::Open,
965                    duplicate: false,
966                })
967            }
968            Err(_) => {
969                // MAX_SUSPEND_MS exceeded: the session is now Expired.
970                self.save_session_to_storage(session).await;
971                self.metrics.record_session_expired(&session.mode);
972                let _ = self
973                    .session_lifecycle_bus
974                    .send(SessionLifecycleEvent::Expired {
975                        session_id: session_id.to_string(),
976                    });
977                Err(MacpError::TtlExpired)
978            }
979        }
980    }
981
982    /// Best-effort log compaction for terminal sessions.
983    /// Returns `true` if compaction succeeded, `false` if skipped or failed.
984    async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
985        // Ordinal accounting for the sequence contract: the checkpoint must
986        // record every accepted ordinal it discards, including any base from
987        // a prior compaction recorded in the current log.
988        let discarded = match self.log_store.get_log(session_id).await {
989            Some(entries) => {
990                let prior_base: u64 = entries
991                    .iter()
992                    .filter(|e| e.entry_kind == EntryKind::Checkpoint)
993                    .map(|e| e.compacted_incoming_ordinals)
994                    .max()
995                    .unwrap_or(0);
996                prior_base
997                    + entries
998                        .iter()
999                        .filter(|e| e.entry_kind == EntryKind::Incoming)
1000                        .count() as u64
1001            }
1002            None => 0,
1003        };
1004        match crate::storage::compaction::compact_session_log(
1005            &*self.storage,
1006            session_id,
1007            session,
1008            discarded,
1009        )
1010        .await
1011        {
1012            Ok(checkpoint) => {
1013                // Keep the in-memory log in step with storage — previously
1014                // only disk was rewritten, so memory and disk diverged and
1015                // post-restart passive-subscribe history vanished silently.
1016                self.log_store
1017                    .replace_session_log(session_id, vec![checkpoint])
1018                    .await;
1019                true
1020            }
1021            Err(e) => {
1022                tracing::debug!(
1023                    session_id,
1024                    error = %e,
1025                    "log compaction skipped (backend may not support it)"
1026                );
1027                false
1028            }
1029        }
1030    }
1031
1032    /// Force a checkpoint entry regardless of interval settings.
1033    /// Used as a fallback when compaction fails on terminal sessions.
1034    async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1035        let persisted = crate::registry::PersistedSession::from(session);
1036        let raw_payload = match serde_json::to_vec(&persisted) {
1037            Ok(bytes) => bytes,
1038            Err(e) => {
1039                tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1040                return;
1041            }
1042        };
1043        let now = Utc::now().timestamp_millis();
1044        let checkpoint = LogEntry {
1045            message_id: String::new(),
1046            received_at_ms: now,
1047            sender: "_runtime".into(),
1048            message_type: "Checkpoint".into(),
1049            raw_payload,
1050            entry_kind: EntryKind::Checkpoint,
1051            session_id: session_id.into(),
1052            mode: session.mode.clone(),
1053            macp_version: String::new(),
1054            timestamp_unix_ms: now,
1055            bound_mode_version: None,
1056            semantics_rev: 0,
1057            bound_max_suspend_ms: None,
1058            compacted_incoming_ordinals: 0,
1059        };
1060        if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1061            tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1062            return;
1063        }
1064        self.log_store.append(session_id, checkpoint).await;
1065        tracing::debug!(
1066            session_id,
1067            "forced checkpoint inserted for terminal session"
1068        );
1069    }
1070
1071    /// Insert a checkpoint entry if the log has reached the configured interval.
1072    async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1073        if self.checkpoint_interval == 0 {
1074            return;
1075        }
1076        let log_len = self
1077            .log_store
1078            .get_log(session_id)
1079            .await
1080            .map(|l| l.len())
1081            .unwrap_or(0);
1082        // Only checkpoint at interval boundaries, and not on the first entry
1083        if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1084            return;
1085        }
1086        self.force_insert_checkpoint(session_id, session).await;
1087        tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1088    }
1089
1090    /// Expire all sessions that have exceeded their TTL.
1091    /// Called by the background cleanup task to proactively transition
1092    /// stale sessions without waiting for the next incoming message.
1093    pub async fn cleanup_expired_sessions(&self) {
1094        let now = Utc::now().timestamp_millis();
1095        // Snapshot the shared handles under a brief map read; never hold the
1096        // map lock across per-session locks or storage I/O. Each session is
1097        // re-checked under its own mutex (it may have been touched since the
1098        // snapshot).
1099        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1100            let guard = self.registry.sessions.read().await;
1101            guard
1102                .iter()
1103                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1104                .collect()
1105        };
1106
1107        let mut expired_count = 0usize;
1108        for (session_id, shared) in candidates {
1109            let mut session = shared.lock().await;
1110            if session.state != SessionState::Open || now <= session.ttl_expiry {
1111                continue;
1112            }
1113            let entry = Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode);
1114            if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1115                tracing::warn!(
1116                    session_id,
1117                    error = %e,
1118                    "failed to write TTL expiry during cleanup"
1119                );
1120                continue;
1121            }
1122            self.log_store.append(&session_id, entry).await;
1123            session.state = SessionState::Expired;
1124            self.metrics.record_session_expired(&session.mode);
1125            self.save_session_to_storage(&session).await;
1126            if !self.maybe_compact_log(&session_id, &session).await {
1127                self.force_insert_checkpoint(&session_id, &session).await;
1128            }
1129            expired_count += 1;
1130            tracing::info!(session_id = %session_id, "session expired via background cleanup");
1131            let _ = self
1132                .session_lifecycle_bus
1133                .send(SessionLifecycleEvent::Expired {
1134                    session_id: session_id.clone(),
1135                });
1136        }
1137
1138        if expired_count > 0 {
1139            tracing::info!(count = expired_count, "background cleanup expired sessions");
1140        }
1141    }
1142
1143    /// Delete terminal sessions' durable data older than `retention_secs`
1144    /// (opt-in via `MACP_SESSION_DISK_RETENTION_SECS`). Before this existed,
1145    /// `storage.delete_session` had no callers at all: disk grew without
1146    /// bound and every restart reloaded every session ever completed.
1147    /// Enumerates STORAGE (not memory — eviction may already have dropped the
1148    /// registry entry), deletes the session's snapshot+log, and clears any
1149    /// in-memory remnants. Returns the number of sessions deleted.
1150    pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1151        let now = Utc::now().timestamp_millis();
1152        let cutoff = now - (retention_secs as i64 * 1000);
1153        let ids = match self.storage.list_session_ids().await {
1154            Ok(ids) => ids,
1155            Err(e) => {
1156                tracing::warn!(error = %e, "disk GC: cannot list sessions");
1157                return 0;
1158            }
1159        };
1160        let mut removed = 0usize;
1161        for id in ids {
1162            // Prefer the in-memory state when present (cheap + current);
1163            // fall back to the stored snapshot for evicted sessions.
1164            let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1165                let s = shared.lock().await;
1166                s.state.is_terminal() && s.started_at_unix_ms < cutoff
1167            } else {
1168                match self.storage.load_session(&id).await {
1169                    Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1170                    // No snapshot (or unreadable): leave it for operator
1171                    // inspection rather than guessing.
1172                    _ => false,
1173                }
1174            };
1175            if !eligible {
1176                continue;
1177            }
1178            match self.storage.delete_session(&id).await {
1179                Ok(()) => {
1180                    {
1181                        let mut guard = self.registry.sessions.write().await;
1182                        guard.remove(&id);
1183                    }
1184                    self.log_store.remove_session_log(&id).await;
1185                    let _ = self.stream_bus.remove_if_unused(&id);
1186                    removed += 1;
1187                }
1188                Err(e) => {
1189                    tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1190                }
1191            }
1192        }
1193        if removed > 0 {
1194            tracing::info!(count = removed, "disk GC removed terminal sessions");
1195        }
1196        removed
1197    }
1198
1199    /// Evict resolved/expired sessions older than `retention_secs` from
1200    /// memory: the registry entry, the in-memory log cache, AND the stream
1201    /// broadcast channel (all three previously grew for the process lifetime;
1202    /// the log cache and stream bus were never evicted at all). Sessions
1203    /// remain queryable from durable storage after eviction.
1204    pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1205        let now = Utc::now().timestamp_millis();
1206        let cutoff = now - (retention_secs as i64 * 1000);
1207
1208        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1209            let guard = self.registry.sessions.read().await;
1210            guard
1211                .iter()
1212                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1213                .collect()
1214        };
1215        let mut evict_ids = Vec::new();
1216        for (id, shared) in candidates {
1217            let session = shared.lock().await;
1218            if matches!(
1219                session.state,
1220                SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1221            ) && session.started_at_unix_ms < cutoff
1222            {
1223                evict_ids.push(id);
1224            }
1225        }
1226
1227        if evict_ids.is_empty() {
1228            return;
1229        }
1230        {
1231            let mut guard = self.registry.sessions.write().await;
1232            for id in &evict_ids {
1233                guard.remove(id);
1234            }
1235        }
1236        for id in &evict_ids {
1237            self.log_store.remove_session_log(id).await;
1238            // Left in place if a subscriber is still attached; retried on the
1239            // next sweep once receivers drop.
1240            let _ = self.stream_bus.remove_if_unused(id);
1241        }
1242        tracing::info!(
1243            count = evict_ids.len(),
1244            "evicted stale sessions from memory (registry + log cache + stream bus)"
1245        );
1246    }
1247}
1248
1249#[cfg(test)]
1250mod tests {
1251    use super::*;
1252    use crate::decision_pb::ProposalPayload;
1253    use crate::pb::{CommitmentPayload, SessionStartPayload};
1254    use prost::Message;
1255
1256    fn new_sid() -> String {
1257        uuid::Uuid::new_v4().as_hyphenated().to_string()
1258    }
1259
1260    fn make_runtime() -> Runtime {
1261        let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1262        let registry = Arc::new(SessionRegistry::new());
1263        let log_store = Arc::new(LogStore::new());
1264        Runtime::new(storage, registry, log_store)
1265    }
1266
1267    fn session_start(participants: Vec<String>) -> Vec<u8> {
1268        SessionStartPayload {
1269            intent: "intent".into(),
1270            participants,
1271            mode_version: "1.0.0".into(),
1272            configuration_version: "cfg-1".into(),
1273            policy_version: String::new(),
1274            ttl_ms: 1_000,
1275            context_id: String::new(),
1276            extensions: std::collections::HashMap::new(),
1277            roots: vec![],
1278            max_suspend_ms: 0,
1279        }
1280        .encode_to_vec()
1281    }
1282
1283    fn env(
1284        mode: &str,
1285        message_type: &str,
1286        message_id: &str,
1287        session_id: &str,
1288        sender: &str,
1289        payload: Vec<u8>,
1290    ) -> Envelope {
1291        Envelope {
1292            macp_version: "1.0".into(),
1293            mode: mode.into(),
1294            message_type: message_type.into(),
1295            message_id: message_id.into(),
1296            session_id: session_id.into(),
1297            sender: sender.into(),
1298            timestamp_unix_ms: Utc::now().timestamp_millis(),
1299            payload,
1300        }
1301    }
1302
1303    #[tokio::test]
1304    async fn standard_session_start_is_strict() {
1305        let rt = make_runtime();
1306        let sid = new_sid();
1307        let bad = SessionStartPayload {
1308            ttl_ms: 0,
1309            ..Default::default()
1310        }
1311        .encode_to_vec();
1312        let err = rt
1313            .process(
1314                &env(
1315                    "macp.mode.decision.v1",
1316                    "SessionStart",
1317                    "m1",
1318                    &sid,
1319                    "agent://orchestrator",
1320                    bad,
1321                ),
1322                None,
1323            )
1324            .await
1325            .unwrap_err();
1326        assert!(matches!(
1327            err,
1328            MacpError::InvalidPayload | MacpError::InvalidTtl
1329        ));
1330    }
1331
1332    /// A **promoted** extension mode keeps the full canonical `SessionStart`
1333    /// contract, including the roster requirement.
1334    ///
1335    /// This pins the *call site's* choice of strictness source, which no other
1336    /// test covers. `ModeRegistry::requires_strict_session_start` reads a
1337    /// per-entry `strict_session_start` flag that `promote_mode` sets to `true`;
1338    /// `macp_core::session::requires_strict_session_start` reads a static name
1339    /// list that has never heard of a promoted mode's name. The two disagree
1340    /// exactly here, so routing this call through
1341    /// `validate_strict_session_start_payload` — which consults the static list
1342    /// — would silently skip canonical validation for every promoted mode while
1343    /// leaving the whole suite green. Measured: it does.
1344    #[tokio::test]
1345    async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1346        let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1347            macp_policy::DefaultPolicyEvaluator,
1348        )));
1349        mode_registry
1350            .register_extension(crate::pb::ModeDescriptor {
1351                mode: "ext.promoted.v1".into(),
1352                mode_version: "1.0.0".into(),
1353                title: "Promoted".into(),
1354                description: "promotion target".into(),
1355                determinism_class: "semantic-deterministic".into(),
1356                participant_model: "declared".into(),
1357                message_types: vec!["SessionStart".into(), "Commitment".into()],
1358                terminal_message_types: vec!["Commitment".into()],
1359                ..Default::default()
1360            })
1361            .expect("register extension");
1362        assert_eq!(
1363            mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1364            "ext.promoted.v1"
1365        );
1366        assert!(
1367            mode_registry.requires_strict_session_start("ext.promoted.v1"),
1368            "promotion must mark the entry strict"
1369        );
1370        assert!(
1371            !crate::session::requires_strict_session_start("ext.promoted.v1"),
1372            "the core's static list must NOT know this name — that disagreement is the point"
1373        );
1374
1375        let rt = Runtime::with_mode_registry(
1376            Arc::new(crate::storage::MemoryBackend),
1377            Arc::new(SessionRegistry::new()),
1378            Arc::new(LogStore::new()),
1379            mode_registry,
1380        );
1381
1382        // An empty roster: refused, because the carve-out names Decision only.
1383        let err = rt
1384            .process(
1385                &env(
1386                    "ext.promoted.v1",
1387                    "SessionStart",
1388                    "m1",
1389                    &new_sid(),
1390                    "agent://orchestrator",
1391                    session_start(vec![]),
1392                ),
1393                None,
1394            )
1395            .await
1396            .unwrap_err();
1397        assert_eq!(err.to_string(), "InvalidPayload");
1398
1399        // And the rest of the canonical contract too, so the assertion above
1400        // cannot be satisfied by a runtime that only kept the roster rule.
1401        let no_versions = SessionStartPayload {
1402            participants: vec!["agent://fraud".into()],
1403            ttl_ms: 1_000,
1404            ..Default::default()
1405        }
1406        .encode_to_vec();
1407        let err = rt
1408            .process(
1409                &env(
1410                    "ext.promoted.v1",
1411                    "SessionStart",
1412                    "m2",
1413                    &new_sid(),
1414                    "agent://orchestrator",
1415                    no_versions,
1416                ),
1417                None,
1418            )
1419            .await
1420            .unwrap_err();
1421        assert_eq!(err.to_string(), "InvalidPayload");
1422
1423        // Positive control: a complete payload is accepted, so the two refusals
1424        // above are the roster and version rules and not a broken mode.
1425        rt.process(
1426            &env(
1427                "ext.promoted.v1",
1428                "SessionStart",
1429                "m3",
1430                &new_sid(),
1431                "agent://orchestrator",
1432                session_start(vec!["agent://fraud".into()]),
1433            ),
1434            None,
1435        )
1436        .await
1437        .expect("a complete SessionStart must still be accepted for a promoted mode");
1438    }
1439
1440    #[tokio::test]
1441    async fn empty_mode_is_rejected() {
1442        let rt = make_runtime();
1443        let sid = new_sid();
1444        let err = rt
1445            .process(
1446                &env(
1447                    "",
1448                    "SessionStart",
1449                    "m1",
1450                    &sid,
1451                    "agent://orchestrator",
1452                    session_start(vec!["agent://fraud".into()]),
1453                ),
1454                None,
1455            )
1456            .await
1457            .unwrap_err();
1458        assert_eq!(err.to_string(), "InvalidEnvelope");
1459    }
1460
1461    #[tokio::test]
1462    async fn rejected_messages_do_not_enter_dedup_state() {
1463        let rt = make_runtime();
1464        let sid = new_sid();
1465        rt.process(
1466            &env(
1467                "macp.mode.decision.v1",
1468                "SessionStart",
1469                "m1",
1470                &sid,
1471                "agent://orchestrator",
1472                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1473            ),
1474            None,
1475        )
1476        .await
1477        .unwrap();
1478
1479        let bad = rt
1480            .process(
1481                &env(
1482                    "macp.mode.decision.v1",
1483                    "Proposal",
1484                    "m2",
1485                    &sid,
1486                    "agent://fraud",
1487                    b"not-protobuf".to_vec(),
1488                ),
1489                None,
1490            )
1491            .await
1492            .unwrap_err();
1493        assert_eq!(bad.to_string(), "InvalidPayload");
1494
1495        let good = ProposalPayload {
1496            proposal_id: "p1".into(),
1497            option: "step-up".into(),
1498            rationale: "risk".into(),
1499            supporting_data: vec![],
1500        }
1501        .encode_to_vec();
1502        let result = rt
1503            .process(
1504                &env(
1505                    "macp.mode.decision.v1",
1506                    "Proposal",
1507                    "m2",
1508                    &sid,
1509                    "agent://orchestrator",
1510                    good,
1511                ),
1512                None,
1513            )
1514            .await
1515            .unwrap();
1516        assert!(!result.duplicate);
1517    }
1518
1519    #[tokio::test]
1520    async fn get_session_transitions_expired_sessions() {
1521        let rt = make_runtime();
1522        let sid = new_sid();
1523        let payload = SessionStartPayload {
1524            intent: "intent".into(),
1525            participants: vec!["agent://fraud".into()],
1526            mode_version: "1.0.0".into(),
1527            configuration_version: "cfg-1".into(),
1528            policy_version: String::new(),
1529            ttl_ms: 1,
1530            context_id: String::new(),
1531            extensions: std::collections::HashMap::new(),
1532            roots: vec![],
1533            max_suspend_ms: 0,
1534        }
1535        .encode_to_vec();
1536        rt.process(
1537            &env(
1538                "macp.mode.decision.v1",
1539                "SessionStart",
1540                "m1",
1541                &sid,
1542                "agent://orchestrator",
1543                payload,
1544            ),
1545            None,
1546        )
1547        .await
1548        .unwrap();
1549        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1550        let session = rt.get_session_checked(&sid).await.unwrap();
1551        assert_eq!(session.state, SessionState::Expired);
1552    }
1553
1554    #[tokio::test]
1555    async fn multi_round_requires_standard_session_start() {
1556        let rt = make_runtime();
1557        let sid = new_sid();
1558        // multi-round is now standards-track: empty mode_version should fail
1559        let payload = SessionStartPayload {
1560            participants: vec!["creator".into(), "other".into()],
1561            ..Default::default()
1562        }
1563        .encode_to_vec();
1564        let err = rt
1565            .process(
1566                &env(
1567                    "ext.multi_round.v1",
1568                    "SessionStart",
1569                    "m1",
1570                    &sid,
1571                    "creator",
1572                    payload,
1573                ),
1574                None,
1575            )
1576            .await
1577            .unwrap_err();
1578        assert!(matches!(
1579            err,
1580            MacpError::InvalidPayload | MacpError::InvalidTtl
1581        ));
1582    }
1583
1584    #[tokio::test]
1585    async fn multi_round_valid_session_start() {
1586        let rt = make_runtime();
1587        let sid = new_sid();
1588        let payload = session_start(vec!["alice".into(), "bob".into()]);
1589        rt.process(
1590            &env(
1591                "ext.multi_round.v1",
1592                "SessionStart",
1593                "m1",
1594                &sid,
1595                "coordinator",
1596                payload,
1597            ),
1598            None,
1599        )
1600        .await
1601        .unwrap();
1602        let session = rt.get_session_checked(&sid).await.unwrap();
1603        assert_eq!(session.mode, "ext.multi_round.v1");
1604        assert_eq!(session.participants, vec!["alice", "bob"]);
1605    }
1606
1607    #[tokio::test]
1608    async fn duplicate_session_start_message_id_returns_duplicate() {
1609        let rt = make_runtime();
1610        let sid = new_sid();
1611        let payload = session_start(vec!["agent://fraud".into()]);
1612        rt.process(
1613            &env(
1614                "macp.mode.decision.v1",
1615                "SessionStart",
1616                "m1",
1617                &sid,
1618                "agent://orchestrator",
1619                payload.clone(),
1620            ),
1621            None,
1622        )
1623        .await
1624        .unwrap();
1625
1626        let result = rt
1627            .process(
1628                &env(
1629                    "macp.mode.decision.v1",
1630                    "SessionStart",
1631                    "m1",
1632                    &sid,
1633                    "agent://orchestrator",
1634                    payload,
1635                ),
1636                None,
1637            )
1638            .await
1639            .unwrap();
1640        assert!(result.duplicate);
1641    }
1642
1643    #[tokio::test]
1644    async fn non_start_mode_mismatch_rejected() {
1645        let rt = make_runtime();
1646        let sid = new_sid();
1647        rt.process(
1648            &env(
1649                "macp.mode.decision.v1",
1650                "SessionStart",
1651                "m1",
1652                &sid,
1653                "agent://orchestrator",
1654                session_start(vec!["agent://fraud".into()]),
1655            ),
1656            None,
1657        )
1658        .await
1659        .unwrap();
1660
1661        let proposal = ProposalPayload {
1662            proposal_id: "p1".into(),
1663            option: "step-up".into(),
1664            rationale: "risk".into(),
1665            supporting_data: vec![],
1666        }
1667        .encode_to_vec();
1668        let err = rt
1669            .process(
1670                &env(
1671                    "macp.mode.task.v1",
1672                    "Proposal",
1673                    "m2",
1674                    &sid,
1675                    "agent://orchestrator",
1676                    proposal,
1677                ),
1678                None,
1679            )
1680            .await
1681            .unwrap_err();
1682        assert_eq!(err.to_string(), "InvalidEnvelope");
1683    }
1684
1685    #[tokio::test]
1686    async fn cancel_idempotent_on_already_expired() {
1687        let rt = make_runtime();
1688        let sid = new_sid();
1689        let payload = SessionStartPayload {
1690            intent: "intent".into(),
1691            participants: vec!["agent://fraud".into()],
1692            mode_version: "1.0.0".into(),
1693            configuration_version: "cfg-1".into(),
1694            policy_version: String::new(),
1695            ttl_ms: 1,
1696            context_id: String::new(),
1697            extensions: std::collections::HashMap::new(),
1698            roots: vec![],
1699            max_suspend_ms: 0,
1700        }
1701        .encode_to_vec();
1702        rt.process(
1703            &env(
1704                "macp.mode.decision.v1",
1705                "SessionStart",
1706                "m1",
1707                &sid,
1708                "agent://orchestrator",
1709                payload,
1710            ),
1711            None,
1712        )
1713        .await
1714        .unwrap();
1715        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1716        let result = rt
1717            .cancel_session(&sid, "cleanup", "agent://orchestrator")
1718            .await
1719            .unwrap();
1720        assert_eq!(result.session_state, SessionState::Expired);
1721    }
1722
1723    #[tokio::test]
1724    async fn accepted_envelopes_are_published_in_order() {
1725        let rt = make_runtime();
1726        let sid = new_sid();
1727        let mut events = rt.subscribe_session_stream(&sid);
1728
1729        let start = env(
1730            "macp.mode.decision.v1",
1731            "SessionStart",
1732            "m1",
1733            &sid,
1734            "agent://orchestrator",
1735            session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1736        );
1737        rt.process(&start, None).await.unwrap();
1738        let first = events.recv().await.unwrap();
1739        assert_eq!(first.message_id, "m1");
1740        assert_eq!(first.message_type, "SessionStart");
1741
1742        let proposal = ProposalPayload {
1743            proposal_id: "p1".into(),
1744            option: "step-up".into(),
1745            rationale: "risk".into(),
1746            supporting_data: vec![],
1747        }
1748        .encode_to_vec();
1749        let proposal_env = env(
1750            "macp.mode.decision.v1",
1751            "Proposal",
1752            "m2",
1753            &sid,
1754            "agent://orchestrator",
1755            proposal,
1756        );
1757        rt.process(&proposal_env, None).await.unwrap();
1758        let second = events.recv().await.unwrap();
1759        assert_eq!(second.message_id, "m2");
1760        assert_eq!(second.message_type, "Proposal");
1761    }
1762
1763    #[tokio::test]
1764    async fn commitment_versions_are_carried_into_resolution() {
1765        let rt = make_runtime();
1766        let sid = new_sid();
1767        rt.process(
1768            &env(
1769                "macp.mode.proposal.v1",
1770                "SessionStart",
1771                "m1",
1772                &sid,
1773                "agent://buyer",
1774                session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
1775            ),
1776            None,
1777        )
1778        .await
1779        .unwrap();
1780
1781        let proposal = crate::proposal_pb::ProposalPayload {
1782            proposal_id: "p1".into(),
1783            title: "offer".into(),
1784            summary: "summary".into(),
1785            details: vec![],
1786            tags: vec![],
1787        }
1788        .encode_to_vec();
1789        rt.process(
1790            &env(
1791                "macp.mode.proposal.v1",
1792                "Proposal",
1793                "m2",
1794                &sid,
1795                "agent://seller",
1796                proposal,
1797            ),
1798            None,
1799        )
1800        .await
1801        .unwrap();
1802        let accept = crate::proposal_pb::AcceptPayload {
1803            proposal_id: "p1".into(),
1804            reason: String::new(),
1805        }
1806        .encode_to_vec();
1807        rt.process(
1808            &env(
1809                "macp.mode.proposal.v1",
1810                "Accept",
1811                "m3",
1812                &sid,
1813                "agent://seller",
1814                accept.clone(),
1815            ),
1816            None,
1817        )
1818        .await
1819        .unwrap();
1820        rt.process(
1821            &env(
1822                "macp.mode.proposal.v1",
1823                "Accept",
1824                "m4",
1825                &sid,
1826                "agent://buyer",
1827                accept,
1828            ),
1829            None,
1830        )
1831        .await
1832        .unwrap();
1833        let commitment = CommitmentPayload {
1834            commitment_id: "c1".into(),
1835            action: "proposal.accepted".into(),
1836            authority_scope: "commercial".into(),
1837            reason: "bound".into(),
1838            mode_version: "1.0.0".into(),
1839            policy_version: "policy.default".into(),
1840            configuration_version: "cfg-1".into(),
1841            outcome_positive: true,
1842            supersedes: None,
1843        }
1844        .encode_to_vec();
1845        let result = rt
1846            .process(
1847                &env(
1848                    "macp.mode.proposal.v1",
1849                    "Commitment",
1850                    "m5",
1851                    &sid,
1852                    "agent://buyer",
1853                    commitment,
1854                ),
1855                None,
1856            )
1857            .await
1858            .unwrap();
1859        assert_eq!(result.session_state, SessionState::Resolved);
1860    }
1861
1862    #[tokio::test]
1863    async fn max_open_sessions_enforced_under_write_lock() {
1864        let rt = make_runtime();
1865        let sid1 = new_sid();
1866        let sid2 = new_sid();
1867        let sid3 = new_sid();
1868        rt.process(
1869            &env(
1870                "macp.mode.decision.v1",
1871                "SessionStart",
1872                "m1",
1873                &sid1,
1874                "agent://orchestrator",
1875                session_start(vec!["agent://fraud".into()]),
1876            ),
1877            Some(1),
1878        )
1879        .await
1880        .unwrap();
1881
1882        let err = rt
1883            .process(
1884                &env(
1885                    "macp.mode.decision.v1",
1886                    "SessionStart",
1887                    "m2",
1888                    &sid2,
1889                    "agent://orchestrator",
1890                    session_start(vec!["agent://fraud".into()]),
1891                ),
1892                Some(1),
1893            )
1894            .await
1895            .unwrap_err();
1896        assert!(matches!(err, MacpError::RateLimited));
1897
1898        rt.process(
1899            &env(
1900                "macp.mode.decision.v1",
1901                "SessionStart",
1902                "m3",
1903                &sid3,
1904                "agent://other",
1905                session_start(vec!["agent://fraud".into()]),
1906            ),
1907            Some(1),
1908        )
1909        .await
1910        .unwrap();
1911    }
1912
1913    #[tokio::test]
1914    async fn weak_session_id_rejected() {
1915        let rt = make_runtime();
1916        let err = rt
1917            .process(
1918                &env(
1919                    "macp.mode.decision.v1",
1920                    "SessionStart",
1921                    "m1",
1922                    "s1",
1923                    "agent://orchestrator",
1924                    session_start(vec!["agent://fraud".into()]),
1925                ),
1926                None,
1927            )
1928            .await
1929            .unwrap_err();
1930        assert_eq!(err.to_string(), "InvalidSessionId");
1931    }
1932
1933    #[tokio::test]
1934    async fn log_append_failure_rejects_session_start() {
1935        use std::io;
1936        struct FailingBackend;
1937        #[async_trait::async_trait]
1938        impl StorageBackend for FailingBackend {
1939            async fn save_session(&self, _: &Session) -> io::Result<()> {
1940                Ok(())
1941            }
1942            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1943                Ok(None)
1944            }
1945            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1946                Ok(vec![])
1947            }
1948            async fn delete_session(&self, _: &str) -> io::Result<()> {
1949                Ok(())
1950            }
1951            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1952                Ok(vec![])
1953            }
1954            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1955                Err(io::Error::other("disk full"))
1956            }
1957            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1958                Ok(vec![])
1959            }
1960            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1961                Ok(())
1962            }
1963        }
1964
1965        let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
1966        let registry = Arc::new(SessionRegistry::new());
1967        let log_store = Arc::new(LogStore::new());
1968        let rt = Runtime::new(storage, registry, log_store);
1969        let sid = new_sid();
1970
1971        let err = rt
1972            .process(
1973                &env(
1974                    "macp.mode.decision.v1",
1975                    "SessionStart",
1976                    "m1",
1977                    &sid,
1978                    "agent://orchestrator",
1979                    session_start(vec!["agent://fraud".into()]),
1980                ),
1981                None,
1982            )
1983            .await
1984            .unwrap_err();
1985        assert_eq!(err.to_string(), "StorageFailed");
1986    }
1987
1988    #[tokio::test]
1989    async fn log_append_failure_rejects_in_session_message() {
1990        use std::io;
1991        use std::sync::atomic::{AtomicUsize, Ordering};
1992
1993        struct FailOnSecondAppend {
1994            count: AtomicUsize,
1995        }
1996        #[async_trait::async_trait]
1997        impl StorageBackend for FailOnSecondAppend {
1998            async fn save_session(&self, _: &Session) -> io::Result<()> {
1999                Ok(())
2000            }
2001            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2002                Ok(None)
2003            }
2004            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2005                Ok(vec![])
2006            }
2007            async fn delete_session(&self, _: &str) -> io::Result<()> {
2008                Ok(())
2009            }
2010            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2011                Ok(vec![])
2012            }
2013            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2014                let n = self.count.fetch_add(1, Ordering::SeqCst);
2015                if n >= 1 {
2016                    Err(io::Error::other("disk full"))
2017                } else {
2018                    Ok(())
2019                }
2020            }
2021            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2022                Ok(vec![])
2023            }
2024            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2025                Ok(())
2026            }
2027        }
2028
2029        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2030            count: AtomicUsize::new(0),
2031        });
2032        let registry = Arc::new(SessionRegistry::new());
2033        let log_store = Arc::new(LogStore::new());
2034        let rt = Runtime::new(storage, registry, log_store);
2035        let sid = new_sid();
2036
2037        // SessionStart succeeds (first append)
2038        rt.process(
2039            &env(
2040                "macp.mode.decision.v1",
2041                "SessionStart",
2042                "m1",
2043                &sid,
2044                "agent://orchestrator",
2045                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2046            ),
2047            None,
2048        )
2049        .await
2050        .unwrap();
2051
2052        // Proposal fails (second append)
2053        let proposal = ProposalPayload {
2054            proposal_id: "p1".into(),
2055            option: "step-up".into(),
2056            rationale: "risk".into(),
2057            supporting_data: vec![],
2058        }
2059        .encode_to_vec();
2060        let err = rt
2061            .process(
2062                &env(
2063                    "macp.mode.decision.v1",
2064                    "Proposal",
2065                    "m2",
2066                    &sid,
2067                    "agent://orchestrator",
2068                    proposal,
2069                ),
2070                None,
2071            )
2072            .await
2073            .unwrap_err();
2074        assert_eq!(err.to_string(), "StorageFailed");
2075
2076        // Verify the message was not added to dedup state
2077        let session = rt.get_session_checked(&sid).await.unwrap();
2078        assert!(!session.seen_message_ids.contains("m2"));
2079    }
2080
2081    #[tokio::test]
2082    async fn cancel_session_fails_if_log_append_fails() {
2083        use std::io;
2084        use std::sync::atomic::{AtomicUsize, Ordering};
2085
2086        struct FailOnSecondAppend {
2087            count: AtomicUsize,
2088        }
2089        #[async_trait::async_trait]
2090        impl StorageBackend for FailOnSecondAppend {
2091            async fn save_session(&self, _: &Session) -> io::Result<()> {
2092                Ok(())
2093            }
2094            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2095                Ok(None)
2096            }
2097            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2098                Ok(vec![])
2099            }
2100            async fn delete_session(&self, _: &str) -> io::Result<()> {
2101                Ok(())
2102            }
2103            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2104                Ok(vec![])
2105            }
2106            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2107                let n = self.count.fetch_add(1, Ordering::SeqCst);
2108                if n >= 1 {
2109                    Err(io::Error::other("disk full"))
2110                } else {
2111                    Ok(())
2112                }
2113            }
2114            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2115                Ok(vec![])
2116            }
2117            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2118                Ok(())
2119            }
2120        }
2121
2122        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2123            count: AtomicUsize::new(0),
2124        });
2125        let registry = Arc::new(SessionRegistry::new());
2126        let log_store = Arc::new(LogStore::new());
2127        let rt = Runtime::new(storage, registry, log_store);
2128        let sid = new_sid();
2129
2130        rt.process(
2131            &env(
2132                "macp.mode.decision.v1",
2133                "SessionStart",
2134                "m1",
2135                &sid,
2136                "agent://orchestrator",
2137                session_start(vec!["agent://fraud".into()]),
2138            ),
2139            None,
2140        )
2141        .await
2142        .unwrap();
2143
2144        let err = rt
2145            .cancel_session(&sid, "test cancel", "agent://orchestrator")
2146            .await
2147            .unwrap_err();
2148        assert_eq!(err.to_string(), "StorageFailed");
2149    }
2150
2151    #[tokio::test]
2152    async fn ttl_expiration_rejects_message() {
2153        let rt = make_runtime();
2154        let sid = new_sid();
2155        let payload = SessionStartPayload {
2156            intent: "intent".into(),
2157            participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2158            mode_version: "1.0.0".into(),
2159            configuration_version: "cfg-1".into(),
2160            policy_version: String::new(),
2161            ttl_ms: 1,
2162            context_id: String::new(),
2163            extensions: std::collections::HashMap::new(),
2164            roots: vec![],
2165            max_suspend_ms: 0,
2166        }
2167        .encode_to_vec();
2168        rt.process(
2169            &env(
2170                "macp.mode.decision.v1",
2171                "SessionStart",
2172                "m1",
2173                &sid,
2174                "agent://orchestrator",
2175                payload,
2176            ),
2177            None,
2178        )
2179        .await
2180        .unwrap();
2181        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2182        let proposal = ProposalPayload {
2183            proposal_id: "p1".into(),
2184            option: "step-up".into(),
2185            rationale: "risk".into(),
2186            supporting_data: vec![],
2187        }
2188        .encode_to_vec();
2189        let err = rt
2190            .process(
2191                &env(
2192                    "macp.mode.decision.v1",
2193                    "Proposal",
2194                    "m2",
2195                    &sid,
2196                    "agent://orchestrator",
2197                    proposal,
2198                ),
2199                None,
2200            )
2201            .await
2202            .unwrap_err();
2203        assert_eq!(err.to_string(), "TtlExpired");
2204    }
2205
2206    #[tokio::test]
2207    async fn cleanup_expired_sessions_marks_expired() {
2208        let rt = make_runtime();
2209        let sid = new_sid();
2210        let payload = SessionStartPayload {
2211            intent: "intent".into(),
2212            participants: vec!["agent://fraud".into()],
2213            mode_version: "1.0.0".into(),
2214            configuration_version: "cfg-1".into(),
2215            policy_version: String::new(),
2216            ttl_ms: 1,
2217            context_id: String::new(),
2218            extensions: std::collections::HashMap::new(),
2219            roots: vec![],
2220            max_suspend_ms: 0,
2221        }
2222        .encode_to_vec();
2223        rt.process(
2224            &env(
2225                "macp.mode.decision.v1",
2226                "SessionStart",
2227                "m1",
2228                &sid,
2229                "agent://orchestrator",
2230                payload,
2231            ),
2232            None,
2233        )
2234        .await
2235        .unwrap();
2236        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2237        rt.cleanup_expired_sessions().await;
2238        let session = rt.get_session_checked(&sid).await.unwrap();
2239        assert_eq!(session.state, SessionState::Expired);
2240    }
2241
2242    #[tokio::test]
2243    async fn evict_stale_sessions_removes_resolved() {
2244        let rt = make_runtime();
2245        let sid = new_sid();
2246        // Start a decision session
2247        rt.process(
2248            &env(
2249                "macp.mode.decision.v1",
2250                "SessionStart",
2251                "m1",
2252                &sid,
2253                "agent://orchestrator",
2254                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2255            ),
2256            None,
2257        )
2258        .await
2259        .unwrap();
2260        // Send a Proposal
2261        let proposal = ProposalPayload {
2262            proposal_id: "p1".into(),
2263            option: "step-up".into(),
2264            rationale: "risk".into(),
2265            supporting_data: vec![],
2266        }
2267        .encode_to_vec();
2268        rt.process(
2269            &env(
2270                "macp.mode.decision.v1",
2271                "Proposal",
2272                "m2",
2273                &sid,
2274                "agent://orchestrator",
2275                proposal,
2276            ),
2277            None,
2278        )
2279        .await
2280        .unwrap();
2281        // Commit to resolve the session
2282        let commitment = CommitmentPayload {
2283            commitment_id: "c1".into(),
2284            action: "decision.selected".into(),
2285            authority_scope: "payments".into(),
2286            reason: "bound".into(),
2287            mode_version: "1.0.0".into(),
2288            policy_version: "policy.default".into(),
2289            configuration_version: "cfg-1".into(),
2290            outcome_positive: true,
2291            supersedes: None,
2292        }
2293        .encode_to_vec();
2294        let result = rt
2295            .process(
2296                &env(
2297                    "macp.mode.decision.v1",
2298                    "Commitment",
2299                    "m3",
2300                    &sid,
2301                    "agent://orchestrator",
2302                    commitment,
2303                ),
2304                None,
2305            )
2306            .await
2307            .unwrap();
2308        assert_eq!(result.session_state, SessionState::Resolved);
2309        // Wait a moment so the session's started_at_unix_ms is strictly in the past
2310        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2311        // Evict with retention = 0 (evict immediately)
2312        rt.evict_stale_sessions(0).await;
2313        // Session should no longer be in the in-memory registry
2314        assert!(rt.registry.get_session(&sid).await.is_none());
2315    }
2316
2317    #[tokio::test]
2318    async fn session_start_with_wrong_mode_version_rejected() {
2319        let rt = make_runtime();
2320        let sid = new_sid();
2321        let payload = SessionStartPayload {
2322            intent: "test".into(),
2323            participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2324            mode_version: "99.0.0".into(), // wrong version
2325            configuration_version: "cfg-1".into(),
2326            policy_version: String::new(),
2327            ttl_ms: 60_000,
2328            context_id: String::new(),
2329            extensions: std::collections::HashMap::new(),
2330            roots: vec![],
2331            max_suspend_ms: 0,
2332        }
2333        .encode_to_vec();
2334
2335        let err = rt
2336            .process(
2337                &env(
2338                    "macp.mode.decision.v1",
2339                    "SessionStart",
2340                    "m1",
2341                    &sid,
2342                    "agent://orchestrator",
2343                    payload,
2344                ),
2345                None,
2346            )
2347            .await
2348            .unwrap_err();
2349        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2350    }
2351
2352    #[tokio::test]
2353    async fn signal_empty_signal_type_rejected() {
2354        let rt = make_runtime();
2355        // Use non-default data so proto3 serializes a non-empty payload
2356        let signal_payload = crate::pb::SignalPayload {
2357            signal_type: String::new(),
2358            data: b"some data".to_vec(),
2359            confidence: 0.0,
2360            correlation_session_id: String::new(),
2361        }
2362        .encode_to_vec();
2363        let signal = Envelope {
2364            macp_version: "1.0".into(),
2365            mode: String::new(),
2366            message_type: "Signal".into(),
2367            message_id: "sig-1".into(),
2368            session_id: String::new(),
2369            sender: "agent://a".into(),
2370            timestamp_unix_ms: 0,
2371            payload: signal_payload,
2372        };
2373        let err = rt.process_signal(&signal).await.unwrap_err();
2374        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2375    }
2376
2377    #[tokio::test]
2378    async fn signal_valid_payload_accepted() {
2379        let rt = make_runtime();
2380        let signal_payload = crate::pb::SignalPayload {
2381            signal_type: "heartbeat".into(),
2382            data: vec![],
2383            confidence: 0.8,
2384            correlation_session_id: String::new(),
2385        }
2386        .encode_to_vec();
2387        let signal = Envelope {
2388            macp_version: "1.0".into(),
2389            mode: String::new(),
2390            message_type: "Signal".into(),
2391            message_id: "sig-2".into(),
2392            session_id: String::new(),
2393            sender: "agent://a".into(),
2394            timestamp_unix_ms: 0,
2395            payload: signal_payload,
2396        };
2397        rt.process_signal(&signal).await.unwrap();
2398    }
2399
2400    #[tokio::test]
2401    async fn signal_empty_payload_accepted() {
2402        let rt = make_runtime();
2403        let signal = Envelope {
2404            macp_version: "1.0".into(),
2405            mode: String::new(),
2406            message_type: "Signal".into(),
2407            message_id: "sig-3".into(),
2408            session_id: String::new(),
2409            sender: "agent://a".into(),
2410            timestamp_unix_ms: 0,
2411            payload: vec![],
2412        };
2413        rt.process_signal(&signal).await.unwrap();
2414    }
2415
2416    /// Freeze invariant: CommitmentPayload version fields must match the
2417    /// session-bound versions — for extension modes too. When a non-strict ext
2418    /// mode's SessionStart omits mode_version, the runtime binds the registered
2419    /// descriptor's version; a Commitment carrying "" must no longer match
2420    /// vacuously.
2421    #[tokio::test]
2422    async fn ext_mode_empty_version_binds_descriptor_version() {
2423        let rt = make_runtime();
2424        rt.register_extension(ModeDescriptor {
2425            mode: "ext.dyn.v1".into(),
2426            mode_version: "2.5.0".into(),
2427            message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2428            terminal_message_types: vec!["Commitment".into()],
2429            ..Default::default()
2430        })
2431        .unwrap();
2432
2433        let sid = new_sid();
2434        let payload = SessionStartPayload {
2435            participants: vec!["alice".into()],
2436            configuration_version: "cfg-1".into(),
2437            ttl_ms: 60_000,
2438            ..Default::default()
2439        }
2440        .encode_to_vec();
2441        rt.process(
2442            &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2443            None,
2444        )
2445        .await
2446        .unwrap();
2447
2448        // The session is bound to the descriptor's version, not "".
2449        let session = rt.get_session_checked(&sid).await.unwrap();
2450        assert_eq!(session.mode_version, "2.5.0");
2451
2452        // Commitment with empty mode_version: rejected (no vacuous match).
2453        let bad = CommitmentPayload {
2454            commitment_id: "c1".into(),
2455            action: "work.completed".into(),
2456            authority_scope: "test".into(),
2457            reason: "done".into(),
2458            mode_version: String::new(),
2459            policy_version: "policy.default".into(),
2460            configuration_version: "cfg-1".into(),
2461            outcome_positive: true,
2462            supersedes: None,
2463        }
2464        .encode_to_vec();
2465        let err = rt
2466            .process(
2467                &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2468                None,
2469            )
2470            .await
2471            .unwrap_err();
2472        assert_eq!(err.to_string(), "InvalidPayload");
2473
2474        // Commitment echoing the bound descriptor version: accepted, resolves.
2475        let good = CommitmentPayload {
2476            commitment_id: "c1".into(),
2477            action: "work.completed".into(),
2478            authority_scope: "test".into(),
2479            reason: "done".into(),
2480            mode_version: "2.5.0".into(),
2481            policy_version: "policy.default".into(),
2482            configuration_version: "cfg-1".into(),
2483            outcome_positive: true,
2484            supersedes: None,
2485        }
2486        .encode_to_vec();
2487        let result = rt
2488            .process(
2489                &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2490                None,
2491            )
2492            .await
2493            .unwrap();
2494        assert_eq!(result.session_state, SessionState::Resolved);
2495    }
2496
2497    /// The binding must be recorded on the SessionStart log entry (replay reads
2498    /// it from there), and only when the payload actually omitted the version.
2499    #[tokio::test]
2500    async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2501        let rt = make_runtime();
2502        rt.register_extension(ModeDescriptor {
2503            mode: "ext.dyn2.v1".into(),
2504            mode_version: "3.0.0".into(),
2505            message_types: vec!["SessionStart".into(), "Commitment".into()],
2506            terminal_message_types: vec!["Commitment".into()],
2507            ..Default::default()
2508        })
2509        .unwrap();
2510
2511        let sid = new_sid();
2512        let payload = SessionStartPayload {
2513            participants: vec!["alice".into()],
2514            configuration_version: "cfg-1".into(),
2515            ttl_ms: 60_000,
2516            ..Default::default()
2517        }
2518        .encode_to_vec();
2519        rt.process(
2520            &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2521            None,
2522        )
2523        .await
2524        .unwrap();
2525
2526        let log = rt.log_store.get_log(&sid).await.unwrap();
2527        assert_eq!(log[0].message_type, "SessionStart");
2528        assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2529
2530        // A payload that carries the version explicitly records no binding.
2531        let sid2 = new_sid();
2532        let payload2 = SessionStartPayload {
2533            participants: vec!["alice".into()],
2534            mode_version: "3.0.0".into(),
2535            configuration_version: "cfg-1".into(),
2536            ttl_ms: 60_000,
2537            ..Default::default()
2538        }
2539        .encode_to_vec();
2540        rt.process(
2541            &env(
2542                "ext.dyn2.v1",
2543                "SessionStart",
2544                "m1",
2545                &sid2,
2546                "alice",
2547                payload2,
2548            ),
2549            None,
2550        )
2551        .await
2552        .unwrap();
2553        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2554        assert_eq!(log2[0].bound_mode_version, None);
2555    }
2556
2557    /// The RESOLVED suspension cap is bound on the session and recorded on
2558    /// the SessionStart log entry (RFC-MACP-0001 §7.5, RFC-MACP-0003 §2):
2559    /// the payload's positive value verbatim, or the runtime default when
2560    /// the payload carried 0 — never left unrecorded on new sessions.
2561    #[tokio::test]
2562    async fn session_start_binds_and_records_max_suspend_cap() {
2563        let rt = make_runtime();
2564
2565        // Explicit cap: recorded verbatim.
2566        let sid = new_sid();
2567        let payload = SessionStartPayload {
2568            participants: vec!["alice".into(), "bob".into()],
2569            mode_version: "1.0.0".into(),
2570            configuration_version: "cfg-1".into(),
2571            ttl_ms: 60_000,
2572            max_suspend_ms: 12_345,
2573            ..Default::default()
2574        }
2575        .encode_to_vec();
2576        rt.process(
2577            &env(
2578                "macp.mode.decision.v1",
2579                "SessionStart",
2580                "m1",
2581                &sid,
2582                "alice",
2583                payload,
2584            ),
2585            None,
2586        )
2587        .await
2588        .unwrap();
2589        let log = rt.log_store.get_log(&sid).await.unwrap();
2590        assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2591
2592        // Payload 0: the runtime default is resolved and recorded.
2593        let sid2 = new_sid();
2594        let payload2 = SessionStartPayload {
2595            participants: vec!["alice".into(), "bob".into()],
2596            mode_version: "1.0.0".into(),
2597            configuration_version: "cfg-1".into(),
2598            ttl_ms: 60_000,
2599            max_suspend_ms: 0,
2600            ..Default::default()
2601        }
2602        .encode_to_vec();
2603        rt.process(
2604            &env(
2605                "macp.mode.decision.v1",
2606                "SessionStart",
2607                "m2",
2608                &sid2,
2609                "alice",
2610                payload2,
2611            ),
2612            None,
2613        )
2614        .await
2615        .unwrap();
2616        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2617        assert_eq!(
2618            log2[0].bound_max_suspend_ms,
2619            Some(macp_core::session::MAX_SUSPEND_MS)
2620        );
2621    }
2622
2623    #[test]
2624    fn audit_verbosity_reads_policy_rules() {
2625        let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2626        assert!(!Runtime::audit_verbose(&session));
2627
2628        session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2629            policy_id: "policy.test.audit".into(),
2630            mode: "*".into(),
2631            description: "audited".into(),
2632            rules: serde_json::json!({ "audit": { "level": "info" } }),
2633            schema_version: 1,
2634        });
2635        assert!(Runtime::audit_verbose(&session));
2636
2637        session.policy_definition.as_mut().unwrap().rules =
2638            serde_json::json!({ "audit": { "level": "debug" } });
2639        assert!(!Runtime::audit_verbose(&session));
2640    }
2641
2642    /// Post-commit-point coherence: once the SessionStart log entry is
2643    /// durable, a snapshot failure must NOT fail (or roll back) the start —
2644    /// the previous fatal path left the durable entry behind, so the
2645    /// "failed" session resurrected on restart and a same-id retry appended
2646    /// a second SessionStart that made the log unreplayable.
2647    #[tokio::test]
2648    async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
2649        use std::io;
2650
2651        struct FailSnapshotBackend;
2652        #[async_trait::async_trait]
2653        impl StorageBackend for FailSnapshotBackend {
2654            async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
2655                Ok(())
2656            }
2657            async fn save_session(&self, _s: &Session) -> io::Result<()> {
2658                Err(io::Error::other("snapshot disk full"))
2659            }
2660            async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
2661                Ok(None)
2662            }
2663            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2664                Ok(vec![])
2665            }
2666            async fn delete_session(&self, _s: &str) -> io::Result<()> {
2667                Ok(())
2668            }
2669            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2670                Ok(vec![])
2671            }
2672            async fn append_log_entry(
2673                &self,
2674                _s: &str,
2675                _e: &crate::log_store::LogEntry,
2676            ) -> io::Result<()> {
2677                Ok(())
2678            }
2679            async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
2680                Ok(vec![])
2681            }
2682        }
2683
2684        let rt = Runtime::new(
2685            Arc::new(FailSnapshotBackend),
2686            Arc::new(SessionRegistry::new()),
2687            Arc::new(LogStore::new()),
2688        );
2689        let sid = new_sid();
2690        let result = rt
2691            .process(
2692                &env(
2693                    "macp.mode.decision.v1",
2694                    "SessionStart",
2695                    "m1",
2696                    &sid,
2697                    "agent://orchestrator",
2698                    session_start(vec!["agent://orchestrator".into()]),
2699                ),
2700                None,
2701            )
2702            .await
2703            .expect("start must succeed: the log append (commit point) succeeded");
2704        assert!(!result.duplicate);
2705        // The session exists and is usable.
2706        assert!(rt.get_session_checked(&sid).await.is_some());
2707    }
2708}