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