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