Skip to main content

macp_runtime/
server.rs

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