1use macp_core::error::MacpError;
2use std::collections::{HashMap, HashSet, VecDeque};
3use std::fs;
4use std::path::PathBuf;
5use std::sync::Arc;
6use std::time::{Duration, Instant};
7use tokio::sync::Mutex;
8use tonic::metadata::MetadataMap;
9
10#[derive(Clone, Debug)]
11pub struct AuthIdentity {
12 pub sender: String,
13 pub allowed_modes: Option<HashSet<String>>,
14 pub can_start_sessions: bool,
15 pub max_open_sessions: Option<usize>,
16 pub can_manage_mode_registry: bool,
17 pub is_observer: bool,
18}
19
20#[derive(Clone, Debug, serde::Deserialize)]
21struct RawIdentity {
22 token: String,
23 sender: String,
24 #[serde(default)]
25 allowed_modes: Vec<String>,
26 #[serde(default = "default_true")]
27 can_start_sessions: bool,
28 max_open_sessions: Option<usize>,
29 #[serde(default)]
30 can_manage_mode_registry: bool,
31 #[serde(default)]
32 is_observer: bool,
33}
34
35#[derive(Clone, Debug, serde::Deserialize)]
36#[serde(untagged)]
37enum RawConfig {
38 List(Vec<RawIdentity>),
39 Wrapped { tokens: Vec<RawIdentity> },
40}
41
42fn default_true() -> bool {
43 true
44}
45
46#[derive(Clone, Debug)]
47pub struct RateLimitConfig {
48 pub limit: usize,
49 pub window: Duration,
50}
51
52#[derive(Default)]
53struct RateBucket {
54 start_events: Mutex<HashMap<String, VecDeque<Instant>>>,
55 message_events: Mutex<HashMap<String, VecDeque<Instant>>>,
56 start_sweep_counter: std::sync::atomic::AtomicU64,
61 message_sweep_counter: std::sync::atomic::AtomicU64,
62}
63
64#[derive(Clone)]
65pub struct SecurityLayer {
66 identities: Arc<HashMap<String, AuthIdentity>>,
67 rate_bucket: Arc<RateBucket>,
68 auth_chain: Option<Arc<crate::auth::AuthResolverChain>>,
69 pub max_payload_bytes: usize,
70 session_start_rate: RateLimitConfig,
71 message_rate: RateLimitConfig,
72}
73
74impl SecurityLayer {
75 pub fn dev_mode() -> Self {
79 Self {
80 identities: Arc::new(HashMap::new()),
81 rate_bucket: Arc::new(RateBucket::default()),
82 auth_chain: None,
83 max_payload_bytes: 1_048_576,
84 session_start_rate: RateLimitConfig {
85 limit: usize::MAX,
86 window: Duration::from_secs(60),
87 },
88 message_rate: RateLimitConfig {
89 limit: usize::MAX,
90 window: Duration::from_secs(60),
91 },
92 }
93 }
94
95 fn dev_authenticate(&self, metadata: &MetadataMap) -> Result<AuthIdentity, MacpError> {
105 if let Some(token) = Self::bearer_token(metadata) {
106 return Ok(AuthIdentity {
107 sender: token,
108 allowed_modes: None,
109 can_start_sessions: true,
110 max_open_sessions: None,
111 can_manage_mode_registry: true,
112 is_observer: false,
113 });
114 }
115 Err(MacpError::Unauthenticated)
116 }
117
118 pub fn from_env() -> Result<Self, Box<dyn std::error::Error>> {
119 let max_payload_bytes = std::env::var("MACP_MAX_PAYLOAD_BYTES")
120 .ok()
121 .and_then(|v| v.parse::<usize>().ok())
122 .unwrap_or(1_048_576);
123
124 let session_start_rate = RateLimitConfig {
125 limit: std::env::var("MACP_SESSION_START_LIMIT_PER_MINUTE")
126 .ok()
127 .and_then(|v| v.parse::<usize>().ok())
128 .unwrap_or(60),
129 window: Duration::from_secs(60),
130 };
131 let message_rate = RateLimitConfig {
132 limit: std::env::var("MACP_MESSAGE_LIMIT_PER_MINUTE")
133 .ok()
134 .and_then(|v| v.parse::<usize>().ok())
135 .unwrap_or(600),
136 window: Duration::from_secs(60),
137 };
138
139 let raw = if let Ok(json) = std::env::var("MACP_AUTH_TOKENS_JSON") {
140 Some(json)
141 } else if let Ok(path) = std::env::var("MACP_AUTH_TOKENS_FILE") {
142 Some(fs::read_to_string(PathBuf::from(path))?)
143 } else {
144 None
145 };
146
147 let identities = raw
148 .as_ref()
149 .map(|json| Self::parse_identities(json))
150 .transpose()?
151 .unwrap_or_default();
152
153 let mut resolvers: Vec<Box<dyn crate::auth::AuthResolver>> = Vec::new();
155
156 if let Ok(issuer) = std::env::var("MACP_AUTH_ISSUER") {
158 let audience =
159 std::env::var("MACP_AUTH_AUDIENCE").unwrap_or_else(|_| "macp-runtime".into());
160 let cache_ttl = std::env::var("MACP_AUTH_JWKS_TTL_SECS")
161 .ok()
162 .and_then(|v| v.parse().ok())
163 .unwrap_or(300u64);
164 let algorithms = std::env::var("MACP_AUTH_JWT_ALGS")
169 .ok()
170 .map(|raw| {
171 raw.split(',')
172 .filter_map(|a| match a.trim().to_uppercase().as_str() {
173 "RS256" => Some(jsonwebtoken::Algorithm::RS256),
174 "ES256" => Some(jsonwebtoken::Algorithm::ES256),
175 "HS256" => Some(jsonwebtoken::Algorithm::HS256),
176 other => {
177 tracing::warn!(alg = other, "ignoring unsupported JWT algorithm");
178 None
179 }
180 })
181 .collect::<Vec<_>>()
182 })
183 .filter(|algs| !algs.is_empty())
184 .unwrap_or_else(|| {
185 vec![
186 jsonwebtoken::Algorithm::RS256,
187 jsonwebtoken::Algorithm::ES256,
188 ]
189 });
190 let config = crate::auth::resolvers::jwt_bearer::JwtConfig {
191 issuer,
192 audience,
193 algorithms,
194 };
195 if let Ok(jwks_json) = std::env::var("MACP_AUTH_JWKS_JSON") {
196 match crate::auth::resolvers::JwtBearerResolver::from_inline_json(
197 config, &jwks_json,
198 ) {
199 Ok(resolver) => resolvers.push(Box::new(resolver)),
200 Err(e) => {
201 tracing::error!("failed to create JWT resolver from inline JWKS: {e}")
202 }
203 }
204 } else if let Ok(jwks_url) = std::env::var("MACP_AUTH_JWKS_URL") {
205 resolvers.push(Box::new(
206 crate::auth::resolvers::JwtBearerResolver::from_url(
207 config, jwks_url, cache_ttl,
208 ),
209 ));
210 }
211 }
212
213 if !identities.is_empty() {
215 resolvers.push(Box::new(crate::auth::resolvers::StaticBearerResolver::new(
216 identities.clone(),
217 )));
218 }
219
220 let auth_chain = if resolvers.is_empty() {
221 None
222 } else {
223 Some(Arc::new(crate::auth::AuthResolverChain::new(resolvers)))
224 };
225
226 Ok(Self {
227 identities: Arc::new(identities),
228 rate_bucket: Arc::new(RateBucket::default()),
229 auth_chain,
230 max_payload_bytes,
231 session_start_rate,
232 message_rate,
233 })
234 }
235
236 fn parse_identities(
237 json: &str,
238 ) -> Result<HashMap<String, AuthIdentity>, Box<dyn std::error::Error>> {
239 let parsed: RawConfig = serde_json::from_str(json)?;
240 let items = match parsed {
241 RawConfig::List(items) => items,
242 RawConfig::Wrapped { tokens } => tokens,
243 };
244 let mut identities = HashMap::new();
245 for item in items {
246 identities.insert(
247 item.token,
248 AuthIdentity {
249 sender: item.sender,
250 allowed_modes: if item.allowed_modes.is_empty() {
251 None
252 } else {
253 Some(item.allowed_modes.into_iter().collect())
254 },
255 can_start_sessions: item.can_start_sessions,
256 max_open_sessions: item.max_open_sessions,
257 can_manage_mode_registry: item.can_manage_mode_registry,
258 is_observer: item.is_observer,
259 },
260 );
261 }
262 Ok(identities)
263 }
264
265 fn bearer_token(metadata: &MetadataMap) -> Option<String> {
266 metadata
267 .get("authorization")
268 .and_then(|value| value.to_str().ok())
269 .and_then(|value| value.strip_prefix("Bearer "))
270 .map(str::to_string)
271 .or_else(|| {
272 metadata
273 .get("x-macp-token")
274 .and_then(|value| value.to_str().ok())
275 .map(str::to_string)
276 })
277 }
278
279 pub async fn authenticate_metadata(
280 &self,
281 metadata: &MetadataMap,
282 ) -> Result<AuthIdentity, MacpError> {
283 if let Some(chain) = &self.auth_chain {
287 return chain.authenticate(metadata).await;
288 }
289
290 if !self.identities.is_empty() {
292 if let Some(token) = Self::bearer_token(metadata) {
293 return self
294 .identities
295 .get(&token)
296 .cloned()
297 .ok_or(MacpError::Unauthenticated);
298 }
299 return Err(MacpError::Unauthenticated);
300 }
301
302 self.dev_authenticate(metadata)
304 }
305
306 pub fn authorize_mode(
307 &self,
308 identity: &AuthIdentity,
309 mode: &str,
310 is_session_start: bool,
311 ) -> Result<(), MacpError> {
312 if is_session_start && !identity.can_start_sessions {
313 return Err(MacpError::Forbidden);
314 }
315 if let Some(allowed_modes) = &identity.allowed_modes {
316 if !allowed_modes.contains(mode) {
317 return Err(MacpError::Forbidden);
318 }
319 }
320 Ok(())
321 }
322
323 pub fn authorize_mode_registry(&self, identity: &AuthIdentity) -> Result<(), MacpError> {
324 if identity.can_manage_mode_registry {
325 Ok(())
326 } else {
327 Err(MacpError::Forbidden)
328 }
329 }
330
331 async fn check_bucket(
332 bucket: &Mutex<HashMap<String, VecDeque<Instant>>>,
333 sweep_counter: &std::sync::atomic::AtomicU64,
334 sender: &str,
335 config: &RateLimitConfig,
336 ) -> Result<(), MacpError> {
337 let now = Instant::now();
338 let mut guard = bucket.lock().await;
339
340 const SWEEP_EVERY: u64 = 128;
348 let tick = sweep_counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
349 if tick.is_multiple_of(SWEEP_EVERY) {
350 guard.retain(|_, deque| {
351 deque
352 .back()
353 .map(|last| now.duration_since(*last) <= config.window)
354 .unwrap_or(false)
355 });
356 }
357
358 let deque = guard.entry(sender.to_string()).or_default();
359 while deque
360 .front()
361 .map(|instant| now.duration_since(*instant) > config.window)
362 .unwrap_or(false)
363 {
364 deque.pop_front();
365 }
366 if deque.len() >= config.limit {
367 return Err(MacpError::RateLimited);
368 }
369 deque.push_back(now);
370 Ok(())
371 }
372
373 pub fn has_configured_auth(&self) -> bool {
377 self.auth_chain.is_some()
378 }
379
380 pub async fn enforce_rate_limit(
381 &self,
382 sender: &str,
383 is_session_start: bool,
384 ) -> Result<(), MacpError> {
385 if is_session_start {
386 Self::check_bucket(
387 &self.rate_bucket.start_events,
388 &self.rate_bucket.start_sweep_counter,
389 sender,
390 &self.session_start_rate,
391 )
392 .await
393 } else {
394 Self::check_bucket(
395 &self.rate_bucket.message_events,
396 &self.rate_bucket.message_sweep_counter,
397 sender,
398 &self.message_rate,
399 )
400 .await
401 }
402 }
403}
404
405#[cfg(test)]
406mod tests {
407 use super::*;
408 use std::io::Write;
409 use tempfile::NamedTempFile;
410 use tonic::metadata::MetadataMap;
411
412 fn layer_with_tokens(json: &str) -> SecurityLayer {
415 let identities = SecurityLayer::parse_identities(json).expect("valid JSON");
416 SecurityLayer {
417 identities: Arc::new(identities),
418 rate_bucket: Arc::new(RateBucket::default()),
419 auth_chain: None,
420 max_payload_bytes: 1_048_576,
421 session_start_rate: RateLimitConfig {
422 limit: usize::MAX,
423 window: Duration::from_secs(60),
424 },
425 message_rate: RateLimitConfig {
426 limit: usize::MAX,
427 window: Duration::from_secs(60),
428 },
429 }
430 }
431
432 fn insecure_layer() -> SecurityLayer {
434 SecurityLayer {
435 identities: Arc::new(HashMap::new()),
436 rate_bucket: Arc::new(RateBucket::default()),
437 auth_chain: None,
438 max_payload_bytes: 1_048_576,
439 session_start_rate: RateLimitConfig {
440 limit: usize::MAX,
441 window: Duration::from_secs(60),
442 },
443 message_rate: RateLimitConfig {
444 limit: usize::MAX,
445 window: Duration::from_secs(60),
446 },
447 }
448 }
449
450 #[tokio::test]
455 async fn dev_mode_requires_dev_header() {
456 let layer = SecurityLayer::dev_mode();
457 let meta = MetadataMap::new();
458 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
459 assert!(matches!(err, MacpError::Unauthenticated));
460 }
461
462 #[tokio::test]
463 async fn dev_mode_rejects_dev_sender_header() {
464 let layer = SecurityLayer::dev_mode();
465 let mut meta = MetadataMap::new();
466 meta.insert("x-macp-agent-id", "agent://dev-bot".parse().unwrap());
467 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
468 assert!(matches!(err, MacpError::Unauthenticated));
469 }
470
471 #[test]
472 fn dev_mode_has_unlimited_rate_limits() {
473 let layer = SecurityLayer::dev_mode();
474 assert_eq!(layer.session_start_rate.limit, usize::MAX);
475 assert_eq!(layer.message_rate.limit, usize::MAX);
476 }
477
478 #[test]
483 fn from_env_defaults_without_env_vars() {
484 let layer = insecure_layer();
486 assert_eq!(layer.max_payload_bytes, 1_048_576);
487 }
488
489 #[tokio::test]
494 async fn bearer_token_authentication_via_authorization_header() {
495 let json = r#"[{"token":"tok-abc","sender":"agent://alice","allowed_modes":[],"can_start_sessions":true}]"#;
496 let layer = layer_with_tokens(json);
497
498 let mut meta = MetadataMap::new();
499 meta.insert("authorization", "Bearer tok-abc".parse().unwrap());
500
501 let id = layer
502 .authenticate_metadata(&meta)
503 .await
504 .expect("should authenticate");
505 assert_eq!(id.sender, "agent://alice");
506 assert!(id.allowed_modes.is_none()); assert!(id.can_start_sessions);
508 }
509
510 #[tokio::test]
511 async fn bearer_token_authentication_via_x_macp_token_header() {
512 let json = r#"[{"token":"tok-xyz","sender":"agent://bob"}]"#;
513 let layer = layer_with_tokens(json);
514
515 let mut meta = MetadataMap::new();
516 meta.insert("x-macp-token", "tok-xyz".parse().unwrap());
517
518 let id = layer
519 .authenticate_metadata(&meta)
520 .await
521 .expect("should authenticate");
522 assert_eq!(id.sender, "agent://bob");
523 }
524
525 #[tokio::test]
526 async fn invalid_bearer_token_returns_unauthenticated() {
527 let json = r#"[{"token":"tok-real","sender":"agent://alice"}]"#;
528 let layer = layer_with_tokens(json);
529
530 let mut meta = MetadataMap::new();
531 meta.insert("authorization", "Bearer tok-fake".parse().unwrap());
532
533 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
534 assert!(matches!(err, MacpError::Unauthenticated));
535 }
536
537 #[tokio::test]
538 async fn no_token_when_auth_required_returns_unauthenticated() {
539 let json = r#"[{"token":"tok-only","sender":"agent://sole"}]"#;
540 let layer = layer_with_tokens(json);
541
542 let meta = MetadataMap::new(); let err = layer.authenticate_metadata(&meta).await.unwrap_err();
544 assert!(matches!(err, MacpError::Unauthenticated));
545 }
546
547 #[tokio::test]
548 async fn parse_identities_wrapped_format() {
549 let json = r#"{"tokens":[{"token":"t1","sender":"agent://wrapped"}]}"#;
550 let layer = layer_with_tokens(json);
551
552 let mut meta = MetadataMap::new();
553 meta.insert("authorization", "Bearer t1".parse().unwrap());
554 let id = layer
555 .authenticate_metadata(&meta)
556 .await
557 .expect("should authenticate");
558 assert_eq!(id.sender, "agent://wrapped");
559 }
560
561 #[tokio::test]
562 async fn parse_identities_with_allowed_modes() {
563 let json = r#"[{"token":"t-modes","sender":"agent://limited","allowed_modes":["macp.mode.decision.v1","macp.mode.task.v1"],"can_start_sessions":false,"max_open_sessions":5}]"#;
564 let layer = layer_with_tokens(json);
565
566 let mut meta = MetadataMap::new();
567 meta.insert("authorization", "Bearer t-modes".parse().unwrap());
568 let id = layer
569 .authenticate_metadata(&meta)
570 .await
571 .expect("should authenticate");
572
573 assert_eq!(id.sender, "agent://limited");
574 assert!(!id.can_start_sessions);
575 assert_eq!(id.max_open_sessions, Some(5));
576 let modes = id
577 .allowed_modes
578 .as_ref()
579 .expect("should have allowed_modes");
580 assert!(modes.contains("macp.mode.decision.v1"));
581 assert!(modes.contains("macp.mode.task.v1"));
582 assert!(!modes.contains("macp.mode.proposal.v1"));
583 }
584
585 #[tokio::test]
586 async fn authorization_header_takes_priority_over_x_macp_token() {
587 let json = r#"[
588 {"token":"bearer-tok","sender":"agent://bearer-user"},
589 {"token":"header-tok","sender":"agent://header-user"}
590 ]"#;
591 let layer = layer_with_tokens(json);
592
593 let mut meta = MetadataMap::new();
594 meta.insert("authorization", "Bearer bearer-tok".parse().unwrap());
595 meta.insert("x-macp-token", "header-tok".parse().unwrap());
596
597 let id = layer
598 .authenticate_metadata(&meta)
599 .await
600 .expect("should authenticate");
601 assert_eq!(id.sender, "agent://bearer-user");
603 }
604
605 #[tokio::test]
610 async fn dev_sender_header_rejected_without_chain() {
611 let layer = SecurityLayer {
612 identities: Arc::new(HashMap::new()),
613 rate_bucket: Arc::new(RateBucket::default()),
614 auth_chain: None,
615 max_payload_bytes: 1_048_576,
616 session_start_rate: RateLimitConfig {
617 limit: usize::MAX,
618 window: Duration::from_secs(60),
619 },
620 message_rate: RateLimitConfig {
621 limit: usize::MAX,
622 window: Duration::from_secs(60),
623 },
624 };
625
626 let mut meta = MetadataMap::new();
627 meta.insert("x-macp-agent-id", "agent://dev-agent".parse().unwrap());
628
629 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
630 assert!(matches!(err, MacpError::Unauthenticated));
631 }
632
633 #[tokio::test]
634 async fn dev_sender_header_ignored_when_not_allowed() {
635 let layer = SecurityLayer {
637 identities: Arc::new(HashMap::new()),
638 rate_bucket: Arc::new(RateBucket::default()),
639 auth_chain: None,
640 max_payload_bytes: 1_048_576,
641 session_start_rate: RateLimitConfig {
642 limit: usize::MAX,
643 window: Duration::from_secs(60),
644 },
645 message_rate: RateLimitConfig {
646 limit: usize::MAX,
647 window: Duration::from_secs(60),
648 },
649 };
650
651 let mut meta = MetadataMap::new();
652 meta.insert("x-macp-agent-id", "agent://sneaky".parse().unwrap());
653
654 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
655 assert!(matches!(err, MacpError::Unauthenticated));
656 }
657
658 #[tokio::test]
659 async fn bearer_token_takes_priority_over_dev_header() {
660 let json = r#"[{"token":"real-tok","sender":"agent://real"}]"#;
661 let identities = SecurityLayer::parse_identities(json).unwrap();
662
663 let layer = SecurityLayer {
664 identities: Arc::new(identities),
665 rate_bucket: Arc::new(RateBucket::default()),
666 auth_chain: None,
667 max_payload_bytes: 1_048_576,
668 session_start_rate: RateLimitConfig {
669 limit: usize::MAX,
670 window: Duration::from_secs(60),
671 },
672 message_rate: RateLimitConfig {
673 limit: usize::MAX,
674 window: Duration::from_secs(60),
675 },
676 };
677
678 let mut meta = MetadataMap::new();
679 meta.insert("authorization", "Bearer real-tok".parse().unwrap());
680 meta.insert("x-macp-agent-id", "agent://dev-override".parse().unwrap());
681
682 let id = layer
683 .authenticate_metadata(&meta)
684 .await
685 .expect("should authenticate via bearer");
686 assert_eq!(id.sender, "agent://real");
687 }
688
689 #[test]
694 fn authorize_mode_allows_any_mode_when_no_restriction() {
695 let layer = SecurityLayer::dev_mode();
696 let id = AuthIdentity {
697 sender: "agent://any".into(),
698 allowed_modes: None,
699 can_start_sessions: true,
700 max_open_sessions: None,
701 can_manage_mode_registry: false,
702 is_observer: false,
703 };
704 assert!(layer
705 .authorize_mode(&id, "macp.mode.decision.v1", false)
706 .is_ok());
707 assert!(layer.authorize_mode(&id, "macp.mode.task.v1", true).is_ok());
708 assert!(layer.authorize_mode(&id, "arbitrary.mode", false).is_ok());
709 }
710
711 #[test]
712 fn authorize_mode_rejects_unlisted_mode() {
713 let layer = SecurityLayer::dev_mode();
714 let mut allowed = HashSet::new();
715 allowed.insert("macp.mode.decision.v1".to_string());
716
717 let id = AuthIdentity {
718 sender: "agent://restricted".into(),
719 allowed_modes: Some(allowed),
720 can_start_sessions: true,
721 max_open_sessions: None,
722 can_manage_mode_registry: false,
723 is_observer: false,
724 };
725 assert!(layer
726 .authorize_mode(&id, "macp.mode.decision.v1", false)
727 .is_ok());
728 let err = layer
729 .authorize_mode(&id, "macp.mode.task.v1", false)
730 .unwrap_err();
731 assert!(matches!(err, MacpError::Forbidden));
732 }
733
734 #[test]
735 fn authorize_mode_rejects_session_start_when_not_allowed() {
736 let layer = SecurityLayer::dev_mode();
737 let id = AuthIdentity {
738 sender: "agent://no-start".into(),
739 allowed_modes: None,
740 can_start_sessions: false,
741 max_open_sessions: None,
742 can_manage_mode_registry: false,
743 is_observer: false,
744 };
745 let err = layer
746 .authorize_mode(&id, "macp.mode.decision.v1", true)
747 .unwrap_err();
748 assert!(matches!(err, MacpError::Forbidden));
749 }
750
751 #[test]
752 fn authorize_mode_allows_non_session_start_even_when_start_forbidden() {
753 let layer = SecurityLayer::dev_mode();
754 let id = AuthIdentity {
755 sender: "agent://no-start".into(),
756 allowed_modes: None,
757 can_start_sessions: false,
758 max_open_sessions: None,
759 can_manage_mode_registry: false,
760 is_observer: false,
761 };
762 assert!(layer
764 .authorize_mode(&id, "macp.mode.decision.v1", false)
765 .is_ok());
766 }
767
768 #[test]
769 fn authorize_mode_checks_both_can_start_and_allowed_modes() {
770 let layer = SecurityLayer::dev_mode();
771 let mut allowed = HashSet::new();
772 allowed.insert("macp.mode.decision.v1".to_string());
773
774 let id = AuthIdentity {
775 sender: "agent://double-check".into(),
776 allowed_modes: Some(allowed),
777 can_start_sessions: false,
778 max_open_sessions: None,
779 can_manage_mode_registry: false,
780 is_observer: false,
781 };
782
783 let err = layer
785 .authorize_mode(&id, "macp.mode.decision.v1", true)
786 .unwrap_err();
787 assert!(matches!(err, MacpError::Forbidden));
788
789 let err = layer
791 .authorize_mode(&id, "macp.mode.task.v1", false)
792 .unwrap_err();
793 assert!(matches!(err, MacpError::Forbidden));
794
795 assert!(layer
797 .authorize_mode(&id, "macp.mode.decision.v1", false)
798 .is_ok());
799 }
800
801 #[test]
802 fn authorize_mode_registry_requires_explicit_privilege() {
803 let layer = SecurityLayer::dev_mode();
804 let id = AuthIdentity {
805 sender: "agent://no-admin".into(),
806 allowed_modes: None,
807 can_start_sessions: true,
808 max_open_sessions: None,
809 is_observer: false,
810 can_manage_mode_registry: false,
811 };
812 let err = layer.authorize_mode_registry(&id).unwrap_err();
813 assert!(matches!(err, MacpError::Forbidden));
814 }
815
816 #[tokio::test]
817 async fn bearer_token_can_manage_mode_registry() {
818 let json =
819 r#"[{"token":"admin-tok","sender":"agent://admin","can_manage_mode_registry":true}]"#;
820 let layer = layer_with_tokens(json);
821 let mut meta = MetadataMap::new();
822 meta.insert("authorization", "Bearer admin-tok".parse().unwrap());
823 let id = layer.authenticate_metadata(&meta).await.unwrap();
824 assert!(layer.authorize_mode_registry(&id).is_ok());
825 }
826
827 #[tokio::test]
832 async fn rate_limit_session_start_enforced() {
833 let layer = SecurityLayer {
834 identities: Arc::new(HashMap::new()),
835 rate_bucket: Arc::new(RateBucket::default()),
836 auth_chain: None,
837 max_payload_bytes: 1_048_576,
838 session_start_rate: RateLimitConfig {
839 limit: 3,
840 window: Duration::from_secs(60),
841 },
842 message_rate: RateLimitConfig {
843 limit: usize::MAX,
844 window: Duration::from_secs(60),
845 },
846 };
847
848 let sender = "agent://rate-test";
849 for _ in 0..3 {
851 assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
852 }
853 let err = layer.enforce_rate_limit(sender, true).await.unwrap_err();
855 assert!(matches!(err, MacpError::RateLimited));
856
857 assert!(layer.enforce_rate_limit(sender, false).await.is_ok());
859 }
860
861 #[tokio::test]
862 async fn rate_limit_message_enforced() {
863 let layer = SecurityLayer {
864 identities: Arc::new(HashMap::new()),
865 rate_bucket: Arc::new(RateBucket::default()),
866 auth_chain: None,
867 max_payload_bytes: 1_048_576,
868 session_start_rate: RateLimitConfig {
869 limit: usize::MAX,
870 window: Duration::from_secs(60),
871 },
872 message_rate: RateLimitConfig {
873 limit: 2,
874 window: Duration::from_secs(60),
875 },
876 };
877
878 let sender = "agent://msg-test";
879 assert!(layer.enforce_rate_limit(sender, false).await.is_ok());
880 assert!(layer.enforce_rate_limit(sender, false).await.is_ok());
881 let err = layer.enforce_rate_limit(sender, false).await.unwrap_err();
882 assert!(matches!(err, MacpError::RateLimited));
883
884 assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
886 }
887
888 #[tokio::test]
889 async fn rate_limit_per_sender_isolation() {
890 let layer = SecurityLayer {
891 identities: Arc::new(HashMap::new()),
892 rate_bucket: Arc::new(RateBucket::default()),
893 auth_chain: None,
894 max_payload_bytes: 1_048_576,
895 session_start_rate: RateLimitConfig {
896 limit: 1,
897 window: Duration::from_secs(60),
898 },
899 message_rate: RateLimitConfig {
900 limit: usize::MAX,
901 window: Duration::from_secs(60),
902 },
903 };
904
905 assert!(layer.enforce_rate_limit("agent://a", true).await.is_ok());
907 assert!(layer.enforce_rate_limit("agent://a", true).await.is_err());
908
909 assert!(layer.enforce_rate_limit("agent://b", true).await.is_ok());
911 }
912
913 #[tokio::test]
914 async fn rate_limit_window_expiry() {
915 let layer = SecurityLayer {
916 identities: Arc::new(HashMap::new()),
917 rate_bucket: Arc::new(RateBucket::default()),
918 auth_chain: None,
919 max_payload_bytes: 1_048_576,
920 session_start_rate: RateLimitConfig {
921 limit: 1,
922 window: Duration::from_millis(1), },
924 message_rate: RateLimitConfig {
925 limit: usize::MAX,
926 window: Duration::from_secs(60),
927 },
928 };
929
930 let sender = "agent://expiry-test";
931 assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
932
933 tokio::time::sleep(Duration::from_millis(5)).await;
935
936 assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
938 }
939
940 #[tokio::test]
945 async fn no_anonymous_fallback_even_when_auth_not_required() {
946 let layer = insecure_layer();
947 let meta = MetadataMap::new();
948 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
949 assert!(matches!(err, MacpError::Unauthenticated));
950 }
951
952 #[tokio::test]
953 async fn no_anonymous_fallback_when_auth_required() {
954 let json = r#"[{"token":"t","sender":"agent://real"}]"#;
955 let layer = layer_with_tokens(json);
956
957 let meta = MetadataMap::new();
958 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
959 assert!(matches!(err, MacpError::Unauthenticated));
960 }
961
962 #[tokio::test]
963 async fn dev_mode_no_fallback_with_empty_metadata() {
964 let layer = SecurityLayer::dev_mode();
967 let meta = MetadataMap::new();
968 let err = layer.authenticate_metadata(&meta).await.unwrap_err();
969 assert!(matches!(err, MacpError::Unauthenticated));
970 }
971
972 #[test]
977 fn token_file_loading_via_parse_identities() {
978 let json = r#"[
981 {"token":"file-tok-1","sender":"agent://file-alice","allowed_modes":["macp.mode.decision.v1"]},
982 {"token":"file-tok-2","sender":"agent://file-bob","can_start_sessions":false}
983 ]"#;
984 let mut tmp = NamedTempFile::new().expect("create temp file");
985 write!(tmp, "{}", json).expect("write temp file");
986
987 let contents = fs::read_to_string(tmp.path()).expect("read temp file");
988 let identities = SecurityLayer::parse_identities(&contents).expect("parse identities");
989
990 assert_eq!(identities.len(), 2);
991
992 let alice = identities.get("file-tok-1").expect("alice entry");
993 assert_eq!(alice.sender, "agent://file-alice");
994 let alice_modes = alice.allowed_modes.as_ref().expect("should have modes");
995 assert!(alice_modes.contains("macp.mode.decision.v1"));
996 assert!(alice.can_start_sessions); let bob = identities.get("file-tok-2").expect("bob entry");
999 assert_eq!(bob.sender, "agent://file-bob");
1000 assert!(!bob.can_start_sessions);
1001 assert!(bob.allowed_modes.is_none()); }
1003
1004 #[tokio::test]
1005 async fn token_file_end_to_end_via_layer() {
1006 let json = r#"[{"token":"e2e-tok","sender":"agent://e2e-agent"}]"#;
1008 let mut tmp = NamedTempFile::new().expect("create temp file");
1009 write!(tmp, "{}", json).expect("write temp file");
1010
1011 let contents = fs::read_to_string(tmp.path()).expect("read temp file");
1012 let layer = layer_with_tokens(&contents);
1013
1014 let mut meta = MetadataMap::new();
1015 meta.insert("authorization", "Bearer e2e-tok".parse().unwrap());
1016 let id = layer
1017 .authenticate_metadata(&meta)
1018 .await
1019 .expect("should authenticate");
1020 assert_eq!(id.sender, "agent://e2e-agent");
1021 }
1022
1023 #[test]
1024 fn parse_identities_invalid_json_returns_error() {
1025 let result = SecurityLayer::parse_identities("not valid json");
1026 assert!(result.is_err());
1027 }
1028
1029 #[test]
1030 fn parse_identities_empty_list() {
1031 let identities = SecurityLayer::parse_identities("[]").expect("valid empty list");
1032 assert!(identities.is_empty());
1033 }
1034
1035 #[test]
1036 fn parse_identities_wrapped_empty() {
1037 let identities =
1038 SecurityLayer::parse_identities(r#"{"tokens":[]}"#).expect("valid wrapped empty");
1039 assert!(identities.is_empty());
1040 }
1041
1042 #[tokio::test]
1046 async fn rate_bucket_sweep_removes_stale_senders() {
1047 let bucket = RateBucket::default();
1048 let config = RateLimitConfig {
1049 limit: 10,
1050 window: Duration::from_millis(1),
1051 };
1052
1053 for i in 0..50 {
1055 SecurityLayer::check_bucket(
1056 &bucket.start_events,
1057 &bucket.start_sweep_counter,
1058 &format!("agent://stale-{i}"),
1059 &config,
1060 )
1061 .await
1062 .unwrap();
1063 }
1064 assert!(bucket.start_events.lock().await.len() >= 49);
1065
1066 tokio::time::sleep(Duration::from_millis(5)).await;
1068
1069 for _ in 0..80 {
1072 SecurityLayer::check_bucket(
1073 &bucket.start_events,
1074 &bucket.start_sweep_counter,
1075 "agent://fresh",
1076 &config,
1077 )
1078 .await
1079 .ok(); }
1081 let map = bucket.start_events.lock().await;
1082 assert!(
1083 map.len() <= 2,
1084 "stale senders must be swept; map still has {} entries",
1085 map.len()
1086 );
1087 assert!(map.contains_key("agent://fresh"));
1088 }
1089}