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/// A group message goes once to the relay, which copies it to everyone else in the room or
189/// to the slots it names, marked as the group's; the relay counts what it passed on.
190#[test]
191fn group_messages_are_copied() {
192 let relay = Relay::start(&[]);
193 let mut peers: Vec<Peer> = (1..=3)
194 .map(|_| relay.join("/v1/room/0123abcd").unwrap())
195 .collect();
196 assert_eq!(peers[0].text(), "welcome 1");
197 assert_eq!(peers[0].text(), "joined 2");
198 assert_eq!(peers[0].text(), "joined 3");
199 assert_eq!(peers[1].text(), "welcome 2 1");
200 assert_eq!(peers[1].text(), "joined 3");
201 assert_eq!(peers[2].text(), "welcome 3 1 2");
202 let from = |slot: u32, bytes: &[u8]| {
203 Message::Binary([&(slot | relay::GROUP).to_be_bytes()[..], bytes].concat())
204 };
205 peers[0].send_to(relay::BROADCAST, b"everyone");
206 assert_eq!(peers[1].hear(), from(1, b"everyone"));
207 assert_eq!(peers[2].hear(), from(1, b"everyone"));
208 let named = [&3u32.to_be_bytes()[..], b"just you"].concat();
209 peers[1].send_to(relay::GROUP | 1, &named);
210 assert_eq!(peers[2].hear(), from(2, b"just you"));
211 peers[0].send_to(2, b"stream");
212 assert_eq!(
213 peers[1].hear(),
214 Message::Binary([&1u32.to_be_bytes()[..], b"stream"].concat())
215 );
216 let health = relay.get("/health");
217 let given = 2 * (4 + 8) + (4 + 8) + (4 + 6);
218 assert!(
219 health.contains(&format!("\"bytes_out\":{given}")),
220 "{health}"
221 );
222}
223
224/// One joining a code's room hears only its owner until the owner says it met it; then the
225/// others in the room, and they it.
226#[test]
227fn a_code_admits_whom_its_owner_met() {
228 let relay = Relay::start(&[]);
229 let (mut owner, number) = claim(&relay);
230 let path = format!("/v1/room/code-{number}");
231 let mut ada = relay.join(&path).unwrap();
232 assert_eq!(ada.text(), "welcome 2 1");
233 assert_eq!(owner.text(), "joined 2");
234 owner.say("met 2");
235 let mut grace = relay.join(&path).unwrap();
236 assert_eq!(grace.text(), "welcome 3 1");
237 assert_eq!(owner.text(), "joined 3");
238 // Ada can't reach Grace before the owner meets her.
239 ada.send_to(3, b"psst");
240 grace.send_to(1, b"hello");
241 assert_eq!(
242 owner.hear(),
243 Message::Binary([&3u32.to_be_bytes()[..], b"hello"].concat())
244 );
245 owner.say("met 3");
246 assert_eq!(grace.text(), "joined 2");
247 assert_eq!(ada.text(), "joined 3");
248 ada.send_to(3, b"hi");
249 assert_eq!(
250 grace.hear(),
251 Message::Binary([&2u32.to_be_bytes()[..], b"hi"].concat())
252 );
253 // Only the owner's word counts.
254 let mut mallory = relay.join(&path).unwrap();
255 assert_eq!(mallory.text(), "welcome 4 1");
256 ada.say("met 4");
257 grace.send_to(4, b"anyone?");
258 owner.say("failed 4");
259 assert!(mallory.closed(Duration::from_secs(5)));
260 assert_eq!(owner.text(), "joined 4");
261 assert_eq!(owner.text(), "left 4");
262}
263
264#[test]
265fn an_unknown_code_or_room_is_refused() {
266 let relay = Relay::start(&[]);
267 assert_eq!(
268 status(relay.join("/v1/room/code-5")),
269 "HTTP/1.1 404 Not Found"
270 );
271 assert_eq!(
272 status(relay.join("/v1/room/UPPER")),
273 "HTTP/1.1 404 Not Found"
274 );
275 assert_eq!(status(relay.join("/v1/room/")), "HTTP/1.1 404 Not Found");
276}
277
278/// Five wrong codes burn a code: it admits no one new, and its owner hears so.
279#[test]
280fn wrong_codes_burn_a_code() {
281 let relay = Relay::start(&["--failures-per-minute", "100"]);
282 let (mut owner, number) = claim(&relay);
283 let path = format!("/v1/room/code-{number}");
284 for slot in 2..7 {
285 let mut guess = relay.join(&path).unwrap();
286 assert_eq!(guess.text(), format!("welcome {slot} 1"));
287 assert_eq!(owner.text(), format!("joined {slot}"));
288 owner.say(&format!("failed {slot}"));
289 assert!(guess.closed(Duration::from_secs(5)));
290 assert_eq!(owner.text(), format!("left {slot}"));
291 }
292 assert_eq!(owner.text(), "burned");
293 assert_eq!(status(relay.join(&path)), "HTTP/1.1 410 Gone");
294}
295
296/// An owner coming back asks for its code's number again, and gets it while no one else
297/// has claimed it.
298#[test]
299fn an_owner_comes_back_to_its_number() {
300 let relay = Relay::start(&[]);
301 let (owner, number) = claim(&relay);
302 let mut ada = relay.join("/v1/room/0123abcd").unwrap();
303 assert_eq!(ada.text(), "welcome 1");
304 let mut taken = relay
305 .join(&format!("/v1/claim?nameplate={number}"))
306 .unwrap();
307 assert_ne!(taken.text(), format!("nameplate {number}"));
308 drop(owner);
309 thread::sleep(Duration::from_millis(200));
310 let mut again = relay
311 .join(&format!("/v1/claim?nameplate={number}"))
312 .unwrap();
313 assert_eq!(again.text(), format!("nameplate {number}"));
314}
315
316/// Silence, or leaving before the owner says it met, counts as a wrong code; too many from
317/// an address in a minute lock it out, while another address is still let in.
318#[test]
319fn an_address_trying_too_many_codes_is_locked_out() {
320 let relay = Relay::start(&[
321 "--trust-forwarded",
322 "true",
323 "--failures-per-minute",
324 "3",
325 "--pending",
326 "1",
327 "--burn-after",
328 "100",
329 ]);
330 let (mut owner, number) = claim(&relay);
331 let path = format!("/v1/room/code-{number}");
332 let mallory = Some("203.0.113.9");
333 // Silent: the relay gives up on it.
334 let mut silent = relay.join_from(&path, mallory).unwrap();
335 assert!(silent.closed(Duration::from_secs(5)));
336 // Gone before the owner's word.
337 drop(relay.join_from(&path, mallory).unwrap());
338 // Refused by the owner.
339 let mut refused = relay.join_from(&path, mallory).unwrap();
340 assert_eq!(refused.text(), "welcome 4 1");
341 owner.say("failed 4");
342 assert!(refused.closed(Duration::from_secs(5)));
343 let locked = status(relay.join_from(&path, mallory));
344 assert_eq!(locked, "HTTP/1.1 429 Too Many Requests");
345 let Err(again) = relay.join_from(&path, mallory) else {
346 panic!("let in while locked out");
347 };
348 assert!(
349 again.contains("Retry-After: 60"),
350 "a minute's lockout: {again}"
351 );
352 // Its /64 neighbour is the same subscriber; another address isn't.
353 let mut ada = relay.join_from(&path, Some("198.51.100.1")).unwrap();
354 assert_eq!(ada.text(), "welcome 5 1");
355}
356
357/// At most three waiting at once from an address allowed three wrong codes a minute: it
358/// can't run many guesses side by side.
359#[test]
360fn guesses_waiting_at_once_count_against_the_limit() {
361 let relay = Relay::start(&["--failures-per-minute", "3"]);
362 let (_owner, number) = claim(&relay);
363 let path = format!("/v1/room/code-{number}");
364 let _waiting: Vec<Peer> = (0..3).map(|_| relay.join(&path).unwrap()).collect();
365 assert_eq!(status(relay.join(&path)), "HTTP/1.1 429 Too Many Requests");
366}
367
368#[test]
369fn joins_rooms_and_messages_are_capped() {
370 let relay = Relay::start(&[
371 "--max-room-peers",
372 "2",
373 "--joins-per-minute",
374 "5",
375 "--max-message",
376 "1000",
377 ]);
378 let mut ada = relay.join("/v1/room/aaaa").unwrap();
379 let _grace = relay.join("/v1/room/aaaa").unwrap();
380 assert_eq!(
381 status(relay.join("/v1/room/aaaa")),
382 "HTTP/1.1 503 Service Unavailable"
383 );
384 let _bbbb = relay.join("/v1/room/bbbb").unwrap();
385 let _fifth = relay.join("/v1/room/cccc").unwrap();
386 assert_eq!(
387 status(relay.join("/v1/room/dddd")),
388 "HTTP/1.1 429 Too Many Requests"
389 );
390 assert_eq!(ada.text(), "welcome 1");
391 ada.send_to(2, &[0; 2000]);
392 assert!(ada.closed(Duration::from_secs(5)));
393}
394
395#[test]
396fn a_silent_connection_is_closed() {
397 let relay = Relay::start(&["--idle", "1"]);
398 let mut quiet = relay.join("/v1/room/bbbb").unwrap();
399 assert_eq!(quiet.text(), "welcome 1");
400 assert!(quiet.closed(Duration::from_secs(5)));
401}