1//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying
2//! streams to peers in it, each opened with SPAKE2 and sealed end to end as on a LAN, so the
3//! relay sees only who talks to whom, when, and how much. In a code's room every two peers
4//! have a stream. In a notebook's room only a host and each guest do; everything else goes
5//! to the room's group (`group`), sent once and copied by the relay. A stream or group frame
6//! that breaks makes this end join the room again, which ends every stream it had there;
7//! peers then meet afresh.
8
9use super::{
10 Event, Hello, OPENING, PATIENCE, Paced, Peer, Pipe, Presence, Relayed as Answer, Room, Shared,
11 Side, code_parts, group,
12 transport::{self, Address, Failure, parse},
13 wire::kind,
14};
15use ::relay::{BROADCAST, GROUP, Notice, SLOT, Verdict, ws};
16use std::{
17 collections::{HashMap, HashSet, hash_map::Entry},
18 io::{self, BufReader, Read, Write},
19 sync::{Arc, Mutex, atomic::Ordering, mpsc},
20 thread,
21 time::{Duration, Instant},
22};
23
24/// How often a quiet connection pings the relay.
25const KEEPALIVE: Duration = Duration::from_secs(30);
26/// A connection that lasted this long was no failure, so the next waits only a second.
27const STEADY: Duration = Duration::from_secs(60);
28/// The most sent in one message, well under any relay's cap.
29const CHUNK: usize = 64 << 10;
30/// The largest message read from the relay.
31const MOST: usize = 1 << 20;
32
33type Reader = ws::Reader<BufReader<Box<dyn Read + Send>>>;
34
35/// Joins `shared`'s room at the relay at `url`, and keeps joining whenever it falls out;
36/// `port` is where it listens, to advertise once the relay numbers its code.
37pub(super) fn join(shared: &Arc<Shared>, url: &str, port: u16) -> io::Result<()> {
38 let address = parse(url)?;
39 let shared = Arc::clone(shared);
40 thread::Builder::new()
41 .name("live relay".into())
42 .spawn(move || keep(&shared, &address, port))?;
43 Ok(())
44}
45
46/// Records how the relay answered, telling `shared`'s events where it changed.
47fn answered(shared: &Shared, answer: Answer) {
48 let mut state = shared.state.lock().unwrap();
49 if state.relayed.as_ref() != Some(&answer) {
50 state.relayed = Some(answer);
51 drop(state);
52 (shared.events)(Event::Changed);
53 }
54}
55
56/// Stays in the room at `address`, joining again whenever the connection ends, until
57/// `shared` stops.
58fn keep(shared: &Arc<Shared>, address: &Address, port: u16) {
59 let mut wait = Duration::from_secs(1);
60 let owner = shared.room.owner();
61 // The number of this end's code, asked for again on joining again.
62 let mut nameplate = owner
63 .then(|| code_parts(shared.state.lock().unwrap().code.as_deref()?).0)
64 .flatten();
65 while !shared.stopped.load(Ordering::Acquire) {
66 let began = Instant::now();
67 let path = match (owner, shared.tag()) {
68 (false, Some(tag)) => format!("{}/v1/room/{tag}", address.path),
69 (false, None) => return,
70 (true, _) => match nameplate {
71 Some(number) => format!("{}/v1/claim?nameplate={number}", address.path),
72 None => format!("{}/v1/claim", address.path),
73 },
74 };
75 match connect(shared, address, &path, owner) {
76 Ok((socket, reader)) => {
77 shared.state.lock().unwrap().relay = Some(Arc::clone(&socket));
78 answered(shared, Answer::Joined);
79 if !shared.stopped.load(Ordering::Acquire) {
80 session(shared, &socket, reader, &mut nameplate, port);
81 }
82 shared.state.lock().unwrap().relay = None;
83 socket.hang_up();
84 }
85 Err(Failure::Refused(status, retry)) => {
86 eprintln!("Live: the relay refused to let this end in ({status})");
87 answered(shared, Answer::Refused(status, retry));
88 if status == 410 {
89 return;
90 }
91 wait = wait.max(retry.unwrap_or_default());
92 }
93 Err(Failure::Trouble(trouble, error)) => {
94 eprintln!("Live: no relay at {}: {error}", address.authority);
95 answered(shared, Answer::Unreachable(trouble));
96 }
97 }
98 if began.elapsed() >= STEADY {
99 wait = Duration::from_secs(1);
100 }
101 thread::sleep(wait);
102 wait = (wait * 2).min(PATIENCE);
103 }
104}
105
106/// Opens the relay's `path` at `address`: the connection, and what reads its messages.
107fn connect(
108 shared: &Shared,
109 address: &Address,
110 path: &str,
111 owner: bool,
112) -> Result<(Arc<Socket>, Reader), Failure> {
113 let connection = transport::connect(address, path)?;
114 let reader = ws::Reader::new(BufReader::new(connection.reader), MOST, false);
115 let group = matches!(shared.room, Room::Notebook(_)).then(|| Group {
116 keys: group::Keys::new(&shared.secret),
117 members: Mutex::default(),
118 out: Mutex::default(),
119 });
120 let socket = Arc::new(Socket {
121 send: Mutex::new(connection.writer),
122 close: connection.close,
123 owner,
124 links: Mutex::default(),
125 group,
126 });
127 Ok((socket, reader))
128}
129
130/// One connection to a relay's room.
131pub(super) struct Socket {
132 send: Mutex<Box<dyn Write + Send>>,
133 close: Box<dyn Fn() + Send + Sync>,
134 /// Whether this end claimed the room for its code, and so tells the relay who knew it.
135 owner: bool,
136 links: Mutex<Links>,
137 /// A notebook's room's group.
138 pub(super) group: Option<Group>,
139}
140
141/// A notebook's room's group: its key, the peers heard in it, and the way to this end's
142/// thread sending to it.
143pub(super) struct Group {
144 keys: group::Keys,
145 /// The peers heard in the group, by slot.
146 pub(super) members: Mutex<HashMap<u32, group::Member>>,
147 /// To the thread sending this end's group frames, while connected.
148 out: Mutex<Option<mpsc::Sender<Out>>>,
149}
150
151/// What this end sends the group next.
152pub(super) enum Out {
153 /// The newest presence, to everyone.
154 Presence,
155 /// A message kind and its encoded body, to everyone or to the peers in these slots.
156 Frame(u16, Vec<u8>, Option<Vec<u32>>),
157}
158
159impl Group {
160 pub(super) fn send(&self, out: Out) {
161 if let Some(sender) = &*self.out.lock().unwrap() {
162 let _ = sender.send(out);
163 }
164 }
165
166 /// The slot of the member that is peer `id`.
167 pub(super) fn slot(&self, id: &[u8; 16]) -> Option<u32> {
168 let members = self.members.lock().unwrap();
169 members
170 .iter()
171 .find_map(|(slot, member)| (member.hello()?.peer == *id).then_some(*slot))
172 }
173}
174
175/// The streams a socket carries, by the slot of the peer at the other end.
176#[derive(Default)]
177struct Links {
178 /// This end's slot.
179 me: u32,
180 inboxes: HashMap<u32, mpsc::Sender<Vec<u8>>>,
181 /// Slots whose stream ended; a peer that comes back has a new slot.
182 ended: HashSet<u32>,
183}
184
185impl Socket {
186 fn send(&self, opcode: u8, payload: &[u8]) -> io::Result<()> {
187 let mut mask = [0; 4];
188 getrandom::fill(&mut mask).map_err(|_| io::Error::other("System random source failed"))?;
189 let frame = ws::frame(opcode, payload, Some(mask));
190 self.send.lock().unwrap().write_all(&frame)
191 }
192
193 pub(super) fn hang_up(&self) {
194 (self.close)();
195 }
196
197 fn forget(&self, slot: u32) {
198 let mut links = self.links.lock().unwrap();
199 links.inboxes.remove(&slot);
200 links.ended.insert(slot);
201 }
202}
203
204/// Reads notices and peers' bytes from the relay until the connection ends, opening a stream
205/// to each peer in the room.
206fn session(
207 shared: &Arc<Shared>,
208 socket: &Arc<Socket>,
209 mut reader: Reader,
210 nameplate: &mut Option<u32>,
211 port: u16,
212) {
213 let mut tag = shared.tag();
214 // Pings while the session lasts: dropping `_beat` at its end stops them.
215 let (_beat, beats) = mpsc::channel::<()>();
216 let pinging = Arc::clone(socket);
217 thread::spawn(move || {
218 while let Err(mpsc::RecvTimeoutError::Timeout) = beats.recv_timeout(KEEPALIVE) {
219 if pinging.send(ws::PING, &[]).is_err() {
220 return;
221 }
222 }
223 });
224 if let Some(group) = &socket.group {
225 let (out, outgoing) = mpsc::channel();
226 *group.out.lock().unwrap() = Some(out);
227 let (shared, socket) = (Arc::clone(shared), Arc::clone(socket));
228 thread::spawn(move || speak(&shared, &socket, outgoing));
229 }
230 let notebook = socket.group.is_some();
231 while let Ok(message) = reader.read() {
232 match message {
233 ws::Message::Text(text) => match text.parse() {
234 Ok(Notice::Nameplate(number)) => {
235 *nameplate = Some(number);
236 tag = Some(format!("code-{number}"));
237 let super::Room::Code { code, .. } = &shared.room else {
238 continue;
239 };
240 let Some(code) = super::code::format(number, &code_parts(code).1) else {
241 continue;
242 };
243 let mut state = shared.state.lock().unwrap();
244 if state.code.as_ref() != Some(&code) {
245 state.code = Some(code);
246 drop(state);
247 shared.advertise(port);
248 (shared.events)(Event::Changed);
249 }
250 }
251 Ok(Notice::Welcome { you, members }) => {
252 socket.links.lock().unwrap().me = you;
253 for slot in members.into_iter().filter(|_| !notebook) {
254 meet(shared, socket, tag.as_deref(), slot, None);
255 }
256 }
257 Ok(Notice::Joined(slot)) if !notebook => {
258 meet(shared, socket, tag.as_deref(), slot, None);
259 }
260 Ok(Notice::Joined(_)) => {}
261 Ok(Notice::Left(slot)) => {
262 socket.forget(slot);
263 let left = (socket.group.as_ref())
264 .and_then(|group| group.members.lock().unwrap().remove(&slot));
265 if left.is_some() {
266 (shared.events)(Event::Changed);
267 }
268 }
269 Ok(Notice::Burned) => {
270 // Coming back, ask for another number.
271 *nameplate = None;
272 shared.state.lock().unwrap().burned = true;
273 (shared.events)(Event::Changed);
274 }
275 Err(()) => {}
276 },
277 ws::Message::Binary(data) => match data.split_first_chunk::<SLOT>() {
278 Some((slot, frame)) if u32::from_be_bytes(*slot) & GROUP != 0 => {
279 let slot = u32::from_be_bytes(*slot) & !GROUP;
280 heard(shared, socket, tag.as_deref(), slot, frame);
281 }
282 Some((slot, bytes)) => {
283 let slot = u32::from_be_bytes(*slot);
284 meet(shared, socket, tag.as_deref(), slot, Some(bytes.to_vec()));
285 }
286 None => {}
287 },
288 ws::Message::Ping(payload) => {
289 let _ = socket.send(ws::PONG, &payload);
290 }
291 ws::Message::Pong => {}
292 ws::Message::Close => break,
293 }
294 }
295 socket.links.lock().unwrap().inboxes.clear();
296 if let Some(group) = &socket.group {
297 *group.out.lock().unwrap() = None;
298 group.members.lock().unwrap().clear();
299 (shared.events)(Event::Changed);
300 }
301}
302
303/// Sends this end's group frames: its hello, asking everyone for theirs, then its presence
304/// at most every `PRESENCE_EVERY`, and frames as they come.
305fn speak(shared: &Shared, socket: &Socket, outgoing: mpsc::Receiver<Out>) {
306 let Some(group) = &socket.group else {
307 return;
308 };
309 let Ok(mut sealer) = group::Sealer::new(&group.keys) else {
310 return;
311 };
312 let mut send = |kind: u16, body: &[u8], to: Option<&[u32]>| -> io::Result<()> {
313 let sealed = sealer.seal(kind, body, to.is_none())?;
314 let mut message = match to {
315 None => BROADCAST.to_be_bytes().to_vec(),
316 Some(slots) => {
317 let mut message = (GROUP | slots.len() as u32).to_be_bytes().to_vec();
318 slots
319 .iter()
320 .for_each(|slot| message.extend(slot.to_be_bytes()));
321 message
322 }
323 };
324 message.extend(sealed);
325 socket.send(ws::BINARY, &message)
326 };
327 let hello = minicbor::to_vec(&shared.me).unwrap_or_default();
328 if send(kind::HELLO, &hello, None).is_err() {
329 return;
330 }
331 let mut paced = Paced::default();
332 paced.changed();
333 loop {
334 let next = match paced.due {
335 Some(_) => outgoing.recv_timeout(paced.wait(Duration::ZERO)),
336 None => outgoing
337 .recv()
338 .map_err(|_| mpsc::RecvTimeoutError::Disconnected),
339 };
340 let result = match next {
341 Ok(Out::Presence) => {
342 paced.changed();
343 Ok(())
344 }
345 Ok(Out::Frame(kind, body, to)) => send(kind, &body, to.as_deref()),
346 Err(mpsc::RecvTimeoutError::Timeout) => Ok(()),
347 Err(mpsc::RecvTimeoutError::Disconnected) => return,
348 };
349 let result = result.and_then(|()| match paced.due(&shared.state) {
350 Some(presence) => send(
351 kind::PRESENCE,
352 &minicbor::to_vec(&presence).unwrap_or_default(),
353 None,
354 ),
355 None => Ok(()),
356 });
357 if result.is_err() {
358 return;
359 }
360 }
361}
362
363/// Takes in a group frame from the peer in `slot`: its hello, answered with this end's and
364/// meeting it where it serves; its presence; or what else it says.
365fn heard(shared: &Arc<Shared>, socket: &Arc<Socket>, tag: Option<&str>, slot: u32, frame: &[u8]) {
366 let Some(group) = &socket.group else {
367 return;
368 };
369 let opened = {
370 let mut members = group.members.lock().unwrap();
371 let member = match members.entry(slot) {
372 Entry::Occupied(member) => Ok(member.into_mut()),
373 Entry::Vacant(vacant) => {
374 group::Member::new(&group.keys, frame).map(|m| vacant.insert(m))
375 }
376 };
377 member.and_then(|member| {
378 let (kind, body) = member.open(frame)?;
379 Ok((kind, body, member.hello().cloned()))
380 })
381 };
382 let (kind, body, hello) = match opened {
383 Ok(opened) => opened,
384 Err(error) => {
385 eprintln!("Live: the relay broke the room's frames ({error}); meeting again");
386 socket.hang_up();
387 return;
388 }
389 };
390 match kind {
391 kind::HELLO | kind::HELLO_BACK => {
392 let Ok(hello) = minicbor::decode::<Hello>(&body) else {
393 return;
394 };
395 if hello.peer == shared.me.peer {
396 return;
397 }
398 let serves = hello.serves.is_some() && shared.me.serves.is_none();
399 if let Some(member) = group.members.lock().unwrap().get_mut(&slot) {
400 member.peer = Some(Peer {
401 hello: Arc::new(hello),
402 presence: None,
403 });
404 }
405 if kind == kind::HELLO {
406 let presence = minicbor::to_vec(&shared.state.lock().unwrap().presence);
407 let presence = presence.unwrap_or_default();
408 let me = minicbor::to_vec(&shared.me).unwrap_or_default();
409 group.send(Out::Frame(kind::HELLO_BACK, me, Some(vec![slot])));
410 group.send(Out::Frame(kind::PRESENCE, presence, Some(vec![slot])));
411 }
412 if serves {
413 meet(shared, socket, tag, slot, None);
414 }
415 (shared.events)(Event::Changed);
416 }
417 kind::PRESENCE => {
418 let Ok(presence) = minicbor::decode::<Presence>(&body) else {
419 return;
420 };
421 let mut members = group.members.lock().unwrap();
422 if let Some(peer) = members.get_mut(&slot).and_then(|m| m.peer.as_mut()) {
423 peer.presence = Some(presence);
424 drop(members);
425 (shared.events)(Event::Changed);
426 }
427 }
428 // A peer met directly says the rest there.
429 kind if kind < 256 => {
430 if let Some(hello) = hello.filter(|hello| !shared.direct(&hello.peer)) {
431 (shared.events)(Event::Frame {
432 from: &hello,
433 kind,
434 body: &body,
435 });
436 }
437 }
438 _ => {}
439 }
440}
441
442/// Hands `bytes` from the peer in `slot` to its stream, or opens one. In a code's room this
443/// end opens a stream to a peer with a lower slot when told of it, and answers one with a
444/// higher slot when its first bytes come; in a notebook's room a guest opens one to its host.
445fn meet(
446 shared: &Arc<Shared>,
447 socket: &Arc<Socket>,
448 tag: Option<&str>,
449 slot: u32,
450 bytes: Option<Vec<u8>>,
451) {
452 let mut links = socket.links.lock().unwrap();
453 if let Some(inbox) = links.inboxes.get(&slot) {
454 if let Some(bytes) = bytes {
455 let _ = inbox.send(bytes);
456 }
457 return;
458 }
459 let notebook = socket.group.is_some();
460 let side = match bytes {
461 None if notebook || slot < links.me => Side::Initiator,
462 Some(_) if (notebook && shared.me.serves.is_some()) || (!notebook && slot > links.me) => {
463 Side::Responder
464 }
465 _ => return,
466 };
467 let Some(tag) = tag.filter(|_| !links.ended.contains(&slot)) else {
468 return;
469 };
470 let (inbox, arriving) = mpsc::channel();
471 if let Some(bytes) = bytes {
472 let _ = inbox.send(bytes);
473 }
474 links.inboxes.insert(slot, inbox);
475 let pipe = Arc::new(Relayed {
476 socket: Arc::clone(socket),
477 slot,
478 inbox: Mutex::new(Inbox {
479 arriving,
480 chunk: Vec::new(),
481 at: 0,
482 }),
483 timeout: Mutex::new(OPENING),
484 });
485 let shared = Arc::clone(shared);
486 let tag = tag.to_owned();
487 thread::spawn(move || shared.run(pipe, side, &tag));
488}
489
490/// The stream to one peer through the relay.
491struct Relayed {
492 socket: Arc<Socket>,
493 slot: u32,
494 inbox: Mutex<Inbox>,
495 timeout: Mutex<Duration>,
496}
497
498struct Inbox {
499 arriving: mpsc::Receiver<Vec<u8>>,
500 chunk: Vec<u8>,
501 at: usize,
502}
503
504impl Pipe for Relayed {
505 fn read(&self, buffer: &mut [u8]) -> io::Result<usize> {
506 let mut inbox = self.inbox.lock().unwrap();
507 while inbox.at == inbox.chunk.len() {
508 let timeout = *self.timeout.lock().unwrap();
509 inbox.chunk = match inbox.arriving.recv_timeout(timeout) {
510 Ok(chunk) => chunk,
511 Err(mpsc::RecvTimeoutError::Timeout) => return Err(io::ErrorKind::TimedOut.into()),
512 Err(mpsc::RecvTimeoutError::Disconnected) => return Ok(0),
513 };
514 inbox.at = 0;
515 }
516 let at = inbox.at;
517 let length = buffer.len().min(inbox.chunk.len() - at);
518 buffer[..length].copy_from_slice(&inbox.chunk[at..at + length]);
519 inbox.at += length;
520 Ok(length)
521 }
522
523 fn write(&self, bytes: &[u8]) -> io::Result<usize> {
524 if self.socket.links.lock().unwrap().ended.contains(&self.slot) {
525 return Err(io::ErrorKind::BrokenPipe.into());
526 }
527 let length = bytes.len().min(CHUNK);
528 let message = [&self.slot.to_be_bytes()[..], &bytes[..length]].concat();
529 self.socket.send(ws::BINARY, &message)?;
530 Ok(length)
531 }
532
533 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> {
534 *self.timeout.lock().unwrap() = timeout;
535 Ok(())
536 }
537
538 fn shutdown(&self) {
539 self.socket.forget(self.slot);
540 }
541
542 fn direct(&self) -> bool {
543 false
544 }
545
546 fn met(&self, met: bool) {
547 if self.socket.owner {
548 let verdict = if met { Verdict::Met } else { Verdict::Failed }(self.slot);
549 let _ = self.socket.send(ws::TEXT, verdict.to_string().as_bytes());
550 }
551 }
552
553 fn broken(&self) {
554 // A deliberate shutdown can wake the reader with EOF.
555 if !self.socket.links.lock().unwrap().ended.contains(&self.slot) {
556 self.socket.hang_up();
557 }
558 }
559}
560
561#[cfg(test)]
562mod tests {
563 use super::*;
564
565 #[test]
566 fn relay_addresses_parse() {
567 let address = parse("wss://live.example.net").unwrap();
568 assert_eq!(
569 address,
570 Address {
571 tls: true,
572 authority: "live.example.net".into(),
573 host: "live.example.net".into(),
574 port: 443,
575 path: String::new(),
576 }
577 );
578 let address = parse("ws://[::1]:7650/snowbound/").unwrap();
579 assert_eq!((address.host.as_str(), address.port), ("::1", 7650));
580 assert_eq!(address.path, "/snowbound");
581 assert!(parse("https://live.example.net").is_err());
582 assert!(parse("ws://:80").is_err());
583 }
584}