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