mx_remote/runtime/mod.rs
1// Author: Lars Op den Kamp (lars@opdenkamp-it.nl)
2// Copyright (c) 2026 Op den Kamp IT Solutions
3
4//! The running client: the socket, the two threads that drive it, and the
5//! read surface over what they discover.
6
7mod control;
8mod info;
9mod schedule;
10
11#[cfg(test)]
12mod control_tests;
13#[cfg(test)]
14mod tests;
15
16// The reuseport check reads a Linux socket option, and pinning an addressless
17// interface is Linux-only, so the whole file is.
18#[cfg(all(test, target_os = "linux"))]
19mod socket;
20
21use std::io;
22use std::net::Ipv4Addr;
23use std::path::PathBuf;
24use std::sync::atomic::{AtomicBool, Ordering};
25use std::sync::{Arc, Mutex, MutexGuard};
26use std::thread::JoinHandle;
27use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
28
29use crate::event::{Event, EventHandler};
30use crate::rx::process_frame;
31use crate::state::{Device, State, CONFIG_GRACE};
32use crate::types::*;
33use crate::wire::{
34 build_hello, op, Addressee, BayUid, Conn, DeviceFeature, DeviceUid, FirmwareType, Opcode,
35 SendError, Tx, V2ipFpgaFeature, MULTICAST_IP, MULTICAST_PORT, PROTOCOL_VERSION, VERSION,
36};
37
38pub use control::ControlError;
39pub use info::{BayInfo, DeviceInfo};
40use schedule::Schedule;
41
42/// The name this client advertises when the caller sets none.
43const DEFAULT_NAME: &str = "MXR Rust";
44
45/// The serial this client advertises.
46///
47/// A client is not a unit with a serial number, but the field is fixed-width
48/// and every device fills it, so it carries a constant rather than a blank.
49const CLIENT_SERIAL: &str = "P9SN00000000";
50
51/// The file the client's identifier is kept in, under the user's home
52/// directory. The identifier has to survive a restart, or every peer would see
53/// each run as a new client.
54const UID_FILE: &str = ".mxr-uid";
55
56/// Largest datagram the receive buffer accepts.
57const RECV_BUFFER: usize = 65535;
58
59/// How often the background thread re-examines the network.
60const PROBE_TICK: Duration = Duration::from_secs(1);
61
62/// How often that thread looks up from waiting to see whether it should stop.
63const SHUTDOWN_POLL: Duration = Duration::from_millis(50);
64
65/// Shortest gap between two discovery requests.
66const DISCOVER_INTERVAL: Duration = Duration::from_secs(5);
67
68/// How long a discover to every device stands in for asking any one of them.
69///
70/// Each device answers it with everything it holds, after a jitter of up to a
71/// few seconds, so asking one again inside this window repeats an answer
72/// already on its way.
73const DISCOVER_COVERS: Duration = Duration::from_secs(60);
74
75/// How a [`Remote`] finds the network.
76///
77/// [`Config::default`] discovers over multicast on the interface the host
78/// picks, which is the right answer on a single-homed machine and arbitrary on
79/// any other.
80#[derive(Clone, Debug, Default)]
81#[non_exhaustive]
82pub struct Config {
83 /// Where to send. Unset means the multicast group, or the interface's
84 /// broadcast address when [`Config::broadcast`] is set.
85 pub target_ip: Option<Ipv4Addr>,
86 /// UDP port. Unset means the default for the selected mode.
87 pub port: Option<u16>,
88 /// Selects the interface by address.
89 ///
90 /// It becomes both the multicast egress interface and the membership
91 /// interface, so it decides which NIC frames leave by and which one they
92 /// are accepted on. Getting it wrong on a multi-homed host is one-sided:
93 /// periodic broadcasts still arrive, so discovery looks healthy while
94 /// every request this client sends leaves by the wrong NIC and is never
95 /// answered.
96 pub local_ip: Option<Ipv4Addr>,
97 /// Selects the interface by name, taking precedence over
98 /// [`Config::local_ip`].
99 ///
100 /// An interface with no address of its own - a tagged VLAN - can only be
101 /// named this way, and only on Linux.
102 pub interface: Option<String>,
103 /// Use broadcast rather than multicast.
104 pub broadcast: bool,
105 /// The name this client advertises. Unset means `MXR Rust`.
106 pub name: Option<String>,
107 /// This client's identifier. Unset loads it from
108 /// [`Config::uid_path`], generating and storing one on first run.
109 pub uid: Option<DeviceUid>,
110 /// Where the identifier is kept. Unset means `.mxr-uid` in the user's home
111 /// directory.
112 pub uid_path: Option<PathBuf>,
113}
114
115/// Everything the threads share.
116struct Shared {
117 uid: DeviceUid,
118 name: String,
119 handler: Arc<dyn EventHandler>,
120 state: Mutex<State>,
121 tx: Mutex<Tx>,
122 schedule: Mutex<Schedule>,
123 network: Mutex<Network>,
124 closing: AtomicBool,
125}
126
127/// The network parameters the socket was opened with, so it can be reopened.
128#[derive(Clone, Debug)]
129struct Network {
130 target_ip: Option<Ipv4Addr>,
131 port: Option<u16>,
132 local_ip: Option<Ipv4Addr>,
133 interface: Option<String>,
134 broadcast: bool,
135}
136
137impl Network {
138 /// The address to send to: what the caller asked for, else the multicast
139 /// group, else - in broadcast mode - the chosen interface's own broadcast
140 /// address.
141 fn target(&self) -> io::Result<Ipv4Addr> {
142 if let Some(ip) = self.target_ip {
143 return Ok(ip);
144 }
145 if !self.broadcast {
146 return Ok(MULTICAST_IP);
147 }
148 Ok(crate::wire::broadcast_address(self.local_ip).unwrap_or(MULTICAST_IP))
149 }
150
151 fn port(&self) -> u16 {
152 self.port.unwrap_or(if self.broadcast {
153 crate::wire::BROADCAST_PORT
154 } else {
155 MULTICAST_PORT
156 })
157 }
158
159 fn open(&self) -> io::Result<Conn> {
160 Conn::open(
161 self.target()?,
162 self.port(),
163 self.local_ip,
164 self.interface.as_deref(),
165 )
166 }
167}
168
169/// A client on the MX Remote network.
170///
171/// Create one with [`Remote::new`], then [`Remote::start`] it. Discovery runs
172/// on threads this client owns until [`Remote::close`], or until the `Remote`
173/// is dropped.
174pub struct Remote {
175 shared: Arc<Shared>,
176 workers: Mutex<Vec<JoinHandle<()>>>,
177}
178
179impl Remote {
180 /// Builds a client, loading or generating its identifier.
181 ///
182 /// Nothing is sent and no socket is opened until [`Remote::start`].
183 pub fn new(config: Config, handler: Arc<dyn EventHandler>) -> io::Result<Self> {
184 let uid = match config.uid {
185 Some(uid) => uid,
186 None => load_uid(config.uid_path.clone())?,
187 };
188 let name = config.name.unwrap_or_else(|| DEFAULT_NAME.to_owned());
189 Ok(Self {
190 shared: Arc::new(Shared {
191 uid,
192 name,
193 handler,
194 state: Mutex::new(State::new(uid)),
195 tx: Mutex::new(Tx::default()),
196 schedule: Mutex::new(Schedule::new()),
197 network: Mutex::new(Network {
198 target_ip: config.target_ip,
199 port: config.port,
200 local_ip: config.local_ip,
201 interface: config.interface,
202 broadcast: config.broadcast,
203 }),
204 closing: AtomicBool::new(false),
205 }),
206 workers: Mutex::new(Vec::new()),
207 })
208 }
209
210 /// Opens the socket, announces this client and begins discovery.
211 ///
212 /// Returns once the receive thread is running; discovery continues in the
213 /// background.
214 pub fn start(&self) -> io::Result<()> {
215 let conn = lock(&self.shared.network).open()?;
216 lock(&self.shared.tx).set_conn(Some(conn));
217 self.shared.closing.store(false, Ordering::SeqCst);
218 self.spawn_workers()?;
219 self.shared.announce();
220 let _ = self.shared.discover();
221 Ok(())
222 }
223
224 /// Stops discovery, closes the socket and waits for the threads to finish.
225 ///
226 /// Calling it more than once, or before [`Remote::start`], does nothing.
227 pub fn close(&self) {
228 self.shared.closing.store(true, Ordering::SeqCst);
229 for worker in std::mem::take(&mut *lock(&self.workers)) {
230 let _ = worker.join();
231 }
232 // Only once no thread can still be reading from it: a descriptor
233 // released while another thread is parked on it can be reissued to
234 // something else before that thread returns.
235 lock(&self.shared.tx).set_conn(None);
236 }
237
238 fn spawn_workers(&self) -> io::Result<()> {
239 let mut workers = lock(&self.workers);
240 if !workers.is_empty() {
241 return Ok(());
242 }
243 for (name, body) in [
244 ("mxr-rx", Shared::receive_loop as fn(&Shared)),
245 ("mxr-probe", Shared::probe_loop as fn(&Shared)),
246 ] {
247 let shared = Arc::clone(&self.shared);
248 workers.push(
249 std::thread::Builder::new()
250 .name(name.to_owned())
251 .spawn(move || body(&shared))?,
252 );
253 }
254 Ok(())
255 }
256
257 // ---- identity ----
258
259 /// This client's identifier, as peers see it.
260 pub fn uid(&self) -> DeviceUid {
261 self.shared.uid
262 }
263
264 /// The name this client advertises.
265 pub fn name(&self) -> &str {
266 &self.shared.name
267 }
268
269 /// The address frames are being sent to, once started.
270 pub fn target(&self) -> Option<std::net::SocketAddrV4> {
271 lock(&self.shared.tx).conn().map(|conn| conn.target())
272 }
273
274 // ---- reading the registry ----
275
276 /// Every device heard from, in no particular order.
277 pub fn devices(&self) -> Vec<DeviceUid> {
278 self.shared
279 .read(|state| state.devices.keys().copied().collect())
280 }
281
282 /// A snapshot of one device.
283 pub fn device(&self, uid: DeviceUid) -> Option<DeviceInfo> {
284 let now = Instant::now();
285 self.shared
286 .read(|state| state.device(uid).map(|d| DeviceInfo::of(d, now)))
287 }
288
289 /// The device with the given serial number.
290 pub fn device_by_serial(&self, serial: &str) -> Option<DeviceUid> {
291 self.shared
292 .read(|state| state.device_by_serial(serial).map(|d| d.uid))
293 }
294
295 /// Resolves a device from its dotted-hex identifier, falling back to a
296 /// serial-number match.
297 pub fn resolve_device(&self, name: &str) -> Option<DeviceUid> {
298 if let Ok(uid) = name.parse::<DeviceUid>() {
299 if self.shared.read(|state| state.device(uid).is_some()) {
300 return Some(uid);
301 }
302 }
303 self.device_by_serial(name)
304 }
305
306 /// A snapshot of one bay.
307 pub fn bay(&self, uid: BayUid) -> Option<BayInfo> {
308 self.shared
309 .read(|state| state.bay(uid).map(|bay| BayInfo::of(state, bay)))
310 }
311
312 /// The bay on `device` with the given port name, such as `Output 1`.
313 pub fn bay_by_name(&self, device: DeviceUid, port_name: &str) -> Option<BayUid> {
314 self.shared.read(|state| {
315 state
316 .device(device)?
317 .bay_by_name(port_name)
318 .map(crate::state::Bay::uid)
319 })
320 }
321
322 /// The source bay advertising the given multicast group, for the video or
323 /// the audio stream.
324 pub fn bay_by_stream_ip(&self, ip: Ipv4Addr, audio: bool) -> Option<BayUid> {
325 self.shared.read(|state| state.bay_by_stream_ip(ip, audio))
326 }
327
328 /// The V2IP streams a device advertises.
329 pub fn v2ip_sources(&self, uid: DeviceUid) -> Option<Vec<V2ipStreamSources>> {
330 self.shared
331 .read(|state| state.device(uid)?.v2ip_sources.clone())
332 }
333
334 /// A V2IP device's own encoder configuration.
335 pub fn v2ip_details(&self, uid: DeviceUid) -> Option<DeviceV2ipDetails> {
336 self.shared.read(|state| state.device(uid)?.v2ip_details)
337 }
338
339 /// The streams a V2IP sink is subscribed to.
340 pub fn v2ip_sink(&self, uid: DeviceUid) -> Option<DeviceV2ipSink> {
341 self.shared.read(|state| state.device(uid)?.v2ip_sink)
342 }
343
344 /// Transport statistics a V2IP device reports.
345 pub fn v2ip_stats(&self, uid: DeviceUid) -> Option<V2ipDeviceStats> {
346 self.shared.read(|state| state.device(uid)?.v2ip_stats)
347 }
348
349 /// What a V2IP device's video processor supports.
350 ///
351 /// `None` until the device has reported a non-empty mask about itself.
352 /// A device reports nothing at all while its processor has yet to answer,
353 /// and an older processor answers with none of the optional commands, so
354 /// "no features" and "not yet known" are the same bytes on the wire and
355 /// both read as `None` here. A mask that has arrived only gains bits.
356 pub fn v2ip_features(&self, uid: DeviceUid) -> Option<V2ipFpgaFeature> {
357 self.shared.read(|state| state.device(uid)?.v2ip_features)
358 }
359
360 /// A V2IP device's settings, `None` until it has reported any.
361 ///
362 /// A device reports only the settings it has, so one missing from
363 /// [`V2ipDeviceSettings::valid`] once the rest have arrived is one the
364 /// device does not have.
365 pub fn v2ip_device_settings(&self, uid: DeviceUid) -> Option<V2ipDeviceSettings> {
366 self.shared.read(|state| state.device(uid)?.v2ip_settings)
367 }
368
369 /// The VLAN configuration a V2IP device last reported about itself,
370 /// `None` until it has reported one.
371 pub fn v2ip_vlan(&self, uid: DeviceUid) -> Option<V2ipVlan> {
372 self.shared.read(|state| state.device(uid)?.v2ip_vlan)
373 }
374
375 /// The health and IP configuration a unit last reported, `None` until it
376 /// has reported them.
377 pub fn unit_status(&self, uid: DeviceUid) -> Option<UnitStatus> {
378 self.shared
379 .read(|state| state.device(uid)?.unit_status.clone())
380 }
381
382 /// The PTP state a V2IP device last reported about itself, `None` until
383 /// it has reported one, and while it has no PTP to report: none on its
384 /// video processor, or in power save.
385 pub fn v2ip_ptp(&self, uid: DeviceUid) -> Option<V2ipPtpState> {
386 self.shared.read(|state| state.device(uid)?.v2ip_ptp)
387 }
388
389 /// The test card a V2IP sink last reported, `None` until it has answered
390 /// [`Remote::request_v2ip_testcard`] or a change.
391 pub fn v2ip_testcard(&self, uid: DeviceUid) -> Option<V2ipTestcard> {
392 self.shared.read(|state| state.device(uid)?.v2ip_testcard)
393 }
394
395 /// The time zone a device announced for its mesh, `None` until it has.
396 ///
397 /// The mesh controller announces it with every periodic broadcast, empty
398 /// when it has none.
399 pub fn time_zone(&self, uid: DeviceUid) -> Option<TimeZone> {
400 self.shared
401 .read(|state| state.device(uid)?.time_zone.clone())
402 }
403
404 /// A device's clock as of now: the time it last announced, advanced by
405 /// how long ago that arrived. `None` until it has announced one.
406 ///
407 /// The mesh controller announces its clock with every periodic broadcast,
408 /// and only once it has been set.
409 pub fn device_clock(&self, uid: DeviceUid) -> Option<SystemTime> {
410 let (utc, at) = self.shared.read(|state| state.device(uid)?.clock)?;
411 let announced = UNIX_EPOCH.checked_add(Duration::from_secs(u64::from(utc)))?;
412 announced.checked_add(at.elapsed())
413 }
414
415 /// The video-wall tiling a V2IP device is configured for.
416 pub fn v2ip_tiling(&self, uid: DeviceUid) -> Option<V2ipTilingConfig> {
417 self.shared.read(|state| state.device(uid)?.tiling)
418 }
419
420 /// The audio endpoint tree a device exposes.
421 pub fn audio_endpoints(&self, uid: DeviceUid) -> Option<AudioEndpoints> {
422 self.shared.read(|state| state.device(uid)?.audio.clone())
423 }
424
425 /// The multiviewer layout a device is showing.
426 pub fn multiviewer_status(&self, uid: DeviceUid) -> Option<MultiviewerStatus> {
427 self.shared
428 .read(|state| state.device(uid)?.multiviewer.clone())
429 }
430
431 /// An amplifier's Dolby decoder settings.
432 pub fn dolby_settings(&self, uid: DeviceUid) -> Option<AmpDolbySettings> {
433 self.shared.read(|state| state.device(uid)?.dolby_settings)
434 }
435
436 /// A device's remote-control settings.
437 pub fn rc_settings(&self, uid: DeviceUid) -> Option<RcSettings> {
438 self.shared
439 .read(|state| state.device(uid)?.rc_settings.clone())
440 }
441
442 /// Every network port a device reports, in port order.
443 pub fn network_status(&self, uid: DeviceUid) -> Vec<NetworkPortStatus> {
444 self.shared.read(|state| {
445 state
446 .device(uid)
447 .map(|d| d.network.values().cloned().collect())
448 .unwrap_or_default()
449 })
450 }
451
452 /// The mesh topology a device reports.
453 pub fn topology(&self, uid: DeviceUid) -> Vec<TopologyEntry> {
454 self.shared.read(|state| {
455 state
456 .device(uid)
457 .map(|d| d.topology.clone())
458 .unwrap_or_default()
459 })
460 }
461
462 /// The EDID a device last reported: the display on its output, or the one
463 /// it presents to the source on its input.
464 ///
465 /// Filled in by a device's answer to [`Remote::request_edid`], and by any
466 /// answer to a peer's request that this client happened to hear.
467 pub fn edid(&self, uid: DeviceUid, output: bool) -> Option<Vec<u8>> {
468 self.shared
469 .read(|state| state.device(uid)?.edid(output).map(<[u8]>::to_vec))
470 }
471
472 /// How many frames from other senders have parsed since this client
473 /// started.
474 ///
475 /// It separates a mesh with nothing on it from an interface nothing is on:
476 /// a client that has discovered no device but is counting frames is
477 /// hearing traffic it cannot get answers from, which on a multi-homed host
478 /// is what a wrong [`Config::local_ip`] looks like. Frames this client
479 /// sent are not counted, because the host loops its own multicast back
480 /// whichever interface was selected.
481 pub fn frames_received(&self) -> u64 {
482 self.shared.read(|state| state.frames_received)
483 }
484
485 /// Every firmware image a device reports a version for.
486 pub fn firmware(&self, uid: DeviceUid) -> Vec<(FirmwareType, FirmwareVersion)> {
487 self.shared.read(|state| {
488 state
489 .device(uid)
490 .map(|d| d.firmware.iter().map(|(k, v)| (*k, v.clone())).collect())
491 .unwrap_or_default()
492 })
493 }
494
495 // ---- reconfiguring ----
496
497 /// Changes the interface and the multicast/broadcast mode while running,
498 /// reopening the socket when either differs from what is in use.
499 pub fn update_config(&self, local_ip: Option<Ipv4Addr>, broadcast: bool) -> io::Result<()> {
500 let network = {
501 let mut network = lock(&self.shared.network);
502 if network.local_ip == local_ip && network.broadcast == broadcast {
503 return Ok(());
504 }
505 network.local_ip = local_ip;
506 network.broadcast = broadcast;
507 network.clone()
508 };
509 // Opened before the old one is dropped, so a bind that fails leaves the
510 // client on the socket it had rather than on none.
511 let conn = network.open()?;
512 lock(&self.shared.tx).set_conn(Some(conn));
513 self.shared.announce();
514 let _ = self.shared.discover();
515 Ok(())
516 }
517
518 /// Asks every device on the network to announce itself.
519 pub fn discover(&self) -> Result<(), SendError> {
520 self.shared.discover()
521 }
522}
523
524impl Drop for Remote {
525 fn drop(&mut self) {
526 self.close();
527 }
528}
529
530impl Shared {
531 /// Reads the registry.
532 fn read<R>(&self, f: impl FnOnce(&State) -> R) -> R {
533 f(&lock(&self.state))
534 }
535
536 /// Mutates the registry, then delivers what changed.
537 ///
538 /// The queue is drained after the lock is released, so an event handler
539 /// may call back into the library.
540 fn mutate<R>(&self, f: impl FnOnce(&mut State, &mut Vec<Event>) -> R) -> R {
541 let mut events = Vec::new();
542 let result = f(&mut lock(&self.state), &mut events);
543 self.dispatch(events);
544 result
545 }
546
547 fn dispatch(&self, events: Vec<Event>) {
548 for event in events {
549 event.dispatch(&*self.handler);
550 }
551 }
552
553 /// Decodes one datagram and delivers what it changed.
554 ///
555 /// This is the receive entry point. Keeping it distinct from the decode
556 /// below matters even though the wrapper is thin: announcing is driven by a
557 /// clock, and the one frame that may add a hello here is a ping addressed
558 /// to this client, which asks for exactly that.
559 fn process_datagram(&self, data: &[u8], from: Ipv4Addr) {
560 let (events, hello_requested, state_requests) = {
561 let mut state = lock(&self.state);
562 let events = process_frame(&mut state, data, Some(from), Instant::now());
563 (
564 events,
565 std::mem::take(&mut state.hello_requested),
566 std::mem::take(&mut state.state_requests),
567 )
568 };
569 self.dispatch(events);
570 if hello_requested {
571 self.announce();
572 }
573 for device in state_requests {
574 self.request_state(device);
575 }
576 }
577
578 /// Sends a frame, refusing one the addressee cannot decode.
579 fn send(&self, to: &Addressee, opcode: Opcode, payload: &[u8]) -> Result<usize, SendError> {
580 lock(&self.tx).send(to, self.uid, opcode, payload)
581 }
582
583 fn discover(&self) -> Result<(), SendError> {
584 lock(&self.schedule).discovered(Instant::now());
585 self.send(&Addressee::Broadcast, op::SYS_DISCOVER, &[])?;
586 Ok(())
587 }
588
589 /// Asks one device for everything it holds, with a discover naming it.
590 ///
591 /// Skipped while a discover to every device is recent enough to have
592 /// covered it. A receiver that predates the named form answers it as a
593 /// discover to every device, so each one asked costs the whole mesh's
594 /// state on such a network.
595 fn request_state(&self, device: DeviceUid) {
596 if lock(&self.schedule).discover_covers(Instant::now()) {
597 return;
598 }
599 let _ = self.send(&Addressee::Broadcast, op::SYS_DISCOVER, device.as_bytes());
600 }
601
602 /// Announces this client, and re-arms the announcement timer only once the
603 /// frame is away.
604 ///
605 /// The firmware resets its own hello timeout inside the branch where the
606 /// transmit succeeded. A send that fails is then retried on the next tick
607 /// rather than costing a whole interval of silence, which matters most at
608 /// startup and after a network blip: exactly when being heard is worth the
609 /// most.
610 fn announce(&self) {
611 let payload = build_hello(
612 PROTOCOL_VERSION,
613 &self.name,
614 CLIENT_SERIAL,
615 VERSION,
616 DeviceFeature::MANAGER.bits(),
617 );
618 match self.send(&Addressee::Broadcast, op::SYS_HELLO, &payload) {
619 Ok(n) if n > 0 => lock(&self.schedule).announced(Instant::now()),
620 _ => {}
621 }
622 }
623
624 /// Waits out one tick, reporting whether the client is still running.
625 ///
626 /// Slept in short steps rather than one, so closing does not have to wait
627 /// out a whole tick before the thread can be joined.
628 fn sleep_until_next_tick(&self) -> bool {
629 let deadline = Instant::now() + PROBE_TICK;
630 while Instant::now() < deadline {
631 if self.closing.load(Ordering::SeqCst) {
632 return false;
633 }
634 std::thread::sleep(SHUTDOWN_POLL);
635 }
636 !self.closing.load(Ordering::SeqCst)
637 }
638
639 /// Whether it is time to announce again.
640 ///
641 /// This is a timer, not a response to traffic: a device announces itself on
642 /// a schedule whether or not anything is talking to it, and a client that
643 /// only re-announced when a datagram arrived would go silent on a quiet
644 /// network and stay unknown to every peer that started after it.
645 fn announce_due(&self, now: Instant) -> bool {
646 !self.closing.load(Ordering::SeqCst) && lock(&self.schedule).announce_due(now)
647 }
648
649 /// Reads datagrams until the client is closing.
650 fn receive_loop(&self) {
651 let mut buf = vec![0u8; RECV_BUFFER];
652 while !self.closing.load(Ordering::SeqCst) {
653 let Some(conn) = lock(&self.tx).conn() else {
654 break;
655 };
656 match conn.recv(&mut buf) {
657 Ok(Some((data, from))) => self.process_datagram(data, from),
658 Ok(None) => {}
659 Err(_) => break,
660 }
661 }
662 }
663
664 /// Runs [`Self::probe_once`] once a second until the client is closing.
665 fn probe_loop(&self) {
666 while self.sleep_until_next_tick() {
667 self.probe_once(Instant::now());
668 }
669 }
670
671 /// One pass: liveness, completion, the announcement timer and discovery.
672 pub(super) fn probe_once(&self, now: Instant) {
673 let want_discover = self.mutate(|state, ev| {
674 let mut incomplete = false;
675 let mut any_complete = false;
676 for device in state.devices.values_mut() {
677 device.check_online(now, ev);
678 // Re-tested here rather than only where a frame lands: a device
679 // stops waiting for its links on the clock alone, and no frame
680 // arrives to announce that.
681 device.check_config_complete(now, ev);
682 if device.configuration_complete(now) {
683 any_complete = true;
684 } else if now.saturating_duration_since(device.first_seen) > CONFIG_GRACE {
685 // Past the grace period a device has said nothing more, so
686 // ask the network again rather than wait forever.
687 incomplete = true;
688 }
689 }
690 // Nothing has finished describing itself, so nothing has been
691 // discovered yet at all.
692 incomplete || !any_complete
693 });
694
695 let discover_due = lock(&self.schedule).discover_due(now);
696 if self.announce_due(now) {
697 self.announce();
698 }
699 if want_discover && discover_due {
700 let _ = self.discover();
701 }
702 }
703}
704
705/// Takes a lock, continuing through a poisoning.
706///
707/// A poisoned lock here means a panic somewhere that held it. The state behind
708/// it is a cache of what devices have reported and is rebuilt by the next
709/// frame from each, so refusing to touch it again would retire a working
710/// client over one bad datagram.
711fn lock<T>(m: &Mutex<T>) -> MutexGuard<'_, T> {
712 m.lock().unwrap_or_else(|e| e.into_inner())
713}
714
715/// Loads the client's identifier, generating and storing one on first run.
716///
717/// A generated identifier that cannot be stored is still used: a client with a
718/// new identity each run is worse than one that works today, and the failure is
719/// the caller's home directory, not the network.
720fn load_uid(path: Option<PathBuf>) -> io::Result<DeviceUid> {
721 let path = path.or_else(|| {
722 std::env::var_os("HOME")
723 .or_else(|| std::env::var_os("USERPROFILE"))
724 .map(|home| PathBuf::from(home).join(UID_FILE))
725 });
726 if let Some(path) = &path {
727 if let Ok(bytes) = std::fs::read(path) {
728 if let Ok(array) = <[u8; 16]>::try_from(bytes.get(..16).unwrap_or_default()) {
729 return Ok(DeviceUid::from_array(array));
730 }
731 }
732 }
733 let mut bytes = [0u8; 16];
734 getrandom::getrandom(&mut bytes).map_err(|e| io::Error::other(e.to_string()))?;
735 if let Some(path) = &path {
736 let _ = std::fs::write(path, bytes);
737 }
738 Ok(DeviceUid::from_array(bytes))
739}
740
741/// The protocol floor is checked against what the device says it can decode.
742impl crate::wire::ProtocolTarget for Device {
743 fn serial(&self) -> &str {
744 Device::serial(self)
745 }
746
747 fn supported_protocol(&self) -> u16 {
748 self.hello.supported_protocol
749 }
750}