Skip to main content

toolkit/runtime/
gear_manager.rs

1//! Gear Manager - tracks and manages all live gear instances in the runtime
2
3use dashmap::DashMap;
4use std::collections::{BTreeMap, HashMap};
5use std::sync::Arc;
6use std::time::{Duration, Instant};
7use uuid::Uuid;
8
9/// Represents an endpoint where a gear instance can be reached
10#[derive(Clone, Debug, PartialEq, Eq, Hash)]
11pub struct Endpoint {
12    pub uri: String,
13}
14
15/// Typed view of an endpoint for parsing and matching
16#[derive(Clone, Debug, PartialEq, Eq)]
17pub enum EndpointKind {
18    /// TCP endpoint with resolved socket address
19    Tcp(std::net::SocketAddr),
20    /// Unix domain socket with file path
21    Uds(std::path::PathBuf),
22    /// Other/unparsed endpoint URI
23    Other(String),
24}
25
26impl Endpoint {
27    pub fn from_uri<S: Into<String>>(s: S) -> Self {
28        Self { uri: s.into() }
29    }
30
31    pub fn uds(path: impl AsRef<std::path::Path>) -> Self {
32        Self {
33            uri: format!("unix://{}", path.as_ref().display()),
34        }
35    }
36
37    #[must_use]
38    pub fn http(host: &str, port: u16) -> Self {
39        Self {
40            uri: format!("http://{host}:{port}"),
41        }
42    }
43
44    #[must_use]
45    pub fn https(host: &str, port: u16) -> Self {
46        Self {
47            uri: format!("https://{host}:{port}"),
48        }
49    }
50
51    /// Parse the endpoint URI into a typed view
52    #[must_use]
53    pub fn kind(&self) -> EndpointKind {
54        if let Some(rest) = self.uri.strip_prefix("unix://") {
55            return EndpointKind::Uds(std::path::PathBuf::from(rest));
56        }
57        if let Some(rest) = self.uri.strip_prefix("http://")
58            && let Ok(addr) = rest.parse::<std::net::SocketAddr>()
59        {
60            return EndpointKind::Tcp(addr);
61        }
62        if let Some(rest) = self.uri.strip_prefix("https://")
63            && let Ok(addr) = rest.parse::<std::net::SocketAddr>()
64        {
65            return EndpointKind::Tcp(addr);
66        }
67        EndpointKind::Other(self.uri.clone())
68    }
69}
70
71#[derive(Clone, Copy, Debug, PartialEq, Eq)]
72pub enum InstanceState {
73    Registered,
74    Ready,
75    Healthy,
76    Quarantined,
77    Draining,
78}
79
80/// Runtime state of an instance (guarded by `RwLock` for safe mutation)
81#[derive(Clone, Debug)]
82pub struct InstanceRuntimeState {
83    pub last_heartbeat: Instant,
84    pub state: InstanceState,
85}
86
87/// Represents a single instance of a gear
88#[derive(Debug)]
89#[must_use]
90pub struct GearInstance {
91    pub gear: String,
92    pub instance_id: Uuid,
93    pub control: Option<Endpoint>,
94    pub grpc_services: HashMap<String, Endpoint>,
95    pub version: Option<String>,
96    pub rest_endpoint: Option<Endpoint>,
97    pub openapi_spec: Option<String>,
98    /// Stable addressing labels (k8s `matchLabels` style) advertised by this
99    /// instance, used for label-based instance selection.
100    pub labels: BTreeMap<String, String>,
101    inner: Arc<parking_lot::RwLock<InstanceRuntimeState>>,
102}
103
104impl Clone for GearInstance {
105    fn clone(&self) -> Self {
106        Self {
107            gear: self.gear.clone(),
108            instance_id: self.instance_id,
109            control: self.control.clone(),
110            grpc_services: self.grpc_services.clone(),
111            version: self.version.clone(),
112            rest_endpoint: self.rest_endpoint.clone(),
113            openapi_spec: self.openapi_spec.clone(),
114            labels: self.labels.clone(),
115            inner: Arc::clone(&self.inner),
116        }
117    }
118}
119
120impl GearInstance {
121    /// Build a new instance record using `other`'s metadata but preserving
122    /// `self`'s `inner` runtime-state lock. This makes re-registration atomic:
123    /// concurrent `update_heartbeat` calls continue to write to the same state
124    /// object instead of racing with a copied/swapped `InstanceRuntimeState`.
125    fn with_metadata_of(&self, other: &GearInstance) -> GearInstance {
126        GearInstance {
127            gear: other.gear.clone(),
128            instance_id: other.instance_id,
129            control: other.control.clone(),
130            grpc_services: other.grpc_services.clone(),
131            version: other.version.clone(),
132            rest_endpoint: other.rest_endpoint.clone(),
133            openapi_spec: other.openapi_spec.clone(),
134            // Empty incoming labels mean "keep the stored set", not "clear"; a
135            // non-empty set replaces it. See `RegisterInstanceInfo::with_labels`
136            // for why (self-heal / REST-augmentation re-registers omit labels).
137            labels: if other.labels.is_empty() {
138                self.labels.clone()
139            } else {
140                other.labels.clone()
141            },
142            inner: Arc::clone(&self.inner),
143        }
144    }
145}
146
147impl GearInstance {
148    pub fn new(gear: impl Into<String>, instance_id: Uuid) -> Self {
149        Self {
150            gear: gear.into(),
151            instance_id,
152            control: None,
153            grpc_services: HashMap::new(),
154            version: None,
155            rest_endpoint: None,
156            openapi_spec: None,
157            labels: BTreeMap::new(),
158            inner: Arc::new(parking_lot::RwLock::new(InstanceRuntimeState {
159                last_heartbeat: Instant::now(),
160                state: InstanceState::Registered,
161            })),
162        }
163    }
164
165    pub fn with_control(mut self, ep: Endpoint) -> Self {
166        self.control = Some(ep);
167        self
168    }
169
170    pub fn with_version(mut self, v: impl Into<String>) -> Self {
171        self.version = Some(v.into());
172        self
173    }
174
175    pub fn with_grpc_service(mut self, name: impl Into<String>, ep: Endpoint) -> Self {
176        self.grpc_services.insert(name.into(), ep);
177        self
178    }
179
180    pub fn with_rest_endpoint(mut self, ep: Endpoint) -> Self {
181        self.rest_endpoint = Some(ep);
182        self
183    }
184
185    pub fn with_openapi_spec(mut self, spec: impl Into<String>) -> Self {
186        self.openapi_spec = Some(spec.into());
187        self
188    }
189
190    pub fn with_labels(mut self, labels: BTreeMap<String, String>) -> Self {
191        self.labels = labels;
192        self
193    }
194
195    /// Get the current state of this instance
196    #[must_use]
197    pub fn state(&self) -> InstanceState {
198        self.inner.read().state
199    }
200
201    /// Get the last heartbeat timestamp
202    #[must_use]
203    pub fn last_heartbeat(&self) -> Instant {
204        self.inner.read().last_heartbeat
205    }
206}
207
208/// Central registry that tracks all running gear instances in the system.
209/// Provides discovery, health tracking, and round-robin load balancing.
210///
211/// Always shared as an `Arc<GearManager>` and NOT `Clone`: a by-value clone
212/// would deep-copy the `DashMap` directories while sharing the guards, so the
213/// registration lock would no longer serialize the map it guards.
214#[must_use]
215pub struct GearManager {
216    inner: DashMap<String, Vec<Arc<GearInstance>>>,
217    rr_counters: DashMap<String, usize>,
218    hb_ttl: Duration,
219    hb_grace: Duration,
220    /// Serializes gRPC-service-name ownership check + insert in
221    /// [`register_instance`](Self::register_instance) so two
222    /// gears cannot both pass an "unowned" check and then both commit the same
223    /// name. `DashMap` locks per entry, giving no boundary across the
224    /// cross-gear ownership scan and the single-entry write; this gate provides
225    /// it.
226    ///
227    /// Contract: every operation that *creates* a name→gear binding must hold
228    /// this lock — today `register_instance`, `set_grpc_service_owners`, and
229    /// `merge_authoritative_grpc_service_owners`. `deregister` / `evict_stale`
230    /// only *free* names and `mark_*` / `update_heartbeat` only touch instance
231    /// state, so they cannot create a second owner and need not take it. Any
232    /// future mutator that can bind a name must.
233    reg_lock: parking_lot::Mutex<()>,
234    /// Authoritative gRPC-service-name -> owning-gear map. When a name appears
235    /// here, ownership is fixed to the named gear: only that gear may advertise
236    /// it and no other gear can ever claim it, regardless of which registered
237    /// first — the stable ownership source the dynamic registration scan is
238    /// not. Populated by the runtime at startup (see
239    /// [`set_grpc_service_owners`](Self::set_grpc_service_owners)); names absent
240    /// from the map fall back to first-registration ownership.
241    service_owners: parking_lot::RwLock<HashMap<String, String>>,
242}
243
244/// Returned by [`GearManager::register_instance`] when an instance
245/// advertises a gRPC service name already owned by a *different* gear. The
246/// store is left unchanged; the caller maps this onto a gRPC status (see
247/// `DirectoryServiceNameConflict`), whose code depends on [`recoverable`].
248///
249/// [`recoverable`]: Self::recoverable
250#[derive(Debug, Clone, PartialEq, Eq)]
251pub struct GrpcServiceNameConflict {
252    /// The gRPC service name that is already owned.
253    pub service_name: String,
254    /// The gear that currently owns `service_name`.
255    pub owner: String,
256    /// Whether waiting could clear the conflict: `false` when the name is pinned
257    /// to `owner` by the authoritative ownership map (permanent), `true` when
258    /// `owner` merely currently advertises it (clears when it deregisters). Set
259    /// at the two rejection sites in [`GearManager::register_instance`].
260    pub recoverable: bool,
261}
262
263impl std::fmt::Display for GrpcServiceNameConflict {
264    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
265        write!(
266            f,
267            "gRPC service name '{}' is already owned by gear '{}'",
268            self.service_name, self.owner
269        )
270    }
271}
272
273impl std::error::Error for GrpcServiceNameConflict {}
274
275impl std::fmt::Debug for GearManager {
276    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
277        let gears: Vec<String> = self.inner.iter().map(|e| e.key().clone()).collect();
278        f.debug_struct("GearManager")
279            .field("instances_count", &self.inner.len())
280            .field("gears", &gears)
281            .field("heartbeat_ttl", &self.hb_ttl)
282            .field("heartbeat_grace", &self.hb_grace)
283            .finish_non_exhaustive()
284    }
285}
286
287impl GearManager {
288    pub fn new() -> Self {
289        Self {
290            inner: DashMap::new(),
291            rr_counters: DashMap::new(),
292            hb_ttl: Duration::from_secs(15),
293            hb_grace: Duration::from_secs(30),
294            reg_lock: parking_lot::Mutex::new(()),
295            service_owners: parking_lot::RwLock::new(HashMap::new()),
296        }
297    }
298
299    pub fn with_heartbeat_policy(mut self, ttl: Duration, grace: Duration) -> Self {
300        self.hb_ttl = ttl;
301        self.hb_grace = grace;
302        self
303    }
304
305    /// Install the authoritative gRPC-service-name -> owning-gear map.
306    ///
307    /// A listed name may be advertised only by its owner (enforced in
308    /// [`register_instance`](Self::register_instance)), closing the hole where a
309    /// gear could squat another's service name by registering first; unlisted
310    /// names keep first-registration ownership. Seeded from operator config, then
311    /// layered with the compiled-in owners via
312    /// [`merge_authoritative_grpc_service_owners`](Self::merge_authoritative_grpc_service_owners) —
313    /// both at startup, before any gear self-registers.
314    ///
315    /// Taken under `reg_lock` so it can't interleave with an in-flight
316    /// [`register_instance`](Self::register_instance) check + commit.
317    pub fn set_grpc_service_owners(&self, owners: HashMap<String, String>) {
318        let _gate = self.reg_lock.lock();
319        *self.service_owners.write() = owners;
320    }
321
322    /// Merge a compiled-in `service_name -> gear` map into the ownership map,
323    /// overriding any conflicting operator-configured entry.
324    ///
325    /// The compiled binary is ground truth — a gear physically *provides* the
326    /// service — so config must not reassign that name elsewhere (that would
327    /// re-open the squat); a disagreeing entry is overridden and warned. Names
328    /// absent here (e.g. remote / out-of-process gears) keep their configured
329    /// owner.
330    ///
331    /// Runs once in the runtime's gRPC phase before self-registration, under
332    /// `reg_lock` like [`set_grpc_service_owners`](Self::set_grpc_service_owners).
333    pub fn merge_authoritative_grpc_service_owners(&self, authoritative: HashMap<String, String>) {
334        let _gate = self.reg_lock.lock();
335        let mut owners = self.service_owners.write();
336        for (service_name, gear) in authoritative {
337            if let Some(configured) = owners.get(&service_name)
338                && configured != &gear
339            {
340                tracing::warn!(
341                    service = %service_name,
342                    configured_owner = %configured,
343                    compiled_owner = %gear,
344                    "grpc-service-owners: operator config assigned a compiled-in service name to a \
345                     different gear; the compiled-in provider is authoritative and overrides it"
346                );
347            }
348            owners.insert(service_name, gear);
349        }
350    }
351
352    /// Register or update a gear instance, enforcing single-gear ownership of
353    /// every gRPC service name it advertises — atomically.
354    ///
355    /// The ownership check ([`check_grpc_service_ownership`](Self::check_grpc_service_ownership))
356    /// and the insert run under one registration lock, so two gears cannot both
357    /// pass an "unowned" check and then both commit the same name — a
358    /// check-then-write race the per-entry `DashMap` locks do not prevent. On a
359    /// conflict the store is left untouched.
360    ///
361    /// Re-registering an existing `instance_id` is an idempotent endpoint /
362    /// metadata refresh, NOT a liveness reset: the existing runtime state
363    /// (`last_heartbeat` + [`InstanceState`]) is carried over onto the new
364    /// record by preserving the same `Arc<RwLock<InstanceRuntimeState>>`.
365    /// Without this, a periodic self-heal re-registration (see
366    /// `oop_registration::presence_loop`) would knock a `Healthy` instance back
367    /// to `Registered`, dropping it out of gRPC round-robin until the next
368    /// heartbeat. It would also race with concurrent `update_heartbeat` calls.
369    ///
370    /// # Errors
371    /// Returns [`GrpcServiceNameConflict`] if an advertised gRPC service name is
372    /// owned by, or currently advertised by, another gear.
373    pub fn register_instance(
374        &self,
375        instance: Arc<GearInstance>,
376    ) -> Result<(), GrpcServiceNameConflict> {
377        let _gate = self.reg_lock.lock();
378        self.check_grpc_service_ownership(&instance)?;
379
380        let gear = instance.gear.clone();
381        let mut vec = self.inner.entry(gear).or_default();
382        // replace by instance_id if it already exists
383        if let Some(pos) = vec
384            .iter()
385            .position(|i| i.instance_id == instance.instance_id)
386        {
387            vec[pos] = Arc::new(vec[pos].with_metadata_of(&instance));
388        } else {
389            vec.push(instance);
390        }
391        Ok(())
392    }
393
394    /// Authorize every gRPC service name `instance` advertises, resolving
395    /// ownership in two tiers:
396    ///
397    /// 1. **Declared (authoritative).** If a name appears in the configured
398    ///    ownership map (see
399    ///    [`set_grpc_service_owners`](Self::set_grpc_service_owners)), only the
400    ///    declared owner may advertise it — any other gear is rejected even if
401    ///    it registers first. This stops a gear from squatting a name it was
402    ///    never assigned.
403    /// 2. **First-registration (fallback).** For names absent from that map, the
404    ///    first gear to register the name owns it.
405    ///
406    /// Authorization is necessary but not sufficient: a name currently
407    /// advertised by a *different* gear always conflicts, so even the declared
408    /// owner is refused while a stale advertiser still holds it (e.g. one that
409    /// registered under a previous ownership map before
410    /// [`set_grpc_service_owners`](Self::set_grpc_service_owners) reassigned the
411    /// name). This keeps ownership from splitting across two advertisers; the
412    /// declared owner succeeds once the stale holder deregisters or expires. A
413    /// name held by the *same* gear is always fine (re-register / additional
414    /// instance).
415    ///
416    /// Must be called with `reg_lock` held so the check stays atomic with the
417    /// caller's subsequent insert.
418    fn check_grpc_service_ownership(
419        &self,
420        instance: &GearInstance,
421    ) -> Result<(), GrpcServiceNameConflict> {
422        // Tier 1: reject a declared name claimed by a non-owner; the rest (the
423        // declared owner, and undeclared names) still need the tier-2 scan.
424        let declared = self.service_owners.read();
425        let mut to_scan: Vec<&str> = Vec::new();
426        for service_name in instance.grpc_services.keys() {
427            match declared.get(service_name) {
428                Some(owner) if *owner != instance.gear => {
429                    // Tier 1: the name is pinned to another gear by the
430                    // authoritative map — permanent, waiting cannot clear it.
431                    return Err(GrpcServiceNameConflict {
432                        service_name: service_name.clone(),
433                        owner: owner.clone(),
434                        recoverable: false,
435                    });
436                }
437                _ => to_scan.push(service_name),
438            }
439        }
440        drop(declared);
441        if to_scan.is_empty() {
442            return Ok(());
443        }
444
445        // Tier 2: one pass over the directory, checking every candidate name per
446        // instance visited rather than re-scanning once per name. A name already
447        // held by a *different* gear conflicts.
448        for entry in &self.inner {
449            if entry.key() == &instance.gear {
450                continue; // same gear: re-register / additional instance is fine
451            }
452            for existing in entry.value() {
453                if let Some(&name) = to_scan
454                    .iter()
455                    .find(|&&n| existing.grpc_services.contains_key(n))
456                {
457                    // Tier 2: a *different* gear currently advertises the name —
458                    // recoverable, it clears when that holder deregisters/expires.
459                    return Err(GrpcServiceNameConflict {
460                        service_name: name.to_owned(),
461                        owner: entry.key().clone(),
462                        recoverable: true,
463                    });
464                }
465            }
466        }
467        Ok(())
468    }
469
470    /// Mark an instance as ready
471    pub fn mark_ready(&self, gear: &str, instance_id: Uuid) {
472        if let Some(mut vec) = self.inner.get_mut(gear)
473            && let Some(inst) = vec.iter_mut().find(|i| i.instance_id == instance_id)
474        {
475            let mut state = inst.inner.write();
476            state.state = InstanceState::Ready;
477        }
478    }
479
480    /// Update the heartbeat timestamp for an instance
481    pub fn update_heartbeat(&self, gear: &str, instance_id: Uuid, at: Instant) {
482        if let Some(mut vec) = self.inner.get_mut(gear)
483            && let Some(inst) = vec.iter_mut().find(|i| i.instance_id == instance_id)
484        {
485            let mut state = inst.inner.write();
486            state.last_heartbeat = at;
487            // Transition Registered -> Healthy on first heartbeat
488            if state.state == InstanceState::Registered {
489                state.state = InstanceState::Healthy;
490            }
491        }
492    }
493
494    /// Mark an instance as quarantined
495    pub fn mark_quarantined(&self, gear: &str, instance_id: Uuid) {
496        if let Some(mut vec) = self.inner.get_mut(gear)
497            && let Some(inst) = vec.iter_mut().find(|i| i.instance_id == instance_id)
498        {
499            inst.inner.write().state = InstanceState::Quarantined;
500        }
501    }
502
503    /// Mark an instance as draining (graceful shutdown in progress)
504    pub fn mark_draining(&self, gear: &str, instance_id: Uuid) {
505        if let Some(mut vec) = self.inner.get_mut(gear)
506            && let Some(inst) = vec.iter_mut().find(|i| i.instance_id == instance_id)
507        {
508            inst.inner.write().state = InstanceState::Draining;
509        }
510    }
511
512    /// Remove an instance from the directory
513    pub fn deregister(&self, gear: &str, instance_id: Uuid) {
514        let mut remove_gear = false;
515        {
516            if let Some(mut vec) = self.inner.get_mut(gear) {
517                let list = vec.value_mut();
518                list.retain(|inst| inst.instance_id != instance_id);
519                if list.is_empty() {
520                    remove_gear = true;
521                }
522            }
523        }
524
525        if remove_gear {
526            self.inner.remove(gear);
527            self.rr_counters.remove(gear);
528            self.rr_counters.remove(&format!("rest:{gear}"));
529        }
530    }
531
532    /// Get all instances of a specific gear
533    #[must_use]
534    pub fn instances_of(&self, gear: &str) -> Vec<Arc<GearInstance>> {
535        self.inner.get(gear).map(|v| v.clone()).unwrap_or_default()
536    }
537
538    /// Get all instances across all gears
539    #[must_use]
540    pub fn all_instances(&self) -> Vec<Arc<GearInstance>> {
541        self.inner
542            .iter()
543            .flat_map(|entry| entry.value().clone())
544            .collect()
545    }
546
547    /// Quarantine or evict stale instances based on heartbeat policy
548    pub fn evict_stale(&self, now: Instant) {
549        use InstanceState::{Draining, Quarantined};
550        let mut empty_gears = Vec::new();
551
552        for mut entry in self.inner.iter_mut() {
553            let gear = entry.key().clone();
554            let vec = entry.value_mut();
555            vec.retain(|inst| {
556                let state = inst.inner.read();
557                let age = now.saturating_duration_since(state.last_heartbeat);
558
559                // Quarantine instances that have exceeded TTL
560                if age >= self.hb_ttl && !matches!(state.state, Quarantined | Draining) {
561                    drop(state); // Release read lock before write
562                    inst.inner.write().state = Quarantined;
563                    return true; // Keep quarantined instances for now
564                }
565
566                // Evict quarantined instances that exceed grace period
567                if state.state == Quarantined && age >= self.hb_ttl + self.hb_grace {
568                    return false; // Remove from directory
569                }
570
571                true
572            });
573
574            if vec.is_empty() {
575                empty_gears.push(gear);
576            }
577        }
578
579        for gear in empty_gears {
580            self.inner.remove(&gear);
581            self.rr_counters.remove(&gear);
582            self.rr_counters.remove(&format!("rest:{gear}"));
583        }
584    }
585
586    /// From a candidate set, return the serving subset (`Healthy`/`Ready`); if
587    /// none are serving, return the whole set so resolution falls back to the
588    /// not-ready instances (`Registered`/`Quarantined`/`Draining`) rather than
589    /// failing closed. `context` names the pool for the cold-path `debug!`.
590    fn prefer_serving(candidates: Vec<Arc<GearInstance>>, context: &str) -> Vec<Arc<GearInstance>> {
591        let serving: Vec<Arc<GearInstance>> = candidates
592            .iter()
593            .filter(|inst| matches!(inst.state(), InstanceState::Healthy | InstanceState::Ready))
594            .cloned()
595            .collect();
596        if serving.is_empty() {
597            tracing::debug!(
598                context,
599                "no serving (Ready/Healthy) instance available; round-robining over the \
600                 not-ready set instead of returning None"
601            );
602            candidates
603        } else {
604            serving
605        }
606    }
607
608    /// Pick an instance using round-robin selection, preferring healthy instances
609    #[must_use]
610    pub fn pick_instance_round_robin(&self, gear: &str) -> Option<Arc<GearInstance>> {
611        let instances_entry = self.inner.get(gear)?;
612        let instances = instances_entry.value();
613
614        if instances.is_empty() {
615            return None;
616        }
617
618        // Non-empty in => non-empty out, so `candidates` is guaranteed non-empty.
619        let candidates = Self::prefer_serving(instances.clone(), gear);
620
621        let len = candidates.len();
622        let mut counter = self.rr_counters.entry(gear.to_owned()).or_insert(0);
623        let idx = *counter % len;
624        *counter = (*counter + 1) % len;
625
626        candidates.get(idx).cloned()
627    }
628
629    /// Pick a service endpoint using round-robin, returning (gear, instance, endpoint).
630    /// Prefers healthy/ready instances and automatically rotates among them.
631    #[must_use]
632    pub fn pick_service_round_robin(
633        &self,
634        service_name: &str,
635    ) -> Option<(String, Arc<GearInstance>, Endpoint)> {
636        // Collect the `Arc` of every instance that provides this service, then
637        // apply the shared serving-preference fallback so gRPC resolution
638        // survives an all-not-ready gear instead of failing closed. Cloning an
639        // `Arc` is a refcount bump; the gear name and endpoint are cloned once,
640        // for the winner, after the index is chosen.
641        let mut providing: Vec<Arc<GearInstance>> = Vec::new();
642        for entry in &self.inner {
643            for inst in entry.value() {
644                if inst.grpc_services.contains_key(service_name) {
645                    providing.push(inst.clone());
646                }
647            }
648        }
649
650        if providing.is_empty() {
651            return None;
652        }
653
654        let mut candidates = Self::prefer_serving(providing, service_name);
655
656        // Use a counter keyed by service name for round-robin
657        let len = candidates.len();
658        let service_key = service_name.to_owned();
659        let mut counter = self.rr_counters.entry(service_key).or_insert(0);
660        let idx = *counter % len;
661        *counter = (*counter + 1) % len;
662
663        let inst = candidates.swap_remove(idx);
664        let endpoint = inst.grpc_services.get(service_name)?.clone();
665        let gear = inst.gear.clone();
666        Some((gear, inst, endpoint))
667    }
668
669    /// Resolve a REST endpoint for a gear using round-robin over instances that
670    /// expose one, preferring healthy/ready instances.
671    #[must_use]
672    pub fn pick_rest_endpoint_round_robin(&self, gear: &str) -> Option<Endpoint> {
673        let instances_entry = self.inner.get(gear)?;
674        let instances = instances_entry.value();
675
676        // Only instances that actually expose a REST endpoint are candidates.
677        let with_rest: Vec<_> = instances
678            .iter()
679            .filter(|inst| inst.rest_endpoint.is_some())
680            .cloned()
681            .collect();
682
683        if with_rest.is_empty() {
684            return None;
685        }
686
687        let candidates = Self::prefer_serving(with_rest, gear);
688
689        let len = candidates.len();
690        let rr_key = format!("rest:{gear}");
691        let mut counter = self.rr_counters.entry(rr_key).or_insert(0);
692        let idx = *counter % len;
693        *counter = (*counter + 1) % len;
694
695        candidates
696            .get(idx)
697            .and_then(|inst| inst.rest_endpoint.clone())
698    }
699
700    /// Return a gear that currently advertises the gRPC service `service_name`,
701    /// if any (any instance state counts); in the misconfigured multi-owner case
702    /// the first match in iteration order.
703    ///
704    /// Test-only: [`register_instance`](Self::register_instance) enforces single
705    /// ownership with its own inline scan, so this is not part of the public
706    /// surface — it exists purely to let the tests below assert who owns a name.
707    #[cfg(test)]
708    #[must_use]
709    fn grpc_service_owner(&self, service_name: &str) -> Option<String> {
710        for entry in &self.inner {
711            if entry
712                .value()
713                .iter()
714                .any(|inst| inst.grpc_services.contains_key(service_name))
715            {
716                return Some(entry.key().clone());
717            }
718        }
719        None
720    }
721
722    /// Retrieve the `OpenAPI` spec of a gear, taken from the first registered
723    /// instance that published one.
724    #[must_use]
725    pub fn openapi_spec_of(&self, gear: &str) -> Option<String> {
726        let instances_entry = self.inner.get(gear)?;
727        instances_entry
728            .value()
729            .iter()
730            .find_map(|inst| inst.openapi_spec.clone())
731    }
732}
733
734impl Default for GearManager {
735    fn default() -> Self {
736        Self::new()
737    }
738}
739
740#[cfg(test)]
741#[cfg_attr(coverage_nightly, coverage(off))]
742mod tests {
743    use super::*;
744    use std::thread::sleep;
745    use std::time::Duration;
746
747    #[test]
748    fn test_register_and_retrieve_instances() {
749        let dir = GearManager::new();
750        let instance_id = Uuid::new_v4();
751        let instance = Arc::new(
752            GearInstance::new("test_gear", instance_id)
753                .with_control(Endpoint::http("localhost", 8080))
754                .with_version("1.0.0"),
755        );
756
757        seed(&dir, instance);
758
759        let instances = dir.instances_of("test_gear");
760        assert_eq!(instances.len(), 1);
761        assert_eq!(instances[0].instance_id, instance_id);
762        assert_eq!(instances[0].gear, "test_gear");
763        assert_eq!(instances[0].version, Some("1.0.0".to_owned()));
764    }
765
766    /// Seed the registry in a test. The fixture is assumed to advertise no
767    /// conflicting gRPC service name, so a conflict is a bug in the test, not
768    /// an expected outcome — surface it loudly rather than swallowing it.
769    #[track_caller]
770    fn seed(mgr: &GearManager, instance: Arc<GearInstance>) {
771        mgr.register_instance(instance)
772            .expect("test fixture must not create a gRPC service-name conflict");
773    }
774
775    fn labels(pairs: &[(&str, &str)]) -> BTreeMap<String, String> {
776        pairs
777            .iter()
778            .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
779            .collect()
780    }
781
782    #[test]
783    fn reregister_without_labels_preserves_stored_labels() {
784        let dir = GearManager::new();
785        let instance_id = Uuid::new_v4();
786
787        // Initial registration carries a shard label.
788        seed(
789            &dir,
790            Arc::new(
791                GearInstance::new("shard-gear", instance_id).with_labels(labels(&[("shard", "7")])),
792            ),
793        );
794
795        // A label-less re-registration (e.g. periodic self-heal / REST
796        // augmentation) must NOT wipe the stored labels.
797        seed(
798            &dir,
799            Arc::new(GearInstance::new("shard-gear", instance_id).with_version("2.0.0")),
800        );
801
802        let registered = dir.instances_of("shard-gear");
803        assert_eq!(registered.len(), 1);
804        assert_eq!(
805            registered[0].labels.get("shard"),
806            Some(&"7".to_owned()),
807            "label-less re-registration must preserve stored labels"
808        );
809        assert_eq!(registered[0].version, Some("2.0.0".to_owned()));
810    }
811
812    #[test]
813    fn reregister_with_labels_replaces_stored_labels() {
814        let dir = GearManager::new();
815        let instance_id = Uuid::new_v4();
816
817        seed(
818            &dir,
819            Arc::new(
820                GearInstance::new("shard-gear", instance_id).with_labels(labels(&[("shard", "7")])),
821            ),
822        );
823        // An explicit non-empty set replaces the stored one wholesale.
824        seed(
825            &dir,
826            Arc::new(
827                GearInstance::new("shard-gear", instance_id).with_labels(labels(&[("shard", "8")])),
828            ),
829        );
830
831        let registered = dir.instances_of("shard-gear");
832        assert_eq!(registered.len(), 1);
833        assert_eq!(registered[0].labels.get("shard"), Some(&"8".to_owned()));
834    }
835
836    #[test]
837    fn test_register_multiple_instances() {
838        let dir = GearManager::new();
839
840        let id1 = Uuid::new_v4();
841        let id2 = Uuid::new_v4();
842        let instance1 = Arc::new(GearInstance::new("test_gear", id1));
843        let instance2 = Arc::new(GearInstance::new("test_gear", id2));
844
845        seed(&dir, instance1);
846        seed(&dir, instance2);
847
848        let registered = dir.instances_of("test_gear");
849        assert_eq!(registered.len(), 2);
850
851        let ids: Vec<_> = registered.iter().map(|i| i.instance_id).collect();
852        assert!(ids.contains(&id1));
853        assert!(ids.contains(&id2));
854    }
855
856    #[test]
857    fn test_update_existing_instance() {
858        let dir = GearManager::new();
859        let instance_id = Uuid::new_v4();
860
861        let initial_instance =
862            Arc::new(GearInstance::new("test_gear", instance_id).with_version("1.0.0"));
863        seed(&dir, initial_instance);
864
865        let updated_instance =
866            Arc::new(GearInstance::new("test_gear", instance_id).with_version("2.0.0"));
867        seed(&dir, updated_instance);
868
869        let registered = dir.instances_of("test_gear");
870        assert_eq!(registered.len(), 1, "Should not duplicate instance");
871        assert_eq!(registered[0].version, Some("2.0.0".to_owned()));
872    }
873
874    #[test]
875    fn test_reregistration_preserves_liveness_state() {
876        let dir = GearManager::new();
877        let instance_id = Uuid::new_v4();
878
879        // Register, then heartbeat so the instance is Healthy (routable).
880        seed(
881            &dir,
882            Arc::new(GearInstance::new("test_gear", instance_id).with_version("1.0.0")),
883        );
884        dir.update_heartbeat("test_gear", instance_id, Instant::now());
885        assert!(matches!(
886            dir.instances_of("test_gear")[0].state(),
887            InstanceState::Healthy
888        ));
889
890        // A periodic self-heal re-registration must NOT reset liveness back to
891        // Registered — otherwise the instance drops out of gRPC round-robin
892        // until the next heartbeat (the "split-brain" flap).
893        seed(
894            &dir,
895            Arc::new(GearInstance::new("test_gear", instance_id).with_version("2.0.0")),
896        );
897
898        let instances = dir.instances_of("test_gear");
899        assert_eq!(instances.len(), 1);
900        assert!(
901            matches!(instances[0].state(), InstanceState::Healthy),
902            "re-registration must preserve the Healthy state"
903        );
904        assert_eq!(
905            instances[0].version,
906            Some("2.0.0".to_owned()),
907            "re-registration must still refresh metadata/endpoints"
908        );
909    }
910
911    #[test]
912    fn test_concurrent_reregister_and_heartbeat_preserves_state() {
913        let dir = GearManager::new();
914        let instance_id = Uuid::new_v4();
915
916        let initial = Arc::new(GearInstance::new("test_gear", instance_id).with_version("1.0.0"));
917        seed(&dir, initial);
918        dir.update_heartbeat("test_gear", instance_id, Instant::now());
919        assert!(matches!(
920            dir.instances_of("test_gear")[0].state(),
921            InstanceState::Healthy
922        ));
923
924        let start = Instant::now();
925
926        std::thread::scope(|s| {
927            s.spawn(|| {
928                for _ in 0..1000 {
929                    dir.update_heartbeat("test_gear", instance_id, Instant::now());
930                }
931            });
932            s.spawn(|| {
933                for i in 0..1000 {
934                    let version = if i % 2 == 0 { "2.0.0" } else { "3.0.0" };
935                    let reinst = Arc::new(
936                        GearInstance::new("test_gear", instance_id)
937                            .with_version(version)
938                            .with_rest_endpoint(Endpoint::http(
939                                "127.0.0.1",
940                                8000u16 + u16::try_from(i % 10).expect("i % 10 fits in u16"),
941                            )),
942                    );
943                    seed(&dir, reinst);
944                }
945            });
946        });
947
948        let instances = dir.instances_of("test_gear");
949        assert_eq!(instances.len(), 1);
950        assert!(
951            matches!(instances[0].state(), InstanceState::Healthy),
952            "concurrent re-registration must not reset Healthy state"
953        );
954        assert!(
955            instances[0].last_heartbeat() >= start,
956            "concurrent re-registration must not lose heartbeat updates"
957        );
958    }
959
960    #[test]
961    fn test_mark_ready() {
962        let dir = GearManager::new();
963        let instance_id = Uuid::new_v4();
964        let instance = Arc::new(GearInstance::new("test_gear", instance_id));
965
966        seed(&dir, instance);
967
968        dir.mark_ready("test_gear", instance_id);
969
970        let instances = dir.instances_of("test_gear");
971        assert_eq!(instances.len(), 1);
972        assert!(matches!(instances[0].state(), InstanceState::Ready));
973    }
974
975    #[test]
976    fn test_update_heartbeat() {
977        let dir = GearManager::new();
978        let instance_id = Uuid::new_v4();
979        let instance = Arc::new(GearInstance::new("test_gear", instance_id));
980        let initial_heartbeat = instance.last_heartbeat();
981
982        seed(&dir, instance);
983
984        // Sleep to ensure time difference
985        sleep(Duration::from_millis(10));
986
987        let new_heartbeat = Instant::now();
988        dir.update_heartbeat("test_gear", instance_id, new_heartbeat);
989
990        let instances = dir.instances_of("test_gear");
991        assert!(instances[0].last_heartbeat() > initial_heartbeat);
992        assert!(matches!(instances[0].state(), InstanceState::Healthy));
993    }
994
995    #[test]
996    fn test_all_instances() {
997        let dir = GearManager::new();
998
999        let instance1 = Arc::new(GearInstance::new("gear_a", Uuid::new_v4()));
1000        let instance2 = Arc::new(GearInstance::new("gear_b", Uuid::new_v4()));
1001        let instance3 = Arc::new(GearInstance::new("gear_a", Uuid::new_v4()));
1002
1003        seed(&dir, instance1);
1004        seed(&dir, instance2);
1005        seed(&dir, instance3);
1006
1007        let all = dir.all_instances();
1008        assert_eq!(all.len(), 3);
1009
1010        let gears: Vec<_> = all.iter().map(|i| i.gear.as_str()).collect();
1011        assert_eq!(gears.iter().filter(|&m| *m == "gear_a").count(), 2);
1012        assert_eq!(gears.iter().filter(|&m| *m == "gear_b").count(), 1);
1013    }
1014
1015    #[test]
1016    fn test_pick_instance_round_robin() {
1017        let dir = GearManager::new();
1018
1019        let id1 = Uuid::new_v4();
1020        let id2 = Uuid::new_v4();
1021        let instance1 = Arc::new(GearInstance::new("test_gear", id1));
1022        let instance2 = Arc::new(GearInstance::new("test_gear", id2));
1023
1024        seed(&dir, instance1);
1025        seed(&dir, instance2);
1026
1027        // Pick three times to verify round-robin behavior
1028        let picked1 = dir.pick_instance_round_robin("test_gear").unwrap();
1029        let picked2 = dir.pick_instance_round_robin("test_gear").unwrap();
1030        let picked3 = dir.pick_instance_round_robin("test_gear").unwrap();
1031
1032        let ids = [
1033            picked1.instance_id,
1034            picked2.instance_id,
1035            picked3.instance_id,
1036        ];
1037
1038        // With 2 instances, we expect round-robin pattern like A, B, A
1039        // Check that both instance IDs appear and that at least one repeats
1040        assert!(ids.contains(&id1));
1041        assert!(ids.contains(&id2));
1042        // First and third pick should be the same (round-robin wraps)
1043        assert_eq!(picked1.instance_id, picked3.instance_id);
1044        // Second pick should be different from the first
1045        assert_ne!(picked1.instance_id, picked2.instance_id);
1046    }
1047
1048    #[test]
1049    fn test_pick_instance_none_available() {
1050        let dir = GearManager::new();
1051        let picked = dir.pick_instance_round_robin("nonexistent_gear");
1052        assert!(picked.is_none());
1053    }
1054
1055    #[test]
1056    fn test_endpoint_creation() {
1057        let plain_ep = Endpoint::http("localhost", 8080);
1058        assert_eq!(plain_ep.uri, "http://localhost:8080");
1059
1060        let secure_ep = Endpoint::https("localhost", 8443);
1061        assert_eq!(secure_ep.uri, "https://localhost:8443");
1062
1063        let uds_ep = Endpoint::uds("/tmp/socket.sock");
1064        assert!(uds_ep.uri.starts_with("unix://"));
1065        assert!(uds_ep.uri.contains("socket.sock"));
1066
1067        let custom_ep = Endpoint::from_uri("http://example.com");
1068        assert_eq!(custom_ep.uri, "http://example.com");
1069    }
1070
1071    #[test]
1072    fn test_endpoint_kind() {
1073        let plain_ep = Endpoint::http("127.0.0.1", 8080);
1074        match plain_ep.kind() {
1075            EndpointKind::Tcp(addr) => {
1076                assert_eq!(addr.ip().to_string(), "127.0.0.1");
1077                assert_eq!(addr.port(), 8080);
1078            }
1079            _ => panic!("Expected TCP endpoint for http"),
1080        }
1081
1082        let secure_ep = Endpoint::https("127.0.0.1", 8443);
1083        match secure_ep.kind() {
1084            EndpointKind::Tcp(addr) => {
1085                assert_eq!(addr.ip().to_string(), "127.0.0.1");
1086                assert_eq!(addr.port(), 8443);
1087            }
1088            _ => panic!("Expected TCP endpoint for https"),
1089        }
1090
1091        let uds_ep = Endpoint::uds("/tmp/test.sock");
1092        match uds_ep.kind() {
1093            EndpointKind::Uds(path) => {
1094                assert!(path.to_string_lossy().contains("test.sock"));
1095            }
1096            _ => panic!("Expected UDS endpoint"),
1097        }
1098
1099        let other_ep = Endpoint::from_uri("grpc://example.com");
1100        match other_ep.kind() {
1101            EndpointKind::Other(uri) => {
1102                assert_eq!(uri, "grpc://example.com");
1103            }
1104            _ => panic!("Expected Other endpoint"),
1105        }
1106    }
1107
1108    #[test]
1109    fn test_gear_instance_builder() {
1110        let instance_id = Uuid::new_v4();
1111        let instance = GearInstance::new("test_gear", instance_id)
1112            .with_control(Endpoint::http("localhost", 8080))
1113            .with_version("1.2.3")
1114            .with_grpc_service("service1", Endpoint::http("localhost", 8082))
1115            .with_grpc_service("service2", Endpoint::http("localhost", 8083));
1116
1117        assert_eq!(instance.gear, "test_gear");
1118        assert_eq!(instance.instance_id, instance_id);
1119        assert!(instance.control.is_some());
1120        assert_eq!(instance.version, Some("1.2.3".to_owned()));
1121        assert_eq!(instance.grpc_services.len(), 2);
1122        assert!(instance.grpc_services.contains_key("service1"));
1123        assert!(instance.grpc_services.contains_key("service2"));
1124        assert!(matches!(instance.state(), InstanceState::Registered));
1125    }
1126
1127    #[test]
1128    fn test_quarantine_and_evict() {
1129        let ttl = Duration::from_millis(50);
1130        let grace = Duration::from_millis(50);
1131        let dir = GearManager::new().with_heartbeat_policy(ttl, grace);
1132
1133        let now = Instant::now();
1134        let instance = GearInstance::new("test_gear", Uuid::new_v4());
1135        // Set the last heartbeat to be stale
1136        instance.inner.write().last_heartbeat = now
1137            .checked_sub(ttl)
1138            .and_then(|t| t.checked_sub(Duration::from_millis(10)))
1139            .expect("test duration subtraction should not underflow");
1140
1141        seed(&dir, Arc::new(instance));
1142
1143        dir.evict_stale(now);
1144        let instances = dir.instances_of("test_gear");
1145        assert_eq!(instances.len(), 1);
1146        assert!(matches!(instances[0].state(), InstanceState::Quarantined));
1147
1148        let later = now + grace + Duration::from_millis(10);
1149        dir.evict_stale(later);
1150
1151        let instances_after = dir.instances_of("test_gear");
1152        assert!(instances_after.is_empty());
1153    }
1154
1155    #[test]
1156    fn test_instances_of_empty() {
1157        let dir = GearManager::new();
1158        let instances = dir.instances_of("nonexistent");
1159        assert!(instances.is_empty());
1160    }
1161
1162    #[test]
1163    fn test_rr_prefers_healthy() {
1164        let dir = GearManager::new();
1165
1166        // Create two instances: one healthy, one quarantined
1167        let healthy_id = Uuid::new_v4();
1168        let healthy = Arc::new(GearInstance::new("test_gear", healthy_id));
1169        seed(&dir, healthy);
1170        dir.update_heartbeat("test_gear", healthy_id, Instant::now());
1171
1172        let quarantined_id = Uuid::new_v4();
1173        let quarantined = Arc::new(GearInstance::new("test_gear", quarantined_id));
1174        seed(&dir, quarantined);
1175        dir.mark_quarantined("test_gear", quarantined_id);
1176
1177        // RR should only pick the healthy instance
1178        for _ in 0..5 {
1179            let picked = dir.pick_instance_round_robin("test_gear").unwrap();
1180            assert_eq!(picked.instance_id, healthy_id);
1181        }
1182    }
1183
1184    #[test]
1185    fn test_pick_rest_endpoint_and_openapi() {
1186        let dir = GearManager::new();
1187        let id = Uuid::new_v4();
1188        let instance = Arc::new(
1189            GearInstance::new("billing", id)
1190                .with_rest_endpoint(Endpoint::http("billing", 8080))
1191                .with_openapi_spec("{\"openapi\":\"3.1.0\"}"),
1192        );
1193        seed(&dir, instance);
1194
1195        let rest = dir.pick_rest_endpoint_round_robin("billing").unwrap();
1196        assert_eq!(rest.uri, "http://billing:8080");
1197
1198        let spec = dir.openapi_spec_of("billing").unwrap();
1199        assert!(spec.contains("openapi"));
1200    }
1201
1202    #[test]
1203    fn test_pick_rest_endpoint_none_when_absent() {
1204        let dir = GearManager::new();
1205        let id = Uuid::new_v4();
1206        // Instance exposes only a gRPC service, no REST endpoint.
1207        let instance = Arc::new(
1208            GearInstance::new("grpc_only", id)
1209                .with_grpc_service("some.Service", Endpoint::http("127.0.0.1", 9000)),
1210        );
1211        seed(&dir, instance);
1212
1213        assert!(dir.pick_rest_endpoint_round_robin("grpc_only").is_none());
1214        assert!(dir.openapi_spec_of("grpc_only").is_none());
1215        assert!(dir.pick_rest_endpoint_round_robin("missing").is_none());
1216    }
1217
1218    #[test]
1219    fn test_pick_rest_endpoint_round_robin_rotates() {
1220        let dir = GearManager::new();
1221        let id1 = Uuid::new_v4();
1222        let id2 = Uuid::new_v4();
1223        let inst1 = Arc::new(
1224            GearInstance::new("web", id1).with_rest_endpoint(Endpoint::http("127.0.0.1", 8001)),
1225        );
1226        let inst2 = Arc::new(
1227            GearInstance::new("web", id2).with_rest_endpoint(Endpoint::http("127.0.0.1", 8002)),
1228        );
1229        seed(&dir, inst1);
1230        seed(&dir, inst2);
1231        dir.update_heartbeat("web", id1, Instant::now());
1232        dir.update_heartbeat("web", id2, Instant::now());
1233
1234        let ep1 = dir.pick_rest_endpoint_round_robin("web").unwrap();
1235        let ep2 = dir.pick_rest_endpoint_round_robin("web").unwrap();
1236        let ep3 = dir.pick_rest_endpoint_round_robin("web").unwrap();
1237
1238        assert_ne!(ep1, ep2);
1239        assert_eq!(ep1, ep3);
1240    }
1241
1242    #[test]
1243    fn test_pick_service_round_robin() {
1244        let dir = GearManager::new();
1245
1246        let id1 = Uuid::new_v4();
1247        let id2 = Uuid::new_v4();
1248        // Register two instances providing the same service
1249        let inst1 = Arc::new(
1250            GearInstance::new("test_gear", id1)
1251                .with_grpc_service("test.Service", Endpoint::http("127.0.0.1", 8001)),
1252        );
1253        let inst2 = Arc::new(
1254            GearInstance::new("test_gear", id2)
1255                .with_grpc_service("test.Service", Endpoint::http("127.0.0.1", 8002)),
1256        );
1257
1258        seed(&dir, inst1);
1259        seed(&dir, inst2);
1260
1261        // Mark both as healthy
1262        dir.update_heartbeat("test_gear", id1, Instant::now());
1263        dir.update_heartbeat("test_gear", id2, Instant::now());
1264
1265        // Pick should rotate between instances
1266        let pick1 = dir.pick_service_round_robin("test.Service");
1267        let pick2 = dir.pick_service_round_robin("test.Service");
1268        let pick3 = dir.pick_service_round_robin("test.Service");
1269
1270        assert!(pick1.is_some());
1271        assert!(pick2.is_some());
1272        assert!(pick3.is_some());
1273
1274        let (_, inst1, ep1) = pick1.unwrap();
1275        let (_, inst2, ep2) = pick2.unwrap();
1276        let (_, inst3, _) = pick3.unwrap();
1277
1278        // First and third should be the same (round-robin)
1279        assert_eq!(inst1.instance_id, inst3.instance_id);
1280        // First and second should be different
1281        assert_ne!(inst1.instance_id, inst2.instance_id);
1282        // Endpoints should differ
1283        assert_ne!(ep1, ep2);
1284    }
1285
1286    #[test]
1287    fn pick_service_falls_back_to_not_ready() {
1288        // A gear whose only service-providing instance is not yet serving
1289        // (Registered, no heartbeat) must still resolve, mirroring the REST and
1290        // instance pickers' not-ready fallback rather than failing closed.
1291        let dir = GearManager::new();
1292        let id = Uuid::new_v4();
1293        seed(
1294            &dir,
1295            Arc::new(
1296                GearInstance::new("worker", id)
1297                    .with_grpc_service("worker.Svc", Endpoint::http("127.0.0.1", 9000)),
1298            ),
1299        );
1300        // No heartbeat: the instance stays Registered (not serving).
1301        assert!(matches!(
1302            dir.instances_of("worker")[0].state(),
1303            InstanceState::Registered
1304        ));
1305
1306        let picked = dir.pick_service_round_robin("worker.Svc");
1307        assert!(
1308            picked.is_some(),
1309            "gRPC service resolution must fall back to the not-ready instance"
1310        );
1311        let (gear, _, ep) = picked.unwrap();
1312        assert_eq!(gear, "worker");
1313        assert_eq!(ep, Endpoint::http("127.0.0.1", 9000));
1314    }
1315
1316    #[test]
1317    fn grpc_service_owner_reports_owning_gear_regardless_of_health() {
1318        let dir = GearManager::new();
1319
1320        // A freshly-registered (not-yet-serving) instance still owns its name.
1321        seed(
1322            &dir,
1323            Arc::new(
1324                GearInstance::new("authz-resolver", Uuid::new_v4()).with_grpc_service(
1325                    "cf.authz.v1.AuthzService",
1326                    Endpoint::http("127.0.0.1", 9000),
1327                ),
1328            ),
1329        );
1330
1331        assert_eq!(
1332            dir.grpc_service_owner("cf.authz.v1.AuthzService")
1333                .as_deref(),
1334            Some("authz-resolver"),
1335            "the advertising gear owns the name even while only Registered"
1336        );
1337        assert!(
1338            dir.grpc_service_owner("unowned.Service").is_none(),
1339            "an unadvertised name has no owner"
1340        );
1341    }
1342
1343    /// The production invariant behind the above: a `Registered`-but-not-healthy
1344    /// owner still blocks another gear's claim of the same name through
1345    /// `register_instance` (ownership does not wait on the owner becoming healthy).
1346    #[test]
1347    fn registered_but_unhealthy_owner_blocks_another_gears_claim() {
1348        let dir = GearManager::new();
1349        let owner_id = Uuid::new_v4();
1350
1351        // Owner claims the name and never heartbeats -> stays Registered.
1352        dir.register_instance(Arc::new(
1353            GearInstance::new("authz-resolver", owner_id).with_grpc_service(
1354                "cf.authz.v1.AuthzService",
1355                Endpoint::http("127.0.0.1", 9000),
1356            ),
1357        ))
1358        .expect("first claim of an unowned name succeeds");
1359        assert!(matches!(
1360            dir.instances_of("authz-resolver")[0].state(),
1361            InstanceState::Registered
1362        ));
1363
1364        // A different gear is rejected via the production path, health regardless.
1365        let conflict = dir
1366            .register_instance(Arc::new(
1367                GearInstance::new("impostor", Uuid::new_v4()).with_grpc_service(
1368                    "cf.authz.v1.AuthzService",
1369                    Endpoint::http("127.0.0.1", 9001),
1370                ),
1371            ))
1372            .unwrap_err();
1373        assert_eq!(conflict.owner, "authz-resolver");
1374    }
1375
1376    #[test]
1377    fn deregister_releases_grpc_service_name_to_another_gear() {
1378        let dir = GearManager::new();
1379        let service = "cf.authz.v1.AuthzService";
1380        let owner_id = Uuid::new_v4();
1381
1382        seed(
1383            &dir,
1384            Arc::new(
1385                GearInstance::new("authz-resolver", owner_id)
1386                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000)),
1387            ),
1388        );
1389
1390        // While the owner is registered, a different gear cannot claim the name.
1391        dir.register_instance(Arc::new(
1392            GearInstance::new("successor", Uuid::new_v4())
1393                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1394        ))
1395        .unwrap_err();
1396
1397        // Deregistering the owner releases the name...
1398        dir.deregister("authz-resolver", owner_id);
1399        assert!(
1400            dir.grpc_service_owner(service).is_none(),
1401            "the name is unowned once its owner deregisters"
1402        );
1403
1404        // ...so a different gear may now claim it.
1405        dir.register_instance(Arc::new(
1406            GearInstance::new("successor", Uuid::new_v4())
1407                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1408        ))
1409        .expect("a released gRPC service name must be claimable by another gear");
1410        assert_eq!(
1411            dir.grpc_service_owner(service).as_deref(),
1412            Some("successor")
1413        );
1414    }
1415
1416    #[test]
1417    fn eviction_releases_grpc_service_name_only_after_full_evict() {
1418        let ttl = Duration::from_millis(50);
1419        let grace = Duration::from_millis(50);
1420        let dir = GearManager::new().with_heartbeat_policy(ttl, grace);
1421        let service = "cf.authz.v1.AuthzService";
1422
1423        let now = Instant::now();
1424        let owner = GearInstance::new("authz-resolver", Uuid::new_v4())
1425            .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000));
1426        // Stale from the start so the first eviction pass quarantines it.
1427        owner.inner.write().last_heartbeat = now
1428            .checked_sub(ttl)
1429            .and_then(|t| t.checked_sub(Duration::from_millis(10)))
1430            .expect("test duration subtraction should not underflow");
1431        seed(&dir, Arc::new(owner));
1432
1433        // First pass only *quarantines* the dead gear — it is not yet evicted,
1434        // so it still owns the name and a successor is rejected. A dead gear
1435        // keeps blocking its gRPC name (regardless of health) until fully
1436        // evicted, i.e. for hb_ttl + hb_grace.
1437        dir.evict_stale(now);
1438        assert!(matches!(
1439            dir.instances_of("authz-resolver")[0].state(),
1440            InstanceState::Quarantined
1441        ));
1442        let conflict = dir
1443            .register_instance(Arc::new(
1444                GearInstance::new("successor", Uuid::new_v4())
1445                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1446            ))
1447            .unwrap_err();
1448        assert_eq!(
1449            conflict.owner, "authz-resolver",
1450            "a quarantined-but-not-evicted owner still holds the name"
1451        );
1452
1453        // Second pass past the grace period fully evicts it, releasing the name
1454        // so another gear may take it over.
1455        dir.evict_stale(now + grace + Duration::from_millis(10));
1456        assert!(
1457            dir.grpc_service_owner(service).is_none(),
1458            "eviction hands the name back once the grace period lapses"
1459        );
1460        dir.register_instance(Arc::new(
1461            GearInstance::new("successor", Uuid::new_v4())
1462                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9002)),
1463        ))
1464        .expect("an evicted owner's gRPC name must be claimable by another gear");
1465        assert_eq!(
1466            dir.grpc_service_owner(service).as_deref(),
1467            Some("successor")
1468        );
1469    }
1470
1471    #[test]
1472    fn test_deregister_clears_rr_counters() {
1473        let dir = GearManager::new();
1474        let id = Uuid::new_v4();
1475        let instance = Arc::new(
1476            GearInstance::new("web", id).with_rest_endpoint(Endpoint::http("127.0.0.1", 8001)),
1477        );
1478        seed(&dir, instance);
1479        dir.update_heartbeat("web", id, Instant::now());
1480
1481        // Exercise both round-robin counters so the keys are created.
1482        assert!(dir.pick_instance_round_robin("web").is_some());
1483        assert!(dir.pick_rest_endpoint_round_robin("web").is_some());
1484
1485        assert!(dir.rr_counters.contains_key("web"));
1486        assert!(dir.rr_counters.contains_key("rest:web"));
1487
1488        dir.deregister("web", id);
1489
1490        assert!(!dir.rr_counters.contains_key("web"));
1491        assert!(!dir.rr_counters.contains_key("rest:web"));
1492    }
1493
1494    #[test]
1495    fn test_evict_stale_clears_rr_counters() {
1496        let ttl = Duration::from_millis(50);
1497        let grace = Duration::from_millis(50);
1498        let dir = GearManager::new().with_heartbeat_policy(ttl, grace);
1499
1500        let now = Instant::now();
1501        let id = Uuid::new_v4();
1502        let instance = Arc::new(
1503            GearInstance::new("web", id).with_rest_endpoint(Endpoint::http("127.0.0.1", 8001)),
1504        );
1505        // Set the last heartbeat to be stale so the instance is quarantined then evicted.
1506        instance.inner.write().last_heartbeat = now
1507            .checked_sub(ttl)
1508            .and_then(|t| t.checked_sub(Duration::from_millis(10)))
1509            .expect("test duration subtraction should not underflow");
1510
1511        seed(&dir, instance);
1512        assert!(dir.pick_rest_endpoint_round_robin("web").is_some());
1513
1514        assert!(dir.rr_counters.contains_key("rest:web"));
1515
1516        // First eviction pass quarantines the stale instance.
1517        dir.evict_stale(now);
1518        let instances = dir.instances_of("web");
1519        assert_eq!(instances.len(), 1);
1520        assert!(matches!(instances[0].state(), InstanceState::Quarantined));
1521
1522        // Second pass evicts quarantined instances that exceeded the grace period.
1523        let later = now + grace + Duration::from_millis(10);
1524        dir.evict_stale(later);
1525
1526        assert!(dir.instances_of("web").is_empty());
1527        assert!(!dir.rr_counters.contains_key("web"));
1528        assert!(!dir.rr_counters.contains_key("rest:web"));
1529    }
1530
1531    #[test]
1532    fn register_instance_rejects_cross_gear_grpc_name() {
1533        let mgr = GearManager::new();
1534        let service = "cf.authz.v1.AuthzService";
1535
1536        mgr.register_instance(Arc::new(
1537            GearInstance::new("authz-resolver", Uuid::new_v4())
1538                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000)),
1539        ))
1540        .expect("claiming an unowned service name must succeed");
1541
1542        // A different gear claiming the same name is rejected, store untouched.
1543        let conflict = mgr
1544            .register_instance(Arc::new(
1545                GearInstance::new("evil", Uuid::new_v4())
1546                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1547            ))
1548            .unwrap_err();
1549        assert_eq!(conflict.service_name, service);
1550        assert_eq!(conflict.owner, "authz-resolver");
1551        assert!(mgr.instances_of("evil").is_empty());
1552
1553        // The owning gear may add another instance under the same name.
1554        mgr.register_instance(Arc::new(
1555            GearInstance::new("authz-resolver", Uuid::new_v4())
1556                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9002)),
1557        ))
1558        .expect("the owning gear may add another instance for its own service");
1559    }
1560
1561    #[test]
1562    fn declared_ownership_beats_registration_order() {
1563        let mgr = GearManager::new();
1564        let service = "cf.authz.v1.AuthzService";
1565
1566        // Config declares the name is owned by `authz`.
1567        mgr.set_grpc_service_owners(HashMap::from([(service.to_owned(), "authz".to_owned())]));
1568
1569        // A squatter registering *first* cannot claim a declared name; it is
1570        // rejected with the declared owner, and the store is untouched.
1571        let conflict = mgr
1572            .register_instance(Arc::new(
1573                GearInstance::new("evil", Uuid::new_v4())
1574                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1575            ))
1576            .unwrap_err();
1577        assert_eq!(conflict.service_name, service);
1578        assert_eq!(conflict.owner, "authz");
1579        assert!(
1580            !conflict.recoverable,
1581            "a declared-owner conflict is pinned by config (permanent)"
1582        );
1583        assert!(mgr.instances_of("evil").is_empty());
1584
1585        // The declared owner is always admitted, even coming in second.
1586        mgr.register_instance(Arc::new(
1587            GearInstance::new("authz", Uuid::new_v4())
1588                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000)),
1589        ))
1590        .expect("the declared owner must be admitted");
1591        assert_eq!(mgr.grpc_service_owner(service).as_deref(), Some("authz"));
1592    }
1593
1594    /// A gear can own an undeclared name via first-registration, after which a
1595    /// `set_grpc_service_owners` call reassigns that name to a different declared
1596    /// owner. The declared owner is authorized, but must not begin advertising
1597    /// the name while the original holder is still serving it — that would split
1598    /// ownership between two advertisers. It is refused until the stale holder is
1599    /// gone, then admitted.
1600    #[test]
1601    fn declared_owner_refused_while_a_stale_advertiser_still_holds_the_name() {
1602        let mgr = GearManager::new();
1603        let service = "cf.authz.v1.AuthzService";
1604        let holder = Uuid::new_v4();
1605
1606        // `first` claims the name via first-registration (not yet declared).
1607        mgr.register_instance(Arc::new(
1608            GearInstance::new("first", holder)
1609                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000)),
1610        ))
1611        .expect("claiming an undeclared, unowned name must succeed");
1612
1613        // The map is (re)installed, declaring the name is owned by `authz`.
1614        mgr.set_grpc_service_owners(HashMap::from([(service.to_owned(), "authz".to_owned())]));
1615
1616        // The declared owner is authorized, but `first` is still advertising the
1617        // name: admitting `authz` now would split ownership, so it is rejected.
1618        let conflict = mgr
1619            .register_instance(Arc::new(
1620                GearInstance::new("authz", Uuid::new_v4())
1621                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1622            ))
1623            .unwrap_err();
1624        assert_eq!(conflict.service_name, service);
1625        assert_eq!(
1626            conflict.owner, "first",
1627            "the current holder is the conflicting owner"
1628        );
1629        assert!(
1630            conflict.recoverable,
1631            "a stale-advertiser conflict clears once the holder deregisters"
1632        );
1633        assert!(mgr.instances_of("authz").is_empty());
1634
1635        // Once the stale holder is gone, the declared owner is admitted and owns
1636        // the name outright — no split.
1637        mgr.deregister("first", holder);
1638        mgr.register_instance(Arc::new(
1639            GearInstance::new("authz", Uuid::new_v4())
1640                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9002)),
1641        ))
1642        .expect("the declared owner must be admitted once the stale holder is gone");
1643        assert_eq!(mgr.grpc_service_owner(service).as_deref(), Some("authz"));
1644    }
1645
1646    /// The compiled-in (authoritative) map overrides a conflicting operator
1647    /// entry — the binary that *provides* a service name wins over config that
1648    /// tries to reassign it — while operator entries for names the binary does
1649    /// not know about are preserved. A subsequent registration is governed by
1650    /// the merged result.
1651    #[test]
1652    fn authoritative_owners_override_conflicting_config_and_keep_the_rest() {
1653        let mgr = GearManager::new();
1654
1655        // Operator config: one name that a compiled gear also provides (a
1656        // reassignment attempt), and one the binary does not compile in.
1657        mgr.set_grpc_service_owners(HashMap::from([
1658            ("cf.authz.v1.AuthzService".to_owned(), "impostor".to_owned()),
1659            ("remote.only.Service".to_owned(), "remote-gear".to_owned()),
1660        ]));
1661
1662        // Compiled registry says authz-resolver provides the authz service.
1663        mgr.merge_authoritative_grpc_service_owners(HashMap::from([(
1664            "cf.authz.v1.AuthzService".to_owned(),
1665            "authz-resolver".to_owned(),
1666        )]));
1667
1668        // Compiled-in provider wins the conflict; the remote-only entry stays.
1669        let conflict = mgr
1670            .register_instance(Arc::new(
1671                GearInstance::new("impostor", Uuid::new_v4()).with_grpc_service(
1672                    "cf.authz.v1.AuthzService",
1673                    Endpoint::http("127.0.0.1", 9000),
1674                ),
1675            ))
1676            .unwrap_err();
1677        assert_eq!(conflict.owner, "authz-resolver");
1678        assert!(
1679            !conflict.recoverable,
1680            "a pinned compiled-in owner is permanent"
1681        );
1682
1683        mgr.register_instance(Arc::new(
1684            GearInstance::new("authz-resolver", Uuid::new_v4()).with_grpc_service(
1685                "cf.authz.v1.AuthzService",
1686                Endpoint::http("127.0.0.1", 9001),
1687            ),
1688        ))
1689        .expect("the compiled-in provider owns its service name");
1690
1691        // The operator entry for the name absent from the compiled map is intact.
1692        let remote_conflict = mgr
1693            .register_instance(Arc::new(
1694                GearInstance::new("squatter", Uuid::new_v4())
1695                    .with_grpc_service("remote.only.Service", Endpoint::http("127.0.0.1", 9002)),
1696            ))
1697            .unwrap_err();
1698        assert_eq!(remote_conflict.owner, "remote-gear");
1699    }
1700
1701    #[test]
1702    fn names_absent_from_declared_map_keep_first_registration_ownership() {
1703        let mgr = GearManager::new();
1704        // The map declares one name but says nothing about `worker.Svc`.
1705        mgr.set_grpc_service_owners(HashMap::from([(
1706            "cf.authz.v1.AuthzService".to_owned(),
1707            "authz".to_owned(),
1708        )]));
1709        let service = "worker.Svc";
1710
1711        mgr.register_instance(Arc::new(
1712            GearInstance::new("worker-a", Uuid::new_v4())
1713                .with_grpc_service(service, Endpoint::http("127.0.0.1", 9000)),
1714        ))
1715        .expect("claiming an undeclared, unowned name must succeed");
1716
1717        let conflict = mgr
1718            .register_instance(Arc::new(
1719                GearInstance::new("worker-b", Uuid::new_v4())
1720                    .with_grpc_service(service, Endpoint::http("127.0.0.1", 9001)),
1721            ))
1722            .unwrap_err();
1723        assert_eq!(conflict.owner, "worker-a");
1724        assert!(
1725            conflict.recoverable,
1726            "a first-registration conflict clears if the holder leaves"
1727        );
1728    }
1729
1730    /// Under contention, the atomic check+insert admits exactly one owner for a
1731    /// gRPC service name: N gears race to claim the same unowned name and only
1732    /// one commits, closing the check-then-write race the per-entry `DashMap`
1733    /// locks leave open.
1734    #[test]
1735    fn concurrent_register_admits_exactly_one_owner() {
1736        use std::sync::Barrier;
1737
1738        let mgr = Arc::new(GearManager::new());
1739        let service = "cf.authz.v1.AuthzService";
1740        let racers = 16;
1741        let barrier = Arc::new(Barrier::new(racers));
1742
1743        #[expect(
1744            clippy::needless_collect,
1745            reason = "a lazy iterator would join before all racers spawn, deadlocking the Barrier"
1746        )]
1747        let handles: Vec<_> = (0..racers)
1748            .map(|i| {
1749                let mgr = Arc::clone(&mgr);
1750                let barrier = Arc::clone(&barrier);
1751                std::thread::spawn(move || {
1752                    let gear = format!("gear-{i}");
1753                    let instance = Arc::new(
1754                        GearInstance::new(gear.clone(), Uuid::new_v4()).with_grpc_service(
1755                            service,
1756                            Endpoint::http("127.0.0.1", 9000 + u16::try_from(i).unwrap()),
1757                        ),
1758                    );
1759                    // Release all threads at once to maximise contention.
1760                    barrier.wait();
1761                    (gear, mgr.register_instance(instance))
1762                })
1763            })
1764            .collect();
1765
1766        let results: Vec<(String, Result<(), GrpcServiceNameConflict>)> = handles
1767            .into_iter()
1768            .map(|h| h.join().expect("registration thread must not panic"))
1769            .collect();
1770
1771        // Exactly one gear wins the race.
1772        let winners: Vec<&str> = results
1773            .iter()
1774            .filter(|(_, r)| r.is_ok())
1775            .map(|(gear, _)| gear.as_str())
1776            .collect();
1777        assert_eq!(
1778            winners.len(),
1779            1,
1780            "exactly one competing registration must win"
1781        );
1782        let winner = winners[0];
1783
1784        // Every loser is rejected with a conflict naming the winning gear — the
1785        // `owner` surfaced in the server-side warn log and the operator's
1786        // remediation hint, so it must be the gear that actually committed, not
1787        // just any error.
1788        for (gear, result) in &results {
1789            if gear == winner {
1790                continue;
1791            }
1792            let conflict = result
1793                .as_ref()
1794                .expect_err("a losing registration must be rejected, not silently dropped");
1795            assert_eq!(conflict.service_name, service);
1796            assert_eq!(
1797                conflict.owner.as_str(),
1798                winner,
1799                "every loser's conflict must name the single winning gear as owner"
1800            );
1801        }
1802
1803        // ...and the store resolves the contested name to exactly that winner.
1804        assert_eq!(
1805            mgr.grpc_service_owner(service).as_deref(),
1806            Some(winner),
1807            "the contested name must resolve to the one gear that won"
1808        );
1809    }
1810}