From 9c5719fd7c0049468dded1316fcb4e9b7cf2bd13 Mon Sep 17 00:00:00 2001 From: reverse Date: Wed, 19 Aug 2026 03:28:51 +0000 Subject: [PATCH] =?UTF-8?q?socketengine:=20cap=20bytes=20drained=20from=20?= =?UTF-8?q?one=20socket=20per=20readable=20event=20(MAX=5FREAD=5FPER=5FTUR?= =?UTF-8?q?N=3D64KiB),=20then=20re-arm=20epoll=20and=20yield=20=E2=80=94?= =?UTF-8?q?=20bounds=20the=20per-turn=20line=20buffer=20and=20stops=20one?= =?UTF-8?q?=20flooding=20client=20from=20monopolising=20the=20reactor;=20t?= =?UTF-8?q?he=20leftover=20waits=20in=20the=20kernel=20buffer=20and=20is?= =?UTF-8?q?=20re-delivered=20next=20turn=20(verified:=20a=20133KB=20single?= =?UTF-8?q?-write=20burst=20gets=20every=20reply=20back)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/socketengine.rs | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) diff --git a/src/socketengine.rs b/src/socketengine.rs index 09c2988..7c11427 100644 --- a/src/socketengine.rs +++ b/src/socketengine.rs @@ -784,6 +784,11 @@ fn try_handshake( true } +/// Most bytes drained from one socket per readable event before we stop, re-arm and +/// yield: bounds the per-turn line buffer and stops one flooding client from +/// monopolising the reactor (the rest waits in the kernel buffer for the next turn). +const MAX_READ_PER_TURN: usize = 64 * 1024; + /// Drain readable bytes from `t` (edge-triggered: read until WouldBlock), frame /// complete lines and forward them to the core; close on EOF/error. fn read_conn(poll: &mut Poll, conns: &mut HashMap, t: usize, core: &Sender) { @@ -796,6 +801,8 @@ fn read_conn(poll: &mut Poll, conns: &mut HashMap, t: usize, core: // a deferred Connect (PROXY conn) to emit, before any lines from the same read let mut connect: Option<(Uid, SocketAddr, u16, bool, Option, OutSink)> = None; let mut close = false; + let mut read_total = 0usize; + let mut capped = false; if let Some(c) = conns.get_mut(&t) { loop { match c.sock.read(&mut chunk) { @@ -805,6 +812,7 @@ fn read_conn(poll: &mut Poll, conns: &mut HashMap, t: usize, core: } Ok(n) => { c.rbuf.extend_from_slice(&chunk[..n]); + read_total += n; if c.proxy_pending { // a v2 header from a TLS-terminating proxy can forward the // client's TLS status + cert fingerprint (see modules::proxy) @@ -860,6 +868,12 @@ fn read_conn(poll: &mut Poll, conns: &mut HashMap, t: usize, core: c.rbuf.clear(); // overlong line with no newline: drop it } } + // fairness + memory bound: after MAX_READ_PER_TURN bytes stop and + // re-arm, so one flooding client can't monopolise this reactor turn + if read_total >= MAX_READ_PER_TURN { + capped = true; + break; + } } Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => break, Err(ref e) if e.kind() == io::ErrorKind::Interrupted => continue, @@ -870,6 +884,18 @@ fn read_conn(poll: &mut Poll, conns: &mut HashMap, t: usize, core: } } } + // if we stopped early for fairness, force epoll to re-deliver the still-readable + // socket next turn (edge-triggered MOD re-reports a ready fd) so no bytes stall. + if capped && !close { + if let Some(c) = conns.get_mut(&t) { + let interest = match (c.want_read, c.want_write) { + (true, true) => Interest::READABLE | Interest::WRITABLE, + (false, true) => Interest::WRITABLE, + _ => Interest::READABLE, + }; + let _ = poll.registry().reregister(c.sock.source(), Token(t), interest); + } + } if let Some((uid, addr, local_port, secure, certfp, out)) = connect { if core .send(Event::Connect {