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