Skip to main content

macp_core/
session.rs

1use crate::error::MacpError;
2use crate::mode::ModeResponse;
3use crate::policy::PolicyDefinition;
4use macp_pb::pb::SessionStartPayload;
5use prost::Message;
6use std::collections::{HashMap, HashSet};
7
8pub const MAX_TTL_MS: i64 = 24 * 60 * 60 * 1000;
9
10/// Default cap on the cumulative time a session may spend `Suspended` before
11/// it is force-expired (RFC-MACP-0001 §7.5). Bounds indefinite
12/// human-in-the-loop holds. Sessions may bind their own cap via
13/// `SessionStartPayload.max_suspend_ms` (0 selects this default); the
14/// resolved cap is recorded at SessionStart and used by replay
15/// (RFC-MACP-0003 §2) — see [`Session::effective_max_suspend_ms`].
16pub const MAX_SUSPEND_MS: i64 = 7 * 24 * 60 * 60 * 1000;
17
18/// Cap on the number of completed suspend/resume *cycles* a session may
19/// record in [`Session::suspension_intervals`]. Enforced in
20/// [`Session::resume`] at **every** revision, but with revision-dependent
21/// behavior (see below), so legacy histories replay bit-identically.
22///
23/// Why a *count* cap is needed even though [`MAX_SUSPEND_MS`] exists:
24/// `MAX_SUSPEND_MS` bounds accumulated suspended *duration* (7 days by
25/// default), not the number of cycles — N one-millisecond suspend/resume
26/// cycles accrue ~0 against that budget, and there is no other cycle counter
27/// anywhere in the session model. `SuspendSession`/`ResumeSession` are also
28/// un-rate-limited RPCs (unlike `Send`), so the cycle count is attacker- (or
29/// bug-) controlled by the session initiator alone.
30///
31/// What the cap closes: each cycle writes two full `PersistedSession`
32/// snapshots, and `suspension_intervals` is part of that snapshot — so an
33/// unbounded vec turns two constant-size writes per cycle into O(N) writes,
34/// i.e. O(N²) total snapshot bytes over a session's life. Bounding the vec
35/// bounds the amplification — and the vec must be bounded at *every*
36/// revision, because a session replayed from a legacy log loads at rev 0/1
37/// and can still be suspended and resumed through the un-rate-limited
38/// `SuspendSession`/`ResumeSession` RPCs.
39///
40/// Two behaviors, split on the revision:
41///
42/// - **`semantics_rev >= 2`** — over the cap, `resume` takes the posture it
43///   already takes for `MAX_SUSPEND_MS`: force-expire the session and return
44///   [`MacpError::TtlExpired`].
45/// - **`semantics_rev <= 1`** — `resume` keeps succeeding but simply **stops
46///   recording** once the vec reaches the cap. Nothing reads
47///   `suspension_intervals` below rev 2 ([`Session::unsuspended_deadline`] is
48///   a rev >= 2 path), so dropping the overflow keeps legacy replay
49///   bit-identical while still bounding memory and snapshot size. A legacy
50///   session is never force-expired by a rule that did not exist when it was
51///   accepted.
52pub const MAX_SUSPENSION_CYCLES: usize = 1024;
53
54/// Current session-semantics revision. Recorded at SessionStart (on the
55/// session and its log entry) and consulted wherever acceptance-time behavior
56/// changed across releases, so legacy histories replay under the semantics
57/// they were accepted with (RFC-MACP-0003 §1).
58///
59/// Revisions:
60/// - 0 — legacy: Handoff implicit-accept timed against the client-supplied
61///   envelope timestamp.
62/// - 1 — Handoff implicit-accept times against the runtime acceptance clock
63///   (`MessageContext::accepted_at_ms`).
64/// - 2 — the Handoff implicit accept becomes a **recorded event** rather than
65///   an inference, and its deadline excludes suspended time (RFC-MACP-0010
66///   §5.1). Two changes, one revision:
67///   - *The synthetic entry.* Once an outstanding offer's
68///     `implicit_accept_timeout_ms` has elapsed, the runtime appends a
69///     synthetic `HandoffAccept` envelope to accepted history — sender = the
70///     offer's target, `implicit = true`, deterministic `message_id`
71///     `implicit-accept:<handoff_id>`, and both clocks
72///     (`timestamp_unix_ms` / the entry's `received_at_ms`) fixed at the
73///     computed deadline `D`, never at observation time — *before* evaluating
74///     any subsequent message against the offer's acceptance state (§5.1(2)).
75///     Clients may not submit that shape: the reserved id namespace and
76///     `implicit = true` are rejected at the client boundary (§5.1(3)).
77///     Revisions 0 and 1 record no such entry and instead infer the accept
78///     inside `Commitment` handling, so their histories replay to the outcome
79///     they were accepted with; at rev >= 2 a `Commitment` on a history
80///     lacking the synthetic entry fails loudly instead.
81///   - *The suspension-corrected deadline* (§5.1(1)): time the session spends
82///     `Suspended` no longer counts toward `implicit_accept_timeout_ms`. The
83///     offer record snapshots `accumulated_suspended_ms` at offer time and the
84///     timeout arithmetic subtracts the suspension accrued since the offer.
85///     Revisions 0 and 1 keep counting suspended time.
86pub const CURRENT_SEMANTICS_REV: u32 = 2;
87
88#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
89pub enum SessionState {
90    Open,
91    /// Non-terminal pause of an `Open` session (RFC-MACP-0001 §7.5). TTL is
92    /// banked while suspended; only `Open`<->`Suspended` and `Suspended`->
93    /// `Expired`/`Cancelled` transitions are permitted.
94    Suspended,
95    Resolved,
96    Expired,
97    /// Terminal: ended by an accepted `CancelSession` (distinct from `Expired`).
98    Cancelled,
99}
100
101impl SessionState {
102    /// Terminal states accept no further transitions.
103    pub fn is_terminal(&self) -> bool {
104        matches!(
105            self,
106            SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
107        )
108    }
109}
110
111/// Session model. Fields are public for reads, but the struct is
112/// `#[non_exhaustive]`: construct via [`Session::builder`]. This lets the model
113/// gain fields without breaking every constructor in downstream crates
114/// (pre-1.0 freeze requirement).
115#[non_exhaustive]
116#[derive(Clone, Debug)]
117pub struct Session {
118    pub session_id: String,
119    pub state: SessionState,
120    pub ttl_expiry: i64,
121    pub ttl_ms: i64,
122    pub started_at_unix_ms: i64,
123    pub resolution: Option<Vec<u8>>,
124    pub mode: String,
125    pub mode_state: Vec<u8>,
126    pub participants: Vec<String>,
127    pub seen_message_ids: HashSet<String>,
128    pub intent: String,
129    pub mode_version: String,
130    pub configuration_version: String,
131    pub policy_version: String,
132    pub context_id: String,
133    pub extensions: HashMap<String, Vec<u8>>,
134    pub roots: Vec<macp_pb::pb::Root>,
135    pub initiator_sender: String,
136    pub participant_message_counts: HashMap<String, u32>,
137    pub participant_last_seen: HashMap<String, i64>,
138    pub policy_definition: Option<PolicyDefinition>,
139    /// Wall-clock (session-timeline) ms at which the session was suspended, or
140    /// `None` when not suspended. Used to bank TTL across a suspension (§7.5).
141    pub suspended_at_ms: Option<i64>,
142    /// Cumulative ms the session has spent suspended across all suspend/resume
143    /// cycles. Drives the `MAX_SUSPEND_MS` cap.
144    pub accumulated_suspended_ms: i64,
145    /// Completed `(suspended_at, resumed_at)` pairs on the session timeline
146    /// (ms), in the order they completed. An *in-progress* suspension is
147    /// deliberately **not** here — the pair is pushed by [`Session::resume`],
148    /// so the vec always describes finished pauses only.
149    ///
150    /// Recorded at **every** `semantics_rev` (so a session that later matters
151    /// has the history), but read only at rev >= 2, by
152    /// [`Session::unsuspended_deadline`] (RFC-MACP-0010 §5.1). Legacy
153    /// snapshots and checkpoints deserialize this as empty; see
154    /// `unsuspended_deadline`'s under-count invariant for why that is safe.
155    ///
156    /// Bounded at every revision by [`MAX_SUSPENSION_CYCLES`] — rev >= 2
157    /// force-expires the session over the cap, rev <= 1 stops recording.
158    pub suspension_intervals: Vec<(i64, i64)>,
159    /// Session-semantics revision this session was accepted under. See
160    /// [`CURRENT_SEMANTICS_REV`]. Legacy persisted sessions load as `0`.
161    pub semantics_rev: u32,
162    /// Maximum-suspension cap bound at SessionStart (RFC-MACP-0001 §7.5,
163    /// RFC-MACP-0003 §2). `0` = unbound (legacy sessions and library
164    /// defaults) — the [`MAX_SUSPEND_MS`] default applies; see
165    /// [`Session::effective_max_suspend_ms`]. The kernel records the
166    /// *resolved* value here for new sessions.
167    pub max_suspend_ms: i64,
168}
169
170impl Session {
171    /// Start building a `Session`. The three arguments are the fields with no
172    /// meaningful default; everything else starts from documented defaults
173    /// (see [`SessionBuilder`]) and is set with the builder's methods.
174    pub fn builder(
175        session_id: impl Into<String>,
176        mode: impl Into<String>,
177        initiator_sender: impl Into<String>,
178    ) -> SessionBuilder {
179        SessionBuilder {
180            inner: Session {
181                session_id: session_id.into(),
182                state: SessionState::Open,
183                // Never-expires by default: every kernel path overrides this
184                // from the validated SessionStart payload; library/test
185                // consumers get a session that behaves until told otherwise.
186                ttl_expiry: i64::MAX,
187                ttl_ms: 0,
188                started_at_unix_ms: 0,
189                resolution: None,
190                mode: mode.into(),
191                mode_state: vec![],
192                participants: vec![],
193                seen_message_ids: HashSet::new(),
194                intent: String::new(),
195                mode_version: String::new(),
196                configuration_version: String::new(),
197                policy_version: String::new(),
198                context_id: String::new(),
199                extensions: HashMap::new(),
200                roots: vec![],
201                initiator_sender: initiator_sender.into(),
202                participant_message_counts: HashMap::new(),
203                participant_last_seen: HashMap::new(),
204                policy_definition: None,
205                suspended_at_ms: None,
206                accumulated_suspended_ms: 0,
207                suspension_intervals: vec![],
208                semantics_rev: CURRENT_SEMANTICS_REV,
209                max_suspend_ms: 0,
210            },
211        }
212    }
213
214    pub fn record_participant_activity(&mut self, sender: &str, timestamp_ms: i64) {
215        *self
216            .participant_message_counts
217            .entry(sender.to_string())
218            .or_insert(0) += 1;
219        self.participant_last_seen
220            .insert(sender.to_string(), timestamp_ms);
221    }
222
223    /// Suspend an `Open` session (RFC-MACP-0001 §7.5). Records the suspend time
224    /// so TTL can be banked on resume. Pure: no clock, no I/O — the caller
225    /// injects `now_ms`.
226    pub fn suspend(&mut self, now_ms: i64) -> Result<(), MacpError> {
227        if self.state != SessionState::Open {
228            return Err(MacpError::SessionNotOpen);
229        }
230        self.state = SessionState::Suspended;
231        self.suspended_at_ms = Some(now_ms);
232        Ok(())
233    }
234
235    /// The suspension cap governing this session: the value bound at
236    /// SessionStart, or the [`MAX_SUSPEND_MS`] default when unbound (0).
237    pub fn effective_max_suspend_ms(&self) -> i64 {
238        if self.max_suspend_ms > 0 {
239            self.max_suspend_ms
240        } else {
241            MAX_SUSPEND_MS
242        }
243    }
244
245    /// Resume a `Suspended` session, banking the suspended duration into the
246    /// TTL deadline (`ttl_expiry += now - suspended_at`) and recording the
247    /// completed pause in [`Session::suspension_intervals`].
248    ///
249    /// Force-expires the session (state `Expired`, `Err(TtlExpired)`) when
250    /// either suspension cap is exceeded: cumulative duration past
251    /// [`MAX_SUSPEND_MS`] (every revision), or completed cycle count past
252    /// [`MAX_SUSPENSION_CYCLES`] (`semantics_rev >= 2` only — at rev <= 1 the
253    /// cap instead stops the recording, see [`MAX_SUSPENSION_CYCLES`]). The
254    /// pair is pushed before either check, so history is recorded even on the
255    /// expiring call.
256    ///
257    /// Pure: no clock, no I/O — the caller injects `now_ms`.
258    pub fn resume(&mut self, now_ms: i64) -> Result<(), MacpError> {
259        if self.state != SessionState::Suspended {
260            return Err(MacpError::SessionNotOpen);
261        }
262        let suspended_at = self.suspended_at_ms.unwrap_or(now_ms);
263        let banked = (now_ms - suspended_at).max(0);
264        self.accumulated_suspended_ms = self.accumulated_suspended_ms.saturating_add(banked);
265        self.suspended_at_ms = None;
266        // Record the completed pair at EVERY revision and BEFORE either cap
267        // check can return: the vec is history, not a decision input, and a
268        // pause that force-expires the session is still a pause that happened.
269        // (Phase-9 precedent: record everywhere, read only under rev >= 2.)
270        //
271        // The one exception is the rev <= 1 overflow below: a legacy session
272        // is reachable through the un-rate-limited Suspend/Resume RPCs just
273        // like a current one, so its vec has to be bounded too — but it must
274        // not be force-expired by a rule postdating its acceptance. Since
275        // nothing reads the vec below rev 2, dropping the overflow bounds
276        // memory and snapshot size while leaving legacy replay bit-identical.
277        if self.semantics_rev >= 2 || self.suspension_intervals.len() < MAX_SUSPENSION_CYCLES {
278            self.suspension_intervals.push((suspended_at, now_ms));
279        }
280        // Cycle-count cap, rev >= 2 only — see `MAX_SUSPENSION_CYCLES`. Gated
281        // on the revision so rev <= 1 histories replay bit-identically.
282        if self.semantics_rev >= 2 && self.suspension_intervals.len() > MAX_SUSPENSION_CYCLES {
283            self.state = SessionState::Expired;
284            return Err(MacpError::TtlExpired);
285        }
286        if self.accumulated_suspended_ms > self.effective_max_suspend_ms() {
287            self.state = SessionState::Expired;
288            return Err(MacpError::TtlExpired);
289        }
290        self.ttl_expiry = self.ttl_expiry.saturating_add(banked);
291        self.state = SessionState::Open;
292        Ok(())
293    }
294
295    /// Cancel an `Open` or `Suspended` session into the terminal `Cancelled`
296    /// state (RFC-MACP-0001 §7.3). Returns an error if already terminal.
297    pub fn cancel(&mut self) -> Result<(), MacpError> {
298        if self.state.is_terminal() {
299            return Err(MacpError::SessionNotOpen);
300        }
301        self.state = SessionState::Cancelled;
302        self.suspended_at_ms = None;
303        Ok(())
304    }
305
306    /// Whether a currently-`Suspended` session has exceeded `MAX_SUSPEND_MS` as
307    /// of `now_ms` (cumulative banked plus the in-progress suspension).
308    pub fn suspend_cap_exceeded(&self, now_ms: i64) -> bool {
309        match self.suspended_at_ms {
310            Some(at) => {
311                self.accumulated_suspended_ms
312                    .saturating_add((now_ms - at).max(0))
313                    > self.effective_max_suspend_ms()
314            }
315            None => self.accumulated_suspended_ms > self.effective_max_suspend_ms(),
316        }
317    }
318
319    /// The session-timeline instant at which `duration_ms` of **unsuspended**
320    /// time has elapsed since `from_ms` — i.e. the earliest `T` with
321    /// `(T - from_ms) - suspended_in[from_ms, T] >= duration_ms`.
322    ///
323    /// This is RFC-MACP-0010 §5.1(3)'s own formula for the synthetic implicit
324    /// accept's timestamp ("offer acceptance time + timeout + suspended time
325    /// **within the window**"), evaluated on the recorded timeline required by
326    /// §5.1(1). It is deliberately not the naive
327    /// `from_ms + duration_ms + banked_since(from_ms)`: that counts pauses
328    /// that begin *after* the true deadline, so it is wrong whenever a
329    /// suspend/resume pair lands between the deadline and the observation —
330    /// fully reachable, since `SuspendSession`/`ResumeSession` are RPCs that
331    /// need no session-scoped message. It would also make the timestamp
332    /// depend on *when* it was computed, which is unacceptable for a value
333    /// baked into permanent history.
334    ///
335    /// The walk: start at `from_ms` with the full `duration_ms` remaining;
336    /// for each completed pause `(s, e)` starting at or after `from_ms`, if
337    /// the unsuspended run up to `s` already covers what remains, stop inside
338    /// that run; otherwise consume it and jump to `e`. A pause starting
339    /// exactly at the returned deadline does not extend it — the offer's
340    /// unsuspended time had already hit the timeout at that instant.
341    ///
342    /// Pure and saturating: no clock read, no I/O, no panics on overflow.
343    ///
344    /// **Under-count invariant.** The walk's contribution satisfies
345    /// `walk_sum <= accumulated_suspended_ms - offer.suspended_ms_at_offer`:
346    /// [`Session::suspension_intervals`] may *under*-report completed pauses
347    /// but can never over-report them. Three distinct sources produce a short
348    /// or empty vec, each permanent for the life of that log:
349    ///
350    /// 1. A snapshot or checkpoint written *before the field existed*
351    ///    deserializes it as empty (`#[serde(default)]`) while
352    ///    `accumulated_suspended_ms` is already positive.
353    /// 2. Any replay — including a post-11b one — that resumes from a
354    ///    *pre-11b mid-session checkpoint*: the fast path replays only
355    ///    `&log_entries[idx + 1..]`, so every pause that completed before that
356    ///    checkpoint is gone and can never be recovered, no matter how many
357    ///    times the log is replayed afterwards.
358    /// 3. A `semantics_rev <= 1` session that exceeded
359    ///    [`MAX_SUSPENSION_CYCLES`]: recording stops rather than
360    ///    force-expiring, so pairs past the cap are dropped. This one is
361    ///    outside this function's read domain by construction — the walk is
362    ///    only consulted at `semantics_rev >= 2`, where the cap force-expires
363    ///    instead of dropping — so the enumeration above is exhaustive for
364    ///    every vec this function can actually be asked to walk.
365    ///
366    /// An under-count only moves the returned deadline *earlier*, never later
367    /// — the safe direction, so callers may rely on `deadline <= now_ms` once
368    /// the scalar arithmetic has already decided the timeout elapsed.
369    ///
370    /// **Why a short vec cannot break replay determinism.** This describes the
371    /// design this function exists to serve; the synthetic entry itself
372    /// arrives in a later phase. Replay is never to *recompute* a synthetic
373    /// implicit-accept timestamp — it replays the recorded synthetic entry as
374    /// data. Only the live emitter computes a deadline, exactly once, at
375    /// *emission* time, from the offer's recorded `offered_at_ms`. So a short
376    /// vec can make a
377    /// *newly* emitted implicit accept land earlier than a fully-recorded one
378    /// would have, but it can never make a *replayed* one disagree with the
379    /// live value already baked into history — byte-identical replay is not
380    /// foreclosed.
381    pub fn unsuspended_deadline(&self, from_ms: i64, duration_ms: i64) -> i64 {
382        let mut cur = from_ms;
383        let mut remaining = duration_ms;
384        for &(s, e) in self
385            .suspension_intervals
386            .iter()
387            .filter(|(s, _)| *s >= from_ms)
388        {
389            // Clamp both the run and the jump. `Utc::now()` is not monotonic
390            // (an NTP step back between suspend and resume yields `e < s`),
391            // pairs can overlap or nest, and the public
392            // `SessionBuilder::suspension_intervals` setter lets a library
393            // consumer supply anything at all. Without the clamps a negative
394            // `run` would *grow* `remaining` and a backwards `e` would rewind
395            // `cur`, i.e. the walk would OVER-report and violate the
396            // under-count invariant above.
397            let run = s.saturating_sub(cur).max(0);
398            if run >= remaining {
399                return cur.saturating_add(remaining);
400            }
401            remaining = remaining.saturating_sub(run);
402            // `s` is in the max as well as `e`: the run up to `s` was just
403            // consumed, so `cur` must end at or after `s` even when the pair
404            // is backwards. Jumping only to `e` there would spend the run
405            // without advancing the cursor, returning a deadline *before*
406            // `from_ms + duration_ms` — safe against over-reporting, but a
407            // timeout that fires before it nominally elapsed is its own bug.
408            // A degenerate pair is therefore treated as a zero-width pause.
409            cur = e.max(s).max(cur);
410        }
411        cur.saturating_add(remaining)
412    }
413
414    pub fn apply_mode_response(&mut self, response: ModeResponse) {
415        match response {
416            ModeResponse::NoOp => {}
417            ModeResponse::PersistState(state) => self.mode_state = state,
418            ModeResponse::Resolve(resolution) => {
419                self.state = SessionState::Resolved;
420                self.resolution = Some(resolution);
421            }
422            ModeResponse::PersistAndResolve { state, resolution } => {
423                self.mode_state = state;
424                self.state = SessionState::Resolved;
425                self.resolution = Some(resolution);
426            }
427        }
428    }
429}
430
431/// Builder for [`Session`] — the only construction path outside `macp-core`
432/// (the struct is `#[non_exhaustive]`).
433///
434/// Defaults: `state: Open`, `ttl_expiry: i64::MAX` (never expires until set),
435/// numeric fields `0`, everything else empty/`None`.
436#[derive(Clone, Debug)]
437pub struct SessionBuilder {
438    inner: Session,
439}
440
441macro_rules! builder_setters {
442    ($($(#[$doc:meta])* $name:ident: $ty:ty),* $(,)?) => {
443        $(
444            $(#[$doc])*
445            pub fn $name(mut self, value: $ty) -> Self {
446                self.inner.$name = value;
447                self
448            }
449        )*
450    };
451}
452
453impl SessionBuilder {
454    builder_setters! {
455        state: SessionState,
456        ttl_expiry: i64,
457        ttl_ms: i64,
458        started_at_unix_ms: i64,
459        resolution: Option<Vec<u8>>,
460        mode_state: Vec<u8>,
461        participants: Vec<String>,
462        seen_message_ids: HashSet<String>,
463        extensions: HashMap<String, Vec<u8>>,
464        roots: Vec<macp_pb::pb::Root>,
465        participant_message_counts: HashMap<String, u32>,
466        participant_last_seen: HashMap<String, i64>,
467        policy_definition: Option<crate::policy::PolicyDefinition>,
468        suspended_at_ms: Option<i64>,
469        accumulated_suspended_ms: i64,
470        /// Completed suspend/resume pairs (see
471        /// [`Session::suspension_intervals`]). Needed so persistence layers
472        /// can restore the field — `Session` is `#[non_exhaustive]`, so they
473        /// cannot use a struct literal.
474        suspension_intervals: Vec<(i64, i64)>,
475        semantics_rev: u32,
476        /// Suspension cap bound at SessionStart; 0 = use the
477        /// [`MAX_SUSPEND_MS`] default (legacy sessions, library consumers).
478        max_suspend_ms: i64,
479    }
480
481    pub fn intent(mut self, value: impl Into<String>) -> Self {
482        self.inner.intent = value.into();
483        self
484    }
485
486    pub fn mode_version(mut self, value: impl Into<String>) -> Self {
487        self.inner.mode_version = value.into();
488        self
489    }
490
491    pub fn configuration_version(mut self, value: impl Into<String>) -> Self {
492        self.inner.configuration_version = value.into();
493        self
494    }
495
496    pub fn policy_version(mut self, value: impl Into<String>) -> Self {
497        self.inner.policy_version = value.into();
498        self
499    }
500
501    pub fn context_id(mut self, value: impl Into<String>) -> Self {
502        self.inner.context_id = value.into();
503        self
504    }
505
506    pub fn build(self) -> Session {
507        self.inner
508    }
509}
510
511pub fn requires_strict_session_start(mode: &str) -> bool {
512    matches!(
513        mode,
514        "macp.mode.decision.v1"
515            | "macp.mode.proposal.v1"
516            | "macp.mode.task.v1"
517            | "macp.mode.handoff.v1"
518            | "macp.mode.quorum.v1"
519            | "ext.multi_round.v1"
520    )
521}
522
523/// Parse a protobuf-encoded SessionStartPayload from raw bytes.
524pub fn parse_session_start_payload(payload: &[u8]) -> Result<SessionStartPayload, MacpError> {
525    if payload.is_empty() {
526        return Err(MacpError::InvalidPayload);
527    }
528    SessionStartPayload::decode(payload).map_err(|_| MacpError::InvalidPayload)
529}
530
531/// Extract and validate TTL from a parsed SessionStartPayload.
532pub fn extract_ttl_ms(payload: &SessionStartPayload) -> Result<i64, MacpError> {
533    if !(1..=MAX_TTL_MS).contains(&payload.ttl_ms) {
534        return Err(MacpError::InvalidTtl);
535    }
536    Ok(payload.ttl_ms)
537}
538
539/// Modes whose canonical `SessionStart` may bind an **empty** `participants`
540/// list.
541///
542/// Decision alone. RFC-MACP-0001 §7.1 requires `participants` only "when
543/// required by the Mode", and RFC-MACP-0007 makes the initiator's authority
544/// role-based rather than membership-based, so a Decision session with no
545/// declared participants is well-defined: nobody — the initiator included — can
546/// emit a `Proposal`, `Evaluation`, `Objection` or `Vote`, because
547/// `DecisionMode::authorize_sender` routes all four through
548/// `is_declared_participant`, which is `false` over an empty list. Such a
549/// session can therefore only expire or be cancelled. Spec #99 removed
550/// `minItems: 1` from the conformance fixture schema on exactly that reasoning
551/// and added `decision_zero_participants.json` to pin it.
552///
553/// **Deliberately a positive allowlist of one, checked in one place.** The
554/// other four standards-track modes each re-reject an insufficient roster in
555/// their own `on_session_start`, but those are five independent
556/// implementations: if the rule lived only there, deleting any one guard would
557/// silently remove the guarantee with nothing at the core level left to notice.
558/// Keeping the rule here means the exception is named once and every other
559/// mode — including a promoted extension mode the list below has never heard
560/// of — keeps the full canonical contract by default.
561fn allows_empty_participants(mode: &str) -> bool {
562    mode == "macp.mode.decision.v1"
563}
564
565/// Validate the complete canonical SessionStart binding contract.
566///
567/// Mode-independent, and therefore holds the roster non-emptiness rule for
568/// **every** mode. Prefer
569/// [`validate_canonical_session_start_payload_for_mode`] on any path that knows
570/// the mode name; this entry point is retained with its original signature and
571/// its original behaviour.
572pub fn validate_canonical_session_start_payload(
573    payload: &SessionStartPayload,
574) -> Result<(), MacpError> {
575    validate_canonical_start(payload, false)
576}
577
578/// The canonical SessionStart binding contract, with the roster rule scoped to
579/// the mode.
580///
581/// Identical to [`validate_canonical_session_start_payload`] in every respect
582/// except one: an empty `participants` list is accepted for the modes
583/// `allows_empty_participants` names (Decision, and only Decision) and
584/// rejected for all others, promoted extension modes included.
585///
586/// This is additive rather than a new parameter on
587/// [`validate_canonical_session_start_payload`] on purpose — changing that
588/// function's signature would be a `macp-core` API break, and every crate in
589/// this workspace shares one version.
590pub fn validate_canonical_session_start_payload_for_mode(
591    mode: &str,
592    payload: &SessionStartPayload,
593) -> Result<(), MacpError> {
594    validate_canonical_start(payload, allows_empty_participants(mode))
595}
596
597fn validate_canonical_start(
598    payload: &SessionStartPayload,
599    allow_empty_participants: bool,
600) -> Result<(), MacpError> {
601    extract_ttl_ms(payload)?;
602
603    if payload.mode_version.trim().is_empty() || payload.configuration_version.trim().is_empty() {
604        return Err(MacpError::InvalidPayload);
605    }
606
607    if payload.participants.is_empty() && !allow_empty_participants {
608        return Err(MacpError::InvalidPayload);
609    }
610
611    // Safety limit: prevent resource exhaustion from excessively large participant lists.
612    const MAX_PARTICIPANTS: usize = 1000;
613    if payload.participants.len() > MAX_PARTICIPANTS {
614        return Err(MacpError::InvalidPayload);
615    }
616
617    let mut seen = HashSet::new();
618    for participant in &payload.participants {
619        let participant = participant.trim();
620        if participant.is_empty() || !seen.insert(participant.to_string()) {
621            return Err(MacpError::InvalidPayload);
622        }
623    }
624
625    // max_suspend_ms: 0 selects the runtime default; a positive value binds a
626    // session-specific cap (RFC-MACP-0001 §7.1). Negative is meaningless.
627    if payload.max_suspend_ms < 0 {
628        return Err(MacpError::InvalidPayload);
629    }
630
631    Ok(())
632}
633
634/// Enforce the strict SessionStart binding contract for standards-track and qualifying extension modes.
635pub fn validate_strict_session_start_payload(
636    mode: &str,
637    payload: &SessionStartPayload,
638) -> Result<(), MacpError> {
639    if !requires_strict_session_start(mode) {
640        return Ok(());
641    }
642
643    validate_canonical_session_start_payload_for_mode(mode, payload)
644}
645
646/// Validate that a session ID meets the acceptance policy.
647///
648/// Accepts:
649/// - UUID v4/v7 in hyphenated lowercase canonical form (36 chars)
650/// - base64url tokens of 22+ chars (`[A-Za-z0-9_-]`)
651///
652/// Rejects everything else (empty, short human-readable, uppercase UUID, etc.).
653pub fn validate_session_id_for_acceptance(session_id: &str) -> Result<(), MacpError> {
654    if session_id.is_empty() {
655        return Err(MacpError::InvalidSessionId);
656    }
657
658    // UUID-shaped ids (36 chars, parseable) are held to strict UUID rules with no
659    // fall-through: canonical lowercase hyphenated form, version v4 or v7. This
660    // keeps non-canonical forms (e.g. uppercase) of the same UUID from being
661    // admitted as distinct base64url tokens. Only strings that do not parse as a
662    // UUID at all fall through to the base64url rule.
663    if session_id.len() == 36 && session_id.contains('-') {
664        if let Ok(parsed) = uuid::Uuid::parse_str(session_id) {
665            if parsed.as_hyphenated().to_string() == session_id {
666                match parsed.get_version() {
667                    Some(uuid::Version::Random) | Some(uuid::Version::SortRand) => {
668                        return Ok(());
669                    }
670                    _ => {}
671                }
672            }
673            return Err(MacpError::InvalidSessionId);
674        }
675    }
676
677    // Try base64url: at least 22 chars, only [A-Za-z0-9_-]
678    if session_id.len() >= 22
679        && session_id
680            .chars()
681            .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
682    {
683        return Ok(());
684    }
685
686    Err(MacpError::InvalidSessionId)
687}
688
689#[cfg(test)]
690mod tests {
691    use super::*;
692    use prost::Message;
693
694    fn encode_payload(ttl_ms: i64, participants: Vec<String>) -> Vec<u8> {
695        let payload = SessionStartPayload {
696            intent: String::new(),
697            participants,
698            mode_version: "1.0.0".into(),
699            configuration_version: "cfg-1".into(),
700            policy_version: String::new(),
701            ttl_ms,
702            context_id: String::new(),
703            extensions: std::collections::HashMap::new(),
704            roots: vec![],
705            max_suspend_ms: 0,
706        };
707        payload.encode_to_vec()
708    }
709
710    #[test]
711    fn parse_empty_payload_is_invalid() {
712        let err = parse_session_start_payload(b"").unwrap_err();
713        assert_eq!(err.to_string(), "InvalidPayload");
714    }
715
716    #[test]
717    fn parse_valid_protobuf_payload() {
718        let bytes = encode_payload(5000, vec!["alice".into(), "bob".into()]);
719        let result = parse_session_start_payload(&bytes).unwrap();
720        assert_eq!(result.ttl_ms, 5000);
721        assert_eq!(result.participants, vec!["alice", "bob"]);
722    }
723
724    #[test]
725    fn extract_ttl_requires_explicit_positive_value() {
726        let payload = SessionStartPayload::default();
727        assert_eq!(
728            extract_ttl_ms(&payload).unwrap_err().to_string(),
729            "InvalidTtl"
730        );
731
732        let payload = SessionStartPayload {
733            ttl_ms: 5000,
734            ..Default::default()
735        };
736        assert_eq!(extract_ttl_ms(&payload).unwrap(), 5000);
737    }
738
739    #[test]
740    fn standard_mode_requires_explicit_versions() {
741        // Renamed from `standard_mode_requires_explicit_versions_and_participants`
742        // and narrowed: the empty-`participants` half moved out, because
743        // Decision now accepts an empty roster. The version and TTL halves of
744        // the strict contract are kept verbatim so the rest stays pinned; the
745        // roster rule is pinned for every other mode by
746        // `every_standard_mode_except_decision_rejects_an_empty_roster` below.
747        let payload = SessionStartPayload {
748            participants: vec!["alice".into()],
749            mode_version: String::new(),
750            configuration_version: "cfg-1".into(),
751            ttl_ms: 1000,
752            ..Default::default()
753        };
754        assert_eq!(
755            validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
756                .unwrap_err()
757                .to_string(),
758            "InvalidPayload"
759        );
760
761        let payload = SessionStartPayload {
762            participants: vec!["alice".into()],
763            mode_version: "1.0.0".into(),
764            configuration_version: String::new(),
765            ttl_ms: 1000,
766            ..Default::default()
767        };
768        assert_eq!(
769            validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
770                .unwrap_err()
771                .to_string(),
772            "InvalidPayload"
773        );
774
775        let payload = SessionStartPayload {
776            participants: vec!["alice".into()],
777            mode_version: "1.0.0".into(),
778            configuration_version: "cfg-1".into(),
779            ttl_ms: 0,
780            ..Default::default()
781        };
782        assert_eq!(
783            validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
784                .unwrap_err()
785                .to_string(),
786            "InvalidTtl"
787        );
788    }
789
790    /// Build an otherwise-valid strict payload with an empty roster.
791    fn empty_roster_payload() -> SessionStartPayload {
792        SessionStartPayload {
793            participants: vec![],
794            mode_version: "1.0.0".into(),
795            configuration_version: "cfg-1".into(),
796            ttl_ms: 1000,
797            ..Default::default()
798        }
799    }
800
801    #[test]
802    fn decision_accepts_an_empty_participant_list() {
803        assert!(
804            validate_strict_session_start_payload("macp.mode.decision.v1", &empty_roster_payload())
805                .is_ok(),
806            "RFC-MACP-0001 §7.1 requires participants only when the mode does, and \
807             RFC-MACP-0007 makes Decision authority role-based; spec #99's \
808             decision_zero_participants.json pins the accepted SessionStart"
809        );
810        // The relaxation must be the *only* thing that moved: the same payload
811        // with everything else intact still fails on a missing version.
812        let mut broken = empty_roster_payload();
813        broken.mode_version = String::new();
814        assert!(
815            validate_strict_session_start_payload("macp.mode.decision.v1", &broken).is_err(),
816            "an empty roster must not waive the rest of the strict contract"
817        );
818    }
819
820    #[test]
821    fn every_standard_mode_except_decision_rejects_an_empty_roster() {
822        // The structural guard, and the reason the rule lives here rather than
823        // in five `on_session_start` implementations. Iterating the strict-mode
824        // list means a mode *added* to it inherits the roster requirement, and
825        // a future change that widened the carve-out would have to edit this
826        // table to stay green — neither is true of a per-mode guard, which a
827        // PR can delete alongside its own test.
828        for mode in [
829            "macp.mode.proposal.v1",
830            "macp.mode.task.v1",
831            "macp.mode.handoff.v1",
832            "macp.mode.quorum.v1",
833            "ext.multi_round.v1",
834        ] {
835            assert!(
836                requires_strict_session_start(mode),
837                "{mode} must be strict for this table to mean anything"
838            );
839            assert_eq!(
840                validate_strict_session_start_payload(mode, &empty_roster_payload())
841                    .unwrap_err()
842                    .to_string(),
843                "InvalidPayload",
844                "{mode} must still reject an empty participant list at the core layer"
845            );
846        }
847
848        // And a *promoted* extension mode — a name the static carve-out list
849        // has never heard of, which `ModeRegistry::promote_mode` can mark
850        // strict at runtime — keeps the full canonical contract. This is the
851        // trap the two call sites had to avoid: swapping the strictness source
852        // for the core's static list would have dropped canonical validation
853        // for exactly these names.
854        assert_eq!(
855            validate_canonical_session_start_payload_for_mode(
856                "ext.promoted.v1",
857                &empty_roster_payload()
858            )
859            .unwrap_err()
860            .to_string(),
861            "InvalidPayload",
862            "an unrecognised (e.g. promoted) mode must default to the strict roster rule"
863        );
864
865        // The mode-independent entry point keeps its original behaviour, so a
866        // caller that cannot supply a mode name is never silently relaxed.
867        assert_eq!(
868            validate_canonical_session_start_payload(&empty_roster_payload())
869                .unwrap_err()
870                .to_string(),
871            "InvalidPayload"
872        );
873    }
874
875    fn open_session(ttl_expiry: i64) -> Session {
876        Session {
877            session_id: "s1".into(),
878            state: SessionState::Open,
879            ttl_expiry,
880            ttl_ms: 60_000,
881            started_at_unix_ms: 0,
882            resolution: None,
883            mode: "macp.mode.decision.v1".into(),
884            mode_state: vec![],
885            participants: vec![],
886            seen_message_ids: HashSet::new(),
887            intent: String::new(),
888            mode_version: "1.0.0".into(),
889            configuration_version: "cfg-1".into(),
890            policy_version: String::new(),
891            context_id: String::new(),
892            extensions: HashMap::new(),
893            roots: vec![],
894            initiator_sender: "agent://a".into(),
895            participant_message_counts: HashMap::new(),
896            participant_last_seen: HashMap::new(),
897            policy_definition: None,
898            suspended_at_ms: None,
899            accumulated_suspended_ms: 0,
900            suspension_intervals: vec![],
901            semantics_rev: CURRENT_SEMANTICS_REV,
902            max_suspend_ms: 0,
903        }
904    }
905
906    #[test]
907    fn suspend_then_resume_banks_ttl() {
908        let mut s = open_session(10_000);
909        s.suspend(2_000).unwrap();
910        assert_eq!(s.state, SessionState::Suspended);
911        assert_eq!(s.suspended_at_ms, Some(2_000));
912        // Resume 3_000ms later: banked 3_000 is added to the deadline.
913        s.resume(5_000).unwrap();
914        assert_eq!(s.state, SessionState::Open);
915        assert_eq!(s.ttl_expiry, 13_000);
916        assert_eq!(s.accumulated_suspended_ms, 3_000);
917        assert_eq!(s.suspended_at_ms, None);
918    }
919
920    #[test]
921    fn suspend_requires_open_and_resume_requires_suspended() {
922        let mut s = open_session(10_000);
923        // resume on an Open session is rejected
924        assert!(matches!(
925            s.resume(1).unwrap_err(),
926            MacpError::SessionNotOpen
927        ));
928        s.suspend(1).unwrap();
929        // double-suspend rejected
930        assert!(matches!(
931            s.suspend(2).unwrap_err(),
932            MacpError::SessionNotOpen
933        ));
934    }
935
936    #[test]
937    fn resume_exceeding_max_suspend_expires() {
938        let mut s = open_session(10_000);
939        s.suspend(0).unwrap();
940        // Resume after more than MAX_SUSPEND_MS: force-expired.
941        let err = s.resume(MAX_SUSPEND_MS + 1).unwrap_err();
942        assert!(matches!(err, MacpError::TtlExpired));
943        assert_eq!(s.state, SessionState::Expired);
944    }
945
946    /// A session-bound cap (SessionStartPayload.max_suspend_ms) overrides the
947    /// default: resuming past the BOUND cap force-expires even though the
948    /// default cap is nowhere near exceeded.
949    #[test]
950    fn bound_cap_overrides_default_on_resume() {
951        let mut s = open_session(10_000);
952        s.max_suspend_ms = 500;
953        s.suspend(0).unwrap();
954        let err = s.resume(501).unwrap_err();
955        assert!(matches!(err, MacpError::TtlExpired));
956        assert_eq!(s.state, SessionState::Expired);
957    }
958
959    #[test]
960    fn bound_cap_within_limit_resumes_and_banks_ttl() {
961        let mut s = open_session(10_000);
962        s.max_suspend_ms = 500;
963        s.suspend(0).unwrap();
964        s.resume(400).unwrap();
965        assert_eq!(s.state, SessionState::Open);
966        assert_eq!(s.ttl_expiry, 10_400);
967    }
968
969    /// The cap is cumulative across pauses, not per-pause: two 300ms
970    /// suspensions each fit under a 500ms bound cap on their own, but the
971    /// second resume sees the 600ms total and force-expires. Pins that
972    /// `resume` accumulates `accumulated_suspended_ms` rather than
973    /// overwriting it with the latest pause — the invariant the rev-2 handoff
974    /// deadline reads (`HandoffMode::rev2_elapsed_ms`).
975    #[test]
976    fn bound_cap_counts_suspension_cumulatively_across_pauses() {
977        let mut s = open_session(10_000);
978        s.max_suspend_ms = 500;
979        s.suspend(0).unwrap();
980        s.resume(300).unwrap();
981        assert_eq!(s.accumulated_suspended_ms, 300);
982        assert_eq!(s.ttl_expiry, 10_300);
983        s.suspend(400).unwrap();
984        // 300 + 300 = 600 > the 500ms cap, though neither pause alone is.
985        let err = s.resume(700).unwrap_err();
986        assert!(matches!(err, MacpError::TtlExpired));
987        assert_eq!(s.state, SessionState::Expired);
988        assert_eq!(s.accumulated_suspended_ms, 600);
989    }
990
991    #[test]
992    fn suspend_cap_exceeded_uses_bound_cap() {
993        let mut s = open_session(10_000);
994        s.max_suspend_ms = 500;
995        s.suspend(0).unwrap();
996        assert!(!s.suspend_cap_exceeded(400));
997        assert!(s.suspend_cap_exceeded(501));
998    }
999
1000    #[test]
1001    fn unbound_session_uses_default_cap() {
1002        let s = open_session(10_000);
1003        assert_eq!(s.max_suspend_ms, 0);
1004        assert_eq!(s.effective_max_suspend_ms(), MAX_SUSPEND_MS);
1005    }
1006
1007    #[test]
1008    fn negative_max_suspend_ms_rejected_in_canonical_payload() {
1009        let payload = SessionStartPayload {
1010            participants: vec!["a".into()],
1011            mode_version: "1.0.0".into(),
1012            configuration_version: "cfg-1".into(),
1013            ttl_ms: 60_000,
1014            max_suspend_ms: -1,
1015            ..Default::default()
1016        };
1017        assert_eq!(
1018            validate_canonical_session_start_payload(&payload)
1019                .unwrap_err()
1020                .to_string(),
1021            "InvalidPayload"
1022        );
1023        // 0 (runtime default) and positive values are both valid.
1024        let ok0 = SessionStartPayload {
1025            max_suspend_ms: 0,
1026            ..payload.clone()
1027        };
1028        validate_canonical_session_start_payload(&ok0).unwrap();
1029        let ok_pos = SessionStartPayload {
1030            max_suspend_ms: 60_000,
1031            ..payload
1032        };
1033        validate_canonical_session_start_payload(&ok_pos).unwrap();
1034    }
1035
1036    #[test]
1037    fn cancel_from_open_or_suspended_then_terminal_is_rejected() {
1038        let mut s = open_session(10_000);
1039        s.suspend(1).unwrap();
1040        s.cancel().unwrap();
1041        assert_eq!(s.state, SessionState::Cancelled);
1042        assert_eq!(s.suspended_at_ms, None);
1043        // Already terminal: further cancel is rejected.
1044        assert!(matches!(s.cancel().unwrap_err(), MacpError::SessionNotOpen));
1045
1046        let mut open = open_session(10_000);
1047        open.cancel().unwrap();
1048        assert_eq!(open.state, SessionState::Cancelled);
1049    }
1050
1051    #[test]
1052    fn standard_mode_rejects_duplicate_participants() {
1053        let payload = SessionStartPayload {
1054            participants: vec!["alice".into(), "alice".into()],
1055            mode_version: "1.0.0".into(),
1056            configuration_version: "cfg-1".into(),
1057            ttl_ms: 1000,
1058            ..Default::default()
1059        };
1060        assert_eq!(
1061            validate_strict_session_start_payload("macp.mode.proposal.v1", &payload)
1062                .unwrap_err()
1063                .to_string(),
1064            "InvalidPayload"
1065        );
1066    }
1067
1068    #[test]
1069    fn multi_round_requires_strict_session_start() {
1070        let payload = SessionStartPayload::default();
1071        assert!(validate_strict_session_start_payload("ext.multi_round.v1", &payload).is_err());
1072    }
1073
1074    #[test]
1075    fn valid_uuid_v4_accepted() {
1076        let id = uuid::Uuid::new_v4().as_hyphenated().to_string();
1077        validate_session_id_for_acceptance(&id).unwrap();
1078    }
1079
1080    #[test]
1081    fn valid_base64url_accepted() {
1082        // 22-char base64url token
1083        validate_session_id_for_acceptance("abcdefghijklmnopqrstuv").unwrap();
1084        // longer base64url with underscore and hyphen
1085        validate_session_id_for_acceptance("abc-def_ghi-jkl_mno-pqr").unwrap();
1086    }
1087
1088    #[test]
1089    fn empty_id_rejected() {
1090        assert_eq!(
1091            validate_session_id_for_acceptance("")
1092                .unwrap_err()
1093                .to_string(),
1094            "InvalidSessionId"
1095        );
1096    }
1097
1098    #[test]
1099    fn short_weak_id_rejected() {
1100        assert_eq!(
1101            validate_session_id_for_acceptance("s1")
1102                .unwrap_err()
1103                .to_string(),
1104            "InvalidSessionId"
1105        );
1106        assert_eq!(
1107            validate_session_id_for_acceptance("decision-demo-1")
1108                .unwrap_err()
1109                .to_string(),
1110            "InvalidSessionId"
1111        );
1112    }
1113
1114    #[test]
1115    fn uppercase_uuid_rejected() {
1116        let id = uuid::Uuid::new_v4()
1117            .as_hyphenated()
1118            .to_string()
1119            .to_uppercase();
1120        assert_eq!(
1121            validate_session_id_for_acceptance(&id)
1122                .unwrap_err()
1123                .to_string(),
1124            "InvalidSessionId"
1125        );
1126    }
1127
1128    #[test]
1129    fn base64url_36_chars_with_hyphen_accepted() {
1130        // 36-char base64url token containing '-' that is NOT UUID-shaped: must be
1131        // accepted via the base64url rule, not rejected by the UUID branch.
1132        // (Regression test for the hard-routing bug: len==36 && contains('-')
1133        // previously returned Err without trying the base64url rule.)
1134        let id = "Zx-abcdefghijklmnopqrstuvwxyz_ABCDE-";
1135        assert_eq!(id.len(), 36);
1136        assert!(uuid::Uuid::parse_str(id).is_err());
1137        validate_session_id_for_acceptance(id).unwrap();
1138    }
1139
1140    #[test]
1141    fn uuid_shaped_but_wrong_version_does_not_fall_through() {
1142        // A canonical v1 UUID is also 36 chars of valid base64url charset; it must
1143        // still be rejected (UUID rules apply, no fall-through to base64url).
1144        let v4 = uuid::Uuid::new_v4();
1145        let mut bytes = *v4.as_bytes();
1146        bytes[6] = (bytes[6] & 0x0F) | 0x10;
1147        bytes[8] = (bytes[8] & 0x3F) | 0x80;
1148        let v1_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
1149        assert!(validate_session_id_for_acceptance(&v1_id).is_err());
1150    }
1151
1152    #[test]
1153    fn base64url_too_short_rejected() {
1154        assert_eq!(
1155            validate_session_id_for_acceptance("abcdefghij")
1156                .unwrap_err()
1157                .to_string(),
1158            "InvalidSessionId"
1159        );
1160    }
1161
1162    #[test]
1163    fn valid_uuid_v7_accepted() {
1164        // Construct a v7 UUID by patching the version nibble of a v4 UUID
1165        let v4 = uuid::Uuid::new_v4();
1166        let mut bytes = *v4.as_bytes();
1167        // Set version nibble (bits 48-51) to 0b0111 (v7)
1168        bytes[6] = (bytes[6] & 0x0F) | 0x70;
1169        // Keep variant bits valid (RFC 4122: 0b10xx)
1170        bytes[8] = (bytes[8] & 0x3F) | 0x80;
1171        let v7_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
1172        assert!(validate_session_id_for_acceptance(&v7_id).is_ok());
1173    }
1174
1175    #[test]
1176    fn uuid_v1_rejected() {
1177        // Construct a v1 UUID by patching the version nibble of a v4 UUID
1178        let v4 = uuid::Uuid::new_v4();
1179        let mut bytes = *v4.as_bytes();
1180        // Set version nibble (bits 48-51) to 0b0001 (v1)
1181        bytes[6] = (bytes[6] & 0x0F) | 0x10;
1182        // Keep variant bits valid (RFC 4122: 0b10xx)
1183        bytes[8] = (bytes[8] & 0x3F) | 0x80;
1184        let v1_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
1185        assert_eq!(
1186            validate_session_id_for_acceptance(&v1_id)
1187                .unwrap_err()
1188                .to_string(),
1189            "InvalidSessionId"
1190        );
1191    }
1192
1193    #[test]
1194    fn too_many_participants_rejected() {
1195        let participants: Vec<String> = (0..1001).map(|i| format!("agent://p{i}")).collect();
1196        let bytes = encode_payload(5000, participants);
1197        let payload = parse_session_start_payload(&bytes).unwrap();
1198        assert_eq!(
1199            validate_canonical_session_start_payload(&payload)
1200                .unwrap_err()
1201                .to_string(),
1202            "InvalidPayload"
1203        );
1204    }
1205
1206    #[test]
1207    fn max_participants_accepted() {
1208        let participants: Vec<String> = (0..1000).map(|i| format!("agent://p{i}")).collect();
1209        let bytes = encode_payload(5000, participants);
1210        let payload = parse_session_start_payload(&bytes).unwrap();
1211        validate_canonical_session_start_payload(&payload).unwrap();
1212    }
1213
1214    // ---- Phase 11b: suspension intervals + the unsuspended deadline ----
1215
1216    /// `resume` records the completed pause, in order, at every revision.
1217    #[test]
1218    fn resume_records_completed_suspension_intervals() {
1219        let mut s = open_session(100_000);
1220        assert!(s.suspension_intervals.is_empty());
1221        s.suspend(2_000).unwrap();
1222        // An in-progress suspension is NOT in the vec — completed pairs only.
1223        assert!(s.suspension_intervals.is_empty());
1224        s.resume(5_000).unwrap();
1225        s.suspend(6_000).unwrap();
1226        s.resume(6_500).unwrap();
1227        assert_eq!(s.suspension_intervals, vec![(2_000, 5_000), (6_000, 6_500)]);
1228
1229        // Recorded at legacy revisions too (read only at rev >= 2).
1230        let mut legacy = open_session(100_000);
1231        legacy.semantics_rev = 0;
1232        legacy.suspend(1_000).unwrap();
1233        legacy.resume(1_400).unwrap();
1234        assert_eq!(legacy.suspension_intervals, vec![(1_000, 1_400)]);
1235    }
1236
1237    /// The pause is recorded even when the resume force-expires the session:
1238    /// the push happens before either cap check can return.
1239    #[test]
1240    fn resume_records_the_pause_even_when_it_force_expires() {
1241        let mut s = open_session(10_000);
1242        s.suspend(0).unwrap();
1243        assert!(s.resume(MAX_SUSPEND_MS + 1).is_err());
1244        assert_eq!(s.state, SessionState::Expired);
1245        assert_eq!(s.suspension_intervals, vec![(0, MAX_SUSPEND_MS + 1)]);
1246    }
1247
1248    /// Acceptance criterion 1 — the `unsuspended_deadline` matrix
1249    /// (RFC-MACP-0010 §5.1(3): "offer acceptance time + timeout + suspended
1250    /// time *within the window*").
1251    #[test]
1252    fn unsuspended_deadline_walks_the_suspension_intervals() {
1253        let base = open_session(1_000_000);
1254        let with = |pairs: Vec<(i64, i64)>| {
1255            let mut s = base.clone();
1256            s.suspension_intervals = pairs;
1257            s
1258        };
1259
1260        // No pauses: the raw deadline.
1261        assert_eq!(with(vec![]).unsuspended_deadline(1_000, 100), 1_100);
1262
1263        // One pause that starts inside the window: extends by its full width.
1264        assert_eq!(
1265            with(vec![(1_050, 1_200)]).unsuspended_deadline(1_000, 100),
1266            1_250
1267        );
1268
1269        // One pause starting after the raw deadline: does NOT extend it.
1270        assert_eq!(
1271            with(vec![(1_500, 1_600)]).unsuspended_deadline(1_000, 100),
1272            1_100
1273        );
1274
1275        // Boundary: a pause starting exactly at the returned deadline does not
1276        // extend it — the timeout had already elapsed at that instant.
1277        assert_eq!(
1278            with(vec![(1_100, 1_300)]).unsuspended_deadline(1_000, 100),
1279            1_100
1280        );
1281
1282        // Two pauses, both inside the window: both are added.
1283        assert_eq!(
1284            with(vec![(1_050, 1_200), (1_230, 1_300)]).unsuspended_deadline(1_000, 100),
1285            1_320
1286        );
1287
1288        // A pause predating `from_ms` is ignored entirely.
1289        assert_eq!(
1290            with(vec![(500, 700)]).unsuspended_deadline(1_000, 100),
1291            1_100
1292        );
1293        assert_eq!(
1294            with(vec![(500, 700), (1_050, 1_200)]).unsuspended_deadline(1_000, 100),
1295            1_250
1296        );
1297    }
1298
1299    /// The walk must never OVER-report, whatever shape the vec is in. The
1300    /// pairs are not guaranteed sorted, disjoint, or even monotone: `Utc::now`
1301    /// is not monotonic (an NTP step back between suspend and resume yields
1302    /// `e < s`, which `resume` records as-is — it clamps only `banked`), and
1303    /// the public `SessionBuilder::suspension_intervals` setter lets a library
1304    /// consumer supply arbitrary pairs. Each case below drove `remaining`
1305    /// upward or `cur` backwards before the clamps in `unsuspended_deadline`.
1306    #[test]
1307    fn unsuspended_deadline_never_over_reports_on_adversarial_pairs() {
1308        let base = open_session(1_000_000);
1309        let with = |pairs: Vec<(i64, i64)>| {
1310            let mut s = base.clone();
1311            s.suspension_intervals = pairs;
1312            s
1313        };
1314
1315        // Overlapping pairs, second starting inside the first. The 100 ms of
1316        // unsuspended time is exhausted at 1_050 + the union of the pauses
1317        // (1_050..1_400) + the remaining 50 ms => 1_450. Before the clamp the
1318        // negative run (1_100 - 1_400) GREW `remaining` to 350 and returned
1319        // 1_550 — a deadline 100 ms LATER than the truth.
1320        assert_eq!(
1321            with(vec![(1_050, 1_400), (1_100, 1_200)]).unsuspended_deadline(1_000, 100),
1322            1_450
1323        );
1324
1325        // A fully-nested pair: the inner pause contributes nothing.
1326        assert_eq!(
1327            with(vec![(1_050, 1_400), (1_200, 1_300)]).unsuspended_deadline(1_000, 100),
1328            1_450
1329        );
1330
1331        // A backwards pair (e < s), as an NTP step back would record it. It
1332        // must neither rewind `cur` nor inflate `remaining`; the 50 ms of run
1333        // before it is consumed and nothing is added.
1334        assert_eq!(
1335            with(vec![(1_050, 1_000)]).unsuspended_deadline(1_000, 100),
1336            1_100
1337        );
1338        // ... and the same pair ahead of a real one extends by the real
1339        // pause's width only (60 ms of run, then the 1_060..1_200 pause, then
1340        // the remaining 40 ms), the degenerate pair counting as zero-width.
1341        assert_eq!(
1342            with(vec![(1_050, 1_000), (1_060, 1_200)]).unsuspended_deadline(1_000, 100),
1343            1_240
1344        );
1345
1346        // Unsorted pairs: the later pause listed first. The walk's answer must
1347        // still not exceed the sorted-order answer.
1348        let sorted = with(vec![(1_050, 1_200), (1_230, 1_300)]).unsuspended_deadline(1_000, 100);
1349        let unsorted = with(vec![(1_230, 1_300), (1_050, 1_200)]).unsuspended_deadline(1_000, 100);
1350        assert_eq!(sorted, 1_320);
1351        assert!(
1352            unsorted <= sorted,
1353            "an unsorted vec must under-report, never over-report \
1354             ({unsorted} vs {sorted})"
1355        );
1356
1357        // The invariant stated in the rustdoc, checked directly: the walk's
1358        // own contribution never exceeds the sum of the normalized pair widths (an upper bound on their union).
1359        for pairs in [
1360            vec![(1_050, 1_400), (1_100, 1_200)],
1361            vec![(1_050, 1_400), (1_200, 1_300)],
1362            vec![(1_050, 1_000)],
1363            vec![(1_230, 1_300), (1_050, 1_200)],
1364        ] {
1365            let union: i64 = pairs.iter().map(|(s, e)| (e - s).max(0)).sum();
1366            let walk = with(pairs.clone()).unsuspended_deadline(1_000, 100) - 1_100;
1367            assert!(
1368                (0..=union).contains(&walk),
1369                "walk contribution {walk} outside [0, {union}] for {pairs:?}"
1370            );
1371        }
1372    }
1373
1374    /// Acceptance criterion 5 — the cycle cap fires on *count*, at rev 2 only.
1375    ///
1376    /// Every pause here is 0 ms wide, so `accumulated_suspended_ms` stays at 0
1377    /// and `MAX_SUSPEND_MS` is nowhere near exhausted: only the count cap can
1378    /// be what expires the session.
1379    #[test]
1380    fn suspension_cycle_cap_force_expires_at_rev2() {
1381        let mut s = open_session(1_000_000_000);
1382        assert_eq!(s.semantics_rev, CURRENT_SEMANTICS_REV);
1383        assert!(CURRENT_SEMANTICS_REV >= 2);
1384        for i in 0..MAX_SUSPENSION_CYCLES as i64 {
1385            s.suspend(i).unwrap();
1386            s.resume(i).unwrap();
1387        }
1388        assert_eq!(s.suspension_intervals.len(), MAX_SUSPENSION_CYCLES);
1389        assert_eq!(s.accumulated_suspended_ms, 0);
1390        assert_eq!(s.state, SessionState::Open);
1391
1392        // One cycle past the cap force-expires.
1393        s.suspend(MAX_SUSPENSION_CYCLES as i64).unwrap();
1394        let err = s.resume(MAX_SUSPENSION_CYCLES as i64).unwrap_err();
1395        assert!(matches!(err, MacpError::TtlExpired));
1396        assert_eq!(s.state, SessionState::Expired);
1397        assert_eq!(
1398            s.accumulated_suspended_ms, 0,
1399            "the duration cap must be nowhere near exhausted, else the count \
1400             cap is not what fired"
1401        );
1402
1403        // The rev gate: the identical sequence on a rev-1 session keeps
1404        // succeeding, so legacy histories replay bit-identically.
1405        let mut legacy = open_session(1_000_000_000);
1406        legacy.semantics_rev = 1;
1407        for i in 0..(MAX_SUSPENSION_CYCLES as i64 + 10) {
1408            legacy.suspend(i).unwrap();
1409            legacy.resume(i).unwrap();
1410        }
1411        assert_eq!(legacy.state, SessionState::Open);
1412    }
1413
1414    /// The other half of the cap: a legacy (rev <= 1) session is reachable
1415    /// through the same un-rate-limited `SuspendSession`/`ResumeSession` RPCs
1416    /// as a current one, so its `suspension_intervals` must be bounded too —
1417    /// but it must NOT be force-expired by a rule postdating its acceptance.
1418    /// So the cap stops the *recording* instead: the session keeps cycling and
1419    /// the vec stays pinned at [`MAX_SUSPENSION_CYCLES`].
1420    #[test]
1421    fn suspension_cycle_cap_stops_recording_at_rev1_without_expiring() {
1422        for rev in [0u32, 1] {
1423            let mut s = open_session(1_000_000_000);
1424            s.semantics_rev = rev;
1425            for i in 0..(MAX_SUSPENSION_CYCLES as i64 * 2) {
1426                s.suspend(i).unwrap();
1427                assert_eq!(s.state, SessionState::Suspended, "rev {rev}, cycle {i}");
1428                s.resume(i)
1429                    .unwrap_or_else(|e| panic!("rev {rev}, cycle {i} must resume, got {e:?}"));
1430                assert_eq!(s.state, SessionState::Open, "rev {rev}, cycle {i}");
1431            }
1432            assert_eq!(
1433                s.suspension_intervals.len(),
1434                MAX_SUSPENSION_CYCLES,
1435                "rev {rev}: the vec must stay pinned at the cap, not grow \
1436                 unbounded — an unbounded vec is the O(N²) snapshot \
1437                 amplification MAX_SUSPENSION_CYCLES exists to close"
1438            );
1439            // The recorded prefix is the FIRST `MAX_SUSPENSION_CYCLES` pauses:
1440            // overflow is dropped, not rotated, so nothing already written
1441            // ever changes.
1442            assert_eq!(s.suspension_intervals[0], (0, 0));
1443            assert_eq!(
1444                s.suspension_intervals[MAX_SUSPENSION_CYCLES - 1],
1445                (
1446                    MAX_SUSPENSION_CYCLES as i64 - 1,
1447                    MAX_SUSPENSION_CYCLES as i64 - 1
1448                )
1449            );
1450        }
1451    }
1452
1453    /// Acceptance criterion 6 — a rev-2 session carrying the pre-11b artifact
1454    /// shape (positive `accumulated_suspended_ms`, empty
1455    /// `suspension_intervals`, e.g. a snapshot or mid-session checkpoint
1456    /// written before the field existed) walks to a deadline at or *earlier*
1457    /// than the fully-recorded one. Under-counting is the safe direction.
1458    #[test]
1459    fn legacy_rev2_snapshot_without_intervals_walks_early_not_late() {
1460        let mut recorded = open_session(1_000_000);
1461        recorded.semantics_rev = 2;
1462        recorded.accumulated_suspended_ms = 150 + 70;
1463        recorded.suspension_intervals = vec![(1_050, 1_200), (1_230, 1_300)];
1464
1465        let mut legacy = recorded.clone();
1466        legacy.suspension_intervals.clear();
1467
1468        let recorded_deadline = recorded.unsuspended_deadline(1_000, 100);
1469        let legacy_deadline = legacy.unsuspended_deadline(1_000, 100);
1470        assert_eq!(recorded_deadline, 1_320);
1471        assert_eq!(legacy_deadline, 1_100);
1472        assert!(
1473            legacy_deadline <= recorded_deadline,
1474            "an under-reported interval vec must move the deadline EARLIER \
1475             ({legacy_deadline} vs {recorded_deadline}), never later"
1476        );
1477
1478        // The invariant the walk relies on: the walk's own contribution never
1479        // exceeds the scalar it is refining.
1480        let walk_sum = recorded_deadline - (1_000 + 100);
1481        assert!(walk_sum <= recorded.accumulated_suspended_ms);
1482        let legacy_walk_sum = legacy_deadline - (1_000 + 100);
1483        assert!(legacy_walk_sum <= legacy.accumulated_suspended_ms);
1484    }
1485}