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