running_process/broker/server/
backend_registry.rs1use std::collections::HashMap;
4
5use crate::broker::backend_handle::BackendHandle;
6use crate::broker::protocol::ServiceDefinition;
7use crate::broker::server::hello_handler::RegisteredBackend;
8use crate::broker::server::instance::BrokerInstanceKey;
9
10#[derive(Clone, Debug, PartialEq, Eq, Hash)]
19pub struct BackendKey {
20 pub instance: BrokerInstanceKey,
22 pub service_name: String,
24 pub service_version: String,
26 pub exe_hash: String,
33}
34
35impl BackendKey {
36 pub fn new(
38 instance: BrokerInstanceKey,
39 service_name: impl Into<String>,
40 service_version: impl Into<String>,
41 exe_hash: impl Into<String>,
42 ) -> Self {
43 Self {
44 instance,
45 service_name: service_name.into(),
46 service_version: service_version.into(),
47 exe_hash: exe_hash.into(),
48 }
49 }
50}
51
52#[derive(Default)]
54pub struct BackendRegistry {
55 entries: HashMap<BackendKey, BackendHandle>,
56}
57
58impl BackendRegistry {
59 pub fn new() -> Self {
61 Self {
62 entries: HashMap::new(),
63 }
64 }
65
66 pub fn len(&self) -> usize {
68 self.entries.len()
69 }
70
71 pub fn is_empty(&self) -> bool {
73 self.entries.is_empty()
74 }
75
76 pub fn insert(
82 &mut self,
83 instance: BrokerInstanceKey,
84 handle: BackendHandle,
85 ) -> Option<BackendHandle> {
86 let key = BackendKey::new(
87 instance,
88 handle.service_name.clone(),
89 handle.service_version.clone(),
90 hex_lower(&handle.daemon_process.exe_hash),
91 );
92 self.entries.insert(key, handle)
93 }
94
95 pub fn get(
101 &self,
102 instance: &BrokerInstanceKey,
103 service_name: &str,
104 service_version: &str,
105 exe_hash: &str,
106 ) -> Option<&BackendHandle> {
107 self.entries.get(&BackendKey::new(
108 instance.clone(),
109 service_name,
110 service_version,
111 exe_hash,
112 ))
113 }
114
115 pub fn get_any_build(
124 &self,
125 instance: &BrokerInstanceKey,
126 service_name: &str,
127 service_version: &str,
128 ) -> Option<&BackendHandle> {
129 self.entries.iter().find_map(|(key, handle)| {
130 (key.instance == *instance
131 && key.service_name == service_name
132 && key.service_version == service_version)
133 .then_some(handle)
134 })
135 }
136
137 pub fn iter(&self) -> impl Iterator<Item = (&BackendKey, &BackendHandle)> {
139 self.entries.iter()
140 }
141
142 pub fn prune_stale(&mut self) -> Vec<BackendKey> {
147 let mut removed = Vec::new();
148 self.entries.retain(|key, handle| {
149 let alive = handle.is_alive();
150 if !alive {
151 removed.push(key.clone());
152 }
153 alive
154 });
155 removed
156 }
157
158 pub fn registered_backend_for(
165 &self,
166 instance: &BrokerInstanceKey,
167 service_definition: &ServiceDefinition,
168 service_version: &str,
169 expected_exe_hash: &str,
170 ) -> Option<RegisteredBackend> {
171 let handle = self.get(
172 instance,
173 &service_definition.service_name,
174 service_version,
175 expected_exe_hash,
176 )?;
177 Some(RegisteredBackend {
178 service_definition: service_definition.clone(),
179 daemon_version: handle.service_version.clone(),
180 backend_pipe: handle.daemon_process.ipc_endpoint.path.clone(),
181 server_capabilities: 0,
182 })
183 }
184
185 pub fn registered_backend_for_any_build(
188 &self,
189 instance: &BrokerInstanceKey,
190 service_definition: &ServiceDefinition,
191 service_version: &str,
192 ) -> Option<RegisteredBackend> {
193 let handle =
194 self.get_any_build(instance, &service_definition.service_name, service_version)?;
195 Some(RegisteredBackend {
196 service_definition: service_definition.clone(),
197 daemon_version: handle.service_version.clone(),
198 backend_pipe: handle.daemon_process.ipc_endpoint.path.clone(),
199 server_capabilities: 0,
200 })
201 }
202}
203
204pub(crate) fn hex_lower(bytes: &[u8; 32]) -> String {
207 use std::fmt::Write as _;
208 let mut out = String::with_capacity(64);
209 for b in bytes {
210 let _ = write!(out, "{b:02x}");
211 }
212 out
213}
214
215#[cfg(test)]
216mod tests {
217 use crate::broker::backend_handle::{BackendHandle, DaemonProcess};
218 use crate::broker::protocol::Endpoint;
219
220 use super::*;
221
222 fn handle(service_name: &str, version: &str, pid: u32) -> BackendHandle {
223 let endpoint = Endpoint {
224 namespace_id: "shared".into(),
225 path: format!("rpb-v1-test-{service_name}-{version}"),
226 };
227 let mut daemon = DaemonProcess::current_process(endpoint, Some(30)).unwrap();
228 daemon.pid = pid;
229
230 BackendHandle {
231 service_name: service_name.into(),
232 service_version: version.into(),
233 daemon_process: daemon,
234 process_handle: None,
235 }
236 }
237
238 fn test_exe_hash() -> String {
241 hex_lower(
242 &handle("probe", "0.0.0", std::process::id())
243 .daemon_process
244 .exe_hash,
245 )
246 }
247
248 #[test]
249 fn prune_stale_removes_dead_handles_and_keeps_live_ones() {
250 let mut registry = BackendRegistry::new();
251 let exe = test_exe_hash();
252 let live_key = BackendKey::new(BrokerInstanceKey::Shared, "zccache", "1.11.20", &exe);
253 let dead_key = BackendKey::new(BrokerInstanceKey::Shared, "zccache", "1.11.21", &exe);
254
255 registry.insert(
256 live_key.instance.clone(),
257 handle(
258 &live_key.service_name,
259 &live_key.service_version,
260 std::process::id(),
261 ),
262 );
263 registry.insert(
264 dead_key.instance.clone(),
265 handle(&dead_key.service_name, &dead_key.service_version, u32::MAX),
266 );
267
268 let removed = registry.prune_stale();
269
270 assert_eq!(removed, vec![dead_key.clone()]);
271 assert!(registry
272 .get(
273 &live_key.instance,
274 &live_key.service_name,
275 &live_key.service_version,
276 &exe,
277 )
278 .is_some());
279 assert!(registry
280 .get(
281 &dead_key.instance,
282 &dead_key.service_name,
283 &dead_key.service_version,
284 &exe,
285 )
286 .is_none());
287 }
288
289 #[test]
290 fn same_version_different_build_is_a_distinct_entry() {
291 let mut registry = BackendRegistry::new();
294
295 let mut a = handle("zccache", "1.11.20", std::process::id());
296 a.daemon_process.exe_hash = [0xAA; 32];
297 let mut b = handle("zccache", "1.11.20", std::process::id());
298 b.daemon_process.exe_hash = [0xBB; 32];
299 let a_pipe = a.daemon_process.ipc_endpoint.path.clone();
300
301 registry.insert(BrokerInstanceKey::Shared, a);
302 let replaced = registry.insert(BrokerInstanceKey::Shared, b);
304 assert!(
305 replaced.is_none(),
306 "a different exe hash must be a new registry entry, not a replacement"
307 );
308 assert_eq!(registry.len(), 2, "both builds coexist");
309
310 let got = registry
312 .get(
313 &BrokerInstanceKey::Shared,
314 "zccache",
315 "1.11.20",
316 &hex_lower(&[0xAA; 32]),
317 )
318 .expect("build A is reachable by its own hash");
319 assert_eq!(got.daemon_process.ipc_endpoint.path, a_pipe);
320
321 assert!(registry
323 .get(
324 &BrokerInstanceKey::Shared,
325 "zccache",
326 "1.11.20",
327 &hex_lower(&[0xCC; 32]),
328 )
329 .is_none());
330 }
331}