Skip to main content

basil_core/service/
broker.rs

1// SPDX-FileCopyrightText: 2026 OpenBasil Contributors
2//
3// SPDX-License-Identifier: Apache-2.0
4
5//! gRPC broker service adapters.
6//!
7//! This module is intentionally a thin transport adapter over the existing
8//! [`BrokerState`] + [`BackendManager`] core. It preserves the same PDP gating
9//! model as the JSON handler: every key-scoped RPC authorizes the
10//! kernel-attested peer uid before dispatching.
11
12#![allow(clippy::result_large_err)]
13
14use std::pin::Pin;
15use std::sync::Arc;
16
17use futures::Stream;
18use std::sync::Mutex;
19use tonic::{Code, Request, Response, Status};
20
21use crate::actor::AuthenticatedActor;
22use crate::catalog::policy::Op;
23use crate::state::BrokerState;
24use crate::transport::{authorize, broker_status};
25
26pub(super) type GrpcResult<T> = Result<Response<T>, Status>;
27pub(super) type BoxStream<T> = Pin<Box<dyn Stream<Item = Result<T, Status>> + Send + 'static>>;
28
29/// Broker identity and response-signing key settings for sealed invocation.
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub struct BrokerIdentityRuntimeConfig {
32    /// Stable broker audience / identity URI.
33    pub id: String,
34    /// Catalog key id used to sign invocation responses.
35    pub response_signing_key_id: String,
36}
37
38/// Runtime settings for the sealed invocation service.
39#[derive(Debug, Clone, PartialEq, Eq)]
40pub struct InvocationRuntimeConfig {
41    /// Whether `InvocationService.Invoke` accepts requests.
42    pub enabled: bool,
43    /// Broker identity and response-signing key. Required when enabled.
44    pub broker_identity: Option<BrokerIdentityRuntimeConfig>,
45    /// Accepted broker audiences. An omitted header audience may be derived only
46    /// when this contains exactly one value.
47    pub audiences: Vec<String>,
48    /// Catalog key id whose public half receives sealed invocation requests.
49    pub request_encryption_key_id: Option<String>,
50    /// Maximum accepted signed request TTL in seconds.
51    pub max_ttl_secs: u32,
52    /// Allowed clock skew in seconds for issue and expiry timestamps.
53    pub clock_skew_secs: u32,
54    /// Maximum replay-cache entries retained in memory.
55    pub replay_cache_capacity: usize,
56    /// Fixed current time override for deterministic tests. Leave unset in
57    /// production.
58    pub now_unix_override: Option<u32>,
59}
60
61impl Default for InvocationRuntimeConfig {
62    fn default() -> Self {
63        Self {
64            enabled: false,
65            broker_identity: None,
66            audiences: Vec::new(),
67            request_encryption_key_id: None,
68            max_ttl_secs: basil_proto::invocation::DEFAULT_EXPIRES_AFTER_SECS,
69            clock_skew_secs: 30,
70            replay_cache_capacity: 4096,
71            now_unix_override: None,
72        }
73    }
74}
75
76/// Shared implementation for all broker gRPC services.
77#[derive(Debug, Clone)]
78pub struct BrokerGrpc {
79    pub(super) state: Arc<BrokerState>,
80    pub(super) invocation: InvocationRuntimeConfig,
81    pub(super) invocation_replay_cache:
82        Arc<Mutex<crate::service::invocation::InvocationReplayCache>>,
83}
84
85impl BrokerGrpc {
86    /// Build a gRPC service adapter over shared broker state.
87    #[must_use]
88    pub fn new(state: Arc<BrokerState>) -> Self {
89        Self::new_with_invocation(state, false)
90    }
91
92    /// Build a gRPC service adapter and explicitly configure invocation serving.
93    #[must_use]
94    pub fn new_with_invocation(state: Arc<BrokerState>, invocation_enabled: bool) -> Self {
95        Self::new_with_invocation_config(
96            state,
97            InvocationRuntimeConfig {
98                enabled: invocation_enabled,
99                ..InvocationRuntimeConfig::default()
100            },
101        )
102    }
103
104    /// Build a gRPC service adapter with full invocation runtime settings.
105    #[must_use]
106    pub fn new_with_invocation_config(
107        state: Arc<BrokerState>,
108        invocation: InvocationRuntimeConfig,
109    ) -> Self {
110        let capacity = invocation.replay_cache_capacity;
111        Self {
112            state,
113            invocation,
114            invocation_replay_cache: Arc::new(Mutex::new(
115                crate::service::invocation::InvocationReplayCache::new(capacity),
116            )),
117        }
118    }
119
120    pub(super) fn invocation_now_unix(&self) -> u32 {
121        if let Some(now) = self.invocation.now_unix_override {
122            return now;
123        }
124        std::time::SystemTime::now()
125            .duration_since(std::time::UNIX_EPOCH)
126            .map_or(0, |duration| {
127                u32::try_from(duration.as_secs()).unwrap_or(u32::MAX)
128            })
129    }
130
131    pub(super) fn authorize<T>(
132        &self,
133        request: &Request<T>,
134        op: Op,
135        key: &str,
136    ) -> Result<AuthenticatedActor, Status> {
137        authorize(&self.state, request, op, key)
138    }
139
140    pub(super) fn visible(&self, actor: &AuthenticatedActor, key: &str) -> bool {
141        let generation = self.state.load_generation();
142        let pdp = generation.pdp();
143        pdp.decide(actor, Op::List, key).is_allow()
144            || pdp.decide(actor, Op::GetPublicKey, key).is_allow()
145            || pdp.decide(actor, Op::Get, key).is_allow()
146    }
147
148    pub(super) fn require_unix_uid(
149        actor: &AuthenticatedActor,
150        op: &'static str,
151    ) -> Result<u32, Status> {
152        actor.unix_uid().ok_or_else(|| {
153            broker_status(
154                Code::Unauthenticated,
155                "UNAUTHENTICATED",
156                op,
157                "operation requires local peer credentials",
158            )
159        })
160    }
161}
162
163#[cfg(test)]
164mod tests {
165    use std::collections::BTreeMap;
166    use std::sync::Arc;
167
168    use async_trait::async_trait;
169    use base64::Engine as _;
170    use basil_proto::KeyType;
171    use basil_proto::broker::v1 as pb;
172    use basil_proto::broker::v1::BrokerErrorInfo;
173    use basil_proto::broker::v1::admin_service_server::AdminService;
174    use basil_proto::broker::v1::minting_service_server::MintingService;
175    use basil_proto::broker::v1::nats_service_server::NatsService;
176    use basil_proto::broker::v1::signing_service_server::SigningService;
177    use basil_proto::google::rpc::Status as RpcStatus;
178    use nkeys::{KeyPair, XKey};
179    use prost::Message;
180    use prost_types::{Duration, Struct, Value as ProstValue, value::Kind as ProstKind};
181    use serde_json::Value as JsonValue;
182    use tonic::{Code, Request, Status};
183
184    use super::BrokerGrpc;
185    use crate::backend::{Backend, BackendError, KvValue, NewKey, PublicKey};
186    use crate::catalog::load;
187    use crate::manager::{BackendManager, ManagerError};
188    use crate::peer::PeerInfo;
189    use crate::service::minting::nats_mint_status;
190    use crate::service::shared::*;
191    use crate::state::BrokerState;
192    use zeroize::Zeroizing;
193
194    fn error_info(status: &Status) -> BrokerErrorInfo {
195        let rpc = RpcStatus::decode(status.details()).expect("status details decode");
196        let detail = rpc.details.first().expect("broker detail present");
197        BrokerErrorInfo::decode(detail.value.as_slice()).expect("broker error info decodes")
198    }
199
200    fn assert_status_omits(status: &Status, canaries: &[&str]) {
201        let info = error_info(status);
202        let visible = format!(
203            "{} {} {} {:?}",
204            status.message(),
205            info.reason,
206            info.op,
207            status.details()
208        );
209        for canary in canaries {
210            assert!(
211                !visible.contains(canary),
212                "status leaked secret canary `{canary}` in `{visible}`"
213            );
214        }
215    }
216
217    #[test]
218    fn unsupported_algorithm_maps_to_unimplemented_detail() {
219        let status = backend_status(
220            "encrypt",
221            &BackendError::UnsupportedAlgorithm(basil_proto::AeadAlgorithm::Aes256Gcm),
222        );
223        let info = error_info(&status);
224        assert_eq!(status.code(), Code::Unimplemented);
225        assert_eq!(info.reason, "UNSUPPORTED_ALGORITHM");
226        assert_eq!(info.op, "encrypt");
227    }
228
229    #[test]
230    fn invalid_request_maps_to_invalid_argument_detail() {
231        let status = key_type(0, "new_key").expect_err("unspecified key type is invalid");
232        let info = error_info(&status);
233        assert_eq!(status.code(), Code::InvalidArgument);
234        assert_eq!(info.reason, "INVALID_REQUEST");
235        assert_eq!(info.op, "new_key");
236    }
237
238    #[test]
239    fn pqc_key_types_map_to_unimplemented_detail() {
240        let status = key_type(pb::KeyType::MlDsa65.into(), "new_key")
241            .expect_err("ML-DSA generation is not implemented yet");
242        let info = error_info(&status);
243        assert_eq!(status.code(), Code::Unimplemented);
244        assert_eq!(info.reason, "UNSUPPORTED_ALGORITHM");
245        assert_eq!(info.op, "new_key");
246
247        let status = key_type(pb::KeyType::MlKem768.into(), "new_key")
248            .expect_err("ML-KEM generation is not implemented yet");
249        let info = error_info(&status);
250        assert_eq!(status.code(), Code::Unimplemented);
251        assert_eq!(info.reason, "UNSUPPORTED_ALGORITHM");
252        assert_eq!(info.op, "new_key");
253    }
254
255    #[test]
256    fn signing_algorithms_are_accepted_at_the_wire_layer() {
257        ensure_supported_signing_algorithm(pb::SigningAlgorithm::Ed25519.into(), "sign")
258            .expect("legacy signing remains supported");
259        // ML-DSA is now serviceable: the wire validation accepts it and the
260        // manager dispatches a software-custodied ML-DSA key through the provider.
261        for algorithm in [
262            pb::SigningAlgorithm::MlDsa44,
263            pb::SigningAlgorithm::MlDsa65,
264            pb::SigningAlgorithm::MlDsa87,
265        ] {
266            ensure_supported_signing_algorithm(algorithm.into(), "sign")
267                .expect("ML-DSA signing is serviceable");
268        }
269    }
270
271    #[test]
272    fn kem_envelope_algorithms_are_validated() {
273        // X25519 and ML-KEM are recognized KEMs for envelope unwrap. Broker-side
274        // wrap remains X25519-only and is rejected in the AEAD service.
275        ensure_supported_kem_algorithm(pb::KemAlgorithm::X25519.into(), "wrap_envelope")
276            .expect("X25519 KEM is serviceable");
277        ensure_supported_kem_algorithm(pb::KemAlgorithm::MlKem1024.into(), "unwrap_envelope")
278            .expect("ML-KEM unwrap is serviceable");
279        ensure_supported_envelope_algorithm(
280            pb::EnvelopeAlgorithm::Aes256Gcm.into(),
281            "wrap_envelope",
282        )
283        .expect("envelope AEAD contract value is recognized");
284
285        let status = ensure_supported_kem_algorithm(0, "wrap_envelope")
286            .expect_err("missing KEM algorithm is invalid");
287        let info = error_info(&status);
288        assert_eq!(status.code(), Code::InvalidArgument);
289        assert_eq!(info.reason, "INVALID_REQUEST");
290        assert_eq!(info.op, "wrap_envelope");
291    }
292
293    #[test]
294    fn payload_cap_maps_to_resource_exhausted_detail() {
295        let status = payload_too_large("set", "too large");
296        let info = error_info(&status);
297        assert_eq!(status.code(), Code::ResourceExhausted);
298        assert_eq!(info.reason, "PAYLOAD_TOO_LARGE");
299        assert_eq!(info.op, "set");
300    }
301
302    #[test]
303    fn backend_transport_maps_to_unavailable_detail() {
304        let status = backend_status("sign", &BackendError::Transport("down".to_string()));
305        let info = error_info(&status);
306        assert_eq!(status.code(), Code::Unavailable);
307        assert_eq!(status.message(), "backend unavailable");
308        assert_eq!(info.reason, "BACKEND_UNAVAILABLE");
309        assert_eq!(info.op, "sign");
310    }
311
312    #[test]
313    fn backend_status_omits_secret_bearing_upstream_details() {
314        let canaries = [
315            "vault-token-s.123",
316            "Authorization: Bearer secret",
317            "/run/credentials/basil/passphrase",
318            "-----BEGIN PRIVATE KEY-----",
319            "upstream-response-body-with-credential",
320        ];
321        for err in [
322            BackendError::Transport(canaries[1].to_string()),
323            BackendError::Backend(canaries[4].to_string()),
324            BackendError::Protocol(canaries[3].to_string()),
325        ] {
326            let status = backend_status("sign", &err);
327            assert_status_omits(&status, &canaries);
328        }
329    }
330
331    #[test]
332    fn decrypt_failure_is_fixed_opaque_invalid_argument() {
333        let status = backend_status("decrypt", &BackendError::DecryptFailed);
334        let info = error_info(&status);
335        assert_eq!(status.code(), Code::InvalidArgument);
336        assert_eq!(status.message(), "decrypt failed");
337        assert_eq!(info.reason, "DECRYPT_FAILED");
338        assert_eq!(info.op, "decrypt");
339    }
340
341    #[test]
342    fn unknown_key_on_gated_manager_path_is_unauthorized_not_not_found() {
343        let status = manager_status("sign", &ManagerError::UnknownKey("hidden.key".to_string()));
344        let info = error_info(&status);
345        assert_eq!(status.code(), Code::PermissionDenied);
346        assert_eq!(status.message(), "not authorized");
347        assert_eq!(info.reason, "UNAUTHORIZED");
348        assert_eq!(info.op, "sign");
349    }
350
351    #[test]
352    fn duration_conversion_rejects_negative_and_rounds_fraction_up() {
353        let ttl = ttl_seconds(
354            Some(&Duration {
355                seconds: 4,
356                nanos: 1,
357            }),
358            "mint",
359        )
360        .expect("ttl converts");
361        assert_eq!(ttl, Some(5));
362
363        let status = ttl_seconds(
364            Some(&Duration {
365                seconds: -1,
366                nanos: 0,
367            }),
368            "mint",
369        )
370        .expect_err("negative ttl rejected");
371        assert_eq!(status.code(), Code::InvalidArgument);
372        assert_eq!(error_info(&status).reason, "INVALID_REQUEST");
373    }
374
375    #[test]
376    fn struct_claims_convert_to_json_and_reject_non_finite_numbers() {
377        let claims = Struct {
378            fields: [
379                (
380                    "scope".to_string(),
381                    ProstValue {
382                        kind: Some(ProstKind::StringValue("read".to_string())),
383                    },
384                ),
385                (
386                    "version".to_string(),
387                    ProstValue {
388                        kind: Some(ProstKind::NumberValue(2.0)),
389                    },
390                ),
391                (
392                    "ratio".to_string(),
393                    ProstValue {
394                        kind: Some(ProstKind::NumberValue(2.5)),
395                    },
396                ),
397            ]
398            .into(),
399        };
400        let json = claims_json(Some(&claims), "mint").expect("claims convert");
401        assert_eq!(json["scope"], JsonValue::String("read".to_string()));
402        assert_eq!(json["version"], JsonValue::Number(2.into()));
403        assert_eq!(
404            json["ratio"],
405            JsonValue::Number(serde_json::Number::from_f64(2.5).expect("finite"))
406        );
407
408        let claims = Struct {
409            fields: [(
410                "bad".to_string(),
411                ProstValue {
412                    kind: Some(ProstKind::NumberValue(f64::NAN)),
413                },
414            )]
415            .into(),
416        };
417        let status = claims_json(Some(&claims), "mint").expect_err("nan rejected");
418        assert_eq!(status.code(), Code::InvalidArgument);
419        assert_eq!(error_info(&status).reason, "INVALID_REQUEST");
420    }
421
422    #[test]
423    fn bad_subject_nkey_maps_to_invalid_argument() {
424        let status = nats_mint_status(
425            "mint_nats_user",
426            &BackendError::Protocol("invalid subject user nkey: bad".to_string()),
427        );
428        let info = error_info(&status);
429        assert_eq!(status.code(), Code::InvalidArgument);
430        assert_eq!(info.reason, "INVALID_REQUEST");
431        assert_eq!(info.op, "mint_nats_user");
432    }
433
434    struct MintBackend;
435
436    #[async_trait]
437    impl Backend for MintBackend {
438        fn kind(&self) -> &'static str {
439            "mint-test"
440        }
441
442        async fn new_key(&self, key_type: KeyType) -> Result<NewKey, BackendError> {
443            let _ = key_type;
444            Err(BackendError::Unsupported("new_key"))
445        }
446
447        async fn public_key(&self, key_id: &str) -> Result<Vec<u8>, BackendError> {
448            let _ = key_id;
449            Ok(vec![7; 32])
450        }
451
452        async fn sign(&self, key_id: &str, message: &[u8]) -> Result<Vec<u8>, BackendError> {
453            let _ = (key_id, message);
454            Ok(vec![9; 64])
455        }
456
457        async fn verify(
458            &self,
459            key_id: &str,
460            message: &[u8],
461            signature: &[u8],
462        ) -> Result<bool, BackendError> {
463            let _ = (key_id, message, signature);
464            Ok(true)
465        }
466
467        async fn kv_get(
468            &self,
469            key_id: &str,
470            version: Option<u32>,
471        ) -> Result<KvValue, BackendError> {
472            let _ = version;
473            match key_id {
474                "nats/curve-box-public" => {
475                    let private = Zeroizing::new([0x55; 32]);
476                    Ok(KvValue {
477                        value: basil_nats::xkey_public_from_private(&private).to_vec(),
478                        version: 1,
479                    })
480                }
481                _ => Err(BackendError::KeyNotFound(key_id.to_string())),
482            }
483        }
484
485        async fn kv_get_secret(
486            &self,
487            key_id: &str,
488            version: Option<u32>,
489        ) -> Result<crate::backend::KvSecret, BackendError> {
490            let _ = version;
491            match key_id {
492                "nats/curve-box" => Ok(crate::backend::KvSecret {
493                    value: Zeroizing::new(vec![0x55; 32]),
494                    version: 1,
495                }),
496                _ => Err(BackendError::KeyNotFound(key_id.to_string())),
497            }
498        }
499    }
500
501    /// Real NATS-key signer used to prove `SigningService.Sign` can complete a
502    /// caller-assembled rich NATS JWT without exposing the issuer seed.
503    struct NatsSignBackend(KeyPair);
504
505    #[async_trait]
506    impl Backend for NatsSignBackend {
507        fn kind(&self) -> &'static str {
508            "nkey-sign-test"
509        }
510
511        async fn new_key(&self, key_type: KeyType) -> Result<NewKey, BackendError> {
512            let _ = key_type;
513            Err(BackendError::Unsupported("new_key"))
514        }
515
516        async fn public_key(&self, key_id: &str) -> Result<Vec<u8>, BackendError> {
517            let _ = key_id;
518            let (_, public) = basil_nats::decode_public(&self.0.public_key())
519                .map_err(|e| BackendError::Protocol(e.to_string()))?;
520            Ok(public.to_vec())
521        }
522
523        async fn public_key_with_meta(&self, key_id: &str) -> Result<PublicKey, BackendError> {
524            Ok(PublicKey {
525                public_key: self.public_key(key_id).await?,
526                key_type: KeyType::Ed25519Nkey,
527                version: 1,
528            })
529        }
530
531        async fn sign(&self, key_id: &str, message: &[u8]) -> Result<Vec<u8>, BackendError> {
532            let _ = key_id;
533            self.0
534                .sign(message)
535                .map_err(|e| BackendError::Backend(e.to_string()))
536        }
537
538        async fn verify(
539            &self,
540            key_id: &str,
541            message: &[u8],
542            signature: &[u8],
543        ) -> Result<bool, BackendError> {
544            let _ = key_id;
545            Ok(self.0.verify(message, signature).is_ok())
546        }
547    }
548
549    fn state_with_backend(backend: Box<dyn Backend>) -> Arc<BrokerState> {
550        let catalog = r#"{
551          "schemaVersion": 1,
552          "backends": { "bao": { "kind": "vault", "addr": "https://127.0.0.1:8200" } },
553          "keys": {
554            "issuer.account": {
555              "class": "asymmetric", "keyType": "ed25519-nkey", "backend": "bao",
556              "path": "issuer/account", "writable": true, "missing": "error",
557              "labels": ["nats_type=A"], "description": "NATS account issuer"
558            },
559            "issuer.operator": {
560              "class": "asymmetric", "keyType": "ed25519-nkey", "backend": "bao",
561              "path": "issuer/operator", "writable": true, "missing": "error",
562              "labels": ["nats_type=O"], "description": "NATS operator issuer"
563            },
564            "issuer.server": {
565              "class": "asymmetric", "keyType": "ed25519-nkey", "backend": "bao",
566              "path": "issuer/server", "writable": true, "missing": "error",
567              "labels": ["nats_type=N"], "description": "NATS server issuer"
568            },
569            "issuer.curve": {
570              "class": "asymmetric", "keyType": "ed25519-nkey", "backend": "bao",
571              "path": "issuer/curve", "writable": true, "missing": "error",
572              "labels": ["nats_type=X"], "description": "NATS curve issuer"
573            },
574            "nats.curve_box": {
575              "class": "sealing", "keyType": "x25519", "backend": "bao",
576              "path": "nats/curve-box", "publicPath": "nats/curve-box-public",
577              "writable": false, "missing": "error",
578              "description": "NATS xkey box custody"
579            }
580          }
581        }"#;
582        let policy = r#"{
583          "schemaVersion": 2,
584          "subjects": {
585            "svc.mint": { "allOf": [ { "kind": "unix", "uid": 42 } ] }
586          },
587          "roles": {
588            "minter": ["mint", "sign_nats_jwt", "validate_nats_jwt"],
589            "reader": ["get", "list", "get_public_key"],
590            "signer": ["sign", "verify"],
591            "nats_box": ["encrypt_nats_curve", "decrypt_nats_curve"]
592          },
593          "rules": [
594            { "id": "mint", "subjects": ["svc.mint"], "action": ["role:minter", "role:reader", "role:signer"], "target": ["issuer.*"] },
595            { "id": "nats-box", "subjects": ["svc.mint"], "action": ["role:nats_box"], "target": ["nats.curve_box"] },
596            { "id": "watch", "subjects": ["svc.mint"], "action": ["op:watch"], "target": ["broker.watch"] }
597          ],
598          "config": {
599            "names": { "users": { "42": "svc-mint" }, "groups": {} },
600            "memberships": { "42": [42] }
601          }
602        }"#;
603        let (catalog, policy, config, warnings) = load(catalog, policy).expect("fixture loads");
604        assert!(warnings.is_empty());
605        let mut backends: BTreeMap<String, Box<dyn Backend>> = BTreeMap::new();
606        backends.insert("bao".to_string(), backend);
607        let manager = BackendManager::new(catalog.clone(), backends).expect("manager builds");
608        Arc::new(BrokerState::new(
609            catalog,
610            policy,
611            config,
612            manager,
613            "mint-test",
614        ))
615    }
616
617    fn mint_state() -> Arc<BrokerState> {
618        state_with_backend(Box::new(MintBackend))
619    }
620
621    fn authed_request<T>(body: T) -> Request<T> {
622        let mut request = Request::new(body);
623        request.extensions_mut().insert(PeerInfo {
624            uid: Some(42),
625            ..PeerInfo::default()
626        });
627        request
628    }
629
630    fn ttl() -> Duration {
631        Duration {
632            seconds: 60,
633            nanos: 0,
634        }
635    }
636
637    fn rich_account_import(exporting: &KeyPair) -> basil_nats::AccountImport {
638        basil_nats::AccountImport {
639            name: "control-read".to_string(),
640            subject: "$JS.API.CONSUMER.MSG.NEXT.control_delivery".to_string(),
641            account: exporting.public_key(),
642            token: String::new(),
643            to: String::new(),
644            local_subject: "R3.$JS.API.CONSUMER.MSG.NEXT.control_delivery".to_string(),
645            kind: basil_nats::ExportType::Service,
646            share: true,
647            allow_trace: true,
648        }
649    }
650
651    fn rich_account_export(revocations: BTreeMap<String, i64>) -> basil_nats::AccountExport {
652        basil_nats::AccountExport {
653            name: "device-messages".to_string(),
654            subject: "dev.*.*.>".to_string(),
655            kind: basil_nats::ExportType::Stream,
656            token_req: true,
657            revocations,
658            response_type: None,
659            response_threshold: 0,
660            service_latency: None,
661            account_token_position: 2,
662            advertise: true,
663            allow_trace: true,
664            description: "realm device delivery".to_string(),
665            info_url: "https://basil.example.test/nats".to_string(),
666        }
667    }
668
669    fn rich_account_limits() -> basil_nats::OperatorLimits {
670        basil_nats::OperatorLimits {
671            nats: basil_nats::NatsLimits {
672                subs: 512,
673                data: 1_048_576,
674                payload: 262_144,
675            },
676            account: basil_nats::AccountLimits {
677                imports: 32,
678                exports: 32,
679                wildcards: true,
680                disallow_bearer: true,
681                conn: 256,
682                leaf: 8,
683            },
684            jetstream: basil_nats::JetStreamLimits {
685                mem_storage: 67_108_864,
686                disk_storage: 1_073_741_824,
687                streams: 64,
688                consumer: 512,
689                max_ack_pending: 10_000,
690                mem_max_stream_bytes: 8_388_608,
691                disk_max_stream_bytes: 134_217_728,
692                max_bytes_required: true,
693            },
694            tiered_limits: BTreeMap::new(),
695        }
696    }
697
698    fn rich_account_default_permissions() -> basil_nats::Permissions {
699        basil_nats::Permissions {
700            publish: basil_nats::Permission {
701                allow: vec!["dev.realm.device.>".to_string()],
702                deny: vec!["dev.realm.device.private.>".to_string()],
703            },
704            sub: basil_nats::Permission {
705                allow: vec!["dist.>".to_string(), "dev.realm.device.>".to_string()],
706                deny: Vec::new(),
707            },
708            resp: Some(basil_nats::ResponsePermission {
709                max: 1,
710                ttl: 2_000_000_000,
711            }),
712        }
713    }
714
715    fn rich_account_claims(exporting: &KeyPair, user: &KeyPair) -> basil_nats::AccountClaims {
716        let mut revocations = BTreeMap::new();
717        revocations.insert(user.public_key(), 1_782_000_001);
718        let mut export_revocations = BTreeMap::new();
719        export_revocations.insert("*".to_string(), 1_782_000_002);
720        basil_nats::AccountClaims {
721            imports: vec![rich_account_import(exporting)],
722            exports: vec![rich_account_export(export_revocations)],
723            limits: rich_account_limits(),
724            signing_keys: Vec::new(),
725            revocations,
726            default_permissions: Some(rich_account_default_permissions()),
727            mappings: BTreeMap::new(),
728            authorization: basil_nats::ExternalAuthorization::default(),
729            trace: Some(basil_nats::MsgTrace {
730                dest: "trace.realm".to_string(),
731                sampling: 25,
732            }),
733            cluster_traffic: Some(basil_nats::ClusterTraffic::System),
734        }
735    }
736
737    fn rich_account_jwt(
738        issuer_nkey: String,
739        account: &KeyPair,
740        signing: &KeyPair,
741        exporting: &KeyPair,
742        user: &KeyPair,
743    ) -> basil_nats::AccountJwt {
744        basil_nats::AccountJwt {
745            issuer: issuer_nkey,
746            subject_account: account.public_key(),
747            name: "basil-realm".to_string(),
748            issued_at: 1_782_000_000,
749            expires: None,
750            signing_keys: vec![signing.public_key()],
751            claims: rich_account_claims(exporting, user),
752        }
753    }
754
755    #[tokio::test]
756    #[allow(clippy::too_many_lines)]
757    async fn grpc_minting_methods_return_credentials() {
758        let service = BrokerGrpc::new(mint_state());
759        let generic = service
760            .mint_jwt(authed_request(pb::MintJwtRequest {
761                key_id: "issuer.account".to_string(),
762                subject: Some("subject".to_string()),
763                ttl: Some(ttl()),
764                claims: None,
765            }))
766            .await
767            .expect("generic mint succeeds")
768            .into_inner();
769        assert!(!generic.token.is_empty());
770        assert!(generic.expires_at.is_some());
771
772        let user_nkey = basil_nats::encode_public(basil_nats::NkeyType::User, &[2; 32])
773            .expect("test user public key encodes");
774        let user = service
775            .mint_nats_user(authed_request(pb::MintNatsUserRequest {
776                key_id: "issuer.account".to_string(),
777                subject_user_nkey: user_nkey,
778                issuer_account: None,
779                name: "user".to_string(),
780                ttl: Some(ttl()),
781                pub_allow: Vec::new(),
782                pub_deny: Vec::new(),
783                sub_allow: Vec::new(),
784                sub_deny: Vec::new(),
785            }))
786            .await
787            .expect("user mint succeeds")
788            .into_inner();
789        assert!(!user.token.is_empty());
790
791        let signed_nats = service
792            .sign_nats_jwt(authed_request(pb::SignNatsJwtRequest {
793                key_id: "issuer.account".to_string(),
794                claims_json: serde_json::to_vec(&serde_json::json!({
795                    "sub": basil_nats::encode_public(basil_nats::NkeyType::User, &[8; 32])
796                        .expect("test user public key encodes"),
797                    "name": "rich-user",
798                    "nats": { "type": "user", "version": 2 }
799                }))
800                .expect("claims json"),
801                expected_type: pb::NatsJwtType::User.into(),
802                ttl: Some(ttl()),
803                expires_at: None,
804                issued_at: None,
805                jti_mode: pb::NatsJtiMode::RequireValid.into(),
806            }))
807            .await
808            .expect("validated nats jwt sign succeeds")
809            .into_inner();
810        assert!(!signed_nats.token.is_empty());
811        assert!(signed_nats.expires_at.is_some());
812
813        let account_nkey = basil_nats::encode_public(basil_nats::NkeyType::Account, &[3; 32])
814            .expect("test account public key encodes");
815        let account = service
816            .mint_nats_account(authed_request(pb::MintNatsAccountRequest {
817                key_id: "issuer.operator".to_string(),
818                subject_account_nkey: account_nkey,
819                name: "account".to_string(),
820                ttl: Some(ttl()),
821                signing_keys: Vec::new(),
822            }))
823            .await
824            .expect("account mint succeeds")
825            .into_inner();
826        assert!(!account.token.is_empty());
827
828        let operator = service
829            .mint_nats_operator(authed_request(pb::MintNatsOperatorRequest {
830                key_id: "issuer.operator".to_string(),
831                subject_operator_nkey: None,
832                name: "operator".to_string(),
833                ttl: Some(ttl()),
834                signing_keys: Vec::new(),
835                account_server_url: None,
836                system_account: None,
837            }))
838            .await
839            .expect("operator mint succeeds")
840            .into_inner();
841        assert!(!operator.token.is_empty());
842
843        let signer = service
844            .mint_nats_signer(authed_request(pb::MintNatsSignerRequest {
845                key_id: "issuer.account".to_string(),
846                subject_nkey: basil_nats::encode_public(basil_nats::NkeyType::Account, &[4; 32])
847                    .expect("test account signer public key encodes"),
848                name: "signer".to_string(),
849                ttl: Some(ttl()),
850            }))
851            .await
852            .expect("signer mint succeeds")
853            .into_inner();
854        assert!(!signer.token.is_empty());
855
856        let server = service
857            .mint_nats_server(authed_request(pb::MintNatsServerRequest {
858                key_id: "issuer.server".to_string(),
859                subject_server_nkey: basil_nats::encode_public(
860                    basil_nats::NkeyType::Server,
861                    &[5; 32],
862                )
863                .expect("test server public key encodes"),
864                name: "server".to_string(),
865                ttl: Some(ttl()),
866            }))
867            .await
868            .expect("server mint succeeds")
869            .into_inner();
870        assert!(!server.token.is_empty());
871
872        let curve = service
873            .mint_nats_curve(authed_request(pb::MintNatsCurveRequest {
874                key_id: "issuer.curve".to_string(),
875                subject_curve_nkey: basil_nats::encode_public(
876                    basil_nats::NkeyType::Curve,
877                    &[6; 32],
878                )
879                .expect("test curve public key encodes"),
880                name: "curve".to_string(),
881                ttl: Some(ttl()),
882            }))
883            .await
884            .expect("curve mint succeeds")
885            .into_inner();
886        assert!(!curve.token.is_empty());
887    }
888
889    #[tokio::test]
890    async fn sign_nats_jwt_ttl_uses_claim_iat_as_base() {
891        let account = KeyPair::new_account();
892        let user = KeyPair::new_user();
893        let service = BrokerGrpc::new(state_with_backend(Box::new(NatsSignBackend(account))));
894        let token = service
895            .sign_nats_jwt(authed_request(pb::SignNatsJwtRequest {
896                key_id: "issuer.account".to_string(),
897                claims_json: serde_json::to_vec(&serde_json::json!({
898                    "iat": 1_700_000_000_u64,
899                    "sub": user.public_key(),
900                    "nats": { "type": "user", "version": 2 }
901                }))
902                .expect("claims json"),
903                expected_type: pb::NatsJwtType::User.into(),
904                ttl: Some(Duration {
905                    seconds: 60,
906                    nanos: 0,
907                }),
908                expires_at: None,
909                issued_at: None,
910                jti_mode: pb::NatsJtiMode::RequireValid.into(),
911            }))
912            .await
913            .expect("sign succeeds")
914            .into_inner()
915            .token;
916
917        let parts: Vec<&str> = token.split('.').collect();
918        let claims: JsonValue = serde_json::from_slice(
919            &base64::engine::general_purpose::URL_SAFE_NO_PAD
920                .decode(parts[1])
921                .expect("claims decode"),
922        )
923        .expect("claims json");
924        assert_eq!(claims["iat"], 1_700_000_000_u64);
925        assert_eq!(claims["exp"], 1_700_000_060_u64);
926    }
927
928    #[tokio::test]
929    async fn grpc_nats_curve_encrypt_decrypt_interops_with_nkeys() {
930        let service = BrokerGrpc::new(mint_state());
931        let broker_private = [0x55; 32];
932        let broker_xkey = XKey::new_from_raw(broker_private);
933        let peer = XKey::new_from_raw([0x66; 32]);
934
935        let encrypted = service
936            .encrypt_nats_curve(authed_request(pb::EncryptNatsCurveRequest {
937                key_id: "nats.curve_box".to_string(),
938                recipient_public_xkey: peer.public_key(),
939                plaintext: b"broker-to-peer".to_vec(),
940            }))
941            .await
942            .expect("encrypt succeeds")
943            .into_inner()
944            .ciphertext;
945        let opened_by_peer = peer
946            .open(&encrypted, &broker_xkey)
947            .expect("nkeys opens broker ciphertext");
948        assert_eq!(opened_by_peer, b"broker-to-peer");
949
950        let peer_box = peer
951            .seal(b"peer-to-broker", &broker_xkey)
952            .expect("nkeys seals");
953        let decrypted = service
954            .decrypt_nats_curve(authed_request(pb::DecryptNatsCurveRequest {
955                key_id: "nats.curve_box".to_string(),
956                sender_public_xkey: peer.public_key(),
957                ciphertext: peer_box,
958            }))
959            .await
960            .expect("decrypt succeeds")
961            .into_inner()
962            .plaintext;
963        assert_eq!(decrypted, b"peer-to-broker");
964
965        let denied = service
966            .encrypt_nats_curve(authed_request(pb::EncryptNatsCurveRequest {
967                key_id: "issuer.account".to_string(),
968                recipient_public_xkey: peer.public_key(),
969                plaintext: b"wrong class".to_vec(),
970            }))
971            .await
972            .expect_err("policy denies ungranted target");
973        assert_eq!(denied.code(), Code::PermissionDenied);
974    }
975
976    #[tokio::test]
977    async fn grpc_sign_on_issuer_key_is_hard_capped_but_rich_jwt_mints_via_sign_nats_jwt() {
978        // Raw `sign` on a NATS operator/account key would let a caller assemble
979        // the JWT signing input off-broker and bypass every validation the
980        // dedicated `sign_nats_jwt` op enforces (kind, jti mode, ttl/expiry),
981        // so the PDP hard-caps it even though the policy grants role:signer
982        // over `issuer.*`. The same rich claims still mint through the
983        // validated op: rejecting the bypass loses no legitimate capability.
984        let operator = KeyPair::new_operator();
985        let operator_public = operator.public_key();
986        let account = KeyPair::new_account();
987        let signing = KeyPair::new_account();
988        let exporting = KeyPair::new_account();
989        let user = KeyPair::new_user();
990        let service = BrokerGrpc::new(state_with_backend(Box::new(NatsSignBackend(operator))));
991
992        let public = service
993            .get_public_key(authed_request(pb::GetPublicKeyRequest {
994                key_id: "issuer.operator".to_string(),
995                version: None,
996            }))
997            .await
998            .expect("public key read succeeds")
999            .into_inner()
1000            .public_key;
1001        let issuer_nkey = basil_nats::encode_public(basil_nats::NkeyType::Operator, &public)
1002            .expect("issuer public encodes as operator nkey");
1003        assert_eq!(issuer_nkey, operator_public);
1004
1005        let jwt = rich_account_jwt(issuer_nkey, &account, &signing, &exporting, &user);
1006        let signing_input = jwt.signing_input().expect("rich signing input builds");
1007
1008        let denied = service
1009            .sign(authed_request(pb::SignRequest {
1010                key_id: "issuer.operator".to_string(),
1011                message: signing_input.as_bytes().to_vec(),
1012                algorithm: pb::SigningAlgorithm::Ed25519Nkey.into(),
1013            }))
1014            .await
1015            .expect_err("raw sign on an issuer key is hard-capped");
1016        assert_eq!(denied.code(), Code::PermissionDenied);
1017
1018        // Route the identical rich claims through the validated minting op.
1019        let claims_json = base64::engine::general_purpose::URL_SAFE_NO_PAD
1020            .decode(signing_input.split('.').nth(1).expect("claims part"))
1021            .expect("claims decode");
1022        let token = service
1023            .sign_nats_jwt(authed_request(pb::SignNatsJwtRequest {
1024                key_id: "issuer.operator".to_string(),
1025                claims_json,
1026                expected_type: pb::NatsJwtType::Account.into(),
1027                ttl: Some(ttl()),
1028                expires_at: None,
1029                issued_at: None,
1030                jti_mode: pb::NatsJtiMode::Rewrite.into(),
1031            }))
1032            .await
1033            .expect("rich account JWT mints through the validated op")
1034            .into_inner()
1035            .token;
1036
1037        let parts: Vec<&str> = token.split('.').collect();
1038        assert_eq!(parts.len(), 3);
1039        let claims: JsonValue = serde_json::from_slice(
1040            &base64::engine::general_purpose::URL_SAFE_NO_PAD
1041                .decode(parts[1])
1042                .expect("claims decode"),
1043        )
1044        .expect("claims json");
1045        let nats = &claims["nats"];
1046        assert_eq!(
1047            nats["imports"][0]["local_subject"],
1048            "R3.$JS.API.CONSUMER.MSG.NEXT.control_delivery"
1049        );
1050        assert_eq!(nats["exports"][0]["token_req"], true);
1051        assert_eq!(nats["limits"]["disk_storage"], 1_073_741_824);
1052        assert_eq!(nats["default_permissions"]["resp"]["max"], 1);
1053        assert_eq!(nats["trace"]["dest"], "trace.realm");
1054
1055        let issuer = KeyPair::from_public_key(claims["iss"].as_str().expect("issuer claim"))
1056            .expect("issuer public key parses");
1057        let signature = base64::engine::general_purpose::URL_SAFE_NO_PAD
1058            .decode(parts[2])
1059            .expect("signature decodes");
1060        issuer
1061            .verify(format!("{}.{}", parts[0], parts[1]).as_bytes(), &signature)
1062            .expect("Basil signature verifies under issuer nkey");
1063    }
1064
1065    #[tokio::test]
1066    #[allow(clippy::too_many_lines)]
1067    async fn grpc_validate_nats_jwt_reports_authoritative_reasons() {
1068        let account = KeyPair::new_account();
1069        let account_public = account.public_key();
1070        let user = KeyPair::new_user();
1071        let service = BrokerGrpc::new(state_with_backend(Box::new(NatsSignBackend(account))));
1072
1073        let token = service
1074            .sign_nats_jwt(authed_request(pb::SignNatsJwtRequest {
1075                key_id: "issuer.account".to_string(),
1076                claims_json: serde_json::to_vec(&serde_json::json!({
1077                    "sub": user.public_key(),
1078                    "name": "valid-user",
1079                    "nats": { "type": "user", "version": 2 }
1080                }))
1081                .expect("claims json"),
1082                expected_type: pb::NatsJwtType::User.into(),
1083                ttl: Some(ttl()),
1084                expires_at: None,
1085                issued_at: None,
1086                jti_mode: pb::NatsJtiMode::RequireValid.into(),
1087            }))
1088            .await
1089            .expect("sign succeeds")
1090            .into_inner()
1091            .token;
1092
1093        let by_key = service
1094            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1095                jwt: token.clone(),
1096                allowed_signers: vec![pb::AllowedNatsSigner {
1097                    signer: Some(pb::allowed_nats_signer::Signer::KeyId(
1098                        "issuer.account".to_string(),
1099                    )),
1100                }],
1101                expected_type: pb::NatsJwtType::User.into(),
1102            }))
1103            .await
1104            .expect("validate by key succeeds")
1105            .into_inner();
1106        assert!(by_key.valid);
1107        assert_eq!(by_key.reason, i32::from(pb::NatsJwtValidationReason::Valid));
1108        assert_eq!(by_key.issuer, account_public.as_str());
1109        assert_eq!(by_key.subject, user.public_key());
1110        assert_eq!(by_key.matched_signer_key_id, "issuer.account");
1111
1112        let by_public_nkey = service
1113            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1114                jwt: token.clone(),
1115                allowed_signers: vec![pb::AllowedNatsSigner {
1116                    signer: Some(pb::allowed_nats_signer::Signer::NatsPublicKey(
1117                        account_public.clone(),
1118                    )),
1119                }],
1120                expected_type: pb::NatsJwtType::User.into(),
1121            }))
1122            .await
1123            .expect("validate by nkey succeeds")
1124            .into_inner();
1125        assert!(by_public_nkey.valid);
1126        assert!(by_public_nkey.matched_signer_key_id.is_empty());
1127
1128        let wrong_type = service
1129            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1130                jwt: token.clone(),
1131                allowed_signers: Vec::new(),
1132                expected_type: pb::NatsJwtType::Account.into(),
1133            }))
1134            .await
1135            .expect("wrong type is authoritative")
1136            .into_inner();
1137        assert!(!wrong_type.valid);
1138        assert_eq!(
1139            wrong_type.reason,
1140            i32::from(pb::NatsJwtValidationReason::WrongType)
1141        );
1142
1143        let unknown = service
1144            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1145                jwt: token,
1146                allowed_signers: vec![pb::AllowedNatsSigner {
1147                    signer: Some(pb::allowed_nats_signer::Signer::NatsPublicKey(
1148                        KeyPair::new_operator().public_key(),
1149                    )),
1150                }],
1151                expected_type: pb::NatsJwtType::User.into(),
1152            }))
1153            .await
1154            .expect("unknown signer is authoritative")
1155            .into_inner();
1156        assert!(!unknown.valid);
1157        assert_eq!(
1158            unknown.reason,
1159            i32::from(pb::NatsJwtValidationReason::UnknownSigner)
1160        );
1161
1162        let malformed = service
1163            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1164                jwt: "not-a-jwt".to_string(),
1165                allowed_signers: Vec::new(),
1166                expected_type: pb::NatsJwtType::Unspecified.into(),
1167            }))
1168            .await
1169            .expect("malformed is authoritative")
1170            .into_inner();
1171        assert!(!malformed.valid);
1172        assert_eq!(
1173            malformed.reason,
1174            i32::from(pb::NatsJwtValidationReason::Malformed)
1175        );
1176    }
1177
1178    #[tokio::test]
1179    async fn grpc_validate_nats_jwt_requires_a_resolved_subject() {
1180        let service = BrokerGrpc::new(mint_state());
1181        // uid 7 resolves to no policy subject: the RPC must fail closed at
1182        // entry, before the caller-supplied-nkey arm can run (finding 16).
1183        let mut request = Request::new(pb::ValidateNatsJwtRequest {
1184            jwt: "not-a-jwt".to_string(),
1185            allowed_signers: vec![pb::AllowedNatsSigner {
1186                signer: Some(pb::allowed_nats_signer::Signer::NatsPublicKey(
1187                    KeyPair::new_account().public_key(),
1188                )),
1189            }],
1190            expected_type: pb::NatsJwtType::Unspecified.into(),
1191        });
1192        request.extensions_mut().insert(PeerInfo {
1193            uid: Some(7),
1194            ..PeerInfo::default()
1195        });
1196        let status = service
1197            .validate_nats_jwt(request)
1198            .await
1199            .expect_err("unresolved peer rejected");
1200        assert_eq!(status.code(), Code::Unauthenticated);
1201    }
1202
1203    #[tokio::test]
1204    async fn grpc_validate_nats_jwt_caps_jwt_length() {
1205        let service = BrokerGrpc::new(mint_state());
1206        let oversized = "a".repeat(service.state.limits().max_payload_size + 1);
1207        let status = service
1208            .validate_nats_jwt(authed_request(pb::ValidateNatsJwtRequest {
1209                jwt: oversized,
1210                allowed_signers: Vec::new(),
1211                expected_type: pb::NatsJwtType::Unspecified.into(),
1212            }))
1213            .await
1214            .expect_err("oversized jwt rejected");
1215        assert_eq!(status.code(), Code::ResourceExhausted);
1216    }
1217
1218    #[tokio::test]
1219    async fn grpc_sign_verify_and_nats_curve_enforce_payload_caps() {
1220        let service = BrokerGrpc::new(mint_state());
1221        let over_payload = vec![0u8; service.state.limits().max_payload_size + 1];
1222        let over_encrypt = vec![0u8; service.state.limits().max_encrypt_size + 1];
1223
1224        // `issuer.server` (nats_type=N): not a credential issuer, so raw
1225        // sign/verify pass the PDP hard cap and reach the payload cap.
1226        let status = service
1227            .sign(authed_request(pb::SignRequest {
1228                key_id: "issuer.server".to_string(),
1229                message: over_payload.clone(),
1230                algorithm: pb::SigningAlgorithm::Ed25519Nkey.into(),
1231            }))
1232            .await
1233            .expect_err("oversized sign message rejected");
1234        assert_eq!(status.code(), Code::ResourceExhausted);
1235
1236        let status = service
1237            .verify(authed_request(pb::VerifyRequest {
1238                key_id: "issuer.server".to_string(),
1239                message: Vec::new(),
1240                signature: over_payload,
1241                algorithm: pb::SigningAlgorithm::Ed25519Nkey.into(),
1242            }))
1243            .await
1244            .expect_err("oversized verify signature rejected");
1245        assert_eq!(status.code(), Code::ResourceExhausted);
1246
1247        let status = service
1248            .encrypt_nats_curve(authed_request(pb::EncryptNatsCurveRequest {
1249                key_id: "nats.curve_box".to_string(),
1250                recipient_public_xkey: String::new(),
1251                plaintext: over_encrypt.clone(),
1252            }))
1253            .await
1254            .expect_err("oversized curve plaintext rejected");
1255        assert_eq!(status.code(), Code::ResourceExhausted);
1256
1257        let status = service
1258            .decrypt_nats_curve(authed_request(pb::DecryptNatsCurveRequest {
1259                key_id: "nats.curve_box".to_string(),
1260                sender_public_xkey: String::new(),
1261                ciphertext: over_encrypt,
1262            }))
1263            .await
1264            .expect_err("oversized curve ciphertext rejected");
1265        assert_eq!(status.code(), Code::ResourceExhausted);
1266    }
1267
1268    #[tokio::test]
1269    async fn grpc_nats_unsupported_issuer_role_is_invalid_argument() {
1270        let service = BrokerGrpc::new(mint_state());
1271        let status = service
1272            .mint_nats_server(authed_request(pb::MintNatsServerRequest {
1273                key_id: "issuer.operator".to_string(),
1274                subject_server_nkey: basil_nats::encode_public(
1275                    basil_nats::NkeyType::Server,
1276                    &[5; 32],
1277                )
1278                .expect("test server public key encodes"),
1279                name: "server".to_string(),
1280                ttl: Some(ttl()),
1281            }))
1282            .await
1283            .expect_err("operator issuer is invalid for server mint");
1284        let info = error_info(&status);
1285        assert_eq!(status.code(), Code::InvalidArgument);
1286        assert_eq!(info.reason, "INVALID_REQUEST");
1287        assert_eq!(info.op, "mint_nats_server");
1288    }
1289
1290    #[tokio::test]
1291    async fn watch_requires_peer_uid() {
1292        let service = BrokerGrpc::new(mint_state());
1293        let Err(status) = service
1294            .watch(Request::new(pb::WatchRequest { kinds: Vec::new() }))
1295            .await
1296        else {
1297            panic!("missing uid should be rejected");
1298        };
1299        assert_eq!(status.code(), Code::Unauthenticated);
1300    }
1301
1302    #[tokio::test]
1303    async fn watch_filters_key_rotation_and_allows_public_events() {
1304        use futures::StreamExt as _;
1305
1306        let state = mint_state();
1307        let service = BrokerGrpc::new(Arc::clone(&state));
1308        let mut stream = service
1309            .watch(authed_request(pb::WatchRequest { kinds: Vec::new() }))
1310            .await
1311            .expect("watch opens")
1312            .into_inner();
1313
1314        state.events().key_rotated("hidden.key", 2);
1315        state.events().bundle_changed("example.org");
1316
1317        let event = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next())
1318            .await
1319            .expect("event arrives")
1320            .expect("stream item")
1321            .expect("event ok");
1322        assert_eq!(event.kind, i32::from(pb::EventKind::BundleChanged));
1323        assert!(matches!(
1324            event.detail,
1325            Some(pb::event::Detail::BundleChanged(_))
1326        ));
1327
1328        state.events().key_rotated("issuer.account", 3);
1329        let event = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next())
1330            .await
1331            .expect("event arrives")
1332            .expect("stream item")
1333            .expect("event ok");
1334        assert_eq!(event.kind, i32::from(pb::EventKind::KeyRotated));
1335        assert!(matches!(
1336            event.detail,
1337            Some(pb::event::Detail::KeyRotated(_))
1338        ));
1339    }
1340
1341    #[tokio::test]
1342    async fn watch_emits_revocation_events() {
1343        use futures::StreamExt as _;
1344
1345        let state = mint_state();
1346        let service = BrokerGrpc::new(Arc::clone(&state));
1347        let mut stream = service
1348            .watch(authed_request(pb::WatchRequest {
1349                kinds: vec![i32::from(pb::EventKind::Revoked)],
1350            }))
1351            .await
1352            .expect("watch opens")
1353            .into_inner();
1354
1355        state
1356            .revoke_jwt_svid(
1357                "example.org",
1358                "test-jti",
1359                jsonwebtoken::get_current_timestamp().saturating_add(300),
1360            )
1361            .await
1362            .expect("revoked jti stored");
1363
1364        let event = tokio::time::timeout(std::time::Duration::from_secs(1), stream.next())
1365            .await
1366            .expect("event arrives")
1367            .expect("stream item")
1368            .expect("event ok");
1369        assert_eq!(event.kind, i32::from(pb::EventKind::Revoked));
1370        let Some(pb::event::Detail::Revoked(revoked)) = event.detail else {
1371            panic!("revoked detail expected");
1372        };
1373        assert_eq!(revoked.trust_domain, "example.org");
1374        assert_eq!(revoked.id, "test-jti");
1375    }
1376}