1use super::*;
2use std::time::Instant;
3
4/// The code for room `number` and `secret`.
5fn code(number: u32, secret: &str) -> String {
6 super::code::format(number, secret).unwrap()
7}
8
9fn 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.
14fn 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
26fn 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]
50fn 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]
75fn 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]
99fn 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"]
152fn 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.
161fn 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.
178fn 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]
192fn 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.
254fn 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.
262fn 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]
277fn 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/// The end sharing a code asks the relay to number it; the other types the whole code. One
296/// with the wrong words never meets, and the relay, told so by the end sharing, burns the
297/// code once too many have tried.
298#[test]
299fn a_relay_numbers_a_code_and_burns_it_after_wrong_tries() {
300 let (url, address) = relay(::relay::server::Config {
301 burn_after: 2,
302 ..Default::default()
303 });
304 let host = Live::start(
305 hello("Ada"),
306 &Room::share("ABCDEF", ""),
307 None,
308 Some(&url),
309 |_| {},
310 )
311 .unwrap();
312 let deadline = Instant::now() + Duration::from_secs(10);
313 let code = loop {
314 if let Some(code) = host.code() {
315 break code;
316 }
317 assert!(Instant::now() < deadline, "no code");
318 thread::sleep(Duration::from_millis(20));
319 };
320 let (number, secret) = super::code::parse(&code).unwrap();
321 assert_eq!(secret, "ABCDEF");
322 let guest = Live::start(
323 hello("Grace"),
324 &Room::join(&code, ""),
325 None,
326 Some(&url),
327 |_| {},
328 )
329 .unwrap();
330 until(&host, |peers| peers.len() == 1);
331 until(&guest, |peers| peers.len() == 1);
332
333 let path = format!("/v1/room/code-{number}");
334 let wrong = Room::join(&self::code(number, "ABCDEG"), "");
335 // An end that typed a wrong code gives up at once; Mallory tries twice, and the second
336 // wrong try burns the code.
337 for _ in 0..2 {
338 let mallory = Live::start(hello("Mallory"), &wrong, None, Some(&url), |_| {}).unwrap();
339 let deadline = Instant::now() + Duration::from_secs(10);
340 while mallory.failed() == 0 {
341 assert!(Instant::now() < deadline, "Mallory never tried");
342 thread::sleep(Duration::from_millis(20));
343 }
344 assert!(mallory.peers().is_empty());
345 }
346 until(&host, |_| host.burned());
347 assert_eq!(status(address, &path), "HTTP/1.1 410 Gone");
348 assert_eq!(host.peers().len(), 1, "Grace stays");
349}
350
351/// What a malicious relay does to one message on its way.
352#[derive(Clone, Copy, Debug)]
353enum Tamper {
354 Drop,
355 Repeat,
356 Reorder,
357 Alter,
358 Inject,
359}
360
361/// A relay in the middle of Grace's connection that passes on what the real one at
362/// `upstream` says, except the first message to her from a peer once `armed`, which it
363/// tampers with. Her next connection waits for `release`.
364struct Malicious {
365 url: String,
366 armed: Arc<std::sync::atomic::AtomicBool>,
367 tampered: mpsc::Receiver<()>,
368 rejoined: mpsc::Receiver<()>,
369 release: mpsc::Sender<()>,
370}
371
372fn malicious(upstream: SocketAddr, tamper: Tamper) -> Malicious {
373 use ::relay::ws::{self, Message};
374 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
375 let url = format!("ws://{}", listener.local_addr().unwrap());
376 let (tampered, told) = mpsc::channel();
377 let (rejoined, heard) = mpsc::channel();
378 let (release, released) = mpsc::channel();
379 let armed = Arc::new(std::sync::atomic::AtomicBool::new(false));
380 let arming = Arc::clone(&armed);
381 thread::spawn(move || {
382 for (index, client) in listener.incoming().enumerate() {
383 let mut client = client.unwrap();
384 if index > 0 {
385 let _ = rejoined.send(());
386 let _ = released.recv();
387 }
388 let server = TcpStream::connect(upstream).unwrap();
389 let (mut up, mut to_server) =
390 (client.try_clone().unwrap(), server.try_clone().unwrap());
391 thread::spawn(move || {
392 let _ = io::copy(&mut up, &mut to_server);
393 let _ = to_server.shutdown(Shutdown::Both);
394 });
395 let (tampered, armed) = (tampered.clone(), Arc::clone(&arming));
396 thread::spawn(move || {
397 let mut reading = io::BufReader::new(server);
398 let head = ws::head(&mut reading).unwrap();
399 client.write_all(head.as_bytes()).unwrap();
400 let mut messages = ws::Reader::new(reading, 1 << 20, false);
401 let mut held = None;
402 while let Ok(message) = messages.read() {
403 let frames: Vec<Vec<u8>> = match message {
404 Message::Binary(mut data) => {
405 let mut out = vec![];
406 if index == 0 && armed.swap(false, std::sync::atomic::Ordering::AcqRel)
407 {
408 match tamper {
409 Tamper::Drop => {}
410 Tamper::Repeat => out = vec![data.clone(), data],
411 Tamper::Reorder => held = Some(data),
412 Tamper::Alter => {
413 *data.last_mut().unwrap() ^= 1;
414 out = vec![data];
415 }
416 Tamper::Inject => {
417 // A slot and part of the sender's id, then
418 // made-up bytes.
419 let mut forged = data.clone();
420 forged[16..].fill(7);
421 out = vec![forged, data];
422 }
423 }
424 let _ = tampered.send(());
425 } else {
426 out.push(data);
427 out.extend(held.take());
428 }
429 out.into_iter()
430 .map(|data| ws::frame(ws::BINARY, &data, None))
431 .collect()
432 }
433 Message::Text(text) => vec![ws::frame(ws::TEXT, text.as_bytes(), None)],
434 Message::Ping(payload) => vec![ws::frame(ws::PING, &payload, None)],
435 Message::Pong => vec![ws::frame(ws::PONG, &[], None)],
436 Message::Close => break,
437 };
438 if frames.iter().any(|frame| client.write_all(frame).is_err()) {
439 break;
440 }
441 }
442 let _ = client.shutdown(Shutdown::Both);
443 });
444 }
445 });
446 Malicious {
447 url,
448 armed,
449 tampered: told,
450 rejoined: heard,
451 release,
452 }
453}
454
455/// Starts `name`, recording each presence of its peer as applied, and `None` when the peer
456/// goes.
457fn recording(name: &str, room: &Room, relay: &str) -> (Live, Arc<Mutex<Vec<Option<Presence>>>>) {
458 let heard = Arc::new(Mutex::new(Vec::new()));
459 let shared: Arc<std::sync::OnceLock<std::sync::Weak<Shared>>> = Arc::default();
460 let (recorded, watched) = (Arc::clone(&heard), Arc::clone(&shared));
461 let live = Live::start(hello(name), room, None, Some(relay), move |event| {
462 if !matches!(event, Event::Changed) {
463 return;
464 }
465 let Some(shared) = watched.get().and_then(std::sync::Weak::upgrade) else {
466 return;
467 };
468 let entry = match shared.peers().first() {
469 None => Some(None),
470 Some(peer) => peer.presence.clone().map(Some),
471 };
472 recorded.lock().unwrap().extend(entry);
473 })
474 .unwrap();
475 shared.set(Arc::downgrade(&live.shared)).ok().unwrap();
476 (live, heard)
477}
478
479/// Whatever a malicious relay does to a frame, the end it was for never applies it or
480/// anything after it, drops the connection as broken, joins again and hears the newest
481/// presence from scratch.
482#[test]
483fn a_malicious_relay_is_caught() {
484 let room = Room::Notebook([4; 16]);
485 for tamper in [
486 Tamper::Drop,
487 Tamper::Repeat,
488 Tamper::Reorder,
489 Tamper::Alter,
490 Tamper::Inject,
491 ] {
492 let (url, address) = relay(Default::default());
493 let ada = Live::start(hello("Ada"), &room, None, Some(&url), |_| {}).unwrap();
494 ada.set_presence(caret(1));
495 let relay = malicious(address, tamper);
496 let (grace, heard) = recording("Grace", &room, &relay.url);
497 until(&grace, |peers| {
498 peers
499 .first()
500 .is_some_and(|peer| peer.presence == Some(caret(1)))
501 });
502 // Ada's next presence goes to everyone, so its loss is caught too.
503 relay
504 .armed
505 .store(true, std::sync::atomic::Ordering::Release);
506 ada.set_presence(caret(2));
507 relay
508 .tampered
509 .recv_timeout(Duration::from_secs(10))
510 .unwrap();
511 ada.set_presence(caret(3));
512 relay
513 .rejoined
514 .recv_timeout(Duration::from_secs(10))
515 .unwrap();
516 let before: Vec<_> = heard.lock().unwrap().clone();
517 let gone = before.iter().position(Option::is_none).unwrap();
518 assert!(
519 !before[..gone].contains(&Some(caret(3))),
520 "{tamper:?}: applied after the tampering: {before:?}"
521 );
522 assert!(grace.peers().is_empty(), "{tamper:?}");
523 relay.release.send(()).unwrap();
524 until(&grace, |peers| {
525 peers
526 .first()
527 .is_some_and(|peer| peer.presence == Some(caret(3)))
528 });
529 }
530}
531
532/// Ends of two versions each learn the other's version from the opening, before any key is
533/// agreed, so each can say which should update.
534#[test]
535fn another_version_is_named_before_any_key() {
536 let listener = TcpListener::bind("127.0.0.1:0").unwrap();
537 let address = listener.local_addr().unwrap();
538 let older = thread::spawn(move || {
539 let mut stream = TcpStream::connect(address).unwrap();
540 #[derive(minicbor::Encode)]
541 #[cbor(map)]
542 struct Open {
543 #[n(0)]
544 version: u16,
545 #[n(1)]
546 room: String,
547 #[cbor(n(2), with = "minicbor::bytes")]
548 pake: Vec<u8>,
549 }
550 let open = minicbor::to_vec(Open {
551 version: 1,
552 room: "room".into(),
553 pake: vec![0; 33],
554 })
555 .unwrap();
556 stream
557 .write_all(&[&(open.len() as u32).to_be_bytes()[..], &open].concat())
558 .unwrap();
559 let mut length = [0; 4];
560 stream.read_exact(&mut length).unwrap();
561 let mut answer = vec![0; u32::from_be_bytes(length) as usize];
562 stream.read_exact(&mut answer).unwrap();
563 answer
564 });
565 let (mut stream, _) = listener.accept().unwrap();
566 let Err(error) = wire::open(&mut stream, Side::Responder, "room", b"secret") else {
567 panic!("met a peer of another version");
568 };
569 assert_eq!(error.kind(), io::ErrorKind::Unsupported);
570 let version = error.get_ref().unwrap().downcast_ref::<wire::Version>();
571 assert_eq!(version, Some(&wire::Version(1)));
572 // The older end heard this one's opening, and with it its version.
573 assert!(!older.join().unwrap().is_empty());
574}