From 7e428759e27de7e50a06ca5afadfbfe0f76c42f0 Mon Sep 17 00:00:00 2001 From: clover caruso Date: Fri, 2 Oct 2026 19:12:37 -0700 Subject: [PATCH] fix: a Live Share host serves each guest through two workers and hangs up on one that floods it A host started a thread for every request a guest sent, so a guest could make it start as many as it liked. Each guest now has two workers and room for 30 requests waiting; it may start 100 requests a second, 200 at once, besides a read's later chunks and an upload's. A guest past either limit hears Bye "flooded" and is hung up on, and a peer's Bye now ends its connection at once rather than after its silence. Assisted-by: claude-opus-5.5 --- crates/notebook/src/live.rs | 27 +++++++-- crates/notebook/src/live/share.rs | 94 +++++++++++++++++++++++------ crates/notebook/tests/live_share.rs | 77 +++++++++++++++++++++++ 3 files changed, 174 insertions(+), 24 deletions(-) diff --git a/crates/notebook/src/live.rs b/crates/notebook/src/live.rs index 5fdce7106330e72794b9ddb5f5d7427945b45de4..0e6f01143e0eb4e1c55a7ca88c1051311c4f76f0 100644 --- a/crates/notebook/src/live.rs +++ b/crates/notebook/src/live.rs @@ -222,6 +222,15 @@ impl Line { .send(Out::Frame(kind, body)) .map_err(|_| io::ErrorKind::NotConnected.into()) } + + /// Says `reason` after the frames sent before, then hangs up. + pub fn hang_up(&self, reason: &str) { + let bye = minicbor::to_vec(wire::Bye { + reason: reason.into(), + }) + .unwrap_or_default(); + let _ = self.0.send(Out::Bye(bye)); + } } /// A stream to one peer, read by one thread and written by another: a TCP connection, or @@ -597,12 +606,18 @@ impl Shared { } (self.events)(Event::Changed); } - kind => (self.events)(Event::Frame { - from: &hello, - kind, - body: &body, - line: &line, - }), + kind => { + (self.events)(Event::Frame { + from: &hello, + kind, + body: &body, + line: &line, + }); + // The peer hangs up after its bye. + if kind == kind::BYE { + break None; + } + } } }; if let Some(error) = ended.filter(|error| error.kind() == io::ErrorKind::InvalidData) { diff --git a/crates/notebook/src/live/share.rs b/crates/notebook/src/live/share.rs index 87c001fd3fd56d1be675a1200300aba76ed662c6..1c0ca9025c5b18a4e40a5f7eb1b7d3463e6c00de 100644 --- a/crates/notebook/src/live/share.rs +++ b/crates/notebook/src/live/share.rs @@ -44,6 +44,16 @@ const SNAPSHOT_AGE: Duration = Duration::from_secs(120); const IMAGES: usize = 4; /// What a host says leaving as it stops sharing. const STOPPED: &str = "stopped"; +/// What a host says hanging up on a guest that asks too much too fast. +const FLOODED: &str = "flooded"; +/// A guest's requests a host holds at once, waiting and in hand, past which it hangs up. +const QUEUED: usize = 32; +/// The threads working through each guest's requests. +const WORKERS: usize = 2; +/// Requests a guest may start each second, and at once: every request but a read's later +/// chunks and an upload's, which come from memory or go to it. +const STARTS: f64 = 100.0; +const BURST: f64 = 200.0; const WORDS: &str = include_str!("words.txt"); @@ -231,7 +241,7 @@ impl Host { images: Mutex::default(), snapshots: Mutex::default(), puts: Mutex::default(), - lines: Mutex::default(), + guests: Mutex::default(), }); let serving = Hello { serves: Some(sharing.share), @@ -244,22 +254,12 @@ impl Host { reach, relay, move |event| match event { - Event::Met(hello, line) => { - heard.lines.lock().unwrap().insert(hello.peer, line.clone()); - } + Event::Met(hello, line) => heard.admit(hello.peer, line), Event::Left(hello) => heard.forget(&hello.peer), Event::Frame { - from, - kind, - body, - line, + from, kind, body, .. } if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => { - let (served, line, peer) = (Arc::clone(&heard), line.clone(), from.peer); - let body = body.to_vec(); - thread::spawn(move || { - let reply = served.handle(&peer, kind, &body); - let _ = line.send(kind::REPLY, &reply); - }); + heard.queue(&from.peer, kind, body); } Event::Changed => told(), _ => {} @@ -379,7 +379,16 @@ struct Served { snapshots: Mutex>, Instant)>>, /// Bytes a later request carries, by guest and upload. puts: Mutex>>, - lines: Mutex>, + guests: Mutex>, +} + +/// A guest as its host serves it: the line to it, its requests waiting for its workers, and +/// how many more it may start now. +struct Admitted { + line: Line, + queue: mpsc::SyncSender<(u16, Vec)>, + starts: f64, + counted: Instant, } /// A section's path, stamp and image. @@ -406,8 +415,57 @@ fn refused(kind: io::ErrorKind, message: &str) -> Error { } impl Served { + /// Serves `peer` on `line`, through workers of its own that end as it leaves. + fn admit(self: &Arc, peer: [u8; 16], line: &Line) { + let (queue, waiting) = mpsc::sync_channel::<(u16, Vec)>(QUEUED - WORKERS); + let waiting = Arc::new(Mutex::new(waiting)); + for _ in 0..WORKERS { + let (served, waiting, line) = + (Arc::downgrade(self), Arc::clone(&waiting), line.clone()); + thread::spawn(move || { + loop { + let next = waiting.lock().unwrap().recv(); + let (Ok((kind, body)), Some(served)) = (next, served.upgrade()) else { + return; + }; + let _ = line.send(kind::REPLY, &served.handle(&peer, kind, &body)); + } + }); + } + let guest = Admitted { + line: line.clone(), + queue, + starts: BURST, + counted: Instant::now(), + }; + self.guests.lock().unwrap().insert(peer, guest); + } + + /// Hands a request from `peer` to its workers, or hangs up on a guest that has too many + /// waiting or starts them too fast. + fn queue(&self, peer: &[u8; 16], kind: u16, body: &[u8]) { + let mut guests = self.guests.lock().unwrap(); + let Some(guest) = guests.get_mut(peer) else { + return; + }; + let continued = kind == kind::PUT + || matches!(kind, kind::READ | kind::READ_FILE) + && minicbor::decode::(body).is_ok_and(|request| request.handle.is_some()); + let now = Instant::now(); + guest.starts = + (guest.starts + now.duration_since(guest.counted).as_secs_f64() * STARTS).min(BURST); + guest.counted = now; + if !continued { + guest.starts -= 1.0; + } + if guest.starts < 0.0 || guest.queue.try_send((kind, body.to_vec())).is_err() { + guest.line.hang_up(FLOODED); + guests.remove(peer); + } + } + fn forget(&self, peer: &[u8; 16]) { - self.lines.lock().unwrap().remove(peer); + self.guests.lock().unwrap().remove(peer); self.snapshots .lock() .unwrap() @@ -423,8 +481,8 @@ impl Served { let touched = Touched { paths: paths.to_vec(), }; - for line in self.lines.lock().unwrap().values() { - let _ = line.send(kind::TOUCHED, &touched); + for guest in self.guests.lock().unwrap().values() { + let _ = guest.line.send(kind::TOUCHED, &touched); } } diff --git a/crates/notebook/tests/live_share.rs b/crates/notebook/tests/live_share.rs index e24cf06409e17d365addb1dd8a5c7729817dbf23..268c76faa5fcc94d61a0d14d799182a10089aad2 100644 --- a/crates/notebook/tests/live_share.rs +++ b/crates/notebook/tests/live_share.rs @@ -429,3 +429,80 @@ fn large_files_travel_in_chunks() { .is_err() ); } + +/// A guest that floods its host with requests is hung up on once too many wait, having had +/// answers to few of them, and the host goes on serving the others. +#[test] +fn a_flooding_guest_is_hung_up_on() { + use notebook::live::{ + Event, Live, Room, + wire::{Bye, Request, kind}, + }; + use std::sync::{ + Mutex, + atomic::{AtomicUsize, Ordering}, + }; + let directory = tempfile::tempdir().unwrap(); + let folder = notebook(directory.path()); + let url = relay(Default::default()); + let sharing = Sharing::new("").unwrap(); + let host = host(&folder, &directory.path().join("host"), &sharing, &url); + let code = code(&host); + let (grace, notebook) = guest("Grace", &code, &url, &directory.path().join("grace")); + let welcome = share::join(hello("Mallory"), &code, "", None, Some(&url)).unwrap(); + let replies = Arc::new(AtomicUsize::new(0)); + let bye = Arc::new(Mutex::new(None)); + let (counted, said) = (Arc::clone(&replies), Arc::clone(&bye)); + let mallory = Live::start( + hello("Mallory"), + &Room::Notebook(welcome.secret), + None, + Some(&url), + move |event| match event { + Event::Frame { + kind: kind::REPLY, .. + } => { + counted.fetch_add(1, Ordering::Relaxed); + } + Event::Frame { + kind: kind::BYE, + body, + .. + } => *said.lock().unwrap() = minicbor::decode::(body).ok(), + _ => {} + }, + ) + .unwrap(); + let served = |live: &Live| { + live.peers() + .into_iter() + .find(|peer| peer.hello.serves == Some(welcome.share)) + }; + until("Mallory never met the host", || served(&mallory).is_some()); + let line = mallory.line(&served(&mallory).unwrap().hello.peer).unwrap(); + const SENT: u64 = 5000; + for id in 0..SENT { + let request = Request { + id, + path: "Garden.one".into(), + ..Request::default() + }; + if line.send(kind::STAMP, &request).is_err() { + break; + } + } + until("the host never hung up on Mallory", || { + bye.lock().unwrap().is_some() && served(&mallory).is_none() + }); + assert_eq!(bye.lock().unwrap().as_ref().unwrap().reason, "flooded"); + let answered = replies.load(Ordering::Relaxed); + // At most the burst a guest may start at once, what waits for the workers, and what the + // rate refills while the flood arrives. + assert!(answered < 300, "Mallory had {answered} answers of {SENT}"); + // Grace, asking at her own pace, is served as before. + assert!(grace.host().is_some()); + assert_eq!( + notebook.read_section("Garden.one").unwrap(), + std::fs::read(folder.join("Garden.one")).unwrap() + ); +} -- 2.54.0