authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-02 03:24:43-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-02 03:45:09-07:00
loge2e4ea381f12668cf85d93e6afeaf73b824a56e9
tree345d3cb73b88b410e683586e4f70966c3a9a3279
parentbb9f2fba980bc24c51881194b08afaa2ef2f0bc3
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

feat: Live Share meets through a relay off the LAN, and no relay can tamper with it unnoticed

- snowbound-relay (crates/relay): a WebSocket rendezvous and relay, one static Linux executable for x86_64 and aarch64 (tools/release_relay.py), with a systemd unit and a deploy README for Caddy or nginx in front - The relay limits wrong codes without learning them: the code's owner tells it who met, ten wrong codes a minute lock an address out for longer each time, and five burn a code - Every live frame carries its number inside the seal, so a frame dropped, repeated, reordered or forged on the way is caught, never applied, and the peers meet again from scratch - Presence reaches peers through a relay named by SNOWBOUND_LIVE_RELAY, by notebook or by a code the relay numbers (feature live) Assisted-by: claude-opus-5.5

18 files changed, 3114 insertions(+), 104 deletions(-)

AGENTS.md+2-1
......@@ -13,7 +13,8 @@ novel UI kit, each new platform port is extremely lightweight.
1313| Path | Owns | Depends on |
1414| --- | --- | --- |
1515| `crates/onestore` | The file format: revision stores (`.one`, `.onetoc2`), the page model, ops, the commit protocol. No network, no SQLite, no `unsafe`. | none |
16| `crates/notebook` | From editor to disk or share: discovery, notebook structure, sessions, the SQLite replica, sync and merging, conflict pages, the embedded SMB client (feature `smb`), live presence between peers (feature `live`). | onestore |
16| `crates/notebook` | From editor to disk or share: discovery, notebook structure, sessions, the SQLite replica, sync and merging, conflict pages, the embedded SMB client (feature `smb`), live presence between peers (feature `live`). | onestore, relay |
17| `crates/relay` | `snowbound-relay`, the WebSocket relay Live Share meets through off the LAN, and the framing and notices its clients share with it. Sees only sealed frames. | none |
1718| `crates/draw` | The wgpu renderer that page and chrome both paint through, and the text-editing core (keys, chords, carets) they share. | none |
1819| `crates/canvas` | The page: editor, OneNote-faithful layout, page scene, interaction, the page's accessibility tree. | onestore, draw |
1920| `crates/ui` | The immediate-mode interface kit and OneNote's chrome controls. Knows nothing of notebooks. | draw |
Cargo.lock+13
......@@ -2189,8 +2189,11 @@ dependencies = [
21892189 "minicbor",
21902190 "nix",
21912191 "onestore",
2192 "relay",
21922193 "rsqlite-vfs",
21932194 "rusqlite",
2195 "rustls",
2196 "rustls-native-certs",
21942197 "serde",
21952198 "serde_json",
21962199 "sha2",
......@@ -2202,6 +2205,7 @@ dependencies = [
22022205 "wasm-bindgen",
22032206 "wasm-bindgen-futures",
22042207 "web-time",
2208 "webpki-root-certs",
22052209 "windows-sys 0.61.2",
22062210 "zeroize",
22072211]
......@@ -3132,6 +3136,15 @@ dependencies = [
31323136 "bitflags 2.13.1",
31333137]
31343138
3139[[package]]
3140name = "relay"
3141version = "0.1.0"
3142dependencies = [
3143 "base64",
3144 "getrandom 0.4.3",
3145 "sha1",
3146]
3147
31353148[[package]]
31363149name = "renderdoc-sys"
31373150version = "1.1.0"
crates/notebook/Cargo.toml+8-3
......@@ -6,9 +6,9 @@ publish = false
66
77[features]
88smb = ["dep:smb2", "dep:tokio"]
9# Live presence: peers found on the network, met through a shared secret, over an encrypted
10# stream (resources/live-share.md).
11live = ["dep:mdns-sd", "dep:minicbor", "dep:spake2"]
9# Live presence: peers found on the network or met through a relay, through a shared secret,
10# over an encrypted stream (resources/live-share.md).
11live = ["dep:mdns-sd", "dep:minicbor", "dep:spake2", "dep:relay", "dep:rustls", "dep:rustls-native-certs", "dep:webpki-root-certs"]
1212
1313[dependencies]
1414onestore = { path = "../onestore" }
......@@ -29,6 +29,11 @@ hmac = "0.13.0"
2929mdns-sd = { version = "0.21.4", default-features = false, optional = true }
3030minicbor = { version = "2.3.0", features = ["derive", "alloc"], optional = true }
3131spake2 = { version = "=0.5.0-pre.0", features = ["getrandom"], optional = true }
32relay = { path = "../relay", optional = true }
33# A relay's TLS, trusting what the app's updates trust: the system's authorities, then Mozilla's.
34rustls = { version = "0.23", default-features = false, features = ["ring", "std", "tls12"], optional = true }
35rustls-native-certs = { version = "0.8", optional = true }
36webpki-root-certs = { version = "1.0", optional = true }
3237# OneNote packages (.onepkg) are cabinet files.
3338cab = "0.6"
3439# std's clock where there is one; the browser's on wasm32-unknown-unknown, whose std has none.
crates/notebook/README.md+16-10
......@@ -531,16 +531,22 @@ physical power-loss durability.
531531
532532## Live presence (feature `live`)
533533
534`live::Live::start(hello, room, reach, notify)` listens on a TCP port and, with a `Reach`,
535advertises `_snowbound._tcp` by mDNS on every network or on loopback alone, connecting to the
536peers in the same `Room` that it finds: a notebook's identity, or a code typed on both
537(`Room::Code("4-violet-otter")`). `connect(address)` meets a peer discovery did not find.
538Peers meet through SPAKE2 on the room's secret, then every frame is AES-256-GCM under the keys
539it agreed: a message kind and a CBOR map (`live::wire`). A reader skips kinds and map keys it
540doesn't know, so later versions add both freely. `set_presence` says which section, page and
541caret this end has (text object and UTF-16 offset, as ops address text); a connection sends
542only the newest. `peers()` lists each connected peer's `Hello` (name, picture) and presence,
543and `notify` runs whenever that changes. Dropping the `Live` leaves.
534`live::Live::start(hello, room, reach, relay, notify)` listens on a TCP port and, with a
535`Reach`, advertises `_snowbound._tcp` by mDNS on every network or on loopback alone,
536connecting to the peers in the same `Room` that it finds: a notebook's identity, or a code
537typed on both (`Room::Code("4-violet-otter")`). With a `relay` (`wss://live.example.net`,
538`crates/relay`) it also joins the room there and meets its peers through it; a code's words
539alone (`Room::Code("violet-otter")`) ask the relay for a number, and `code()` then has the
540whole code. `connect(address)` meets a peer discovery did not find. Peers meet through
541SPAKE2 on the room's secret, then every frame is AES-256-GCM under the keys it agreed: its
542number, which is also its nonce, then a message kind and a CBOR map (`live::wire`). A frame
543lost, repeated, reordered or forged on the way fails where it lands; the connection is
544dropped as broken, nothing from it after the fault is applied, and the ends meet again from
545scratch. A reader skips kinds and map keys it doesn't know, so later versions add both
546freely. `set_presence` says which section, page and caret this end has (text object and
547UTF-16 offset, as ops address text); a connection sends only the newest. `peers()` lists
548each connected peer's `Hello` (name, picture) and presence, and `notify` runs whenever that
549changes. Dropping the `Live` leaves.
544550
545551## Queue measurement
546552
crates/notebook/src/live.rs+204-60
......@@ -1,9 +1,12 @@
11//! Live presence: who else has the notebook open, the page they are on and their caret,
2//! straight from one Snowbound to another. Peers find each other with mDNS
3//! (`_snowbound._tcp`) and meet through a secret both hold, a notebook's identity or a code
4//! typed on both, which SPAKE2 turns into the keys every frame after the opening is sealed
5//! with. The lower peer id connects; each connection has a thread reading and one writing.
2//! straight from one Snowbound to another, or through a relay (`crates/relay`) where they
3//! aren't on one network. Peers find each other with mDNS (`_snowbound._tcp`) or in the
4//! relay's room, and meet through a secret both hold, a notebook's identity or a code typed
5//! on both, which SPAKE2 turns into the keys every frame after the opening is sealed with.
6//! Each connection has a thread reading and one writing. A connection whose frames arrive
7//! out of order is dropped and met again from scratch.
68
9mod relay;
710pub mod wire;
811pub use wire::{Caret, Guid, Hello, Presence, Spot};
912
......@@ -11,7 +14,7 @@ use mdns_sd::{IfKind, ServiceDaemon, ServiceEvent, ServiceInfo};
1114use sha2::{Digest, Sha256};
1215use std::{
1316 collections::BTreeMap,
14 io,
17 io::{self, Read, Write},
1518 net::{IpAddr, Ipv4Addr, Shutdown, SocketAddr, TcpListener, TcpStream},
1619 sync::{
1720 Arc, Mutex,
......@@ -28,6 +31,8 @@ const SERVICE: &str = "_snowbound._tcp.local.";
2831const PING: Duration = Duration::from_secs(15);
2932const GONE: Duration = Duration::from_secs(45);
3033const OPENING: Duration = Duration::from_secs(5);
34/// The longest wait before meeting again.
35const PATIENCE: Duration = Duration::from_secs(30);
3136
3237/// The secret peers meet through.
3338#[derive(Clone, Debug, PartialEq, Eq)]
......@@ -36,27 +41,40 @@ pub enum Room {
3641 /// hold.
3742 Notebook([u8; 16]),
3843 /// A code typed on both: `7-violet-otter`, whose number names it on the network and whose
39 /// words only the two people know.
44 /// words only the two people know. Its words alone (`violet-otter`) ask the relay for a
45 /// free number, and `Live::code` then has the whole code.
4046 Code(String),
4147}
4248
4349impl Room {
44 /// What names the room in the clear: a hash of a notebook's identity, a code's number.
45 fn tag(&self) -> String {
50 /// What names the room in the clear: a hash of a notebook's identity, a code's number;
51 /// none for a code the relay hasn't numbered.
52 fn tag(&self) -> Option<String> {
4653 match self {
47 Room::Notebook(id) => hex(&Sha256::digest([&b"Snowbound room "[..], id].concat())[..8]),
48 Room::Code(code) => format!("code-{}", code.split('-').next().unwrap_or_default()),
54 Room::Notebook(id) => Some(hex(&Sha256::digest(
55 [&b"Snowbound room "[..], id].concat(),
56 )[..8])),
57 Room::Code(code) => code_parts(code).0.map(|number| format!("code-{number}")),
4958 }
5059 }
5160
5261 fn secret(&self) -> Vec<u8> {
5362 match self {
5463 Room::Notebook(id) => id.to_vec(),
55 Room::Code(code) => code.trim().to_lowercase().into_bytes(),
64 Room::Code(code) => code_parts(code).1.to_lowercase().into_bytes(),
5665 }
5766 }
5867}
5968
69/// A code's number, if it has one, and its words.
70fn code_parts(code: &str) -> (Option<u32>, &str) {
71 let code = code.trim();
72 match code.split_once('-') {
73 Some((number, words)) if number.parse::<u32>().is_ok() => (number.parse().ok(), words),
74 _ => (None, code),
75 }
76}
77
6078/// Where peers are looked for.
6179#[derive(Clone, Copy, Debug, PartialEq, Eq)]
6280pub enum Reach {
......@@ -82,7 +100,7 @@ pub struct Live {
82100
83101struct Shared {
84102 me: Hello,
85 tag: String,
103 room: Room,
86104 secret: Vec<u8>,
87105 state: Mutex<State>,
88106 notify: Box<dyn Fn() + Send + Sync>,
......@@ -96,25 +114,93 @@ struct State {
96114 /// Counts changes to `presence`, so a writer sends only the newest.
97115 generation: u64,
98116 peers: BTreeMap<[u8; 16], Link>,
117 /// The code others type, once known.
118 code: Option<String>,
119 /// The relay connection open now, to hang up on leaving.
120 relay: Option<Arc<relay::Socket>>,
99121}
100122
101123struct Link {
102124 connection: u64,
103125 peer: Peer,
104126 wake: mpsc::Sender<()>,
105 stream: TcpStream,
127 pipe: Arc<dyn Pipe>,
128}
129
130/// A stream to one peer, read by one thread and written by another: a TCP connection, or
131/// one carried through a relay.
132trait Pipe: Send + Sync {
133 fn read(&self, buffer: &mut [u8]) -> io::Result<usize>;
134 fn write(&self, bytes: &[u8]) -> io::Result<usize>;
135 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()>;
136 /// Hangs up, ending the thread reading.
137 fn shutdown(&self);
138 /// Whether it goes straight to the peer, which is better than through a relay.
139 fn direct(&self) -> bool;
140 /// Hears whether the peer at the other end knew the secret.
141 fn met(&self, _met: bool) {}
142 /// Hears that frames arrived lost, repeated, reordered or forged, after which this end
143 /// meets its peers again.
144 fn broken(&self) {}
145}
146
147impl Pipe for TcpStream {
148 fn read(&self, buffer: &mut [u8]) -> io::Result<usize> {
149 Read::read(&mut &*self, buffer)
150 }
151
152 fn write(&self, bytes: &[u8]) -> io::Result<usize> {
153 Write::write(&mut &*self, bytes)
154 }
155
156 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> {
157 TcpStream::set_read_timeout(self, Some(timeout))
158 }
159
160 fn shutdown(&self) {
161 let _ = TcpStream::shutdown(self, Shutdown::Both);
162 }
163
164 fn direct(&self) -> bool {
165 true
166 }
167}
168
169impl Read for &dyn Pipe {
170 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
171 Pipe::read(*self, buffer)
172 }
173}
174
175impl Write for &dyn Pipe {
176 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
177 Pipe::write(*self, bytes)
178 }
179
180 fn flush(&mut self) -> io::Result<()> {
181 Ok(())
182 }
106183}
107184
108185impl Live {
109186 /// Starts listening as `me` in `room`, advertised and looked for where `reach` says, or
110 /// not at all with `None`, leaving peers to `connect`. `notify` runs on a network thread
111 /// whenever `peers` changes.
187 /// not at all with `None`, leaving peers to `connect`, and in the room at `relay`
188 /// (`wss://live.example.net`) where given. `notify` runs on a network thread whenever
189 /// `peers` or `code` changes.
112190 pub fn start(
113191 me: Hello,
114192 room: &Room,
115193 reach: Option<Reach>,
194 relay: Option<&str>,
116195 notify: impl Fn() + Send + Sync + 'static,
117196 ) -> io::Result<Live> {
197 let tag = room.tag();
198 if tag.is_none() && relay.is_none() {
199 return Err(io::Error::new(
200 io::ErrorKind::InvalidInput,
201 "Only a relay can number a code",
202 ));
203 }
118204 let host = match reach {
119205 Some(Reach::Network) => IpAddr::V4(Ipv4Addr::UNSPECIFIED),
120206 Some(Reach::Loopback) | None => IpAddr::V4(Ipv4Addr::LOCALHOST),
......@@ -124,16 +210,27 @@ impl Live {
124210 if address.ip().is_unspecified() {
125211 address.set_ip(IpAddr::V4(Ipv4Addr::LOCALHOST));
126212 }
213 let code = match room {
214 Room::Code(code) if tag.is_some() => Some(code.trim().to_owned()),
215 _ => None,
216 };
127217 let shared = Arc::new(Shared {
128218 me,
129 tag: room.tag(),
219 room: room.clone(),
130220 secret: room.secret(),
131 state: Mutex::default(),
221 state: Mutex::new(State {
222 code,
223 ..State::default()
224 }),
132225 notify: Box::new(notify),
133226 stopped: AtomicBool::new(false),
134227 connections: AtomicU64::new(0),
135228 });
229 if let Some(relay) = relay {
230 relay::join(&shared, relay)?;
231 }
136232 let accepting = Arc::clone(&shared);
233 let accepted = tag.clone();
137234 thread::Builder::new()
138235 .name("live accept".into())
139236 .spawn(move || {
......@@ -141,17 +238,17 @@ impl Live {
141238 if accepting.stopped.load(Ordering::Acquire) {
142239 return;
143240 }
144 if let Ok(stream) = stream {
241 if let (Ok(stream), Some(tag)) = (stream, accepted.clone()) {
145242 let shared = Arc::clone(&accepting);
146 thread::spawn(move || shared.run(stream, Side::Responder));
243 thread::spawn(move || shared.run(Arc::new(stream), Side::Responder, &tag));
147244 }
148245 }
149246 })?;
150 let daemon = match reach {
151 Some(reach) => {
152 Some(advertise(&shared, reach, address.port()).map_err(io::Error::other)?)
247 let daemon = match (reach, tag) {
248 (Some(reach), Some(tag)) => {
249 Some(advertise(&shared, reach, tag, address.port()).map_err(io::Error::other)?)
153250 }
154 None => None,
251 _ => None,
155252 };
156253 Ok(Live {
157254 shared,
......@@ -165,6 +262,12 @@ impl Live {
165262 self.address
166263 }
167264
265 /// The code others type to meet this end: the one it was given, or the one the relay
266 /// numbered.
267 pub fn code(&self) -> Option<String> {
268 self.shared.state.lock().unwrap().code.clone()
269 }
270
168271 /// Connects to a peer at `address` that discovery did not find.
169272 pub fn connect(&self, address: SocketAddr) {
170273 let shared = Arc::clone(&self.shared);
......@@ -199,22 +302,27 @@ impl Drop for Live {
199302 }
200303 // Wakes the accepting thread to see it has stopped.
201304 let _ = TcpStream::connect_timeout(&self.address, OPENING);
202 for link in self.shared.state.lock().unwrap().peers.values() {
203 let _ = link.stream.shutdown(Shutdown::Both);
305 let state = self.shared.state.lock().unwrap();
306 for link in state.peers.values() {
307 link.pipe.shutdown();
308 }
309 if let Some(relay) = &state.relay {
310 relay.hang_up();
204311 }
205312 }
206313}
207314
208/// Advertises `shared` on `port` and connects to the peers in its room that discovery finds
209/// with a higher id than its own, which leave the connecting to it.
210fn advertise(shared: &Arc<Shared>, reach: Reach, port: u16) -> mdns_sd::Result<ServiceDaemon> {
315/// Advertises `shared` as in room `tag` on `port` and connects to the peers in its room that
316/// discovery finds with a higher id than its own, which leave the connecting to it.
317fn advertise(
318 shared: &Arc<Shared>,
319 reach: Reach,
320 tag: String,
321 port: u16,
322) -> mdns_sd::Result<ServiceDaemon> {
211323 let daemon = ServiceDaemon::new()?;
212324 let id = hex(&shared.me.peer);
213 let properties = [
214 ("v", "1"),
215 ("room", shared.tag.as_str()),
216 ("peer", id.as_str()),
217 ];
325 let properties = [("v", "1"), ("room", tag.as_str()), ("peer", id.as_str())];
218326 let host = format!("snowbound-{id}.local.");
219327 let info = match reach {
220328 Reach::Network => {
......@@ -249,7 +357,7 @@ fn advertise(shared: &Arc<Shared>, reach: Reach, port: u16) -> mdns_sd::Result<S
249357 ) else {
250358 continue;
251359 };
252 if room != shared.tag || peer <= id.as_str() || shared.knows(peer) {
360 if room != tag || peer <= id.as_str() || shared.knows(peer) {
253361 continue;
254362 }
255363 let mut addresses: Vec<IpAddr> =
......@@ -272,21 +380,35 @@ impl Shared {
272380 state.peers.keys().any(|id| hex(id) == peer)
273381 }
274382
383 /// Connects to `address`, and again while the peer is there and the connection was the
384 /// one this end kept.
275385 fn dial(self: Arc<Self>, address: SocketAddr) {
276 if let Ok(stream) = TcpStream::connect_timeout(&address, OPENING) {
277 self.run(stream, Side::Initiator);
386 let Some(tag) = self.room.tag() else {
387 return;
388 };
389 let mut wait = Duration::from_secs(1);
390 while let Ok(stream) = TcpStream::connect_timeout(&address, OPENING) {
391 if !Arc::clone(&self).run(Arc::new(stream), Side::Initiator, &tag)
392 || self.stopped.load(Ordering::Acquire)
393 {
394 return;
395 }
396 thread::sleep(wait);
397 wait = (wait * 2).min(PATIENCE);
278398 }
279399 }
280400
281 /// Meets the peer at the other end of `stream`, then reads from it until it goes.
282 fn run(self: Arc<Self>, mut stream: TcpStream, side: Side) {
401 /// Meets the peer at the other end of `pipe` in room `tag`, then reads from it until it
402 /// goes: whether it was the connection kept to that peer.
403 fn run(self: Arc<Self>, pipe: Arc<dyn Pipe>, side: Side, tag: &str) -> bool {
283404 if self.stopped.load(Ordering::Acquire) {
284 return;
405 pipe.shutdown();
406 return false;
285407 }
408 let mut stream: &dyn Pipe = &*pipe;
286409 let met = (|| {
287 stream.set_read_timeout(Some(OPENING))?;
288 stream.set_nodelay(true)?;
289 let (mut send, mut receive) = wire::open(&mut stream, side, &self.tag, &self.secret)?;
410 stream.set_read_timeout(OPENING)?;
411 let (mut send, mut receive) = wire::open(&mut stream, side, tag, &self.secret)?;
290412 send.send(&mut stream, kind::HELLO, &self.me)?;
291413 let (first, body) = receive.receive(&mut stream)?;
292414 if first != kind::HELLO {
......@@ -294,27 +416,40 @@ impl Shared {
294416 }
295417 let hello: Hello = minicbor::decode(&body)
296418 .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "A malformed hello"))?;
297 stream.set_read_timeout(Some(GONE))?;
419 stream.set_read_timeout(GONE)?;
298420 Ok((send, receive, hello))
299421 })();
422 pipe.met(met.is_ok());
300423 let (send, mut receive, hello) = match met {
301424 Ok(met) => met,
302425 Err(error) => {
303 eprintln!("Live: no meeting with {:?}: {error}", stream.peer_addr());
304 return;
426 eprintln!("Live: no meeting in {tag}: {error}");
427 pipe.shutdown();
428 return false;
305429 }
306430 };
307431 let peer = hello.peer;
432 let name = hello.name.clone();
308433 let connection = self.connections.fetch_add(1, Ordering::Relaxed);
309434 let (wake, woken) = mpsc::channel();
310435 {
311436 let mut state = self.state.lock().unwrap();
312 if peer == self.me.peer || state.peers.contains_key(&peer) {
313 return;
314 }
315 let Ok(writing) = stream.try_clone() else {
316 return;
437 // A peer met both directly and through a relay keeps the direct connection, as
438 // both ends then agree.
439 let kept = match state.peers.get(&peer) {
440 _ if peer == self.me.peer => false,
441 Some(link) if link.pipe.direct() || !pipe.direct() => false,
442 Some(link) => {
443 link.pipe.shutdown();
444 true
445 }
446 None => true,
317447 };
448 if !kept {
449 drop(state);
450 pipe.shutdown();
451 return false;
452 }
318453 let _ = wake.send(());
319454 state.peers.insert(
320455 peer,
......@@ -325,28 +460,35 @@ impl Shared {
325460 presence: None,
326461 },
327462 wake,
328 stream: writing,
463 pipe: Arc::clone(&pipe),
329464 },
330465 );
331466 }
332467 (self.notify)();
333 if let Ok(writing) = stream.try_clone() {
334 let shared = Arc::clone(&self);
335 thread::spawn(move || shared.write(writing, send, woken));
336 }
337 while let Ok((message, body)) = receive.receive(&mut stream) {
468 let writing = Arc::clone(&pipe);
469 let shared = Arc::clone(&self);
470 thread::spawn(move || shared.write(&*writing, send, woken));
471 let ended = loop {
472 let (message, body) = match receive.receive(&mut stream) {
473 Ok(frame) => frame,
474 Err(error) => break Some(error),
475 };
338476 if message != kind::PRESENCE {
339477 continue;
340478 }
341479 let Ok(presence) = minicbor::decode::<Presence>(&body) else {
342 break;
480 break None;
343481 };
344482 if let Some(link) = self.state.lock().unwrap().peers.get_mut(&peer) {
345483 link.peer.presence = Some(presence);
346484 }
347485 (self.notify)();
486 };
487 if let Some(error) = ended.filter(|error| error.kind() == io::ErrorKind::InvalidData) {
488 eprintln!("Live: the connection to {name} broke ({error}); meeting again");
489 pipe.broken();
348490 }
349 let _ = stream.shutdown(Shutdown::Both);
491 pipe.shutdown();
350492 let mut state = self.state.lock().unwrap();
351493 if state
352494 .peers
......@@ -357,10 +499,12 @@ impl Shared {
357499 drop(state);
358500 (self.notify)();
359501 }
502 true
360503 }
361504
362505 /// Sends the newest presence whenever woken, and a ping when quiet.
363 fn write(&self, mut stream: TcpStream, mut send: Sealer, woken: mpsc::Receiver<()>) {
506 fn write(&self, pipe: &dyn Pipe, mut send: Sealer, woken: mpsc::Receiver<()>) {
507 let mut stream = pipe;
364508 let mut sent = None;
365509 loop {
366510 let result = match woken.recv_timeout(PING) {
......@@ -379,7 +523,7 @@ impl Shared {
379523 Err(mpsc::RecvTimeoutError::Disconnected) => return,
380524 };
381525 if result.is_err() {
382 let _ = stream.shutdown(Shutdown::Both);
526 pipe.shutdown();
383527 return;
384528 }
385529 }
crates/notebook/src/live/relay.rs created+537
......@@ -0,0 +1,537 @@
1//! Meeting peers through a relay (`crates/relay`): one WebSocket to the room, carrying a
2//! stream to each peer in it, each opened with SPAKE2 and sealed end to end as on a LAN, so
3//! the relay sees only who talks to whom, when, and how much. A stream that breaks makes this
4//! end join the room again, which ends every stream it had there; peers then meet afresh.
5
6use super::{OPENING, PATIENCE, Pipe, Shared, Side, code_parts};
7use ::relay::{Notice, SLOT, Verdict, ws};
8use base64::Engine;
9use rustls::{ClientConfig, ClientConnection, RootCertStore, pki_types::ServerName};
10use std::{
11 collections::{HashMap, HashSet},
12 io::{self, BufReader, Read, Write},
13 net::{Shutdown, TcpStream, ToSocketAddrs},
14 sync::{Arc, Mutex, OnceLock, atomic::Ordering, mpsc},
15 thread,
16 time::{Duration, Instant},
17};
18
19const CONNECT: Duration = Duration::from_secs(10);
20/// How often a quiet connection pings the relay, and how long it waits to hear anything.
21const KEEPALIVE: Duration = Duration::from_secs(30);
22const QUIET: Duration = Duration::from_secs(75);
23/// A connection that lasted this long was no failure, so the next waits only a second.
24const STEADY: Duration = Duration::from_secs(60);
25/// The most sent in one message, well under any relay's cap.
26const CHUNK: usize = 64 << 10;
27/// The largest message read from the relay.
28const MOST: usize = 1 << 20;
29
30type Reader = ws::Reader<BufReader<Box<dyn Read + Send>>>;
31
32/// Joins `shared`'s room at the relay at `url`, and keeps joining whenever it falls out.
33pub(super) fn join(shared: &Arc<Shared>, url: &str) -> io::Result<()> {
34 let address = parse(url)?;
35 let shared = Arc::clone(shared);
36 thread::Builder::new()
37 .name("live relay".into())
38 .spawn(move || keep(&shared, &address))?;
39 Ok(())
40}
41
42/// A relay's address: `ws://` or `wss://`, a host, and the path the relay's `/v1/` follows.
43#[derive(Debug, PartialEq)]
44struct Address {
45 tls: bool,
46 /// The host and port as the URL gave them, for the `Host` header.
47 authority: String,
48 host: String,
49 port: u16,
50 path: String,
51}
52
53fn parse(url: &str) -> io::Result<Address> {
54 let bad = || io::Error::new(io::ErrorKind::InvalidInput, format!("Not a relay: {url}"));
55 let (tls, rest) = match url.split_once("://") {
56 Some(("wss", rest)) => (true, rest),
57 Some(("ws", rest)) => (false, rest),
58 _ => return Err(bad()),
59 };
60 let (authority, path) = rest.split_at(rest.find('/').unwrap_or(rest.len()));
61 let (host, port) = match authority.rsplit_once(':') {
62 Some((host, port)) if !port.contains(']') => (host, port.parse().map_err(|_| bad())?),
63 _ => (authority, if tls { 443 } else { 80 }),
64 };
65 let host = host.trim_start_matches('[').trim_end_matches(']');
66 if host.is_empty() {
67 return Err(bad());
68 }
69 Ok(Address {
70 tls,
71 authority: authority.into(),
72 host: host.into(),
73 port,
74 path: path.trim_end_matches('/').into(),
75 })
76}
77
78enum Failure {
79 Network(io::Error),
80 /// The relay answered, but with this HTTP status and how long to wait.
81 Refused(u16, Option<Duration>),
82}
83
84impl From<io::Error> for Failure {
85 fn from(error: io::Error) -> Self {
86 Failure::Network(error)
87 }
88}
89
90/// Stays in the room at `address`, joining again whenever the connection ends, until
91/// `shared` stops.
92fn keep(shared: &Arc<Shared>, address: &Address) {
93 let mut wait = Duration::from_secs(1);
94 // The number the relay gave this end's code, asked for again on joining again.
95 let mut nameplate = None;
96 while !shared.stopped.load(Ordering::Acquire) {
97 let began = Instant::now();
98 let tag = shared.room.tag();
99 let path = match &tag {
100 Some(tag) => format!("{}/v1/room/{tag}", address.path),
101 None => match nameplate {
102 Some(number) => format!("{}/v1/claim?nameplate={number}", address.path),
103 None => format!("{}/v1/claim", address.path),
104 },
105 };
106 match connect(address, &path, tag.is_none()) {
107 Ok((socket, reader)) => {
108 shared.state.lock().unwrap().relay = Some(Arc::clone(&socket));
109 if !shared.stopped.load(Ordering::Acquire) {
110 session(shared, &socket, reader, &mut nameplate);
111 }
112 shared.state.lock().unwrap().relay = None;
113 socket.hang_up();
114 }
115 Err(Failure::Refused(410, _)) => {
116 eprintln!("Live: the code has expired; ask for a new one");
117 return;
118 }
119 Err(Failure::Refused(status, retry)) => {
120 eprintln!("Live: the relay refused to let this end in ({status})");
121 wait = wait.max(retry.unwrap_or_default());
122 }
123 Err(Failure::Network(error)) => {
124 eprintln!("Live: no relay at {}: {error}", address.authority);
125 }
126 }
127 if began.elapsed() >= STEADY {
128 wait = Duration::from_secs(1);
129 }
130 thread::sleep(wait);
131 wait = (wait * 2).min(PATIENCE);
132 }
133}
134
135/// Opens a WebSocket to `path` at `address`: the connection, and its reading half.
136fn connect(address: &Address, path: &str, owner: bool) -> Result<(Arc<Socket>, Reader), Failure> {
137 let target = (address.host.as_str(), address.port)
138 .to_socket_addrs()?
139 .next()
140 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "No address for the relay"))?;
141 let tcp = TcpStream::connect_timeout(&target, CONNECT)?;
142 tcp.set_nodelay(true)?;
143 tcp.set_read_timeout(Some(QUIET))?;
144 tcp.set_write_timeout(Some(QUIET))?;
145 let (reading, mut writing): (Box<dyn Read + Send>, Box<dyn Write + Send>) = if address.tls {
146 let name = ServerName::try_from(address.host.clone())
147 .map_err(|error| io::Error::new(io::ErrorKind::InvalidInput, error))?;
148 let mut connection = ClientConnection::new(tls(), name).map_err(io::Error::other)?;
149 while connection.is_handshaking() {
150 connection.complete_io(&mut &tcp)?;
151 }
152 let connection = Arc::new(Mutex::new(connection));
153 (
154 Box::new(TlsReader {
155 tcp: tcp.try_clone()?,
156 tls: Arc::clone(&connection),
157 plain: Vec::new(),
158 at: 0,
159 }),
160 Box::new(TlsWriter {
161 tcp: tcp.try_clone()?,
162 tls: connection,
163 }),
164 )
165 } else {
166 (Box::new(tcp.try_clone()?), Box::new(tcp.try_clone()?))
167 };
168 let mut key = [0; 16];
169 getrandom::fill(&mut key).map_err(|_| io::Error::other("System random source failed"))?;
170 let key = base64::engine::general_purpose::STANDARD.encode(key);
171 write!(
172 writing,
173 "GET {path} HTTP/1.1\r\nHost: {}\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
174 Sec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\nUser-Agent: Snowbound/{}\r\n\r\n",
175 address.authority,
176 env!("CARGO_PKG_VERSION"),
177 )?;
178 let mut reader = BufReader::new(reading);
179 let head = ws::head(&mut reader)?;
180 let status = head
181 .split(' ')
182 .nth(1)
183 .and_then(|status| status.parse().ok())
184 .unwrap_or(0);
185 if status != 101 {
186 let retry = ws::header(&head, "Retry-After")
187 .and_then(|seconds| seconds.parse().ok())
188 .map(Duration::from_secs);
189 return Err(Failure::Refused(status, retry));
190 }
191 if ws::header(&head, "Sec-WebSocket-Accept") != Some(ws::accept(&key).as_str()) {
192 return Err(io::Error::new(io::ErrorKind::InvalidData, "Not a relay").into());
193 }
194 let socket = Arc::new(Socket {
195 send: Mutex::new(writing),
196 tcp,
197 owner,
198 links: Mutex::default(),
199 });
200 Ok((socket, ws::Reader::new(reader, MOST, false)))
201}
202
203/// The certificate authorities the system trusts, then Mozilla's for a system whose store is
204/// missing or stale, as updates trust them.
205fn tls() -> Arc<ClientConfig> {
206 static CONFIG: OnceLock<Arc<ClientConfig>> = OnceLock::new();
207 Arc::clone(CONFIG.get_or_init(|| {
208 let mut roots = RootCertStore::empty();
209 roots.add_parsable_certificates(rustls_native_certs::load_native_certs().certs);
210 roots.add_parsable_certificates(webpki_root_certs::TLS_SERVER_ROOT_CERTS.iter().cloned());
211 let provider = Arc::new(rustls::crypto::ring::default_provider());
212 Arc::new(
213 ClientConfig::builder_with_provider(provider)
214 .with_safe_default_protocol_versions()
215 .expect("ring speaks TLS 1.2 and 1.3")
216 .with_root_certificates(roots)
217 .with_no_client_auth(),
218 )
219 }))
220}
221
222/// TLS read on one thread while another writes: the socket is read without the lock, and
223/// what arrives is decrypted under it.
224struct TlsReader {
225 tcp: TcpStream,
226 tls: Arc<Mutex<ClientConnection>>,
227 plain: Vec<u8>,
228 at: usize,
229}
230
231impl Read for TlsReader {
232 fn read(&mut self, buffer: &mut [u8]) -> io::Result<usize> {
233 while self.at == self.plain.len() {
234 self.plain.clear();
235 self.at = 0;
236 let mut raw = vec![0; 16 << 10];
237 let length = self.tcp.read(&mut raw)?;
238 if length == 0 {
239 return Ok(0);
240 }
241 let mut tls = self.tls.lock().unwrap();
242 let mut arrived = &raw[..length];
243 let mut closed = false;
244 while !arrived.is_empty() {
245 tls.read_tls(&mut arrived)?;
246 tls.process_new_packets()
247 .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?;
248 let mut chunk = [0; 4096];
249 loop {
250 match tls.reader().read(&mut chunk) {
251 Ok(0) => {
252 closed = true;
253 break;
254 }
255 Ok(length) => self.plain.extend_from_slice(&chunk[..length]),
256 Err(error) if error.kind() == io::ErrorKind::WouldBlock => break,
257 Err(error) => return Err(error),
258 }
259 }
260 }
261 while tls.wants_write() {
262 tls.write_tls(&mut &self.tcp)?;
263 }
264 if closed && self.plain.is_empty() {
265 return Ok(0);
266 }
267 }
268 let length = buffer.len().min(self.plain.len() - self.at);
269 buffer[..length].copy_from_slice(&self.plain[self.at..self.at + length]);
270 self.at += length;
271 Ok(length)
272 }
273}
274
275struct TlsWriter {
276 tcp: TcpStream,
277 tls: Arc<Mutex<ClientConnection>>,
278}
279
280impl Write for TlsWriter {
281 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
282 let mut tls = self.tls.lock().unwrap();
283 let length = tls.writer().write(bytes)?;
284 while tls.wants_write() {
285 tls.write_tls(&mut &self.tcp)?;
286 }
287 Ok(length)
288 }
289
290 fn flush(&mut self) -> io::Result<()> {
291 Ok(())
292 }
293}
294
295/// One connection to a relay's room.
296pub(super) struct Socket {
297 send: Mutex<Box<dyn Write + Send>>,
298 tcp: TcpStream,
299 /// Whether this end claimed the room for its code, and so tells the relay who knew it.
300 owner: bool,
301 links: Mutex<Links>,
302}
303
304/// The streams a socket carries, by the slot of the peer at the other end.
305#[derive(Default)]
306struct Links {
307 /// This end's slot.
308 me: u32,
309 inboxes: HashMap<u32, mpsc::Sender<Vec<u8>>>,
310 /// Slots whose stream ended; a peer that comes back has a new slot.
311 ended: HashSet<u32>,
312}
313
314impl Socket {
315 fn send(&self, opcode: u8, payload: &[u8]) -> io::Result<()> {
316 let mut mask = [0; 4];
317 getrandom::fill(&mut mask).map_err(|_| io::Error::other("System random source failed"))?;
318 let frame = ws::frame(opcode, payload, Some(mask));
319 self.send.lock().unwrap().write_all(&frame)
320 }
321
322 pub(super) fn hang_up(&self) {
323 let _ = self.tcp.shutdown(Shutdown::Both);
324 }
325
326 fn forget(&self, slot: u32) {
327 let mut links = self.links.lock().unwrap();
328 links.inboxes.remove(&slot);
329 links.ended.insert(slot);
330 }
331}
332
333/// Reads notices and peers' bytes from the relay until the connection ends, opening a stream
334/// to each peer in the room.
335fn session(
336 shared: &Arc<Shared>,
337 socket: &Arc<Socket>,
338 mut reader: Reader,
339 nameplate: &mut Option<u32>,
340) {
341 let mut tag = shared.room.tag();
342 // Pings while the session lasts: dropping `_beat` at its end stops them.
343 let (_beat, beats) = mpsc::channel::<()>();
344 let pinging = Arc::clone(socket);
345 thread::spawn(move || {
346 while let Err(mpsc::RecvTimeoutError::Timeout) = beats.recv_timeout(KEEPALIVE) {
347 if pinging.send(ws::PING, &[]).is_err() {
348 return;
349 }
350 }
351 });
352 while let Ok(message) = reader.read() {
353 match message {
354 ws::Message::Text(text) => match text.parse() {
355 Ok(Notice::Nameplate(number)) => {
356 *nameplate = Some(number);
357 tag = Some(format!("code-{number}"));
358 let super::Room::Code(words) = &shared.room else {
359 continue;
360 };
361 let code = format!("{number}-{}", code_parts(words).1);
362 let mut state = shared.state.lock().unwrap();
363 if state.code.as_ref() != Some(&code) {
364 eprintln!("Live: the code is {code}");
365 state.code = Some(code);
366 drop(state);
367 (shared.notify)();
368 }
369 }
370 Ok(Notice::Welcome { you, members }) => {
371 socket.links.lock().unwrap().me = you;
372 for slot in members {
373 meet(shared, socket, tag.as_deref(), slot, None);
374 }
375 }
376 Ok(Notice::Joined(slot)) => meet(shared, socket, tag.as_deref(), slot, None),
377 Ok(Notice::Left(slot)) => socket.forget(slot),
378 Ok(Notice::Burned) => {
379 // Coming back, ask for another number.
380 *nameplate = None;
381 eprintln!("Live: too many wrong tries; the code admits no one new");
382 }
383 Err(()) => {}
384 },
385 ws::Message::Binary(data) => {
386 if let Some((slot, bytes)) = data.split_first_chunk::<SLOT>() {
387 let slot = u32::from_be_bytes(*slot);
388 meet(shared, socket, tag.as_deref(), slot, Some(bytes.to_vec()));
389 }
390 }
391 ws::Message::Ping(payload) => {
392 let _ = socket.send(ws::PONG, &payload);
393 }
394 ws::Message::Pong => {}
395 ws::Message::Close => break,
396 }
397 }
398 socket.links.lock().unwrap().inboxes.clear();
399}
400
401/// Hands `bytes` from the peer in `slot` to its stream, or opens one: this end opens a stream
402/// to a peer with a lower slot when told of it, and answers one with a higher slot when its
403/// first bytes come.
404fn meet(
405 shared: &Arc<Shared>,
406 socket: &Arc<Socket>,
407 tag: Option<&str>,
408 slot: u32,
409 bytes: Option<Vec<u8>>,
410) {
411 let mut links = socket.links.lock().unwrap();
412 if let Some(inbox) = links.inboxes.get(&slot) {
413 if let Some(bytes) = bytes {
414 let _ = inbox.send(bytes);
415 }
416 return;
417 }
418 let side = match bytes {
419 None if slot < links.me => Side::Initiator,
420 Some(_) if slot > links.me => Side::Responder,
421 _ => return,
422 };
423 let Some(tag) = tag.filter(|_| !links.ended.contains(&slot)) else {
424 return;
425 };
426 let (inbox, arriving) = mpsc::channel();
427 if let Some(bytes) = bytes {
428 let _ = inbox.send(bytes);
429 }
430 links.inboxes.insert(slot, inbox);
431 let pipe = Arc::new(Relayed {
432 socket: Arc::clone(socket),
433 slot,
434 inbox: Mutex::new(Inbox {
435 arriving,
436 chunk: Vec::new(),
437 at: 0,
438 }),
439 timeout: Mutex::new(OPENING),
440 });
441 let shared = Arc::clone(shared);
442 let tag = tag.to_owned();
443 thread::spawn(move || shared.run(pipe, side, &tag));
444}
445
446/// The stream to one peer through the relay.
447struct Relayed {
448 socket: Arc<Socket>,
449 slot: u32,
450 inbox: Mutex<Inbox>,
451 timeout: Mutex<Duration>,
452}
453
454struct Inbox {
455 arriving: mpsc::Receiver<Vec<u8>>,
456 chunk: Vec<u8>,
457 at: usize,
458}
459
460impl Pipe for Relayed {
461 fn read(&self, buffer: &mut [u8]) -> io::Result<usize> {
462 let mut inbox = self.inbox.lock().unwrap();
463 while inbox.at == inbox.chunk.len() {
464 let timeout = *self.timeout.lock().unwrap();
465 inbox.chunk = match inbox.arriving.recv_timeout(timeout) {
466 Ok(chunk) => chunk,
467 Err(mpsc::RecvTimeoutError::Timeout) => return Err(io::ErrorKind::TimedOut.into()),
468 Err(mpsc::RecvTimeoutError::Disconnected) => return Ok(0),
469 };
470 inbox.at = 0;
471 }
472 let at = inbox.at;
473 let length = buffer.len().min(inbox.chunk.len() - at);
474 buffer[..length].copy_from_slice(&inbox.chunk[at..at + length]);
475 inbox.at += length;
476 Ok(length)
477 }
478
479 fn write(&self, bytes: &[u8]) -> io::Result<usize> {
480 if self.socket.links.lock().unwrap().ended.contains(&self.slot) {
481 return Err(io::ErrorKind::BrokenPipe.into());
482 }
483 let length = bytes.len().min(CHUNK);
484 let message = [&self.slot.to_be_bytes()[..], &bytes[..length]].concat();
485 self.socket.send(ws::BINARY, &message)?;
486 Ok(length)
487 }
488
489 fn set_read_timeout(&self, timeout: Duration) -> io::Result<()> {
490 *self.timeout.lock().unwrap() = timeout;
491 Ok(())
492 }
493
494 fn shutdown(&self) {
495 self.socket.forget(self.slot);
496 }
497
498 fn direct(&self) -> bool {
499 false
500 }
501
502 fn met(&self, met: bool) {
503 if self.socket.owner {
504 let verdict = if met { Verdict::Met } else { Verdict::Failed }(self.slot);
505 let _ = self.socket.send(ws::TEXT, verdict.to_string().as_bytes());
506 }
507 }
508
509 fn broken(&self) {
510 self.socket.hang_up();
511 }
512}
513
514#[cfg(test)]
515mod tests {
516 use super::*;
517
518 #[test]
519 fn relay_addresses_parse() {
520 let address = parse("wss://live.example.net").unwrap();
521 assert_eq!(
522 address,
523 Address {
524 tls: true,
525 authority: "live.example.net".into(),
526 host: "live.example.net".into(),
527 port: 443,
528 path: String::new(),
529 }
530 );
531 let address = parse("ws://[::1]:7650/snowbound/").unwrap();
532 assert_eq!((address.host.as_str(), address.port), ("::1", 7650));
533 assert_eq!(address.path, "/snowbound");
534 assert!(parse("https://live.example.net").is_err());
535 assert!(parse("ws://:80").is_err());
536 }
537}
crates/notebook/src/live/tests.rs+368-7
......@@ -44,8 +44,8 @@ fn caret(offset: u32) -> Presence {
4444#[test]
4545fn peers_meet_and_follow_presence() {
4646 let room = Room::Code("7-violet-otter".into());
47 let ada = Live::start(hello("Ada"), &room, None, || {}).unwrap();
48 let grace = Live::start(hello("Grace"), &room, None, || {}).unwrap();
47 let ada = Live::start(hello("Ada"), &room, None, None, || {}).unwrap();
48 let grace = Live::start(hello("Grace"), &room, None, None, || {}).unwrap();
4949 ada.set_presence(caret(1));
5050 ada.connect(grace.address());
5151 let seen = until(&grace, |peers| {
......@@ -72,6 +72,7 @@ fn another_code_never_meets() {
7272 hello("Ada"),
7373 &Room::Code("7-violet-otter".into()),
7474 None,
75 None,
7576 || {},
7677 )
7778 .unwrap();
......@@ -79,6 +80,7 @@ fn another_code_never_meets() {
7980 hello("Mallory"),
8081 &Room::Code("7-violet-ocelot".into()),
8182 None,
83 None,
8284 || {},
8385 )
8486 .unwrap();
......@@ -114,12 +116,17 @@ fn later_fields_and_kinds_are_skipped() {
114116 );
115117
116118 let room = Room::Code("4-quiet-heron".into());
117 let grace = Live::start(hello("Grace"), &room, None, || {}).unwrap();
119 let grace = Live::start(hello("Grace"), &room, None, None, || {}).unwrap();
118120 // A later version: it greets, says something new, then where it is.
119121 let later = thread::spawn(move || {
120122 let mut stream = TcpStream::connect(grace.address()).unwrap();
121 let (mut send, mut receive) =
122 wire::open(&mut stream, Side::Initiator, &room.tag(), &room.secret()).unwrap();
123 let (mut send, mut receive) = wire::open(
124 &mut stream,
125 Side::Initiator,
126 &room.tag().unwrap(),
127 &room.secret(),
128 )
129 .unwrap();
123130 send.send(&mut stream, kind::HELLO, &hello("Later"))
124131 .unwrap();
125132 receive.receive(&mut stream).unwrap();
......@@ -139,8 +146,362 @@ fn later_fields_and_kinds_are_skipped() {
139146#[ignore = "multicasts mDNS on the loopback interface"]
140147fn peers_find_each_other_on_loopback() {
141148 let room = Room::Notebook([9; 16]);
142 let ada = Live::start(hello("Ada"), &room, Some(Reach::Loopback), || {}).unwrap();
143 let grace = Live::start(hello("Grace"), &room, Some(Reach::Loopback), || {}).unwrap();
149 let ada = Live::start(hello("Ada"), &room, Some(Reach::Loopback), None, || {}).unwrap();
150 let grace = Live::start(hello("Grace"), &room, Some(Reach::Loopback), None, || {}).unwrap();
144151 until(&ada, |peers| peers.len() == 1);
145152 until(&grace, |peers| peers.len() == 1);
146153}
154
155/// Two sealers that met over loopback: Ada's to send with and Grace's to receive with.
156fn sealers() -> (Sealer, Sealer) {
157 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
158 let address = listener.local_addr().unwrap();
159 let ada = thread::spawn(move || {
160 let mut stream = TcpStream::connect(address).unwrap();
161 wire::open(&mut stream, Side::Initiator, "room", b"secret")
162 .unwrap()
163 .0
164 });
165 let (mut stream, _) = listener.accept().unwrap();
166 let grace = wire::open(&mut stream, Side::Responder, "room", b"secret")
167 .unwrap()
168 .1;
169 (ada.join().unwrap(), grace)
170}
171
172/// Ada's carets at offsets `0..count`, each frame as its own block of bytes.
173fn frames(ada: &mut Sealer, count: u32) -> Vec<Vec<u8>> {
174 (0..count)
175 .map(|offset| {
176 let mut block = Vec::new();
177 ada.send(&mut block, kind::PRESENCE, &caret(offset))
178 .unwrap();
179 block
180 })
181 .collect()
182}
183
184/// A frame lost, repeated, reordered, altered, forged or sealed for another meeting is caught
185/// where it lands, before anything in it or after it is read.
186#[test]
187fn frames_out_of_place_are_caught() {
188 let (mut other, _) = sealers();
189 let elsewhere = frames(&mut other, 3).remove(2);
190 type Edit = dyn Fn(&mut Vec<Vec<u8>>);
191 let cases: [(&str, Box<Edit>, &str); 6] = [
192 (
193 "lost",
194 Box::new(|sent| drop(sent.remove(2))),
195 "Frame 3 came where frame 2 was due",
196 ),
197 (
198 "repeated",
199 Box::new(|sent| {
200 let again = sent[1].clone();
201 sent.insert(2, again);
202 }),
203 "Frame 1 came where frame 2 was due",
204 ),
205 (
206 "reordered",
207 Box::new(|sent| sent.swap(2, 3)),
208 "Frame 3 came where frame 2 was due",
209 ),
210 (
211 "altered",
212 Box::new(|sent| *sent[2].last_mut().unwrap() ^= 1),
213 "Frame 2 does not open",
214 ),
215 (
216 "forged",
217 Box::new(|sent| {
218 // Its length and number kept, the rest made up.
219 let mut forged = sent[2].clone();
220 forged[12..].fill(7);
221 sent.insert(2, forged);
222 }),
223 "Frame 2 does not open",
224 ),
225 (
226 "from another meeting",
227 Box::new(move |sent| sent.insert(2, elsewhere.clone())),
228 "Frame 2 does not open",
229 ),
230 ];
231 for (name, tamper, expected) in cases {
232 let (mut ada, mut grace) = sealers();
233 let mut sent = frames(&mut ada, 5);
234 tamper(&mut sent);
235 let stream = sent.concat();
236 let mut arriving = &stream[..];
237 for offset in 0..2 {
238 let (message, body) = grace.receive(&mut arriving).unwrap();
239 assert_eq!(message, kind::PRESENCE);
240 assert_eq!(minicbor::decode::<Presence>(&body).unwrap(), caret(offset));
241 }
242 let error = grace.receive(&mut arriving).unwrap_err();
243 assert_eq!(error.kind(), io::ErrorKind::InvalidData, "{name}");
244 assert!(error.to_string().starts_with(expected), "{name}: {error}");
245 }
246}
247
248/// A relay on this computer, with `config`'s limits: its URL.
249fn relay(config: ::relay::server::Config) -> (String, SocketAddr) {
250 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
251 let address = listener.local_addr().unwrap();
252 thread::spawn(move || ::relay::server::serve(listener, config));
253 (format!("ws://{address}"), address)
254}
255
256/// The status a relay at `address` answers a WebSocket to `path` with.
257fn status(address: SocketAddr, path: &str) -> String {
258 let mut stream = TcpStream::connect(address).unwrap();
259 write!(
260 stream,
261 "GET {path} HTTP/1.1\r\nHost: relay\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
262 Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Version: 13\r\n\r\n"
263 )
264 .unwrap();
265 let head = ::relay::ws::head(&mut stream).unwrap();
266 head.lines().next().unwrap().to_owned()
267}
268
269/// Two ends of one notebook that share no network meet in the relay's room, and see each
270/// other leave.
271#[test]
272fn peers_meet_through_a_relay() {
273 let (url, _) = relay(Default::default());
274 let room = Room::Notebook([3; 16]);
275 let ada = Live::start(hello("Ada"), &room, None, Some(&url), || {}).unwrap();
276 ada.set_presence(caret(1));
277 let grace = Live::start(hello("Grace"), &room, None, Some(&url), || {}).unwrap();
278 until(&grace, |peers| {
279 peers.len() == 1 && peers[0].presence == Some(caret(1))
280 });
281 until(&ada, |peers| {
282 peers.len() == 1 && peers[0].hello.name == "Grace"
283 });
284 ada.set_presence(caret(2));
285 until(&grace, |peers| peers[0].presence == Some(caret(2)));
286 drop(ada);
287 until(&grace, <[Peer]>::is_empty);
288}
289
290/// The end sharing a code asks the relay to number it; the other types the whole code. One
291/// with the wrong words never meets, and the relay, told so by the end sharing, burns the
292/// code once too many have tried.
293#[test]
294fn a_relay_numbers_a_code_and_burns_it_after_wrong_tries() {
295 let (url, address) = relay(::relay::server::Config {
296 burn_after: 2,
297 ..Default::default()
298 });
299 let host = Live::start(
300 hello("Ada"),
301 &Room::Code("violet-otter".into()),
302 None,
303 Some(&url),
304 || {},
305 )
306 .unwrap();
307 let deadline = Instant::now() + Duration::from_secs(10);
308 let code = loop {
309 if let Some(code) = host.code() {
310 break code;
311 }
312 assert!(Instant::now() < deadline, "no code");
313 thread::sleep(Duration::from_millis(20));
314 };
315 let (number, words) = code.split_once('-').unwrap();
316 assert_eq!(words, "violet-otter");
317 let guest = Live::start(
318 hello("Grace"),
319 &Room::Code(code.clone()),
320 None,
321 Some(&url),
322 || {},
323 )
324 .unwrap();
325 until(&host, |peers| peers.len() == 1);
326 until(&guest, |peers| peers.len() == 1);
327
328 let path = format!("/v1/room/code-{number}");
329 let wrong = Room::Code(format!("{number}-violet-ocelot"));
330 let mallory = Live::start(hello("Mallory"), &wrong, None, Some(&url), || {}).unwrap();
331 // Mallory tries again a second later, and the second wrong try burns the code. Asking
332 // sooner would count as a try itself.
333 thread::sleep(Duration::from_secs(5));
334 assert_eq!(status(address, &path), "HTTP/1.1 410 Gone");
335 assert!(mallory.peers().is_empty());
336 assert_eq!(host.peers().len(), 1, "Grace stays");
337}
338
339/// What a malicious relay does to one message on its way.
340#[derive(Clone, Copy, Debug)]
341enum Tamper {
342 Drop,
343 Repeat,
344 Reorder,
345 Alter,
346 Inject,
347}
348
349/// A relay in the middle of Grace's connection that passes on what the real one at
350/// `upstream` says, except the message to her numbered `at` among those from peers, which it
351/// tampers with. Her next connection waits for `release`.
352struct Malicious {
353 url: String,
354 tampered: mpsc::Receiver<()>,
355 rejoined: mpsc::Receiver<()>,
356 release: mpsc::Sender<()>,
357}
358
359fn malicious(upstream: SocketAddr, tamper: Tamper, at: usize) -> Malicious {
360 use ::relay::ws::{self, Message};
361 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
362 let url = format!("ws://{}", listener.local_addr().unwrap());
363 let (tampered, told) = mpsc::channel();
364 let (rejoined, heard) = mpsc::channel();
365 let (release, released) = mpsc::channel();
366 thread::spawn(move || {
367 for (index, client) in listener.incoming().enumerate() {
368 let mut client = client.unwrap();
369 if index > 0 {
370 let _ = rejoined.send(());
371 let _ = released.recv();
372 }
373 let server = TcpStream::connect(upstream).unwrap();
374 let (mut up, mut to_server) =
375 (client.try_clone().unwrap(), server.try_clone().unwrap());
376 thread::spawn(move || {
377 let _ = io::copy(&mut up, &mut to_server);
378 let _ = to_server.shutdown(Shutdown::Both);
379 });
380 let tampered = tampered.clone();
381 thread::spawn(move || {
382 let mut reading = io::BufReader::new(server);
383 let head = ws::head(&mut reading).unwrap();
384 client.write_all(head.as_bytes()).unwrap();
385 let mut messages = ws::Reader::new(reading, 1 << 20, false);
386 let (mut count, mut held) = (0, None);
387 while let Ok(message) = messages.read() {
388 let frames: Vec<Vec<u8>> = match message {
389 Message::Binary(mut data) => {
390 count += 1;
391 let mut out = vec![];
392 if index == 0 && count - 1 == at {
393 match tamper {
394 Tamper::Drop => {}
395 Tamper::Repeat => out = vec![data.clone(), data],
396 Tamper::Reorder => held = Some(data),
397 Tamper::Alter => {
398 *data.last_mut().unwrap() ^= 1;
399 out = vec![data];
400 }
401 Tamper::Inject => {
402 // A slot, length and number, then made-up bytes.
403 let mut forged = data.clone();
404 forged[16..].fill(7);
405 out = vec![forged, data];
406 }
407 }
408 let _ = tampered.send(());
409 } else {
410 out.push(data);
411 out.extend(held.take());
412 }
413 out.into_iter()
414 .map(|data| ws::frame(ws::BINARY, &data, None))
415 .collect()
416 }
417 Message::Text(text) => vec![ws::frame(ws::TEXT, text.as_bytes(), None)],
418 Message::Ping(payload) => vec![ws::frame(ws::PING, &payload, None)],
419 Message::Pong => vec![ws::frame(ws::PONG, &[], None)],
420 Message::Close => break,
421 };
422 if frames.iter().any(|frame| client.write_all(frame).is_err()) {
423 break;
424 }
425 }
426 let _ = client.shutdown(Shutdown::Both);
427 });
428 }
429 });
430 Malicious {
431 url,
432 tampered: told,
433 rejoined: heard,
434 release,
435 }
436}
437
438/// Starts `name`, recording each presence of its peer as applied, and `None` when the peer
439/// goes.
440fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc<Mutex<Vec<Option<Presence>>>>) {
441 let heard = Arc::new(Mutex::new(Vec::new()));
442 let shared: Arc<std::sync::OnceLock<std::sync::Weak<Shared>>> = Arc::default();
443 let (recorded, watched) = (Arc::clone(&heard), Arc::clone(&shared));
444 let live = Live::start(hello(name), room, None, Some(relay), move || {
445 let Some(shared) = watched.get().and_then(std::sync::Weak::upgrade) else {
446 return;
447 };
448 let entry = match shared.state.lock().unwrap().peers.values().next() {
449 None => Some(None),
450 Some(link) => link.peer.presence.clone().map(Some),
451 };
452 recorded.lock().unwrap().extend(entry);
453 })
454 .unwrap();
455 shared.set(Arc::downgrade(&live.shared)).ok().unwrap();
456 (live, heard)
457}
458
459/// Whatever a malicious relay does to a frame, the end it was for never applies it or
460/// anything after it, drops the connection as broken, joins again and hears the newest
461/// presence from scratch.
462#[test]
463fn a_malicious_relay_is_caught() {
464 let room = Room::Notebook([4; 16]);
465 for tamper in [
466 Tamper::Drop,
467 Tamper::Repeat,
468 Tamper::Reorder,
469 Tamper::Alter,
470 Tamper::Inject,
471 ] {
472 let (url, address) = relay(Default::default());
473 let ada = Live::start(hello("Ada"), &room, None, Some(&url), || {}).unwrap();
474 ada.set_presence(caret(1));
475 // Ada's opening, hello and first presence reach Grace; the next is tampered with.
476 let relay = malicious(address, tamper, 3);
477 let (grace, heard) = recording("Grace", &room, &relay.url);
478 until(&grace, |peers| {
479 peers
480 .first()
481 .is_some_and(|peer| peer.presence == Some(caret(1)))
482 });
483 ada.set_presence(caret(2));
484 relay
485 .tampered
486 .recv_timeout(Duration::from_secs(10))
487 .unwrap();
488 ada.set_presence(caret(3));
489 relay
490 .rejoined
491 .recv_timeout(Duration::from_secs(10))
492 .unwrap();
493 let before: Vec<_> = heard.lock().unwrap().clone();
494 let gone = before.iter().position(Option::is_none).unwrap();
495 assert!(
496 !before[..gone].contains(&Some(caret(3))),
497 "{tamper:?}: applied after the tampering: {before:?}"
498 );
499 assert!(grace.peers().is_empty(), "{tamper:?}");
500 relay.release.send(()).unwrap();
501 until(&grace, |peers| {
502 peers
503 .first()
504 .is_some_and(|peer| peer.presence == Some(caret(3)))
505 });
506 }
507}
crates/notebook/src/live/wire.rs+37-16
......@@ -136,7 +136,9 @@ struct Open {
136136 pake: Vec<u8>,
137137}
138138
139/// One direction's AEAD key and the count of frames sealed under it, which is each frame's nonce.
139/// One direction's AEAD key and the count of frames sealed under it. Each frame carries its
140/// number, which is also its nonce, so a frame lost, repeated, reordered or forged on the way
141/// is caught before anything in it or after it is read.
140142pub struct Sealer {
141143 cipher: Aes256Gcm,
142144 count: u64,
......@@ -152,13 +154,6 @@ impl Sealer {
152154 }
153155 }
154156
155 fn nonce(&mut self) -> [u8; 12] {
156 let mut nonce = [0; 12];
157 nonce[4..].copy_from_slice(&self.count.to_be_bytes());
158 self.count += 1;
159 nonce
160 }
161
162157 /// Writes message `kind` holding `body`.
163158 pub fn send(
164159 &mut self,
......@@ -168,28 +163,54 @@ impl Sealer {
168163 ) -> io::Result<()> {
169164 let mut clear = kind.to_be_bytes().to_vec();
170165 minicbor::encode(body, &mut clear).map_err(io::Error::other)?;
171 let nonce = self.nonce();
166 let number = self.count.to_be_bytes();
172167 let sealed = self
173168 .cipher
174 .encrypt(&nonce.into(), clear.as_slice())
169 .encrypt(&nonce(number).into(), clear.as_slice())
175170 .map_err(|_| io::Error::other("A frame could not be sealed"))?;
176 write_block(to, &sealed)
171 self.count += 1;
172 write_block(to, &[&number[..], &sealed].concat())
177173 }
178174
179 /// Reads the next message: its kind and body.
175 /// Reads the next message: its kind and body. An error of kind `InvalidData` means the
176 /// stream broke, and nothing more on it can be trusted.
180177 pub fn receive(&mut self, from: &mut impl Read) -> io::Result<(u16, Vec<u8>)> {
181 let sealed = read_block(from)?;
182 let nonce = self.nonce();
178 let block = read_block(from)?;
179 let (number, sealed) = block
180 .split_first_chunk::<8>()
181 .ok_or_else(|| invalid("A frame without its number"))?;
182 let due = self.count;
183 if u64::from_be_bytes(*number) != due {
184 return Err(io::Error::new(
185 io::ErrorKind::InvalidData,
186 format!(
187 "Frame {} came where frame {due} was due",
188 u64::from_be_bytes(*number)
189 ),
190 ));
191 }
183192 let mut clear = self
184193 .cipher
185 .decrypt(&nonce.into(), sealed.as_slice())
186 .map_err(|_| invalid("A frame does not open under the agreed key"))?;
194 .decrypt(&nonce(*number).into(), sealed)
195 .map_err(|_| {
196 io::Error::new(
197 io::ErrorKind::InvalidData,
198 format!("Frame {due} does not open under the agreed key"),
199 )
200 })?;
201 self.count += 1;
187202 let body = clear.split_off(2.min(clear.len()));
188203 let kind = u16::from_be_bytes(clear.try_into().map_err(|_| invalid("An empty frame"))?);
189204 Ok((kind, body))
190205 }
191206}
192207
208fn nonce(number: [u8; 8]) -> [u8; 12] {
209 let mut nonce = [0; 12];
210 nonce[4..].copy_from_slice(&number);
211 nonce
212}
213
193214/// Which end of the connection this is.
194215#[derive(Clone, Copy, PartialEq, Eq)]
195216pub enum Side {
crates/relay/Cargo.toml created+16
......@@ -0,0 +1,16 @@
1[package]
2name = "relay"
3version = "0.1.0"
4edition = "2024"
5publish = false
6
7# The relay Live Share meets through off the LAN (resources/live-share.md), and the WebSocket
8# framing and notices its clients share with it.
9[[bin]]
10name = "snowbound-relay"
11path = "src/main.rs"
12
13[dependencies]
14base64 = { version = "0.23.1", default-features = false, features = ["std"] }
15getrandom = "0.4.3"
16sha1 = "0.11.0"
crates/relay/README.md created+95
......@@ -0,0 +1,95 @@
1# snowbound-relay
2
3The relay Live Share meets through when two Snowbounds aren't on one network. Peers join a
4room named by a tag (a hash of the notebook's secret, or a code's number), and the relay
5passes their messages between them. Every message after the opening is sealed end to end
6and numbered inside the seal, so the relay can't read, alter, drop or reorder one unseen;
7it learns who talks to whom, when, and how much. `src/lib.rs` describes the protocol.
8
9One static Linux executable with no configuration file and nothing on disk. It speaks plain
10HTTP and WebSocket; a proxy in front of it terminates TLS.
11
12## Build
13
14```sh
15python3 tools/release_relay.py # target/relay/snowbound-relay-linux-{x86_64,aarch64}
16```
17
18It links with Rust's own lld against Rust's own musl, so it needs only `rustup`. The folder
19it makes holds both executables, `SHA256SUMS`, this README and `snowbound-relay.service`.
20
21## Deploy
22
23```sh
24scp target/relay/snowbound-relay-linux-x86_64 vps:/tmp/snowbound-relay
25scp target/relay/snowbound-relay.service vps:/tmp/
26ssh vps
27sudo install -m 755 /tmp/snowbound-relay /usr/local/bin/snowbound-relay
28sudo install -m 644 /tmp/snowbound-relay.service /etc/systemd/system/
29sudo systemctl daemon-reload
30sudo systemctl enable --now snowbound-relay
31curl -s http://127.0.0.1:7650/health # {"rooms":0,"peers":0,"connections":1,"seconds":3}
32```
33
34Then point a name at the server and put a TLS proxy in front. Caddy fetches its own
35certificate and passes WebSocket upgrades through as they are:
36
37```text
38live.example.net {
39 reverse_proxy 127.0.0.1:7650
40}
41```
42
43nginx, with a certificate from certbot:
44
45```nginx
46server {
47 listen 443 ssl;
48 server_name live.example.net;
49 ssl_certificate /etc/letsencrypt/live/live.example.net/fullchain.pem;
50 ssl_certificate_key /etc/letsencrypt/live/live.example.net/privkey.pem;
51 location / {
52 proxy_pass http://127.0.0.1:7650;
53 proxy_http_version 1.1;
54 proxy_set_header Upgrade $http_upgrade;
55 proxy_set_header Connection "upgrade";
56 # Replaced, not appended to, so a client can't name its own address.
57 proxy_set_header X-Forwarded-For $remote_addr;
58 proxy_read_timeout 1h;
59 }
60}
61```
62
63The app then uses `wss://live.example.net` (for now, `SNOWBOUND_LIVE_RELAY=wss://live.example.net`
64with a build that has the `live` feature).
65
66`--trust-forwarded true`, as the unit sets it, counts each peer by the last
67`X-Forwarded-For` entry, the one the proxy added. Without a proxy, leave it off: a client
68could otherwise claim any address. Listening on a public address without TLS works but lets
69anyone on the path see room tags and nameplates.
70
71## Options
72
73`snowbound-relay --help` lists every option and its default. Each can also be set in the
74unit's environment as `SNOWBOUND_RELAY_<OPTION>` (`SNOWBOUND_RELAY_MAX_ROOMS=64`).
75
76Memory stays under about `max-connections × (max-message + queue)` plus two small thread
77stacks per connection: some 320 MiB at the defaults, and a few MB in use by a handful of
78people. The unit caps it at 512 MiB.
79
80## Abuse
81
82- **Wrong codes.** A peer joining a code's room hears only the room's owner until the owner
83 tells the relay it met it (`met`), so the relay never needs the code. A peer the owner
84 says `failed`, one that leaves first, and one silent past `--pending` each count a wrong
85 code against its address (an IPv4 address or an IPv6 /64) and against the code. Ten in a
86 minute lock the address out for a minute, then two, four, up to an hour, forgotten after
87 a quiet hour; tries still waiting count toward the ten, so they can't run side by side.
88 Five burn the code: it admits no one new, and the owner makes another.
89- **Load.** Joins per address and per room are rate limited; connections in all, per
90 address, rooms, and peers per room are capped; a room's bytes per second are throttled by
91 slowing its senders; a message over `--max-message` or a peer too slow to take what waits
92 for it (`--queue`) is hung up on; a silent connection closes after `--idle`.
93
94The relay logs lockouts and burned codes to stderr (`journalctl -u snowbound-relay`), and
95nothing else.
crates/relay/deploy/snowbound-relay.service created+33
......@@ -0,0 +1,33 @@
1# /etc/systemd/system/snowbound-relay.service: the Live Share relay behind a TLS proxy on
2# this machine (README.md). Options go on ExecStart, or as SNOWBOUND_RELAY_* in Environment.
3[Unit]
4Description=Snowbound Live Share relay
5After=network-online.target
6Wants=network-online.target
7
8[Service]
9ExecStart=/usr/local/bin/snowbound-relay --listen 127.0.0.1:7650 --trust-forwarded true
10Restart=always
11RestartSec=2
12DynamicUser=yes
13# Two threads a connection, and a few more.
14TasksMax=600
15LimitNOFILE=1024
16MemoryMax=512M
17NoNewPrivileges=yes
18ProtectSystem=strict
19ProtectHome=yes
20PrivateTmp=yes
21PrivateDevices=yes
22ProtectKernelTunables=yes
23ProtectKernelModules=yes
24ProtectControlGroups=yes
25RestrictAddressFamilies=AF_INET AF_INET6
26RestrictNamespaces=yes
27LockPersonality=yes
28MemoryDenyWriteExecute=yes
29SystemCallArchitectures=native
30CapabilityBoundingSet=
31
32[Install]
33WantedBy=multi-user.target
crates/relay/src/lib.rs created+137
......@@ -0,0 +1,137 @@
1//! Snowbound's Live Share relay, and what its clients share with it.
2//!
3//! A peer opens a WebSocket to `/v1/room/<tag>`, or to `/v1/claim` for a code's room, whose
4//! number the relay chooses. Every peer in a room has a slot, a number the relay gives it
5//! that is never reused in that room. A binary message is a slot (`u32`, big-endian) and
6//! bytes: sent, the slot it goes to; received, the slot it came from. What the bytes are
7//! the relay never knows; peers seal them end to end. Text messages are [`Notice`]s from the
8//! relay and [`Verdict`]s to it.
9//!
10//! Of each pair of peers the one with the higher slot opens the stream between them.
11//! A peer joining a code's room waits, hearing only the room's owner, until the owner says
12//! its opening `met`; a peer that `failed`, left first, or stayed silent too long counts a
13//! wrong code against its address and the code. See `resources/live-share.md`.
14
15pub mod server;
16pub mod ws;
17
18use std::{fmt, str::FromStr};
19
20/// The bytes before a binary message's payload: the slot it goes to or came from.
21pub const SLOT: usize = 4;
22
23/// What the relay tells a peer.
24#[derive(Clone, Debug, PartialEq, Eq)]
25pub enum Notice {
26 /// The number of the code whose room this peer claimed; sent before `Welcome`.
27 Nameplate(u32),
28 /// This peer's slot, and the slots of those already in the room it may talk to.
29 Welcome {
30 you: u32,
31 members: Vec<u32>,
32 },
33 /// Someone this peer may now talk to.
34 Joined(u32),
35 Left(u32),
36 /// The code had too many wrong tries and admits no one new; sent to its owner.
37 Burned,
38}
39
40/// What a code's owner tells the relay of the peer in a slot.
41#[derive(Clone, Copy, Debug, PartialEq, Eq)]
42pub enum Verdict {
43 /// The peer knew the code.
44 Met(u32),
45 /// The peer's opening failed: a wrong code.
46 Failed(u32),
47}
48
49impl fmt::Display for Notice {
50 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
51 match self {
52 Notice::Nameplate(number) => write!(f, "nameplate {number}"),
53 Notice::Welcome { you, members } => {
54 write!(f, "welcome {you}")?;
55 members.iter().try_for_each(|slot| write!(f, " {slot}"))
56 }
57 Notice::Joined(slot) => write!(f, "joined {slot}"),
58 Notice::Left(slot) => write!(f, "left {slot}"),
59 Notice::Burned => f.write_str("burned"),
60 }
61 }
62}
63
64impl FromStr for Notice {
65 type Err = ();
66
67 fn from_str(text: &str) -> Result<Self, ()> {
68 let mut words = text.split(' ');
69 let name = words.next().ok_or(())?;
70 let numbers: Vec<u32> = words
71 .map(|word| word.parse().map_err(|_| ()))
72 .collect::<Result<_, _>>()?;
73 Ok(match (name, &numbers[..]) {
74 ("nameplate", [number]) => Notice::Nameplate(*number),
75 ("welcome", [you, members @ ..]) => Notice::Welcome {
76 you: *you,
77 members: members.to_vec(),
78 },
79 ("joined", [slot]) => Notice::Joined(*slot),
80 ("left", [slot]) => Notice::Left(*slot),
81 ("burned", []) => Notice::Burned,
82 _ => return Err(()),
83 })
84 }
85}
86
87impl fmt::Display for Verdict {
88 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
89 match self {
90 Verdict::Met(slot) => write!(f, "met {slot}"),
91 Verdict::Failed(slot) => write!(f, "failed {slot}"),
92 }
93 }
94}
95
96impl FromStr for Verdict {
97 type Err = ();
98
99 fn from_str(text: &str) -> Result<Self, ()> {
100 let (name, slot) = text.split_once(' ').ok_or(())?;
101 let slot = slot.parse().map_err(|_| ())?;
102 match name {
103 "met" => Ok(Verdict::Met(slot)),
104 "failed" => Ok(Verdict::Failed(slot)),
105 _ => Err(()),
106 }
107 }
108}
109
110#[cfg(test)]
111mod tests {
112 use super::*;
113
114 #[test]
115 fn notices_and_verdicts_read_back() {
116 for notice in [
117 Notice::Nameplate(412),
118 Notice::Welcome {
119 you: 3,
120 members: vec![1, 2],
121 },
122 Notice::Welcome {
123 you: 1,
124 members: vec![],
125 },
126 Notice::Joined(4),
127 Notice::Left(4),
128 Notice::Burned,
129 ] {
130 assert_eq!(notice.to_string().parse(), Ok(notice));
131 }
132 for verdict in [Verdict::Met(2), Verdict::Failed(9)] {
133 assert_eq!(verdict.to_string().parse(), Ok(verdict));
134 }
135 assert_eq!("shrug 1".parse::<Notice>(), Err(()));
136 }
137}
crates/relay/src/main.rs created+120
......@@ -0,0 +1,120 @@
1//! `snowbound-relay`: see `crates/relay/README.md`.
2
3use relay::server::{self, Config};
4use std::{net::TcpListener, process::ExitCode, time::Duration};
5
6const USAGE: &str = "\
7usage: snowbound-relay [OPTION VALUE]...
8
9Each option may also come from the environment as SNOWBOUND_RELAY_<OPTION>, upper case with
10underscores: SNOWBOUND_RELAY_LISTEN=127.0.0.1:7650. Flags win.
11
12 --listen ADDRESS where to listen (127.0.0.1:7650)
13 --trust-forwarded true|false count peers by X-Forwarded-For, behind a proxy (false)
14 --max-connections N (256)
15 --max-connections-per-address N per IPv4 address or IPv6 /64 (16)
16 --max-rooms N (128)
17 --max-room-peers N (16)
18 --max-message BYTES (262144)
19 --queue BYTES waiting to go to one peer before it is dropped (1048576)
20 --idle SECONDS silence before a connection is closed (600)
21 --room-bytes-per-second BYTES (4194304)
22 --joins-per-minute N per address (30)
23 --room-joins-per-minute N (30)
24 --failures-per-minute N wrong codes per address before a lockout (10)
25 --burn-after N wrong codes before a code admits no one new (5)
26 --pending SECONDS for a peer joining a code to meet its owner (20)
27";
28
29fn main() -> ExitCode {
30 let mut listen = String::from("127.0.0.1:7650");
31 let mut config = Config::default();
32 let from_environment = std::env::vars().filter_map(|(key, value)| {
33 let option = key.strip_prefix("SNOWBOUND_RELAY_")?;
34 Some((option.to_lowercase().replace('_', "-"), value))
35 });
36 let mut arguments = std::env::args().skip(1);
37 let mut from_flags = Vec::new();
38 while let Some(flag) = arguments.next() {
39 if flag == "--help" {
40 print!("{USAGE}");
41 return ExitCode::SUCCESS;
42 }
43 let Some(option) = flag.strip_prefix("--") else {
44 eprint!("{USAGE}");
45 return ExitCode::from(2);
46 };
47 let (option, value) = match option.split_once('=') {
48 Some((option, value)) => (option.to_owned(), value.to_owned()),
49 None => match arguments.next() {
50 Some(value) => (option.to_owned(), value),
51 None => {
52 eprintln!("--{option} needs a value\n\n{USAGE}");
53 return ExitCode::from(2);
54 }
55 },
56 };
57 from_flags.push((option, value));
58 }
59 for (option, value) in from_environment.chain(from_flags) {
60 let set = match option.as_str() {
61 "listen" => {
62 listen = value.clone();
63 Ok(())
64 }
65 _ => set(&mut config, &option, &value),
66 };
67 if let Err(error) = set {
68 eprintln!("--{option} {value}: {error}\n\n{USAGE}");
69 return ExitCode::from(2);
70 }
71 }
72 let listener = match TcpListener::bind(&listen) {
73 Ok(listener) => listener,
74 Err(error) => {
75 eprintln!("Cannot listen on {listen}: {error}");
76 return ExitCode::FAILURE;
77 }
78 };
79 match listener.local_addr() {
80 Ok(address) => println!(
81 "snowbound-relay {} listening on {address}",
82 env!("CARGO_PKG_VERSION")
83 ),
84 Err(error) => eprintln!("{error}"),
85 }
86 match server::serve(listener, config) {
87 Ok(()) => ExitCode::SUCCESS,
88 Err(error) => {
89 eprintln!("{error}");
90 ExitCode::FAILURE
91 }
92 }
93}
94
95fn set(config: &mut Config, option: &str, value: &str) -> Result<(), String> {
96 let number = || value.parse::<u64>().map_err(|error| error.to_string());
97 let count = || number().map(|number| number as usize);
98 let small =
99 || number().and_then(|number| u32::try_from(number).map_err(|error| error.to_string()));
100 match option {
101 "trust-forwarded" => {
102 config.trust_forwarded = value.parse().map_err(|_| "true or false".to_owned())?
103 }
104 "max-connections" => config.max_connections = count()?,
105 "max-connections-per-address" => config.max_connections_per_address = count()?,
106 "max-rooms" => config.max_rooms = count()?,
107 "max-room-peers" => config.max_room_peers = count()?,
108 "max-message" => config.max_message = count()?,
109 "queue" => config.queue = count()?,
110 "idle" => config.idle = Duration::from_secs(number()?),
111 "room-bytes-per-second" => config.room_bytes_per_second = number()?.max(1),
112 "joins-per-minute" => config.joins_per_minute = small()?.max(1),
113 "room-joins-per-minute" => config.room_joins_per_minute = small()?.max(1),
114 "failures-per-minute" => config.failures_per_minute = small()?.max(1),
115 "burn-after" => config.burn_after = small()?.max(1),
116 "pending" => config.pending = Duration::from_secs(number()?),
117 _ => return Err("no such option".into()),
118 }
119 Ok(())
120}
crates/relay/src/server.rs created+881
......@@ -0,0 +1,881 @@
1//! The relay: rooms of WebSocket peers, each message passed on to the peer it names, with
2//! limits on everything a stranger can make it hold. A thread reads each connection and
3//! another writes it, from a queue capped in bytes.
4
5use crate::{Notice, SLOT, Verdict, ws};
6use std::{
7 collections::{BTreeMap, HashMap, VecDeque},
8 io::{BufReader, Write},
9 net::{IpAddr, Ipv4Addr, Ipv6Addr, Shutdown, TcpListener, TcpStream},
10 ops::RangeInclusive,
11 sync::{
12 Arc, Condvar, Mutex,
13 atomic::{AtomicUsize, Ordering},
14 },
15 thread,
16 time::{Duration, Instant},
17};
18
19/// What a relay allows; `snowbound-relay --help` explains each.
20#[derive(Clone, Debug)]
21pub struct Config {
22 /// Counts a peer by the last `X-Forwarded-For` entry, the one its own proxy added.
23 pub trust_forwarded: bool,
24 pub max_connections: usize,
25 pub max_connections_per_address: usize,
26 pub max_rooms: usize,
27 pub max_room_peers: usize,
28 /// The largest message, in bytes.
29 pub max_message: usize,
30 /// The most bytes waiting to go to one peer before the relay hangs up on it.
31 pub queue: usize,
32 /// How long a connection may send nothing.
33 pub idle: Duration,
34 pub room_bytes_per_second: u64,
35 pub joins_per_minute: u32,
36 pub room_joins_per_minute: u32,
37 /// Wrong codes an address may try in a minute before it is locked out, a minute the first
38 /// time and twice as long each time after, up to an hour.
39 pub failures_per_minute: u32,
40 /// Wrong codes after which a code admits no one new.
41 pub burn_after: u32,
42 /// How long a peer joining a code's room has to meet its owner.
43 pub pending: Duration,
44}
45
46impl Default for Config {
47 fn default() -> Self {
48 Self {
49 trust_forwarded: false,
50 max_connections: 256,
51 max_connections_per_address: 16,
52 max_rooms: 128,
53 max_room_peers: 16,
54 max_message: 256 << 10,
55 queue: 1 << 20,
56 idle: Duration::from_secs(600),
57 room_bytes_per_second: 4 << 20,
58 joins_per_minute: 30,
59 room_joins_per_minute: 30,
60 failures_per_minute: 10,
61 burn_after: 5,
62 pending: Duration::from_secs(20),
63 }
64 }
65}
66
67const HANDSHAKE: Duration = Duration::from_secs(10);
68const WRITE: Duration = Duration::from_secs(30);
69const MINUTE: Duration = Duration::from_secs(60);
70const HOUR: Duration = Duration::from_secs(3600);
71/// Addresses remembered at once; past it, new ones wait.
72const ADDRESSES: usize = 1 << 16;
73const NAMEPLATES: RangeInclusive<u32> = 1..=999;
74const STACK: usize = 256 << 10;
75
76/// Serves `listener` until it fails.
77pub fn serve(listener: TcpListener, config: Config) -> std::io::Result<()> {
78 let relay = Arc::new(Relay {
79 config,
80 state: Mutex::default(),
81 connections: AtomicUsize::new(0),
82 started: Instant::now(),
83 });
84 let sweeping = Arc::clone(&relay);
85 thread::Builder::new().name("sweep".into()).spawn(move || {
86 loop {
87 thread::sleep(Duration::from_secs(1));
88 sweeping.sweep(Instant::now());
89 }
90 })?;
91 for stream in listener.incoming() {
92 let Ok(stream) = stream else {
93 // Out of descriptors, most likely: let some close.
94 thread::sleep(Duration::from_millis(50));
95 continue;
96 };
97 if relay.connections.fetch_add(1, Ordering::AcqRel) >= relay.config.max_connections {
98 relay.connections.fetch_sub(1, Ordering::AcqRel);
99 continue;
100 }
101 let serving = Arc::clone(&relay);
102 let spawned = thread::Builder::new().stack_size(STACK).spawn(move || {
103 serving.connection(stream);
104 serving.connections.fetch_sub(1, Ordering::AcqRel);
105 });
106 if spawned.is_err() {
107 relay.connections.fetch_sub(1, Ordering::AcqRel);
108 }
109 }
110 Ok(())
111}
112
113struct Relay {
114 config: Config,
115 state: Mutex<State>,
116 /// Connections open, joined or not.
117 connections: AtomicUsize,
118 started: Instant,
119}
120
121#[derive(Default)]
122struct State {
123 rooms: HashMap<String, Room>,
124 addresses: HashMap<IpAddr, Address>,
125}
126
127struct Room {
128 /// The slot of the peer that claimed a code's room, while it is there.
129 owner: Option<u32>,
130 next: u32,
131 members: BTreeMap<u32, Member>,
132 joins: Bucket,
133 bytes: Bucket,
134 failures: u32,
135 burned: bool,
136}
137
138struct Member {
139 outbox: Arc<Outbox>,
140 address: IpAddr,
141 /// When a peer joining a code's room joined, until its owner says it met it.
142 pending: Option<Instant>,
143}
144
145/// What one address (or IPv6 /64) is doing.
146struct Address {
147 connections: usize,
148 pending: usize,
149 joins: Bucket,
150 /// Wrong codes in the last minute.
151 failures: VecDeque<Instant>,
152 /// Lockouts in a row, each twice as long.
153 strikes: u32,
154 struck: Option<Instant>,
155 locked: Option<Instant>,
156}
157
158enum Ask {
159 /// A code's room, numbered by the relay or, coming back, as it was.
160 Claim(Option<u32>),
161 Room(String),
162}
163
164enum Refusal {
165 NotFound,
166 Gone,
167 Wait(Duration),
168 Full,
169}
170
171impl Relay {
172 fn connection(&self, stream: TcpStream) {
173 let _ = stream.set_nodelay(true);
174 let _ = stream.set_read_timeout(Some(HANDSHAKE));
175 let _ = stream.set_write_timeout(Some(WRITE));
176 let Ok(reading) = stream.try_clone() else {
177 return;
178 };
179 let mut reader = BufReader::new(reading);
180 let Ok(head) = ws::head(&mut reader) else {
181 return;
182 };
183 let address = self.address(&stream, &head);
184 let target = head.split(' ').nth(1).unwrap_or_default();
185 let (path, query) = target.split_once('?').unwrap_or((target, ""));
186 let ask = match path {
187 "/health" => {
188 respond(&stream, "200 OK", "application/json", "", &self.health());
189 return;
190 }
191 "/v1/claim" => Ask::Claim(
192 query
193 .split('&')
194 .find_map(|pair| pair.strip_prefix("nameplate="))
195 .and_then(|number| number.parse().ok()),
196 ),
197 _ => match path.strip_prefix("/v1/room/").filter(|tag| valid(tag)) {
198 Some(tag) => Ask::Room(tag.into()),
199 None => {
200 respond(&stream, "404 Not Found", "text/plain", "", "No such page\n");
201 return;
202 }
203 },
204 };
205 let upgrade = ws::header(&head, "Upgrade")
206 .is_some_and(|value| value.eq_ignore_ascii_case("websocket"));
207 let (true, Some(key)) = (
208 upgrade && head.starts_with("GET "),
209 ws::header(&head, "Sec-WebSocket-Key"),
210 ) else {
211 respond(
212 &stream,
213 "400 Bad Request",
214 "text/plain",
215 "",
216 "A WebSocket only\n",
217 );
218 return;
219 };
220 let Ok(writing) = stream.try_clone() else {
221 return;
222 };
223 let outbox = Arc::new(Outbox::new(writing, self.config.queue));
224 let (tag, slot) = match self.join(address, ask, &outbox, Instant::now()) {
225 Ok(joined) => joined,
226 Err(refusal) => {
227 let (status, headers, body) = match refusal {
228 Refusal::NotFound => ("404 Not Found", String::new(), "No such code\n"),
229 Refusal::Gone => ("410 Gone", String::new(), "The code has expired\n"),
230 Refusal::Wait(wait) => (
231 "429 Too Many Requests",
232 // Rounded up, so that a client waiting so long finds it over.
233 format!(
234 "Retry-After: {}\r\n",
235 wait.as_secs() + u64::from(wait.subsec_nanos() > 0)
236 ),
237 "Too many tries\n",
238 ),
239 Refusal::Full => (
240 "503 Service Unavailable",
241 "Retry-After: 30\r\n".into(),
242 "The relay is full\n",
243 ),
244 };
245 respond(&stream, status, "text/plain", &headers, body);
246 return;
247 }
248 };
249 let switching = format!(
250 "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
251 Sec-WebSocket-Accept: {}\r\n\r\n",
252 ws::accept(key)
253 );
254 let writer = Arc::clone(&outbox);
255 if (&stream).write_all(switching.as_bytes()).is_ok()
256 && thread::Builder::new()
257 .stack_size(STACK)
258 .spawn(move || writer.drain())
259 .is_ok()
260 {
261 let _ = stream.set_read_timeout(Some(self.config.idle));
262 let mut reader = ws::Reader::new(reader, self.config.max_message, true);
263 while let Ok(message) = reader.read() {
264 if !self.heard(&tag, slot, &outbox, message) {
265 break;
266 }
267 }
268 }
269 let mut state = self.state.lock().unwrap();
270 depart(&mut state, &self.config, &tag, slot, Instant::now());
271 }
272
273 /// Where a peer counts: the address it connected from, or its proxy says it did.
274 fn address(&self, stream: &TcpStream, head: &str) -> IpAddr {
275 let forwarded = ws::header(head, "X-Forwarded-For")
276 .filter(|_| self.config.trust_forwarded)
277 .and_then(|value| value.rsplit(',').next()?.trim().parse().ok());
278 let ip = forwarded
279 .or_else(|| stream.peer_addr().ok().map(|address| address.ip()))
280 .unwrap_or(IpAddr::V4(Ipv4Addr::UNSPECIFIED));
281 match ip {
282 IpAddr::V6(v6) => match v6.to_ipv4_mapped() {
283 Some(v4) => IpAddr::V4(v4),
284 // One subscriber is usually given a whole /64.
285 None => {
286 let [a, b, c, d, ..] = v6.segments();
287 IpAddr::V6(Ipv6Addr::new(a, b, c, d, 0, 0, 0, 0))
288 }
289 },
290 v4 => v4,
291 }
292 }
293
294 fn health(&self) -> String {
295 let state = self.state.lock().unwrap();
296 let peers: usize = state.rooms.values().map(|room| room.members.len()).sum();
297 format!(
298 "{{\"rooms\":{},\"peers\":{peers},\"connections\":{},\"seconds\":{}}}\n",
299 state.rooms.len(),
300 self.connections.load(Ordering::Acquire),
301 self.started.elapsed().as_secs()
302 )
303 }
304
305 /// Puts a peer from `address` in the room it asks for, telling it and those it may talk
306 /// to: its room's tag and its slot.
307 fn join(
308 &self,
309 address: IpAddr,
310 ask: Ask,
311 outbox: &Arc<Outbox>,
312 now: Instant,
313 ) -> Result<(String, u32), Refusal> {
314 let config = &self.config;
315 let mut state = self.state.lock().unwrap();
316 let State { rooms, addresses } = &mut *state;
317 if !addresses.contains_key(&address) && addresses.len() >= ADDRESSES {
318 return Err(Refusal::Full);
319 }
320 let client = addresses
321 .entry(address)
322 .or_insert_with(|| Address::new(config, now));
323 client.refresh(now);
324 if let Some(until) = client.locked {
325 return Err(Refusal::Wait(until - now));
326 }
327 if client.connections >= config.max_connections_per_address {
328 return Err(Refusal::Wait(Duration::from_secs(10)));
329 }
330 let per_minute = f64::from(config.joins_per_minute);
331 client
332 .joins
333 .take(per_minute, per_minute / 60.0, now)
334 .map_err(Refusal::Wait)?;
335 let (tag, nameplate) = match ask {
336 Ask::Claim(back) => {
337 let free = |number: u32| {
338 rooms
339 .get(&code(number))
340 .is_none_or(|room| room.owner.is_none() && !room.burned)
341 };
342 let number = back
343 .filter(|number| NAMEPLATES.contains(number) && free(*number))
344 .or_else(|| {
345 (0..32)
346 .filter_map(|_| {
347 let mut bytes = [0; 4];
348 getrandom::fill(&mut bytes).ok()?;
349 Some(1 + u32::from_le_bytes(bytes) % NAMEPLATES.end())
350 })
351 .chain(NAMEPLATES)
352 .find(|number| free(*number))
353 })
354 .ok_or(Refusal::Full)?;
355 (code(number), Some(number))
356 }
357 Ask::Room(tag) if tag.starts_with("code-") => {
358 let room = rooms
359 .get(&tag)
360 .filter(|room| room.owner.is_some())
361 .ok_or(Refusal::NotFound)?;
362 if room.burned {
363 return Err(Refusal::Gone);
364 }
365 if client.failures.len() + client.pending >= config.failures_per_minute as usize {
366 let wait = client.failures.front().map_or(config.pending, |first| {
367 (*first + MINUTE).saturating_duration_since(now)
368 });
369 return Err(Refusal::Wait(wait));
370 }
371 (tag, None)
372 }
373 Ask::Room(tag) => (tag, None),
374 };
375 if !rooms.contains_key(&tag) && rooms.len() >= config.max_rooms {
376 return Err(Refusal::Full);
377 }
378 let room = rooms.entry(tag.clone()).or_insert_with(|| Room {
379 owner: None,
380 next: 1,
381 members: BTreeMap::new(),
382 joins: Bucket::full(f64::from(config.room_joins_per_minute), now),
383 bytes: Bucket::full(config.room_bytes_per_second as f64, now),
384 failures: 0,
385 burned: false,
386 });
387 let per_minute = f64::from(config.room_joins_per_minute);
388 let admitted = room.members.len() < config.max_room_peers;
389 let joined = admitted.then(|| room.joins.take(per_minute, per_minute / 60.0, now));
390 let refusal = match joined {
391 None => Some(Refusal::Full),
392 Some(Err(wait)) => Some(Refusal::Wait(wait)),
393 Some(Ok(())) => None,
394 };
395 if let Some(refusal) = refusal {
396 if room.members.is_empty() {
397 rooms.remove(&tag);
398 }
399 return Err(refusal);
400 }
401 let slot = room.next;
402 room.next += 1;
403 let pending = room.owner.is_some() && nameplate.is_none();
404 let visible: Vec<u32> = match room.owner {
405 Some(owner) if pending => vec![owner],
406 _ => (room.members.iter())
407 .filter(|(_, member)| member.pending.is_none())
408 .map(|(slot, _)| *slot)
409 .collect(),
410 };
411 if let Some(number) = nameplate {
412 room.owner = Some(slot);
413 outbox.push(notice(Notice::Nameplate(number)));
414 }
415 outbox.push(notice(Notice::Welcome {
416 you: slot,
417 members: visible.clone(),
418 }));
419 for other in visible {
420 room.members[&other]
421 .outbox
422 .push(notice(Notice::Joined(slot)));
423 }
424 room.members.insert(
425 slot,
426 Member {
427 outbox: Arc::clone(outbox),
428 address,
429 pending: pending.then_some(now),
430 },
431 );
432 client.connections += 1;
433 client.pending += usize::from(pending);
434 Ok((tag, slot))
435 }
436
437 /// Acts on a message from `slot`: false once it is gone or broke the protocol.
438 fn heard(&self, tag: &str, slot: u32, outbox: &Outbox, message: ws::Message) -> bool {
439 match message {
440 ws::Message::Binary(mut data) => {
441 if data.len() < SLOT {
442 return false;
443 }
444 let rate = self.config.room_bytes_per_second as f64;
445 let wait = match self.state.lock().unwrap().rooms.get_mut(tag) {
446 Some(room) => room.bytes.spend(data.len() as f64, rate, Instant::now()),
447 None => return false,
448 };
449 // Waiting here slows the sender alone, as its socket fills.
450 thread::sleep(wait);
451 let to = u32::from_be_bytes(data[..SLOT].try_into().expect("a slot"));
452 let state = self.state.lock().unwrap();
453 let Some(room) = state.rooms.get(tag) else {
454 return false;
455 };
456 let (Some(from), Some(target)) = (room.members.get(&slot), room.members.get(&to))
457 else {
458 return room.members.contains_key(&slot);
459 };
460 // One waiting for the code's owner talks to the owner alone.
461 if (from.pending.is_none() || room.owner == Some(to))
462 && (target.pending.is_none() || room.owner == Some(slot))
463 {
464 data[..SLOT].copy_from_slice(&slot.to_be_bytes());
465 target.outbox.push(ws::frame(ws::BINARY, &data, None));
466 }
467 true
468 }
469 ws::Message::Text(text) => {
470 if let Ok(verdict) = text.parse() {
471 self.judge(tag, slot, verdict);
472 }
473 true
474 }
475 ws::Message::Ping(payload) => {
476 outbox.push(ws::frame(ws::PONG, &payload, None));
477 true
478 }
479 ws::Message::Pong => true,
480 ws::Message::Close => false,
481 }
482 }
483
484 /// Takes a code's owner's word on the peer in a slot waiting to meet it.
485 fn judge(&self, tag: &str, from: u32, verdict: Verdict) {
486 let mut state = self.state.lock().unwrap();
487 let State { rooms, addresses } = &mut *state;
488 let Some(room) = rooms.get_mut(tag).filter(|room| room.owner == Some(from)) else {
489 return;
490 };
491 match verdict {
492 Verdict::Met(slot) => {
493 let Some(member) = room.members.get_mut(&slot) else {
494 return;
495 };
496 if member.pending.take().is_none() {
497 return;
498 }
499 if let Some(address) = addresses.get_mut(&member.address) {
500 address.pending -= 1;
501 }
502 let newcomer = Arc::clone(&member.outbox);
503 for (other, peer) in &room.members {
504 if *other != slot && *other != from && peer.pending.is_none() {
505 newcomer.push(notice(Notice::Joined(*other)));
506 peer.outbox.push(notice(Notice::Joined(slot)));
507 }
508 }
509 }
510 Verdict::Failed(slot) => {
511 if room
512 .members
513 .get(&slot)
514 .is_some_and(|member| member.pending.is_some())
515 {
516 depart(&mut state, &self.config, tag, slot, Instant::now());
517 }
518 }
519 }
520 }
521
522 /// Fails the peers that waited too long to meet a code's owner, and forgets quiet addresses.
523 fn sweep(&self, now: Instant) {
524 let config = &self.config;
525 let mut state = self.state.lock().unwrap();
526 let late: Vec<(String, u32)> = (state.rooms.iter())
527 .flat_map(|(tag, room)| {
528 (room.members.iter())
529 .filter(|(_, member)| {
530 member.pending.is_some_and(|since| {
531 now.saturating_duration_since(since) >= config.pending
532 })
533 })
534 .map(|(slot, _)| (tag.clone(), *slot))
535 })
536 .collect();
537 for (tag, slot) in late {
538 depart(&mut state, config, &tag, slot, now);
539 }
540 state.addresses.retain(|_, address| {
541 address.refresh(now);
542 !address.quiet(config, now)
543 });
544 }
545}
546
547/// Takes `slot` out of room `tag` and hangs up on it. One still waiting to meet a code's
548/// owner tried a wrong code; the owner leaving excuses those waiting for it.
549fn depart(state: &mut State, config: &Config, tag: &str, slot: u32, now: Instant) {
550 let State { rooms, addresses } = state;
551 let Some(room) = rooms.get_mut(tag) else {
552 return;
553 };
554 let Some(member) = room.members.remove(&slot) else {
555 return;
556 };
557 member.outbox.close();
558 if let Some(address) = addresses.get_mut(&member.address) {
559 address.connections -= 1;
560 if member.pending.is_some() {
561 address.pending -= 1;
562 if let Some(lockout) = address.fail(now, config.failures_per_minute) {
563 eprintln!(
564 "{}: locked out for {} s after too many wrong codes",
565 member.address,
566 lockout.as_secs()
567 );
568 }
569 }
570 }
571 for peer in room.members.values() {
572 peer.outbox.push(notice(Notice::Left(slot)));
573 }
574 if member.pending.is_some() {
575 room.failures += 1;
576 if room.failures >= config.burn_after && !room.burned {
577 room.burned = true;
578 eprintln!("{tag}: burned after {} wrong codes", room.failures);
579 if let Some(owner) = room.owner.and_then(|owner| room.members.get(&owner)) {
580 owner.outbox.push(notice(Notice::Burned));
581 }
582 }
583 }
584 if room.owner == Some(slot) {
585 room.owner = None;
586 for member in room.members.values_mut() {
587 if member.pending.take().is_some() {
588 if let Some(address) = addresses.get_mut(&member.address) {
589 address.pending -= 1;
590 }
591 member.outbox.close();
592 }
593 }
594 }
595 if room.members.is_empty() {
596 rooms.remove(tag);
597 }
598}
599
600fn notice(notice: Notice) -> Vec<u8> {
601 ws::frame(ws::TEXT, notice.to_string().as_bytes(), None)
602}
603
604fn code(number: u32) -> String {
605 format!("code-{number}")
606}
607
608/// A room tag as clients make them: a hash in hex, or `code-` and a number.
609fn valid(tag: &str) -> bool {
610 (1..=64).contains(&tag.len())
611 && tag
612 .bytes()
613 .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-')
614}
615
616fn respond(mut stream: &TcpStream, status: &str, kind: &str, headers: &str, body: &str) {
617 let _ = write!(
618 stream,
619 "HTTP/1.1 {status}\r\nContent-Type: {kind}\r\nContent-Length: {}\r\nConnection: close\r\n\
620 {headers}\r\n{body}",
621 body.len()
622 );
623}
624
625impl Address {
626 fn new(config: &Config, now: Instant) -> Self {
627 Self {
628 connections: 0,
629 pending: 0,
630 joins: Bucket::full(f64::from(config.joins_per_minute), now),
631 failures: VecDeque::new(),
632 strikes: 0,
633 struck: None,
634 locked: None,
635 }
636 }
637
638 /// Forgets failures over a minute old, a lockout that has ended, and strikes after a
639 /// quiet hour.
640 fn refresh(&mut self, now: Instant) {
641 while self
642 .failures
643 .front()
644 .is_some_and(|failed| now.saturating_duration_since(*failed) >= MINUTE)
645 {
646 self.failures.pop_front();
647 }
648 self.locked = self.locked.filter(|until| *until > now);
649 if self
650 .struck
651 .is_some_and(|struck| now.saturating_duration_since(struck) >= HOUR)
652 {
653 self.strikes = 0;
654 self.struck = None;
655 }
656 }
657
658 /// Counts a wrong code: at `limit` in a minute, the lockout it starts.
659 fn fail(&mut self, now: Instant, limit: u32) -> Option<Duration> {
660 self.refresh(now);
661 self.failures.push_back(now);
662 if self.failures.len() < limit as usize {
663 return None;
664 }
665 self.failures.clear();
666 self.strikes += 1;
667 self.struck = Some(now);
668 let lockout = MINUTE
669 .saturating_mul(1 << (self.strikes - 1).min(6))
670 .min(HOUR);
671 self.locked = Some(now + lockout);
672 Some(lockout)
673 }
674
675 /// Whether nothing about this address needs remembering.
676 fn quiet(&self, config: &Config, now: Instant) -> bool {
677 let mut joins = self.joins;
678 let per_minute = f64::from(config.joins_per_minute);
679 self.connections == 0
680 && self.pending == 0
681 && self.failures.is_empty()
682 && self.locked.is_none()
683 && self.strikes == 0
684 && joins.take(per_minute, per_minute / 60.0, now).is_ok()
685 }
686}
687
688/// A token bucket.
689#[derive(Clone, Copy)]
690struct Bucket {
691 tokens: f64,
692 at: Instant,
693}
694
695impl Bucket {
696 fn full(capacity: f64, now: Instant) -> Self {
697 Self {
698 tokens: capacity,
699 at: now,
700 }
701 }
702
703 fn refill(&mut self, capacity: f64, rate: f64, now: Instant) {
704 let elapsed = now.saturating_duration_since(self.at).as_secs_f64();
705 self.tokens = (self.tokens + rate * elapsed).min(capacity);
706 self.at = self.at.max(now);
707 }
708
709 /// Takes one token from a bucket of `capacity` refilling at `rate` a second; else how
710 /// long until one is there.
711 fn take(&mut self, capacity: f64, rate: f64, now: Instant) -> Result<(), Duration> {
712 self.refill(capacity, rate, now);
713 if self.tokens >= 1.0 {
714 self.tokens -= 1.0;
715 Ok(())
716 } else {
717 Err(Duration::from_secs_f64((1.0 - self.tokens) / rate))
718 }
719 }
720
721 /// Spends `amount` from a bucket holding a second's worth at `rate`, into debt if need
722 /// be: how long until the debt is repaid.
723 fn spend(&mut self, amount: f64, rate: f64, now: Instant) -> Duration {
724 self.refill(rate, rate, now);
725 self.tokens -= amount;
726 Duration::from_secs_f64((-self.tokens / rate).max(0.0))
727 }
728}
729
730/// What waits to be written to one peer, capped in bytes.
731struct Outbox {
732 stream: TcpStream,
733 queue: Mutex<Queue>,
734 ready: Condvar,
735 most: usize,
736}
737
738#[derive(Default)]
739struct Queue {
740 frames: VecDeque<Vec<u8>>,
741 bytes: usize,
742 closed: bool,
743}
744
745impl Outbox {
746 fn new(stream: TcpStream, most: usize) -> Self {
747 Self {
748 stream,
749 queue: Mutex::default(),
750 ready: Condvar::new(),
751 most,
752 }
753 }
754
755 /// Queues `frame`, or hangs up on a peer that reads too slowly to take it.
756 fn push(&self, frame: Vec<u8>) {
757 let mut queue = self.queue.lock().unwrap();
758 if queue.closed {
759 return;
760 }
761 if queue.bytes + frame.len() > self.most {
762 drop(queue);
763 self.close();
764 return;
765 }
766 queue.bytes += frame.len();
767 queue.frames.push_back(frame);
768 self.ready.notify_one();
769 }
770
771 fn close(&self) {
772 let mut queue = self.queue.lock().unwrap();
773 *queue = Queue {
774 closed: true,
775 ..Queue::default()
776 };
777 self.ready.notify_one();
778 drop(queue);
779 let _ = self.stream.shutdown(Shutdown::Both);
780 }
781
782 /// Writes what is queued until closed.
783 fn drain(&self) {
784 loop {
785 let frame = {
786 let mut queue = self.queue.lock().unwrap();
787 loop {
788 if queue.closed {
789 return;
790 }
791 if let Some(frame) = queue.frames.pop_front() {
792 queue.bytes -= frame.len();
793 break frame;
794 }
795 queue = self.ready.wait(queue).unwrap();
796 }
797 };
798 if (&self.stream).write_all(&frame).is_err() {
799 self.close();
800 return;
801 }
802 }
803 }
804}
805
806#[cfg(test)]
807mod tests {
808 use super::*;
809
810 /// Ten wrong codes in a minute lock an address out for a minute, then two, then four,
811 /// and a quiet hour forgives it.
812 #[test]
813 fn wrong_codes_lock_out_for_longer_each_time() {
814 let config = Config::default();
815 let start = Instant::now();
816 let mut address = Address::new(&config, start);
817 let mut now = start;
818 for expected in [1, 2, 4] {
819 for _ in 0..9 {
820 assert_eq!(address.fail(now, 10), None);
821 }
822 assert_eq!(address.fail(now, 10), Some(MINUTE * expected));
823 now += MINUTE * expected;
824 }
825 now += HOUR;
826 address.refresh(now);
827 assert_eq!(address.strikes, 0);
828 assert!(address.quiet(&config, now));
829 }
830
831 /// Failures spread out over more than a minute never add up to a lockout.
832 #[test]
833 fn slow_wrong_codes_never_lock_out() {
834 let start = Instant::now();
835 let mut address = Address::new(&Config::default(), start);
836 for minute in 0..30 {
837 for second in [0, 30] {
838 let now = start + MINUTE * minute + Duration::from_secs(second);
839 assert_eq!(address.fail(now, 10), None);
840 }
841 }
842 }
843
844 #[test]
845 fn buckets_refill_and_debts_are_waited_out() {
846 let now = Instant::now();
847 let mut joins = Bucket::full(2.0, now);
848 assert!(joins.take(2.0, 1.0, now).is_ok());
849 assert!(joins.take(2.0, 1.0, now).is_ok());
850 assert_eq!(joins.take(2.0, 1.0, now), Err(Duration::from_secs(1)));
851 assert!(joins.take(2.0, 1.0, now + Duration::from_secs(1)).is_ok());
852
853 let mut bytes = Bucket::full(100.0, now);
854 assert_eq!(bytes.spend(100.0, 100.0, now), Duration::ZERO);
855 assert_eq!(bytes.spend(50.0, 100.0, now), Duration::from_millis(500));
856 }
857
858 #[test]
859 fn ipv6_addresses_count_by_their_64() {
860 let relay = Relay {
861 config: Config {
862 trust_forwarded: true,
863 ..Config::default()
864 },
865 state: Mutex::default(),
866 connections: AtomicUsize::new(0),
867 started: Instant::now(),
868 };
869 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
870 let stream = TcpStream::connect(listener.local_addr().unwrap()).unwrap();
871 let head = "GET / HTTP/1.1\r\nX-Forwarded-For: 1.2.3.4, 2001:db8:1:2:3:4:5:6\r\n\r\n";
872 assert_eq!(
873 relay.address(&stream, head),
874 "2001:db8:1:2::".parse::<IpAddr>().unwrap()
875 );
876 assert_eq!(
877 relay.address(&stream, "GET / HTTP/1.1\r\n\r\n"),
878 "127.0.0.1".parse::<IpAddr>().unwrap()
879 );
880 }
881}
crates/relay/src/ws.rs created+224
......@@ -0,0 +1,224 @@
1//! The WebSocket framing both ends speak (RFC 6455): whole text and binary messages, pings and
2//! closes, each message capped in size.
3
4use base64::Engine;
5use sha1::{Digest, Sha1};
6use std::io::{self, Read};
7
8pub const TEXT: u8 = 1;
9pub const BINARY: u8 = 2;
10pub const CLOSE: u8 = 8;
11pub const PING: u8 = 9;
12pub const PONG: u8 = 10;
13/// The longest HTTP head either end reads.
14pub const HEAD: usize = 8 << 10;
15
16#[derive(Debug, PartialEq, Eq)]
17pub enum Message {
18 Text(String),
19 Binary(Vec<u8>),
20 Ping(Vec<u8>),
21 Pong,
22 Close,
23}
24
25/// The `Sec-WebSocket-Accept` answering a `Sec-WebSocket-Key` of `key`.
26pub fn accept(key: &str) -> String {
27 let digest = Sha1::digest(format!("{key}258EAFA5-E914-47DA-95CA-C5AB0DC85B11"));
28 base64::engine::general_purpose::STANDARD.encode(digest)
29}
30
31/// One unfragmented frame holding `payload`, masked with `mask` as a client's must be.
32pub fn frame(opcode: u8, payload: &[u8], mask: Option<[u8; 4]>) -> Vec<u8> {
33 let mut frame = Vec::with_capacity(payload.len() + 14);
34 frame.push(0x80 | opcode);
35 let masked = if mask.is_some() { 0x80 } else { 0 };
36 match payload.len() {
37 length @ 0..=125 => frame.push(masked | length as u8),
38 length @ 126..=0xffff => {
39 frame.push(masked | 126);
40 frame.extend_from_slice(&(length as u16).to_be_bytes());
41 }
42 length => {
43 frame.push(masked | 127);
44 frame.extend_from_slice(&(length as u64).to_be_bytes());
45 }
46 }
47 match mask {
48 Some(mask) => {
49 frame.extend_from_slice(&mask);
50 frame.extend(payload.iter().zip(mask.iter().cycle()).map(|(b, m)| b ^ m));
51 }
52 None => frame.extend_from_slice(payload),
53 }
54 frame
55}
56
57/// An HTTP head, through the blank line that ends it, read a byte at a time so that nothing
58/// after it is consumed.
59pub fn head(from: &mut impl Read) -> io::Result<String> {
60 let mut head = Vec::new();
61 let mut byte = [0];
62 while !head.ends_with(b"\r\n\r\n") {
63 if head.len() >= HEAD {
64 return Err(invalid("An HTTP head too long"));
65 }
66 from.read_exact(&mut byte)?;
67 head.push(byte[0]);
68 }
69 String::from_utf8(head).map_err(|_| invalid("An HTTP head that isn't UTF-8"))
70}
71
72/// The value of header `name` in `head`.
73pub fn header<'a>(head: &'a str, name: &str) -> Option<&'a str> {
74 head.lines().skip(1).find_map(|line| {
75 let (key, value) = line.split_once(':')?;
76 key.trim().eq_ignore_ascii_case(name).then(|| value.trim())
77 })
78}
79
80/// Reads one direction of a connection as whole messages, joining fragments.
81pub struct Reader<R> {
82 inner: R,
83 /// The largest message, in bytes.
84 most: usize,
85 /// Whether frames arrive masked, as a server reads them.
86 masked: bool,
87 /// A message's opcode and its fragments so far.
88 partial: Option<(u8, Vec<u8>)>,
89}
90
91impl<R: Read> Reader<R> {
92 pub fn new(inner: R, most: usize, masked: bool) -> Self {
93 Self {
94 inner,
95 most,
96 masked,
97 partial: None,
98 }
99 }
100
101 pub fn get_mut(&mut self) -> &mut R {
102 &mut self.inner
103 }
104
105 /// The next message, or an error of kind `InvalidData` for one breaking the protocol or
106 /// the cap.
107 pub fn read(&mut self) -> io::Result<Message> {
108 loop {
109 let mut head = [0; 2];
110 self.inner.read_exact(&mut head)?;
111 let (last, opcode) = (head[0] & 0x80 != 0, head[0] & 0x0f);
112 if head[0] & 0x70 != 0 || (head[1] & 0x80 != 0) != self.masked {
113 return Err(invalid(
114 "A WebSocket frame with reserved bits or the wrong mask",
115 ));
116 }
117 let length = match head[1] & 0x7f {
118 126 => {
119 let mut length = [0; 2];
120 self.inner.read_exact(&mut length)?;
121 u64::from(u16::from_be_bytes(length))
122 }
123 127 => {
124 let mut length = [0; 8];
125 self.inner.read_exact(&mut length)?;
126 u64::from_be_bytes(length)
127 }
128 length => u64::from(length),
129 };
130 let control = opcode & 0x08 != 0;
131 let room = match (&self.partial, control) {
132 (_, true) if !last => return Err(invalid("A fragmented control frame")),
133 (_, true) => 125,
134 (Some((_, so_far)), false) => self.most - so_far.len(),
135 (None, false) => self.most,
136 };
137 if length > room as u64 {
138 return Err(invalid("A WebSocket message too large"));
139 }
140 let mut mask = [0; 4];
141 if self.masked {
142 self.inner.read_exact(&mut mask)?;
143 }
144 let mut payload = vec![0; length as usize];
145 self.inner.read_exact(&mut payload)?;
146 if self.masked {
147 for (byte, mask) in payload.iter_mut().zip(mask.iter().cycle()) {
148 *byte ^= mask;
149 }
150 }
151 let (opcode, payload) = match (opcode, &mut self.partial) {
152 (PING, _) => return Ok(Message::Ping(payload)),
153 (PONG, _) => return Ok(Message::Pong),
154 (CLOSE, _) => return Ok(Message::Close),
155 (TEXT | BINARY, None) if last => (opcode, payload),
156 (TEXT | BINARY, None) => {
157 self.partial = Some((opcode, payload));
158 continue;
159 }
160 (0, Some((_, so_far))) => {
161 so_far.extend_from_slice(&payload);
162 if !last {
163 continue;
164 }
165 self.partial.take().expect("a partial message")
166 }
167 _ => return Err(invalid("An unexpected WebSocket frame")),
168 };
169 return Ok(if opcode == TEXT {
170 Message::Text(
171 String::from_utf8(payload).map_err(|_| invalid("Text that isn't UTF-8"))?,
172 )
173 } else {
174 Message::Binary(payload)
175 });
176 }
177 }
178}
179
180pub(crate) fn invalid(message: &'static str) -> io::Error {
181 io::Error::new(io::ErrorKind::InvalidData, message)
182}
183
184#[cfg(test)]
185mod tests {
186 use super::*;
187
188 #[test]
189 fn frames_read_back_masked_or_not_and_fragments_join() {
190 let mut bytes = frame(TEXT, b"welcome 1", Some([1, 2, 3, 4]));
191 bytes.extend(frame(BINARY, &[7; 70_000], Some([9, 8, 7, 6])));
192 bytes.extend(frame(PING, b"hi", Some([0; 4])));
193 let mut reader = Reader::new(&bytes[..], 1 << 20, true);
194 assert_eq!(reader.read().unwrap(), Message::Text("welcome 1".into()));
195 assert_eq!(reader.read().unwrap(), Message::Binary(vec![7; 70_000]));
196 assert_eq!(reader.read().unwrap(), Message::Ping(b"hi".to_vec()));
197
198 // "ab" then "cd" as a fragmented binary message, a ping between them.
199 let mut fragments = vec![BINARY, 2, b'a', b'b'];
200 fragments.extend(frame(PING, b"", None));
201 fragments.extend([0x80, 2, b'c', b'd']);
202 let mut reader = Reader::new(&fragments[..], 4, false);
203 assert_eq!(reader.read().unwrap(), Message::Ping(vec![]));
204 assert_eq!(reader.read().unwrap(), Message::Binary(b"abcd".to_vec()));
205 }
206
207 #[test]
208 fn a_message_over_the_cap_or_masked_wrongly_is_refused() {
209 let big = frame(BINARY, &[0; 10], None);
210 let error = Reader::new(&big[..], 9, false).read().unwrap_err();
211 assert_eq!(error.kind(), io::ErrorKind::InvalidData);
212 let unmasked = frame(BINARY, b"x", None);
213 assert!(Reader::new(&unmasked[..], 9, true).read().is_err());
214 }
215
216 /// RFC 6455's own example.
217 #[test]
218 fn accepts_the_rfc_key() {
219 assert_eq!(
220 accept("dGhlIHNhbXBsZSBub25jZQ=="),
221 "s3pPLMBiTxaQ9kYGzzhZRbK+xOo="
222 );
223 }
224}
crates/relay/tests/relay.rs created+365
......@@ -0,0 +1,365 @@
1//! The relay as shipped, run as its own process and spoken to as a WebSocket.
2
3use relay::ws::{self, Message};
4use std::{
5 io::{BufRead, BufReader, Read, Write},
6 net::{SocketAddr, TcpStream},
7 process::{Child, Command, Stdio},
8 thread,
9 time::Duration,
10};
11
12/// A relay process, killed when dropped.
13struct Relay {
14 child: Child,
15 address: SocketAddr,
16}
17
18impl Relay {
19 fn start(options: &[&str]) -> Relay {
20 let mut child = Command::new(env!("CARGO_BIN_EXE_snowbound-relay"))
21 .args(["--listen", "127.0.0.1:0"])
22 .args(options)
23 .env_clear()
24 .stdout(Stdio::piped())
25 .spawn()
26 .unwrap();
27 let mut line = String::new();
28 BufReader::new(child.stdout.take().unwrap())
29 .read_line(&mut line)
30 .unwrap();
31 let address = line.trim().rsplit(' ').next().unwrap().parse().unwrap();
32 Relay { child, address }
33 }
34
35 fn get(&self, path: &str) -> String {
36 let mut stream = TcpStream::connect(self.address).unwrap();
37 write!(stream, "GET {path} HTTP/1.1\r\nHost: relay\r\n\r\n").unwrap();
38 let mut response = String::new();
39 stream.read_to_string(&mut response).unwrap();
40 response
41 }
42
43 fn join(&self, path: &str) -> Result<Peer, String> {
44 self.join_from(path, None)
45 }
46
47 /// Joins as if from `address`, as a proxy would say.
48 fn join_from(&self, path: &str, address: Option<&str>) -> Result<Peer, String> {
49 let mut stream = TcpStream::connect(self.address).unwrap();
50 stream
51 .set_read_timeout(Some(Duration::from_secs(10)))
52 .unwrap();
53 let forwarded = address
54 .map(|address| format!("X-Forwarded-For: {address}\r\n"))
55 .unwrap_or_default();
56 write!(
57 stream,
58 "GET {path} HTTP/1.1\r\nHost: relay\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\
59 Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Version: 13\r\n{forwarded}\r\n"
60 )
61 .unwrap();
62 let mut reader = BufReader::new(stream.try_clone().unwrap());
63 let head = ws::head(&mut reader).map_err(|error| error.to_string())?;
64 if !head.starts_with("HTTP/1.1 101") {
65 return Err(head);
66 }
67 assert!(head.contains("s3pPLMBiTxaQ9kYGzzhZRbK+xOo="));
68 Ok(Peer {
69 stream,
70 reader: ws::Reader::new(reader, 1 << 20, false),
71 })
72 }
73}
74
75impl Drop for Relay {
76 fn drop(&mut self) {
77 let _ = self.child.kill();
78 let _ = self.child.wait();
79 }
80}
81
82struct Peer {
83 stream: TcpStream,
84 reader: ws::Reader<BufReader<TcpStream>>,
85}
86
87impl Peer {
88 fn send(&mut self, opcode: u8, payload: &[u8]) {
89 self.stream
90 .write_all(&ws::frame(opcode, payload, Some([1, 2, 3, 4])))
91 .unwrap();
92 }
93
94 fn say(&mut self, text: &str) {
95 self.send(ws::TEXT, text.as_bytes());
96 }
97
98 fn send_to(&mut self, slot: u32, bytes: &[u8]) {
99 self.send(ws::BINARY, &[&slot.to_be_bytes()[..], bytes].concat());
100 }
101
102 fn hear(&mut self) -> Message {
103 self.reader.read().unwrap()
104 }
105
106 fn text(&mut self) -> String {
107 match self.hear() {
108 Message::Text(text) => text,
109 other => panic!("heard {other:?}"),
110 }
111 }
112
113 /// Whether the relay has hung up, waiting at most `patience`.
114 fn closed(&mut self, patience: Duration) -> bool {
115 self.stream.set_read_timeout(Some(patience)).unwrap();
116 loop {
117 match self.reader.read() {
118 Ok(_) => continue,
119 Err(error) => {
120 return !matches!(
121 error.kind(),
122 std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
123 );
124 }
125 }
126 }
127 }
128}
129
130fn status(refusal: Result<Peer, String>) -> String {
131 match refusal {
132 Ok(_) => panic!("joined"),
133 Err(head) => head.lines().next().unwrap().to_owned(),
134 }
135}
136
137/// The number of the code a claim was given.
138fn claim(relay: &Relay) -> (Peer, u32) {
139 let mut owner = relay.join("/v1/claim").unwrap();
140 let number = owner
141 .text()
142 .strip_prefix("nameplate ")
143 .unwrap()
144 .parse()
145 .unwrap();
146 assert_eq!(owner.text(), "welcome 1");
147 (owner, number)
148}
149
150#[test]
151fn health_answers() {
152 let relay = Relay::start(&[]);
153 let response = relay.get("/health");
154 assert!(response.starts_with("HTTP/1.1 200"), "{response}");
155 assert!(response.contains("\"rooms\":0"), "{response}");
156 assert!(relay.get("/nothing").starts_with("HTTP/1.1 404"));
157 assert!(relay.get("/v1/room/abc").starts_with("HTTP/1.1 400"));
158}
159
160/// Peers in a room hear who comes and goes and pass each other messages, which arrive
161/// labelled with the sender's slot.
162#[test]
163fn peers_in_a_room_pass_messages() {
164 let relay = Relay::start(&[]);
165 let mut ada = relay.join("/v1/room/0123abcd").unwrap();
166 assert_eq!(ada.text(), "welcome 1");
167 let mut grace = relay.join("/v1/room/0123abcd").unwrap();
168 assert_eq!(grace.text(), "welcome 2 1");
169 assert_eq!(ada.text(), "joined 2");
170 grace.send_to(1, b"hello");
171 assert_eq!(
172 ada.hear(),
173 Message::Binary([&2u32.to_be_bytes()[..], b"hello"].concat())
174 );
175 ada.send_to(2, b"hi");
176 assert_eq!(
177 grace.hear(),
178 Message::Binary([&1u32.to_be_bytes()[..], b"hi"].concat())
179 );
180 ada.send(ws::PING, b"p");
181 assert_eq!(ada.hear(), Message::Pong);
182 let health = relay.get("/health");
183 assert!(health.contains("\"rooms\":1,\"peers\":2"), "{health}");
184 drop(grace);
185 assert_eq!(ada.text(), "left 2");
186}
187
188/// One joining a code's room hears only its owner until the owner says it met it; then the
189/// others in the room, and they it.
190#[test]
191fn a_code_admits_whom_its_owner_met() {
192 let relay = Relay::start(&[]);
193 let (mut owner, number) = claim(&relay);
194 let path = format!("/v1/room/code-{number}");
195 let mut ada = relay.join(&path).unwrap();
196 assert_eq!(ada.text(), "welcome 2 1");
197 assert_eq!(owner.text(), "joined 2");
198 owner.say("met 2");
199 let mut grace = relay.join(&path).unwrap();
200 assert_eq!(grace.text(), "welcome 3 1");
201 assert_eq!(owner.text(), "joined 3");
202 // Ada can't reach Grace before the owner meets her.
203 ada.send_to(3, b"psst");
204 grace.send_to(1, b"hello");
205 assert_eq!(
206 owner.hear(),
207 Message::Binary([&3u32.to_be_bytes()[..], b"hello"].concat())
208 );
209 owner.say("met 3");
210 assert_eq!(grace.text(), "joined 2");
211 assert_eq!(ada.text(), "joined 3");
212 ada.send_to(3, b"hi");
213 assert_eq!(
214 grace.hear(),
215 Message::Binary([&2u32.to_be_bytes()[..], b"hi"].concat())
216 );
217 // Only the owner's word counts.
218 let mut mallory = relay.join(&path).unwrap();
219 assert_eq!(mallory.text(), "welcome 4 1");
220 ada.say("met 4");
221 grace.send_to(4, b"anyone?");
222 owner.say("failed 4");
223 assert!(mallory.closed(Duration::from_secs(5)));
224 assert_eq!(owner.text(), "joined 4");
225 assert_eq!(owner.text(), "left 4");
226}
227
228#[test]
229fn an_unknown_code_or_room_is_refused() {
230 let relay = Relay::start(&[]);
231 assert_eq!(
232 status(relay.join("/v1/room/code-5")),
233 "HTTP/1.1 404 Not Found"
234 );
235 assert_eq!(
236 status(relay.join("/v1/room/UPPER")),
237 "HTTP/1.1 404 Not Found"
238 );
239 assert_eq!(status(relay.join("/v1/room/")), "HTTP/1.1 404 Not Found");
240}
241
242/// Five wrong codes burn a code: it admits no one new, and its owner hears so.
243#[test]
244fn wrong_codes_burn_a_code() {
245 let relay = Relay::start(&["--failures-per-minute", "100"]);
246 let (mut owner, number) = claim(&relay);
247 let path = format!("/v1/room/code-{number}");
248 for slot in 2..7 {
249 let mut guess = relay.join(&path).unwrap();
250 assert_eq!(guess.text(), format!("welcome {slot} 1"));
251 assert_eq!(owner.text(), format!("joined {slot}"));
252 owner.say(&format!("failed {slot}"));
253 assert!(guess.closed(Duration::from_secs(5)));
254 assert_eq!(owner.text(), format!("left {slot}"));
255 }
256 assert_eq!(owner.text(), "burned");
257 assert_eq!(status(relay.join(&path)), "HTTP/1.1 410 Gone");
258}
259
260/// An owner coming back asks for its code's number again, and gets it while no one else
261/// has claimed it.
262#[test]
263fn an_owner_comes_back_to_its_number() {
264 let relay = Relay::start(&[]);
265 let (owner, number) = claim(&relay);
266 let mut ada = relay.join("/v1/room/0123abcd").unwrap();
267 assert_eq!(ada.text(), "welcome 1");
268 let mut taken = relay
269 .join(&format!("/v1/claim?nameplate={number}"))
270 .unwrap();
271 assert_ne!(taken.text(), format!("nameplate {number}"));
272 drop(owner);
273 thread::sleep(Duration::from_millis(200));
274 let mut again = relay
275 .join(&format!("/v1/claim?nameplate={number}"))
276 .unwrap();
277 assert_eq!(again.text(), format!("nameplate {number}"));
278}
279
280/// Silence, or leaving before the owner says it met, counts as a wrong code; too many from
281/// an address in a minute lock it out, while another address is still let in.
282#[test]
283fn an_address_trying_too_many_codes_is_locked_out() {
284 let relay = Relay::start(&[
285 "--trust-forwarded",
286 "true",
287 "--failures-per-minute",
288 "3",
289 "--pending",
290 "1",
291 "--burn-after",
292 "100",
293 ]);
294 let (mut owner, number) = claim(&relay);
295 let path = format!("/v1/room/code-{number}");
296 let mallory = Some("203.0.113.9");
297 // Silent: the relay gives up on it.
298 let mut silent = relay.join_from(&path, mallory).unwrap();
299 assert!(silent.closed(Duration::from_secs(5)));
300 // Gone before the owner's word.
301 drop(relay.join_from(&path, mallory).unwrap());
302 // Refused by the owner.
303 let mut refused = relay.join_from(&path, mallory).unwrap();
304 assert_eq!(refused.text(), "welcome 4 1");
305 owner.say("failed 4");
306 assert!(refused.closed(Duration::from_secs(5)));
307 let locked = status(relay.join_from(&path, mallory));
308 assert_eq!(locked, "HTTP/1.1 429 Too Many Requests");
309 let Err(again) = relay.join_from(&path, mallory) else {
310 panic!("let in while locked out");
311 };
312 assert!(
313 again.contains("Retry-After: 60"),
314 "a minute's lockout: {again}"
315 );
316 // Its /64 neighbour is the same subscriber; another address isn't.
317 let mut ada = relay.join_from(&path, Some("198.51.100.1")).unwrap();
318 assert_eq!(ada.text(), "welcome 5 1");
319}
320
321/// At most three waiting at once from an address allowed three wrong codes a minute: it
322/// can't run many guesses side by side.
323#[test]
324fn guesses_waiting_at_once_count_against_the_limit() {
325 let relay = Relay::start(&["--failures-per-minute", "3"]);
326 let (_owner, number) = claim(&relay);
327 let path = format!("/v1/room/code-{number}");
328 let _waiting: Vec<Peer> = (0..3).map(|_| relay.join(&path).unwrap()).collect();
329 assert_eq!(status(relay.join(&path)), "HTTP/1.1 429 Too Many Requests");
330}
331
332#[test]
333fn joins_rooms_and_messages_are_capped() {
334 let relay = Relay::start(&[
335 "--max-room-peers",
336 "2",
337 "--joins-per-minute",
338 "5",
339 "--max-message",
340 "1000",
341 ]);
342 let mut ada = relay.join("/v1/room/aaaa").unwrap();
343 let _grace = relay.join("/v1/room/aaaa").unwrap();
344 assert_eq!(
345 status(relay.join("/v1/room/aaaa")),
346 "HTTP/1.1 503 Service Unavailable"
347 );
348 let _bbbb = relay.join("/v1/room/bbbb").unwrap();
349 let _fifth = relay.join("/v1/room/cccc").unwrap();
350 assert_eq!(
351 status(relay.join("/v1/room/dddd")),
352 "HTTP/1.1 429 Too Many Requests"
353 );
354 assert_eq!(ada.text(), "welcome 1");
355 ada.send_to(2, &[0; 2000]);
356 assert!(ada.closed(Duration::from_secs(5)));
357}
358
359#[test]
360fn a_silent_connection_is_closed() {
361 let relay = Relay::start(&["--idle", "1"]);
362 let mut quiet = relay.join("/v1/room/bbbb").unwrap();
363 assert_eq!(quiet.text(), "welcome 1");
364 assert!(quiet.closed(Duration::from_secs(5)));
365}
crates/snowbound/src/live.rs+11-7
......@@ -1,8 +1,9 @@
11//! Live presence in the window (feature `live`, on where `SNOWBOUND_LIVE` is `network` or
2//! `loopback`): the others with the notebook open as avatars beside the search box, a dot on
3//! the tab of the page each has open, and their carets on the page, each in a colour of their
4//! own. They meet through the notebook's identity, or through `SNOWBOUND_LIVE_CODE` where it
5//! is set (`resources/live-share.md`).
2//! `loopback`, or `SNOWBOUND_LIVE_RELAY` names a relay, `wss://live.example.net`): the others
3//! with the notebook open as avatars beside the search box, a dot on the tab of the page each
4//! has open, and their carets on the page, each in a colour of their own. They meet through
5//! the notebook's identity, or through `SNOWBOUND_LIVE_CODE` where it is set; through a relay,
6//! its words alone ask the relay to number it (`resources/live-share.md`).
67
78use crate::{Command, Session, State, TAB_ROW, page, platform};
89use canvas::{document::TextPosition, editor::TextOutline};
......@@ -37,9 +38,10 @@ fn reach() -> Option<Reach> {
3738impl State {
3839 /// Joins the open notebook's room, leaving any other, and says where this window is.
3940 pub(crate) fn follow_peers(&mut self) {
40 let Some(reach) = reach() else {
41 let (reach, relay) = (reach(), std::env::var("SNOWBOUND_LIVE_RELAY").ok());
42 if reach.is_none() && relay.is_none() {
4143 return;
42 };
44 }
4345 let room = match std::env::var("SNOWBOUND_LIVE_CODE") {
4446 Ok(code) => Some(Room::Code(code)),
4547 Err(_) => self
......@@ -57,7 +59,9 @@ impl State {
5759 .flatten();
5860 let live = Hello::new(self.author.clone(), picture)
5961 .and_then(|me| {
60 live::Live::start(me, &room, Some(reach), move || redraw.wake_by_ref())
62 live::Live::start(me, &room, reach, relay.as_deref(), move || {
63 redraw.wake_by_ref()
64 })
6165 })
6266 .inspect_err(|error| eprintln!("Live presence: {error}"))
6367 .ok();
tools/release_relay.py created+47
......@@ -0,0 +1,47 @@
1#!/usr/bin/env python3
2"""Builds snowbound-relay, the Live Share relay (crates/relay), as one static executable for each
3Linux architecture a server runs, with the files that deploy it, into target/relay/ by default.
4It ships on its own, outside the app's builds and latest.json; see crates/relay/README.md."""
5import argparse
6import hashlib
7import os
8from pathlib import Path
9import shutil
10import subprocess
11
12ROOT = Path(__file__).resolve().parents[1]
13ARCHITECTURES = ['x86_64', 'aarch64']
14
15
16def build(arch, output):
17 """Linked against Rust's own musl by its own lld, so the file needs nothing of the Linux it
18 runs on."""
19 triple = f'{arch}-unknown-linux-musl'
20 subprocess.run(['rustup', 'target', 'add', triple], check=True, capture_output=True)
21 environment = {**os.environ, f'CARGO_TARGET_{triple.upper().replace("-", "_")}_LINKER': 'rust-lld',
22 'CARGO_PROFILE_RELEASE_STRIP': 'symbols'}
23 subprocess.run(['cargo', 'build', '--locked', '--release', '-p', 'relay', '--target', triple],
24 cwd=ROOT, env=environment, check=True)
25 target = Path(os.environ.get('CARGO_TARGET_DIR') or ROOT / 'target')
26 file = output / f'snowbound-relay-linux-{arch}'
27 shutil.copy(target / triple / 'release/snowbound-relay', file)
28 return file
29
30
31def main():
32 parser = argparse.ArgumentParser(description=__doc__)
33 parser.add_argument('--output', type=Path, default=ROOT / 'target/relay')
34 parser.add_argument('--architectures', nargs='+', choices=ARCHITECTURES, default=ARCHITECTURES)
35 args = parser.parse_args()
36 args.output.mkdir(parents=True, exist_ok=True)
37 files = [build(arch, args.output) for arch in args.architectures]
38 for deployed in (ROOT / 'crates/relay/deploy').iterdir():
39 shutil.copy(deployed, args.output)
40 shutil.copy(ROOT / 'crates/relay/README.md', args.output)
41 sums = ''.join(f'{hashlib.sha256(file.read_bytes()).hexdigest()} {file.name}\n' for file in files)
42 (args.output / 'SHA256SUMS').write_text(sums)
43 print(f'{args.output}:\n{sums}', end='')
44
45
46if __name__ == '__main__':
47 main()