1//! Live presence and Live Share: who else has the notebook open, the page they are on and
2//! their caret, straight from one Snowbound to another, or through a relay (`crates/relay`)
3//! where they aren't on one network; and a notebook one machine holds opened on another
4//! (`share`). Peers find each other with mDNS (`_snowbound._tcp`) or in the relay's room, and
5//! meet through a secret both hold, a room's or a code typed on both, which SPAKE2 turns into
6//! the keys every frame after the opening is sealed with. Each connection has a thread reading
7//! and one writing. A connection whose frames arrive out of order is dropped and met again
8//! from scratch.
9
10pub use ::relay::code;
11mod group;
12pub mod proxy;
13mod relay;
14pub mod share;
15mod transport;
16pub use transport::Trouble;
17pub mod wire;
18pub use wire::{Caret, Guid, Hello, Presence, Spot};
19
20use mdns_sd::{IfKind, ServiceDaemon, ServiceEvent, ServiceInfo};
21use minicbor::Encode;
22use std::{
23 collections::BTreeMap,
24 io::{self, Read, Write},
25 net::{IpAddr, Ipv4Addr, Shutdown, SocketAddr, TcpListener, TcpStream},
26 sync::{
27 Arc, Mutex,
28 atomic::{AtomicBool, AtomicU64, Ordering},
29 mpsc,
30 },
31 thread,
32 time::{Duration, Instant},
33};
34use wire::{Sealer, Side, kind};
35
36const SERVICE: &str = "_snowbound._tcp.local.";
37/// How long a connection may be quiet before a ping, and before the peer counts as gone.
38const PING: Duration = Duration::from_secs(15);
39const GONE: Duration = Duration::from_secs(45);
40const OPENING: Duration = Duration::from_secs(5);
41/// The most often presence goes to a peer.
42const PRESENCE_EVERY: Duration = Duration::from_millis(100);
43/// The longest wait before meeting again.
44const PATIENCE: Duration = Duration::from_secs(30);
45/// Wrong tries of a code met off any relay before it admits no one new, as a relay burns one.
46const TRIES: u32 = 5;
47
48mod model;
49pub use model::{Event, Peer, Reach, Relayed, Room};
50use model::{code_parts, hex};
51
52/// Presence on the network while it lives; dropping it leaves.
53pub struct Live {
54 shared: Arc<Shared>,
55 address: SocketAddr,
56}
57
58struct Shared {
59 me: Hello,
60 room: Room,
61 secret: Vec<u8>,
62 reach: Option<Reach>,
63 state: Mutex<State>,
64 events: Box<dyn Fn(Event) + Send + Sync>,
65 stopped: AtomicBool,
66 connections: AtomicU64,
67}
68
69#[derive(Default)]
70struct State {
71 presence: Presence,
72 /// Counts changes to `presence`, so a writer sends only the newest.
73 generation: u64,
74 peers: BTreeMap<[u8; 16], Link>,
75 /// The code others type, once known.
76 code: Option<String>,
77 /// The relay connection open now, to hang up on leaving.
78 relay: Option<Arc<relay::Socket>>,
79 relayed: Option<Relayed>,
80 /// Meetings that failed on the secret: wrong codes tried here, or this end's.
81 failed: u32,
82 /// The code admits no one new.
83 burned: bool,
84 /// The Live Share version of the last peer met that speaks another.
85 outdated: Option<u16>,
86 daemon: Option<ServiceDaemon>,
87}
88
89impl State {
90 /// The relay's group, in a notebook's room while the relay is reached.
91 fn group(&self) -> Option<&relay::Group> {
92 self.relay.as_ref()?.group.as_ref()
93 }
94}
95
96/// Presence sent at most every `PRESENCE_EVERY`, and only the newest.
97#[derive(Default)]
98struct Paced {
99 due: Option<Instant>,
100 last: Option<Instant>,
101 sent: Option<u64>,
102}
103
104impl Paced {
105 /// How long to wait for news: until presence is due, else `idle`.
106 fn wait(&self, idle: Duration) -> Duration {
107 self.due
108 .map_or(idle, |due| due.saturating_duration_since(Instant::now()))
109 }
110
111 /// Hears that presence changed.
112 fn changed(&mut self) {
113 let last = self.last;
114 self.due
115 .get_or_insert_with(|| last.map_or_else(Instant::now, |at| at + PRESENCE_EVERY));
116 }
117
118 /// The presence to send now, where it is due and new.
119 fn due(&mut self, state: &Mutex<State>) -> Option<Presence> {
120 self.due.filter(|due| *due <= Instant::now())?;
121 self.due = None;
122 let state = state.lock().unwrap();
123 if self.sent == Some(state.generation) {
124 return None;
125 }
126 (self.sent, self.last) = (Some(state.generation), Some(Instant::now()));
127 Some(state.presence.clone())
128 }
129}
130
131struct Link {
132 connection: u64,
133 peer: Peer,
134 line: Line,
135 pipe: Arc<dyn Pipe>,
136}
137
138/// What a connection's writer sends next.
139enum Out {
140 /// The newest presence.
141 Presence,
142 Frame(u16, Vec<u8>),
143 /// A `Bye`, after which it hangs up.
144 Bye(Vec<u8>),
145}
146
147/// The way to one connected peer: frames sent on it go after those sent before.
148#[derive(Clone)]
149pub struct Line(mpsc::Sender<Out>);
150
151impl Line {
152 pub fn send(&self, kind: u16, body: &impl Encode<()>) -> io::Result<()> {
153 let body = minicbor::to_vec(body).map_err(io::Error::other)?;
154 self.0
155 .send(Out::Frame(kind, body))
156 .map_err(|_| io::ErrorKind::NotConnected.into())
157 }
158
159 /// Says `reason` after the frames sent before, then hangs up.
160 pub fn hang_up(&self, reason: &str) {
161 let bye = minicbor::to_vec(wire::Bye {
162 reason: reason.into(),
163 })
164 .unwrap_or_default();
165 let _ = self.0.send(Out::Bye(bye));
166 }
167}
168
169/// A stream to one peer, read by one thread and written by another: a TCP connection, or
170/// one carried through a relay.
171trait Pipe: Send + Sync {
172 fn read(&self, buffer: &mut [u8]) -> io::Result<usize>;
173 fn write(&self, bytes: &[u8]) -> io::Result<usize>;
174 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()>;
175 /// Hangs up, ending the thread reading.
176 fn shutdown(&self);
177 /// Whether it goes straight to the peer, which is better than through a relay.
178 fn direct(&self) -> bool;
179 /// Hears whether the peer at the other end knew the secret.
180 fn met(&self, _met: bool) {}
181 /// Hears that the stream failed, after which this end meets its peers again.
182 fn broken(&self) {}
183}
184
185impl Pipe for TcpStream {
186 fn read(&self, buffer: &mut [u8]) -> io::Result<usize> {
187 Read::read(&mut &*self, buffer)
188 }
189
190 fn write(&self, bytes: &[u8]) -> io::Result<usize> {
191 Write::write(&mut &*self, bytes)
192 }
193
194 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> {
195 TcpStream::set_read_timeout(self, Some(timeout))
196 }
197
198 fn shutdown(&self) {
199 let _ = TcpStream::shutdown(self, Shutdown::Both);
200 }
201
202 fn direct(&self) -> bool {
203 true
204 }
205}
206
207impl Read for &dyn Pipe {
208 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
209 Pipe::read(*self, buffer)
210 }
211}
212
213impl Write for &dyn Pipe {
214 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
215 Pipe::write(*self, bytes)
216 }
217
218 fn flush(&mut self) -> io::Result<()> {
219 Ok(())
220 }
221}
222
223impl Live {
224 /// Starts listening as `me` in `room`, advertised and looked for where `reach` says, or
225 /// not at all with `None`, leaving peers to `connect`, and in the room at `relay`
226 /// (`wss://live.example.net`) where given. `events` runs on a network thread.
227 pub fn start(
228 me: Hello,
229 room: &Room,
230 reach: Option<Reach>,
231 relay: Option<&str>,
232 events: impl Fn(Event) + Send + Sync + 'static,
233 ) -> io::Result<Live> {
234 let host = match reach {
235 Some(Reach::Network) => IpAddr::V4(Ipv4Addr::UNSPECIFIED),
236 Some(Reach::Loopback) | None => IpAddr::V4(Ipv4Addr::LOCALHOST),
237 };
238 let listener = TcpListener::bind((host, 0))?;
239 let mut address = listener.local_addr()?;
240 if address.ip().is_unspecified() {
241 address.set_ip(IpAddr::V4(Ipv4Addr::LOCALHOST));
242 }
243 let code = match room {
244 Room::Code { code, .. } => match code_parts(code) {
245 (Some(number), secret) => code::format(number, &secret),
246 // Off any relay, the end sharing a secret alone numbers it itself.
247 (None, secret) if relay.is_none() => {
248 let mut number = [0; 4];
249 getrandom::fill(&mut number)
250 .map_err(|_| io::Error::other("System random source failed"))?;
251 code::format(u32::from_le_bytes(number) % code::NAMEPLATES, &secret)
252 }
253 (None, _) => None,
254 },
255 Room::Notebook(_) => None,
256 };
257 let shared = Arc::new(Shared {
258 me,
259 room: room.clone(),
260 secret: room.secret(),
261 reach,
262 state: Mutex::new(State {
263 code,
264 ..State::default()
265 }),
266 events: Box::new(events),
267 stopped: AtomicBool::new(false),
268 connections: AtomicU64::new(0),
269 });
270 if let Some(relay) = relay {
271 relay::join(&shared, relay, address.port())?;
272 }
273 let accepting = Arc::clone(&shared);
274 thread::Builder::new()
275 .name("live accept".into())
276 .spawn(move || {
277 for stream in listener.incoming() {
278 if accepting.stopped.load(Ordering::Acquire) {
279 return;
280 }
281 if let (Ok(stream), Some(tag)) = (stream, accepting.tag()) {
282 let shared = Arc::clone(&accepting);
283 thread::spawn(move || shared.run(Arc::new(stream), Side::Responder, &tag));
284 }
285 }
286 })?;
287 shared.advertise(address.port());
288 Ok(Live { shared, address })
289 }
290
291 /// Where this end listens.
292 pub fn address(&self) -> SocketAddr {
293 self.address
294 }
295
296 /// The code others type to meet this end: the one it was given, or the one the relay
297 /// or this end numbered.
298 pub fn code(&self) -> Option<String> {
299 self.shared.state.lock().unwrap().code.clone()
300 }
301
302 /// How the relay last answered.
303 pub fn relayed(&self) -> Relayed {
304 let state = self.shared.state.lock().unwrap();
305 state.relayed.clone().unwrap_or(Relayed::Unknown)
306 }
307
308 /// Meetings that failed on the secret: wrong tries of this end's code, or this end's own
309 /// wrong code.
310 pub fn failed(&self) -> u32 {
311 self.shared.state.lock().unwrap().failed
312 }
313
314 /// The Live Share version of the last peer met that speaks another, which one of the two
315 /// must update to meet.
316 pub fn other_version(&self) -> Option<u16> {
317 self.shared.state.lock().unwrap().outdated
318 }
319
320 /// Whether this end's code had too many wrong tries and admits no one new.
321 pub fn burned(&self) -> bool {
322 self.shared.state.lock().unwrap().burned
323 }
324
325 /// Connects to a peer at `address` that discovery did not find.
326 pub fn connect(&self, address: SocketAddr) {
327 let shared = Arc::clone(&self.shared);
328 thread::spawn(move || shared.dial(address));
329 }
330
331 /// Says where this end is now; peers hear only the newest of quick changes, at most every
332 /// tenth of a second.
333 pub fn set_presence(&self, presence: Presence) {
334 let mut state = self.shared.state.lock().unwrap();
335 if state.presence == presence {
336 return;
337 }
338 state.presence = presence;
339 state.generation += 1;
340 for link in state.peers.values() {
341 if self.shared.carries_presence(&*link.pipe) {
342 let _ = link.line.0.send(Out::Presence);
343 }
344 }
345 if let Some(group) = state.group() {
346 group.send(relay::Out::Presence);
347 }
348 }
349
350 /// The peers in the room now, by id.
351 pub fn peers(&self) -> Vec<Peer> {
352 self.shared.peers()
353 }
354
355 /// A way to send to this room's peers that doesn't keep it open.
356 pub fn sender(&self) -> Sender {
357 Sender(Arc::downgrade(&self.shared))
358 }
359
360 /// The line to `peer`, while it is connected.
361 pub fn line(&self, peer: &[u8; 16]) -> Option<Line> {
362 let state = self.shared.state.lock().unwrap();
363 state.peers.get(peer).map(|link| link.line.clone())
364 }
365
366 /// Says `reason` to every peer and leaves, waiting a moment for them to hear it.
367 pub fn leave(self, reason: &str) {
368 let bye = minicbor::to_vec(wire::Bye {
369 reason: reason.into(),
370 })
371 .unwrap_or_default();
372 for link in self.shared.state.lock().unwrap().peers.values() {
373 let _ = link.line.0.send(Out::Bye(bye.clone()));
374 }
375 let deadline = Instant::now() + Duration::from_secs(1);
376 while !self.shared.state.lock().unwrap().peers.is_empty() && Instant::now() < deadline {
377 thread::sleep(Duration::from_millis(20));
378 }
379 }
380}
381
382/// Sends to a room's peers while the room is open.
383#[derive(Clone)]
384pub struct Sender(std::sync::Weak<Shared>);
385
386impl Sender {
387 /// Sends message `kind` holding `body` to the peers `to` names, or to everyone: once to
388 /// the room's group through the relay, and to each peer met directly.
389 pub fn send(&self, kind: u16, body: &impl Encode<()>, to: Option<&[[u8; 16]]>) {
390 let (Some(shared), Ok(body)) = (self.0.upgrade(), minicbor::to_vec(body)) else {
391 return;
392 };
393 let state = shared.state.lock().unwrap();
394 let group = state.group();
395 let named = |id: &[u8; 16]| to.is_none_or(|to| to.contains(id));
396 let mut slots = Vec::new();
397 for (id, link) in state.peers.iter().filter(|(id, _)| named(id)) {
398 match group {
399 Some(group) if !link.pipe.direct() => slots.extend(group.slot(id)),
400 _ => {
401 let _ = link.line.0.send(Out::Frame(kind, body.clone()));
402 }
403 }
404 }
405 let Some(group) = group else {
406 return;
407 };
408 match to {
409 None => group.send(relay::Out::Frame(kind, body, None)),
410 Some(to) => {
411 let unlinked = to.iter().filter(|id| !state.peers.contains_key(*id));
412 slots.extend(unlinked.filter_map(|id| group.slot(id)));
413 if !slots.is_empty() {
414 group.send(relay::Out::Frame(kind, body, Some(slots)));
415 }
416 }
417 }
418 }
419}
420
421impl Drop for Live {
422 fn drop(&mut self) {
423 self.shared.stopped.store(true, Ordering::Release);
424 // Wakes the accepting thread to see it has stopped.
425 let _ = TcpStream::connect_timeout(&self.address, OPENING);
426 let mut state = self.shared.state.lock().unwrap();
427 if let Some(daemon) = state.daemon.take() {
428 let _ = daemon.shutdown();
429 }
430 for link in state.peers.values() {
431 link.pipe.shutdown();
432 }
433 if let Some(relay) = &state.relay {
434 relay.hang_up();
435 }
436 }
437}
438
439impl Shared {
440 /// The peers in the room: those met directly, and the rest as the relay's group or a
441 /// stream through the relay last heard of them.
442 fn peers(&self) -> Vec<Peer> {
443 let state = self.state.lock().unwrap();
444 let mut peers = BTreeMap::new();
445 if let Some(group) = state.group() {
446 for member in group.members.lock().unwrap().values() {
447 if let Some(peer) = &member.peer {
448 peers.insert(peer.hello.peer, peer.clone());
449 }
450 }
451 }
452 for (id, link) in &state.peers {
453 if link.pipe.direct() || !peers.contains_key(id) {
454 peers.insert(*id, link.peer.clone());
455 }
456 }
457 peers.into_values().collect()
458 }
459
460 /// Whether peer `id` is met directly, which is where it says everything.
461 fn direct(&self, id: &[u8; 16]) -> bool {
462 let state = self.state.lock().unwrap();
463 state.peers.get(id).is_some_and(|link| link.pipe.direct())
464 }
465
466 /// Whether presence goes on `pipe`: not on a stream through the relay in a notebook's
467 /// room, whose group carries it.
468 fn carries_presence(&self, pipe: &dyn Pipe) -> bool {
469 pipe.direct() || matches!(self.room, Room::Code { .. })
470 }
471
472 /// The room's tag as it stands: a code's, once numbered.
473 fn tag(&self) -> Option<String> {
474 match &self.room {
475 Room::Code { .. } => code_parts(self.state.lock().unwrap().code.as_deref()?)
476 .0
477 .map(|number| format!("code-{number}")),
478 room => room.tag(),
479 }
480 }
481
482 /// Advertises this end on `port` once its room has a tag, where its reach says, and
483 /// connects to the peers in its room that discovery finds with a higher id than its own,
484 /// which leave the connecting to it.
485 fn advertise(self: &Arc<Self>, port: u16) {
486 let (Some(reach), Some(tag)) = (self.reach, self.tag()) else {
487 return;
488 };
489 let mut state = self.state.lock().unwrap();
490 if state.daemon.is_some() || self.stopped.load(Ordering::Acquire) {
491 return;
492 }
493 match discover(self, reach, tag, port) {
494 Ok(daemon) => state.daemon = Some(daemon),
495 Err(error) => eprintln!("Live: no discovery: {error}"),
496 }
497 }
498
499 fn knows(&self, peer: &str) -> bool {
500 let state = self.state.lock().unwrap();
501 state.peers.keys().any(|id| hex(id) == peer)
502 }
503
504 /// Connects to `address`, and again while the peer is there and the connection was the
505 /// one this end kept.
506 fn dial(self: Arc<Self>, address: SocketAddr) {
507 let mut wait = Duration::from_secs(1);
508 while let Ok(stream) = TcpStream::connect_timeout(&address, OPENING) {
509 let Some(tag) = self.tag() else {
510 return;
511 };
512 if !Arc::clone(&self).run(Arc::new(stream), Side::Initiator, &tag)
513 || self.stopped.load(Ordering::Acquire)
514 {
515 return;
516 }
517 thread::sleep(wait);
518 wait = (wait * 2).min(PATIENCE);
519 }
520 }
521
522 /// Counts a meeting that failed on the secret: the end sharing a code burns it after too
523 /// many, and an end that typed one gives up at once.
524 fn failed(&self) {
525 let mut state = self.state.lock().unwrap();
526 state.failed += 1;
527 match &self.room {
528 Room::Code { owner: true, .. } if state.failed >= TRIES => state.burned = true,
529 Room::Code { owner: false, .. } => self.stopped.store(true, Ordering::Release),
530 _ => {}
531 }
532 drop(state);
533 (self.events)(Event::Changed);
534 }
535
536 /// Meets the peer at the other end of `pipe` in room `tag`, then reads from it until it
537 /// goes: whether it was the connection kept to that peer.
538 fn run(self: Arc<Self>, pipe: Arc<dyn Pipe>, side: Side, tag: &str) -> bool {
539 if self.stopped.load(Ordering::Acquire) || self.state.lock().unwrap().burned {
540 pipe.shutdown();
541 return false;
542 }
543 let mut stream: &dyn Pipe = &*pipe;
544 let met = (|| {
545 stream.set_read_timeout(OPENING)?;
546 let (mut send, mut receive) = wire::open(&mut stream, side, tag, &self.secret)?;
547 send.send(&mut stream, kind::HELLO, &self.me)?;
548 let (first, body) = receive.receive(&mut stream)?;
549 if first != kind::HELLO {
550 return Err(io::Error::new(io::ErrorKind::InvalidData, "No hello"));
551 }
552 let hello: Hello = minicbor::decode(&body)
553 .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "A malformed hello"))?;
554 stream.set_read_timeout(GONE)?;
555 Ok((send, receive, hello))
556 })();
557 let other = (met.as_ref().err())
558 .and_then(|error| error.get_ref()?.downcast_ref::<wire::Version>())
559 .copied();
560 let (send, mut receive, hello) = match met {
561 Ok(met) => {
562 pipe.met(true);
563 met
564 }
565 Err(error) => {
566 eprintln!("Live: no meeting in {tag}: {error}");
567 if let Some(wire::Version(version)) = other {
568 // Another version guessed nothing: no keys were agreed.
569 pipe.met(true);
570 self.state.lock().unwrap().outdated = Some(version);
571 if matches!(self.room, Room::Code { owner: false, .. }) {
572 self.stopped.store(true, Ordering::Release);
573 }
574 (self.events)(Event::Changed);
575 } else if error.kind() == io::ErrorKind::InvalidData {
576 pipe.met(false);
577 self.failed();
578 } else {
579 pipe.broken();
580 }
581 pipe.shutdown();
582 return false;
583 }
584 };
585 let peer = hello.peer;
586 let name = hello.name.clone();
587 let hello = Arc::new(hello);
588 let connection = self.connections.fetch_add(1, Ordering::Relaxed);
589 let (out, outgoing) = mpsc::channel();
590 let line = Line(out);
591 {
592 let mut state = self.state.lock().unwrap();
593 // A peer met both directly and through a relay keeps the direct connection, as
594 // both ends then agree.
595 let kept = match state.peers.get(&peer) {
596 _ if peer == self.me.peer => false,
597 Some(link) if link.pipe.direct() || !pipe.direct() => false,
598 Some(link) => {
599 link.pipe.shutdown();
600 true
601 }
602 None => true,
603 };
604 if !kept {
605 drop(state);
606 pipe.shutdown();
607 return false;
608 }
609 if self.carries_presence(&*pipe) {
610 let _ = line.0.send(Out::Presence);
611 }
612 state.peers.insert(
613 peer,
614 Link {
615 connection,
616 peer: Peer {
617 hello: Arc::clone(&hello),
618 presence: None,
619 },
620 line: line.clone(),
621 pipe: Arc::clone(&pipe),
622 },
623 );
624 }
625 let writing = Arc::clone(&pipe);
626 let shared = Arc::clone(&self);
627 thread::spawn(move || shared.write(&*writing, send, outgoing));
628 (self.events)(Event::Met(&hello, &line));
629 (self.events)(Event::Changed);
630 let ended = loop {
631 let (message, body) = match receive.receive(&mut stream) {
632 Ok(frame) => frame,
633 Err(error) => break Some(error),
634 };
635 match message {
636 kind::HELLO | kind::PING => {}
637 kind::PRESENCE => {
638 let Ok(presence) = minicbor::decode::<Presence>(&body) else {
639 break None;
640 };
641 if let Some(link) = self.state.lock().unwrap().peers.get_mut(&peer) {
642 link.peer.presence = Some(presence);
643 }
644 (self.events)(Event::Changed);
645 }
646 kind => {
647 (self.events)(Event::Frame {
648 from: &hello,
649 kind,
650 body: &body,
651 });
652 // The peer hangs up after its bye.
653 if kind == kind::BYE {
654 break None;
655 }
656 }
657 }
658 };
659 if let Some(error) = ended {
660 eprintln!("Live: the connection to {name} broke ({error}); meeting again");
661 pipe.broken();
662 }
663 pipe.shutdown();
664 let mut state = self.state.lock().unwrap();
665 if state
666 .peers
667 .get(&peer)
668 .is_some_and(|link| link.connection == connection)
669 {
670 state.peers.remove(&peer);
671 drop(state);
672 (self.events)(Event::Left(&hello));
673 (self.events)(Event::Changed);
674 }
675 true
676 }
677
678 /// Sends the newest presence when woken, at most every `PRESENCE_EVERY`, frames as they
679 /// come, and a ping when quiet.
680 fn write(&self, pipe: &dyn Pipe, mut send: Sealer, outgoing: mpsc::Receiver<Out>) {
681 let mut stream = pipe;
682 let mut paced = Paced::default();
683 loop {
684 let result = match outgoing.recv_timeout(paced.wait(PING)) {
685 Ok(Out::Presence) => {
686 paced.changed();
687 Ok(())
688 }
689 Ok(Out::Frame(kind, body)) => send.send_encoded(&mut stream, kind, &body),
690 Ok(Out::Bye(body)) => {
691 let _ = send.send_encoded(&mut stream, kind::BYE, &body);
692 pipe.shutdown();
693 return;
694 }
695 Err(mpsc::RecvTimeoutError::Timeout) if paced.due.is_none() => {
696 send.send(&mut stream, kind::PING, &())
697 }
698 Err(mpsc::RecvTimeoutError::Timeout) => Ok(()),
699 Err(mpsc::RecvTimeoutError::Disconnected) => return,
700 };
701 let result = result.and_then(|()| match paced.due(&self.state) {
702 Some(presence) => send.send(&mut stream, kind::PRESENCE, &presence),
703 None => Ok(()),
704 });
705 if result.is_err() {
706 pipe.broken();
707 pipe.shutdown();
708 return;
709 }
710 }
711 }
712}
713
714/// Advertises `shared` as in room `tag` on `port` and connects to the peers in its room that
715/// discovery finds with a higher id than its own.
716fn discover(
717 shared: &Arc<Shared>,
718 reach: Reach,
719 tag: String,
720 port: u16,
721) -> mdns_sd::Result<ServiceDaemon> {
722 let daemon = ServiceDaemon::new()?;
723 let id = hex(&shared.me.peer);
724 let properties = [("v", "1"), ("room", tag.as_str()), ("peer", id.as_str())];
725 let host = format!("snowbound-{id}.local.");
726 let info = match reach {
727 Reach::Network => {
728 ServiceInfo::new(SERVICE, &id, &host, "", port, &properties[..])?.enable_addr_auto()
729 }
730 Reach::Loopback => {
731 daemon.disable_interface(IfKind::All)?;
732 daemon.enable_interface(IfKind::LoopbackV4)?;
733 ServiceInfo::new(
734 SERVICE,
735 &id,
736 &host,
737 IpAddr::V4(Ipv4Addr::LOCALHOST),
738 port,
739 &properties[..],
740 )?
741 }
742 };
743 daemon.register(info)?;
744 let found = daemon.browse(SERVICE)?;
745 let shared = Arc::downgrade(shared);
746 thread::Builder::new()
747 .name("live discovery".into())
748 .spawn(move || {
749 while let Ok(event) = found.recv() {
750 let ServiceEvent::ServiceResolved(service) = event else {
751 continue;
752 };
753 let Some(shared) = shared.upgrade() else {
754 return;
755 };
756 let (Some(room), Some(peer)) = (
757 service.get_property_val_str("room"),
758 service.get_property_val_str("peer"),
759 ) else {
760 continue;
761 };
762 if room != tag || peer <= id.as_str() || shared.knows(peer) {
763 continue;
764 }
765 let mut addresses: Vec<IpAddr> =
766 service.addresses.iter().map(|ip| ip.to_ip_addr()).collect();
767 addresses.sort_by_key(|ip| (!ip.is_ipv4(), !ip.is_loopback()));
768 if let Some(ip) = addresses.first() {
769 let address = SocketAddr::new(*ip, service.port);
770 thread::spawn(move || shared.dial(address));
771 }
772 }
773 })
774 .map_err(|error| mdns_sd::Error::Msg(error.to_string()))?;
775 Ok(daemon)
776}
777
778/// `text` with `%XX` escapes decoded, as a URL's name and password.
779fn decode(text: &str) -> String {
780 let bytes = text.as_bytes();
781 let mut decoded = Vec::with_capacity(bytes.len());
782 let mut at = 0;
783 while at < bytes.len() {
784 match text
785 .get(at + 1..at + 3)
786 .filter(|_| bytes[at] == b'%')
787 .and_then(|hex| u8::from_str_radix(hex, 16).ok())
788 {
789 Some(byte) => {
790 decoded.push(byte);
791 at += 3;
792 }
793 None => {
794 decoded.push(bytes[at]);
795 at += 1;
796 }
797 }
798 }
799 String::from_utf8_lossy(&decoded).into_owned()
800}
801
802#[cfg(test)]
803mod tests;