Skip to main content

macp_runtime/
runtime.rs

1use chrono::Utc;
2use std::sync::Arc;
3
4use crate::error::MacpError;
5use crate::extensions::ExtensionProviderRegistry;
6use crate::log_store::{EntryKind, LogEntry, LogStore};
7use crate::metrics::RuntimeMetrics;
8use crate::mode_registry::ModeRegistry;
9use crate::pb::{Envelope, ModeDescriptor};
10use crate::policy::registry::PolicyRegistry;
11use crate::policy::PolicyDefinition;
12use crate::registry::SessionRegistry;
13use crate::session::{
14    extract_ttl_ms, parse_session_start_payload, validate_canonical_session_start_payload_for_mode,
15    validate_session_id_for_acceptance, Session, SessionState,
16};
17use crate::storage::StorageBackend;
18use crate::stream_bus::SessionStreamBus;
19
20#[derive(Debug)]
21pub struct ProcessResult {
22    pub session_state: SessionState,
23    pub duplicate: bool,
24}
25
26#[derive(Clone, Debug)]
27pub enum SessionLifecycleEvent {
28    Created { session_id: String },
29    Resolved { session_id: String },
30    Expired { session_id: String },
31    Suspended { session_id: String },
32    Resumed { session_id: String },
33    Cancelled { session_id: String },
34}
35
36pub struct Runtime {
37    pub storage: Arc<dyn StorageBackend>,
38    pub registry: Arc<SessionRegistry>,
39    pub log_store: Arc<LogStore>,
40    stream_bus: Arc<SessionStreamBus>,
41    signal_bus: tokio::sync::broadcast::Sender<Envelope>,
42    session_lifecycle_bus: tokio::sync::broadcast::Sender<SessionLifecycleEvent>,
43    mode_registry: Arc<ModeRegistry>,
44    policy_registry: Arc<PolicyRegistry>,
45    #[allow(dead_code)] // plumbed for future session-extension providers; register API TBD
46    extensions: Arc<ExtensionProviderRegistry>,
47    metrics: Arc<RuntimeMetrics>,
48    checkpoint_interval: usize,
49}
50
51impl Runtime {
52    pub fn new(
53        storage: Arc<dyn StorageBackend>,
54        registry: Arc<SessionRegistry>,
55        log_store: Arc<LogStore>,
56    ) -> Self {
57        Self::with_mode_registry(
58            storage,
59            registry,
60            log_store,
61            Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
62                macp_policy::DefaultPolicyEvaluator,
63            ))),
64        )
65    }
66
67    pub fn with_mode_registry(
68        storage: Arc<dyn StorageBackend>,
69        registry: Arc<SessionRegistry>,
70        log_store: Arc<LogStore>,
71        mode_registry: Arc<ModeRegistry>,
72    ) -> Self {
73        Self::with_registries(
74            storage,
75            registry,
76            log_store,
77            mode_registry,
78            Arc::new(PolicyRegistry::new()),
79        )
80    }
81
82    pub fn with_registries(
83        storage: Arc<dyn StorageBackend>,
84        registry: Arc<SessionRegistry>,
85        log_store: Arc<LogStore>,
86        mode_registry: Arc<ModeRegistry>,
87        policy_registry: Arc<PolicyRegistry>,
88    ) -> Self {
89        let checkpoint_interval = std::env::var("MACP_CHECKPOINT_INTERVAL")
90            .ok()
91            .and_then(|v| v.parse().ok())
92            .unwrap_or(0); // 0 = disabled by default
93        let (signal_tx, _) = tokio::sync::broadcast::channel(256);
94        let (session_lifecycle_tx, _) = tokio::sync::broadcast::channel(64);
95        Self {
96            storage,
97            registry,
98            log_store,
99            stream_bus: Arc::new(SessionStreamBus::default()),
100            signal_bus: signal_tx,
101            session_lifecycle_bus: session_lifecycle_tx,
102            mode_registry,
103            policy_registry,
104            extensions: Arc::new(ExtensionProviderRegistry::new()),
105            metrics: Arc::new(RuntimeMetrics::new()),
106            checkpoint_interval,
107        }
108    }
109
110    /// Returns all mode names the runtime can handle (standards-track + extensions).
111    /// Used by Initialize and GetManifest to advertise full capability.
112    pub fn registered_mode_names(&self) -> Vec<String> {
113        self.mode_registry.all_mode_names()
114    }
115
116    /// Returns only standards-track mode descriptors for ListModes.
117    pub fn standard_mode_descriptors(&self) -> Vec<ModeDescriptor> {
118        self.mode_registry.standard_mode_descriptors()
119    }
120
121    /// Returns only extension mode descriptors for ListExtModes.
122    pub fn extension_mode_descriptors(&self) -> Vec<ModeDescriptor> {
123        self.mode_registry.extension_mode_descriptors()
124    }
125
126    pub fn register_extension(&self, descriptor: ModeDescriptor) -> Result<(), String> {
127        self.mode_registry.register_extension(descriptor)
128    }
129
130    pub fn unregister_extension(&self, mode: &str) -> Result<(), String> {
131        self.mode_registry.unregister_extension(mode)
132    }
133
134    pub fn promote_mode(&self, mode: &str, new_name: Option<&str>) -> Result<String, String> {
135        self.mode_registry.promote_mode(mode, new_name)
136    }
137
138    pub fn subscribe_mode_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
139        self.mode_registry.subscribe_changes()
140    }
141
142    pub fn mode_registry(&self) -> &Arc<ModeRegistry> {
143        &self.mode_registry
144    }
145
146    // ── Policy registry delegation ──────────────────────────────────
147
148    pub fn register_policy(&self, definition: PolicyDefinition) -> Result<(), String> {
149        self.policy_registry.register(definition)
150    }
151
152    pub fn unregister_policy(&self, policy_id: &str) -> Result<(), String> {
153        self.policy_registry.unregister(policy_id)
154    }
155
156    pub fn get_policy(&self, policy_id: &str) -> Option<PolicyDefinition> {
157        self.policy_registry.get(policy_id)
158    }
159
160    pub fn list_policies(&self, mode_filter: Option<&str>) -> Vec<PolicyDefinition> {
161        self.policy_registry.list(mode_filter)
162    }
163
164    pub fn subscribe_policy_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
165        self.policy_registry.subscribe_changes()
166    }
167
168    pub fn policy_registry(&self) -> &Arc<PolicyRegistry> {
169        &self.policy_registry
170    }
171
172    pub fn metrics(&self) -> &Arc<RuntimeMetrics> {
173        &self.metrics
174    }
175
176    pub fn subscribe_session_stream(
177        &self,
178        session_id: &str,
179    ) -> tokio::sync::broadcast::Receiver<Envelope> {
180        self.stream_bus.subscribe(session_id)
181    }
182
183    pub fn subscribe_signals(&self) -> tokio::sync::broadcast::Receiver<Envelope> {
184        self.signal_bus.subscribe()
185    }
186
187    pub fn subscribe_session_lifecycle(
188        &self,
189    ) -> tokio::sync::broadcast::Receiver<SessionLifecycleEvent> {
190        self.session_lifecycle_bus.subscribe()
191    }
192
193    /// RFC-MACP-0006 §3.2: Replay accepted envelopes from the session log for
194    /// passive subscribe, strictly after `after_sequence` (1-based accepted
195    /// ordinal, exclusive; 0 = from the start). `Err(base)` when the
196    /// requested range was discarded by log compaction — the caller must
197    /// surface an explicit error, not silently skip missing history.
198    pub async fn get_session_envelopes_after(
199        &self,
200        session_id: &str,
201        after_sequence: u64,
202    ) -> Result<Vec<Envelope>, u64> {
203        Ok(self
204            .log_store
205            .get_incoming_after(session_id, after_sequence)
206            .await?
207            .into_iter()
208            .map(|(_idx, entry)| Envelope {
209                macp_version: if entry.macp_version.is_empty() {
210                    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.
1197                self.save_session_to_storage(session).await;
1198                self.metrics.record_session_expired(&session.mode);
1199                let _ = self
1200                    .session_lifecycle_bus
1201                    .send(SessionLifecycleEvent::Expired {
1202                        session_id: session_id.to_string(),
1203                    });
1204                Err(MacpError::TtlExpired)
1205            }
1206        }
1207    }
1208
1209    /// Best-effort log compaction for terminal sessions.
1210    /// Returns `true` if compaction succeeded, `false` if skipped or failed.
1211    async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
1212        // Ordinal accounting for the sequence contract: the checkpoint must
1213        // record every accepted ordinal it discards, including any base from
1214        // a prior compaction recorded in the current log.
1215        let discarded = match self.log_store.get_log(session_id).await {
1216            Some(entries) => {
1217                let prior_base: u64 = entries
1218                    .iter()
1219                    .filter(|e| e.entry_kind == EntryKind::Checkpoint)
1220                    .map(|e| e.compacted_incoming_ordinals)
1221                    .max()
1222                    .unwrap_or(0);
1223                prior_base
1224                    + entries
1225                        .iter()
1226                        .filter(|e| e.entry_kind == EntryKind::Incoming)
1227                        .count() as u64
1228            }
1229            None => 0,
1230        };
1231        match crate::storage::compaction::compact_session_log(
1232            &*self.storage,
1233            session_id,
1234            session,
1235            discarded,
1236        )
1237        .await
1238        {
1239            Ok(checkpoint) => {
1240                // Keep the in-memory log in step with storage — previously
1241                // only disk was rewritten, so memory and disk diverged and
1242                // post-restart passive-subscribe history vanished silently.
1243                self.log_store
1244                    .replace_session_log(session_id, vec![checkpoint])
1245                    .await;
1246                true
1247            }
1248            Err(e) => {
1249                tracing::debug!(
1250                    session_id,
1251                    error = %e,
1252                    "log compaction skipped (backend may not support it)"
1253                );
1254                false
1255            }
1256        }
1257    }
1258
1259    /// Force a checkpoint entry regardless of interval settings.
1260    /// Used as a fallback when compaction fails on terminal sessions.
1261    async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1262        let persisted = crate::registry::PersistedSession::from(session);
1263        let raw_payload = match serde_json::to_vec(&persisted) {
1264            Ok(bytes) => bytes,
1265            Err(e) => {
1266                tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1267                return;
1268            }
1269        };
1270        let now = Utc::now().timestamp_millis();
1271        let checkpoint = LogEntry {
1272            message_id: String::new(),
1273            received_at_ms: now,
1274            sender: "_runtime".into(),
1275            message_type: "Checkpoint".into(),
1276            raw_payload,
1277            entry_kind: EntryKind::Checkpoint,
1278            session_id: session_id.into(),
1279            mode: session.mode.clone(),
1280            macp_version: String::new(),
1281            timestamp_unix_ms: now,
1282            bound_mode_version: None,
1283            semantics_rev: 0,
1284            bound_max_suspend_ms: None,
1285            compacted_incoming_ordinals: 0,
1286        };
1287        if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1288            tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1289            return;
1290        }
1291        self.log_store.append(session_id, checkpoint).await;
1292        tracing::debug!(
1293            session_id,
1294            "forced checkpoint inserted for terminal session"
1295        );
1296    }
1297
1298    /// Insert a checkpoint entry if the log has reached the configured interval.
1299    ///
1300    /// Called after **every** append that can cross a boundary: the tail of
1301    /// `process_message`, and [`Runtime::synthesize_due_accept`] for the
1302    /// synthetic entry it writes. The second call site is what keeps the eager
1303    /// sweep and the lazy trigger placing checkpoints identically — see the
1304    /// note at that call. The decision is a pure function of the current
1305    /// `log_len`, so asking twice within one `process_message` (at two
1306    /// different lengths) is not a repeat.
1307    async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1308        if self.checkpoint_interval == 0 {
1309            return;
1310        }
1311        let log_len = self
1312            .log_store
1313            .get_log(session_id)
1314            .await
1315            .map(|l| l.len())
1316            .unwrap_or(0);
1317        // Only checkpoint at interval boundaries, and not on the first entry
1318        if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1319            return;
1320        }
1321        self.force_insert_checkpoint(session_id, session).await;
1322        tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1323    }
1324
1325    /// Expire all sessions that have exceeded their TTL.
1326    /// Called by the background cleanup task to proactively transition
1327    /// stale sessions without waiting for the next incoming message.
1328    pub async fn cleanup_expired_sessions(&self) {
1329        let now = Utc::now().timestamp_millis();
1330        // Snapshot the shared handles under a brief map read; never hold the
1331        // map lock across per-session locks or storage I/O. Each session is
1332        // re-checked under its own mutex (it may have been touched since the
1333        // snapshot).
1334        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1335            let guard = self.registry.sessions.read().await;
1336            guard
1337                .iter()
1338                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1339                .collect()
1340        };
1341
1342        let mut expired_count = 0usize;
1343        for (session_id, shared) in candidates {
1344            let mut session = shared.lock().await;
1345            if session.state != SessionState::Open || now <= session.ttl_expiry {
1346                continue;
1347            }
1348            let entry =
1349                Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode, now);
1350            if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1351                tracing::warn!(
1352                    session_id,
1353                    error = %e,
1354                    "failed to write TTL expiry during cleanup"
1355                );
1356                continue;
1357            }
1358            self.log_store.append(&session_id, entry).await;
1359            session.state = SessionState::Expired;
1360            self.metrics.record_session_expired(&session.mode);
1361            self.save_session_to_storage(&session).await;
1362            if !self.maybe_compact_log(&session_id, &session).await {
1363                self.force_insert_checkpoint(&session_id, &session).await;
1364            }
1365            expired_count += 1;
1366            tracing::info!(session_id = %session_id, "session expired via background cleanup");
1367            let _ = self
1368                .session_lifecycle_bus
1369                .send(SessionLifecycleEvent::Expired {
1370                    session_id: session_id.clone(),
1371                });
1372        }
1373
1374        if expired_count > 0 {
1375            tracing::info!(count = expired_count, "background cleanup expired sessions");
1376        }
1377    }
1378
1379    /// Emit every synthetic envelope that has become due, across all open
1380    /// sessions — RFC-MACP-0010 §5.1(2)'s **eager** observation of the
1381    /// implicit-accept deadline.
1382    ///
1383    /// Lazy observation (`process_message` -> `synthesize_due_accept`) is the
1384    /// MUST and ships on its own: it guarantees no message is ever evaluated
1385    /// against a stale offer. This is the SHOULD on top of it — without it a
1386    /// session where nobody speaks again keeps an accepted offer out of history
1387    /// indefinitely, and `GetSession` / `StreamSession` show an offer that the
1388    /// protocol says was accepted at `D`. Called from the background
1389    /// maintenance loop in `src/main.rs`, so `MACP_CLEANUP_INTERVAL_SECS` is
1390    /// the latency bound on the observation (never on the recorded timestamp,
1391    /// which is `D` whichever path emits — see below).
1392    ///
1393    /// # Not a variant of `cleanup_expired_sessions`
1394    ///
1395    /// Different predicate (a mode-computed deadline inside `mode_state`, not
1396    /// `ttl_expiry`) and a different product: an `EntryKind::Incoming` entry
1397    /// that consumes an accepted ordinal and is published to `StreamSession`,
1398    /// where `TtlExpired` is `Internal` and is published to neither. It is
1399    /// deliberately ordered *after* `cleanup_expired_sessions` in that loop so
1400    /// a TTL-expired session is already non-`Open` when the sweep reaches it —
1401    /// which is exactly the precedence the lazy path gives (`Precheck::Expired`
1402    /// returns before `synthesize_due_accept` is ever called), so the two paths
1403    /// cannot disagree about a session whose TTL and implicit-accept deadline
1404    /// both passed unobserved.
1405    ///
1406    /// # Locking
1407    ///
1408    /// The registry map lock is held only for the snapshot of `(id, Arc)` pairs
1409    /// and is released before any session mutex is taken or any I/O happens —
1410    /// the lock-ordering contract on `SessionRegistry`. The snapshot fixes the
1411    /// *set*, not the state.
1412    ///
1413    /// # Why no snapshot-to-append race can leave an orphan entry
1414    ///
1415    /// A session can resolve, cancel or expire between the snapshot and the
1416    /// moment this loop reaches it, and the `Arc` keeps it alive (and writable)
1417    /// regardless. Nothing here reads state at snapshot time: every decision is
1418    /// made *under the session mutex*, which is the same mutex every writer —
1419    /// `process_message`, `cancel_session`, `suspend_session`, `resume_session`,
1420    /// `cleanup_expired_sessions` — holds across its own validate-append-commit.
1421    /// So the `state != Open` re-check inside `synthesize_due_accept` observes
1422    /// the session as the last writer left it, and a session that terminated
1423    /// after the snapshot is skipped. Two further guards make it belt and
1424    /// braces: the mode returns `None` once the offer's disposition is no
1425    /// longer `Offered` (so an explicit `HandoffAccept` that won the race
1426    /// disarms the synthesis), and the deterministic `message_id` is already in
1427    /// `seen_message_ids` after any emission. Eviction cannot orphan one
1428    /// either: `evict_stale_sessions` and `gc_disk_sessions` only ever drop
1429    /// *terminal* sessions, which fail the `Open` check.
1430    ///
1431    /// # Cost
1432    ///
1433    /// One `mode_state` decode per open session whose mode implements the hook
1434    /// and whose policy binds a timeout, per tick; everything else short-circuits
1435    /// before decoding (`Mode::due_synthetic_envelope`'s own cost note).
1436    ///
1437    /// Returns the number of synthetic entries appended.
1438    pub async fn sweep_due_synthetic_accepts(&self) -> usize {
1439        let now = Utc::now().timestamp_millis();
1440        // Snapshot the shared handles under a brief map read; never hold the
1441        // map lock across per-session locks or storage I/O. Mirrors
1442        // `cleanup_expired_sessions`.
1443        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1444            let guard = self.registry.sessions.read().await;
1445            guard
1446                .iter()
1447                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1448                .collect()
1449        };
1450
1451        let mut emitted = 0usize;
1452        for (session_id, shared) in candidates {
1453            let mut session = shared.lock().await;
1454            // `Suspended` (and every other non-`Open` state) is skipped here
1455            // *and* inside `synthesize_due_accept`. Measured: removing this
1456            // check alone changes no behaviour — the seam's own filter still
1457            // declines — so treat it as the documented precondition of the
1458            // call below rather than as the enforcement. The enforcement
1459            // matters: both computations behind the synthetic entry ignore an
1460            // in-flight pause, so asking about a paused session over-counts
1461            // elapsed time and bakes a wrong `D` into permanent history, and
1462            // this is the only caller that can ever be handed one (no message
1463            // path reaches a non-`Open` session at all).
1464            if session.state != SessionState::Open {
1465                continue;
1466            }
1467            match self
1468                .synthesize_due_accept(&session_id, &mut session, now)
1469                .await
1470            {
1471                Ok(true) => emitted += 1,
1472                Ok(false) => {}
1473                // No triggering message to reject: a failed append (or a mode
1474                // that refused its own synthetic) is logged and the offer stays
1475                // outstanding, to be retried on the next tick or settled by the
1476                // lazy path. The same posture `cleanup_expired_sessions` takes
1477                // on a failed `TtlExpired` append.
1478                Err(e) => {
1479                    tracing::warn!(
1480                        session_id = %session_id,
1481                        error = %e,
1482                        "eager sweep could not emit a due synthetic envelope"
1483                    );
1484                }
1485            }
1486        }
1487
1488        if emitted > 0 {
1489            tracing::info!(
1490                count = emitted,
1491                "eager sweep appended due synthetic envelopes"
1492            );
1493        }
1494        emitted
1495    }
1496
1497    /// Delete terminal sessions' durable data older than `retention_secs`
1498    /// (opt-in via `MACP_SESSION_DISK_RETENTION_SECS`). Before this existed,
1499    /// `storage.delete_session` had no callers at all: disk grew without
1500    /// bound and every restart reloaded every session ever completed.
1501    /// Enumerates STORAGE (not memory — eviction may already have dropped the
1502    /// registry entry), deletes the session's snapshot+log, and clears any
1503    /// in-memory remnants. Returns the number of sessions deleted.
1504    pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1505        let now = Utc::now().timestamp_millis();
1506        let cutoff = now - (retention_secs as i64 * 1000);
1507        let ids = match self.storage.list_session_ids().await {
1508            Ok(ids) => ids,
1509            Err(e) => {
1510                tracing::warn!(error = %e, "disk GC: cannot list sessions");
1511                return 0;
1512            }
1513        };
1514        let mut removed = 0usize;
1515        for id in ids {
1516            // Prefer the in-memory state when present (cheap + current);
1517            // fall back to the stored snapshot for evicted sessions.
1518            let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1519                let s = shared.lock().await;
1520                s.state.is_terminal() && s.started_at_unix_ms < cutoff
1521            } else {
1522                match self.storage.load_session(&id).await {
1523                    Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1524                    // No snapshot (or unreadable): leave it for operator
1525                    // inspection rather than guessing.
1526                    _ => false,
1527                }
1528            };
1529            if !eligible {
1530                continue;
1531            }
1532            match self.storage.delete_session(&id).await {
1533                Ok(()) => {
1534                    {
1535                        let mut guard = self.registry.sessions.write().await;
1536                        guard.remove(&id);
1537                    }
1538                    self.log_store.remove_session_log(&id).await;
1539                    let _ = self.stream_bus.remove_if_unused(&id);
1540                    removed += 1;
1541                }
1542                Err(e) => {
1543                    tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1544                }
1545            }
1546        }
1547        if removed > 0 {
1548            tracing::info!(count = removed, "disk GC removed terminal sessions");
1549        }
1550        removed
1551    }
1552
1553    /// Evict resolved/expired sessions older than `retention_secs` from
1554    /// memory: the registry entry, the in-memory log cache, AND the stream
1555    /// broadcast channel (all three previously grew for the process lifetime;
1556    /// the log cache and stream bus were never evicted at all). Sessions
1557    /// remain queryable from durable storage after eviction.
1558    pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1559        let now = Utc::now().timestamp_millis();
1560        let cutoff = now - (retention_secs as i64 * 1000);
1561
1562        let candidates: Vec<(String, crate::registry::SharedSession)> = {
1563            let guard = self.registry.sessions.read().await;
1564            guard
1565                .iter()
1566                .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1567                .collect()
1568        };
1569        let mut evict_ids = Vec::new();
1570        for (id, shared) in candidates {
1571            let session = shared.lock().await;
1572            if matches!(
1573                session.state,
1574                SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1575            ) && session.started_at_unix_ms < cutoff
1576            {
1577                evict_ids.push(id);
1578            }
1579        }
1580
1581        if evict_ids.is_empty() {
1582            return;
1583        }
1584        {
1585            let mut guard = self.registry.sessions.write().await;
1586            for id in &evict_ids {
1587                guard.remove(id);
1588            }
1589        }
1590        for id in &evict_ids {
1591            self.log_store.remove_session_log(id).await;
1592            // Left in place if a subscriber is still attached; retried on the
1593            // next sweep once receivers drop.
1594            let _ = self.stream_bus.remove_if_unused(id);
1595        }
1596        tracing::info!(
1597            count = evict_ids.len(),
1598            "evicted stale sessions from memory (registry + log cache + stream bus)"
1599        );
1600    }
1601}
1602
1603#[cfg(test)]
1604mod tests {
1605    use super::*;
1606    use crate::decision_pb::ProposalPayload;
1607    use crate::pb::{CommitmentPayload, SessionStartPayload};
1608    use prost::Message;
1609
1610    fn new_sid() -> String {
1611        uuid::Uuid::new_v4().as_hyphenated().to_string()
1612    }
1613
1614    fn make_runtime() -> Runtime {
1615        let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1616        let registry = Arc::new(SessionRegistry::new());
1617        let log_store = Arc::new(LogStore::new());
1618        Runtime::new(storage, registry, log_store)
1619    }
1620
1621    fn session_start(participants: Vec<String>) -> Vec<u8> {
1622        SessionStartPayload {
1623            intent: "intent".into(),
1624            participants,
1625            mode_version: "1.0.0".into(),
1626            configuration_version: "cfg-1".into(),
1627            policy_version: String::new(),
1628            ttl_ms: 1_000,
1629            context_id: String::new(),
1630            extensions: std::collections::HashMap::new(),
1631            roots: vec![],
1632            max_suspend_ms: 0,
1633        }
1634        .encode_to_vec()
1635    }
1636
1637    fn env(
1638        mode: &str,
1639        message_type: &str,
1640        message_id: &str,
1641        session_id: &str,
1642        sender: &str,
1643        payload: Vec<u8>,
1644    ) -> Envelope {
1645        Envelope {
1646            macp_version: "1.0".into(),
1647            mode: mode.into(),
1648            message_type: message_type.into(),
1649            message_id: message_id.into(),
1650            session_id: session_id.into(),
1651            sender: sender.into(),
1652            timestamp_unix_ms: Utc::now().timestamp_millis(),
1653            payload,
1654        }
1655    }
1656
1657    #[tokio::test]
1658    async fn standard_session_start_is_strict() {
1659        let rt = make_runtime();
1660        let sid = new_sid();
1661        let bad = SessionStartPayload {
1662            ttl_ms: 0,
1663            ..Default::default()
1664        }
1665        .encode_to_vec();
1666        let err = rt
1667            .process(
1668                &env(
1669                    "macp.mode.decision.v1",
1670                    "SessionStart",
1671                    "m1",
1672                    &sid,
1673                    "agent://orchestrator",
1674                    bad,
1675                ),
1676                None,
1677            )
1678            .await
1679            .unwrap_err();
1680        assert!(matches!(
1681            err,
1682            MacpError::InvalidPayload | MacpError::InvalidTtl
1683        ));
1684    }
1685
1686    /// A **promoted** extension mode keeps the full canonical `SessionStart`
1687    /// contract, including the roster requirement.
1688    ///
1689    /// This pins the *call site's* choice of strictness source, which no other
1690    /// test covers. `ModeRegistry::requires_strict_session_start` reads a
1691    /// per-entry `strict_session_start` flag that `promote_mode` sets to `true`;
1692    /// `macp_core::session::requires_strict_session_start` reads a static name
1693    /// list that has never heard of a promoted mode's name. The two disagree
1694    /// exactly here, so routing this call through
1695    /// `validate_strict_session_start_payload` — which consults the static list
1696    /// — would silently skip canonical validation for every promoted mode while
1697    /// leaving the whole suite green. Measured: it does.
1698    #[tokio::test]
1699    async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1700        let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1701            macp_policy::DefaultPolicyEvaluator,
1702        )));
1703        mode_registry
1704            .register_extension(crate::pb::ModeDescriptor {
1705                mode: "ext.promoted.v1".into(),
1706                mode_version: "1.0.0".into(),
1707                title: "Promoted".into(),
1708                description: "promotion target".into(),
1709                determinism_class: "semantic-deterministic".into(),
1710                participant_model: "declared".into(),
1711                message_types: vec!["SessionStart".into(), "Commitment".into()],
1712                terminal_message_types: vec!["Commitment".into()],
1713                ..Default::default()
1714            })
1715            .expect("register extension");
1716        assert_eq!(
1717            mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1718            "ext.promoted.v1"
1719        );
1720        assert!(
1721            mode_registry.requires_strict_session_start("ext.promoted.v1"),
1722            "promotion must mark the entry strict"
1723        );
1724        assert!(
1725            !crate::session::requires_strict_session_start("ext.promoted.v1"),
1726            "the core's static list must NOT know this name — that disagreement is the point"
1727        );
1728
1729        let rt = Runtime::with_mode_registry(
1730            Arc::new(crate::storage::MemoryBackend),
1731            Arc::new(SessionRegistry::new()),
1732            Arc::new(LogStore::new()),
1733            mode_registry,
1734        );
1735
1736        // An empty roster: refused, because the carve-out names Decision only.
1737        let err = rt
1738            .process(
1739                &env(
1740                    "ext.promoted.v1",
1741                    "SessionStart",
1742                    "m1",
1743                    &new_sid(),
1744                    "agent://orchestrator",
1745                    session_start(vec![]),
1746                ),
1747                None,
1748            )
1749            .await
1750            .unwrap_err();
1751        assert_eq!(err.to_string(), "InvalidPayload");
1752
1753        // And the rest of the canonical contract too, so the assertion above
1754        // cannot be satisfied by a runtime that only kept the roster rule.
1755        let no_versions = SessionStartPayload {
1756            participants: vec!["agent://fraud".into()],
1757            ttl_ms: 1_000,
1758            ..Default::default()
1759        }
1760        .encode_to_vec();
1761        let err = rt
1762            .process(
1763                &env(
1764                    "ext.promoted.v1",
1765                    "SessionStart",
1766                    "m2",
1767                    &new_sid(),
1768                    "agent://orchestrator",
1769                    no_versions,
1770                ),
1771                None,
1772            )
1773            .await
1774            .unwrap_err();
1775        assert_eq!(err.to_string(), "InvalidPayload");
1776
1777        // Positive control: a complete payload is accepted, so the two refusals
1778        // above are the roster and version rules and not a broken mode.
1779        rt.process(
1780            &env(
1781                "ext.promoted.v1",
1782                "SessionStart",
1783                "m3",
1784                &new_sid(),
1785                "agent://orchestrator",
1786                session_start(vec!["agent://fraud".into()]),
1787            ),
1788            None,
1789        )
1790        .await
1791        .expect("a complete SessionStart must still be accepted for a promoted mode");
1792    }
1793
1794    #[tokio::test]
1795    async fn empty_mode_is_rejected() {
1796        let rt = make_runtime();
1797        let sid = new_sid();
1798        let err = rt
1799            .process(
1800                &env(
1801                    "",
1802                    "SessionStart",
1803                    "m1",
1804                    &sid,
1805                    "agent://orchestrator",
1806                    session_start(vec!["agent://fraud".into()]),
1807                ),
1808                None,
1809            )
1810            .await
1811            .unwrap_err();
1812        assert_eq!(err.to_string(), "InvalidEnvelope");
1813    }
1814
1815    #[tokio::test]
1816    async fn rejected_messages_do_not_enter_dedup_state() {
1817        let rt = make_runtime();
1818        let sid = new_sid();
1819        rt.process(
1820            &env(
1821                "macp.mode.decision.v1",
1822                "SessionStart",
1823                "m1",
1824                &sid,
1825                "agent://orchestrator",
1826                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1827            ),
1828            None,
1829        )
1830        .await
1831        .unwrap();
1832
1833        let bad = rt
1834            .process(
1835                &env(
1836                    "macp.mode.decision.v1",
1837                    "Proposal",
1838                    "m2",
1839                    &sid,
1840                    "agent://fraud",
1841                    b"not-protobuf".to_vec(),
1842                ),
1843                None,
1844            )
1845            .await
1846            .unwrap_err();
1847        assert_eq!(bad.to_string(), "InvalidPayload");
1848
1849        let good = ProposalPayload {
1850            proposal_id: "p1".into(),
1851            option: "step-up".into(),
1852            rationale: "risk".into(),
1853            supporting_data: vec![],
1854        }
1855        .encode_to_vec();
1856        let result = rt
1857            .process(
1858                &env(
1859                    "macp.mode.decision.v1",
1860                    "Proposal",
1861                    "m2",
1862                    &sid,
1863                    "agent://orchestrator",
1864                    good,
1865                ),
1866                None,
1867            )
1868            .await
1869            .unwrap();
1870        assert!(!result.duplicate);
1871    }
1872
1873    #[tokio::test]
1874    async fn get_session_transitions_expired_sessions() {
1875        let rt = make_runtime();
1876        let sid = new_sid();
1877        let payload = SessionStartPayload {
1878            intent: "intent".into(),
1879            participants: vec!["agent://fraud".into()],
1880            mode_version: "1.0.0".into(),
1881            configuration_version: "cfg-1".into(),
1882            policy_version: String::new(),
1883            ttl_ms: 1,
1884            context_id: String::new(),
1885            extensions: std::collections::HashMap::new(),
1886            roots: vec![],
1887            max_suspend_ms: 0,
1888        }
1889        .encode_to_vec();
1890        rt.process(
1891            &env(
1892                "macp.mode.decision.v1",
1893                "SessionStart",
1894                "m1",
1895                &sid,
1896                "agent://orchestrator",
1897                payload,
1898            ),
1899            None,
1900        )
1901        .await
1902        .unwrap();
1903        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1904        let session = rt.get_session_checked(&sid).await.unwrap();
1905        assert_eq!(session.state, SessionState::Expired);
1906    }
1907
1908    #[tokio::test]
1909    async fn multi_round_requires_standard_session_start() {
1910        let rt = make_runtime();
1911        let sid = new_sid();
1912        // multi-round is now standards-track: empty mode_version should fail
1913        let payload = SessionStartPayload {
1914            participants: vec!["creator".into(), "other".into()],
1915            ..Default::default()
1916        }
1917        .encode_to_vec();
1918        let err = rt
1919            .process(
1920                &env(
1921                    "ext.multi_round.v1",
1922                    "SessionStart",
1923                    "m1",
1924                    &sid,
1925                    "creator",
1926                    payload,
1927                ),
1928                None,
1929            )
1930            .await
1931            .unwrap_err();
1932        assert!(matches!(
1933            err,
1934            MacpError::InvalidPayload | MacpError::InvalidTtl
1935        ));
1936    }
1937
1938    #[tokio::test]
1939    async fn multi_round_valid_session_start() {
1940        let rt = make_runtime();
1941        let sid = new_sid();
1942        let payload = session_start(vec!["alice".into(), "bob".into()]);
1943        rt.process(
1944            &env(
1945                "ext.multi_round.v1",
1946                "SessionStart",
1947                "m1",
1948                &sid,
1949                "coordinator",
1950                payload,
1951            ),
1952            None,
1953        )
1954        .await
1955        .unwrap();
1956        let session = rt.get_session_checked(&sid).await.unwrap();
1957        assert_eq!(session.mode, "ext.multi_round.v1");
1958        assert_eq!(session.participants, vec!["alice", "bob"]);
1959    }
1960
1961    #[tokio::test]
1962    async fn duplicate_session_start_message_id_returns_duplicate() {
1963        let rt = make_runtime();
1964        let sid = new_sid();
1965        let payload = session_start(vec!["agent://fraud".into()]);
1966        rt.process(
1967            &env(
1968                "macp.mode.decision.v1",
1969                "SessionStart",
1970                "m1",
1971                &sid,
1972                "agent://orchestrator",
1973                payload.clone(),
1974            ),
1975            None,
1976        )
1977        .await
1978        .unwrap();
1979
1980        let result = rt
1981            .process(
1982                &env(
1983                    "macp.mode.decision.v1",
1984                    "SessionStart",
1985                    "m1",
1986                    &sid,
1987                    "agent://orchestrator",
1988                    payload,
1989                ),
1990                None,
1991            )
1992            .await
1993            .unwrap();
1994        assert!(result.duplicate);
1995    }
1996
1997    #[tokio::test]
1998    async fn non_start_mode_mismatch_rejected() {
1999        let rt = make_runtime();
2000        let sid = new_sid();
2001        rt.process(
2002            &env(
2003                "macp.mode.decision.v1",
2004                "SessionStart",
2005                "m1",
2006                &sid,
2007                "agent://orchestrator",
2008                session_start(vec!["agent://fraud".into()]),
2009            ),
2010            None,
2011        )
2012        .await
2013        .unwrap();
2014
2015        let proposal = ProposalPayload {
2016            proposal_id: "p1".into(),
2017            option: "step-up".into(),
2018            rationale: "risk".into(),
2019            supporting_data: vec![],
2020        }
2021        .encode_to_vec();
2022        let err = rt
2023            .process(
2024                &env(
2025                    "macp.mode.task.v1",
2026                    "Proposal",
2027                    "m2",
2028                    &sid,
2029                    "agent://orchestrator",
2030                    proposal,
2031                ),
2032                None,
2033            )
2034            .await
2035            .unwrap_err();
2036        assert_eq!(err.to_string(), "InvalidEnvelope");
2037    }
2038
2039    #[tokio::test]
2040    async fn cancel_idempotent_on_already_expired() {
2041        let rt = make_runtime();
2042        let sid = new_sid();
2043        let payload = SessionStartPayload {
2044            intent: "intent".into(),
2045            participants: vec!["agent://fraud".into()],
2046            mode_version: "1.0.0".into(),
2047            configuration_version: "cfg-1".into(),
2048            policy_version: String::new(),
2049            ttl_ms: 1,
2050            context_id: String::new(),
2051            extensions: std::collections::HashMap::new(),
2052            roots: vec![],
2053            max_suspend_ms: 0,
2054        }
2055        .encode_to_vec();
2056        rt.process(
2057            &env(
2058                "macp.mode.decision.v1",
2059                "SessionStart",
2060                "m1",
2061                &sid,
2062                "agent://orchestrator",
2063                payload,
2064            ),
2065            None,
2066        )
2067        .await
2068        .unwrap();
2069        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2070        let result = rt
2071            .cancel_session(&sid, "cleanup", "agent://orchestrator")
2072            .await
2073            .unwrap();
2074        assert_eq!(result.session_state, SessionState::Expired);
2075    }
2076
2077    #[tokio::test]
2078    async fn accepted_envelopes_are_published_in_order() {
2079        let rt = make_runtime();
2080        let sid = new_sid();
2081        let mut events = rt.subscribe_session_stream(&sid);
2082
2083        let start = env(
2084            "macp.mode.decision.v1",
2085            "SessionStart",
2086            "m1",
2087            &sid,
2088            "agent://orchestrator",
2089            session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2090        );
2091        rt.process(&start, None).await.unwrap();
2092        let first = events.recv().await.unwrap();
2093        assert_eq!(first.message_id, "m1");
2094        assert_eq!(first.message_type, "SessionStart");
2095
2096        let proposal = ProposalPayload {
2097            proposal_id: "p1".into(),
2098            option: "step-up".into(),
2099            rationale: "risk".into(),
2100            supporting_data: vec![],
2101        }
2102        .encode_to_vec();
2103        let proposal_env = env(
2104            "macp.mode.decision.v1",
2105            "Proposal",
2106            "m2",
2107            &sid,
2108            "agent://orchestrator",
2109            proposal,
2110        );
2111        rt.process(&proposal_env, None).await.unwrap();
2112        let second = events.recv().await.unwrap();
2113        assert_eq!(second.message_id, "m2");
2114        assert_eq!(second.message_type, "Proposal");
2115    }
2116
2117    #[tokio::test]
2118    async fn commitment_versions_are_carried_into_resolution() {
2119        let rt = make_runtime();
2120        let sid = new_sid();
2121        rt.process(
2122            &env(
2123                "macp.mode.proposal.v1",
2124                "SessionStart",
2125                "m1",
2126                &sid,
2127                "agent://buyer",
2128                session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
2129            ),
2130            None,
2131        )
2132        .await
2133        .unwrap();
2134
2135        let proposal = crate::proposal_pb::ProposalPayload {
2136            proposal_id: "p1".into(),
2137            title: "offer".into(),
2138            summary: "summary".into(),
2139            details: vec![],
2140            tags: vec![],
2141        }
2142        .encode_to_vec();
2143        rt.process(
2144            &env(
2145                "macp.mode.proposal.v1",
2146                "Proposal",
2147                "m2",
2148                &sid,
2149                "agent://seller",
2150                proposal,
2151            ),
2152            None,
2153        )
2154        .await
2155        .unwrap();
2156        let accept = crate::proposal_pb::AcceptPayload {
2157            proposal_id: "p1".into(),
2158            reason: String::new(),
2159        }
2160        .encode_to_vec();
2161        rt.process(
2162            &env(
2163                "macp.mode.proposal.v1",
2164                "Accept",
2165                "m3",
2166                &sid,
2167                "agent://seller",
2168                accept.clone(),
2169            ),
2170            None,
2171        )
2172        .await
2173        .unwrap();
2174        rt.process(
2175            &env(
2176                "macp.mode.proposal.v1",
2177                "Accept",
2178                "m4",
2179                &sid,
2180                "agent://buyer",
2181                accept,
2182            ),
2183            None,
2184        )
2185        .await
2186        .unwrap();
2187        let commitment = CommitmentPayload {
2188            commitment_id: "c1".into(),
2189            action: "proposal.accepted".into(),
2190            authority_scope: "commercial".into(),
2191            reason: "bound".into(),
2192            mode_version: "1.0.0".into(),
2193            policy_version: "policy.default".into(),
2194            configuration_version: "cfg-1".into(),
2195            outcome_positive: true,
2196            supersedes: None,
2197        }
2198        .encode_to_vec();
2199        let result = rt
2200            .process(
2201                &env(
2202                    "macp.mode.proposal.v1",
2203                    "Commitment",
2204                    "m5",
2205                    &sid,
2206                    "agent://buyer",
2207                    commitment,
2208                ),
2209                None,
2210            )
2211            .await
2212            .unwrap();
2213        assert_eq!(result.session_state, SessionState::Resolved);
2214    }
2215
2216    #[tokio::test]
2217    async fn max_open_sessions_enforced_under_write_lock() {
2218        let rt = make_runtime();
2219        let sid1 = new_sid();
2220        let sid2 = new_sid();
2221        let sid3 = new_sid();
2222        rt.process(
2223            &env(
2224                "macp.mode.decision.v1",
2225                "SessionStart",
2226                "m1",
2227                &sid1,
2228                "agent://orchestrator",
2229                session_start(vec!["agent://fraud".into()]),
2230            ),
2231            Some(1),
2232        )
2233        .await
2234        .unwrap();
2235
2236        let err = rt
2237            .process(
2238                &env(
2239                    "macp.mode.decision.v1",
2240                    "SessionStart",
2241                    "m2",
2242                    &sid2,
2243                    "agent://orchestrator",
2244                    session_start(vec!["agent://fraud".into()]),
2245                ),
2246                Some(1),
2247            )
2248            .await
2249            .unwrap_err();
2250        assert!(matches!(err, MacpError::RateLimited));
2251
2252        rt.process(
2253            &env(
2254                "macp.mode.decision.v1",
2255                "SessionStart",
2256                "m3",
2257                &sid3,
2258                "agent://other",
2259                session_start(vec!["agent://fraud".into()]),
2260            ),
2261            Some(1),
2262        )
2263        .await
2264        .unwrap();
2265    }
2266
2267    #[tokio::test]
2268    async fn weak_session_id_rejected() {
2269        let rt = make_runtime();
2270        let err = rt
2271            .process(
2272                &env(
2273                    "macp.mode.decision.v1",
2274                    "SessionStart",
2275                    "m1",
2276                    "s1",
2277                    "agent://orchestrator",
2278                    session_start(vec!["agent://fraud".into()]),
2279                ),
2280                None,
2281            )
2282            .await
2283            .unwrap_err();
2284        assert_eq!(err.to_string(), "InvalidSessionId");
2285    }
2286
2287    #[tokio::test]
2288    async fn log_append_failure_rejects_session_start() {
2289        use std::io;
2290        struct FailingBackend;
2291        #[async_trait::async_trait]
2292        impl StorageBackend for FailingBackend {
2293            async fn save_session(&self, _: &Session) -> io::Result<()> {
2294                Ok(())
2295            }
2296            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2297                Ok(None)
2298            }
2299            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2300                Ok(vec![])
2301            }
2302            async fn delete_session(&self, _: &str) -> io::Result<()> {
2303                Ok(())
2304            }
2305            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2306                Ok(vec![])
2307            }
2308            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2309                Err(io::Error::other("disk full"))
2310            }
2311            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2312                Ok(vec![])
2313            }
2314            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2315                Ok(())
2316            }
2317        }
2318
2319        let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
2320        let registry = Arc::new(SessionRegistry::new());
2321        let log_store = Arc::new(LogStore::new());
2322        let rt = Runtime::new(storage, registry, log_store);
2323        let sid = new_sid();
2324
2325        let err = rt
2326            .process(
2327                &env(
2328                    "macp.mode.decision.v1",
2329                    "SessionStart",
2330                    "m1",
2331                    &sid,
2332                    "agent://orchestrator",
2333                    session_start(vec!["agent://fraud".into()]),
2334                ),
2335                None,
2336            )
2337            .await
2338            .unwrap_err();
2339        assert_eq!(err.to_string(), "StorageFailed");
2340    }
2341
2342    #[tokio::test]
2343    async fn log_append_failure_rejects_in_session_message() {
2344        use std::io;
2345        use std::sync::atomic::{AtomicUsize, Ordering};
2346
2347        struct FailOnSecondAppend {
2348            count: AtomicUsize,
2349        }
2350        #[async_trait::async_trait]
2351        impl StorageBackend for FailOnSecondAppend {
2352            async fn save_session(&self, _: &Session) -> io::Result<()> {
2353                Ok(())
2354            }
2355            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2356                Ok(None)
2357            }
2358            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2359                Ok(vec![])
2360            }
2361            async fn delete_session(&self, _: &str) -> io::Result<()> {
2362                Ok(())
2363            }
2364            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2365                Ok(vec![])
2366            }
2367            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2368                let n = self.count.fetch_add(1, Ordering::SeqCst);
2369                if n >= 1 {
2370                    Err(io::Error::other("disk full"))
2371                } else {
2372                    Ok(())
2373                }
2374            }
2375            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2376                Ok(vec![])
2377            }
2378            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2379                Ok(())
2380            }
2381        }
2382
2383        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2384            count: AtomicUsize::new(0),
2385        });
2386        let registry = Arc::new(SessionRegistry::new());
2387        let log_store = Arc::new(LogStore::new());
2388        let rt = Runtime::new(storage, registry, log_store);
2389        let sid = new_sid();
2390
2391        // SessionStart succeeds (first append)
2392        rt.process(
2393            &env(
2394                "macp.mode.decision.v1",
2395                "SessionStart",
2396                "m1",
2397                &sid,
2398                "agent://orchestrator",
2399                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2400            ),
2401            None,
2402        )
2403        .await
2404        .unwrap();
2405
2406        // Proposal fails (second append)
2407        let proposal = ProposalPayload {
2408            proposal_id: "p1".into(),
2409            option: "step-up".into(),
2410            rationale: "risk".into(),
2411            supporting_data: vec![],
2412        }
2413        .encode_to_vec();
2414        let err = rt
2415            .process(
2416                &env(
2417                    "macp.mode.decision.v1",
2418                    "Proposal",
2419                    "m2",
2420                    &sid,
2421                    "agent://orchestrator",
2422                    proposal,
2423                ),
2424                None,
2425            )
2426            .await
2427            .unwrap_err();
2428        assert_eq!(err.to_string(), "StorageFailed");
2429
2430        // Verify the message was not added to dedup state
2431        let session = rt.get_session_checked(&sid).await.unwrap();
2432        assert!(!session.seen_message_ids.contains("m2"));
2433    }
2434
2435    #[tokio::test]
2436    async fn cancel_session_fails_if_log_append_fails() {
2437        use std::io;
2438        use std::sync::atomic::{AtomicUsize, Ordering};
2439
2440        struct FailOnSecondAppend {
2441            count: AtomicUsize,
2442        }
2443        #[async_trait::async_trait]
2444        impl StorageBackend for FailOnSecondAppend {
2445            async fn save_session(&self, _: &Session) -> io::Result<()> {
2446                Ok(())
2447            }
2448            async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2449                Ok(None)
2450            }
2451            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2452                Ok(vec![])
2453            }
2454            async fn delete_session(&self, _: &str) -> io::Result<()> {
2455                Ok(())
2456            }
2457            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2458                Ok(vec![])
2459            }
2460            async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2461                let n = self.count.fetch_add(1, Ordering::SeqCst);
2462                if n >= 1 {
2463                    Err(io::Error::other("disk full"))
2464                } else {
2465                    Ok(())
2466                }
2467            }
2468            async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2469                Ok(vec![])
2470            }
2471            async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2472                Ok(())
2473            }
2474        }
2475
2476        let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2477            count: AtomicUsize::new(0),
2478        });
2479        let registry = Arc::new(SessionRegistry::new());
2480        let log_store = Arc::new(LogStore::new());
2481        let rt = Runtime::new(storage, registry, log_store);
2482        let sid = new_sid();
2483
2484        rt.process(
2485            &env(
2486                "macp.mode.decision.v1",
2487                "SessionStart",
2488                "m1",
2489                &sid,
2490                "agent://orchestrator",
2491                session_start(vec!["agent://fraud".into()]),
2492            ),
2493            None,
2494        )
2495        .await
2496        .unwrap();
2497
2498        let err = rt
2499            .cancel_session(&sid, "test cancel", "agent://orchestrator")
2500            .await
2501            .unwrap_err();
2502        assert_eq!(err.to_string(), "StorageFailed");
2503    }
2504
2505    #[tokio::test]
2506    async fn ttl_expiration_rejects_message() {
2507        let rt = make_runtime();
2508        let sid = new_sid();
2509        let payload = SessionStartPayload {
2510            intent: "intent".into(),
2511            participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2512            mode_version: "1.0.0".into(),
2513            configuration_version: "cfg-1".into(),
2514            policy_version: String::new(),
2515            ttl_ms: 1,
2516            context_id: String::new(),
2517            extensions: std::collections::HashMap::new(),
2518            roots: vec![],
2519            max_suspend_ms: 0,
2520        }
2521        .encode_to_vec();
2522        rt.process(
2523            &env(
2524                "macp.mode.decision.v1",
2525                "SessionStart",
2526                "m1",
2527                &sid,
2528                "agent://orchestrator",
2529                payload,
2530            ),
2531            None,
2532        )
2533        .await
2534        .unwrap();
2535        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2536        let proposal = ProposalPayload {
2537            proposal_id: "p1".into(),
2538            option: "step-up".into(),
2539            rationale: "risk".into(),
2540            supporting_data: vec![],
2541        }
2542        .encode_to_vec();
2543        let err = rt
2544            .process(
2545                &env(
2546                    "macp.mode.decision.v1",
2547                    "Proposal",
2548                    "m2",
2549                    &sid,
2550                    "agent://orchestrator",
2551                    proposal,
2552                ),
2553                None,
2554            )
2555            .await
2556            .unwrap_err();
2557        assert_eq!(err.to_string(), "TtlExpired");
2558    }
2559
2560    #[tokio::test]
2561    async fn cleanup_expired_sessions_marks_expired() {
2562        let rt = make_runtime();
2563        let sid = new_sid();
2564        let payload = SessionStartPayload {
2565            intent: "intent".into(),
2566            participants: vec!["agent://fraud".into()],
2567            mode_version: "1.0.0".into(),
2568            configuration_version: "cfg-1".into(),
2569            policy_version: String::new(),
2570            ttl_ms: 1,
2571            context_id: String::new(),
2572            extensions: std::collections::HashMap::new(),
2573            roots: vec![],
2574            max_suspend_ms: 0,
2575        }
2576        .encode_to_vec();
2577        rt.process(
2578            &env(
2579                "macp.mode.decision.v1",
2580                "SessionStart",
2581                "m1",
2582                &sid,
2583                "agent://orchestrator",
2584                payload,
2585            ),
2586            None,
2587        )
2588        .await
2589        .unwrap();
2590        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2591        rt.cleanup_expired_sessions().await;
2592        let session = rt.get_session_checked(&sid).await.unwrap();
2593        assert_eq!(session.state, SessionState::Expired);
2594    }
2595
2596    #[tokio::test]
2597    async fn evict_stale_sessions_removes_resolved() {
2598        let rt = make_runtime();
2599        let sid = new_sid();
2600        // Start a decision session
2601        rt.process(
2602            &env(
2603                "macp.mode.decision.v1",
2604                "SessionStart",
2605                "m1",
2606                &sid,
2607                "agent://orchestrator",
2608                session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2609            ),
2610            None,
2611        )
2612        .await
2613        .unwrap();
2614        // Send a Proposal
2615        let proposal = ProposalPayload {
2616            proposal_id: "p1".into(),
2617            option: "step-up".into(),
2618            rationale: "risk".into(),
2619            supporting_data: vec![],
2620        }
2621        .encode_to_vec();
2622        rt.process(
2623            &env(
2624                "macp.mode.decision.v1",
2625                "Proposal",
2626                "m2",
2627                &sid,
2628                "agent://orchestrator",
2629                proposal,
2630            ),
2631            None,
2632        )
2633        .await
2634        .unwrap();
2635        // Commit to resolve the session
2636        let commitment = CommitmentPayload {
2637            commitment_id: "c1".into(),
2638            action: "decision.selected".into(),
2639            authority_scope: "payments".into(),
2640            reason: "bound".into(),
2641            mode_version: "1.0.0".into(),
2642            policy_version: "policy.default".into(),
2643            configuration_version: "cfg-1".into(),
2644            outcome_positive: true,
2645            supersedes: None,
2646        }
2647        .encode_to_vec();
2648        let result = rt
2649            .process(
2650                &env(
2651                    "macp.mode.decision.v1",
2652                    "Commitment",
2653                    "m3",
2654                    &sid,
2655                    "agent://orchestrator",
2656                    commitment,
2657                ),
2658                None,
2659            )
2660            .await
2661            .unwrap();
2662        assert_eq!(result.session_state, SessionState::Resolved);
2663        // Wait a moment so the session's started_at_unix_ms is strictly in the past
2664        tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2665        // Evict with retention = 0 (evict immediately)
2666        rt.evict_stale_sessions(0).await;
2667        // Session should no longer be in the in-memory registry
2668        assert!(rt.registry.get_session(&sid).await.is_none());
2669    }
2670
2671    #[tokio::test]
2672    async fn session_start_with_wrong_mode_version_rejected() {
2673        let rt = make_runtime();
2674        let sid = new_sid();
2675        let payload = SessionStartPayload {
2676            intent: "test".into(),
2677            participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2678            mode_version: "99.0.0".into(), // wrong version
2679            configuration_version: "cfg-1".into(),
2680            policy_version: String::new(),
2681            ttl_ms: 60_000,
2682            context_id: String::new(),
2683            extensions: std::collections::HashMap::new(),
2684            roots: vec![],
2685            max_suspend_ms: 0,
2686        }
2687        .encode_to_vec();
2688
2689        let err = rt
2690            .process(
2691                &env(
2692                    "macp.mode.decision.v1",
2693                    "SessionStart",
2694                    "m1",
2695                    &sid,
2696                    "agent://orchestrator",
2697                    payload,
2698                ),
2699                None,
2700            )
2701            .await
2702            .unwrap_err();
2703        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2704    }
2705
2706    #[tokio::test]
2707    async fn signal_empty_signal_type_rejected() {
2708        let rt = make_runtime();
2709        // Use non-default data so proto3 serializes a non-empty payload
2710        let signal_payload = crate::pb::SignalPayload {
2711            signal_type: String::new(),
2712            data: b"some data".to_vec(),
2713            confidence: 0.0,
2714            correlation_session_id: String::new(),
2715        }
2716        .encode_to_vec();
2717        let signal = Envelope {
2718            macp_version: "1.0".into(),
2719            mode: String::new(),
2720            message_type: "Signal".into(),
2721            message_id: "sig-1".into(),
2722            session_id: String::new(),
2723            sender: "agent://a".into(),
2724            timestamp_unix_ms: 0,
2725            payload: signal_payload,
2726        };
2727        let err = rt.process_signal(&signal).await.unwrap_err();
2728        assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2729    }
2730
2731    #[tokio::test]
2732    async fn signal_valid_payload_accepted() {
2733        let rt = make_runtime();
2734        let signal_payload = crate::pb::SignalPayload {
2735            signal_type: "heartbeat".into(),
2736            data: vec![],
2737            confidence: 0.8,
2738            correlation_session_id: String::new(),
2739        }
2740        .encode_to_vec();
2741        let signal = Envelope {
2742            macp_version: "1.0".into(),
2743            mode: String::new(),
2744            message_type: "Signal".into(),
2745            message_id: "sig-2".into(),
2746            session_id: String::new(),
2747            sender: "agent://a".into(),
2748            timestamp_unix_ms: 0,
2749            payload: signal_payload,
2750        };
2751        rt.process_signal(&signal).await.unwrap();
2752    }
2753
2754    #[tokio::test]
2755    async fn signal_empty_payload_accepted() {
2756        let rt = make_runtime();
2757        let signal = Envelope {
2758            macp_version: "1.0".into(),
2759            mode: String::new(),
2760            message_type: "Signal".into(),
2761            message_id: "sig-3".into(),
2762            session_id: String::new(),
2763            sender: "agent://a".into(),
2764            timestamp_unix_ms: 0,
2765            payload: vec![],
2766        };
2767        rt.process_signal(&signal).await.unwrap();
2768    }
2769
2770    /// Freeze invariant: CommitmentPayload version fields must match the
2771    /// session-bound versions — for extension modes too. When a non-strict ext
2772    /// mode's SessionStart omits mode_version, the runtime binds the registered
2773    /// descriptor's version; a Commitment carrying "" must no longer match
2774    /// vacuously.
2775    #[tokio::test]
2776    async fn ext_mode_empty_version_binds_descriptor_version() {
2777        let rt = make_runtime();
2778        rt.register_extension(ModeDescriptor {
2779            mode: "ext.dyn.v1".into(),
2780            mode_version: "2.5.0".into(),
2781            message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2782            terminal_message_types: vec!["Commitment".into()],
2783            ..Default::default()
2784        })
2785        .unwrap();
2786
2787        let sid = new_sid();
2788        let payload = SessionStartPayload {
2789            participants: vec!["alice".into()],
2790            configuration_version: "cfg-1".into(),
2791            ttl_ms: 60_000,
2792            ..Default::default()
2793        }
2794        .encode_to_vec();
2795        rt.process(
2796            &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2797            None,
2798        )
2799        .await
2800        .unwrap();
2801
2802        // The session is bound to the descriptor's version, not "".
2803        let session = rt.get_session_checked(&sid).await.unwrap();
2804        assert_eq!(session.mode_version, "2.5.0");
2805
2806        // Commitment with empty mode_version: rejected (no vacuous match).
2807        let bad = CommitmentPayload {
2808            commitment_id: "c1".into(),
2809            action: "work.completed".into(),
2810            authority_scope: "test".into(),
2811            reason: "done".into(),
2812            mode_version: String::new(),
2813            policy_version: "policy.default".into(),
2814            configuration_version: "cfg-1".into(),
2815            outcome_positive: true,
2816            supersedes: None,
2817        }
2818        .encode_to_vec();
2819        let err = rt
2820            .process(
2821                &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2822                None,
2823            )
2824            .await
2825            .unwrap_err();
2826        assert_eq!(err.to_string(), "InvalidPayload");
2827
2828        // Commitment echoing the bound descriptor version: accepted, resolves.
2829        let good = CommitmentPayload {
2830            commitment_id: "c1".into(),
2831            action: "work.completed".into(),
2832            authority_scope: "test".into(),
2833            reason: "done".into(),
2834            mode_version: "2.5.0".into(),
2835            policy_version: "policy.default".into(),
2836            configuration_version: "cfg-1".into(),
2837            outcome_positive: true,
2838            supersedes: None,
2839        }
2840        .encode_to_vec();
2841        let result = rt
2842            .process(
2843                &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2844                None,
2845            )
2846            .await
2847            .unwrap();
2848        assert_eq!(result.session_state, SessionState::Resolved);
2849    }
2850
2851    /// The binding must be recorded on the SessionStart log entry (replay reads
2852    /// it from there), and only when the payload actually omitted the version.
2853    #[tokio::test]
2854    async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2855        let rt = make_runtime();
2856        rt.register_extension(ModeDescriptor {
2857            mode: "ext.dyn2.v1".into(),
2858            mode_version: "3.0.0".into(),
2859            message_types: vec!["SessionStart".into(), "Commitment".into()],
2860            terminal_message_types: vec!["Commitment".into()],
2861            ..Default::default()
2862        })
2863        .unwrap();
2864
2865        let sid = new_sid();
2866        let payload = SessionStartPayload {
2867            participants: vec!["alice".into()],
2868            configuration_version: "cfg-1".into(),
2869            ttl_ms: 60_000,
2870            ..Default::default()
2871        }
2872        .encode_to_vec();
2873        rt.process(
2874            &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2875            None,
2876        )
2877        .await
2878        .unwrap();
2879
2880        let log = rt.log_store.get_log(&sid).await.unwrap();
2881        assert_eq!(log[0].message_type, "SessionStart");
2882        assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2883
2884        // A payload that carries the version explicitly records no binding.
2885        let sid2 = new_sid();
2886        let payload2 = SessionStartPayload {
2887            participants: vec!["alice".into()],
2888            mode_version: "3.0.0".into(),
2889            configuration_version: "cfg-1".into(),
2890            ttl_ms: 60_000,
2891            ..Default::default()
2892        }
2893        .encode_to_vec();
2894        rt.process(
2895            &env(
2896                "ext.dyn2.v1",
2897                "SessionStart",
2898                "m1",
2899                &sid2,
2900                "alice",
2901                payload2,
2902            ),
2903            None,
2904        )
2905        .await
2906        .unwrap();
2907        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2908        assert_eq!(log2[0].bound_mode_version, None);
2909    }
2910
2911    /// The RESOLVED suspension cap is bound on the session and recorded on
2912    /// the SessionStart log entry (RFC-MACP-0001 §7.5, RFC-MACP-0003 §2):
2913    /// the payload's positive value verbatim, or the runtime default when
2914    /// the payload carried 0 — never left unrecorded on new sessions.
2915    #[tokio::test]
2916    async fn session_start_binds_and_records_max_suspend_cap() {
2917        let rt = make_runtime();
2918
2919        // Explicit cap: recorded verbatim.
2920        let sid = new_sid();
2921        let payload = SessionStartPayload {
2922            participants: vec!["alice".into(), "bob".into()],
2923            mode_version: "1.0.0".into(),
2924            configuration_version: "cfg-1".into(),
2925            ttl_ms: 60_000,
2926            max_suspend_ms: 12_345,
2927            ..Default::default()
2928        }
2929        .encode_to_vec();
2930        rt.process(
2931            &env(
2932                "macp.mode.decision.v1",
2933                "SessionStart",
2934                "m1",
2935                &sid,
2936                "alice",
2937                payload,
2938            ),
2939            None,
2940        )
2941        .await
2942        .unwrap();
2943        let log = rt.log_store.get_log(&sid).await.unwrap();
2944        assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2945
2946        // Payload 0: the runtime default is resolved and recorded.
2947        let sid2 = new_sid();
2948        let payload2 = SessionStartPayload {
2949            participants: vec!["alice".into(), "bob".into()],
2950            mode_version: "1.0.0".into(),
2951            configuration_version: "cfg-1".into(),
2952            ttl_ms: 60_000,
2953            max_suspend_ms: 0,
2954            ..Default::default()
2955        }
2956        .encode_to_vec();
2957        rt.process(
2958            &env(
2959                "macp.mode.decision.v1",
2960                "SessionStart",
2961                "m2",
2962                &sid2,
2963                "alice",
2964                payload2,
2965            ),
2966            None,
2967        )
2968        .await
2969        .unwrap();
2970        let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2971        assert_eq!(
2972            log2[0].bound_max_suspend_ms,
2973            Some(macp_core::session::MAX_SUSPEND_MS)
2974        );
2975    }
2976
2977    #[test]
2978    fn audit_verbosity_reads_policy_rules() {
2979        let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2980        assert!(!Runtime::audit_verbose(&session));
2981
2982        session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2983            policy_id: "policy.test.audit".into(),
2984            mode: "*".into(),
2985            description: "audited".into(),
2986            rules: serde_json::json!({ "audit": { "level": "info" } }),
2987            schema_version: 1,
2988        });
2989        assert!(Runtime::audit_verbose(&session));
2990
2991        session.policy_definition.as_mut().unwrap().rules =
2992            serde_json::json!({ "audit": { "level": "debug" } });
2993        assert!(!Runtime::audit_verbose(&session));
2994    }
2995
2996    /// Post-commit-point coherence: once the SessionStart log entry is
2997    /// durable, a snapshot failure must NOT fail (or roll back) the start —
2998    /// the previous fatal path left the durable entry behind, so the
2999    /// "failed" session resurrected on restart and a same-id retry appended
3000    /// a second SessionStart that made the log unreplayable.
3001    #[tokio::test]
3002    async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
3003        use std::io;
3004
3005        struct FailSnapshotBackend;
3006        #[async_trait::async_trait]
3007        impl StorageBackend for FailSnapshotBackend {
3008            async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
3009                Ok(())
3010            }
3011            async fn save_session(&self, _s: &Session) -> io::Result<()> {
3012                Err(io::Error::other("snapshot disk full"))
3013            }
3014            async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
3015                Ok(None)
3016            }
3017            async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
3018                Ok(vec![])
3019            }
3020            async fn delete_session(&self, _s: &str) -> io::Result<()> {
3021                Ok(())
3022            }
3023            async fn list_session_ids(&self) -> io::Result<Vec<String>> {
3024                Ok(vec![])
3025            }
3026            async fn append_log_entry(
3027                &self,
3028                _s: &str,
3029                _e: &crate::log_store::LogEntry,
3030            ) -> io::Result<()> {
3031                Ok(())
3032            }
3033            async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
3034                Ok(vec![])
3035            }
3036        }
3037
3038        let rt = Runtime::new(
3039            Arc::new(FailSnapshotBackend),
3040            Arc::new(SessionRegistry::new()),
3041            Arc::new(LogStore::new()),
3042        );
3043        let sid = new_sid();
3044        let result = rt
3045            .process(
3046                &env(
3047                    "macp.mode.decision.v1",
3048                    "SessionStart",
3049                    "m1",
3050                    &sid,
3051                    "agent://orchestrator",
3052                    session_start(vec!["agent://orchestrator".into()]),
3053                ),
3054                None,
3055            )
3056            .await
3057            .expect("start must succeed: the log append (commit point) succeeded");
3058        assert!(!result.duplicate);
3059        // The session exists and is usable.
3060        assert!(rt.get_session_checked(&sid).await.is_some());
3061    }
3062
3063    // --- The client boundary at the kernel's two live entry points (11c) ---
3064    //
3065    // RFC-MACP-0010 §5.1(3). These are the runtime-level halves; the
3066    // mode-level rules are unit-tested in
3067    // `crates/macp-modes/src/mode/handoff.rs`, `step::validate_message` in
3068    // `crates/macp-modes/src/step.rs`, and the proof that the hook is NOT on
3069    // the replay path in `src/replay.rs`
3070    // (`reserved_prefix_entry_replays_at_every_rev`).
3071    //
3072    // Both entry points matter independently: `process_message` and
3073    // `process_session_start` each call the hook themselves, because the
3074    // runtime deliberately bypasses `macp_modes::step::validate_message` (it
3075    // interposes its durable append between validation and commit), so neither
3076    // call site is covered by the other.
3077
3078    const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
3079    const OWNER: &str = "agent://owner";
3080    const TARGET: &str = "agent://target";
3081
3082    fn reserved_id(handoff_id: &str) -> String {
3083        format!(
3084            "{}{handoff_id}",
3085            crate::mode::handoff::IMPLICIT_ACCEPT_MESSAGE_ID_PREFIX
3086        )
3087    }
3088
3089    fn handoff_start_payload() -> Vec<u8> {
3090        session_start(vec![OWNER.into(), TARGET.into()])
3091    }
3092
3093    /// A `Commitment` that echoes the versions `handoff_session_with_offer`
3094    /// binds, so the only thing left to reject it is the `message_id`.
3095    fn handoff_commitment_payload() -> Vec<u8> {
3096        CommitmentPayload {
3097            commitment_id: "c1".into(),
3098            action: "handoff.accepted".into(),
3099            authority_scope: "support".into(),
3100            reason: "bound".into(),
3101            mode_version: "1.0.0".into(),
3102            policy_version: "policy.default".into(),
3103            configuration_version: "cfg-1".into(),
3104            outcome_positive: true,
3105            supersedes: None,
3106        }
3107        .encode_to_vec()
3108    }
3109
3110    fn handoff_offer(handoff_id: &str) -> Vec<u8> {
3111        crate::handoff_pb::HandoffOfferPayload {
3112            handoff_id: handoff_id.into(),
3113            target_participant: TARGET.into(),
3114            scope: "support".into(),
3115            reason: "escalate".into(),
3116        }
3117        .encode_to_vec()
3118    }
3119
3120    fn handoff_context(handoff_id: &str) -> Vec<u8> {
3121        crate::handoff_pb::HandoffContextPayload {
3122            handoff_id: handoff_id.into(),
3123            content_type: "text/plain".into(),
3124            context: b"background".to_vec(),
3125        }
3126        .encode_to_vec()
3127    }
3128
3129    fn handoff_accept(handoff_id: &str, implicit: bool) -> Vec<u8> {
3130        crate::handoff_pb::HandoffAcceptPayload {
3131            handoff_id: handoff_id.into(),
3132            accepted_by: TARGET.into(),
3133            reason: "ready".into(),
3134            implicit,
3135        }
3136        .encode_to_vec()
3137    }
3138
3139    /// An Open handoff session at the current semantics revision with one
3140    /// outstanding offer `h1`. Returns the session id.
3141    async fn handoff_session_with_offer(rt: &Runtime) -> String {
3142        let sid = new_sid();
3143        rt.process(
3144            &env(
3145                HANDOFF_MODE,
3146                "SessionStart",
3147                "start-1",
3148                &sid,
3149                OWNER,
3150                handoff_start_payload(),
3151            ),
3152            None,
3153        )
3154        .await
3155        .expect("handoff session start");
3156        rt.process(
3157            &env(
3158                HANDOFF_MODE,
3159                "HandoffOffer",
3160                "offer-1",
3161                &sid,
3162                OWNER,
3163                handoff_offer("h1"),
3164            ),
3165            None,
3166        )
3167        .await
3168        .expect("handoff offer");
3169        assert_eq!(
3170            rt.get_session_checked(&sid).await.unwrap().semantics_rev,
3171            macp_core::session::CURRENT_SEMANTICS_REV
3172        );
3173        sid
3174    }
3175
3176    /// Phase 11c criterion 1, message path: at the current revision a client
3177    /// envelope whose `message_id` is in the reserved `implicit-accept:`
3178    /// namespace is rejected `InvalidEnvelope`, and the rejection mutates
3179    /// nothing — neither accepted history nor `seen_message_ids` (CLAUDE.md §8
3180    /// dedup invariant).
3181    ///
3182    /// `HandoffContext` is the interesting carrier: it is the one mode message
3183    /// the offerer may send at any disposition, and it is rejected here
3184    /// **before** dispatch, which is why the check cannot live in the mode's
3185    /// own rules. `Commitment` is included because the initiator can squat the
3186    /// id that way too, and because it proves the reservation is not scoped to
3187    /// `HandoffAccept`.
3188    #[tokio::test]
3189    async fn reserved_message_id_namespace_is_rejected_at_rev2() {
3190        let rt = make_runtime();
3191        let sid = handoff_session_with_offer(&rt).await;
3192
3193        let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3194        let dedup_before = rt
3195            .get_session_checked(&sid)
3196            .await
3197            .unwrap()
3198            .seen_message_ids
3199            .clone();
3200
3201        for (message_type, sender, payload) in [
3202            ("HandoffContext", OWNER, handoff_context("h1")),
3203            ("Commitment", OWNER, handoff_commitment_payload()),
3204            ("HandoffAccept", TARGET, handoff_accept("h1", false)),
3205        ] {
3206            let err = rt
3207                .process(
3208                    &env(
3209                        HANDOFF_MODE,
3210                        message_type,
3211                        &reserved_id("h1"),
3212                        &sid,
3213                        sender,
3214                        payload,
3215                    ),
3216                    None,
3217                )
3218                .await
3219                .unwrap_err();
3220            assert!(
3221                matches!(err, MacpError::InvalidEnvelope),
3222                "{message_type} with a reserved id must be InvalidEnvelope, got {err}"
3223            );
3224        }
3225
3226        // Nothing was appended and no dedup slot was consumed.
3227        let session = rt.get_session_checked(&sid).await.unwrap();
3228        assert_eq!(
3229            rt.log_store.get_log(&sid).await.unwrap().len(),
3230            history_before
3231        );
3232        assert_eq!(session.seen_message_ids, dedup_before);
3233        assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3234        assert_eq!(session.state, SessionState::Open);
3235
3236        // Control: the same `HandoffContext` with an ordinary id is accepted,
3237        // so the rejections above are the id's doing and not the payload's.
3238        rt.process(
3239            &env(
3240                HANDOFF_MODE,
3241                "HandoffContext",
3242                "ctx-1",
3243                &sid,
3244                OWNER,
3245                handoff_context("h1"),
3246            ),
3247            None,
3248        )
3249        .await
3250        .expect("an ordinary id is accepted");
3251        assert_eq!(
3252            rt.log_store.get_log(&sid).await.unwrap().len(),
3253            history_before + 1
3254        );
3255    }
3256
3257    /// Phase 11c criterion 1, start path. This call site exists precisely
3258    /// because the namespace is squattable through `SessionStart`, whose
3259    /// `message_id` enters `seen_message_ids` at the commit point — after
3260    /// which the runtime's own later synthesis would be silently skipped.
3261    ///
3262    /// The rejection must leave **no** trace: no session in the registry, no
3263    /// session log, and no reservation of the session id — proved by starting
3264    /// the same session id again with a clean `message_id` and having it
3265    /// succeed (`SessionAlreadyExists` would be the failure signature of a
3266    /// leaked reservation).
3267    #[tokio::test]
3268    async fn reserved_message_id_is_rejected_on_the_session_start_path() {
3269        let rt = make_runtime();
3270        let sid = new_sid();
3271
3272        let err = rt
3273            .process(
3274                &env(
3275                    HANDOFF_MODE,
3276                    "SessionStart",
3277                    &reserved_id("h1"),
3278                    &sid,
3279                    OWNER,
3280                    handoff_start_payload(),
3281                ),
3282                None,
3283            )
3284            .await
3285            .unwrap_err();
3286        assert!(
3287            matches!(err, MacpError::InvalidEnvelope),
3288            "reserved id on SessionStart must be InvalidEnvelope, got {err}"
3289        );
3290
3291        // No session, no log, nothing to roll back.
3292        assert!(rt.get_session_checked(&sid).await.is_none());
3293        assert!(rt.log_store.get_log(&sid).await.is_none());
3294        assert!(!rt.registry.sessions.read().await.contains_key(&sid));
3295
3296        // The session id was never reserved and the id never consumed a dedup
3297        // slot: the same session starts cleanly.
3298        rt.process(
3299            &env(
3300                HANDOFF_MODE,
3301                "SessionStart",
3302                "start-1",
3303                &sid,
3304                OWNER,
3305                handoff_start_payload(),
3306            ),
3307            None,
3308        )
3309        .await
3310        .expect("a rejected SessionStart must not reserve the session id");
3311        let session = rt.get_session_checked(&sid).await.unwrap();
3312        assert!(session.seen_message_ids.contains("start-1"));
3313        assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3314    }
3315
3316    /// Phase 11c criterion 2, runtime path: a client `HandoffAccept` carrying
3317    /// `implicit = true` never enters history.
3318    ///
3319    /// Two envelopes, as the criterion requires:
3320    /// (a) `implicit = true` with an ordinary `message_id` — `InvalidPayload`.
3321    ///     **This assertion is double-guarded today** and so does not isolate
3322    ///     the hook: `handle_message` rejects the same envelope with the same
3323    ///     code (deliberately kept until 11d restructures that arm). What it
3324    ///     does pin is the criterion's actual requirement — that the rev-2
3325    ///     error *surface* through `Send` is unchanged — and it becomes the
3326    ///     only guard once 11d teaches dispatch to accept the shape.
3327    /// (b) `implicit = true` with the **reserved** `message_id` and the
3328    ///     correct sender — `InvalidEnvelope`, which only the boundary can
3329    ///     produce (dispatch would say `InvalidPayload`). The envelope is
3330    ///     byte-shaped exactly like the one the runtime will synthesize
3331    ///     from 11e.
3332    ///
3333    /// Measured, **neither half pins the `implicit` rule**: (b) is killed by
3334    /// the *reserved-prefix* rule, which fires first and returns
3335    /// `InvalidEnvelope` whatever the flag says, so deleting the `implicit`
3336    /// rule leaves this whole test green. The `implicit` rule's only
3337    /// non-vacuous guard is the mode-level unit test
3338    /// `handoff::tests::client_implicit_accept_rejected_at_the_boundary`.
3339    /// What this test pins is the runtime-level *error surface* at rev 2 —
3340    /// which is the criterion's requirement, and which 11d must keep in view
3341    /// when it restructures the dispatch arm.
3342    #[tokio::test]
3343    async fn client_implicit_accept_rejected_through_the_runtime() {
3344        let rt = make_runtime();
3345        let sid = handoff_session_with_offer(&rt).await;
3346        let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3347
3348        // (a) ordinary id.
3349        let err = rt
3350            .process(
3351                &env(
3352                    HANDOFF_MODE,
3353                    "HandoffAccept",
3354                    "accept-1",
3355                    &sid,
3356                    TARGET,
3357                    handoff_accept("h1", true),
3358                ),
3359                None,
3360            )
3361            .await
3362            .unwrap_err();
3363        assert!(matches!(err, MacpError::InvalidPayload), "got {err}");
3364
3365        // (b) the runtime's own synthetic shape, submitted by the target.
3366        let err = rt
3367            .process(
3368                &env(
3369                    HANDOFF_MODE,
3370                    "HandoffAccept",
3371                    &reserved_id("h1"),
3372                    &sid,
3373                    TARGET,
3374                    handoff_accept("h1", true),
3375                ),
3376                None,
3377            )
3378            .await
3379            .unwrap_err();
3380        assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3381
3382        // The offer is still outstanding and history is untouched.
3383        let session = rt.get_session_checked(&sid).await.unwrap();
3384        assert_eq!(
3385            rt.log_store.get_log(&sid).await.unwrap().len(),
3386            history_before
3387        );
3388        assert!(session.seen_message_ids.is_disjoint(
3389            &["accept-1".to_string(), reserved_id("h1")]
3390                .into_iter()
3391                .collect()
3392        ));
3393        let mode_state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
3394        assert_eq!(mode_state["offers"]["h1"]["disposition"], "Offered");
3395
3396        // Control: the explicit accept (`implicit = false`, ordinary id) is
3397        // accepted, so the rejections above are not the envelope's other
3398        // fields.
3399        rt.process(
3400            &env(
3401                HANDOFF_MODE,
3402                "HandoffAccept",
3403                "accept-2",
3404                &sid,
3405                TARGET,
3406                handoff_accept("h1", false),
3407            ),
3408            None,
3409        )
3410        .await
3411        .expect("an explicit accept is still accepted");
3412    }
3413
3414    /// Error-code ordering at the kernel is unchanged by the phase: an
3415    /// envelope that is **both** unauthorized and carries a reserved id
3416    /// reports `Forbidden`, because the hook is called after
3417    /// `mode.authorize_sender`. Rev <= 1's Forbidden-before-payload ordering
3418    /// therefore does not shift at rev 2.
3419    #[tokio::test]
3420    async fn client_boundary_error_ordering_is_unchanged_at_rev2() {
3421        let rt = make_runtime();
3422        let sid = handoff_session_with_offer(&rt).await;
3423
3424        // `agent://stranger` is not a declared participant.
3425        let err = rt
3426            .process(
3427                &env(
3428                    HANDOFF_MODE,
3429                    "HandoffAccept",
3430                    &reserved_id("h1"),
3431                    &sid,
3432                    "agent://stranger",
3433                    handoff_accept("h1", true),
3434                ),
3435                None,
3436            )
3437            .await
3438            .unwrap_err();
3439        assert!(
3440            matches!(err, MacpError::Forbidden),
3441            "authorization must be reported before the client boundary, got {err}"
3442        );
3443
3444        // Same envelope from the authorized sender: now the boundary speaks.
3445        let err = rt
3446            .process(
3447                &env(
3448                    HANDOFF_MODE,
3449                    "HandoffAccept",
3450                    &reserved_id("h1"),
3451                    &sid,
3452                    TARGET,
3453                    handoff_accept("h1", true),
3454                ),
3455                None,
3456            )
3457            .await
3458            .unwrap_err();
3459        assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3460    }
3461
3462    // --- The synthesis seam's own guards (11e) ---
3463
3464    /// A handoff session with a bound `implicit_accept_timeout_ms` and an
3465    /// outstanding offer. Unlike [`handoff_session_with_offer`] this one binds
3466    /// a policy, so an implicit accept can actually become due.
3467    async fn handoff_session_with_timed_offer(rt: &Runtime, timeout_ms: i64) -> String {
3468        rt.register_policy(macp_core::policy::PolicyDefinition {
3469            policy_id: "handoff-timed".into(),
3470            mode: HANDOFF_MODE.into(),
3471            description: "implicit accept".into(),
3472            rules: serde_json::json!({
3473                "acceptance": { "implicit_accept_timeout_ms": timeout_ms },
3474                "commitment": { "authority": "initiator_only" }
3475            }),
3476            schema_version: 1,
3477        })
3478        .expect("policy registers");
3479
3480        let sid = new_sid();
3481        let start = SessionStartPayload {
3482            intent: "escalate".into(),
3483            participants: vec![OWNER.into(), TARGET.into()],
3484            mode_version: "1.0.0".into(),
3485            configuration_version: "cfg-1".into(),
3486            policy_version: "handoff-timed".into(),
3487            ttl_ms: 60_000,
3488            context_id: String::new(),
3489            extensions: std::collections::HashMap::new(),
3490            roots: vec![],
3491            max_suspend_ms: 0,
3492        }
3493        .encode_to_vec();
3494        rt.process(
3495            &env(HANDOFF_MODE, "SessionStart", "start-1", &sid, OWNER, start),
3496            None,
3497        )
3498        .await
3499        .expect("session start");
3500        rt.process(
3501            &env(
3502                HANDOFF_MODE,
3503                "HandoffOffer",
3504                "offer-1",
3505                &sid,
3506                OWNER,
3507                handoff_offer("h1"),
3508            ),
3509            None,
3510        )
3511        .await
3512        .expect("offer");
3513        sid
3514    }
3515
3516    /// Phase 11e criterion 11, first half: `synthesize_due_accept` must refuse
3517    /// to synthesize for a session that is not `Open`.
3518    ///
3519    /// This is a **correctness** gate, not hygiene, and it is the only one
3520    /// that ships. Both computations behind the synthetic entry ignore an
3521    /// in-flight pause by design — `Session::unsuspended_deadline` walks
3522    /// completed intervals only, and `HandoffMode::rev2_elapsed_ms` has no
3523    /// in-flight term — so asking a `Suspended` session over-counts elapsed
3524    /// time (the accept can be judged due when it is not) and under-computes
3525    /// `D` (a wrong `timestamp_unix_ms` baked into permanent history, in
3526    /// silence). The mode carries a `debug_assert!` for the same thing, which
3527    /// compiles out in release builds and therefore guards nothing where it
3528    /// matters.
3529    ///
3530    /// The method is called directly because no message path can reach it with
3531    /// a suspended session — `step::check_preconditions` rejects every message
3532    /// to a non-`Open` session first. Phase 12's eager sweep is the caller that
3533    /// *will* see suspended sessions, which is exactly why the filter has to
3534    /// be in the seam rather than at its current call site.
3535    #[tokio::test]
3536    async fn synthesis_is_skipped_for_a_non_open_session() {
3537        let rt = make_runtime();
3538        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3539        rt.suspend_session(&sid, "hold", OWNER)
3540            .await
3541            .expect("suspend");
3542
3543        let shared = rt.registry.get_shared(&sid).await.unwrap();
3544        let mut guard = shared.lock().await;
3545        let session = &mut *guard;
3546        assert_eq!(session.state, SessionState::Suspended);
3547        assert!(session.suspended_at_ms.is_some());
3548
3549        // Far past the deadline on wall time — the only thing stopping a
3550        // synthesis here is the state filter.
3551        let long_after = session.suspended_at_ms.unwrap() + 10_000;
3552        let log_before = rt.log_store.get_log(&sid).await.unwrap_or_default().len();
3553        let dedup_before = session.seen_message_ids.len();
3554        let mode_state_before = session.mode_state.clone();
3555
3556        rt.synthesize_due_accept(&sid, session, long_after)
3557            .await
3558            .expect("the filter is a skip, not an error");
3559
3560        assert_eq!(
3561            rt.log_store.get_log(&sid).await.unwrap_or_default().len(),
3562            log_before,
3563            "a suspended session must not gain a synthetic entry"
3564        );
3565        assert_eq!(session.seen_message_ids.len(), dedup_before);
3566        assert_eq!(session.mode_state, mode_state_before);
3567
3568        // Control: the filter is what declined, not the arithmetic — and what
3569        // it kept out of history was wrong, not merely early.
3570        //
3571        // The mode cannot be asked about the genuinely suspended session in a
3572        // debug build: `due_synthetic_envelope` carries a `debug_assert!` on
3573        // `suspended_at_ms.is_none()` and would panic. So the probe is a clone
3574        // that differs *only* by that tripwire — same offer, same banked
3575        // suspension (none: the pause is still in flight and banks on resume),
3576        // same clock. That is exactly the state the mode would see if the
3577        // kernel filter were removed and the assertion compiled out, which is
3578        // what a release build does.
3579        let mut without_the_tripwire = session.clone();
3580        without_the_tripwire.state = SessionState::Open;
3581        without_the_tripwire.suspended_at_ms = None;
3582        let mode = rt.mode_registry.get_mode(&session.mode).unwrap();
3583        let would_have_emitted = mode
3584            .due_synthetic_envelope(&without_the_tripwire, long_after)
3585            .expect("the mode would have synthesized; only the kernel filter stopped it");
3586        // The harm, concretely: the D it would have recorded falls *inside* the
3587        // pause that is still running — a moment at which the session was not
3588        // ticking at all. `unsuspended_deadline` walks completed intervals
3589        // only, and this pause is not completed, so it is invisible to the
3590        // walk. Recorded, that timestamp would be permanent and wrong.
3591        let suspended_at = session.suspended_at_ms.unwrap();
3592        assert!(
3593            would_have_emitted.timestamp_unix_ms >= suspended_at
3594                && would_have_emitted.timestamp_unix_ms < long_after,
3595            "D {} must fall inside the still-open pause starting at {suspended_at}",
3596            would_have_emitted.timestamp_unix_ms
3597        );
3598
3599        // And once the session is Open again the seam works normally, so the
3600        // filter is a skip rather than a permanent disable.
3601        drop(guard);
3602        rt.resume_session(&sid, "go", OWNER).await.expect("resume");
3603        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3604        let shared = rt.registry.get_shared(&sid).await.unwrap();
3605        let mut guard = shared.lock().await;
3606        let session = &mut *guard;
3607        let now = Utc::now().timestamp_millis();
3608        rt.synthesize_due_accept(&sid, session, now).await.unwrap();
3609        assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3610    }
3611
3612    /// Phase 11e criterion 11, second half: the appended entry stamps
3613    /// `received_at_ms` with the envelope's own `timestamp_unix_ms` (the
3614    /// deadline `D`), never wall-clock.
3615    ///
3616    /// Asserted on the stamped value directly, at the seam, rather than
3617    /// inferred from a replay — a replay would pass under a wall-clock stamp
3618    /// too, because handoff's accept arm is time-blind. The contract is
3619    /// general: replay derives its dispatch clock from `received_at_ms`, so
3620    /// the next mode to use this hook would silently diverge.
3621    #[tokio::test]
3622    async fn synthetic_entry_stamps_received_at_with_the_deadline() {
3623        let rt = make_runtime();
3624        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3625        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3626
3627        let offer_received_at = rt
3628            .log_store
3629            .get_log(&sid)
3630            .await
3631            .unwrap()
3632            .iter()
3633            .find(|e| e.message_type == "HandoffOffer")
3634            .expect("offer entry")
3635            .received_at_ms;
3636        let expected_d = offer_received_at + 20;
3637
3638        let shared = rt.registry.get_shared(&sid).await.unwrap();
3639        let mut guard = shared.lock().await;
3640        let session = &mut *guard;
3641        // Observed far later than D, so a wall-clock stamp would be obvious.
3642        let observed = Utc::now().timestamp_millis();
3643        assert!(observed > expected_d);
3644        rt.synthesize_due_accept(&sid, session, observed)
3645            .await
3646            .unwrap();
3647        drop(guard);
3648
3649        let entry = rt
3650            .log_store
3651            .get_log(&sid)
3652            .await
3653            .unwrap()
3654            .into_iter()
3655            .find(|e| e.message_id == reserved_id("h1"))
3656            .expect("the synthetic entry");
3657        assert_eq!(entry.timestamp_unix_ms, expected_d, "envelope clock is D");
3658        assert_eq!(entry.received_at_ms, expected_d, "entry clock is D");
3659        assert_ne!(
3660            entry.received_at_ms, observed,
3661            "received_at_ms must not be the observation time"
3662        );
3663        assert_eq!(entry.entry_kind, EntryKind::Incoming);
3664    }
3665
3666    /// The synthetic accept is deliberately NOT credited as participant
3667    /// activity: `replay_entry` never records activity for any entry kind, so
3668    /// calling `record_participant_activity` live would guarantee a
3669    /// live/replay divergence in `participant_message_counts` — and the target
3670    /// did not, in fact, send anything.
3671    ///
3672    /// The consequence is user-visible (`SessionMetadata.participant_activity`
3673    /// via `server::session_to_metadata`), so it is pinned rather than left as
3674    /// a comment.
3675    #[tokio::test]
3676    async fn synthetic_accept_is_not_credited_as_participant_activity() {
3677        let rt = make_runtime();
3678        let sid = handoff_session_with_timed_offer(&rt, 20).await;
3679        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3680
3681        let shared = rt.registry.get_shared(&sid).await.unwrap();
3682        let mut guard = shared.lock().await;
3683        let session = &mut *guard;
3684        let before = session.participant_message_counts.get(TARGET).copied();
3685        rt.synthesize_due_accept(&sid, session, Utc::now().timestamp_millis())
3686            .await
3687            .unwrap();
3688        assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3689        assert_eq!(
3690            session.participant_message_counts.get(TARGET).copied(),
3691            before,
3692            "the target must not be credited with a message they did not send"
3693        );
3694    }
3695    /// The checkpoint-interval check belongs to the *append*, so the eager
3696    /// sweep and a lazy trigger place the same checkpoint in the same
3697    /// position.
3698    ///
3699    /// `maybe_insert_checkpoint` used to be called only from the tail of
3700    /// `process_message`. A synthetic entry advances `log_len` without ever
3701    /// reaching it: the lazy path checked only the length *after* the
3702    /// trigger's own append (one greater), and the eager path — where
3703    /// `sweep_due_synthetic_accepts` is the whole call — checked nothing at
3704    /// all. With `MACP_CHECKPOINT_INTERVAL > 0` that is a genuine eager/lazy
3705    /// divergence: identical accepted histories, different checkpoint
3706    /// placement, and a boundary the synthetic crossed silently skipped.
3707    /// Calling it from inside `synthesize_due_accept` makes the two agree by
3708    /// construction, which is what this pins.
3709    ///
3710    /// Checkpoints are a replay optimization rather than a correctness
3711    /// property, which is exactly why this needs a test: nothing else would
3712    /// ever notice.
3713    ///
3714    /// `checkpoint_interval` is set on the struct rather than through
3715    /// `MACP_CHECKPOINT_INTERVAL`, because that variable is read once in
3716    /// `Runtime::with_registries` and this binary runs its tests in parallel —
3717    /// setting it here would leak into every other runtime built concurrently.
3718    ///
3719    /// Interval 3, with `SessionStart` + `HandoffOffer` already logged, puts
3720    /// the synthetic accept exactly on the boundary: the skipped case.
3721    #[tokio::test]
3722    async fn a_synthetic_entry_on_the_checkpoint_boundary_checkpoints_either_path() {
3723        async fn shape(rt: &Runtime, sid: &str) -> Vec<(EntryKind, String)> {
3724            rt.log_store
3725                .get_log(sid)
3726                .await
3727                .expect("log")
3728                .iter()
3729                .map(|e| (e.entry_kind.clone(), e.message_type.clone()))
3730                .collect()
3731        }
3732
3733        // Eager: the sweep is the only thing that touches the session.
3734        let mut eager = make_runtime();
3735        eager.checkpoint_interval = 3;
3736        let eager_sid = handoff_session_with_timed_offer(&eager, 20).await;
3737        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3738        assert_eq!(
3739            eager.sweep_due_synthetic_accepts().await,
3740            1,
3741            "sweep emitted"
3742        );
3743
3744        // Lazy: a trigger message arrives instead. `HandoffContext` is the one
3745        // mode message the offerer may send at any disposition, so it provokes
3746        // the synthesis without being an accept itself.
3747        let mut lazy = make_runtime();
3748        lazy.checkpoint_interval = 3;
3749        let lazy_sid = handoff_session_with_timed_offer(&lazy, 20).await;
3750        tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3751        lazy.process(
3752            &env(
3753                HANDOFF_MODE,
3754                "HandoffContext",
3755                "ctx-1",
3756                &lazy_sid,
3757                OWNER,
3758                handoff_context("h1"),
3759            ),
3760            None,
3761        )
3762        .await
3763        .expect("context accepted");
3764
3765        let expected = vec![
3766            (EntryKind::Incoming, "SessionStart".to_string()),
3767            (EntryKind::Incoming, "HandoffOffer".to_string()),
3768            (EntryKind::Incoming, "HandoffAccept".to_string()),
3769            (EntryKind::Checkpoint, "Checkpoint".to_string()),
3770        ];
3771        assert_eq!(shape(&eager, &eager_sid).await, expected, "eager sweep");
3772
3773        let lazy_shape = shape(&lazy, &lazy_sid).await;
3774        assert_eq!(
3775            lazy_shape[..4],
3776            expected[..],
3777            "the lazy path must checkpoint in the same place as the eager one"
3778        );
3779        // ... and exactly once: `process_message` runs its own check after
3780        // appending the trigger, at a length the synthesis never tested.
3781        assert_eq!(
3782            lazy_shape[4..],
3783            [(EntryKind::Incoming, "HandoffContext".to_string())],
3784            "no second checkpoint for the trigger's own append"
3785        );
3786    }
3787}