1use dashmap::DashMap;
4use std::collections::{BTreeMap, HashMap};
5use std::sync::Arc;
6use std::time::{Duration, Instant};
7use uuid::Uuid;
8
9#[derive(Clone, Debug, PartialEq, Eq, Hash)]
11pub struct Endpoint {
12 pub uri: String,
13}
14
15#[derive(Clone, Debug, PartialEq, Eq)]
17pub enum EndpointKind {
18 Tcp(std::net::SocketAddr),
20 Uds(std::path::PathBuf),
22 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 #[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#[derive(Clone, Debug)]
82pub struct InstanceRuntimeState {
83 pub last_heartbeat: Instant,
84 pub state: InstanceState,
85}
86
87#[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 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 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 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 #[must_use]
197 pub fn state(&self) -> InstanceState {
198 self.inner.read().state
199 }
200
201 #[must_use]
203 pub fn last_heartbeat(&self) -> Instant {
204 self.inner.read().last_heartbeat
205 }
206}
207
208#[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 reg_lock: parking_lot::Mutex<()>,
234 service_owners: parking_lot::RwLock<HashMap<String, String>>,
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
251pub struct GrpcServiceNameConflict {
252 pub service_name: String,
254 pub owner: String,
256 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 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 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 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 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 fn check_grpc_service_ownership(
419 &self,
420 instance: &GearInstance,
421 ) -> Result<(), GrpcServiceNameConflict> {
422 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 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 for entry in &self.inner {
449 if entry.key() == &instance.gear {
450 continue; }
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 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 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 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 if state.state == InstanceState::Registered {
489 state.state = InstanceState::Healthy;
490 }
491 }
492 }
493
494 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 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 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 #[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 #[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 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 if age >= self.hb_ttl && !matches!(state.state, Quarantined | Draining) {
561 drop(state); inst.inner.write().state = Quarantined;
563 return true; }
565
566 if state.state == Quarantined && age >= self.hb_ttl + self.hb_grace {
568 return false; }
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 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 #[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 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 #[must_use]
632 pub fn pick_service_round_robin(
633 &self,
634 service_name: &str,
635 ) -> Option<(String, Arc<GearInstance>, Endpoint)> {
636 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 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 #[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 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 #[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 #[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 #[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 seed(
789 &dir,
790 Arc::new(
791 GearInstance::new("shard-gear", instance_id).with_labels(labels(&[("shard", "7")])),
792 ),
793 );
794
795 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 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 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 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(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 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 assert!(ids.contains(&id1));
1041 assert!(ids.contains(&id2));
1042 assert_eq!(picked1.instance_id, picked3.instance_id);
1044 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 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 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 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 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 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 dir.update_heartbeat("test_gear", id1, Instant::now());
1263 dir.update_heartbeat("test_gear", id2, Instant::now());
1264
1265 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 assert_eq!(inst1.instance_id, inst3.instance_id);
1280 assert_ne!(inst1.instance_id, inst2.instance_id);
1282 assert_ne!(ep1, ep2);
1284 }
1285
1286 #[test]
1287 fn pick_service_falls_back_to_not_ready() {
1288 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 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 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 mgr.set_grpc_service_owners(HashMap::from([(service.to_owned(), "authz".to_owned())]));
1568
1569 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 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 #[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 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 mgr.set_grpc_service_owners(HashMap::from([(service.to_owned(), "authz".to_owned())]));
1615
1616 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 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 #[test]
1652 fn authoritative_owners_override_conflicting_config_and_keep_the_rest() {
1653 let mgr = GearManager::new();
1654
1655 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 mgr.merge_authoritative_grpc_service_owners(HashMap::from([(
1664 "cf.authz.v1.AuthzService".to_owned(),
1665 "authz-resolver".to_owned(),
1666 )]));
1667
1668 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 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 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 #[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 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 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 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 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}