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 /// Optional `OpenAPI` spec (JSON) this instance published, if any.
140 pub openapi_spec: Option<String>,
141 /// Stable content token for the published `OpenAPI` spec, if any.
142 pub openapi_spec_hash: Option<String>,
143 /// Published gRPC services as `(service name, endpoint)` pairs.
144 ///
145 /// Kept as an ordered list rather than a map so it round-trips a repeated
146 /// proto field without reordering; service names are unique per instance.
147 /// Carried back by `list_instances` so a subsequent `register_instance`
148 /// (which replaces the entry wholesale) can augment — rather than clobber —
149 /// the previously-registered gRPC services when adding a REST endpoint.
150 pub grpc_services: Vec<(String, ServiceEndpoint)>,
151 /// Stable addressing labels published by this instance (k8s `matchLabels`
152 /// style). Used by [`DirectoryClient::resolve_by_labels`] to select shards
153 /// or peers within a directory name.
154 pub labels: BTreeMap<String, String>,
155 /// Live serving state, so a `resolve_by_labels` caller can apply its own
156 /// health policy without a second call. Defaults to
157 /// [`InstanceState::Registered`].
158 pub state: InstanceState,
159}
160
161impl ServiceInstanceInfo {
162 /// Start a projection for `gear` / `instance_id` with no primary endpoint
163 /// and no optional metadata. Chain the `with_*` builders to populate it.
164 /// `#[non_exhaustive]` means callers must use this builder rather than a
165 /// struct literal, so future fields stay additive.
166 #[must_use]
167 pub fn new(gear: impl Into<String>, instance_id: impl Into<String>) -> Self {
168 Self {
169 gear: gear.into(),
170 instance_id: instance_id.into(),
171 ..Self::default()
172 }
173 }
174
175 /// Set the primary endpoint (`None` for an instance with no primary
176 /// endpoint).
177 #[must_use]
178 pub fn with_endpoint(mut self, endpoint: Option<ServiceEndpoint>) -> Self {
179 self.endpoint = endpoint;
180 self
181 }
182
183 /// Set the optional version string.
184 #[must_use]
185 pub fn with_version(mut self, version: Option<String>) -> Self {
186 self.version = version;
187 self
188 }
189
190 /// Set the optional REST endpoint.
191 #[must_use]
192 pub fn with_rest_endpoint(mut self, rest_endpoint: Option<ServiceEndpoint>) -> Self {
193 self.rest_endpoint = rest_endpoint;
194 self
195 }
196
197 /// Set the optional `OpenAPI` spec (JSON).
198 #[must_use]
199 pub fn with_openapi_spec(mut self, openapi_spec: Option<String>) -> Self {
200 self.openapi_spec = openapi_spec;
201 self
202 }
203
204 /// Set the optional `OpenAPI` spec content token.
205 #[must_use]
206 pub fn with_openapi_spec_hash(mut self, openapi_spec_hash: Option<String>) -> Self {
207 self.openapi_spec_hash = openapi_spec_hash;
208 self
209 }
210
211 /// Set the published gRPC services.
212 #[must_use]
213 pub fn with_grpc_services(mut self, grpc_services: Vec<(String, ServiceEndpoint)>) -> Self {
214 self.grpc_services = grpc_services;
215 self
216 }
217
218 /// Set the stable addressing labels.
219 #[must_use]
220 pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
221 self.labels = labels;
222 self
223 }
224
225 /// Set the live serving state.
226 #[must_use]
227 pub fn with_state(mut self, state: InstanceState) -> Self {
228 self.state = state;
229 self
230 }
231
232 /// Drop the stable addressing labels, returning `self`.
233 ///
234 /// The shared transform every [`DirectoryClient::list_all_instances`]
235 /// implementation applies so the cross-gear snapshot omits labels
236 /// identically across transports — see that method's doc for why.
237 #[must_use]
238 pub fn without_labels(mut self) -> Self {
239 self.labels.clear();
240 self
241 }
242}
243
244/// Information for registering a new gear instance
245#[derive(Debug, Clone, Default)]
246#[non_exhaustive]
247pub struct RegisterInstanceInfo {
248 /// Gear name
249 pub gear: String,
250 /// Unique instance identifier
251 pub instance_id: String,
252 /// Published gRPC services as `(service name, endpoint)` pairs (service
253 /// names are unique per instance; ordered to mirror the proto repeated
254 /// field rather than reordered into a map).
255 pub grpc_services: Vec<(String, ServiceEndpoint)>,
256 /// Optional version string
257 pub version: Option<String>,
258 /// Optional REST endpoint (HTTP base URL) exposed by the gear.
259 pub rest_endpoint: Option<ServiceEndpoint>,
260 /// Optional `OpenAPI` spec (JSON) published by the gear.
261 pub openapi_spec: Option<String>,
262 /// Stable addressing labels for this instance (k8s `matchLabels` style).
263 pub labels: BTreeMap<String, String>,
264}
265
266impl RegisterInstanceInfo {
267 /// Start a registration for `gear` / `instance_id` with no endpoints,
268 /// version, spec, or labels. Chain the `with_*` builders to populate it.
269 #[must_use]
270 pub fn new(gear: impl Into<String>, instance_id: impl Into<String>) -> Self {
271 Self {
272 gear: gear.into(),
273 instance_id: instance_id.into(),
274 ..Self::default()
275 }
276 }
277
278 /// Set the published gRPC services.
279 #[must_use]
280 pub fn with_grpc_services(mut self, grpc_services: Vec<(String, ServiceEndpoint)>) -> Self {
281 self.grpc_services = grpc_services;
282 self
283 }
284
285 /// Set the version string.
286 #[must_use]
287 pub fn with_version(mut self, version: impl Into<String>) -> Self {
288 self.version = Some(version.into());
289 self
290 }
291
292 /// Set the REST endpoint (HTTP base URL).
293 #[must_use]
294 pub fn with_rest_endpoint(mut self, rest_endpoint: ServiceEndpoint) -> Self {
295 self.rest_endpoint = Some(rest_endpoint);
296 self
297 }
298
299 /// Set the published `OpenAPI` spec (JSON).
300 #[must_use]
301 pub fn with_openapi_spec(mut self, openapi_spec: impl Into<String>) -> Self {
302 self.openapi_spec = Some(openapi_spec.into());
303 self
304 }
305
306 /// Set the stable addressing labels.
307 ///
308 /// Re-registration replaces the directory entry wholesale, but an **empty**
309 /// label set is treated as "preserve the stored labels", not "clear them"
310 /// (see `GearManager::with_metadata_of`): the periodic self-heal and REST
311 /// augmentation re-register without labels and must not erase a shard's
312 /// identity. A non-empty set replaces the stored one. Clearing labels
313 /// outright is therefore intentionally not expressible — an instance's
314 /// labels are its static shard/peer identity for its lifetime; changing
315 /// them means restarting with new configuration.
316 #[must_use]
317 pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
318 self.labels = labels;
319 self
320 }
321}
322
323/// A resolved gRPC service and the endpoint it is reachable at.
324#[derive(Clone, Debug, PartialEq, Eq, Hash)]
325pub struct GrpcServiceInfo {
326 /// Fully-qualified gRPC service name (e.g. `payment.v1.PaymentApi`).
327 pub service_name: String,
328 /// Endpoint the service is reachable at.
329 pub endpoint: ServiceEndpoint,
330}
331
332impl GrpcServiceInfo {
333 pub fn new(service_name: impl Into<String>, endpoint: ServiceEndpoint) -> Self {
334 Self {
335 service_name: service_name.into(),
336 endpoint,
337 }
338 }
339}
340
341/// Sentinel error wrapped via `anyhow::Error` to signal "the requested gear or
342/// service is not registered (or has no live instance)" through the
343/// [`DirectoryClient`] trait. Consumers downcast to this type to distinguish a
344/// not-ready provider (eventual readiness) from a directory-backend failure —
345/// see `toolkit::discovery::DirectoryEndpointResolver`.
346#[derive(Debug, Clone, PartialEq, Eq)]
347pub struct DirectoryNotFound {
348 /// What was being looked up — e.g. `"gear foo"` or `"service foo.Bar"`.
349 pub resource: String,
350}
351
352impl DirectoryNotFound {
353 pub fn new(resource: impl Into<String>) -> Self {
354 Self {
355 resource: resource.into(),
356 }
357 }
358}
359
360impl std::fmt::Display for DirectoryNotFound {
361 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
362 write!(f, "directory: not found: {}", self.resource)
363 }
364}
365
366impl std::error::Error for DirectoryNotFound {}
367
368/// Sentinel error wrapped via `anyhow::Error` to signal "client-supplied
369/// argument is malformed" (e.g. invalid UUID) through the [`DirectoryClient`]
370/// trait. Allows the gRPC server boundary to return `Status::invalid_argument`
371/// instead of mislabeling a client bug as an internal failure.
372#[derive(Debug, Clone, PartialEq, Eq)]
373pub struct DirectoryInvalidArgument {
374 /// Human-readable description of what was invalid.
375 pub message: String,
376}
377
378impl DirectoryInvalidArgument {
379 pub fn new(message: impl Into<String>) -> Self {
380 Self {
381 message: message.into(),
382 }
383 }
384}
385
386impl std::fmt::Display for DirectoryInvalidArgument {
387 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
388 write!(f, "directory: invalid argument: {}", self.message)
389 }
390}
391
392impl std::error::Error for DirectoryInvalidArgument {}
393
394/// Directory API trait for service discovery and instance management
395///
396/// This trait defines the contract for interacting with the gear directory.
397/// It can be implemented by:
398/// - A local implementation that delegates to `GearManager`
399/// - A gRPC client for out-of-process gears
400#[async_trait]
401pub trait DirectoryClient: Send + Sync {
402 /// Resolve a gRPC service by its logical name to an endpoint
403 async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint>;
404
405 /// Resolve a REST endpoint (HTTP base URL) for a gear by its name.
406 ///
407 /// Returns the base URL (e.g. `http://billing:8080`) that callers use to
408 /// make REST requests to the resolved gear.
409 async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint>;
410
411 /// Retrieve the `OpenAPI` spec (JSON) published by a gear.
412 async fn get_openapi_spec(&self, gear_name: &str) -> Result<String>;
413
414 /// List all service instances for a given gear.
415 ///
416 /// Entries are **spec-free**: only the `openapi_spec_hash` is carried, never
417 /// the full `openapi_spec` document, so a multi-instance response stays
418 /// bounded regardless of spec size. Callers fetch the document out-of-band
419 /// via [`get_openapi_spec`](Self::get_openapi_spec); the hash lets a consumer
420 /// detect a spec change without inlining N copies of the document.
421 ///
422 /// An endpoint-less instance is still returned (`endpoint`/`rest_endpoint` =
423 /// `None`, never an empty-URI sentinel): *no-endpoint, not an error*. Only
424 /// [`list_all_instances`](Self::list_all_instances), which needs a dialable
425 /// URI, drops such entries.
426 async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>>;
427
428 /// Resolve instances of `gear` whose labels satisfy `selector`.
429 ///
430 /// Matching is equality-AND (k8s `matchLabels`): an instance matches iff
431 /// it carries **every** `(key, value)` pair in the selector; an empty
432 /// selector matches every instance of the name.
433 ///
434 /// The returned set is **not** health-filtered — matched instances are
435 /// returned regardless of readiness/heartbeat state, each carrying its
436 /// `instance_id`, `labels`, endpoints, and live
437 /// [`state`](ServiceInstanceInfo::state). Applying liveness policy (e.g.
438 /// keeping only [`InstanceState::is_serving`]) is the caller's
439 /// responsibility.
440 ///
441 /// Nor is it **endpoint-filtered**: a matched instance that advertises no
442 /// endpoint yet is still returned (`endpoint`/`rest_endpoint` = `None`) —
443 /// *no-endpoint, not an error* — so the caller can fall back rather than
444 /// have the match silently hidden.
445 ///
446 /// Returned entries are **spec-free** on every implementation — the full
447 /// `openapi_spec` document is omitted (only its hash is carried), so the
448 /// polled resolve payload stays bounded; the caller fetches documents via
449 /// [`get_openapi_spec`](Self::get_openapi_spec). This holds regardless of
450 /// whether the selector is empty. The default body enumerates via
451 /// [`list_instances`](Self::list_instances), filters in-process, and drops
452 /// the spec; transport-backed clients (e.g. the gRPC client) override it to
453 /// push the selector server-side and request spec-free entries over the
454 /// wire.
455 ///
456 /// Label-based selection is effectively **out-of-process only**. Labels are
457 /// published from the `OoP` serve path's configuration (`oop_http.labels`);
458 /// the in-process registration path (grpc-hub start phase /
459 /// `run_directory_register_phase`) has no label source, so an in-process
460 /// gear carries no labels and a non-empty selector never matches it. This
461 /// matches the deployment model: shards/peers are distinct instances
462 /// (separate pods/processes), which is inherently the `OoP` topology.
463 async fn resolve_by_labels(
464 &self,
465 gear: &str,
466 selector: &LabelSelector,
467 ) -> Result<Vec<ServiceInstanceInfo>> {
468 // cancel-safe: the single await precedes any mutation; cancelling here
469 // just drops the in-flight list and leaves no partial state.
470 let instances = self.list_instances(gear).await?;
471 Ok(instances
472 .into_iter()
473 .filter(|i| selector.matches(&i.labels))
474 // Spec-free entries, matching the documented contract and the
475 // transport-backed overrides: the label-resolve path never inlines
476 // the OpenAPI document (the hash still rides along).
477 .map(|mut i| {
478 i.openapi_spec = None;
479 i
480 })
481 .collect())
482 }
483
484 /// List every service instance across all registered gears.
485 ///
486 /// Used by the edge gateway to discover which gears (and their REST
487 /// endpoints) to reverse-proxy. This is a lightweight discovery snapshot:
488 /// the returned instances do **not** carry `openapi_spec` — even when the
489 /// backing store holds a stored specification. The edge fetches a gear's
490 /// document once, on first discovery, via
491 /// [`get_openapi_spec`](Self::get_openapi_spec).
492 ///
493 /// The returned instances also do **not** carry `labels`: every
494 /// implementation applies [`ServiceInstanceInfo::without_labels`] so the
495 /// snapshot omits them identically regardless of transport. Labels drive
496 /// the targeted [`resolve_by_labels`](Self::resolve_by_labels) path, not
497 /// this snapshot.
498 async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>>;
499
500 /// Register a new gear instance with the directory
501 async fn register_instance(&self, info: RegisterInstanceInfo) -> Result<()>;
502
503 /// Deregister a gear instance (for graceful shutdown)
504 async fn deregister_instance(&self, gear: &str, instance_id: &str) -> Result<()>;
505
506 /// Send a heartbeat for a gear instance to indicate it's still alive
507 async fn send_heartbeat(&self, gear: &str, instance_id: &str) -> Result<()>;
508}
509
510#[cfg(test)]
511#[cfg_attr(coverage_nightly, coverage(off))]
512mod tests {
513 use super::*;
514
515 #[test]
516 fn test_service_endpoint_creation() {
517 let http_ep = ServiceEndpoint::http("localhost", 8080);
518 assert_eq!(http_ep.uri, concat!("http", "://localhost:8080"));
519
520 let https_endpoint = ServiceEndpoint::https("localhost", 8443);
521 assert_eq!(https_endpoint.uri, "https://localhost:8443");
522
523 let uds_ep = ServiceEndpoint::uds("/tmp/socket.sock");
524 assert!(uds_ep.uri.starts_with("unix://"));
525 assert!(uds_ep.uri.contains("socket.sock"));
526
527 let custom_ep = ServiceEndpoint::new(concat!("http", "://example.com"));
528 assert_eq!(custom_ep.uri, concat!("http", "://example.com"));
529 }
530
531 #[test]
532 fn test_register_instance_info() {
533 let info = RegisterInstanceInfo::new("test_gear", "instance1")
534 .with_grpc_services(vec![(
535 "test.Service".to_owned(),
536 ServiceEndpoint::http("127.0.0.1", 8001),
537 )])
538 .with_version("1.0.0");
539
540 assert_eq!(info.gear, "test_gear");
541 assert_eq!(info.instance_id, "instance1");
542 assert_eq!(info.grpc_services.len(), 1);
543 assert!(info.rest_endpoint.is_none());
544 assert!(info.openapi_spec.is_none());
545 assert!(info.labels.is_empty());
546 }
547
548 #[test]
549 fn test_register_instance_info_with_rest() {
550 let info = RegisterInstanceInfo {
551 gear: "billing".to_owned(),
552 instance_id: "instance1".to_owned(),
553 grpc_services: vec![],
554 version: Some("2.0.0".to_owned()),
555 labels: BTreeMap::new(),
556 rest_endpoint: Some(ServiceEndpoint::http("billing", 8080)),
557 openapi_spec: Some("{\"openapi\":\"3.1.0\"}".to_owned()),
558 };
559
560 assert_eq!(info.gear, "billing");
561 assert_eq!(
562 info.rest_endpoint.as_ref().unwrap().uri,
563 concat!("http", "://billing:8080")
564 );
565 assert!(info.openapi_spec.is_some());
566 }
567
568 fn labels(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
569 pairs
570 .iter()
571 .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
572 .collect()
573 }
574
575 #[test]
576 fn label_selector_equality_and_matches() {
577 let selector = LabelSelector::new()
578 .with("shard", "7")
579 .with("role", "ingest");
580
581 // Superset of the required labels matches.
582 assert!(selector.matches(&labels(&[
583 ("shard", "7"),
584 ("role", "ingest"),
585 ("extra", "x"),
586 ])));
587 // Missing one required key fails.
588 assert!(!selector.matches(&labels(&[("shard", "7")])));
589 // Wrong value on a required key fails.
590 assert!(!selector.matches(&labels(&[("shard", "8"), ("role", "ingest")])));
591 }
592
593 #[test]
594 fn empty_label_selector_matches_everything() {
595 let selector = LabelSelector::new();
596 assert!(selector.is_empty());
597 assert!(selector.matches(&BTreeMap::new()));
598 assert!(selector.matches(&labels(&[("shard", "7")])));
599 }
600
601 struct StaticDirectory {
602 instances: Vec<ServiceInstanceInfo>,
603 }
604
605 #[async_trait]
606 impl DirectoryClient for StaticDirectory {
607 async fn resolve_grpc_service(&self, service_name: &str) -> Result<ServiceEndpoint> {
608 Err(DirectoryNotFound::new(format!("service {service_name}")).into())
609 }
610 async fn resolve_rest_service(&self, gear_name: &str) -> Result<ServiceEndpoint> {
611 Err(DirectoryNotFound::new(format!("gear {gear_name}")).into())
612 }
613 async fn get_openapi_spec(&self, gear_name: &str) -> Result<String> {
614 Err(DirectoryNotFound::new(format!("spec {gear_name}")).into())
615 }
616 async fn list_instances(&self, gear: &str) -> Result<Vec<ServiceInstanceInfo>> {
617 Ok(self
618 .instances
619 .iter()
620 .filter(|i| i.gear == gear)
621 .cloned()
622 .collect())
623 }
624 async fn list_all_instances(&self) -> Result<Vec<ServiceInstanceInfo>> {
625 Ok(self.instances.clone())
626 }
627 async fn register_instance(&self, _info: RegisterInstanceInfo) -> Result<()> {
628 Ok(())
629 }
630 async fn deregister_instance(&self, _gear: &str, _instance_id: &str) -> Result<()> {
631 Ok(())
632 }
633 async fn send_heartbeat(&self, _gear: &str, _instance_id: &str) -> Result<()> {
634 Ok(())
635 }
636 }
637
638 fn instance(id: &str, labels_pairs: &[(&str, &str)]) -> ServiceInstanceInfo {
639 ServiceInstanceInfo {
640 gear: "worker".to_owned(),
641 instance_id: id.to_owned(),
642 labels: labels(labels_pairs),
643 ..Default::default()
644 }
645 }
646
647 #[tokio::test]
648 async fn resolve_by_labels_default_filters_by_equality_and() {
649 let dir = StaticDirectory {
650 instances: vec![
651 instance("a", &[("shard", "7")]),
652 instance("b", &[("shard", "8")]),
653 instance("c", &[("shard", "7"), ("role", "ingest")]),
654 ],
655 };
656
657 let selector = LabelSelector::new().with("shard", "7");
658 let mut matched = dir
659 .resolve_by_labels("worker", &selector)
660 .await
661 .unwrap()
662 .into_iter()
663 .map(|i| i.instance_id)
664 .collect::<Vec<_>>();
665 matched.sort();
666 assert_eq!(matched, vec!["a".to_owned(), "c".to_owned()]);
667
668 // An empty selector returns every instance of the name.
669 let all = dir
670 .resolve_by_labels("worker", &LabelSelector::new())
671 .await
672 .unwrap();
673 assert_eq!(all.len(), 3);
674 }
675}