Skip to main content

gonzalo_server/
grpc.rs

1//! gRPC transport: adapts the generated `Gonzalo` service to the shared
2//! `Service`, carrying `gonzalo-core` types as JSON payloads.
3
4use crate::Service;
5use crate::auth::{Access, Auth, Principal};
6use gonzalo_core::{
7    ContentHash, DeleteResult, Identity, KeyPrefix, PutResult, Record, RecordKey, Revision,
8};
9use gonzalo_proto::v1::{
10    DeleteBlobRequest, DeleteBlobResponse, DeleteRequest, DeleteResponse, GetBlobRequest,
11    GetBlobResponse, GetRequest, GetResponse, GraphLocatedResponse, GraphNamesResponse,
12    GraphQueryRequest, ListBlobsRequest, ListBlobsResponse, ListRequest, ListResponse,
13    PutBlobRequest, PutBlobResponse, PutRequest, PutResponse, TicketSyncRequest,
14    TicketSyncResponse,
15    gonzalo_server::{Gonzalo, GonzaloServer},
16};
17use serde::Serialize;
18use std::sync::Arc;
19use tonic::metadata::MetadataMap;
20use tonic::{Request, Response, Status};
21
22/// Reserved authz namespace for namespace-agnostic blob ops (ADR 0015), matching
23/// the HTTP transport.
24const BLOB_NS: &str = "_blobs";
25
26/// Adapts [`Service`] to the generated gRPC trait, enforcing namespace-scoped
27/// auth (ADR 0015) per call from the request's bearer metadata.
28pub struct GrpcAdapter {
29    service: Service,
30    auth: Arc<Auth>,
31}
32
33impl GrpcAdapter {
34    /// Adapter with auth disabled (open) — used by tests and open deployments.
35    pub fn new(service: Service) -> Self {
36        Self::with_auth(service, Arc::new(Auth::Disabled))
37    }
38
39    /// Adapter enforcing `auth`.
40    pub fn with_auth(service: Service, auth: Arc<Auth>) -> Self {
41        Self { service, auth }
42    }
43
44    /// Authenticate the call's bearer token and authorize `access` on
45    /// `namespace`. Returns the [`Principal`] (for author stamping on writes).
46    #[allow(clippy::result_large_err)]
47    fn authorize(
48        &self,
49        metadata: &MetadataMap,
50        access: Access,
51        namespace: &str,
52    ) -> Result<Principal, Status> {
53        let principal = self.authenticate(metadata)?;
54        self.check_access(&principal, access, namespace)?;
55        Ok(principal)
56    }
57
58    /// Authenticate the call's bearer token into a [`Principal`], independent of
59    /// any namespace. Split out from [`authorize`] so a handler can reject an
60    /// unauthenticated caller *before* deserializing attacker-controlled JSON
61    /// (#146) and only then authorize against a namespace parsed from the body.
62    #[allow(clippy::result_large_err)]
63    fn authenticate(&self, metadata: &MetadataMap) -> Result<Principal, Status> {
64        self.auth
65            .authenticate(bearer(metadata))
66            .ok_or_else(|| Status::unauthenticated("invalid or missing token"))
67    }
68
69    /// Authorize an already-authenticated `principal` for `access` on `namespace`.
70    #[allow(clippy::result_large_err)]
71    fn check_access(
72        &self,
73        principal: &Principal,
74        access: Access,
75        namespace: &str,
76    ) -> Result<(), Status> {
77        if principal.allows(access, namespace) {
78            Ok(())
79        } else {
80            Err(Status::permission_denied(format!(
81                "principal {:?} lacks {access:?} on namespace {namespace:?}",
82                principal.name()
83            )))
84        }
85    }
86}
87
88/// The bearer token from gRPC `authorization: Bearer <token>` metadata.
89fn bearer(metadata: &MetadataMap) -> Option<&str> {
90    metadata
91        .get("authorization")?
92        .to_str()
93        .ok()?
94        .strip_prefix("Bearer ")
95}
96
97/// Map a backend failure to an opaque `Internal` status. The full error is
98/// logged server-side; the client sees only "internal error" so on-disk graph
99/// paths, SQLite text, and S3 endpoint/bucket detail never leak to the network
100/// (#148).
101fn internal<E: std::fmt::Display>(e: E) -> Status {
102    eprintln!("gonzalod: internal error: {e}");
103    Status::internal("internal error")
104}
105
106#[tonic::async_trait]
107impl Gonzalo for GrpcAdapter {
108    async fn get(&self, req: Request<GetRequest>) -> Result<Response<GetResponse>, Status> {
109        let (metadata, _ext, r) = req.into_parts();
110        self.authorize(&metadata, Access::Read, &r.namespace)?;
111        let key = RecordKey::new(r.namespace, r.collection, r.id);
112        let rec = self.service.get(&key).await.map_err(internal)?;
113        let resp = match rec {
114            Some(rec) => GetResponse {
115                found: true,
116                record_json: serde_json::to_vec(&rec).map_err(internal)?,
117            },
118            None => GetResponse {
119                found: false,
120                record_json: Vec::new(),
121            },
122        };
123        Ok(Response::new(resp))
124    }
125
126    async fn put(&self, req: Request<PutRequest>) -> Result<Response<PutResponse>, Status> {
127        let (metadata, _ext, r) = req.into_parts();
128        // Authenticate BEFORE deserializing attacker-controlled JSON (#146): an
129        // unauthenticated caller is rejected without ever feeding its body to
130        // serde. Only then parse the body (malformed input is the caller's
131        // error → invalid_argument, not internal) and authorize the write
132        // against the namespace named in the record's key.
133        let principal = self.authenticate(&metadata)?;
134        let mut record: Record = serde_json::from_slice(&r.record_json)
135            .map_err(|e| Status::invalid_argument(e.to_string()))?;
136        let expected: Option<Revision> = serde_json::from_slice(&r.expected_json)
137            .map_err(|e| Status::invalid_argument(e.to_string()))?;
138        self.check_access(&principal, Access::Write, &record.key.namespace)?;
139        // Stamp the author from the authenticated principal — unforgeable (ADR
140        // 0015). Open mode (no auth) leaves the record's author untouched.
141        if principal.is_authenticated() {
142            record.meta.author = Identity::new(principal.name());
143        }
144        let outcome = self.service.put(record, expected).await.map_err(internal)?;
145        let resp = match outcome {
146            PutResult::Committed(rev) => PutResponse {
147                outcome: "committed".into(),
148                payload_json: serde_json::to_vec(&rev).map_err(internal)?,
149            },
150            PutResult::Conflict(c) => PutResponse {
151                outcome: "conflict".into(),
152                payload_json: serde_json::to_vec(&*c).map_err(internal)?,
153            },
154        };
155        Ok(Response::new(resp))
156    }
157
158    async fn delete(
159        &self,
160        req: Request<DeleteRequest>,
161    ) -> Result<Response<DeleteResponse>, Status> {
162        let (metadata, _ext, r) = req.into_parts();
163        // Authenticate BEFORE deserializing attacker-controlled JSON (#146), then
164        // parse the precondition (malformed input is the caller's error →
165        // invalid_argument), authorize the write against the path's namespace,
166        // and build the key from (namespace, collection, id).
167        let principal = self.authenticate(&metadata)?;
168        let expected: Option<Revision> = serde_json::from_slice(&r.expected_json)
169            .map_err(|e| Status::invalid_argument(e.to_string()))?;
170        self.check_access(&principal, Access::Write, &r.namespace)?;
171        let key = RecordKey::new(r.namespace, r.collection, r.id);
172        let outcome = self
173            .service
174            .delete(&key, expected)
175            .await
176            .map_err(internal)?;
177        let resp = match outcome {
178            DeleteResult::Deleted => DeleteResponse {
179                outcome: "deleted".into(),
180                payload_json: Vec::new(),
181            },
182            DeleteResult::Conflict(c) => DeleteResponse {
183                outcome: "conflict".into(),
184                payload_json: serde_json::to_vec(&*c).map_err(internal)?,
185            },
186        };
187        Ok(Response::new(resp))
188    }
189
190    async fn list(&self, req: Request<ListRequest>) -> Result<Response<ListResponse>, Status> {
191        let (metadata, _ext, r) = req.into_parts();
192        // Listing without a namespace spans all namespaces → requires admin
193        // (`read` on `"*"`); a namespaced list needs read on that namespace.
194        self.authorize(
195            &metadata,
196            Access::Read,
197            r.namespace.as_deref().unwrap_or("*"),
198        )?;
199        let prefix = KeyPrefix {
200            namespace: r.namespace,
201            collection: r.collection,
202        };
203        let keys = self.service.list(&prefix).await.map_err(internal)?;
204        let keys_json = keys
205            .iter()
206            .map(serde_json::to_vec)
207            .collect::<std::result::Result<Vec<_>, _>>()
208            .map_err(internal)?;
209        Ok(Response::new(ListResponse { keys_json }))
210    }
211
212    async fn ticket_sync(
213        &self,
214        req: Request<TicketSyncRequest>,
215    ) -> Result<Response<TicketSyncResponse>, Status> {
216        let (metadata, _ext, r) = req.into_parts();
217        // Ticket sync writes records in the `tickets` namespace.
218        self.authorize(&metadata, Access::Write, "tickets")?;
219        let conn: gonzalo_ticket_config::Connection = serde_json::from_slice(&r.connection_json)
220            .map_err(|e| Status::invalid_argument(e.to_string()))?;
221        let summary = self
222            .service
223            .ticket_sync(&conn, "gonzalod")
224            .await
225            .map_err(|e| match e {
226                // A misconfigured request is the caller's own input → safe to
227                // echo. An internal failure goes through `internal` so its
228                // detail is logged, not leaked (#148).
229                crate::service::TicketSyncError::BadRequest(m) => Status::invalid_argument(m),
230                crate::service::TicketSyncError::Internal(m) => internal(m),
231            })?;
232        Ok(Response::new(TicketSyncResponse {
233            imported: summary.imported as u64,
234            updated: summary.updated as u64,
235            unchanged: summary.unchanged as u64,
236        }))
237    }
238
239    async fn graph_definitions(
240        &self,
241        req: Request<GraphQueryRequest>,
242    ) -> Result<Response<GraphLocatedResponse>, Status> {
243        let (metadata, _ext, r) = req.into_parts();
244        self.authorize(&metadata, Access::Read, &r.repo)?;
245        let items = self
246            .service
247            .graph_definitions(&r.repo, &r.view_id, &r.name)
248            .await
249            .map_err(internal)?;
250        Ok(Response::new(located_response(&items)?))
251    }
252
253    async fn graph_references_to(
254        &self,
255        req: Request<GraphQueryRequest>,
256    ) -> Result<Response<GraphLocatedResponse>, Status> {
257        let (metadata, _ext, r) = req.into_parts();
258        self.authorize(&metadata, Access::Read, &r.repo)?;
259        let items = self
260            .service
261            .graph_references_to(&r.repo, &r.view_id, &r.name)
262            .await
263            .map_err(internal)?;
264        Ok(Response::new(located_response(&items)?))
265    }
266
267    async fn graph_callers_of(
268        &self,
269        req: Request<GraphQueryRequest>,
270    ) -> Result<Response<GraphNamesResponse>, Status> {
271        let (metadata, _ext, r) = req.into_parts();
272        self.authorize(&metadata, Access::Read, &r.repo)?;
273        let names = self
274            .service
275            .graph_callers_of(&r.repo, &r.view_id, &r.name)
276            .await
277            .map_err(internal)?;
278        Ok(Response::new(GraphNamesResponse { names }))
279    }
280
281    async fn graph_callees(
282        &self,
283        req: Request<GraphQueryRequest>,
284    ) -> Result<Response<GraphNamesResponse>, Status> {
285        let (metadata, _ext, r) = req.into_parts();
286        self.authorize(&metadata, Access::Read, &r.repo)?;
287        let names = self
288            .service
289            .graph_callees(&r.repo, &r.view_id, &r.name)
290            .await
291            .map_err(internal)?;
292        Ok(Response::new(GraphNamesResponse { names }))
293    }
294
295    async fn graph_impact(
296        &self,
297        req: Request<GraphQueryRequest>,
298    ) -> Result<Response<GraphNamesResponse>, Status> {
299        let (metadata, _ext, r) = req.into_parts();
300        self.authorize(&metadata, Access::Read, &r.repo)?;
301        let names = self
302            .service
303            .graph_impact_names(&r.repo, &r.view_id, &r.name)
304            .await
305            .map_err(internal)?;
306        Ok(Response::new(GraphNamesResponse { names }))
307    }
308
309    async fn put_blob(
310        &self,
311        req: Request<PutBlobRequest>,
312    ) -> Result<Response<PutBlobResponse>, Status> {
313        let (metadata, _ext, r) = req.into_parts();
314        self.authorize(&metadata, Access::Write, BLOB_NS)?;
315        // Verify the content hashes to the advertised value before writing —
316        // same integrity check as the HTTP hash-addressed PUT.
317        let computed = ContentHash::of(&r.content);
318        if computed.0 != r.hash {
319            return Err(Status::invalid_argument(
320                "blob content does not match the advertised hash",
321            ));
322        }
323        let hash = self.service.put_blob(&r.content).await.map_err(internal)?;
324        Ok(Response::new(PutBlobResponse { hash: hash.0 }))
325    }
326
327    async fn get_blob(
328        &self,
329        req: Request<GetBlobRequest>,
330    ) -> Result<Response<GetBlobResponse>, Status> {
331        let (metadata, _ext, r) = req.into_parts();
332        self.authorize(&metadata, Access::Read, BLOB_NS)?;
333        let found = self
334            .service
335            .get_blob(&ContentHash(r.hash))
336            .await
337            .map_err(internal)?;
338        let resp = match found {
339            Some(content) => GetBlobResponse {
340                found: true,
341                content,
342            },
343            None => GetBlobResponse {
344                found: false,
345                content: Vec::new(),
346            },
347        };
348        Ok(Response::new(resp))
349    }
350
351    async fn list_blobs(
352        &self,
353        req: Request<ListBlobsRequest>,
354    ) -> Result<Response<ListBlobsResponse>, Status> {
355        let (metadata, _ext, _r) = req.into_parts();
356        self.authorize(&metadata, Access::Read, BLOB_NS)?;
357        let hashes = self.service.list_blobs().await.map_err(internal)?;
358        Ok(Response::new(ListBlobsResponse {
359            hashes: hashes.into_iter().map(|h| h.0).collect(),
360        }))
361    }
362
363    async fn delete_blob(
364        &self,
365        req: Request<DeleteBlobRequest>,
366    ) -> Result<Response<DeleteBlobResponse>, Status> {
367        let (metadata, _ext, r) = req.into_parts();
368        self.authorize(&metadata, Access::Write, BLOB_NS)?;
369        self.service
370            .delete_blob(&ContentHash(r.hash))
371            .await
372            .map_err(internal)?;
373        Ok(Response::new(DeleteBlobResponse {}))
374    }
375}
376
377/// JSON-encode each located item into a `GraphLocatedResponse` (the shared
378/// JSON-in-bytes convention).
379// `Status` is large but fixed by tonic's API, so the large-err lint can't be
380// acted on (same as `serve_grpc`).
381#[allow(clippy::result_large_err)]
382fn located_response<T: Serialize>(items: &[T]) -> Result<GraphLocatedResponse, Status> {
383    let items_json = items
384        .iter()
385        .map(serde_json::to_vec)
386        .collect::<std::result::Result<Vec<_>, _>>()
387        .map_err(internal)?;
388    Ok(GraphLocatedResponse { items_json })
389}
390
391/// Serve gRPC on an already-bound listener until the process ends. `auth`
392/// governs per-call namespace authorization (ADR 0015); `Auth::Disabled` serves
393/// open.
394pub async fn serve_grpc(
395    listener: tokio::net::TcpListener,
396    service: Service,
397    auth: Arc<Auth>,
398) -> Result<(), tonic::transport::Error> {
399    let max_blob = service.max_blob_size();
400    let adapter = GrpcAdapter::with_auth(service, auth);
401    let incoming = tokio_stream::wrappers::TcpListenerStream::new(listener);
402    tonic::transport::Server::builder()
403        .add_service(GonzaloServer::new(adapter).max_decoding_message_size(max_blob))
404        .serve_with_incoming(incoming)
405        .await
406}
407
408#[cfg(test)]
409mod tests {
410    use super::*;
411    use gonzalo_core::{BlobStore, Identity, Manifest, Meta, RecordKind, Store};
412    use gonzalo_graph::{Located, Symbol, build_rust};
413    use gonzalo_store_fs::FsStore;
414    use std::collections::BTreeMap;
415    use std::sync::Arc;
416
417    /// Seed one view (`r`/`main`) with two slices and return a gRPC adapter over it.
418    async fn seeded_adapter() -> GrpcAdapter {
419        let fs = Arc::new(FsStore::new(tempfile::tempdir().unwrap().keep()));
420        let mut manifest = Manifest::new();
421        for (path, src) in [
422            ("lib.rs", "fn helper() {}"),
423            ("main.rs", "fn main() { helper(); }"),
424        ] {
425            let hash = fs
426                .put_blob(&build_rust(src).to_slice_bytes())
427                .await
428                .unwrap();
429            manifest.insert(path, hash);
430        }
431        let body = manifest.to_body();
432        let record = Record {
433            revision: Revision::initial(body.bytes()),
434            parent: None,
435            body,
436            kind: RecordKind::GraphManifest,
437            meta: Meta {
438                author: Identity::new("tester"),
439                origin_system: "test".into(),
440                created: 0,
441                updated: 0,
442                labels: BTreeMap::new(),
443            },
444            links: Vec::new(),
445            key: Manifest::key("r", "main"),
446        };
447        let outcome = fs.put(record, None).await.unwrap();
448        assert!(matches!(outcome, PutResult::Committed(_)));
449        GrpcAdapter::new(Service::new(fs.clone(), fs))
450    }
451
452    fn query(name: &str) -> Request<GraphQueryRequest> {
453        Request::new(GraphQueryRequest {
454            repo: "r".into(),
455            view_id: "main".into(),
456            name: name.into(),
457        })
458    }
459
460    #[tokio::test]
461    async fn graph_definitions_returns_located_json() {
462        let adapter = seeded_adapter().await;
463        let resp = adapter
464            .graph_definitions(query("helper"))
465            .await
466            .unwrap()
467            .into_inner();
468        assert_eq!(resp.items_json.len(), 1);
469        let located: Located<Symbol> = serde_json::from_slice(&resp.items_json[0]).unwrap();
470        assert_eq!(located.path, "lib.rs");
471        assert_eq!(located.item.name, "helper");
472    }
473
474    // --- namespace-scoped auth (ADR 0015) ---
475
476    use std::collections::HashMap;
477
478    /// A `writer` principal scoped to the `memory` namespace, plus an `admin`.
479    fn scoped_auth() -> Arc<Auth> {
480        Arc::new(Auth::Enabled(HashMap::from([
481            (
482                "wtok".to_string(),
483                Principal::new("writer", vec!["memory".into()], vec!["memory".into()]),
484            ),
485            ("atok".to_string(), Principal::admin("admin")),
486        ])))
487    }
488
489    fn with_token<T>(msg: T, token: &str) -> Request<T> {
490        let mut req = Request::new(msg);
491        req.metadata_mut()
492            .insert("authorization", format!("Bearer {token}").parse().unwrap());
493        req
494    }
495
496    fn fs_adapter(auth: Arc<Auth>) -> GrpcAdapter {
497        let fs = Arc::new(FsStore::new(tempfile::tempdir().unwrap().keep()));
498        GrpcAdapter::with_auth(Service::new(fs.clone(), fs), auth)
499    }
500
501    fn get_req(namespace: &str) -> GetRequest {
502        GetRequest {
503            namespace: namespace.into(),
504            collection: "col".into(),
505            id: "x".into(),
506        }
507    }
508
509    fn put_req(namespace: &str, author: &str) -> PutRequest {
510        let record = Record {
511            revision: Revision::initial(b"{}"),
512            parent: None,
513            body: gonzalo_core::Body::Inline(b"{}".to_vec()),
514            kind: RecordKind::MemoryTier,
515            meta: Meta {
516                author: Identity::new(author),
517                origin_system: "test".into(),
518                created: 0,
519                updated: 0,
520                labels: BTreeMap::new(),
521            },
522            links: Vec::new(),
523            key: RecordKey::new(namespace, "col", "x"),
524        };
525        PutRequest {
526            record_json: serde_json::to_vec(&record).unwrap(),
527            expected_json: serde_json::to_vec(&Option::<Revision>::None).unwrap(),
528        }
529    }
530
531    /// A well-formed `PutRequest` whose `record_json` is not valid JSON.
532    fn malformed_put_req() -> PutRequest {
533        PutRequest {
534            record_json: b"definitely not a record".to_vec(),
535            expected_json: serde_json::to_vec(&Option::<Revision>::None).unwrap(),
536        }
537    }
538
539    /// A store whose every op fails, to force the `internal` (server-error) path.
540    struct DownStore;
541
542    #[async_trait::async_trait]
543    impl Store for DownStore {
544        async fn get(&self, _key: &RecordKey) -> gonzalo_core::Result<Option<Record>> {
545            Err(gonzalo_core::CoreError::Backend(
546                "/var/lib/gonzalo/graphs/secret.sqlite unreachable".into(),
547            ))
548        }
549        async fn put(
550            &self,
551            _record: Record,
552            _expected: Option<Revision>,
553        ) -> gonzalo_core::Result<PutResult> {
554            Err(gonzalo_core::CoreError::Backend(
555                "s3://secret-bucket".into(),
556            ))
557        }
558        async fn list(
559            &self,
560            _prefix: &gonzalo_core::KeyPrefix,
561        ) -> gonzalo_core::Result<Vec<RecordKey>> {
562            Err(gonzalo_core::CoreError::Backend(
563                "s3://secret-bucket".into(),
564            ))
565        }
566        async fn delete(
567            &self,
568            _key: &RecordKey,
569            _expected: Option<Revision>,
570        ) -> gonzalo_core::Result<DeleteResult> {
571            Err(gonzalo_core::CoreError::Backend(
572                "s3://secret-bucket".into(),
573            ))
574        }
575    }
576
577    #[tokio::test]
578    async fn put_authenticates_before_deserializing() {
579        // A malformed body with NO token is rejected at authentication, before
580        // serde ever runs on the attacker-controlled JSON (#146).
581        let adapter = fs_adapter(scoped_auth());
582        let err = adapter
583            .put(Request::new(malformed_put_req()))
584            .await
585            .unwrap_err();
586        assert_eq!(err.code(), tonic::Code::Unauthenticated);
587    }
588
589    #[tokio::test]
590    async fn put_malformed_body_is_invalid_argument_not_internal() {
591        // An authorized caller sending a malformed body gets InvalidArgument —
592        // the caller's own bad input — not Internal (#146).
593        let adapter = fs_adapter(scoped_auth());
594        let err = adapter
595            .put(with_token(malformed_put_req(), "wtok"))
596            .await
597            .unwrap_err();
598        assert_eq!(err.code(), tonic::Code::InvalidArgument);
599    }
600
601    #[tokio::test]
602    async fn backend_error_is_opaque() {
603        // A forced backend failure yields an opaque Internal status; the leaky
604        // path/bucket detail never reaches the client (#148).
605        let fs = Arc::new(FsStore::new(tempfile::tempdir().unwrap().keep()));
606        let adapter = GrpcAdapter::new(Service::new(Arc::new(DownStore), fs));
607        let err = adapter.get(Request::new(get_req("any"))).await.unwrap_err();
608        assert_eq!(err.code(), tonic::Code::Internal);
609        assert_eq!(err.message(), "internal error");
610        assert!(!err.message().contains("secret"));
611    }
612
613    #[tokio::test]
614    async fn missing_token_is_unauthenticated() {
615        let adapter = fs_adapter(scoped_auth());
616        let err = adapter
617            .get(Request::new(get_req("memory")))
618            .await
619            .unwrap_err();
620        assert_eq!(err.code(), tonic::Code::Unauthenticated);
621    }
622
623    #[tokio::test]
624    async fn read_is_allowed_in_scope_denied_out_of_scope() {
625        let adapter = fs_adapter(scoped_auth());
626        // In-scope read succeeds (record absent → found:false, but authorized).
627        assert!(
628            adapter
629                .get(with_token(get_req("memory"), "wtok"))
630                .await
631                .is_ok()
632        );
633        // Out-of-scope read is denied.
634        let err = adapter
635            .get(with_token(get_req("secrets"), "wtok"))
636            .await
637            .unwrap_err();
638        assert_eq!(err.code(), tonic::Code::PermissionDenied);
639    }
640
641    #[tokio::test]
642    async fn write_is_scoped_and_author_is_stamped() {
643        let adapter = fs_adapter(scoped_auth());
644        // Write outside scope is denied.
645        let err = adapter
646            .put(with_token(put_req("secrets", "writer"), "wtok"))
647            .await
648            .unwrap_err();
649        assert_eq!(err.code(), tonic::Code::PermissionDenied);
650
651        // In-scope write commits — even though the client claimed author
652        // "forged", the daemon stamps the authenticated principal.
653        adapter
654            .put(with_token(put_req("memory", "forged"), "wtok"))
655            .await
656            .unwrap();
657        let resp = adapter
658            .get(with_token(get_req("memory"), "wtok"))
659            .await
660            .unwrap()
661            .into_inner();
662        let record: Record = serde_json::from_slice(&resp.record_json).unwrap();
663        assert_eq!(record.meta.author, Identity::new("writer"));
664    }
665
666    #[tokio::test]
667    async fn list_without_namespace_requires_admin() {
668        let adapter = fs_adapter(scoped_auth());
669        let scoped = adapter
670            .list(with_token(ListRequest::default(), "wtok"))
671            .await
672            .unwrap_err();
673        assert_eq!(scoped.code(), tonic::Code::PermissionDenied);
674        // Admin (wildcard) may list across all namespaces.
675        assert!(
676            adapter
677                .list(with_token(ListRequest::default(), "atok"))
678                .await
679                .is_ok()
680        );
681    }
682
683    // --- blob RPCs (#184) ---
684
685    #[tokio::test]
686    async fn grpc_blob_roundtrip_open() {
687        let fs = Arc::new(FsStore::new(tempfile::tempdir().unwrap().keep()));
688        let adapter = GrpcAdapter::new(Service::new(fs.clone(), fs));
689        let content = b"grpc blob body".to_vec();
690        let hash = gonzalo_core::ContentHash::of(&content).0;
691
692        let put = adapter
693            .put_blob(Request::new(PutBlobRequest {
694                hash: hash.clone(),
695                content: content.clone(),
696            }))
697            .await
698            .unwrap()
699            .into_inner();
700        assert_eq!(put.hash, hash);
701
702        let got = adapter
703            .get_blob(Request::new(GetBlobRequest { hash: hash.clone() }))
704            .await
705            .unwrap()
706            .into_inner();
707        assert!(got.found);
708        assert_eq!(got.content, content);
709
710        let listed = adapter
711            .list_blobs(Request::new(ListBlobsRequest {}))
712            .await
713            .unwrap()
714            .into_inner();
715        assert_eq!(listed.hashes, vec![hash.clone()]);
716
717        adapter
718            .delete_blob(Request::new(DeleteBlobRequest { hash: hash.clone() }))
719            .await
720            .unwrap();
721        let gone = adapter
722            .get_blob(Request::new(GetBlobRequest { hash }))
723            .await
724            .unwrap()
725            .into_inner();
726        assert!(!gone.found);
727    }
728
729    #[tokio::test]
730    async fn grpc_put_blob_hash_mismatch_is_invalid_argument() {
731        let fs = Arc::new(FsStore::new(tempfile::tempdir().unwrap().keep()));
732        let adapter = GrpcAdapter::new(Service::new(fs.clone(), fs));
733        let err = adapter
734            .put_blob(Request::new(PutBlobRequest {
735                hash: gonzalo_core::ContentHash::of(b"not the body").0,
736                content: b"the body".to_vec(),
737            }))
738            .await
739            .unwrap_err();
740        assert_eq!(err.code(), tonic::Code::InvalidArgument);
741    }
742
743    #[tokio::test]
744    async fn grpc_blob_ops_require_blobs_scope() {
745        // `scoped_auth()` grants `memory` only, not `_blobs`.
746        let adapter = fs_adapter(scoped_auth());
747        let content = b"scoped".to_vec();
748        let hash = gonzalo_core::ContentHash::of(&content).0;
749
750        let denied = adapter
751            .put_blob(with_token(
752                PutBlobRequest {
753                    hash: hash.clone(),
754                    content: content.clone(),
755                },
756                "wtok",
757            ))
758            .await
759            .unwrap_err();
760        assert_eq!(denied.code(), tonic::Code::PermissionDenied);
761
762        // Admin token succeeds.
763        adapter
764            .put_blob(with_token(PutBlobRequest { hash, content }, "atok"))
765            .await
766            .unwrap();
767    }
768
769    #[tokio::test]
770    async fn graph_name_queries_return_names() {
771        let adapter = seeded_adapter().await;
772        let impact = adapter
773            .graph_impact(query("helper"))
774            .await
775            .unwrap()
776            .into_inner();
777        assert_eq!(impact.names, vec!["main".to_string()]);
778        let callees = adapter
779            .graph_callees(query("main"))
780            .await
781            .unwrap()
782            .into_inner();
783        assert_eq!(callees.names, vec!["helper".to_string()]);
784        let callers = adapter
785            .graph_callers_of(query("helper"))
786            .await
787            .unwrap()
788            .into_inner();
789        assert_eq!(callers.names, vec!["main".to_string()]);
790    }
791}