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 != macp_core::MACP_VERSION {
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
804            .supported_protocol_versions
805            .iter()
806            .any(|v| v == macp_core::MACP_VERSION)
807        {
808            return Err(Status::failed_precondition(
809                "UNSUPPORTED_PROTOCOL_VERSION: no mutually supported protocol version",
810            ));
811        }
812
813        Ok(Response::new(InitializeResponse {
814            selected_protocol_version: macp_core::MACP_VERSION.into(),
815            runtime_info: Some(RuntimeInfo {
816                name: "macp-runtime".into(),
817                title: "MACP Reference Runtime".into(),
818                // Tracks the workspace version; a literal here drifted once
819                // (still advertising 0.4.0 after the 0.5.0 release).
820                version: env!("CARGO_PKG_VERSION").into(),
821                description: "Reference implementation of the Multi-Agent Coordination Protocol"
822                    .into(),
823                website_url: String::new(),
824            }),
825            capabilities: Some(Capabilities {
826                sessions: Some(SessionsCapability { stream: true, list_sessions: true, watch_sessions: true }),
827                cancellation: Some(CancellationCapability {
828                    cancel_session: true,
829                }),
830                progress: Some(ProgressCapability { progress: true }),
831                manifest: Some(ManifestCapability { get_manifest: true }),
832                mode_registry: Some(ModeRegistryCapability {
833                    list_modes: true,
834                    list_changed: true,
835                }),
836                roots: Some(RootsCapability {
837                    // ListRoots is answerable (the root set is empty — a valid
838                    // state), but this runtime has no roots provider, so the
839                    // set never changes: do not advertise change notifications
840                    // (RFC-MACP-0006 §3.3 gates WatchRoots on list_changed).
841                    // Revisit when a roots provider lands (plans E2).
842                    list_roots: true,
843                    list_changed: false,
844                }),
845                policy_registry: Some(PolicyRegistryCapability {
846                    register_policy: !self.policies_read_only,
847                    list_policies: true,
848                    list_changed: true,
849                }),
850                experimental: Some(crate::pb::ExperimentalCapabilities {
851                    features: HashMap::from([
852                        ("ext_mode_lifecycle".into(), "true".into()),
853                    ]),
854                }),
855            }),
856            supported_modes: self.runtime.registered_mode_names(),
857            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(),
858        }))
859    }
860
861    async fn send(&self, request: Request<SendRequest>) -> Result<Response<SendResponse>, Status> {
862        let env = request
863            .get_ref()
864            .envelope
865            .clone()
866            .ok_or_else(|| Status::invalid_argument("SendRequest must contain an envelope"))?;
867
868        let result = async {
869            self.validate_envelope_shape(&env)?;
870            let (env, max_open) = self.authenticate_send_request(&request, env).await?;
871            self.runtime
872                .process(&env, max_open)
873                .await
874                .map(|process_result| (env, process_result))
875        }
876        .await;
877
878        let ack = match result {
879            Ok((env, process_result)) => Ack {
880                ok: true,
881                duplicate: process_result.duplicate,
882                message_id: env.message_id.clone(),
883                session_id: env.session_id.clone(),
884                accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
885                session_state: Self::session_state_to_pb(&process_result.session_state),
886                error: None,
887            },
888            Err(err) => {
889                let env = request.get_ref().envelope.clone().unwrap_or_default();
890                // Rejection counters were collected but never recorded before
891                // (permanently zero). Session-scoped rejections are counted
892                // per mode; commitments additionally under their own counter.
893                if !env.session_id.is_empty() {
894                    self.runtime.metrics().record_message_rejected(&env.mode);
895                    if env.message_type == "Commitment" {
896                        self.runtime.metrics().record_commitment_rejected(&env.mode);
897                    }
898                }
899                Self::make_error_ack(&err, &env)
900            }
901        };
902
903        Ok(Response::new(SendResponse { ack: Some(ack) }))
904    }
905
906    async fn get_session(
907        &self,
908        request: Request<GetSessionRequest>,
909    ) -> Result<Response<GetSessionResponse>, Status> {
910        let session_id = request.get_ref().session_id.clone();
911        let _identity = self
912            .authenticate_session_access(&request, &session_id)
913            .await?;
914        let session = self
915            .runtime
916            .get_session_checked(&session_id)
917            .await
918            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
919
920        Ok(Response::new(GetSessionResponse {
921            metadata: Some(Self::session_to_metadata(&session)),
922        }))
923    }
924
925    async fn cancel_session(
926        &self,
927        request: Request<CancelSessionRequest>,
928    ) -> Result<Response<CancelSessionResponse>, Status> {
929        let session_id = request.get_ref().session_id.clone();
930        let identity = self
931            .security
932            .authenticate_metadata(request.metadata())
933            .await
934            .map_err(Self::status_from_error)?;
935        let session = self
936            .runtime
937            .get_session_checked(&session_id)
938            .await
939            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
940        // RFC-MACP-0001: "Only the initiator and policy-delegated roles may cancel."
941        // CancelSession is a Core control-plane message — mode authorization does not apply.
942        if identity.sender != session.initiator_sender
943            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
944        {
945            return Err(Status::permission_denied(
946                "FORBIDDEN: only the session initiator or policy-delegated roles can cancel",
947            ));
948        }
949        let sender = identity.sender.clone();
950        let req = request.into_inner();
951        match self
952            .runtime
953            .cancel_session(&req.session_id, &req.reason, &sender)
954            .await
955        {
956            Ok(result) => Ok(Response::new(CancelSessionResponse {
957                ack: Some(Ack {
958                    ok: true,
959                    duplicate: false,
960                    message_id: String::new(),
961                    session_id: req.session_id,
962                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
963                    session_state: Self::session_state_to_pb(&result.session_state),
964                    error: None,
965                }),
966            })),
967            Err(err) => Ok(Response::new(CancelSessionResponse {
968                ack: Some(Ack {
969                    ok: false,
970                    duplicate: false,
971                    message_id: String::new(),
972                    session_id: req.session_id.clone(),
973                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
974                    session_state: PbSessionState::Unspecified.into(),
975                    error: Some(PbMacpError {
976                        code: err.error_code().into(),
977                        message: err.to_string(),
978                        session_id: req.session_id,
979                        message_id: String::new(),
980                        details: vec![],
981                    }),
982                }),
983            })),
984        }
985    }
986
987    async fn suspend_session(
988        &self,
989        request: Request<SuspendSessionRequest>,
990    ) -> Result<Response<SuspendSessionResponse>, Status> {
991        let session_id = request.get_ref().session_id.clone();
992        let identity = self
993            .security
994            .authenticate_metadata(request.metadata())
995            .await
996            .map_err(Self::status_from_error)?;
997        let session = self
998            .runtime
999            .get_session_checked(&session_id)
1000            .await
1001            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1002        // RFC-MACP-0001 §7.5: same authority model as CancelSession — initiator
1003        // or policy-delegated roles only; mode authorization does not apply.
1004        if identity.sender != session.initiator_sender
1005            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1006        {
1007            return Err(Status::permission_denied(
1008                "FORBIDDEN: only the session initiator or policy-delegated roles can suspend",
1009            ));
1010        }
1011        let sender = identity.sender.clone();
1012        let req = request.into_inner();
1013        match self
1014            .runtime
1015            .suspend_session(&req.session_id, &req.reason, &sender)
1016            .await
1017        {
1018            Ok(result) => Ok(Response::new(SuspendSessionResponse {
1019                ack: Some(Ack {
1020                    ok: true,
1021                    duplicate: false,
1022                    message_id: String::new(),
1023                    session_id: req.session_id,
1024                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1025                    session_state: Self::session_state_to_pb(&result.session_state),
1026                    error: None,
1027                }),
1028            })),
1029            Err(err) => Ok(Response::new(SuspendSessionResponse {
1030                ack: Some(Ack {
1031                    ok: false,
1032                    duplicate: false,
1033                    message_id: String::new(),
1034                    session_id: req.session_id.clone(),
1035                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1036                    session_state: PbSessionState::Unspecified.into(),
1037                    error: Some(PbMacpError {
1038                        code: err.error_code().into(),
1039                        message: err.to_string(),
1040                        session_id: req.session_id,
1041                        message_id: String::new(),
1042                        details: vec![],
1043                    }),
1044                }),
1045            })),
1046        }
1047    }
1048
1049    async fn resume_session(
1050        &self,
1051        request: Request<ResumeSessionRequest>,
1052    ) -> Result<Response<ResumeSessionResponse>, Status> {
1053        let session_id = request.get_ref().session_id.clone();
1054        let identity = self
1055            .security
1056            .authenticate_metadata(request.metadata())
1057            .await
1058            .map_err(Self::status_from_error)?;
1059        let session = self
1060            .runtime
1061            .get_session_checked(&session_id)
1062            .await
1063            .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1064        if identity.sender != session.initiator_sender
1065            && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1066        {
1067            return Err(Status::permission_denied(
1068                "FORBIDDEN: only the session initiator or policy-delegated roles can resume",
1069            ));
1070        }
1071        let sender = identity.sender.clone();
1072        let req = request.into_inner();
1073        match self
1074            .runtime
1075            .resume_session(&req.session_id, &req.reason, &sender)
1076            .await
1077        {
1078            Ok(result) => Ok(Response::new(ResumeSessionResponse {
1079                ack: Some(Ack {
1080                    ok: true,
1081                    duplicate: false,
1082                    message_id: String::new(),
1083                    session_id: req.session_id,
1084                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1085                    session_state: Self::session_state_to_pb(&result.session_state),
1086                    error: None,
1087                }),
1088            })),
1089            Err(err) => Ok(Response::new(ResumeSessionResponse {
1090                ack: Some(Ack {
1091                    ok: false,
1092                    duplicate: false,
1093                    message_id: String::new(),
1094                    session_id: req.session_id.clone(),
1095                    accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1096                    session_state: PbSessionState::Unspecified.into(),
1097                    error: Some(PbMacpError {
1098                        code: err.error_code().into(),
1099                        message: err.to_string(),
1100                        session_id: req.session_id,
1101                        message_id: String::new(),
1102                        details: vec![],
1103                    }),
1104                }),
1105            })),
1106        }
1107    }
1108
1109    async fn get_manifest(
1110        &self,
1111        request: Request<GetManifestRequest>,
1112    ) -> Result<Response<GetManifestResponse>, Status> {
1113        let req = request.into_inner();
1114        if !req.agent_id.is_empty() && req.agent_id != "macp-runtime" {
1115            return Err(Status::not_found(format!(
1116                "Agent '{}' not found",
1117                req.agent_id
1118            )));
1119        }
1120
1121        Ok(Response::new(GetManifestResponse {
1122            manifest: Some(crate::pb::AgentManifest {
1123                agent_id: "macp-runtime".into(),
1124                title: "MACP Reference Runtime".into(),
1125                description: "Reference implementation of MACP".into(),
1126                supported_modes: self.runtime.registered_mode_names(),
1127                input_content_types: vec!["application/macp-envelope+proto".into()],
1128                output_content_types: vec!["application/macp-envelope+proto".into()],
1129                metadata: HashMap::new(),
1130                // Empty: unary-first profile has no dedicated transport endpoints.
1131                transport_endpoints: vec![],
1132            }),
1133        }))
1134    }
1135
1136    async fn list_modes(
1137        &self,
1138        _request: Request<ListModesRequest>,
1139    ) -> Result<Response<ListModesResponse>, Status> {
1140        Ok(Response::new(ListModesResponse {
1141            modes: self.runtime.standard_mode_descriptors(),
1142        }))
1143    }
1144
1145    async fn list_roots(
1146        &self,
1147        _request: Request<ListRootsRequest>,
1148    ) -> Result<Response<ListRootsResponse>, Status> {
1149        Ok(Response::new(ListRootsResponse { roots: vec![] }))
1150    }
1151
1152    type StreamSessionStream = SessionResponseStream;
1153
1154    async fn stream_session(
1155        &self,
1156        request: Request<tonic::Streaming<StreamSessionRequest>>,
1157    ) -> Result<Response<Self::StreamSessionStream>, Status> {
1158        let identity = self
1159            .security
1160            .authenticate_metadata(request.metadata())
1161            .await
1162            .map_err(Self::status_from_error)?;
1163        let inbound = request.into_inner();
1164        Ok(Response::new(
1165            self.build_stream_session_stream(identity, inbound),
1166        ))
1167    }
1168
1169    type WatchModeRegistryStream = std::pin::Pin<
1170        Box<dyn futures_core::Stream<Item = Result<WatchModeRegistryResponse, Status>> + Send>,
1171    >;
1172
1173    async fn watch_mode_registry(
1174        &self,
1175        _request: Request<WatchModeRegistryRequest>,
1176    ) -> Result<Response<Self::WatchModeRegistryStream>, Status> {
1177        let mut rx = self.runtime.subscribe_mode_changes();
1178        let stream = async_stream::try_stream! {
1179            // Send initial state
1180            yield WatchModeRegistryResponse {
1181                change: Some(crate::pb::RegistryChanged {
1182                    registry: "modes".into(),
1183                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1184                }),
1185            };
1186            // Wait for changes from register/unregister/promote
1187            while rx.recv().await.is_ok() {
1188                yield WatchModeRegistryResponse {
1189                    change: Some(crate::pb::RegistryChanged {
1190                        registry: "modes".into(),
1191                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1192                    }),
1193                };
1194            }
1195        };
1196        Ok(Response::new(Box::pin(stream)))
1197    }
1198
1199    type WatchRootsStream = std::pin::Pin<
1200        Box<dyn futures_core::Stream<Item = Result<WatchRootsResponse, Status>> + Send>,
1201    >;
1202
1203    async fn watch_roots(
1204        &self,
1205        _request: Request<WatchRootsRequest>,
1206    ) -> Result<Response<Self::WatchRootsStream>, Status> {
1207        let initial = WatchRootsResponse {
1208            change: Some(crate::pb::RootsChanged {
1209                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1210            }),
1211        };
1212        let stream = async_stream::try_stream! {
1213            yield initial;
1214            // Roots are static — keep the stream open but idle.
1215            std::future::pending::<()>().await;
1216        };
1217        Ok(Response::new(Box::pin(stream)))
1218    }
1219
1220    type WatchSignalsStream = std::pin::Pin<
1221        Box<dyn futures_core::Stream<Item = Result<WatchSignalsResponse, Status>> + Send>,
1222    >;
1223
1224    type WatchSessionsStream = std::pin::Pin<
1225        Box<dyn futures_core::Stream<Item = Result<WatchSessionsResponse, Status>> + Send>,
1226    >;
1227
1228    async fn watch_signals(
1229        &self,
1230        request: Request<WatchSignalsRequest>,
1231    ) -> Result<Response<Self::WatchSignalsStream>, Status> {
1232        // Ambient signals carry agent-generated payload data; subscribing is
1233        // gated on authentication like the session-observation surfaces.
1234        // (RFC-0004 §4.1 constrains unauthenticated *producers*; requiring
1235        // authenticated subscribers is this runtime's hardening posture.)
1236        let _identity = self
1237            .security
1238            .authenticate_metadata(request.metadata())
1239            .await
1240            .map_err(Self::status_from_error)?;
1241        let mut rx = self.runtime.subscribe_signals();
1242        let stream = async_stream::try_stream! {
1243            loop {
1244                match rx.recv().await {
1245                    Ok(envelope) => {
1246                        yield WatchSignalsResponse {
1247                            envelope: Some(envelope),
1248                        };
1249                    }
1250                    // Surface lag instead of silently ending the stream: a
1251                    // slow consumer must be able to distinguish "no traffic"
1252                    // from "events dropped" (mirrors StreamSession).
1253                    Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1254                        Err(Status::resource_exhausted(format!(
1255                            "WatchSignals receiver fell behind by {skipped} signals"
1256                        )))?;
1257                    }
1258                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1259                }
1260            }
1261        };
1262        Ok(Response::new(Box::pin(stream)))
1263    }
1264
1265    // Session lifecycle observation RPCs
1266
1267    async fn list_sessions(
1268        &self,
1269        request: Request<ListSessionsRequest>,
1270    ) -> Result<Response<ListSessionsResponse>, Status> {
1271        // Authentication stays the FIRST statement: an unauthenticated caller
1272        // must see UNAUTHENTICATED, never a request-shape error, which would
1273        // confirm the endpoint is live and leak validation order pre-auth.
1274        let _identity = self
1275            .security
1276            .authenticate_metadata(request.metadata())
1277            .await
1278            .map_err(Self::status_from_error)?;
1279        let req = request.into_inner();
1280
1281        // Paged: a keyset scan over session IDs in ascending byte order, with
1282        // `next_page_token` carrying the last emitted ID as the cursor. Per
1283        // `macp-proto` core.proto:411-426, `page_size = 0` means the
1284        // server-chosen default, the server MAY cap the effective size, and a
1285        // response is complete only when `next_page_token` is empty — a page
1286        // may be short while more results remain.
1287        if req.page_size < 0 {
1288            return Err(Status::invalid_argument(
1289                "INVALID_ARGUMENT: page_size must not be negative",
1290            ));
1291        }
1292        // The cast is safe only below the negative guard above: `page_size` is
1293        // an int32, and `-1 as usize` is astronomically large on 64-bit.
1294        let effective = if req.page_size == 0 {
1295            self.security.list_sessions_default_page_size
1296        } else {
1297            (req.page_size as usize).min(self.security.list_sessions_max_page_size)
1298        };
1299        // Floor, not redundant: at `effective == 0`, `page_ids` is `&ids[..0]`,
1300        // so `page_ids.last()` is `None` and the `next_page_token` match below
1301        // falls to its `_` arm — an empty page paired with an empty token,
1302        // which looks complete. The traversal terminates immediately and
1303        // `ListSessions` silently returns nothing, forever. The configured
1304        // values are `>= 1` four ways today, but both fields are deliberately
1305        // `pub`, so a library consumer or a test can set 0.
1306        let effective = effective.max(1);
1307
1308        let cursor = if req.page_token.is_empty() {
1309            None
1310        } else {
1311            // Every decode failure collapses to one opaque message: the token
1312            // must not be an oracle for which check rejected it. The error
1313            // discriminant is dropped here and never logged (see
1314            // `crate::pagination` — the token is attacker-chosen bytes).
1315            Some(
1316                crate::pagination::decode_page_token(&req.page_token).map_err(|_| {
1317                    Status::invalid_argument(
1318                        "INVALID_ARGUMENT: page_token is not a valid continuation token",
1319                    )
1320                })?,
1321            )
1322        };
1323
1324        // Fetch one extra ID than the page holds: its presence is an exact
1325        // "more results exist" signal, so the terminal page carries an empty
1326        // token with no extra round trip. `saturating_add` because `effective`
1327        // is operator-controlled.
1328        let ids = self
1329            .runtime
1330            .registry
1331            .session_ids_after(cursor.as_deref(), effective.saturating_add(1))
1332            .await;
1333        let has_more = ids.len() > effective;
1334        let page_ids = &ids[..effective.min(ids.len())];
1335
1336        // The cursor comes from the ID list, NOT from the materialized
1337        // sessions below. A session can be removed between the ID scan and the
1338        // fetch; deriving the cursor from what survived would stall the cursor
1339        // (re-emitting the same page) or, if the whole page vanished, drop
1340        // every remaining session by terminating the traversal early.
1341        let next_page_token = match (has_more, page_ids.last()) {
1342            (true, Some(last)) => crate::pagination::encode_page_token(last),
1343            _ => String::new(),
1344        };
1345
1346        let mut metadata: Vec<SessionMetadata> = Vec::with_capacity(page_ids.len());
1347        for id in page_ids {
1348            // A `None` here means the session was removed between the scan and
1349            // this fetch; skipping it is correct. The resulting short page with
1350            // a non-empty token is explicitly permitted by core.proto:411-414.
1351            if let Some(session) = self.runtime.registry.get_session(id).await {
1352                debug_assert_eq!(
1353                    session.session_id, *id,
1354                    "registry map key must equal Session::session_id — paging orders \
1355                     by the key but emits the field"
1356                );
1357                metadata.push(Self::session_to_metadata(&session));
1358            }
1359        }
1360
1361        Ok(Response::new(ListSessionsResponse {
1362            sessions: metadata,
1363            next_page_token,
1364        }))
1365    }
1366
1367    async fn watch_sessions(
1368        &self,
1369        request: Request<WatchSessionsRequest>,
1370    ) -> Result<Response<Self::WatchSessionsStream>, Status> {
1371        let _identity = self
1372            .security
1373            .authenticate_metadata(request.metadata())
1374            .await
1375            .map_err(Self::status_from_error)?;
1376        // Subscribed HERE, during the unary call — strictly before the generator
1377        // below is first polled, which does not happen until the client reads.
1378        // An event published in that gap must be buffered by an existing
1379        // subscription, not missed; moving this inside the generator
1380        // reintroduces exactly that race.
1381        let mut rx = self.runtime.subscribe_session_lifecycle();
1382        let runtime = Arc::clone(&self.runtime);
1383        let stream = async_stream::try_stream! {
1384            // Initial sync: emit all current sessions as CREATED events.
1385            //
1386            // The traversal snapshots the registry's shared session handles
1387            // ONCE and then locks and clones one session at a time (see
1388            // `watch_sync`), so peak resident `Session` clones is one rather
1389            // than the registry size: this generator is paced by the client's
1390            // reads and there can be `MACP_MAX_CONCURRENT_STREAMS` of them at
1391            // once.
1392            let mut sync = crate::watch_sync::InitialSync::begin(&runtime.registry).await;
1393            // The IDs this sync emits. A Created event buffered in the
1394            // subscribe→sync window would duplicate one of them — the session
1395            // was already registered when the snapshot was taken, but
1396            // `process_session_start` inserts into the registry BEFORE
1397            // publishing Created, so that event can still arrive afterwards.
1398            // Session IDs are create-once and `runtime.rs` holds the only
1399            // `Created` publisher, so one `send` per session start means a
1400            // *live* Created can never repeat for an ID already in this set —
1401            // membership is read, never extended, past the sync. That keeps the
1402            // set bounded by the registry size at subscribe time instead of
1403            // growing with every session the stream ever observes.
1404            let mut synced: std::collections::HashSet<String> =
1405                std::collections::HashSet::with_capacity(sync.remaining());
1406            // Lifecycle events that arrive while the sync is still emitting.
1407            // The sync loop cannot `recv().await` (it has its own output to
1408            // produce) but must not ignore the bus either: the bus holds 64
1409            // events, so a slow sync would otherwise make the first post-sync
1410            // `recv()` return `Lagged` and kill the stream. Bounded by
1411            // `PENDING_EVENT_LIMIT`; on overflow the client gets the same
1412            // `RESOURCE_EXHAUSTED` it gets for bus lag.
1413            let mut pending: std::collections::VecDeque<crate::runtime::SessionLifecycleEvent> =
1414                std::collections::VecDeque::new();
1415            loop {
1416                // Drained BEFORE the next session is fetched, so the bus is
1417                // relieved once per emitted session rather than once for the
1418                // whole sync.
1419                if let Err(drain_err) = crate::watch_sync::drain_lifecycle_events(
1420                    &mut rx,
1421                    &mut pending,
1422                    crate::watch_sync::PENDING_EVENT_LIMIT,
1423                ) {
1424                    Err(Status::resource_exhausted(drain_err.message()))?;
1425                    break;
1426                }
1427                // One session, emitted and dropped before the next is asked
1428                // for, which is what keeps residency bounded.
1429                let Some(session) = sync.next_session().await else { break };
1430                synced.insert(session.session_id.clone());
1431                yield WatchSessionsResponse {
1432                    event: Some(SessionLifecycleEvent {
1433                        event_type: session_lifecycle_event::EventType::Created.into(),
1434                        session: Some(Self::session_to_metadata(&session)),
1435                        observed_at_unix_ms: session.started_at_unix_ms,
1436                    }),
1437                };
1438            }
1439            // Stream lifecycle transitions: first the ones buffered during the
1440            // sync (in bus order), then live ones.
1441            loop {
1442                let event = match pending.pop_front() {
1443                    Some(event) => event,
1444                    None => match rx.recv().await {
1445                        Ok(event) => event,
1446                        Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1447                            Err(Status::resource_exhausted(format!(
1448                                "WatchSessions receiver fell behind by {skipped} events"
1449                            )))?;
1450                            break;
1451                        }
1452                        Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1453                    },
1454                };
1455                let (event_type, sid) = match &event {
1456                    crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1457                        (session_lifecycle_event::EventType::Created, session_id.clone()),
1458                    crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1459                        (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1460                    crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1461                        (session_lifecycle_event::EventType::Expired, session_id.clone()),
1462                    crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1463                        (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1464                    crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1465                        (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1466                    crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1467                        (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1468                };
1469                // Skip the buffered duplicate of an initial-sync entry;
1470                // non-Created events for synced sessions are new information
1471                // and pass through.
1472                //
1473                // `contains`, not `insert`: a live Created for a session the
1474                // sync never saw is emitted as-is and NOT recorded. Recording
1475                // it would be the only thing making this set grow with the
1476                // stream's lifetime, and it would buy nothing — `runtime.rs`
1477                // is the single `Created` publisher and sends once per session
1478                // start, so no live Created can repeat.
1479                if event_type == session_lifecycle_event::EventType::Created
1480                    && synced.contains(&sid)
1481                {
1482                    continue;
1483                }
1484                let session_meta = runtime.registry.get_session(&sid).await
1485                    .map(|s| Self::session_to_metadata(&s));
1486                yield WatchSessionsResponse {
1487                    event: Some(SessionLifecycleEvent {
1488                        event_type: event_type.into(),
1489                        session: session_meta,
1490                        observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1491                    }),
1492                };
1493            }
1494        };
1495        Ok(Response::new(Box::pin(stream)))
1496    }
1497
1498    // Extension mode lifecycle RPCs
1499
1500    async fn list_ext_modes(
1501        &self,
1502        _request: Request<ListExtModesRequest>,
1503    ) -> Result<Response<ListExtModesResponse>, Status> {
1504        Ok(Response::new(ListExtModesResponse {
1505            modes: self.runtime.extension_mode_descriptors(),
1506        }))
1507    }
1508
1509    async fn register_ext_mode(
1510        &self,
1511        request: Request<RegisterExtModeRequest>,
1512    ) -> Result<Response<RegisterExtModeResponse>, Status> {
1513        let identity = self
1514            .security
1515            .authenticate_metadata(request.metadata())
1516            .await
1517            .map_err(Self::status_from_error)?;
1518        self.security
1519            .authorize_mode_registry(&identity)
1520            .map_err(Self::status_from_error)?;
1521        let req = request.into_inner();
1522        let descriptor = req
1523            .mode_descriptor
1524            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1525        match self.runtime.register_extension(descriptor) {
1526            Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1527                ok: true,
1528                error: String::new(),
1529            })),
1530            Err(e) => Ok(Response::new(RegisterExtModeResponse {
1531                ok: false,
1532                error: e,
1533            })),
1534        }
1535    }
1536
1537    async fn unregister_ext_mode(
1538        &self,
1539        request: Request<UnregisterExtModeRequest>,
1540    ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1541        let identity = self
1542            .security
1543            .authenticate_metadata(request.metadata())
1544            .await
1545            .map_err(Self::status_from_error)?;
1546        self.security
1547            .authorize_mode_registry(&identity)
1548            .map_err(Self::status_from_error)?;
1549        let req = request.into_inner();
1550        match self.runtime.unregister_extension(&req.mode) {
1551            Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1552                ok: true,
1553                error: String::new(),
1554            })),
1555            Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1556                ok: false,
1557                error: e,
1558            })),
1559        }
1560    }
1561
1562    async fn promote_mode(
1563        &self,
1564        request: Request<PromoteModeRequest>,
1565    ) -> Result<Response<PromoteModeResponse>, Status> {
1566        let identity = self
1567            .security
1568            .authenticate_metadata(request.metadata())
1569            .await
1570            .map_err(Self::status_from_error)?;
1571        self.security
1572            .authorize_mode_registry(&identity)
1573            .map_err(Self::status_from_error)?;
1574        let req = request.into_inner();
1575        let new_name = if req.promoted_mode_name.is_empty() {
1576            None
1577        } else {
1578            Some(req.promoted_mode_name.as_str())
1579        };
1580        match self.runtime.promote_mode(&req.mode, new_name) {
1581            Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1582                ok: true,
1583                error: String::new(),
1584                mode: final_name,
1585            })),
1586            Err(e) => Ok(Response::new(PromoteModeResponse {
1587                ok: false,
1588                error: e,
1589                mode: String::new(),
1590            })),
1591        }
1592    }
1593
1594    // ── Governance policy lifecycle RPCs (RFC-MACP-0012) ────────────
1595
1596    async fn register_policy(
1597        &self,
1598        request: Request<RegisterPolicyRequest>,
1599    ) -> Result<Response<RegisterPolicyResponse>, Status> {
1600        if self.policies_read_only {
1601            return Err(Status::failed_precondition(
1602                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1603            ));
1604        }
1605        let identity = self
1606            .security
1607            .authenticate_metadata(request.metadata())
1608            .await
1609            .map_err(Self::status_from_error)?;
1610        self.security
1611            .authorize_mode_registry(&identity)
1612            .map_err(Self::status_from_error)?;
1613        let req = request.into_inner();
1614        let descriptor = req
1615            .policy_descriptor
1616            .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1617        let definition = Self::policy_descriptor_to_definition(&descriptor);
1618        match self.runtime.register_policy(definition) {
1619            Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1620                ok: true,
1621                error: String::new(),
1622            })),
1623            Err(e) => Ok(Response::new(RegisterPolicyResponse {
1624                ok: false,
1625                error: e,
1626            })),
1627        }
1628    }
1629
1630    async fn unregister_policy(
1631        &self,
1632        request: Request<UnregisterPolicyRequest>,
1633    ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1634        if self.policies_read_only {
1635            return Err(Status::failed_precondition(
1636                "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1637            ));
1638        }
1639        let identity = self
1640            .security
1641            .authenticate_metadata(request.metadata())
1642            .await
1643            .map_err(Self::status_from_error)?;
1644        self.security
1645            .authorize_mode_registry(&identity)
1646            .map_err(Self::status_from_error)?;
1647        let req = request.into_inner();
1648        match self.runtime.unregister_policy(&req.policy_id) {
1649            Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1650                ok: true,
1651                error: String::new(),
1652            })),
1653            Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1654                ok: false,
1655                error: e,
1656            })),
1657        }
1658    }
1659
1660    async fn get_policy(
1661        &self,
1662        request: Request<GetPolicyRequest>,
1663    ) -> Result<Response<GetPolicyResponse>, Status> {
1664        let _identity = self
1665            .security
1666            .authenticate_metadata(request.metadata())
1667            .await
1668            .map_err(Self::status_from_error)?;
1669        let req = request.into_inner();
1670        let policy = self
1671            .runtime
1672            .get_policy(&req.policy_id)
1673            .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1674        Ok(Response::new(GetPolicyResponse {
1675            policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1676        }))
1677    }
1678
1679    async fn list_policies(
1680        &self,
1681        request: Request<ListPoliciesRequest>,
1682    ) -> Result<Response<ListPoliciesResponse>, Status> {
1683        let _identity = self
1684            .security
1685            .authenticate_metadata(request.metadata())
1686            .await
1687            .map_err(Self::status_from_error)?;
1688        let req = request.into_inner();
1689        let mode_filter = if req.mode.is_empty() {
1690            None
1691        } else {
1692            Some(req.mode.as_str())
1693        };
1694        let policies = self.runtime.list_policies(mode_filter);
1695        let descriptors = policies
1696            .iter()
1697            .map(Self::policy_definition_to_descriptor)
1698            .collect();
1699        Ok(Response::new(ListPoliciesResponse { descriptors }))
1700    }
1701
1702    type WatchPoliciesStream = std::pin::Pin<
1703        Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1704    >;
1705
1706    async fn watch_policies(
1707        &self,
1708        _request: Request<WatchPoliciesRequest>,
1709    ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1710        let mut rx = self.runtime.subscribe_policy_changes();
1711        let runtime = Arc::clone(&self.runtime);
1712        let stream = async_stream::try_stream! {
1713            // Send initial state
1714            let policies = runtime.list_policies(None);
1715            let descriptors: Vec<PolicyDescriptor> = policies
1716                .iter()
1717                .map(MacpServer::policy_definition_to_descriptor)
1718                .collect();
1719            yield WatchPoliciesResponse {
1720                descriptors,
1721                observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1722            };
1723            // Wait for changes
1724            while rx.recv().await.is_ok() {
1725                let policies = runtime.list_policies(None);
1726                let descriptors: Vec<PolicyDescriptor> = policies
1727                    .iter()
1728                    .map(MacpServer::policy_definition_to_descriptor)
1729                    .collect();
1730                yield WatchPoliciesResponse {
1731                    descriptors,
1732                    observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1733                };
1734            }
1735        };
1736        Ok(Response::new(Box::pin(stream)))
1737    }
1738}
1739
1740// ── Policy type conversion helpers ──────────────────────────────────
1741
1742impl MacpServer {
1743    fn policy_descriptor_to_definition(
1744        descriptor: &PolicyDescriptor,
1745    ) -> crate::policy::PolicyDefinition {
1746        let rules: serde_json::Value = if descriptor.rules.is_empty() {
1747            serde_json::json!({})
1748        } else {
1749            serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1750        };
1751        crate::policy::PolicyDefinition {
1752            policy_id: descriptor.policy_id.clone(),
1753            mode: descriptor.mode.clone(),
1754            description: descriptor.description.clone(),
1755            rules,
1756            schema_version: descriptor.schema_version,
1757        }
1758    }
1759
1760    fn policy_definition_to_descriptor(
1761        definition: &crate::policy::PolicyDefinition,
1762    ) -> PolicyDescriptor {
1763        PolicyDescriptor {
1764            policy_id: definition.policy_id.clone(),
1765            mode: definition.mode.clone(),
1766            description: definition.description.clone(),
1767            rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1768            schema_version: definition.schema_version,
1769            registered_at_unix_ms: 0,
1770        }
1771    }
1772}
1773
1774#[cfg(test)]
1775mod tests {
1776    use super::*;
1777    use crate::log_store::LogStore;
1778    use crate::pb::SessionStartPayload;
1779    use crate::registry::SessionRegistry;
1780    use chrono::Utc;
1781    use prost::Message;
1782
1783    fn new_sid() -> String {
1784        uuid::Uuid::new_v4().as_hyphenated().to_string()
1785    }
1786
1787    fn make_server() -> (MacpServer, Arc<Runtime>) {
1788        make_server_with_security(SecurityLayer::dev_mode())
1789    }
1790
1791    /// Same harness, with the `SecurityLayer` supplied by the caller so a test
1792    /// can pin `list_sessions_{default,max}_page_size` without touching process
1793    /// env (which is not deterministic under `cargo test`'s thread pool).
1794    fn make_server_with_security(security: SecurityLayer) -> (MacpServer, Arc<Runtime>) {
1795        let storage: Arc<dyn crate::storage::StorageBackend> =
1796            Arc::new(crate::storage::MemoryBackend);
1797        let registry = Arc::new(SessionRegistry::new());
1798        let log_store = Arc::new(LogStore::new());
1799        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1800        let server = MacpServer::new(runtime.clone(), security);
1801        (server, runtime)
1802    }
1803
1804    fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1805        let mut req = Request::new(SendRequest {
1806            envelope: Some(env),
1807        });
1808        req.metadata_mut()
1809            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1810        req
1811    }
1812
1813    async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1814        let resp = server.send(send_req(sender, env)).await.unwrap();
1815        resp.into_inner().ack.unwrap()
1816    }
1817
1818    fn start_payload() -> Vec<u8> {
1819        SessionStartPayload {
1820            intent: "intent".into(),
1821            participants: vec!["agent://fraud".into()],
1822            mode_version: "1.0.0".into(),
1823            configuration_version: "cfg-1".into(),
1824            policy_version: String::new(),
1825            ttl_ms: 1000,
1826            context_id: String::new(),
1827            extensions: std::collections::HashMap::new(),
1828            roots: vec![],
1829            max_suspend_ms: 0,
1830        }
1831        .encode_to_vec()
1832    }
1833
1834    #[tokio::test]
1835    async fn sender_is_derived_from_authenticated_metadata() {
1836        let (server, runtime) = make_server();
1837        let sid = new_sid();
1838        let ack = do_send(
1839            &server,
1840            "agent://orchestrator",
1841            Envelope {
1842                macp_version: "1.0".into(),
1843                mode: "macp.mode.decision.v1".into(),
1844                message_type: "SessionStart".into(),
1845                message_id: "m1".into(),
1846                session_id: sid.clone(),
1847                sender: String::new(),
1848                timestamp_unix_ms: Utc::now().timestamp_millis(),
1849                payload: start_payload(),
1850            },
1851        )
1852        .await;
1853        assert!(ack.ok);
1854        let session = runtime.get_session_checked(&sid).await.unwrap();
1855        assert_eq!(session.initiator_sender, "agent://orchestrator");
1856    }
1857
1858    #[tokio::test]
1859    async fn spoofed_sender_is_rejected() {
1860        let (server, _) = make_server();
1861        let sid = new_sid();
1862        let ack = do_send(
1863            &server,
1864            "agent://orchestrator",
1865            Envelope {
1866                macp_version: "1.0".into(),
1867                mode: "macp.mode.decision.v1".into(),
1868                message_type: "SessionStart".into(),
1869                message_id: "m1".into(),
1870                session_id: sid,
1871                sender: "agent://spoof".into(),
1872                timestamp_unix_ms: Utc::now().timestamp_millis(),
1873                payload: start_payload(),
1874            },
1875        )
1876        .await;
1877        assert!(!ack.ok);
1878        assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1879    }
1880
1881    #[tokio::test]
1882    async fn get_session_requires_session_membership() {
1883        let (server, _) = make_server();
1884        let sid = new_sid();
1885        let ack = do_send(
1886            &server,
1887            "agent://orchestrator",
1888            Envelope {
1889                macp_version: "1.0".into(),
1890                mode: "macp.mode.decision.v1".into(),
1891                message_type: "SessionStart".into(),
1892                message_id: "m1".into(),
1893                session_id: sid.clone(),
1894                sender: String::new(),
1895                timestamp_unix_ms: Utc::now().timestamp_millis(),
1896                payload: start_payload(),
1897            },
1898        )
1899        .await;
1900        assert!(ack.ok);
1901
1902        let mut req = Request::new(GetSessionRequest { session_id: sid });
1903        req.metadata_mut().insert(
1904            "authorization",
1905            format!("Bearer {}", "agent://outsider").parse().unwrap(),
1906        );
1907        let err = server.get_session(req).await.unwrap_err();
1908        assert_eq!(err.code(), tonic::Code::PermissionDenied);
1909    }
1910
1911    #[tokio::test]
1912    async fn register_ext_mode_requires_authenticated_registry_permission() {
1913        let storage: Arc<dyn crate::storage::StorageBackend> =
1914            Arc::new(crate::storage::MemoryBackend);
1915        let registry = Arc::new(SessionRegistry::new());
1916        let log_store = Arc::new(LogStore::new());
1917        let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1918        let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1919        let server = MacpServer::new(runtime, security);
1920
1921        let req = Request::new(RegisterExtModeRequest {
1922            mode_descriptor: Some(crate::pb::ModeDescriptor {
1923                mode: "ext.custom.v1".into(),
1924                mode_version: "1.0.0".into(),
1925                message_types: vec!["SessionStart".into(), "Commitment".into()],
1926                ..Default::default()
1927            }),
1928        });
1929        let err = server.register_ext_mode(req).await.unwrap_err();
1930        assert_eq!(err.code(), tonic::Code::Unauthenticated);
1931    }
1932
1933    fn stream_identity(sender: &str) -> AuthIdentity {
1934        AuthIdentity {
1935            sender: sender.into(),
1936            allowed_modes: None,
1937            can_start_sessions: true,
1938            max_open_sessions: None,
1939            can_manage_mode_registry: false,
1940            is_observer: false,
1941        }
1942    }
1943
1944    #[tokio::test]
1945    async fn stream_session_emits_accepted_envelopes_only() {
1946        use tokio_stream::{iter, StreamExt};
1947
1948        let (server, _) = make_server();
1949        let sid = new_sid();
1950        let requests = iter(vec![Ok(StreamSessionRequest {
1951            subscribe_session_id: String::new(),
1952            after_sequence: 0,
1953            envelope: Some(Envelope {
1954                macp_version: "1.0".into(),
1955                mode: "macp.mode.decision.v1".into(),
1956                message_type: "SessionStart".into(),
1957                message_id: "m1".into(),
1958                session_id: sid.clone(),
1959                sender: String::new(),
1960                timestamp_unix_ms: Utc::now().timestamp_millis(),
1961                payload: start_payload(),
1962            }),
1963        })]);
1964
1965        let mut stream =
1966            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1967
1968        let response = stream.next().await.unwrap().unwrap();
1969        let envelope = match response.response.unwrap() {
1970            crate::pb::stream_session_response::Response::Envelope(e) => e,
1971            _ => panic!("expected envelope"),
1972        };
1973        assert_eq!(envelope.message_type, "SessionStart");
1974        assert_eq!(envelope.message_id, "m1");
1975        assert!(stream.next().await.is_none());
1976    }
1977
1978    #[tokio::test]
1979    async fn stream_session_rejects_mixed_session_ids() {
1980        use tokio_stream::{iter, StreamExt};
1981
1982        let (server, _) = make_server();
1983        let sid1 = new_sid();
1984        let sid2 = new_sid();
1985        let requests = iter(vec![
1986            Ok(StreamSessionRequest {
1987                subscribe_session_id: String::new(),
1988                after_sequence: 0,
1989                envelope: Some(Envelope {
1990                    macp_version: "1.0".into(),
1991                    mode: "macp.mode.decision.v1".into(),
1992                    message_type: "SessionStart".into(),
1993                    message_id: "m1".into(),
1994                    session_id: sid1.clone(),
1995                    sender: String::new(),
1996                    timestamp_unix_ms: Utc::now().timestamp_millis(),
1997                    payload: start_payload(),
1998                }),
1999            }),
2000            Ok(StreamSessionRequest {
2001                subscribe_session_id: String::new(),
2002                after_sequence: 0,
2003                envelope: Some(Envelope {
2004                    macp_version: "1.0".into(),
2005                    mode: "macp.mode.decision.v1".into(),
2006                    message_type: "SessionStart".into(),
2007                    message_id: "m2".into(),
2008                    session_id: sid2,
2009                    sender: String::new(),
2010                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2011                    payload: start_payload(),
2012                }),
2013            }),
2014        ]);
2015
2016        let mut stream =
2017            server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
2018
2019        let first = stream.next().await.unwrap().unwrap();
2020        let first_env = match first.response.unwrap() {
2021            crate::pb::stream_session_response::Response::Envelope(e) => e,
2022            _ => panic!("expected envelope"),
2023        };
2024        assert_eq!(first_env.session_id, sid1);
2025        let err = stream.next().await.unwrap().unwrap_err();
2026        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2027    }
2028
2029    #[tokio::test]
2030    async fn list_modes_returns_standard_modes() {
2031        let (server, _) = make_server();
2032        let resp = server
2033            .list_modes(Request::new(ListModesRequest {}))
2034            .await
2035            .unwrap();
2036        let names: Vec<String> = resp
2037            .into_inner()
2038            .modes
2039            .iter()
2040            .map(|m| m.mode.clone())
2041            .collect();
2042        assert_eq!(names.len(), 5);
2043        assert!(names.contains(&"macp.mode.decision.v1".to_string()));
2044        assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
2045        assert!(names.contains(&"macp.mode.task.v1".to_string()));
2046        assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
2047        assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
2048        // multi_round is now an extension, not in ListModes
2049        assert!(!names.contains(&"ext.multi_round.v1".to_string()));
2050    }
2051
2052    #[tokio::test]
2053    async fn list_ext_modes_returns_extensions() {
2054        let (server, _) = make_server();
2055        let resp = server
2056            .list_ext_modes(Request::new(ListExtModesRequest {}))
2057            .await
2058            .unwrap();
2059        let names: Vec<String> = resp
2060            .into_inner()
2061            .modes
2062            .iter()
2063            .map(|m| m.mode.clone())
2064            .collect();
2065        assert_eq!(names.len(), 1);
2066        assert!(names.contains(&"ext.multi_round.v1".to_string()));
2067    }
2068
2069    #[tokio::test]
2070    async fn get_manifest_includes_all_modes() {
2071        let (server, _) = make_server();
2072        let resp = server
2073            .get_manifest(Request::new(crate::pb::GetManifestRequest {
2074                agent_id: String::new(),
2075            }))
2076            .await
2077            .unwrap();
2078        let manifest = resp.into_inner().manifest.unwrap();
2079        assert_eq!(manifest.supported_modes.len(), 6);
2080        assert!(manifest
2081            .supported_modes
2082            .contains(&"ext.multi_round.v1".to_string()));
2083    }
2084
2085    #[tokio::test]
2086    async fn get_session_returns_metadata() {
2087        let (server, _) = make_server();
2088        let sid = new_sid();
2089        let ack = do_send(
2090            &server,
2091            "agent://orchestrator",
2092            Envelope {
2093                macp_version: "1.0".into(),
2094                mode: "macp.mode.decision.v1".into(),
2095                message_type: "SessionStart".into(),
2096                message_id: "m1".into(),
2097                session_id: sid.clone(),
2098                sender: String::new(),
2099                timestamp_unix_ms: Utc::now().timestamp_millis(),
2100                payload: start_payload(),
2101            },
2102        )
2103        .await;
2104        assert!(ack.ok);
2105
2106        let mut req = Request::new(GetSessionRequest {
2107            session_id: sid.clone(),
2108        });
2109        req.metadata_mut().insert(
2110            "authorization",
2111            format!("Bearer {}", "agent://orchestrator")
2112                .parse()
2113                .unwrap(),
2114        );
2115        let resp = server.get_session(req).await.unwrap();
2116        let meta = resp.into_inner().metadata.unwrap();
2117        assert_eq!(meta.session_id, sid);
2118        assert_eq!(meta.mode, "macp.mode.decision.v1");
2119        assert_eq!(meta.mode_version, "1.0.0");
2120        assert_eq!(meta.configuration_version, "cfg-1");
2121    }
2122
2123    #[tokio::test]
2124    async fn cancel_session_transitions_to_cancelled() {
2125        let (server, _) = make_server();
2126        let sid = new_sid();
2127        let ack = do_send(
2128            &server,
2129            "agent://orchestrator",
2130            Envelope {
2131                macp_version: "1.0".into(),
2132                mode: "macp.mode.decision.v1".into(),
2133                message_type: "SessionStart".into(),
2134                message_id: "m1".into(),
2135                session_id: sid.clone(),
2136                sender: String::new(),
2137                timestamp_unix_ms: Utc::now().timestamp_millis(),
2138                payload: start_payload(),
2139            },
2140        )
2141        .await;
2142        assert!(ack.ok);
2143
2144        let mut req = Request::new(CancelSessionRequest {
2145            session_id: sid,
2146            reason: "no longer needed".into(),
2147        });
2148        req.metadata_mut().insert(
2149            "authorization",
2150            format!("Bearer {}", "agent://orchestrator")
2151                .parse()
2152                .unwrap(),
2153        );
2154        let resp = server.cancel_session(req).await.unwrap();
2155        let ack = resp.into_inner().ack.unwrap();
2156        assert!(ack.ok);
2157        // RFC-MACP-0001 §7.3: cancellation now yields the distinct CANCELLED state.
2158        assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2159    }
2160
2161    #[tokio::test]
2162    async fn participant_cannot_cancel_session() {
2163        let (server, _) = make_server();
2164        let sid = new_sid();
2165        let ack = do_send(
2166            &server,
2167            "agent://orchestrator",
2168            Envelope {
2169                macp_version: "1.0".into(),
2170                mode: "macp.mode.decision.v1".into(),
2171                message_type: "SessionStart".into(),
2172                message_id: "m1".into(),
2173                session_id: sid.clone(),
2174                sender: String::new(),
2175                timestamp_unix_ms: Utc::now().timestamp_millis(),
2176                payload: start_payload(),
2177            },
2178        )
2179        .await;
2180        assert!(ack.ok);
2181
2182        let mut req = Request::new(CancelSessionRequest {
2183            session_id: sid,
2184            reason: "I want to cancel".into(),
2185        });
2186        req.metadata_mut().insert(
2187            "authorization",
2188            format!("Bearer {}", "agent://fraud").parse().unwrap(),
2189        );
2190        let err = server.cancel_session(req).await.unwrap_err();
2191        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2192    }
2193
2194    #[tokio::test]
2195    async fn cancel_session_unknown_session_returns_error() {
2196        let (server, _) = make_server();
2197        let mut req = Request::new(CancelSessionRequest {
2198            session_id: "nonexistent".into(),
2199            reason: "test".into(),
2200        });
2201        req.metadata_mut().insert(
2202            "authorization",
2203            format!("Bearer {}", "agent://orchestrator")
2204                .parse()
2205                .unwrap(),
2206        );
2207        let err = server.cancel_session(req).await.unwrap_err();
2208        assert_eq!(err.code(), tonic::Code::NotFound);
2209    }
2210
2211    #[tokio::test]
2212    async fn ambient_signal_accepted() {
2213        let (server, _) = make_server();
2214        let ack = do_send(
2215            &server,
2216            "agent://orchestrator",
2217            Envelope {
2218                macp_version: "1.0".into(),
2219                mode: String::new(),
2220                message_type: "Signal".into(),
2221                message_id: "sig-1".into(),
2222                session_id: String::new(),
2223                sender: String::new(),
2224                timestamp_unix_ms: Utc::now().timestamp_millis(),
2225                payload: vec![],
2226            },
2227        )
2228        .await;
2229        assert!(ack.ok);
2230    }
2231
2232    #[tokio::test]
2233    async fn signal_with_session_id_rejected() {
2234        let (server, _) = make_server();
2235        let ack = do_send(
2236            &server,
2237            "agent://orchestrator",
2238            Envelope {
2239                macp_version: "1.0".into(),
2240                mode: String::new(),
2241                message_type: "Signal".into(),
2242                message_id: "sig-2".into(),
2243                session_id: "some-session".into(),
2244                sender: String::new(),
2245                timestamp_unix_ms: Utc::now().timestamp_millis(),
2246                payload: vec![],
2247            },
2248        )
2249        .await;
2250        assert!(!ack.ok);
2251        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2252    }
2253
2254    #[tokio::test]
2255    async fn signal_with_mode_rejected() {
2256        let (server, _) = make_server();
2257        let ack = do_send(
2258            &server,
2259            "agent://orchestrator",
2260            Envelope {
2261                macp_version: "1.0".into(),
2262                mode: "macp.mode.decision.v1".into(),
2263                message_type: "Signal".into(),
2264                message_id: "sig-3".into(),
2265                session_id: String::new(),
2266                sender: String::new(),
2267                timestamp_unix_ms: Utc::now().timestamp_millis(),
2268                payload: vec![],
2269            },
2270        )
2271        .await;
2272        assert!(!ack.ok);
2273        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2274    }
2275
2276    #[tokio::test]
2277    async fn ambient_progress_accepted() {
2278        let (server, _) = make_server();
2279        let ack = do_send(
2280            &server,
2281            "agent://orchestrator",
2282            Envelope {
2283                macp_version: "1.0".into(),
2284                mode: String::new(),
2285                message_type: "Progress".into(),
2286                message_id: "prog-1".into(),
2287                session_id: String::new(),
2288                sender: String::new(),
2289                timestamp_unix_ms: Utc::now().timestamp_millis(),
2290                payload: vec![],
2291            },
2292        )
2293        .await;
2294        assert!(ack.ok);
2295    }
2296
2297    #[tokio::test]
2298    async fn ambient_progress_with_mode_rejected() {
2299        let (server, _) = make_server();
2300        let ack = do_send(
2301            &server,
2302            "agent://orchestrator",
2303            Envelope {
2304                macp_version: "1.0".into(),
2305                mode: "macp.mode.decision.v1".into(),
2306                message_type: "Progress".into(),
2307                message_id: "prog-2".into(),
2308                session_id: String::new(),
2309                sender: String::new(),
2310                timestamp_unix_ms: Utc::now().timestamp_millis(),
2311                payload: vec![],
2312            },
2313        )
2314        .await;
2315        assert!(!ack.ok);
2316        assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2317    }
2318
2319    #[tokio::test]
2320    async fn manifest_advertises_stream_enabled() {
2321        let (server, _) = make_server();
2322        let resp = server
2323            .initialize(Request::new(InitializeRequest {
2324                supported_protocol_versions: vec!["1.0".into()],
2325                client_info: None,
2326                capabilities: None,
2327            }))
2328            .await
2329            .unwrap();
2330        let caps = resp.into_inner().capabilities.unwrap();
2331        assert!(caps.sessions.unwrap().stream);
2332    }
2333
2334    #[tokio::test]
2335    async fn initialize_empty_versions_rejected() {
2336        let (server, _) = make_server();
2337        let err = server
2338            .initialize(Request::new(InitializeRequest {
2339                supported_protocol_versions: vec![],
2340                client_info: None,
2341                capabilities: None,
2342            }))
2343            .await
2344            .unwrap_err();
2345        assert_eq!(err.code(), tonic::Code::InvalidArgument);
2346    }
2347
2348    #[tokio::test]
2349    async fn initialize_unsupported_version_rejected() {
2350        let (server, _) = make_server();
2351        let err = server
2352            .initialize(Request::new(InitializeRequest {
2353                supported_protocol_versions: vec!["2.0".into()],
2354                client_info: None,
2355                capabilities: None,
2356            }))
2357            .await
2358            .unwrap_err();
2359        assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2360    }
2361
2362    // ── RFC-MACP-0006-A1: passive subscribe tests ──────────────────────
2363
2364    fn observer_identity(sender: &str) -> AuthIdentity {
2365        AuthIdentity {
2366            sender: sender.into(),
2367            allowed_modes: None,
2368            can_start_sessions: false,
2369            max_open_sessions: None,
2370            can_manage_mode_registry: false,
2371            is_observer: true,
2372        }
2373    }
2374
2375    fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2376        StreamSessionRequest {
2377            subscribe_session_id: session_id.into(),
2378            after_sequence: after,
2379            envelope: None,
2380        }
2381    }
2382
2383    fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2384        SessionStartPayload {
2385            intent: "intent".into(),
2386            participants,
2387            mode_version: "1.0.0".into(),
2388            configuration_version: "cfg-1".into(),
2389            policy_version: String::new(),
2390            ttl_ms: 60_000,
2391            context_id: String::new(),
2392            extensions: std::collections::HashMap::new(),
2393            roots: vec![],
2394            max_suspend_ms: 0,
2395        }
2396        .encode_to_vec()
2397    }
2398
2399    async fn start_session(
2400        server: &MacpServer,
2401        initiator: &str,
2402        sid: &str,
2403        participants: Vec<String>,
2404    ) {
2405        let ack = do_send(
2406            server,
2407            initiator,
2408            Envelope {
2409                macp_version: "1.0".into(),
2410                mode: "macp.mode.decision.v1".into(),
2411                message_type: "SessionStart".into(),
2412                message_id: "start".into(),
2413                session_id: sid.into(),
2414                sender: String::new(),
2415                timestamp_unix_ms: Utc::now().timestamp_millis(),
2416                payload: start_multi_participant(participants),
2417            },
2418        )
2419        .await;
2420        assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2421    }
2422
2423    async fn send_proposal(
2424        server: &MacpServer,
2425        sender: &str,
2426        sid: &str,
2427        message_id: &str,
2428        proposal_id: &str,
2429    ) {
2430        let payload = crate::decision_pb::ProposalPayload {
2431            proposal_id: proposal_id.into(),
2432            option: "opt".into(),
2433            rationale: "r".into(),
2434            supporting_data: vec![],
2435        }
2436        .encode_to_vec();
2437        let ack = do_send(
2438            server,
2439            sender,
2440            Envelope {
2441                macp_version: "1.0".into(),
2442                mode: "macp.mode.decision.v1".into(),
2443                message_type: "Proposal".into(),
2444                message_id: message_id.into(),
2445                session_id: sid.into(),
2446                sender: String::new(),
2447                timestamp_unix_ms: Utc::now().timestamp_millis(),
2448                payload,
2449            },
2450        )
2451        .await;
2452        assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2453    }
2454
2455    #[tokio::test]
2456    async fn subscribe_replays_session_history_from_zero() {
2457        let (server, _) = make_server();
2458        let sid = new_sid();
2459        let initiator = "agent://orchestrator";
2460        let peer = "agent://fraud";
2461        start_session(
2462            &server,
2463            initiator,
2464            &sid,
2465            vec![initiator.into(), peer.into()],
2466        )
2467        .await;
2468        send_proposal(&server, peer, &sid, "m2", "p1").await;
2469
2470        let mut bound = None;
2471        let mut events = None;
2472        let replay = server
2473            .process_stream_request(
2474                &stream_identity(peer),
2475                subscribe_frame(&sid, 0),
2476                &mut bound,
2477                &mut events,
2478            )
2479            .await
2480            .unwrap();
2481
2482        assert_eq!(replay.len(), 2);
2483        assert_eq!(replay[0].message_type, "SessionStart");
2484        assert_eq!(replay[0].message_id, "start");
2485        assert_eq!(replay[1].message_type, "Proposal");
2486        assert_eq!(replay[1].message_id, "m2");
2487        assert_eq!(bound.as_deref(), Some(sid.as_str()));
2488        assert!(events.is_some());
2489    }
2490
2491    #[tokio::test]
2492    async fn subscribe_after_sequence_filters_history() {
2493        let (server, _) = make_server();
2494        let sid = new_sid();
2495        let initiator = "agent://orchestrator";
2496        let peer = "agent://fraud";
2497        start_session(
2498            &server,
2499            initiator,
2500            &sid,
2501            vec![initiator.into(), peer.into()],
2502        )
2503        .await;
2504        send_proposal(&server, peer, &sid, "m2", "p1").await;
2505        send_proposal(&server, peer, &sid, "m3", "p2").await;
2506
2507        let mut bound = None;
2508        let mut events = None;
2509        let replay = server
2510            .process_stream_request(
2511                &stream_identity(peer),
2512                subscribe_frame(&sid, 2),
2513                &mut bound,
2514                &mut events,
2515            )
2516            .await
2517            .unwrap();
2518
2519        assert_eq!(replay.len(), 1);
2520        assert_eq!(replay[0].message_id, "m3");
2521    }
2522
2523    #[tokio::test]
2524    async fn subscribe_unknown_session_returns_not_found() {
2525        let (server, _) = make_server();
2526        let mut bound = None;
2527        let mut events = None;
2528        let status = server
2529            .process_stream_request(
2530                &stream_identity("agent://orchestrator"),
2531                subscribe_frame("missing-session", 0),
2532                &mut bound,
2533                &mut events,
2534            )
2535            .await
2536            .unwrap_err();
2537        assert_eq!(status.code(), tonic::Code::NotFound);
2538        assert!(bound.is_none());
2539        assert!(events.is_none());
2540    }
2541
2542    #[tokio::test]
2543    async fn subscribe_non_participant_is_forbidden() {
2544        let (server, _) = make_server();
2545        let sid = new_sid();
2546        start_session(
2547            &server,
2548            "agent://orchestrator",
2549            &sid,
2550            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2551        )
2552        .await;
2553
2554        let mut bound = None;
2555        let mut events = None;
2556        let status = server
2557            .process_stream_request(
2558                &stream_identity("agent://outsider"),
2559                subscribe_frame(&sid, 0),
2560                &mut bound,
2561                &mut events,
2562            )
2563            .await
2564            .unwrap_err();
2565        assert_eq!(status.code(), tonic::Code::PermissionDenied);
2566    }
2567
2568    #[tokio::test]
2569    async fn subscribe_observer_identity_allowed() {
2570        let (server, _) = make_server();
2571        let sid = new_sid();
2572        start_session(
2573            &server,
2574            "agent://orchestrator",
2575            &sid,
2576            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2577        )
2578        .await;
2579
2580        let mut bound = None;
2581        let mut events = None;
2582        let replay = server
2583            .process_stream_request(
2584                &observer_identity("agent://auditor"),
2585                subscribe_frame(&sid, 0),
2586                &mut bound,
2587                &mut events,
2588            )
2589            .await
2590            .unwrap();
2591        assert_eq!(replay.len(), 1);
2592        assert_eq!(replay[0].message_type, "SessionStart");
2593    }
2594
2595    #[tokio::test]
2596    async fn subscribe_initiator_allowed_even_when_not_listed() {
2597        // Per RFC-MACP-0007, the initiator is always authorized for session
2598        // access, even if not present in the participants list.
2599        let (server, _) = make_server();
2600        let sid = new_sid();
2601        start_session(
2602            &server,
2603            "agent://orchestrator",
2604            &sid,
2605            vec!["agent://fraud".into()],
2606        )
2607        .await;
2608
2609        let mut bound = None;
2610        let mut events = None;
2611        let replay = server
2612            .process_stream_request(
2613                &stream_identity("agent://orchestrator"),
2614                subscribe_frame(&sid, 0),
2615                &mut bound,
2616                &mut events,
2617            )
2618            .await
2619            .unwrap();
2620        assert_eq!(replay.len(), 1);
2621    }
2622
2623    #[tokio::test]
2624    async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2625        let (server, _) = make_server();
2626        let sid = new_sid();
2627        let req = StreamSessionRequest {
2628            subscribe_session_id: sid.clone(),
2629            after_sequence: 0,
2630            envelope: Some(Envelope {
2631                macp_version: "1.0".into(),
2632                mode: "macp.mode.decision.v1".into(),
2633                message_type: "SessionStart".into(),
2634                message_id: "m1".into(),
2635                session_id: sid,
2636                sender: String::new(),
2637                timestamp_unix_ms: Utc::now().timestamp_millis(),
2638                payload: start_payload(),
2639            }),
2640        };
2641
2642        let mut bound = None;
2643        let mut events = None;
2644        let status = server
2645            .process_stream_request(
2646                &stream_identity("agent://orchestrator"),
2647                req,
2648                &mut bound,
2649                &mut events,
2650            )
2651            .await
2652            .unwrap_err();
2653        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2654    }
2655
2656    #[tokio::test]
2657    async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2658        let (server, _) = make_server();
2659        let sid1 = new_sid();
2660        let sid2 = new_sid();
2661        start_session(
2662            &server,
2663            "agent://orchestrator",
2664            &sid1,
2665            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2666        )
2667        .await;
2668        start_session(
2669            &server,
2670            "agent://orchestrator",
2671            &sid2,
2672            vec!["agent://orchestrator".into(), "agent://fraud".into()],
2673        )
2674        .await;
2675
2676        // First subscribe binds the stream to sid1
2677        let identity = stream_identity("agent://fraud");
2678        let mut bound = None;
2679        let mut events = None;
2680        server
2681            .process_stream_request(
2682                &identity,
2683                subscribe_frame(&sid1, 0),
2684                &mut bound,
2685                &mut events,
2686            )
2687            .await
2688            .unwrap();
2689        assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2690
2691        // Second subscribe to sid2 on the same stream must be rejected
2692        let status = server
2693            .process_stream_request(
2694                &identity,
2695                subscribe_frame(&sid2, 0),
2696                &mut bound,
2697                &mut events,
2698            )
2699            .await
2700            .unwrap_err();
2701        assert_eq!(status.code(), tonic::Code::InvalidArgument);
2702    }
2703
2704    /// E3: an injected ingress engine gates session start, messages, and
2705    /// session reads — deny-one-sender double proves all three hooks fire and
2706    /// that denial surfaces as POLICY_DENIED / PermissionDenied (fail closed).
2707    struct DenySenderEngine {
2708        denied: String,
2709    }
2710
2711    #[async_trait::async_trait]
2712    impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2713        async fn evaluate_session_start(
2714            &self,
2715            identity: &crate::security::AuthIdentity,
2716            _mode: &str,
2717            _env: &Envelope,
2718        ) -> macp_core::policy::PolicyDecision {
2719            if identity.sender == self.denied {
2720                macp_core::policy::PolicyDecision::Deny {
2721                    reasons: vec!["sender embargoed".into()],
2722                }
2723            } else {
2724                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2725            }
2726        }
2727
2728        async fn evaluate_message(
2729            &self,
2730            identity: &crate::security::AuthIdentity,
2731            _session: &macp_core::session::Session,
2732            _env: &Envelope,
2733        ) -> macp_core::policy::PolicyDecision {
2734            if identity.sender == self.denied {
2735                macp_core::policy::PolicyDecision::Deny {
2736                    reasons: vec!["sender embargoed".into()],
2737                }
2738            } else {
2739                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2740            }
2741        }
2742
2743        async fn evaluate_session_access(
2744            &self,
2745            identity: &crate::security::AuthIdentity,
2746            _session: &macp_core::session::Session,
2747        ) -> macp_core::policy::PolicyDecision {
2748            if identity.sender == self.denied {
2749                macp_core::policy::PolicyDecision::Deny {
2750                    reasons: vec!["sender embargoed".into()],
2751                }
2752            } else {
2753                macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2754            }
2755        }
2756    }
2757
2758    #[tokio::test]
2759    async fn policy_engine_gates_all_three_ingress_points() {
2760        let (server, _runtime) = make_server();
2761        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2762            denied: "agent://embargoed".into(),
2763        }));
2764
2765        let sid = new_sid();
2766        let start_payload = SessionStartPayload {
2767            intent: "e3".into(),
2768            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2769            mode_version: "1.0.0".into(),
2770            configuration_version: "cfg-1".into(),
2771            policy_version: String::new(),
2772            ttl_ms: 60_000,
2773            context_id: String::new(),
2774            extensions: Default::default(),
2775            roots: vec![],
2776            max_suspend_ms: 0,
2777        }
2778        .encode_to_vec();
2779        let start_env = |sender: &str, sid: &str| Envelope {
2780            macp_version: "1.0".into(),
2781            mode: "macp.mode.decision.v1".into(),
2782            message_type: "SessionStart".into(),
2783            message_id: new_sid(),
2784            session_id: sid.into(),
2785            sender: sender.into(),
2786            timestamp_unix_ms: Utc::now().timestamp_millis(),
2787            payload: start_payload.clone(),
2788        };
2789
2790        // 1. Embargoed sender cannot start a session.
2791        let ack = server
2792            .send(send_req(
2793                "agent://embargoed",
2794                start_env("agent://embargoed", &sid),
2795            ))
2796            .await
2797            .unwrap()
2798            .into_inner()
2799            .ack
2800            .unwrap();
2801        assert!(!ack.ok);
2802        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2803
2804        // Allowed sender starts it.
2805        let ack = server
2806            .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2807            .await
2808            .unwrap()
2809            .into_inner()
2810            .ack
2811            .unwrap();
2812        assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2813
2814        // 2. Embargoed sender cannot send into the session.
2815        let proposal = crate::decision_pb::ProposalPayload {
2816            proposal_id: "p1".into(),
2817            option: "x".into(),
2818            rationale: "r".into(),
2819            supporting_data: vec![],
2820        }
2821        .encode_to_vec();
2822        let msg_env = Envelope {
2823            macp_version: "1.0".into(),
2824            mode: "macp.mode.decision.v1".into(),
2825            message_type: "Proposal".into(),
2826            message_id: new_sid(),
2827            session_id: sid.clone(),
2828            sender: "agent://embargoed".into(),
2829            timestamp_unix_ms: Utc::now().timestamp_millis(),
2830            payload: proposal,
2831        };
2832        let ack = server
2833            .send(send_req("agent://embargoed", msg_env))
2834            .await
2835            .unwrap()
2836            .into_inner()
2837            .ack
2838            .unwrap();
2839        assert!(!ack.ok);
2840        assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2841
2842        // 3. Embargoed sender cannot read the session.
2843        let mut req = Request::new(crate::pb::GetSessionRequest {
2844            session_id: sid.clone(),
2845        });
2846        req.metadata_mut()
2847            .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2848        let err = server
2849            .get_session(req)
2850            .await
2851            .expect_err("embargoed read must be denied");
2852        assert_eq!(err.code(), tonic::Code::PermissionDenied);
2853    }
2854
2855    /// E3 transport-parity: the ingress engine gates the STREAM path too — a
2856    /// denied sender must not be able to bypass the engine by switching from
2857    /// unary Send to StreamSession (envelope frames or subscribe frames).
2858    #[tokio::test]
2859    async fn policy_engine_gates_stream_path() {
2860        let (server, runtime) = make_server();
2861        let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2862            denied: "agent://embargoed".into(),
2863        }));
2864
2865        // Session started by an allowed sender (participants include the
2866        // embargoed agent so built-in membership checks pass — only the
2867        // engine denies it).
2868        let sid = new_sid();
2869        let payload = SessionStartPayload {
2870            intent: "e3-stream".into(),
2871            participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2872            mode_version: "1.0.0".into(),
2873            configuration_version: "cfg-1".into(),
2874            policy_version: String::new(),
2875            ttl_ms: 60_000,
2876            context_id: String::new(),
2877            extensions: Default::default(),
2878            roots: vec![],
2879            max_suspend_ms: 0,
2880        }
2881        .encode_to_vec();
2882        runtime
2883            .process(
2884                &Envelope {
2885                    macp_version: "1.0".into(),
2886                    mode: "macp.mode.decision.v1".into(),
2887                    message_type: "SessionStart".into(),
2888                    message_id: new_sid(),
2889                    session_id: sid.clone(),
2890                    sender: "agent://ok".into(),
2891                    timestamp_unix_ms: Utc::now().timestamp_millis(),
2892                    payload,
2893                },
2894                None,
2895            )
2896            .await
2897            .unwrap();
2898
2899        let embargoed = crate::security::AuthIdentity {
2900            sender: "agent://embargoed".into(),
2901            allowed_modes: None,
2902            can_start_sessions: true,
2903            max_open_sessions: None,
2904            can_manage_mode_registry: false,
2905            is_observer: false,
2906        };
2907        let mut bound = None;
2908        let mut events = None;
2909
2910        // 1. Stream envelope frame from the embargoed sender: denied.
2911        let proposal = crate::decision_pb::ProposalPayload {
2912            proposal_id: "p1".into(),
2913            option: "x".into(),
2914            rationale: "r".into(),
2915            supporting_data: vec![],
2916        }
2917        .encode_to_vec();
2918        let req = StreamSessionRequest {
2919            envelope: Some(Envelope {
2920                macp_version: "1.0".into(),
2921                mode: "macp.mode.decision.v1".into(),
2922                message_type: "Proposal".into(),
2923                message_id: new_sid(),
2924                session_id: sid.clone(),
2925                sender: "agent://embargoed".into(),
2926                timestamp_unix_ms: Utc::now().timestamp_millis(),
2927                payload: proposal,
2928            }),
2929            subscribe_session_id: String::new(),
2930            after_sequence: 0,
2931        };
2932        let err = server
2933            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2934            .await
2935            .expect_err("stream envelope from embargoed sender must be denied");
2936        // PolicyDenied maps to FailedPrecondition on the transport (same
2937        // error the unary path expresses as a POLICY_DENIED ack).
2938        assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2939        assert!(err.message().contains("PolicyDenied"), "{err:?}");
2940
2941        // 2. Passive-subscribe frame (history read) from the embargoed
2942        //    sender: denied even though membership would allow it.
2943        let req = StreamSessionRequest {
2944            envelope: None,
2945            subscribe_session_id: sid.clone(),
2946            after_sequence: 0,
2947        };
2948        let err = server
2949            .process_stream_request(&embargoed, req, &mut bound, &mut events)
2950            .await
2951            .expect_err("stream subscribe from embargoed sender must be denied");
2952        assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2953    }
2954    // ── ListSessions pagination (core.proto:411-426) ───────────────────
2955
2956    fn paged_session(id: &str) -> crate::session::Session {
2957        crate::session::Session::builder(id, "macp.mode.decision.v1", "agent://initiator")
2958            .participants(vec!["agent://a".into()])
2959            .mode_version("1.0.0")
2960            .configuration_version("cfg-1")
2961            .started_at_unix_ms(1)
2962            .build()
2963    }
2964
2965    /// Insert sessions straight into the registry, with each `Session`'s
2966    /// `session_id` equal to its map key (the Phase 1 `debug_assert_eq!`
2967    /// enforces the pair).
2968    async fn seed_sessions(runtime: &Arc<Runtime>, ids: &[String]) {
2969        for id in ids {
2970            runtime
2971                .registry
2972                .insert_recovered_session(id.clone(), paged_session(id))
2973                .await;
2974        }
2975    }
2976
2977    fn list_sessions_req(page_size: i32, page_token: &str) -> Request<ListSessionsRequest> {
2978        let mut req = Request::new(ListSessionsRequest {
2979            page_size,
2980            page_token: page_token.to_string(),
2981        });
2982        req.metadata_mut()
2983            .insert("authorization", "Bearer agent://observer".parse().unwrap());
2984        req
2985    }
2986
2987    fn page_size_security(default: usize, max: usize) -> SecurityLayer {
2988        let mut security = SecurityLayer::dev_mode();
2989        security.list_sessions_default_page_size = default;
2990        security.list_sessions_max_page_size = max;
2991        security
2992    }
2993
2994    fn seed_ids(n: usize) -> Vec<String> {
2995        (0..n).map(|i| format!("session-{i:03}")).collect()
2996    }
2997
2998    #[tokio::test]
2999    async fn list_sessions_applies_default_page_size_when_zero() {
3000        let (server, runtime) = make_server_with_security(page_size_security(3, 1000));
3001        seed_sessions(&runtime, &seed_ids(10)).await;
3002
3003        let resp = server
3004            .list_sessions(list_sessions_req(0, ""))
3005            .await
3006            .unwrap()
3007            .into_inner();
3008        assert_eq!(resp.sessions.len(), 3);
3009        assert!(!resp.next_page_token.is_empty());
3010    }
3011
3012    #[tokio::test]
3013    async fn list_sessions_honors_explicit_page_size() {
3014        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3015        seed_sessions(&runtime, &seed_ids(10)).await;
3016
3017        let resp = server
3018            .list_sessions(list_sessions_req(4, ""))
3019            .await
3020            .unwrap()
3021            .into_inner();
3022        assert_eq!(resp.sessions.len(), 4);
3023        assert!(!resp.next_page_token.is_empty());
3024    }
3025
3026    #[tokio::test]
3027    async fn list_sessions_clamps_page_size_above_max() {
3028        let (server, runtime) = make_server_with_security(page_size_security(100, 3));
3029        seed_sessions(&runtime, &seed_ids(10)).await;
3030
3031        let resp = server
3032            .list_sessions(list_sessions_req(1000, ""))
3033            .await
3034            .unwrap()
3035            .into_inner();
3036        assert_eq!(resp.sessions.len(), 3);
3037        assert!(!resp.next_page_token.is_empty());
3038    }
3039
3040    #[tokio::test]
3041    async fn list_sessions_rejects_negative_page_size() {
3042        let (server, runtime) = make_server();
3043        seed_sessions(&runtime, &seed_ids(3)).await;
3044
3045        let err = server
3046            .list_sessions(list_sessions_req(-1, ""))
3047            .await
3048            .unwrap_err();
3049        assert_eq!(err.code(), tonic::Code::InvalidArgument, "{err:?}");
3050        assert!(err.message().contains("page_size"), "{err:?}");
3051    }
3052
3053    #[tokio::test]
3054    async fn list_sessions_rejects_garbage_page_token() {
3055        use base64::Engine;
3056        let (server, runtime) = make_server();
3057        seed_sessions(&runtime, &seed_ids(3)).await;
3058
3059        let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD;
3060        let valid = engine.encode("v1:session-000");
3061        let tokens = vec![
3062            // not base64url
3063            "not-a-token!".to_string(),
3064            // wrong version prefix
3065            engine.encode("v2:session-000"),
3066            // prefix present, cursor empty
3067            engine.encode("v1:"),
3068            // truncated *through* the version prefix. Note that lopping bytes
3069            // off the end of an encoded token instead yields a valid, shorter
3070            // cursor — harmless, since a cursor is a position, not a handle —
3071            // so the truncation that must be rejected is the one that damages
3072            // the prefix.
3073            engine.encode("v1"),
3074            // front-truncated: the leading base64 character is gone, so the
3075            // decoded bytes are no longer valid UTF-8 (and could not carry the
3076            // prefix regardless).
3077            valid[1..].to_string(),
3078            // oversized: rejected by the length branch, before any decode
3079            "A".repeat(2 * 1024 * 1024),
3080        ];
3081        for token in tokens {
3082            let err = server
3083                .list_sessions(list_sessions_req(0, &token))
3084                .await
3085                .unwrap_err();
3086            assert_eq!(err.code(), tonic::Code::InvalidArgument);
3087            // One opaque message for every rejection reason.
3088            assert_eq!(
3089                err.message(),
3090                "INVALID_ARGUMENT: page_token is not a valid continuation token"
3091            );
3092        }
3093    }
3094
3095    #[tokio::test]
3096    async fn list_sessions_full_traversal_visits_every_session_exactly_once() {
3097        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3098        let ids = seed_ids(25);
3099        seed_sessions(&runtime, &ids).await;
3100
3101        let mut collected: Vec<String> = Vec::new();
3102        let mut token = String::new();
3103        for _ in 0..100 {
3104            let resp = server
3105                .list_sessions(list_sessions_req(4, &token))
3106                .await
3107                .unwrap()
3108                .into_inner();
3109            collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3110            token = resp.next_page_token;
3111            if token.is_empty() {
3112                break;
3113            }
3114        }
3115        assert!(token.is_empty(), "traversal did not terminate");
3116        let unique: std::collections::HashSet<&String> = collected.iter().collect();
3117        // Both assertions: the set alone would hide duplicates, the total
3118        // alone would hide a duplicate paired with a drop.
3119        assert_eq!(unique.len(), 25, "sessions were dropped or duplicated");
3120        assert_eq!(collected.len(), 25, "sessions were duplicated");
3121    }
3122
3123    #[tokio::test]
3124    async fn list_sessions_terminal_page_has_empty_next_page_token() {
3125        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3126        seed_sessions(&runtime, &seed_ids(10)).await;
3127
3128        let mut tokens: Vec<String> = Vec::new();
3129        let mut token = String::new();
3130        for _ in 0..20 {
3131            let resp = server
3132                .list_sessions(list_sessions_req(5, &token))
3133                .await
3134                .unwrap()
3135                .into_inner();
3136            token = resp.next_page_token;
3137            tokens.push(token.clone());
3138            if token.is_empty() {
3139                break;
3140            }
3141        }
3142        // 10 sessions at 5/page: exactly two pages, and only the last one
3143        // carries the empty token.
3144        assert_eq!(tokens.len(), 2, "{tokens:?}");
3145        assert!(!tokens[0].is_empty());
3146        assert!(tokens[1].is_empty());
3147    }
3148
3149    #[tokio::test]
3150    async fn list_sessions_orders_by_session_id_ascending() {
3151        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3152        // Insertion order deliberately unrelated to sort order.
3153        let ids: Vec<String> = ["delta", "alpha", "echo", "charlie", "bravo"]
3154            .iter()
3155            .map(|s| s.to_string())
3156            .collect();
3157        seed_sessions(&runtime, &ids).await;
3158
3159        let mut collected: Vec<String> = Vec::new();
3160        let mut token = String::new();
3161        loop {
3162            let resp = server
3163                .list_sessions(list_sessions_req(2, &token))
3164                .await
3165                .unwrap()
3166                .into_inner();
3167            collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3168            token = resp.next_page_token;
3169            if token.is_empty() {
3170                break;
3171            }
3172        }
3173        // Ascending across the whole traversal, not merely within a page.
3174        assert_eq!(
3175            collected,
3176            vec!["alpha", "bravo", "charlie", "delta", "echo"]
3177        );
3178    }
3179
3180    #[tokio::test]
3181    async fn list_sessions_still_requires_authentication() {
3182        let (server, runtime) = make_server();
3183        seed_sessions(&runtime, &seed_ids(3)).await;
3184
3185        // No authorization metadata, and a request body that would otherwise
3186        // be INVALID_ARGUMENT: authentication must still be what answers.
3187        let req = Request::new(ListSessionsRequest {
3188            page_size: -1,
3189            page_token: String::new(),
3190        });
3191        let err = server.list_sessions(req).await.unwrap_err();
3192        assert_eq!(err.code(), tonic::Code::Unauthenticated, "{err:?}");
3193    }
3194
3195    #[tokio::test]
3196    async fn list_sessions_tolerates_cursor_for_removed_session() {
3197        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3198        let ids = seed_ids(4);
3199        seed_sessions(&runtime, &ids).await;
3200
3201        let first = server
3202            .list_sessions(list_sessions_req(1, ""))
3203            .await
3204            .unwrap()
3205            .into_inner();
3206        assert_eq!(first.sessions[0].session_id, "session-000");
3207        assert!(!first.next_page_token.is_empty());
3208
3209        // Delete the very session the cursor names. A keyset cursor is a
3210        // position, not a handle, so paging must continue undisturbed.
3211        runtime
3212            .registry
3213            .sessions
3214            .write()
3215            .await
3216            .remove("session-000");
3217
3218        let second = server
3219            .list_sessions(list_sessions_req(1, &first.next_page_token))
3220            .await
3221            .unwrap()
3222            .into_inner();
3223        assert_eq!(second.sessions[0].session_id, "session-001");
3224    }
3225
3226    #[tokio::test]
3227    async fn list_sessions_cursor_comes_from_the_id_list_not_the_returned_sessions() {
3228        // The cursor must be the last *candidate ID*, not the last *returned
3229        // session*. The two differ only when the registry is mutated between
3230        // the ID scan and the per-ID fetch, so this test manufactures exactly
3231        // that window: the handler parks on the first page entry's session
3232        // mutex, and while it is parked the last page entry is removed.
3233        //
3234        // Deriving the cursor from the returned sessions instead would move it
3235        // backwards, re-scanning IDs the page already accounted for — and, when
3236        // an entire page vanishes, would emit an empty token and silently
3237        // terminate the traversal, dropping every remaining session.
3238        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3239        seed_sessions(&runtime, &seed_ids(6)).await;
3240
3241        // Hold the first page entry's session mutex: the handler's fetch loop
3242        // parks there, which is the only deterministic yield point between the
3243        // ID scan and the rest of the fetches.
3244        let first = runtime.registry.get_shared("session-000").await.unwrap();
3245        let guard = first.lock().await;
3246
3247        let handler = server.list_sessions(list_sessions_req(3, ""));
3248        let mutator = async {
3249            // The handler clones the Arc in `get_shared` before parking on the
3250            // mutex, so a strong count of 3 (map + this test + handler) means
3251            // it is parked. Bounded so a missed interleave fails loudly rather
3252            // than hanging.
3253            let mut spins = 0;
3254            while Arc::strong_count(&first) < 3 {
3255                assert!(spins < 10_000, "handler never parked on the session mutex");
3256                spins += 1;
3257                tokio::task::yield_now().await;
3258            }
3259            runtime
3260                .registry
3261                .sessions
3262                .write()
3263                .await
3264                .remove("session-002");
3265            drop(guard);
3266        };
3267        let (resp, ()) = tokio::join!(handler, mutator);
3268        let resp = resp.unwrap().into_inner();
3269
3270        // The interleave actually happened: the last candidate was skipped.
3271        assert_eq!(
3272            resp.sessions.len(),
3273            2,
3274            "expected session-002 to vanish between the scan and the fetch"
3275        );
3276        assert_eq!(resp.sessions[1].session_id, "session-001");
3277        // ...yet the cursor is the last candidate, not the last survivor.
3278        assert_eq!(
3279            crate::pagination::decode_page_token(&resp.next_page_token),
3280            Ok("session-002".to_string()),
3281            "cursor was derived from the returned sessions, not the ID list"
3282        );
3283
3284        // Observable through the API too: put session-002 back and page on.
3285        runtime
3286            .registry
3287            .insert_recovered_session("session-002".to_string(), paged_session("session-002"))
3288            .await;
3289        let second = server
3290            .list_sessions(list_sessions_req(3, &resp.next_page_token))
3291            .await
3292            .unwrap()
3293            .into_inner();
3294        assert_eq!(
3295            second.sessions[0].session_id, "session-003",
3296            "the cursor moved backwards past an ID the page had already accounted for"
3297        );
3298    }
3299
3300    #[tokio::test]
3301    async fn list_sessions_replaying_a_token_returns_the_identical_page() {
3302        let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3303        seed_sessions(&runtime, &seed_ids(10)).await;
3304
3305        let first = server
3306            .list_sessions(list_sessions_req(3, ""))
3307            .await
3308            .unwrap()
3309            .into_inner();
3310        let token = first.next_page_token;
3311        assert!(!token.is_empty());
3312
3313        let page_a = server
3314            .list_sessions(list_sessions_req(3, &token))
3315            .await
3316            .unwrap()
3317            .into_inner();
3318        let page_b = server
3319            .list_sessions(list_sessions_req(3, &token))
3320            .await
3321            .unwrap()
3322            .into_inner();
3323
3324        let ids_a: Vec<&str> = page_a.sessions.iter().map(|s| &*s.session_id).collect();
3325        let ids_b: Vec<&str> = page_b.sessions.iter().map(|s| &*s.session_id).collect();
3326        assert_eq!(ids_a, ids_b);
3327        assert_eq!(page_a.next_page_token, page_b.next_page_token);
3328    }
3329
3330    #[tokio::test]
3331    async fn list_sessions_survives_zero_effective_page_size() {
3332        // Both fields are `pub`, so a consumer can reach 0. The floor in the
3333        // handler must keep the response well-formed: never an empty page
3334        // paired with a non-empty token (which would never advance).
3335        let (server, runtime) = make_server_with_security(page_size_security(0, 0));
3336        seed_sessions(&runtime, &seed_ids(3)).await;
3337
3338        let resp = server
3339            .list_sessions(list_sessions_req(0, ""))
3340            .await
3341            .unwrap()
3342            .into_inner();
3343        assert!(
3344            !resp.sessions.is_empty(),
3345            "empty page with token {:?} — the traversal terminates and ListSessions returns nothing",
3346            resp.next_page_token
3347        );
3348        assert_eq!(resp.sessions.len(), 1);
3349        assert!(!resp.next_page_token.is_empty());
3350
3351        // And it actually advances.
3352        let next = server
3353            .list_sessions(list_sessions_req(0, &resp.next_page_token))
3354            .await
3355            .unwrap()
3356            .into_inner();
3357        assert_eq!(next.sessions.len(), 1);
3358        assert_ne!(next.sessions[0].session_id, resp.sessions[0].session_id);
3359    }
3360
3361    fn watch_sessions_req(sender: &str) -> Request<WatchSessionsRequest> {
3362        let mut req = Request::new(WatchSessionsRequest {});
3363        req.metadata_mut()
3364            .insert("authorization", format!("Bearer {sender}").parse().unwrap());
3365        req
3366    }
3367
3368    /// Read the next event, failing (rather than hanging) if none arrives.
3369    async fn next_lifecycle_event(
3370        stream: &mut <MacpServer as MacpRuntimeService>::WatchSessionsStream,
3371    ) -> crate::pb::SessionLifecycleEvent {
3372        use tokio_stream::StreamExt;
3373        let resp = tokio::time::timeout(std::time::Duration::from_secs(5), stream.next())
3374            .await
3375            .expect("WatchSessions produced no event within 5s")
3376            .expect("stream ended")
3377            .expect("stream errored");
3378        resp.event.expect("event present")
3379    }
3380
3381    /// Criterion 1 through the handler: N sessions in the registry produce
3382    /// exactly N `Created` events, one per session, now that the sync
3383    /// materializes them one at a time instead of deep-cloning the registry.
3384    #[tokio::test]
3385    async fn watch_sessions_initial_sync_emits_each_session_exactly_once() {
3386        let (server, runtime) = make_server();
3387        let ids = seed_ids(24);
3388        seed_sessions(&runtime, &ids).await;
3389
3390        let mut stream = server
3391            .watch_sessions(watch_sessions_req("agent://observer"))
3392            .await
3393            .unwrap()
3394            .into_inner();
3395
3396        let mut counts: HashMap<String, usize> = HashMap::new();
3397        for _ in 0..ids.len() {
3398            let event = next_lifecycle_event(&mut stream).await;
3399            assert_eq!(
3400                event.event_type,
3401                session_lifecycle_event::EventType::Created as i32
3402            );
3403            let session = event.session.expect("initial sync always carries metadata");
3404            *counts.entry(session.session_id).or_default() += 1;
3405        }
3406        assert_eq!(counts.len(), ids.len(), "sync emitted the wrong set");
3407        for id in &ids {
3408            assert_eq!(
3409                counts.get(id).copied(),
3410                Some(1),
3411                "{id} was not emitted exactly once"
3412            );
3413        }
3414    }
3415
3416    /// Criterion 4: the lifecycle subscription is taken during the unary call,
3417    /// not lazily inside the generator.
3418    ///
3419    /// The generator is not polled until the client reads, so this drives a
3420    /// session terminal *after* the response is returned but *before* the first
3421    /// read. The `Cancelled` event is published while nothing is polling the
3422    /// stream — it can only be delivered because the subscription already
3423    /// existed. Move `subscribe_session_lifecycle()` inside the `try_stream!`
3424    /// and the second read here finds nothing and times out.
3425    #[tokio::test]
3426    async fn watch_sessions_subscribes_before_the_generator_is_polled() {
3427        let (server, _runtime) = make_server();
3428        let initiator = "agent://orchestrator";
3429        let sid = new_sid();
3430        start_session(&server, initiator, &sid, vec![initiator.into()]).await;
3431
3432        let mut stream = server
3433            .watch_sessions(watch_sessions_req("agent://observer"))
3434            .await
3435            .unwrap()
3436            .into_inner();
3437
3438        // Nothing has polled the stream yet; this event has only the
3439        // subscription taken above to land in.
3440        let mut cancel = Request::new(CancelSessionRequest {
3441            session_id: sid.clone(),
3442            reason: "test".into(),
3443        });
3444        cancel.metadata_mut().insert(
3445            "authorization",
3446            format!("Bearer {initiator}").parse().unwrap(),
3447        );
3448        let ack = server
3449            .cancel_session(cancel)
3450            .await
3451            .unwrap()
3452            .into_inner()
3453            .ack
3454            .unwrap();
3455        assert!(ack.ok);
3456
3457        // The initial sync comes first...
3458        let first = next_lifecycle_event(&mut stream).await;
3459        assert_eq!(
3460            first.event_type,
3461            session_lifecycle_event::EventType::Created as i32
3462        );
3463        assert_eq!(first.session.unwrap().session_id, sid);
3464
3465        // ...then the event buffered while the generator was still unpolled.
3466        let second = next_lifecycle_event(&mut stream).await;
3467        assert_eq!(
3468            second.event_type,
3469            session_lifecycle_event::EventType::Cancelled as i32,
3470            "the event published before the first poll was lost — the \
3471             subscription must be taken in the unary call"
3472        );
3473        assert_eq!(second.session.unwrap().session_id, sid);
3474    }
3475
3476    /// Both arms of the `Created` dedup, with the race that makes it necessary
3477    /// driven deterministically.
3478    ///
3479    /// `process_session_start` registers the session BEFORE it publishes
3480    /// `Created`, so a session started after the subscribe but before the first
3481    /// poll is in the sync snapshot *and* has a `Created` sitting on the bus.
3482    /// The sync must emit it once and the buffered copy must be dropped. A
3483    /// session started after the sync is not in the snapshot, and its live
3484    /// `Created` must pass through — the handler tests membership of the sync
3485    /// set with `contains` rather than `insert`, so that pass-through does not
3486    /// extend the set; the set stays bounded by the registry size at subscribe
3487    /// time. That bound is structural and has no observable signal, so what is
3488    /// asserted here is the exactly-once contract it must not break.
3489    #[tokio::test]
3490    async fn watch_sessions_emits_created_once_for_synced_and_live_sessions() {
3491        let (server, _runtime) = make_server();
3492        let initiator = "agent://orchestrator";
3493        let synced_sid = new_sid();
3494
3495        let mut stream = server
3496            .watch_sessions(watch_sessions_req("agent://observer"))
3497            .await
3498            .unwrap()
3499            .into_inner();
3500
3501        // Registered and published while nothing is polling: this lands in the
3502        // subscription AND in the snapshot the first poll takes.
3503        start_session(&server, initiator, &synced_sid, vec![initiator.into()]).await;
3504
3505        let from_sync = next_lifecycle_event(&mut stream).await;
3506        assert_eq!(
3507            from_sync.event_type,
3508            session_lifecycle_event::EventType::Created as i32
3509        );
3510        assert_eq!(from_sync.session.unwrap().session_id, synced_sid);
3511
3512        // Started after the sync, so it is absent from the snapshot.
3513        let live_sid = new_sid();
3514        start_session(&server, initiator, &live_sid, vec![initiator.into()]).await;
3515
3516        // The next event must be the live session's Created. If the buffered
3517        // duplicate leaked through, this is `synced_sid` a second time.
3518        let live = next_lifecycle_event(&mut stream).await;
3519        assert_eq!(
3520            live.event_type,
3521            session_lifecycle_event::EventType::Created as i32
3522        );
3523        assert_eq!(
3524            live.session.unwrap().session_id,
3525            live_sid,
3526            "the sync entry's buffered Created must be suppressed, and the \
3527             live session's must not be"
3528        );
3529
3530        // And nothing further: neither Created repeats.
3531        use tokio_stream::StreamExt;
3532        let extra =
3533            tokio::time::timeout(std::time::Duration::from_millis(300), stream.next()).await;
3534        assert!(
3535            extra.is_err(),
3536            "unexpected extra lifecycle event: {extra:?}"
3537        );
3538    }
3539
3540    /// Criterion 3 end to end through the real handler, which the unit tests of
3541    /// `drain_lifecycle_events` cannot reach: a client that reads slowly while
3542    /// lifecycle events arrive throughout a long sync must not be killed with
3543    /// `RESOURCE_EXHAUSTED`, and must still see every `Created` exactly once.
3544    ///
3545    /// The shape matters. The lifecycle bus holds 64 events, and far more than
3546    /// that arrive here — but they arrive *interleaved* with the reads, which is
3547    /// what a slow consumer actually looks like. The handler drains the bus
3548    /// before fetching each session, so the receiver never falls 64 behind.
3549    /// Delete that drain (or hoist it out of the loop) and the bus overruns
3550    /// mid-sync, the first post-sync `recv()` returns `Lagged`, and this test
3551    /// fails on the stream error.
3552    #[tokio::test]
3553    async fn watch_sessions_survives_a_slow_consumer_during_a_long_sync() {
3554        use tokio_stream::StreamExt;
3555
3556        let (server, runtime) = make_server();
3557        let seeded = seed_ids(200);
3558        seed_sessions(&runtime, &seeded).await;
3559
3560        let mut stream = server
3561            .watch_sessions(watch_sessions_req("agent://observer"))
3562            .await
3563            .unwrap()
3564            .into_inner();
3565
3566        let mut counts: HashMap<String, usize> = HashMap::new();
3567        // Reads the next event, failing loudly if the stream errored — that is
3568        // the RESOURCE_EXHAUSTED this test exists to rule out.
3569        async fn read(
3570            stream: &mut <MacpServer as MacpRuntimeService>::WatchSessionsStream,
3571            counts: &mut HashMap<String, usize>,
3572        ) {
3573            let resp = tokio::time::timeout(std::time::Duration::from_secs(10), stream.next())
3574                .await
3575                .expect("WatchSessions stalled")
3576                .expect("stream ended early")
3577                .expect("stream must not be terminated (RESOURCE_EXHAUSTED)");
3578            let event = resp.event.expect("event present");
3579            if event.event_type == session_lifecycle_event::EventType::Created as i32 {
3580                let session = event.session.expect("Created always carries metadata");
3581                *counts.entry(session.session_id).or_default() += 1;
3582            }
3583        }
3584
3585        // One read starts the sync and suspends the generator at its first
3586        // yield, with 199 sessions still to emit.
3587        read(&mut stream, &mut counts).await;
3588
3589        // 70 live starts — more than the bus capacity — spread across the sync,
3590        // two reads per start so the consumer stays behind the producer without
3591        // ever stopping.
3592        let initiator = "agent://orchestrator";
3593        let mut live = Vec::new();
3594        for _ in 0..70 {
3595            let sid = new_sid();
3596            start_session(&server, initiator, &sid, vec![initiator.into()]).await;
3597            live.push(sid);
3598            read(&mut stream, &mut counts).await;
3599            read(&mut stream, &mut counts).await;
3600        }
3601
3602        // Drain the rest of the sync plus the buffered live events.
3603        let expected = seeded.len() + live.len();
3604        while counts.len() < expected {
3605            read(&mut stream, &mut counts).await;
3606        }
3607
3608        for id in seeded.iter().chain(live.iter()) {
3609            assert_eq!(
3610                counts.get(id).copied(),
3611                Some(1),
3612                "{id} was not emitted exactly once"
3613            );
3614        }
3615        assert_eq!(counts.len(), expected, "unexpected extra sessions emitted");
3616    }
3617}