| author | |
| committer | |
| log | 7e428759e27de7e50a06ca5afadfbfe0f76c42f0 |
| tree | 8d03cd96c7b5e3ee0f7a7cc0f4e462cbd0b20e6e |
| parent | 3aa700c6c057e00e4eaf633161e28c5611b8b687 |
| signature | Signed by SSH key SHA256:52mNGHRsVFBDED9IAX5pe+LRWUefqTbxEReunq21QvU |
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.53 files changed, 174 insertions(+), 24 deletions(-)
crates/notebook/src/live.rs+21-6| ... | ... | @@ -222,6 +222,15 @@ impl Line { |
| 222 | 222 | .send(Out::Frame(kind, body)) |
| 223 | 223 | .map_err(|_| io::ErrorKind::NotConnected.into()) |
| 224 | 224 | } |
| 225 | ||
| 226 | /// Says `reason` after the frames sent before, then hangs up. | |
| 227 | pub fn hang_up(&self, reason: &str) { | |
| 228 | let bye = minicbor::to_vec(wire::Bye { | |
| 229 | reason: reason.into(), | |
| 230 | }) | |
| 231 | .unwrap_or_default(); | |
| 232 | let _ = self.0.send(Out::Bye(bye)); | |
| 233 | } | |
| 225 | 234 | } |
| 226 | 235 | |
| 227 | 236 | /// A stream to one peer, read by one thread and written by another: a TCP connection, or |
| ... | ... | @@ -597,12 +606,18 @@ impl Shared { |
| 597 | 606 | } |
| 598 | 607 | (self.events)(Event::Changed); |
| 599 | 608 | } |
| 600 | kind => (self.events)(Event::Frame { | |
| 601 | from: &hello, | |
| 602 | kind, | |
| 603 | body: &body, | |
| 604 | line: &line, | |
| 605 | }), | |
| 609 | kind => { | |
| 610 | (self.events)(Event::Frame { | |
| 611 | from: &hello, | |
| 612 | kind, | |
| 613 | body: &body, | |
| 614 | line: &line, | |
| 615 | }); | |
| 616 | // The peer hangs up after its bye. | |
| 617 | if kind == kind::BYE { | |
| 618 | break None; | |
| 619 | } | |
| 620 | } | |
| 606 | 621 | } |
| 607 | 622 | }; |
| 608 | 623 | if let Some(error) = ended.filter(|error| error.kind() == io::ErrorKind::InvalidData) { |
crates/notebook/src/live/share.rs+76-18| ... | ... | @@ -44,6 +44,16 @@ const SNAPSHOT_AGE: Duration = Duration::from_secs(120); |
| 44 | 44 | const IMAGES: usize = 4; |
| 45 | 45 | /// What a host says leaving as it stops sharing. |
| 46 | 46 | const STOPPED: &str = "stopped"; |
| 47 | /// What a host says hanging up on a guest that asks too much too fast. | |
| 48 | const FLOODED: &str = "flooded"; | |
| 49 | /// A guest's requests a host holds at once, waiting and in hand, past which it hangs up. | |
| 50 | const QUEUED: usize = 32; | |
| 51 | /// The threads working through each guest's requests. | |
| 52 | const WORKERS: usize = 2; | |
| 53 | /// Requests a guest may start each second, and at once: every request but a read's later | |
| 54 | /// chunks and an upload's, which come from memory or go to it. | |
| 55 | const STARTS: f64 = 100.0; | |
| 56 | const BURST: f64 = 200.0; | |
| 47 | 57 | |
| 48 | 58 | const WORDS: &str = include_str!("words.txt"); |
| 49 | 59 | |
| ... | ... | @@ -231,7 +241,7 @@ impl Host { |
| 231 | 241 | images: Mutex::default(), |
| 232 | 242 | snapshots: Mutex::default(), |
| 233 | 243 | puts: Mutex::default(), |
| 234 | lines: Mutex::default(), | |
| 244 | guests: Mutex::default(), | |
| 235 | 245 | }); |
| 236 | 246 | let serving = Hello { |
| 237 | 247 | serves: Some(sharing.share), |
| ... | ... | @@ -244,22 +254,12 @@ impl Host { |
| 244 | 254 | reach, |
| 245 | 255 | relay, |
| 246 | 256 | move |event| match event { |
| 247 | Event::Met(hello, line) => { | |
| 248 | heard.lines.lock().unwrap().insert(hello.peer, line.clone()); | |
| 249 | } | |
| 257 | Event::Met(hello, line) => heard.admit(hello.peer, line), | |
| 250 | 258 | Event::Left(hello) => heard.forget(&hello.peer), |
| 251 | 259 | Event::Frame { |
| 252 | from, | |
| 253 | kind, | |
| 254 | body, | |
| 255 | line, | |
| 260 | from, kind, body, .. | |
| 256 | 261 | } if wire::KNOWN.contains(&kind) && kind > 256 && kind != kind::REPLY => { |
| 257 | let (served, line, peer) = (Arc::clone(&heard), line.clone(), from.peer); | |
| 258 | let body = body.to_vec(); | |
| 259 | thread::spawn(move || { | |
| 260 | let reply = served.handle(&peer, kind, &body); | |
| 261 | let _ = line.send(kind::REPLY, &reply); | |
| 262 | }); | |
| 262 | heard.queue(&from.peer, kind, body); | |
| 263 | 263 | } |
| 264 | 264 | Event::Changed => told(), |
| 265 | 265 | _ => {} |
| ... | ... | @@ -379,7 +379,16 @@ struct Served { |
| 379 | 379 | snapshots: Mutex<ByGuest<(Arc<Vec<u8>>, Instant)>>, |
| 380 | 380 | /// Bytes a later request carries, by guest and upload. |
| 381 | 381 | puts: Mutex<ByGuest<Vec<u8>>>, |
| 382 | lines: Mutex<BTreeMap<[u8; 16], Line>>, | |
| 382 | guests: Mutex<BTreeMap<[u8; 16], Admitted>>, | |
| 383 | } | |
| 384 | ||
| 385 | /// A guest as its host serves it: the line to it, its requests waiting for its workers, and | |
| 386 | /// how many more it may start now. | |
| 387 | struct Admitted { | |
| 388 | line: Line, | |
| 389 | queue: mpsc::SyncSender<(u16, Vec<u8>)>, | |
| 390 | starts: f64, | |
| 391 | counted: Instant, | |
| 383 | 392 | } |
| 384 | 393 | |
| 385 | 394 | /// A section's path, stamp and image. |
| ... | ... | @@ -406,8 +415,57 @@ fn refused(kind: io::ErrorKind, message: &str) -> Error { |
| 406 | 415 | } |
| 407 | 416 | |
| 408 | 417 | impl Served { |
| 418 | /// Serves `peer` on `line`, through workers of its own that end as it leaves. | |
| 419 | fn admit(self: &Arc<Self>, peer: [u8; 16], line: &Line) { | |
| 420 | let (queue, waiting) = mpsc::sync_channel::<(u16, Vec<u8>)>(QUEUED - WORKERS); | |
| 421 | let waiting = Arc::new(Mutex::new(waiting)); | |
| 422 | for _ in 0..WORKERS { | |
| 423 | let (served, waiting, line) = | |
| 424 | (Arc::downgrade(self), Arc::clone(&waiting), line.clone()); | |
| 425 | thread::spawn(move || { | |
| 426 | loop { | |
| 427 | let next = waiting.lock().unwrap().recv(); | |
| 428 | let (Ok((kind, body)), Some(served)) = (next, served.upgrade()) else { | |
| 429 | return; | |
| 430 | }; | |
| 431 | let _ = line.send(kind::REPLY, &served.handle(&peer, kind, &body)); | |
| 432 | } | |
| 433 | }); | |
| 434 | } | |
| 435 | let guest = Admitted { | |
| 436 | line: line.clone(), | |
| 437 | queue, | |
| 438 | starts: BURST, | |
| 439 | counted: Instant::now(), | |
| 440 | }; | |
| 441 | self.guests.lock().unwrap().insert(peer, guest); | |
| 442 | } | |
| 443 | ||
| 444 | /// Hands a request from `peer` to its workers, or hangs up on a guest that has too many | |
| 445 | /// waiting or starts them too fast. | |
| 446 | fn queue(&self, peer: &[u8; 16], kind: u16, body: &[u8]) { | |
| 447 | let mut guests = self.guests.lock().unwrap(); | |
| 448 | let Some(guest) = guests.get_mut(peer) else { | |
| 449 | return; | |
| 450 | }; | |
| 451 | let continued = kind == kind::PUT | |
| 452 | || matches!(kind, kind::READ | kind::READ_FILE) | |
| 453 | && minicbor::decode::<Request>(body).is_ok_and(|request| request.handle.is_some()); | |
| 454 | let now = Instant::now(); | |
| 455 | guest.starts = | |
| 456 | (guest.starts + now.duration_since(guest.counted).as_secs_f64() * STARTS).min(BURST); | |
| 457 | guest.counted = now; | |
| 458 | if !continued { | |
| 459 | guest.starts -= 1.0; | |
| 460 | } | |
| 461 | if guest.starts < 0.0 || guest.queue.try_send((kind, body.to_vec())).is_err() { | |
| 462 | guest.line.hang_up(FLOODED); | |
| 463 | guests.remove(peer); | |
| 464 | } | |
| 465 | } | |
| 466 | ||
| 409 | 467 | fn forget(&self, peer: &[u8; 16]) { |
| 410 | self.lines.lock().unwrap().remove(peer); | |
| 468 | self.guests.lock().unwrap().remove(peer); | |
| 411 | 469 | self.snapshots |
| 412 | 470 | .lock() |
| 413 | 471 | .unwrap() |
| ... | ... | @@ -423,8 +481,8 @@ impl Served { |
| 423 | 481 | let touched = Touched { |
| 424 | 482 | paths: paths.to_vec(), |
| 425 | 483 | }; |
| 426 | for line in self.lines.lock().unwrap().values() { | |
| 427 | let _ = line.send(kind::TOUCHED, &touched); | |
| 484 | for guest in self.guests.lock().unwrap().values() { | |
| 485 | let _ = guest.line.send(kind::TOUCHED, &touched); | |
| 428 | 486 | } |
| 429 | 487 | } |
| 430 | 488 |
crates/notebook/tests/live_share.rs+77| ... | ... | @@ -429,3 +429,80 @@ fn large_files_travel_in_chunks() { |
| 429 | 429 | .is_err() |
| 430 | 430 | ); |
| 431 | 431 | } |
| 432 | ||
| 433 | /// A guest that floods its host with requests is hung up on once too many wait, having had | |
| 434 | /// answers to few of them, and the host goes on serving the others. | |
| 435 | #[test] | |
| 436 | fn a_flooding_guest_is_hung_up_on() { | |
| 437 | use notebook::live::{ | |
| 438 | Event, Live, Room, | |
| 439 | wire::{Bye, Request, kind}, | |
| 440 | }; | |
| 441 | use std::sync::{ | |
| 442 | Mutex, | |
| 443 | atomic::{AtomicUsize, Ordering}, | |
| 444 | }; | |
| 445 | let directory = tempfile::tempdir().unwrap(); | |
| 446 | let folder = notebook(directory.path()); | |
| 447 | let url = relay(Default::default()); | |
| 448 | let sharing = Sharing::new("").unwrap(); | |
| 449 | let host = host(&folder, &directory.path().join("host"), &sharing, &url); | |
| 450 | let code = code(&host); | |
| 451 | let (grace, notebook) = guest("Grace", &code, &url, &directory.path().join("grace")); | |
| 452 | let welcome = share::join(hello("Mallory"), &code, "", None, Some(&url)).unwrap(); | |
| 453 | let replies = Arc::new(AtomicUsize::new(0)); | |
| 454 | let bye = Arc::new(Mutex::new(None)); | |
| 455 | let (counted, said) = (Arc::clone(&replies), Arc::clone(&bye)); | |
| 456 | let mallory = Live::start( | |
| 457 | hello("Mallory"), | |
| 458 | &Room::Notebook(welcome.secret), | |
| 459 | None, | |
| 460 | Some(&url), | |
| 461 | move |event| match event { | |
| 462 | Event::Frame { | |
| 463 | kind: kind::REPLY, .. | |
| 464 | } => { | |
| 465 | counted.fetch_add(1, Ordering::Relaxed); | |
| 466 | } | |
| 467 | Event::Frame { | |
| 468 | kind: kind::BYE, | |
| 469 | body, | |
| 470 | .. | |
| 471 | } => *said.lock().unwrap() = minicbor::decode::<Bye>(body).ok(), | |
| 472 | _ => {} | |
| 473 | }, | |
| 474 | ) | |
| 475 | .unwrap(); | |
| 476 | let served = |live: &Live| { | |
| 477 | live.peers() | |
| 478 | .into_iter() | |
| 479 | .find(|peer| peer.hello.serves == Some(welcome.share)) | |
| 480 | }; | |
| 481 | until("Mallory never met the host", || served(&mallory).is_some()); | |
| 482 | let line = mallory.line(&served(&mallory).unwrap().hello.peer).unwrap(); | |
| 483 | const SENT: u64 = 5000; | |
| 484 | for id in 0..SENT { | |
| 485 | let request = Request { | |
| 486 | id, | |
| 487 | path: "Garden.one".into(), | |
| 488 | ..Request::default() | |
| 489 | }; | |
| 490 | if line.send(kind::STAMP, &request).is_err() { | |
| 491 | break; | |
| 492 | } | |
| 493 | } | |
| 494 | until("the host never hung up on Mallory", || { | |
| 495 | bye.lock().unwrap().is_some() && served(&mallory).is_none() | |
| 496 | }); | |
| 497 | assert_eq!(bye.lock().unwrap().as_ref().unwrap().reason, "flooded"); | |
| 498 | let answered = replies.load(Ordering::Relaxed); | |
| 499 | // At most the burst a guest may start at once, what waits for the workers, and what the | |
| 500 | // rate refills while the flood arrives. | |
| 501 | assert!(answered < 300, "Mallory had {answered} answers of {SENT}"); | |
| 502 | // Grace, asking at her own pace, is served as before. | |
| 503 | assert!(grace.host().is_some()); | |
| 504 | assert_eq!( | |
| 505 | notebook.read_section("Garden.one").unwrap(), | |
| 506 | std::fs::read(folder.join("Garden.one")).unwrap() | |
| 507 | ); | |
| 508 | } |