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 policies_read_only: bool,
42 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 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 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 match self.runtime.get_session_checked(&env.session_id).await {
94 Some(session) => engine.evaluate_message(identity, &session, env).await,
95 None => return Ok(()),
96 }
97 } else {
98 return Ok(());
99 };
100 match decision {
101 macp_core::policy::PolicyDecision::Allow { .. } => Ok(()),
102 macp_core::policy::PolicyDecision::Deny { reasons } => {
103 Err(MacpError::PolicyDenied { reasons })
104 }
105 other => {
106 tracing::warn!(decision = ?other, "unrecognized ingress policy decision");
107 Err(MacpError::PolicyDenied {
108 reasons: vec!["unrecognized policy decision".into()],
109 })
110 }
111 }
112 }
113
114 fn validate_envelope_shape(&self, env: &Envelope) -> Result<(), MacpError> {
115 if env.macp_version != "1.0" {
116 return Err(MacpError::InvalidMacpVersion);
117 }
118 if env.message_type.is_empty() || env.message_id.is_empty() {
119 return Err(MacpError::InvalidEnvelope);
120 }
121 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 let val = tonic::metadata::MetadataValue::from_bytes(&details);
780 status
781 .metadata_mut()
782 .insert_bin("macp-error-details-bin", val);
783 }
784 status
785 }
786 _ => Status::failed_precondition(err.to_string()),
787 }
788 }
789}
790
791#[tonic::async_trait]
792impl MacpRuntimeService for MacpServer {
793 async fn initialize(
794 &self,
795 request: Request<InitializeRequest>,
796 ) -> Result<Response<InitializeResponse>, Status> {
797 let req = request.into_inner();
798 if req.supported_protocol_versions.is_empty() {
799 return Err(Status::invalid_argument(
800 "INVALID_REQUEST: supported_protocol_versions must not be empty",
801 ));
802 }
803 if !req.supported_protocol_versions.iter().any(|v| v == "1.0") {
804 return Err(Status::failed_precondition(
805 "UNSUPPORTED_PROTOCOL_VERSION: no mutually supported protocol version",
806 ));
807 }
808
809 Ok(Response::new(InitializeResponse {
810 selected_protocol_version: "1.0".into(),
811 runtime_info: Some(RuntimeInfo {
812 name: "macp-runtime".into(),
813 title: "MACP Reference Runtime".into(),
814 version: env!("CARGO_PKG_VERSION").into(),
817 description: "Reference implementation of the Multi-Agent Coordination Protocol"
818 .into(),
819 website_url: String::new(),
820 }),
821 capabilities: Some(Capabilities {
822 sessions: Some(SessionsCapability { stream: true, list_sessions: true, watch_sessions: true }),
823 cancellation: Some(CancellationCapability {
824 cancel_session: true,
825 }),
826 progress: Some(ProgressCapability { progress: true }),
827 manifest: Some(ManifestCapability { get_manifest: true }),
828 mode_registry: Some(ModeRegistryCapability {
829 list_modes: true,
830 list_changed: true,
831 }),
832 roots: Some(RootsCapability {
833 list_roots: true,
839 list_changed: false,
840 }),
841 policy_registry: Some(PolicyRegistryCapability {
842 register_policy: !self.policies_read_only,
843 list_policies: true,
844 list_changed: true,
845 }),
846 experimental: Some(crate::pb::ExperimentalCapabilities {
847 features: HashMap::from([
848 ("ext_mode_lifecycle".into(), "true".into()),
849 ]),
850 }),
851 }),
852 supported_modes: self.runtime.registered_mode_names(),
853 instructions: "Authenticate requests with Authorization: Bearer <token>. Use the unary Send RPC for all session messaging. For local development only, x-macp-agent-id may be enabled by configuration.".into(),
854 }))
855 }
856
857 async fn send(&self, request: Request<SendRequest>) -> Result<Response<SendResponse>, Status> {
858 let env = request
859 .get_ref()
860 .envelope
861 .clone()
862 .ok_or_else(|| Status::invalid_argument("SendRequest must contain an envelope"))?;
863
864 let result = async {
865 self.validate_envelope_shape(&env)?;
866 let (env, max_open) = self.authenticate_send_request(&request, env).await?;
867 self.runtime
868 .process(&env, max_open)
869 .await
870 .map(|process_result| (env, process_result))
871 }
872 .await;
873
874 let ack = match result {
875 Ok((env, process_result)) => Ack {
876 ok: true,
877 duplicate: process_result.duplicate,
878 message_id: env.message_id.clone(),
879 session_id: env.session_id.clone(),
880 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
881 session_state: Self::session_state_to_pb(&process_result.session_state),
882 error: None,
883 },
884 Err(err) => {
885 let env = request.get_ref().envelope.clone().unwrap_or_default();
886 if !env.session_id.is_empty() {
890 self.runtime.metrics().record_message_rejected(&env.mode);
891 if env.message_type == "Commitment" {
892 self.runtime.metrics().record_commitment_rejected(&env.mode);
893 }
894 }
895 Self::make_error_ack(&err, &env)
896 }
897 };
898
899 Ok(Response::new(SendResponse { ack: Some(ack) }))
900 }
901
902 async fn get_session(
903 &self,
904 request: Request<GetSessionRequest>,
905 ) -> Result<Response<GetSessionResponse>, Status> {
906 let session_id = request.get_ref().session_id.clone();
907 let _identity = self
908 .authenticate_session_access(&request, &session_id)
909 .await?;
910 let session = self
911 .runtime
912 .get_session_checked(&session_id)
913 .await
914 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
915
916 Ok(Response::new(GetSessionResponse {
917 metadata: Some(Self::session_to_metadata(&session)),
918 }))
919 }
920
921 async fn cancel_session(
922 &self,
923 request: Request<CancelSessionRequest>,
924 ) -> Result<Response<CancelSessionResponse>, Status> {
925 let session_id = request.get_ref().session_id.clone();
926 let identity = self
927 .security
928 .authenticate_metadata(request.metadata())
929 .await
930 .map_err(Self::status_from_error)?;
931 let session = self
932 .runtime
933 .get_session_checked(&session_id)
934 .await
935 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
936 if identity.sender != session.initiator_sender
939 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
940 {
941 return Err(Status::permission_denied(
942 "FORBIDDEN: only the session initiator or policy-delegated roles can cancel",
943 ));
944 }
945 let sender = identity.sender.clone();
946 let req = request.into_inner();
947 match self
948 .runtime
949 .cancel_session(&req.session_id, &req.reason, &sender)
950 .await
951 {
952 Ok(result) => Ok(Response::new(CancelSessionResponse {
953 ack: Some(Ack {
954 ok: true,
955 duplicate: false,
956 message_id: String::new(),
957 session_id: req.session_id,
958 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
959 session_state: Self::session_state_to_pb(&result.session_state),
960 error: None,
961 }),
962 })),
963 Err(err) => Ok(Response::new(CancelSessionResponse {
964 ack: Some(Ack {
965 ok: false,
966 duplicate: false,
967 message_id: String::new(),
968 session_id: req.session_id.clone(),
969 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
970 session_state: PbSessionState::Unspecified.into(),
971 error: Some(PbMacpError {
972 code: err.error_code().into(),
973 message: err.to_string(),
974 session_id: req.session_id,
975 message_id: String::new(),
976 details: vec![],
977 }),
978 }),
979 })),
980 }
981 }
982
983 async fn suspend_session(
984 &self,
985 request: Request<SuspendSessionRequest>,
986 ) -> Result<Response<SuspendSessionResponse>, Status> {
987 let session_id = request.get_ref().session_id.clone();
988 let identity = self
989 .security
990 .authenticate_metadata(request.metadata())
991 .await
992 .map_err(Self::status_from_error)?;
993 let session = self
994 .runtime
995 .get_session_checked(&session_id)
996 .await
997 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
998 if identity.sender != session.initiator_sender
1001 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1002 {
1003 return Err(Status::permission_denied(
1004 "FORBIDDEN: only the session initiator or policy-delegated roles can suspend",
1005 ));
1006 }
1007 let sender = identity.sender.clone();
1008 let req = request.into_inner();
1009 match self
1010 .runtime
1011 .suspend_session(&req.session_id, &req.reason, &sender)
1012 .await
1013 {
1014 Ok(result) => Ok(Response::new(SuspendSessionResponse {
1015 ack: Some(Ack {
1016 ok: true,
1017 duplicate: false,
1018 message_id: String::new(),
1019 session_id: req.session_id,
1020 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1021 session_state: Self::session_state_to_pb(&result.session_state),
1022 error: None,
1023 }),
1024 })),
1025 Err(err) => Ok(Response::new(SuspendSessionResponse {
1026 ack: Some(Ack {
1027 ok: false,
1028 duplicate: false,
1029 message_id: String::new(),
1030 session_id: req.session_id.clone(),
1031 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1032 session_state: PbSessionState::Unspecified.into(),
1033 error: Some(PbMacpError {
1034 code: err.error_code().into(),
1035 message: err.to_string(),
1036 session_id: req.session_id,
1037 message_id: String::new(),
1038 details: vec![],
1039 }),
1040 }),
1041 })),
1042 }
1043 }
1044
1045 async fn resume_session(
1046 &self,
1047 request: Request<ResumeSessionRequest>,
1048 ) -> Result<Response<ResumeSessionResponse>, Status> {
1049 let session_id = request.get_ref().session_id.clone();
1050 let identity = self
1051 .security
1052 .authenticate_metadata(request.metadata())
1053 .await
1054 .map_err(Self::status_from_error)?;
1055 let session = self
1056 .runtime
1057 .get_session_checked(&session_id)
1058 .await
1059 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1060 if identity.sender != session.initiator_sender
1061 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1062 {
1063 return Err(Status::permission_denied(
1064 "FORBIDDEN: only the session initiator or policy-delegated roles can resume",
1065 ));
1066 }
1067 let sender = identity.sender.clone();
1068 let req = request.into_inner();
1069 match self
1070 .runtime
1071 .resume_session(&req.session_id, &req.reason, &sender)
1072 .await
1073 {
1074 Ok(result) => Ok(Response::new(ResumeSessionResponse {
1075 ack: Some(Ack {
1076 ok: true,
1077 duplicate: false,
1078 message_id: String::new(),
1079 session_id: req.session_id,
1080 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1081 session_state: Self::session_state_to_pb(&result.session_state),
1082 error: None,
1083 }),
1084 })),
1085 Err(err) => Ok(Response::new(ResumeSessionResponse {
1086 ack: Some(Ack {
1087 ok: false,
1088 duplicate: false,
1089 message_id: String::new(),
1090 session_id: req.session_id.clone(),
1091 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1092 session_state: PbSessionState::Unspecified.into(),
1093 error: Some(PbMacpError {
1094 code: err.error_code().into(),
1095 message: err.to_string(),
1096 session_id: req.session_id,
1097 message_id: String::new(),
1098 details: vec![],
1099 }),
1100 }),
1101 })),
1102 }
1103 }
1104
1105 async fn get_manifest(
1106 &self,
1107 request: Request<GetManifestRequest>,
1108 ) -> Result<Response<GetManifestResponse>, Status> {
1109 let req = request.into_inner();
1110 if !req.agent_id.is_empty() && req.agent_id != "macp-runtime" {
1111 return Err(Status::not_found(format!(
1112 "Agent '{}' not found",
1113 req.agent_id
1114 )));
1115 }
1116
1117 Ok(Response::new(GetManifestResponse {
1118 manifest: Some(crate::pb::AgentManifest {
1119 agent_id: "macp-runtime".into(),
1120 title: "MACP Reference Runtime".into(),
1121 description: "Reference implementation of MACP".into(),
1122 supported_modes: self.runtime.registered_mode_names(),
1123 input_content_types: vec!["application/macp-envelope+proto".into()],
1124 output_content_types: vec!["application/macp-envelope+proto".into()],
1125 metadata: HashMap::new(),
1126 transport_endpoints: vec![],
1128 }),
1129 }))
1130 }
1131
1132 async fn list_modes(
1133 &self,
1134 _request: Request<ListModesRequest>,
1135 ) -> Result<Response<ListModesResponse>, Status> {
1136 Ok(Response::new(ListModesResponse {
1137 modes: self.runtime.standard_mode_descriptors(),
1138 }))
1139 }
1140
1141 async fn list_roots(
1142 &self,
1143 _request: Request<ListRootsRequest>,
1144 ) -> Result<Response<ListRootsResponse>, Status> {
1145 Ok(Response::new(ListRootsResponse { roots: vec![] }))
1146 }
1147
1148 type StreamSessionStream = SessionResponseStream;
1149
1150 async fn stream_session(
1151 &self,
1152 request: Request<tonic::Streaming<StreamSessionRequest>>,
1153 ) -> Result<Response<Self::StreamSessionStream>, Status> {
1154 let identity = self
1155 .security
1156 .authenticate_metadata(request.metadata())
1157 .await
1158 .map_err(Self::status_from_error)?;
1159 let inbound = request.into_inner();
1160 Ok(Response::new(
1161 self.build_stream_session_stream(identity, inbound),
1162 ))
1163 }
1164
1165 type WatchModeRegistryStream = std::pin::Pin<
1166 Box<dyn futures_core::Stream<Item = Result<WatchModeRegistryResponse, Status>> + Send>,
1167 >;
1168
1169 async fn watch_mode_registry(
1170 &self,
1171 _request: Request<WatchModeRegistryRequest>,
1172 ) -> Result<Response<Self::WatchModeRegistryStream>, Status> {
1173 let mut rx = self.runtime.subscribe_mode_changes();
1174 let stream = async_stream::try_stream! {
1175 yield WatchModeRegistryResponse {
1177 change: Some(crate::pb::RegistryChanged {
1178 registry: "modes".into(),
1179 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1180 }),
1181 };
1182 while rx.recv().await.is_ok() {
1184 yield WatchModeRegistryResponse {
1185 change: Some(crate::pb::RegistryChanged {
1186 registry: "modes".into(),
1187 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1188 }),
1189 };
1190 }
1191 };
1192 Ok(Response::new(Box::pin(stream)))
1193 }
1194
1195 type WatchRootsStream = std::pin::Pin<
1196 Box<dyn futures_core::Stream<Item = Result<WatchRootsResponse, Status>> + Send>,
1197 >;
1198
1199 async fn watch_roots(
1200 &self,
1201 _request: Request<WatchRootsRequest>,
1202 ) -> Result<Response<Self::WatchRootsStream>, Status> {
1203 let initial = WatchRootsResponse {
1204 change: Some(crate::pb::RootsChanged {
1205 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1206 }),
1207 };
1208 let stream = async_stream::try_stream! {
1209 yield initial;
1210 std::future::pending::<()>().await;
1212 };
1213 Ok(Response::new(Box::pin(stream)))
1214 }
1215
1216 type WatchSignalsStream = std::pin::Pin<
1217 Box<dyn futures_core::Stream<Item = Result<WatchSignalsResponse, Status>> + Send>,
1218 >;
1219
1220 type WatchSessionsStream = std::pin::Pin<
1221 Box<dyn futures_core::Stream<Item = Result<WatchSessionsResponse, Status>> + Send>,
1222 >;
1223
1224 async fn watch_signals(
1225 &self,
1226 request: Request<WatchSignalsRequest>,
1227 ) -> Result<Response<Self::WatchSignalsStream>, Status> {
1228 let _identity = self
1233 .security
1234 .authenticate_metadata(request.metadata())
1235 .await
1236 .map_err(Self::status_from_error)?;
1237 let mut rx = self.runtime.subscribe_signals();
1238 let stream = async_stream::try_stream! {
1239 loop {
1240 match rx.recv().await {
1241 Ok(envelope) => {
1242 yield WatchSignalsResponse {
1243 envelope: Some(envelope),
1244 };
1245 }
1246 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1250 Err(Status::resource_exhausted(format!(
1251 "WatchSignals receiver fell behind by {skipped} signals"
1252 )))?;
1253 }
1254 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1255 }
1256 }
1257 };
1258 Ok(Response::new(Box::pin(stream)))
1259 }
1260
1261 async fn list_sessions(
1264 &self,
1265 request: Request<ListSessionsRequest>,
1266 ) -> Result<Response<ListSessionsResponse>, Status> {
1267 let _identity = self
1271 .security
1272 .authenticate_metadata(request.metadata())
1273 .await
1274 .map_err(Self::status_from_error)?;
1275 let req = request.into_inner();
1276
1277 if req.page_size < 0 {
1284 return Err(Status::invalid_argument(
1285 "INVALID_ARGUMENT: page_size must not be negative",
1286 ));
1287 }
1288 let effective = if req.page_size == 0 {
1291 self.security.list_sessions_default_page_size
1292 } else {
1293 (req.page_size as usize).min(self.security.list_sessions_max_page_size)
1294 };
1295 let effective = effective.max(1);
1303
1304 let cursor = if req.page_token.is_empty() {
1305 None
1306 } else {
1307 Some(
1312 crate::pagination::decode_page_token(&req.page_token).map_err(|_| {
1313 Status::invalid_argument(
1314 "INVALID_ARGUMENT: page_token is not a valid continuation token",
1315 )
1316 })?,
1317 )
1318 };
1319
1320 let ids = self
1325 .runtime
1326 .registry
1327 .session_ids_after(cursor.as_deref(), effective.saturating_add(1))
1328 .await;
1329 let has_more = ids.len() > effective;
1330 let page_ids = &ids[..effective.min(ids.len())];
1331
1332 let next_page_token = match (has_more, page_ids.last()) {
1338 (true, Some(last)) => crate::pagination::encode_page_token(last),
1339 _ => String::new(),
1340 };
1341
1342 let mut metadata: Vec<SessionMetadata> = Vec::with_capacity(page_ids.len());
1343 for id in page_ids {
1344 if let Some(session) = self.runtime.registry.get_session(id).await {
1348 debug_assert_eq!(
1349 session.session_id, *id,
1350 "registry map key must equal Session::session_id — paging orders \
1351 by the key but emits the field"
1352 );
1353 metadata.push(Self::session_to_metadata(&session));
1354 }
1355 }
1356
1357 Ok(Response::new(ListSessionsResponse {
1358 sessions: metadata,
1359 next_page_token,
1360 }))
1361 }
1362
1363 async fn watch_sessions(
1364 &self,
1365 request: Request<WatchSessionsRequest>,
1366 ) -> Result<Response<Self::WatchSessionsStream>, Status> {
1367 let _identity = self
1368 .security
1369 .authenticate_metadata(request.metadata())
1370 .await
1371 .map_err(Self::status_from_error)?;
1372 let mut rx = self.runtime.subscribe_session_lifecycle();
1373 let runtime = Arc::clone(&self.runtime);
1374 let stream = async_stream::try_stream! {
1375 let sessions = runtime.registry.get_all_sessions().await;
1381 let mut synced: std::collections::HashSet<String> =
1382 std::collections::HashSet::with_capacity(sessions.len());
1383 for session in &sessions {
1384 synced.insert(session.session_id.clone());
1385 yield WatchSessionsResponse {
1386 event: Some(SessionLifecycleEvent {
1387 event_type: session_lifecycle_event::EventType::Created.into(),
1388 session: Some(Self::session_to_metadata(session)),
1389 observed_at_unix_ms: session.started_at_unix_ms,
1390 }),
1391 };
1392 }
1393 loop {
1395 let event = match rx.recv().await {
1396 Ok(event) => event,
1397 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1398 Err(Status::resource_exhausted(format!(
1399 "WatchSessions receiver fell behind by {skipped} events"
1400 )))?;
1401 break;
1402 }
1403 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1404 };
1405 let (event_type, sid) = match &event {
1406 crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1407 (session_lifecycle_event::EventType::Created, session_id.clone()),
1408 crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1409 (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1410 crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1411 (session_lifecycle_event::EventType::Expired, session_id.clone()),
1412 crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1413 (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1414 crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1415 (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1416 crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1417 (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1418 };
1419 if event_type == session_lifecycle_event::EventType::Created
1423 && !synced.insert(sid.clone())
1424 {
1425 continue;
1426 }
1427 let session_meta = runtime.registry.get_session(&sid).await
1428 .map(|s| Self::session_to_metadata(&s));
1429 yield WatchSessionsResponse {
1430 event: Some(SessionLifecycleEvent {
1431 event_type: event_type.into(),
1432 session: session_meta,
1433 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1434 }),
1435 };
1436 }
1437 };
1438 Ok(Response::new(Box::pin(stream)))
1439 }
1440
1441 async fn list_ext_modes(
1444 &self,
1445 _request: Request<ListExtModesRequest>,
1446 ) -> Result<Response<ListExtModesResponse>, Status> {
1447 Ok(Response::new(ListExtModesResponse {
1448 modes: self.runtime.extension_mode_descriptors(),
1449 }))
1450 }
1451
1452 async fn register_ext_mode(
1453 &self,
1454 request: Request<RegisterExtModeRequest>,
1455 ) -> Result<Response<RegisterExtModeResponse>, Status> {
1456 let identity = self
1457 .security
1458 .authenticate_metadata(request.metadata())
1459 .await
1460 .map_err(Self::status_from_error)?;
1461 self.security
1462 .authorize_mode_registry(&identity)
1463 .map_err(Self::status_from_error)?;
1464 let req = request.into_inner();
1465 let descriptor = req
1466 .mode_descriptor
1467 .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1468 match self.runtime.register_extension(descriptor) {
1469 Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1470 ok: true,
1471 error: String::new(),
1472 })),
1473 Err(e) => Ok(Response::new(RegisterExtModeResponse {
1474 ok: false,
1475 error: e,
1476 })),
1477 }
1478 }
1479
1480 async fn unregister_ext_mode(
1481 &self,
1482 request: Request<UnregisterExtModeRequest>,
1483 ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1484 let identity = self
1485 .security
1486 .authenticate_metadata(request.metadata())
1487 .await
1488 .map_err(Self::status_from_error)?;
1489 self.security
1490 .authorize_mode_registry(&identity)
1491 .map_err(Self::status_from_error)?;
1492 let req = request.into_inner();
1493 match self.runtime.unregister_extension(&req.mode) {
1494 Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1495 ok: true,
1496 error: String::new(),
1497 })),
1498 Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1499 ok: false,
1500 error: e,
1501 })),
1502 }
1503 }
1504
1505 async fn promote_mode(
1506 &self,
1507 request: Request<PromoteModeRequest>,
1508 ) -> Result<Response<PromoteModeResponse>, Status> {
1509 let identity = self
1510 .security
1511 .authenticate_metadata(request.metadata())
1512 .await
1513 .map_err(Self::status_from_error)?;
1514 self.security
1515 .authorize_mode_registry(&identity)
1516 .map_err(Self::status_from_error)?;
1517 let req = request.into_inner();
1518 let new_name = if req.promoted_mode_name.is_empty() {
1519 None
1520 } else {
1521 Some(req.promoted_mode_name.as_str())
1522 };
1523 match self.runtime.promote_mode(&req.mode, new_name) {
1524 Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1525 ok: true,
1526 error: String::new(),
1527 mode: final_name,
1528 })),
1529 Err(e) => Ok(Response::new(PromoteModeResponse {
1530 ok: false,
1531 error: e,
1532 mode: String::new(),
1533 })),
1534 }
1535 }
1536
1537 async fn register_policy(
1540 &self,
1541 request: Request<RegisterPolicyRequest>,
1542 ) -> Result<Response<RegisterPolicyResponse>, Status> {
1543 if self.policies_read_only {
1544 return Err(Status::failed_precondition(
1545 "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1546 ));
1547 }
1548 let identity = self
1549 .security
1550 .authenticate_metadata(request.metadata())
1551 .await
1552 .map_err(Self::status_from_error)?;
1553 self.security
1554 .authorize_mode_registry(&identity)
1555 .map_err(Self::status_from_error)?;
1556 let req = request.into_inner();
1557 let descriptor = req
1558 .policy_descriptor
1559 .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1560 let definition = Self::policy_descriptor_to_definition(&descriptor);
1561 match self.runtime.register_policy(definition) {
1562 Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1563 ok: true,
1564 error: String::new(),
1565 })),
1566 Err(e) => Ok(Response::new(RegisterPolicyResponse {
1567 ok: false,
1568 error: e,
1569 })),
1570 }
1571 }
1572
1573 async fn unregister_policy(
1574 &self,
1575 request: Request<UnregisterPolicyRequest>,
1576 ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1577 if self.policies_read_only {
1578 return Err(Status::failed_precondition(
1579 "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1580 ));
1581 }
1582 let identity = self
1583 .security
1584 .authenticate_metadata(request.metadata())
1585 .await
1586 .map_err(Self::status_from_error)?;
1587 self.security
1588 .authorize_mode_registry(&identity)
1589 .map_err(Self::status_from_error)?;
1590 let req = request.into_inner();
1591 match self.runtime.unregister_policy(&req.policy_id) {
1592 Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1593 ok: true,
1594 error: String::new(),
1595 })),
1596 Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1597 ok: false,
1598 error: e,
1599 })),
1600 }
1601 }
1602
1603 async fn get_policy(
1604 &self,
1605 request: Request<GetPolicyRequest>,
1606 ) -> Result<Response<GetPolicyResponse>, Status> {
1607 let _identity = self
1608 .security
1609 .authenticate_metadata(request.metadata())
1610 .await
1611 .map_err(Self::status_from_error)?;
1612 let req = request.into_inner();
1613 let policy = self
1614 .runtime
1615 .get_policy(&req.policy_id)
1616 .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1617 Ok(Response::new(GetPolicyResponse {
1618 policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1619 }))
1620 }
1621
1622 async fn list_policies(
1623 &self,
1624 request: Request<ListPoliciesRequest>,
1625 ) -> Result<Response<ListPoliciesResponse>, Status> {
1626 let _identity = self
1627 .security
1628 .authenticate_metadata(request.metadata())
1629 .await
1630 .map_err(Self::status_from_error)?;
1631 let req = request.into_inner();
1632 let mode_filter = if req.mode.is_empty() {
1633 None
1634 } else {
1635 Some(req.mode.as_str())
1636 };
1637 let policies = self.runtime.list_policies(mode_filter);
1638 let descriptors = policies
1639 .iter()
1640 .map(Self::policy_definition_to_descriptor)
1641 .collect();
1642 Ok(Response::new(ListPoliciesResponse { descriptors }))
1643 }
1644
1645 type WatchPoliciesStream = std::pin::Pin<
1646 Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1647 >;
1648
1649 async fn watch_policies(
1650 &self,
1651 _request: Request<WatchPoliciesRequest>,
1652 ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1653 let mut rx = self.runtime.subscribe_policy_changes();
1654 let runtime = Arc::clone(&self.runtime);
1655 let stream = async_stream::try_stream! {
1656 let policies = runtime.list_policies(None);
1658 let descriptors: Vec<PolicyDescriptor> = policies
1659 .iter()
1660 .map(MacpServer::policy_definition_to_descriptor)
1661 .collect();
1662 yield WatchPoliciesResponse {
1663 descriptors,
1664 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1665 };
1666 while rx.recv().await.is_ok() {
1668 let policies = runtime.list_policies(None);
1669 let descriptors: Vec<PolicyDescriptor> = policies
1670 .iter()
1671 .map(MacpServer::policy_definition_to_descriptor)
1672 .collect();
1673 yield WatchPoliciesResponse {
1674 descriptors,
1675 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1676 };
1677 }
1678 };
1679 Ok(Response::new(Box::pin(stream)))
1680 }
1681}
1682
1683impl MacpServer {
1686 fn policy_descriptor_to_definition(
1687 descriptor: &PolicyDescriptor,
1688 ) -> crate::policy::PolicyDefinition {
1689 let rules: serde_json::Value = if descriptor.rules.is_empty() {
1690 serde_json::json!({})
1691 } else {
1692 serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1693 };
1694 crate::policy::PolicyDefinition {
1695 policy_id: descriptor.policy_id.clone(),
1696 mode: descriptor.mode.clone(),
1697 description: descriptor.description.clone(),
1698 rules,
1699 schema_version: descriptor.schema_version,
1700 }
1701 }
1702
1703 fn policy_definition_to_descriptor(
1704 definition: &crate::policy::PolicyDefinition,
1705 ) -> PolicyDescriptor {
1706 PolicyDescriptor {
1707 policy_id: definition.policy_id.clone(),
1708 mode: definition.mode.clone(),
1709 description: definition.description.clone(),
1710 rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1711 schema_version: definition.schema_version,
1712 registered_at_unix_ms: 0,
1713 }
1714 }
1715}
1716
1717#[cfg(test)]
1718mod tests {
1719 use super::*;
1720 use crate::log_store::LogStore;
1721 use crate::pb::SessionStartPayload;
1722 use crate::registry::SessionRegistry;
1723 use chrono::Utc;
1724 use prost::Message;
1725
1726 fn new_sid() -> String {
1727 uuid::Uuid::new_v4().as_hyphenated().to_string()
1728 }
1729
1730 fn make_server() -> (MacpServer, Arc<Runtime>) {
1731 make_server_with_security(SecurityLayer::dev_mode())
1732 }
1733
1734 fn make_server_with_security(security: SecurityLayer) -> (MacpServer, Arc<Runtime>) {
1738 let storage: Arc<dyn crate::storage::StorageBackend> =
1739 Arc::new(crate::storage::MemoryBackend);
1740 let registry = Arc::new(SessionRegistry::new());
1741 let log_store = Arc::new(LogStore::new());
1742 let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1743 let server = MacpServer::new(runtime.clone(), security);
1744 (server, runtime)
1745 }
1746
1747 fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1748 let mut req = Request::new(SendRequest {
1749 envelope: Some(env),
1750 });
1751 req.metadata_mut()
1752 .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1753 req
1754 }
1755
1756 async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1757 let resp = server.send(send_req(sender, env)).await.unwrap();
1758 resp.into_inner().ack.unwrap()
1759 }
1760
1761 fn start_payload() -> Vec<u8> {
1762 SessionStartPayload {
1763 intent: "intent".into(),
1764 participants: vec!["agent://fraud".into()],
1765 mode_version: "1.0.0".into(),
1766 configuration_version: "cfg-1".into(),
1767 policy_version: String::new(),
1768 ttl_ms: 1000,
1769 context_id: String::new(),
1770 extensions: std::collections::HashMap::new(),
1771 roots: vec![],
1772 max_suspend_ms: 0,
1773 }
1774 .encode_to_vec()
1775 }
1776
1777 #[tokio::test]
1778 async fn sender_is_derived_from_authenticated_metadata() {
1779 let (server, runtime) = make_server();
1780 let sid = new_sid();
1781 let ack = do_send(
1782 &server,
1783 "agent://orchestrator",
1784 Envelope {
1785 macp_version: "1.0".into(),
1786 mode: "macp.mode.decision.v1".into(),
1787 message_type: "SessionStart".into(),
1788 message_id: "m1".into(),
1789 session_id: sid.clone(),
1790 sender: String::new(),
1791 timestamp_unix_ms: Utc::now().timestamp_millis(),
1792 payload: start_payload(),
1793 },
1794 )
1795 .await;
1796 assert!(ack.ok);
1797 let session = runtime.get_session_checked(&sid).await.unwrap();
1798 assert_eq!(session.initiator_sender, "agent://orchestrator");
1799 }
1800
1801 #[tokio::test]
1802 async fn spoofed_sender_is_rejected() {
1803 let (server, _) = make_server();
1804 let sid = new_sid();
1805 let ack = do_send(
1806 &server,
1807 "agent://orchestrator",
1808 Envelope {
1809 macp_version: "1.0".into(),
1810 mode: "macp.mode.decision.v1".into(),
1811 message_type: "SessionStart".into(),
1812 message_id: "m1".into(),
1813 session_id: sid,
1814 sender: "agent://spoof".into(),
1815 timestamp_unix_ms: Utc::now().timestamp_millis(),
1816 payload: start_payload(),
1817 },
1818 )
1819 .await;
1820 assert!(!ack.ok);
1821 assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1822 }
1823
1824 #[tokio::test]
1825 async fn get_session_requires_session_membership() {
1826 let (server, _) = make_server();
1827 let sid = new_sid();
1828 let ack = do_send(
1829 &server,
1830 "agent://orchestrator",
1831 Envelope {
1832 macp_version: "1.0".into(),
1833 mode: "macp.mode.decision.v1".into(),
1834 message_type: "SessionStart".into(),
1835 message_id: "m1".into(),
1836 session_id: sid.clone(),
1837 sender: String::new(),
1838 timestamp_unix_ms: Utc::now().timestamp_millis(),
1839 payload: start_payload(),
1840 },
1841 )
1842 .await;
1843 assert!(ack.ok);
1844
1845 let mut req = Request::new(GetSessionRequest { session_id: sid });
1846 req.metadata_mut().insert(
1847 "authorization",
1848 format!("Bearer {}", "agent://outsider").parse().unwrap(),
1849 );
1850 let err = server.get_session(req).await.unwrap_err();
1851 assert_eq!(err.code(), tonic::Code::PermissionDenied);
1852 }
1853
1854 #[tokio::test]
1855 async fn register_ext_mode_requires_authenticated_registry_permission() {
1856 let storage: Arc<dyn crate::storage::StorageBackend> =
1857 Arc::new(crate::storage::MemoryBackend);
1858 let registry = Arc::new(SessionRegistry::new());
1859 let log_store = Arc::new(LogStore::new());
1860 let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1861 let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1862 let server = MacpServer::new(runtime, security);
1863
1864 let req = Request::new(RegisterExtModeRequest {
1865 mode_descriptor: Some(crate::pb::ModeDescriptor {
1866 mode: "ext.custom.v1".into(),
1867 mode_version: "1.0.0".into(),
1868 message_types: vec!["SessionStart".into(), "Commitment".into()],
1869 ..Default::default()
1870 }),
1871 });
1872 let err = server.register_ext_mode(req).await.unwrap_err();
1873 assert_eq!(err.code(), tonic::Code::Unauthenticated);
1874 }
1875
1876 fn stream_identity(sender: &str) -> AuthIdentity {
1877 AuthIdentity {
1878 sender: sender.into(),
1879 allowed_modes: None,
1880 can_start_sessions: true,
1881 max_open_sessions: None,
1882 can_manage_mode_registry: false,
1883 is_observer: false,
1884 }
1885 }
1886
1887 #[tokio::test]
1888 async fn stream_session_emits_accepted_envelopes_only() {
1889 use tokio_stream::{iter, StreamExt};
1890
1891 let (server, _) = make_server();
1892 let sid = new_sid();
1893 let requests = iter(vec![Ok(StreamSessionRequest {
1894 subscribe_session_id: String::new(),
1895 after_sequence: 0,
1896 envelope: Some(Envelope {
1897 macp_version: "1.0".into(),
1898 mode: "macp.mode.decision.v1".into(),
1899 message_type: "SessionStart".into(),
1900 message_id: "m1".into(),
1901 session_id: sid.clone(),
1902 sender: String::new(),
1903 timestamp_unix_ms: Utc::now().timestamp_millis(),
1904 payload: start_payload(),
1905 }),
1906 })]);
1907
1908 let mut stream =
1909 server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1910
1911 let response = stream.next().await.unwrap().unwrap();
1912 let envelope = match response.response.unwrap() {
1913 crate::pb::stream_session_response::Response::Envelope(e) => e,
1914 _ => panic!("expected envelope"),
1915 };
1916 assert_eq!(envelope.message_type, "SessionStart");
1917 assert_eq!(envelope.message_id, "m1");
1918 assert!(stream.next().await.is_none());
1919 }
1920
1921 #[tokio::test]
1922 async fn stream_session_rejects_mixed_session_ids() {
1923 use tokio_stream::{iter, StreamExt};
1924
1925 let (server, _) = make_server();
1926 let sid1 = new_sid();
1927 let sid2 = new_sid();
1928 let requests = iter(vec![
1929 Ok(StreamSessionRequest {
1930 subscribe_session_id: String::new(),
1931 after_sequence: 0,
1932 envelope: Some(Envelope {
1933 macp_version: "1.0".into(),
1934 mode: "macp.mode.decision.v1".into(),
1935 message_type: "SessionStart".into(),
1936 message_id: "m1".into(),
1937 session_id: sid1.clone(),
1938 sender: String::new(),
1939 timestamp_unix_ms: Utc::now().timestamp_millis(),
1940 payload: start_payload(),
1941 }),
1942 }),
1943 Ok(StreamSessionRequest {
1944 subscribe_session_id: String::new(),
1945 after_sequence: 0,
1946 envelope: Some(Envelope {
1947 macp_version: "1.0".into(),
1948 mode: "macp.mode.decision.v1".into(),
1949 message_type: "SessionStart".into(),
1950 message_id: "m2".into(),
1951 session_id: sid2,
1952 sender: String::new(),
1953 timestamp_unix_ms: Utc::now().timestamp_millis(),
1954 payload: start_payload(),
1955 }),
1956 }),
1957 ]);
1958
1959 let mut stream =
1960 server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1961
1962 let first = stream.next().await.unwrap().unwrap();
1963 let first_env = match first.response.unwrap() {
1964 crate::pb::stream_session_response::Response::Envelope(e) => e,
1965 _ => panic!("expected envelope"),
1966 };
1967 assert_eq!(first_env.session_id, sid1);
1968 let err = stream.next().await.unwrap().unwrap_err();
1969 assert_eq!(err.code(), tonic::Code::InvalidArgument);
1970 }
1971
1972 #[tokio::test]
1973 async fn list_modes_returns_standard_modes() {
1974 let (server, _) = make_server();
1975 let resp = server
1976 .list_modes(Request::new(ListModesRequest {}))
1977 .await
1978 .unwrap();
1979 let names: Vec<String> = resp
1980 .into_inner()
1981 .modes
1982 .iter()
1983 .map(|m| m.mode.clone())
1984 .collect();
1985 assert_eq!(names.len(), 5);
1986 assert!(names.contains(&"macp.mode.decision.v1".to_string()));
1987 assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
1988 assert!(names.contains(&"macp.mode.task.v1".to_string()));
1989 assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
1990 assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
1991 assert!(!names.contains(&"ext.multi_round.v1".to_string()));
1993 }
1994
1995 #[tokio::test]
1996 async fn list_ext_modes_returns_extensions() {
1997 let (server, _) = make_server();
1998 let resp = server
1999 .list_ext_modes(Request::new(ListExtModesRequest {}))
2000 .await
2001 .unwrap();
2002 let names: Vec<String> = resp
2003 .into_inner()
2004 .modes
2005 .iter()
2006 .map(|m| m.mode.clone())
2007 .collect();
2008 assert_eq!(names.len(), 1);
2009 assert!(names.contains(&"ext.multi_round.v1".to_string()));
2010 }
2011
2012 #[tokio::test]
2013 async fn get_manifest_includes_all_modes() {
2014 let (server, _) = make_server();
2015 let resp = server
2016 .get_manifest(Request::new(crate::pb::GetManifestRequest {
2017 agent_id: String::new(),
2018 }))
2019 .await
2020 .unwrap();
2021 let manifest = resp.into_inner().manifest.unwrap();
2022 assert_eq!(manifest.supported_modes.len(), 6);
2023 assert!(manifest
2024 .supported_modes
2025 .contains(&"ext.multi_round.v1".to_string()));
2026 }
2027
2028 #[tokio::test]
2029 async fn get_session_returns_metadata() {
2030 let (server, _) = make_server();
2031 let sid = new_sid();
2032 let ack = do_send(
2033 &server,
2034 "agent://orchestrator",
2035 Envelope {
2036 macp_version: "1.0".into(),
2037 mode: "macp.mode.decision.v1".into(),
2038 message_type: "SessionStart".into(),
2039 message_id: "m1".into(),
2040 session_id: sid.clone(),
2041 sender: String::new(),
2042 timestamp_unix_ms: Utc::now().timestamp_millis(),
2043 payload: start_payload(),
2044 },
2045 )
2046 .await;
2047 assert!(ack.ok);
2048
2049 let mut req = Request::new(GetSessionRequest {
2050 session_id: sid.clone(),
2051 });
2052 req.metadata_mut().insert(
2053 "authorization",
2054 format!("Bearer {}", "agent://orchestrator")
2055 .parse()
2056 .unwrap(),
2057 );
2058 let resp = server.get_session(req).await.unwrap();
2059 let meta = resp.into_inner().metadata.unwrap();
2060 assert_eq!(meta.session_id, sid);
2061 assert_eq!(meta.mode, "macp.mode.decision.v1");
2062 assert_eq!(meta.mode_version, "1.0.0");
2063 assert_eq!(meta.configuration_version, "cfg-1");
2064 }
2065
2066 #[tokio::test]
2067 async fn cancel_session_transitions_to_cancelled() {
2068 let (server, _) = make_server();
2069 let sid = new_sid();
2070 let ack = do_send(
2071 &server,
2072 "agent://orchestrator",
2073 Envelope {
2074 macp_version: "1.0".into(),
2075 mode: "macp.mode.decision.v1".into(),
2076 message_type: "SessionStart".into(),
2077 message_id: "m1".into(),
2078 session_id: sid.clone(),
2079 sender: String::new(),
2080 timestamp_unix_ms: Utc::now().timestamp_millis(),
2081 payload: start_payload(),
2082 },
2083 )
2084 .await;
2085 assert!(ack.ok);
2086
2087 let mut req = Request::new(CancelSessionRequest {
2088 session_id: sid,
2089 reason: "no longer needed".into(),
2090 });
2091 req.metadata_mut().insert(
2092 "authorization",
2093 format!("Bearer {}", "agent://orchestrator")
2094 .parse()
2095 .unwrap(),
2096 );
2097 let resp = server.cancel_session(req).await.unwrap();
2098 let ack = resp.into_inner().ack.unwrap();
2099 assert!(ack.ok);
2100 assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2102 }
2103
2104 #[tokio::test]
2105 async fn participant_cannot_cancel_session() {
2106 let (server, _) = make_server();
2107 let sid = new_sid();
2108 let ack = do_send(
2109 &server,
2110 "agent://orchestrator",
2111 Envelope {
2112 macp_version: "1.0".into(),
2113 mode: "macp.mode.decision.v1".into(),
2114 message_type: "SessionStart".into(),
2115 message_id: "m1".into(),
2116 session_id: sid.clone(),
2117 sender: String::new(),
2118 timestamp_unix_ms: Utc::now().timestamp_millis(),
2119 payload: start_payload(),
2120 },
2121 )
2122 .await;
2123 assert!(ack.ok);
2124
2125 let mut req = Request::new(CancelSessionRequest {
2126 session_id: sid,
2127 reason: "I want to cancel".into(),
2128 });
2129 req.metadata_mut().insert(
2130 "authorization",
2131 format!("Bearer {}", "agent://fraud").parse().unwrap(),
2132 );
2133 let err = server.cancel_session(req).await.unwrap_err();
2134 assert_eq!(err.code(), tonic::Code::PermissionDenied);
2135 }
2136
2137 #[tokio::test]
2138 async fn cancel_session_unknown_session_returns_error() {
2139 let (server, _) = make_server();
2140 let mut req = Request::new(CancelSessionRequest {
2141 session_id: "nonexistent".into(),
2142 reason: "test".into(),
2143 });
2144 req.metadata_mut().insert(
2145 "authorization",
2146 format!("Bearer {}", "agent://orchestrator")
2147 .parse()
2148 .unwrap(),
2149 );
2150 let err = server.cancel_session(req).await.unwrap_err();
2151 assert_eq!(err.code(), tonic::Code::NotFound);
2152 }
2153
2154 #[tokio::test]
2155 async fn ambient_signal_accepted() {
2156 let (server, _) = make_server();
2157 let ack = do_send(
2158 &server,
2159 "agent://orchestrator",
2160 Envelope {
2161 macp_version: "1.0".into(),
2162 mode: String::new(),
2163 message_type: "Signal".into(),
2164 message_id: "sig-1".into(),
2165 session_id: String::new(),
2166 sender: String::new(),
2167 timestamp_unix_ms: Utc::now().timestamp_millis(),
2168 payload: vec![],
2169 },
2170 )
2171 .await;
2172 assert!(ack.ok);
2173 }
2174
2175 #[tokio::test]
2176 async fn signal_with_session_id_rejected() {
2177 let (server, _) = make_server();
2178 let ack = do_send(
2179 &server,
2180 "agent://orchestrator",
2181 Envelope {
2182 macp_version: "1.0".into(),
2183 mode: String::new(),
2184 message_type: "Signal".into(),
2185 message_id: "sig-2".into(),
2186 session_id: "some-session".into(),
2187 sender: String::new(),
2188 timestamp_unix_ms: Utc::now().timestamp_millis(),
2189 payload: vec![],
2190 },
2191 )
2192 .await;
2193 assert!(!ack.ok);
2194 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2195 }
2196
2197 #[tokio::test]
2198 async fn signal_with_mode_rejected() {
2199 let (server, _) = make_server();
2200 let ack = do_send(
2201 &server,
2202 "agent://orchestrator",
2203 Envelope {
2204 macp_version: "1.0".into(),
2205 mode: "macp.mode.decision.v1".into(),
2206 message_type: "Signal".into(),
2207 message_id: "sig-3".into(),
2208 session_id: String::new(),
2209 sender: String::new(),
2210 timestamp_unix_ms: Utc::now().timestamp_millis(),
2211 payload: vec![],
2212 },
2213 )
2214 .await;
2215 assert!(!ack.ok);
2216 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2217 }
2218
2219 #[tokio::test]
2220 async fn ambient_progress_accepted() {
2221 let (server, _) = make_server();
2222 let ack = do_send(
2223 &server,
2224 "agent://orchestrator",
2225 Envelope {
2226 macp_version: "1.0".into(),
2227 mode: String::new(),
2228 message_type: "Progress".into(),
2229 message_id: "prog-1".into(),
2230 session_id: String::new(),
2231 sender: String::new(),
2232 timestamp_unix_ms: Utc::now().timestamp_millis(),
2233 payload: vec![],
2234 },
2235 )
2236 .await;
2237 assert!(ack.ok);
2238 }
2239
2240 #[tokio::test]
2241 async fn ambient_progress_with_mode_rejected() {
2242 let (server, _) = make_server();
2243 let ack = do_send(
2244 &server,
2245 "agent://orchestrator",
2246 Envelope {
2247 macp_version: "1.0".into(),
2248 mode: "macp.mode.decision.v1".into(),
2249 message_type: "Progress".into(),
2250 message_id: "prog-2".into(),
2251 session_id: String::new(),
2252 sender: String::new(),
2253 timestamp_unix_ms: Utc::now().timestamp_millis(),
2254 payload: vec![],
2255 },
2256 )
2257 .await;
2258 assert!(!ack.ok);
2259 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2260 }
2261
2262 #[tokio::test]
2263 async fn manifest_advertises_stream_enabled() {
2264 let (server, _) = make_server();
2265 let resp = server
2266 .initialize(Request::new(InitializeRequest {
2267 supported_protocol_versions: vec!["1.0".into()],
2268 client_info: None,
2269 capabilities: None,
2270 }))
2271 .await
2272 .unwrap();
2273 let caps = resp.into_inner().capabilities.unwrap();
2274 assert!(caps.sessions.unwrap().stream);
2275 }
2276
2277 #[tokio::test]
2278 async fn initialize_empty_versions_rejected() {
2279 let (server, _) = make_server();
2280 let err = server
2281 .initialize(Request::new(InitializeRequest {
2282 supported_protocol_versions: vec![],
2283 client_info: None,
2284 capabilities: None,
2285 }))
2286 .await
2287 .unwrap_err();
2288 assert_eq!(err.code(), tonic::Code::InvalidArgument);
2289 }
2290
2291 #[tokio::test]
2292 async fn initialize_unsupported_version_rejected() {
2293 let (server, _) = make_server();
2294 let err = server
2295 .initialize(Request::new(InitializeRequest {
2296 supported_protocol_versions: vec!["2.0".into()],
2297 client_info: None,
2298 capabilities: None,
2299 }))
2300 .await
2301 .unwrap_err();
2302 assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2303 }
2304
2305 fn observer_identity(sender: &str) -> AuthIdentity {
2308 AuthIdentity {
2309 sender: sender.into(),
2310 allowed_modes: None,
2311 can_start_sessions: false,
2312 max_open_sessions: None,
2313 can_manage_mode_registry: false,
2314 is_observer: true,
2315 }
2316 }
2317
2318 fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2319 StreamSessionRequest {
2320 subscribe_session_id: session_id.into(),
2321 after_sequence: after,
2322 envelope: None,
2323 }
2324 }
2325
2326 fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2327 SessionStartPayload {
2328 intent: "intent".into(),
2329 participants,
2330 mode_version: "1.0.0".into(),
2331 configuration_version: "cfg-1".into(),
2332 policy_version: String::new(),
2333 ttl_ms: 60_000,
2334 context_id: String::new(),
2335 extensions: std::collections::HashMap::new(),
2336 roots: vec![],
2337 max_suspend_ms: 0,
2338 }
2339 .encode_to_vec()
2340 }
2341
2342 async fn start_session(
2343 server: &MacpServer,
2344 initiator: &str,
2345 sid: &str,
2346 participants: Vec<String>,
2347 ) {
2348 let ack = do_send(
2349 server,
2350 initiator,
2351 Envelope {
2352 macp_version: "1.0".into(),
2353 mode: "macp.mode.decision.v1".into(),
2354 message_type: "SessionStart".into(),
2355 message_id: "start".into(),
2356 session_id: sid.into(),
2357 sender: String::new(),
2358 timestamp_unix_ms: Utc::now().timestamp_millis(),
2359 payload: start_multi_participant(participants),
2360 },
2361 )
2362 .await;
2363 assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2364 }
2365
2366 async fn send_proposal(
2367 server: &MacpServer,
2368 sender: &str,
2369 sid: &str,
2370 message_id: &str,
2371 proposal_id: &str,
2372 ) {
2373 let payload = crate::decision_pb::ProposalPayload {
2374 proposal_id: proposal_id.into(),
2375 option: "opt".into(),
2376 rationale: "r".into(),
2377 supporting_data: vec![],
2378 }
2379 .encode_to_vec();
2380 let ack = do_send(
2381 server,
2382 sender,
2383 Envelope {
2384 macp_version: "1.0".into(),
2385 mode: "macp.mode.decision.v1".into(),
2386 message_type: "Proposal".into(),
2387 message_id: message_id.into(),
2388 session_id: sid.into(),
2389 sender: String::new(),
2390 timestamp_unix_ms: Utc::now().timestamp_millis(),
2391 payload,
2392 },
2393 )
2394 .await;
2395 assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2396 }
2397
2398 #[tokio::test]
2399 async fn subscribe_replays_session_history_from_zero() {
2400 let (server, _) = make_server();
2401 let sid = new_sid();
2402 let initiator = "agent://orchestrator";
2403 let peer = "agent://fraud";
2404 start_session(
2405 &server,
2406 initiator,
2407 &sid,
2408 vec![initiator.into(), peer.into()],
2409 )
2410 .await;
2411 send_proposal(&server, peer, &sid, "m2", "p1").await;
2412
2413 let mut bound = None;
2414 let mut events = None;
2415 let replay = server
2416 .process_stream_request(
2417 &stream_identity(peer),
2418 subscribe_frame(&sid, 0),
2419 &mut bound,
2420 &mut events,
2421 )
2422 .await
2423 .unwrap();
2424
2425 assert_eq!(replay.len(), 2);
2426 assert_eq!(replay[0].message_type, "SessionStart");
2427 assert_eq!(replay[0].message_id, "start");
2428 assert_eq!(replay[1].message_type, "Proposal");
2429 assert_eq!(replay[1].message_id, "m2");
2430 assert_eq!(bound.as_deref(), Some(sid.as_str()));
2431 assert!(events.is_some());
2432 }
2433
2434 #[tokio::test]
2435 async fn subscribe_after_sequence_filters_history() {
2436 let (server, _) = make_server();
2437 let sid = new_sid();
2438 let initiator = "agent://orchestrator";
2439 let peer = "agent://fraud";
2440 start_session(
2441 &server,
2442 initiator,
2443 &sid,
2444 vec![initiator.into(), peer.into()],
2445 )
2446 .await;
2447 send_proposal(&server, peer, &sid, "m2", "p1").await;
2448 send_proposal(&server, peer, &sid, "m3", "p2").await;
2449
2450 let mut bound = None;
2451 let mut events = None;
2452 let replay = server
2453 .process_stream_request(
2454 &stream_identity(peer),
2455 subscribe_frame(&sid, 2),
2456 &mut bound,
2457 &mut events,
2458 )
2459 .await
2460 .unwrap();
2461
2462 assert_eq!(replay.len(), 1);
2463 assert_eq!(replay[0].message_id, "m3");
2464 }
2465
2466 #[tokio::test]
2467 async fn subscribe_unknown_session_returns_not_found() {
2468 let (server, _) = make_server();
2469 let mut bound = None;
2470 let mut events = None;
2471 let status = server
2472 .process_stream_request(
2473 &stream_identity("agent://orchestrator"),
2474 subscribe_frame("missing-session", 0),
2475 &mut bound,
2476 &mut events,
2477 )
2478 .await
2479 .unwrap_err();
2480 assert_eq!(status.code(), tonic::Code::NotFound);
2481 assert!(bound.is_none());
2482 assert!(events.is_none());
2483 }
2484
2485 #[tokio::test]
2486 async fn subscribe_non_participant_is_forbidden() {
2487 let (server, _) = make_server();
2488 let sid = new_sid();
2489 start_session(
2490 &server,
2491 "agent://orchestrator",
2492 &sid,
2493 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2494 )
2495 .await;
2496
2497 let mut bound = None;
2498 let mut events = None;
2499 let status = server
2500 .process_stream_request(
2501 &stream_identity("agent://outsider"),
2502 subscribe_frame(&sid, 0),
2503 &mut bound,
2504 &mut events,
2505 )
2506 .await
2507 .unwrap_err();
2508 assert_eq!(status.code(), tonic::Code::PermissionDenied);
2509 }
2510
2511 #[tokio::test]
2512 async fn subscribe_observer_identity_allowed() {
2513 let (server, _) = make_server();
2514 let sid = new_sid();
2515 start_session(
2516 &server,
2517 "agent://orchestrator",
2518 &sid,
2519 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2520 )
2521 .await;
2522
2523 let mut bound = None;
2524 let mut events = None;
2525 let replay = server
2526 .process_stream_request(
2527 &observer_identity("agent://auditor"),
2528 subscribe_frame(&sid, 0),
2529 &mut bound,
2530 &mut events,
2531 )
2532 .await
2533 .unwrap();
2534 assert_eq!(replay.len(), 1);
2535 assert_eq!(replay[0].message_type, "SessionStart");
2536 }
2537
2538 #[tokio::test]
2539 async fn subscribe_initiator_allowed_even_when_not_listed() {
2540 let (server, _) = make_server();
2543 let sid = new_sid();
2544 start_session(
2545 &server,
2546 "agent://orchestrator",
2547 &sid,
2548 vec!["agent://fraud".into()],
2549 )
2550 .await;
2551
2552 let mut bound = None;
2553 let mut events = None;
2554 let replay = server
2555 .process_stream_request(
2556 &stream_identity("agent://orchestrator"),
2557 subscribe_frame(&sid, 0),
2558 &mut bound,
2559 &mut events,
2560 )
2561 .await
2562 .unwrap();
2563 assert_eq!(replay.len(), 1);
2564 }
2565
2566 #[tokio::test]
2567 async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2568 let (server, _) = make_server();
2569 let sid = new_sid();
2570 let req = StreamSessionRequest {
2571 subscribe_session_id: sid.clone(),
2572 after_sequence: 0,
2573 envelope: Some(Envelope {
2574 macp_version: "1.0".into(),
2575 mode: "macp.mode.decision.v1".into(),
2576 message_type: "SessionStart".into(),
2577 message_id: "m1".into(),
2578 session_id: sid,
2579 sender: String::new(),
2580 timestamp_unix_ms: Utc::now().timestamp_millis(),
2581 payload: start_payload(),
2582 }),
2583 };
2584
2585 let mut bound = None;
2586 let mut events = None;
2587 let status = server
2588 .process_stream_request(
2589 &stream_identity("agent://orchestrator"),
2590 req,
2591 &mut bound,
2592 &mut events,
2593 )
2594 .await
2595 .unwrap_err();
2596 assert_eq!(status.code(), tonic::Code::InvalidArgument);
2597 }
2598
2599 #[tokio::test]
2600 async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2601 let (server, _) = make_server();
2602 let sid1 = new_sid();
2603 let sid2 = new_sid();
2604 start_session(
2605 &server,
2606 "agent://orchestrator",
2607 &sid1,
2608 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2609 )
2610 .await;
2611 start_session(
2612 &server,
2613 "agent://orchestrator",
2614 &sid2,
2615 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2616 )
2617 .await;
2618
2619 let identity = stream_identity("agent://fraud");
2621 let mut bound = None;
2622 let mut events = None;
2623 server
2624 .process_stream_request(
2625 &identity,
2626 subscribe_frame(&sid1, 0),
2627 &mut bound,
2628 &mut events,
2629 )
2630 .await
2631 .unwrap();
2632 assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2633
2634 let status = server
2636 .process_stream_request(
2637 &identity,
2638 subscribe_frame(&sid2, 0),
2639 &mut bound,
2640 &mut events,
2641 )
2642 .await
2643 .unwrap_err();
2644 assert_eq!(status.code(), tonic::Code::InvalidArgument);
2645 }
2646
2647 struct DenySenderEngine {
2651 denied: String,
2652 }
2653
2654 #[async_trait::async_trait]
2655 impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2656 async fn evaluate_session_start(
2657 &self,
2658 identity: &crate::security::AuthIdentity,
2659 _mode: &str,
2660 _env: &Envelope,
2661 ) -> macp_core::policy::PolicyDecision {
2662 if identity.sender == self.denied {
2663 macp_core::policy::PolicyDecision::Deny {
2664 reasons: vec!["sender embargoed".into()],
2665 }
2666 } else {
2667 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2668 }
2669 }
2670
2671 async fn evaluate_message(
2672 &self,
2673 identity: &crate::security::AuthIdentity,
2674 _session: &macp_core::session::Session,
2675 _env: &Envelope,
2676 ) -> macp_core::policy::PolicyDecision {
2677 if identity.sender == self.denied {
2678 macp_core::policy::PolicyDecision::Deny {
2679 reasons: vec!["sender embargoed".into()],
2680 }
2681 } else {
2682 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2683 }
2684 }
2685
2686 async fn evaluate_session_access(
2687 &self,
2688 identity: &crate::security::AuthIdentity,
2689 _session: &macp_core::session::Session,
2690 ) -> macp_core::policy::PolicyDecision {
2691 if identity.sender == self.denied {
2692 macp_core::policy::PolicyDecision::Deny {
2693 reasons: vec!["sender embargoed".into()],
2694 }
2695 } else {
2696 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2697 }
2698 }
2699 }
2700
2701 #[tokio::test]
2702 async fn policy_engine_gates_all_three_ingress_points() {
2703 let (server, _runtime) = make_server();
2704 let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2705 denied: "agent://embargoed".into(),
2706 }));
2707
2708 let sid = new_sid();
2709 let start_payload = SessionStartPayload {
2710 intent: "e3".into(),
2711 participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2712 mode_version: "1.0.0".into(),
2713 configuration_version: "cfg-1".into(),
2714 policy_version: String::new(),
2715 ttl_ms: 60_000,
2716 context_id: String::new(),
2717 extensions: Default::default(),
2718 roots: vec![],
2719 max_suspend_ms: 0,
2720 }
2721 .encode_to_vec();
2722 let start_env = |sender: &str, sid: &str| Envelope {
2723 macp_version: "1.0".into(),
2724 mode: "macp.mode.decision.v1".into(),
2725 message_type: "SessionStart".into(),
2726 message_id: new_sid(),
2727 session_id: sid.into(),
2728 sender: sender.into(),
2729 timestamp_unix_ms: Utc::now().timestamp_millis(),
2730 payload: start_payload.clone(),
2731 };
2732
2733 let ack = server
2735 .send(send_req(
2736 "agent://embargoed",
2737 start_env("agent://embargoed", &sid),
2738 ))
2739 .await
2740 .unwrap()
2741 .into_inner()
2742 .ack
2743 .unwrap();
2744 assert!(!ack.ok);
2745 assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2746
2747 let ack = server
2749 .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2750 .await
2751 .unwrap()
2752 .into_inner()
2753 .ack
2754 .unwrap();
2755 assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2756
2757 let proposal = crate::decision_pb::ProposalPayload {
2759 proposal_id: "p1".into(),
2760 option: "x".into(),
2761 rationale: "r".into(),
2762 supporting_data: vec![],
2763 }
2764 .encode_to_vec();
2765 let msg_env = Envelope {
2766 macp_version: "1.0".into(),
2767 mode: "macp.mode.decision.v1".into(),
2768 message_type: "Proposal".into(),
2769 message_id: new_sid(),
2770 session_id: sid.clone(),
2771 sender: "agent://embargoed".into(),
2772 timestamp_unix_ms: Utc::now().timestamp_millis(),
2773 payload: proposal,
2774 };
2775 let ack = server
2776 .send(send_req("agent://embargoed", msg_env))
2777 .await
2778 .unwrap()
2779 .into_inner()
2780 .ack
2781 .unwrap();
2782 assert!(!ack.ok);
2783 assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2784
2785 let mut req = Request::new(crate::pb::GetSessionRequest {
2787 session_id: sid.clone(),
2788 });
2789 req.metadata_mut()
2790 .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2791 let err = server
2792 .get_session(req)
2793 .await
2794 .expect_err("embargoed read must be denied");
2795 assert_eq!(err.code(), tonic::Code::PermissionDenied);
2796 }
2797
2798 #[tokio::test]
2802 async fn policy_engine_gates_stream_path() {
2803 let (server, runtime) = make_server();
2804 let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2805 denied: "agent://embargoed".into(),
2806 }));
2807
2808 let sid = new_sid();
2812 let payload = SessionStartPayload {
2813 intent: "e3-stream".into(),
2814 participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2815 mode_version: "1.0.0".into(),
2816 configuration_version: "cfg-1".into(),
2817 policy_version: String::new(),
2818 ttl_ms: 60_000,
2819 context_id: String::new(),
2820 extensions: Default::default(),
2821 roots: vec![],
2822 max_suspend_ms: 0,
2823 }
2824 .encode_to_vec();
2825 runtime
2826 .process(
2827 &Envelope {
2828 macp_version: "1.0".into(),
2829 mode: "macp.mode.decision.v1".into(),
2830 message_type: "SessionStart".into(),
2831 message_id: new_sid(),
2832 session_id: sid.clone(),
2833 sender: "agent://ok".into(),
2834 timestamp_unix_ms: Utc::now().timestamp_millis(),
2835 payload,
2836 },
2837 None,
2838 )
2839 .await
2840 .unwrap();
2841
2842 let embargoed = crate::security::AuthIdentity {
2843 sender: "agent://embargoed".into(),
2844 allowed_modes: None,
2845 can_start_sessions: true,
2846 max_open_sessions: None,
2847 can_manage_mode_registry: false,
2848 is_observer: false,
2849 };
2850 let mut bound = None;
2851 let mut events = None;
2852
2853 let proposal = crate::decision_pb::ProposalPayload {
2855 proposal_id: "p1".into(),
2856 option: "x".into(),
2857 rationale: "r".into(),
2858 supporting_data: vec![],
2859 }
2860 .encode_to_vec();
2861 let req = StreamSessionRequest {
2862 envelope: Some(Envelope {
2863 macp_version: "1.0".into(),
2864 mode: "macp.mode.decision.v1".into(),
2865 message_type: "Proposal".into(),
2866 message_id: new_sid(),
2867 session_id: sid.clone(),
2868 sender: "agent://embargoed".into(),
2869 timestamp_unix_ms: Utc::now().timestamp_millis(),
2870 payload: proposal,
2871 }),
2872 subscribe_session_id: String::new(),
2873 after_sequence: 0,
2874 };
2875 let err = server
2876 .process_stream_request(&embargoed, req, &mut bound, &mut events)
2877 .await
2878 .expect_err("stream envelope from embargoed sender must be denied");
2879 assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2882 assert!(err.message().contains("PolicyDenied"), "{err:?}");
2883
2884 let req = StreamSessionRequest {
2887 envelope: None,
2888 subscribe_session_id: sid.clone(),
2889 after_sequence: 0,
2890 };
2891 let err = server
2892 .process_stream_request(&embargoed, req, &mut bound, &mut events)
2893 .await
2894 .expect_err("stream subscribe from embargoed sender must be denied");
2895 assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2896 }
2897 fn paged_session(id: &str) -> crate::session::Session {
2900 crate::session::Session::builder(id, "macp.mode.decision.v1", "agent://initiator")
2901 .participants(vec!["agent://a".into()])
2902 .mode_version("1.0.0")
2903 .configuration_version("cfg-1")
2904 .started_at_unix_ms(1)
2905 .build()
2906 }
2907
2908 async fn seed_sessions(runtime: &Arc<Runtime>, ids: &[String]) {
2912 for id in ids {
2913 runtime
2914 .registry
2915 .insert_recovered_session(id.clone(), paged_session(id))
2916 .await;
2917 }
2918 }
2919
2920 fn list_sessions_req(page_size: i32, page_token: &str) -> Request<ListSessionsRequest> {
2921 let mut req = Request::new(ListSessionsRequest {
2922 page_size,
2923 page_token: page_token.to_string(),
2924 });
2925 req.metadata_mut()
2926 .insert("authorization", "Bearer agent://observer".parse().unwrap());
2927 req
2928 }
2929
2930 fn page_size_security(default: usize, max: usize) -> SecurityLayer {
2931 let mut security = SecurityLayer::dev_mode();
2932 security.list_sessions_default_page_size = default;
2933 security.list_sessions_max_page_size = max;
2934 security
2935 }
2936
2937 fn seed_ids(n: usize) -> Vec<String> {
2938 (0..n).map(|i| format!("session-{i:03}")).collect()
2939 }
2940
2941 #[tokio::test]
2942 async fn list_sessions_applies_default_page_size_when_zero() {
2943 let (server, runtime) = make_server_with_security(page_size_security(3, 1000));
2944 seed_sessions(&runtime, &seed_ids(10)).await;
2945
2946 let resp = server
2947 .list_sessions(list_sessions_req(0, ""))
2948 .await
2949 .unwrap()
2950 .into_inner();
2951 assert_eq!(resp.sessions.len(), 3);
2952 assert!(!resp.next_page_token.is_empty());
2953 }
2954
2955 #[tokio::test]
2956 async fn list_sessions_honors_explicit_page_size() {
2957 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
2958 seed_sessions(&runtime, &seed_ids(10)).await;
2959
2960 let resp = server
2961 .list_sessions(list_sessions_req(4, ""))
2962 .await
2963 .unwrap()
2964 .into_inner();
2965 assert_eq!(resp.sessions.len(), 4);
2966 assert!(!resp.next_page_token.is_empty());
2967 }
2968
2969 #[tokio::test]
2970 async fn list_sessions_clamps_page_size_above_max() {
2971 let (server, runtime) = make_server_with_security(page_size_security(100, 3));
2972 seed_sessions(&runtime, &seed_ids(10)).await;
2973
2974 let resp = server
2975 .list_sessions(list_sessions_req(1000, ""))
2976 .await
2977 .unwrap()
2978 .into_inner();
2979 assert_eq!(resp.sessions.len(), 3);
2980 assert!(!resp.next_page_token.is_empty());
2981 }
2982
2983 #[tokio::test]
2984 async fn list_sessions_rejects_negative_page_size() {
2985 let (server, runtime) = make_server();
2986 seed_sessions(&runtime, &seed_ids(3)).await;
2987
2988 let err = server
2989 .list_sessions(list_sessions_req(-1, ""))
2990 .await
2991 .unwrap_err();
2992 assert_eq!(err.code(), tonic::Code::InvalidArgument, "{err:?}");
2993 assert!(err.message().contains("page_size"), "{err:?}");
2994 }
2995
2996 #[tokio::test]
2997 async fn list_sessions_rejects_garbage_page_token() {
2998 use base64::Engine;
2999 let (server, runtime) = make_server();
3000 seed_sessions(&runtime, &seed_ids(3)).await;
3001
3002 let engine = base64::engine::general_purpose::URL_SAFE_NO_PAD;
3003 let valid = engine.encode("v1:session-000");
3004 let tokens = vec![
3005 "not-a-token!".to_string(),
3007 engine.encode("v2:session-000"),
3009 engine.encode("v1:"),
3011 engine.encode("v1"),
3017 valid[1..].to_string(),
3021 "A".repeat(2 * 1024 * 1024),
3023 ];
3024 for token in tokens {
3025 let err = server
3026 .list_sessions(list_sessions_req(0, &token))
3027 .await
3028 .unwrap_err();
3029 assert_eq!(err.code(), tonic::Code::InvalidArgument);
3030 assert_eq!(
3032 err.message(),
3033 "INVALID_ARGUMENT: page_token is not a valid continuation token"
3034 );
3035 }
3036 }
3037
3038 #[tokio::test]
3039 async fn list_sessions_full_traversal_visits_every_session_exactly_once() {
3040 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3041 let ids = seed_ids(25);
3042 seed_sessions(&runtime, &ids).await;
3043
3044 let mut collected: Vec<String> = Vec::new();
3045 let mut token = String::new();
3046 for _ in 0..100 {
3047 let resp = server
3048 .list_sessions(list_sessions_req(4, &token))
3049 .await
3050 .unwrap()
3051 .into_inner();
3052 collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3053 token = resp.next_page_token;
3054 if token.is_empty() {
3055 break;
3056 }
3057 }
3058 assert!(token.is_empty(), "traversal did not terminate");
3059 let unique: std::collections::HashSet<&String> = collected.iter().collect();
3060 assert_eq!(unique.len(), 25, "sessions were dropped or duplicated");
3063 assert_eq!(collected.len(), 25, "sessions were duplicated");
3064 }
3065
3066 #[tokio::test]
3067 async fn list_sessions_terminal_page_has_empty_next_page_token() {
3068 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3069 seed_sessions(&runtime, &seed_ids(10)).await;
3070
3071 let mut tokens: Vec<String> = Vec::new();
3072 let mut token = String::new();
3073 for _ in 0..20 {
3074 let resp = server
3075 .list_sessions(list_sessions_req(5, &token))
3076 .await
3077 .unwrap()
3078 .into_inner();
3079 token = resp.next_page_token;
3080 tokens.push(token.clone());
3081 if token.is_empty() {
3082 break;
3083 }
3084 }
3085 assert_eq!(tokens.len(), 2, "{tokens:?}");
3088 assert!(!tokens[0].is_empty());
3089 assert!(tokens[1].is_empty());
3090 }
3091
3092 #[tokio::test]
3093 async fn list_sessions_orders_by_session_id_ascending() {
3094 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3095 let ids: Vec<String> = ["delta", "alpha", "echo", "charlie", "bravo"]
3097 .iter()
3098 .map(|s| s.to_string())
3099 .collect();
3100 seed_sessions(&runtime, &ids).await;
3101
3102 let mut collected: Vec<String> = Vec::new();
3103 let mut token = String::new();
3104 loop {
3105 let resp = server
3106 .list_sessions(list_sessions_req(2, &token))
3107 .await
3108 .unwrap()
3109 .into_inner();
3110 collected.extend(resp.sessions.iter().map(|s| s.session_id.clone()));
3111 token = resp.next_page_token;
3112 if token.is_empty() {
3113 break;
3114 }
3115 }
3116 assert_eq!(
3118 collected,
3119 vec!["alpha", "bravo", "charlie", "delta", "echo"]
3120 );
3121 }
3122
3123 #[tokio::test]
3124 async fn list_sessions_still_requires_authentication() {
3125 let (server, runtime) = make_server();
3126 seed_sessions(&runtime, &seed_ids(3)).await;
3127
3128 let req = Request::new(ListSessionsRequest {
3131 page_size: -1,
3132 page_token: String::new(),
3133 });
3134 let err = server.list_sessions(req).await.unwrap_err();
3135 assert_eq!(err.code(), tonic::Code::Unauthenticated, "{err:?}");
3136 }
3137
3138 #[tokio::test]
3139 async fn list_sessions_tolerates_cursor_for_removed_session() {
3140 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3141 let ids = seed_ids(4);
3142 seed_sessions(&runtime, &ids).await;
3143
3144 let first = server
3145 .list_sessions(list_sessions_req(1, ""))
3146 .await
3147 .unwrap()
3148 .into_inner();
3149 assert_eq!(first.sessions[0].session_id, "session-000");
3150 assert!(!first.next_page_token.is_empty());
3151
3152 runtime
3155 .registry
3156 .sessions
3157 .write()
3158 .await
3159 .remove("session-000");
3160
3161 let second = server
3162 .list_sessions(list_sessions_req(1, &first.next_page_token))
3163 .await
3164 .unwrap()
3165 .into_inner();
3166 assert_eq!(second.sessions[0].session_id, "session-001");
3167 }
3168
3169 #[tokio::test]
3170 async fn list_sessions_cursor_comes_from_the_id_list_not_the_returned_sessions() {
3171 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3182 seed_sessions(&runtime, &seed_ids(6)).await;
3183
3184 let first = runtime.registry.get_shared("session-000").await.unwrap();
3188 let guard = first.lock().await;
3189
3190 let handler = server.list_sessions(list_sessions_req(3, ""));
3191 let mutator = async {
3192 let mut spins = 0;
3197 while Arc::strong_count(&first) < 3 {
3198 assert!(spins < 10_000, "handler never parked on the session mutex");
3199 spins += 1;
3200 tokio::task::yield_now().await;
3201 }
3202 runtime
3203 .registry
3204 .sessions
3205 .write()
3206 .await
3207 .remove("session-002");
3208 drop(guard);
3209 };
3210 let (resp, ()) = tokio::join!(handler, mutator);
3211 let resp = resp.unwrap().into_inner();
3212
3213 assert_eq!(
3215 resp.sessions.len(),
3216 2,
3217 "expected session-002 to vanish between the scan and the fetch"
3218 );
3219 assert_eq!(resp.sessions[1].session_id, "session-001");
3220 assert_eq!(
3222 crate::pagination::decode_page_token(&resp.next_page_token),
3223 Ok("session-002".to_string()),
3224 "cursor was derived from the returned sessions, not the ID list"
3225 );
3226
3227 runtime
3229 .registry
3230 .insert_recovered_session("session-002".to_string(), paged_session("session-002"))
3231 .await;
3232 let second = server
3233 .list_sessions(list_sessions_req(3, &resp.next_page_token))
3234 .await
3235 .unwrap()
3236 .into_inner();
3237 assert_eq!(
3238 second.sessions[0].session_id, "session-003",
3239 "the cursor moved backwards past an ID the page had already accounted for"
3240 );
3241 }
3242
3243 #[tokio::test]
3244 async fn list_sessions_replaying_a_token_returns_the_identical_page() {
3245 let (server, runtime) = make_server_with_security(page_size_security(100, 1000));
3246 seed_sessions(&runtime, &seed_ids(10)).await;
3247
3248 let first = server
3249 .list_sessions(list_sessions_req(3, ""))
3250 .await
3251 .unwrap()
3252 .into_inner();
3253 let token = first.next_page_token;
3254 assert!(!token.is_empty());
3255
3256 let page_a = server
3257 .list_sessions(list_sessions_req(3, &token))
3258 .await
3259 .unwrap()
3260 .into_inner();
3261 let page_b = server
3262 .list_sessions(list_sessions_req(3, &token))
3263 .await
3264 .unwrap()
3265 .into_inner();
3266
3267 let ids_a: Vec<&str> = page_a.sessions.iter().map(|s| &*s.session_id).collect();
3268 let ids_b: Vec<&str> = page_b.sessions.iter().map(|s| &*s.session_id).collect();
3269 assert_eq!(ids_a, ids_b);
3270 assert_eq!(page_a.next_page_token, page_b.next_page_token);
3271 }
3272
3273 #[tokio::test]
3274 async fn list_sessions_survives_zero_effective_page_size() {
3275 let (server, runtime) = make_server_with_security(page_size_security(0, 0));
3279 seed_sessions(&runtime, &seed_ids(3)).await;
3280
3281 let resp = server
3282 .list_sessions(list_sessions_req(0, ""))
3283 .await
3284 .unwrap()
3285 .into_inner();
3286 assert!(
3287 !resp.sessions.is_empty(),
3288 "empty page with token {:?} — the traversal terminates and ListSessions returns nothing",
3289 resp.next_page_token
3290 );
3291 assert_eq!(resp.sessions.len(), 1);
3292 assert!(!resp.next_page_token.is_empty());
3293
3294 let next = server
3296 .list_sessions(list_sessions_req(0, &resp.next_page_token))
3297 .await
3298 .unwrap()
3299 .into_inner();
3300 assert_eq!(next.sessions.len(), 1);
3301 assert_ne!(next.sessions[0].session_id, resp.sessions[0].session_id);
3302 }
3303}