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: "0.4.0".into(),
813 description: "Reference implementation of the Multi-Agent Coordination Protocol"
814 .into(),
815 website_url: String::new(),
816 }),
817 capabilities: Some(Capabilities {
818 sessions: Some(SessionsCapability { stream: true, list_sessions: true, watch_sessions: true }),
819 cancellation: Some(CancellationCapability {
820 cancel_session: true,
821 }),
822 progress: Some(ProgressCapability { progress: true }),
823 manifest: Some(ManifestCapability { get_manifest: true }),
824 mode_registry: Some(ModeRegistryCapability {
825 list_modes: true,
826 list_changed: true,
827 }),
828 roots: Some(RootsCapability {
829 list_roots: true,
835 list_changed: false,
836 }),
837 policy_registry: Some(PolicyRegistryCapability {
838 register_policy: !self.policies_read_only,
839 list_policies: true,
840 list_changed: true,
841 }),
842 experimental: Some(crate::pb::ExperimentalCapabilities {
843 features: HashMap::from([
844 ("ext_mode_lifecycle".into(), "true".into()),
845 ]),
846 }),
847 }),
848 supported_modes: self.runtime.registered_mode_names(),
849 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(),
850 }))
851 }
852
853 async fn send(&self, request: Request<SendRequest>) -> Result<Response<SendResponse>, Status> {
854 let env = request
855 .get_ref()
856 .envelope
857 .clone()
858 .ok_or_else(|| Status::invalid_argument("SendRequest must contain an envelope"))?;
859
860 let result = async {
861 self.validate_envelope_shape(&env)?;
862 let (env, max_open) = self.authenticate_send_request(&request, env).await?;
863 self.runtime
864 .process(&env, max_open)
865 .await
866 .map(|process_result| (env, process_result))
867 }
868 .await;
869
870 let ack = match result {
871 Ok((env, process_result)) => Ack {
872 ok: true,
873 duplicate: process_result.duplicate,
874 message_id: env.message_id.clone(),
875 session_id: env.session_id.clone(),
876 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
877 session_state: Self::session_state_to_pb(&process_result.session_state),
878 error: None,
879 },
880 Err(err) => {
881 let env = request.get_ref().envelope.clone().unwrap_or_default();
882 if !env.session_id.is_empty() {
886 self.runtime.metrics().record_message_rejected(&env.mode);
887 if env.message_type == "Commitment" {
888 self.runtime.metrics().record_commitment_rejected(&env.mode);
889 }
890 }
891 Self::make_error_ack(&err, &env)
892 }
893 };
894
895 Ok(Response::new(SendResponse { ack: Some(ack) }))
896 }
897
898 async fn get_session(
899 &self,
900 request: Request<GetSessionRequest>,
901 ) -> Result<Response<GetSessionResponse>, Status> {
902 let session_id = request.get_ref().session_id.clone();
903 let _identity = self
904 .authenticate_session_access(&request, &session_id)
905 .await?;
906 let session = self
907 .runtime
908 .get_session_checked(&session_id)
909 .await
910 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
911
912 Ok(Response::new(GetSessionResponse {
913 metadata: Some(Self::session_to_metadata(&session)),
914 }))
915 }
916
917 async fn cancel_session(
918 &self,
919 request: Request<CancelSessionRequest>,
920 ) -> Result<Response<CancelSessionResponse>, Status> {
921 let session_id = request.get_ref().session_id.clone();
922 let identity = self
923 .security
924 .authenticate_metadata(request.metadata())
925 .await
926 .map_err(Self::status_from_error)?;
927 let session = self
928 .runtime
929 .get_session_checked(&session_id)
930 .await
931 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
932 if identity.sender != session.initiator_sender
935 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
936 {
937 return Err(Status::permission_denied(
938 "FORBIDDEN: only the session initiator or policy-delegated roles can cancel",
939 ));
940 }
941 let sender = identity.sender.clone();
942 let req = request.into_inner();
943 match self
944 .runtime
945 .cancel_session(&req.session_id, &req.reason, &sender)
946 .await
947 {
948 Ok(result) => Ok(Response::new(CancelSessionResponse {
949 ack: Some(Ack {
950 ok: true,
951 duplicate: false,
952 message_id: String::new(),
953 session_id: req.session_id,
954 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
955 session_state: Self::session_state_to_pb(&result.session_state),
956 error: None,
957 }),
958 })),
959 Err(err) => Ok(Response::new(CancelSessionResponse {
960 ack: Some(Ack {
961 ok: false,
962 duplicate: false,
963 message_id: String::new(),
964 session_id: req.session_id.clone(),
965 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
966 session_state: PbSessionState::Unspecified.into(),
967 error: Some(PbMacpError {
968 code: err.error_code().into(),
969 message: err.to_string(),
970 session_id: req.session_id,
971 message_id: String::new(),
972 details: vec![],
973 }),
974 }),
975 })),
976 }
977 }
978
979 async fn suspend_session(
980 &self,
981 request: Request<SuspendSessionRequest>,
982 ) -> Result<Response<SuspendSessionResponse>, Status> {
983 let session_id = request.get_ref().session_id.clone();
984 let identity = self
985 .security
986 .authenticate_metadata(request.metadata())
987 .await
988 .map_err(Self::status_from_error)?;
989 let session = self
990 .runtime
991 .get_session_checked(&session_id)
992 .await
993 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
994 if identity.sender != session.initiator_sender
997 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
998 {
999 return Err(Status::permission_denied(
1000 "FORBIDDEN: only the session initiator or policy-delegated roles can suspend",
1001 ));
1002 }
1003 let sender = identity.sender.clone();
1004 let req = request.into_inner();
1005 match self
1006 .runtime
1007 .suspend_session(&req.session_id, &req.reason, &sender)
1008 .await
1009 {
1010 Ok(result) => Ok(Response::new(SuspendSessionResponse {
1011 ack: Some(Ack {
1012 ok: true,
1013 duplicate: false,
1014 message_id: String::new(),
1015 session_id: req.session_id,
1016 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1017 session_state: Self::session_state_to_pb(&result.session_state),
1018 error: None,
1019 }),
1020 })),
1021 Err(err) => Ok(Response::new(SuspendSessionResponse {
1022 ack: Some(Ack {
1023 ok: false,
1024 duplicate: false,
1025 message_id: String::new(),
1026 session_id: req.session_id.clone(),
1027 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1028 session_state: PbSessionState::Unspecified.into(),
1029 error: Some(PbMacpError {
1030 code: err.error_code().into(),
1031 message: err.to_string(),
1032 session_id: req.session_id,
1033 message_id: String::new(),
1034 details: vec![],
1035 }),
1036 }),
1037 })),
1038 }
1039 }
1040
1041 async fn resume_session(
1042 &self,
1043 request: Request<ResumeSessionRequest>,
1044 ) -> Result<Response<ResumeSessionResponse>, Status> {
1045 let session_id = request.get_ref().session_id.clone();
1046 let identity = self
1047 .security
1048 .authenticate_metadata(request.metadata())
1049 .await
1050 .map_err(Self::status_from_error)?;
1051 let session = self
1052 .runtime
1053 .get_session_checked(&session_id)
1054 .await
1055 .ok_or_else(|| Status::not_found(format!("Session '{}' not found", session_id)))?;
1056 if identity.sender != session.initiator_sender
1057 && crate::mode::util::check_commitment_authority(&session, &identity.sender).is_err()
1058 {
1059 return Err(Status::permission_denied(
1060 "FORBIDDEN: only the session initiator or policy-delegated roles can resume",
1061 ));
1062 }
1063 let sender = identity.sender.clone();
1064 let req = request.into_inner();
1065 match self
1066 .runtime
1067 .resume_session(&req.session_id, &req.reason, &sender)
1068 .await
1069 {
1070 Ok(result) => Ok(Response::new(ResumeSessionResponse {
1071 ack: Some(Ack {
1072 ok: true,
1073 duplicate: false,
1074 message_id: String::new(),
1075 session_id: req.session_id,
1076 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1077 session_state: Self::session_state_to_pb(&result.session_state),
1078 error: None,
1079 }),
1080 })),
1081 Err(err) => Ok(Response::new(ResumeSessionResponse {
1082 ack: Some(Ack {
1083 ok: false,
1084 duplicate: false,
1085 message_id: String::new(),
1086 session_id: req.session_id.clone(),
1087 accepted_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1088 session_state: PbSessionState::Unspecified.into(),
1089 error: Some(PbMacpError {
1090 code: err.error_code().into(),
1091 message: err.to_string(),
1092 session_id: req.session_id,
1093 message_id: String::new(),
1094 details: vec![],
1095 }),
1096 }),
1097 })),
1098 }
1099 }
1100
1101 async fn get_manifest(
1102 &self,
1103 request: Request<GetManifestRequest>,
1104 ) -> Result<Response<GetManifestResponse>, Status> {
1105 let req = request.into_inner();
1106 if !req.agent_id.is_empty() && req.agent_id != "macp-runtime" {
1107 return Err(Status::not_found(format!(
1108 "Agent '{}' not found",
1109 req.agent_id
1110 )));
1111 }
1112
1113 Ok(Response::new(GetManifestResponse {
1114 manifest: Some(crate::pb::AgentManifest {
1115 agent_id: "macp-runtime".into(),
1116 title: "MACP Reference Runtime".into(),
1117 description: "Reference implementation of MACP".into(),
1118 supported_modes: self.runtime.registered_mode_names(),
1119 input_content_types: vec!["application/macp-envelope+proto".into()],
1120 output_content_types: vec!["application/macp-envelope+proto".into()],
1121 metadata: HashMap::new(),
1122 transport_endpoints: vec![],
1124 }),
1125 }))
1126 }
1127
1128 async fn list_modes(
1129 &self,
1130 _request: Request<ListModesRequest>,
1131 ) -> Result<Response<ListModesResponse>, Status> {
1132 Ok(Response::new(ListModesResponse {
1133 modes: self.runtime.standard_mode_descriptors(),
1134 }))
1135 }
1136
1137 async fn list_roots(
1138 &self,
1139 _request: Request<ListRootsRequest>,
1140 ) -> Result<Response<ListRootsResponse>, Status> {
1141 Ok(Response::new(ListRootsResponse { roots: vec![] }))
1142 }
1143
1144 type StreamSessionStream = SessionResponseStream;
1145
1146 async fn stream_session(
1147 &self,
1148 request: Request<tonic::Streaming<StreamSessionRequest>>,
1149 ) -> Result<Response<Self::StreamSessionStream>, Status> {
1150 let identity = self
1151 .security
1152 .authenticate_metadata(request.metadata())
1153 .await
1154 .map_err(Self::status_from_error)?;
1155 let inbound = request.into_inner();
1156 Ok(Response::new(
1157 self.build_stream_session_stream(identity, inbound),
1158 ))
1159 }
1160
1161 type WatchModeRegistryStream = std::pin::Pin<
1162 Box<dyn futures_core::Stream<Item = Result<WatchModeRegistryResponse, Status>> + Send>,
1163 >;
1164
1165 async fn watch_mode_registry(
1166 &self,
1167 _request: Request<WatchModeRegistryRequest>,
1168 ) -> Result<Response<Self::WatchModeRegistryStream>, Status> {
1169 let mut rx = self.runtime.subscribe_mode_changes();
1170 let stream = async_stream::try_stream! {
1171 yield WatchModeRegistryResponse {
1173 change: Some(crate::pb::RegistryChanged {
1174 registry: "modes".into(),
1175 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1176 }),
1177 };
1178 while rx.recv().await.is_ok() {
1180 yield WatchModeRegistryResponse {
1181 change: Some(crate::pb::RegistryChanged {
1182 registry: "modes".into(),
1183 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1184 }),
1185 };
1186 }
1187 };
1188 Ok(Response::new(Box::pin(stream)))
1189 }
1190
1191 type WatchRootsStream = std::pin::Pin<
1192 Box<dyn futures_core::Stream<Item = Result<WatchRootsResponse, Status>> + Send>,
1193 >;
1194
1195 async fn watch_roots(
1196 &self,
1197 _request: Request<WatchRootsRequest>,
1198 ) -> Result<Response<Self::WatchRootsStream>, Status> {
1199 let initial = WatchRootsResponse {
1200 change: Some(crate::pb::RootsChanged {
1201 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1202 }),
1203 };
1204 let stream = async_stream::try_stream! {
1205 yield initial;
1206 std::future::pending::<()>().await;
1208 };
1209 Ok(Response::new(Box::pin(stream)))
1210 }
1211
1212 type WatchSignalsStream = std::pin::Pin<
1213 Box<dyn futures_core::Stream<Item = Result<WatchSignalsResponse, Status>> + Send>,
1214 >;
1215
1216 type WatchSessionsStream = std::pin::Pin<
1217 Box<dyn futures_core::Stream<Item = Result<WatchSessionsResponse, Status>> + Send>,
1218 >;
1219
1220 async fn watch_signals(
1221 &self,
1222 request: Request<WatchSignalsRequest>,
1223 ) -> Result<Response<Self::WatchSignalsStream>, Status> {
1224 let _identity = self
1229 .security
1230 .authenticate_metadata(request.metadata())
1231 .await
1232 .map_err(Self::status_from_error)?;
1233 let mut rx = self.runtime.subscribe_signals();
1234 let stream = async_stream::try_stream! {
1235 loop {
1236 match rx.recv().await {
1237 Ok(envelope) => {
1238 yield WatchSignalsResponse {
1239 envelope: Some(envelope),
1240 };
1241 }
1242 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1246 Err(Status::resource_exhausted(format!(
1247 "WatchSignals receiver fell behind by {skipped} signals"
1248 )))?;
1249 }
1250 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1251 }
1252 }
1253 };
1254 Ok(Response::new(Box::pin(stream)))
1255 }
1256
1257 async fn list_sessions(
1260 &self,
1261 request: Request<ListSessionsRequest>,
1262 ) -> Result<Response<ListSessionsResponse>, Status> {
1263 let _identity = self
1264 .security
1265 .authenticate_metadata(request.metadata())
1266 .await
1267 .map_err(Self::status_from_error)?;
1268 let sessions = self.runtime.registry.get_all_sessions().await;
1269 let metadata: Vec<SessionMetadata> =
1270 sessions.iter().map(Self::session_to_metadata).collect();
1271 Ok(Response::new(ListSessionsResponse { sessions: metadata }))
1272 }
1273
1274 async fn watch_sessions(
1275 &self,
1276 request: Request<WatchSessionsRequest>,
1277 ) -> Result<Response<Self::WatchSessionsStream>, Status> {
1278 let _identity = self
1279 .security
1280 .authenticate_metadata(request.metadata())
1281 .await
1282 .map_err(Self::status_from_error)?;
1283 let mut rx = self.runtime.subscribe_session_lifecycle();
1284 let runtime = Arc::clone(&self.runtime);
1285 let stream = async_stream::try_stream! {
1286 let sessions = runtime.registry.get_all_sessions().await;
1292 let mut synced: std::collections::HashSet<String> =
1293 std::collections::HashSet::with_capacity(sessions.len());
1294 for session in &sessions {
1295 synced.insert(session.session_id.clone());
1296 yield WatchSessionsResponse {
1297 event: Some(SessionLifecycleEvent {
1298 event_type: session_lifecycle_event::EventType::Created.into(),
1299 session: Some(Self::session_to_metadata(session)),
1300 observed_at_unix_ms: session.started_at_unix_ms,
1301 }),
1302 };
1303 }
1304 loop {
1306 let event = match rx.recv().await {
1307 Ok(event) => event,
1308 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
1309 Err(Status::resource_exhausted(format!(
1310 "WatchSessions receiver fell behind by {skipped} events"
1311 )))?;
1312 break;
1313 }
1314 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
1315 };
1316 let (event_type, sid) = match &event {
1317 crate::runtime::SessionLifecycleEvent::Created { session_id } =>
1318 (session_lifecycle_event::EventType::Created, session_id.clone()),
1319 crate::runtime::SessionLifecycleEvent::Resolved { session_id } =>
1320 (session_lifecycle_event::EventType::Resolved, session_id.clone()),
1321 crate::runtime::SessionLifecycleEvent::Expired { session_id } =>
1322 (session_lifecycle_event::EventType::Expired, session_id.clone()),
1323 crate::runtime::SessionLifecycleEvent::Suspended { session_id } =>
1324 (session_lifecycle_event::EventType::Suspended, session_id.clone()),
1325 crate::runtime::SessionLifecycleEvent::Resumed { session_id } =>
1326 (session_lifecycle_event::EventType::Resumed, session_id.clone()),
1327 crate::runtime::SessionLifecycleEvent::Cancelled { session_id } =>
1328 (session_lifecycle_event::EventType::Cancelled, session_id.clone()),
1329 };
1330 if event_type == session_lifecycle_event::EventType::Created
1334 && !synced.insert(sid.clone())
1335 {
1336 continue;
1337 }
1338 let session_meta = runtime.registry.get_session(&sid).await
1339 .map(|s| Self::session_to_metadata(&s));
1340 yield WatchSessionsResponse {
1341 event: Some(SessionLifecycleEvent {
1342 event_type: event_type.into(),
1343 session: session_meta,
1344 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1345 }),
1346 };
1347 }
1348 };
1349 Ok(Response::new(Box::pin(stream)))
1350 }
1351
1352 async fn list_ext_modes(
1355 &self,
1356 _request: Request<ListExtModesRequest>,
1357 ) -> Result<Response<ListExtModesResponse>, Status> {
1358 Ok(Response::new(ListExtModesResponse {
1359 modes: self.runtime.extension_mode_descriptors(),
1360 }))
1361 }
1362
1363 async fn register_ext_mode(
1364 &self,
1365 request: Request<RegisterExtModeRequest>,
1366 ) -> Result<Response<RegisterExtModeResponse>, Status> {
1367 let identity = self
1368 .security
1369 .authenticate_metadata(request.metadata())
1370 .await
1371 .map_err(Self::status_from_error)?;
1372 self.security
1373 .authorize_mode_registry(&identity)
1374 .map_err(Self::status_from_error)?;
1375 let req = request.into_inner();
1376 let descriptor = req
1377 .mode_descriptor
1378 .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1379 match self.runtime.register_extension(descriptor) {
1380 Ok(()) => Ok(Response::new(RegisterExtModeResponse {
1381 ok: true,
1382 error: String::new(),
1383 })),
1384 Err(e) => Ok(Response::new(RegisterExtModeResponse {
1385 ok: false,
1386 error: e,
1387 })),
1388 }
1389 }
1390
1391 async fn unregister_ext_mode(
1392 &self,
1393 request: Request<UnregisterExtModeRequest>,
1394 ) -> Result<Response<UnregisterExtModeResponse>, Status> {
1395 let identity = self
1396 .security
1397 .authenticate_metadata(request.metadata())
1398 .await
1399 .map_err(Self::status_from_error)?;
1400 self.security
1401 .authorize_mode_registry(&identity)
1402 .map_err(Self::status_from_error)?;
1403 let req = request.into_inner();
1404 match self.runtime.unregister_extension(&req.mode) {
1405 Ok(()) => Ok(Response::new(UnregisterExtModeResponse {
1406 ok: true,
1407 error: String::new(),
1408 })),
1409 Err(e) => Ok(Response::new(UnregisterExtModeResponse {
1410 ok: false,
1411 error: e,
1412 })),
1413 }
1414 }
1415
1416 async fn promote_mode(
1417 &self,
1418 request: Request<PromoteModeRequest>,
1419 ) -> Result<Response<PromoteModeResponse>, Status> {
1420 let identity = self
1421 .security
1422 .authenticate_metadata(request.metadata())
1423 .await
1424 .map_err(Self::status_from_error)?;
1425 self.security
1426 .authorize_mode_registry(&identity)
1427 .map_err(Self::status_from_error)?;
1428 let req = request.into_inner();
1429 let new_name = if req.promoted_mode_name.is_empty() {
1430 None
1431 } else {
1432 Some(req.promoted_mode_name.as_str())
1433 };
1434 match self.runtime.promote_mode(&req.mode, new_name) {
1435 Ok(final_name) => Ok(Response::new(PromoteModeResponse {
1436 ok: true,
1437 error: String::new(),
1438 mode: final_name,
1439 })),
1440 Err(e) => Ok(Response::new(PromoteModeResponse {
1441 ok: false,
1442 error: e,
1443 mode: String::new(),
1444 })),
1445 }
1446 }
1447
1448 async fn register_policy(
1451 &self,
1452 request: Request<RegisterPolicyRequest>,
1453 ) -> Result<Response<RegisterPolicyResponse>, Status> {
1454 if self.policies_read_only {
1455 return Err(Status::failed_precondition(
1456 "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1457 ));
1458 }
1459 let identity = self
1460 .security
1461 .authenticate_metadata(request.metadata())
1462 .await
1463 .map_err(Self::status_from_error)?;
1464 self.security
1465 .authorize_mode_registry(&identity)
1466 .map_err(Self::status_from_error)?;
1467 let req = request.into_inner();
1468 let descriptor = req
1469 .policy_descriptor
1470 .ok_or_else(|| Status::invalid_argument("descriptor required"))?;
1471 let definition = Self::policy_descriptor_to_definition(&descriptor);
1472 match self.runtime.register_policy(definition) {
1473 Ok(()) => Ok(Response::new(RegisterPolicyResponse {
1474 ok: true,
1475 error: String::new(),
1476 })),
1477 Err(e) => Ok(Response::new(RegisterPolicyResponse {
1478 ok: false,
1479 error: e,
1480 })),
1481 }
1482 }
1483
1484 async fn unregister_policy(
1485 &self,
1486 request: Request<UnregisterPolicyRequest>,
1487 ) -> Result<Response<UnregisterPolicyResponse>, Status> {
1488 if self.policies_read_only {
1489 return Err(Status::failed_precondition(
1490 "policy registry is read-only: policies are file-loaded via MACP_POLICIES_DIR",
1491 ));
1492 }
1493 let identity = self
1494 .security
1495 .authenticate_metadata(request.metadata())
1496 .await
1497 .map_err(Self::status_from_error)?;
1498 self.security
1499 .authorize_mode_registry(&identity)
1500 .map_err(Self::status_from_error)?;
1501 let req = request.into_inner();
1502 match self.runtime.unregister_policy(&req.policy_id) {
1503 Ok(()) => Ok(Response::new(UnregisterPolicyResponse {
1504 ok: true,
1505 error: String::new(),
1506 })),
1507 Err(e) => Ok(Response::new(UnregisterPolicyResponse {
1508 ok: false,
1509 error: e,
1510 })),
1511 }
1512 }
1513
1514 async fn get_policy(
1515 &self,
1516 request: Request<GetPolicyRequest>,
1517 ) -> Result<Response<GetPolicyResponse>, Status> {
1518 let _identity = self
1519 .security
1520 .authenticate_metadata(request.metadata())
1521 .await
1522 .map_err(Self::status_from_error)?;
1523 let req = request.into_inner();
1524 let policy = self
1525 .runtime
1526 .get_policy(&req.policy_id)
1527 .ok_or_else(|| Status::not_found(format!("Policy '{}' not found", req.policy_id)))?;
1528 Ok(Response::new(GetPolicyResponse {
1529 policy_descriptor: Some(Self::policy_definition_to_descriptor(&policy)),
1530 }))
1531 }
1532
1533 async fn list_policies(
1534 &self,
1535 request: Request<ListPoliciesRequest>,
1536 ) -> Result<Response<ListPoliciesResponse>, Status> {
1537 let _identity = self
1538 .security
1539 .authenticate_metadata(request.metadata())
1540 .await
1541 .map_err(Self::status_from_error)?;
1542 let req = request.into_inner();
1543 let mode_filter = if req.mode.is_empty() {
1544 None
1545 } else {
1546 Some(req.mode.as_str())
1547 };
1548 let policies = self.runtime.list_policies(mode_filter);
1549 let descriptors = policies
1550 .iter()
1551 .map(Self::policy_definition_to_descriptor)
1552 .collect();
1553 Ok(Response::new(ListPoliciesResponse { descriptors }))
1554 }
1555
1556 type WatchPoliciesStream = std::pin::Pin<
1557 Box<dyn futures_core::Stream<Item = Result<WatchPoliciesResponse, Status>> + Send>,
1558 >;
1559
1560 async fn watch_policies(
1561 &self,
1562 _request: Request<WatchPoliciesRequest>,
1563 ) -> Result<Response<Self::WatchPoliciesStream>, Status> {
1564 let mut rx = self.runtime.subscribe_policy_changes();
1565 let runtime = Arc::clone(&self.runtime);
1566 let stream = async_stream::try_stream! {
1567 let policies = runtime.list_policies(None);
1569 let descriptors: Vec<PolicyDescriptor> = policies
1570 .iter()
1571 .map(MacpServer::policy_definition_to_descriptor)
1572 .collect();
1573 yield WatchPoliciesResponse {
1574 descriptors,
1575 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1576 };
1577 while rx.recv().await.is_ok() {
1579 let policies = runtime.list_policies(None);
1580 let descriptors: Vec<PolicyDescriptor> = policies
1581 .iter()
1582 .map(MacpServer::policy_definition_to_descriptor)
1583 .collect();
1584 yield WatchPoliciesResponse {
1585 descriptors,
1586 observed_at_unix_ms: chrono::Utc::now().timestamp_millis(),
1587 };
1588 }
1589 };
1590 Ok(Response::new(Box::pin(stream)))
1591 }
1592}
1593
1594impl MacpServer {
1597 fn policy_descriptor_to_definition(
1598 descriptor: &PolicyDescriptor,
1599 ) -> crate::policy::PolicyDefinition {
1600 let rules: serde_json::Value = if descriptor.rules.is_empty() {
1601 serde_json::json!({})
1602 } else {
1603 serde_json::from_str(&descriptor.rules).unwrap_or_else(|_| serde_json::json!({}))
1604 };
1605 crate::policy::PolicyDefinition {
1606 policy_id: descriptor.policy_id.clone(),
1607 mode: descriptor.mode.clone(),
1608 description: descriptor.description.clone(),
1609 rules,
1610 schema_version: descriptor.schema_version,
1611 }
1612 }
1613
1614 fn policy_definition_to_descriptor(
1615 definition: &crate::policy::PolicyDefinition,
1616 ) -> PolicyDescriptor {
1617 PolicyDescriptor {
1618 policy_id: definition.policy_id.clone(),
1619 mode: definition.mode.clone(),
1620 description: definition.description.clone(),
1621 rules: serde_json::to_string(&definition.rules).unwrap_or_default(),
1622 schema_version: definition.schema_version,
1623 registered_at_unix_ms: 0,
1624 }
1625 }
1626}
1627
1628#[cfg(test)]
1629mod tests {
1630 use super::*;
1631 use crate::log_store::LogStore;
1632 use crate::pb::SessionStartPayload;
1633 use crate::registry::SessionRegistry;
1634 use chrono::Utc;
1635 use prost::Message;
1636
1637 fn new_sid() -> String {
1638 uuid::Uuid::new_v4().as_hyphenated().to_string()
1639 }
1640
1641 fn make_server() -> (MacpServer, Arc<Runtime>) {
1642 let storage: Arc<dyn crate::storage::StorageBackend> =
1643 Arc::new(crate::storage::MemoryBackend);
1644 let registry = Arc::new(SessionRegistry::new());
1645 let log_store = Arc::new(LogStore::new());
1646 let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1647 let server = MacpServer::new(runtime.clone(), SecurityLayer::dev_mode());
1648 (server, runtime)
1649 }
1650
1651 fn send_req(sender: &str, env: Envelope) -> Request<SendRequest> {
1652 let mut req = Request::new(SendRequest {
1653 envelope: Some(env),
1654 });
1655 req.metadata_mut()
1656 .insert("authorization", format!("Bearer {sender}").parse().unwrap());
1657 req
1658 }
1659
1660 async fn do_send(server: &MacpServer, sender: &str, env: Envelope) -> Ack {
1661 let resp = server.send(send_req(sender, env)).await.unwrap();
1662 resp.into_inner().ack.unwrap()
1663 }
1664
1665 fn start_payload() -> Vec<u8> {
1666 SessionStartPayload {
1667 intent: "intent".into(),
1668 participants: vec!["agent://fraud".into()],
1669 mode_version: "1.0.0".into(),
1670 configuration_version: "cfg-1".into(),
1671 policy_version: String::new(),
1672 ttl_ms: 1000,
1673 context_id: String::new(),
1674 extensions: std::collections::HashMap::new(),
1675 roots: vec![],
1676 max_suspend_ms: 0,
1677 }
1678 .encode_to_vec()
1679 }
1680
1681 #[tokio::test]
1682 async fn sender_is_derived_from_authenticated_metadata() {
1683 let (server, runtime) = make_server();
1684 let sid = new_sid();
1685 let ack = do_send(
1686 &server,
1687 "agent://orchestrator",
1688 Envelope {
1689 macp_version: "1.0".into(),
1690 mode: "macp.mode.decision.v1".into(),
1691 message_type: "SessionStart".into(),
1692 message_id: "m1".into(),
1693 session_id: sid.clone(),
1694 sender: String::new(),
1695 timestamp_unix_ms: Utc::now().timestamp_millis(),
1696 payload: start_payload(),
1697 },
1698 )
1699 .await;
1700 assert!(ack.ok);
1701 let session = runtime.get_session_checked(&sid).await.unwrap();
1702 assert_eq!(session.initiator_sender, "agent://orchestrator");
1703 }
1704
1705 #[tokio::test]
1706 async fn spoofed_sender_is_rejected() {
1707 let (server, _) = make_server();
1708 let sid = new_sid();
1709 let ack = do_send(
1710 &server,
1711 "agent://orchestrator",
1712 Envelope {
1713 macp_version: "1.0".into(),
1714 mode: "macp.mode.decision.v1".into(),
1715 message_type: "SessionStart".into(),
1716 message_id: "m1".into(),
1717 session_id: sid,
1718 sender: "agent://spoof".into(),
1719 timestamp_unix_ms: Utc::now().timestamp_millis(),
1720 payload: start_payload(),
1721 },
1722 )
1723 .await;
1724 assert!(!ack.ok);
1725 assert_eq!(ack.error.as_ref().unwrap().code, "UNAUTHENTICATED");
1726 }
1727
1728 #[tokio::test]
1729 async fn get_session_requires_session_membership() {
1730 let (server, _) = make_server();
1731 let sid = new_sid();
1732 let ack = do_send(
1733 &server,
1734 "agent://orchestrator",
1735 Envelope {
1736 macp_version: "1.0".into(),
1737 mode: "macp.mode.decision.v1".into(),
1738 message_type: "SessionStart".into(),
1739 message_id: "m1".into(),
1740 session_id: sid.clone(),
1741 sender: String::new(),
1742 timestamp_unix_ms: Utc::now().timestamp_millis(),
1743 payload: start_payload(),
1744 },
1745 )
1746 .await;
1747 assert!(ack.ok);
1748
1749 let mut req = Request::new(GetSessionRequest { session_id: sid });
1750 req.metadata_mut().insert(
1751 "authorization",
1752 format!("Bearer {}", "agent://outsider").parse().unwrap(),
1753 );
1754 let err = server.get_session(req).await.unwrap_err();
1755 assert_eq!(err.code(), tonic::Code::PermissionDenied);
1756 }
1757
1758 #[tokio::test]
1759 async fn register_ext_mode_requires_authenticated_registry_permission() {
1760 let storage: Arc<dyn crate::storage::StorageBackend> =
1761 Arc::new(crate::storage::MemoryBackend);
1762 let registry = Arc::new(SessionRegistry::new());
1763 let log_store = Arc::new(LogStore::new());
1764 let runtime = Arc::new(Runtime::new(storage, registry, log_store));
1765 let security = SecurityLayer::from_env().unwrap_or_else(|_| SecurityLayer::dev_mode());
1766 let server = MacpServer::new(runtime, security);
1767
1768 let req = Request::new(RegisterExtModeRequest {
1769 mode_descriptor: Some(crate::pb::ModeDescriptor {
1770 mode: "ext.custom.v1".into(),
1771 mode_version: "1.0.0".into(),
1772 message_types: vec!["SessionStart".into(), "Commitment".into()],
1773 ..Default::default()
1774 }),
1775 });
1776 let err = server.register_ext_mode(req).await.unwrap_err();
1777 assert_eq!(err.code(), tonic::Code::Unauthenticated);
1778 }
1779
1780 fn stream_identity(sender: &str) -> AuthIdentity {
1781 AuthIdentity {
1782 sender: sender.into(),
1783 allowed_modes: None,
1784 can_start_sessions: true,
1785 max_open_sessions: None,
1786 can_manage_mode_registry: false,
1787 is_observer: false,
1788 }
1789 }
1790
1791 #[tokio::test]
1792 async fn stream_session_emits_accepted_envelopes_only() {
1793 use tokio_stream::{iter, StreamExt};
1794
1795 let (server, _) = make_server();
1796 let sid = new_sid();
1797 let requests = iter(vec![Ok(StreamSessionRequest {
1798 subscribe_session_id: String::new(),
1799 after_sequence: 0,
1800 envelope: Some(Envelope {
1801 macp_version: "1.0".into(),
1802 mode: "macp.mode.decision.v1".into(),
1803 message_type: "SessionStart".into(),
1804 message_id: "m1".into(),
1805 session_id: sid.clone(),
1806 sender: String::new(),
1807 timestamp_unix_ms: Utc::now().timestamp_millis(),
1808 payload: start_payload(),
1809 }),
1810 })]);
1811
1812 let mut stream =
1813 server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1814
1815 let response = stream.next().await.unwrap().unwrap();
1816 let envelope = match response.response.unwrap() {
1817 crate::pb::stream_session_response::Response::Envelope(e) => e,
1818 _ => panic!("expected envelope"),
1819 };
1820 assert_eq!(envelope.message_type, "SessionStart");
1821 assert_eq!(envelope.message_id, "m1");
1822 assert!(stream.next().await.is_none());
1823 }
1824
1825 #[tokio::test]
1826 async fn stream_session_rejects_mixed_session_ids() {
1827 use tokio_stream::{iter, StreamExt};
1828
1829 let (server, _) = make_server();
1830 let sid1 = new_sid();
1831 let sid2 = new_sid();
1832 let requests = iter(vec![
1833 Ok(StreamSessionRequest {
1834 subscribe_session_id: String::new(),
1835 after_sequence: 0,
1836 envelope: Some(Envelope {
1837 macp_version: "1.0".into(),
1838 mode: "macp.mode.decision.v1".into(),
1839 message_type: "SessionStart".into(),
1840 message_id: "m1".into(),
1841 session_id: sid1.clone(),
1842 sender: String::new(),
1843 timestamp_unix_ms: Utc::now().timestamp_millis(),
1844 payload: start_payload(),
1845 }),
1846 }),
1847 Ok(StreamSessionRequest {
1848 subscribe_session_id: String::new(),
1849 after_sequence: 0,
1850 envelope: Some(Envelope {
1851 macp_version: "1.0".into(),
1852 mode: "macp.mode.decision.v1".into(),
1853 message_type: "SessionStart".into(),
1854 message_id: "m2".into(),
1855 session_id: sid2,
1856 sender: String::new(),
1857 timestamp_unix_ms: Utc::now().timestamp_millis(),
1858 payload: start_payload(),
1859 }),
1860 }),
1861 ]);
1862
1863 let mut stream =
1864 server.build_stream_session_stream(stream_identity("agent://orchestrator"), requests);
1865
1866 let first = stream.next().await.unwrap().unwrap();
1867 let first_env = match first.response.unwrap() {
1868 crate::pb::stream_session_response::Response::Envelope(e) => e,
1869 _ => panic!("expected envelope"),
1870 };
1871 assert_eq!(first_env.session_id, sid1);
1872 let err = stream.next().await.unwrap().unwrap_err();
1873 assert_eq!(err.code(), tonic::Code::InvalidArgument);
1874 }
1875
1876 #[tokio::test]
1877 async fn list_modes_returns_standard_modes() {
1878 let (server, _) = make_server();
1879 let resp = server
1880 .list_modes(Request::new(ListModesRequest {}))
1881 .await
1882 .unwrap();
1883 let names: Vec<String> = resp
1884 .into_inner()
1885 .modes
1886 .iter()
1887 .map(|m| m.mode.clone())
1888 .collect();
1889 assert_eq!(names.len(), 5);
1890 assert!(names.contains(&"macp.mode.decision.v1".to_string()));
1891 assert!(names.contains(&"macp.mode.proposal.v1".to_string()));
1892 assert!(names.contains(&"macp.mode.task.v1".to_string()));
1893 assert!(names.contains(&"macp.mode.handoff.v1".to_string()));
1894 assert!(names.contains(&"macp.mode.quorum.v1".to_string()));
1895 assert!(!names.contains(&"ext.multi_round.v1".to_string()));
1897 }
1898
1899 #[tokio::test]
1900 async fn list_ext_modes_returns_extensions() {
1901 let (server, _) = make_server();
1902 let resp = server
1903 .list_ext_modes(Request::new(ListExtModesRequest {}))
1904 .await
1905 .unwrap();
1906 let names: Vec<String> = resp
1907 .into_inner()
1908 .modes
1909 .iter()
1910 .map(|m| m.mode.clone())
1911 .collect();
1912 assert_eq!(names.len(), 1);
1913 assert!(names.contains(&"ext.multi_round.v1".to_string()));
1914 }
1915
1916 #[tokio::test]
1917 async fn get_manifest_includes_all_modes() {
1918 let (server, _) = make_server();
1919 let resp = server
1920 .get_manifest(Request::new(crate::pb::GetManifestRequest {
1921 agent_id: String::new(),
1922 }))
1923 .await
1924 .unwrap();
1925 let manifest = resp.into_inner().manifest.unwrap();
1926 assert_eq!(manifest.supported_modes.len(), 6);
1927 assert!(manifest
1928 .supported_modes
1929 .contains(&"ext.multi_round.v1".to_string()));
1930 }
1931
1932 #[tokio::test]
1933 async fn get_session_returns_metadata() {
1934 let (server, _) = make_server();
1935 let sid = new_sid();
1936 let ack = do_send(
1937 &server,
1938 "agent://orchestrator",
1939 Envelope {
1940 macp_version: "1.0".into(),
1941 mode: "macp.mode.decision.v1".into(),
1942 message_type: "SessionStart".into(),
1943 message_id: "m1".into(),
1944 session_id: sid.clone(),
1945 sender: String::new(),
1946 timestamp_unix_ms: Utc::now().timestamp_millis(),
1947 payload: start_payload(),
1948 },
1949 )
1950 .await;
1951 assert!(ack.ok);
1952
1953 let mut req = Request::new(GetSessionRequest {
1954 session_id: sid.clone(),
1955 });
1956 req.metadata_mut().insert(
1957 "authorization",
1958 format!("Bearer {}", "agent://orchestrator")
1959 .parse()
1960 .unwrap(),
1961 );
1962 let resp = server.get_session(req).await.unwrap();
1963 let meta = resp.into_inner().metadata.unwrap();
1964 assert_eq!(meta.session_id, sid);
1965 assert_eq!(meta.mode, "macp.mode.decision.v1");
1966 assert_eq!(meta.mode_version, "1.0.0");
1967 assert_eq!(meta.configuration_version, "cfg-1");
1968 }
1969
1970 #[tokio::test]
1971 async fn cancel_session_transitions_to_cancelled() {
1972 let (server, _) = make_server();
1973 let sid = new_sid();
1974 let ack = do_send(
1975 &server,
1976 "agent://orchestrator",
1977 Envelope {
1978 macp_version: "1.0".into(),
1979 mode: "macp.mode.decision.v1".into(),
1980 message_type: "SessionStart".into(),
1981 message_id: "m1".into(),
1982 session_id: sid.clone(),
1983 sender: String::new(),
1984 timestamp_unix_ms: Utc::now().timestamp_millis(),
1985 payload: start_payload(),
1986 },
1987 )
1988 .await;
1989 assert!(ack.ok);
1990
1991 let mut req = Request::new(CancelSessionRequest {
1992 session_id: sid,
1993 reason: "no longer needed".into(),
1994 });
1995 req.metadata_mut().insert(
1996 "authorization",
1997 format!("Bearer {}", "agent://orchestrator")
1998 .parse()
1999 .unwrap(),
2000 );
2001 let resp = server.cancel_session(req).await.unwrap();
2002 let ack = resp.into_inner().ack.unwrap();
2003 assert!(ack.ok);
2004 assert_eq!(ack.session_state, PbSessionState::Cancelled as i32);
2006 }
2007
2008 #[tokio::test]
2009 async fn participant_cannot_cancel_session() {
2010 let (server, _) = make_server();
2011 let sid = new_sid();
2012 let ack = do_send(
2013 &server,
2014 "agent://orchestrator",
2015 Envelope {
2016 macp_version: "1.0".into(),
2017 mode: "macp.mode.decision.v1".into(),
2018 message_type: "SessionStart".into(),
2019 message_id: "m1".into(),
2020 session_id: sid.clone(),
2021 sender: String::new(),
2022 timestamp_unix_ms: Utc::now().timestamp_millis(),
2023 payload: start_payload(),
2024 },
2025 )
2026 .await;
2027 assert!(ack.ok);
2028
2029 let mut req = Request::new(CancelSessionRequest {
2030 session_id: sid,
2031 reason: "I want to cancel".into(),
2032 });
2033 req.metadata_mut().insert(
2034 "authorization",
2035 format!("Bearer {}", "agent://fraud").parse().unwrap(),
2036 );
2037 let err = server.cancel_session(req).await.unwrap_err();
2038 assert_eq!(err.code(), tonic::Code::PermissionDenied);
2039 }
2040
2041 #[tokio::test]
2042 async fn cancel_session_unknown_session_returns_error() {
2043 let (server, _) = make_server();
2044 let mut req = Request::new(CancelSessionRequest {
2045 session_id: "nonexistent".into(),
2046 reason: "test".into(),
2047 });
2048 req.metadata_mut().insert(
2049 "authorization",
2050 format!("Bearer {}", "agent://orchestrator")
2051 .parse()
2052 .unwrap(),
2053 );
2054 let err = server.cancel_session(req).await.unwrap_err();
2055 assert_eq!(err.code(), tonic::Code::NotFound);
2056 }
2057
2058 #[tokio::test]
2059 async fn ambient_signal_accepted() {
2060 let (server, _) = make_server();
2061 let ack = do_send(
2062 &server,
2063 "agent://orchestrator",
2064 Envelope {
2065 macp_version: "1.0".into(),
2066 mode: String::new(),
2067 message_type: "Signal".into(),
2068 message_id: "sig-1".into(),
2069 session_id: String::new(),
2070 sender: String::new(),
2071 timestamp_unix_ms: Utc::now().timestamp_millis(),
2072 payload: vec![],
2073 },
2074 )
2075 .await;
2076 assert!(ack.ok);
2077 }
2078
2079 #[tokio::test]
2080 async fn signal_with_session_id_rejected() {
2081 let (server, _) = make_server();
2082 let ack = do_send(
2083 &server,
2084 "agent://orchestrator",
2085 Envelope {
2086 macp_version: "1.0".into(),
2087 mode: String::new(),
2088 message_type: "Signal".into(),
2089 message_id: "sig-2".into(),
2090 session_id: "some-session".into(),
2091 sender: String::new(),
2092 timestamp_unix_ms: Utc::now().timestamp_millis(),
2093 payload: vec![],
2094 },
2095 )
2096 .await;
2097 assert!(!ack.ok);
2098 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2099 }
2100
2101 #[tokio::test]
2102 async fn signal_with_mode_rejected() {
2103 let (server, _) = make_server();
2104 let ack = do_send(
2105 &server,
2106 "agent://orchestrator",
2107 Envelope {
2108 macp_version: "1.0".into(),
2109 mode: "macp.mode.decision.v1".into(),
2110 message_type: "Signal".into(),
2111 message_id: "sig-3".into(),
2112 session_id: String::new(),
2113 sender: String::new(),
2114 timestamp_unix_ms: Utc::now().timestamp_millis(),
2115 payload: vec![],
2116 },
2117 )
2118 .await;
2119 assert!(!ack.ok);
2120 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2121 }
2122
2123 #[tokio::test]
2124 async fn ambient_progress_accepted() {
2125 let (server, _) = make_server();
2126 let ack = do_send(
2127 &server,
2128 "agent://orchestrator",
2129 Envelope {
2130 macp_version: "1.0".into(),
2131 mode: String::new(),
2132 message_type: "Progress".into(),
2133 message_id: "prog-1".into(),
2134 session_id: String::new(),
2135 sender: String::new(),
2136 timestamp_unix_ms: Utc::now().timestamp_millis(),
2137 payload: vec![],
2138 },
2139 )
2140 .await;
2141 assert!(ack.ok);
2142 }
2143
2144 #[tokio::test]
2145 async fn ambient_progress_with_mode_rejected() {
2146 let (server, _) = make_server();
2147 let ack = do_send(
2148 &server,
2149 "agent://orchestrator",
2150 Envelope {
2151 macp_version: "1.0".into(),
2152 mode: "macp.mode.decision.v1".into(),
2153 message_type: "Progress".into(),
2154 message_id: "prog-2".into(),
2155 session_id: String::new(),
2156 sender: String::new(),
2157 timestamp_unix_ms: Utc::now().timestamp_millis(),
2158 payload: vec![],
2159 },
2160 )
2161 .await;
2162 assert!(!ack.ok);
2163 assert_eq!(ack.error.as_ref().unwrap().code, "INVALID_ENVELOPE");
2164 }
2165
2166 #[tokio::test]
2167 async fn manifest_advertises_stream_enabled() {
2168 let (server, _) = make_server();
2169 let resp = server
2170 .initialize(Request::new(InitializeRequest {
2171 supported_protocol_versions: vec!["1.0".into()],
2172 client_info: None,
2173 capabilities: None,
2174 }))
2175 .await
2176 .unwrap();
2177 let caps = resp.into_inner().capabilities.unwrap();
2178 assert!(caps.sessions.unwrap().stream);
2179 }
2180
2181 #[tokio::test]
2182 async fn initialize_empty_versions_rejected() {
2183 let (server, _) = make_server();
2184 let err = server
2185 .initialize(Request::new(InitializeRequest {
2186 supported_protocol_versions: vec![],
2187 client_info: None,
2188 capabilities: None,
2189 }))
2190 .await
2191 .unwrap_err();
2192 assert_eq!(err.code(), tonic::Code::InvalidArgument);
2193 }
2194
2195 #[tokio::test]
2196 async fn initialize_unsupported_version_rejected() {
2197 let (server, _) = make_server();
2198 let err = server
2199 .initialize(Request::new(InitializeRequest {
2200 supported_protocol_versions: vec!["2.0".into()],
2201 client_info: None,
2202 capabilities: None,
2203 }))
2204 .await
2205 .unwrap_err();
2206 assert_eq!(err.code(), tonic::Code::FailedPrecondition);
2207 }
2208
2209 fn observer_identity(sender: &str) -> AuthIdentity {
2212 AuthIdentity {
2213 sender: sender.into(),
2214 allowed_modes: None,
2215 can_start_sessions: false,
2216 max_open_sessions: None,
2217 can_manage_mode_registry: false,
2218 is_observer: true,
2219 }
2220 }
2221
2222 fn subscribe_frame(session_id: &str, after: u64) -> StreamSessionRequest {
2223 StreamSessionRequest {
2224 subscribe_session_id: session_id.into(),
2225 after_sequence: after,
2226 envelope: None,
2227 }
2228 }
2229
2230 fn start_multi_participant(participants: Vec<String>) -> Vec<u8> {
2231 SessionStartPayload {
2232 intent: "intent".into(),
2233 participants,
2234 mode_version: "1.0.0".into(),
2235 configuration_version: "cfg-1".into(),
2236 policy_version: String::new(),
2237 ttl_ms: 60_000,
2238 context_id: String::new(),
2239 extensions: std::collections::HashMap::new(),
2240 roots: vec![],
2241 max_suspend_ms: 0,
2242 }
2243 .encode_to_vec()
2244 }
2245
2246 async fn start_session(
2247 server: &MacpServer,
2248 initiator: &str,
2249 sid: &str,
2250 participants: Vec<String>,
2251 ) {
2252 let ack = do_send(
2253 server,
2254 initiator,
2255 Envelope {
2256 macp_version: "1.0".into(),
2257 mode: "macp.mode.decision.v1".into(),
2258 message_type: "SessionStart".into(),
2259 message_id: "start".into(),
2260 session_id: sid.into(),
2261 sender: String::new(),
2262 timestamp_unix_ms: Utc::now().timestamp_millis(),
2263 payload: start_multi_participant(participants),
2264 },
2265 )
2266 .await;
2267 assert!(ack.ok, "SessionStart failed: {:?}", ack.error);
2268 }
2269
2270 async fn send_proposal(
2271 server: &MacpServer,
2272 sender: &str,
2273 sid: &str,
2274 message_id: &str,
2275 proposal_id: &str,
2276 ) {
2277 let payload = crate::decision_pb::ProposalPayload {
2278 proposal_id: proposal_id.into(),
2279 option: "opt".into(),
2280 rationale: "r".into(),
2281 supporting_data: vec![],
2282 }
2283 .encode_to_vec();
2284 let ack = do_send(
2285 server,
2286 sender,
2287 Envelope {
2288 macp_version: "1.0".into(),
2289 mode: "macp.mode.decision.v1".into(),
2290 message_type: "Proposal".into(),
2291 message_id: message_id.into(),
2292 session_id: sid.into(),
2293 sender: String::new(),
2294 timestamp_unix_ms: Utc::now().timestamp_millis(),
2295 payload,
2296 },
2297 )
2298 .await;
2299 assert!(ack.ok, "Proposal failed: {:?}", ack.error);
2300 }
2301
2302 #[tokio::test]
2303 async fn subscribe_replays_session_history_from_zero() {
2304 let (server, _) = make_server();
2305 let sid = new_sid();
2306 let initiator = "agent://orchestrator";
2307 let peer = "agent://fraud";
2308 start_session(
2309 &server,
2310 initiator,
2311 &sid,
2312 vec![initiator.into(), peer.into()],
2313 )
2314 .await;
2315 send_proposal(&server, peer, &sid, "m2", "p1").await;
2316
2317 let mut bound = None;
2318 let mut events = None;
2319 let replay = server
2320 .process_stream_request(
2321 &stream_identity(peer),
2322 subscribe_frame(&sid, 0),
2323 &mut bound,
2324 &mut events,
2325 )
2326 .await
2327 .unwrap();
2328
2329 assert_eq!(replay.len(), 2);
2330 assert_eq!(replay[0].message_type, "SessionStart");
2331 assert_eq!(replay[0].message_id, "start");
2332 assert_eq!(replay[1].message_type, "Proposal");
2333 assert_eq!(replay[1].message_id, "m2");
2334 assert_eq!(bound.as_deref(), Some(sid.as_str()));
2335 assert!(events.is_some());
2336 }
2337
2338 #[tokio::test]
2339 async fn subscribe_after_sequence_filters_history() {
2340 let (server, _) = make_server();
2341 let sid = new_sid();
2342 let initiator = "agent://orchestrator";
2343 let peer = "agent://fraud";
2344 start_session(
2345 &server,
2346 initiator,
2347 &sid,
2348 vec![initiator.into(), peer.into()],
2349 )
2350 .await;
2351 send_proposal(&server, peer, &sid, "m2", "p1").await;
2352 send_proposal(&server, peer, &sid, "m3", "p2").await;
2353
2354 let mut bound = None;
2355 let mut events = None;
2356 let replay = server
2357 .process_stream_request(
2358 &stream_identity(peer),
2359 subscribe_frame(&sid, 2),
2360 &mut bound,
2361 &mut events,
2362 )
2363 .await
2364 .unwrap();
2365
2366 assert_eq!(replay.len(), 1);
2367 assert_eq!(replay[0].message_id, "m3");
2368 }
2369
2370 #[tokio::test]
2371 async fn subscribe_unknown_session_returns_not_found() {
2372 let (server, _) = make_server();
2373 let mut bound = None;
2374 let mut events = None;
2375 let status = server
2376 .process_stream_request(
2377 &stream_identity("agent://orchestrator"),
2378 subscribe_frame("missing-session", 0),
2379 &mut bound,
2380 &mut events,
2381 )
2382 .await
2383 .unwrap_err();
2384 assert_eq!(status.code(), tonic::Code::NotFound);
2385 assert!(bound.is_none());
2386 assert!(events.is_none());
2387 }
2388
2389 #[tokio::test]
2390 async fn subscribe_non_participant_is_forbidden() {
2391 let (server, _) = make_server();
2392 let sid = new_sid();
2393 start_session(
2394 &server,
2395 "agent://orchestrator",
2396 &sid,
2397 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2398 )
2399 .await;
2400
2401 let mut bound = None;
2402 let mut events = None;
2403 let status = server
2404 .process_stream_request(
2405 &stream_identity("agent://outsider"),
2406 subscribe_frame(&sid, 0),
2407 &mut bound,
2408 &mut events,
2409 )
2410 .await
2411 .unwrap_err();
2412 assert_eq!(status.code(), tonic::Code::PermissionDenied);
2413 }
2414
2415 #[tokio::test]
2416 async fn subscribe_observer_identity_allowed() {
2417 let (server, _) = make_server();
2418 let sid = new_sid();
2419 start_session(
2420 &server,
2421 "agent://orchestrator",
2422 &sid,
2423 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2424 )
2425 .await;
2426
2427 let mut bound = None;
2428 let mut events = None;
2429 let replay = server
2430 .process_stream_request(
2431 &observer_identity("agent://auditor"),
2432 subscribe_frame(&sid, 0),
2433 &mut bound,
2434 &mut events,
2435 )
2436 .await
2437 .unwrap();
2438 assert_eq!(replay.len(), 1);
2439 assert_eq!(replay[0].message_type, "SessionStart");
2440 }
2441
2442 #[tokio::test]
2443 async fn subscribe_initiator_allowed_even_when_not_listed() {
2444 let (server, _) = make_server();
2447 let sid = new_sid();
2448 start_session(
2449 &server,
2450 "agent://orchestrator",
2451 &sid,
2452 vec!["agent://fraud".into()],
2453 )
2454 .await;
2455
2456 let mut bound = None;
2457 let mut events = None;
2458 let replay = server
2459 .process_stream_request(
2460 &stream_identity("agent://orchestrator"),
2461 subscribe_frame(&sid, 0),
2462 &mut bound,
2463 &mut events,
2464 )
2465 .await
2466 .unwrap();
2467 assert_eq!(replay.len(), 1);
2468 }
2469
2470 #[tokio::test]
2471 async fn stream_request_with_envelope_and_subscribe_is_rejected() {
2472 let (server, _) = make_server();
2473 let sid = new_sid();
2474 let req = StreamSessionRequest {
2475 subscribe_session_id: sid.clone(),
2476 after_sequence: 0,
2477 envelope: Some(Envelope {
2478 macp_version: "1.0".into(),
2479 mode: "macp.mode.decision.v1".into(),
2480 message_type: "SessionStart".into(),
2481 message_id: "m1".into(),
2482 session_id: sid,
2483 sender: String::new(),
2484 timestamp_unix_ms: Utc::now().timestamp_millis(),
2485 payload: start_payload(),
2486 }),
2487 };
2488
2489 let mut bound = None;
2490 let mut events = None;
2491 let status = server
2492 .process_stream_request(
2493 &stream_identity("agent://orchestrator"),
2494 req,
2495 &mut bound,
2496 &mut events,
2497 )
2498 .await
2499 .unwrap_err();
2500 assert_eq!(status.code(), tonic::Code::InvalidArgument);
2501 }
2502
2503 #[tokio::test]
2504 async fn subscribe_to_different_session_on_bound_stream_is_rejected() {
2505 let (server, _) = make_server();
2506 let sid1 = new_sid();
2507 let sid2 = new_sid();
2508 start_session(
2509 &server,
2510 "agent://orchestrator",
2511 &sid1,
2512 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2513 )
2514 .await;
2515 start_session(
2516 &server,
2517 "agent://orchestrator",
2518 &sid2,
2519 vec!["agent://orchestrator".into(), "agent://fraud".into()],
2520 )
2521 .await;
2522
2523 let identity = stream_identity("agent://fraud");
2525 let mut bound = None;
2526 let mut events = None;
2527 server
2528 .process_stream_request(
2529 &identity,
2530 subscribe_frame(&sid1, 0),
2531 &mut bound,
2532 &mut events,
2533 )
2534 .await
2535 .unwrap();
2536 assert_eq!(bound.as_deref(), Some(sid1.as_str()));
2537
2538 let status = server
2540 .process_stream_request(
2541 &identity,
2542 subscribe_frame(&sid2, 0),
2543 &mut bound,
2544 &mut events,
2545 )
2546 .await
2547 .unwrap_err();
2548 assert_eq!(status.code(), tonic::Code::InvalidArgument);
2549 }
2550
2551 struct DenySenderEngine {
2555 denied: String,
2556 }
2557
2558 #[async_trait::async_trait]
2559 impl crate::policy_engine::PolicyEngine for DenySenderEngine {
2560 async fn evaluate_session_start(
2561 &self,
2562 identity: &crate::security::AuthIdentity,
2563 _mode: &str,
2564 _env: &Envelope,
2565 ) -> macp_core::policy::PolicyDecision {
2566 if identity.sender == self.denied {
2567 macp_core::policy::PolicyDecision::Deny {
2568 reasons: vec!["sender embargoed".into()],
2569 }
2570 } else {
2571 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2572 }
2573 }
2574
2575 async fn evaluate_message(
2576 &self,
2577 identity: &crate::security::AuthIdentity,
2578 _session: &macp_core::session::Session,
2579 _env: &Envelope,
2580 ) -> macp_core::policy::PolicyDecision {
2581 if identity.sender == self.denied {
2582 macp_core::policy::PolicyDecision::Deny {
2583 reasons: vec!["sender embargoed".into()],
2584 }
2585 } else {
2586 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2587 }
2588 }
2589
2590 async fn evaluate_session_access(
2591 &self,
2592 identity: &crate::security::AuthIdentity,
2593 _session: &macp_core::session::Session,
2594 ) -> macp_core::policy::PolicyDecision {
2595 if identity.sender == self.denied {
2596 macp_core::policy::PolicyDecision::Deny {
2597 reasons: vec!["sender embargoed".into()],
2598 }
2599 } else {
2600 macp_core::policy::PolicyDecision::Allow { reasons: vec![] }
2601 }
2602 }
2603 }
2604
2605 #[tokio::test]
2606 async fn policy_engine_gates_all_three_ingress_points() {
2607 let (server, _runtime) = make_server();
2608 let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2609 denied: "agent://embargoed".into(),
2610 }));
2611
2612 let sid = new_sid();
2613 let start_payload = SessionStartPayload {
2614 intent: "e3".into(),
2615 participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2616 mode_version: "1.0.0".into(),
2617 configuration_version: "cfg-1".into(),
2618 policy_version: String::new(),
2619 ttl_ms: 60_000,
2620 context_id: String::new(),
2621 extensions: Default::default(),
2622 roots: vec![],
2623 max_suspend_ms: 0,
2624 }
2625 .encode_to_vec();
2626 let start_env = |sender: &str, sid: &str| Envelope {
2627 macp_version: "1.0".into(),
2628 mode: "macp.mode.decision.v1".into(),
2629 message_type: "SessionStart".into(),
2630 message_id: new_sid(),
2631 session_id: sid.into(),
2632 sender: sender.into(),
2633 timestamp_unix_ms: Utc::now().timestamp_millis(),
2634 payload: start_payload.clone(),
2635 };
2636
2637 let ack = server
2639 .send(send_req(
2640 "agent://embargoed",
2641 start_env("agent://embargoed", &sid),
2642 ))
2643 .await
2644 .unwrap()
2645 .into_inner()
2646 .ack
2647 .unwrap();
2648 assert!(!ack.ok);
2649 assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2650
2651 let ack = server
2653 .send(send_req("agent://ok", start_env("agent://ok", &sid)))
2654 .await
2655 .unwrap()
2656 .into_inner()
2657 .ack
2658 .unwrap();
2659 assert!(ack.ok, "allowed sender must start: {:?}", ack.error);
2660
2661 let proposal = crate::decision_pb::ProposalPayload {
2663 proposal_id: "p1".into(),
2664 option: "x".into(),
2665 rationale: "r".into(),
2666 supporting_data: vec![],
2667 }
2668 .encode_to_vec();
2669 let msg_env = Envelope {
2670 macp_version: "1.0".into(),
2671 mode: "macp.mode.decision.v1".into(),
2672 message_type: "Proposal".into(),
2673 message_id: new_sid(),
2674 session_id: sid.clone(),
2675 sender: "agent://embargoed".into(),
2676 timestamp_unix_ms: Utc::now().timestamp_millis(),
2677 payload: proposal,
2678 };
2679 let ack = server
2680 .send(send_req("agent://embargoed", msg_env))
2681 .await
2682 .unwrap()
2683 .into_inner()
2684 .ack
2685 .unwrap();
2686 assert!(!ack.ok);
2687 assert_eq!(ack.error.unwrap().code, "POLICY_DENIED");
2688
2689 let mut req = Request::new(crate::pb::GetSessionRequest {
2691 session_id: sid.clone(),
2692 });
2693 req.metadata_mut()
2694 .insert("authorization", "Bearer agent://embargoed".parse().unwrap());
2695 let err = server
2696 .get_session(req)
2697 .await
2698 .expect_err("embargoed read must be denied");
2699 assert_eq!(err.code(), tonic::Code::PermissionDenied);
2700 }
2701
2702 #[tokio::test]
2706 async fn policy_engine_gates_stream_path() {
2707 let (server, runtime) = make_server();
2708 let server = server.with_policy_engine(Arc::new(DenySenderEngine {
2709 denied: "agent://embargoed".into(),
2710 }));
2711
2712 let sid = new_sid();
2716 let payload = SessionStartPayload {
2717 intent: "e3-stream".into(),
2718 participants: vec!["agent://ok".into(), "agent://embargoed".into()],
2719 mode_version: "1.0.0".into(),
2720 configuration_version: "cfg-1".into(),
2721 policy_version: String::new(),
2722 ttl_ms: 60_000,
2723 context_id: String::new(),
2724 extensions: Default::default(),
2725 roots: vec![],
2726 max_suspend_ms: 0,
2727 }
2728 .encode_to_vec();
2729 runtime
2730 .process(
2731 &Envelope {
2732 macp_version: "1.0".into(),
2733 mode: "macp.mode.decision.v1".into(),
2734 message_type: "SessionStart".into(),
2735 message_id: new_sid(),
2736 session_id: sid.clone(),
2737 sender: "agent://ok".into(),
2738 timestamp_unix_ms: Utc::now().timestamp_millis(),
2739 payload,
2740 },
2741 None,
2742 )
2743 .await
2744 .unwrap();
2745
2746 let embargoed = crate::security::AuthIdentity {
2747 sender: "agent://embargoed".into(),
2748 allowed_modes: None,
2749 can_start_sessions: true,
2750 max_open_sessions: None,
2751 can_manage_mode_registry: false,
2752 is_observer: false,
2753 };
2754 let mut bound = None;
2755 let mut events = None;
2756
2757 let proposal = crate::decision_pb::ProposalPayload {
2759 proposal_id: "p1".into(),
2760 option: "x".into(),
2761 rationale: "r".into(),
2762 supporting_data: vec![],
2763 }
2764 .encode_to_vec();
2765 let req = StreamSessionRequest {
2766 envelope: Some(Envelope {
2767 macp_version: "1.0".into(),
2768 mode: "macp.mode.decision.v1".into(),
2769 message_type: "Proposal".into(),
2770 message_id: new_sid(),
2771 session_id: sid.clone(),
2772 sender: "agent://embargoed".into(),
2773 timestamp_unix_ms: Utc::now().timestamp_millis(),
2774 payload: proposal,
2775 }),
2776 subscribe_session_id: String::new(),
2777 after_sequence: 0,
2778 };
2779 let err = server
2780 .process_stream_request(&embargoed, req, &mut bound, &mut events)
2781 .await
2782 .expect_err("stream envelope from embargoed sender must be denied");
2783 assert_eq!(err.code(), tonic::Code::FailedPrecondition, "{err:?}");
2786 assert!(err.message().contains("PolicyDenied"), "{err:?}");
2787
2788 let req = StreamSessionRequest {
2791 envelope: None,
2792 subscribe_session_id: sid.clone(),
2793 after_sequence: 0,
2794 };
2795 let err = server
2796 .process_stream_request(&embargoed, req, &mut bound, &mut events)
2797 .await
2798 .expect_err("stream subscribe from embargoed sender must be denied");
2799 assert_eq!(err.code(), tonic::Code::PermissionDenied, "{err:?}");
2800 }
2801}