Skip to main content

macp_modes/
step.rs

1//! The per-message coordination step — the pure, I/O-free kernel invariants.
2//!
3//! Every accepted MACP message passes the same per-message invariants: dedup
4//! (RFC-MACP-0001 §8 idempotency), mode-binding, TTL, and the monotonic OPEN
5//! gate (§7.2/§7.3), then mode validation, then commit. Historically these
6//! lived welded into the gRPC server's `process_message`, so any other consumer
7//! of the coordination core (e.g. an embedding library) had to re-implement
8//! them and risk drift. This module hosts them once — synchronous and free of
9//! tokio, storage, transport, and the wall clock (the caller injects `now_ms`).
10//!
11//! Two ways to drive it:
12//! - [`step`] — all-in-one, for in-memory consumers that do not interpose
13//!   durable storage between validation and commit.
14//! - [`check_preconditions`] + [`validate_message`] + [`commit`] — the phases,
15//!   for a durable consumer (the runtime) that must write the message to its
16//!   append-only log *between* validation and commit, so a failed write never
17//!   consumes a dedup slot.
18
19use crate::mode::{Mode, ModeResponse};
20use macp_core::error::MacpError;
21use macp_core::session::{Session, SessionState};
22use macp_pb::pb::Envelope;
23
24/// Outcome of the mode-independent precondition checks.
25#[derive(Debug, Clone, PartialEq, Eq)]
26pub enum Precheck {
27    /// `message_id` already accepted — idempotent no-op.
28    Duplicate,
29    /// The session's TTL has elapsed; the caller must expire the session.
30    Expired,
31    /// Preconditions satisfied — proceed to mode validation.
32    Proceed,
33}
34
35/// Mode-independent per-message invariants. Pure: no mutation, no I/O, no clock.
36///
37/// Order mirrors the runtime's `process_message` exactly: dedup → mode-binding
38/// → TTL → the monotonic OPEN gate. `now_ms` is the injected clock (the
39/// envelope/replay timestamp). The TTL check uses a strict `>` and is guarded
40/// on `Open`, matching the runtime's `maybe_expire_session`: a message arriving
41/// exactly at `ttl_expiry` does not expire, and a non-`Open` session is never
42/// re-expired (it falls through to [`MacpError::SessionNotOpen`]).
43pub 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
63/// Mode-dependent validation: sender authorization + the client boundary +
64/// mode rules. Pure — returns the [`ModeResponse`] to apply and mutates
65/// nothing. Call only after [`check_preconditions`] returns
66/// [`Precheck::Proceed`].
67///
68/// The middle phase is [`Mode::validate_client_envelope`], which rejects
69/// envelope shapes that are legal as *recorded history* but illegal as *client
70/// submissions* (see that method's contract). It runs after
71/// [`Mode::authorize_sender`] so the existing authorization-before-payload
72/// error ordering is untouched. Because it is a client-boundary check, replay
73/// must not go through this function — this runtime's replay calls
74/// `authorize_sender`/`on_message_at` directly and never reaches here.
75///
76/// A durable consumer that bypasses this helper (as the runtime does, to
77/// interpose its append between validation and commit) MUST call
78/// [`Mode::validate_client_envelope`] itself.
79pub 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
89/// Commit a validated message into the session: consume the dedup slot, record
90/// participant activity, and apply the mode response. Returns the resulting
91/// session state.
92///
93/// A durable consumer MUST call this only after the message has been durably
94/// recorded, so a failed write never consumes a dedup slot. Because nothing
95/// here mutates the session until validation has already succeeded, a rejected
96/// message likewise leaves `seen_message_ids` untouched.
97pub 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/// Outcome of [`step`].
110#[derive(Debug, Clone, PartialEq)]
111pub enum StepOutcome {
112    /// `message_id` already accepted — nothing changed.
113    Duplicate,
114    /// Message validated, committed, and applied; carries the resulting state.
115    Accepted { state: SessionState },
116}
117
118/// All-in-one per-message step for in-memory consumers: preconditions → mode
119/// validation → commit, mirroring the runtime's external contract. Expiry marks
120/// the session `Expired` and returns [`MacpError::TtlExpired`]; a duplicate is
121/// reported as [`StepOutcome::Duplicate`]; any other rejection returns its error
122/// without consuming a dedup slot or applying state.
123///
124/// A durable consumer should instead use [`check_preconditions`],
125/// [`validate_message`], and [`commit`] so it can interpose its append-only
126/// write between validation and commit (see the runtime's `process_message`).
127pub 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    // A trivial mode: every participant may send; a `Commitment` resolves the
154    // session, anything else just persists. Lets us exercise the step invariants
155    // without policy or protobuf payloads.
156    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        // default authorize_sender: sender must be a declared participant.
172    }
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        // now == ttl_expiry: NOT expired (strict `>`), message is accepted.
223        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        // now > ttl_expiry: expired, session marked Expired, dedup untouched.
234        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        // A resolved session past its original ttl must report SessionNotOpen,
245        // not flip to Expired or return TtlExpired (matches maybe_expire_session).
246        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        // Dedup invariant (CLAUDE.md §8): a message rejected by mode validation
298        // must NOT consume its dedup slot — a later valid message with the same
299        // id is accepted normally.
300        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        // Same id, now from an authorized participant: accepted.
305        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        // Expiry is decided purely by the injected now_ms, independent of the
318        // wall clock — a far-future deadline never expires, a past one does.
319        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    // A mode with a client boundary: it refuses one reserved `message_id` and
334    // records whether dispatch was reached, so a test can tell "rejected at
335    // the boundary" from "rejected by the mode's own rules".
336    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        // default authorize_sender: sender must be a declared participant.
365    }
366
367    /// `validate_message` enforces the client boundary (Phase 11c criterion 4),
368    /// so a library consumer driving the phases through this helper inherits it
369    /// without knowing the hook exists.
370    ///
371    /// Asserts the *position* too, not just the error: dispatch must not have
372    /// been reached, and the envelope must be rejected **after**
373    /// `authorize_sender` so the pre-existing Forbidden-before-payload error
374    /// ordering is untouched. `step` inherits it by construction (it calls
375    /// `validate_message`), and a rejection consumes no dedup slot.
376    #[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        // Ordering: unauthorized *and* reserved reports the authorization
389        // error, because the hook runs after `authorize_sender`.
390        let err =
391            validate_message(&s, &env("agent://stranger", "Msg", "reserved"), &mode).unwrap_err();
392        assert!(matches!(err, MacpError::Forbidden));
393
394        // A non-reserved id from the same sender passes the boundary and
395        // reaches dispatch, so the rejection above was the hook's doing.
396        validate_message(&s, &env("agent://a", "Msg", "m1"), &mode).unwrap();
397        assert!(mode.dispatched());
398
399        // Through `step`: rejected, and no dedup slot consumed.
400        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    /// The default hook is `Ok(())`, so every mode that has no client-boundary
408    /// rules is unaffected by the phase (semver-minor, no behavior change).
409    #[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        // The runtime path: check_preconditions -> validate_message -> commit.
425        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        // Nothing applied until commit.
430        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}