Skip to main content

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}