1mod control;
8mod info;
9mod schedule;
10
11#[cfg(test)]
12mod control_tests;
13#[cfg(test)]
14mod tests;
15
16#[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
42const DEFAULT_NAME: &str = "MXR Rust";
44
45const CLIENT_SERIAL: &str = "P9SN00000000";
50
51const UID_FILE: &str = ".mxr-uid";
55
56const RECV_BUFFER: usize = 65535;
58
59const PROBE_TICK: Duration = Duration::from_secs(1);
61
62const SHUTDOWN_POLL: Duration = Duration::from_millis(50);
64
65const DISCOVER_INTERVAL: Duration = Duration::from_secs(5);
67
68#[derive(Clone, Debug, Default)]
74#[non_exhaustive]
75pub struct Config {
76 pub target_ip: Option<Ipv4Addr>,
79 pub port: Option<u16>,
81 pub local_ip: Option<Ipv4Addr>,
90 pub interface: Option<String>,
96 pub broadcast: bool,
98 pub name: Option<String>,
100 pub uid: Option<DeviceUid>,
103 pub uid_path: Option<PathBuf>,
106}
107
108struct Shared {
110 uid: DeviceUid,
111 name: String,
112 handler: Arc<dyn EventHandler>,
113 state: Mutex<State>,
114 tx: Mutex<Tx>,
115 schedule: Mutex<Schedule>,
116 network: Mutex<Network>,
117 closing: AtomicBool,
118}
119
120#[derive(Clone, Debug)]
122struct Network {
123 target_ip: Option<Ipv4Addr>,
124 port: Option<u16>,
125 local_ip: Option<Ipv4Addr>,
126 interface: Option<String>,
127 broadcast: bool,
128}
129
130impl Network {
131 fn target(&self) -> io::Result<Ipv4Addr> {
135 if let Some(ip) = self.target_ip {
136 return Ok(ip);
137 }
138 if !self.broadcast {
139 return Ok(MULTICAST_IP);
140 }
141 Ok(crate::wire::broadcast_address(self.local_ip).unwrap_or(MULTICAST_IP))
142 }
143
144 fn port(&self) -> u16 {
145 self.port.unwrap_or(if self.broadcast {
146 crate::wire::BROADCAST_PORT
147 } else {
148 MULTICAST_PORT
149 })
150 }
151
152 fn open(&self) -> io::Result<Conn> {
153 Conn::open(
154 self.target()?,
155 self.port(),
156 self.local_ip,
157 self.interface.as_deref(),
158 )
159 }
160}
161
162pub struct Remote {
168 shared: Arc<Shared>,
169 workers: Mutex<Vec<JoinHandle<()>>>,
170}
171
172impl Remote {
173 pub fn new(config: Config, handler: Arc<dyn EventHandler>) -> io::Result<Self> {
177 let uid = match config.uid {
178 Some(uid) => uid,
179 None => load_uid(config.uid_path.clone())?,
180 };
181 let name = config.name.unwrap_or_else(|| DEFAULT_NAME.to_owned());
182 Ok(Self {
183 shared: Arc::new(Shared {
184 uid,
185 name,
186 handler,
187 state: Mutex::new(State::new(uid)),
188 tx: Mutex::new(Tx::default()),
189 schedule: Mutex::new(Schedule::new()),
190 network: Mutex::new(Network {
191 target_ip: config.target_ip,
192 port: config.port,
193 local_ip: config.local_ip,
194 interface: config.interface,
195 broadcast: config.broadcast,
196 }),
197 closing: AtomicBool::new(false),
198 }),
199 workers: Mutex::new(Vec::new()),
200 })
201 }
202
203 pub fn start(&self) -> io::Result<()> {
208 let conn = lock(&self.shared.network).open()?;
209 lock(&self.shared.tx).set_conn(Some(conn));
210 self.shared.closing.store(false, Ordering::SeqCst);
211 self.spawn_workers()?;
212 self.shared.announce();
213 let _ = self.shared.discover();
214 Ok(())
215 }
216
217 pub fn close(&self) {
221 self.shared.closing.store(true, Ordering::SeqCst);
222 for worker in std::mem::take(&mut *lock(&self.workers)) {
223 let _ = worker.join();
224 }
225 lock(&self.shared.tx).set_conn(None);
229 }
230
231 fn spawn_workers(&self) -> io::Result<()> {
232 let mut workers = lock(&self.workers);
233 if !workers.is_empty() {
234 return Ok(());
235 }
236 for (name, body) in [
237 ("mxr-rx", Shared::receive_loop as fn(&Shared)),
238 ("mxr-probe", Shared::probe_loop as fn(&Shared)),
239 ] {
240 let shared = Arc::clone(&self.shared);
241 workers.push(
242 std::thread::Builder::new()
243 .name(name.to_owned())
244 .spawn(move || body(&shared))?,
245 );
246 }
247 Ok(())
248 }
249
250 pub fn uid(&self) -> DeviceUid {
254 self.shared.uid
255 }
256
257 pub fn name(&self) -> &str {
259 &self.shared.name
260 }
261
262 pub fn target(&self) -> Option<std::net::SocketAddrV4> {
264 lock(&self.shared.tx).conn().map(|conn| conn.target())
265 }
266
267 pub fn devices(&self) -> Vec<DeviceUid> {
271 self.shared
272 .read(|state| state.devices.keys().copied().collect())
273 }
274
275 pub fn device(&self, uid: DeviceUid) -> Option<DeviceInfo> {
277 let now = Instant::now();
278 self.shared
279 .read(|state| state.device(uid).map(|d| DeviceInfo::of(d, now)))
280 }
281
282 pub fn device_by_serial(&self, serial: &str) -> Option<DeviceUid> {
284 self.shared
285 .read(|state| state.device_by_serial(serial).map(|d| d.uid))
286 }
287
288 pub fn resolve_device(&self, name: &str) -> Option<DeviceUid> {
291 if let Ok(uid) = name.parse::<DeviceUid>() {
292 if self.shared.read(|state| state.device(uid).is_some()) {
293 return Some(uid);
294 }
295 }
296 self.device_by_serial(name)
297 }
298
299 pub fn bay(&self, uid: BayUid) -> Option<BayInfo> {
301 self.shared
302 .read(|state| state.bay(uid).map(|bay| BayInfo::of(state, bay)))
303 }
304
305 pub fn bay_by_name(&self, device: DeviceUid, port_name: &str) -> Option<BayUid> {
307 self.shared.read(|state| {
308 state
309 .device(device)?
310 .bay_by_name(port_name)
311 .map(crate::state::Bay::uid)
312 })
313 }
314
315 pub fn bay_by_stream_ip(&self, ip: Ipv4Addr, audio: bool) -> Option<BayUid> {
318 self.shared.read(|state| state.bay_by_stream_ip(ip, audio))
319 }
320
321 pub fn v2ip_sources(&self, uid: DeviceUid) -> Option<Vec<V2ipStreamSources>> {
323 self.shared
324 .read(|state| state.device(uid)?.v2ip_sources.clone())
325 }
326
327 pub fn v2ip_details(&self, uid: DeviceUid) -> Option<DeviceV2ipDetails> {
329 self.shared.read(|state| state.device(uid)?.v2ip_details)
330 }
331
332 pub fn v2ip_sink(&self, uid: DeviceUid) -> Option<DeviceV2ipSink> {
334 self.shared.read(|state| state.device(uid)?.v2ip_sink)
335 }
336
337 pub fn v2ip_stats(&self, uid: DeviceUid) -> Option<V2ipDeviceStats> {
339 self.shared.read(|state| state.device(uid)?.v2ip_stats)
340 }
341
342 pub fn v2ip_features(&self, uid: DeviceUid) -> Option<V2ipFpgaFeature> {
350 self.shared.read(|state| state.device(uid)?.v2ip_features)
351 }
352
353 pub fn v2ip_device_settings(&self, uid: DeviceUid) -> Option<V2ipDeviceSettings> {
359 self.shared.read(|state| state.device(uid)?.v2ip_settings)
360 }
361
362 pub fn time_zone(&self, uid: DeviceUid) -> Option<TimeZone> {
366 self.shared
367 .read(|state| state.device(uid)?.time_zone.clone())
368 }
369
370 pub fn device_clock(&self, uid: DeviceUid) -> Option<SystemTime> {
376 let (utc, at) = self.shared.read(|state| state.device(uid)?.clock)?;
377 let announced = UNIX_EPOCH.checked_add(Duration::from_secs(u64::from(utc)))?;
378 announced.checked_add(at.elapsed())
379 }
380
381 pub fn v2ip_tiling(&self, uid: DeviceUid) -> Option<V2ipTilingConfig> {
383 self.shared.read(|state| state.device(uid)?.tiling)
384 }
385
386 pub fn audio_endpoints(&self, uid: DeviceUid) -> Option<AudioEndpoints> {
388 self.shared.read(|state| state.device(uid)?.audio.clone())
389 }
390
391 pub fn multiviewer_status(&self, uid: DeviceUid) -> Option<MultiviewerStatus> {
393 self.shared
394 .read(|state| state.device(uid)?.multiviewer.clone())
395 }
396
397 pub fn dolby_settings(&self, uid: DeviceUid) -> Option<AmpDolbySettings> {
399 self.shared.read(|state| state.device(uid)?.dolby_settings)
400 }
401
402 pub fn rc_settings(&self, uid: DeviceUid) -> Option<RcSettings> {
404 self.shared
405 .read(|state| state.device(uid)?.rc_settings.clone())
406 }
407
408 pub fn network_status(&self, uid: DeviceUid) -> Vec<NetworkPortStatus> {
410 self.shared.read(|state| {
411 state
412 .device(uid)
413 .map(|d| d.network.values().cloned().collect())
414 .unwrap_or_default()
415 })
416 }
417
418 pub fn topology(&self, uid: DeviceUid) -> Vec<TopologyEntry> {
420 self.shared.read(|state| {
421 state
422 .device(uid)
423 .map(|d| d.topology.clone())
424 .unwrap_or_default()
425 })
426 }
427
428 pub fn edid(&self, uid: DeviceUid, output: bool) -> Option<Vec<u8>> {
434 self.shared
435 .read(|state| state.device(uid)?.edid(output).map(<[u8]>::to_vec))
436 }
437
438 pub fn frames_received(&self) -> u64 {
448 self.shared.read(|state| state.frames_received)
449 }
450
451 pub fn firmware(&self, uid: DeviceUid) -> Vec<(FirmwareType, FirmwareVersion)> {
453 self.shared.read(|state| {
454 state
455 .device(uid)
456 .map(|d| d.firmware.iter().map(|(k, v)| (*k, v.clone())).collect())
457 .unwrap_or_default()
458 })
459 }
460
461 pub fn update_config(&self, local_ip: Option<Ipv4Addr>, broadcast: bool) -> io::Result<()> {
466 let network = {
467 let mut network = lock(&self.shared.network);
468 if network.local_ip == local_ip && network.broadcast == broadcast {
469 return Ok(());
470 }
471 network.local_ip = local_ip;
472 network.broadcast = broadcast;
473 network.clone()
474 };
475 let conn = network.open()?;
478 lock(&self.shared.tx).set_conn(Some(conn));
479 self.shared.announce();
480 let _ = self.shared.discover();
481 Ok(())
482 }
483
484 pub fn discover(&self) -> Result<(), SendError> {
486 self.shared.discover()
487 }
488}
489
490impl Drop for Remote {
491 fn drop(&mut self) {
492 self.close();
493 }
494}
495
496impl Shared {
497 fn read<R>(&self, f: impl FnOnce(&State) -> R) -> R {
499 f(&lock(&self.state))
500 }
501
502 fn mutate<R>(&self, f: impl FnOnce(&mut State, &mut Vec<Event>) -> R) -> R {
507 let mut events = Vec::new();
508 let result = f(&mut lock(&self.state), &mut events);
509 self.dispatch(events);
510 result
511 }
512
513 fn dispatch(&self, events: Vec<Event>) {
514 for event in events {
515 event.dispatch(&*self.handler);
516 }
517 }
518
519 fn process_datagram(&self, data: &[u8], from: Ipv4Addr) {
526 let (events, hello_requested) = {
527 let mut state = lock(&self.state);
528 let events = process_frame(&mut state, data, Some(from), Instant::now());
529 (events, std::mem::take(&mut state.hello_requested))
530 };
531 self.dispatch(events);
532 if hello_requested {
533 self.announce();
534 }
535 }
536
537 fn send(&self, to: &Addressee, opcode: Opcode, payload: &[u8]) -> Result<usize, SendError> {
539 lock(&self.tx).send(to, self.uid, opcode, payload)
540 }
541
542 fn discover(&self) -> Result<(), SendError> {
543 lock(&self.schedule).discovered(Instant::now());
544 self.send(&Addressee::Broadcast, op::SYS_DISCOVER, &[])?;
545 Ok(())
546 }
547
548 fn announce(&self) {
557 let payload = build_hello(
558 PROTOCOL_VERSION,
559 &self.name,
560 CLIENT_SERIAL,
561 VERSION,
562 DeviceFeature::MANAGER.bits(),
563 );
564 match self.send(&Addressee::Broadcast, op::SYS_HELLO, &payload) {
565 Ok(n) if n > 0 => lock(&self.schedule).announced(Instant::now()),
566 _ => {}
567 }
568 }
569
570 fn sleep_until_next_tick(&self) -> bool {
575 let deadline = Instant::now() + PROBE_TICK;
576 while Instant::now() < deadline {
577 if self.closing.load(Ordering::SeqCst) {
578 return false;
579 }
580 std::thread::sleep(SHUTDOWN_POLL);
581 }
582 !self.closing.load(Ordering::SeqCst)
583 }
584
585 fn announce_due(&self, now: Instant) -> bool {
592 !self.closing.load(Ordering::SeqCst) && lock(&self.schedule).announce_due(now)
593 }
594
595 fn receive_loop(&self) {
597 let mut buf = vec![0u8; RECV_BUFFER];
598 while !self.closing.load(Ordering::SeqCst) {
599 let Some(conn) = lock(&self.tx).conn() else {
600 break;
601 };
602 match conn.recv(&mut buf) {
603 Ok(Some((data, from))) => self.process_datagram(data, from),
604 Ok(None) => {}
605 Err(_) => break,
606 }
607 }
608 }
609
610 fn probe_loop(&self) {
612 while self.sleep_until_next_tick() {
613 self.probe_once(Instant::now());
614 }
615 }
616
617 pub(super) fn probe_once(&self, now: Instant) {
619 let want_discover = self.mutate(|state, ev| {
620 let mut incomplete = false;
621 let mut any_complete = false;
622 for device in state.devices.values_mut() {
623 device.check_online(now, ev);
624 device.check_config_complete(now, ev);
628 if device.configuration_complete(now) {
629 any_complete = true;
630 } else if now.saturating_duration_since(device.first_seen) > CONFIG_GRACE {
631 incomplete = true;
634 }
635 }
636 incomplete || !any_complete
639 });
640
641 let discover_due = lock(&self.schedule).discover_due(now);
642 if self.announce_due(now) {
643 self.announce();
644 }
645 if want_discover && discover_due {
646 let _ = self.discover();
647 }
648 }
649}
650
651fn lock<T>(m: &Mutex<T>) -> MutexGuard<'_, T> {
658 m.lock().unwrap_or_else(|e| e.into_inner())
659}
660
661fn load_uid(path: Option<PathBuf>) -> io::Result<DeviceUid> {
667 let path = path.or_else(|| {
668 std::env::var_os("HOME")
669 .or_else(|| std::env::var_os("USERPROFILE"))
670 .map(|home| PathBuf::from(home).join(UID_FILE))
671 });
672 if let Some(path) = &path {
673 if let Ok(bytes) = std::fs::read(path) {
674 if let Ok(array) = <[u8; 16]>::try_from(bytes.get(..16).unwrap_or_default()) {
675 return Ok(DeviceUid::from_array(array));
676 }
677 }
678 }
679 let mut bytes = [0u8; 16];
680 getrandom::getrandom(&mut bytes).map_err(|e| io::Error::other(e.to_string()))?;
681 if let Some(path) = &path {
682 let _ = std::fs::write(path, bytes);
683 }
684 Ok(DeviceUid::from_array(bytes))
685}
686
687impl crate::wire::ProtocolTarget for Device {
689 fn serial(&self) -> &str {
690 Device::serial(self)
691 }
692
693 fn supported_protocol(&self) -> u16 {
694 self.hello.supported_protocol
695 }
696}