Skip to main content

cf_system_sdk_directory/grpc/
client.rs

1//! gRPC client implementation of Directory API
2//!
3//! This client allows remote gears to discover and resolve services via gRPC.
4
5use anyhow::Result;
6use async_trait::async_trait;
7use tonic::service::interceptor::InterceptedService;
8use tonic::transport::Channel;
9
10use crate::ProtoInstanceState;
11use crate::api::{
12    DirectoryClient, DirectoryInvalidArgument, DirectoryNotFound, DirectoryPermissionDenied,
13    InstanceState, LabelSelector, RegisterInstanceInfo, ServiceEndpoint, ServiceInstanceInfo,
14};
15use std::collections::BTreeMap;
16use toolkit_transport_grpc::InternalAuthInterceptor;
17use toolkit_transport_grpc::client::{GrpcClientConfig, connect_lazy, connect_with_retry};
18
19use crate::{
20    DeregisterInstanceRequest, DirectoryServiceClient, GetOpenApiSpecRequest, GrpcServiceEndpoint,
21    HeartbeatRequest, InstanceInfo, ListAllInstancesRequest, ListInstancesRequest,
22    RegisterInstanceRequest, ResolveGrpcServiceRequest, ResolveRestServiceRequest,
23};
24
25/// The directory channel wrapped with the platform-plane
26/// [`InternalAuthInterceptor`], which attaches the gear's internal token
27/// (`x-toolkit-internal-token`) to every outbound system call. A
28/// [`disabled`](InternalAuthInterceptor::disabled) interceptor attaches
29/// nothing (Profile 1 / no platform-plane credential).
30type AuthedChannel = InterceptedService<Channel, InternalAuthInterceptor>;
31
32/// Map a lookup RPC's `tonic::Status` onto the directory's typed sentinels.
33///
34/// The status code is the only thing distinguishing "this name is not
35/// registered" from "the directory is unreachable", and stringifying the status
36/// throws it away. Callers downcast the result: `DirectoryEndpointResolver`
37/// turns `DirectoryNotFound` into `Ok(None)` (a provider that has not come up
38/// yet — routine during startup) and anything else into a real error.
39fn lookup_error(resource: &str, status: &tonic::Status) -> anyhow::Error {
40    match status.code() {
41        tonic::Code::NotFound => DirectoryNotFound::new(resource.to_owned()).into(),
42        tonic::Code::InvalidArgument => {
43            DirectoryInvalidArgument::new(status.message().to_owned()).into()
44        }
45        code => anyhow::anyhow!(
46            "directory lookup for {resource} failed: gRPC {code:?}: {}",
47            status.message()
48        ),
49    }
50}
51
52/// Map a mutating RPC's `tonic::Status`, preserving the code in the message.
53///
54/// A bare `"gRPC call failed"` hides whether the directory was unreachable
55/// (transient) or rejected the request permanently. The two permanent classes
56/// are typed so a caller retrying a mutation (e.g. the presence loop) can stop
57/// rather than spin forever:
58/// - `InvalidArgument` → [`DirectoryInvalidArgument`] (malformed request);
59/// - `PermissionDenied` → [`DirectoryPermissionDenied`] (a permanent decision:
60///   peer not authorized, namespace / trust domain not allowlisted, or a gRPC
61///   service name *pinned* to another gear by the ownership map).
62///
63/// Every other code stays an opaque, transient-by-default error (code kept in
64/// the message) and is retried:
65/// - `Unauthenticated` — a stale or rotated token is valid again on the next
66///   attempt, so a retry recovers what the permanent sentinel would strand;
67/// - `FailedPrecondition` — a *recoverable* gRPC service-name conflict (the name
68///   is merely currently advertised by another gear and frees up when it
69///   deregisters), retried rather than mislabeled as a permanent denial;
70/// - the transport-level transient codes (`Unavailable`, …).
71fn call_error(op: &str, status: &tonic::Status) -> anyhow::Error {
72    match status.code() {
73        tonic::Code::InvalidArgument => {
74            DirectoryInvalidArgument::new(status.message().to_owned()).into()
75        }
76        tonic::Code::PermissionDenied => {
77            DirectoryPermissionDenied::new(status.message().to_owned()).into()
78        }
79        code => anyhow::anyhow!("directory {op} failed: gRPC {code:?}: {}", status.message()),
80    }
81}
82
83/// gRPC client for Directory API
84///
85/// This client connects to a remote `DirectoryService` via gRPC and provides
86/// typed access to service discovery functionality. It includes:
87/// - Configurable timeouts and retries via transport stack
88/// - Automatic proto ↔ domain type conversions
89/// - Distributed tracing and metrics
90/// - Platform-plane credential attachment via an [`InternalAuthInterceptor`]
91///   (defaults to attaching nothing; supply one via the `*_with_interceptor`
92///   constructors for Profile-3 / shared-secret deployments)
93pub struct DirectoryGrpcClient {
94    inner: DirectoryServiceClient<AuthedChannel>,
95}
96
97impl DirectoryGrpcClient {
98    /// Connect to a directory service using default configuration with retries.
99    ///
100    /// Uses exponential backoff retry logic for reliable connection establishment.
101    /// This is the recommended method for `OoP` gears connecting to the master host.
102    ///
103    /// # Errors
104    /// It will return an error when it fails
105    pub async fn connect(uri: impl Into<String>) -> Result<Self> {
106        let cfg = GrpcClientConfig::new("directory");
107        Self::connect_with_retry(uri, &cfg).await
108    }
109
110    /// Connect with default configuration + retries, attaching `interceptor`'s
111    /// platform-plane credential to every outbound call.
112    ///
113    /// This is the Profile-3 / shared-secret entry point: the interceptor is
114    /// typically built from a
115    /// [`ServiceAccountTokenReader`](toolkit_transport_grpc::ServiceAccountTokenReader)
116    /// (rotating SA token) or
117    /// [`InternalAuthInterceptor::from_token`] (static shared secret).
118    ///
119    /// # Errors
120    /// It will return an error when it fails
121    pub async fn connect_with_interceptor(
122        uri: impl Into<String>,
123        interceptor: InternalAuthInterceptor,
124    ) -> Result<Self> {
125        let cfg = GrpcClientConfig::new("directory");
126        let channel: Channel = connect_with_retry(uri, &cfg).await?;
127        Ok(Self::from_channel_with_interceptor(channel, interceptor))
128    }
129
130    /// Create a directory client with a **lazily-connecting** channel.
131    ///
132    /// Performs **no** eager connection: the channel connects on the first RPC
133    /// and transparently reconnects on failure. This is the eventual-readiness
134    /// entry point (`cpt-cf-adr-eventual-readiness`) for `OoP` bootstrap — the
135    /// process starts even when the `DirectoryService` is not yet reachable, and
136    /// the presence loop's backoff retry absorbs the startup window instead of
137    /// the process crashing (which would offload retries onto a k8s
138    /// `CrashLoopBackOff`).
139    ///
140    /// # Runtime context
141    /// Must be called from within a Tokio runtime context: building the lazy
142    /// channel initialises the hyper reactor. Calling it outside a runtime
143    /// returns an error rather than panicking (it still does not connect).
144    ///
145    /// # Errors
146    /// Returns an error if called outside a Tokio runtime context, or if `uri`
147    /// is malformed — never for an unreachable peer.
148    pub fn connect_lazy(uri: impl Into<String>) -> Result<Self> {
149        let cfg = GrpcClientConfig::new("directory");
150        let channel: Channel = connect_lazy(uri, &cfg)?;
151        Ok(Self::from_channel(channel))
152    }
153
154    /// Create a directory client with a **lazily-connecting** channel, attaching
155    /// `interceptor`'s platform-plane credential to every outbound call.
156    ///
157    /// The lazy counterpart of [`connect_with_interceptor`](Self::connect_with_interceptor);
158    /// see [`connect_lazy`](Self::connect_lazy) for the connection semantics.
159    /// The URI is validated before `interceptor` is consumed, so the credential
160    /// is only moved into the client on success.
161    ///
162    /// # Errors
163    /// Returns an error only if `uri` is malformed — never for an unreachable
164    /// peer.
165    pub fn connect_lazy_with_interceptor(
166        uri: impl Into<String>,
167        interceptor: InternalAuthInterceptor,
168    ) -> Result<Self> {
169        let cfg = GrpcClientConfig::new("directory");
170        // Validate the URI (build the channel) before consuming `interceptor`.
171        let channel: Channel = connect_lazy(uri, &cfg)?;
172        Ok(Self::from_channel_with_interceptor(channel, interceptor))
173    }
174
175    /// Connect to a directory service with custom configuration and retry logic.
176    ///
177    /// Uses exponential backoff based on `cfg.max_retries`, `cfg.base_backoff`,
178    /// and `cfg.max_backoff` settings.
179    ///
180    /// # Errors
181    /// It will return an error when it fails
182    pub async fn connect_with_retry(
183        uri: impl Into<String>,
184        cfg: &GrpcClientConfig,
185    ) -> Result<Self> {
186        let channel: Channel = connect_with_retry(uri, cfg).await?;
187        Ok(Self::from_channel(channel))
188    }
189
190    /// Connect to a directory service without retry logic.
191    ///
192    /// This method attempts a single connection. Use `connect` or `connect_with_retry`
193    /// for production scenarios where the directory service may not be immediately available.
194    ///
195    /// # Errors
196    /// It will return an error when it fails
197    pub async fn connect_no_retry(uri: impl Into<String>, cfg: &GrpcClientConfig) -> Result<Self> {
198        let uri_string = uri.into();
199
200        // Create endpoint with timeouts from config
201        let endpoint = tonic::transport::Endpoint::from_shared(uri_string)?
202            .connect_timeout(cfg.connect_timeout)
203            .timeout(cfg.rpc_timeout);
204
205        // Connect to the service
206        let channel = endpoint.connect().await?;
207
208        if cfg.enable_tracing {
209            tracing::debug!(
210                service_name = cfg.service_name,
211                connect_timeout_ms = cfg.connect_timeout.as_millis(),
212                rpc_timeout_ms = cfg.rpc_timeout.as_millis(),
213                "directory gRPC client connected"
214            );
215        }
216
217        Ok(Self::from_channel(channel))
218    }
219
220    /// Create from an existing channel (useful for testing or custom setup).
221    ///
222    /// Attaches no platform-plane credential; use
223    /// [`from_channel_with_interceptor`](Self::from_channel_with_interceptor)
224    /// to attach one.
225    #[must_use]
226    pub fn from_channel(channel: Channel) -> Self {
227        Self::from_channel_with_interceptor(channel, InternalAuthInterceptor::disabled())
228    }
229
230    /// Create from an existing channel, attaching `interceptor`'s platform-plane
231    /// credential to every outbound call.
232    #[must_use]
233    pub fn from_channel_with_interceptor(
234        channel: Channel,
235        interceptor: InternalAuthInterceptor,
236    ) -> Self {
237        Self {
238            inner: DirectoryServiceClient::with_interceptor(channel, interceptor),
239        }
240    }
241}
242
243#[async_trait]
244impl DirectoryClient for DirectoryGrpcClient {
245    async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint> {
246        let mut client = self.inner.clone();
247        let request = tonic::Request::new(ResolveGrpcServiceRequest {
248            service_name: service_name.to_owned(),
249        });
250
251        let response = client
252            .resolve_grpc_service(request)
253            .await
254            .map_err(|e| lookup_error(&format!("service {service_name}"), &e))?;
255
256        let proto_response = response.into_inner();
257        Ok(ServiceEndpoint::new(proto_response.endpoint_uri))
258    }
259
260    async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint> {
261        let mut client = self.inner.clone();
262        let request = tonic::Request::new(ResolveRestServiceRequest {
263            gear_name: gear_name.to_owned(),
264        });
265
266        let response = client
267            .resolve_rest_service(request)
268            .await
269            .map_err(|e| lookup_error(&format!("gear {gear_name}"), &e))?;
270
271        let proto_response = response.into_inner();
272        Ok(ServiceEndpoint::new(proto_response.endpoint_uri))
273    }
274
275    async fn get_openapi_spec(&self, gear_name: &str) -> Result<String> {
276        let mut client = self.inner.clone();
277        let request = tonic::Request::new(GetOpenApiSpecRequest {
278            gear_name: gear_name.to_owned(),
279        });
280
281        let response = client
282            .get_open_api_spec(request)
283            .await
284            .map_err(|e| lookup_error(&format!("openapi spec for gear {gear_name}"), &e))?;
285
286        Ok(response.into_inner().openapi_spec)
287    }
288
289    async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>> {
290        let mut client = self.inner.clone();
291        let request = tonic::Request::new(ListInstancesRequest {
292            gear_name: gear.to_owned(),
293            match_labels: std::collections::HashMap::new(),
294        });
295
296        let response = client
297            .list_instances(request)
298            .await
299            .map_err(|e| lookup_error(&format!("instances of gear {gear}"), &e))?;
300
301        let instances = response
302            .into_inner()
303            .instances
304            .into_iter()
305            .map(proto_instance_to_domain)
306            .collect();
307
308        Ok(instances)
309    }
310
311    async fn resolve_by_labels(
312        &self,
313        gear: &str,
314        selector: &LabelSelector,
315    ) -> Result<Vec<ServiceInstanceInfo>> {
316        // Push the selector server-side so the directory returns only matching
317        // instances. Every `list_instances` response is spec-free (only the
318        // `openapi_spec_hash` rides along; the document is fetched via
319        // `GetOpenApiSpec`), so an empty (match-all) selector never leaks a full
320        // OpenAPI document over the wire.
321        //
322        // cancel-safe: the single await is the unary `list_instances` RPC, which
323        // precedes any local state change; cancelling it just drops the in-flight
324        // response and leaves nothing partially applied.
325        let mut client = self.inner.clone();
326        let match_labels = selector
327            .match_labels
328            .iter()
329            .map(|(k, v)| (k.clone(), v.clone()))
330            .collect();
331        let request = tonic::Request::new(ListInstancesRequest {
332            gear_name: gear.to_owned(),
333            match_labels,
334        });
335
336        let response = client
337            .list_instances(request)
338            .await
339            .map_err(|e| lookup_error(&format!("instances of gear {gear}"), &e))?;
340
341        let instances = response
342            .into_inner()
343            .instances
344            .into_iter()
345            .map(proto_instance_to_domain)
346            .filter(|i| selector.matches(&i.labels))
347            .collect();
348
349        Ok(instances)
350    }
351
352    async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>> {
353        let mut client = self.inner.clone();
354        let response = client
355            .list_all_instances(tonic::Request::new(ListAllInstancesRequest {}))
356            .await
357            .map_err(|e| lookup_error("all instances", &e))?;
358
359        let instances = response
360            .into_inner()
361            .instances
362            .into_iter()
363            .map(|proto| proto_instance_to_domain(proto).without_labels())
364            .collect();
365
366        Ok(instances)
367    }
368
369    async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()> {
370        let mut client = self.inner.clone();
371
372        // Convert gRPC service endpoints
373        let grpc_services = info
374            .grpc_services
375            .into_iter()
376            .map(|(name, ep)| GrpcServiceEndpoint {
377                service_name: name,
378                endpoint_uri: ep.uri,
379            })
380            .collect();
381
382        let req = RegisterInstanceRequest {
383            gear_name: info.gear,
384            instance_id: info.instance_id,
385            grpc_services,
386            version: info.version.unwrap_or_default(),
387            rest_endpoint_uri: info.rest_endpoint.map(|ep| ep.uri),
388            openapi_spec: info.openapi_spec,
389            labels: info.labels.into_iter().collect(),
390        };
391
392        client
393            .register_instance(tonic::Request::new(req))
394            .await
395            .map_err(|e| call_error("register_instance", &e))?;
396
397        Ok(())
398    }
399
400    async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()> {
401        let mut client = self.inner.clone();
402
403        let req = DeregisterInstanceRequest {
404            gear_name: gear.to_owned(),
405            instance_id: instance_id.to_owned(),
406        };
407
408        client
409            .deregister_instance(tonic::Request::new(req))
410            .await
411            .map_err(|e| call_error("deregister_instance", &e))?;
412
413        Ok(())
414    }
415
416    async fn send_heartbeat(&self, gear: &str, instance_id: &str) -> Result<()> {
417        let mut client = self.inner.clone();
418
419        let req = HeartbeatRequest {
420            gear_name: gear.to_owned(),
421            instance_id: instance_id.to_owned(),
422        };
423
424        client
425            .heartbeat(tonic::Request::new(req))
426            .await
427            .map_err(|e| call_error("heartbeat", &e))?;
428
429        Ok(())
430    }
431}
432
433/// Convert a proto `InstanceInfo` into the domain [`ServiceInstanceInfo`].
434fn proto_instance_to_domain(proto: InstanceInfo) -> ServiceInstanceInfo {
435    ServiceInstanceInfo {
436        gear: proto.gear_name,
437        instance_id: proto.instance_id,
438        endpoint: if proto.endpoint_uri.is_empty() {
439            None
440        } else {
441            Some(ServiceEndpoint::new(proto.endpoint_uri))
442        },
443        version: if proto.version.is_empty() {
444            None
445        } else {
446            Some(proto.version)
447        },
448        rest_endpoint: proto.rest_endpoint_uri.map(ServiceEndpoint::new),
449        openapi_spec_hash: proto.openapi_spec_hash,
450        // The `InstanceInfo` proto message carries no per-service gRPC
451        // breakdown, so nothing to reconstruct over the OoP directory transport;
452        // only the in-process `LocalDirectoryClient` populates this (from the
453        // live `GearInstance`). No consumer reads `grpc_services` off a
454        // gRPC-obtained instance today — a single gRPC endpoint is available via
455        // the primary `endpoint`. When a label-targeted gRPC client needs the
456        // full per-service map remotely (the TopologyView work), add a
457        // `repeated GrpcServiceEndpoint grpc_services` field to `InstanceInfo`.
458        grpc_services: Vec::new(),
459        // Stable addressing labels cross the wire (ordering is normalized into a
460        // BTreeMap for deterministic selector matching).
461        labels: proto.labels.into_iter().collect::<BTreeMap<_, _>>(),
462        // Live serving state so label-targeted callers can filter on health.
463        state: proto_state_to_domain(proto.state),
464    }
465}
466
467/// Map the proto `InstanceState` (an open enum carried as `i32`) onto the
468/// domain [`InstanceState`].
469///
470/// `UNSPECIFIED` (a peer that never set the field, e.g. an older server) and
471/// any discriminant this build does not recognise map to the non-serving
472/// [`InstanceState::Unknown`] — kept distinct from
473/// [`InstanceState::Registered`] so "the state is unknown" is not silently read
474/// as a known pre-serving baseline. An unrecognised discriminant is logged so
475/// the version skew is observable rather than swallowed.
476fn proto_state_to_domain(state: i32) -> InstanceState {
477    match ProtoInstanceState::try_from(state) {
478        Ok(ProtoInstanceState::Ready) => InstanceState::Ready,
479        Ok(ProtoInstanceState::Healthy) => InstanceState::Healthy,
480        Ok(ProtoInstanceState::Quarantined) => InstanceState::Quarantined,
481        Ok(ProtoInstanceState::Draining) => InstanceState::Draining,
482        Ok(ProtoInstanceState::Registered) => InstanceState::Registered,
483        Ok(ProtoInstanceState::Unspecified) => InstanceState::Unknown,
484        Err(_) => {
485            tracing::warn!(
486                raw_state = state,
487                "directory returned an unrecognised InstanceState discriminant; \
488                 treating as Unknown (non-serving)"
489            );
490            InstanceState::Unknown
491        }
492    }
493}
494
495#[cfg(test)]
496#[cfg_attr(coverage_nightly, coverage(off))]
497mod tests {
498    use super::*;
499
500    #[tokio::test]
501    async fn test_grpc_client_can_be_constructed() {
502        // Smoke test to ensure types compile and connect
503        let endpoint = tonic::transport::Endpoint::from_static("http://[::1]:50051");
504
505        // We can't actually connect without a server, but we can construct the client type
506        // This ensures the API is correct
507        let channel_result = endpoint.connect().await;
508
509        // It's expected to fail since there's no server, but if it does somehow succeed:
510        if let Ok(channel) = channel_result {
511            let _client = DirectoryGrpcClient::from_channel(channel);
512        }
513    }
514
515    #[tokio::test]
516    async fn from_channel_constructs_without_connecting() {
517        // `connect_lazy` yields a Channel without a live server, so both the
518        // default (no-credential) and interceptor-bearing constructors can be
519        // exercised offline.
520        let channel = Channel::from_static("http://[::1]:50051").connect_lazy();
521        let _default = DirectoryGrpcClient::from_channel(channel.clone());
522        let _authed = DirectoryGrpcClient::from_channel_with_interceptor(
523            channel,
524            InternalAuthInterceptor::disabled(),
525        );
526    }
527
528    #[tokio::test]
529    async fn connect_lazy_succeeds_against_unreachable_peer() {
530        // The lazy constructor performs no eager connect, so an OoP gear can
531        // build its directory client before the `DirectoryService` is up
532        // (`cpt-cf-adr-eventual-readiness`). Nothing is listening on port 1, yet
533        // both the plain and interceptor-bearing constructors return `Ok`.
534        let plain = DirectoryGrpcClient::connect_lazy("http://127.0.0.1:1");
535        assert!(
536            plain.is_ok(),
537            "connect_lazy must not eagerly connect (unreachable peer -> Ok)"
538        );
539
540        let authed = DirectoryGrpcClient::connect_lazy_with_interceptor(
541            "http://127.0.0.1:1",
542            InternalAuthInterceptor::disabled(),
543        );
544        assert!(
545            authed.is_ok(),
546            "connect_lazy_with_interceptor must not eagerly connect (unreachable peer -> Ok)"
547        );
548    }
549
550    #[tokio::test]
551    async fn connect_lazy_rejects_malformed_uri() {
552        // A malformed endpoint is a static misconfiguration worth failing fast
553        // on — the only error path of the lazy constructors.
554        assert!(
555            DirectoryGrpcClient::connect_lazy(String::new()).is_err(),
556            "connect_lazy should fail on a malformed URI"
557        );
558        assert!(
559            DirectoryGrpcClient::connect_lazy_with_interceptor(
560                String::new(),
561                InternalAuthInterceptor::disabled(),
562            )
563            .is_err(),
564            "connect_lazy_with_interceptor should fail on a malformed URI"
565        );
566    }
567
568    #[tokio::test]
569    async fn resolve_grpc_service_through_lazy_client_errors_not_hangs() {
570        // A lazy client builds against an unreachable directory; the first RPC
571        // returns a lookup/call error rather than hanging (outer timeout proves
572        // non-hang; nothing is listening on port 1).
573        let client =
574            DirectoryGrpcClient::connect_lazy("http://127.0.0.1:1").expect("lazy build ok");
575        let outcome = tokio::time::timeout(
576            std::time::Duration::from_secs(5),
577            client.resolve_grpc_service("cf.directory.v1.DirectoryService"),
578        )
579        .await;
580        assert!(
581            outcome.is_ok(),
582            "resolve_grpc_service through a lazy client must not hang against an unreachable peer"
583        );
584        assert!(
585            outcome.unwrap().is_err(),
586            "resolve_grpc_service against an unreachable directory must return Err"
587        );
588    }
589
590    #[test]
591    fn proto_instance_maps_all_fields_to_domain() {
592        let proto = InstanceInfo {
593            gear_name: "calc".to_owned(),
594            instance_id: "calc-1".to_owned(),
595            endpoint_uri: "http://calc:8080".to_owned(),
596            version: "1.2.3".to_owned(),
597            rest_endpoint_uri: Some("http://calc:8080".to_owned()),
598            openapi_spec_hash: Some("1a2b3c4d5e6f7a8b".to_owned()),
599            labels: [("shard".to_owned(), "7".to_owned())].into_iter().collect(),
600            state: ProtoInstanceState::Healthy as i32,
601        };
602        let domain = proto_instance_to_domain(proto);
603        assert_eq!(domain.gear, "calc");
604        assert_eq!(domain.state, InstanceState::Healthy);
605        assert_eq!(domain.instance_id, "calc-1");
606        assert_eq!(
607            domain.endpoint.as_ref().map(|e| e.uri.as_str()),
608            Some("http://calc:8080")
609        );
610        assert_eq!(domain.version.as_deref(), Some("1.2.3"));
611        assert_eq!(
612            domain.rest_endpoint.map(|e| e.uri),
613            Some("http://calc:8080".to_owned())
614        );
615        // Enumeration is spec-free: only the hash crosses the wire.
616        assert_eq!(
617            domain.openapi_spec_hash.as_deref(),
618            Some("1a2b3c4d5e6f7a8b")
619        );
620        // Labels cross the wire and land in a BTreeMap for deterministic matching.
621        assert_eq!(domain.labels.get("shard"), Some(&"7".to_owned()));
622    }
623
624    #[test]
625    fn proto_instance_maps_empty_version_to_none() {
626        let proto = InstanceInfo {
627            gear_name: "worker".to_owned(),
628            instance_id: "worker-1".to_owned(),
629            endpoint_uri: "http://worker:7000".to_owned(),
630            version: String::new(),
631            rest_endpoint_uri: None,
632            openapi_spec_hash: None,
633            labels: std::collections::HashMap::new(),
634            state: ProtoInstanceState::Unspecified as i32,
635        };
636        let domain = proto_instance_to_domain(proto);
637        // A non-empty proto endpoint_uri maps to `Some`.
638        assert_eq!(
639            domain.endpoint.as_ref().map(|e| e.uri.as_str()),
640            Some("http://worker:7000")
641        );
642        // An empty proto version string maps to `None` rather than an empty string.
643        assert!(domain.version.is_none());
644        assert!(domain.rest_endpoint.is_none());
645        assert!(domain.openapi_spec_hash.is_none());
646        assert!(domain.labels.is_empty());
647        // An unset proto state (`UNSPECIFIED`) maps to the non-serving Unknown
648        // sentinel — distinct from the pre-serving Registered baseline.
649        assert_eq!(domain.state, InstanceState::Unknown);
650        assert!(!domain.state.is_serving());
651    }
652
653    #[test]
654    fn proto_instance_maps_empty_endpoint_to_none() {
655        // proto3 carries an absent primary endpoint as an empty string; it must
656        // map back to `None`, not an empty-URI sentinel a dialer could mistake
657        // for a real address.
658        let proto = InstanceInfo {
659            gear_name: "grpc-only".to_owned(),
660            instance_id: "g-1".to_owned(),
661            endpoint_uri: String::new(),
662            version: String::new(),
663            rest_endpoint_uri: None,
664            openapi_spec_hash: None,
665            labels: std::collections::HashMap::new(),
666            state: ProtoInstanceState::Ready as i32,
667        };
668        let domain = proto_instance_to_domain(proto);
669        assert!(
670            domain.endpoint.is_none(),
671            "an empty proto endpoint_uri must map to None"
672        );
673    }
674
675    #[test]
676    fn proto_state_unspecified_and_unrecognised_map_to_unknown() {
677        // `UNSPECIFIED` (peer never set the field) and any discriminant this
678        // build does not know (e.g. a newer server) both collapse to the
679        // non-serving Unknown sentinel rather than being read as Registered.
680        assert_eq!(
681            proto_state_to_domain(ProtoInstanceState::Unspecified as i32),
682            InstanceState::Unknown
683        );
684        assert_eq!(proto_state_to_domain(9999), InstanceState::Unknown);
685        assert!(!proto_state_to_domain(9999).is_serving());
686
687        // Known discriminants still map through unchanged.
688        assert_eq!(
689            proto_state_to_domain(ProtoInstanceState::Registered as i32),
690            InstanceState::Registered
691        );
692        assert_eq!(
693            proto_state_to_domain(ProtoInstanceState::Healthy as i32),
694            InstanceState::Healthy
695        );
696    }
697
698    #[test]
699    fn call_error_types_invalid_argument_as_permanent() {
700        let err = call_error(
701            "register_instance",
702            &tonic::Status::invalid_argument("bad label"),
703        );
704        assert!(
705            err.downcast_ref::<DirectoryInvalidArgument>().is_some(),
706            "InvalidArgument must map to the typed permanent DirectoryInvalidArgument"
707        );
708    }
709
710    #[test]
711    fn call_error_types_permission_denied_as_permanent() {
712        // A true authorization decision (peer not authorized, or namespace /
713        // trust domain not allowlisted) is permanent, so the presence loop stops
714        // retrying.
715        let err = call_error(
716            "register_instance",
717            &tonic::Status::permission_denied("peer not authorized"),
718        );
719        assert!(
720            err.downcast_ref::<DirectoryPermissionDenied>().is_some(),
721            "PermissionDenied must map to the typed permanent DirectoryPermissionDenied"
722        );
723    }
724
725    #[test]
726    fn call_error_keeps_transient_codes_opaque() {
727        // Transient codes stay untyped (code preserved in the message) so the
728        // presence loop keeps retrying:
729        // - `Unauthenticated` — a stale/rotated credential must not collapse to
730        //   the permanent sentinel;
731        // - `FailedPrecondition` — a recoverable gRPC service-name conflict must
732        //   not be mislabeled as a permanent authorization denial.
733        for status in [
734            tonic::Status::unavailable("connection reset"),
735            tonic::Status::unauthenticated("invalid internal token"),
736            tonic::Status::failed_precondition("service name already owned by another gear"),
737        ] {
738            let err = call_error("heartbeat", &status);
739            assert!(err.downcast_ref::<DirectoryInvalidArgument>().is_none());
740            assert!(
741                err.downcast_ref::<DirectoryPermissionDenied>().is_none(),
742                "{:?} must stay transient/retryable, not map to DirectoryPermissionDenied",
743                status.code()
744            );
745            assert!(err.to_string().contains(&format!("{:?}", status.code())));
746        }
747    }
748}