use crate::error::MacpError;
use crate::mode::ModeResponse;
use crate::policy::PolicyDefinition;
use macp_pb::pb::SessionStartPayload;
use prost::Message;
use std::collections::{HashMap, HashSet};
#[doc(hidden)]
pub const DEFAULT_MODE_VERSION: &str = "1.0.0";
#[doc(hidden)]
pub const DEFAULT_CONFIGURATION_VERSION: &str = "config.default";
pub const MAX_TTL_MS: i64 = 24 * 60 * 60 * 1000;
pub const MAX_SUSPEND_MS: i64 = 7 * 24 * 60 * 60 * 1000;
pub const MAX_SUSPENSION_CYCLES: usize = 1024;
pub const CURRENT_SEMANTICS_REV: u32 = 3;
#[derive(Clone, Debug, PartialEq, serde::Serialize, serde::Deserialize)]
pub enum SessionState {
Open,
Suspended,
Resolved,
Expired,
Cancelled,
}
impl SessionState {
pub fn is_terminal(&self) -> bool {
matches!(
self,
SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
)
}
}
#[non_exhaustive]
#[derive(Clone, Debug)]
pub struct Session {
pub session_id: String,
pub state: SessionState,
pub ttl_expiry: i64,
pub ttl_ms: i64,
pub started_at_unix_ms: i64,
pub resolution: Option<Vec<u8>>,
pub mode: String,
pub mode_state: Vec<u8>,
pub participants: Vec<String>,
pub seen_message_ids: HashSet<String>,
pub intent: String,
pub mode_version: String,
pub configuration_version: String,
pub policy_version: String,
pub context_id: String,
pub extensions: HashMap<String, Vec<u8>>,
pub roots: Vec<macp_pb::pb::Root>,
pub initiator_sender: String,
pub participant_message_counts: HashMap<String, u32>,
pub participant_last_seen: HashMap<String, i64>,
pub policy_definition: Option<PolicyDefinition>,
pub suspended_at_ms: Option<i64>,
pub accumulated_suspended_ms: i64,
pub suspension_intervals: Vec<(i64, i64)>,
pub semantics_rev: u32,
pub max_suspend_ms: i64,
}
impl Session {
pub fn builder(
session_id: impl Into<String>,
mode: impl Into<String>,
initiator_sender: impl Into<String>,
) -> SessionBuilder {
SessionBuilder {
inner: Session {
session_id: session_id.into(),
state: SessionState::Open,
ttl_expiry: i64::MAX,
ttl_ms: 0,
started_at_unix_ms: 0,
resolution: None,
mode: mode.into(),
mode_state: vec![],
participants: vec![],
seen_message_ids: HashSet::new(),
intent: String::new(),
mode_version: String::new(),
configuration_version: String::new(),
policy_version: String::new(),
context_id: String::new(),
extensions: HashMap::new(),
roots: vec![],
initiator_sender: initiator_sender.into(),
participant_message_counts: HashMap::new(),
participant_last_seen: HashMap::new(),
policy_definition: None,
suspended_at_ms: None,
accumulated_suspended_ms: 0,
suspension_intervals: vec![],
semantics_rev: CURRENT_SEMANTICS_REV,
max_suspend_ms: 0,
},
}
}
pub fn record_participant_activity(&mut self, sender: &str, timestamp_ms: i64) {
*self
.participant_message_counts
.entry(sender.to_string())
.or_insert(0) += 1;
self.participant_last_seen
.insert(sender.to_string(), timestamp_ms);
}
pub fn suspend(&mut self, now_ms: i64) -> Result<(), MacpError> {
if self.state != SessionState::Open {
return Err(MacpError::SessionNotOpen);
}
self.state = SessionState::Suspended;
self.suspended_at_ms = Some(now_ms);
Ok(())
}
pub fn effective_max_suspend_ms(&self) -> i64 {
if self.max_suspend_ms > 0 {
self.max_suspend_ms
} else {
MAX_SUSPEND_MS
}
}
pub fn resume(&mut self, now_ms: i64) -> Result<(), MacpError> {
if self.state != SessionState::Suspended {
return Err(MacpError::SessionNotOpen);
}
let suspended_at = self.suspended_at_ms.unwrap_or(now_ms);
let banked = (now_ms - suspended_at).max(0);
self.accumulated_suspended_ms = self.accumulated_suspended_ms.saturating_add(banked);
self.suspended_at_ms = None;
if self.semantics_rev >= 2 || self.suspension_intervals.len() < MAX_SUSPENSION_CYCLES {
self.suspension_intervals.push((suspended_at, now_ms));
}
if self.semantics_rev >= 2 && self.suspension_intervals.len() > MAX_SUSPENSION_CYCLES {
self.state = SessionState::Expired;
return Err(MacpError::TtlExpired);
}
if self.accumulated_suspended_ms > self.effective_max_suspend_ms() {
self.state = SessionState::Expired;
return Err(MacpError::TtlExpired);
}
self.ttl_expiry = self.ttl_expiry.saturating_add(banked);
self.state = SessionState::Open;
Ok(())
}
pub fn cancel(&mut self) -> Result<(), MacpError> {
if self.state.is_terminal() {
return Err(MacpError::SessionNotOpen);
}
self.state = SessionState::Cancelled;
self.suspended_at_ms = None;
Ok(())
}
pub fn suspend_cap_exceeded(&self, now_ms: i64) -> bool {
match self.suspended_at_ms {
Some(at) => {
self.accumulated_suspended_ms
.saturating_add((now_ms - at).max(0))
> self.effective_max_suspend_ms()
}
None => self.accumulated_suspended_ms > self.effective_max_suspend_ms(),
}
}
pub fn unsuspended_deadline(&self, from_ms: i64, duration_ms: i64) -> i64 {
let mut cur = from_ms;
let mut remaining = duration_ms;
for &(s, e) in self
.suspension_intervals
.iter()
.filter(|(s, _)| *s >= from_ms)
{
let run = s.saturating_sub(cur).max(0);
if run >= remaining {
return cur.saturating_add(remaining);
}
remaining = remaining.saturating_sub(run);
cur = e.max(s).max(cur);
}
cur.saturating_add(remaining)
}
pub fn apply_mode_response(&mut self, response: ModeResponse) {
match response {
ModeResponse::NoOp => {}
ModeResponse::PersistState(state) => self.mode_state = state,
ModeResponse::Resolve(resolution) => {
self.state = SessionState::Resolved;
self.resolution = Some(resolution);
}
ModeResponse::PersistAndResolve { state, resolution } => {
self.mode_state = state;
self.state = SessionState::Resolved;
self.resolution = Some(resolution);
}
}
}
}
#[derive(Clone, Debug)]
pub struct SessionBuilder {
inner: Session,
}
macro_rules! builder_setters {
($($(#[$doc:meta])* $name:ident: $ty:ty),* $(,)?) => {
$(
$(#[$doc])*
pub fn $name(mut self, value: $ty) -> Self {
self.inner.$name = value;
self
}
)*
};
}
impl SessionBuilder {
builder_setters! {
state: SessionState,
ttl_expiry: i64,
ttl_ms: i64,
started_at_unix_ms: i64,
resolution: Option<Vec<u8>>,
mode_state: Vec<u8>,
participants: Vec<String>,
seen_message_ids: HashSet<String>,
extensions: HashMap<String, Vec<u8>>,
roots: Vec<macp_pb::pb::Root>,
participant_message_counts: HashMap<String, u32>,
participant_last_seen: HashMap<String, i64>,
policy_definition: Option<crate::policy::PolicyDefinition>,
suspended_at_ms: Option<i64>,
accumulated_suspended_ms: i64,
suspension_intervals: Vec<(i64, i64)>,
semantics_rev: u32,
max_suspend_ms: i64,
}
pub fn intent(mut self, value: impl Into<String>) -> Self {
self.inner.intent = value.into();
self
}
pub fn mode_version(mut self, value: impl Into<String>) -> Self {
self.inner.mode_version = value.into();
self
}
pub fn configuration_version(mut self, value: impl Into<String>) -> Self {
self.inner.configuration_version = value.into();
self
}
pub fn policy_version(mut self, value: impl Into<String>) -> Self {
self.inner.policy_version = value.into();
self
}
pub fn context_id(mut self, value: impl Into<String>) -> Self {
self.inner.context_id = value.into();
self
}
pub fn build(self) -> Session {
self.inner
}
}
pub fn requires_strict_session_start(mode: &str) -> bool {
matches!(
mode,
"macp.mode.decision.v1"
| "macp.mode.proposal.v1"
| "macp.mode.task.v1"
| "macp.mode.handoff.v1"
| "macp.mode.quorum.v1"
| "ext.multi_round.v1"
)
}
pub fn parse_session_start_payload(payload: &[u8]) -> Result<SessionStartPayload, MacpError> {
if payload.is_empty() {
return Err(MacpError::InvalidPayload);
}
SessionStartPayload::decode(payload).map_err(|_| MacpError::InvalidPayload)
}
pub fn extract_ttl_ms(payload: &SessionStartPayload) -> Result<i64, MacpError> {
if !(1..=MAX_TTL_MS).contains(&payload.ttl_ms) {
return Err(MacpError::InvalidTtl);
}
Ok(payload.ttl_ms)
}
fn allows_empty_participants(mode: &str) -> bool {
mode == "macp.mode.decision.v1"
}
pub fn validate_canonical_session_start_payload(
payload: &SessionStartPayload,
) -> Result<(), MacpError> {
validate_canonical_start(payload, false)
}
pub fn validate_canonical_session_start_payload_for_mode(
mode: &str,
payload: &SessionStartPayload,
) -> Result<(), MacpError> {
validate_canonical_start(payload, allows_empty_participants(mode))
}
fn validate_canonical_start(
payload: &SessionStartPayload,
allow_empty_participants: bool,
) -> Result<(), MacpError> {
extract_ttl_ms(payload)?;
if payload.mode_version.trim().is_empty() || payload.configuration_version.trim().is_empty() {
return Err(MacpError::InvalidPayload);
}
if payload.participants.is_empty() && !allow_empty_participants {
return Err(MacpError::InvalidPayload);
}
const MAX_PARTICIPANTS: usize = 1000;
if payload.participants.len() > MAX_PARTICIPANTS {
return Err(MacpError::InvalidPayload);
}
let mut seen = HashSet::new();
for participant in &payload.participants {
let participant = participant.trim();
if participant.is_empty() || !seen.insert(participant.to_string()) {
return Err(MacpError::InvalidPayload);
}
}
if payload.max_suspend_ms < 0 {
return Err(MacpError::InvalidPayload);
}
Ok(())
}
pub fn validate_strict_session_start_payload(
mode: &str,
payload: &SessionStartPayload,
) -> Result<(), MacpError> {
if !requires_strict_session_start(mode) {
return Ok(());
}
validate_canonical_session_start_payload_for_mode(mode, payload)
}
pub fn validate_session_id_for_acceptance(session_id: &str) -> Result<(), MacpError> {
if session_id.is_empty() {
return Err(MacpError::InvalidSessionId);
}
if session_id.len() == 36 && session_id.contains('-') {
if let Ok(parsed) = uuid::Uuid::parse_str(session_id) {
if parsed.as_hyphenated().to_string() == session_id {
match parsed.get_version() {
Some(uuid::Version::Random) | Some(uuid::Version::SortRand) => {
return Ok(());
}
_ => {}
}
}
return Err(MacpError::InvalidSessionId);
}
}
if session_id.len() >= 22
&& session_id
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
{
return Ok(());
}
Err(MacpError::InvalidSessionId)
}
#[cfg(test)]
mod tests {
use super::*;
use prost::Message;
#[test]
fn default_mode_and_configuration_version_values() {
assert_eq!(DEFAULT_MODE_VERSION, "1.0.0");
assert_eq!(DEFAULT_CONFIGURATION_VERSION, "config.default");
}
fn encode_payload(ttl_ms: i64, participants: Vec<String>) -> Vec<u8> {
let payload = SessionStartPayload {
intent: String::new(),
participants,
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
policy_version: String::new(),
ttl_ms,
context_id: String::new(),
extensions: std::collections::HashMap::new(),
roots: vec![],
max_suspend_ms: 0,
};
payload.encode_to_vec()
}
#[test]
fn parse_empty_payload_is_invalid() {
let err = parse_session_start_payload(b"").unwrap_err();
assert_eq!(err.to_string(), "InvalidPayload");
}
#[test]
fn parse_valid_protobuf_payload() {
let bytes = encode_payload(5000, vec!["alice".into(), "bob".into()]);
let result = parse_session_start_payload(&bytes).unwrap();
assert_eq!(result.ttl_ms, 5000);
assert_eq!(result.participants, vec!["alice", "bob"]);
}
#[test]
fn extract_ttl_requires_explicit_positive_value() {
let payload = SessionStartPayload::default();
assert_eq!(
extract_ttl_ms(&payload).unwrap_err().to_string(),
"InvalidTtl"
);
let payload = SessionStartPayload {
ttl_ms: 5000,
..Default::default()
};
assert_eq!(extract_ttl_ms(&payload).unwrap(), 5000);
}
#[test]
fn standard_mode_requires_explicit_versions() {
let payload = SessionStartPayload {
participants: vec!["alice".into()],
mode_version: String::new(),
configuration_version: "cfg-1".into(),
ttl_ms: 1000,
..Default::default()
};
assert_eq!(
validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
.unwrap_err()
.to_string(),
"InvalidPayload"
);
let payload = SessionStartPayload {
participants: vec!["alice".into()],
mode_version: "1.0.0".into(),
configuration_version: String::new(),
ttl_ms: 1000,
..Default::default()
};
assert_eq!(
validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
.unwrap_err()
.to_string(),
"InvalidPayload"
);
let payload = SessionStartPayload {
participants: vec!["alice".into()],
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
ttl_ms: 0,
..Default::default()
};
assert_eq!(
validate_strict_session_start_payload("macp.mode.decision.v1", &payload)
.unwrap_err()
.to_string(),
"InvalidTtl"
);
}
fn empty_roster_payload() -> SessionStartPayload {
SessionStartPayload {
participants: vec![],
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
ttl_ms: 1000,
..Default::default()
}
}
#[test]
fn decision_accepts_an_empty_participant_list() {
assert!(
validate_strict_session_start_payload("macp.mode.decision.v1", &empty_roster_payload())
.is_ok(),
"RFC-MACP-0001 §7.1 requires participants only when the mode does, and \
RFC-MACP-0007 makes Decision authority role-based; spec #99's \
decision_zero_participants.json pins the accepted SessionStart"
);
let mut broken = empty_roster_payload();
broken.mode_version = String::new();
assert!(
validate_strict_session_start_payload("macp.mode.decision.v1", &broken).is_err(),
"an empty roster must not waive the rest of the strict contract"
);
}
#[test]
fn every_standard_mode_except_decision_rejects_an_empty_roster() {
for mode in [
"macp.mode.proposal.v1",
"macp.mode.task.v1",
"macp.mode.handoff.v1",
"macp.mode.quorum.v1",
"ext.multi_round.v1",
] {
assert!(
requires_strict_session_start(mode),
"{mode} must be strict for this table to mean anything"
);
assert_eq!(
validate_strict_session_start_payload(mode, &empty_roster_payload())
.unwrap_err()
.to_string(),
"InvalidPayload",
"{mode} must still reject an empty participant list at the core layer"
);
}
assert_eq!(
validate_canonical_session_start_payload_for_mode(
"ext.promoted.v1",
&empty_roster_payload()
)
.unwrap_err()
.to_string(),
"InvalidPayload",
"an unrecognised (e.g. promoted) mode must default to the strict roster rule"
);
assert_eq!(
validate_canonical_session_start_payload(&empty_roster_payload())
.unwrap_err()
.to_string(),
"InvalidPayload"
);
}
fn open_session(ttl_expiry: i64) -> Session {
Session {
session_id: "s1".into(),
state: SessionState::Open,
ttl_expiry,
ttl_ms: 60_000,
started_at_unix_ms: 0,
resolution: None,
mode: "macp.mode.decision.v1".into(),
mode_state: vec![],
participants: vec![],
seen_message_ids: HashSet::new(),
intent: String::new(),
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
policy_version: String::new(),
context_id: String::new(),
extensions: HashMap::new(),
roots: vec![],
initiator_sender: "agent://a".into(),
participant_message_counts: HashMap::new(),
participant_last_seen: HashMap::new(),
policy_definition: None,
suspended_at_ms: None,
accumulated_suspended_ms: 0,
suspension_intervals: vec![],
semantics_rev: CURRENT_SEMANTICS_REV,
max_suspend_ms: 0,
}
}
#[test]
fn suspend_then_resume_banks_ttl() {
let mut s = open_session(10_000);
s.suspend(2_000).unwrap();
assert_eq!(s.state, SessionState::Suspended);
assert_eq!(s.suspended_at_ms, Some(2_000));
s.resume(5_000).unwrap();
assert_eq!(s.state, SessionState::Open);
assert_eq!(s.ttl_expiry, 13_000);
assert_eq!(s.accumulated_suspended_ms, 3_000);
assert_eq!(s.suspended_at_ms, None);
}
#[test]
fn suspend_requires_open_and_resume_requires_suspended() {
let mut s = open_session(10_000);
assert!(matches!(
s.resume(1).unwrap_err(),
MacpError::SessionNotOpen
));
s.suspend(1).unwrap();
assert!(matches!(
s.suspend(2).unwrap_err(),
MacpError::SessionNotOpen
));
}
#[test]
fn resume_exceeding_max_suspend_expires() {
let mut s = open_session(10_000);
s.suspend(0).unwrap();
let err = s.resume(MAX_SUSPEND_MS + 1).unwrap_err();
assert!(matches!(err, MacpError::TtlExpired));
assert_eq!(s.state, SessionState::Expired);
}
#[test]
fn bound_cap_overrides_default_on_resume() {
let mut s = open_session(10_000);
s.max_suspend_ms = 500;
s.suspend(0).unwrap();
let err = s.resume(501).unwrap_err();
assert!(matches!(err, MacpError::TtlExpired));
assert_eq!(s.state, SessionState::Expired);
}
#[test]
fn bound_cap_within_limit_resumes_and_banks_ttl() {
let mut s = open_session(10_000);
s.max_suspend_ms = 500;
s.suspend(0).unwrap();
s.resume(400).unwrap();
assert_eq!(s.state, SessionState::Open);
assert_eq!(s.ttl_expiry, 10_400);
}
#[test]
fn bound_cap_counts_suspension_cumulatively_across_pauses() {
let mut s = open_session(10_000);
s.max_suspend_ms = 500;
s.suspend(0).unwrap();
s.resume(300).unwrap();
assert_eq!(s.accumulated_suspended_ms, 300);
assert_eq!(s.ttl_expiry, 10_300);
s.suspend(400).unwrap();
let err = s.resume(700).unwrap_err();
assert!(matches!(err, MacpError::TtlExpired));
assert_eq!(s.state, SessionState::Expired);
assert_eq!(s.accumulated_suspended_ms, 600);
}
#[test]
fn suspend_cap_exceeded_uses_bound_cap() {
let mut s = open_session(10_000);
s.max_suspend_ms = 500;
s.suspend(0).unwrap();
assert!(!s.suspend_cap_exceeded(400));
assert!(s.suspend_cap_exceeded(501));
}
#[test]
fn unbound_session_uses_default_cap() {
let s = open_session(10_000);
assert_eq!(s.max_suspend_ms, 0);
assert_eq!(s.effective_max_suspend_ms(), MAX_SUSPEND_MS);
}
#[test]
fn negative_max_suspend_ms_rejected_in_canonical_payload() {
let payload = SessionStartPayload {
participants: vec!["a".into()],
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
ttl_ms: 60_000,
max_suspend_ms: -1,
..Default::default()
};
assert_eq!(
validate_canonical_session_start_payload(&payload)
.unwrap_err()
.to_string(),
"InvalidPayload"
);
let ok0 = SessionStartPayload {
max_suspend_ms: 0,
..payload.clone()
};
validate_canonical_session_start_payload(&ok0).unwrap();
let ok_pos = SessionStartPayload {
max_suspend_ms: 60_000,
..payload
};
validate_canonical_session_start_payload(&ok_pos).unwrap();
}
#[test]
fn cancel_from_open_or_suspended_then_terminal_is_rejected() {
let mut s = open_session(10_000);
s.suspend(1).unwrap();
s.cancel().unwrap();
assert_eq!(s.state, SessionState::Cancelled);
assert_eq!(s.suspended_at_ms, None);
assert!(matches!(s.cancel().unwrap_err(), MacpError::SessionNotOpen));
let mut open = open_session(10_000);
open.cancel().unwrap();
assert_eq!(open.state, SessionState::Cancelled);
}
#[test]
fn standard_mode_rejects_duplicate_participants() {
let payload = SessionStartPayload {
participants: vec!["alice".into(), "alice".into()],
mode_version: "1.0.0".into(),
configuration_version: "cfg-1".into(),
ttl_ms: 1000,
..Default::default()
};
assert_eq!(
validate_strict_session_start_payload("macp.mode.proposal.v1", &payload)
.unwrap_err()
.to_string(),
"InvalidPayload"
);
}
#[test]
fn multi_round_requires_strict_session_start() {
let payload = SessionStartPayload::default();
assert!(validate_strict_session_start_payload("ext.multi_round.v1", &payload).is_err());
}
#[test]
fn valid_uuid_v4_accepted() {
let id = uuid::Uuid::new_v4().as_hyphenated().to_string();
validate_session_id_for_acceptance(&id).unwrap();
}
#[test]
fn valid_base64url_accepted() {
validate_session_id_for_acceptance("abcdefghijklmnopqrstuv").unwrap();
validate_session_id_for_acceptance("abc-def_ghi-jkl_mno-pqr").unwrap();
}
#[test]
fn empty_id_rejected() {
assert_eq!(
validate_session_id_for_acceptance("")
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
}
#[test]
fn short_weak_id_rejected() {
assert_eq!(
validate_session_id_for_acceptance("s1")
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
assert_eq!(
validate_session_id_for_acceptance("decision-demo-1")
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
}
#[test]
fn uppercase_uuid_rejected() {
let id = uuid::Uuid::new_v4()
.as_hyphenated()
.to_string()
.to_uppercase();
assert_eq!(
validate_session_id_for_acceptance(&id)
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
}
#[test]
fn base64url_36_chars_with_hyphen_accepted() {
let id = "Zx-abcdefghijklmnopqrstuvwxyz_ABCDE-";
assert_eq!(id.len(), 36);
assert!(uuid::Uuid::parse_str(id).is_err());
validate_session_id_for_acceptance(id).unwrap();
}
#[test]
fn uuid_shaped_but_wrong_version_does_not_fall_through() {
let v4 = uuid::Uuid::new_v4();
let mut bytes = *v4.as_bytes();
bytes[6] = (bytes[6] & 0x0F) | 0x10;
bytes[8] = (bytes[8] & 0x3F) | 0x80;
let v1_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
assert!(validate_session_id_for_acceptance(&v1_id).is_err());
}
#[test]
fn base64url_too_short_rejected() {
assert_eq!(
validate_session_id_for_acceptance("abcdefghij")
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
}
#[test]
fn valid_uuid_v7_accepted() {
let v4 = uuid::Uuid::new_v4();
let mut bytes = *v4.as_bytes();
bytes[6] = (bytes[6] & 0x0F) | 0x70;
bytes[8] = (bytes[8] & 0x3F) | 0x80;
let v7_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
assert!(validate_session_id_for_acceptance(&v7_id).is_ok());
}
#[test]
fn uuid_v1_rejected() {
let v4 = uuid::Uuid::new_v4();
let mut bytes = *v4.as_bytes();
bytes[6] = (bytes[6] & 0x0F) | 0x10;
bytes[8] = (bytes[8] & 0x3F) | 0x80;
let v1_id = uuid::Uuid::from_bytes(bytes).as_hyphenated().to_string();
assert_eq!(
validate_session_id_for_acceptance(&v1_id)
.unwrap_err()
.to_string(),
"InvalidSessionId"
);
}
#[test]
fn too_many_participants_rejected() {
let participants: Vec<String> = (0..1001).map(|i| format!("agent://p{i}")).collect();
let bytes = encode_payload(5000, participants);
let payload = parse_session_start_payload(&bytes).unwrap();
assert_eq!(
validate_canonical_session_start_payload(&payload)
.unwrap_err()
.to_string(),
"InvalidPayload"
);
}
#[test]
fn max_participants_accepted() {
let participants: Vec<String> = (0..1000).map(|i| format!("agent://p{i}")).collect();
let bytes = encode_payload(5000, participants);
let payload = parse_session_start_payload(&bytes).unwrap();
validate_canonical_session_start_payload(&payload).unwrap();
}
#[test]
fn resume_records_completed_suspension_intervals() {
let mut s = open_session(100_000);
assert!(s.suspension_intervals.is_empty());
s.suspend(2_000).unwrap();
assert!(s.suspension_intervals.is_empty());
s.resume(5_000).unwrap();
s.suspend(6_000).unwrap();
s.resume(6_500).unwrap();
assert_eq!(s.suspension_intervals, vec![(2_000, 5_000), (6_000, 6_500)]);
let mut legacy = open_session(100_000);
legacy.semantics_rev = 0;
legacy.suspend(1_000).unwrap();
legacy.resume(1_400).unwrap();
assert_eq!(legacy.suspension_intervals, vec![(1_000, 1_400)]);
}
#[test]
fn resume_records_the_pause_even_when_it_force_expires() {
let mut s = open_session(10_000);
s.suspend(0).unwrap();
assert!(s.resume(MAX_SUSPEND_MS + 1).is_err());
assert_eq!(s.state, SessionState::Expired);
assert_eq!(s.suspension_intervals, vec![(0, MAX_SUSPEND_MS + 1)]);
}
#[test]
fn unsuspended_deadline_walks_the_suspension_intervals() {
let base = open_session(1_000_000);
let with = |pairs: Vec<(i64, i64)>| {
let mut s = base.clone();
s.suspension_intervals = pairs;
s
};
assert_eq!(with(vec![]).unsuspended_deadline(1_000, 100), 1_100);
assert_eq!(
with(vec![(1_050, 1_200)]).unsuspended_deadline(1_000, 100),
1_250
);
assert_eq!(
with(vec![(1_500, 1_600)]).unsuspended_deadline(1_000, 100),
1_100
);
assert_eq!(
with(vec![(1_100, 1_300)]).unsuspended_deadline(1_000, 100),
1_100
);
assert_eq!(
with(vec![(1_050, 1_200), (1_230, 1_300)]).unsuspended_deadline(1_000, 100),
1_320
);
assert_eq!(
with(vec![(500, 700)]).unsuspended_deadline(1_000, 100),
1_100
);
assert_eq!(
with(vec![(500, 700), (1_050, 1_200)]).unsuspended_deadline(1_000, 100),
1_250
);
}
#[test]
fn unsuspended_deadline_never_over_reports_on_adversarial_pairs() {
let base = open_session(1_000_000);
let with = |pairs: Vec<(i64, i64)>| {
let mut s = base.clone();
s.suspension_intervals = pairs;
s
};
assert_eq!(
with(vec![(1_050, 1_400), (1_100, 1_200)]).unsuspended_deadline(1_000, 100),
1_450
);
assert_eq!(
with(vec![(1_050, 1_400), (1_200, 1_300)]).unsuspended_deadline(1_000, 100),
1_450
);
assert_eq!(
with(vec![(1_050, 1_000)]).unsuspended_deadline(1_000, 100),
1_100
);
assert_eq!(
with(vec![(1_050, 1_000), (1_060, 1_200)]).unsuspended_deadline(1_000, 100),
1_240
);
let sorted = with(vec![(1_050, 1_200), (1_230, 1_300)]).unsuspended_deadline(1_000, 100);
let unsorted = with(vec![(1_230, 1_300), (1_050, 1_200)]).unsuspended_deadline(1_000, 100);
assert_eq!(sorted, 1_320);
assert!(
unsorted <= sorted,
"an unsorted vec must under-report, never over-report \
({unsorted} vs {sorted})"
);
for pairs in [
vec![(1_050, 1_400), (1_100, 1_200)],
vec![(1_050, 1_400), (1_200, 1_300)],
vec![(1_050, 1_000)],
vec![(1_230, 1_300), (1_050, 1_200)],
] {
let union: i64 = pairs.iter().map(|(s, e)| (e - s).max(0)).sum();
let walk = with(pairs.clone()).unsuspended_deadline(1_000, 100) - 1_100;
assert!(
(0..=union).contains(&walk),
"walk contribution {walk} outside [0, {union}] for {pairs:?}"
);
}
}
#[test]
fn suspension_cycle_cap_force_expires_at_rev2() {
let mut s = open_session(1_000_000_000);
assert_eq!(s.semantics_rev, CURRENT_SEMANTICS_REV);
#[allow(clippy::assertions_on_constants)]
{
assert!(CURRENT_SEMANTICS_REV >= 2);
}
for i in 0..MAX_SUSPENSION_CYCLES as i64 {
s.suspend(i).unwrap();
s.resume(i).unwrap();
}
assert_eq!(s.suspension_intervals.len(), MAX_SUSPENSION_CYCLES);
assert_eq!(s.accumulated_suspended_ms, 0);
assert_eq!(s.state, SessionState::Open);
s.suspend(MAX_SUSPENSION_CYCLES as i64).unwrap();
let err = s.resume(MAX_SUSPENSION_CYCLES as i64).unwrap_err();
assert!(matches!(err, MacpError::TtlExpired));
assert_eq!(s.state, SessionState::Expired);
assert_eq!(
s.accumulated_suspended_ms, 0,
"the duration cap must be nowhere near exhausted, else the count \
cap is not what fired"
);
let mut legacy = open_session(1_000_000_000);
legacy.semantics_rev = 1;
for i in 0..(MAX_SUSPENSION_CYCLES as i64 + 10) {
legacy.suspend(i).unwrap();
legacy.resume(i).unwrap();
}
assert_eq!(legacy.state, SessionState::Open);
}
#[test]
fn suspension_cycle_cap_stops_recording_at_rev1_without_expiring() {
for rev in [0u32, 1] {
let mut s = open_session(1_000_000_000);
s.semantics_rev = rev;
for i in 0..(MAX_SUSPENSION_CYCLES as i64 * 2) {
s.suspend(i).unwrap();
assert_eq!(s.state, SessionState::Suspended, "rev {rev}, cycle {i}");
s.resume(i)
.unwrap_or_else(|e| panic!("rev {rev}, cycle {i} must resume, got {e:?}"));
assert_eq!(s.state, SessionState::Open, "rev {rev}, cycle {i}");
}
assert_eq!(
s.suspension_intervals.len(),
MAX_SUSPENSION_CYCLES,
"rev {rev}: the vec must stay pinned at the cap, not grow \
unbounded — an unbounded vec is the O(N²) snapshot \
amplification MAX_SUSPENSION_CYCLES exists to close"
);
assert_eq!(s.suspension_intervals[0], (0, 0));
assert_eq!(
s.suspension_intervals[MAX_SUSPENSION_CYCLES - 1],
(
MAX_SUSPENSION_CYCLES as i64 - 1,
MAX_SUSPENSION_CYCLES as i64 - 1
)
);
}
}
#[test]
fn legacy_rev2_snapshot_without_intervals_walks_early_not_late() {
let mut recorded = open_session(1_000_000);
recorded.semantics_rev = 2;
recorded.accumulated_suspended_ms = 150 + 70;
recorded.suspension_intervals = vec![(1_050, 1_200), (1_230, 1_300)];
let mut legacy = recorded.clone();
legacy.suspension_intervals.clear();
let recorded_deadline = recorded.unsuspended_deadline(1_000, 100);
let legacy_deadline = legacy.unsuspended_deadline(1_000, 100);
assert_eq!(recorded_deadline, 1_320);
assert_eq!(legacy_deadline, 1_100);
assert!(
legacy_deadline <= recorded_deadline,
"an under-reported interval vec must move the deadline EARLIER \
({legacy_deadline} vs {recorded_deadline}), never later"
);
let walk_sum = recorded_deadline - (1_000 + 100);
assert!(walk_sum <= recorded.accumulated_suspended_ms);
let legacy_walk_sum = legacy_deadline - (1_000 + 100);
assert!(legacy_walk_sum <= legacy.accumulated_suspended_ms);
}
}