authorgravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-05 22:29:47-07:00
committergravatar for git@paperclover.netclover caruso <git@paperclover.net> 2026-10-06 00:27:38-07:00
logfccad14f94e3d22b269957a42ae9c4d3f615edb0
tree7c14efd08ce8d47f6463800eb1229f549167b776
parent2c99e71ccf170b16592314fda8f6de56584acd3a
signature Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU

fix: reopen Live Share after relay stream failures

Retry transient handshake and established-stream I/O failures with fresh keys without counting rejected authentication. Preserve explicit peer departure without retiring the room socket. fixes #95 Assisted-by: gpt-6.1-sol

3 files changed, 203 insertions(+), 7 deletions(-)

crates/notebook/src/live.rs+12-6
......@@ -178,8 +178,7 @@ trait Pipe: Send + Sync {
178178 fn direct(&self) -> bool;
179179 /// Hears whether the peer at the other end knew the secret.
180180 fn met(&self, _met: bool) {}
181 /// Hears that frames arrived lost, repeated, reordered or forged, after which this end
182 /// meets its peers again.
181 /// Hears that the stream failed, after which this end meets its peers again.
183182 fn broken(&self) {}
184183}
185184
......@@ -558,20 +557,26 @@ impl Shared {
558557 let other = (met.as_ref().err())
559558 .and_then(|error| error.get_ref()?.downcast_ref::<wire::Version>())
560559 .copied();
561 // A peer of another version guessed nothing, the keys never being agreed.
562 pipe.met(met.is_ok() || other.is_some());
563560 let (send, mut receive, hello) = match met {
564 Ok(met) => met,
561 Ok(met) => {
562 pipe.met(true);
563 met
564 }
565565 Err(error) => {
566566 eprintln!("Live: no meeting in {tag}: {error}");
567567 if let Some(wire::Version(version)) = other {
568 // Another version guessed nothing: no keys were agreed.
569 pipe.met(true);
568570 self.state.lock().unwrap().outdated = Some(version);
569571 if matches!(self.room, Room::Code { owner: false, .. }) {
570572 self.stopped.store(true, Ordering::Release);
571573 }
572574 (self.events)(Event::Changed);
573575 } else if error.kind() == io::ErrorKind::InvalidData {
576 pipe.met(false);
574577 self.failed();
578 } else {
579 pipe.broken();
575580 }
576581 pipe.shutdown();
577582 return false;
......@@ -651,7 +656,7 @@ impl Shared {
651656 }
652657 }
653658 };
654 if let Some(error) = ended.filter(|error| error.kind() == io::ErrorKind::InvalidData) {
659 if let Some(error) = ended {
655660 eprintln!("Live: the connection to {name} broke ({error}); meeting again");
656661 pipe.broken();
657662 }
......@@ -698,6 +703,7 @@ impl Shared {
698703 None => Ok(()),
699704 });
700705 if result.is_err() {
706 pipe.broken();
701707 pipe.shutdown();
702708 return;
703709 }
crates/notebook/src/live/relay.rs+4-1
......@@ -551,7 +551,10 @@ impl Pipe for Relayed {
551551 }
552552
553553 fn broken(&self) {
554 self.socket.hang_up();
554 // A deliberate shutdown can wake the reader with EOF.
555 if !self.socket.links.lock().unwrap().ended.contains(&self.slot) {
556 self.socket.hang_up();
557 }
555558 }
556559}
557560
crates/notebook/src/live/tests.rs+187
......@@ -292,6 +292,193 @@ fn peers_meet_through_a_relay() {
292292 until(&grace, <[Peer]>::is_empty);
293293}
294294
295#[test]
296fn 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
295482/// The end sharing a code asks the relay to number it; the other types the whole code. One
296483/// with the wrong words never meets, and the relay, told so by the end sharing, burns the
297484/// code once too many have tried.