Skip to main content

cf_system_sdk_directory/
api.rs

1//! Directory API - contract for service discovery and instance resolution
2//!
3//! This gear defines the core traits and types for the directory service API.
4
5use anyhow::Result;
6use async_trait::async_trait;
7use std::collections::BTreeMap;
8
9/// Represents an endpoint where a service can be reached
10#[derive(Clone, Debug, Default, PartialEq, Eq, Hash)]
11pub struct ServiceEndpoint {
12    pub uri: String,
13}
14
15impl ServiceEndpoint {
16    pub fn new(uri: impl Into<String>) -> Self {
17        Self { uri: uri.into() }
18    }
19
20    #[must_use]
21    pub fn http(host: &str, port: u16) -> Self {
22        Self {
23            uri: format!("{}://{}:{}", "http", host, port),
24        }
25    }
26
27    #[must_use]
28    pub fn https(host: &str, port: u16) -> Self {
29        Self {
30            uri: format!("https://{host}:{port}"),
31        }
32    }
33
34    pub fn uds(path: impl AsRef<std::path::Path>) -> Self {
35        Self {
36            uri: format!("unix://{}", path.as_ref().display()),
37        }
38    }
39}
40
41/// Label-equality selector over instance labels (k8s `matchLabels` style).
42///
43/// Matching is equality-AND: an instance matches iff it carries **every**
44/// `(key, value)` pair in `match_labels`. An empty selector matches every
45/// instance. This is the selector consumed by
46/// [`DirectoryClient::resolve_by_labels`].
47#[derive(Clone, Debug, Default, PartialEq, Eq)]
48pub struct LabelSelector {
49    /// Required label equalities; all must match.
50    pub match_labels: BTreeMap<String, String>,
51}
52
53impl LabelSelector {
54    /// An empty selector (matches every instance).
55    #[must_use]
56    pub fn new() -> Self {
57        Self::default()
58    }
59
60    /// Build a selector from an existing `matchLabels` map.
61    #[must_use]
62    pub fn from_match_labels(match_labels: BTreeMap<String, String>) -> Self {
63        Self { match_labels }
64    }
65
66    /// Add one required label equality, returning `self` for chaining.
67    #[must_use]
68    pub fn with(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
69        self.match_labels.insert(key.into(), value.into());
70        self
71    }
72
73    /// Whether this selector constrains nothing (matches every instance).
74    #[must_use]
75    pub fn is_empty(&self) -> bool {
76        self.match_labels.is_empty()
77    }
78
79    /// Whether `labels` satisfies every required equality in this selector.
80    #[must_use]
81    pub fn matches(&self, labels: &BTreeMap<String, String>) -> bool {
82        self.match_labels
83            .iter()
84            .all(|(k, v)| labels.get(k) == Some(v))
85    }
86}
87
88/// Serving state of a directory instance, projected from the directory's live
89/// registration/heartbeat tracking.
90///
91/// Only [`Ready`](Self::Ready) / [`Healthy`](Self::Healthy) are **serving**;
92/// [`Draining`](Self::Draining) is a graceful shutdown (callers should avoid
93/// sending new work) and [`Quarantined`](Self::Quarantined) is stale (missed
94/// heartbeats). [`Registered`](Self::Registered) is the pre-serving baseline.
95/// Exposed on [`ServiceInstanceInfo`] so [`DirectoryClient::resolve_by_labels`]
96/// callers can apply their own health policy from a single resolve.
97#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
98pub enum InstanceState {
99    /// Registered but not yet serving (no heartbeat / readiness signal).
100    #[default]
101    Registered,
102    /// Dependencies resolved and marked ready.
103    Ready,
104    /// Serving and heartbeating.
105    Healthy,
106    /// Missed its heartbeat deadline; not serving, pending eviction.
107    Quarantined,
108    /// Gracefully shutting down; upstreams should stop routing new work.
109    Draining,
110    /// State could not be determined.
111    Unknown,
112}
113
114impl InstanceState {
115    /// Whether an instance in this state should receive new traffic
116    /// (`Ready` / `Healthy`). All other states are non-serving.
117    #[must_use]
118    pub fn is_serving(self) -> bool {
119        matches!(self, InstanceState::Ready | InstanceState::Healthy)
120    }
121}
122
123/// Information about a service instance
124#[derive(Debug, Clone, Default)]
125#[non_exhaustive]
126pub struct ServiceInstanceInfo {
127    /// Gear name this instance belongs to
128    pub gear: String,
129    /// Unique instance identifier
130    pub instance_id: String,
131    /// Primary endpoint for the instance
132    pub endpoint: Option<ServiceEndpoint>,
133    /// Optional version string
134    pub version: Option<String>,
135    /// REST endpoint (HTTP base URL) for this instance, or `None` if the gear
136    /// exposes no REST API. This is the address the edge reverse-proxy actually
137    /// dials for REST routing.
138    pub rest_endpoint: Option<ServiceEndpoint>,
139    /// Stable content token for the published `OpenAPI` spec, if any.
140    ///
141    /// Enumeration is deliberately spec-free: only this hash rides along, never
142    /// the full document (fetched out-of-band via
143    /// [`DirectoryClient::get_openapi_spec`]). The hash lets a consumer detect a
144    /// spec change without inlining N copies of the document.
145    pub openapi_spec_hash: Option<String>,
146    /// Published gRPC services as `(service name, endpoint)` pairs.
147    ///
148    /// Kept as an ordered list rather than a map so it round-trips a repeated
149    /// proto field without reordering; service names are unique per instance.
150    /// Carried back by `list_instances` so a subsequent `register_instance`
151    /// (which replaces the entry wholesale) can augment — rather than clobber —
152    /// the previously-registered gRPC services when adding a REST endpoint.
153    pub grpc_services: Vec<(String, ServiceEndpoint)>,
154    /// Stable addressing labels published by this instance (k8s `matchLabels`
155    /// style). Used by [`DirectoryClient::resolve_by_labels`] to select shards
156    /// or peers within a directory name.
157    pub labels: BTreeMap<String, String>,
158    /// Live serving state, so a `resolve_by_labels` caller can apply its own
159    /// health policy without a second call. Defaults to
160    /// [`InstanceState::Registered`].
161    pub state: InstanceState,
162}
163
164impl ServiceInstanceInfo {
165    /// Start a projection for `gear` / `instance_id` with no primary endpoint
166    /// and no optional metadata. Chain the `with_*` builders to populate it.
167    /// `#[non_exhaustive]` means callers must use this builder rather than a
168    /// struct literal, so future fields stay additive.
169    #[must_use]
170    pub fn new(gear: impl Into<String>, instance_id: impl Into<String>) -> Self {
171        Self {
172            gear: gear.into(),
173            instance_id: instance_id.into(),
174            ..Self::default()
175        }
176    }
177
178    /// Set the primary endpoint (`None` for an instance with no primary
179    /// endpoint).
180    #[must_use]
181    pub fn with_endpoint(mut self, endpoint: Option<ServiceEndpoint>) -> Self {
182        self.endpoint = endpoint;
183        self
184    }
185
186    /// Set the optional version string.
187    #[must_use]
188    pub fn with_version(mut self, version: Option<String>) -> Self {
189        self.version = version;
190        self
191    }
192
193    /// Set the optional REST endpoint.
194    #[must_use]
195    pub fn with_rest_endpoint(mut self, rest_endpoint: Option<ServiceEndpoint>) -> Self {
196        self.rest_endpoint = rest_endpoint;
197        self
198    }
199
200    /// Set the optional `OpenAPI` spec content token.
201    #[must_use]
202    pub fn with_openapi_spec_hash(mut self, openapi_spec_hash: Option<String>) -> Self {
203        self.openapi_spec_hash = openapi_spec_hash;
204        self
205    }
206
207    /// Set the published gRPC services.
208    #[must_use]
209    pub fn with_grpc_services(mut self, grpc_services: Vec<(String, ServiceEndpoint)>) -> Self {
210        self.grpc_services = grpc_services;
211        self
212    }
213
214    /// Set the stable addressing labels.
215    #[must_use]
216    pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
217        self.labels = labels;
218        self
219    }
220
221    /// Set the live serving state.
222    #[must_use]
223    pub fn with_state(mut self, state: InstanceState) -> Self {
224        self.state = state;
225        self
226    }
227
228    /// Drop the stable addressing labels, returning `self`.
229    ///
230    /// The shared transform every [`DirectoryClient::list_all_instances`]
231    /// implementation applies so the cross-gear snapshot omits labels
232    /// identically across transports — see that method's doc for why.
233    #[must_use]
234    pub fn without_labels(mut self) -> Self {
235        self.labels.clear();
236        self
237    }
238}
239
240/// Information for registering a new gear instance
241#[derive(Debug, Clone, Default)]
242#[non_exhaustive]
243pub struct RegisterInstanceInfo {
244    /// Gear name
245    pub gear: String,
246    /// Unique instance identifier
247    pub instance_id: String,
248    /// Published gRPC services as `(service name, endpoint)` pairs (service
249    /// names are unique per instance; ordered to mirror the proto repeated
250    /// field rather than reordered into a map).
251    pub grpc_services: Vec<(String, ServiceEndpoint)>,
252    /// Optional version string
253    pub version: Option<String>,
254    /// Optional REST endpoint (HTTP base URL) exposed by the gear.
255    pub rest_endpoint: Option<ServiceEndpoint>,
256    /// Optional `OpenAPI` spec (JSON) published by the gear.
257    pub openapi_spec: Option<String>,
258    /// Stable addressing labels for this instance (k8s `matchLabels` style).
259    pub labels: BTreeMap<String, String>,
260}
261
262impl RegisterInstanceInfo {
263    /// Start a registration for `gear` / `instance_id` with no endpoints,
264    /// version, spec, or labels. Chain the `with_*` builders to populate it.
265    #[must_use]
266    pub fn new(gear: impl Into<String>, instance_id: impl Into<String>) -> Self {
267        Self {
268            gear: gear.into(),
269            instance_id: instance_id.into(),
270            ..Self::default()
271        }
272    }
273
274    /// Set the published gRPC services.
275    #[must_use]
276    pub fn with_grpc_services(mut self, grpc_services: Vec<(String, ServiceEndpoint)>) -> Self {
277        self.grpc_services = grpc_services;
278        self
279    }
280
281    /// Set the version string.
282    #[must_use]
283    pub fn with_version(mut self, version: impl Into<String>) -> Self {
284        self.version = Some(version.into());
285        self
286    }
287
288    /// Set the REST endpoint (HTTP base URL).
289    #[must_use]
290    pub fn with_rest_endpoint(mut self, rest_endpoint: ServiceEndpoint) -> Self {
291        self.rest_endpoint = Some(rest_endpoint);
292        self
293    }
294
295    /// Set the published `OpenAPI` spec (JSON).
296    #[must_use]
297    pub fn with_openapi_spec(mut self, openapi_spec: impl Into<String>) -> Self {
298        self.openapi_spec = Some(openapi_spec.into());
299        self
300    }
301
302    /// Set the stable addressing labels.
303    ///
304    /// Re-registration replaces the directory entry wholesale, but an **empty**
305    /// label set is treated as "preserve the stored labels", not "clear them"
306    /// (see `GearManager::with_metadata_of`): the periodic self-heal and REST
307    /// augmentation re-register without labels and must not erase a shard's
308    /// identity. A non-empty set replaces the stored one. Clearing labels
309    /// outright is therefore intentionally not expressible — an instance's
310    /// labels are its static shard/peer identity for its lifetime; changing
311    /// them means restarting with new configuration.
312    #[must_use]
313    pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
314        self.labels = labels;
315        self
316    }
317}
318
319/// A resolved gRPC service and the endpoint it is reachable at.
320#[derive(Clone, Debug, PartialEq, Eq, Hash)]
321pub struct GrpcServiceInfo {
322    /// Fully-qualified gRPC service name (e.g. `payment.v1.PaymentApi`).
323    pub service_name: String,
324    /// Endpoint the service is reachable at.
325    pub endpoint: ServiceEndpoint,
326}
327
328impl GrpcServiceInfo {
329    pub fn new(service_name: impl Into<String>, endpoint: ServiceEndpoint) -> Self {
330        Self {
331            service_name: service_name.into(),
332            endpoint,
333        }
334    }
335}
336
337/// Sentinel error wrapped via `anyhow::Error` to signal "the requested gear or
338/// service is not registered (or has no live instance)" through the
339/// [`DirectoryClient`] trait. Consumers downcast to this type to distinguish a
340/// not-ready provider (eventual readiness) from a directory-backend failure —
341/// see `toolkit::discovery::DirectoryEndpointResolver`.
342#[derive(Debug, Clone, PartialEq, Eq)]
343pub struct DirectoryNotFound {
344    /// What was being looked up — e.g. `"gear foo"` or `"service foo.Bar"`.
345    pub resource: String,
346}
347
348impl DirectoryNotFound {
349    pub fn new(resource: impl Into<String>) -> Self {
350        Self {
351            resource: resource.into(),
352        }
353    }
354}
355
356impl std::fmt::Display for DirectoryNotFound {
357    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
358        write!(f, "directory: not found: {}", self.resource)
359    }
360}
361
362impl std::error::Error for DirectoryNotFound {}
363
364/// Sentinel error wrapped via `anyhow::Error` to signal "client-supplied
365/// argument is malformed" (e.g. invalid UUID) through the [`DirectoryClient`]
366/// trait. Allows the gRPC server boundary to return `Status::invalid_argument`
367/// instead of mislabeling a client bug as an internal failure.
368#[derive(Debug, Clone, PartialEq, Eq)]
369pub struct DirectoryInvalidArgument {
370    /// Human-readable description of what was invalid.
371    pub message: String,
372}
373
374impl DirectoryInvalidArgument {
375    pub fn new(message: impl Into<String>) -> Self {
376        Self {
377            message: message.into(),
378        }
379    }
380}
381
382impl std::fmt::Display for DirectoryInvalidArgument {
383    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
384        write!(f, "directory: invalid argument: {}", self.message)
385    }
386}
387
388impl std::error::Error for DirectoryInvalidArgument {}
389
390/// Sentinel error (wrapped via `anyhow::Error`) signalling a **permanent**
391/// authorization refusal: retrying the identical request can never turn a "no"
392/// into a "yes". Carries the gRPC `PermissionDenied` code across the
393/// [`DirectoryClient`] boundary (otherwise lost when a `tonic::Status` is
394/// stringified) so the presence loop stops retrying and logs loudly instead of
395/// spinning at `warn!`. Reached when the peer is not authorized for the gear, its
396/// namespace / trust domain is not allowlisted, or a service name is *pinned* to
397/// another gear (a non-recoverable [`DirectoryServiceNameConflict`]; a recoverable
398/// one is `FailedPrecondition`).
399#[derive(Debug, Clone, PartialEq, Eq)]
400pub struct DirectoryPermissionDenied {
401    /// Human-readable description of why the call was refused.
402    pub message: String,
403}
404
405impl DirectoryPermissionDenied {
406    pub fn new(message: impl Into<String>) -> Self {
407        Self {
408            message: message.into(),
409        }
410    }
411}
412
413impl std::fmt::Display for DirectoryPermissionDenied {
414    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
415        write!(f, "directory: permission denied: {}", self.message)
416    }
417}
418
419impl std::error::Error for DirectoryPermissionDenied {}
420
421/// Sentinel error signalling "a gRPC service name in this registration is already
422/// owned by a *different* gear" (single-gear ownership is enforced atomically in
423/// `GearManager::register_instance`). The gRPC boundary logs the conflicting
424/// `service_name` / `owner` server-side and returns a static-message status whose
425/// code depends on [`recoverable`](Self::recoverable):
426///
427/// - `recoverable` → `Status::failed_precondition`: another gear merely
428///   *currently advertises* the name; it clears when that gear deregisters, so
429///   the registrant retries.
430/// - not `recoverable` → `Status::permission_denied`: the name is pinned to
431///   another gear by the ownership map; retrying can never reassign it.
432#[derive(Debug, Clone, PartialEq, Eq)]
433pub struct DirectoryServiceNameConflict {
434    /// The gRPC service name that is already owned.
435    pub service_name: String,
436    /// The gear that currently owns `service_name`.
437    pub owner: String,
438    /// Whether waiting could clear the conflict (see the type docs). `true` for
439    /// a current-advertiser conflict, `false` for a pinned-ownership conflict.
440    pub recoverable: bool,
441}
442
443impl std::fmt::Display for DirectoryServiceNameConflict {
444    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
445        write!(
446            f,
447            "directory: gRPC service name '{}' already owned by gear '{}'",
448            self.service_name, self.owner
449        )
450    }
451}
452
453impl std::error::Error for DirectoryServiceNameConflict {}
454
455/// Directory API trait for service discovery and instance management
456///
457/// This trait defines the contract for interacting with the gear directory.
458/// It can be implemented by:
459/// - A local implementation that delegates to `GearManager`
460/// - A gRPC client for out-of-process gears
461#[async_trait]
462pub trait DirectoryClient: Send + Sync {
463    /// Resolve a gRPC service by its logical name to an endpoint
464    async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint>;
465
466    /// Resolve a REST endpoint (HTTP base URL) for a gear by its name.
467    ///
468    /// Returns the base URL (e.g. `http://billing:8080`) that callers use to
469    /// make REST requests to the resolved gear.
470    async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint>;
471
472    /// Retrieve the `OpenAPI` spec (JSON) published by a gear.
473    async fn get_openapi_spec(&self, gear_name: &str) -> Result<String>;
474
475    /// List all service instances for a given gear.
476    ///
477    /// Entries are **spec-free**: only the `openapi_spec_hash` is carried, never
478    /// the full `OpenAPI` document, so a multi-instance response stays
479    /// bounded regardless of spec size. Callers fetch the document out-of-band
480    /// via [`get_openapi_spec`](Self::get_openapi_spec); the hash lets a consumer
481    /// detect a spec change without inlining N copies of the document.
482    ///
483    /// An endpoint-less instance is still returned (`endpoint`/`rest_endpoint` =
484    /// `None`, never an empty-URI sentinel): *no-endpoint, not an error*. Only
485    /// [`list_all_instances`](Self::list_all_instances), which needs a dialable
486    /// URI, drops such entries.
487    async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>>;
488
489    /// Resolve instances of `gear` whose labels satisfy `selector`.
490    ///
491    /// Matching is equality-AND (k8s `matchLabels`): an instance matches iff
492    /// it carries **every** `(key, value)` pair in the selector; an empty
493    /// selector matches every instance of the name.
494    ///
495    /// The returned set is **not** health-filtered — matched instances are
496    /// returned regardless of readiness/heartbeat state, each carrying its
497    /// `instance_id`, `labels`, endpoints, and live
498    /// [`state`](ServiceInstanceInfo::state). Applying liveness policy (e.g.
499    /// keeping only [`InstanceState::is_serving`]) is the caller's
500    /// responsibility.
501    ///
502    /// Nor is it **endpoint-filtered**: a matched instance that advertises no
503    /// endpoint yet is still returned (`endpoint`/`rest_endpoint` = `None`) —
504    /// *no-endpoint, not an error* — so the caller can fall back rather than
505    /// have the match silently hidden.
506    ///
507    /// Returned entries are **spec-free** on every implementation — the full
508    /// `OpenAPI` document is omitted (only its hash is carried), so the
509    /// polled resolve payload stays bounded; the caller fetches documents via
510    /// [`get_openapi_spec`](Self::get_openapi_spec). This holds regardless of
511    /// whether the selector is empty. The default body enumerates via
512    /// [`list_instances`](Self::list_instances) and filters in-process;
513    /// transport-backed clients (e.g. the gRPC client) override it to push the
514    /// selector server-side.
515    ///
516    /// Label-based selection is effectively **out-of-process only**. Labels are
517    /// published from the `OoP` serve path's configuration (`oop_http.labels`);
518    /// the in-process registration path (grpc-hub start phase /
519    /// `run_directory_register_phase`) has no label source, so an in-process
520    /// gear carries no labels and a non-empty selector never matches it. This
521    /// matches the deployment model: shards/peers are distinct instances
522    /// (separate pods/processes), which is inherently the `OoP` topology.
523    async fn resolve_by_labels(
524        &self,
525        gear: &str,
526        selector: &LabelSelector,
527    ) -> Result<Vec<ServiceInstanceInfo>> {
528        // cancel-safe: the single await precedes any mutation; cancelling here
529        // just drops the in-flight list and leaves no partial state.
530        let instances = self.list_instances(gear).await?;
531        // Entries are spec-free by construction — `ServiceInstanceInfo` carries
532        // only the `openapi_spec_hash`, never the full OpenAPI document.
533        Ok(instances
534            .into_iter()
535            .filter(|i| selector.matches(&i.labels))
536            .collect())
537    }
538
539    /// List every service instance across all registered gears.
540    ///
541    /// Used by the edge gateway to discover which gears (and their REST
542    /// endpoints) to reverse-proxy. This is a lightweight discovery snapshot.
543    /// Like every enumeration path, it is **spec-free**: entries carry only the
544    /// `openapi_spec_hash`, never the full `OpenAPI` document, even when the
545    /// backing store holds a stored specification. The edge fetches a gear's
546    /// document once, on first discovery, via
547    /// [`get_openapi_spec`](Self::get_openapi_spec).
548    ///
549    /// The returned instances also do **not** carry `labels`: every
550    /// implementation applies [`ServiceInstanceInfo::without_labels`] so the
551    /// snapshot omits them identically regardless of transport. Labels drive
552    /// the targeted [`resolve_by_labels`](Self::resolve_by_labels) path, not
553    /// this snapshot.
554    async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>>;
555
556    /// Register a new gear instance with the directory
557    async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()>;
558
559    /// Deregister a gear instance (for graceful shutdown)
560    async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()>;
561
562    /// Send a heartbeat for a gear instance to indicate it's still alive
563    async fn send_heartbeat(&self, gear: &str, instance_id: &str) -> Result<()>;
564}
565
566#[cfg(test)]
567#[cfg_attr(coverage_nightly, coverage(off))]
568mod tests {
569    use super::*;
570
571    #[test]
572    fn test_service_endpoint_creation() {
573        let http_ep = ServiceEndpoint::http("localhost", 8080);
574        assert_eq!(http_ep.uri, concat!("http", "://localhost:8080"));
575
576        let https_endpoint = ServiceEndpoint::https("localhost", 8443);
577        assert_eq!(https_endpoint.uri, "https://localhost:8443");
578
579        let uds_ep = ServiceEndpoint::uds("/tmp/socket.sock");
580        assert!(uds_ep.uri.starts_with("unix://"));
581        assert!(uds_ep.uri.contains("socket.sock"));
582
583        let custom_ep = ServiceEndpoint::new(concat!("http", "://example.com"));
584        assert_eq!(custom_ep.uri, concat!("http", "://example.com"));
585    }
586
587    #[test]
588    fn test_register_instance_info() {
589        let info = RegisterInstanceInfo::new("test_gear", "instance1")
590            .with_grpc_services(vec![(
591                "test.Service".to_owned(),
592                ServiceEndpoint::http("127.0.0.1", 8001),
593            )])
594            .with_version("1.0.0");
595
596        assert_eq!(info.gear, "test_gear");
597        assert_eq!(info.instance_id, "instance1");
598        assert_eq!(info.grpc_services.len(), 1);
599        assert!(info.rest_endpoint.is_none());
600        assert!(info.openapi_spec.is_none());
601        assert!(info.labels.is_empty());
602    }
603
604    #[test]
605    fn test_register_instance_info_with_rest() {
606        let info = RegisterInstanceInfo {
607            gear: "billing".to_owned(),
608            instance_id: "instance1".to_owned(),
609            grpc_services: vec![],
610            version: Some("2.0.0".to_owned()),
611            labels: BTreeMap::new(),
612            rest_endpoint: Some(ServiceEndpoint::http("billing", 8080)),
613            openapi_spec: Some("{\"openapi\":\"3.1.0\"}".to_owned()),
614        };
615
616        assert_eq!(info.gear, "billing");
617        assert_eq!(
618            info.rest_endpoint.as_ref().unwrap().uri,
619            concat!("http", "://billing:8080")
620        );
621        assert!(info.openapi_spec.is_some());
622    }
623
624    fn labels(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
625        pairs
626            .iter()
627            .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
628            .collect()
629    }
630
631    #[test]
632    fn label_selector_equality_and_matches() {
633        let selector = LabelSelector::new()
634            .with("shard", "7")
635            .with("role", "ingest");
636
637        // Superset of the required labels matches.
638        assert!(selector.matches(&labels(&[
639            ("shard", "7"),
640            ("role", "ingest"),
641            ("extra", "x"),
642        ])));
643        // Missing one required key fails.
644        assert!(!selector.matches(&labels(&[("shard", "7")])));
645        // Wrong value on a required key fails.
646        assert!(!selector.matches(&labels(&[("shard", "8"), ("role", "ingest")])));
647    }
648
649    #[test]
650    fn empty_label_selector_matches_everything() {
651        let selector = LabelSelector::new();
652        assert!(selector.is_empty());
653        assert!(selector.matches(&BTreeMap::new()));
654        assert!(selector.matches(&labels(&[("shard", "7")])));
655    }
656
657    struct StaticDirectory {
658        instances: Vec<ServiceInstanceInfo>,
659    }
660
661    #[async_trait]
662    impl DirectoryClient for StaticDirectory {
663        async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint> {
664            Err(DirectoryNotFound::new(format!("service {service_name}")).into())
665        }
666        async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint> {
667            Err(DirectoryNotFound::new(format!("gear {gear_name}")).into())
668        }
669        async fn get_openapi_spec(&self, gear_name: &str) -> Result<String> {
670            Err(DirectoryNotFound::new(format!("spec {gear_name}")).into())
671        }
672        async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>> {
673            Ok(self
674                .instances
675                .iter()
676                .filter(|i| i.gear == gear)
677                .cloned()
678                .collect())
679        }
680        async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>> {
681            Ok(self.instances.clone())
682        }
683        async fn register_instance(&self, _info: RegisterInstanceInfo) -> Result<()> {
684            Ok(())
685        }
686        async fn deregister_instance(&self, _gear: &str, _instance_id: &str) -> Result<()> {
687            Ok(())
688        }
689        async fn send_heartbeat(&self, _gear: &str, _instance_id: &str) -> Result<()> {
690            Ok(())
691        }
692    }
693
694    fn instance(id: &str, labels_pairs: &[(&str, &str)]) -> ServiceInstanceInfo {
695        ServiceInstanceInfo {
696            gear: "worker".to_owned(),
697            instance_id: id.to_owned(),
698            labels: labels(labels_pairs),
699            ..Default::default()
700        }
701    }
702
703    #[tokio::test]
704    async fn resolve_by_labels_default_filters_by_equality_and() {
705        let dir = StaticDirectory {
706            instances: vec![
707                instance("a", &[("shard", "7")]),
708                instance("b", &[("shard", "8")]),
709                instance("c", &[("shard", "7"), ("role", "ingest")]),
710            ],
711        };
712
713        let selector = LabelSelector::new().with("shard", "7");
714        let mut matched = dir
715            .resolve_by_labels("worker", &selector)
716            .await
717            .unwrap()
718            .into_iter()
719            .map(|i| i.instance_id)
720            .collect::<Vec<_>>();
721        matched.sort();
722        assert_eq!(matched, vec!["a".to_owned(), "c".to_owned()]);
723
724        // An empty selector returns every instance of the name.
725        let all = dir
726            .resolve_by_labels("worker", &LabelSelector::new())
727            .await
728            .unwrap();
729        assert_eq!(all.len(), 3);
730    }
731}