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