Skip to main content

macp_runtime/
runtime.rs

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