Skip to main content

macp_auth/
security.rs

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    /// Requests since the last full stale-sweep of each map. Full sweeps are
57    /// amortized (every `SWEEP_EVERY` requests) so no single request pays a
58    /// scan proportional to total sender cardinality, while the maps still
59    /// get fully cleaned on a bounded cadence.
60    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    /// Creates a test-friendly SecurityLayer that maps any bearer token
76    /// `"tok-<sender>"` to an identity with `sender = <token-value>`.
77    /// For tests, use `Authorization: Bearer agent://name` to authenticate as `agent://name`.
78    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    /// Dev-mode authenticate: accepts any bearer token as a FULLY-PRIVILEGED
96    /// identity (can start sessions, can manage the mode registry).
97    ///
98    /// Reached whenever no auth is configured — both by `dev_mode()` in tests
99    /// AND by `from_env()` when the operator sets no tokens/issuer. Startup
100    /// therefore refuses to run without configured auth unless
101    /// `MACP_ALLOW_INSECURE=1` (see `has_configured_auth` and `src/main.rs`);
102    /// an operator who forgets auth env vars must not silently run an
103    /// any-token-is-admin server.
104    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        // Build auth resolver chain
154        let mut resolvers: Vec<Box<dyn crate::auth::AuthResolver>> = Vec::new();
155
156        // JWT resolver (if configured)
157        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            // Asymmetric algorithms only by default. HS256 in a default
165            // allowlist is a latent confusion risk: if the JWKS ever contains
166            // an `oct` key, symmetric tokens become verifiable. Operators who
167            // genuinely use shared-secret JWTs must opt in explicitly.
168            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        // Static bearer resolver (always present if tokens are configured)
214        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        // Production path: use the auth resolver chain. Fully async — the
284        // previous block_in_place/block_on bridge parked a worker thread for
285        // the whole JWKS fetch and panicked on a current-thread runtime.
286        if let Some(chain) = &self.auth_chain {
287            return chain.authenticate(metadata).await;
288        }
289
290        // Explicit identity map (layer_with_tokens in tests)
291        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        // Dev-mode: any bearer token → identity (for tests only)
303        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        // Amortized stale-sender sweep. A per-request full scan is O(total
341        // senders) — and sender cardinality is attacker-controllable via
342        // distinct authenticated identities — so the full sweep runs only
343        // every SWEEP_EVERY requests. Between sweeps a request touches only
344        // its own deque. The map therefore stays bounded (a full clean every
345        // SWEEP_EVERY requests) without any request paying the whole scan
346        // more than 1/SWEEP_EVERY of the time.
347        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    /// Whether any real authentication is configured (static tokens and/or a
374    /// JWT resolver). When false, `authenticate_metadata` falls through to
375    /// the any-token-is-admin dev path — callers gate startup on this.
376    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    /// Build a SecurityLayer with bearer token identities loaded from a JSON string.
413    /// This avoids touching environment variables (safe for parallel tests).
414    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    /// Build a SecurityLayer with no tokens that does not require auth.
433    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    // ---------------------------------------------------------------
451    // 1. dev_mode() creates a SecurityLayer that doesn't require auth
452    // ---------------------------------------------------------------
453
454    #[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    // ---------------------------------------------------------------
479    // 2. from_env() with no env vars creates an insecure layer
480    // ---------------------------------------------------------------
481
482    #[test]
483    fn from_env_defaults_without_env_vars() {
484        // Verify default configuration via direct construction.
485        let layer = insecure_layer();
486        assert_eq!(layer.max_payload_bytes, 1_048_576);
487    }
488
489    // ---------------------------------------------------------------
490    // 3. Bearer token auth: loading tokens and authenticating
491    // ---------------------------------------------------------------
492
493    #[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()); // empty vec -> None
507        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(); // no auth header at all
543        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        // Authorization header should take priority
602        assert_eq!(id.sender, "agent://bearer-user");
603    }
604
605    // ---------------------------------------------------------------
606    // 4. Dev header extraction: x-macp-agent-id
607    // ---------------------------------------------------------------
608
609    #[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        // allow_dev_sender_header=false, no tokens
636        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    // ---------------------------------------------------------------
690    // 5. authorize_mode() with allowed modes and without
691    // ---------------------------------------------------------------
692
693    #[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        // Regular messages (not session start) should succeed
763        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        // Cannot start sessions (checked first)
784        let err = layer
785            .authorize_mode(&id, "macp.mode.decision.v1", true)
786            .unwrap_err();
787        assert!(matches!(err, MacpError::Forbidden));
788
789        // Cannot use unlisted mode
790        let err = layer
791            .authorize_mode(&id, "macp.mode.task.v1", false)
792            .unwrap_err();
793        assert!(matches!(err, MacpError::Forbidden));
794
795        // Can send non-start message on allowed mode
796        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    // ---------------------------------------------------------------
828    // 6. enforce_rate_limit() with session_start and message categories
829    // ---------------------------------------------------------------
830
831    #[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        // First 3 should succeed
850        for _ in 0..3 {
851            assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
852        }
853        // 4th should be rate limited
854        let err = layer.enforce_rate_limit(sender, true).await.unwrap_err();
855        assert!(matches!(err, MacpError::RateLimited));
856
857        // Regular messages should still be fine (separate bucket)
858        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        // Session starts should still be fine (separate bucket)
885        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        // Sender A exhausts limit
906        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        // Sender B should still be able to start sessions
910        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), // very short window
923            },
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        // Wait for the window to expire
934        tokio::time::sleep(Duration::from_millis(5)).await;
935
936        // Should succeed again after window expiry
937        assert!(layer.enforce_rate_limit(sender, true).await.is_ok());
938    }
939
940    // ---------------------------------------------------------------
941    // 7. Anonymous fallback behavior
942    // ---------------------------------------------------------------
943
944    #[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        // dev_mode: allow_dev_sender_header=true
965        // With no headers at all, returns Unauthenticated (no anonymous fallback)
966        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    // ---------------------------------------------------------------
973    // 8. Token file loading via MACP_AUTH_TOKENS_FILE
974    // ---------------------------------------------------------------
975
976    #[test]
977    fn token_file_loading_via_parse_identities() {
978        // Test the parse_identities path that from_env uses after reading the file.
979        // We write a temp file and then read + parse it the same way from_env would.
980        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); // default_true
997
998        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()); // empty vec -> None
1002    }
1003
1004    #[tokio::test]
1005    async fn token_file_end_to_end_via_layer() {
1006        // Build a layer as if loaded from a token file, then authenticate with it.
1007        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    /// The amortized sweep must actually bound the sender map: stale senders
1043    /// are fully removed when the periodic full sweep fires, so the map does
1044    /// not grow with total distinct-sender cardinality forever.
1045    #[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        // Tick 0 sweeps the (empty) map; ticks 1..=49 add 49 distinct senders.
1054        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        // Let every recorded event age out of the window.
1067        tokio::time::sleep(Duration::from_millis(5)).await;
1068
1069        // Drive the counter across the next sweep boundary (tick 128) with a
1070        // single fresh sender. After the sweep only the fresh sender remains.
1071        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(); // fresh sender may hit its own limit; irrelevant here
1080        }
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}