| 1 | //! The relay as shipped, run as its own process and spoken to as a WebSocket. |
| 2 | |
| 3 | use relay::ws::{self, Message}; |
| 4 | use 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. |
| 13 | struct Relay { |
| 14 | child: Child, |
| 15 | address: SocketAddr, |
| 16 | } |
| 17 | |
| 18 | impl 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 | |
| 75 | impl Drop for Relay { |
| 76 | fn drop(&mut self) { |
| 77 | let _ = self.child.kill(); |
| 78 | let _ = self.child.wait(); |
| 79 | } |
| 80 | } |
| 81 | |
| 82 | struct Peer { |
| 83 | stream: TcpStream, |
| 84 | reader: ws::Reader<BufReader<TcpStream>>, |
| 85 | } |
| 86 | |
| 87 | impl 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 | |
| 130 | fn 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. |
| 138 | fn 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] |
| 151 | fn 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] |
| 163 | fn 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] |
| 191 | fn 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] |
| 227 | fn 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] |
| 265 | fn 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] |
| 280 | fn 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] |
| 299 | fn 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] |
| 319 | fn 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] |
| 360 | fn 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] |
| 369 | fn 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] |
| 396 | fn 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 | } |