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, MAX_SUSPENSION_CYCLES,
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                    macp_core::MACP_VERSION.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    /// Build a runtime-authored (`EntryKind::Internal`) log entry stamped with
267    /// `at_ms`.
268    ///
269    /// The clock is **injected, never read here**, mirroring
270    /// [`Self::make_incoming_entry`]'s `received_at_ms`. Replay reconstructs
271    /// suspension state from these recorded stamps (`SessionSuspend` /
272    /// `SessionResume` in `replay::replay_entry`), so the entry must carry the
273    /// *same* instant the caller used to mutate the live `Session`. When this
274    /// helper read `Utc::now()` itself, `suspend_session` / `resume_session`
275    /// read the clock twice — once for `Session::suspend`/`resume`, once here —
276    /// and a tick landing between the two reads made the live
277    /// `accumulated_suspended_ms` differ from the replayed one by ~1 ms. That
278    /// value gates the handoff implicit-accept decision, so a flip within 1 ms
279    /// of the deadline could make a live-`Resolved` session fail replay
280    /// entirely and be skipped at startup (`src/main.rs`, "failed to replay
281    /// session; skipping"). One clock read per entry removes the whole class.
282    fn make_internal_entry(
283        message_type: &str,
284        payload: &[u8],
285        session_id: &str,
286        mode: &str,
287        at_ms: i64,
288    ) -> LogEntry {
289        LogEntry {
290            message_id: String::new(),
291            received_at_ms: at_ms,
292            sender: "_runtime".into(),
293            message_type: message_type.into(),
294            raw_payload: payload.to_vec(),
295            entry_kind: EntryKind::Internal,
296            session_id: session_id.into(),
297            mode: mode.into(),
298            macp_version: macp_core::MACP_VERSION.into(),
299            timestamp_unix_ms: at_ms,
300            bound_mode_version: None,
301            semantics_rev: 0,
302            bound_max_suspend_ms: None,
303            compacted_incoming_ordinals: 0,
304        }
305    }
306
307    async fn save_session_to_storage(&self, session: &Session) {
308        if let Err(err) = self.storage.save_session(session).await {
309            tracing::warn!(
310                session_id = %session.session_id,
311                error = %err,
312                "failed to persist session snapshot"
313            );
314        }
315    }
316
317    async fn maybe_expire_session(
318        &self,
319        session_id: &str,
320        session: &mut Session,
321    ) -> Result<bool, MacpError> {
322        let now = Utc::now().timestamp_millis();
323        // An Open session past its deadline, or a Suspended session that has
324        // exceeded the MAX_SUSPEND_MS cap (RFC-MACP-0001 §7.5), expires.
325        let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
326            || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
327        if expires {
328            let entry =
329                Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode, now);
330            self.storage
331                .append_log_entry(session_id, &entry)
332                .await
333                .map_err(|_| MacpError::StorageFailed)?;
334            self.log_store.append(session_id, entry).await;
335            session.state = SessionState::Expired;
336            session.suspended_at_ms = None;
337            self.metrics.record_session_expired(&session.mode);
338            tracing::info!(session_id, "session expired via TTL");
339            let _ = self
340                .session_lifecycle_bus
341                .send(SessionLifecycleEvent::Expired {
342                    session_id: session_id.to_string(),
343                });
344            return Ok(true);
345        }
346        Ok(false)
347    }
348
349    pub async fn process(
350        &self,
351        env: &Envelope,
352        max_open_sessions: Option<usize>,
353    ) -> Result<ProcessResult, MacpError> {
354        match env.message_type.as_str() {
355            "SessionStart" => self.process_session_start(env, max_open_sessions).await,
356            "Signal" | "Progress" => self.process_signal(env).await,
357            _ => self.process_message(env).await,
358        }
359    }
360
361    async fn process_session_start(
362        &self,
363        env: &Envelope,
364        max_open_sessions: Option<usize>,
365    ) -> Result<ProcessResult, MacpError> {
366        if env.mode.trim().is_empty() {
367            return Err(MacpError::InvalidEnvelope);
368        }
369        validate_session_id_for_acceptance(&env.session_id)?;
370        let mode_name = env.mode.as_str();
371        let mode = self
372            .mode_registry
373            .get_mode(mode_name)
374            .ok_or(MacpError::UnknownMode)?;
375
376        let start_payload = parse_session_start_payload(&env.payload)?;
377        // `requires_strict_session_start` stays the source of *whether* the
378        // canonical contract applies: it reads the registry's per-entry
379        // `strict_session_start` flag, which `promote_mode` sets for modes the
380        // core's static name list has never heard of. The `_for_mode` validator
381        // decides only *which* roster rule applies within that contract, and
382        // defaults to the strict one for any mode it does not recognise — so a
383        // promoted mode keeps full canonical validation.
384        let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
385        if require_complete_start {
386            validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
387        }
388
389        // Validate mode_version matches the registered descriptor's version.
390        // When the payload omits mode_version (only possible for non-strict
391        // extension modes), bind the descriptor's version instead of leaving the
392        // session bound to "" — an empty binding makes the Commitment version
393        // check vacuous (any commitment with mode_version "" would match).
394        // The bound value is recorded on the SessionStart log entry so replay
395        // uses the recorded binding, never the live registry.
396        let descriptor_version = self.mode_registry.get_mode_version(mode_name);
397        if let Some(descriptor_version) = &descriptor_version {
398            if !start_payload.mode_version.is_empty()
399                && &start_payload.mode_version != descriptor_version
400            {
401                tracing::warn!(
402                    mode = mode_name,
403                    payload_version = %start_payload.mode_version,
404                    descriptor_version = %descriptor_version,
405                    "mode_version mismatch"
406                );
407                return Err(MacpError::InvalidEnvelope);
408            }
409        }
410        let bound_mode_version: Option<String> = if start_payload.mode_version.is_empty() {
411            descriptor_version
412        } else {
413            None
414        };
415        let effective_mode_version = bound_mode_version
416            .clone()
417            .unwrap_or_else(|| start_payload.mode_version.clone());
418
419        let ttl_ms = extract_ttl_ms(&start_payload)?;
420
421        // Existing-session path: duplicate SessionStart handling. Take the
422        // shared handle under a brief map read, then check dedup under the
423        // session's own mutex (never await a session mutex while holding the
424        // map lock).
425        if let Some(existing) = self.registry.get_shared(&env.session_id).await {
426            let existing = existing.lock().await;
427            if existing.seen_message_ids.contains(&env.message_id) {
428                return Ok(ProcessResult {
429                    session_state: existing.state.clone(),
430                    duplicate: true,
431                });
432            }
433            return Err(MacpError::SessionAlreadyExists);
434        }
435
436        // Resolve the governance policy for this session.
437        // RFC-MACP-0012 §6.1: policy_version is resolved at SessionStart; empty
438        // resolves to "policy.default". The resolved PolicyDescriptor is stored
439        // immutably on the session for deterministic replay (RFC-MACP-0003 §3).
440        let effective_policy_version = if start_payload.policy_version.is_empty() {
441            crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
442        } else {
443            start_payload.policy_version.clone()
444        };
445        let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
446            Ok(policy) => {
447                // RFC 6.1: reject if policy mode doesn't match session mode
448                if policy.mode != "*" && policy.mode != mode_name {
449                    return Err(MacpError::InvalidPolicyDefinition);
450                }
451                Some(policy)
452            }
453            Err(_) => {
454                return Err(MacpError::UnknownPolicyVersion);
455            }
456        };
457
458        let accepted_at = Utc::now().timestamp_millis();
459        // RFC-MACP-0003 §2: TTL deadline is computed from the SessionStart
460        // envelope's timestamp_unix_ms, not wall-clock time. This ensures
461        // deterministic replay. Fall back to accepted_at if envelope has no timestamp.
462        let ttl_base = if env.timestamp_unix_ms > 0 {
463            env.timestamp_unix_ms
464        } else {
465            accepted_at
466        };
467        let ttl_expiry = ttl_base.saturating_add(ttl_ms);
468        // Resolve the suspension cap (RFC-MACP-0001 §7.5): the payload's
469        // positive value, else the runtime default. The RESOLVED value is
470        // bound on the session and recorded on the SessionStart log entry so
471        // replay uses it — never live configuration (RFC-MACP-0003 §2).
472        let bound_max_suspend_ms = if start_payload.max_suspend_ms > 0 {
473            start_payload.max_suspend_ms
474        } else {
475            macp_core::session::MAX_SUSPEND_MS
476        };
477        let session = Session::builder(env.session_id.clone(), mode_name, env.sender.clone())
478            .ttl_expiry(ttl_expiry)
479            .ttl_ms(ttl_ms)
480            .max_suspend_ms(bound_max_suspend_ms)
481            .started_at_unix_ms(accepted_at)
482            .participants(start_payload.participants.clone())
483            .intent(start_payload.intent.clone())
484            .mode_version(effective_mode_version)
485            .configuration_version(start_payload.configuration_version.clone())
486            .policy_version(effective_policy_version)
487            .context_id(start_payload.context_id.clone())
488            .extensions(start_payload.extensions.clone())
489            .roots(start_payload.roots.clone())
490            .policy_definition(policy_definition)
491            .build();
492
493        let response = mode.on_session_start(&session, env)?;
494        // The client boundary, on the start path too: the reserved
495        // `message_id` namespace is squattable through `SessionStart`, whose
496        // id enters `seen_message_ids` below. Runs after `on_session_start`
497        // (so existing roster rejections keep their error) and well before the
498        // commit-point append, and before the registry reservation — so a
499        // rejection needs no rollback and leaves no session behind.
500        mode.validate_client_envelope(&session, env)?;
501        let semantics_rev = session.semantics_rev;
502
503        // Reserve the session id atomically (dedup + max_open TOCTOU safety),
504        // then do the storage I/O with the map lock RELEASED and only this
505        // session's mutex held — a slow fsync on one SessionStart no longer
506        // stalls every other session.
507        let shared = std::sync::Arc::new(tokio::sync::Mutex::new(session));
508        // Lock our own reservation BEFORE publishing it, so any concurrent
509        // access to this session id blocks until start completes or rolls back.
510        let mut session_guard = shared
511            .clone()
512            .try_lock_owned()
513            .expect("freshly created mutex is uncontended");
514        {
515            let mut map = self.registry.sessions.write().await;
516            if map.contains_key(&env.session_id) {
517                // Lost a same-id race after the earlier existence check.
518                return Err(MacpError::SessionAlreadyExists);
519            }
520            if let Some(max_open) = max_open_sessions {
521                let now = Utc::now().timestamp_millis();
522                let mut count = 0usize;
523                for arc in map.values() {
524                    // Never await a session mutex under the map lock: a
525                    // locked entry is in-flight and therefore Open —
526                    // counting it is the conservative direction for a
527                    // rate limit.
528                    let counts = match arc.try_lock() {
529                        Ok(s) => {
530                            s.initiator_sender == env.sender
531                                && s.state == SessionState::Open
532                                && now <= s.ttl_expiry
533                        }
534                        Err(_) => true,
535                    };
536                    if counts {
537                        count += 1;
538                    }
539                }
540                if count >= max_open {
541                    return Err(MacpError::RateLimited);
542                }
543            }
544            map.insert(env.session_id.clone(), std::sync::Arc::clone(&shared));
545        }
546
547        // Roll back the reservation on any storage failure: poison the
548        // placeholder (non-Open) BEFORE removing it so a waiter that already
549        // cloned the Arc fails the OPEN gate instead of processing a message
550        // for a session whose SessionStart never committed.
551        let rollback = |runtime: &Self, session_guard: &mut Session| {
552            session_guard.state = SessionState::Expired;
553            let registry = std::sync::Arc::clone(&runtime.registry);
554            let sid = env.session_id.clone();
555            async move {
556                let mut map = registry.sessions.write().await;
557                map.remove(&sid);
558            }
559        };
560
561        // 1. Create storage directory and write log entry (COMMIT POINT)
562        if self
563            .storage
564            .create_session_storage(&env.session_id)
565            .await
566            .is_err()
567        {
568            rollback(self, &mut session_guard).await;
569            return Err(MacpError::StorageFailed);
570        }
571        let mut incoming_entry = Self::make_incoming_entry(env, accepted_at);
572        incoming_entry.bound_mode_version = bound_mode_version;
573        incoming_entry.semantics_rev = semantics_rev;
574        incoming_entry.bound_max_suspend_ms = Some(bound_max_suspend_ms);
575        if self
576            .storage
577            .append_log_entry(&env.session_id, &incoming_entry)
578            .await
579            .is_err()
580        {
581            rollback(self, &mut session_guard).await;
582            return Err(MacpError::StorageFailed);
583        }
584
585        // 2. Update in-memory caches
586        self.log_store.create_session_log(&env.session_id).await;
587        self.log_store.append(&env.session_id, incoming_entry).await;
588
589        session_guard
590            .seen_message_ids
591            .insert(env.message_id.clone());
592        session_guard.apply_mode_response(response);
593
594        let result_state = session_guard.state.clone();
595        // 3. Session snapshot — best-effort AFTER the durable append. The log
596        // entry above is the COMMIT POINT: once it is durable, the session
597        // exists and replay reconstructs it, so a snapshot failure must NOT
598        // fail (or roll back) the start. The previous fatal+rollback here was
599        // incoherent past the commit point — it could not un-append the
600        // durable SessionStart, so the "failed" session resurrected on
601        // restart, and a same-id client retry appended a SECOND SessionStart
602        // that made the log unreplayable.
603        if let Err(err) = self.storage.save_session(&session_guard).await {
604            tracing::warn!(
605                session_id = %session_guard.session_id,
606                error = %err,
607                "failed to persist session snapshot at SessionStart (recoverable via replay)"
608            );
609        }
610        self.metrics.record_session_start(mode_name);
611        tracing::info!(
612            session_id = %env.session_id,
613            mode = mode_name,
614            sender = %env.sender,
615            "session started"
616        );
617        // Publish while still holding the session mutex — publish order must
618        // equal acceptance order (process_message publishes under the mutex
619        // too). Publishing after the drop let a subscriber observe a later
620        // message's broadcast BEFORE this SessionStart's, breaking the FIFO
621        // premise the subscribe-window dedupe relies on.
622        self.publish_accepted_envelope(env);
623        drop(session_guard);
624        let _ = self
625            .session_lifecycle_bus
626            .send(SessionLifecycleEvent::Created {
627                session_id: env.session_id.clone(),
628            });
629
630        Ok(ProcessResult {
631            session_state: result_state,
632            duplicate: false,
633        })
634    }
635
636    /// Emit the mode's due synthetic envelope, if one is due, into accepted
637    /// history — the kernel half of [`macp_modes::mode::Mode::due_synthetic_envelope`]'s
638    /// contract (RFC-MACP-0010 §5.1(2) for handoff's implicit accept).
639    ///
640    /// Called from `process_message` *before* the triggering message is
641    /// dispatched, sharing the trigger's single clock read. The second caller
642    /// is [`Runtime::sweep_due_synthetic_accepts`], the eager sweep driven by
643    /// the background maintenance loop in `src/main.rs` (alongside
644    /// `cleanup_expired_sessions`, `evict_stale_sessions` and
645    /// `gc_disk_sessions`): it calls this for every open session so an offer is
646    /// settled on time rather than only when the next message happens to
647    /// arrive. Both callers must hold the session mutex.
648    ///
649    /// Returns `true` when an entry was appended — the sweep's emission count,
650    /// and the only honest one: every skip below is a silent `Ok`.
651    ///
652    /// # The non-`Open` filter is a correctness gate, not hygiene
653    ///
654    /// Both computations feeding the synthetic entry deliberately ignore an
655    /// *in-flight* pause: `Session::unsuspended_deadline` walks completed
656    /// intervals only, and `rev2_elapsed_ms` has no in-flight term. Asking a
657    /// `Suspended` session therefore over-counts elapsed time (the accept can
658    /// be judged due when it is not) and under-computes `D` (a wrong
659    /// `timestamp_unix_ms` baked into permanent history). The mode carries a
660    /// `debug_assert!` for this, which compiles out in release builds; this
661    /// filter is the enforcement that ships.
662    ///
663    /// # Clocks
664    ///
665    /// The entry is dispatched *and* stamped with the envelope's own
666    /// `timestamp_unix_ms` (the computed deadline `D`), never wall-clock:
667    /// `received_at_ms == timestamp_unix_ms == D`. Replay derives its dispatch
668    /// clock from `received_at_ms`, so any other value forks the live session
669    /// from what the log rebuilds.
670    ///
671    /// # No `step::commit`
672    ///
673    /// The commit here is `seen_message_ids` + `apply_mode_response` only.
674    /// `step::commit` would also call `record_participant_activity`, which
675    /// replay never does for any entry kind — so calling it live would
676    /// guarantee a live/replay divergence in `participant_message_counts` /
677    /// `participant_last_seen` — and crediting the target with activity they
678    /// did not perform would be false besides. The consequence is deliberate
679    /// and documented: the target's `message_count` in
680    /// `SessionMetadata.participant_activity` does not include the synthetic
681    /// accept.
682    ///
683    /// # The save is load-bearing
684    ///
685    /// `save_session_to_storage` is called *here*, not left to the caller: the
686    /// trigger's own save is downstream of its `on_message_at`, which returns
687    /// early when the mode rejects the trigger. Without this call, a synthesis
688    /// followed by a rejected trigger would leave the in-memory session
689    /// carrying new `mode_state` and a new dedup id that the on-disk snapshot
690    /// lacks — and `replay::validate_replay_consistency` would warn about it on
691    /// the next startup. The in-tree precedent is the `Precheck::Expired` arm
692    /// below, which saves before returning `Err` for the same reason.
693    ///
694    /// # Freeze-profile carve-out
695    ///
696    /// The tracked invariant "rejected messages don't consume dedup slots or
697    /// mutate history" (`CONTRIBUTING.md`) is amended, not broken, by this
698    /// method: a *rejected* trigger can now leave a runtime-originated entry in
699    /// accepted history. Three arguments, recorded here because this is what
700    /// stops the next reader from reverting it:
701    ///
702    /// 1. **The precedent is already shipped.** `Precheck::Expired` does all of
703    ///    this on a message it then rejects — `maybe_expire_session` appends a
704    ///    durable `TtlExpired` entry and mutates `session.state`, then
705    ///    `process_message` saves the snapshot and returns `Err`. The honest
706    ///    delta is only that this observation lands in **accepted** history
707    ///    (`EntryKind::Incoming`, so it consumes an accepted ordinal) and is
708    ///    **published to `StreamSession`** subscribers; `TtlExpired` is
709    ///    `EntryKind::Internal` and does neither.
710    /// 2. **The "conservative" alternative is the non-conformant one.**
711    ///    Synthesizing only for *accepted* triggers inverts RFC-MACP-0010
712    ///    §5.1(4): a late explicit `HandoffAccept` would then be validated
713    ///    against a still-unaccepted offer, pass, be accepted — and the
714    ///    synthetic would never be emitted at all. §5.1(2) requires the
715    ///    synthetic in history *before* any subsequent message is evaluated
716    ///    against the offer's acceptance state.
717    /// 3. **The dedup half is preserved exactly.** The rejected trigger's own
718    ///    `message_id` is never inserted — only the synthetic's deterministic
719    ///    id is — so re-sending a corrected message with that same id is still
720    ///    accepted. No existing dedup-invariant test needed weakening.
721    async fn synthesize_due_accept(
722        &self,
723        session_id: &str,
724        session: &mut Session,
725        now_ms: i64,
726    ) -> Result<bool, MacpError> {
727        if session.state != SessionState::Open {
728            return Ok(false);
729        }
730        let Some(mode) = self.mode_registry.get_mode(&session.mode) else {
731            return Ok(false);
732        };
733        let Some(syn) = mode.due_synthetic_envelope(session, now_ms) else {
734            return Ok(false);
735        };
736        // Idempotence backstop. The mode's own contract already returns `None`
737        // once the envelope has been applied; this makes a mode that forgets
738        // append a duplicate entry impossible rather than merely unlikely.
739        if session.seen_message_ids.contains(&syn.message_id) {
740            return Ok(false);
741        }
742        // Mirrors replay, which authorizes every `Incoming` entry. The target
743        // is a declared participant by offer validation, so this always passes
744        // for handoff — but the seam is general.
745        mode.authorize_sender(session, &syn)?;
746        let response = mode.on_message_at(
747            session,
748            &syn,
749            &macp_core::mode::MessageContext::new(syn.timestamp_unix_ms),
750        )?;
751
752        // COMMIT POINT: nothing above has mutated the session, so a failed
753        // append rejects the *triggering* message with `StorageFailed` before
754        // it consumed a dedup slot. RFC-MACP-0010 §5.1(2) forbids evaluating
755        // that message without the accept in history, and this runtime never
756        // acknowledges what it could not persist.
757        let entry = Self::make_incoming_entry(&syn, syn.timestamp_unix_ms);
758        self.storage
759            .append_log_entry(session_id, &entry)
760            .await
761            .map_err(|_| MacpError::StorageFailed)?;
762        self.log_store.append(session_id, entry).await;
763
764        session.seen_message_ids.insert(syn.message_id.clone());
765        session.apply_mode_response(response);
766        self.metrics.record_message_accepted(&session.mode);
767
768        tracing::info!(
769            session_id = %session_id,
770            message_type = %syn.message_type,
771            message_id = %syn.message_id,
772            sender = %syn.sender,
773            deadline_ms = syn.timestamp_unix_ms,
774            "synthetic envelope appended to accepted history"
775        );
776
777        self.save_session_to_storage(session).await;
778        // The checkpoint-interval check belongs to *the append*, not to the
779        // caller. It lives inside this seam rather than in either caller so the
780        // two paths cannot drift: a synthetic entry advances `log_len` exactly
781        // the same way from `process_message` and from
782        // `sweep_due_synthetic_accepts`, so it must cross an interval boundary
783        // the same way too. Put it in the sweep instead and the eager path
784        // would silently skip every boundary a synthetic crossed — checkpoints
785        // are only a replay optimization, but the divergence would be real and
786        // invisible.
787        //
788        // No double-checkpoint on the lazy path. `process_message` runs its own
789        // `maybe_insert_checkpoint` after appending the triggering message, by
790        // which point the log is at least one entry longer than it is here, so
791        // the two calls never test the same `log_len`. The check is a pure
792        // function of that length, so re-asking at a different length is not a
793        // repeat.
794        self.maybe_insert_checkpoint(session_id, session).await;
795        self.publish_accepted_envelope(&syn);
796        Ok(true)
797    }
798
799    /// Process a session-scoped message following the RFC-MACP-0001 Section 7.3
800    /// terminal-state transition order:
801    /// 1. Check session OPEN
802    /// 2. Validate message (mode.authorize_sender + mode.on_message)
803    /// 3. Accept into history (log_store.append)
804    /// 4. Transition to RESOLVED (session.apply_mode_response)
805    /// 5. Reject subsequent messages (enforced by step 1 on next call)
806    async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
807        // Per-session serialization (RFC-0001 §8.1): clone the shared handle
808        // under a brief map read, then hold ONLY this session's mutex across
809        // validate + append (fsync) + commit. Different sessions' appends
810        // proceed in parallel; the same session's appends stay strictly
811        // ordered (which also keeps RocksDB's per-session next_seq
812        // read-modify-write safe).
813        let shared = self
814            .registry
815            .get_shared(&env.session_id)
816            .await
817            .ok_or(MacpError::UnknownSession)?;
818        let mut session_guard = shared.lock().await;
819        let session = &mut *session_guard;
820
821        // Per-message kernel invariants (dedup, mode-binding, TTL, the monotonic
822        // OPEN gate) live in `macp_modes::step` so any consumer of the
823        // coordination core runs the identical checks. The runtime is the first
824        // caller: it drives the phases here so it can interpose its append-only
825        // write between validation and commit (a failed write must not consume a
826        // dedup slot) — which a single all-in-one step could not preserve.
827        let now_ms = chrono::Utc::now().timestamp_millis();
828        match macp_modes::step::check_preconditions(session, env, now_ms)? {
829            macp_modes::step::Precheck::Duplicate => {
830                return Ok(ProcessResult {
831                    session_state: session.state.clone(),
832                    duplicate: true,
833                });
834            }
835            macp_modes::step::Precheck::Expired => {
836                // Durable expiry via the existing path: it appends the
837                // `TtlExpired` log entry, updates metrics/lifecycle, and marks
838                // the session Expired. `check_preconditions` and
839                // `maybe_expire_session` share the same strict `>`, OPEN-guarded
840                // rule, so this always expires.
841                let expired = self.maybe_expire_session(&env.session_id, session).await?;
842                debug_assert!(expired, "check_preconditions reported Expired");
843                self.save_session_to_storage(session).await;
844                return Err(MacpError::TtlExpired);
845            }
846            macp_modes::step::Precheck::Proceed => {}
847        }
848
849        let mode = self
850            .mode_registry
851            .get_mode(&session.mode)
852            .ok_or(MacpError::UnknownMode)?;
853        mode.authorize_sender(session, env)?;
854        // The client boundary (RFC-MACP-0010 §5.1(3)): shapes that are legal
855        // as recorded history but illegal as client submissions. Called here
856        // and NEVER on replay — replay re-reads entries that already passed
857        // this check when they were first accepted, and a runtime-synthesized
858        // envelope is not a client submission either. Placed after
859        // `authorize_sender` so the pre-existing Forbidden-before-InvalidPayload
860        // ordering is unchanged. This call is not optional plumbing: the
861        // runtime deliberately bypasses `macp_modes::step::validate_message`
862        // (which also calls the hook) so it can interpose the durable append
863        // between validation and commit, so without this line the runtime
864        // would have no client boundary at all.
865        mode.validate_client_envelope(session, env)?;
866        // One acceptance clock for both the mode call and the log entry, so
867        // replay (which re-reads received_at_ms) observes the identical time.
868        let accepted_at_ms = Utc::now().timestamp_millis();
869        // RFC-MACP-0010 §5.1(2): anything the mode says is already due MUST
870        // enter accepted history BEFORE this message is evaluated against the
871        // state it changes. Lazy emission — the same clock reading the trigger
872        // is about to be dispatched with. Deliberately after the prechecks (a
873        // duplicate, TTL-expired, suspended or unauthorized trigger
874        // synthesizes nothing) and after the client boundary, so a forged
875        // envelope can never provoke an append.
876        self.synthesize_due_accept(&env.session_id, session, accepted_at_ms)
877            .await?;
878        let response = mode.on_message_at(
879            session,
880            env,
881            &macp_core::mode::MessageContext::new(accepted_at_ms),
882        )?;
883
884        // 1. COMMIT POINT: write log entry to disk
885        let incoming_entry = Self::make_incoming_entry(env, accepted_at_ms);
886        self.storage
887            .append_log_entry(&env.session_id, &incoming_entry)
888            .await
889            .map_err(|_| MacpError::StorageFailed)?;
890
891        // 2. Update in-memory state via the shared commit phase (consume dedup
892        //    slot, record participant activity, apply mode response) — the exact
893        //    sequence a library consumer runs through `macp_modes::step`.
894        self.log_store.append(&env.session_id, incoming_entry).await;
895        let result_state = macp_modes::step::commit(session, env, response, now_ms);
896
897        self.metrics.record_message_accepted(&session.mode);
898        if env.message_type == "Commitment" {
899            self.metrics.record_commitment_accepted(&session.mode);
900        }
901
902        // Policy-driven audit verbosity (E3b): a bound policy may request
903        // per-message audit lines at info level via an `audit.level` rules
904        // block ("info"); default stays debug. Mode rule schemas ignore
905        // unknown blocks, so `audit` composes with any mode's rules.
906        if Self::audit_verbose(session) {
907            tracing::info!(
908                session_id = %env.session_id,
909                message_type = %env.message_type,
910                sender = %env.sender,
911                state = ?result_state,
912                "message accepted (audit)"
913            );
914        } else {
915            tracing::debug!(
916                session_id = %env.session_id,
917                message_type = %env.message_type,
918                sender = %env.sender,
919                state = ?result_state,
920                "message accepted"
921            );
922        }
923
924        if result_state == SessionState::Resolved {
925            self.metrics.record_session_resolved(&session.mode);
926            tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
927            let _ = self
928                .session_lifecycle_bus
929                .send(SessionLifecycleEvent::Resolved {
930                    session_id: env.session_id.clone(),
931                });
932        }
933
934        // 3. Best-effort session save + checkpoint
935        self.save_session_to_storage(session).await;
936        if result_state == SessionState::Resolved {
937            if !self.maybe_compact_log(&env.session_id, session).await {
938                self.force_insert_checkpoint(&env.session_id, session).await;
939            }
940        } else {
941            self.maybe_insert_checkpoint(&env.session_id, session).await;
942        }
943        self.publish_accepted_envelope(env);
944
945        Ok(ProcessResult {
946            session_state: result_state,
947            duplicate: false,
948        })
949    }
950
951    /// Process a Signal or Progress envelope. Signals are informational out-of-band
952    /// notifications. Progress messages carry structured ProgressPayload.
953    /// Neither mutates session state — both are broadcast to subscribers.
954    async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
955        // RFC-MACP-0001 §4 / RFC-MACP-0010: validate SignalPayload structure.
956        // signal_type must be non-empty when a payload is present.
957        if env.message_type == "Signal" && !env.payload.is_empty() {
958            let signal: crate::pb::SignalPayload =
959                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
960            if signal.signal_type.trim().is_empty() {
961                return Err(MacpError::InvalidPayload);
962            }
963        }
964        // RFC-MACP-0001: validate ProgressPayload structure for Progress messages.
965        if env.message_type == "Progress" && !env.payload.is_empty() {
966            let _: crate::pb::ProgressPayload =
967                prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
968        }
969        tracing::debug!(
970            sender = %env.sender,
971            message_id = %env.message_id,
972            message_type = %env.message_type,
973            "signal received"
974        );
975        let _ = self.signal_bus.send(env.clone());
976        Ok(ProcessResult {
977            session_state: SessionState::Open,
978            duplicate: false,
979        })
980    }
981
982    pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
983        let shared = self.registry.get_shared(session_id).await?;
984        let mut session = shared.lock().await;
985        let changed = self
986            .maybe_expire_session(session_id, &mut session)
987            .await
988            .unwrap_or(false);
989        if changed {
990            self.save_session_to_storage(&session).await;
991        }
992        Some(session.clone())
993    }
994
995    /// Cancel a session. The `cancelled_by` parameter MUST be the authenticated
996    /// sender of the CancelSession RPC (RFC-MACP-0001 Section 7.3: CancelSession
997    /// is a Core control-plane message; mode authorization does not apply).
998    pub async fn cancel_session(
999        &self,
1000        session_id: &str,
1001        reason: &str,
1002        cancelled_by: &str,
1003    ) -> Result<ProcessResult, MacpError> {
1004        let shared = self
1005            .registry
1006            .get_shared(session_id)
1007            .await
1008            .ok_or(MacpError::UnknownSession)?;
1009        let mut session_guard = shared.lock().await;
1010        let session = &mut *session_guard;
1011
1012        self.maybe_expire_session(session_id, session).await?;
1013
1014        // Already terminal (Resolved/Expired/Cancelled): nothing to do. An Open
1015        // or Suspended session can still be cancelled (RFC-MACP-0001 §7.2/§7.3).
1016        if session.state.is_terminal() {
1017            let result_state = session.state.clone();
1018            self.save_session_to_storage(session).await;
1019            return Ok(ProcessResult {
1020                session_state: result_state,
1021                duplicate: false,
1022            });
1023        }
1024
1025        // RFC-MACP-0001: runtime encodes a proper SessionCancelPayload with
1026        // `cancelled_by` set to the authenticated sender identity.
1027        let now_ms = Utc::now().timestamp_millis();
1028        let cancel_payload = crate::pb::SessionCancelPayload {
1029            reason: reason.to_string(),
1030            cancelled_by: cancelled_by.to_string(),
1031        };
1032        let cancel_entry = Self::make_internal_entry(
1033            "SessionCancel",
1034            &prost::Message::encode_to_vec(&cancel_payload),
1035            session_id,
1036            &session.mode,
1037            now_ms,
1038        );
1039        self.storage
1040            .append_log_entry(session_id, &cancel_entry)
1041            .await
1042            .map_err(|_| MacpError::StorageFailed)?;
1043        self.log_store.append(session_id, cancel_entry).await;
1044        // RFC-MACP-0001 §7.3: cancellation terminates as CANCELLED (distinct
1045        // from EXPIRED) — `cancel()` also clears any suspension marker.
1046        let _ = session.cancel();
1047        self.save_session_to_storage(session).await;
1048        if !self.maybe_compact_log(session_id, session).await {
1049            self.force_insert_checkpoint(session_id, session).await;
1050        }
1051        self.metrics.record_session_cancelled(&session.mode);
1052        tracing::info!(session_id, reason, "session cancelled");
1053        let _ = self
1054            .session_lifecycle_bus
1055            .send(SessionLifecycleEvent::Cancelled {
1056                session_id: session_id.to_string(),
1057            });
1058
1059        Ok(ProcessResult {
1060            session_state: SessionState::Cancelled,
1061            duplicate: false,
1062        })
1063    }
1064
1065    /// Suspend an `Open` session (RFC-MACP-0001 §7.5). Appends a `SessionSuspend`
1066    /// annotation, transitions Open -> Suspended, and emits a lifecycle event.
1067    /// The session's TTL is banked and restored on resume.
1068    pub async fn suspend_session(
1069        &self,
1070        session_id: &str,
1071        reason: &str,
1072        suspended_by: &str,
1073    ) -> Result<ProcessResult, MacpError> {
1074        let shared = self
1075            .registry
1076            .get_shared(session_id)
1077            .await
1078            .ok_or(MacpError::UnknownSession)?;
1079        let mut session_guard = shared.lock().await;
1080        let session = &mut *session_guard;
1081
1082        self.maybe_expire_session(session_id, session).await?;
1083        if session.state != SessionState::Open {
1084            return Err(MacpError::SessionNotOpen);
1085        }
1086
1087        let now_ms = chrono::Utc::now().timestamp_millis();
1088        let payload = crate::pb::SessionSuspendPayload {
1089            reason: reason.to_string(),
1090            suspended_by: suspended_by.to_string(),
1091        };
1092        let entry = Self::make_internal_entry(
1093            "SessionSuspend",
1094            &prost::Message::encode_to_vec(&payload),
1095            session_id,
1096            &session.mode,
1097            now_ms,
1098        );
1099        self.storage
1100            .append_log_entry(session_id, &entry)
1101            .await
1102            .map_err(|_| MacpError::StorageFailed)?;
1103        self.log_store.append(session_id, entry).await;
1104        session.suspend(now_ms)?;
1105        self.save_session_to_storage(session).await;
1106        self.metrics.record_session_suspended(&session.mode);
1107        tracing::info!(session_id, reason, "session suspended");
1108        let _ = self
1109            .session_lifecycle_bus
1110            .send(SessionLifecycleEvent::Suspended {
1111                session_id: session_id.to_string(),
1112            });
1113
1114        Ok(ProcessResult {
1115            session_state: SessionState::Suspended,
1116            duplicate: false,
1117        })
1118    }
1119
1120    /// Resume a `Suspended` session (RFC-MACP-0001 §7.5), banking the suspended
1121    /// duration into the TTL deadline. If the `MAX_SUSPEND_MS` cap is exceeded,
1122    /// the session is force-expired instead.
1123    pub async fn resume_session(
1124        &self,
1125        session_id: &str,
1126        reason: &str,
1127        resumed_by: &str,
1128    ) -> Result<ProcessResult, MacpError> {
1129        let shared = self
1130            .registry
1131            .get_shared(session_id)
1132            .await
1133            .ok_or(MacpError::UnknownSession)?;
1134        let mut session_guard = shared.lock().await;
1135        let session = &mut *session_guard;
1136
1137        if session.state != SessionState::Suspended {
1138            return Err(MacpError::SessionNotOpen);
1139        }
1140
1141        let now_ms = chrono::Utc::now().timestamp_millis();
1142        // `banked_ms` on the wire payload is **informational only**. Both this
1143        // value and the `SessionResume` entry's `received_at_ms` now come from
1144        // the single `now_ms` read above, so it is exactly
1145        // `resume_entry.received_at_ms - suspend_entry.received_at_ms` — i.e.
1146        // equal to the value replay derives by construction. Replay still
1147        // ignores it and re-derives the banked duration from the two entry
1148        // timestamps (see the `SessionSuspend`/`SessionResume` arms of
1149        // `replay::replay_entry`). Keep it that way: starting to consume
1150        // `banked_ms` would change how *legacy* logs replay, because entries
1151        // written before this change recorded a second, independently-read
1152        // clock and their `banked_ms` can disagree with their timestamps.
1153        let banked_before = session
1154            .suspended_at_ms
1155            .map(|at| (now_ms - at).max(0))
1156            .unwrap_or(0);
1157        let payload = crate::pb::SessionResumePayload {
1158            reason: reason.to_string(),
1159            resumed_by: resumed_by.to_string(),
1160            banked_ms: banked_before,
1161        };
1162        let entry = Self::make_internal_entry(
1163            "SessionResume",
1164            &prost::Message::encode_to_vec(&payload),
1165            session_id,
1166            &session.mode,
1167            now_ms,
1168        );
1169        self.storage
1170            .append_log_entry(session_id, &entry)
1171            .await
1172            .map_err(|_| MacpError::StorageFailed)?;
1173        self.log_store.append(session_id, entry).await;
1174
1175        // `resume` banks the TTL; if the suspend cap is exceeded it force-expires.
1176        match session.resume(now_ms) {
1177            Ok(()) => {
1178                self.save_session_to_storage(session).await;
1179                self.metrics.record_session_resumed(&session.mode);
1180                tracing::info!(session_id, reason, "session resumed");
1181                let _ = self
1182                    .session_lifecycle_bus
1183                    .send(SessionLifecycleEvent::Resumed {
1184                        session_id: session_id.to_string(),
1185                    });
1186                Ok(ProcessResult {
1187                    session_state: SessionState::Open,
1188                    duplicate: false,
1189                })
1190            }
1191            Err(_) => {
1192                // A suspension cap was exceeded — either MAX_SUSPEND_MS
1193                // (cumulative suspended duration) or, at semantics_rev >= 2,
1194                // MAX_SUSPENSION_CYCLES (completed suspend/resume cycles).
1195                // Both take the same posture in `Session::resume`: the
1196                // session is now Expired. `resume` mutates `session` before
1197                // returning `Err`, so both fields are readable here to name
1198                // which cap actually fired instead of leaving an operator to
1199                // guess between two causes that share one error variant.
1200                let cycle_cap_exceeded = session.suspension_intervals.len() > MAX_SUSPENSION_CYCLES;
1201                let duration_cap_exceeded =
1202                    session.accumulated_suspended_ms > session.effective_max_suspend_ms();
1203                tracing::warn!(
1204                    session_id,
1205                    cycle_cap_exceeded,
1206                    duration_cap_exceeded,
1207                    suspension_cycles = session.suspension_intervals.len(),
1208                    accumulated_suspended_ms = session.accumulated_suspended_ms,
1209                    "session force-expired: suspension cap exceeded"
1210                );
1211                self.save_session_to_storage(session).await;
1212                self.metrics.record_session_expired(&session.mode);
1213                let _ = self
1214                    .session_lifecycle_bus
1215                    .send(SessionLifecycleEvent::Expired {
1216                        session_id: session_id.to_string(),
1217                    });
1218                Err(MacpError::TtlExpired)
1219            }
1220        }
1221    }
1222
1223    /// Best-effort log compaction for terminal sessions.
1224    /// Returns `true` if compaction succeeded, `false` if skipped or failed.
1225    async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
1226        // Ordinal accounting for the sequence contract: the checkpoint must
1227        // record every accepted ordinal it discards, including any base from
1228        // a prior compaction recorded in the current log.
1229        let discarded = match self.log_store.get_log(session_id).await {
1230            Some(entries) => {
1231                let prior_base: u64 = entries
1232                    .iter()
1233                    .filter(|e| e.entry_kind == EntryKind::Checkpoint)
1234                    .map(|e| e.compacted_incoming_ordinals)
1235                    .max()
1236                    .unwrap_or(0);
1237                prior_base
1238                    + entries
1239                        .iter()
1240                        .filter(|e| e.entry_kind == EntryKind::Incoming)
1241                        .count() as u64
1242            }
1243            None => 0,
1244        };
1245        match crate::storage::compaction::compact_session_log(
1246            &*self.storage,
1247            session_id,
1248            session,
1249            discarded,
1250        )
1251        .await
1252        {
1253            Ok(checkpoint) => {
1254                // Keep the in-memory log in step with storage — previously
1255                // only disk was rewritten, so memory and disk diverged and
1256                // post-restart passive-subscribe history vanished silently.
1257                self.log_store
1258                    .replace_session_log(session_id, vec![checkpoint])
1259                    .await;
1260                true
1261            }
1262            Err(e) => {
1263                tracing::debug!(
1264                    session_id,
1265                    error = %e,
1266                    "log compaction skipped (backend may not support it)"
1267                );
1268                false
1269            }
1270        }
1271    }
1272
1273    /// Force a checkpoint entry regardless of interval settings.
1274    /// Used as a fallback when compaction fails on terminal sessions.
1275    async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1276        let persisted = crate::registry::PersistedSession::from(session);
1277        let raw_payload = match serde_json::to_vec(&persisted) {
1278            Ok(bytes) => bytes,
1279            Err(e) => {
1280                tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1281                return;
1282            }
1283        };
1284        let now = Utc::now().timestamp_millis();
1285        let checkpoint = LogEntry {
1286            message_id: String::new(),
1287            received_at_ms: now,
1288            sender: "_runtime".into(),
1289            message_type: "Checkpoint".into(),
1290            raw_payload,
1291            entry_kind: EntryKind::Checkpoint,
1292            session_id: session_id.into(),
1293            mode: session.mode.clone(),
1294            macp_version: String::new(),
1295            timestamp_unix_ms: now,
1296            bound_mode_version: None,
1297            semantics_rev: 0,
1298            bound_max_suspend_ms: None,
1299            compacted_incoming_ordinals: 0,
1300        };
1301        if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1302            tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1303            return;
1304        }
1305        self.log_store.append(session_id, checkpoint).await;
1306        tracing::debug!(
1307            session_id,
1308            "forced checkpoint inserted for terminal session"
1309        );
1310    }
1311
1312    /// Insert a checkpoint entry if the log has reached the configured interval.
1313    ///
1314    /// Called after **every** append that can cross a boundary: the tail of
1315    /// `process_message`, and [`Runtime::synthesize_due_accept`] for the
1316    /// synthetic entry it writes. The second call site is what keeps the eager
1317    /// sweep and the lazy trigger placing checkpoints identically — see the
1318    /// note at that call. The decision is a pure function of the current
1319    /// `log_len`, so asking twice within one `process_message` (at two
1320    /// different lengths) is not a repeat.
1321    async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1322        if self.checkpoint_interval == 0 {
1323            return;
1324        }
1325        let log_len = self
1326            .log_store
1327            .get_log(session_id)
1328            .await
1329            .map(|l| l.len())
1330            .unwrap_or(0);
1331        // Only checkpoint at interval boundaries, and not on the first entry
1332        if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1333            return;
1334        }
1335        self.force_insert_checkpoint(session_id, session).await;
1336        tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1337    }
1338
1339    /// Expire all sessions that have exceeded their TTL.
1340    /// Called by the background cleanup task to proactively transition
1341    /// stale sessions without waiting for the next incoming message.
1342    pub async fn cleanup_expired_sessions(&self) {
1343        let now = Utc::now().timestamp_millis();
1344        // Snapshot the shared handles under a brief map read; never hold the
1345        // map lock across per-session locks or storage I/O. Each session is
1346        // re-checked under its own mutex (it may have been touched since the
1347        // snapshot).
1348        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1349            let guard = self.registry.sessions.read().await;
1350            guard
1351                .iter()
1352                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1353                .collect()
1354        };
1355
1356        let mut expired_count = 0usize;
1357        for (session_id, shared) in candidates {
1358            let mut session = shared.lock().await;
1359            if session.state != SessionState::Open || now <= session.ttl_expiry {
1360                continue;
1361            }
1362            let entry =
1363                Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode, now);
1364            if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1365                tracing::warn!(
1366                    session_id,
1367                    error = %e,
1368                    "failed to write TTL expiry during cleanup"
1369                );
1370                continue;
1371            }
1372            self.log_store.append(&session_id, entry).await;
1373            session.state = SessionState::Expired;
1374            self.metrics.record_session_expired(&session.mode);
1375            self.save_session_to_storage(&session).await;
1376            if !self.maybe_compact_log(&session_id, &session).await {
1377                self.force_insert_checkpoint(&session_id, &session).await;
1378            }
1379            expired_count += 1;
1380            tracing::info!(session_id = %session_id, "session expired via background cleanup");
1381            let _ = self
1382                .session_lifecycle_bus
1383                .send(SessionLifecycleEvent::Expired {
1384                    session_id: session_id.clone(),
1385                });
1386        }
1387
1388        if expired_count > 0 {
1389            tracing::info!(count = expired_count, "background cleanup expired sessions");
1390        }
1391    }
1392
1393    /// Emit every synthetic envelope that has become due, across all open
1394    /// sessions — RFC-MACP-0010 §5.1(2)'s **eager** observation of the
1395    /// implicit-accept deadline.
1396    ///
1397    /// Lazy observation (`process_message` -> `synthesize_due_accept`) is the
1398    /// MUST and ships on its own: it guarantees no message is ever evaluated
1399    /// against a stale offer. This is the SHOULD on top of it — without it a
1400    /// session where nobody speaks again keeps an accepted offer out of history
1401    /// indefinitely, and `GetSession` / `StreamSession` show an offer that the
1402    /// protocol says was accepted at `D`. Called from the background
1403    /// maintenance loop in `src/main.rs`, so `MACP_CLEANUP_INTERVAL_SECS` is
1404    /// the latency bound on the observation (never on the recorded timestamp,
1405    /// which is `D` whichever path emits — see below).
1406    ///
1407    /// # Not a variant of `cleanup_expired_sessions`
1408    ///
1409    /// Different predicate (a mode-computed deadline inside `mode_state`, not
1410    /// `ttl_expiry`) and a different product: an `EntryKind::Incoming` entry
1411    /// that consumes an accepted ordinal and is published to `StreamSession`,
1412    /// where `TtlExpired` is `Internal` and is published to neither. It is
1413    /// deliberately ordered *after* `cleanup_expired_sessions` in that loop so
1414    /// a TTL-expired session is already non-`Open` when the sweep reaches it —
1415    /// which is exactly the precedence the lazy path gives (`Precheck::Expired`
1416    /// returns before `synthesize_due_accept` is ever called), so the two paths
1417    /// cannot disagree about a session whose TTL and implicit-accept deadline
1418    /// both passed unobserved.
1419    ///
1420    /// # Locking
1421    ///
1422    /// The registry map lock is held only for the snapshot of `(id, Arc)` pairs
1423    /// and is released before any session mutex is taken or any I/O happens —
1424    /// the lock-ordering contract on `SessionRegistry`. The snapshot fixes the
1425    /// *set*, not the state.
1426    ///
1427    /// # Why no snapshot-to-append race can leave an orphan entry
1428    ///
1429    /// A session can resolve, cancel or expire between the snapshot and the
1430    /// moment this loop reaches it, and the `Arc` keeps it alive (and writable)
1431    /// regardless. Nothing here reads state at snapshot time: every decision is
1432    /// made *under the session mutex*, which is the same mutex every writer —
1433    /// `process_message`, `cancel_session`, `suspend_session`, `resume_session`,
1434    /// `cleanup_expired_sessions` — holds across its own validate-append-commit.
1435    /// So the `state != Open` re-check inside `synthesize_due_accept` observes
1436    /// the session as the last writer left it, and a session that terminated
1437    /// after the snapshot is skipped. Two further guards make it belt and
1438    /// braces: the mode returns `None` once the offer's disposition is no
1439    /// longer `Offered` (so an explicit `HandoffAccept` that won the race
1440    /// disarms the synthesis), and the deterministic `message_id` is already in
1441    /// `seen_message_ids` after any emission. Eviction cannot orphan one
1442    /// either: `evict_stale_sessions` and `gc_disk_sessions` only ever drop
1443    /// *terminal* sessions, which fail the `Open` check.
1444    ///
1445    /// # Cost
1446    ///
1447    /// One `mode_state` decode per open session whose mode implements the hook
1448    /// and whose policy binds a timeout, per tick; everything else short-circuits
1449    /// before decoding (`Mode::due_synthetic_envelope`'s own cost note).
1450    ///
1451    /// Returns the number of synthetic entries appended.
1452    pub async fn sweep_due_synthetic_accepts(&self) -> usize {
1453        let now = Utc::now().timestamp_millis();
1454        // Snapshot the shared handles under a brief map read; never hold the
1455        // map lock across per-session locks or storage I/O. Mirrors
1456        // `cleanup_expired_sessions`.
1457        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1458            let guard = self.registry.sessions.read().await;
1459            guard
1460                .iter()
1461                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1462                .collect()
1463        };
1464
1465        let mut emitted = 0usize;
1466        for (session_id, shared) in candidates {
1467            let mut session = shared.lock().await;
1468            // `Suspended` (and every other non-`Open` state) is skipped here
1469            // *and* inside `synthesize_due_accept`. Measured: removing this
1470            // check alone changes no behaviour — the seam's own filter still
1471            // declines — so treat it as the documented precondition of the
1472            // call below rather than as the enforcement. The enforcement
1473            // matters: both computations behind the synthetic entry ignore an
1474            // in-flight pause, so asking about a paused session over-counts
1475            // elapsed time and bakes a wrong `D` into permanent history, and
1476            // this is the only caller that can ever be handed one (no message
1477            // path reaches a non-`Open` session at all).
1478            if session.state != SessionState::Open {
1479                continue;
1480            }
1481            match self
1482                .synthesize_due_accept(&session_id, &mut session, now)
1483                .await
1484            {
1485                Ok(true) => emitted += 1,
1486                Ok(false) => {}
1487                // No triggering message to reject: a failed append (or a mode
1488                // that refused its own synthetic) is logged and the offer stays
1489                // outstanding, to be retried on the next tick or settled by the
1490                // lazy path. The same posture `cleanup_expired_sessions` takes
1491                // on a failed `TtlExpired` append.
1492                Err(e) => {
1493                    tracing::warn!(
1494                        session_id = %session_id,
1495                        error = %e,
1496                        "eager sweep could not emit a due synthetic envelope"
1497                    );
1498                }
1499            }
1500        }
1501
1502        if emitted > 0 {
1503            tracing::info!(
1504                count = emitted,
1505                "eager sweep appended due synthetic envelopes"
1506            );
1507        }
1508        emitted
1509    }
1510
1511    /// Delete terminal sessions' durable data older than `retention_secs`
1512    /// (opt-in via `MACP_SESSION_DISK_RETENTION_SECS`). Before this existed,
1513    /// `storage.delete_session` had no callers at all: disk grew without
1514    /// bound and every restart reloaded every session ever completed.
1515    /// Enumerates STORAGE (not memory — eviction may already have dropped the
1516    /// registry entry), deletes the session's snapshot+log, and clears any
1517    /// in-memory remnants. Returns the number of sessions deleted.
1518    pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1519        let now = Utc::now().timestamp_millis();
1520        let cutoff = now - (retention_secs as i64 * 1000);
1521        let ids = match self.storage.list_session_ids().await {
1522            Ok(ids) => ids,
1523            Err(e) => {
1524                tracing::warn!(error = %e, "disk GC: cannot list sessions");
1525                return 0;
1526            }
1527        };
1528        let mut removed = 0usize;
1529        for id in ids {
1530            // Prefer the in-memory state when present (cheap + current);
1531            // fall back to the stored snapshot for evicted sessions.
1532            let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1533                let s = shared.lock().await;
1534                s.state.is_terminal() && s.started_at_unix_ms < cutoff
1535            } else {
1536                match self.storage.load_session(&id).await {
1537                    Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1538                    // No snapshot (or unreadable): leave it for operator
1539                    // inspection rather than guessing.
1540                    _ => false,
1541                }
1542            };
1543            if !eligible {
1544                continue;
1545            }
1546            match self.storage.delete_session(&id).await {
1547                Ok(()) => {
1548                    {
1549                        let mut guard = self.registry.sessions.write().await;
1550                        guard.remove(&id);
1551                    }
1552                    self.log_store.remove_session_log(&id).await;
1553                    let _ = self.stream_bus.remove_if_unused(&id);
1554                    removed += 1;
1555                }
1556                Err(e) => {
1557                    tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1558                }
1559            }
1560        }
1561        if removed > 0 {
1562            tracing::info!(count = removed, "disk GC removed terminal sessions");
1563        }
1564        removed
1565    }
1566
1567    /// Evict resolved/expired sessions older than `retention_secs` from
1568    /// memory: the registry entry, the in-memory log cache, AND the stream
1569    /// broadcast channel (all three previously grew for the process lifetime;
1570    /// the log cache and stream bus were never evicted at all). Sessions
1571    /// remain queryable from durable storage after eviction.
1572    pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1573        let now = Utc::now().timestamp_millis();
1574        let cutoff = now - (retention_secs as i64 * 1000);
1575
1576        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1577            let guard = self.registry.sessions.read().await;
1578            guard
1579                .iter()
1580                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1581                .collect()
1582        };
1583        let mut evict_ids = Vec::new();
1584        for (id, shared) in candidates {
1585            let session = shared.lock().await;
1586            if matches!(
1587                session.state,
1588                SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1589            ) && session.started_at_unix_ms < cutoff
1590            {
1591                evict_ids.push(id);
1592            }
1593        }
1594
1595        if evict_ids.is_empty() {
1596            return;
1597        }
1598        {
1599            let mut guard = self.registry.sessions.write().await;
1600            for id in &evict_ids {
1601                guard.remove(id);
1602            }
1603        }
1604        for id in &evict_ids {
1605            self.log_store.remove_session_log(id).await;
1606            // Left in place if a subscriber is still attached; retried on the
1607            // next sweep once receivers drop.
1608            let _ = self.stream_bus.remove_if_unused(id);
1609        }
1610        tracing::info!(
1611            count = evict_ids.len(),
1612            "evicted stale sessions from memory (registry + log cache + stream bus)"
1613        );
1614    }
1615}
1616
1617#[cfg(test)]
1618mod tests {
1619    use super::*;
1620    use crate::decision_pb::ProposalPayload;
1621    use crate::pb::{CommitmentPayload, SessionStartPayload};
1622    use prost::Message;
1623
1624    fn new_sid() -> String {
1625        uuid::Uuid::new_v4().as_hyphenated().to_string()
1626    }
1627
1628    fn make_runtime() -> Runtime {
1629        let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1630        let registry = Arc::new(SessionRegistry::new());
1631        let log_store = Arc::new(LogStore::new());
1632        Runtime::new(storage, registry, log_store)
1633    }
1634
1635    fn session_start(participants: Vec<String>) -> Vec<u8> {
1636        SessionStartPayload {
1637            intent: "intent".into(),
1638            participants,
1639            mode_version: "1.0.0".into(),
1640            configuration_version: "cfg-1".into(),
1641            policy_version: String::new(),
1642            ttl_ms: 1_000,
1643            context_id: String::new(),
1644            extensions: std::collections::HashMap::new(),
1645            roots: vec![],
1646            max_suspend_ms: 0,
1647        }
1648        .encode_to_vec()
1649    }
1650
1651    fn env(
1652        mode: &str,
1653        message_type: &str,
1654        message_id: &str,
1655        session_id: &str,
1656        sender: &str,
1657        payload: Vec<u8>,
1658    ) -> Envelope {
1659        Envelope {
1660            macp_version: "1.0".into(),
1661            mode: mode.into(),
1662            message_type: message_type.into(),
1663            message_id: message_id.into(),
1664            session_id: session_id.into(),
1665            sender: sender.into(),
1666            timestamp_unix_ms: Utc::now().timestamp_millis(),
1667            payload,
1668        }
1669    }
1670
1671    #[tokio::test]
1672    async fn standard_session_start_is_strict() {
1673        let rt = make_runtime();
1674        let sid = new_sid();
1675        let bad = SessionStartPayload {
1676            ttl_ms: 0,
1677            ..Default::default()
1678        }
1679        .encode_to_vec();
1680        let err = rt
1681            .process(
1682                &env(
1683                    "macp.mode.decision.v1",
1684                    "SessionStart",
1685                    "m1",
1686                    &sid,
1687                    "agent://orchestrator",
1688                    bad,
1689                ),
1690                None,
1691            )
1692            .await
1693            .unwrap_err();
1694        assert!(matches!(
1695            err,
1696            MacpError::InvalidPayload | MacpError::InvalidTtl
1697        ));
1698    }
1699
1700    /// A **promoted** extension mode keeps the full canonical `SessionStart`
1701    /// contract, including the roster requirement.
1702    ///
1703    /// This pins the *call site's* choice of strictness source, which no other
1704    /// test covers. `ModeRegistry::requires_strict_session_start` reads a
1705    /// per-entry `strict_session_start` flag that `promote_mode` sets to `true`;
1706    /// `macp_core::session::requires_strict_session_start` reads a static name
1707    /// list that has never heard of a promoted mode's name. The two disagree
1708    /// exactly here, so routing this call through
1709    /// `validate_strict_session_start_payload` — which consults the static list
1710    /// — would silently skip canonical validation for every promoted mode while
1711    /// leaving the whole suite green. Measured: it does.
1712    #[tokio::test]
1713    async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1714        let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1715            macp_policy::DefaultPolicyEvaluator,
1716        )));
1717        mode_registry
1718            .register_extension(crate::pb::ModeDescriptor {
1719                mode: "ext.promoted.v1".into(),
1720                mode_version: "1.0.0".into(),
1721                title: "Promoted".into(),
1722                description: "promotion target".into(),
1723                determinism_class: "semantic-deterministic".into(),
1724                participant_model: "declared".into(),
1725                message_types: vec!["SessionStart".into(), "Commitment".into()],
1726                terminal_message_types: vec!["Commitment".into()],
1727                ..Default::default()
1728            })
1729            .expect("register extension");
1730        assert_eq!(
1731            mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1732            "ext.promoted.v1"
1733        );
1734        assert!(
1735            mode_registry.requires_strict_session_start("ext.promoted.v1"),
1736            "promotion must mark the entry strict"
1737        );
1738        assert!(
1739            !crate::session::requires_strict_session_start("ext.promoted.v1"),
1740            "the core's static list must NOT know this name — that disagreement is the point"
1741        );
1742
1743        let rt = Runtime::with_mode_registry(
1744            Arc::new(crate::storage::MemoryBackend),
1745            Arc::new(SessionRegistry::new()),
1746            Arc::new(LogStore::new()),
1747            mode_registry,
1748        );
1749
1750        // An empty roster: refused, because the carve-out names Decision only.
1751        let err = rt
1752            .process(
1753                &env(
1754                    "ext.promoted.v1",
1755                    "SessionStart",
1756                    "m1",
1757                    &new_sid(),
1758                    "agent://orchestrator",
1759                    session_start(vec![]),
1760                ),
1761                None,
1762            )
1763            .await
1764            .unwrap_err();
1765        assert_eq!(err.to_string(), "InvalidPayload");
1766
1767        // And the rest of the canonical contract too, so the assertion above
1768        // cannot be satisfied by a runtime that only kept the roster rule.
1769        let no_versions = SessionStartPayload {
1770            participants: vec!["agent://fraud".into()],
1771            ttl_ms: 1_000,
1772            ..Default::default()
1773        }
1774        .encode_to_vec();
1775        let err = rt
1776            .process(
1777                &env(
1778                    "ext.promoted.v1",
1779                    "SessionStart",
1780                    "m2",
1781                    &new_sid(),
1782                    "agent://orchestrator",
1783                    no_versions,
1784                ),
1785                None,
1786            )
1787            .await
1788            .unwrap_err();
1789        assert_eq!(err.to_string(), "InvalidPayload");
1790
1791        // Positive control: a complete payload is accepted, so the two refusals
1792        // above are the roster and version rules and not a broken mode.
1793        rt.process(
1794            &env(
1795                "ext.promoted.v1",
1796                "SessionStart",
1797                "m3",
1798                &new_sid(),
1799                "agent://orchestrator",
1800                session_start(vec!["agent://fraud".into()]),
1801            ),
1802            None,
1803        )
1804        .await
1805        .expect("a complete SessionStart must still be accepted for a promoted mode");
1806    }
1807
1808    #[tokio::test]
1809    async fn empty_mode_is_rejected() {
1810        let rt = make_runtime();
1811        let sid = new_sid();
1812        let err = rt
1813            .process(
1814                &env(
1815                    "",
1816                    "SessionStart",
1817                    "m1",
1818                    &sid,
1819                    "agent://orchestrator",
1820                    session_start(vec!["agent://fraud".into()]),
1821                ),
1822                None,
1823            )
1824            .await
1825            .unwrap_err();
1826        assert_eq!(err.to_string(), "InvalidEnvelope");
1827    }
1828
1829    #[tokio::test]
1830    async fn rejected_messages_do_not_enter_dedup_state() {
1831        let rt = make_runtime();
1832        let sid = new_sid();
1833        rt.process(
1834            &env(
1835                "macp.mode.decision.v1",
1836                "SessionStart",
1837                "m1",
1838                &sid,
1839                "agent://orchestrator",
1840                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1841            ),
1842            None,
1843        )
1844        .await
1845        .unwrap();
1846
1847        let bad = rt
1848            .process(
1849                &env(
1850                    "macp.mode.decision.v1",
1851                    "Proposal",
1852                    "m2",
1853                    &sid,
1854                    "agent://fraud",
1855                    b"not-protobuf".to_vec(),
1856                ),
1857                None,
1858            )
1859            .await
1860            .unwrap_err();
1861        assert_eq!(bad.to_string(), "InvalidPayload");
1862
1863        let good = ProposalPayload {
1864            proposal_id: "p1".into(),
1865            option: "step-up".into(),
1866            rationale: "risk".into(),
1867            supporting_data: vec![],
1868        }
1869        .encode_to_vec();
1870        let result = rt
1871            .process(
1872                &env(
1873                    "macp.mode.decision.v1",
1874                    "Proposal",
1875                    "m2",
1876                    &sid,
1877                    "agent://orchestrator",
1878                    good,
1879                ),
1880                None,
1881            )
1882            .await
1883            .unwrap();
1884        assert!(!result.duplicate);
1885    }
1886
1887    #[tokio::test]
1888    async fn get_session_transitions_expired_sessions() {
1889        let rt = make_runtime();
1890        let sid = new_sid();
1891        let payload = SessionStartPayload {
1892            intent: "intent".into(),
1893            participants: vec!["agent://fraud".into()],
1894            mode_version: "1.0.0".into(),
1895            configuration_version: "cfg-1".into(),
1896            policy_version: String::new(),
1897            ttl_ms: 1,
1898            context_id: String::new(),
1899            extensions: std::collections::HashMap::new(),
1900            roots: vec![],
1901            max_suspend_ms: 0,
1902        }
1903        .encode_to_vec();
1904        rt.process(
1905            &env(
1906                "macp.mode.decision.v1",
1907                "SessionStart",
1908                "m1",
1909                &sid,
1910                "agent://orchestrator",
1911                payload,
1912            ),
1913            None,
1914        )
1915        .await
1916        .unwrap();
1917        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1918        let session = rt.get_session_checked(&sid).await.unwrap();
1919        assert_eq!(session.state, SessionState::Expired);
1920    }
1921
1922    #[tokio::test]
1923    async fn multi_round_requires_standard_session_start() {
1924        let rt = make_runtime();
1925        let sid = new_sid();
1926        // multi-round is now standards-track: empty mode_version should fail
1927        let payload = SessionStartPayload {
1928            participants: vec!["creator".into(), "other".into()],
1929            ..Default::default()
1930        }
1931        .encode_to_vec();
1932        let err = rt
1933            .process(
1934                &env(
1935                    "ext.multi_round.v1",
1936                    "SessionStart",
1937                    "m1",
1938                    &sid,
1939                    "creator",
1940                    payload,
1941                ),
1942                None,
1943            )
1944            .await
1945            .unwrap_err();
1946        assert!(matches!(
1947            err,
1948            MacpError::InvalidPayload | MacpError::InvalidTtl
1949        ));
1950    }
1951
1952    #[tokio::test]
1953    async fn multi_round_valid_session_start() {
1954        let rt = make_runtime();
1955        let sid = new_sid();
1956        let payload = session_start(vec!["alice".into(), "bob".into()]);
1957        rt.process(
1958            &env(
1959                "ext.multi_round.v1",
1960                "SessionStart",
1961                "m1",
1962                &sid,
1963                "coordinator",
1964                payload,
1965            ),
1966            None,
1967        )
1968        .await
1969        .unwrap();
1970        let session = rt.get_session_checked(&sid).await.unwrap();
1971        assert_eq!(session.mode, "ext.multi_round.v1");
1972        assert_eq!(session.participants, vec!["alice", "bob"]);
1973    }
1974
1975    #[tokio::test]
1976    async fn duplicate_session_start_message_id_returns_duplicate() {
1977        let rt = make_runtime();
1978        let sid = new_sid();
1979        let payload = session_start(vec!["agent://fraud".into()]);
1980        rt.process(
1981            &env(
1982                "macp.mode.decision.v1",
1983                "SessionStart",
1984                "m1",
1985                &sid,
1986                "agent://orchestrator",
1987                payload.clone(),
1988            ),
1989            None,
1990        )
1991        .await
1992        .unwrap();
1993
1994        let result = rt
1995            .process(
1996                &env(
1997                    "macp.mode.decision.v1",
1998                    "SessionStart",
1999                    "m1",
2000                    &sid,
2001                    "agent://orchestrator",
2002                    payload,
2003                ),
2004                None,
2005            )
2006            .await
2007            .unwrap();
2008        assert!(result.duplicate);
2009    }
2010
2011    #[tokio::test]
2012    async fn non_start_mode_mismatch_rejected() {
2013        let rt = make_runtime();
2014        let sid = new_sid();
2015        rt.process(
2016            &env(
2017                "macp.mode.decision.v1",
2018                "SessionStart",
2019                "m1",
2020                &sid,
2021                "agent://orchestrator",
2022                session_start(vec!["agent://fraud".into()]),
2023            ),
2024            None,
2025        )
2026        .await
2027        .unwrap();
2028
2029        let proposal = ProposalPayload {
2030            proposal_id: "p1".into(),
2031            option: "step-up".into(),
2032            rationale: "risk".into(),
2033            supporting_data: vec![],
2034        }
2035        .encode_to_vec();
2036        let err = rt
2037            .process(
2038                &env(
2039                    "macp.mode.task.v1",
2040                    "Proposal",
2041                    "m2",
2042                    &sid,
2043                    "agent://orchestrator",
2044                    proposal,
2045                ),
2046                None,
2047            )
2048            .await
2049            .unwrap_err();
2050        assert_eq!(err.to_string(), "InvalidEnvelope");
2051    }
2052
2053    #[tokio::test]
2054    async fn cancel_idempotent_on_already_expired() {
2055        let rt = make_runtime();
2056        let sid = new_sid();
2057        let payload = SessionStartPayload {
2058            intent: "intent".into(),
2059            participants: vec!["agent://fraud".into()],
2060            mode_version: "1.0.0".into(),
2061            configuration_version: "cfg-1".into(),
2062            policy_version: String::new(),
2063            ttl_ms: 1,
2064            context_id: String::new(),
2065            extensions: std::collections::HashMap::new(),
2066            roots: vec![],
2067            max_suspend_ms: 0,
2068        }
2069        .encode_to_vec();
2070        rt.process(
2071            &env(
2072                "macp.mode.decision.v1",
2073                "SessionStart",
2074                "m1",
2075                &sid,
2076                "agent://orchestrator",
2077                payload,
2078            ),
2079            None,
2080        )
2081        .await
2082        .unwrap();
2083        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2084        let result = rt
2085            .cancel_session(&sid, "cleanup", "agent://orchestrator")
2086            .await
2087            .unwrap();
2088        assert_eq!(result.session_state, SessionState::Expired);
2089    }
2090
2091    #[tokio::test]
2092    async fn accepted_envelopes_are_published_in_order() {
2093        let rt = make_runtime();
2094        let sid = new_sid();
2095        let mut events = rt.subscribe_session_stream(&sid);
2096
2097        let start = env(
2098            "macp.mode.decision.v1",
2099            "SessionStart",
2100            "m1",
2101            &sid,
2102            "agent://orchestrator",
2103            session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2104        );
2105        rt.process(&start, None).await.unwrap();
2106        let first = events.recv().await.unwrap();
2107        assert_eq!(first.message_id, "m1");
2108        assert_eq!(first.message_type, "SessionStart");
2109
2110        let proposal = ProposalPayload {
2111            proposal_id: "p1".into(),
2112            option: "step-up".into(),
2113            rationale: "risk".into(),
2114            supporting_data: vec![],
2115        }
2116        .encode_to_vec();
2117        let proposal_env = env(
2118            "macp.mode.decision.v1",
2119            "Proposal",
2120            "m2",
2121            &sid,
2122            "agent://orchestrator",
2123            proposal,
2124        );
2125        rt.process(&proposal_env, None).await.unwrap();
2126        let second = events.recv().await.unwrap();
2127        assert_eq!(second.message_id, "m2");
2128        assert_eq!(second.message_type, "Proposal");
2129    }
2130
2131    #[tokio::test]
2132    async fn commitment_versions_are_carried_into_resolution() {
2133        let rt = make_runtime();
2134        let sid = new_sid();
2135        rt.process(
2136            &env(
2137                "macp.mode.proposal.v1",
2138                "SessionStart",
2139                "m1",
2140                &sid,
2141                "agent://buyer",
2142                session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
2143            ),
2144            None,
2145        )
2146        .await
2147        .unwrap();
2148
2149        let proposal = crate::proposal_pb::ProposalPayload {
2150            proposal_id: "p1".into(),
2151            title: "offer".into(),
2152            summary: "summary".into(),
2153            details: vec![],
2154            tags: vec![],
2155        }
2156        .encode_to_vec();
2157        rt.process(
2158            &env(
2159                "macp.mode.proposal.v1",
2160                "Proposal",
2161                "m2",
2162                &sid,
2163                "agent://seller",
2164                proposal,
2165            ),
2166            None,
2167        )
2168        .await
2169        .unwrap();
2170        let accept = crate::proposal_pb::AcceptPayload {
2171            proposal_id: "p1".into(),
2172            reason: String::new(),
2173        }
2174        .encode_to_vec();
2175        rt.process(
2176            &env(
2177                "macp.mode.proposal.v1",
2178                "Accept",
2179                "m3",
2180                &sid,
2181                "agent://seller",
2182                accept.clone(),
2183            ),
2184            None,
2185        )
2186        .await
2187        .unwrap();
2188        rt.process(
2189            &env(
2190                "macp.mode.proposal.v1",
2191                "Accept",
2192                "m4",
2193                &sid,
2194                "agent://buyer",
2195                accept,
2196            ),
2197            None,
2198        )
2199        .await
2200        .unwrap();
2201        let commitment = CommitmentPayload {
2202            commitment_id: "c1".into(),
2203            action: "proposal.accepted".into(),
2204            authority_scope: "commercial".into(),
2205            reason: "bound".into(),
2206            mode_version: "1.0.0".into(),
2207            policy_version: "policy.default".into(),
2208            configuration_version: "cfg-1".into(),
2209            outcome_positive: true,
2210            supersedes: None,
2211        }
2212        .encode_to_vec();
2213        let result = rt
2214            .process(
2215                &env(
2216                    "macp.mode.proposal.v1",
2217                    "Commitment",
2218                    "m5",
2219                    &sid,
2220                    "agent://buyer",
2221                    commitment,
2222                ),
2223                None,
2224            )
2225            .await
2226            .unwrap();
2227        assert_eq!(result.session_state, SessionState::Resolved);
2228    }
2229
2230    #[tokio::test]
2231    async fn max_open_sessions_enforced_under_write_lock() {
2232        let rt = make_runtime();
2233        let sid1 = new_sid();
2234        let sid2 = new_sid();
2235        let sid3 = new_sid();
2236        rt.process(
2237            &env(
2238                "macp.mode.decision.v1",
2239                "SessionStart",
2240                "m1",
2241                &sid1,
2242                "agent://orchestrator",
2243                session_start(vec!["agent://fraud".into()]),
2244            ),
2245            Some(1),
2246        )
2247        .await
2248        .unwrap();
2249
2250        let err = rt
2251            .process(
2252                &env(
2253                    "macp.mode.decision.v1",
2254                    "SessionStart",
2255                    "m2",
2256                    &sid2,
2257                    "agent://orchestrator",
2258                    session_start(vec!["agent://fraud".into()]),
2259                ),
2260                Some(1),
2261            )
2262            .await
2263            .unwrap_err();
2264        assert!(matches!(err, MacpError::RateLimited));
2265
2266        rt.process(
2267            &env(
2268                "macp.mode.decision.v1",
2269                "SessionStart",
2270                "m3",
2271                &sid3,
2272                "agent://other",
2273                session_start(vec!["agent://fraud".into()]),
2274            ),
2275            Some(1),
2276        )
2277        .await
2278        .unwrap();
2279    }
2280
2281    #[tokio::test]
2282    async fn weak_session_id_rejected() {
2283        let rt = make_runtime();
2284        let err = rt
2285            .process(
2286                &env(
2287                    "macp.mode.decision.v1",
2288                    "SessionStart",
2289                    "m1",
2290                    "s1",
2291                    "agent://orchestrator",
2292                    session_start(vec!["agent://fraud".into()]),
2293                ),
2294                None,
2295            )
2296            .await
2297            .unwrap_err();
2298        assert_eq!(err.to_string(), "InvalidSessionId");
2299    }
2300
2301    #[tokio::test]
2302    async fn log_append_failure_rejects_session_start() {
2303        use std::io;
2304        struct FailingBackend;
2305        #[async_trait::async_trait]
2306        impl StorageBackend for FailingBackend {
2307            async fn save_session(&self, _: &Session) -> io::Result<()> {
2308                Ok(())
2309            }
2310            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2311                Ok(None)
2312            }
2313            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2314                Ok(vec![])
2315            }
2316            async fn delete_session(&self, _: &str) -> io::Result<()> {
2317                Ok(())
2318            }
2319            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2320                Ok(vec![])
2321            }
2322            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2323                Err(io::Error::other("disk full"))
2324            }
2325            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2326                Ok(vec![])
2327            }
2328            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2329                Ok(())
2330            }
2331        }
2332
2333        let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
2334        let registry = Arc::new(SessionRegistry::new());
2335        let log_store = Arc::new(LogStore::new());
2336        let rt = Runtime::new(storage, registry, log_store);
2337        let sid = new_sid();
2338
2339        let err = rt
2340            .process(
2341                &env(
2342                    "macp.mode.decision.v1",
2343                    "SessionStart",
2344                    "m1",
2345                    &sid,
2346                    "agent://orchestrator",
2347                    session_start(vec!["agent://fraud".into()]),
2348                ),
2349                None,
2350            )
2351            .await
2352            .unwrap_err();
2353        assert_eq!(err.to_string(), "StorageFailed");
2354    }
2355
2356    #[tokio::test]
2357    async fn log_append_failure_rejects_in_session_message() {
2358        use std::io;
2359        use std::sync::atomic::{AtomicUsize, Ordering};
2360
2361        struct FailOnSecondAppend {
2362            count: AtomicUsize,
2363        }
2364        #[async_trait::async_trait]
2365        impl StorageBackend for FailOnSecondAppend {
2366            async fn save_session(&self, _: &Session) -> io::Result<()> {
2367                Ok(())
2368            }
2369            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2370                Ok(None)
2371            }
2372            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2373                Ok(vec![])
2374            }
2375            async fn delete_session(&self, _: &str) -> io::Result<()> {
2376                Ok(())
2377            }
2378            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2379                Ok(vec![])
2380            }
2381            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2382                let n = self.count.fetch_add(1, Ordering::SeqCst);
2383                if n >= 1 {
2384                    Err(io::Error::other("disk full"))
2385                } else {
2386                    Ok(())
2387                }
2388            }
2389            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2390                Ok(vec![])
2391            }
2392            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2393                Ok(())
2394            }
2395        }
2396
2397        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2398            count: AtomicUsize::new(0),
2399        });
2400        let registry = Arc::new(SessionRegistry::new());
2401        let log_store = Arc::new(LogStore::new());
2402        let rt = Runtime::new(storage, registry, log_store);
2403        let sid = new_sid();
2404
2405        // SessionStart succeeds (first append)
2406        rt.process(
2407            &env(
2408                "macp.mode.decision.v1",
2409                "SessionStart",
2410                "m1",
2411                &sid,
2412                "agent://orchestrator",
2413                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2414            ),
2415            None,
2416        )
2417        .await
2418        .unwrap();
2419
2420        // Proposal fails (second append)
2421        let proposal = ProposalPayload {
2422            proposal_id: "p1".into(),
2423            option: "step-up".into(),
2424            rationale: "risk".into(),
2425            supporting_data: vec![],
2426        }
2427        .encode_to_vec();
2428        let err = rt
2429            .process(
2430                &env(
2431                    "macp.mode.decision.v1",
2432                    "Proposal",
2433                    "m2",
2434                    &sid,
2435                    "agent://orchestrator",
2436                    proposal,
2437                ),
2438                None,
2439            )
2440            .await
2441            .unwrap_err();
2442        assert_eq!(err.to_string(), "StorageFailed");
2443
2444        // Verify the message was not added to dedup state
2445        let session = rt.get_session_checked(&sid).await.unwrap();
2446        assert!(!session.seen_message_ids.contains("m2"));
2447    }
2448
2449    #[tokio::test]
2450    async fn cancel_session_fails_if_log_append_fails() {
2451        use std::io;
2452        use std::sync::atomic::{AtomicUsize, Ordering};
2453
2454        struct FailOnSecondAppend {
2455            count: AtomicUsize,
2456        }
2457        #[async_trait::async_trait]
2458        impl StorageBackend for FailOnSecondAppend {
2459            async fn save_session(&self, _: &Session) -> io::Result<()> {
2460                Ok(())
2461            }
2462            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2463                Ok(None)
2464            }
2465            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2466                Ok(vec![])
2467            }
2468            async fn delete_session(&self, _: &str) -> io::Result<()> {
2469                Ok(())
2470            }
2471            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2472                Ok(vec![])
2473            }
2474            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2475                let n = self.count.fetch_add(1, Ordering::SeqCst);
2476                if n >= 1 {
2477                    Err(io::Error::other("disk full"))
2478                } else {
2479                    Ok(())
2480                }
2481            }
2482            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2483                Ok(vec![])
2484            }
2485            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2486                Ok(())
2487            }
2488        }
2489
2490        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2491            count: AtomicUsize::new(0),
2492        });
2493        let registry = Arc::new(SessionRegistry::new());
2494        let log_store = Arc::new(LogStore::new());
2495        let rt = Runtime::new(storage, registry, log_store);
2496        let sid = new_sid();
2497
2498        rt.process(
2499            &env(
2500                "macp.mode.decision.v1",
2501                "SessionStart",
2502                "m1",
2503                &sid,
2504                "agent://orchestrator",
2505                session_start(vec!["agent://fraud".into()]),
2506            ),
2507            None,
2508        )
2509        .await
2510        .unwrap();
2511
2512        let err = rt
2513            .cancel_session(&sid, "test cancel", "agent://orchestrator")
2514            .await
2515            .unwrap_err();
2516        assert_eq!(err.to_string(), "StorageFailed");
2517    }
2518
2519    #[tokio::test]
2520    async fn ttl_expiration_rejects_message() {
2521        let rt = make_runtime();
2522        let sid = new_sid();
2523        let payload = SessionStartPayload {
2524            intent: "intent".into(),
2525            participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2526            mode_version: "1.0.0".into(),
2527            configuration_version: "cfg-1".into(),
2528            policy_version: String::new(),
2529            ttl_ms: 1,
2530            context_id: String::new(),
2531            extensions: std::collections::HashMap::new(),
2532            roots: vec![],
2533            max_suspend_ms: 0,
2534        }
2535        .encode_to_vec();
2536        rt.process(
2537            &env(
2538                "macp.mode.decision.v1",
2539                "SessionStart",
2540                "m1",
2541                &sid,
2542                "agent://orchestrator",
2543                payload,
2544            ),
2545            None,
2546        )
2547        .await
2548        .unwrap();
2549        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2550        let proposal = ProposalPayload {
2551            proposal_id: "p1".into(),
2552            option: "step-up".into(),
2553            rationale: "risk".into(),
2554            supporting_data: vec![],
2555        }
2556        .encode_to_vec();
2557        let err = rt
2558            .process(
2559                &env(
2560                    "macp.mode.decision.v1",
2561                    "Proposal",
2562                    "m2",
2563                    &sid,
2564                    "agent://orchestrator",
2565                    proposal,
2566                ),
2567                None,
2568            )
2569            .await
2570            .unwrap_err();
2571        assert_eq!(err.to_string(), "TtlExpired");
2572    }
2573
2574    #[tokio::test]
2575    async fn cleanup_expired_sessions_marks_expired() {
2576        let rt = make_runtime();
2577        let sid = new_sid();
2578        let payload = SessionStartPayload {
2579            intent: "intent".into(),
2580            participants: vec!["agent://fraud".into()],
2581            mode_version: "1.0.0".into(),
2582            configuration_version: "cfg-1".into(),
2583            policy_version: String::new(),
2584            ttl_ms: 1,
2585            context_id: String::new(),
2586            extensions: std::collections::HashMap::new(),
2587            roots: vec![],
2588            max_suspend_ms: 0,
2589        }
2590        .encode_to_vec();
2591        rt.process(
2592            &env(
2593                "macp.mode.decision.v1",
2594                "SessionStart",
2595                "m1",
2596                &sid,
2597                "agent://orchestrator",
2598                payload,
2599            ),
2600            None,
2601        )
2602        .await
2603        .unwrap();
2604        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2605        rt.cleanup_expired_sessions().await;
2606        let session = rt.get_session_checked(&sid).await.unwrap();
2607        assert_eq!(session.state, SessionState::Expired);
2608    }
2609
2610    #[tokio::test]
2611    async fn evict_stale_sessions_removes_resolved() {
2612        let rt = make_runtime();
2613        let sid = new_sid();
2614        // Start a decision session
2615        rt.process(
2616            &env(
2617                "macp.mode.decision.v1",
2618                "SessionStart",
2619                "m1",
2620                &sid,
2621                "agent://orchestrator",
2622                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2623            ),
2624            None,
2625        )
2626        .await
2627        .unwrap();
2628        // Send a Proposal
2629        let proposal = ProposalPayload {
2630            proposal_id: "p1".into(),
2631            option: "step-up".into(),
2632            rationale: "risk".into(),
2633            supporting_data: vec![],
2634        }
2635        .encode_to_vec();
2636        rt.process(
2637            &env(
2638                "macp.mode.decision.v1",
2639                "Proposal",
2640                "m2",
2641                &sid,
2642                "agent://orchestrator",
2643                proposal,
2644            ),
2645            None,
2646        )
2647        .await
2648        .unwrap();
2649        // Commit to resolve the session
2650        let commitment = CommitmentPayload {
2651            commitment_id: "c1".into(),
2652            action: "decision.selected".into(),
2653            authority_scope: "payments".into(),
2654            reason: "bound".into(),
2655            mode_version: "1.0.0".into(),
2656            policy_version: "policy.default".into(),
2657            configuration_version: "cfg-1".into(),
2658            outcome_positive: true,
2659            supersedes: None,
2660        }
2661        .encode_to_vec();
2662        let result = rt
2663            .process(
2664                &env(
2665                    "macp.mode.decision.v1",
2666                    "Commitment",
2667                    "m3",
2668                    &sid,
2669                    "agent://orchestrator",
2670                    commitment,
2671                ),
2672                None,
2673            )
2674            .await
2675            .unwrap();
2676        assert_eq!(result.session_state, SessionState::Resolved);
2677        // Wait a moment so the session's started_at_unix_ms is strictly in the past
2678        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2679        // Evict with retention = 0 (evict immediately)
2680        rt.evict_stale_sessions(0).await;
2681        // Session should no longer be in the in-memory registry
2682        assert!(rt.registry.get_session(&sid).await.is_none());
2683    }
2684
2685    #[tokio::test]
2686    async fn session_start_with_wrong_mode_version_rejected() {
2687        let rt = make_runtime();
2688        let sid = new_sid();
2689        let payload = SessionStartPayload {
2690            intent: "test".into(),
2691            participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2692            mode_version: "99.0.0".into(), // wrong version
2693            configuration_version: "cfg-1".into(),
2694            policy_version: String::new(),
2695            ttl_ms: 60_000,
2696            context_id: String::new(),
2697            extensions: std::collections::HashMap::new(),
2698            roots: vec![],
2699            max_suspend_ms: 0,
2700        }
2701        .encode_to_vec();
2702
2703        let err = rt
2704            .process(
2705                &env(
2706                    "macp.mode.decision.v1",
2707                    "SessionStart",
2708                    "m1",
2709                    &sid,
2710                    "agent://orchestrator",
2711                    payload,
2712                ),
2713                None,
2714            )
2715            .await
2716            .unwrap_err();
2717        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2718    }
2719
2720    #[tokio::test]
2721    async fn signal_empty_signal_type_rejected() {
2722        let rt = make_runtime();
2723        // Use non-default data so proto3 serializes a non-empty payload
2724        let signal_payload = crate::pb::SignalPayload {
2725            signal_type: String::new(),
2726            data: b"some data".to_vec(),
2727            confidence: 0.0,
2728            correlation_session_id: String::new(),
2729        }
2730        .encode_to_vec();
2731        let signal = Envelope {
2732            macp_version: "1.0".into(),
2733            mode: String::new(),
2734            message_type: "Signal".into(),
2735            message_id: "sig-1".into(),
2736            session_id: String::new(),
2737            sender: "agent://a".into(),
2738            timestamp_unix_ms: 0,
2739            payload: signal_payload,
2740        };
2741        let err = rt.process_signal(&signal).await.unwrap_err();
2742        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2743    }
2744
2745    #[tokio::test]
2746    async fn signal_valid_payload_accepted() {
2747        let rt = make_runtime();
2748        let signal_payload = crate::pb::SignalPayload {
2749            signal_type: "heartbeat".into(),
2750            data: vec![],
2751            confidence: 0.8,
2752            correlation_session_id: String::new(),
2753        }
2754        .encode_to_vec();
2755        let signal = Envelope {
2756            macp_version: "1.0".into(),
2757            mode: String::new(),
2758            message_type: "Signal".into(),
2759            message_id: "sig-2".into(),
2760            session_id: String::new(),
2761            sender: "agent://a".into(),
2762            timestamp_unix_ms: 0,
2763            payload: signal_payload,
2764        };
2765        rt.process_signal(&signal).await.unwrap();
2766    }
2767
2768    #[tokio::test]
2769    async fn signal_empty_payload_accepted() {
2770        let rt = make_runtime();
2771        let signal = Envelope {
2772            macp_version: "1.0".into(),
2773            mode: String::new(),
2774            message_type: "Signal".into(),
2775            message_id: "sig-3".into(),
2776            session_id: String::new(),
2777            sender: "agent://a".into(),
2778            timestamp_unix_ms: 0,
2779            payload: vec![],
2780        };
2781        rt.process_signal(&signal).await.unwrap();
2782    }
2783
2784    /// Freeze invariant: CommitmentPayload version fields must match the
2785    /// session-bound versions — for extension modes too. When a non-strict ext
2786    /// mode's SessionStart omits mode_version, the runtime binds the registered
2787    /// descriptor's version; a Commitment carrying "" must no longer match
2788    /// vacuously.
2789    #[tokio::test]
2790    async fn ext_mode_empty_version_binds_descriptor_version() {
2791        let rt = make_runtime();
2792        rt.register_extension(ModeDescriptor {
2793            mode: "ext.dyn.v1".into(),
2794            mode_version: "2.5.0".into(),
2795            message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2796            terminal_message_types: vec!["Commitment".into()],
2797            ..Default::default()
2798        })
2799        .unwrap();
2800
2801        let sid = new_sid();
2802        let payload = SessionStartPayload {
2803            participants: vec!["alice".into()],
2804            configuration_version: "cfg-1".into(),
2805            ttl_ms: 60_000,
2806            ..Default::default()
2807        }
2808        .encode_to_vec();
2809        rt.process(
2810            &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2811            None,
2812        )
2813        .await
2814        .unwrap();
2815
2816        // The session is bound to the descriptor's version, not "".
2817        let session = rt.get_session_checked(&sid).await.unwrap();
2818        assert_eq!(session.mode_version, "2.5.0");
2819
2820        // Commitment with empty mode_version: rejected (no vacuous match).
2821        let bad = CommitmentPayload {
2822            commitment_id: "c1".into(),
2823            action: "work.completed".into(),
2824            authority_scope: "test".into(),
2825            reason: "done".into(),
2826            mode_version: String::new(),
2827            policy_version: "policy.default".into(),
2828            configuration_version: "cfg-1".into(),
2829            outcome_positive: true,
2830            supersedes: None,
2831        }
2832        .encode_to_vec();
2833        let err = rt
2834            .process(
2835                &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2836                None,
2837            )
2838            .await
2839            .unwrap_err();
2840        assert_eq!(err.to_string(), "InvalidPayload");
2841
2842        // Commitment echoing the bound descriptor version: accepted, resolves.
2843        let good = CommitmentPayload {
2844            commitment_id: "c1".into(),
2845            action: "work.completed".into(),
2846            authority_scope: "test".into(),
2847            reason: "done".into(),
2848            mode_version: "2.5.0".into(),
2849            policy_version: "policy.default".into(),
2850            configuration_version: "cfg-1".into(),
2851            outcome_positive: true,
2852            supersedes: None,
2853        }
2854        .encode_to_vec();
2855        let result = rt
2856            .process(
2857                &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2858                None,
2859            )
2860            .await
2861            .unwrap();
2862        assert_eq!(result.session_state, SessionState::Resolved);
2863    }
2864
2865    /// The binding must be recorded on the SessionStart log entry (replay reads
2866    /// it from there), and only when the payload actually omitted the version.
2867    #[tokio::test]
2868    async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2869        let rt = make_runtime();
2870        rt.register_extension(ModeDescriptor {
2871            mode: "ext.dyn2.v1".into(),
2872            mode_version: "3.0.0".into(),
2873            message_types: vec!["SessionStart".into(), "Commitment".into()],
2874            terminal_message_types: vec!["Commitment".into()],
2875            ..Default::default()
2876        })
2877        .unwrap();
2878
2879        let sid = new_sid();
2880        let payload = SessionStartPayload {
2881            participants: vec!["alice".into()],
2882            configuration_version: "cfg-1".into(),
2883            ttl_ms: 60_000,
2884            ..Default::default()
2885        }
2886        .encode_to_vec();
2887        rt.process(
2888            &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2889            None,
2890        )
2891        .await
2892        .unwrap();
2893
2894        let log = rt.log_store.get_log(&sid).await.unwrap();
2895        assert_eq!(log[0].message_type, "SessionStart");
2896        assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2897
2898        // A payload that carries the version explicitly records no binding.
2899        let sid2 = new_sid();
2900        let payload2 = SessionStartPayload {
2901            participants: vec!["alice".into()],
2902            mode_version: "3.0.0".into(),
2903            configuration_version: "cfg-1".into(),
2904            ttl_ms: 60_000,
2905            ..Default::default()
2906        }
2907        .encode_to_vec();
2908        rt.process(
2909            &env(
2910                "ext.dyn2.v1",
2911                "SessionStart",
2912                "m1",
2913                &sid2,
2914                "alice",
2915                payload2,
2916            ),
2917            None,
2918        )
2919        .await
2920        .unwrap();
2921        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2922        assert_eq!(log2[0].bound_mode_version, None);
2923    }
2924
2925    /// The RESOLVED suspension cap is bound on the session and recorded on
2926    /// the SessionStart log entry (RFC-MACP-0001 §7.5, RFC-MACP-0003 §2):
2927    /// the payload's positive value verbatim, or the runtime default when
2928    /// the payload carried 0 — never left unrecorded on new sessions.
2929    #[tokio::test]
2930    async fn session_start_binds_and_records_max_suspend_cap() {
2931        let rt = make_runtime();
2932
2933        // Explicit cap: recorded verbatim.
2934        let sid = new_sid();
2935        let payload = SessionStartPayload {
2936            participants: vec!["alice".into(), "bob".into()],
2937            mode_version: "1.0.0".into(),
2938            configuration_version: "cfg-1".into(),
2939            ttl_ms: 60_000,
2940            max_suspend_ms: 12_345,
2941            ..Default::default()
2942        }
2943        .encode_to_vec();
2944        rt.process(
2945            &env(
2946                "macp.mode.decision.v1",
2947                "SessionStart",
2948                "m1",
2949                &sid,
2950                "alice",
2951                payload,
2952            ),
2953            None,
2954        )
2955        .await
2956        .unwrap();
2957        let log = rt.log_store.get_log(&sid).await.unwrap();
2958        assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2959
2960        // Payload 0: the runtime default is resolved and recorded.
2961        let sid2 = new_sid();
2962        let payload2 = SessionStartPayload {
2963            participants: vec!["alice".into(), "bob".into()],
2964            mode_version: "1.0.0".into(),
2965            configuration_version: "cfg-1".into(),
2966            ttl_ms: 60_000,
2967            max_suspend_ms: 0,
2968            ..Default::default()
2969        }
2970        .encode_to_vec();
2971        rt.process(
2972            &env(
2973                "macp.mode.decision.v1",
2974                "SessionStart",
2975                "m2",
2976                &sid2,
2977                "alice",
2978                payload2,
2979            ),
2980            None,
2981        )
2982        .await
2983        .unwrap();
2984        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2985        assert_eq!(
2986            log2[0].bound_max_suspend_ms,
2987            Some(macp_core::session::MAX_SUSPEND_MS)
2988        );
2989    }
2990
2991    #[test]
2992    fn audit_verbosity_reads_policy_rules() {
2993        let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2994        assert!(!Runtime::audit_verbose(&session));
2995
2996        session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2997            policy_id: "policy.test.audit".into(),
2998            mode: "*".into(),
2999            description: "audited".into(),
3000            rules: serde_json::json!({ "audit": { "level": "info" } }),
3001            schema_version: 1,
3002        });
3003        assert!(Runtime::audit_verbose(&session));
3004
3005        session.policy_definition.as_mut().unwrap().rules =
3006            serde_json::json!({ "audit": { "level": "debug" } });
3007        assert!(!Runtime::audit_verbose(&session));
3008    }
3009
3010    /// Post-commit-point coherence: once the SessionStart log entry is
3011    /// durable, a snapshot failure must NOT fail (or roll back) the start —
3012    /// the previous fatal path left the durable entry behind, so the
3013    /// "failed" session resurrected on restart and a same-id retry appended
3014    /// a second SessionStart that made the log unreplayable.
3015    #[tokio::test]
3016    async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
3017        use std::io;
3018
3019        struct FailSnapshotBackend;
3020        #[async_trait::async_trait]
3021        impl StorageBackend for FailSnapshotBackend {
3022            async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
3023                Ok(())
3024            }
3025            async fn save_session(&self, _s: &Session) -> io::Result<()> {
3026                Err(io::Error::other("snapshot disk full"))
3027            }
3028            async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
3029                Ok(None)
3030            }
3031            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
3032                Ok(vec![])
3033            }
3034            async fn delete_session(&self, _s: &str) -> io::Result<()> {
3035                Ok(())
3036            }
3037            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
3038                Ok(vec![])
3039            }
3040            async fn append_log_entry(
3041                &self,
3042                _s: &str,
3043                _e: &crate::log_store::LogEntry,
3044            ) -> io::Result<()> {
3045                Ok(())
3046            }
3047            async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
3048                Ok(vec![])
3049            }
3050        }
3051
3052        let rt = Runtime::new(
3053            Arc::new(FailSnapshotBackend),
3054            Arc::new(SessionRegistry::new()),
3055            Arc::new(LogStore::new()),
3056        );
3057        let sid = new_sid();
3058        let result = rt
3059            .process(
3060                &env(
3061                    "macp.mode.decision.v1",
3062                    "SessionStart",
3063                    "m1",
3064                    &sid,
3065                    "agent://orchestrator",
3066                    session_start(vec!["agent://orchestrator".into()]),
3067                ),
3068                None,
3069            )
3070            .await
3071            .expect("start must succeed: the log append (commit point) succeeded");
3072        assert!(!result.duplicate);
3073        // The session exists and is usable.
3074        assert!(rt.get_session_checked(&sid).await.is_some());
3075    }
3076
3077    // --- The client boundary at the kernel's two live entry points (11c) ---
3078    //
3079    // RFC-MACP-0010 §5.1(3). These are the runtime-level halves; the
3080    // mode-level rules are unit-tested in
3081    // `crates/macp-modes/src/mode/handoff.rs`, `step::validate_message` in
3082    // `crates/macp-modes/src/step.rs`, and the proof that the hook is NOT on
3083    // the replay path in `src/replay.rs`
3084    // (`reserved_prefix_entry_replays_at_every_rev`).
3085    //
3086    // Both entry points matter independently: `process_message` and
3087    // `process_session_start` each call the hook themselves, because the
3088    // runtime deliberately bypasses `macp_modes::step::validate_message` (it
3089    // interposes its durable append between validation and commit), so neither
3090    // call site is covered by the other.
3091
3092    const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
3093    const OWNER: &str = "agent://owner";
3094    const TARGET: &str = "agent://target";
3095
3096    fn reserved_id(handoff_id: &str) -> String {
3097        format!(
3098            "{}{handoff_id}",
3099            crate::mode::handoff::IMPLICIT_ACCEPT_MESSAGE_ID_PREFIX
3100        )
3101    }
3102
3103    fn handoff_start_payload() -> Vec<u8> {
3104        session_start(vec![OWNER.into(), TARGET.into()])
3105    }
3106
3107    /// A `Commitment` that echoes the versions `handoff_session_with_offer`
3108    /// binds, so the only thing left to reject it is the `message_id`.
3109    fn handoff_commitment_payload() -> Vec<u8> {
3110        CommitmentPayload {
3111            commitment_id: "c1".into(),
3112            action: "handoff.accepted".into(),
3113            authority_scope: "support".into(),
3114            reason: "bound".into(),
3115            mode_version: "1.0.0".into(),
3116            policy_version: "policy.default".into(),
3117            configuration_version: "cfg-1".into(),
3118            outcome_positive: true,
3119            supersedes: None,
3120        }
3121        .encode_to_vec()
3122    }
3123
3124    fn handoff_offer(handoff_id: &str) -> Vec<u8> {
3125        crate::handoff_pb::HandoffOfferPayload {
3126            handoff_id: handoff_id.into(),
3127            target_participant: TARGET.into(),
3128            scope: "support".into(),
3129            reason: "escalate".into(),
3130        }
3131        .encode_to_vec()
3132    }
3133
3134    fn handoff_context(handoff_id: &str) -> Vec<u8> {
3135        crate::handoff_pb::HandoffContextPayload {
3136            handoff_id: handoff_id.into(),
3137            content_type: "text/plain".into(),
3138            context: b"background".to_vec(),
3139        }
3140        .encode_to_vec()
3141    }
3142
3143    fn handoff_accept(handoff_id: &str, implicit: bool) -> Vec<u8> {
3144        crate::handoff_pb::HandoffAcceptPayload {
3145            handoff_id: handoff_id.into(),
3146            accepted_by: TARGET.into(),
3147            reason: "ready".into(),
3148            implicit,
3149        }
3150        .encode_to_vec()
3151    }
3152
3153    /// An Open handoff session at the current semantics revision with one
3154    /// outstanding offer `h1`. Returns the session id.
3155    async fn handoff_session_with_offer(rt: &Runtime) -> String {
3156        let sid = new_sid();
3157        rt.process(
3158            &env(
3159                HANDOFF_MODE,
3160                "SessionStart",
3161                "start-1",
3162                &sid,
3163                OWNER,
3164                handoff_start_payload(),
3165            ),
3166            None,
3167        )
3168        .await
3169        .expect("handoff session start");
3170        rt.process(
3171            &env(
3172                HANDOFF_MODE,
3173                "HandoffOffer",
3174                "offer-1",
3175                &sid,
3176                OWNER,
3177                handoff_offer("h1"),
3178            ),
3179            None,
3180        )
3181        .await
3182        .expect("handoff offer");
3183        assert_eq!(
3184            rt.get_session_checked(&sid).await.unwrap().semantics_rev,
3185            macp_core::session::CURRENT_SEMANTICS_REV
3186        );
3187        sid
3188    }
3189
3190    /// Phase 11c criterion 1, message path: at the current revision a client
3191    /// envelope whose `message_id` is in the reserved `implicit-accept:`
3192    /// namespace is rejected `InvalidEnvelope`, and the rejection mutates
3193    /// nothing — neither accepted history nor `seen_message_ids` (CLAUDE.md §8
3194    /// dedup invariant).
3195    ///
3196    /// `HandoffContext` is the interesting carrier: it is the one mode message
3197    /// the offerer may send at any disposition, and it is rejected here
3198    /// **before** dispatch, which is why the check cannot live in the mode's
3199    /// own rules. `Commitment` is included because the initiator can squat the
3200    /// id that way too, and because it proves the reservation is not scoped to
3201    /// `HandoffAccept`.
3202    #[tokio::test]
3203    async fn reserved_message_id_namespace_is_rejected_at_rev2() {
3204        let rt = make_runtime();
3205        let sid = handoff_session_with_offer(&rt).await;
3206
3207        let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3208        let dedup_before = rt
3209            .get_session_checked(&sid)
3210            .await
3211            .unwrap()
3212            .seen_message_ids
3213            .clone();
3214
3215        for (message_type, sender, payload) in [
3216            ("HandoffContext", OWNER, handoff_context("h1")),
3217            ("Commitment", OWNER, handoff_commitment_payload()),
3218            ("HandoffAccept", TARGET, handoff_accept("h1", false)),
3219        ] {
3220            let err = rt
3221                .process(
3222                    &env(
3223                        HANDOFF_MODE,
3224                        message_type,
3225                        &reserved_id("h1"),
3226                        &sid,
3227                        sender,
3228                        payload,
3229                    ),
3230                    None,
3231                )
3232                .await
3233                .unwrap_err();
3234            assert!(
3235                matches!(err, MacpError::InvalidEnvelope),
3236                "{message_type} with a reserved id must be InvalidEnvelope, got {err}"
3237            );
3238        }
3239
3240        // Nothing was appended and no dedup slot was consumed.
3241        let session = rt.get_session_checked(&sid).await.unwrap();
3242        assert_eq!(
3243            rt.log_store.get_log(&sid).await.unwrap().len(),
3244            history_before
3245        );
3246        assert_eq!(session.seen_message_ids, dedup_before);
3247        assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3248        assert_eq!(session.state, SessionState::Open);
3249
3250        // Control: the same `HandoffContext` with an ordinary id is accepted,
3251        // so the rejections above are the id's doing and not the payload's.
3252        rt.process(
3253            &env(
3254                HANDOFF_MODE,
3255                "HandoffContext",
3256                "ctx-1",
3257                &sid,
3258                OWNER,
3259                handoff_context("h1"),
3260            ),
3261            None,
3262        )
3263        .await
3264        .expect("an ordinary id is accepted");
3265        assert_eq!(
3266            rt.log_store.get_log(&sid).await.unwrap().len(),
3267            history_before + 1
3268        );
3269    }
3270
3271    /// Phase 11c criterion 1, start path. This call site exists precisely
3272    /// because the namespace is squattable through `SessionStart`, whose
3273    /// `message_id` enters `seen_message_ids` at the commit point — after
3274    /// which the runtime's own later synthesis would be silently skipped.
3275    ///
3276    /// The rejection must leave **no** trace: no session in the registry, no
3277    /// session log, and no reservation of the session id — proved by starting
3278    /// the same session id again with a clean `message_id` and having it
3279    /// succeed (`SessionAlreadyExists` would be the failure signature of a
3280    /// leaked reservation).
3281    #[tokio::test]
3282    async fn reserved_message_id_is_rejected_on_the_session_start_path() {
3283        let rt = make_runtime();
3284        let sid = new_sid();
3285
3286        let err = rt
3287            .process(
3288                &env(
3289                    HANDOFF_MODE,
3290                    "SessionStart",
3291                    &reserved_id("h1"),
3292                    &sid,
3293                    OWNER,
3294                    handoff_start_payload(),
3295                ),
3296                None,
3297            )
3298            .await
3299            .unwrap_err();
3300        assert!(
3301            matches!(err, MacpError::InvalidEnvelope),
3302            "reserved id on SessionStart must be InvalidEnvelope, got {err}"
3303        );
3304
3305        // No session, no log, nothing to roll back.
3306        assert!(rt.get_session_checked(&sid).await.is_none());
3307        assert!(rt.log_store.get_log(&sid).await.is_none());
3308        assert!(!rt.registry.sessions.read().await.contains_key(&sid));
3309
3310        // The session id was never reserved and the id never consumed a dedup
3311        // slot: the same session starts cleanly.
3312        rt.process(
3313            &env(
3314                HANDOFF_MODE,
3315                "SessionStart",
3316                "start-1",
3317                &sid,
3318                OWNER,
3319                handoff_start_payload(),
3320            ),
3321            None,
3322        )
3323        .await
3324        .expect("a rejected SessionStart must not reserve the session id");
3325        let session = rt.get_session_checked(&sid).await.unwrap();
3326        assert!(session.seen_message_ids.contains("start-1"));
3327        assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3328    }
3329
3330    /// Phase 11c criterion 2, runtime path: a client `HandoffAccept` carrying
3331    /// `implicit = true` never enters history.
3332    ///
3333    /// Two envelopes, as the criterion requires:
3334    /// (a) `implicit = true` with an ordinary `message_id` — `InvalidPayload`.
3335    ///     **This assertion is double-guarded**: `handle_message` rejects the
3336    ///     envelope at the client boundary before dispatch ever sees it
3337    ///     (`crates/macp-modes/src/mode/handoff.rs`, the client-envelope
3338    ///     validation that refuses a client-submitted `implicit = true`), and
3339    ///     — since 11d — `dispatch_implicit_accept`'s own reserved-`message_id`
3340    ///     check (`handoff.rs:733-740`) would refuse the same envelope on its
3341    ///     `message_id` alone even if the boundary check were removed. What it
3342    ///     pins is the criterion's actual requirement — that the rev-2 error
3343    ///     *surface* through `Send` is unchanged — not which single guard is
3344    ///     load-bearing.
3345    /// (b) `implicit = true` with the **reserved** `message_id` and the
3346    ///     correct sender — `InvalidEnvelope`, which only the boundary can
3347    ///     produce (dispatch would say `InvalidPayload`). The envelope is
3348    ///     byte-shaped exactly like the one the runtime synthesizes in
3349    ///     `synthesize_due_accept`.
3350    ///
3351    /// Measured, **neither half pins the `implicit` rule in isolation**: (b)
3352    /// is killed by the *reserved-prefix* rule, which fires first and returns
3353    /// `InvalidEnvelope` whatever the flag says, so deleting the `implicit`
3354    /// rule leaves this whole test green. The `implicit` rule's only
3355    /// non-vacuous guard is the mode-level unit test
3356    /// `handoff::tests::client_implicit_accept_rejected_at_the_boundary`.
3357    /// What this test pins is the runtime-level *error surface* at rev 2 —
3358    /// which is the criterion's requirement.
3359    #[tokio::test]
3360    async fn client_implicit_accept_rejected_through_the_runtime() {
3361        let rt = make_runtime();
3362        let sid = handoff_session_with_offer(&rt).await;
3363        let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3364
3365        // (a) ordinary id.
3366        let err = rt
3367            .process(
3368                &env(
3369                    HANDOFF_MODE,
3370                    "HandoffAccept",
3371                    "accept-1",
3372                    &sid,
3373                    TARGET,
3374                    handoff_accept("h1", true),
3375                ),
3376                None,
3377            )
3378            .await
3379            .unwrap_err();
3380        assert!(matches!(err, MacpError::InvalidPayload), "got {err}");
3381
3382        // (b) the runtime's own synthetic shape, submitted by the target.
3383        let err = rt
3384            .process(
3385                &env(
3386                    HANDOFF_MODE,
3387                    "HandoffAccept",
3388                    &reserved_id("h1"),
3389                    &sid,
3390                    TARGET,
3391                    handoff_accept("h1", true),
3392                ),
3393                None,
3394            )
3395            .await
3396            .unwrap_err();
3397        assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3398
3399        // The offer is still outstanding and history is untouched.
3400        let session = rt.get_session_checked(&sid).await.unwrap();
3401        assert_eq!(
3402            rt.log_store.get_log(&sid).await.unwrap().len(),
3403            history_before
3404        );
3405        assert!(session.seen_message_ids.is_disjoint(
3406            &["accept-1".to_string(), reserved_id("h1")]
3407                .into_iter()
3408                .collect()
3409        ));
3410        let mode_state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
3411        assert_eq!(mode_state["offers"]["h1"]["disposition"], "Offered");
3412
3413        // Control: the explicit accept (`implicit = false`, ordinary id) is
3414        // accepted, so the rejections above are not the envelope's other
3415        // fields.
3416        rt.process(
3417            &env(
3418                HANDOFF_MODE,
3419                "HandoffAccept",
3420                "accept-2",
3421                &sid,
3422                TARGET,
3423                handoff_accept("h1", false),
3424            ),
3425            None,
3426        )
3427        .await
3428        .expect("an explicit accept is still accepted");
3429    }
3430
3431    /// Error-code ordering at the kernel is unchanged by the phase: an
3432    /// envelope that is **both** unauthorized and carries a reserved id
3433    /// reports `Forbidden`, because the hook is called after
3434    /// `mode.authorize_sender`. Rev <= 1's Forbidden-before-payload ordering
3435    /// therefore does not shift at rev 2.
3436    #[tokio::test]
3437    async fn client_boundary_error_ordering_is_unchanged_at_rev2() {
3438        let rt = make_runtime();
3439        let sid = handoff_session_with_offer(&rt).await;
3440
3441        // `agent://stranger` is not a declared participant.
3442        let err = rt
3443            .process(
3444                &env(
3445                    HANDOFF_MODE,
3446                    "HandoffAccept",
3447                    &reserved_id("h1"),
3448                    &sid,
3449                    "agent://stranger",
3450                    handoff_accept("h1", true),
3451                ),
3452                None,
3453            )
3454            .await
3455            .unwrap_err();
3456        assert!(
3457            matches!(err, MacpError::Forbidden),
3458            "authorization must be reported before the client boundary, got {err}"
3459        );
3460
3461        // Same envelope from the authorized sender: now the boundary speaks.
3462        let err = rt
3463            .process(
3464                &env(
3465                    HANDOFF_MODE,
3466                    "HandoffAccept",
3467                    &reserved_id("h1"),
3468                    &sid,
3469                    TARGET,
3470                    handoff_accept("h1", true),
3471                ),
3472                None,
3473            )
3474            .await
3475            .unwrap_err();
3476        assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3477    }
3478
3479    // --- The synthesis seam's own guards (11e) ---
3480
3481    /// A handoff session with a bound `implicit_accept_timeout_ms` and an
3482    /// outstanding offer. Unlike [`handoff_session_with_offer`] this one binds
3483    /// a policy, so an implicit accept can actually become due.
3484    async fn handoff_session_with_timed_offer(rt: &Runtime, timeout_ms: i64) -> String {
3485        rt.register_policy(macp_core::policy::PolicyDefinition {
3486            policy_id: "handoff-timed".into(),
3487            mode: HANDOFF_MODE.into(),
3488            description: "implicit accept".into(),
3489            rules: serde_json::json!({
3490                "acceptance": { "implicit_accept_timeout_ms": timeout_ms },
3491                "commitment": { "authority": "initiator_only" }
3492            }),
3493            schema_version: 1,
3494        })
3495        .expect("policy registers");
3496
3497        let sid = new_sid();
3498        let start = SessionStartPayload {
3499            intent: "escalate".into(),
3500            participants: vec![OWNER.into(), TARGET.into()],
3501            mode_version: "1.0.0".into(),
3502            configuration_version: "cfg-1".into(),
3503            policy_version: "handoff-timed".into(),
3504            ttl_ms: 60_000,
3505            context_id: String::new(),
3506            extensions: std::collections::HashMap::new(),
3507            roots: vec![],
3508            max_suspend_ms: 0,
3509        }
3510        .encode_to_vec();
3511        rt.process(
3512            &env(HANDOFF_MODE, "SessionStart", "start-1", &sid, OWNER, start),
3513            None,
3514        )
3515        .await
3516        .expect("session start");
3517        rt.process(
3518            &env(
3519                HANDOFF_MODE,
3520                "HandoffOffer",
3521                "offer-1",
3522                &sid,
3523                OWNER,
3524                handoff_offer("h1"),
3525            ),
3526            None,
3527        )
3528        .await
3529        .expect("offer");
3530        sid
3531    }
3532
3533    /// Phase 11e criterion 11, first half: `synthesize_due_accept` must refuse
3534    /// to synthesize for a session that is not `Open`.
3535    ///
3536    /// This is a **correctness** gate, not hygiene, and it is the only one
3537    /// that ships. Both computations behind the synthetic entry ignore an
3538    /// in-flight pause by design — `Session::unsuspended_deadline` walks
3539    /// completed intervals only, and `HandoffMode::rev2_elapsed_ms` has no
3540    /// in-flight term — so asking a `Suspended` session over-counts elapsed
3541    /// time (the accept can be judged due when it is not) and under-computes
3542    /// `D` (a wrong `timestamp_unix_ms` baked into permanent history, in
3543    /// silence). The mode carries a `debug_assert!` for the same thing, which
3544    /// compiles out in release builds and therefore guards nothing where it
3545    /// matters.
3546    ///
3547    /// The method is called directly because no message path can reach it with
3548    /// a suspended session — `step::check_preconditions` rejects every message
3549    /// to a non-`Open` session first. The eager sweep (`sweep_due_synthetic_accepts`)
3550    /// is the caller that *does* see suspended sessions, which is exactly why
3551    /// this filter has to live in the seam rather than at a single call site.
3552    #[tokio::test]
3553    async fn synthesis_is_skipped_for_a_non_open_session() {
3554        let rt = make_runtime();
3555        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3556        rt.suspend_session(&sid, "hold", OWNER)
3557            .await
3558            .expect("suspend");
3559
3560        let shared = rt.registry.get_shared(&sid).await.unwrap();
3561        let mut guard = shared.lock().await;
3562        let session = &mut *guard;
3563        assert_eq!(session.state, SessionState::Suspended);
3564        assert!(session.suspended_at_ms.is_some());
3565
3566        // Far past the deadline on wall time — the only thing stopping a
3567        // synthesis here is the state filter.
3568        let long_after = session.suspended_at_ms.unwrap() + 10_000;
3569        let log_before = rt.log_store.get_log(&sid).await.unwrap_or_default().len();
3570        let dedup_before = session.seen_message_ids.len();
3571        let mode_state_before = session.mode_state.clone();
3572
3573        rt.synthesize_due_accept(&sid, session, long_after)
3574            .await
3575            .expect("the filter is a skip, not an error");
3576
3577        assert_eq!(
3578            rt.log_store.get_log(&sid).await.unwrap_or_default().len(),
3579            log_before,
3580            "a suspended session must not gain a synthetic entry"
3581        );
3582        assert_eq!(session.seen_message_ids.len(), dedup_before);
3583        assert_eq!(session.mode_state, mode_state_before);
3584
3585        // Control: the filter is what declined, not the arithmetic — and what
3586        // it kept out of history was wrong, not merely early.
3587        //
3588        // The mode cannot be asked about the genuinely suspended session in a
3589        // debug build: `due_synthetic_envelope` carries a `debug_assert!` on
3590        // `suspended_at_ms.is_none()` and would panic. So the probe is a clone
3591        // that differs *only* by that tripwire — same offer, same banked
3592        // suspension (none: the pause is still in flight and banks on resume),
3593        // same clock. That is exactly the state the mode would see if the
3594        // kernel filter were removed and the assertion compiled out, which is
3595        // what a release build does.
3596        let mut without_the_tripwire = session.clone();
3597        without_the_tripwire.state = SessionState::Open;
3598        without_the_tripwire.suspended_at_ms = None;
3599        let mode = rt.mode_registry.get_mode(&session.mode).unwrap();
3600        let would_have_emitted = mode
3601            .due_synthetic_envelope(&without_the_tripwire, long_after)
3602            .expect("the mode would have synthesized; only the kernel filter stopped it");
3603        // The harm, concretely: the D it would have recorded falls *inside* the
3604        // pause that is still running — a moment at which the session was not
3605        // ticking at all. `unsuspended_deadline` walks completed intervals
3606        // only, and this pause is not completed, so it is invisible to the
3607        // walk. Recorded, that timestamp would be permanent and wrong.
3608        let suspended_at = session.suspended_at_ms.unwrap();
3609        assert!(
3610            would_have_emitted.timestamp_unix_ms >= suspended_at
3611                && would_have_emitted.timestamp_unix_ms < long_after,
3612            "D {} must fall inside the still-open pause starting at {suspended_at}",
3613            would_have_emitted.timestamp_unix_ms
3614        );
3615
3616        // And once the session is Open again the seam works normally, so the
3617        // filter is a skip rather than a permanent disable.
3618        drop(guard);
3619        rt.resume_session(&sid, "go", OWNER).await.expect("resume");
3620        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3621        let shared = rt.registry.get_shared(&sid).await.unwrap();
3622        let mut guard = shared.lock().await;
3623        let session = &mut *guard;
3624        let now = Utc::now().timestamp_millis();
3625        rt.synthesize_due_accept(&sid, session, now).await.unwrap();
3626        assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3627    }
3628
3629    /// Phase 11e criterion 11, second half: the appended entry stamps
3630    /// `received_at_ms` with the envelope's own `timestamp_unix_ms` (the
3631    /// deadline `D`), never wall-clock.
3632    ///
3633    /// Asserted on the stamped value directly, at the seam, rather than
3634    /// inferred from a replay — a replay would pass under a wall-clock stamp
3635    /// too, because handoff's accept arm is time-blind. The contract is
3636    /// general: replay derives its dispatch clock from `received_at_ms`, so
3637    /// the next mode to use this hook would silently diverge.
3638    #[tokio::test]
3639    async fn synthetic_entry_stamps_received_at_with_the_deadline() {
3640        let rt = make_runtime();
3641        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3642        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3643
3644        let offer_received_at = rt
3645            .log_store
3646            .get_log(&sid)
3647            .await
3648            .unwrap()
3649            .iter()
3650            .find(|e| e.message_type == "HandoffOffer")
3651            .expect("offer entry")
3652            .received_at_ms;
3653        let expected_d = offer_received_at + 20;
3654
3655        let shared = rt.registry.get_shared(&sid).await.unwrap();
3656        let mut guard = shared.lock().await;
3657        let session = &mut *guard;
3658        // Observed far later than D, so a wall-clock stamp would be obvious.
3659        let observed = Utc::now().timestamp_millis();
3660        assert!(observed > expected_d);
3661        rt.synthesize_due_accept(&sid, session, observed)
3662            .await
3663            .unwrap();
3664        drop(guard);
3665
3666        let entry = rt
3667            .log_store
3668            .get_log(&sid)
3669            .await
3670            .unwrap()
3671            .into_iter()
3672            .find(|e| e.message_id == reserved_id("h1"))
3673            .expect("the synthetic entry");
3674        assert_eq!(entry.timestamp_unix_ms, expected_d, "envelope clock is D");
3675        assert_eq!(entry.received_at_ms, expected_d, "entry clock is D");
3676        assert_ne!(
3677            entry.received_at_ms, observed,
3678            "received_at_ms must not be the observation time"
3679        );
3680        assert_eq!(entry.entry_kind, EntryKind::Incoming);
3681    }
3682
3683    /// The synthetic accept is deliberately NOT credited as participant
3684    /// activity: `replay_entry` never records activity for any entry kind, so
3685    /// calling `record_participant_activity` live would guarantee a
3686    /// live/replay divergence in `participant_message_counts` — and the target
3687    /// did not, in fact, send anything.
3688    ///
3689    /// The consequence is user-visible (`SessionMetadata.participant_activity`
3690    /// via `server::session_to_metadata`), so it is pinned rather than left as
3691    /// a comment.
3692    #[tokio::test]
3693    async fn synthetic_accept_is_not_credited_as_participant_activity() {
3694        let rt = make_runtime();
3695        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3696        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3697
3698        let shared = rt.registry.get_shared(&sid).await.unwrap();
3699        let mut guard = shared.lock().await;
3700        let session = &mut *guard;
3701        let before = session.participant_message_counts.get(TARGET).copied();
3702        rt.synthesize_due_accept(&sid, session, Utc::now().timestamp_millis())
3703            .await
3704            .unwrap();
3705        assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3706        assert_eq!(
3707            session.participant_message_counts.get(TARGET).copied(),
3708            before,
3709            "the target must not be credited with a message they did not send"
3710        );
3711    }
3712    /// The checkpoint-interval check belongs to the *append*, so the eager
3713    /// sweep and a lazy trigger place the same checkpoint in the same
3714    /// position.
3715    ///
3716    /// `maybe_insert_checkpoint` used to be called only from the tail of
3717    /// `process_message`. A synthetic entry advances `log_len` without ever
3718    /// reaching it: the lazy path checked only the length *after* the
3719    /// trigger's own append (one greater), and the eager path — where
3720    /// `sweep_due_synthetic_accepts` is the whole call — checked nothing at
3721    /// all. With `MACP_CHECKPOINT_INTERVAL > 0` that is a genuine eager/lazy
3722    /// divergence: identical accepted histories, different checkpoint
3723    /// placement, and a boundary the synthetic crossed silently skipped.
3724    /// Calling it from inside `synthesize_due_accept` makes the two agree by
3725    /// construction, which is what this pins.
3726    ///
3727    /// Checkpoints are a replay optimization rather than a correctness
3728    /// property, which is exactly why this needs a test: nothing else would
3729    /// ever notice.
3730    ///
3731    /// `checkpoint_interval` is set on the struct rather than through
3732    /// `MACP_CHECKPOINT_INTERVAL`, because that variable is read once in
3733    /// `Runtime::with_registries` and this binary runs its tests in parallel —
3734    /// setting it here would leak into every other runtime built concurrently.
3735    ///
3736    /// Interval 3, with `SessionStart` + `HandoffOffer` already logged, puts
3737    /// the synthetic accept exactly on the boundary: the skipped case.
3738    #[tokio::test]
3739    async fn a_synthetic_entry_on_the_checkpoint_boundary_checkpoints_either_path() {
3740        async fn shape(rt: &Runtime, sid: &str) -> Vec<(EntryKind, String)> {
3741            rt.log_store
3742                .get_log(sid)
3743                .await
3744                .expect("log")
3745                .iter()
3746                .map(|e| (e.entry_kind.clone(), e.message_type.clone()))
3747                .collect()
3748        }
3749
3750        // Eager: the sweep is the only thing that touches the session.
3751        let mut eager = make_runtime();
3752        eager.checkpoint_interval = 3;
3753        let eager_sid = handoff_session_with_timed_offer(&eager, 20).await;
3754        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3755        assert_eq!(
3756            eager.sweep_due_synthetic_accepts().await,
3757            1,
3758            "sweep emitted"
3759        );
3760
3761        // Lazy: a trigger message arrives instead. `HandoffContext` is the one
3762        // mode message the offerer may send at any disposition, so it provokes
3763        // the synthesis without being an accept itself.
3764        let mut lazy = make_runtime();
3765        lazy.checkpoint_interval = 3;
3766        let lazy_sid = handoff_session_with_timed_offer(&lazy, 20).await;
3767        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3768        lazy.process(
3769            &env(
3770                HANDOFF_MODE,
3771                "HandoffContext",
3772                "ctx-1",
3773                &lazy_sid,
3774                OWNER,
3775                handoff_context("h1"),
3776            ),
3777            None,
3778        )
3779        .await
3780        .expect("context accepted");
3781
3782        let expected = vec![
3783            (EntryKind::Incoming, "SessionStart".to_string()),
3784            (EntryKind::Incoming, "HandoffOffer".to_string()),
3785            (EntryKind::Incoming, "HandoffAccept".to_string()),
3786            (EntryKind::Checkpoint, "Checkpoint".to_string()),
3787        ];
3788        assert_eq!(shape(&eager, &eager_sid).await, expected, "eager sweep");
3789
3790        let lazy_shape = shape(&lazy, &lazy_sid).await;
3791        assert_eq!(
3792            lazy_shape[..4],
3793            expected[..],
3794            "the lazy path must checkpoint in the same place as the eager one"
3795        );
3796        // ... and exactly once: `process_message` runs its own check after
3797        // appending the trigger, at a length the synthesis never tested.
3798        assert_eq!(
3799            lazy_shape[4..],
3800            [(EntryKind::Incoming, "HandoffContext".to_string())],
3801            "no second checkpoint for the trigger's own append"
3802        );
3803    }
3804}