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