macp-runtime 0.8.5

MACP reference runtime: a coordination kernel and gRPC server enforcing session boundaries, message validation, append-only history, modes, and governance policy.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
use macp_runtime::log_store::LogStore;
use macp_runtime::pb::{Envelope, SessionStartPayload};
use macp_runtime::registry::SessionRegistry;
use macp_runtime::runtime::Runtime;
use macp_runtime::storage::MemoryBackend;
use prost::Message;
use std::sync::Arc;

fn new_sid() -> String {
    uuid::Uuid::new_v4().as_hyphenated().to_string()
}

fn make_runtime() -> Runtime {
    let storage: Arc<dyn macp_runtime::storage::StorageBackend> = Arc::new(MemoryBackend);
    let registry = Arc::new(SessionRegistry::new());
    let log_store = Arc::new(LogStore::new());
    Runtime::new(storage, registry, log_store)
}

fn session_start(participants: Vec<String>) -> Vec<u8> {
    SessionStartPayload {
        intent: "stream-test".into(),
        participants,
        mode_version: "1.0.0".into(),
        configuration_version: "cfg-1".into(),
        policy_version: String::new(),
        ttl_ms: 60_000,
        context_id: String::new(),
        extensions: std::collections::HashMap::new(),
        roots: vec![],
        max_suspend_ms: 0,
    }
    .encode_to_vec()
}

fn envelope(
    mode: &str,
    message_type: &str,
    message_id: &str,
    session_id: &str,
    sender: &str,
    payload: Vec<u8>,
) -> Envelope {
    Envelope {
        macp_version: "1.0".into(),
        mode: mode.into(),
        message_type: message_type.into(),
        message_id: message_id.into(),
        session_id: session_id.into(),
        sender: sender.into(),
        timestamp_unix_ms: chrono::Utc::now().timestamp_millis(),
        payload,
    }
}

#[tokio::test]
async fn stream_receives_accepted_envelopes() {
    let rt = make_runtime();
    let sid = new_sid();
    let mode = "macp.mode.decision.v1";

    let mut rx = rt.subscribe_session_stream(&sid);

    rt.process(
        &envelope(
            mode,
            "SessionStart",
            "m1",
            &sid,
            "agent://orchestrator",
            session_start(vec!["agent://orchestrator".into(), "agent://a".into()]),
        ),
        None,
    )
    .await
    .unwrap();

    let env = rx.recv().await.unwrap();
    assert_eq!(env.message_id, "m1");
    assert_eq!(env.message_type, "SessionStart");
}

#[tokio::test]
async fn stream_ordering_matches_processing_order() {
    let rt = make_runtime();
    let sid = new_sid();
    let mode = "macp.mode.decision.v1";

    let mut rx = rt.subscribe_session_stream(&sid);

    rt.process(
        &envelope(
            mode,
            "SessionStart",
            "m1",
            &sid,
            "agent://orchestrator",
            session_start(vec!["agent://orchestrator".into(), "agent://a".into()]),
        ),
        None,
    )
    .await
    .unwrap();

    let proposal = macp_runtime::decision_pb::ProposalPayload {
        proposal_id: "p1".into(),
        option: "deploy".into(),
        rationale: "ready".into(),
        supporting_data: vec![],
    }
    .encode_to_vec();
    rt.process(
        &envelope(
            mode,
            "Proposal",
            "m2",
            &sid,
            "agent://orchestrator",
            proposal,
        ),
        None,
    )
    .await
    .unwrap();

    let first = rx.recv().await.unwrap();
    let second = rx.recv().await.unwrap();
    assert_eq!(first.message_id, "m1");
    assert_eq!(second.message_id, "m2");
}

#[tokio::test]
async fn concurrent_subscribers_both_receive_events() {
    let rt = make_runtime();
    let sid = new_sid();
    let mode = "macp.mode.decision.v1";

    let mut rx1 = rt.subscribe_session_stream(&sid);
    let mut rx2 = rt.subscribe_session_stream(&sid);

    rt.process(
        &envelope(
            mode,
            "SessionStart",
            "m1",
            &sid,
            "agent://orchestrator",
            session_start(vec!["agent://orchestrator".into(), "agent://a".into()]),
        ),
        None,
    )
    .await
    .unwrap();

    let env1 = rx1.recv().await.unwrap();
    let env2 = rx2.recv().await.unwrap();
    assert_eq!(env1.message_id, "m1");
    assert_eq!(env2.message_id, "m1");
}

// ---------------------------------------------------------------------------
// RFC-MACP-0010 ยง5.1(2) โ€” the synthetic implicit accept on the wire.
// ---------------------------------------------------------------------------
//
// The synthetic accept is `EntryKind::Incoming`, so unlike the runtime's
// internal lifecycle entries it is an *accepted* envelope and reaches
// `StreamSession` subscribers like any other. That is one of the two honest
// deltas the freeze-profile carve-out names (the other being that it consumes
// an accepted ordinal), so it is pinned here rather than left implicit.

const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
const HANDOFF_OWNER: &str = "agent://owner";
const HANDOFF_TARGET: &str = "agent://target";
const HANDOFF_POLICY: &str = "handoff-auto-accept";
const HANDOFF_TIMEOUT_MS: i64 = 60;

fn handoff_runtime() -> Runtime {
    let rt = make_runtime();
    rt.register_policy(macp_runtime::macp_core::policy::PolicyDefinition {
        policy_id: HANDOFF_POLICY.into(),
        mode: HANDOFF_MODE.into(),
        description: "implicit accept after a short timeout".into(),
        rules: serde_json::json!({
            "acceptance": { "implicit_accept_timeout_ms": HANDOFF_TIMEOUT_MS },
            "commitment": { "authority": "initiator_only" }
        }),
        schema_version: 1,
    })
    .expect("policy registers");
    rt
}

fn handoff_start() -> Vec<u8> {
    SessionStartPayload {
        intent: "escalate".into(),
        participants: vec![HANDOFF_OWNER.into(), HANDOFF_TARGET.into()],
        mode_version: "1.0.0".into(),
        configuration_version: "cfg-1".into(),
        policy_version: HANDOFF_POLICY.into(),
        ttl_ms: 60_000,
        context_id: String::new(),
        extensions: std::collections::HashMap::new(),
        roots: vec![],
        max_suspend_ms: 0,
    }
    .encode_to_vec()
}

/// Criterion 8: a subscriber sees the synthetic envelope in emission order โ€”
/// after the last message accepted before the deadline, and before the trigger
/// that provoked it.
#[tokio::test]
async fn stream_subscribers_see_the_synthetic_envelope_in_order() {
    let rt = handoff_runtime();
    let sid = new_sid();
    let mut rx = rt.subscribe_session_stream(&sid);

    rt.process(
        &envelope(
            HANDOFF_MODE,
            "SessionStart",
            "start-1",
            &sid,
            HANDOFF_OWNER,
            handoff_start(),
        ),
        None,
    )
    .await
    .unwrap();
    rt.process(
        &envelope(
            HANDOFF_MODE,
            "HandoffOffer",
            "offer-1",
            &sid,
            HANDOFF_OWNER,
            macp_runtime::handoff_pb::HandoffOfferPayload {
                handoff_id: "h1".into(),
                target_participant: HANDOFF_TARGET.into(),
                scope: "support".into(),
                reason: "escalate".into(),
            }
            .encode_to_vec(),
        ),
        None,
    )
    .await
    .unwrap();
    // The pre-deadline message the synthetic must come *after*.
    rt.process(
        &envelope(
            HANDOFF_MODE,
            "HandoffContext",
            "ctx-1",
            &sid,
            HANDOFF_OWNER,
            macp_runtime::handoff_pb::HandoffContextPayload {
                handoff_id: "h1".into(),
                content_type: "text/plain".into(),
                context: b"background".to_vec(),
            }
            .encode_to_vec(),
        ),
        None,
    )
    .await
    .unwrap();

    tokio::time::sleep(std::time::Duration::from_millis(
        HANDOFF_TIMEOUT_MS as u64 + 40,
    ))
    .await;

    rt.process(
        &envelope(
            HANDOFF_MODE,
            "Commitment",
            "commit-1",
            &sid,
            HANDOFF_OWNER,
            macp_runtime::pb::CommitmentPayload {
                commitment_id: "c1".into(),
                action: "handoff.accepted".into(),
                authority_scope: "support".into(),
                reason: "bound".into(),
                mode_version: "1.0.0".into(),
                policy_version: HANDOFF_POLICY.into(),
                configuration_version: "cfg-1".into(),
                outcome_positive: true,
                supersedes: None,
            }
            .encode_to_vec(),
        ),
        None,
    )
    .await
    .expect("resolves");

    let mut seen = Vec::new();
    for _ in 0..4 {
        seen.push(rx.recv().await.expect("subscriber must not lag"));
    }
    let ids: Vec<&str> = seen.iter().map(|e| e.message_id.as_str()).collect();
    assert_eq!(
        ids,
        vec!["start-1", "offer-1", "ctx-1", "implicit-accept:h1"],
        "the synthetic must follow the last pre-deadline message"
    );

    // The fifth is the trigger, strictly after the synthetic.
    let commit = rx.recv().await.unwrap();
    assert_eq!(commit.message_id, "commit-1");

    // The synthetic on the wire is the target's accept, not a runtime message.
    let syn = &seen[3];
    assert_eq!(syn.sender, HANDOFF_TARGET);
    assert_eq!(syn.message_type, "HandoffAccept");
    assert_eq!(syn.mode, HANDOFF_MODE);
    assert_eq!(syn.session_id, sid);
    let payload = macp_runtime::handoff_pb::HandoffAcceptPayload::decode(&*syn.payload).unwrap();
    assert!(payload.implicit);
    assert_eq!(payload.accepted_by, HANDOFF_TARGET);
    // Its envelope clock is the deadline, not the publish time.
    assert!(syn.timestamp_unix_ms < commit.timestamp_unix_ms);
}

// ---------------------------------------------------------------------------
// RFC-MACP-0006 ยง3.2 โ€” the live-StreamSession half of the ordinal/delivery
// invariant for SessionSuspend/SessionResume/SessionCancel.
// ---------------------------------------------------------------------------
//
// Contrast with `stream_subscribers_see_the_synthetic_envelope_in_order`
// above: that test pins the *opposite* answer for the handoff synthetic
// accept, which is a mode message with a spec-pinned sender and so DOES
// consume an ordinal and DOES publish. These three entry kinds are runtime
// bookkeeping (`sender == "_runtime"`, `EntryKind::Internal`) and must not.
//
// Citation scoping: `:117` ("MUST NOT consume ordinals") is unrestricted and
// names all three types. `:122` ("MUST NOT deliver an internal annotation on
// this stream") is textually scoped to "a subscribe stream" โ€” the passive
// `StreamSession` replay path pinned by
// `subscribe_never_delivers_lifecycle_annotations` in
// `integration_tests/tests/tier1_protocol/test_suspend_resume.rs`, which is
// the clause that governs literally. The live path asserted here follows *a
// fortiori* from `:117` plus `:120`'s counting argument ("a client can
// determine its position only by counting the distinct accepted envelopes it
// has been delivered"): this runtime's live and passive paths share one
// `stream_bus` and one `get_incoming_after`, so a live subscriber counting its
// position the same way must see the same exclusion.
#[tokio::test]
async fn stream_subscriber_never_sees_lifecycle_annotations() {
    let rt = make_runtime();
    let sid = new_sid();
    let mode = "macp.mode.decision.v1";

    let mut rx = rt.subscribe_session_stream(&sid);

    rt.process(
        &envelope(
            mode,
            "SessionStart",
            "m1",
            &sid,
            "agent://orchestrator",
            session_start(vec!["agent://orchestrator".into(), "agent://a".into()]),
        ),
        None,
    )
    .await
    .unwrap();

    rt.suspend_session(&sid, "pause", "agent://orchestrator")
        .await
        .unwrap();
    rt.resume_session(&sid, "carry on", "agent://orchestrator")
        .await
        .unwrap();

    let proposal = macp_runtime::decision_pb::ProposalPayload {
        proposal_id: "p1".into(),
        option: "deploy".into(),
        rationale: "after resume".into(),
        supporting_data: vec![],
    }
    .encode_to_vec();
    rt.process(
        &envelope(
            mode,
            "Proposal",
            "m2",
            &sid,
            "agent://orchestrator",
            proposal,
        ),
        None,
    )
    .await
    .unwrap();

    rt.cancel_session(&sid, "done", "agent://orchestrator")
        .await
        .unwrap();

    let mut received = Vec::new();
    loop {
        match rx.try_recv() {
            Ok(env) => received.push(env.message_type),
            Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break,
            // Treat a lag as a hard failure rather than folding it into
            // `Empty` โ€” these tests publish single-digit envelopes into a
            // 256-capacity channel, so `Lagged` is unreachable today; a
            // future capacity change should surface here as red, not a false
            // pass.
            Err(tokio::sync::broadcast::error::TryRecvError::Lagged(n)) => {
                panic!("subscriber lagged by {n} envelopes")
            }
            Err(tokio::sync::broadcast::error::TryRecvError::Closed) => break,
        }
    }

    assert_eq!(
        received,
        vec!["SessionStart".to_string(), "Proposal".to_string()],
        "a live StreamSession subscriber must receive exactly the ordinal-\
         consuming client envelopes, in order, and none of \
         SessionSuspend/SessionResume/SessionCancel"
    );
}