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