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                version: "0.4.0".into(),
813                description: "Reference implementation of the Multi-Agent Coordination Protocol"
814                    .into(),
815                website_url: String::new(),
816            }),
817            capabilities: Some(Capabilities {
818                sessions: Some(SessionsCapability { stream: true, list_sessions: true, watch_sessions: true }),
819                cancellation: Some(CancellationCapability {
820                    cancel_session: true,
821                }),
822                progress: Some(ProgressCapability { progress: true }),
823                manifest: Some(ManifestCapability { get_manifest: true }),
824                mode_registry: Some(ModeRegistryCapability {
825                    list_modes: true,
826                    list_changed: true,
827                }),
828                roots: Some(RootsCapability {
829                    // ListRoots is answerable (the root set is empty — a valid
830                    // state), but this runtime has no roots provider, so the
831                    // set never changes: do not advertise change notifications
832                    // (RFC-MACP-0006 §3.3 gates WatchRoots on list_changed).
833                    // Revisit when a roots provider lands (plans E2).
834                    list_roots: true,
835                    list_changed: false,
836                }),
837                policy_registry: Some(PolicyRegistryCapability {
838                    register_policy: !self.policies_read_only,
839                    list_policies: true,
840                    list_changed: true,
841                }),
842                experimental: Some(crate::pb::ExperimentalCapabilities {
843                    features: HashMap::from([
844                        ("ext_mode_lifecycle".into(), "true".into()),
845                    ]),
846                }),
847            }),
848            supported_modes: self.runtime.registered_mode_names(),
849            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(),
850        }))
851    }
852
853    async fn send(&self, request: Request<SendRequest>) -> Result<Response<SendResponse>, Status> {
854        let env = request
855            .get_ref()
856            .envelope
857            .clone()
858            .ok_or_else(|| Status::invalid_argument("SendRequest must contain an envelope"))?;
859
860        let result = async {
861            self.validate_envelope_shape(&env)?;
862            let (env, max_open) = self.authenticate_send_request(&request, env).await?;
863            self.runtime
864                .process(&env, max_open)
865                .await
866                .map(|process_result| (env, process_result))
867        }
868        .await;
869
870        let ack = match result {
871            Ok((env, process_result)) => Ack {
872                ok: true,
873                duplicate: process_result.duplicate,
874                message_id: env.message_id.clone(),
875                session_id: env.session_id.clone(),
876                accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
877                session_state: Self::session_state_to_pb(&process_result.session_state),
878                error: None,
879            },
880            Err(err) => {
881                let env = request.get_ref().envelope.clone().unwrap_or_default();
882                // Rejection counters were collected but never recorded before
883                // (permanently zero). Session-scoped rejections are counted
884                // per mode; commitments additionally under their own counter.
885                if !env.session_id.is_empty() {
886                    self.runtime.metrics().record_message_rejected(&env.mode);
887                    if env.message_type == "Commitment" {
888                        self.runtime.metrics().record_commitment_rejected(&env.mode);
889                    }
890                }
891                Self::make_error_ack(&err, &env)
892            }
893        };
894
895        Ok(Response::new(SendResponse { ack: Some(ack) }))
896    }
897
898    async fn get_session(
899        &self,
900        request: Request<GetSessionRequest>,
901    ) -> Result<Response<GetSessionResponse>, Status> {
902        let session_id = request.get_ref().session_id.clone();
903        let _identity = self
904            .authenticate_session_access(&request, &session_id)
905            .await?;
906        let session = self
907            .runtime
908            .get_session_checked(&session_id)
909            .await
910            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
911
912        Ok(Response::new(GetSessionResponse {
913            metadata: Some(Self::session_to_metadata(&session)),
914        }))
915    }
916
917    async fn cancel_session(
918        &self,
919        request: Request<CancelSessionRequest>,
920    ) -> Result<Response<CancelSessionResponse>, Status> {
921        let session_id = request.get_ref().session_id.clone();
922        let identity = self
923            .security
924            .authenticate_metadata(request.metadata())
925            .await
926            .map_err(Self::status_from_error)?;
927        let session = self
928            .runtime
929            .get_session_checked(&session_id)
930            .await
931            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
932        // RFC-MACP-0001: "Only the initiator and policy-delegated roles may cancel."
933        // CancelSession is a Core control-plane message — mode authorization does not apply.
934        if identity.sender != session.initiator_sender
935            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
936        {
937            return Err(Status::permission_denied(
938                "FORBIDDEN: only the session initiator or policy-delegated roles can cancel",
939            ));
940        }
941        let sender = identity.sender.clone();
942        let req = request.into_inner();
943        match self
944            .runtime
945            .cancel_session(&req.session_id, &req.reason, &sender)
946            .await
947        {
948            Ok(result) => Ok(Response::new(CancelSessionResponse {
949                ack: Some(Ack {
950                    ok: true,
951                    duplicate: false,
952                    message_id: String::new(),
953                    session_id: req.session_id,
954                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
955                    session_state: Self::session_state_to_pb(&result.session_state),
956                    error: None,
957                }),
958            })),
959            Err(err) => Ok(Response::new(CancelSessionResponse {
960                ack: Some(Ack {
961                    ok: false,
962                    duplicate: false,
963                    message_id: String::new(),
964                    session_id: req.session_id.clone(),
965                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
966                    session_state: PbSessionState::Unspecified.into(),
967                    error: Some(PbMacpError {
968                        code: err.error_code().into(),
969                        message: err.to_string(),
970                        session_id: req.session_id,
971                        message_id: String::new(),
972                        details: vec![],
973                    }),
974                }),
975            })),
976        }
977    }
978
979    async fn suspend_session(
980        &self,
981        request: Request<SuspendSessionRequest>,
982    ) -> Result<Response<SuspendSessionResponse>, Status> {
983        let session_id = request.get_ref().session_id.clone();
984        let identity = self
985            .security
986            .authenticate_metadata(request.metadata())
987            .await
988            .map_err(Self::status_from_error)?;
989        let session = self
990            .runtime
991            .get_session_checked(&session_id)
992            .await
993            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
994        // RFC-MACP-0001 §7.5: same authority model as CancelSession — initiator
995        // or policy-delegated roles only; mode authorization does not apply.
996        if identity.sender != session.initiator_sender
997            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
998        {
999            return Err(Status::permission_denied(
1000                "FORBIDDEN: only the session initiator or policy-delegated roles can suspend",
1001            ));
1002        }
1003        let sender = identity.sender.clone();
1004        let req = request.into_inner();
1005        match self
1006            .runtime
1007            .suspend_session(&req.session_id, &req.reason, &sender)
1008            .await
1009        {
1010            Ok(result) => Ok(Response::new(SuspendSessionResponse {
1011                ack: Some(Ack {
1012                    ok: true,
1013                    duplicate: false,
1014                    message_id: String::new(),
1015                    session_id: req.session_id,
1016                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1017                    session_state: Self::session_state_to_pb(&result.session_state),
1018                    error: None,
1019                }),
1020            })),
1021            Err(err) => Ok(Response::new(SuspendSessionResponse {
1022                ack: Some(Ack {
1023                    ok: false,
1024                    duplicate: false,
1025                    message_id: String::new(),
1026                    session_id: req.session_id.clone(),
1027                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1028                    session_state: PbSessionState::Unspecified.into(),
1029                    error: Some(PbMacpError {
1030                        code: err.error_code().into(),
1031                        message: err.to_string(),
1032                        session_id: req.session_id,
1033                        message_id: String::new(),
1034                        details: vec![],
1035                    }),
1036                }),
1037            })),
1038        }
1039    }
1040
1041    async fn resume_session(
1042        &self,
1043        request: Request<ResumeSessionRequest>,
1044    ) -> Result<Response<ResumeSessionResponse>, Status> {
1045        let session_id = request.get_ref().session_id.clone();
1046        let identity = self
1047            .security
1048            .authenticate_metadata(request.metadata())
1049            .await
1050            .map_err(Self::status_from_error)?;
1051        let session = self
1052            .runtime
1053            .get_session_checked(&session_id)
1054            .await
1055            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1056        if identity.sender != session.initiator_sender
1057            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1058        {
1059            return Err(Status::permission_denied(
1060                "FORBIDDEN: only the session initiator or policy-delegated roles can resume",
1061            ));
1062        }
1063        let sender = identity.sender.clone();
1064        let req = request.into_inner();
1065        match self
1066            .runtime
1067            .resume_session(&req.session_id, &req.reason, &sender)
1068            .await
1069        {
1070            Ok(result) => Ok(Response::new(ResumeSessionResponse {
1071                ack: Some(Ack {
1072                    ok: true,
1073                    duplicate: false,
1074                    message_id: String::new(),
1075                    session_id: req.session_id,
1076                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1077                    session_state: Self::session_state_to_pb(&result.session_state),
1078                    error: None,
1079                }),
1080            })),
1081            Err(err) => Ok(Response::new(ResumeSessionResponse {
1082                ack: Some(Ack {
1083                    ok: false,
1084                    duplicate: false,
1085                    message_id: String::new(),
1086                    session_id: req.session_id.clone(),
1087                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1088                    session_state: PbSessionState::Unspecified.into(),
1089                    error: Some(PbMacpError {
1090                        code: err.error_code().into(),
1091                        message: err.to_string(),
1092                        session_id: req.session_id,
1093                        message_id: String::new(),
1094                        details: vec![],
1095                    }),
1096                }),
1097            })),
1098        }
1099    }
1100
1101    async fn get_manifest(
1102        &self,
1103        request: Request<GetManifestRequest>,
1104    ) -> Result<Response<GetManifestResponse>, Status> {
1105        let req = request.into_inner();
1106        if !req.agent_id.is_empty() && req.agent_id != "macp-runtime" {
1107            return Err(Status::not_found(format!(
1108                "Agent '{}' not found",
1109                req.agent_id
1110            )));
1111        }
1112
1113        Ok(Response::new(GetManifestResponse {
1114            manifest: Some(crate::pb::AgentManifest {
1115                agent_id: "macp-runtime".into(),
1116                title: "MACP Reference Runtime".into(),
1117                description: "Reference implementation of MACP".into(),
1118                supported_modes: self.runtime.registered_mode_names(),
1119                input_content_types: vec!["application/macp-envelope+proto".into()],
1120                output_content_types: vec!["application/macp-envelope+proto".into()],
1121                metadata: HashMap::new(),
1122                // Empty: unary-first profile has no dedicated transport endpoints.
1123                transport_endpoints: vec![],
1124            }),
1125        }))
1126    }
1127
1128    async fn list_modes(
1129        &self,
1130        _request: Request<ListModesRequest>,
1131    ) -> Result<Response<ListModesResponse>, Status> {
1132        Ok(Response::new(ListModesResponse {
1133            modes: self.runtime.standard_mode_descriptors(),
1134        }))
1135    }
1136
1137    async fn list_roots(
1138        &self,
1139        _request: Request<ListRootsRequest>,
1140    ) -> Result<Response<ListRootsResponse>, Status> {
1141        Ok(Response::new(ListRootsResponse { roots: vec![] }))
1142    }
1143
1144    type StreamSessionStream = SessionResponseStream;
1145
1146    async fn stream_session(
1147        &self,
1148        request: Request<tonic::Streaming<StreamSessionRequest>>,
1149    ) -> Result<Response<Self::StreamSessionStream>, Status> {
1150        let identity = self
1151            .security
1152            .authenticate_metadata(request.metadata())
1153            .await
1154            .map_err(Self::status_from_error)?;
1155        let inbound = request.into_inner();
1156        Ok(Response::new(
1157            self.build_stream_session_stream(identity, inbound),
1158        ))
1159    }
1160
1161    type WatchModeRegistryStream = std::pin::Pin<
1162        Box<dyn futures_core::Stream<Item = Result<WatchModeRegistryResponse, Status>> + Send>,
1163    >;
1164
1165    async fn watch_mode_registry(
1166        &self,
1167        _request: Request<WatchModeRegistryRequest>,
1168    ) -> Result<Response<Self::WatchModeRegistryStream>, Status> {
1169        let mut rx = self.runtime.subscribe_mode_changes();
1170        let stream = async_stream::try_stream! {
1171            // Send initial state
1172            yield WatchModeRegistryResponse {
1173                change: Some(crate::pb::RegistryChanged {
1174                    registry: "modes".into(),
1175                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1176                }),
1177            };
1178            // Wait for changes from register/unregister/promote
1179            while rx.recv().await.is_ok() {
1180                yield WatchModeRegistryResponse {
1181                    change: Some(crate::pb::RegistryChanged {
1182                        registry: "modes".into(),
1183                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1184                    }),
1185                };
1186            }
1187        };
1188        Ok(Response::new(Box::pin(stream)))
1189    }
1190
1191    type WatchRootsStream = std::pin::Pin<
1192        Box<dyn futures_core::Stream<Item = Result<WatchRootsResponse, Status>> + Send>,
1193    >;
1194
1195    async fn watch_roots(
1196        &self,
1197        _request: Request<WatchRootsRequest>,
1198    ) -> Result<Response<Self::WatchRootsStream>, Status> {
1199        let initial = WatchRootsResponse {
1200            change: Some(crate::pb::RootsChanged {
1201                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1202            }),
1203        };
1204        let stream = async_stream::try_stream! {
1205            yield initial;
1206            // Roots are static — keep the stream open but idle.
1207            std::future::pending::<()>().await;
1208        };
1209        Ok(Response::new(Box::pin(stream)))
1210    }
1211
1212    type WatchSignalsStream = std::pin::Pin<
1213        Box<dyn futures_core::Stream<Item = Result<WatchSignalsResponse, Status>> + Send>,
1214    >;
1215
1216    type WatchSessionsStream = std::pin::Pin<
1217        Box<dyn futures_core::Stream<Item = Result<WatchSessionsResponse, Status>> + Send>,
1218    >;
1219
1220    async fn watch_signals(
1221        &self,
1222        request: Request<WatchSignalsRequest>,
1223    ) -> Result<Response<Self::WatchSignalsStream>, Status> {
1224        // Ambient signals carry agent-generated payload data; subscribing is
1225        // gated on authentication like the session-observation surfaces.
1226        // (RFC-0004 §4.1 constrains unauthenticated *producers*; requiring
1227        // authenticated subscribers is this runtime's hardening posture.)
1228        let _identity = self
1229            .security
1230            .authenticate_metadata(request.metadata())
1231            .await
1232            .map_err(Self::status_from_error)?;
1233        let mut rx = self.runtime.subscribe_signals();
1234        let stream = async_stream::try_stream! {
1235            loop {
1236                match rx.recv().await {
1237                    Ok(envelope) => {
1238                        yield WatchSignalsResponse {
1239                            envelope: Some(envelope),
1240                        };
1241                    }
1242                    // Surface lag instead of silently ending the stream: a
1243                    // slow consumer must be able to distinguish "no traffic"
1244                    // from "events dropped" (mirrors StreamSession).
1245                    Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1246                        Err(Status::resource_exhausted(format!(
1247                            "WatchSignals receiver fell behind by {skipped} signals"
1248                        )))?;
1249                    }
1250                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1251                }
1252            }
1253        };
1254        Ok(Response::new(Box::pin(stream)))
1255    }
1256
1257    // Session lifecycle observation RPCs
1258
1259    async fn list_sessions(
1260        &self,
1261        request: Request<ListSessionsRequest>,
1262    ) -> Result<Response<ListSessionsResponse>, Status> {
1263        let _identity = self
1264            .security
1265            .authenticate_metadata(request.metadata())
1266            .await
1267            .map_err(Self::status_from_error)?;
1268        let sessions = self.runtime.registry.get_all_sessions().await;
1269        let metadata: Vec<SessionMetadata> =
1270            sessions.iter().map(Self::session_to_metadata).collect();
1271        Ok(Response::new(ListSessionsResponse { sessions: metadata }))
1272    }
1273
1274    async fn watch_sessions(
1275        &self,
1276        request: Request<WatchSessionsRequest>,
1277    ) -> Result<Response<Self::WatchSessionsStream>, Status> {
1278        let _identity = self
1279            .security
1280            .authenticate_metadata(request.metadata())
1281            .await
1282            .map_err(Self::status_from_error)?;
1283        let mut rx = self.runtime.subscribe_session_lifecycle();
1284        let runtime = Arc::clone(&self.runtime);
1285        let stream = async_stream::try_stream! {
1286            // Initial sync: emit all current sessions as CREATED events. The
1287            // lifecycle bus was subscribed *before* this snapshot (so no event
1288            // is missed); any Created event buffered in that window would
1289            // duplicate a snapshot entry — session IDs are create-once, so we
1290            // dedupe buffered Created events against the synced set below.
1291            let sessions = runtime.registry.get_all_sessions().await;
1292            let mut synced: std::collections::HashSet<String> =
1293                std::collections::HashSet::with_capacity(sessions.len());
1294            for session in &sessions {
1295                synced.insert(session.session_id.clone());
1296                yield WatchSessionsResponse {
1297                    event: Some(SessionLifecycleEvent {
1298                        event_type: session_lifecycle_event::EventType::Created.into(),
1299                        session: Some(Self::session_to_metadata(session)),
1300                        observed_at_unix_ms: session.started_at_unix_ms,
1301                    }),
1302                };
1303            }
1304            // Stream lifecycle transitions
1305            loop {
1306                let event = match rx.recv().await {
1307                    Ok(event) => event,
1308                    Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1309                        Err(Status::resource_exhausted(format!(
1310                            "WatchSessions receiver fell behind by {skipped} events"
1311                        )))?;
1312                        break;
1313                    }
1314                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1315                };
1316                let (event_type, sid) = match &event {
1317                    crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1318                        (session_lifecycle_event::EventType::Created, session_id.clone()),
1319                    crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1320                        (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1321                    crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1322                        (session_lifecycle_event::EventType::Expired, session_id.clone()),
1323                    crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1324                        (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1325                    crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1326                        (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1327                    crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1328                        (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1329                };
1330                // Skip the buffered duplicate of an initial-sync entry;
1331                // non-Created events for synced sessions are new information
1332                // and pass through.
1333                if event_type == session_lifecycle_event::EventType::Created
1334                    && !synced.insert(sid.clone())
1335                {
1336                    continue;
1337                }
1338                let session_meta = runtime.registry.get_session(&sid).await
1339                    .map(|s| Self::session_to_metadata(&s));
1340                yield WatchSessionsResponse {
1341                    event: Some(SessionLifecycleEvent {
1342                        event_type: event_type.into(),
1343                        session: session_meta,
1344                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1345                    }),
1346                };
1347            }
1348        };
1349        Ok(Response::new(Box::pin(stream)))
1350    }
1351
1352    // Extension mode lifecycle RPCs
1353
1354    async fn list_ext_modes(
1355        &self,
1356        _request: Request<ListExtModesRequest>,
1357    ) -> Result<Response<ListExtModesResponse>, Status> {
1358        Ok(Response::new(ListExtModesResponse {
1359            modes: self.runtime.extension_mode_descriptors(),
1360        }))
1361    }
1362
1363    async fn register_ext_mode(
1364        &self,
1365        request: Request<RegisterExtModeRequest>,
1366    ) -> Result<Response<RegisterExtModeResponse>, Status> {
1367        let identity = self
1368            .security
1369            .authenticate_metadata(request.metadata())
1370            .await
1371            .map_err(Self::status_from_error)?;
1372        self.security
1373            .authorize_mode_registry(&identity)
1374            .map_err(Self::status_from_error)?;
1375        let req = request.into_inner();
1376        let descriptor = req
1377            .mode_descriptor
1378            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1379        match self.runtime.register_extension(descriptor) {
1380            Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1381                ok: true,
1382                error: String::new(),
1383            })),
1384            Err(e) => Ok(Response::new(RegisterExtModeResponse {
1385                ok: false,
1386                error: e,
1387            })),
1388        }
1389    }
1390
1391    async fn unregister_ext_mode(
1392        &self,
1393        request: Request<UnregisterExtModeRequest>,
1394    ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1395        let identity = self
1396            .security
1397            .authenticate_metadata(request.metadata())
1398            .await
1399            .map_err(Self::status_from_error)?;
1400        self.security
1401            .authorize_mode_registry(&identity)
1402            .map_err(Self::status_from_error)?;
1403        let req = request.into_inner();
1404        match self.runtime.unregister_extension(&req.mode) {
1405            Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1406                ok: true,
1407                error: String::new(),
1408            })),
1409            Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1410                ok: false,
1411                error: e,
1412            })),
1413        }
1414    }
1415
1416    async fn promote_mode(
1417        &self,
1418        request: Request<PromoteModeRequest>,
1419    ) -> Result<Response<PromoteModeResponse>, Status> {
1420        let identity = self
1421            .security
1422            .authenticate_metadata(request.metadata())
1423            .await
1424            .map_err(Self::status_from_error)?;
1425        self.security
1426            .authorize_mode_registry(&identity)
1427            .map_err(Self::status_from_error)?;
1428        let req = request.into_inner();
1429        let new_name = if req.promoted_mode_name.is_empty() {
1430            None
1431        } else {
1432            Some(req.promoted_mode_name.as_str())
1433        };
1434        match self.runtime.promote_mode(&req.mode, new_name) {
1435            Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1436                ok: true,
1437                error: String::new(),
1438                mode: final_name,
1439            })),
1440            Err(e) => Ok(Response::new(PromoteModeResponse {
1441                ok: false,
1442                error: e,
1443                mode: String::new(),
1444            })),
1445        }
1446    }
1447
1448    // ── Governance policy lifecycle RPCs (RFC-MACP-0012) ────────────
1449
1450    async fn register_policy(
1451        &self,
1452        request: Request<RegisterPolicyRequest>,
1453    ) -> Result<Response<RegisterPolicyResponse>, Status> {
1454        if self.policies_read_only {
1455            return Err(Status::failed_precondition(
1456                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1457            ));
1458        }
1459        let identity = self
1460            .security
1461            .authenticate_metadata(request.metadata())
1462            .await
1463            .map_err(Self::status_from_error)?;
1464        self.security
1465            .authorize_mode_registry(&identity)
1466            .map_err(Self::status_from_error)?;
1467        let req = request.into_inner();
1468        let descriptor = req
1469            .policy_descriptor
1470            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1471        let definition = Self::policy_descriptor_to_definition(&descriptor);
1472        match self.runtime.register_policy(definition) {
1473            Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1474                ok: true,
1475                error: String::new(),
1476            })),
1477            Err(e) => Ok(Response::new(RegisterPolicyResponse {
1478                ok: false,
1479                error: e,
1480            })),
1481        }
1482    }
1483
1484    async fn unregister_policy(
1485        &self,
1486        request: Request<UnregisterPolicyRequest>,
1487    ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1488        if self.policies_read_only {
1489            return Err(Status::failed_precondition(
1490                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1491            ));
1492        }
1493        let identity = self
1494            .security
1495            .authenticate_metadata(request.metadata())
1496            .await
1497            .map_err(Self::status_from_error)?;
1498        self.security
1499            .authorize_mode_registry(&identity)
1500            .map_err(Self::status_from_error)?;
1501        let req = request.into_inner();
1502        match self.runtime.unregister_policy(&req.policy_id) {
1503            Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1504                ok: true,
1505                error: String::new(),
1506            })),
1507            Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1508                ok: false,
1509                error: e,
1510            })),
1511        }
1512    }
1513
1514    async fn get_policy(
1515        &self,
1516        request: Request<GetPolicyRequest>,
1517    ) -> Result<Response<GetPolicyResponse>, Status> {
1518        let _identity = self
1519            .security
1520            .authenticate_metadata(request.metadata())
1521            .await
1522            .map_err(Self::status_from_error)?;
1523        let req = request.into_inner();
1524        let policy = self
1525            .runtime
1526            .get_policy(&req.policy_id)
1527            .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1528        Ok(Response::new(GetPolicyResponse {
1529            policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1530        }))
1531    }
1532
1533    async fn list_policies(
1534        &self,
1535        request: Request<ListPoliciesRequest>,
1536    ) -> Result<Response<ListPoliciesResponse>, Status> {
1537        let _identity = self
1538            .security
1539            .authenticate_metadata(request.metadata())
1540            .await
1541            .map_err(Self::status_from_error)?;
1542        let req = request.into_inner();
1543        let mode_filter = if req.mode.is_empty() {
1544            None
1545        } else {
1546            Some(req.mode.as_str())
1547        };
1548        let policies = self.runtime.list_policies(mode_filter);
1549        let descriptors = policies
1550            .iter()
1551            .map(Self::policy_definition_to_descriptor)
1552            .collect();
1553        Ok(Response::new(ListPoliciesResponse { descriptors }))
1554    }
1555
1556    type WatchPoliciesStream = std::pin::Pin<
1557        Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1558    >;
1559
1560    async fn watch_policies(
1561        &self,
1562        _request: Request<WatchPoliciesRequest>,
1563    ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1564        let mut rx = self.runtime.subscribe_policy_changes();
1565        let runtime = Arc::clone(&self.runtime);
1566        let stream = async_stream::try_stream! {
1567            // Send initial state
1568            let policies = runtime.list_policies(None);
1569            let descriptors: Vec<PolicyDescriptor> = policies
1570                .iter()
1571                .map(MacpServer::policy_definition_to_descriptor)
1572                .collect();
1573            yield WatchPoliciesResponse {
1574                descriptors,
1575                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1576            };
1577            // Wait for changes
1578            while rx.recv().await.is_ok() {
1579                let policies = runtime.list_policies(None);
1580                let descriptors: Vec<PolicyDescriptor> = policies
1581                    .iter()
1582                    .map(MacpServer::policy_definition_to_descriptor)
1583                    .collect();
1584                yield WatchPoliciesResponse {
1585                    descriptors,
1586                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1587                };
1588            }
1589        };
1590        Ok(Response::new(Box::pin(stream)))
1591    }
1592}
1593
1594// ── Policy type conversion helpers ──────────────────────────────────
1595
1596impl MacpServer {
1597    fn policy_descriptor_to_definition(
1598        descriptor: &PolicyDescriptor,
1599    ) -> crate::policy::PolicyDefinition {
1600        let rules: serde_json::Value = if descriptor.rules.is_empty() {
1601            serde_json::json!({})
1602        } else {
1603            serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1604        };
1605        crate::policy::PolicyDefinition {
1606            policy_id: descriptor.policy_id.clone(),
1607            mode: descriptor.mode.clone(),
1608            description: descriptor.description.clone(),
1609            rules,
1610            schema_version: descriptor.schema_version,
1611        }
1612    }
1613
1614    fn policy_definition_to_descriptor(
1615        definition: &crate::policy::PolicyDefinition,
1616    ) -> PolicyDescriptor {
1617        PolicyDescriptor {
1618            policy_id: definition.policy_id.clone(),
1619            mode: definition.mode.clone(),
1620            description: definition.description.clone(),
1621            rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1622            schema_version: definition.schema_version,
1623            registered_at_unix_ms: 0,
1624        }
1625    }
1626}
1627
1628#[cfg(test)]
1629mod tests {
1630    use super::*;
1631    use crate::log_store::LogStore;
1632    use crate::pb::SessionStartPayload;
1633    use crate::registry::SessionRegistry;
1634    use chrono::Utc;
1635    use prost::Message;
1636
1637    fn new_sid() -> String {
1638        uuid::Uuid::new_v4().as_hyphenated().to_string()
1639    }
1640
1641    fn make_server() -> (MacpServer, Arc<Runtime>) {
1642        let storage: Arc<dyn crate::storage::StorageBackend> =
1643            Arc::new(crate::storage::MemoryBackend);
1644        let registry = Arc::new(SessionRegistry::new());
1645        let log_store = Arc::new(LogStore::new());
1646        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1647        let server = MacpServer::new(runtime.clone(), SecurityLayer::dev_mode());
1648        (server, runtime)
1649    }
1650
1651    fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1652        let mut req = Request::new(SendRequest {
1653            envelope: Some(env),
1654        });
1655        req.metadata_mut()
1656            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1657        req
1658    }
1659
1660    async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1661        let resp = server.send(send_req(sender, env)).await.unwrap();
1662        resp.into_inner().ack.unwrap()
1663    }
1664
1665    fn start_payload() -> Vec<u8> {
1666        SessionStartPayload {
1667            intent: "intent".into(),
1668            participants: vec!["agent://fraud".into()],
1669            mode_version: "1.0.0".into(),
1670            configuration_version: "cfg-1".into(),
1671            policy_version: String::new(),
1672            ttl_ms: 1000,
1673            context_id: String::new(),
1674            extensions: std::collections::HashMap::new(),
1675            roots: vec![],
1676            max_suspend_ms: 0,
1677        }
1678        .encode_to_vec()
1679    }
1680
1681    #[tokio::test]
1682    async fn sender_is_derived_from_authenticated_metadata() {
1683        let (server, runtime) = make_server();
1684        let sid = new_sid();
1685        let ack = do_send(
1686            &server,
1687            "agent://orchestrator",
1688            Envelope {
1689                macp_version: "1.0".into(),
1690                mode: "macp.mode.decision.v1".into(),
1691                message_type: "SessionStart".into(),
1692                message_id: "m1".into(),
1693                session_id: sid.clone(),
1694                sender: String::new(),
1695                timestamp_unix_ms: Utc::now().timestamp_millis(),
1696                payload: start_payload(),
1697            },
1698        )
1699        .await;
1700        assert!(ack.ok);
1701        let session = runtime.get_session_checked(&sid).await.unwrap();
1702        assert_eq!(session.initiator_sender, "agent://orchestrator");
1703    }
1704
1705    #[tokio::test]
1706    async fn spoofed_sender_is_rejected() {
1707        let (server, _) = make_server();
1708        let sid = new_sid();
1709        let ack = do_send(
1710            &server,
1711            "agent://orchestrator",
1712            Envelope {
1713                macp_version: "1.0".into(),
1714                mode: "macp.mode.decision.v1".into(),
1715                message_type: "SessionStart".into(),
1716                message_id: "m1".into(),
1717                session_id: sid,
1718                sender: "agent://spoof".into(),
1719                timestamp_unix_ms: Utc::now().timestamp_millis(),
1720                payload: start_payload(),
1721            },
1722        )
1723        .await;
1724        assert!(!ack.ok);
1725        assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1726    }
1727
1728    #[tokio::test]
1729    async fn get_session_requires_session_membership() {
1730        let (server, _) = make_server();
1731        let sid = new_sid();
1732        let ack = do_send(
1733            &server,
1734            "agent://orchestrator",
1735            Envelope {
1736                macp_version: "1.0".into(),
1737                mode: "macp.mode.decision.v1".into(),
1738                message_type: "SessionStart".into(),
1739                message_id: "m1".into(),
1740                session_id: sid.clone(),
1741                sender: String::new(),
1742                timestamp_unix_ms: Utc::now().timestamp_millis(),
1743                payload: start_payload(),
1744            },
1745        )
1746        .await;
1747        assert!(ack.ok);
1748
1749        let mut req = Request::new(GetSessionRequest { session_id: sid });
1750        req.metadata_mut().insert(
1751            "authorization",
1752            format!("Bearer {}", "agent://outsider").parse().unwrap(),
1753        );
1754        let err = server.get_session(req).await.unwrap_err();
1755        assert_eq!(err.code(), tonic::Code::PermissionDenied);
1756    }
1757
1758    #[tokio::test]
1759    async fn register_ext_mode_requires_authenticated_registry_permission() {
1760        let storage: Arc<dyn crate::storage::StorageBackend> =
1761            Arc::new(crate::storage::MemoryBackend);
1762        let registry = Arc::new(SessionRegistry::new());
1763        let log_store = Arc::new(LogStore::new());
1764        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1765        let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1766        let server = MacpServer::new(runtime, security);
1767
1768        let req = Request::new(RegisterExtModeRequest {
1769            mode_descriptor: Some(crate::pb::ModeDescriptor {
1770                mode: "ext.custom.v1".into(),
1771                mode_version: "1.0.0".into(),
1772                message_types: vec!["SessionStart".into(), "Commitment".into()],
1773                ..Default::default()
1774            }),
1775        });
1776        let err = server.register_ext_mode(req).await.unwrap_err();
1777        assert_eq!(err.code(), tonic::Code::Unauthenticated);
1778    }
1779
1780    fn stream_identity(sender: &str) -> AuthIdentity {
1781        AuthIdentity {
1782            sender: sender.into(),
1783            allowed_modes: None,
1784            can_start_sessions: true,
1785            max_open_sessions: None,
1786            can_manage_mode_registry: false,
1787            is_observer: false,
1788        }
1789    }
1790
1791    #[tokio::test]
1792    async fn stream_session_emits_accepted_envelopes_only() {
1793        use tokio_stream::{iter, StreamExt};
1794
1795        let (server, _) = make_server();
1796        let sid = new_sid();
1797        let requests = iter(vec![Ok(StreamSessionRequest {
1798            subscribe_session_id: String::new(),
1799            after_sequence: 0,
1800            envelope: Some(Envelope {
1801                macp_version: "1.0".into(),
1802                mode: "macp.mode.decision.v1".into(),
1803                message_type: "SessionStart".into(),
1804                message_id: "m1".into(),
1805                session_id: sid.clone(),
1806                sender: String::new(),
1807                timestamp_unix_ms: Utc::now().timestamp_millis(),
1808                payload: start_payload(),
1809            }),
1810        })]);
1811
1812        let mut stream =
1813            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1814
1815        let response = stream.next().await.unwrap().unwrap();
1816        let envelope = match response.response.unwrap() {
1817            crate::pb::stream_session_response::Response::Envelope(e) => e,
1818            _ => panic!("expected envelope"),
1819        };
1820        assert_eq!(envelope.message_type, "SessionStart");
1821        assert_eq!(envelope.message_id, "m1");
1822        assert!(stream.next().await.is_none());
1823    }
1824
1825    #[tokio::test]
1826    async fn stream_session_rejects_mixed_session_ids() {
1827        use tokio_stream::{iter, StreamExt};
1828
1829        let (server, _) = make_server();
1830        let sid1 = new_sid();
1831        let sid2 = new_sid();
1832        let requests = iter(vec![
1833            Ok(StreamSessionRequest {
1834                subscribe_session_id: String::new(),
1835                after_sequence: 0,
1836                envelope: Some(Envelope {
1837                    macp_version: "1.0".into(),
1838                    mode: "macp.mode.decision.v1".into(),
1839                    message_type: "SessionStart".into(),
1840                    message_id: "m1".into(),
1841                    session_id: sid1.clone(),
1842                    sender: String::new(),
1843                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1844                    payload: start_payload(),
1845                }),
1846            }),
1847            Ok(StreamSessionRequest {
1848                subscribe_session_id: String::new(),
1849                after_sequence: 0,
1850                envelope: Some(Envelope {
1851                    macp_version: "1.0".into(),
1852                    mode: "macp.mode.decision.v1".into(),
1853                    message_type: "SessionStart".into(),
1854                    message_id: "m2".into(),
1855                    session_id: sid2,
1856                    sender: String::new(),
1857                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1858                    payload: start_payload(),
1859                }),
1860            }),
1861        ]);
1862
1863        let mut stream =
1864            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1865
1866        let first = stream.next().await.unwrap().unwrap();
1867        let first_env = match first.response.unwrap() {
1868            crate::pb::stream_session_response::Response::Envelope(e) => e,
1869            _ => panic!("expected envelope"),
1870        };
1871        assert_eq!(first_env.session_id, sid1);
1872        let err = stream.next().await.unwrap().unwrap_err();
1873        assert_eq!(err.code(), tonic::Code::InvalidArgument);
1874    }
1875
1876    #[tokio::test]
1877    async fn list_modes_returns_standard_modes() {
1878        let (server, _) = make_server();
1879        let resp = server
1880            .list_modes(Request::new(ListModesRequest {}))
1881            .await
1882            .unwrap();
1883        let names: Vec<String> = resp
1884            .into_inner()
1885            .modes
1886            .iter()
1887            .map(|m| m.mode.clone())
1888            .collect();
1889        assert_eq!(names.len(), 5);
1890        assert!(names.contains(&"macp.mode.decision.v1".to_string()));
1891        assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
1892        assert!(names.contains(&"macp.mode.task.v1".to_string()));
1893        assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
1894        assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
1895        // multi_round is now an extension, not in ListModes
1896        assert!(!names.contains(&"ext.multi_round.v1".to_string()));
1897    }
1898
1899    #[tokio::test]
1900    async fn list_ext_modes_returns_extensions() {
1901        let (server, _) = make_server();
1902        let resp = server
1903            .list_ext_modes(Request::new(ListExtModesRequest {}))
1904            .await
1905            .unwrap();
1906        let names: Vec<String> = resp
1907            .into_inner()
1908            .modes
1909            .iter()
1910            .map(|m| m.mode.clone())
1911            .collect();
1912        assert_eq!(names.len(), 1);
1913        assert!(names.contains(&"ext.multi_round.v1".to_string()));
1914    }
1915
1916    #[tokio::test]
1917    async fn get_manifest_includes_all_modes() {
1918        let (server, _) = make_server();
1919        let resp = server
1920            .get_manifest(Request::new(crate::pb::GetManifestRequest {
1921                agent_id: String::new(),
1922            }))
1923            .await
1924            .unwrap();
1925        let manifest = resp.into_inner().manifest.unwrap();
1926        assert_eq!(manifest.supported_modes.len(), 6);
1927        assert!(manifest
1928            .supported_modes
1929            .contains(&"ext.multi_round.v1".to_string()));
1930    }
1931
1932    #[tokio::test]
1933    async fn get_session_returns_metadata() {
1934        let (server, _) = make_server();
1935        let sid = new_sid();
1936        let ack = do_send(
1937            &server,
1938            "agent://orchestrator",
1939            Envelope {
1940                macp_version: "1.0".into(),
1941                mode: "macp.mode.decision.v1".into(),
1942                message_type: "SessionStart".into(),
1943                message_id: "m1".into(),
1944                session_id: sid.clone(),
1945                sender: String::new(),
1946                timestamp_unix_ms: Utc::now().timestamp_millis(),
1947                payload: start_payload(),
1948            },
1949        )
1950        .await;
1951        assert!(ack.ok);
1952
1953        let mut req = Request::new(GetSessionRequest {
1954            session_id: sid.clone(),
1955        });
1956        req.metadata_mut().insert(
1957            "authorization",
1958            format!("Bearer {}", "agent://orchestrator")
1959                .parse()
1960                .unwrap(),
1961        );
1962        let resp = server.get_session(req).await.unwrap();
1963        let meta = resp.into_inner().metadata.unwrap();
1964        assert_eq!(meta.session_id, sid);
1965        assert_eq!(meta.mode, "macp.mode.decision.v1");
1966        assert_eq!(meta.mode_version, "1.0.0");
1967        assert_eq!(meta.configuration_version, "cfg-1");
1968    }
1969
1970    #[tokio::test]
1971    async fn cancel_session_transitions_to_cancelled() {
1972        let (server, _) = make_server();
1973        let sid = new_sid();
1974        let ack = do_send(
1975            &server,
1976            "agent://orchestrator",
1977            Envelope {
1978                macp_version: "1.0".into(),
1979                mode: "macp.mode.decision.v1".into(),
1980                message_type: "SessionStart".into(),
1981                message_id: "m1".into(),
1982                session_id: sid.clone(),
1983                sender: String::new(),
1984                timestamp_unix_ms: Utc::now().timestamp_millis(),
1985                payload: start_payload(),
1986            },
1987        )
1988        .await;
1989        assert!(ack.ok);
1990
1991        let mut req = Request::new(CancelSessionRequest {
1992            session_id: sid,
1993            reason: "no longer needed".into(),
1994        });
1995        req.metadata_mut().insert(
1996            "authorization",
1997            format!("Bearer {}", "agent://orchestrator")
1998                .parse()
1999                .unwrap(),
2000        );
2001        let resp = server.cancel_session(req).await.unwrap();
2002        let ack = resp.into_inner().ack.unwrap();
2003        assert!(ack.ok);
2004        // RFC-MACP-0001 §7.3: cancellation now yields the distinct CANCELLED state.
2005        assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2006    }
2007
2008    #[tokio::test]
2009    async fn participant_cannot_cancel_session() {
2010        let (server, _) = make_server();
2011        let sid = new_sid();
2012        let ack = do_send(
2013            &server,
2014            "agent://orchestrator",
2015            Envelope {
2016                macp_version: "1.0".into(),
2017                mode: "macp.mode.decision.v1".into(),
2018                message_type: "SessionStart".into(),
2019                message_id: "m1".into(),
2020                session_id: sid.clone(),
2021                sender: String::new(),
2022                timestamp_unix_ms: Utc::now().timestamp_millis(),
2023                payload: start_payload(),
2024            },
2025        )
2026        .await;
2027        assert!(ack.ok);
2028
2029        let mut req = Request::new(CancelSessionRequest {
2030            session_id: sid,
2031            reason: "I want to cancel".into(),
2032        });
2033        req.metadata_mut().insert(
2034            "authorization",
2035            format!("Bearer {}", "agent://fraud").parse().unwrap(),
2036        );
2037        let err = server.cancel_session(req).await.unwrap_err();
2038        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2039    }
2040
2041    #[tokio::test]
2042    async fn cancel_session_unknown_session_returns_error() {
2043        let (server, _) = make_server();
2044        let mut req = Request::new(CancelSessionRequest {
2045            session_id: "nonexistent".into(),
2046            reason: "test".into(),
2047        });
2048        req.metadata_mut().insert(
2049            "authorization",
2050            format!("Bearer {}", "agent://orchestrator")
2051                .parse()
2052                .unwrap(),
2053        );
2054        let err = server.cancel_session(req).await.unwrap_err();
2055        assert_eq!(err.code(), tonic::Code::NotFound);
2056    }
2057
2058    #[tokio::test]
2059    async fn ambient_signal_accepted() {
2060        let (server, _) = make_server();
2061        let ack = do_send(
2062            &server,
2063            "agent://orchestrator",
2064            Envelope {
2065                macp_version: "1.0".into(),
2066                mode: String::new(),
2067                message_type: "Signal".into(),
2068                message_id: "sig-1".into(),
2069                session_id: String::new(),
2070                sender: String::new(),
2071                timestamp_unix_ms: Utc::now().timestamp_millis(),
2072                payload: vec![],
2073            },
2074        )
2075        .await;
2076        assert!(ack.ok);
2077    }
2078
2079    #[tokio::test]
2080    async fn signal_with_session_id_rejected() {
2081        let (server, _) = make_server();
2082        let ack = do_send(
2083            &server,
2084            "agent://orchestrator",
2085            Envelope {
2086                macp_version: "1.0".into(),
2087                mode: String::new(),
2088                message_type: "Signal".into(),
2089                message_id: "sig-2".into(),
2090                session_id: "some-session".into(),
2091                sender: String::new(),
2092                timestamp_unix_ms: Utc::now().timestamp_millis(),
2093                payload: vec![],
2094            },
2095        )
2096        .await;
2097        assert!(!ack.ok);
2098        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2099    }
2100
2101    #[tokio::test]
2102    async fn signal_with_mode_rejected() {
2103        let (server, _) = make_server();
2104        let ack = do_send(
2105            &server,
2106            "agent://orchestrator",
2107            Envelope {
2108                macp_version: "1.0".into(),
2109                mode: "macp.mode.decision.v1".into(),
2110                message_type: "Signal".into(),
2111                message_id: "sig-3".into(),
2112                session_id: String::new(),
2113                sender: String::new(),
2114                timestamp_unix_ms: Utc::now().timestamp_millis(),
2115                payload: vec![],
2116            },
2117        )
2118        .await;
2119        assert!(!ack.ok);
2120        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2121    }
2122
2123    #[tokio::test]
2124    async fn ambient_progress_accepted() {
2125        let (server, _) = make_server();
2126        let ack = do_send(
2127            &server,
2128            "agent://orchestrator",
2129            Envelope {
2130                macp_version: "1.0".into(),
2131                mode: String::new(),
2132                message_type: "Progress".into(),
2133                message_id: "prog-1".into(),
2134                session_id: String::new(),
2135                sender: String::new(),
2136                timestamp_unix_ms: Utc::now().timestamp_millis(),
2137                payload: vec![],
2138            },
2139        )
2140        .await;
2141        assert!(ack.ok);
2142    }
2143
2144    #[tokio::test]
2145    async fn ambient_progress_with_mode_rejected() {
2146        let (server, _) = make_server();
2147        let ack = do_send(
2148            &server,
2149            "agent://orchestrator",
2150            Envelope {
2151                macp_version: "1.0".into(),
2152                mode: "macp.mode.decision.v1".into(),
2153                message_type: "Progress".into(),
2154                message_id: "prog-2".into(),
2155                session_id: String::new(),
2156                sender: String::new(),
2157                timestamp_unix_ms: Utc::now().timestamp_millis(),
2158                payload: vec![],
2159            },
2160        )
2161        .await;
2162        assert!(!ack.ok);
2163        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2164    }
2165
2166    #[tokio::test]
2167    async fn manifest_advertises_stream_enabled() {
2168        let (server, _) = make_server();
2169        let resp = server
2170            .initialize(Request::new(InitializeRequest {
2171                supported_protocol_versions: vec!["1.0".into()],
2172                client_info: None,
2173                capabilities: None,
2174            }))
2175            .await
2176            .unwrap();
2177        let caps = resp.into_inner().capabilities.unwrap();
2178        assert!(caps.sessions.unwrap().stream);
2179    }
2180
2181    #[tokio::test]
2182    async fn initialize_empty_versions_rejected() {
2183        let (server, _) = make_server();
2184        let err = server
2185            .initialize(Request::new(InitializeRequest {
2186                supported_protocol_versions: vec![],
2187                client_info: None,
2188                capabilities: None,
2189            }))
2190            .await
2191            .unwrap_err();
2192        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2193    }
2194
2195    #[tokio::test]
2196    async fn initialize_unsupported_version_rejected() {
2197        let (server, _) = make_server();
2198        let err = server
2199            .initialize(Request::new(InitializeRequest {
2200                supported_protocol_versions: vec!["2.0".into()],
2201                client_info: None,
2202                capabilities: None,
2203            }))
2204            .await
2205            .unwrap_err();
2206        assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2207    }
2208
2209    // ── RFC-MACP-0006-A1: passive subscribe tests ──────────────────────
2210
2211    fn observer_identity(sender: &str) -> AuthIdentity {
2212        AuthIdentity {
2213            sender: sender.into(),
2214            allowed_modes: None,
2215            can_start_sessions: false,
2216            max_open_sessions: None,
2217            can_manage_mode_registry: false,
2218            is_observer: true,
2219        }
2220    }
2221
2222    fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2223        StreamSessionRequest {
2224            subscribe_session_id: session_id.into(),
2225            after_sequence: after,
2226            envelope: None,
2227        }
2228    }
2229
2230    fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2231        SessionStartPayload {
2232            intent: "intent".into(),
2233            participants,
2234            mode_version: "1.0.0".into(),
2235            configuration_version: "cfg-1".into(),
2236            policy_version: String::new(),
2237            ttl_ms: 60_000,
2238            context_id: String::new(),
2239            extensions: std::collections::HashMap::new(),
2240            roots: vec![],
2241            max_suspend_ms: 0,
2242        }
2243        .encode_to_vec()
2244    }
2245
2246    async fn start_session(
2247        server: &MacpServer,
2248        initiator: &str,
2249        sid: &str,
2250        participants: Vec<String>,
2251    ) {
2252        let ack = do_send(
2253            server,
2254            initiator,
2255            Envelope {
2256                macp_version: "1.0".into(),
2257                mode: "macp.mode.decision.v1".into(),
2258                message_type: "SessionStart".into(),
2259                message_id: "start".into(),
2260                session_id: sid.into(),
2261                sender: String::new(),
2262                timestamp_unix_ms: Utc::now().timestamp_millis(),
2263                payload: start_multi_participant(participants),
2264            },
2265        )
2266        .await;
2267        assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2268    }
2269
2270    async fn send_proposal(
2271        server: &MacpServer,
2272        sender: &str,
2273        sid: &str,
2274        message_id: &str,
2275        proposal_id: &str,
2276    ) {
2277        let payload = crate::decision_pb::ProposalPayload {
2278            proposal_id: proposal_id.into(),
2279            option: "opt".into(),
2280            rationale: "r".into(),
2281            supporting_data: vec![],
2282        }
2283        .encode_to_vec();
2284        let ack = do_send(
2285            server,
2286            sender,
2287            Envelope {
2288                macp_version: "1.0".into(),
2289                mode: "macp.mode.decision.v1".into(),
2290                message_type: "Proposal".into(),
2291                message_id: message_id.into(),
2292                session_id: sid.into(),
2293                sender: String::new(),
2294                timestamp_unix_ms: Utc::now().timestamp_millis(),
2295                payload,
2296            },
2297        )
2298        .await;
2299        assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2300    }
2301
2302    #[tokio::test]
2303    async fn subscribe_replays_session_history_from_zero() {
2304        let (server, _) = make_server();
2305        let sid = new_sid();
2306        let initiator = "agent://orchestrator";
2307        let peer = "agent://fraud";
2308        start_session(
2309            &server,
2310            initiator,
2311            &sid,
2312            vec![initiator.into(), peer.into()],
2313        )
2314        .await;
2315        send_proposal(&server, peer, &sid, "m2", "p1").await;
2316
2317        let mut bound = None;
2318        let mut events = None;
2319        let replay = server
2320            .process_stream_request(
2321                &stream_identity(peer),
2322                subscribe_frame(&sid, 0),
2323                &mut bound,
2324                &mut events,
2325            )
2326            .await
2327            .unwrap();
2328
2329        assert_eq!(replay.len(), 2);
2330        assert_eq!(replay[0].message_type, "SessionStart");
2331        assert_eq!(replay[0].message_id, "start");
2332        assert_eq!(replay[1].message_type, "Proposal");
2333        assert_eq!(replay[1].message_id, "m2");
2334        assert_eq!(bound.as_deref(), Some(sid.as_str()));
2335        assert!(events.is_some());
2336    }
2337
2338    #[tokio::test]
2339    async fn subscribe_after_sequence_filters_history() {
2340        let (server, _) = make_server();
2341        let sid = new_sid();
2342        let initiator = "agent://orchestrator";
2343        let peer = "agent://fraud";
2344        start_session(
2345            &server,
2346            initiator,
2347            &sid,
2348            vec![initiator.into(), peer.into()],
2349        )
2350        .await;
2351        send_proposal(&server, peer, &sid, "m2", "p1").await;
2352        send_proposal(&server, peer, &sid, "m3", "p2").await;
2353
2354        let mut bound = None;
2355        let mut events = None;
2356        let replay = server
2357            .process_stream_request(
2358                &stream_identity(peer),
2359                subscribe_frame(&sid, 2),
2360                &mut bound,
2361                &mut events,
2362            )
2363            .await
2364            .unwrap();
2365
2366        assert_eq!(replay.len(), 1);
2367        assert_eq!(replay[0].message_id, "m3");
2368    }
2369
2370    #[tokio::test]
2371    async fn subscribe_unknown_session_returns_not_found() {
2372        let (server, _) = make_server();
2373        let mut bound = None;
2374        let mut events = None;
2375        let status = server
2376            .process_stream_request(
2377                &stream_identity("agent://orchestrator"),
2378                subscribe_frame("missing-session", 0),
2379                &mut bound,
2380                &mut events,
2381            )
2382            .await
2383            .unwrap_err();
2384        assert_eq!(status.code(), tonic::Code::NotFound);
2385        assert!(bound.is_none());
2386        assert!(events.is_none());
2387    }
2388
2389    #[tokio::test]
2390    async fn subscribe_non_participant_is_forbidden() {
2391        let (server, _) = make_server();
2392        let sid = new_sid();
2393        start_session(
2394            &server,
2395            "agent://orchestrator",
2396            &sid,
2397            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2398        )
2399        .await;
2400
2401        let mut bound = None;
2402        let mut events = None;
2403        let status = server
2404            .process_stream_request(
2405                &stream_identity("agent://outsider"),
2406                subscribe_frame(&sid, 0),
2407                &mut bound,
2408                &mut events,
2409            )
2410            .await
2411            .unwrap_err();
2412        assert_eq!(status.code(), tonic::Code::PermissionDenied);
2413    }
2414
2415    #[tokio::test]
2416    async fn subscribe_observer_identity_allowed() {
2417        let (server, _) = make_server();
2418        let sid = new_sid();
2419        start_session(
2420            &server,
2421            "agent://orchestrator",
2422            &sid,
2423            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2424        )
2425        .await;
2426
2427        let mut bound = None;
2428        let mut events = None;
2429        let replay = server
2430            .process_stream_request(
2431                &observer_identity("agent://auditor"),
2432                subscribe_frame(&sid, 0),
2433                &mut bound,
2434                &mut events,
2435            )
2436            .await
2437            .unwrap();
2438        assert_eq!(replay.len(), 1);
2439        assert_eq!(replay[0].message_type, "SessionStart");
2440    }
2441
2442    #[tokio::test]
2443    async fn subscribe_initiator_allowed_even_when_not_listed() {
2444        // Per RFC-MACP-0007, the initiator is always authorized for session
2445        // access, even if not present in the participants list.
2446        let (server, _) = make_server();
2447        let sid = new_sid();
2448        start_session(
2449            &server,
2450            "agent://orchestrator",
2451            &sid,
2452            vec!["agent://fraud".into()],
2453        )
2454        .await;
2455
2456        let mut bound = None;
2457        let mut events = None;
2458        let replay = server
2459            .process_stream_request(
2460                &stream_identity("agent://orchestrator"),
2461                subscribe_frame(&sid, 0),
2462                &mut bound,
2463                &mut events,
2464            )
2465            .await
2466            .unwrap();
2467        assert_eq!(replay.len(), 1);
2468    }
2469
2470    #[tokio::test]
2471    async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2472        let (server, _) = make_server();
2473        let sid = new_sid();
2474        let req = StreamSessionRequest {
2475            subscribe_session_id: sid.clone(),
2476            after_sequence: 0,
2477            envelope: Some(Envelope {
2478                macp_version: "1.0".into(),
2479                mode: "macp.mode.decision.v1".into(),
2480                message_type: "SessionStart".into(),
2481                message_id: "m1".into(),
2482                session_id: sid,
2483                sender: String::new(),
2484                timestamp_unix_ms: Utc::now().timestamp_millis(),
2485                payload: start_payload(),
2486            }),
2487        };
2488
2489        let mut bound = None;
2490        let mut events = None;
2491        let status = server
2492            .process_stream_request(
2493                &stream_identity("agent://orchestrator"),
2494                req,
2495                &mut bound,
2496                &mut events,
2497            )
2498            .await
2499            .unwrap_err();
2500        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2501    }
2502
2503    #[tokio::test]
2504    async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2505        let (server, _) = make_server();
2506        let sid1 = new_sid();
2507        let sid2 = new_sid();
2508        start_session(
2509            &server,
2510            "agent://orchestrator",
2511            &sid1,
2512            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2513        )
2514        .await;
2515        start_session(
2516            &server,
2517            "agent://orchestrator",
2518            &sid2,
2519            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2520        )
2521        .await;
2522
2523        // First subscribe binds the stream to sid1
2524        let identity = stream_identity("agent://fraud");
2525        let mut bound = None;
2526        let mut events = None;
2527        server
2528            .process_stream_request(
2529                &identity,
2530                subscribe_frame(&sid1, 0),
2531                &mut bound,
2532                &mut events,
2533            )
2534            .await
2535            .unwrap();
2536        assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2537
2538        // Second subscribe to sid2 on the same stream must be rejected
2539        let status = server
2540            .process_stream_request(
2541                &identity,
2542                subscribe_frame(&sid2, 0),
2543                &mut bound,
2544                &mut events,
2545            )
2546            .await
2547            .unwrap_err();
2548        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2549    }
2550
2551    /// E3: an injected ingress engine gates session start, messages, and
2552    /// session reads — deny-one-sender double proves all three hooks fire and
2553    /// that denial surfaces as POLICY_DENIED / PermissionDenied (fail closed).
2554    struct DenySenderEngine {
2555        denied: String,
2556    }
2557
2558    #[async_trait::async_trait]
2559    impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2560        async fn evaluate_session_start(
2561            &self,
2562            identity: &crate::security::AuthIdentity,
2563            _mode: &str,
2564            _env: &Envelope,
2565        ) -> macp_core::policy::PolicyDecision {
2566            if identity.sender == self.denied {
2567                macp_core::policy::PolicyDecision::Deny {
2568                    reasons: vec!["sender embargoed".into()],
2569                }
2570            } else {
2571                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2572            }
2573        }
2574
2575        async fn evaluate_message(
2576            &self,
2577            identity: &crate::security::AuthIdentity,
2578            _session: &macp_core::session::Session,
2579            _env: &Envelope,
2580        ) -> macp_core::policy::PolicyDecision {
2581            if identity.sender == self.denied {
2582                macp_core::policy::PolicyDecision::Deny {
2583                    reasons: vec!["sender embargoed".into()],
2584                }
2585            } else {
2586                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2587            }
2588        }
2589
2590        async fn evaluate_session_access(
2591            &self,
2592            identity: &crate::security::AuthIdentity,
2593            _session: &macp_core::session::Session,
2594        ) -> macp_core::policy::PolicyDecision {
2595            if identity.sender == self.denied {
2596                macp_core::policy::PolicyDecision::Deny {
2597                    reasons: vec!["sender embargoed".into()],
2598                }
2599            } else {
2600                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2601            }
2602        }
2603    }
2604
2605    #[tokio::test]
2606    async fn policy_engine_gates_all_three_ingress_points() {
2607        let (server, _runtime) = make_server();
2608        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2609            denied: "agent://embargoed".into(),
2610        }));
2611
2612        let sid = new_sid();
2613        let start_payload = SessionStartPayload {
2614            intent: "e3".into(),
2615            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2616            mode_version: "1.0.0".into(),
2617            configuration_version: "cfg-1".into(),
2618            policy_version: String::new(),
2619            ttl_ms: 60_000,
2620            context_id: String::new(),
2621            extensions: Default::default(),
2622            roots: vec![],
2623            max_suspend_ms: 0,
2624        }
2625        .encode_to_vec();
2626        let start_env = |sender: &str, sid: &str| Envelope {
2627            macp_version: "1.0".into(),
2628            mode: "macp.mode.decision.v1".into(),
2629            message_type: "SessionStart".into(),
2630            message_id: new_sid(),
2631            session_id: sid.into(),
2632            sender: sender.into(),
2633            timestamp_unix_ms: Utc::now().timestamp_millis(),
2634            payload: start_payload.clone(),
2635        };
2636
2637        // 1. Embargoed sender cannot start a session.
2638        let ack = server
2639            .send(send_req(
2640                "agent://embargoed",
2641                start_env("agent://embargoed", &sid),
2642            ))
2643            .await
2644            .unwrap()
2645            .into_inner()
2646            .ack
2647            .unwrap();
2648        assert!(!ack.ok);
2649        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2650
2651        // Allowed sender starts it.
2652        let ack = server
2653            .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2654            .await
2655            .unwrap()
2656            .into_inner()
2657            .ack
2658            .unwrap();
2659        assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2660
2661        // 2. Embargoed sender cannot send into the session.
2662        let proposal = crate::decision_pb::ProposalPayload {
2663            proposal_id: "p1".into(),
2664            option: "x".into(),
2665            rationale: "r".into(),
2666            supporting_data: vec![],
2667        }
2668        .encode_to_vec();
2669        let msg_env = Envelope {
2670            macp_version: "1.0".into(),
2671            mode: "macp.mode.decision.v1".into(),
2672            message_type: "Proposal".into(),
2673            message_id: new_sid(),
2674            session_id: sid.clone(),
2675            sender: "agent://embargoed".into(),
2676            timestamp_unix_ms: Utc::now().timestamp_millis(),
2677            payload: proposal,
2678        };
2679        let ack = server
2680            .send(send_req("agent://embargoed", msg_env))
2681            .await
2682            .unwrap()
2683            .into_inner()
2684            .ack
2685            .unwrap();
2686        assert!(!ack.ok);
2687        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2688
2689        // 3. Embargoed sender cannot read the session.
2690        let mut req = Request::new(crate::pb::GetSessionRequest {
2691            session_id: sid.clone(),
2692        });
2693        req.metadata_mut()
2694            .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2695        let err = server
2696            .get_session(req)
2697            .await
2698            .expect_err("embargoed read must be denied");
2699        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2700    }
2701
2702    /// E3 transport-parity: the ingress engine gates the STREAM path too — a
2703    /// denied sender must not be able to bypass the engine by switching from
2704    /// unary Send to StreamSession (envelope frames or subscribe frames).
2705    #[tokio::test]
2706    async fn policy_engine_gates_stream_path() {
2707        let (server, runtime) = make_server();
2708        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2709            denied: "agent://embargoed".into(),
2710        }));
2711
2712        // Session started by an allowed sender (participants include the
2713        // embargoed agent so built-in membership checks pass — only the
2714        // engine denies it).
2715        let sid = new_sid();
2716        let payload = SessionStartPayload {
2717            intent: "e3-stream".into(),
2718            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2719            mode_version: "1.0.0".into(),
2720            configuration_version: "cfg-1".into(),
2721            policy_version: String::new(),
2722            ttl_ms: 60_000,
2723            context_id: String::new(),
2724            extensions: Default::default(),
2725            roots: vec![],
2726            max_suspend_ms: 0,
2727        }
2728        .encode_to_vec();
2729        runtime
2730            .process(
2731                &Envelope {
2732                    macp_version: "1.0".into(),
2733                    mode: "macp.mode.decision.v1".into(),
2734                    message_type: "SessionStart".into(),
2735                    message_id: new_sid(),
2736                    session_id: sid.clone(),
2737                    sender: "agent://ok".into(),
2738                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2739                    payload,
2740                },
2741                None,
2742            )
2743            .await
2744            .unwrap();
2745
2746        let embargoed = crate::security::AuthIdentity {
2747            sender: "agent://embargoed".into(),
2748            allowed_modes: None,
2749            can_start_sessions: true,
2750            max_open_sessions: None,
2751            can_manage_mode_registry: false,
2752            is_observer: false,
2753        };
2754        let mut bound = None;
2755        let mut events = None;
2756
2757        // 1. Stream envelope frame from the embargoed sender: denied.
2758        let proposal = crate::decision_pb::ProposalPayload {
2759            proposal_id: "p1".into(),
2760            option: "x".into(),
2761            rationale: "r".into(),
2762            supporting_data: vec![],
2763        }
2764        .encode_to_vec();
2765        let req = StreamSessionRequest {
2766            envelope: Some(Envelope {
2767                macp_version: "1.0".into(),
2768                mode: "macp.mode.decision.v1".into(),
2769                message_type: "Proposal".into(),
2770                message_id: new_sid(),
2771                session_id: sid.clone(),
2772                sender: "agent://embargoed".into(),
2773                timestamp_unix_ms: Utc::now().timestamp_millis(),
2774                payload: proposal,
2775            }),
2776            subscribe_session_id: String::new(),
2777            after_sequence: 0,
2778        };
2779        let err = server
2780            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2781            .await
2782            .expect_err("stream envelope from embargoed sender must be denied");
2783        // PolicyDenied maps to FailedPrecondition on the transport (same
2784        // error the unary path expresses as a POLICY_DENIED ack).
2785        assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2786        assert!(err.message().contains("PolicyDenied"), "{err:?}");
2787
2788        // 2. Passive-subscribe frame (history read) from the embargoed
2789        //    sender: denied even though membership would allow it.
2790        let req = StreamSessionRequest {
2791            envelope: None,
2792            subscribe_session_id: sid.clone(),
2793            after_sequence: 0,
2794        };
2795        let err = server
2796            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2797            .await
2798            .expect_err("stream subscribe from embargoed sender must be denied");
2799        assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2800    }
2801}