Skip to main content

macp_runtime/
server.rs

1use crate::error::MacpError;
2use crate::pb::macp_runtime_service_server::MacpRuntimeService;
3use crate::pb::{
4    session_lifecycle_event, Ack, CancelSessionRequest, CancelSessionResponse,
5    CancellationCapability, Capabilities, Envelope, GetManifestRequest, GetManifestResponse,
6    GetPolicyRequest, GetPolicyResponse, GetSessionRequest, GetSessionResponse, InitializeRequest,
7    InitializeResponse, ListExtModesRequest, ListExtModesResponse, ListModesRequest,
8    ListModesResponse, ListPoliciesRequest, ListPoliciesResponse, ListRootsRequest,
9    ListRootsResponse, ListSessionsRequest, ListSessionsResponse, MacpError as PbMacpError,
10    ManifestCapability, ModeRegistryCapability, ParticipantActivity, PolicyDescriptor,
11    PolicyRegistryCapability, ProgressCapability, PromoteModeRequest, PromoteModeResponse,
12    RegisterExtModeRequest, RegisterExtModeResponse, RegisterPolicyRequest, RegisterPolicyResponse,
13    ResumeSessionRequest, ResumeSessionResponse, RootsCapability, RuntimeInfo, SendRequest,
14    SendResponse, SessionLifecycleEvent, SessionMetadata, SessionState as PbSessionState,
15    SessionsCapability, StreamSessionRequest, StreamSessionResponse, SuspendSessionRequest,
16    SuspendSessionResponse, UnregisterExtModeRequest, UnregisterExtModeResponse,
17    UnregisterPolicyRequest, UnregisterPolicyResponse, WatchModeRegistryRequest,
18    WatchModeRegistryResponse, WatchPoliciesRequest, WatchPoliciesResponse, WatchRootsRequest,
19    WatchRootsResponse, WatchSessionsRequest, WatchSessionsResponse, WatchSignalsRequest,
20    WatchSignalsResponse,
21};
22use crate::runtime::Runtime;
23use crate::security::{AuthIdentity, SecurityLayer};
24use crate::session::SessionState;
25use std::collections::HashMap;
26use std::sync::Arc;
27use tonic::{Request, Response, Status};
28
29type SessionResponseStream = std::pin::Pin<
30    Box<dyn futures_core::Stream<Item = Result<StreamSessionResponse, Status>> + Send>,
31>;
32
33#[derive(Clone)]
34pub struct MacpServer {
35    runtime: Arc<Runtime>,
36    security: SecurityLayer,
37    /// RFC-MACP-0012 §9 file-loaded profile: when policies are preloaded from
38    /// `MACP_POLICIES_DIR`, the policy registry is read-only over the wire —
39    /// `register_policy` is advertised `false` and the mutating RPCs return
40    /// `FAILED_PRECONDITION`. Governance then has exactly one source of truth.
41    policies_read_only: bool,
42    /// Optional external ingress policy engine (E3). Consulted after
43    /// authentication, before kernel acceptance; deny-on-error. See
44    /// `crate::policy_engine`.
45    policy_engine: Option<Arc<dyn crate::policy_engine::PolicyEngine>>,
46}
47
48impl MacpServer {
49    pub fn new(runtime: Arc<Runtime>, security: SecurityLayer) -> Self {
50        Self {
51            runtime,
52            security,
53            policies_read_only: false,
54            policy_engine: None,
55        }
56    }
57
58    pub fn with_read_only_policies(mut self) -> Self {
59        self.policies_read_only = true;
60        self
61    }
62
63    /// Install an external ingress policy engine (OPA/Cedar/custom). All
64    /// session starts, session-scoped sends, and session reads are then
65    /// additionally gated on it (fail closed).
66    pub fn with_policy_engine(
67        mut self,
68        engine: Arc<dyn crate::policy_engine::PolicyEngine>,
69    ) -> Self {
70        self.policy_engine = Some(engine);
71        self
72    }
73
74    /// E3 ingress gate for the send path. No-op without an engine. Denials
75    /// surface as `POLICY_DENIED` acks (fail closed, including unrecognized
76    /// decisions — `PolicyDecision` is `#[non_exhaustive]`).
77    async fn enforce_ingress_policy(
78        &self,
79        identity: &crate::security::AuthIdentity,
80        env: &Envelope,
81    ) -> Result<(), MacpError> {
82        let Some(engine) = &self.policy_engine else {
83            return Ok(());
84        };
85        let decision = if env.message_type == "SessionStart" {
86            engine
87                .evaluate_session_start(identity, &env.mode, env)
88                .await
89        } else if !env.session_id.is_empty() {
90            // Session-scoped message: give the engine the session context. A
91            // missing session falls through to the kernel's own
92            // UnknownSession handling.
93            match self.runtime.get_session_checked(&env.session_id).await {
94                Some(session) => engine.evaluate_message(identity, &session, env).await,
95                None => return Ok(()),
96            }
97        } else {
98            return Ok(());
99        };
100        match decision {
101            macp_core::policy::PolicyDecision::Allow { .. } => Ok(()),
102            macp_core::policy::PolicyDecision::Deny { reasons } => {
103                Err(MacpError::PolicyDenied { reasons })
104            }
105            other => {
106                tracing::warn!(decision = ?other, "unrecognized ingress policy decision");
107                Err(MacpError::PolicyDenied {
108                    reasons: vec!["unrecognized policy decision".into()],
109                })
110            }
111        }
112    }
113
114    fn validate_envelope_shape(&self, env: &Envelope) -> Result<(), MacpError> {
115        if env.macp_version != "1.0" {
116            return Err(MacpError::InvalidMacpVersion);
117        }
118        if env.message_type.is_empty() || env.message_id.is_empty() {
119            return Err(MacpError::InvalidEnvelope);
120        }
121        // RFC-MACP-0001: Signals MUST have empty session_id and empty mode.
122        // Progress messages MAY be ambient (empty session_id/mode) or session-scoped.
123        let is_ambient_type = env.message_type == "Signal" || env.message_type == "Progress";
124        if env.message_type == "Signal" {
125            if !env.session_id.is_empty() {
126                return Err(MacpError::InvalidEnvelope);
127            }
128            if !env.mode.trim().is_empty() {
129                return Err(MacpError::InvalidEnvelope);
130            }
131        }
132        if env.message_type == "Progress" && env.session_id.is_empty() {
133            // Ambient Progress: mode must also be empty
134            if !env.mode.trim().is_empty() {
135                return Err(MacpError::InvalidEnvelope);
136            }
137        }
138        if !is_ambient_type && env.session_id.is_empty() {
139            return Err(MacpError::InvalidEnvelope);
140        }
141        if !is_ambient_type && env.mode.trim().is_empty() {
142            return Err(MacpError::InvalidEnvelope);
143        }
144        // Session-scoped Progress must have non-empty mode (enforced above for non-ambient types,
145        // and ambient Progress with non-empty session_id falls through to here naturally)
146        if env.payload.len() > self.security.max_payload_bytes {
147            return Err(MacpError::PayloadTooLarge);
148        }
149        Ok(())
150    }
151
152    fn session_state_to_pb(state: &SessionState) -> i32 {
153        match state {
154            SessionState::Open => PbSessionState::Open.into(),
155            SessionState::Suspended => PbSessionState::Suspended.into(),
156            SessionState::Resolved => PbSessionState::Resolved.into(),
157            SessionState::Expired => PbSessionState::Expired.into(),
158            SessionState::Cancelled => PbSessionState::Cancelled.into(),
159        }
160    }
161
162    fn session_to_metadata(session: &crate::session::Session) -> SessionMetadata {
163        let participant_activity = session
164            .participant_message_counts
165            .iter()
166            .map(|(pid, count)| ParticipantActivity {
167                participant_id: pid.clone(),
168                last_message_at_unix_ms: session
169                    .participant_last_seen
170                    .get(pid)
171                    .copied()
172                    .unwrap_or(0),
173                message_count: *count,
174            })
175            .collect();
176        SessionMetadata {
177            session_id: session.session_id.clone(),
178            mode: session.mode.clone(),
179            state: Self::session_state_to_pb(&session.state),
180            started_at_unix_ms: session.started_at_unix_ms,
181            expires_at_unix_ms: session.ttl_expiry,
182            mode_version: session.mode_version.clone(),
183            configuration_version: session.configuration_version.clone(),
184            policy_version: session.policy_version.clone(),
185            participants: session.participants.clone(),
186            participant_activity,
187            initiator: session.initiator_sender.clone(),
188            context_id: session.context_id.clone(),
189            extension_keys: session.extensions.keys().cloned().collect(),
190        }
191    }
192
193    fn make_error_ack(e: &MacpError, env: &Envelope) -> Ack {
194        let details = Self::error_details_bytes(e);
195        Ack {
196            ok: false,
197            duplicate: false,
198            message_id: env.message_id.clone(),
199            session_id: env.session_id.clone(),
200            accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
201            session_state: PbSessionState::Unspecified.into(),
202            error: Some(PbMacpError {
203                code: e.error_code().into(),
204                message: e.to_string(),
205                session_id: env.session_id.clone(),
206                message_id: env.message_id.clone(),
207                details,
208            }),
209        }
210    }
211
212    /// Serialize structured error details as JSON bytes for the `details` field.
213    /// Currently only `PolicyDenied` carries additional detail (its reasons list).
214    fn error_details_bytes(e: &MacpError) -> Vec<u8> {
215        match e {
216            MacpError::PolicyDenied { reasons } => {
217                serde_json::to_vec(&serde_json::json!({ "reasons": reasons })).unwrap_or_default()
218            }
219            _ => vec![],
220        }
221    }
222
223    fn apply_authenticated_sender(
224        identity: &AuthIdentity,
225        mut env: Envelope,
226    ) -> Result<Envelope, MacpError> {
227        if !env.sender.is_empty() && env.sender != identity.sender {
228            return Err(MacpError::Unauthenticated);
229        }
230        env.sender = identity.sender.clone();
231        Ok(env)
232    }
233
234    async fn authenticate_send_request(
235        &self,
236        request: &Request<SendRequest>,
237        env: Envelope,
238    ) -> Result<(Envelope, Option<usize>), MacpError> {
239        let identity = self
240            .security
241            .authenticate_metadata(request.metadata())
242            .await?;
243        let env = Self::apply_authenticated_sender(&identity, env)?;
244        let is_session_start = env.message_type == "SessionStart";
245        self.security
246            .authorize_mode(&identity, &env.mode, is_session_start)?;
247        self.security
248            .enforce_rate_limit(&identity.sender, is_session_start)
249            .await?;
250        // External ingress policy engine (E3): after authentication and the
251        // built-in security checks, before kernel acceptance.
252        self.enforce_ingress_policy(&identity, &env).await?;
253        let max_open = if is_session_start {
254            identity.max_open_sessions
255        } else {
256            None
257        };
258        Ok((env, max_open))
259    }
260
261    async fn authenticate_session_access<T>(
262        &self,
263        request: &Request<T>,
264        session_id: &str,
265    ) -> Result<AuthIdentity, Status> {
266        let identity = self
267            .security
268            .authenticate_metadata(request.metadata())
269            .await
270            .map_err(Self::status_from_error)?;
271        let session = self
272            .runtime
273            .get_session_checked(session_id)
274            .await
275            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
276        let allowed = identity.is_observer
277            || session.initiator_sender == identity.sender
278            || session.participants.iter().any(|p| p == &identity.sender);
279        if !allowed {
280            return Err(Status::permission_denied(
281                "FORBIDDEN: session access denied",
282            ));
283        }
284        // External ingress policy engine (E3): may additionally restrict
285        // reads beyond the built-in membership check (fail closed).
286        if let Some(engine) = &self.policy_engine {
287            let decision = engine.evaluate_session_access(&identity, &session).await;
288            crate::policy_engine::require_allow(decision, "session access")?;
289        }
290        Ok(identity)
291    }
292
293    /// Subscribe-window dedupe (see `replay_dedup` in the stream loop): drop
294    /// a buffered envelope that was already delivered in the replay batch;
295    /// the first miss disarms the filter (the receiver is FIFO and every
296    /// in-window duplicate precedes the first post-snapshot envelope). Every
297    /// path that yields broadcast envelopes MUST route through this — the
298    /// drain loops bypassing it delivered subscribe-window duplicates.
299    fn should_skip_replayed(
300        replay_dedup: &mut Option<std::collections::HashSet<String>>,
301        envelope: &Envelope,
302    ) -> bool {
303        if let Some(seen) = replay_dedup.as_mut() {
304            if seen.remove(&envelope.message_id) {
305                return true;
306            }
307            *replay_dedup = None;
308        }
309        false
310    }
311
312    fn try_next_stream_event(
313        receiver: &mut Option<tokio::sync::broadcast::Receiver<Envelope>>,
314    ) -> Result<Option<Envelope>, Status> {
315        use tokio::sync::broadcast::error::TryRecvError;
316
317        let rx = match receiver.as_mut() {
318            Some(rx) => rx,
319            None => return Ok(None),
320        };
321
322        match rx.try_recv() {
323            Ok(envelope) => Ok(Some(envelope)),
324            Err(TryRecvError::Empty) => Ok(None),
325            Err(TryRecvError::Closed) => {
326                *receiver = None;
327                Ok(None)
328            }
329            Err(TryRecvError::Lagged(skipped)) => {
330                // Terminate the stream so the client knows it missed messages.
331                // Consistent with the async recv() path which also returns ResourceExhausted.
332                tracing::warn!(
333                    skipped,
334                    "StreamSession receiver fell behind; terminating stream"
335                );
336                Err(Status::resource_exhausted(format!(
337                    "StreamSession receiver fell behind by {skipped} envelopes"
338                )))
339            }
340        }
341    }
342
343    /// Process a single StreamSessionRequest frame.
344    ///
345    /// Returns `Ok(replay_envelopes)` — empty for normal sends, non-empty when
346    /// a subscribe frame triggers history replay (RFC-MACP-0006-A1).
347    async fn process_stream_request(
348        &self,
349        identity: &AuthIdentity,
350        req: StreamSessionRequest,
351        bound_session_id: &mut Option<String>,
352        session_events: &mut Option<tokio::sync::broadcast::Receiver<Envelope>>,
353    ) -> Result<Vec<Envelope>, Status> {
354        // RFC-MACP-0006-A1: Handle subscribe-only frame.
355        // When subscribe_session_id is set and envelope is absent, subscribe to
356        // the session's broadcast channel and replay accepted history.
357        if !req.subscribe_session_id.is_empty() {
358            if req.envelope.is_some() {
359                return Err(Status::invalid_argument(
360                    "StreamSessionRequest must not contain both envelope and subscribe_session_id",
361                ));
362            }
363            return self
364                .process_subscribe_frame(
365                    identity,
366                    &req.subscribe_session_id,
367                    req.after_sequence,
368                    bound_session_id,
369                    session_events,
370                )
371                .await;
372        }
373
374        let envelope = req.envelope.ok_or_else(|| {
375            Status::invalid_argument(
376                "StreamSessionRequest must contain an envelope or subscribe_session_id",
377            )
378        })?;
379
380        self.validate_envelope_shape(&envelope)
381            .map_err(Self::status_from_error)?;
382        if envelope.session_id.trim().is_empty() {
383            return Err(Status::invalid_argument(
384                "StreamSession requires a non-empty session_id",
385            ));
386        }
387        if envelope.mode.trim().is_empty() {
388            return Err(Status::invalid_argument(
389                "StreamSession requires a non-empty mode",
390            ));
391        }
392        if let Some(bound) = bound_session_id.as_ref() {
393            if bound != &envelope.session_id {
394                return Err(Status::invalid_argument(
395                    "StreamSession may only carry envelopes for one session_id",
396                ));
397            }
398        }
399
400        let envelope = Self::apply_authenticated_sender(identity, envelope)
401            .map_err(Self::status_from_error)?;
402        let is_session_start = envelope.message_type == "SessionStart";
403
404        if !is_session_start {
405            if let Some(session) = self.runtime.get_session_checked(&envelope.session_id).await {
406                if envelope.mode != session.mode {
407                    return Err(Status::invalid_argument(
408                        "INVALID_ENVELOPE: envelope mode does not match the bound session mode",
409                    ));
410                }
411                if session.state != SessionState::Open {
412                    return Err(Status::invalid_argument("SESSION_NOT_OPEN"));
413                }
414            } else if envelope.message_type == "Signal" {
415                return Err(Status::not_found(format!(
416                    "Session '{}' not found",
417                    envelope.session_id
418                )));
419            }
420        }
421
422        self.security
423            .authorize_mode(identity, &envelope.mode, is_session_start)
424            .map_err(Self::status_from_error)?;
425        // External ingress policy engine (E3): the stream path must be gated
426        // identically to unary Send — without this, a sender denied by an
427        // installed engine could simply switch transports.
428        self.enforce_ingress_policy(identity, &envelope)
429            .await
430            .map_err(Self::status_from_error)?;
431        self.security
432            .enforce_rate_limit(&identity.sender, is_session_start)
433            .await
434            .map_err(Self::status_from_error)?;
435
436        if session_events.is_none() {
437            *bound_session_id = Some(envelope.session_id.clone());
438            *session_events = Some(self.runtime.subscribe_session_stream(&envelope.session_id));
439        }
440
441        let max_open = if is_session_start {
442            identity.max_open_sessions
443        } else {
444            None
445        };
446        self.runtime
447            .process(&envelope, max_open)
448            .await
449            .map_err(Self::status_from_error)?;
450        Ok(vec![])
451    }
452
453    /// RFC-MACP-0006-A1: Process a subscribe-only frame.
454    /// Subscribes the stream to the session's broadcast channel and replays
455    /// accepted envelope history from `after_sequence` onwards.
456    async fn process_subscribe_frame(
457        &self,
458        identity: &AuthIdentity,
459        session_id: &str,
460        after_sequence: u64,
461        bound_session_id: &mut Option<String>,
462        session_events: &mut Option<tokio::sync::broadcast::Receiver<Envelope>>,
463    ) -> Result<Vec<Envelope>, Status> {
464        // Validate: only one session per stream
465        if let Some(bound) = bound_session_id.as_ref() {
466            if bound != session_id {
467                return Err(Status::invalid_argument(
468                    "StreamSession may only carry envelopes for one session_id",
469                ));
470            }
471        }
472
473        // Validate session exists
474        let session = self
475            .runtime
476            .get_session_checked(session_id)
477            .await
478            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
479
480        // Authorize: caller must be a declared participant, initiator, or observer
481        let allowed = identity.is_observer
482            || session.initiator_sender == identity.sender
483            || session.participants.iter().any(|p| p == &identity.sender);
484        if !allowed {
485            return Err(Status::permission_denied(
486                "FORBIDDEN: caller is not a declared participant or observer for this session",
487            ));
488        }
489        // External ingress policy engine (E3): stream-based history replay is
490        // a read and must be gated like GetSession (fail closed).
491        if let Some(engine) = &self.policy_engine {
492            let decision = engine.evaluate_session_access(identity, &session).await;
493            crate::policy_engine::require_allow(decision, "session access")?;
494        }
495
496        // Subscribe to live broadcast (if not already subscribed)
497        if session_events.is_none() {
498            *bound_session_id = Some(session_id.to_string());
499            *session_events = Some(self.runtime.subscribe_session_stream(session_id));
500        }
501
502        tracing::info!(
503            session_id = %session_id,
504            sender = %identity.sender,
505            after_sequence = after_sequence,
506            "passive subscribe: replaying session history"
507        );
508
509        // Replay accepted envelopes from LogStore
510        let replay = self
511            .runtime
512            .get_session_envelopes_after(session_id, after_sequence)
513            .await
514            .map_err(|base| {
515                Status::failed_precondition(format!(
516                    "session history before ordinal {base} was compacted; \
517                     resume with after_sequence >= {base} or re-read state via GetSession"
518                ))
519            })?;
520
521        Ok(replay)
522    }
523
524    fn build_stream_session_stream<S>(
525        &self,
526        identity: AuthIdentity,
527        inbound: S,
528    ) -> SessionResponseStream
529    where
530        S: futures_core::Stream<Item = Result<StreamSessionRequest, Status>> + Send + 'static,
531    {
532        use tokio::sync::broadcast;
533        use tokio_stream::StreamExt;
534
535        // Actions collected from tokio::select! arms to process outside the
536        // select scope, avoiding borrow and macro-expansion issues with `?`
537        // and `yield` inside select branches within try_stream!.
538        enum StreamAction {
539            ProcessRequest(StreamSessionRequest),
540            EmitEnvelope(Envelope),
541            ClientError(Status),
542            ClientDone,
543            EventsClosed,
544            Lagged(u64),
545        }
546
547        let server = self.clone();
548        let output = async_stream::try_stream! {
549            let mut inbound = Box::pin(inbound);
550            let mut bound_session_id: Option<String> = None;
551            let mut session_events: Option<broadcast::Receiver<Envelope>> = None;
552            // Subscribe-window dedup (RFC-0006 §3.2): the receiver is
553            // subscribed BEFORE the history snapshot is read, so an envelope
554            // accepted in that window arrives twice — once in the replay
555            // batch, once buffered on the receiver. Buffered events are FIFO
556            // and all in-window events precede post-snapshot ones, so we drop
557            // buffered envelopes whose message_id was replayed and disarm on
558            // the first miss.
559            let mut replay_dedup: Option<std::collections::HashSet<String>> = None;
560
561            loop {
562                if session_events.is_some() {
563                    let action = {
564                        let events = session_events.as_mut().unwrap();
565                        tokio::select! {
566                            maybe_req = inbound.next() => {
567                                match maybe_req {
568                                    Some(Ok(req)) => StreamAction::ProcessRequest(req),
569                                    Some(Err(status)) => StreamAction::ClientError(status),
570                                    None => StreamAction::ClientDone,
571                                }
572                            }
573                            recv_result = events.recv() => {
574                                match recv_result {
575                                    Ok(envelope) => StreamAction::EmitEnvelope(envelope),
576                                    Err(broadcast::error::RecvError::Closed) => StreamAction::EventsClosed,
577                                    Err(broadcast::error::RecvError::Lagged(n)) => StreamAction::Lagged(n),
578                                }
579                            }
580                        }
581                    };
582
583                    match action {
584                        StreamAction::ProcessRequest(req) => {
585                            match server
586                                .process_stream_request(
587                                    &identity,
588                                    req,
589                                    &mut bound_session_id,
590                                    &mut session_events,
591                                )
592                                .await
593                            {
594                                Ok(replay) => {
595                                    // RFC-MACP-0006-A1: yield replayed envelopes from subscribe
596                                    if !replay.is_empty() {
597                                        replay_dedup = Some(
598                                            replay.iter().map(|e| e.message_id.clone()).collect(),
599                                        );
600                                    }
601                                    for env in replay {
602                                        yield StreamSessionResponse {
603                                            response: Some(
604                                                crate::pb::stream_session_response::Response::Envelope(env),
605                                            ),
606                                        };
607                                    }
608                                }
609                                Err(status) if Self::is_stream_terminal_error(&status) => {
610                                    Err(status)?;
611                                }
612                                Err(status) => {
613                                    // RFC-MACP-0001: application-level validation errors
614                                    // are sent as inline MACPError; stream remains open.
615                                    yield StreamSessionResponse {
616                                        response: Some(
617                                            crate::pb::stream_session_response::Response::Error(
618                                                PbMacpError {
619                                                    code: status.message().to_string(),
620                                                    message: status.message().to_string(),
621                                                    session_id: bound_session_id.clone().unwrap_or_default(),
622                                                    message_id: String::new(),
623                                                    details: vec![],
624                                                },
625                                            ),
626                                        ),
627                                    };
628                                }
629                            }
630                            while let Some(envelope) = Self::try_next_stream_event(&mut session_events)? {
631                                if Self::should_skip_replayed(&mut replay_dedup, &envelope) {
632                                    continue;
633                                }
634                                yield StreamSessionResponse {
635                                    response: Some(
636                                        crate::pb::stream_session_response::Response::Envelope(envelope),
637                                    ),
638                                };
639                            }
640                        }
641                        StreamAction::EmitEnvelope(envelope) => {
642                            if Self::should_skip_replayed(&mut replay_dedup, &envelope) {
643                                continue;
644                            }
645                            yield StreamSessionResponse {
646                                response: Some(
647                                    crate::pb::stream_session_response::Response::Envelope(envelope),
648                                ),
649                            };
650                        }
651                        StreamAction::ClientError(status) => {
652                            Err(status)?;
653                        }
654                        StreamAction::ClientDone => {
655                            while let Some(envelope) = Self::try_next_stream_event(&mut session_events)? {
656                                if Self::should_skip_replayed(&mut replay_dedup, &envelope) {
657                                    continue;
658                                }
659                                yield StreamSessionResponse {
660                                    response: Some(
661                                        crate::pb::stream_session_response::Response::Envelope(envelope),
662                                    ),
663                                };
664                            }
665                            break;
666                        }
667                        StreamAction::EventsClosed => {
668                            session_events = None;
669                        }
670                        StreamAction::Lagged(skipped) => {
671                            Err(Status::resource_exhausted(format!(
672                                "StreamSession receiver fell behind by {skipped} envelopes"
673                            )))?;
674                        }
675                    }
676                } else {
677                    match inbound.next().await {
678                        Some(Ok(req)) => {
679                            match server
680                                .process_stream_request(
681                                    &identity,
682                                    req,
683                                    &mut bound_session_id,
684                                    &mut session_events,
685                                )
686                                .await
687                            {
688                                Ok(replay) => {
689                                    // RFC-MACP-0006-A1: yield replayed envelopes from subscribe
690                                    if !replay.is_empty() {
691                                        replay_dedup = Some(
692                                            replay.iter().map(|e| e.message_id.clone()).collect(),
693                                        );
694                                    }
695                                    for env in replay {
696                                        yield StreamSessionResponse {
697                                            response: Some(
698                                                crate::pb::stream_session_response::Response::Envelope(env),
699                                            ),
700                                        };
701                                    }
702                                }
703                                Err(status) if Self::is_stream_terminal_error(&status) => {
704                                    Err(status)?;
705                                }
706                                Err(status) => {
707                                    yield StreamSessionResponse {
708                                        response: Some(
709                                            crate::pb::stream_session_response::Response::Error(
710                                                PbMacpError {
711                                                    code: status.message().to_string(),
712                                                    message: status.message().to_string(),
713                                                    session_id: bound_session_id.clone().unwrap_or_default(),
714                                                    message_id: String::new(),
715                                                    details: vec![],
716                                                },
717                                            ),
718                                        ),
719                                    };
720                                }
721                            }
722                            while let Some(envelope) = Self::try_next_stream_event(&mut session_events)? {
723                                if Self::should_skip_replayed(&mut replay_dedup, &envelope) {
724                                    continue;
725                                }
726                                yield StreamSessionResponse {
727                                    response: Some(
728                                        crate::pb::stream_session_response::Response::Envelope(envelope),
729                                    ),
730                                };
731                            }
732                        }
733                        Some(Err(status)) => Err(status)?,
734                        None => break,
735                    }
736                }
737            }
738        };
739        Box::pin(output)
740    }
741
742    /// Returns true if the error should terminate a StreamSession stream.
743    /// Transport and binding errors terminate. Application-level validation
744    /// errors (from `runtime.process()`) are sent as inline MACPError per RFC-0001.
745    fn is_stream_terminal_error(status: &Status) -> bool {
746        matches!(
747            status.code(),
748            tonic::Code::Unauthenticated
749                | tonic::Code::Internal
750                | tonic::Code::ResourceExhausted
751                | tonic::Code::InvalidArgument
752                | tonic::Code::NotFound
753                | tonic::Code::AlreadyExists
754        )
755    }
756
757    fn status_from_error(err: MacpError) -> Status {
758        match err {
759            MacpError::Unauthenticated => Status::unauthenticated(err.to_string()),
760            MacpError::Forbidden => Status::permission_denied(err.to_string()),
761            MacpError::PayloadTooLarge => Status::resource_exhausted(err.to_string()),
762            MacpError::RateLimited => Status::resource_exhausted(err.to_string()),
763            MacpError::StorageFailed => Status::internal(err.to_string()),
764            MacpError::InvalidSessionId => Status::invalid_argument(err.to_string()),
765            MacpError::InvalidPolicyDefinition => Status::invalid_argument(err.to_string()),
766            MacpError::SessionAlreadyExists => Status::already_exists(err.to_string()),
767            MacpError::PolicyDenied { ref reasons } => {
768                let details = Self::error_details_bytes(&err);
769                let msg = if reasons.is_empty() {
770                    "PolicyDenied".to_string()
771                } else {
772                    format!("PolicyDenied: {}", reasons.join("; "))
773                };
774                let mut status = Status::failed_precondition(msg);
775                if !details.is_empty() {
776                    // Attach JSON details as binary metadata so clients can parse structured reasons.
777                    let val = tonic::metadata::MetadataValue::from_bytes(&details);
778                    status
779                        .metadata_mut()
780                        .insert_bin("macp-error-details-bin", val);
781                }
782                status
783            }
784            _ => Status::failed_precondition(err.to_string()),
785        }
786    }
787}
788
789#[tonic::async_trait]
790impl MacpRuntimeService for MacpServer {
791    async fn initialize(
792        &self,
793        request: Request<InitializeRequest>,
794    ) -> Result<Response<InitializeResponse>, Status> {
795        let req = request.into_inner();
796        if req.supported_protocol_versions.is_empty() {
797            return Err(Status::invalid_argument(
798                "INVALID_REQUEST: supported_protocol_versions must not be empty",
799            ));
800        }
801        if !req.supported_protocol_versions.iter().any(|v| v == "1.0") {
802            return Err(Status::failed_precondition(
803                "UNSUPPORTED_PROTOCOL_VERSION: no mutually supported protocol version",
804            ));
805        }
806
807        Ok(Response::new(InitializeResponse {
808            selected_protocol_version: "1.0".into(),
809            runtime_info: Some(RuntimeInfo {
810                name: "macp-runtime".into(),
811                title: "MACP Reference Runtime".into(),
812                // Tracks the workspace version; a literal here drifted once
813                // (still advertising 0.4.0 after the 0.5.0 release).
814                version: env!("CARGO_PKG_VERSION").into(),
815                description: "Reference implementation of the Multi-Agent Coordination Protocol"
816                    .into(),
817                website_url: String::new(),
818            }),
819            capabilities: Some(Capabilities {
820                sessions: Some(SessionsCapability { stream: true, list_sessions: true, watch_sessions: true }),
821                cancellation: Some(CancellationCapability {
822                    cancel_session: true,
823                }),
824                progress: Some(ProgressCapability { progress: true }),
825                manifest: Some(ManifestCapability { get_manifest: true }),
826                mode_registry: Some(ModeRegistryCapability {
827                    list_modes: true,
828                    list_changed: true,
829                }),
830                roots: Some(RootsCapability {
831                    // ListRoots is answerable (the root set is empty — a valid
832                    // state), but this runtime has no roots provider, so the
833                    // set never changes: do not advertise change notifications
834                    // (RFC-MACP-0006 §3.3 gates WatchRoots on list_changed).
835                    // Revisit when a roots provider lands (plans E2).
836                    list_roots: true,
837                    list_changed: false,
838                }),
839                policy_registry: Some(PolicyRegistryCapability {
840                    register_policy: !self.policies_read_only,
841                    list_policies: true,
842                    list_changed: true,
843                }),
844                experimental: Some(crate::pb::ExperimentalCapabilities {
845                    features: HashMap::from([
846                        ("ext_mode_lifecycle".into(), "true".into()),
847                    ]),
848                }),
849            }),
850            supported_modes: self.runtime.registered_mode_names(),
851            instructions: "Authenticate requests with Authorization: Bearer <token>. Use the unary Send RPC for all session messaging. For local development only, x-macp-agent-id may be enabled by configuration.".into(),
852        }))
853    }
854
855    async fn send(&self, request: Request<SendRequest>) -> Result<Response<SendResponse>, Status> {
856        let env = request
857            .get_ref()
858            .envelope
859            .clone()
860            .ok_or_else(|| Status::invalid_argument("SendRequest must contain an envelope"))?;
861
862        let result = async {
863            self.validate_envelope_shape(&env)?;
864            let (env, max_open) = self.authenticate_send_request(&request, env).await?;
865            self.runtime
866                .process(&env, max_open)
867                .await
868                .map(|process_result| (env, process_result))
869        }
870        .await;
871
872        let ack = match result {
873            Ok((env, process_result)) => Ack {
874                ok: true,
875                duplicate: process_result.duplicate,
876                message_id: env.message_id.clone(),
877                session_id: env.session_id.clone(),
878                accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
879                session_state: Self::session_state_to_pb(&process_result.session_state),
880                error: None,
881            },
882            Err(err) => {
883                let env = request.get_ref().envelope.clone().unwrap_or_default();
884                // Rejection counters were collected but never recorded before
885                // (permanently zero). Session-scoped rejections are counted
886                // per mode; commitments additionally under their own counter.
887                if !env.session_id.is_empty() {
888                    self.runtime.metrics().record_message_rejected(&env.mode);
889                    if env.message_type == "Commitment" {
890                        self.runtime.metrics().record_commitment_rejected(&env.mode);
891                    }
892                }
893                Self::make_error_ack(&err, &env)
894            }
895        };
896
897        Ok(Response::new(SendResponse { ack: Some(ack) }))
898    }
899
900    async fn get_session(
901        &self,
902        request: Request<GetSessionRequest>,
903    ) -> Result<Response<GetSessionResponse>, Status> {
904        let session_id = request.get_ref().session_id.clone();
905        let _identity = self
906            .authenticate_session_access(&request, &session_id)
907            .await?;
908        let session = self
909            .runtime
910            .get_session_checked(&session_id)
911            .await
912            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
913
914        Ok(Response::new(GetSessionResponse {
915            metadata: Some(Self::session_to_metadata(&session)),
916        }))
917    }
918
919    async fn cancel_session(
920        &self,
921        request: Request<CancelSessionRequest>,
922    ) -> Result<Response<CancelSessionResponse>, Status> {
923        let session_id = request.get_ref().session_id.clone();
924        let identity = self
925            .security
926            .authenticate_metadata(request.metadata())
927            .await
928            .map_err(Self::status_from_error)?;
929        let session = self
930            .runtime
931            .get_session_checked(&session_id)
932            .await
933            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
934        // RFC-MACP-0001: "Only the initiator and policy-delegated roles may cancel."
935        // CancelSession is a Core control-plane message — mode authorization does not apply.
936        if identity.sender != session.initiator_sender
937            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
938        {
939            return Err(Status::permission_denied(
940                "FORBIDDEN: only the session initiator or policy-delegated roles can cancel",
941            ));
942        }
943        let sender = identity.sender.clone();
944        let req = request.into_inner();
945        match self
946            .runtime
947            .cancel_session(&req.session_id, &req.reason, &sender)
948            .await
949        {
950            Ok(result) => Ok(Response::new(CancelSessionResponse {
951                ack: Some(Ack {
952                    ok: true,
953                    duplicate: false,
954                    message_id: String::new(),
955                    session_id: req.session_id,
956                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
957                    session_state: Self::session_state_to_pb(&result.session_state),
958                    error: None,
959                }),
960            })),
961            Err(err) => Ok(Response::new(CancelSessionResponse {
962                ack: Some(Ack {
963                    ok: false,
964                    duplicate: false,
965                    message_id: String::new(),
966                    session_id: req.session_id.clone(),
967                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
968                    session_state: PbSessionState::Unspecified.into(),
969                    error: Some(PbMacpError {
970                        code: err.error_code().into(),
971                        message: err.to_string(),
972                        session_id: req.session_id,
973                        message_id: String::new(),
974                        details: vec![],
975                    }),
976                }),
977            })),
978        }
979    }
980
981    async fn suspend_session(
982        &self,
983        request: Request<SuspendSessionRequest>,
984    ) -> Result<Response<SuspendSessionResponse>, Status> {
985        let session_id = request.get_ref().session_id.clone();
986        let identity = self
987            .security
988            .authenticate_metadata(request.metadata())
989            .await
990            .map_err(Self::status_from_error)?;
991        let session = self
992            .runtime
993            .get_session_checked(&session_id)
994            .await
995            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
996        // RFC-MACP-0001 §7.5: same authority model as CancelSession — initiator
997        // or policy-delegated roles only; mode authorization does not apply.
998        if identity.sender != session.initiator_sender
999            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1000        {
1001            return Err(Status::permission_denied(
1002                "FORBIDDEN: only the session initiator or policy-delegated roles can suspend",
1003            ));
1004        }
1005        let sender = identity.sender.clone();
1006        let req = request.into_inner();
1007        match self
1008            .runtime
1009            .suspend_session(&req.session_id, &req.reason, &sender)
1010            .await
1011        {
1012            Ok(result) => Ok(Response::new(SuspendSessionResponse {
1013                ack: Some(Ack {
1014                    ok: true,
1015                    duplicate: false,
1016                    message_id: String::new(),
1017                    session_id: req.session_id,
1018                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1019                    session_state: Self::session_state_to_pb(&result.session_state),
1020                    error: None,
1021                }),
1022            })),
1023            Err(err) => Ok(Response::new(SuspendSessionResponse {
1024                ack: Some(Ack {
1025                    ok: false,
1026                    duplicate: false,
1027                    message_id: String::new(),
1028                    session_id: req.session_id.clone(),
1029                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1030                    session_state: PbSessionState::Unspecified.into(),
1031                    error: Some(PbMacpError {
1032                        code: err.error_code().into(),
1033                        message: err.to_string(),
1034                        session_id: req.session_id,
1035                        message_id: String::new(),
1036                        details: vec![],
1037                    }),
1038                }),
1039            })),
1040        }
1041    }
1042
1043    async fn resume_session(
1044        &self,
1045        request: Request<ResumeSessionRequest>,
1046    ) -> Result<Response<ResumeSessionResponse>, Status> {
1047        let session_id = request.get_ref().session_id.clone();
1048        let identity = self
1049            .security
1050            .authenticate_metadata(request.metadata())
1051            .await
1052            .map_err(Self::status_from_error)?;
1053        let session = self
1054            .runtime
1055            .get_session_checked(&session_id)
1056            .await
1057            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1058        if identity.sender != session.initiator_sender
1059            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1060        {
1061            return Err(Status::permission_denied(
1062                "FORBIDDEN: only the session initiator or policy-delegated roles can resume",
1063            ));
1064        }
1065        let sender = identity.sender.clone();
1066        let req = request.into_inner();
1067        match self
1068            .runtime
1069            .resume_session(&req.session_id, &req.reason, &sender)
1070            .await
1071        {
1072            Ok(result) => Ok(Response::new(ResumeSessionResponse {
1073                ack: Some(Ack {
1074                    ok: true,
1075                    duplicate: false,
1076                    message_id: String::new(),
1077                    session_id: req.session_id,
1078                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1079                    session_state: Self::session_state_to_pb(&result.session_state),
1080                    error: None,
1081                }),
1082            })),
1083            Err(err) => Ok(Response::new(ResumeSessionResponse {
1084                ack: Some(Ack {
1085                    ok: false,
1086                    duplicate: false,
1087                    message_id: String::new(),
1088                    session_id: req.session_id.clone(),
1089                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1090                    session_state: PbSessionState::Unspecified.into(),
1091                    error: Some(PbMacpError {
1092                        code: err.error_code().into(),
1093                        message: err.to_string(),
1094                        session_id: req.session_id,
1095                        message_id: String::new(),
1096                        details: vec![],
1097                    }),
1098                }),
1099            })),
1100        }
1101    }
1102
1103    async fn get_manifest(
1104        &self,
1105        request: Request<GetManifestRequest>,
1106    ) -> Result<Response<GetManifestResponse>, Status> {
1107        let req = request.into_inner();
1108        if !req.agent_id.is_empty() && req.agent_id != "macp-runtime" {
1109            return Err(Status::not_found(format!(
1110                "Agent '{}' not found",
1111                req.agent_id
1112            )));
1113        }
1114
1115        Ok(Response::new(GetManifestResponse {
1116            manifest: Some(crate::pb::AgentManifest {
1117                agent_id: "macp-runtime".into(),
1118                title: "MACP Reference Runtime".into(),
1119                description: "Reference implementation of MACP".into(),
1120                supported_modes: self.runtime.registered_mode_names(),
1121                input_content_types: vec!["application/macp-envelope+proto".into()],
1122                output_content_types: vec!["application/macp-envelope+proto".into()],
1123                metadata: HashMap::new(),
1124                // Empty: unary-first profile has no dedicated transport endpoints.
1125                transport_endpoints: vec![],
1126            }),
1127        }))
1128    }
1129
1130    async fn list_modes(
1131        &self,
1132        _request: Request<ListModesRequest>,
1133    ) -> Result<Response<ListModesResponse>, Status> {
1134        Ok(Response::new(ListModesResponse {
1135            modes: self.runtime.standard_mode_descriptors(),
1136        }))
1137    }
1138
1139    async fn list_roots(
1140        &self,
1141        _request: Request<ListRootsRequest>,
1142    ) -> Result<Response<ListRootsResponse>, Status> {
1143        Ok(Response::new(ListRootsResponse { roots: vec![] }))
1144    }
1145
1146    type StreamSessionStream = SessionResponseStream;
1147
1148    async fn stream_session(
1149        &self,
1150        request: Request<tonic::Streaming<StreamSessionRequest>>,
1151    ) -> Result<Response<Self::StreamSessionStream>, Status> {
1152        let identity = self
1153            .security
1154            .authenticate_metadata(request.metadata())
1155            .await
1156            .map_err(Self::status_from_error)?;
1157        let inbound = request.into_inner();
1158        Ok(Response::new(
1159            self.build_stream_session_stream(identity, inbound),
1160        ))
1161    }
1162
1163    type WatchModeRegistryStream = std::pin::Pin<
1164        Box<dyn futures_core::Stream<Item = Result<WatchModeRegistryResponse, Status>> + Send>,
1165    >;
1166
1167    async fn watch_mode_registry(
1168        &self,
1169        _request: Request<WatchModeRegistryRequest>,
1170    ) -> Result<Response<Self::WatchModeRegistryStream>, Status> {
1171        let mut rx = self.runtime.subscribe_mode_changes();
1172        let stream = async_stream::try_stream! {
1173            // Send initial state
1174            yield WatchModeRegistryResponse {
1175                change: Some(crate::pb::RegistryChanged {
1176                    registry: "modes".into(),
1177                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1178                }),
1179            };
1180            // Wait for changes from register/unregister/promote
1181            while rx.recv().await.is_ok() {
1182                yield WatchModeRegistryResponse {
1183                    change: Some(crate::pb::RegistryChanged {
1184                        registry: "modes".into(),
1185                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1186                    }),
1187                };
1188            }
1189        };
1190        Ok(Response::new(Box::pin(stream)))
1191    }
1192
1193    type WatchRootsStream = std::pin::Pin<
1194        Box<dyn futures_core::Stream<Item = Result<WatchRootsResponse, Status>> + Send>,
1195    >;
1196
1197    async fn watch_roots(
1198        &self,
1199        _request: Request<WatchRootsRequest>,
1200    ) -> Result<Response<Self::WatchRootsStream>, Status> {
1201        let initial = WatchRootsResponse {
1202            change: Some(crate::pb::RootsChanged {
1203                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1204            }),
1205        };
1206        let stream = async_stream::try_stream! {
1207            yield initial;
1208            // Roots are static — keep the stream open but idle.
1209            std::future::pending::<()>().await;
1210        };
1211        Ok(Response::new(Box::pin(stream)))
1212    }
1213
1214    type WatchSignalsStream = std::pin::Pin<
1215        Box<dyn futures_core::Stream<Item = Result<WatchSignalsResponse, Status>> + Send>,
1216    >;
1217
1218    type WatchSessionsStream = std::pin::Pin<
1219        Box<dyn futures_core::Stream<Item = Result<WatchSessionsResponse, Status>> + Send>,
1220    >;
1221
1222    async fn watch_signals(
1223        &self,
1224        request: Request<WatchSignalsRequest>,
1225    ) -> Result<Response<Self::WatchSignalsStream>, Status> {
1226        // Ambient signals carry agent-generated payload data; subscribing is
1227        // gated on authentication like the session-observation surfaces.
1228        // (RFC-0004 §4.1 constrains unauthenticated *producers*; requiring
1229        // authenticated subscribers is this runtime's hardening posture.)
1230        let _identity = self
1231            .security
1232            .authenticate_metadata(request.metadata())
1233            .await
1234            .map_err(Self::status_from_error)?;
1235        let mut rx = self.runtime.subscribe_signals();
1236        let stream = async_stream::try_stream! {
1237            loop {
1238                match rx.recv().await {
1239                    Ok(envelope) => {
1240                        yield WatchSignalsResponse {
1241                            envelope: Some(envelope),
1242                        };
1243                    }
1244                    // Surface lag instead of silently ending the stream: a
1245                    // slow consumer must be able to distinguish "no traffic"
1246                    // from "events dropped" (mirrors StreamSession).
1247                    Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1248                        Err(Status::resource_exhausted(format!(
1249                            "WatchSignals receiver fell behind by {skipped} signals"
1250                        )))?;
1251                    }
1252                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1253                }
1254            }
1255        };
1256        Ok(Response::new(Box::pin(stream)))
1257    }
1258
1259    // Session lifecycle observation RPCs
1260
1261    async fn list_sessions(
1262        &self,
1263        request: Request<ListSessionsRequest>,
1264    ) -> Result<Response<ListSessionsResponse>, Status> {
1265        let _identity = self
1266            .security
1267            .authenticate_metadata(request.metadata())
1268            .await
1269            .map_err(Self::status_from_error)?;
1270        let sessions = self.runtime.registry.get_all_sessions().await;
1271        let metadata: Vec<SessionMetadata> =
1272            sessions.iter().map(Self::session_to_metadata).collect();
1273        Ok(Response::new(ListSessionsResponse { sessions: metadata }))
1274    }
1275
1276    async fn watch_sessions(
1277        &self,
1278        request: Request<WatchSessionsRequest>,
1279    ) -> Result<Response<Self::WatchSessionsStream>, Status> {
1280        let _identity = self
1281            .security
1282            .authenticate_metadata(request.metadata())
1283            .await
1284            .map_err(Self::status_from_error)?;
1285        let mut rx = self.runtime.subscribe_session_lifecycle();
1286        let runtime = Arc::clone(&self.runtime);
1287        let stream = async_stream::try_stream! {
1288            // Initial sync: emit all current sessions as CREATED events. The
1289            // lifecycle bus was subscribed *before* this snapshot (so no event
1290            // is missed); any Created event buffered in that window would
1291            // duplicate a snapshot entry — session IDs are create-once, so we
1292            // dedupe buffered Created events against the synced set below.
1293            let sessions = runtime.registry.get_all_sessions().await;
1294            let mut synced: std::collections::HashSet<String> =
1295                std::collections::HashSet::with_capacity(sessions.len());
1296            for session in &sessions {
1297                synced.insert(session.session_id.clone());
1298                yield WatchSessionsResponse {
1299                    event: Some(SessionLifecycleEvent {
1300                        event_type: session_lifecycle_event::EventType::Created.into(),
1301                        session: Some(Self::session_to_metadata(session)),
1302                        observed_at_unix_ms: session.started_at_unix_ms,
1303                    }),
1304                };
1305            }
1306            // Stream lifecycle transitions
1307            loop {
1308                let event = match rx.recv().await {
1309                    Ok(event) => event,
1310                    Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1311                        Err(Status::resource_exhausted(format!(
1312                            "WatchSessions receiver fell behind by {skipped} events"
1313                        )))?;
1314                        break;
1315                    }
1316                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1317                };
1318                let (event_type, sid) = match &event {
1319                    crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1320                        (session_lifecycle_event::EventType::Created, session_id.clone()),
1321                    crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1322                        (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1323                    crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1324                        (session_lifecycle_event::EventType::Expired, session_id.clone()),
1325                    crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1326                        (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1327                    crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1328                        (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1329                    crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1330                        (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1331                };
1332                // Skip the buffered duplicate of an initial-sync entry;
1333                // non-Created events for synced sessions are new information
1334                // and pass through.
1335                if event_type == session_lifecycle_event::EventType::Created
1336                    && !synced.insert(sid.clone())
1337                {
1338                    continue;
1339                }
1340                let session_meta = runtime.registry.get_session(&sid).await
1341                    .map(|s| Self::session_to_metadata(&s));
1342                yield WatchSessionsResponse {
1343                    event: Some(SessionLifecycleEvent {
1344                        event_type: event_type.into(),
1345                        session: session_meta,
1346                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1347                    }),
1348                };
1349            }
1350        };
1351        Ok(Response::new(Box::pin(stream)))
1352    }
1353
1354    // Extension mode lifecycle RPCs
1355
1356    async fn list_ext_modes(
1357        &self,
1358        _request: Request<ListExtModesRequest>,
1359    ) -> Result<Response<ListExtModesResponse>, Status> {
1360        Ok(Response::new(ListExtModesResponse {
1361            modes: self.runtime.extension_mode_descriptors(),
1362        }))
1363    }
1364
1365    async fn register_ext_mode(
1366        &self,
1367        request: Request<RegisterExtModeRequest>,
1368    ) -> Result<Response<RegisterExtModeResponse>, Status> {
1369        let identity = self
1370            .security
1371            .authenticate_metadata(request.metadata())
1372            .await
1373            .map_err(Self::status_from_error)?;
1374        self.security
1375            .authorize_mode_registry(&identity)
1376            .map_err(Self::status_from_error)?;
1377        let req = request.into_inner();
1378        let descriptor = req
1379            .mode_descriptor
1380            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1381        match self.runtime.register_extension(descriptor) {
1382            Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1383                ok: true,
1384                error: String::new(),
1385            })),
1386            Err(e) => Ok(Response::new(RegisterExtModeResponse {
1387                ok: false,
1388                error: e,
1389            })),
1390        }
1391    }
1392
1393    async fn unregister_ext_mode(
1394        &self,
1395        request: Request<UnregisterExtModeRequest>,
1396    ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1397        let identity = self
1398            .security
1399            .authenticate_metadata(request.metadata())
1400            .await
1401            .map_err(Self::status_from_error)?;
1402        self.security
1403            .authorize_mode_registry(&identity)
1404            .map_err(Self::status_from_error)?;
1405        let req = request.into_inner();
1406        match self.runtime.unregister_extension(&req.mode) {
1407            Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1408                ok: true,
1409                error: String::new(),
1410            })),
1411            Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1412                ok: false,
1413                error: e,
1414            })),
1415        }
1416    }
1417
1418    async fn promote_mode(
1419        &self,
1420        request: Request<PromoteModeRequest>,
1421    ) -> Result<Response<PromoteModeResponse>, Status> {
1422        let identity = self
1423            .security
1424            .authenticate_metadata(request.metadata())
1425            .await
1426            .map_err(Self::status_from_error)?;
1427        self.security
1428            .authorize_mode_registry(&identity)
1429            .map_err(Self::status_from_error)?;
1430        let req = request.into_inner();
1431        let new_name = if req.promoted_mode_name.is_empty() {
1432            None
1433        } else {
1434            Some(req.promoted_mode_name.as_str())
1435        };
1436        match self.runtime.promote_mode(&req.mode, new_name) {
1437            Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1438                ok: true,
1439                error: String::new(),
1440                mode: final_name,
1441            })),
1442            Err(e) => Ok(Response::new(PromoteModeResponse {
1443                ok: false,
1444                error: e,
1445                mode: String::new(),
1446            })),
1447        }
1448    }
1449
1450    // ── Governance policy lifecycle RPCs (RFC-MACP-0012) ────────────
1451
1452    async fn register_policy(
1453        &self,
1454        request: Request<RegisterPolicyRequest>,
1455    ) -> Result<Response<RegisterPolicyResponse>, Status> {
1456        if self.policies_read_only {
1457            return Err(Status::failed_precondition(
1458                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1459            ));
1460        }
1461        let identity = self
1462            .security
1463            .authenticate_metadata(request.metadata())
1464            .await
1465            .map_err(Self::status_from_error)?;
1466        self.security
1467            .authorize_mode_registry(&identity)
1468            .map_err(Self::status_from_error)?;
1469        let req = request.into_inner();
1470        let descriptor = req
1471            .policy_descriptor
1472            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1473        let definition = Self::policy_descriptor_to_definition(&descriptor);
1474        match self.runtime.register_policy(definition) {
1475            Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1476                ok: true,
1477                error: String::new(),
1478            })),
1479            Err(e) => Ok(Response::new(RegisterPolicyResponse {
1480                ok: false,
1481                error: e,
1482            })),
1483        }
1484    }
1485
1486    async fn unregister_policy(
1487        &self,
1488        request: Request<UnregisterPolicyRequest>,
1489    ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1490        if self.policies_read_only {
1491            return Err(Status::failed_precondition(
1492                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1493            ));
1494        }
1495        let identity = self
1496            .security
1497            .authenticate_metadata(request.metadata())
1498            .await
1499            .map_err(Self::status_from_error)?;
1500        self.security
1501            .authorize_mode_registry(&identity)
1502            .map_err(Self::status_from_error)?;
1503        let req = request.into_inner();
1504        match self.runtime.unregister_policy(&req.policy_id) {
1505            Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1506                ok: true,
1507                error: String::new(),
1508            })),
1509            Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1510                ok: false,
1511                error: e,
1512            })),
1513        }
1514    }
1515
1516    async fn get_policy(
1517        &self,
1518        request: Request<GetPolicyRequest>,
1519    ) -> Result<Response<GetPolicyResponse>, Status> {
1520        let _identity = self
1521            .security
1522            .authenticate_metadata(request.metadata())
1523            .await
1524            .map_err(Self::status_from_error)?;
1525        let req = request.into_inner();
1526        let policy = self
1527            .runtime
1528            .get_policy(&req.policy_id)
1529            .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1530        Ok(Response::new(GetPolicyResponse {
1531            policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1532        }))
1533    }
1534
1535    async fn list_policies(
1536        &self,
1537        request: Request<ListPoliciesRequest>,
1538    ) -> Result<Response<ListPoliciesResponse>, Status> {
1539        let _identity = self
1540            .security
1541            .authenticate_metadata(request.metadata())
1542            .await
1543            .map_err(Self::status_from_error)?;
1544        let req = request.into_inner();
1545        let mode_filter = if req.mode.is_empty() {
1546            None
1547        } else {
1548            Some(req.mode.as_str())
1549        };
1550        let policies = self.runtime.list_policies(mode_filter);
1551        let descriptors = policies
1552            .iter()
1553            .map(Self::policy_definition_to_descriptor)
1554            .collect();
1555        Ok(Response::new(ListPoliciesResponse { descriptors }))
1556    }
1557
1558    type WatchPoliciesStream = std::pin::Pin<
1559        Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1560    >;
1561
1562    async fn watch_policies(
1563        &self,
1564        _request: Request<WatchPoliciesRequest>,
1565    ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1566        let mut rx = self.runtime.subscribe_policy_changes();
1567        let runtime = Arc::clone(&self.runtime);
1568        let stream = async_stream::try_stream! {
1569            // Send initial state
1570            let policies = runtime.list_policies(None);
1571            let descriptors: Vec<PolicyDescriptor> = policies
1572                .iter()
1573                .map(MacpServer::policy_definition_to_descriptor)
1574                .collect();
1575            yield WatchPoliciesResponse {
1576                descriptors,
1577                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1578            };
1579            // Wait for changes
1580            while rx.recv().await.is_ok() {
1581                let policies = runtime.list_policies(None);
1582                let descriptors: Vec<PolicyDescriptor> = policies
1583                    .iter()
1584                    .map(MacpServer::policy_definition_to_descriptor)
1585                    .collect();
1586                yield WatchPoliciesResponse {
1587                    descriptors,
1588                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1589                };
1590            }
1591        };
1592        Ok(Response::new(Box::pin(stream)))
1593    }
1594}
1595
1596// ── Policy type conversion helpers ──────────────────────────────────
1597
1598impl MacpServer {
1599    fn policy_descriptor_to_definition(
1600        descriptor: &PolicyDescriptor,
1601    ) -> crate::policy::PolicyDefinition {
1602        let rules: serde_json::Value = if descriptor.rules.is_empty() {
1603            serde_json::json!({})
1604        } else {
1605            serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1606        };
1607        crate::policy::PolicyDefinition {
1608            policy_id: descriptor.policy_id.clone(),
1609            mode: descriptor.mode.clone(),
1610            description: descriptor.description.clone(),
1611            rules,
1612            schema_version: descriptor.schema_version,
1613        }
1614    }
1615
1616    fn policy_definition_to_descriptor(
1617        definition: &crate::policy::PolicyDefinition,
1618    ) -> PolicyDescriptor {
1619        PolicyDescriptor {
1620            policy_id: definition.policy_id.clone(),
1621            mode: definition.mode.clone(),
1622            description: definition.description.clone(),
1623            rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1624            schema_version: definition.schema_version,
1625            registered_at_unix_ms: 0,
1626        }
1627    }
1628}
1629
1630#[cfg(test)]
1631mod tests {
1632    use super::*;
1633    use crate::log_store::LogStore;
1634    use crate::pb::SessionStartPayload;
1635    use crate::registry::SessionRegistry;
1636    use chrono::Utc;
1637    use prost::Message;
1638
1639    fn new_sid() -> String {
1640        uuid::Uuid::new_v4().as_hyphenated().to_string()
1641    }
1642
1643    fn make_server() -> (MacpServer, Arc<Runtime>) {
1644        let storage: Arc<dyn crate::storage::StorageBackend> =
1645            Arc::new(crate::storage::MemoryBackend);
1646        let registry = Arc::new(SessionRegistry::new());
1647        let log_store = Arc::new(LogStore::new());
1648        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1649        let server = MacpServer::new(runtime.clone(), SecurityLayer::dev_mode());
1650        (server, runtime)
1651    }
1652
1653    fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1654        let mut req = Request::new(SendRequest {
1655            envelope: Some(env),
1656        });
1657        req.metadata_mut()
1658            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1659        req
1660    }
1661
1662    async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1663        let resp = server.send(send_req(sender, env)).await.unwrap();
1664        resp.into_inner().ack.unwrap()
1665    }
1666
1667    fn start_payload() -> Vec<u8> {
1668        SessionStartPayload {
1669            intent: "intent".into(),
1670            participants: vec!["agent://fraud".into()],
1671            mode_version: "1.0.0".into(),
1672            configuration_version: "cfg-1".into(),
1673            policy_version: String::new(),
1674            ttl_ms: 1000,
1675            context_id: String::new(),
1676            extensions: std::collections::HashMap::new(),
1677            roots: vec![],
1678            max_suspend_ms: 0,
1679        }
1680        .encode_to_vec()
1681    }
1682
1683    #[tokio::test]
1684    async fn sender_is_derived_from_authenticated_metadata() {
1685        let (server, runtime) = make_server();
1686        let sid = new_sid();
1687        let ack = do_send(
1688            &server,
1689            "agent://orchestrator",
1690            Envelope {
1691                macp_version: "1.0".into(),
1692                mode: "macp.mode.decision.v1".into(),
1693                message_type: "SessionStart".into(),
1694                message_id: "m1".into(),
1695                session_id: sid.clone(),
1696                sender: String::new(),
1697                timestamp_unix_ms: Utc::now().timestamp_millis(),
1698                payload: start_payload(),
1699            },
1700        )
1701        .await;
1702        assert!(ack.ok);
1703        let session = runtime.get_session_checked(&sid).await.unwrap();
1704        assert_eq!(session.initiator_sender, "agent://orchestrator");
1705    }
1706
1707    #[tokio::test]
1708    async fn spoofed_sender_is_rejected() {
1709        let (server, _) = make_server();
1710        let sid = new_sid();
1711        let ack = do_send(
1712            &server,
1713            "agent://orchestrator",
1714            Envelope {
1715                macp_version: "1.0".into(),
1716                mode: "macp.mode.decision.v1".into(),
1717                message_type: "SessionStart".into(),
1718                message_id: "m1".into(),
1719                session_id: sid,
1720                sender: "agent://spoof".into(),
1721                timestamp_unix_ms: Utc::now().timestamp_millis(),
1722                payload: start_payload(),
1723            },
1724        )
1725        .await;
1726        assert!(!ack.ok);
1727        assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1728    }
1729
1730    #[tokio::test]
1731    async fn get_session_requires_session_membership() {
1732        let (server, _) = make_server();
1733        let sid = new_sid();
1734        let ack = do_send(
1735            &server,
1736            "agent://orchestrator",
1737            Envelope {
1738                macp_version: "1.0".into(),
1739                mode: "macp.mode.decision.v1".into(),
1740                message_type: "SessionStart".into(),
1741                message_id: "m1".into(),
1742                session_id: sid.clone(),
1743                sender: String::new(),
1744                timestamp_unix_ms: Utc::now().timestamp_millis(),
1745                payload: start_payload(),
1746            },
1747        )
1748        .await;
1749        assert!(ack.ok);
1750
1751        let mut req = Request::new(GetSessionRequest { session_id: sid });
1752        req.metadata_mut().insert(
1753            "authorization",
1754            format!("Bearer {}", "agent://outsider").parse().unwrap(),
1755        );
1756        let err = server.get_session(req).await.unwrap_err();
1757        assert_eq!(err.code(), tonic::Code::PermissionDenied);
1758    }
1759
1760    #[tokio::test]
1761    async fn register_ext_mode_requires_authenticated_registry_permission() {
1762        let storage: Arc<dyn crate::storage::StorageBackend> =
1763            Arc::new(crate::storage::MemoryBackend);
1764        let registry = Arc::new(SessionRegistry::new());
1765        let log_store = Arc::new(LogStore::new());
1766        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1767        let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1768        let server = MacpServer::new(runtime, security);
1769
1770        let req = Request::new(RegisterExtModeRequest {
1771            mode_descriptor: Some(crate::pb::ModeDescriptor {
1772                mode: "ext.custom.v1".into(),
1773                mode_version: "1.0.0".into(),
1774                message_types: vec!["SessionStart".into(), "Commitment".into()],
1775                ..Default::default()
1776            }),
1777        });
1778        let err = server.register_ext_mode(req).await.unwrap_err();
1779        assert_eq!(err.code(), tonic::Code::Unauthenticated);
1780    }
1781
1782    fn stream_identity(sender: &str) -> AuthIdentity {
1783        AuthIdentity {
1784            sender: sender.into(),
1785            allowed_modes: None,
1786            can_start_sessions: true,
1787            max_open_sessions: None,
1788            can_manage_mode_registry: false,
1789            is_observer: false,
1790        }
1791    }
1792
1793    #[tokio::test]
1794    async fn stream_session_emits_accepted_envelopes_only() {
1795        use tokio_stream::{iter, StreamExt};
1796
1797        let (server, _) = make_server();
1798        let sid = new_sid();
1799        let requests = iter(vec![Ok(StreamSessionRequest {
1800            subscribe_session_id: String::new(),
1801            after_sequence: 0,
1802            envelope: Some(Envelope {
1803                macp_version: "1.0".into(),
1804                mode: "macp.mode.decision.v1".into(),
1805                message_type: "SessionStart".into(),
1806                message_id: "m1".into(),
1807                session_id: sid.clone(),
1808                sender: String::new(),
1809                timestamp_unix_ms: Utc::now().timestamp_millis(),
1810                payload: start_payload(),
1811            }),
1812        })]);
1813
1814        let mut stream =
1815            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1816
1817        let response = stream.next().await.unwrap().unwrap();
1818        let envelope = match response.response.unwrap() {
1819            crate::pb::stream_session_response::Response::Envelope(e) => e,
1820            _ => panic!("expected envelope"),
1821        };
1822        assert_eq!(envelope.message_type, "SessionStart");
1823        assert_eq!(envelope.message_id, "m1");
1824        assert!(stream.next().await.is_none());
1825    }
1826
1827    #[tokio::test]
1828    async fn stream_session_rejects_mixed_session_ids() {
1829        use tokio_stream::{iter, StreamExt};
1830
1831        let (server, _) = make_server();
1832        let sid1 = new_sid();
1833        let sid2 = new_sid();
1834        let requests = iter(vec![
1835            Ok(StreamSessionRequest {
1836                subscribe_session_id: String::new(),
1837                after_sequence: 0,
1838                envelope: Some(Envelope {
1839                    macp_version: "1.0".into(),
1840                    mode: "macp.mode.decision.v1".into(),
1841                    message_type: "SessionStart".into(),
1842                    message_id: "m1".into(),
1843                    session_id: sid1.clone(),
1844                    sender: String::new(),
1845                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1846                    payload: start_payload(),
1847                }),
1848            }),
1849            Ok(StreamSessionRequest {
1850                subscribe_session_id: String::new(),
1851                after_sequence: 0,
1852                envelope: Some(Envelope {
1853                    macp_version: "1.0".into(),
1854                    mode: "macp.mode.decision.v1".into(),
1855                    message_type: "SessionStart".into(),
1856                    message_id: "m2".into(),
1857                    session_id: sid2,
1858                    sender: String::new(),
1859                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1860                    payload: start_payload(),
1861                }),
1862            }),
1863        ]);
1864
1865        let mut stream =
1866            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1867
1868        let first = stream.next().await.unwrap().unwrap();
1869        let first_env = match first.response.unwrap() {
1870            crate::pb::stream_session_response::Response::Envelope(e) => e,
1871            _ => panic!("expected envelope"),
1872        };
1873        assert_eq!(first_env.session_id, sid1);
1874        let err = stream.next().await.unwrap().unwrap_err();
1875        assert_eq!(err.code(), tonic::Code::InvalidArgument);
1876    }
1877
1878    #[tokio::test]
1879    async fn list_modes_returns_standard_modes() {
1880        let (server, _) = make_server();
1881        let resp = server
1882            .list_modes(Request::new(ListModesRequest {}))
1883            .await
1884            .unwrap();
1885        let names: Vec<String> = resp
1886            .into_inner()
1887            .modes
1888            .iter()
1889            .map(|m| m.mode.clone())
1890            .collect();
1891        assert_eq!(names.len(), 5);
1892        assert!(names.contains(&"macp.mode.decision.v1".to_string()));
1893        assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
1894        assert!(names.contains(&"macp.mode.task.v1".to_string()));
1895        assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
1896        assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
1897        // multi_round is now an extension, not in ListModes
1898        assert!(!names.contains(&"ext.multi_round.v1".to_string()));
1899    }
1900
1901    #[tokio::test]
1902    async fn list_ext_modes_returns_extensions() {
1903        let (server, _) = make_server();
1904        let resp = server
1905            .list_ext_modes(Request::new(ListExtModesRequest {}))
1906            .await
1907            .unwrap();
1908        let names: Vec<String> = resp
1909            .into_inner()
1910            .modes
1911            .iter()
1912            .map(|m| m.mode.clone())
1913            .collect();
1914        assert_eq!(names.len(), 1);
1915        assert!(names.contains(&"ext.multi_round.v1".to_string()));
1916    }
1917
1918    #[tokio::test]
1919    async fn get_manifest_includes_all_modes() {
1920        let (server, _) = make_server();
1921        let resp = server
1922            .get_manifest(Request::new(crate::pb::GetManifestRequest {
1923                agent_id: String::new(),
1924            }))
1925            .await
1926            .unwrap();
1927        let manifest = resp.into_inner().manifest.unwrap();
1928        assert_eq!(manifest.supported_modes.len(), 6);
1929        assert!(manifest
1930            .supported_modes
1931            .contains(&"ext.multi_round.v1".to_string()));
1932    }
1933
1934    #[tokio::test]
1935    async fn get_session_returns_metadata() {
1936        let (server, _) = make_server();
1937        let sid = new_sid();
1938        let ack = do_send(
1939            &server,
1940            "agent://orchestrator",
1941            Envelope {
1942                macp_version: "1.0".into(),
1943                mode: "macp.mode.decision.v1".into(),
1944                message_type: "SessionStart".into(),
1945                message_id: "m1".into(),
1946                session_id: sid.clone(),
1947                sender: String::new(),
1948                timestamp_unix_ms: Utc::now().timestamp_millis(),
1949                payload: start_payload(),
1950            },
1951        )
1952        .await;
1953        assert!(ack.ok);
1954
1955        let mut req = Request::new(GetSessionRequest {
1956            session_id: sid.clone(),
1957        });
1958        req.metadata_mut().insert(
1959            "authorization",
1960            format!("Bearer {}", "agent://orchestrator")
1961                .parse()
1962                .unwrap(),
1963        );
1964        let resp = server.get_session(req).await.unwrap();
1965        let meta = resp.into_inner().metadata.unwrap();
1966        assert_eq!(meta.session_id, sid);
1967        assert_eq!(meta.mode, "macp.mode.decision.v1");
1968        assert_eq!(meta.mode_version, "1.0.0");
1969        assert_eq!(meta.configuration_version, "cfg-1");
1970    }
1971
1972    #[tokio::test]
1973    async fn cancel_session_transitions_to_cancelled() {
1974        let (server, _) = make_server();
1975        let sid = new_sid();
1976        let ack = do_send(
1977            &server,
1978            "agent://orchestrator",
1979            Envelope {
1980                macp_version: "1.0".into(),
1981                mode: "macp.mode.decision.v1".into(),
1982                message_type: "SessionStart".into(),
1983                message_id: "m1".into(),
1984                session_id: sid.clone(),
1985                sender: String::new(),
1986                timestamp_unix_ms: Utc::now().timestamp_millis(),
1987                payload: start_payload(),
1988            },
1989        )
1990        .await;
1991        assert!(ack.ok);
1992
1993        let mut req = Request::new(CancelSessionRequest {
1994            session_id: sid,
1995            reason: "no longer needed".into(),
1996        });
1997        req.metadata_mut().insert(
1998            "authorization",
1999            format!("Bearer {}", "agent://orchestrator")
2000                .parse()
2001                .unwrap(),
2002        );
2003        let resp = server.cancel_session(req).await.unwrap();
2004        let ack = resp.into_inner().ack.unwrap();
2005        assert!(ack.ok);
2006        // RFC-MACP-0001 §7.3: cancellation now yields the distinct CANCELLED state.
2007        assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2008    }
2009
2010    #[tokio::test]
2011    async fn participant_cannot_cancel_session() {
2012        let (server, _) = make_server();
2013        let sid = new_sid();
2014        let ack = do_send(
2015            &server,
2016            "agent://orchestrator",
2017            Envelope {
2018                macp_version: "1.0".into(),
2019                mode: "macp.mode.decision.v1".into(),
2020                message_type: "SessionStart".into(),
2021                message_id: "m1".into(),
2022                session_id: sid.clone(),
2023                sender: String::new(),
2024                timestamp_unix_ms: Utc::now().timestamp_millis(),
2025                payload: start_payload(),
2026            },
2027        )
2028        .await;
2029        assert!(ack.ok);
2030
2031        let mut req = Request::new(CancelSessionRequest {
2032            session_id: sid,
2033            reason: "I want to cancel".into(),
2034        });
2035        req.metadata_mut().insert(
2036            "authorization",
2037            format!("Bearer {}", "agent://fraud").parse().unwrap(),
2038        );
2039        let err = server.cancel_session(req).await.unwrap_err();
2040        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2041    }
2042
2043    #[tokio::test]
2044    async fn cancel_session_unknown_session_returns_error() {
2045        let (server, _) = make_server();
2046        let mut req = Request::new(CancelSessionRequest {
2047            session_id: "nonexistent".into(),
2048            reason: "test".into(),
2049        });
2050        req.metadata_mut().insert(
2051            "authorization",
2052            format!("Bearer {}", "agent://orchestrator")
2053                .parse()
2054                .unwrap(),
2055        );
2056        let err = server.cancel_session(req).await.unwrap_err();
2057        assert_eq!(err.code(), tonic::Code::NotFound);
2058    }
2059
2060    #[tokio::test]
2061    async fn ambient_signal_accepted() {
2062        let (server, _) = make_server();
2063        let ack = do_send(
2064            &server,
2065            "agent://orchestrator",
2066            Envelope {
2067                macp_version: "1.0".into(),
2068                mode: String::new(),
2069                message_type: "Signal".into(),
2070                message_id: "sig-1".into(),
2071                session_id: String::new(),
2072                sender: String::new(),
2073                timestamp_unix_ms: Utc::now().timestamp_millis(),
2074                payload: vec![],
2075            },
2076        )
2077        .await;
2078        assert!(ack.ok);
2079    }
2080
2081    #[tokio::test]
2082    async fn signal_with_session_id_rejected() {
2083        let (server, _) = make_server();
2084        let ack = do_send(
2085            &server,
2086            "agent://orchestrator",
2087            Envelope {
2088                macp_version: "1.0".into(),
2089                mode: String::new(),
2090                message_type: "Signal".into(),
2091                message_id: "sig-2".into(),
2092                session_id: "some-session".into(),
2093                sender: String::new(),
2094                timestamp_unix_ms: Utc::now().timestamp_millis(),
2095                payload: vec![],
2096            },
2097        )
2098        .await;
2099        assert!(!ack.ok);
2100        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2101    }
2102
2103    #[tokio::test]
2104    async fn signal_with_mode_rejected() {
2105        let (server, _) = make_server();
2106        let ack = do_send(
2107            &server,
2108            "agent://orchestrator",
2109            Envelope {
2110                macp_version: "1.0".into(),
2111                mode: "macp.mode.decision.v1".into(),
2112                message_type: "Signal".into(),
2113                message_id: "sig-3".into(),
2114                session_id: String::new(),
2115                sender: String::new(),
2116                timestamp_unix_ms: Utc::now().timestamp_millis(),
2117                payload: vec![],
2118            },
2119        )
2120        .await;
2121        assert!(!ack.ok);
2122        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2123    }
2124
2125    #[tokio::test]
2126    async fn ambient_progress_accepted() {
2127        let (server, _) = make_server();
2128        let ack = do_send(
2129            &server,
2130            "agent://orchestrator",
2131            Envelope {
2132                macp_version: "1.0".into(),
2133                mode: String::new(),
2134                message_type: "Progress".into(),
2135                message_id: "prog-1".into(),
2136                session_id: String::new(),
2137                sender: String::new(),
2138                timestamp_unix_ms: Utc::now().timestamp_millis(),
2139                payload: vec![],
2140            },
2141        )
2142        .await;
2143        assert!(ack.ok);
2144    }
2145
2146    #[tokio::test]
2147    async fn ambient_progress_with_mode_rejected() {
2148        let (server, _) = make_server();
2149        let ack = do_send(
2150            &server,
2151            "agent://orchestrator",
2152            Envelope {
2153                macp_version: "1.0".into(),
2154                mode: "macp.mode.decision.v1".into(),
2155                message_type: "Progress".into(),
2156                message_id: "prog-2".into(),
2157                session_id: String::new(),
2158                sender: String::new(),
2159                timestamp_unix_ms: Utc::now().timestamp_millis(),
2160                payload: vec![],
2161            },
2162        )
2163        .await;
2164        assert!(!ack.ok);
2165        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2166    }
2167
2168    #[tokio::test]
2169    async fn manifest_advertises_stream_enabled() {
2170        let (server, _) = make_server();
2171        let resp = server
2172            .initialize(Request::new(InitializeRequest {
2173                supported_protocol_versions: vec!["1.0".into()],
2174                client_info: None,
2175                capabilities: None,
2176            }))
2177            .await
2178            .unwrap();
2179        let caps = resp.into_inner().capabilities.unwrap();
2180        assert!(caps.sessions.unwrap().stream);
2181    }
2182
2183    #[tokio::test]
2184    async fn initialize_empty_versions_rejected() {
2185        let (server, _) = make_server();
2186        let err = server
2187            .initialize(Request::new(InitializeRequest {
2188                supported_protocol_versions: vec![],
2189                client_info: None,
2190                capabilities: None,
2191            }))
2192            .await
2193            .unwrap_err();
2194        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2195    }
2196
2197    #[tokio::test]
2198    async fn initialize_unsupported_version_rejected() {
2199        let (server, _) = make_server();
2200        let err = server
2201            .initialize(Request::new(InitializeRequest {
2202                supported_protocol_versions: vec!["2.0".into()],
2203                client_info: None,
2204                capabilities: None,
2205            }))
2206            .await
2207            .unwrap_err();
2208        assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2209    }
2210
2211    // ── RFC-MACP-0006-A1: passive subscribe tests ──────────────────────
2212
2213    fn observer_identity(sender: &str) -> AuthIdentity {
2214        AuthIdentity {
2215            sender: sender.into(),
2216            allowed_modes: None,
2217            can_start_sessions: false,
2218            max_open_sessions: None,
2219            can_manage_mode_registry: false,
2220            is_observer: true,
2221        }
2222    }
2223
2224    fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2225        StreamSessionRequest {
2226            subscribe_session_id: session_id.into(),
2227            after_sequence: after,
2228            envelope: None,
2229        }
2230    }
2231
2232    fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2233        SessionStartPayload {
2234            intent: "intent".into(),
2235            participants,
2236            mode_version: "1.0.0".into(),
2237            configuration_version: "cfg-1".into(),
2238            policy_version: String::new(),
2239            ttl_ms: 60_000,
2240            context_id: String::new(),
2241            extensions: std::collections::HashMap::new(),
2242            roots: vec![],
2243            max_suspend_ms: 0,
2244        }
2245        .encode_to_vec()
2246    }
2247
2248    async fn start_session(
2249        server: &MacpServer,
2250        initiator: &str,
2251        sid: &str,
2252        participants: Vec<String>,
2253    ) {
2254        let ack = do_send(
2255            server,
2256            initiator,
2257            Envelope {
2258                macp_version: "1.0".into(),
2259                mode: "macp.mode.decision.v1".into(),
2260                message_type: "SessionStart".into(),
2261                message_id: "start".into(),
2262                session_id: sid.into(),
2263                sender: String::new(),
2264                timestamp_unix_ms: Utc::now().timestamp_millis(),
2265                payload: start_multi_participant(participants),
2266            },
2267        )
2268        .await;
2269        assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2270    }
2271
2272    async fn send_proposal(
2273        server: &MacpServer,
2274        sender: &str,
2275        sid: &str,
2276        message_id: &str,
2277        proposal_id: &str,
2278    ) {
2279        let payload = crate::decision_pb::ProposalPayload {
2280            proposal_id: proposal_id.into(),
2281            option: "opt".into(),
2282            rationale: "r".into(),
2283            supporting_data: vec![],
2284        }
2285        .encode_to_vec();
2286        let ack = do_send(
2287            server,
2288            sender,
2289            Envelope {
2290                macp_version: "1.0".into(),
2291                mode: "macp.mode.decision.v1".into(),
2292                message_type: "Proposal".into(),
2293                message_id: message_id.into(),
2294                session_id: sid.into(),
2295                sender: String::new(),
2296                timestamp_unix_ms: Utc::now().timestamp_millis(),
2297                payload,
2298            },
2299        )
2300        .await;
2301        assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2302    }
2303
2304    #[tokio::test]
2305    async fn subscribe_replays_session_history_from_zero() {
2306        let (server, _) = make_server();
2307        let sid = new_sid();
2308        let initiator = "agent://orchestrator";
2309        let peer = "agent://fraud";
2310        start_session(
2311            &server,
2312            initiator,
2313            &sid,
2314            vec![initiator.into(), peer.into()],
2315        )
2316        .await;
2317        send_proposal(&server, peer, &sid, "m2", "p1").await;
2318
2319        let mut bound = None;
2320        let mut events = None;
2321        let replay = server
2322            .process_stream_request(
2323                &stream_identity(peer),
2324                subscribe_frame(&sid, 0),
2325                &mut bound,
2326                &mut events,
2327            )
2328            .await
2329            .unwrap();
2330
2331        assert_eq!(replay.len(), 2);
2332        assert_eq!(replay[0].message_type, "SessionStart");
2333        assert_eq!(replay[0].message_id, "start");
2334        assert_eq!(replay[1].message_type, "Proposal");
2335        assert_eq!(replay[1].message_id, "m2");
2336        assert_eq!(bound.as_deref(), Some(sid.as_str()));
2337        assert!(events.is_some());
2338    }
2339
2340    #[tokio::test]
2341    async fn subscribe_after_sequence_filters_history() {
2342        let (server, _) = make_server();
2343        let sid = new_sid();
2344        let initiator = "agent://orchestrator";
2345        let peer = "agent://fraud";
2346        start_session(
2347            &server,
2348            initiator,
2349            &sid,
2350            vec![initiator.into(), peer.into()],
2351        )
2352        .await;
2353        send_proposal(&server, peer, &sid, "m2", "p1").await;
2354        send_proposal(&server, peer, &sid, "m3", "p2").await;
2355
2356        let mut bound = None;
2357        let mut events = None;
2358        let replay = server
2359            .process_stream_request(
2360                &stream_identity(peer),
2361                subscribe_frame(&sid, 2),
2362                &mut bound,
2363                &mut events,
2364            )
2365            .await
2366            .unwrap();
2367
2368        assert_eq!(replay.len(), 1);
2369        assert_eq!(replay[0].message_id, "m3");
2370    }
2371
2372    #[tokio::test]
2373    async fn subscribe_unknown_session_returns_not_found() {
2374        let (server, _) = make_server();
2375        let mut bound = None;
2376        let mut events = None;
2377        let status = server
2378            .process_stream_request(
2379                &stream_identity("agent://orchestrator"),
2380                subscribe_frame("missing-session", 0),
2381                &mut bound,
2382                &mut events,
2383            )
2384            .await
2385            .unwrap_err();
2386        assert_eq!(status.code(), tonic::Code::NotFound);
2387        assert!(bound.is_none());
2388        assert!(events.is_none());
2389    }
2390
2391    #[tokio::test]
2392    async fn subscribe_non_participant_is_forbidden() {
2393        let (server, _) = make_server();
2394        let sid = new_sid();
2395        start_session(
2396            &server,
2397            "agent://orchestrator",
2398            &sid,
2399            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2400        )
2401        .await;
2402
2403        let mut bound = None;
2404        let mut events = None;
2405        let status = server
2406            .process_stream_request(
2407                &stream_identity("agent://outsider"),
2408                subscribe_frame(&sid, 0),
2409                &mut bound,
2410                &mut events,
2411            )
2412            .await
2413            .unwrap_err();
2414        assert_eq!(status.code(), tonic::Code::PermissionDenied);
2415    }
2416
2417    #[tokio::test]
2418    async fn subscribe_observer_identity_allowed() {
2419        let (server, _) = make_server();
2420        let sid = new_sid();
2421        start_session(
2422            &server,
2423            "agent://orchestrator",
2424            &sid,
2425            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2426        )
2427        .await;
2428
2429        let mut bound = None;
2430        let mut events = None;
2431        let replay = server
2432            .process_stream_request(
2433                &observer_identity("agent://auditor"),
2434                subscribe_frame(&sid, 0),
2435                &mut bound,
2436                &mut events,
2437            )
2438            .await
2439            .unwrap();
2440        assert_eq!(replay.len(), 1);
2441        assert_eq!(replay[0].message_type, "SessionStart");
2442    }
2443
2444    #[tokio::test]
2445    async fn subscribe_initiator_allowed_even_when_not_listed() {
2446        // Per RFC-MACP-0007, the initiator is always authorized for session
2447        // access, even if not present in the participants list.
2448        let (server, _) = make_server();
2449        let sid = new_sid();
2450        start_session(
2451            &server,
2452            "agent://orchestrator",
2453            &sid,
2454            vec!["agent://fraud".into()],
2455        )
2456        .await;
2457
2458        let mut bound = None;
2459        let mut events = None;
2460        let replay = server
2461            .process_stream_request(
2462                &stream_identity("agent://orchestrator"),
2463                subscribe_frame(&sid, 0),
2464                &mut bound,
2465                &mut events,
2466            )
2467            .await
2468            .unwrap();
2469        assert_eq!(replay.len(), 1);
2470    }
2471
2472    #[tokio::test]
2473    async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2474        let (server, _) = make_server();
2475        let sid = new_sid();
2476        let req = StreamSessionRequest {
2477            subscribe_session_id: sid.clone(),
2478            after_sequence: 0,
2479            envelope: Some(Envelope {
2480                macp_version: "1.0".into(),
2481                mode: "macp.mode.decision.v1".into(),
2482                message_type: "SessionStart".into(),
2483                message_id: "m1".into(),
2484                session_id: sid,
2485                sender: String::new(),
2486                timestamp_unix_ms: Utc::now().timestamp_millis(),
2487                payload: start_payload(),
2488            }),
2489        };
2490
2491        let mut bound = None;
2492        let mut events = None;
2493        let status = server
2494            .process_stream_request(
2495                &stream_identity("agent://orchestrator"),
2496                req,
2497                &mut bound,
2498                &mut events,
2499            )
2500            .await
2501            .unwrap_err();
2502        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2503    }
2504
2505    #[tokio::test]
2506    async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2507        let (server, _) = make_server();
2508        let sid1 = new_sid();
2509        let sid2 = new_sid();
2510        start_session(
2511            &server,
2512            "agent://orchestrator",
2513            &sid1,
2514            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2515        )
2516        .await;
2517        start_session(
2518            &server,
2519            "agent://orchestrator",
2520            &sid2,
2521            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2522        )
2523        .await;
2524
2525        // First subscribe binds the stream to sid1
2526        let identity = stream_identity("agent://fraud");
2527        let mut bound = None;
2528        let mut events = None;
2529        server
2530            .process_stream_request(
2531                &identity,
2532                subscribe_frame(&sid1, 0),
2533                &mut bound,
2534                &mut events,
2535            )
2536            .await
2537            .unwrap();
2538        assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2539
2540        // Second subscribe to sid2 on the same stream must be rejected
2541        let status = server
2542            .process_stream_request(
2543                &identity,
2544                subscribe_frame(&sid2, 0),
2545                &mut bound,
2546                &mut events,
2547            )
2548            .await
2549            .unwrap_err();
2550        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2551    }
2552
2553    /// E3: an injected ingress engine gates session start, messages, and
2554    /// session reads — deny-one-sender double proves all three hooks fire and
2555    /// that denial surfaces as POLICY_DENIED / PermissionDenied (fail closed).
2556    struct DenySenderEngine {
2557        denied: String,
2558    }
2559
2560    #[async_trait::async_trait]
2561    impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2562        async fn evaluate_session_start(
2563            &self,
2564            identity: &crate::security::AuthIdentity,
2565            _mode: &str,
2566            _env: &Envelope,
2567        ) -> macp_core::policy::PolicyDecision {
2568            if identity.sender == self.denied {
2569                macp_core::policy::PolicyDecision::Deny {
2570                    reasons: vec!["sender embargoed".into()],
2571                }
2572            } else {
2573                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2574            }
2575        }
2576
2577        async fn evaluate_message(
2578            &self,
2579            identity: &crate::security::AuthIdentity,
2580            _session: &macp_core::session::Session,
2581            _env: &Envelope,
2582        ) -> macp_core::policy::PolicyDecision {
2583            if identity.sender == self.denied {
2584                macp_core::policy::PolicyDecision::Deny {
2585                    reasons: vec!["sender embargoed".into()],
2586                }
2587            } else {
2588                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2589            }
2590        }
2591
2592        async fn evaluate_session_access(
2593            &self,
2594            identity: &crate::security::AuthIdentity,
2595            _session: &macp_core::session::Session,
2596        ) -> macp_core::policy::PolicyDecision {
2597            if identity.sender == self.denied {
2598                macp_core::policy::PolicyDecision::Deny {
2599                    reasons: vec!["sender embargoed".into()],
2600                }
2601            } else {
2602                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2603            }
2604        }
2605    }
2606
2607    #[tokio::test]
2608    async fn policy_engine_gates_all_three_ingress_points() {
2609        let (server, _runtime) = make_server();
2610        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2611            denied: "agent://embargoed".into(),
2612        }));
2613
2614        let sid = new_sid();
2615        let start_payload = SessionStartPayload {
2616            intent: "e3".into(),
2617            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2618            mode_version: "1.0.0".into(),
2619            configuration_version: "cfg-1".into(),
2620            policy_version: String::new(),
2621            ttl_ms: 60_000,
2622            context_id: String::new(),
2623            extensions: Default::default(),
2624            roots: vec![],
2625            max_suspend_ms: 0,
2626        }
2627        .encode_to_vec();
2628        let start_env = |sender: &str, sid: &str| Envelope {
2629            macp_version: "1.0".into(),
2630            mode: "macp.mode.decision.v1".into(),
2631            message_type: "SessionStart".into(),
2632            message_id: new_sid(),
2633            session_id: sid.into(),
2634            sender: sender.into(),
2635            timestamp_unix_ms: Utc::now().timestamp_millis(),
2636            payload: start_payload.clone(),
2637        };
2638
2639        // 1. Embargoed sender cannot start a session.
2640        let ack = server
2641            .send(send_req(
2642                "agent://embargoed",
2643                start_env("agent://embargoed", &sid),
2644            ))
2645            .await
2646            .unwrap()
2647            .into_inner()
2648            .ack
2649            .unwrap();
2650        assert!(!ack.ok);
2651        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2652
2653        // Allowed sender starts it.
2654        let ack = server
2655            .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2656            .await
2657            .unwrap()
2658            .into_inner()
2659            .ack
2660            .unwrap();
2661        assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2662
2663        // 2. Embargoed sender cannot send into the session.
2664        let proposal = crate::decision_pb::ProposalPayload {
2665            proposal_id: "p1".into(),
2666            option: "x".into(),
2667            rationale: "r".into(),
2668            supporting_data: vec![],
2669        }
2670        .encode_to_vec();
2671        let msg_env = Envelope {
2672            macp_version: "1.0".into(),
2673            mode: "macp.mode.decision.v1".into(),
2674            message_type: "Proposal".into(),
2675            message_id: new_sid(),
2676            session_id: sid.clone(),
2677            sender: "agent://embargoed".into(),
2678            timestamp_unix_ms: Utc::now().timestamp_millis(),
2679            payload: proposal,
2680        };
2681        let ack = server
2682            .send(send_req("agent://embargoed", msg_env))
2683            .await
2684            .unwrap()
2685            .into_inner()
2686            .ack
2687            .unwrap();
2688        assert!(!ack.ok);
2689        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2690
2691        // 3. Embargoed sender cannot read the session.
2692        let mut req = Request::new(crate::pb::GetSessionRequest {
2693            session_id: sid.clone(),
2694        });
2695        req.metadata_mut()
2696            .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2697        let err = server
2698            .get_session(req)
2699            .await
2700            .expect_err("embargoed read must be denied");
2701        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2702    }
2703
2704    /// E3 transport-parity: the ingress engine gates the STREAM path too — a
2705    /// denied sender must not be able to bypass the engine by switching from
2706    /// unary Send to StreamSession (envelope frames or subscribe frames).
2707    #[tokio::test]
2708    async fn policy_engine_gates_stream_path() {
2709        let (server, runtime) = make_server();
2710        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2711            denied: "agent://embargoed".into(),
2712        }));
2713
2714        // Session started by an allowed sender (participants include the
2715        // embargoed agent so built-in membership checks pass — only the
2716        // engine denies it).
2717        let sid = new_sid();
2718        let payload = SessionStartPayload {
2719            intent: "e3-stream".into(),
2720            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2721            mode_version: "1.0.0".into(),
2722            configuration_version: "cfg-1".into(),
2723            policy_version: String::new(),
2724            ttl_ms: 60_000,
2725            context_id: String::new(),
2726            extensions: Default::default(),
2727            roots: vec![],
2728            max_suspend_ms: 0,
2729        }
2730        .encode_to_vec();
2731        runtime
2732            .process(
2733                &Envelope {
2734                    macp_version: "1.0".into(),
2735                    mode: "macp.mode.decision.v1".into(),
2736                    message_type: "SessionStart".into(),
2737                    message_id: new_sid(),
2738                    session_id: sid.clone(),
2739                    sender: "agent://ok".into(),
2740                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2741                    payload,
2742                },
2743                None,
2744            )
2745            .await
2746            .unwrap();
2747
2748        let embargoed = crate::security::AuthIdentity {
2749            sender: "agent://embargoed".into(),
2750            allowed_modes: None,
2751            can_start_sessions: true,
2752            max_open_sessions: None,
2753            can_manage_mode_registry: false,
2754            is_observer: false,
2755        };
2756        let mut bound = None;
2757        let mut events = None;
2758
2759        // 1. Stream envelope frame from the embargoed sender: denied.
2760        let proposal = crate::decision_pb::ProposalPayload {
2761            proposal_id: "p1".into(),
2762            option: "x".into(),
2763            rationale: "r".into(),
2764            supporting_data: vec![],
2765        }
2766        .encode_to_vec();
2767        let req = StreamSessionRequest {
2768            envelope: Some(Envelope {
2769                macp_version: "1.0".into(),
2770                mode: "macp.mode.decision.v1".into(),
2771                message_type: "Proposal".into(),
2772                message_id: new_sid(),
2773                session_id: sid.clone(),
2774                sender: "agent://embargoed".into(),
2775                timestamp_unix_ms: Utc::now().timestamp_millis(),
2776                payload: proposal,
2777            }),
2778            subscribe_session_id: String::new(),
2779            after_sequence: 0,
2780        };
2781        let err = server
2782            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2783            .await
2784            .expect_err("stream envelope from embargoed sender must be denied");
2785        // PolicyDenied maps to FailedPrecondition on the transport (same
2786        // error the unary path expresses as a POLICY_DENIED ack).
2787        assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2788        assert!(err.message().contains("PolicyDenied"), "{err:?}");
2789
2790        // 2. Passive-subscribe frame (history read) from the embargoed
2791        //    sender: denied even though membership would allow it.
2792        let req = StreamSessionRequest {
2793            envelope: None,
2794            subscribe_session_id: sid.clone(),
2795            after_sequence: 0,
2796        };
2797        let err = server
2798            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2799            .await
2800            .expect_err("stream subscribe from embargoed sender must be denied");
2801        assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2802    }
2803}