| 1 | use super::*; |
| 2 | use std::time::Instant; |
| 3 | |
| 4 | /// The code for room `number` and `secret`. |
| 5 | fn code(number: u32, secret: &str) -> String { |
| 6 | super::code::format(number, secret).unwrap() |
| 7 | } |
| 8 | |
| 9 | fn hello(name: &str) -> Hello { |
| 10 | Hello::new(name.into(), Some(vec![1, 2, 3])).unwrap() |
| 11 | } |
| 12 | |
| 13 | /// Waits until `done` holds of `live`'s peers. |
| 14 | fn until(live: &Live, done: impl Fn(&[Peer]) -> bool) -> Vec<Peer> { |
| 15 | let deadline = Instant::now() + Duration::from_secs(10); |
| 16 | loop { |
| 17 | let peers = live.peers(); |
| 18 | if done(&peers) { |
| 19 | return peers; |
| 20 | } |
| 21 | assert!(Instant::now() < deadline, "peers stayed {peers:?}"); |
| 22 | thread::sleep(Duration::from_millis(20)); |
| 23 | } |
| 24 | } |
| 25 | |
| 26 | fn caret(offset: u32) -> Presence { |
| 27 | let spot = Spot { |
| 28 | text: Guid { |
| 29 | guid: [7; 16], |
| 30 | n: 3, |
| 31 | }, |
| 32 | offset, |
| 33 | }; |
| 34 | Presence { |
| 35 | section: Some([5; 16]), |
| 36 | page: Some(Guid { |
| 37 | guid: [6; 16], |
| 38 | n: 1, |
| 39 | }), |
| 40 | caret: Some(Caret { |
| 41 | anchor: spot, |
| 42 | focus: spot, |
| 43 | }), |
| 44 | } |
| 45 | } |
| 46 | |
| 47 | /// Two ends of one code meet, greet each other by name and picture, hear each other's newest |
| 48 | /// caret, and see the other leave. |
| 49 | #[test] |
| 50 | fn peers_meet_and_follow_presence() { |
| 51 | let room = Room::join(&code(7, "ABCDEF"), ""); |
| 52 | let ada = Live::start(hello("Ada"), &room, None, None, |_| {}).unwrap(); |
| 53 | let grace = Live::start(hello("Grace"), &room, None, None, |_| {}).unwrap(); |
| 54 | ada.set_presence(caret(1)); |
| 55 | ada.connect(grace.address()); |
| 56 | let seen = until(&grace, |peers| { |
| 57 | peers.iter().any(|peer| peer.presence.is_some()) |
| 58 | }); |
| 59 | assert_eq!(seen[0].hello.name, "Ada"); |
| 60 | assert_eq!(seen[0].hello.picture.as_deref(), Some(&[1, 2, 3][..])); |
| 61 | assert_eq!(seen[0].presence, Some(caret(1))); |
| 62 | until(&ada, |peers| { |
| 63 | peers.len() == 1 && peers[0].hello.name == "Grace" |
| 64 | }); |
| 65 | for offset in 2..20 { |
| 66 | ada.set_presence(caret(offset)); |
| 67 | } |
| 68 | until(&grace, |peers| peers[0].presence == Some(caret(19))); |
| 69 | drop(ada); |
| 70 | until(&grace, <[Peer]>::is_empty); |
| 71 | } |
| 72 | |
| 73 | /// A peer holding another code never meets: its first frame does not open. |
| 74 | #[test] |
| 75 | fn another_code_never_meets() { |
| 76 | let ada = Live::start( |
| 77 | hello("Ada"), |
| 78 | &Room::join(&code(7, "ABCDEF"), ""), |
| 79 | None, |
| 80 | None, |
| 81 | |_| {}, |
| 82 | ) |
| 83 | .unwrap(); |
| 84 | let mallory = Live::start( |
| 85 | hello("Mallory"), |
| 86 | &Room::join(&code(7, "ABCDEG"), ""), |
| 87 | None, |
| 88 | None, |
| 89 | |_| {}, |
| 90 | ) |
| 91 | .unwrap(); |
| 92 | mallory.connect(ada.address()); |
| 93 | thread::sleep(Duration::from_millis(500)); |
| 94 | assert!(ada.peers().is_empty() && mallory.peers().is_empty()); |
| 95 | } |
| 96 | |
| 97 | /// A field or a kind a later version adds is skipped by this one. |
| 98 | #[test] |
| 99 | fn later_fields_and_kinds_are_skipped() { |
| 100 | #[derive(Encode)] |
| 101 | #[cbor(map)] |
| 102 | struct Later { |
| 103 | #[cbor(n(0), with = "minicbor::bytes")] |
| 104 | section: Option<[u8; 16]>, |
| 105 | #[n(9)] |
| 106 | mood: String, |
| 107 | } |
| 108 | use minicbor::Encode; |
| 109 | let body = minicbor::to_vec(Later { |
| 110 | section: Some([5; 16]), |
| 111 | mood: "curious".into(), |
| 112 | }) |
| 113 | .unwrap(); |
| 114 | let presence: Presence = minicbor::decode(&body).unwrap(); |
| 115 | assert_eq!( |
| 116 | presence, |
| 117 | Presence { |
| 118 | section: Some([5; 16]), |
| 119 | ..Presence::default() |
| 120 | } |
| 121 | ); |
| 122 | |
| 123 | let room = Room::join(&code(4, "QJETHR"), ""); |
| 124 | let grace = Live::start(hello("Grace"), &room, None, None, |_| {}).unwrap(); |
| 125 | // A later version: it greets, says something new, then where it is. |
| 126 | let later = thread::spawn(move || { |
| 127 | let mut stream = TcpStream::connect(grace.address()).unwrap(); |
| 128 | let (mut send, mut receive) = wire::open( |
| 129 | &mut stream, |
| 130 | Side::Initiator, |
| 131 | &room.tag().unwrap(), |
| 132 | &room.secret(), |
| 133 | ) |
| 134 | .unwrap(); |
| 135 | send.send(&mut stream, kind::HELLO, &hello("Later")) |
| 136 | .unwrap(); |
| 137 | receive.receive(&mut stream).unwrap(); |
| 138 | send.send(&mut stream, 999, &"a chat message").unwrap(); |
| 139 | send.send(&mut stream, kind::PRESENCE, &caret(4)).unwrap(); |
| 140 | until(&grace, |peers| { |
| 141 | peers |
| 142 | .first() |
| 143 | .is_some_and(|peer| peer.presence == Some(caret(4))) |
| 144 | }); |
| 145 | }); |
| 146 | later.join().unwrap(); |
| 147 | } |
| 148 | |
| 149 | /// Two ends of one notebook find each other by mDNS on this computer's loopback. |
| 150 | #[test] |
| 151 | #[ignore = "multicasts mDNS on the loopback interface"] |
| 152 | fn peers_find_each_other_on_loopback() { |
| 153 | let room = Room::Notebook([9; 16]); |
| 154 | let ada = Live::start(hello("Ada"), &room, Some(Reach::Loopback), None, |_| {}).unwrap(); |
| 155 | let grace = Live::start(hello("Grace"), &room, Some(Reach::Loopback), None, |_| {}).unwrap(); |
| 156 | until(&ada, |peers| peers.len() == 1); |
| 157 | until(&grace, |peers| peers.len() == 1); |
| 158 | } |
| 159 | |
| 160 | /// Two sealers that met over loopback: Ada's to send with and Grace's to receive with. |
| 161 | fn sealers() -> (Sealer, Sealer) { |
| 162 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); |
| 163 | let address = listener.local_addr().unwrap(); |
| 164 | let ada = thread::spawn(move || { |
| 165 | let mut stream = TcpStream::connect(address).unwrap(); |
| 166 | wire::open(&mut stream, Side::Initiator, "room", b"secret") |
| 167 | .unwrap() |
| 168 | .0 |
| 169 | }); |
| 170 | let (mut stream, _) = listener.accept().unwrap(); |
| 171 | let grace = wire::open(&mut stream, Side::Responder, "room", b"secret") |
| 172 | .unwrap() |
| 173 | .1; |
| 174 | (ada.join().unwrap(), grace) |
| 175 | } |
| 176 | |
| 177 | /// Ada's carets at offsets `0..count`, each frame as its own block of bytes. |
| 178 | fn frames(ada: &mut Sealer, count: u32) -> Vec<Vec<u8>> { |
| 179 | (0..count) |
| 180 | .map(|offset| { |
| 181 | let mut block = Vec::new(); |
| 182 | ada.send(&mut block, kind::PRESENCE, &caret(offset)) |
| 183 | .unwrap(); |
| 184 | block |
| 185 | }) |
| 186 | .collect() |
| 187 | } |
| 188 | |
| 189 | /// A frame lost, repeated, reordered, altered, forged or sealed for another meeting is caught |
| 190 | /// where it lands, before anything in it or after it is read. |
| 191 | #[test] |
| 192 | fn frames_out_of_place_are_caught() { |
| 193 | let (mut other, _) = sealers(); |
| 194 | let elsewhere = frames(&mut other, 3).remove(2); |
| 195 | type Edit = dyn Fn(&mut Vec<Vec<u8>>); |
| 196 | let cases: [(&str, Box<Edit>, &str); 6] = [ |
| 197 | ( |
| 198 | "lost", |
| 199 | Box::new(|sent| drop(sent.remove(2))), |
| 200 | "Frame 3 came where frame 2 was due", |
| 201 | ), |
| 202 | ( |
| 203 | "repeated", |
| 204 | Box::new(|sent| { |
| 205 | let again = sent[1].clone(); |
| 206 | sent.insert(2, again); |
| 207 | }), |
| 208 | "Frame 1 came where frame 2 was due", |
| 209 | ), |
| 210 | ( |
| 211 | "reordered", |
| 212 | Box::new(|sent| sent.swap(2, 3)), |
| 213 | "Frame 3 came where frame 2 was due", |
| 214 | ), |
| 215 | ( |
| 216 | "altered", |
| 217 | Box::new(|sent| *sent[2].last_mut().unwrap() ^= 1), |
| 218 | "Frame 2 does not open", |
| 219 | ), |
| 220 | ( |
| 221 | "forged", |
| 222 | Box::new(|sent| { |
| 223 | // Its length and number kept, the rest made up. |
| 224 | let mut forged = sent[2].clone(); |
| 225 | forged[12..].fill(7); |
| 226 | sent.insert(2, forged); |
| 227 | }), |
| 228 | "Frame 2 does not open", |
| 229 | ), |
| 230 | ( |
| 231 | "from another meeting", |
| 232 | Box::new(move |sent| sent.insert(2, elsewhere.clone())), |
| 233 | "Frame 2 does not open", |
| 234 | ), |
| 235 | ]; |
| 236 | for (name, tamper, expected) in cases { |
| 237 | let (mut ada, mut grace) = sealers(); |
| 238 | let mut sent = frames(&mut ada, 5); |
| 239 | tamper(&mut sent); |
| 240 | let stream = sent.concat(); |
| 241 | let mut arriving = &stream[..]; |
| 242 | for offset in 0..2 { |
| 243 | let (message, body) = grace.receive(&mut arriving).unwrap(); |
| 244 | assert_eq!(message, kind::PRESENCE); |
| 245 | assert_eq!(minicbor::decode::<Presence>(&body).unwrap(), caret(offset)); |
| 246 | } |
| 247 | let error = grace.receive(&mut arriving).unwrap_err(); |
| 248 | assert_eq!(error.kind(), io::ErrorKind::InvalidData, "{name}"); |
| 249 | assert!(error.to_string().starts_with(expected), "{name}: {error}"); |
| 250 | } |
| 251 | } |
| 252 | |
| 253 | /// A relay on this computer, with `config`'s limits: its URL. |
| 254 | fn relay(config: ::relay::server::Config) -> (String, SocketAddr) { |
| 255 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); |
| 256 | let address = listener.local_addr().unwrap(); |
| 257 | thread::spawn(move || ::relay::server::serve(listener, config)); |
| 258 | (format!("ws://{address}"), address) |
| 259 | } |
| 260 | |
| 261 | /// The status a relay at `address` answers a WebSocket to `path` with. |
| 262 | fn status(address: SocketAddr, path: &str) -> String { |
| 263 | let mut stream = TcpStream::connect(address).unwrap(); |
| 264 | write!( |
| 265 | stream, |
| 266 | "GET {path} HTTP/1.1\r\nHost: relay\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n\ |
| 267 | Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==\r\nSec-WebSocket-Version: 13\r\n\r\n" |
| 268 | ) |
| 269 | .unwrap(); |
| 270 | let head = ::relay::ws::head(&mut stream).unwrap(); |
| 271 | head.lines().next().unwrap().to_owned() |
| 272 | } |
| 273 | |
| 274 | /// Two ends of one notebook that share no network meet in the relay's room, and see each |
| 275 | /// other leave. |
| 276 | #[test] |
| 277 | fn peers_meet_through_a_relay() { |
| 278 | let (url, _) = relay(Default::default()); |
| 279 | let room = Room::Notebook([3; 16]); |
| 280 | let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap(); |
| 281 | ada.set_presence(caret(1)); |
| 282 | let grace = Live::start(hello("Grace"), &room, None, Some(&url), |_| {}).unwrap(); |
| 283 | until(&grace, |peers| { |
| 284 | peers.len() == 1 && peers[0].presence == Some(caret(1)) |
| 285 | }); |
| 286 | until(&ada, |peers| { |
| 287 | peers.len() == 1 && peers[0].hello.name == "Grace" |
| 288 | }); |
| 289 | ada.set_presence(caret(2)); |
| 290 | until(&grace, |peers| peers[0].presence == Some(caret(2))); |
| 291 | drop(ada); |
| 292 | until(&grace, <[Peer]>::is_empty); |
| 293 | } |
| 294 | |
| 295 | #[test] |
| 296 | fn relay_timeouts_reopen_with_fresh_keys() { |
| 297 | use ::relay::ws::{self, Message}; |
| 298 | |
| 299 | let (url, upstream) = relay(Default::default()); |
| 300 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); |
| 301 | let delayed = format!("ws://{}", listener.local_addr().unwrap()); |
| 302 | listener.set_nonblocking(true).unwrap(); |
| 303 | let (opening, openings) = mpsc::channel(); |
| 304 | let proxy = thread::spawn(move || { |
| 305 | for attempt in 0..3 { |
| 306 | let deadline = Instant::now() + Duration::from_secs(15); |
| 307 | let mut client = loop { |
| 308 | match listener.accept() { |
| 309 | Ok((client, _)) => break client, |
| 310 | Err(error) if error.kind() == io::ErrorKind::WouldBlock => { |
| 311 | assert!(Instant::now() < deadline, "the relay room never reopened"); |
| 312 | thread::sleep(Duration::from_millis(20)); |
| 313 | } |
| 314 | Err(error) => panic!("{error}"), |
| 315 | } |
| 316 | }; |
| 317 | client.set_nonblocking(false).unwrap(); |
| 318 | client |
| 319 | .set_read_timeout(Some(Duration::from_secs(15))) |
| 320 | .unwrap(); |
| 321 | let mut server = TcpStream::connect(upstream).unwrap(); |
| 322 | let head = ws::head(&mut client).unwrap(); |
| 323 | server.write_all(head.as_bytes()).unwrap(); |
| 324 | let (mut down, mut to_client) = |
| 325 | (server.try_clone().unwrap(), client.try_clone().unwrap()); |
| 326 | let replies = thread::spawn(move || { |
| 327 | let _ = io::copy(&mut down, &mut to_client); |
| 328 | let _ = to_client.shutdown(Shutdown::Both); |
| 329 | }); |
| 330 | let mut messages = ws::Reader::new(client, 1 << 20, true); |
| 331 | let mut first = true; |
| 332 | while let Ok(message) = messages.read() { |
| 333 | let (opcode, data) = match message { |
| 334 | Message::Binary(data) => { |
| 335 | let slot = u32::from_be_bytes(data[..4].try_into().unwrap()); |
| 336 | if first && slot > 0 && slot < ::relay::GROUP { |
| 337 | first = false; |
| 338 | opening.send(data[4..].to_vec()).unwrap(); |
| 339 | // Keep the first opening delayed until its stream is abandoned. |
| 340 | if attempt == 0 { |
| 341 | continue; |
| 342 | } |
| 343 | } |
| 344 | (ws::BINARY, data) |
| 345 | } |
| 346 | Message::Text(text) => (ws::TEXT, text.into_bytes()), |
| 347 | Message::Ping(data) => (ws::PING, data), |
| 348 | Message::Pong => (ws::PONG, vec![]), |
| 349 | Message::Close => break, |
| 350 | }; |
| 351 | if server |
| 352 | .write_all(&ws::frame(opcode, &data, Some([1, 2, 3, 4]))) |
| 353 | .is_err() |
| 354 | { |
| 355 | break; |
| 356 | } |
| 357 | } |
| 358 | let _ = server.shutdown(Shutdown::Both); |
| 359 | replies.join().unwrap(); |
| 360 | } |
| 361 | }); |
| 362 | |
| 363 | let room = Room::Notebook([7; 16]); |
| 364 | let mut host = hello("Host"); |
| 365 | host.serves = Some([8; 16]); |
| 366 | let (met, host_met) = mpsc::channel(); |
| 367 | let (frame, frames) = mpsc::channel(); |
| 368 | let host = Live::start(host, &room, None, Some(&url), move |event| match event { |
| 369 | Event::Met(_, line) => { |
| 370 | let _ = met.send(line.clone()); |
| 371 | } |
| 372 | Event::Frame { kind, body, .. } => { |
| 373 | let _ = frame.send((kind, body.to_vec())); |
| 374 | } |
| 375 | _ => {} |
| 376 | }) |
| 377 | .unwrap(); |
| 378 | let (met, guest_met) = mpsc::channel(); |
| 379 | let (left, guest_left) = mpsc::channel(); |
| 380 | let (frame, guest_frames) = mpsc::channel(); |
| 381 | let guest = Live::start( |
| 382 | hello("Guest"), |
| 383 | &room, |
| 384 | None, |
| 385 | Some(&delayed), |
| 386 | move |event| match event { |
| 387 | Event::Met(_, line) => { |
| 388 | let _ = met.send(line.clone()); |
| 389 | } |
| 390 | Event::Left(_) => { |
| 391 | let _ = left.send(()); |
| 392 | } |
| 393 | Event::Frame { kind, body, .. } => { |
| 394 | let _ = frame.send((kind, body.to_vec())); |
| 395 | } |
| 396 | _ => {} |
| 397 | }, |
| 398 | ) |
| 399 | .unwrap(); |
| 400 | let first = openings.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 401 | assert!(matches!( |
| 402 | host_met.try_recv(), |
| 403 | Err(mpsc::TryRecvError::Empty) |
| 404 | )); |
| 405 | assert!(matches!( |
| 406 | guest_met.try_recv(), |
| 407 | Err(mpsc::TryRecvError::Empty) |
| 408 | )); |
| 409 | let fresh = openings.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 410 | assert_ne!(first, fresh, "the retry reused the abandoned key exchange"); |
| 411 | let host_line = host_met.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 412 | guest_met.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 413 | let before_timeout = guest.shared.state.lock().unwrap().relay.clone().unwrap(); |
| 414 | guest |
| 415 | .shared |
| 416 | .state |
| 417 | .lock() |
| 418 | .unwrap() |
| 419 | .peers |
| 420 | .values() |
| 421 | .next() |
| 422 | .unwrap() |
| 423 | .pipe |
| 424 | .set_read_timeout(Duration::from_secs(1)) |
| 425 | .unwrap(); |
| 426 | host_line.send(kind::EDITS, &40_u64).unwrap(); |
| 427 | let (kind, body) = guest_frames.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 428 | assert_eq!(kind, kind::EDITS); |
| 429 | assert_eq!(minicbor::decode::<u64>(&body).unwrap(), 40); |
| 430 | guest_left.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 431 | let reopened = openings.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 432 | assert_ne!( |
| 433 | fresh, reopened, |
| 434 | "the retry reused the timed-out stream's keys" |
| 435 | ); |
| 436 | host_met.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 437 | let line = guest_met.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 438 | assert_eq!((host.failed(), guest.failed()), (0, 0)); |
| 439 | let before = guest.shared.state.lock().unwrap().relay.clone().unwrap(); |
| 440 | assert!(!Arc::ptr_eq(&before_timeout, &before)); |
| 441 | for edit in [41_u64, 42] { |
| 442 | line.send(kind::EDITS, &edit).unwrap(); |
| 443 | } |
| 444 | line.hang_up("left"); |
| 445 | let mut delivered = Vec::new(); |
| 446 | loop { |
| 447 | let (kind, body) = frames.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 448 | match kind { |
| 449 | kind::EDITS => delivered.push(minicbor::decode::<u64>(&body).unwrap()), |
| 450 | kind::BYE => break, |
| 451 | _ => {} |
| 452 | } |
| 453 | } |
| 454 | assert_eq!(delivered, [41, 42]); |
| 455 | guest_left.recv_timeout(Duration::from_secs(10)).unwrap(); |
| 456 | host.set_presence(caret(43)); |
| 457 | until(&guest, |peers| { |
| 458 | peers |
| 459 | .first() |
| 460 | .is_some_and(|peer| peer.presence == Some(caret(43))) |
| 461 | }); |
| 462 | let state = guest.shared.state.lock().unwrap(); |
| 463 | assert!(Arc::ptr_eq(&before, state.relay.as_ref().unwrap())); |
| 464 | drop(state); |
| 465 | assert!(matches!( |
| 466 | host_met.try_recv(), |
| 467 | Err(mpsc::TryRecvError::Empty) |
| 468 | )); |
| 469 | assert!(matches!( |
| 470 | guest_met.try_recv(), |
| 471 | Err(mpsc::TryRecvError::Empty) |
| 472 | )); |
| 473 | assert!(matches!( |
| 474 | guest_frames.try_recv(), |
| 475 | Err(mpsc::TryRecvError::Empty) |
| 476 | )); |
| 477 | drop(guest); |
| 478 | drop(host); |
| 479 | proxy.join().unwrap(); |
| 480 | } |
| 481 | |
| 482 | /// The end sharing a code asks the relay to number it; the other types the whole code. One |
| 483 | /// with the wrong words never meets, and the relay, told so by the end sharing, burns the |
| 484 | /// code once too many have tried. |
| 485 | #[test] |
| 486 | fn a_relay_numbers_a_code_and_burns_it_after_wrong_tries() { |
| 487 | let (url, address) = relay(::relay::server::Config { |
| 488 | burn_after: 2, |
| 489 | ..Default::default() |
| 490 | }); |
| 491 | let host = Live::start( |
| 492 | hello("Ada"), |
| 493 | &Room::share("ABCDEF", ""), |
| 494 | None, |
| 495 | Some(&url), |
| 496 | |_| {}, |
| 497 | ) |
| 498 | .unwrap(); |
| 499 | let deadline = Instant::now() + Duration::from_secs(10); |
| 500 | let code = loop { |
| 501 | if let Some(code) = host.code() { |
| 502 | break code; |
| 503 | } |
| 504 | assert!(Instant::now() < deadline, "no code"); |
| 505 | thread::sleep(Duration::from_millis(20)); |
| 506 | }; |
| 507 | let (number, secret) = super::code::parse(&code).unwrap(); |
| 508 | assert_eq!(secret, "ABCDEF"); |
| 509 | let guest = Live::start( |
| 510 | hello("Grace"), |
| 511 | &Room::join(&code, ""), |
| 512 | None, |
| 513 | Some(&url), |
| 514 | |_| {}, |
| 515 | ) |
| 516 | .unwrap(); |
| 517 | until(&host, |peers| peers.len() == 1); |
| 518 | until(&guest, |peers| peers.len() == 1); |
| 519 | |
| 520 | let path = format!("/v1/room/code-{number}"); |
| 521 | let wrong = Room::join(&self::code(number, "ABCDEG"), ""); |
| 522 | // An end that typed a wrong code gives up at once; Mallory tries twice, and the second |
| 523 | // wrong try burns the code. |
| 524 | for _ in 0..2 { |
| 525 | let mallory = Live::start(hello("Mallory"), &wrong, None, Some(&url), |_| {}).unwrap(); |
| 526 | let deadline = Instant::now() + Duration::from_secs(10); |
| 527 | while mallory.failed() == 0 { |
| 528 | assert!(Instant::now() < deadline, "Mallory never tried"); |
| 529 | thread::sleep(Duration::from_millis(20)); |
| 530 | } |
| 531 | assert!(mallory.peers().is_empty()); |
| 532 | } |
| 533 | until(&host, |_| host.burned()); |
| 534 | assert_eq!(status(address, &path), "HTTP/1.1 410 Gone"); |
| 535 | assert_eq!(host.peers().len(), 1, "Grace stays"); |
| 536 | } |
| 537 | |
| 538 | /// What a malicious relay does to one message on its way. |
| 539 | #[derive(Clone, Copy, Debug)] |
| 540 | enum Tamper { |
| 541 | Drop, |
| 542 | Repeat, |
| 543 | Reorder, |
| 544 | Alter, |
| 545 | Inject, |
| 546 | } |
| 547 | |
| 548 | /// A relay in the middle of Grace's connection that passes on what the real one at |
| 549 | /// `upstream` says, except the first message to her from a peer once `armed`, which it |
| 550 | /// tampers with. Her next connection waits for `release`. |
| 551 | struct Malicious { |
| 552 | url: String, |
| 553 | armed: Arc<std::sync::atomic::AtomicBool>, |
| 554 | tampered: mpsc::Receiver<()>, |
| 555 | rejoined: mpsc::Receiver<()>, |
| 556 | release: mpsc::Sender<()>, |
| 557 | } |
| 558 | |
| 559 | fn malicious(upstream: SocketAddr, tamper: Tamper) -> Malicious { |
| 560 | use ::relay::ws::{self, Message}; |
| 561 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); |
| 562 | let url = format!("ws://{}", listener.local_addr().unwrap()); |
| 563 | let (tampered, told) = mpsc::channel(); |
| 564 | let (rejoined, heard) = mpsc::channel(); |
| 565 | let (release, released) = mpsc::channel(); |
| 566 | let armed = Arc::new(std::sync::atomic::AtomicBool::new(false)); |
| 567 | let arming = Arc::clone(&armed); |
| 568 | thread::spawn(move || { |
| 569 | for (index, client) in listener.incoming().enumerate() { |
| 570 | let mut client = client.unwrap(); |
| 571 | if index > 0 { |
| 572 | let _ = rejoined.send(()); |
| 573 | let _ = released.recv(); |
| 574 | } |
| 575 | let server = TcpStream::connect(upstream).unwrap(); |
| 576 | let (mut up, mut to_server) = |
| 577 | (client.try_clone().unwrap(), server.try_clone().unwrap()); |
| 578 | thread::spawn(move || { |
| 579 | let _ = io::copy(&mut up, &mut to_server); |
| 580 | let _ = to_server.shutdown(Shutdown::Both); |
| 581 | }); |
| 582 | let (tampered, armed) = (tampered.clone(), Arc::clone(&arming)); |
| 583 | thread::spawn(move || { |
| 584 | let mut reading = io::BufReader::new(server); |
| 585 | let head = ws::head(&mut reading).unwrap(); |
| 586 | client.write_all(head.as_bytes()).unwrap(); |
| 587 | let mut messages = ws::Reader::new(reading, 1 << 20, false); |
| 588 | let mut held = None; |
| 589 | while let Ok(message) = messages.read() { |
| 590 | let frames: Vec<Vec<u8>> = match message { |
| 591 | Message::Binary(mut data) => { |
| 592 | let mut out = vec![]; |
| 593 | if index == 0 && armed.swap(false, std::sync::atomic::Ordering::AcqRel) |
| 594 | { |
| 595 | match tamper { |
| 596 | Tamper::Drop => {} |
| 597 | Tamper::Repeat => out = vec![data.clone(), data], |
| 598 | Tamper::Reorder => held = Some(data), |
| 599 | Tamper::Alter => { |
| 600 | *data.last_mut().unwrap() ^= 1; |
| 601 | out = vec![data]; |
| 602 | } |
| 603 | Tamper::Inject => { |
| 604 | // A slot and part of the sender's id, then |
| 605 | // made-up bytes. |
| 606 | let mut forged = data.clone(); |
| 607 | forged[16..].fill(7); |
| 608 | out = vec![forged, data]; |
| 609 | } |
| 610 | } |
| 611 | let _ = tampered.send(()); |
| 612 | } else { |
| 613 | out.push(data); |
| 614 | out.extend(held.take()); |
| 615 | } |
| 616 | out.into_iter() |
| 617 | .map(|data| ws::frame(ws::BINARY, &data, None)) |
| 618 | .collect() |
| 619 | } |
| 620 | Message::Text(text) => vec![ws::frame(ws::TEXT, text.as_bytes(), None)], |
| 621 | Message::Ping(payload) => vec![ws::frame(ws::PING, &payload, None)], |
| 622 | Message::Pong => vec![ws::frame(ws::PONG, &[], None)], |
| 623 | Message::Close => break, |
| 624 | }; |
| 625 | if frames.iter().any(|frame| client.write_all(frame).is_err()) { |
| 626 | break; |
| 627 | } |
| 628 | } |
| 629 | let _ = client.shutdown(Shutdown::Both); |
| 630 | }); |
| 631 | } |
| 632 | }); |
| 633 | Malicious { |
| 634 | url, |
| 635 | armed, |
| 636 | tampered: told, |
| 637 | rejoined: heard, |
| 638 | release, |
| 639 | } |
| 640 | } |
| 641 | |
| 642 | /// Starts `name`, recording each presence of its peer as applied, and `None` when the peer |
| 643 | /// goes. |
| 644 | fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc<Mutex<Vec<Option<Presence>>>>) { |
| 645 | let heard = Arc::new(Mutex::new(Vec::new())); |
| 646 | let shared: Arc<std::sync::OnceLock<std::sync::Weak<Shared>>> = Arc::default(); |
| 647 | let (recorded, watched) = (Arc::clone(&heard), Arc::clone(&shared)); |
| 648 | let live = Live::start(hello(name), room, None, Some(relay), move |event| { |
| 649 | if !matches!(event, Event::Changed) { |
| 650 | return; |
| 651 | } |
| 652 | let Some(shared) = watched.get().and_then(std::sync::Weak::upgrade) else { |
| 653 | return; |
| 654 | }; |
| 655 | let entry = match shared.peers().first() { |
| 656 | None => Some(None), |
| 657 | Some(peer) => peer.presence.clone().map(Some), |
| 658 | }; |
| 659 | recorded.lock().unwrap().extend(entry); |
| 660 | }) |
| 661 | .unwrap(); |
| 662 | shared.set(Arc::downgrade(&live.shared)).ok().unwrap(); |
| 663 | (live, heard) |
| 664 | } |
| 665 | |
| 666 | /// Whatever a malicious relay does to a frame, the end it was for never applies it or |
| 667 | /// anything after it, drops the connection as broken, joins again and hears the newest |
| 668 | /// presence from scratch. |
| 669 | #[test] |
| 670 | fn a_malicious_relay_is_caught() { |
| 671 | let room = Room::Notebook([4; 16]); |
| 672 | for tamper in [ |
| 673 | Tamper::Drop, |
| 674 | Tamper::Repeat, |
| 675 | Tamper::Reorder, |
| 676 | Tamper::Alter, |
| 677 | Tamper::Inject, |
| 678 | ] { |
| 679 | let (url, address) = relay(Default::default()); |
| 680 | let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap(); |
| 681 | ada.set_presence(caret(1)); |
| 682 | let relay = malicious(address, tamper); |
| 683 | let (grace, heard) = recording("Grace", &room, &relay.url); |
| 684 | until(&grace, |peers| { |
| 685 | peers |
| 686 | .first() |
| 687 | .is_some_and(|peer| peer.presence == Some(caret(1))) |
| 688 | }); |
| 689 | // Ada's next presence goes to everyone, so its loss is caught too. |
| 690 | relay |
| 691 | .armed |
| 692 | .store(true, std::sync::atomic::Ordering::Release); |
| 693 | ada.set_presence(caret(2)); |
| 694 | relay |
| 695 | .tampered |
| 696 | .recv_timeout(Duration::from_secs(10)) |
| 697 | .unwrap(); |
| 698 | ada.set_presence(caret(3)); |
| 699 | relay |
| 700 | .rejoined |
| 701 | .recv_timeout(Duration::from_secs(10)) |
| 702 | .unwrap(); |
| 703 | let before: Vec<_> = heard.lock().unwrap().clone(); |
| 704 | let gone = before.iter().position(Option::is_none).unwrap(); |
| 705 | assert!( |
| 706 | !before[..gone].contains(&Some(caret(3))), |
| 707 | "{tamper:?}: applied after the tampering: {before:?}" |
| 708 | ); |
| 709 | assert!(grace.peers().is_empty(), "{tamper:?}"); |
| 710 | relay.release.send(()).unwrap(); |
| 711 | until(&grace, |peers| { |
| 712 | peers |
| 713 | .first() |
| 714 | .is_some_and(|peer| peer.presence == Some(caret(3))) |
| 715 | }); |
| 716 | } |
| 717 | } |
| 718 | |
| 719 | /// Ends of two versions each learn the other's version from the opening, before any key is |
| 720 | /// agreed, so each can say which should update. |
| 721 | #[test] |
| 722 | fn another_version_is_named_before_any_key() { |
| 723 | let listener = TcpListener::bind("127.0.0.1:0").unwrap(); |
| 724 | let address = listener.local_addr().unwrap(); |
| 725 | let older = thread::spawn(move || { |
| 726 | let mut stream = TcpStream::connect(address).unwrap(); |
| 727 | #[derive(minicbor::Encode)] |
| 728 | #[cbor(map)] |
| 729 | struct Open { |
| 730 | #[n(0)] |
| 731 | version: u16, |
| 732 | #[n(1)] |
| 733 | room: String, |
| 734 | #[cbor(n(2), with = "minicbor::bytes")] |
| 735 | pake: Vec<u8>, |
| 736 | } |
| 737 | let open = minicbor::to_vec(Open { |
| 738 | version: 1, |
| 739 | room: "room".into(), |
| 740 | pake: vec![0; 33], |
| 741 | }) |
| 742 | .unwrap(); |
| 743 | stream |
| 744 | .write_all(&[&(open.len() as u32).to_be_bytes()[..], &open].concat()) |
| 745 | .unwrap(); |
| 746 | let mut length = [0; 4]; |
| 747 | stream.read_exact(&mut length).unwrap(); |
| 748 | let mut answer = vec![0; u32::from_be_bytes(length) as usize]; |
| 749 | stream.read_exact(&mut answer).unwrap(); |
| 750 | answer |
| 751 | }); |
| 752 | let (mut stream, _) = listener.accept().unwrap(); |
| 753 | let Err(error) = wire::open(&mut stream, Side::Responder, "room", b"secret") else { |
| 754 | panic!("met a peer of another version"); |
| 755 | }; |
| 756 | assert_eq!(error.kind(), io::ErrorKind::Unsupported); |
| 757 | let version = error.get_ref().unwrap().downcast_ref::<wire::Version>(); |
| 758 | assert_eq!(version, Some(&wire::Version(1))); |
| 759 | // The older end heard this one's opening, and with it its version. |
| 760 | assert!(!older.join().unwrap().is_empty()); |
| 761 | } |