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