1use crate::mode::{Mode, ModeResponse};
20use macp_core::error::MacpError;
21use macp_core::session::{Session, SessionState};
22use macp_pb::pb::Envelope;
23
24#[derive(Debug, Clone, PartialEq, Eq)]
26pub enum Precheck {
27 Duplicate,
29 Expired,
31 Proceed,
33}
34
35pub fn check_preconditions(
44 session: &Session,
45 env: &Envelope,
46 now_ms: i64,
47) -> Result<Precheck, MacpError> {
48 if session.seen_message_ids.contains(&env.message_id) {
49 return Ok(Precheck::Duplicate);
50 }
51 if env.mode != session.mode {
52 return Err(MacpError::InvalidEnvelope);
53 }
54 if session.state == SessionState::Open && now_ms > session.ttl_expiry {
55 return Ok(Precheck::Expired);
56 }
57 if session.state != SessionState::Open {
58 return Err(MacpError::SessionNotOpen);
59 }
60 Ok(Precheck::Proceed)
61}
62
63pub fn validate_message(
80 session: &Session,
81 env: &Envelope,
82 mode: &dyn Mode,
83) -> Result<ModeResponse, MacpError> {
84 mode.authorize_sender(session, env)?;
85 mode.validate_client_envelope(session, env)?;
86 mode.on_message(session, env)
87}
88
89pub fn commit(
98 session: &mut Session,
99 env: &Envelope,
100 response: ModeResponse,
101 now_ms: i64,
102) -> SessionState {
103 session.seen_message_ids.insert(env.message_id.clone());
104 session.record_participant_activity(&env.sender, now_ms);
105 session.apply_mode_response(response);
106 session.state.clone()
107}
108
109#[derive(Debug, Clone, PartialEq)]
111pub enum StepOutcome {
112 Duplicate,
114 Accepted { state: SessionState },
116}
117
118pub fn step(
128 session: &mut Session,
129 env: &Envelope,
130 mode: &dyn Mode,
131 now_ms: i64,
132) -> Result<StepOutcome, MacpError> {
133 match check_preconditions(session, env, now_ms)? {
134 Precheck::Duplicate => Ok(StepOutcome::Duplicate),
135 Precheck::Expired => {
136 session.state = SessionState::Expired;
137 Err(MacpError::TtlExpired)
138 }
139 Precheck::Proceed => {
140 let response = validate_message(session, env, mode)?;
141 let state = commit(session, env, response, now_ms);
142 Ok(StepOutcome::Accepted { state })
143 }
144 }
145}
146
147#[cfg(test)]
148mod tests {
149 use super::*;
150
151 const MODE: &str = "macp.mode.test.v1";
152
153 struct TestMode;
157 impl Mode for TestMode {
158 fn on_session_start(&self, _s: &Session, _e: &Envelope) -> Result<ModeResponse, MacpError> {
159 Ok(ModeResponse::PersistState(vec![]))
160 }
161 fn on_message(&self, _s: &Session, env: &Envelope) -> Result<ModeResponse, MacpError> {
162 if env.message_type == "Commitment" {
163 Ok(ModeResponse::PersistAndResolve {
164 state: vec![1],
165 resolution: vec![2],
166 })
167 } else {
168 Ok(ModeResponse::PersistState(vec![1]))
169 }
170 }
171 }
173
174 fn session() -> Session {
175 Session::builder("11111111-1111-4111-8111-111111111111", MODE, "agent://a")
176 .ttl_expiry(10_000)
177 .ttl_ms(10_000)
178 .participants(vec!["agent://a".into(), "agent://b".into()])
179 .mode_version("1.0.0")
180 .configuration_version("cfg-1")
181 .build()
182 }
183
184 fn env(sender: &str, message_type: &str, message_id: &str) -> Envelope {
185 Envelope {
186 macp_version: "1.0".into(),
187 mode: MODE.into(),
188 message_type: message_type.into(),
189 message_id: message_id.into(),
190 session_id: "11111111-1111-4111-8111-111111111111".into(),
191 sender: sender.into(),
192 timestamp_unix_ms: 0,
193 payload: vec![],
194 }
195 }
196
197 #[test]
198 fn duplicate_is_reported_and_changes_nothing() {
199 let mut s = session();
200 s.seen_message_ids.insert("m1".into());
201 let before = s.seen_message_ids.len();
202 let out = step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, 1).unwrap();
203 assert_eq!(out, StepOutcome::Duplicate);
204 assert_eq!(s.seen_message_ids.len(), before);
205 assert_eq!(s.state, SessionState::Open);
206 }
207
208 #[test]
209 fn mode_binding_mismatch_rejected() {
210 let mut s = session();
211 let mut e = env("agent://a", "Msg", "m1");
212 e.mode = "macp.mode.other.v1".into();
213 assert!(matches!(
214 step(&mut s, &e, &TestMode, 1).unwrap_err(),
215 MacpError::InvalidEnvelope
216 ));
217 assert!(s.seen_message_ids.is_empty());
218 }
219
220 #[test]
221 fn ttl_strict_boundary_does_not_expire_but_past_does() {
222 let mut s = session();
224 let deadline = s.ttl_expiry;
225 let out = step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, deadline).unwrap();
226 assert_eq!(
227 out,
228 StepOutcome::Accepted {
229 state: SessionState::Open
230 }
231 );
232
233 let mut s2 = session();
235 let past = s2.ttl_expiry + 1;
236 let err = step(&mut s2, &env("agent://a", "Msg", "m2"), &TestMode, past).unwrap_err();
237 assert!(matches!(err, MacpError::TtlExpired));
238 assert_eq!(s2.state, SessionState::Expired);
239 assert!(s2.seen_message_ids.is_empty());
240 }
241
242 #[test]
243 fn ttl_does_not_re_expire_a_resolved_session() {
244 let mut s = session();
247 s.state = SessionState::Resolved;
248 let past = s.ttl_expiry + 5_000;
249 let err = step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, past).unwrap_err();
250 assert!(matches!(err, MacpError::SessionNotOpen));
251 assert_eq!(s.state, SessionState::Resolved);
252 }
253
254 #[test]
255 fn non_open_session_rejected() {
256 for st in [SessionState::Resolved, SessionState::Expired] {
257 let mut s = session();
258 s.state = st.clone();
259 assert!(matches!(
260 step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, 1).unwrap_err(),
261 MacpError::SessionNotOpen
262 ));
263 }
264 }
265
266 #[test]
267 fn accepted_consumes_dedup_records_activity_and_applies_state() {
268 let mut s = session();
269 let out = step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, 42).unwrap();
270 assert_eq!(
271 out,
272 StepOutcome::Accepted {
273 state: SessionState::Open
274 }
275 );
276 assert!(s.seen_message_ids.contains("m1"));
277 assert_eq!(s.mode_state, vec![1]);
278 assert_eq!(s.participant_last_seen.get("agent://a"), Some(&42));
279 }
280
281 #[test]
282 fn commitment_resolves() {
283 let mut s = session();
284 let out = step(&mut s, &env("agent://a", "Commitment", "c1"), &TestMode, 1).unwrap();
285 assert_eq!(
286 out,
287 StepOutcome::Accepted {
288 state: SessionState::Resolved
289 }
290 );
291 assert_eq!(s.state, SessionState::Resolved);
292 assert_eq!(s.resolution, Some(vec![2]));
293 }
294
295 #[test]
296 fn rejected_validation_does_not_consume_dedup_slot() {
297 let mut s = session();
301 let err = step(&mut s, &env("agent://stranger", "Msg", "m1"), &TestMode, 1).unwrap_err();
302 assert!(matches!(err, MacpError::Forbidden));
303 assert!(!s.seen_message_ids.contains("m1"));
304 let out = step(&mut s, &env("agent://a", "Msg", "m1"), &TestMode, 1).unwrap();
306 assert_eq!(
307 out,
308 StepOutcome::Accepted {
309 state: SessionState::Open
310 }
311 );
312 assert!(s.seen_message_ids.contains("m1"));
313 }
314
315 #[test]
316 fn clock_is_injected_not_wall_clock() {
317 let mut s = session();
320 s.ttl_expiry = i64::MAX;
321 assert!(matches!(
322 check_preconditions(&s, &env("agent://a", "Msg", "m1"), i64::MAX - 1),
323 Ok(Precheck::Proceed)
324 ));
325 let mut s2 = session();
326 s2.ttl_expiry = 0;
327 assert!(matches!(
328 check_preconditions(&s2, &env("agent://a", "Msg", "m1"), 1),
329 Ok(Precheck::Expired)
330 ));
331 }
332
333 struct BoundaryMode {
337 dispatched: std::sync::atomic::AtomicBool,
338 }
339 impl BoundaryMode {
340 fn new() -> Self {
341 Self {
342 dispatched: std::sync::atomic::AtomicBool::new(false),
343 }
344 }
345 fn dispatched(&self) -> bool {
346 self.dispatched.load(std::sync::atomic::Ordering::SeqCst)
347 }
348 }
349 impl Mode for BoundaryMode {
350 fn on_session_start(&self, _s: &Session, _e: &Envelope) -> Result<ModeResponse, MacpError> {
351 Ok(ModeResponse::PersistState(vec![]))
352 }
353 fn on_message(&self, _s: &Session, _e: &Envelope) -> Result<ModeResponse, MacpError> {
354 self.dispatched
355 .store(true, std::sync::atomic::Ordering::SeqCst);
356 Ok(ModeResponse::PersistState(vec![1]))
357 }
358 fn validate_client_envelope(&self, _s: &Session, env: &Envelope) -> Result<(), MacpError> {
359 if env.message_id == "reserved" {
360 return Err(MacpError::InvalidEnvelope);
361 }
362 Ok(())
363 }
364 }
366
367 #[test]
377 fn validate_message_enforces_the_client_boundary() {
378 let s = session();
379 let mode = BoundaryMode::new();
380
381 let err = validate_message(&s, &env("agent://a", "Msg", "reserved"), &mode).unwrap_err();
382 assert!(matches!(err, MacpError::InvalidEnvelope));
383 assert!(
384 !mode.dispatched(),
385 "the boundary must reject before dispatch"
386 );
387
388 let err =
391 validate_message(&s, &env("agent://stranger", "Msg", "reserved"), &mode).unwrap_err();
392 assert!(matches!(err, MacpError::Forbidden));
393
394 validate_message(&s, &env("agent://a", "Msg", "m1"), &mode).unwrap();
397 assert!(mode.dispatched());
398
399 let mut s2 = session();
401 let err = step(&mut s2, &env("agent://a", "Msg", "reserved"), &mode, 1).unwrap_err();
402 assert!(matches!(err, MacpError::InvalidEnvelope));
403 assert!(s2.seen_message_ids.is_empty());
404 assert_eq!(s2.state, SessionState::Open);
405 }
406
407 #[test]
410 fn default_client_boundary_is_open() {
411 let s = session();
412 for message_id in ["m1", "reserved", "implicit-accept:h1"] {
413 assert!(
414 TestMode
415 .validate_client_envelope(&s, &env("agent://a", "Msg", message_id))
416 .is_ok(),
417 "default hook must accept {message_id}"
418 );
419 }
420 }
421
422 #[test]
423 fn phases_compose_like_step_for_durable_consumers() {
424 let mut s = session();
426 let e = env("agent://b", "Msg", "m1");
427 assert_eq!(check_preconditions(&s, &e, 5).unwrap(), Precheck::Proceed);
428 let resp = validate_message(&s, &e, &TestMode).unwrap();
429 assert!(s.seen_message_ids.is_empty());
431 let state = commit(&mut s, &e, resp, 5);
432 assert_eq!(state, SessionState::Open);
433 assert!(s.seen_message_ids.contains("m1"));
434 }
435}