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        // Subscribed HERE, during the unary call — strictly before the generator
1373        // below is first polled, which does not happen until the client reads.
1374        // An event published in that gap must be buffered by an existing
1375        // subscription, not missed; moving this inside the generator
1376        // reintroduces exactly that race.
1377        let mut rx = self.runtime.subscribe_session_lifecycle();
1378        let runtime = Arc::clone(&self.runtime);
1379        let stream = async_stream::try_stream! {
1380            // Initial sync: emit all current sessions as CREATED events.
1381            //
1382            // The traversal snapshots the registry's shared session handles
1383            // ONCE and then locks and clones one session at a time (see
1384            // `watch_sync`), so peak resident `Session` clones is one rather
1385            // than the registry size: this generator is paced by the client's
1386            // reads and there can be `MACP_MAX_CONCURRENT_STREAMS` of them at
1387            // once.
1388            let mut sync = crate::watch_sync::InitialSync::begin(&runtime.registry).await;
1389            // The IDs this sync emits. A Created event buffered in the
1390            // subscribe→sync window would duplicate one of them — the session
1391            // was already registered when the snapshot was taken, but
1392            // `process_session_start` inserts into the registry BEFORE
1393            // publishing Created, so that event can still arrive afterwards.
1394            // Session IDs are create-once and `runtime.rs` holds the only
1395            // `Created` publisher, so one `send` per session start means a
1396            // *live* Created can never repeat for an ID already in this set —
1397            // membership is read, never extended, past the sync. That keeps the
1398            // set bounded by the registry size at subscribe time instead of
1399            // growing with every session the stream ever observes.
1400            let mut synced: std::collections::HashSet<String> =
1401                std::collections::HashSet::with_capacity(sync.remaining());
1402            // Lifecycle events that arrive while the sync is still emitting.
1403            // The sync loop cannot `recv().await` (it has its own output to
1404            // produce) but must not ignore the bus either: the bus holds 64
1405            // events, so a slow sync would otherwise make the first post-sync
1406            // `recv()` return `Lagged` and kill the stream. Bounded by
1407            // `PENDING_EVENT_LIMIT`; on overflow the client gets the same
1408            // `RESOURCE_EXHAUSTED` it gets for bus lag.
1409            let mut pending: std::collections::VecDeque<crate::runtime::SessionLifecycleEvent> =
1410                std::collections::VecDeque::new();
1411            loop {
1412                // Drained BEFORE the next session is fetched, so the bus is
1413                // relieved once per emitted session rather than once for the
1414                // whole sync.
1415                if let Err(drain_err) = crate::watch_sync::drain_lifecycle_events(
1416                    &mut rx,
1417                    &mut pending,
1418                    crate::watch_sync::PENDING_EVENT_LIMIT,
1419                ) {
1420                    Err(Status::resource_exhausted(drain_err.message()))?;
1421                    break;
1422                }
1423                // One session, emitted and dropped before the next is asked
1424                // for, which is what keeps residency bounded.
1425                let Some(session) = sync.next_session().await else { break };
1426                synced.insert(session.session_id.clone());
1427                yield WatchSessionsResponse {
1428                    event: Some(SessionLifecycleEvent {
1429                        event_type: session_lifecycle_event::EventType::Created.into(),
1430                        session: Some(Self::session_to_metadata(&session)),
1431                        observed_at_unix_ms: session.started_at_unix_ms,
1432                    }),
1433                };
1434            }
1435            // Stream lifecycle transitions: first the ones buffered during the
1436            // sync (in bus order), then live ones.
1437            loop {
1438                let event = match pending.pop_front() {
1439                    Some(event) => event,
1440                    None => match rx.recv().await {
1441                        Ok(event) => event,
1442                        Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1443                            Err(Status::resource_exhausted(format!(
1444                                "WatchSessions receiver fell behind by {skipped} events"
1445                            )))?;
1446                            break;
1447                        }
1448                        Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1449                    },
1450                };
1451                let (event_type, sid) = match &event {
1452                    crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1453                        (session_lifecycle_event::EventType::Created, session_id.clone()),
1454                    crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1455                        (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1456                    crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1457                        (session_lifecycle_event::EventType::Expired, session_id.clone()),
1458                    crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1459                        (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1460                    crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1461                        (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1462                    crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1463                        (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1464                };
1465                // Skip the buffered duplicate of an initial-sync entry;
1466                // non-Created events for synced sessions are new information
1467                // and pass through.
1468                //
1469                // `contains`, not `insert`: a live Created for a session the
1470                // sync never saw is emitted as-is and NOT recorded. Recording
1471                // it would be the only thing making this set grow with the
1472                // stream's lifetime, and it would buy nothing — `runtime.rs`
1473                // is the single `Created` publisher and sends once per session
1474                // start, so no live Created can repeat.
1475                if event_type == session_lifecycle_event::EventType::Created
1476                    && synced.contains(&sid)
1477                {
1478                    continue;
1479                }
1480                let session_meta = runtime.registry.get_session(&sid).await
1481                    .map(|s| Self::session_to_metadata(&s));
1482                yield WatchSessionsResponse {
1483                    event: Some(SessionLifecycleEvent {
1484                        event_type: event_type.into(),
1485                        session: session_meta,
1486                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1487                    }),
1488                };
1489            }
1490        };
1491        Ok(Response::new(Box::pin(stream)))
1492    }
1493
1494    // Extension mode lifecycle RPCs
1495
1496    async fn list_ext_modes(
1497        &self,
1498        _request: Request<ListExtModesRequest>,
1499    ) -> Result<Response<ListExtModesResponse>, Status> {
1500        Ok(Response::new(ListExtModesResponse {
1501            modes: self.runtime.extension_mode_descriptors(),
1502        }))
1503    }
1504
1505    async fn register_ext_mode(
1506        &self,
1507        request: Request<RegisterExtModeRequest>,
1508    ) -> Result<Response<RegisterExtModeResponse>, 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 descriptor = req
1519            .mode_descriptor
1520            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1521        match self.runtime.register_extension(descriptor) {
1522            Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1523                ok: true,
1524                error: String::new(),
1525            })),
1526            Err(e) => Ok(Response::new(RegisterExtModeResponse {
1527                ok: false,
1528                error: e,
1529            })),
1530        }
1531    }
1532
1533    async fn unregister_ext_mode(
1534        &self,
1535        request: Request<UnregisterExtModeRequest>,
1536    ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1537        let identity = self
1538            .security
1539            .authenticate_metadata(request.metadata())
1540            .await
1541            .map_err(Self::status_from_error)?;
1542        self.security
1543            .authorize_mode_registry(&identity)
1544            .map_err(Self::status_from_error)?;
1545        let req = request.into_inner();
1546        match self.runtime.unregister_extension(&req.mode) {
1547            Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1548                ok: true,
1549                error: String::new(),
1550            })),
1551            Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1552                ok: false,
1553                error: e,
1554            })),
1555        }
1556    }
1557
1558    async fn promote_mode(
1559        &self,
1560        request: Request<PromoteModeRequest>,
1561    ) -> Result<Response<PromoteModeResponse>, Status> {
1562        let identity = self
1563            .security
1564            .authenticate_metadata(request.metadata())
1565            .await
1566            .map_err(Self::status_from_error)?;
1567        self.security
1568            .authorize_mode_registry(&identity)
1569            .map_err(Self::status_from_error)?;
1570        let req = request.into_inner();
1571        let new_name = if req.promoted_mode_name.is_empty() {
1572            None
1573        } else {
1574            Some(req.promoted_mode_name.as_str())
1575        };
1576        match self.runtime.promote_mode(&req.mode, new_name) {
1577            Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1578                ok: true,
1579                error: String::new(),
1580                mode: final_name,
1581            })),
1582            Err(e) => Ok(Response::new(PromoteModeResponse {
1583                ok: false,
1584                error: e,
1585                mode: String::new(),
1586            })),
1587        }
1588    }
1589
1590    // ── Governance policy lifecycle RPCs (RFC-MACP-0012) ────────────
1591
1592    async fn register_policy(
1593        &self,
1594        request: Request<RegisterPolicyRequest>,
1595    ) -> Result<Response<RegisterPolicyResponse>, Status> {
1596        if self.policies_read_only {
1597            return Err(Status::failed_precondition(
1598                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1599            ));
1600        }
1601        let identity = self
1602            .security
1603            .authenticate_metadata(request.metadata())
1604            .await
1605            .map_err(Self::status_from_error)?;
1606        self.security
1607            .authorize_mode_registry(&identity)
1608            .map_err(Self::status_from_error)?;
1609        let req = request.into_inner();
1610        let descriptor = req
1611            .policy_descriptor
1612            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1613        let definition = Self::policy_descriptor_to_definition(&descriptor);
1614        match self.runtime.register_policy(definition) {
1615            Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1616                ok: true,
1617                error: String::new(),
1618            })),
1619            Err(e) => Ok(Response::new(RegisterPolicyResponse {
1620                ok: false,
1621                error: e,
1622            })),
1623        }
1624    }
1625
1626    async fn unregister_policy(
1627        &self,
1628        request: Request<UnregisterPolicyRequest>,
1629    ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1630        if self.policies_read_only {
1631            return Err(Status::failed_precondition(
1632                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1633            ));
1634        }
1635        let identity = self
1636            .security
1637            .authenticate_metadata(request.metadata())
1638            .await
1639            .map_err(Self::status_from_error)?;
1640        self.security
1641            .authorize_mode_registry(&identity)
1642            .map_err(Self::status_from_error)?;
1643        let req = request.into_inner();
1644        match self.runtime.unregister_policy(&req.policy_id) {
1645            Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1646                ok: true,
1647                error: String::new(),
1648            })),
1649            Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1650                ok: false,
1651                error: e,
1652            })),
1653        }
1654    }
1655
1656    async fn get_policy(
1657        &self,
1658        request: Request<GetPolicyRequest>,
1659    ) -> Result<Response<GetPolicyResponse>, Status> {
1660        let _identity = self
1661            .security
1662            .authenticate_metadata(request.metadata())
1663            .await
1664            .map_err(Self::status_from_error)?;
1665        let req = request.into_inner();
1666        let policy = self
1667            .runtime
1668            .get_policy(&req.policy_id)
1669            .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1670        Ok(Response::new(GetPolicyResponse {
1671            policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1672        }))
1673    }
1674
1675    async fn list_policies(
1676        &self,
1677        request: Request<ListPoliciesRequest>,
1678    ) -> Result<Response<ListPoliciesResponse>, Status> {
1679        let _identity = self
1680            .security
1681            .authenticate_metadata(request.metadata())
1682            .await
1683            .map_err(Self::status_from_error)?;
1684        let req = request.into_inner();
1685        let mode_filter = if req.mode.is_empty() {
1686            None
1687        } else {
1688            Some(req.mode.as_str())
1689        };
1690        let policies = self.runtime.list_policies(mode_filter);
1691        let descriptors = policies
1692            .iter()
1693            .map(Self::policy_definition_to_descriptor)
1694            .collect();
1695        Ok(Response::new(ListPoliciesResponse { descriptors }))
1696    }
1697
1698    type WatchPoliciesStream = std::pin::Pin<
1699        Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1700    >;
1701
1702    async fn watch_policies(
1703        &self,
1704        _request: Request<WatchPoliciesRequest>,
1705    ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1706        let mut rx = self.runtime.subscribe_policy_changes();
1707        let runtime = Arc::clone(&self.runtime);
1708        let stream = async_stream::try_stream! {
1709            // Send initial state
1710            let policies = runtime.list_policies(None);
1711            let descriptors: Vec<PolicyDescriptor> = policies
1712                .iter()
1713                .map(MacpServer::policy_definition_to_descriptor)
1714                .collect();
1715            yield WatchPoliciesResponse {
1716                descriptors,
1717                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1718            };
1719            // Wait for changes
1720            while rx.recv().await.is_ok() {
1721                let policies = runtime.list_policies(None);
1722                let descriptors: Vec<PolicyDescriptor> = policies
1723                    .iter()
1724                    .map(MacpServer::policy_definition_to_descriptor)
1725                    .collect();
1726                yield WatchPoliciesResponse {
1727                    descriptors,
1728                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1729                };
1730            }
1731        };
1732        Ok(Response::new(Box::pin(stream)))
1733    }
1734}
1735
1736// ── Policy type conversion helpers ──────────────────────────────────
1737
1738impl MacpServer {
1739    fn policy_descriptor_to_definition(
1740        descriptor: &PolicyDescriptor,
1741    ) -> crate::policy::PolicyDefinition {
1742        let rules: serde_json::Value = if descriptor.rules.is_empty() {
1743            serde_json::json!({})
1744        } else {
1745            serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1746        };
1747        crate::policy::PolicyDefinition {
1748            policy_id: descriptor.policy_id.clone(),
1749            mode: descriptor.mode.clone(),
1750            description: descriptor.description.clone(),
1751            rules,
1752            schema_version: descriptor.schema_version,
1753        }
1754    }
1755
1756    fn policy_definition_to_descriptor(
1757        definition: &crate::policy::PolicyDefinition,
1758    ) -> PolicyDescriptor {
1759        PolicyDescriptor {
1760            policy_id: definition.policy_id.clone(),
1761            mode: definition.mode.clone(),
1762            description: definition.description.clone(),
1763            rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1764            schema_version: definition.schema_version,
1765            registered_at_unix_ms: 0,
1766        }
1767    }
1768}
1769
1770#[cfg(test)]
1771mod tests {
1772    use super::*;
1773    use crate::log_store::LogStore;
1774    use crate::pb::SessionStartPayload;
1775    use crate::registry::SessionRegistry;
1776    use chrono::Utc;
1777    use prost::Message;
1778
1779    fn new_sid() -> String {
1780        uuid::Uuid::new_v4().as_hyphenated().to_string()
1781    }
1782
1783    fn make_server() -> (MacpServer, Arc<Runtime>) {
1784        make_server_with_security(SecurityLayer::dev_mode())
1785    }
1786
1787    /// Same harness, with the `SecurityLayer` supplied by the caller so a test
1788    /// can pin `list_sessions_{default,max}_page_size` without touching process
1789    /// env (which is not deterministic under `cargo test`'s thread pool).
1790    fn make_server_with_security(security: SecurityLayer) -> (MacpServer, Arc<Runtime>) {
1791        let storage: Arc<dyn crate::storage::StorageBackend> =
1792            Arc::new(crate::storage::MemoryBackend);
1793        let registry = Arc::new(SessionRegistry::new());
1794        let log_store = Arc::new(LogStore::new());
1795        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1796        let server = MacpServer::new(runtime.clone(), security);
1797        (server, runtime)
1798    }
1799
1800    fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1801        let mut req = Request::new(SendRequest {
1802            envelope: Some(env),
1803        });
1804        req.metadata_mut()
1805            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1806        req
1807    }
1808
1809    async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1810        let resp = server.send(send_req(sender, env)).await.unwrap();
1811        resp.into_inner().ack.unwrap()
1812    }
1813
1814    fn start_payload() -> Vec<u8> {
1815        SessionStartPayload {
1816            intent: "intent".into(),
1817            participants: vec!["agent://fraud".into()],
1818            mode_version: "1.0.0".into(),
1819            configuration_version: "cfg-1".into(),
1820            policy_version: String::new(),
1821            ttl_ms: 1000,
1822            context_id: String::new(),
1823            extensions: std::collections::HashMap::new(),
1824            roots: vec![],
1825            max_suspend_ms: 0,
1826        }
1827        .encode_to_vec()
1828    }
1829
1830    #[tokio::test]
1831    async fn sender_is_derived_from_authenticated_metadata() {
1832        let (server, runtime) = make_server();
1833        let sid = new_sid();
1834        let ack = do_send(
1835            &server,
1836            "agent://orchestrator",
1837            Envelope {
1838                macp_version: "1.0".into(),
1839                mode: "macp.mode.decision.v1".into(),
1840                message_type: "SessionStart".into(),
1841                message_id: "m1".into(),
1842                session_id: sid.clone(),
1843                sender: String::new(),
1844                timestamp_unix_ms: Utc::now().timestamp_millis(),
1845                payload: start_payload(),
1846            },
1847        )
1848        .await;
1849        assert!(ack.ok);
1850        let session = runtime.get_session_checked(&sid).await.unwrap();
1851        assert_eq!(session.initiator_sender, "agent://orchestrator");
1852    }
1853
1854    #[tokio::test]
1855    async fn spoofed_sender_is_rejected() {
1856        let (server, _) = make_server();
1857        let sid = new_sid();
1858        let ack = do_send(
1859            &server,
1860            "agent://orchestrator",
1861            Envelope {
1862                macp_version: "1.0".into(),
1863                mode: "macp.mode.decision.v1".into(),
1864                message_type: "SessionStart".into(),
1865                message_id: "m1".into(),
1866                session_id: sid,
1867                sender: "agent://spoof".into(),
1868                timestamp_unix_ms: Utc::now().timestamp_millis(),
1869                payload: start_payload(),
1870            },
1871        )
1872        .await;
1873        assert!(!ack.ok);
1874        assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1875    }
1876
1877    #[tokio::test]
1878    async fn get_session_requires_session_membership() {
1879        let (server, _) = make_server();
1880        let sid = new_sid();
1881        let ack = do_send(
1882            &server,
1883            "agent://orchestrator",
1884            Envelope {
1885                macp_version: "1.0".into(),
1886                mode: "macp.mode.decision.v1".into(),
1887                message_type: "SessionStart".into(),
1888                message_id: "m1".into(),
1889                session_id: sid.clone(),
1890                sender: String::new(),
1891                timestamp_unix_ms: Utc::now().timestamp_millis(),
1892                payload: start_payload(),
1893            },
1894        )
1895        .await;
1896        assert!(ack.ok);
1897
1898        let mut req = Request::new(GetSessionRequest { session_id: sid });
1899        req.metadata_mut().insert(
1900            "authorization",
1901            format!("Bearer {}", "agent://outsider").parse().unwrap(),
1902        );
1903        let err = server.get_session(req).await.unwrap_err();
1904        assert_eq!(err.code(), tonic::Code::PermissionDenied);
1905    }
1906
1907    #[tokio::test]
1908    async fn register_ext_mode_requires_authenticated_registry_permission() {
1909        let storage: Arc<dyn crate::storage::StorageBackend> =
1910            Arc::new(crate::storage::MemoryBackend);
1911        let registry = Arc::new(SessionRegistry::new());
1912        let log_store = Arc::new(LogStore::new());
1913        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1914        let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1915        let server = MacpServer::new(runtime, security);
1916
1917        let req = Request::new(RegisterExtModeRequest {
1918            mode_descriptor: Some(crate::pb::ModeDescriptor {
1919                mode: "ext.custom.v1".into(),
1920                mode_version: "1.0.0".into(),
1921                message_types: vec!["SessionStart".into(), "Commitment".into()],
1922                ..Default::default()
1923            }),
1924        });
1925        let err = server.register_ext_mode(req).await.unwrap_err();
1926        assert_eq!(err.code(), tonic::Code::Unauthenticated);
1927    }
1928
1929    fn stream_identity(sender: &str) -> AuthIdentity {
1930        AuthIdentity {
1931            sender: sender.into(),
1932            allowed_modes: None,
1933            can_start_sessions: true,
1934            max_open_sessions: None,
1935            can_manage_mode_registry: false,
1936            is_observer: false,
1937        }
1938    }
1939
1940    #[tokio::test]
1941    async fn stream_session_emits_accepted_envelopes_only() {
1942        use tokio_stream::{iter, StreamExt};
1943
1944        let (server, _) = make_server();
1945        let sid = new_sid();
1946        let requests = iter(vec![Ok(StreamSessionRequest {
1947            subscribe_session_id: String::new(),
1948            after_sequence: 0,
1949            envelope: Some(Envelope {
1950                macp_version: "1.0".into(),
1951                mode: "macp.mode.decision.v1".into(),
1952                message_type: "SessionStart".into(),
1953                message_id: "m1".into(),
1954                session_id: sid.clone(),
1955                sender: String::new(),
1956                timestamp_unix_ms: Utc::now().timestamp_millis(),
1957                payload: start_payload(),
1958            }),
1959        })]);
1960
1961        let mut stream =
1962            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1963
1964        let response = stream.next().await.unwrap().unwrap();
1965        let envelope = match response.response.unwrap() {
1966            crate::pb::stream_session_response::Response::Envelope(e) => e,
1967            _ => panic!("expected envelope"),
1968        };
1969        assert_eq!(envelope.message_type, "SessionStart");
1970        assert_eq!(envelope.message_id, "m1");
1971        assert!(stream.next().await.is_none());
1972    }
1973
1974    #[tokio::test]
1975    async fn stream_session_rejects_mixed_session_ids() {
1976        use tokio_stream::{iter, StreamExt};
1977
1978        let (server, _) = make_server();
1979        let sid1 = new_sid();
1980        let sid2 = new_sid();
1981        let requests = iter(vec![
1982            Ok(StreamSessionRequest {
1983                subscribe_session_id: String::new(),
1984                after_sequence: 0,
1985                envelope: Some(Envelope {
1986                    macp_version: "1.0".into(),
1987                    mode: "macp.mode.decision.v1".into(),
1988                    message_type: "SessionStart".into(),
1989                    message_id: "m1".into(),
1990                    session_id: sid1.clone(),
1991                    sender: String::new(),
1992                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1993                    payload: start_payload(),
1994                }),
1995            }),
1996            Ok(StreamSessionRequest {
1997                subscribe_session_id: String::new(),
1998                after_sequence: 0,
1999                envelope: Some(Envelope {
2000                    macp_version: "1.0".into(),
2001                    mode: "macp.mode.decision.v1".into(),
2002                    message_type: "SessionStart".into(),
2003                    message_id: "m2".into(),
2004                    session_id: sid2,
2005                    sender: String::new(),
2006                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2007                    payload: start_payload(),
2008                }),
2009            }),
2010        ]);
2011
2012        let mut stream =
2013            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
2014
2015        let first = stream.next().await.unwrap().unwrap();
2016        let first_env = match first.response.unwrap() {
2017            crate::pb::stream_session_response::Response::Envelope(e) => e,
2018            _ => panic!("expected envelope"),
2019        };
2020        assert_eq!(first_env.session_id, sid1);
2021        let err = stream.next().await.unwrap().unwrap_err();
2022        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2023    }
2024
2025    #[tokio::test]
2026    async fn list_modes_returns_standard_modes() {
2027        let (server, _) = make_server();
2028        let resp = server
2029            .list_modes(Request::new(ListModesRequest {}))
2030            .await
2031            .unwrap();
2032        let names: Vec<String> = resp
2033            .into_inner()
2034            .modes
2035            .iter()
2036            .map(|m| m.mode.clone())
2037            .collect();
2038        assert_eq!(names.len(), 5);
2039        assert!(names.contains(&"macp.mode.decision.v1".to_string()));
2040        assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
2041        assert!(names.contains(&"macp.mode.task.v1".to_string()));
2042        assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
2043        assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
2044        // multi_round is now an extension, not in ListModes
2045        assert!(!names.contains(&"ext.multi_round.v1".to_string()));
2046    }
2047
2048    #[tokio::test]
2049    async fn list_ext_modes_returns_extensions() {
2050        let (server, _) = make_server();
2051        let resp = server
2052            .list_ext_modes(Request::new(ListExtModesRequest {}))
2053            .await
2054            .unwrap();
2055        let names: Vec<String> = resp
2056            .into_inner()
2057            .modes
2058            .iter()
2059            .map(|m| m.mode.clone())
2060            .collect();
2061        assert_eq!(names.len(), 1);
2062        assert!(names.contains(&"ext.multi_round.v1".to_string()));
2063    }
2064
2065    #[tokio::test]
2066    async fn get_manifest_includes_all_modes() {
2067        let (server, _) = make_server();
2068        let resp = server
2069            .get_manifest(Request::new(crate::pb::GetManifestRequest {
2070                agent_id: String::new(),
2071            }))
2072            .await
2073            .unwrap();
2074        let manifest = resp.into_inner().manifest.unwrap();
2075        assert_eq!(manifest.supported_modes.len(), 6);
2076        assert!(manifest
2077            .supported_modes
2078            .contains(&"ext.multi_round.v1".to_string()));
2079    }
2080
2081    #[tokio::test]
2082    async fn get_session_returns_metadata() {
2083        let (server, _) = make_server();
2084        let sid = new_sid();
2085        let ack = do_send(
2086            &server,
2087            "agent://orchestrator",
2088            Envelope {
2089                macp_version: "1.0".into(),
2090                mode: "macp.mode.decision.v1".into(),
2091                message_type: "SessionStart".into(),
2092                message_id: "m1".into(),
2093                session_id: sid.clone(),
2094                sender: String::new(),
2095                timestamp_unix_ms: Utc::now().timestamp_millis(),
2096                payload: start_payload(),
2097            },
2098        )
2099        .await;
2100        assert!(ack.ok);
2101
2102        let mut req = Request::new(GetSessionRequest {
2103            session_id: sid.clone(),
2104        });
2105        req.metadata_mut().insert(
2106            "authorization",
2107            format!("Bearer {}", "agent://orchestrator")
2108                .parse()
2109                .unwrap(),
2110        );
2111        let resp = server.get_session(req).await.unwrap();
2112        let meta = resp.into_inner().metadata.unwrap();
2113        assert_eq!(meta.session_id, sid);
2114        assert_eq!(meta.mode, "macp.mode.decision.v1");
2115        assert_eq!(meta.mode_version, "1.0.0");
2116        assert_eq!(meta.configuration_version, "cfg-1");
2117    }
2118
2119    #[tokio::test]
2120    async fn cancel_session_transitions_to_cancelled() {
2121        let (server, _) = make_server();
2122        let sid = new_sid();
2123        let ack = do_send(
2124            &server,
2125            "agent://orchestrator",
2126            Envelope {
2127                macp_version: "1.0".into(),
2128                mode: "macp.mode.decision.v1".into(),
2129                message_type: "SessionStart".into(),
2130                message_id: "m1".into(),
2131                session_id: sid.clone(),
2132                sender: String::new(),
2133                timestamp_unix_ms: Utc::now().timestamp_millis(),
2134                payload: start_payload(),
2135            },
2136        )
2137        .await;
2138        assert!(ack.ok);
2139
2140        let mut req = Request::new(CancelSessionRequest {
2141            session_id: sid,
2142            reason: "no longer needed".into(),
2143        });
2144        req.metadata_mut().insert(
2145            "authorization",
2146            format!("Bearer {}", "agent://orchestrator")
2147                .parse()
2148                .unwrap(),
2149        );
2150        let resp = server.cancel_session(req).await.unwrap();
2151        let ack = resp.into_inner().ack.unwrap();
2152        assert!(ack.ok);
2153        // RFC-MACP-0001 §7.3: cancellation now yields the distinct CANCELLED state.
2154        assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2155    }
2156
2157    #[tokio::test]
2158    async fn participant_cannot_cancel_session() {
2159        let (server, _) = make_server();
2160        let sid = new_sid();
2161        let ack = do_send(
2162            &server,
2163            "agent://orchestrator",
2164            Envelope {
2165                macp_version: "1.0".into(),
2166                mode: "macp.mode.decision.v1".into(),
2167                message_type: "SessionStart".into(),
2168                message_id: "m1".into(),
2169                session_id: sid.clone(),
2170                sender: String::new(),
2171                timestamp_unix_ms: Utc::now().timestamp_millis(),
2172                payload: start_payload(),
2173            },
2174        )
2175        .await;
2176        assert!(ack.ok);
2177
2178        let mut req = Request::new(CancelSessionRequest {
2179            session_id: sid,
2180            reason: "I want to cancel".into(),
2181        });
2182        req.metadata_mut().insert(
2183            "authorization",
2184            format!("Bearer {}", "agent://fraud").parse().unwrap(),
2185        );
2186        let err = server.cancel_session(req).await.unwrap_err();
2187        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2188    }
2189
2190    #[tokio::test]
2191    async fn cancel_session_unknown_session_returns_error() {
2192        let (server, _) = make_server();
2193        let mut req = Request::new(CancelSessionRequest {
2194            session_id: "nonexistent".into(),
2195            reason: "test".into(),
2196        });
2197        req.metadata_mut().insert(
2198            "authorization",
2199            format!("Bearer {}", "agent://orchestrator")
2200                .parse()
2201                .unwrap(),
2202        );
2203        let err = server.cancel_session(req).await.unwrap_err();
2204        assert_eq!(err.code(), tonic::Code::NotFound);
2205    }
2206
2207    #[tokio::test]
2208    async fn ambient_signal_accepted() {
2209        let (server, _) = make_server();
2210        let ack = do_send(
2211            &server,
2212            "agent://orchestrator",
2213            Envelope {
2214                macp_version: "1.0".into(),
2215                mode: String::new(),
2216                message_type: "Signal".into(),
2217                message_id: "sig-1".into(),
2218                session_id: String::new(),
2219                sender: String::new(),
2220                timestamp_unix_ms: Utc::now().timestamp_millis(),
2221                payload: vec![],
2222            },
2223        )
2224        .await;
2225        assert!(ack.ok);
2226    }
2227
2228    #[tokio::test]
2229    async fn signal_with_session_id_rejected() {
2230        let (server, _) = make_server();
2231        let ack = do_send(
2232            &server,
2233            "agent://orchestrator",
2234            Envelope {
2235                macp_version: "1.0".into(),
2236                mode: String::new(),
2237                message_type: "Signal".into(),
2238                message_id: "sig-2".into(),
2239                session_id: "some-session".into(),
2240                sender: String::new(),
2241                timestamp_unix_ms: Utc::now().timestamp_millis(),
2242                payload: vec![],
2243            },
2244        )
2245        .await;
2246        assert!(!ack.ok);
2247        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2248    }
2249
2250    #[tokio::test]
2251    async fn signal_with_mode_rejected() {
2252        let (server, _) = make_server();
2253        let ack = do_send(
2254            &server,
2255            "agent://orchestrator",
2256            Envelope {
2257                macp_version: "1.0".into(),
2258                mode: "macp.mode.decision.v1".into(),
2259                message_type: "Signal".into(),
2260                message_id: "sig-3".into(),
2261                session_id: String::new(),
2262                sender: String::new(),
2263                timestamp_unix_ms: Utc::now().timestamp_millis(),
2264                payload: vec![],
2265            },
2266        )
2267        .await;
2268        assert!(!ack.ok);
2269        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2270    }
2271
2272    #[tokio::test]
2273    async fn ambient_progress_accepted() {
2274        let (server, _) = make_server();
2275        let ack = do_send(
2276            &server,
2277            "agent://orchestrator",
2278            Envelope {
2279                macp_version: "1.0".into(),
2280                mode: String::new(),
2281                message_type: "Progress".into(),
2282                message_id: "prog-1".into(),
2283                session_id: String::new(),
2284                sender: String::new(),
2285                timestamp_unix_ms: Utc::now().timestamp_millis(),
2286                payload: vec![],
2287            },
2288        )
2289        .await;
2290        assert!(ack.ok);
2291    }
2292
2293    #[tokio::test]
2294    async fn ambient_progress_with_mode_rejected() {
2295        let (server, _) = make_server();
2296        let ack = do_send(
2297            &server,
2298            "agent://orchestrator",
2299            Envelope {
2300                macp_version: "1.0".into(),
2301                mode: "macp.mode.decision.v1".into(),
2302                message_type: "Progress".into(),
2303                message_id: "prog-2".into(),
2304                session_id: String::new(),
2305                sender: String::new(),
2306                timestamp_unix_ms: Utc::now().timestamp_millis(),
2307                payload: vec![],
2308            },
2309        )
2310        .await;
2311        assert!(!ack.ok);
2312        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2313    }
2314
2315    #[tokio::test]
2316    async fn manifest_advertises_stream_enabled() {
2317        let (server, _) = make_server();
2318        let resp = server
2319            .initialize(Request::new(InitializeRequest {
2320                supported_protocol_versions: vec!["1.0".into()],
2321                client_info: None,
2322                capabilities: None,
2323            }))
2324            .await
2325            .unwrap();
2326        let caps = resp.into_inner().capabilities.unwrap();
2327        assert!(caps.sessions.unwrap().stream);
2328    }
2329
2330    #[tokio::test]
2331    async fn initialize_empty_versions_rejected() {
2332        let (server, _) = make_server();
2333        let err = server
2334            .initialize(Request::new(InitializeRequest {
2335                supported_protocol_versions: vec![],
2336                client_info: None,
2337                capabilities: None,
2338            }))
2339            .await
2340            .unwrap_err();
2341        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2342    }
2343
2344    #[tokio::test]
2345    async fn initialize_unsupported_version_rejected() {
2346        let (server, _) = make_server();
2347        let err = server
2348            .initialize(Request::new(InitializeRequest {
2349                supported_protocol_versions: vec!["2.0".into()],
2350                client_info: None,
2351                capabilities: None,
2352            }))
2353            .await
2354            .unwrap_err();
2355        assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2356    }
2357
2358    // ── RFC-MACP-0006-A1: passive subscribe tests ──────────────────────
2359
2360    fn observer_identity(sender: &str) -> AuthIdentity {
2361        AuthIdentity {
2362            sender: sender.into(),
2363            allowed_modes: None,
2364            can_start_sessions: false,
2365            max_open_sessions: None,
2366            can_manage_mode_registry: false,
2367            is_observer: true,
2368        }
2369    }
2370
2371    fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2372        StreamSessionRequest {
2373            subscribe_session_id: session_id.into(),
2374            after_sequence: after,
2375            envelope: None,
2376        }
2377    }
2378
2379    fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2380        SessionStartPayload {
2381            intent: "intent".into(),
2382            participants,
2383            mode_version: "1.0.0".into(),
2384            configuration_version: "cfg-1".into(),
2385            policy_version: String::new(),
2386            ttl_ms: 60_000,
2387            context_id: String::new(),
2388            extensions: std::collections::HashMap::new(),
2389            roots: vec![],
2390            max_suspend_ms: 0,
2391        }
2392        .encode_to_vec()
2393    }
2394
2395    async fn start_session(
2396        server: &MacpServer,
2397        initiator: &str,
2398        sid: &str,
2399        participants: Vec<String>,
2400    ) {
2401        let ack = do_send(
2402            server,
2403            initiator,
2404            Envelope {
2405                macp_version: "1.0".into(),
2406                mode: "macp.mode.decision.v1".into(),
2407                message_type: "SessionStart".into(),
2408                message_id: "start".into(),
2409                session_id: sid.into(),
2410                sender: String::new(),
2411                timestamp_unix_ms: Utc::now().timestamp_millis(),
2412                payload: start_multi_participant(participants),
2413            },
2414        )
2415        .await;
2416        assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2417    }
2418
2419    async fn send_proposal(
2420        server: &MacpServer,
2421        sender: &str,
2422        sid: &str,
2423        message_id: &str,
2424        proposal_id: &str,
2425    ) {
2426        let payload = crate::decision_pb::ProposalPayload {
2427            proposal_id: proposal_id.into(),
2428            option: "opt".into(),
2429            rationale: "r".into(),
2430            supporting_data: vec![],
2431        }
2432        .encode_to_vec();
2433        let ack = do_send(
2434            server,
2435            sender,
2436            Envelope {
2437                macp_version: "1.0".into(),
2438                mode: "macp.mode.decision.v1".into(),
2439                message_type: "Proposal".into(),
2440                message_id: message_id.into(),
2441                session_id: sid.into(),
2442                sender: String::new(),
2443                timestamp_unix_ms: Utc::now().timestamp_millis(),
2444                payload,
2445            },
2446        )
2447        .await;
2448        assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2449    }
2450
2451    #[tokio::test]
2452    async fn subscribe_replays_session_history_from_zero() {
2453        let (server, _) = make_server();
2454        let sid = new_sid();
2455        let initiator = "agent://orchestrator";
2456        let peer = "agent://fraud";
2457        start_session(
2458            &server,
2459            initiator,
2460            &sid,
2461            vec![initiator.into(), peer.into()],
2462        )
2463        .await;
2464        send_proposal(&server, peer, &sid, "m2", "p1").await;
2465
2466        let mut bound = None;
2467        let mut events = None;
2468        let replay = server
2469            .process_stream_request(
2470                &stream_identity(peer),
2471                subscribe_frame(&sid, 0),
2472                &mut bound,
2473                &mut events,
2474            )
2475            .await
2476            .unwrap();
2477
2478        assert_eq!(replay.len(), 2);
2479        assert_eq!(replay[0].message_type, "SessionStart");
2480        assert_eq!(replay[0].message_id, "start");
2481        assert_eq!(replay[1].message_type, "Proposal");
2482        assert_eq!(replay[1].message_id, "m2");
2483        assert_eq!(bound.as_deref(), Some(sid.as_str()));
2484        assert!(events.is_some());
2485    }
2486
2487    #[tokio::test]
2488    async fn subscribe_after_sequence_filters_history() {
2489        let (server, _) = make_server();
2490        let sid = new_sid();
2491        let initiator = "agent://orchestrator";
2492        let peer = "agent://fraud";
2493        start_session(
2494            &server,
2495            initiator,
2496            &sid,
2497            vec![initiator.into(), peer.into()],
2498        )
2499        .await;
2500        send_proposal(&server, peer, &sid, "m2", "p1").await;
2501        send_proposal(&server, peer, &sid, "m3", "p2").await;
2502
2503        let mut bound = None;
2504        let mut events = None;
2505        let replay = server
2506            .process_stream_request(
2507                &stream_identity(peer),
2508                subscribe_frame(&sid, 2),
2509                &mut bound,
2510                &mut events,
2511            )
2512            .await
2513            .unwrap();
2514
2515        assert_eq!(replay.len(), 1);
2516        assert_eq!(replay[0].message_id, "m3");
2517    }
2518
2519    #[tokio::test]
2520    async fn subscribe_unknown_session_returns_not_found() {
2521        let (server, _) = make_server();
2522        let mut bound = None;
2523        let mut events = None;
2524        let status = server
2525            .process_stream_request(
2526                &stream_identity("agent://orchestrator"),
2527                subscribe_frame("missing-session", 0),
2528                &mut bound,
2529                &mut events,
2530            )
2531            .await
2532            .unwrap_err();
2533        assert_eq!(status.code(), tonic::Code::NotFound);
2534        assert!(bound.is_none());
2535        assert!(events.is_none());
2536    }
2537
2538    #[tokio::test]
2539    async fn subscribe_non_participant_is_forbidden() {
2540        let (server, _) = make_server();
2541        let sid = new_sid();
2542        start_session(
2543            &server,
2544            "agent://orchestrator",
2545            &sid,
2546            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2547        )
2548        .await;
2549
2550        let mut bound = None;
2551        let mut events = None;
2552        let status = server
2553            .process_stream_request(
2554                &stream_identity("agent://outsider"),
2555                subscribe_frame(&sid, 0),
2556                &mut bound,
2557                &mut events,
2558            )
2559            .await
2560            .unwrap_err();
2561        assert_eq!(status.code(), tonic::Code::PermissionDenied);
2562    }
2563
2564    #[tokio::test]
2565    async fn subscribe_observer_identity_allowed() {
2566        let (server, _) = make_server();
2567        let sid = new_sid();
2568        start_session(
2569            &server,
2570            "agent://orchestrator",
2571            &sid,
2572            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2573        )
2574        .await;
2575
2576        let mut bound = None;
2577        let mut events = None;
2578        let replay = server
2579            .process_stream_request(
2580                &observer_identity("agent://auditor"),
2581                subscribe_frame(&sid, 0),
2582                &mut bound,
2583                &mut events,
2584            )
2585            .await
2586            .unwrap();
2587        assert_eq!(replay.len(), 1);
2588        assert_eq!(replay[0].message_type, "SessionStart");
2589    }
2590
2591    #[tokio::test]
2592    async fn subscribe_initiator_allowed_even_when_not_listed() {
2593        // Per RFC-MACP-0007, the initiator is always authorized for session
2594        // access, even if not present in the participants list.
2595        let (server, _) = make_server();
2596        let sid = new_sid();
2597        start_session(
2598            &server,
2599            "agent://orchestrator",
2600            &sid,
2601            vec!["agent://fraud".into()],
2602        )
2603        .await;
2604
2605        let mut bound = None;
2606        let mut events = None;
2607        let replay = server
2608            .process_stream_request(
2609                &stream_identity("agent://orchestrator"),
2610                subscribe_frame(&sid, 0),
2611                &mut bound,
2612                &mut events,
2613            )
2614            .await
2615            .unwrap();
2616        assert_eq!(replay.len(), 1);
2617    }
2618
2619    #[tokio::test]
2620    async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2621        let (server, _) = make_server();
2622        let sid = new_sid();
2623        let req = StreamSessionRequest {
2624            subscribe_session_id: sid.clone(),
2625            after_sequence: 0,
2626            envelope: Some(Envelope {
2627                macp_version: "1.0".into(),
2628                mode: "macp.mode.decision.v1".into(),
2629                message_type: "SessionStart".into(),
2630                message_id: "m1".into(),
2631                session_id: sid,
2632                sender: String::new(),
2633                timestamp_unix_ms: Utc::now().timestamp_millis(),
2634                payload: start_payload(),
2635            }),
2636        };
2637
2638        let mut bound = None;
2639        let mut events = None;
2640        let status = server
2641            .process_stream_request(
2642                &stream_identity("agent://orchestrator"),
2643                req,
2644                &mut bound,
2645                &mut events,
2646            )
2647            .await
2648            .unwrap_err();
2649        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2650    }
2651
2652    #[tokio::test]
2653    async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2654        let (server, _) = make_server();
2655        let sid1 = new_sid();
2656        let sid2 = new_sid();
2657        start_session(
2658            &server,
2659            "agent://orchestrator",
2660            &sid1,
2661            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2662        )
2663        .await;
2664        start_session(
2665            &server,
2666            "agent://orchestrator",
2667            &sid2,
2668            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2669        )
2670        .await;
2671
2672        // First subscribe binds the stream to sid1
2673        let identity = stream_identity("agent://fraud");
2674        let mut bound = None;
2675        let mut events = None;
2676        server
2677            .process_stream_request(
2678                &identity,
2679                subscribe_frame(&sid1, 0),
2680                &mut bound,
2681                &mut events,
2682            )
2683            .await
2684            .unwrap();
2685        assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2686
2687        // Second subscribe to sid2 on the same stream must be rejected
2688        let status = server
2689            .process_stream_request(
2690                &identity,
2691                subscribe_frame(&sid2, 0),
2692                &mut bound,
2693                &mut events,
2694            )
2695            .await
2696            .unwrap_err();
2697        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2698    }
2699
2700    /// E3: an injected ingress engine gates session start, messages, and
2701    /// session reads — deny-one-sender double proves all three hooks fire and
2702    /// that denial surfaces as POLICY_DENIED / PermissionDenied (fail closed).
2703    struct DenySenderEngine {
2704        denied: String,
2705    }
2706
2707    #[async_trait::async_trait]
2708    impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2709        async fn evaluate_session_start(
2710            &self,
2711            identity: &crate::security::AuthIdentity,
2712            _mode: &str,
2713            _env: &Envelope,
2714        ) -> macp_core::policy::PolicyDecision {
2715            if identity.sender == self.denied {
2716                macp_core::policy::PolicyDecision::Deny {
2717                    reasons: vec!["sender embargoed".into()],
2718                }
2719            } else {
2720                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2721            }
2722        }
2723
2724        async fn evaluate_message(
2725            &self,
2726            identity: &crate::security::AuthIdentity,
2727            _session: &macp_core::session::Session,
2728            _env: &Envelope,
2729        ) -> macp_core::policy::PolicyDecision {
2730            if identity.sender == self.denied {
2731                macp_core::policy::PolicyDecision::Deny {
2732                    reasons: vec!["sender embargoed".into()],
2733                }
2734            } else {
2735                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2736            }
2737        }
2738
2739        async fn evaluate_session_access(
2740            &self,
2741            identity: &crate::security::AuthIdentity,
2742            _session: &macp_core::session::Session,
2743        ) -> macp_core::policy::PolicyDecision {
2744            if identity.sender == self.denied {
2745                macp_core::policy::PolicyDecision::Deny {
2746                    reasons: vec!["sender embargoed".into()],
2747                }
2748            } else {
2749                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2750            }
2751        }
2752    }
2753
2754    #[tokio::test]
2755    async fn policy_engine_gates_all_three_ingress_points() {
2756        let (server, _runtime) = make_server();
2757        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2758            denied: "agent://embargoed".into(),
2759        }));
2760
2761        let sid = new_sid();
2762        let start_payload = SessionStartPayload {
2763            intent: "e3".into(),
2764            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2765            mode_version: "1.0.0".into(),
2766            configuration_version: "cfg-1".into(),
2767            policy_version: String::new(),
2768            ttl_ms: 60_000,
2769            context_id: String::new(),
2770            extensions: Default::default(),
2771            roots: vec![],
2772            max_suspend_ms: 0,
2773        }
2774        .encode_to_vec();
2775        let start_env = |sender: &str, sid: &str| Envelope {
2776            macp_version: "1.0".into(),
2777            mode: "macp.mode.decision.v1".into(),
2778            message_type: "SessionStart".into(),
2779            message_id: new_sid(),
2780            session_id: sid.into(),
2781            sender: sender.into(),
2782            timestamp_unix_ms: Utc::now().timestamp_millis(),
2783            payload: start_payload.clone(),
2784        };
2785
2786        // 1. Embargoed sender cannot start a session.
2787        let ack = server
2788            .send(send_req(
2789                "agent://embargoed",
2790                start_env("agent://embargoed", &sid),
2791            ))
2792            .await
2793            .unwrap()
2794            .into_inner()
2795            .ack
2796            .unwrap();
2797        assert!(!ack.ok);
2798        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2799
2800        // Allowed sender starts it.
2801        let ack = server
2802            .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2803            .await
2804            .unwrap()
2805            .into_inner()
2806            .ack
2807            .unwrap();
2808        assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2809
2810        // 2. Embargoed sender cannot send into the session.
2811        let proposal = crate::decision_pb::ProposalPayload {
2812            proposal_id: "p1".into(),
2813            option: "x".into(),
2814            rationale: "r".into(),
2815            supporting_data: vec![],
2816        }
2817        .encode_to_vec();
2818        let msg_env = Envelope {
2819            macp_version: "1.0".into(),
2820            mode: "macp.mode.decision.v1".into(),
2821            message_type: "Proposal".into(),
2822            message_id: new_sid(),
2823            session_id: sid.clone(),
2824            sender: "agent://embargoed".into(),
2825            timestamp_unix_ms: Utc::now().timestamp_millis(),
2826            payload: proposal,
2827        };
2828        let ack = server
2829            .send(send_req("agent://embargoed", msg_env))
2830            .await
2831            .unwrap()
2832            .into_inner()
2833            .ack
2834            .unwrap();
2835        assert!(!ack.ok);
2836        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2837
2838        // 3. Embargoed sender cannot read the session.
2839        let mut req = Request::new(crate::pb::GetSessionRequest {
2840            session_id: sid.clone(),
2841        });
2842        req.metadata_mut()
2843            .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2844        let err = server
2845            .get_session(req)
2846            .await
2847            .expect_err("embargoed read must be denied");
2848        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2849    }
2850
2851    /// E3 transport-parity: the ingress engine gates the STREAM path too — a
2852    /// denied sender must not be able to bypass the engine by switching from
2853    /// unary Send to StreamSession (envelope frames or subscribe frames).
2854    #[tokio::test]
2855    async fn policy_engine_gates_stream_path() {
2856        let (server, runtime) = make_server();
2857        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2858            denied: "agent://embargoed".into(),
2859        }));
2860
2861        // Session started by an allowed sender (participants include the
2862        // embargoed agent so built-in membership checks pass — only the
2863        // engine denies it).
2864        let sid = new_sid();
2865        let payload = SessionStartPayload {
2866            intent: "e3-stream".into(),
2867            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2868            mode_version: "1.0.0".into(),
2869            configuration_version: "cfg-1".into(),
2870            policy_version: String::new(),
2871            ttl_ms: 60_000,
2872            context_id: String::new(),
2873            extensions: Default::default(),
2874            roots: vec![],
2875            max_suspend_ms: 0,
2876        }
2877        .encode_to_vec();
2878        runtime
2879            .process(
2880                &Envelope {
2881                    macp_version: "1.0".into(),
2882                    mode: "macp.mode.decision.v1".into(),
2883                    message_type: "SessionStart".into(),
2884                    message_id: new_sid(),
2885                    session_id: sid.clone(),
2886                    sender: "agent://ok".into(),
2887                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2888                    payload,
2889                },
2890                None,
2891            )
2892            .await
2893            .unwrap();
2894
2895        let embargoed = crate::security::AuthIdentity {
2896            sender: "agent://embargoed".into(),
2897            allowed_modes: None,
2898            can_start_sessions: true,
2899            max_open_sessions: None,
2900            can_manage_mode_registry: false,
2901            is_observer: false,
2902        };
2903        let mut bound = None;
2904        let mut events = None;
2905
2906        // 1. Stream envelope frame from the embargoed sender: denied.
2907        let proposal = crate::decision_pb::ProposalPayload {
2908            proposal_id: "p1".into(),
2909            option: "x".into(),
2910            rationale: "r".into(),
2911            supporting_data: vec![],
2912        }
2913        .encode_to_vec();
2914        let req = StreamSessionRequest {
2915            envelope: Some(Envelope {
2916                macp_version: "1.0".into(),
2917                mode: "macp.mode.decision.v1".into(),
2918                message_type: "Proposal".into(),
2919                message_id: new_sid(),
2920                session_id: sid.clone(),
2921                sender: "agent://embargoed".into(),
2922                timestamp_unix_ms: Utc::now().timestamp_millis(),
2923                payload: proposal,
2924            }),
2925            subscribe_session_id: String::new(),
2926            after_sequence: 0,
2927        };
2928        let err = server
2929            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2930            .await
2931            .expect_err("stream envelope from embargoed sender must be denied");
2932        // PolicyDenied maps to FailedPrecondition on the transport (same
2933        // error the unary path expresses as a POLICY_DENIED ack).
2934        assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2935        assert!(err.message().contains("PolicyDenied"), "{err:?}");
2936
2937        // 2. Passive-subscribe frame (history read) from the embargoed
2938        //    sender: denied even though membership would allow it.
2939        let req = StreamSessionRequest {
2940            envelope: None,
2941            subscribe_session_id: sid.clone(),
2942            after_sequence: 0,
2943        };
2944        let err = server
2945            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2946            .await
2947            .expect_err("stream subscribe from embargoed sender must be denied");
2948        assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2949    }
2950    // ── ListSessions pagination (core.proto:411-426) ───────────────────
2951
2952    fn paged_session(id: &str) -> crate::session::Session {
2953        crate::session::Session::builder(id, "macp.mode.decision.v1", "agent://initiator")
2954            .participants(vec!["agent://a".into()])
2955            .mode_version("1.0.0")
2956            .configuration_version("cfg-1")
2957            .started_at_unix_ms(1)
2958            .build()
2959    }
2960
2961    /// Insert sessions straight into the registry, with each `Session`'s
2962    /// `session_id` equal to its map key (the Phase 1 `debug_assert_eq!`
2963    /// enforces the pair).
2964    async fn seed_sessions(runtime: &Arc<Runtime>, ids: &[String]) {
2965        for id in ids {
2966            runtime
2967                .registry
2968                .insert_recovered_session(id.clone(), paged_session(id))
2969                .await;
2970        }
2971    }
2972
2973    fn list_sessions_req(page_size: i32, page_token: &str) -> Request<ListSessionsRequest> {
2974        let mut req = Request::new(ListSessionsRequest {
2975            page_size,
2976            page_token: page_token.to_string(),
2977        });
2978        req.metadata_mut()
2979            .insert("authorization", "Bearer agent://observer".parse().unwrap());
2980        req
2981    }
2982
2983    fn page_size_security(default: usize, max: usize) -> SecurityLayer {
2984        let mut security = SecurityLayer::dev_mode();
2985        security.list_sessions_default_page_size = default;
2986        security.list_sessions_max_page_size = max;
2987        security
2988    }
2989
2990    fn seed_ids(n: usize) -> Vec<String> {
2991        (0..n).map(|i| format!("session-{i:03}")).collect()
2992    }
2993
2994    #[tokio::test]
2995    async fn list_sessions_applies_default_page_size_when_zero() {
2996        let (server, runtime) = make_server_with_security(page_size_security(3, 1000));
2997        seed_sessions(&runtime, &seed_ids(10)).await;
2998
2999        let resp = server
3000            .list_sessions(list_sessions_req(0, ""))
3001            .await
3002            .unwrap()
3003            .into_inner();
3004        assert_eq!(resp.sessions.len(), 3);
3005        assert!(!resp.next_page_token.is_empty());
3006    }
3007
3008    #[tokio::test]
3009    async fn list_sessions_honors_explicit_page_size() {
3010        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3011        seed_sessions(&runtime, &seed_ids(10)).await;
3012
3013        let resp = server
3014            .list_sessions(list_sessions_req(4, ""))
3015            .await
3016            .unwrap()
3017            .into_inner();
3018        assert_eq!(resp.sessions.len(), 4);
3019        assert!(!resp.next_page_token.is_empty());
3020    }
3021
3022    #[tokio::test]
3023    async fn list_sessions_clamps_page_size_above_max() {
3024        let (server, runtime) = make_server_with_security(page_size_security(100, 3));
3025        seed_sessions(&runtime, &seed_ids(10)).await;
3026
3027        let resp = server
3028            .list_sessions(list_sessions_req(1000, ""))
3029            .await
3030            .unwrap()
3031            .into_inner();
3032        assert_eq!(resp.sessions.len(), 3);
3033        assert!(!resp.next_page_token.is_empty());
3034    }
3035
3036    #[tokio::test]
3037    async fn list_sessions_rejects_negative_page_size() {
3038        let (server, runtime) = make_server();
3039        seed_sessions(&runtime, &seed_ids(3)).await;
3040
3041        let err = server
3042            .list_sessions(list_sessions_req(-1, ""))
3043            .await
3044            .unwrap_err();
3045        assert_eq!(err.code(), tonic::Code::InvalidArgument, "{err:?}");
3046        assert!(err.message().contains("page_size"), "{err:?}");
3047    }
3048
3049    #[tokio::test]
3050    async fn list_sessions_rejects_garbage_page_token() {
3051        use base64::Engine;
3052        let (server, runtime) = make_server();
3053        seed_sessions(&runtime, &seed_ids(3)).await;
3054
3055        let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD;
3056        let valid = engine.encode("v1:session-000");
3057        let tokens = vec![
3058            // not base64url
3059            "not-a-token!".to_string(),
3060            // wrong version prefix
3061            engine.encode("v2:session-000"),
3062            // prefix present, cursor empty
3063            engine.encode("v1:"),
3064            // truncated *through* the version prefix. Note that lopping bytes
3065            // off the end of an encoded token instead yields a valid, shorter
3066            // cursor — harmless, since a cursor is a position, not a handle —
3067            // so the truncation that must be rejected is the one that damages
3068            // the prefix.
3069            engine.encode("v1"),
3070            // front-truncated: the leading base64 character is gone, so the
3071            // decoded bytes are no longer valid UTF-8 (and could not carry the
3072            // prefix regardless).
3073            valid[1..].to_string(),
3074            // oversized: rejected by the length branch, before any decode
3075            "A".repeat(2 * 1024 * 1024),
3076        ];
3077        for token in tokens {
3078            let err = server
3079                .list_sessions(list_sessions_req(0, &token))
3080                .await
3081                .unwrap_err();
3082            assert_eq!(err.code(), tonic::Code::InvalidArgument);
3083            // One opaque message for every rejection reason.
3084            assert_eq!(
3085                err.message(),
3086                "INVALID_ARGUMENT: page_token is not a valid continuation token"
3087            );
3088        }
3089    }
3090
3091    #[tokio::test]
3092    async fn list_sessions_full_traversal_visits_every_session_exactly_once() {
3093        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3094        let ids = seed_ids(25);
3095        seed_sessions(&runtime, &ids).await;
3096
3097        let mut collected: Vec<String> = Vec::new();
3098        let mut token = String::new();
3099        for _ in 0..100 {
3100            let resp = server
3101                .list_sessions(list_sessions_req(4, &token))
3102                .await
3103                .unwrap()
3104                .into_inner();
3105            collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3106            token = resp.next_page_token;
3107            if token.is_empty() {
3108                break;
3109            }
3110        }
3111        assert!(token.is_empty(), "traversal did not terminate");
3112        let unique: std::collections::HashSet<&String> = collected.iter().collect();
3113        // Both assertions: the set alone would hide duplicates, the total
3114        // alone would hide a duplicate paired with a drop.
3115        assert_eq!(unique.len(), 25, "sessions were dropped or duplicated");
3116        assert_eq!(collected.len(), 25, "sessions were duplicated");
3117    }
3118
3119    #[tokio::test]
3120    async fn list_sessions_terminal_page_has_empty_next_page_token() {
3121        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3122        seed_sessions(&runtime, &seed_ids(10)).await;
3123
3124        let mut tokens: Vec<String> = Vec::new();
3125        let mut token = String::new();
3126        for _ in 0..20 {
3127            let resp = server
3128                .list_sessions(list_sessions_req(5, &token))
3129                .await
3130                .unwrap()
3131                .into_inner();
3132            token = resp.next_page_token;
3133            tokens.push(token.clone());
3134            if token.is_empty() {
3135                break;
3136            }
3137        }
3138        // 10 sessions at 5/page: exactly two pages, and only the last one
3139        // carries the empty token.
3140        assert_eq!(tokens.len(), 2, "{tokens:?}");
3141        assert!(!tokens[0].is_empty());
3142        assert!(tokens[1].is_empty());
3143    }
3144
3145    #[tokio::test]
3146    async fn list_sessions_orders_by_session_id_ascending() {
3147        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3148        // Insertion order deliberately unrelated to sort order.
3149        let ids: Vec<String> = ["delta", "alpha", "echo", "charlie", "bravo"]
3150            .iter()
3151            .map(|s| s.to_string())
3152            .collect();
3153        seed_sessions(&runtime, &ids).await;
3154
3155        let mut collected: Vec<String> = Vec::new();
3156        let mut token = String::new();
3157        loop {
3158            let resp = server
3159                .list_sessions(list_sessions_req(2, &token))
3160                .await
3161                .unwrap()
3162                .into_inner();
3163            collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3164            token = resp.next_page_token;
3165            if token.is_empty() {
3166                break;
3167            }
3168        }
3169        // Ascending across the whole traversal, not merely within a page.
3170        assert_eq!(
3171            collected,
3172            vec!["alpha", "bravo", "charlie", "delta", "echo"]
3173        );
3174    }
3175
3176    #[tokio::test]
3177    async fn list_sessions_still_requires_authentication() {
3178        let (server, runtime) = make_server();
3179        seed_sessions(&runtime, &seed_ids(3)).await;
3180
3181        // No authorization metadata, and a request body that would otherwise
3182        // be INVALID_ARGUMENT: authentication must still be what answers.
3183        let req = Request::new(ListSessionsRequest {
3184            page_size: -1,
3185            page_token: String::new(),
3186        });
3187        let err = server.list_sessions(req).await.unwrap_err();
3188        assert_eq!(err.code(), tonic::Code::Unauthenticated, "{err:?}");
3189    }
3190
3191    #[tokio::test]
3192    async fn list_sessions_tolerates_cursor_for_removed_session() {
3193        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3194        let ids = seed_ids(4);
3195        seed_sessions(&runtime, &ids).await;
3196
3197        let first = server
3198            .list_sessions(list_sessions_req(1, ""))
3199            .await
3200            .unwrap()
3201            .into_inner();
3202        assert_eq!(first.sessions[0].session_id, "session-000");
3203        assert!(!first.next_page_token.is_empty());
3204
3205        // Delete the very session the cursor names. A keyset cursor is a
3206        // position, not a handle, so paging must continue undisturbed.
3207        runtime
3208            .registry
3209            .sessions
3210            .write()
3211            .await
3212            .remove("session-000");
3213
3214        let second = server
3215            .list_sessions(list_sessions_req(1, &first.next_page_token))
3216            .await
3217            .unwrap()
3218            .into_inner();
3219        assert_eq!(second.sessions[0].session_id, "session-001");
3220    }
3221
3222    #[tokio::test]
3223    async fn list_sessions_cursor_comes_from_the_id_list_not_the_returned_sessions() {
3224        // The cursor must be the last *candidate ID*, not the last *returned
3225        // session*. The two differ only when the registry is mutated between
3226        // the ID scan and the per-ID fetch, so this test manufactures exactly
3227        // that window: the handler parks on the first page entry's session
3228        // mutex, and while it is parked the last page entry is removed.
3229        //
3230        // Deriving the cursor from the returned sessions instead would move it
3231        // backwards, re-scanning IDs the page already accounted for — and, when
3232        // an entire page vanishes, would emit an empty token and silently
3233        // terminate the traversal, dropping every remaining session.
3234        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3235        seed_sessions(&runtime, &seed_ids(6)).await;
3236
3237        // Hold the first page entry's session mutex: the handler's fetch loop
3238        // parks there, which is the only deterministic yield point between the
3239        // ID scan and the rest of the fetches.
3240        let first = runtime.registry.get_shared("session-000").await.unwrap();
3241        let guard = first.lock().await;
3242
3243        let handler = server.list_sessions(list_sessions_req(3, ""));
3244        let mutator = async {
3245            // The handler clones the Arc in `get_shared` before parking on the
3246            // mutex, so a strong count of 3 (map + this test + handler) means
3247            // it is parked. Bounded so a missed interleave fails loudly rather
3248            // than hanging.
3249            let mut spins = 0;
3250            while Arc::strong_count(&first) < 3 {
3251                assert!(spins < 10_000, "handler never parked on the session mutex");
3252                spins += 1;
3253                tokio::task::yield_now().await;
3254            }
3255            runtime
3256                .registry
3257                .sessions
3258                .write()
3259                .await
3260                .remove("session-002");
3261            drop(guard);
3262        };
3263        let (resp, ()) = tokio::join!(handler, mutator);
3264        let resp = resp.unwrap().into_inner();
3265
3266        // The interleave actually happened: the last candidate was skipped.
3267        assert_eq!(
3268            resp.sessions.len(),
3269            2,
3270            "expected session-002 to vanish between the scan and the fetch"
3271        );
3272        assert_eq!(resp.sessions[1].session_id, "session-001");
3273        // ...yet the cursor is the last candidate, not the last survivor.
3274        assert_eq!(
3275            crate::pagination::decode_page_token(&resp.next_page_token),
3276            Ok("session-002".to_string()),
3277            "cursor was derived from the returned sessions, not the ID list"
3278        );
3279
3280        // Observable through the API too: put session-002 back and page on.
3281        runtime
3282            .registry
3283            .insert_recovered_session("session-002".to_string(), paged_session("session-002"))
3284            .await;
3285        let second = server
3286            .list_sessions(list_sessions_req(3, &resp.next_page_token))
3287            .await
3288            .unwrap()
3289            .into_inner();
3290        assert_eq!(
3291            second.sessions[0].session_id, "session-003",
3292            "the cursor moved backwards past an ID the page had already accounted for"
3293        );
3294    }
3295
3296    #[tokio::test]
3297    async fn list_sessions_replaying_a_token_returns_the_identical_page() {
3298        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3299        seed_sessions(&runtime, &seed_ids(10)).await;
3300
3301        let first = server
3302            .list_sessions(list_sessions_req(3, ""))
3303            .await
3304            .unwrap()
3305            .into_inner();
3306        let token = first.next_page_token;
3307        assert!(!token.is_empty());
3308
3309        let page_a = server
3310            .list_sessions(list_sessions_req(3, &token))
3311            .await
3312            .unwrap()
3313            .into_inner();
3314        let page_b = server
3315            .list_sessions(list_sessions_req(3, &token))
3316            .await
3317            .unwrap()
3318            .into_inner();
3319
3320        let ids_a: Vec<&str> = page_a.sessions.iter().map(|s| &*s.session_id).collect();
3321        let ids_b: Vec<&str> = page_b.sessions.iter().map(|s| &*s.session_id).collect();
3322        assert_eq!(ids_a, ids_b);
3323        assert_eq!(page_a.next_page_token, page_b.next_page_token);
3324    }
3325
3326    #[tokio::test]
3327    async fn list_sessions_survives_zero_effective_page_size() {
3328        // Both fields are `pub`, so a consumer can reach 0. The floor in the
3329        // handler must keep the response well-formed: never an empty page
3330        // paired with a non-empty token (which would never advance).
3331        let (server, runtime) = make_server_with_security(page_size_security(0, 0));
3332        seed_sessions(&runtime, &seed_ids(3)).await;
3333
3334        let resp = server
3335            .list_sessions(list_sessions_req(0, ""))
3336            .await
3337            .unwrap()
3338            .into_inner();
3339        assert!(
3340            !resp.sessions.is_empty(),
3341            "empty page with token {:?} — the traversal terminates and ListSessions returns nothing",
3342            resp.next_page_token
3343        );
3344        assert_eq!(resp.sessions.len(), 1);
3345        assert!(!resp.next_page_token.is_empty());
3346
3347        // And it actually advances.
3348        let next = server
3349            .list_sessions(list_sessions_req(0, &resp.next_page_token))
3350            .await
3351            .unwrap()
3352            .into_inner();
3353        assert_eq!(next.sessions.len(), 1);
3354        assert_ne!(next.sessions[0].session_id, resp.sessions[0].session_id);
3355    }
3356
3357    fn watch_sessions_req(sender: &str) -> Request<WatchSessionsRequest> {
3358        let mut req = Request::new(WatchSessionsRequest {});
3359        req.metadata_mut()
3360            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
3361        req
3362    }
3363
3364    /// Read the next event, failing (rather than hanging) if none arrives.
3365    async fn next_lifecycle_event(
3366        stream: &mut <MacpServer as MacpRuntimeService>::WatchSessionsStream,
3367    ) -> crate::pb::SessionLifecycleEvent {
3368        use tokio_stream::StreamExt;
3369        let resp = tokio::time::timeout(std::time::Duration::from_secs(5), stream.next())
3370            .await
3371            .expect("WatchSessions produced no event within 5s")
3372            .expect("stream ended")
3373            .expect("stream errored");
3374        resp.event.expect("event present")
3375    }
3376
3377    /// Criterion 1 through the handler: N sessions in the registry produce
3378    /// exactly N `Created` events, one per session, now that the sync
3379    /// materializes them one at a time instead of deep-cloning the registry.
3380    #[tokio::test]
3381    async fn watch_sessions_initial_sync_emits_each_session_exactly_once() {
3382        let (server, runtime) = make_server();
3383        let ids = seed_ids(24);
3384        seed_sessions(&runtime, &ids).await;
3385
3386        let mut stream = server
3387            .watch_sessions(watch_sessions_req("agent://observer"))
3388            .await
3389            .unwrap()
3390            .into_inner();
3391
3392        let mut counts: HashMap<String, usize> = HashMap::new();
3393        for _ in 0..ids.len() {
3394            let event = next_lifecycle_event(&mut stream).await;
3395            assert_eq!(
3396                event.event_type,
3397                session_lifecycle_event::EventType::Created as i32
3398            );
3399            let session = event.session.expect("initial sync always carries metadata");
3400            *counts.entry(session.session_id).or_default() += 1;
3401        }
3402        assert_eq!(counts.len(), ids.len(), "sync emitted the wrong set");
3403        for id in &ids {
3404            assert_eq!(
3405                counts.get(id).copied(),
3406                Some(1),
3407                "{id} was not emitted exactly once"
3408            );
3409        }
3410    }
3411
3412    /// Criterion 4: the lifecycle subscription is taken during the unary call,
3413    /// not lazily inside the generator.
3414    ///
3415    /// The generator is not polled until the client reads, so this drives a
3416    /// session terminal *after* the response is returned but *before* the first
3417    /// read. The `Cancelled` event is published while nothing is polling the
3418    /// stream — it can only be delivered because the subscription already
3419    /// existed. Move `subscribe_session_lifecycle()` inside the `try_stream!`
3420    /// and the second read here finds nothing and times out.
3421    #[tokio::test]
3422    async fn watch_sessions_subscribes_before_the_generator_is_polled() {
3423        let (server, _runtime) = make_server();
3424        let initiator = "agent://orchestrator";
3425        let sid = new_sid();
3426        start_session(&server, initiator, &sid, vec![initiator.into()]).await;
3427
3428        let mut stream = server
3429            .watch_sessions(watch_sessions_req("agent://observer"))
3430            .await
3431            .unwrap()
3432            .into_inner();
3433
3434        // Nothing has polled the stream yet; this event has only the
3435        // subscription taken above to land in.
3436        let mut cancel = Request::new(CancelSessionRequest {
3437            session_id: sid.clone(),
3438            reason: "test".into(),
3439        });
3440        cancel.metadata_mut().insert(
3441            "authorization",
3442            format!("Bearer {initiator}").parse().unwrap(),
3443        );
3444        let ack = server
3445            .cancel_session(cancel)
3446            .await
3447            .unwrap()
3448            .into_inner()
3449            .ack
3450            .unwrap();
3451        assert!(ack.ok);
3452
3453        // The initial sync comes first...
3454        let first = next_lifecycle_event(&mut stream).await;
3455        assert_eq!(
3456            first.event_type,
3457            session_lifecycle_event::EventType::Created as i32
3458        );
3459        assert_eq!(first.session.unwrap().session_id, sid);
3460
3461        // ...then the event buffered while the generator was still unpolled.
3462        let second = next_lifecycle_event(&mut stream).await;
3463        assert_eq!(
3464            second.event_type,
3465            session_lifecycle_event::EventType::Cancelled as i32,
3466            "the event published before the first poll was lost — the \
3467             subscription must be taken in the unary call"
3468        );
3469        assert_eq!(second.session.unwrap().session_id, sid);
3470    }
3471
3472    /// Both arms of the `Created` dedup, with the race that makes it necessary
3473    /// driven deterministically.
3474    ///
3475    /// `process_session_start` registers the session BEFORE it publishes
3476    /// `Created`, so a session started after the subscribe but before the first
3477    /// poll is in the sync snapshot *and* has a `Created` sitting on the bus.
3478    /// The sync must emit it once and the buffered copy must be dropped. A
3479    /// session started after the sync is not in the snapshot, and its live
3480    /// `Created` must pass through — the handler tests membership of the sync
3481    /// set with `contains` rather than `insert`, so that pass-through does not
3482    /// extend the set; the set stays bounded by the registry size at subscribe
3483    /// time. That bound is structural and has no observable signal, so what is
3484    /// asserted here is the exactly-once contract it must not break.
3485    #[tokio::test]
3486    async fn watch_sessions_emits_created_once_for_synced_and_live_sessions() {
3487        let (server, _runtime) = make_server();
3488        let initiator = "agent://orchestrator";
3489        let synced_sid = new_sid();
3490
3491        let mut stream = server
3492            .watch_sessions(watch_sessions_req("agent://observer"))
3493            .await
3494            .unwrap()
3495            .into_inner();
3496
3497        // Registered and published while nothing is polling: this lands in the
3498        // subscription AND in the snapshot the first poll takes.
3499        start_session(&server, initiator, &synced_sid, vec![initiator.into()]).await;
3500
3501        let from_sync = next_lifecycle_event(&mut stream).await;
3502        assert_eq!(
3503            from_sync.event_type,
3504            session_lifecycle_event::EventType::Created as i32
3505        );
3506        assert_eq!(from_sync.session.unwrap().session_id, synced_sid);
3507
3508        // Started after the sync, so it is absent from the snapshot.
3509        let live_sid = new_sid();
3510        start_session(&server, initiator, &live_sid, vec![initiator.into()]).await;
3511
3512        // The next event must be the live session's Created. If the buffered
3513        // duplicate leaked through, this is `synced_sid` a second time.
3514        let live = next_lifecycle_event(&mut stream).await;
3515        assert_eq!(
3516            live.event_type,
3517            session_lifecycle_event::EventType::Created as i32
3518        );
3519        assert_eq!(
3520            live.session.unwrap().session_id,
3521            live_sid,
3522            "the sync entry's buffered Created must be suppressed, and the \
3523             live session's must not be"
3524        );
3525
3526        // And nothing further: neither Created repeats.
3527        use tokio_stream::StreamExt;
3528        let extra =
3529            tokio::time::timeout(std::time::Duration::from_millis(300), stream.next()).await;
3530        assert!(
3531            extra.is_err(),
3532            "unexpected extra lifecycle event: {extra:?}"
3533        );
3534    }
3535
3536    /// Criterion 3 end to end through the real handler, which the unit tests of
3537    /// `drain_lifecycle_events` cannot reach: a client that reads slowly while
3538    /// lifecycle events arrive throughout a long sync must not be killed with
3539    /// `RESOURCE_EXHAUSTED`, and must still see every `Created` exactly once.
3540    ///
3541    /// The shape matters. The lifecycle bus holds 64 events, and far more than
3542    /// that arrive here — but they arrive *interleaved* with the reads, which is
3543    /// what a slow consumer actually looks like. The handler drains the bus
3544    /// before fetching each session, so the receiver never falls 64 behind.
3545    /// Delete that drain (or hoist it out of the loop) and the bus overruns
3546    /// mid-sync, the first post-sync `recv()` returns `Lagged`, and this test
3547    /// fails on the stream error.
3548    #[tokio::test]
3549    async fn watch_sessions_survives_a_slow_consumer_during_a_long_sync() {
3550        use tokio_stream::StreamExt;
3551
3552        let (server, runtime) = make_server();
3553        let seeded = seed_ids(200);
3554        seed_sessions(&runtime, &seeded).await;
3555
3556        let mut stream = server
3557            .watch_sessions(watch_sessions_req("agent://observer"))
3558            .await
3559            .unwrap()
3560            .into_inner();
3561
3562        let mut counts: HashMap<String, usize> = HashMap::new();
3563        // Reads the next event, failing loudly if the stream errored — that is
3564        // the RESOURCE_EXHAUSTED this test exists to rule out.
3565        async fn read(
3566            stream: &mut <MacpServer as MacpRuntimeService>::WatchSessionsStream,
3567            counts: &mut HashMap<String, usize>,
3568        ) {
3569            let resp = tokio::time::timeout(std::time::Duration::from_secs(10), stream.next())
3570                .await
3571                .expect("WatchSessions stalled")
3572                .expect("stream ended early")
3573                .expect("stream must not be terminated (RESOURCE_EXHAUSTED)");
3574            let event = resp.event.expect("event present");
3575            if event.event_type == session_lifecycle_event::EventType::Created as i32 {
3576                let session = event.session.expect("Created always carries metadata");
3577                *counts.entry(session.session_id).or_default() += 1;
3578            }
3579        }
3580
3581        // One read starts the sync and suspends the generator at its first
3582        // yield, with 199 sessions still to emit.
3583        read(&mut stream, &mut counts).await;
3584
3585        // 70 live starts — more than the bus capacity — spread across the sync,
3586        // two reads per start so the consumer stays behind the producer without
3587        // ever stopping.
3588        let initiator = "agent://orchestrator";
3589        let mut live = Vec::new();
3590        for _ in 0..70 {
3591            let sid = new_sid();
3592            start_session(&server, initiator, &sid, vec![initiator.into()]).await;
3593            live.push(sid);
3594            read(&mut stream, &mut counts).await;
3595            read(&mut stream, &mut counts).await;
3596        }
3597
3598        // Drain the rest of the sync plus the buffered live events.
3599        let expected = seeded.len() + live.len();
3600        while counts.len() < expected {
3601            read(&mut stream, &mut counts).await;
3602        }
3603
3604        for id in seeded.iter().chain(live.iter()) {
3605            assert_eq!(
3606                counts.get(id).copied(),
3607                Some(1),
3608                "{id} was not emitted exactly once"
3609            );
3610        }
3611        assert_eq!(counts.len(), expected, "unexpected extra sessions emitted");
3612    }
3613}