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/// Directory API trait for service discovery and instance management
391///
392/// This trait defines the contract for interacting with the gear directory.
393/// It can be implemented by:
394/// - A local implementation that delegates to `GearManager`
395/// - A gRPC client for out-of-process gears
396#[async_trait]
397pub trait DirectoryClient: Send + Sync {
398    /// Resolve a gRPC service by its logical name to an endpoint
399    async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint>;
400
401    /// Resolve a REST endpoint (HTTP base URL) for a gear by its name.
402    ///
403    /// Returns the base URL (e.g. `http://billing:8080`) that callers use to
404    /// make REST requests to the resolved gear.
405    async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint>;
406
407    /// Retrieve the `OpenAPI` spec (JSON) published by a gear.
408    async fn get_openapi_spec(&self, gear_name: &str) -> Result<String>;
409
410    /// List all service instances for a given gear.
411    ///
412    /// Entries are **spec-free**: only the `openapi_spec_hash` is carried, never
413    /// the full `OpenAPI` document, so a multi-instance response stays
414    /// bounded regardless of spec size. Callers fetch the document out-of-band
415    /// via [`get_openapi_spec`](Self::get_openapi_spec); the hash lets a consumer
416    /// detect a spec change without inlining N copies of the document.
417    ///
418    /// An endpoint-less instance is still returned (`endpoint`/`rest_endpoint` =
419    /// `None`, never an empty-URI sentinel): *no-endpoint, not an error*. Only
420    /// [`list_all_instances`](Self::list_all_instances), which needs a dialable
421    /// URI, drops such entries.
422    async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>>;
423
424    /// Resolve instances of `gear` whose labels satisfy `selector`.
425    ///
426    /// Matching is equality-AND (k8s `matchLabels`): an instance matches iff
427    /// it carries **every** `(key, value)` pair in the selector; an empty
428    /// selector matches every instance of the name.
429    ///
430    /// The returned set is **not** health-filtered — matched instances are
431    /// returned regardless of readiness/heartbeat state, each carrying its
432    /// `instance_id`, `labels`, endpoints, and live
433    /// [`state`](ServiceInstanceInfo::state). Applying liveness policy (e.g.
434    /// keeping only [`InstanceState::is_serving`]) is the caller's
435    /// responsibility.
436    ///
437    /// Nor is it **endpoint-filtered**: a matched instance that advertises no
438    /// endpoint yet is still returned (`endpoint`/`rest_endpoint` = `None`) —
439    /// *no-endpoint, not an error* — so the caller can fall back rather than
440    /// have the match silently hidden.
441    ///
442    /// Returned entries are **spec-free** on every implementation — the full
443    /// `OpenAPI` document is omitted (only its hash is carried), so the
444    /// polled resolve payload stays bounded; the caller fetches documents via
445    /// [`get_openapi_spec`](Self::get_openapi_spec). This holds regardless of
446    /// whether the selector is empty. The default body enumerates via
447    /// [`list_instances`](Self::list_instances) and filters in-process;
448    /// transport-backed clients (e.g. the gRPC client) override it to push the
449    /// selector server-side.
450    ///
451    /// Label-based selection is effectively **out-of-process only**. Labels are
452    /// published from the `OoP` serve path's configuration (`oop_http.labels`);
453    /// the in-process registration path (grpc-hub start phase /
454    /// `run_directory_register_phase`) has no label source, so an in-process
455    /// gear carries no labels and a non-empty selector never matches it. This
456    /// matches the deployment model: shards/peers are distinct instances
457    /// (separate pods/processes), which is inherently the `OoP` topology.
458    async fn resolve_by_labels(
459        &self,
460        gear: &str,
461        selector: &LabelSelector,
462    ) -> Result<Vec<ServiceInstanceInfo>> {
463        // cancel-safe: the single await precedes any mutation; cancelling here
464        // just drops the in-flight list and leaves no partial state.
465        let instances = self.list_instances(gear).await?;
466        // Entries are spec-free by construction — `ServiceInstanceInfo` carries
467        // only the `openapi_spec_hash`, never the full OpenAPI document.
468        Ok(instances
469            .into_iter()
470            .filter(|i| selector.matches(&i.labels))
471            .collect())
472    }
473
474    /// List every service instance across all registered gears.
475    ///
476    /// Used by the edge gateway to discover which gears (and their REST
477    /// endpoints) to reverse-proxy. This is a lightweight discovery snapshot.
478    /// Like every enumeration path, it is **spec-free**: entries carry only the
479    /// `openapi_spec_hash`, never the full `OpenAPI` document, even when the
480    /// backing store holds a stored specification. The edge fetches a gear's
481    /// document once, on first discovery, via
482    /// [`get_openapi_spec`](Self::get_openapi_spec).
483    ///
484    /// The returned instances also do **not** carry `labels`: every
485    /// implementation applies [`ServiceInstanceInfo::without_labels`] so the
486    /// snapshot omits them identically regardless of transport. Labels drive
487    /// the targeted [`resolve_by_labels`](Self::resolve_by_labels) path, not
488    /// this snapshot.
489    async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>>;
490
491    /// Register a new gear instance with the directory
492    async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()>;
493
494    /// Deregister a gear instance (for graceful shutdown)
495    async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()>;
496
497    /// Send a heartbeat for a gear instance to indicate it's still alive
498    async fn send_heartbeat(&self, gear: &str, instance_id: &str) -> Result<()>;
499}
500
501#[cfg(test)]
502#[cfg_attr(coverage_nightly, coverage(off))]
503mod tests {
504    use super::*;
505
506    #[test]
507    fn test_service_endpoint_creation() {
508        let http_ep = ServiceEndpoint::http("localhost", 8080);
509        assert_eq!(http_ep.uri, concat!("http", "://localhost:8080"));
510
511        let https_endpoint = ServiceEndpoint::https("localhost", 8443);
512        assert_eq!(https_endpoint.uri, "https://localhost:8443");
513
514        let uds_ep = ServiceEndpoint::uds("/tmp/socket.sock");
515        assert!(uds_ep.uri.starts_with("unix://"));
516        assert!(uds_ep.uri.contains("socket.sock"));
517
518        let custom_ep = ServiceEndpoint::new(concat!("http", "://example.com"));
519        assert_eq!(custom_ep.uri, concat!("http", "://example.com"));
520    }
521
522    #[test]
523    fn test_register_instance_info() {
524        let info = RegisterInstanceInfo::new("test_gear", "instance1")
525            .with_grpc_services(vec![(
526                "test.Service".to_owned(),
527                ServiceEndpoint::http("127.0.0.1", 8001),
528            )])
529            .with_version("1.0.0");
530
531        assert_eq!(info.gear, "test_gear");
532        assert_eq!(info.instance_id, "instance1");
533        assert_eq!(info.grpc_services.len(), 1);
534        assert!(info.rest_endpoint.is_none());
535        assert!(info.openapi_spec.is_none());
536        assert!(info.labels.is_empty());
537    }
538
539    #[test]
540    fn test_register_instance_info_with_rest() {
541        let info = RegisterInstanceInfo {
542            gear: "billing".to_owned(),
543            instance_id: "instance1".to_owned(),
544            grpc_services: vec![],
545            version: Some("2.0.0".to_owned()),
546            labels: BTreeMap::new(),
547            rest_endpoint: Some(ServiceEndpoint::http("billing", 8080)),
548            openapi_spec: Some("{\"openapi\":\"3.1.0\"}".to_owned()),
549        };
550
551        assert_eq!(info.gear, "billing");
552        assert_eq!(
553            info.rest_endpoint.as_ref().unwrap().uri,
554            concat!("http", "://billing:8080")
555        );
556        assert!(info.openapi_spec.is_some());
557    }
558
559    fn labels(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
560        pairs
561            .iter()
562            .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
563            .collect()
564    }
565
566    #[test]
567    fn label_selector_equality_and_matches() {
568        let selector = LabelSelector::new()
569            .with("shard", "7")
570            .with("role", "ingest");
571
572        // Superset of the required labels matches.
573        assert!(selector.matches(&labels(&[
574            ("shard", "7"),
575            ("role", "ingest"),
576            ("extra", "x"),
577        ])));
578        // Missing one required key fails.
579        assert!(!selector.matches(&labels(&[("shard", "7")])));
580        // Wrong value on a required key fails.
581        assert!(!selector.matches(&labels(&[("shard", "8"), ("role", "ingest")])));
582    }
583
584    #[test]
585    fn empty_label_selector_matches_everything() {
586        let selector = LabelSelector::new();
587        assert!(selector.is_empty());
588        assert!(selector.matches(&BTreeMap::new()));
589        assert!(selector.matches(&labels(&[("shard", "7")])));
590    }
591
592    struct StaticDirectory {
593        instances: Vec<ServiceInstanceInfo>,
594    }
595
596    #[async_trait]
597    impl DirectoryClient for StaticDirectory {
598        async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint> {
599            Err(DirectoryNotFound::new(format!("service {service_name}")).into())
600        }
601        async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint> {
602            Err(DirectoryNotFound::new(format!("gear {gear_name}")).into())
603        }
604        async fn get_openapi_spec(&self, gear_name: &str) -> Result<String> {
605            Err(DirectoryNotFound::new(format!("spec {gear_name}")).into())
606        }
607        async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>> {
608            Ok(self
609                .instances
610                .iter()
611                .filter(|i| i.gear == gear)
612                .cloned()
613                .collect())
614        }
615        async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>> {
616            Ok(self.instances.clone())
617        }
618        async fn register_instance(&self, _info: RegisterInstanceInfo) -> Result<()> {
619            Ok(())
620        }
621        async fn deregister_instance(&self, _gear: &str, _instance_id: &str) -> Result<()> {
622            Ok(())
623        }
624        async fn send_heartbeat(&self, _gear: &str, _instance_id: &str) -> Result<()> {
625            Ok(())
626        }
627    }
628
629    fn instance(id: &str, labels_pairs: &[(&str, &str)]) -> ServiceInstanceInfo {
630        ServiceInstanceInfo {
631            gear: "worker".to_owned(),
632            instance_id: id.to_owned(),
633            labels: labels(labels_pairs),
634            ..Default::default()
635        }
636    }
637
638    #[tokio::test]
639    async fn resolve_by_labels_default_filters_by_equality_and() {
640        let dir = StaticDirectory {
641            instances: vec![
642                instance("a", &[("shard", "7")]),
643                instance("b", &[("shard", "8")]),
644                instance("c", &[("shard", "7"), ("role", "ingest")]),
645            ],
646        };
647
648        let selector = LabelSelector::new().with("shard", "7");
649        let mut matched = dir
650            .resolve_by_labels("worker", &selector)
651            .await
652            .unwrap()
653            .into_iter()
654            .map(|i| i.instance_id)
655            .collect::<Vec<_>>();
656        matched.sort();
657        assert_eq!(matched, vec!["a".to_owned(), "c".to_owned()]);
658
659        // An empty selector returns every instance of the name.
660        let all = dir
661            .resolve_by_labels("worker", &LabelSelector::new())
662            .await
663            .unwrap();
664        assert_eq!(all.len(), 3);
665    }
666}