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}