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