diff --git a/crates/beacon_api/src/server.rs b/crates/beacon_api/src/server.rs index 14b1a407..ca9c7520 100644 --- a/crates/beacon_api/src/server.rs +++ b/crates/beacon_api/src/server.rs @@ -9,8 +9,8 @@ use mio::{Events, Interest, Registry, Token, event::Event}; use silver_beacon_state_data::{BeaconStateReader, SpecConfig}; use silver_common::{Enr, Identify, Keypair}; use silver_httpcore::{ - AfterResponse, Bind, ChunkedResponse, Listener, ParsedRequest, ServerConnection, Stream, - TokenRange, + AfterResponse, Bind, ChunkedResponse, Closed, Listener, ParsedRequest, ServerConnection, + Stream, TokenRange, }; use crate::{ @@ -255,18 +255,7 @@ impl Subscription { return Ok(true); } - if event.is_writable() { - while !self.body.pending_write().is_empty() { - match stream.write(self.body.pending_write()) { - Ok(0) => { - return Err(io::Error::new(io::ErrorKind::WriteZero, "write returned 0")) - } - Ok(n) => self.body.commit_write(n, now), - Err(e) if would_block(&e) => return Ok(false), - Err(e) if interrupted(&e) => continue, - Err(e) => return Err(e), - } - } + if event.is_writable() && self.body.drain_into(stream, now)? { registry.reregister(stream, event.token(), Interest::READABLE)?; } @@ -448,26 +437,36 @@ impl BeaconApi { let Self { connections, registry, .. } = self; let mut pushed = false; connections.retain(|token, conn| { - let State::Subscription(subscription) = &mut conn.state else { return true }; + let Connection { stream, state: State::Subscription(subscription) } = conn else { + return true; + }; if !wants(subscription) { return true; } - if !subscription.body.push(chunk, now) { - tracing::warn!( - "beacon api subscriber would exceed send cap with {} bytes already pending, closing", - subscription.body.pending_write().len() - ); - let _ = registry.deregister(&mut conn.stream); - return false; - } - pushed = true; - let interest = Interest::READABLE | Interest::WRITABLE; - if let Err(e) = registry.reregister(&mut conn.stream, *token, interest) { - tracing::warn!("beacon api subscriber lost: {e}"); - let _ = registry.deregister(&mut conn.stream); - return false; + let outcome = subscription.body.deliver(stream, chunk, now).and_then(|interest| { + pushed = true; + match interest { + Some(interest) => { + registry.reregister(stream, *token, interest).map_err(Closed::Lost) + } + None => Ok(()), + } + }); + match outcome { + Ok(()) => true, + Err(Closed::AtCap { pending }) => { + tracing::warn!( + "beacon api subscriber would exceed send cap with {pending} bytes already pending, closing" + ); + let _ = registry.deregister(stream); + false + } + Err(Closed::Lost(e)) => { + tracing::warn!("beacon api subscriber lost: {e}"); + let _ = registry.deregister(stream); + false + } } - true }); pushed } @@ -757,7 +756,7 @@ mod tests { /// later one. Handing that offset out again replaces the map entry, which /// drops the older connection and closes its socket unannounced. #[test] - fn a_recycled_offset_skips_the_connection_still_holding_it() { + fn recycled_offset_skips_the_connection_still_holding_it() { let span = 3; let tokens = TokenRange::new(64, span); let binds = [Bind::parse("127.0.0.1:0")]; @@ -1072,7 +1071,7 @@ mod tests { /// An operator large enough to declare more body than the read buffer /// holds gets a status back rather than a connection that goes quiet. #[test] - fn a_body_declared_past_the_read_cap_is_answered_with_413() { + fn body_declared_past_the_read_cap_is_answered_with_413() { let mut server = server_with(64, LONG_TIMEOUT); let addr = tcp_addr(&server); @@ -1097,7 +1096,7 @@ mod tests { /// like every other reject, where dropping the socket mid-send would be /// the reset that costs a client its node. #[test] - fn a_head_that_outgrows_the_read_buffer_is_answered_with_431() { + fn head_that_outgrows_the_read_buffer_is_answered_with_431() { let mut server = server_with(64, LONG_TIMEOUT); let addr = tcp_addr(&server); @@ -1129,7 +1128,7 @@ mod tests { /// caller and costs the node its place in the rotation, where a 413 costs /// nothing. #[test] - fn a_client_still_streaming_when_the_413_is_framed_reads_all_of_it() { + fn client_still_streaming_when_the_413_is_framed_reads_all_of_it() { let mut server = server_with(64, LONG_TIMEOUT); let addr = tcp_addr(&server); @@ -1148,7 +1147,7 @@ mod tests { /// Unix sockets take the same half-close, so the drain ends on the peer's /// own close there too rather than running to the linger cap. #[test] - fn a_client_still_streaming_over_uds_reads_all_of_the_413() { + fn client_still_streaming_over_uds_reads_all_of_the_413() { let dir = tempfile::tempdir().unwrap(); let socket = dir.path().join("api.sock"); let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); @@ -1179,7 +1178,7 @@ mod tests { /// Draining an answered connection is bounded: one client cannot hold a /// slot for as long as it cares to keep sending. #[test] - fn a_client_that_never_stops_sending_is_dropped_at_the_linger_cap() { + fn client_that_never_stops_sending_is_dropped_at_the_linger_cap() { let mut server = server_with(64, Duration::from_millis(800)); // A peer that never pauses keeps the wait between reads at zero, so the // total cap is the only one that can end it. @@ -1218,7 +1217,7 @@ mod tests { /// the wait between reads, not for the whole draining window — and not for /// the far longer deadline that keeps a served connection available. #[test] - fn a_lingering_connection_that_goes_quiet_is_dropped_at_the_idle_cap() { + fn lingering_connection_that_goes_quiet_is_dropped_at_the_idle_cap() { let idle_timeout = Duration::from_secs(2); let mut server = server_with(64, idle_timeout); server.api.linger = @@ -1412,7 +1411,7 @@ mod tests { } #[test] - fn a_subscriber_gets_the_head_then_every_block_published_on_its_channel() { + fn subscriber_gets_the_head_then_every_block_published_on_its_channel() { let mut server = server_with(64, LONG_TIMEOUT); let mut client = connect(tcp_addr(&server)); subscribe(&mut client, "block"); @@ -1571,7 +1570,7 @@ mod tests { } #[test] - fn a_topic_silver_does_not_serve_is_refused_on_an_ordinary_connection() { + fn topic_silver_does_not_serve_is_refused_on_an_ordinary_connection() { let mut server = server_with(64, LONG_TIMEOUT); let addr = tcp_addr(&server); let client = std::thread::spawn(move || { @@ -1605,8 +1604,117 @@ mod tests { assert_same_bytes(&got, &expected); } + fn burst_frame(index: usize, len: usize) -> Vec { + let mut frame = format!("event: burst\ndata: {index:02}").into_bytes(); + frame.resize(len - 2, b'c'); + frame.extend_from_slice(b"\n\n"); + frame + } + + /// The burst fits in the application buffer even if the socket initially + /// accepts no bytes. Delivery therefore does not require a particular + /// kernel send-buffer capacity. The frames approximate one block's column + /// events with commitments; any unsent remainder drains through the + /// readiness loop. + #[test] + fn burst_reaches_a_reading_subscriber_in_order() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("api.sock"); + let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); + let mut client = connect_uds(&socket); + subscribe(&mut client, "block"); + pump_until(&mut server, "subscribed and head sent", |server| { + subscribers(server) == 1 && bytes_waiting_for_subscribers(server) == 0 + }); + + let frames: Vec<_> = (0..128).map(|index| burst_frame(index, 2300)).collect(); + let mut expected = SSE_HEAD.to_vec(); + frames.iter().for_each(|frame| expected.extend(chunk(frame))); + let now = Instant::now(); + for frame in &frames { + assert!(server.api.fan_out(|_| true, frame, now), "queued for the subscriber"); + } + assert_eq!(subscribers(&server), 1); + + let got = serve(&mut server, read_exactly(client, expected.len()), "the burst"); + assert_same_bytes(&got, &expected); + } + + /// Checks that publication attempts delivery before the readiness loop, + /// independently of whether a larger burst fits the application buffer. + #[test] + fn a_publish_to_a_reading_subscriber_leaves_nothing_pending() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("api.sock"); + let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); + let mut client = connect_uds(&socket); + subscribe(&mut client, "block"); + pump_until(&mut server, "subscribed and head sent", |server| { + subscribers(server) == 1 && bytes_waiting_for_subscribers(server) == 0 + }); + + for slot in 1..=3 { + server.api.publish_block(slot, &[0x33; 32]); + } + assert_eq!(bytes_waiting_for_subscribers(&server), 0); + assert_eq!(subscribers(&server), 1); + drop(client); + } + + /// On Linux, registering `WRITABLE` on a writable socket queues an epoll + /// event. The small frame keeps this probe independent of burst capacity. + #[cfg(target_os = "linux")] + #[test] + fn frame_written_whole_to_an_idle_subscriber_leaves_nothing_to_report() { + let mut server = server_with(64, LONG_TIMEOUT); + let mut client = connect(tcp_addr(&server)); + subscribe(&mut client, "block"); + pump_until(&mut server, "subscribed and head sent", |server| { + subscribers(server) == 1 && bytes_waiting_for_subscribers(server) == 0 + }); + + server.api.publish_block(7, &[0x77; 32]); + assert_eq!(bytes_waiting_for_subscribers(&server), 0, "the socket took the frame"); + server.readiness.wait(Duration::ZERO); + assert_eq!(server.readiness.events().iter().count(), 0); + drop(client); + } + + /// The queued response head leaves `WRITABLE` registered. On Linux, a + /// publish that drains the head must remove that interest before polling + /// and retain interest in inbound bytes. + #[cfg(target_os = "linux")] + #[test] + fn publish_that_drains_the_unsent_head_leaves_nothing_to_report() { + let dir = tempfile::tempdir().unwrap(); + let socket = dir.path().join("api.sock"); + let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); + let mut client = connect_uds(&socket); + subscribe(&mut client, "block"); + pump_until(&mut server, "subscribed", |server| subscribers(server) == 1); + assert!(bytes_waiting_for_subscribers(&server) > 0, "the head is still queued"); + + server.api.publish_block(5, &[0x55; 32]); + assert_eq!(bytes_waiting_for_subscribers(&server), 0, "the publish drained the head"); + server.readiness.wait(Duration::ZERO); + assert_eq!(server.readiness.events().iter().count(), 0, "no writable event remains"); + + client.write_all(b"GET /eth/v1/node/version HTTP/1.1\r\nHost: x\r\n\r\n").unwrap(); + assert!(server.pump(), "inbound bytes produce a readable event"); + assert_eq!( + subscribers(&server), + 1, + "a request behind the subscribe is dropped, not answered" + ); + + drop(client); + pump_until(&mut server, "hung-up subscriber removed", |server| { + server.api.connections.is_empty() + }); + } + #[test] - fn a_subscriber_that_stops_reading_is_closed_at_the_cap() { + fn subscriber_that_stops_reading_is_closed_at_the_cap() { let dir = tempfile::tempdir().unwrap(); let socket = dir.path().join("api.sock"); let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); @@ -1628,7 +1736,7 @@ mod tests { } #[test] - fn a_subscriber_that_takes_nothing_for_the_send_deadline_is_closed() { + fn subscriber_that_takes_nothing_for_the_send_deadline_is_closed() { let dir = tempfile::tempdir().unwrap(); let socket = dir.path().join("api.sock"); let mut server = server_bound_to(&[Bind::Unix(socket.clone())], 64, LONG_TIMEOUT); @@ -1654,7 +1762,7 @@ mod tests { } #[test] - fn a_quiet_subscriber_outlives_the_idle_timeout() { + fn quiet_subscriber_outlives_the_idle_timeout() { let idle_timeout = Duration::from_millis(100); let mut server = server_with(64, idle_timeout); let mut client = connect(tcp_addr(&server)); @@ -1667,7 +1775,7 @@ mod tests { } #[test] - fn a_subscriber_that_hangs_up_is_forgotten() { + fn subscriber_that_hangs_up_is_forgotten() { let mut server = server_with(64, LONG_TIMEOUT); let mut client = connect(tcp_addr(&server)); subscribe(&mut client, "block"); @@ -1680,7 +1788,7 @@ mod tests { } #[test] - fn a_request_pipelined_behind_the_subscribe_is_never_answered() { + fn request_pipelined_behind_the_subscribe_is_never_answered() { let mut server = server_with(64, LONG_TIMEOUT); let mut client = connect(tcp_addr(&server)); write!( diff --git a/crates/httpcore/src/chunked_response.rs b/crates/httpcore/src/chunked_response.rs index c6dd5b01..9509917e 100644 --- a/crates/httpcore/src/chunked_response.rs +++ b/crates/httpcore/src/chunked_response.rs @@ -1,11 +1,15 @@ use std::{ - io::Write, + io::{self, ErrorKind, Write}, time::{Duration, Instant}, }; +use mio::Interest; + // Reserve the full output allowance at construction so accepted pushes do // not reallocate. The response head and chunk framing count against it. -const PENDING_MAX: usize = 64 << 10; +// Allows one block's 128 column events with 21 commitments each, about +// 290 KiB in total, even when the writer accepts no bytes. +const PENDING_MAX: usize = 512 << 10; const DISCARD_LEN: usize = 4096; /// Does not emit a terminal chunk; the caller ends the stream by closing @@ -57,6 +61,52 @@ impl ChunkedResponse { true } + /// Attempts to drain existing output before checking whether the framed + /// chunk fits under the pending-byte cap. If it fits, queues it and + /// attempts to drain again. Socket write progress can therefore free + /// room for the new chunk. + /// + /// Returns replacement readiness interests for the caller to register. + /// `None` means no output was pending on entry and the new chunk + /// drained completely, so the caller can retain its `READABLE` + /// registration. Errors require the caller to close the connection. + pub fn deliver( + &mut self, + stream: &mut impl Write, + chunk: &[u8], + now: Instant, + ) -> Result, Closed> { + let had_backlog = !self.pending_write().is_empty(); + if had_backlog { + self.drain_into(stream, now).map_err(Closed::Lost)?; + } + if !self.push(chunk, now) { + return Err(Closed::AtCap { pending: self.pending_write().len() }); + } + // An empty buffer already has READABLE alone. A backlog may retain + // WRITABLE from the response head or an earlier blocked write. + Ok(match (self.drain_into(stream, now).map_err(Closed::Lost)?, had_backlog) { + (true, false) => None, + (true, true) => Some(Interest::READABLE), + (false, _) => Some(Interest::READABLE | Interest::WRITABLE), + }) + } + + /// Returns `true` when no output remains, or `false` on `WouldBlock`. + /// On `true`, readiness-driven callers restore `READABLE` alone. + pub fn drain_into(&mut self, stream: &mut impl Write, now: Instant) -> io::Result { + while !self.pending_write().is_empty() { + match stream.write(self.pending_write()) { + Ok(0) => return Err(io::Error::new(ErrorKind::WriteZero, "write returned 0")), + Ok(n) => self.commit_write(n, now), + Err(e) if e.kind() == ErrorKind::WouldBlock => return Ok(false), + Err(e) if e.kind() == ErrorKind::Interrupted => continue, + Err(e) => return Err(e), + } + } + Ok(true) + } + pub fn pending_write(&self) -> &[u8] { &self.pending[self.write_pos..] } @@ -86,6 +136,12 @@ impl ChunkedResponse { } } +#[derive(Debug)] +pub enum Closed { + AtCap { pending: usize }, + Lost(io::Error), +} + fn hex_digits(n: usize) -> usize { (usize::BITS - n.leading_zeros()).div_ceil(4).max(1) as usize } @@ -102,7 +158,7 @@ pub fn frame_chunked_head(out: &mut Vec, content_type: &str, headers: &[(&st #[cfg(test)] mod tests { - use std::cell::Cell; + use std::{cell::Cell, collections::VecDeque}; use super::*; use crate::{ParsedRequest, ServerConnection, frame_response}; @@ -157,8 +213,271 @@ mod tests { vec![b'e'; n] } + fn burst_frame(index: usize, len: usize) -> Vec { + let mut frame = format!("event: burst\ndata: {index:02}").into_bytes(); + frame.resize(len - 2, b'c'); + frame.extend_from_slice(b"\n\n"); + frame + } + + #[derive(Clone, Copy)] + enum Step { + Take(usize), + WouldBlock, + Zero, + Broken, + } + + /// After the scripted steps are consumed, each write follows `then`. + struct ScriptedSocket { + steps: VecDeque, + then: Step, + taken: Vec, + } + + impl ScriptedSocket { + fn taking_everything() -> Self { + Self { steps: VecDeque::new(), then: Step::Take(usize::MAX), taken: Vec::new() } + } + + fn refusing_everything() -> Self { + Self { then: Step::WouldBlock, ..Self::taking_everything() } + } + + fn script(&mut self, steps: impl IntoIterator) { + self.steps.extend(steps); + } + } + + impl Write for ScriptedSocket { + fn write(&mut self, buf: &[u8]) -> io::Result { + match self.steps.pop_front().unwrap_or(self.then) { + Step::Take(n) => { + let n = n.min(buf.len()); + self.taken.extend_from_slice(&buf[..n]); + Ok(n) + } + Step::WouldBlock => Err(ErrorKind::WouldBlock.into()), + Step::Zero => Ok(0), + Step::Broken => Err(ErrorKind::BrokenPipe.into()), + } + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + const BOTH: Interest = Interest::READABLE.add(Interest::WRITABLE); + + #[derive(Default)] + struct BudgetSocket { + remaining: usize, + taken: Vec, + interrupt_once: bool, + } + + impl Write for BudgetSocket { + fn write(&mut self, buf: &[u8]) -> io::Result { + if std::mem::take(&mut self.interrupt_once) { + return Err(ErrorKind::Interrupted.into()); + } + if self.remaining == 0 { + return Err(ErrorKind::WouldBlock.into()); + } + let n = buf.len().min(self.remaining); + self.taken.extend_from_slice(&buf[..n]); + self.remaining -= n; + Ok(n) + } + + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + + /// The queued head supplies the initial backlog. Once drained, later + /// deliveries need no registration change while the writer accepts output. + #[test] + fn burst_past_the_cap_is_written_as_it_is_pushed() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = ScriptedSocket::taking_everything(); + let past_the_cap = PENDING_MAX / framed(&burst_frame(0, 2300)).len() + 1; + let frames: Vec<_> = (0..past_the_cap).map(|index| burst_frame(index, 2300)).collect(); + let mut expected = HEAD.to_vec(); + frames.iter().for_each(|frame| expected.extend(framed(frame))); + assert!(expected.len() - HEAD.len() > PENDING_MAX, "the burst passes the cap"); + + let (first, rest) = frames.split_first().unwrap(); + assert_eq!(stream.deliver(&mut socket, first, t0).unwrap(), Some(Interest::READABLE)); + for frame in rest { + assert_eq!(stream.deliver(&mut socket, frame, t0).unwrap(), None); + } + assert_eq!(socket.taken, expected); + assert!(stream.pending_write().is_empty()); + } + + /// A blocked write requests `WRITABLE`; a delivery that drains the backlog + /// requests `READABLE` alone. The readiness path reports a complete drain + /// for the caller to make the same registration change. + #[test] + fn refused_write_arms_writable_until_the_backlog_drains() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = ScriptedSocket::taking_everything(); + assert!(stream.drain_into(&mut socket, t0).unwrap()); + let frames = [burst_frame(1, 300), burst_frame(2, 300), burst_frame(3, 300)]; + + socket.script([Step::WouldBlock]); + assert_eq!(stream.deliver(&mut socket, &frames[0], t0).unwrap(), Some(BOTH)); + assert_eq!(stream.pending_write(), framed(&frames[0])); + let by_publish = stream.deliver(&mut socket, &frames[1], t0).unwrap(); + assert_eq!(by_publish, Some(Interest::READABLE)); + assert!(stream.pending_write().is_empty()); + + socket.script([Step::WouldBlock]); + assert_eq!(stream.deliver(&mut socket, &frames[2], t0).unwrap(), Some(BOTH)); + assert!(stream.drain_into(&mut socket, t0).unwrap(), "the loop drains the rest"); + + let expected: Vec<_> = frames.iter().flat_map(|frame| framed(frame)).collect(); + assert_eq!(socket.taken[HEAD.len()..], expected); + } + + /// Partial writes reset the stall clock. Interruptions are retried, and + /// draining the remaining bytes clears the clock. + #[test] + fn partial_writes_continue_until_the_socket_refuses() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = BudgetSocket { remaining: usize::MAX, ..BudgetSocket::default() }; + let frame = burst_frame(0, 1000); + let expected = [HEAD, &framed(&frame)].concat(); + + assert!(stream.drain_into(&mut socket, t0).unwrap()); + let already_sent = socket.taken.len(); + let allowance = 200; + socket.remaining = allowance; + let t1 = t0 + Duration::from_secs(1); + assert_eq!(stream.deliver(&mut socket, &frame, t1).unwrap(), Some(BOTH)); + let sent = already_sent + allowance; + assert_eq!(socket.taken, expected[..sent]); + assert_eq!(stream.pending_write(), &expected[sent..]); + assert!(!stream.stalled(t1 + DEADLINE, DEADLINE)); + assert!(stream.stalled(t1 + DEADLINE + Duration::from_millis(1), DEADLINE)); + + socket.interrupt_once = true; + socket.remaining = usize::MAX; + assert!(stream.drain_into(&mut socket, t1).unwrap()); + assert_eq!(socket.taken, expected); + assert!(!stream.stalled(t1 + DEADLINE * 100, DEADLINE)); + } + + /// Synthetic frames approximate column events with 21 commitments. + /// The queued response head also counts against the allowance. + #[test] + fn one_blocks_column_events_fit_the_cap_with_the_socket_taking_nothing() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = ScriptedSocket::refusing_everything(); + let frames: Vec<_> = (0..128).map(|index| burst_frame(index, 2300)).collect(); + + for frame in &frames { + assert_eq!(stream.deliver(&mut socket, frame, t0).unwrap(), Some(BOTH)); + } + let mut expected = HEAD.to_vec(); + frames.iter().for_each(|frame| expected.extend(framed(frame))); + assert_eq!(stream.pending_write(), expected); + } + + /// The cap error reports bytes already pending, including the response + /// head. + #[test] + fn pushes_the_socket_never_takes_close_at_the_cap() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = ScriptedSocket::refusing_everything(); + let frame = burst_frame(0, 2300); + let chunk = framed(&frame).len(); + let fit = (PENDING_MAX - HEAD.len()) / chunk; + + for pushed in 0..fit { + let interest = stream.deliver(&mut socket, &frame, t0).unwrap(); + assert_eq!(interest, Some(BOTH), "push {pushed} waits for the socket"); + } + let closed = stream.deliver(&mut socket, &frame, t0).unwrap_err(); + let at_cap = + matches!(closed, Closed::AtCap { pending } if pending == HEAD.len() + fit * chunk); + assert!(at_cap, "{closed:?}"); + assert!(socket.taken.is_empty()); + } + + #[test] + fn draining_a_full_backlog_makes_room_for_the_next_frame() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = BudgetSocket::default(); + let first = payload_framed_to(PENDING_MAX - HEAD.len()); + stream.deliver(&mut socket, &first, t0).unwrap(); + assert_eq!(stream.pending_write().len(), PENDING_MAX); + + let next = burst_frame(1, 300); + socket.remaining = framed(&next).len(); + let t1 = t0 + Duration::from_secs(1); + stream.deliver(&mut socket, &next, t1).unwrap(); + assert!(stream.pending_write().len() <= PENDING_MAX); + + socket.remaining = usize::MAX; + assert!(stream.drain_into(&mut socket, t1).unwrap()); + assert_eq!(socket.taken, [HEAD, &framed(&first), &framed(&next)].concat()); + } + + #[test] + fn draining_too_little_still_refuses_the_next_frame() { + let t0 = Instant::now(); + let mut stream = subscribed(t0); + let mut socket = BudgetSocket::default(); + let first = payload_framed_to(PENDING_MAX - HEAD.len()); + stream.deliver(&mut socket, &first, t0).unwrap(); + assert_eq!(stream.pending_write().len(), PENDING_MAX); + + let next = burst_frame(1, 300); + socket.remaining = framed(&next).len() - 1; + let t1 = t0 + Duration::from_secs(1); + let Err(Closed::AtCap { pending }) = stream.deliver(&mut socket, &next, t1) else { + panic!("insufficient room after draining must still end delivery") + }; + + let expected = [HEAD, &framed(&first)].concat(); + let sent = socket.taken.len(); + assert!(sent > 0, "delivery must try draining before refusing the frame"); + assert_eq!(pending, expected.len() - sent); + assert_eq!(socket.taken, expected[..sent]); + } + + /// Both zero-length writes and write errors require the caller to close. + #[test] + fn dead_socket_ends_the_stream() { + let t0 = Instant::now(); + let frame = burst_frame(0, 100); + + let mut zero = ScriptedSocket::taking_everything(); + zero.script([Step::Zero]); + let Err(Closed::Lost(e)) = subscribed(t0).deliver(&mut zero, &frame, t0) else { + panic!("a zero-length write is an error") + }; + assert_eq!(e.kind(), ErrorKind::WriteZero); + + let mut broken = ScriptedSocket::taking_everything(); + broken.script([Step::Broken]); + let Err(Closed::Lost(e)) = subscribed(t0).deliver(&mut broken, &frame, t0) else { + panic!("a failed write is an error") + }; + assert_eq!(e.kind(), ErrorKind::BrokenPipe); + } + #[test] - fn the_head_leaves_first_ahead_of_chunks_pushed_before_the_first_drain() { + fn head_leaves_first_ahead_of_chunks_pushed_before_the_first_drain() { let t0 = Instant::now(); let mut stream = subscribed(t0); assert!(stream.push(b"event: block\ndata: {}\n\n", t0)); @@ -185,7 +504,7 @@ mod tests { } #[test] - fn a_push_landing_exactly_on_the_cap_fits_and_the_next_byte_is_refused() { + fn push_landing_exactly_on_the_cap_fits_and_the_next_byte_is_refused() { let t0 = Instant::now(); let mut stream = subscribed(t0); let payload = vec![b'x'; largest_fitting_payload(&stream)]; @@ -203,7 +522,7 @@ mod tests { #[test] #[should_panic(expected = "terminal chunk")] - fn an_empty_chunk_is_refused() { + fn empty_chunk_is_refused() { let t0 = Instant::now(); let mut stream = subscribed(t0); let _ = stream.push(b"", t0); @@ -228,7 +547,7 @@ mod tests { } #[test] - fn the_send_buffer_is_allocated_once_at_the_cap_and_never_moves() { + fn send_buffer_is_allocated_once_at_the_cap_and_never_moves() { let t0 = Instant::now(); let mut stream = subscribed(t0); assert_eq!(stream.pending.capacity(), PENDING_MAX); @@ -288,7 +607,7 @@ mod tests { } #[test] - fn an_exact_fit_behind_a_consumed_prefix_is_appended_without_compaction() { + fn exact_fit_behind_a_consumed_prefix_is_appended_without_compaction() { let t0 = Instant::now(); let mut stream = subscribed(t0); drain_all(&mut stream, t0); @@ -328,7 +647,7 @@ mod tests { } #[test] - fn a_refused_push_moves_nothing_either() { + fn refused_push_moves_nothing_either() { let t0 = Instant::now(); let mut stream = subscribed(t0); drain_all(&mut stream, t0); @@ -344,12 +663,12 @@ mod tests { #[test] #[should_panic(expected = "exceeds the send cap")] - fn a_head_past_the_cap_is_a_bug() { + fn head_past_the_cap_is_a_bug() { let _ = ChunkedResponse::new(vec![b'h'; PENDING_MAX + 1], vec![0; 16], Instant::now()); } #[test] - fn a_read_buffer_grown_by_an_earlier_request_is_cut_back_to_scratch_size() { + fn read_buffer_grown_by_an_earlier_request_is_cut_back_to_scratch_size() { let t0 = Instant::now(); let mut conn = ServerConnection::new(); let body = vec![b'b'; 6000]; @@ -375,7 +694,7 @@ mod tests { } #[test] - fn the_waiting_clock_moves_only_when_bytes_leave() { + fn waiting_clock_moves_only_when_bytes_leave() { let t0 = Instant::now(); let mut stream = subscribed(t0); @@ -461,7 +780,7 @@ mod tests { } #[test] - fn a_materialised_response_still_frames_as_before() { + fn materialised_response_still_frames_as_before() { let mut out = Vec::new(); frame_response(&mut out, "200 OK", None, b"ok"); assert_eq!(out, b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok"); diff --git a/crates/httpcore/src/lib.rs b/crates/httpcore/src/lib.rs index 4aae12e9..b8d8aa63 100644 --- a/crates/httpcore/src/lib.rs +++ b/crates/httpcore/src/lib.rs @@ -6,7 +6,7 @@ mod server; mod stream; mod token_range; -pub use chunked_response::{ChunkedResponse, frame_chunked_head}; +pub use chunked_response::{ChunkedResponse, Closed, frame_chunked_head}; pub use client::{ClientConnection, frame_request}; pub use query::Query; pub use readiness::Readiness; diff --git a/crates/httpcore/src/server.rs b/crates/httpcore/src/server.rs index bdb4ac6c..f6c6f101 100644 --- a/crates/httpcore/src/server.rs +++ b/crates/httpcore/src/server.rs @@ -572,7 +572,7 @@ mod tests { /// the buffer filling up is itself the verdict: the 431 of RFC 6585, /// answered and lingered like every other reject rather than a bare drop. #[test] - fn an_exhausted_buffer_holding_no_request_writes_431_then_lingers() { + fn exhausted_buffer_holding_no_request_writes_431_then_lingers() { let mut conn = ServerConnection::new(); feed(&mut conn, b"GET /eth/v1/node/health HTTP/1.1\r\nCookie: "); loop { @@ -605,7 +605,7 @@ mod tests { /// The body the client is still sending is read for one reason only — to /// keep the answer from being lost — so it must cost nothing to read. #[test] - fn a_lingering_connection_discards_without_growing() { + fn lingering_connection_discards_without_growing() { let mut conn = reject_and_linger(&oversized_post(READ_BUF_MAX), "413 Payload Too Large"); let scratch = conn.read_buf.len(); @@ -865,7 +865,7 @@ mod tests { } #[test] - fn a_large_response_does_not_leave_the_connection_inflated() { + fn large_response_does_not_leave_the_connection_inflated() { let mut conn = ServerConnection::new(); let big = vec![b'x'; 4 << 20]; feed(&mut conn, &get_req("/big", "HTTP/1.1")); diff --git a/docs/adr/0004-sync-materialized-api.md b/docs/adr/0004-sync-materialized-api.md index 4ddf2f00..d9e9c75d 100644 --- a/docs/adr/0004-sync-materialized-api.md +++ b/docs/adr/0004-sync-materialized-api.md @@ -67,10 +67,11 @@ head, then the connection replaces its `ServerConnection` with this machine. Each SSE frame occupies one HTTP chunk. Closing the connection ends the stream without a terminal chunk. -Subscriptions are exempt from the request idle timeout. A push that would take -pending output past 64 KiB closes the connection. The expiry sweep also closes -connections whose pending output has made no socket write progress for over -12 seconds. This measures writes accepted by the socket, not reads by the peer. +Subscriptions are exempt from the request idle timeout. After attempting to drain +existing output, a push that would take pending output past 512 KiB closes the +connection. The expiry sweep also closes connections whose pending output has +made no socket write progress for over 30 seconds. This measures writes accepted +by the socket, not reads by the peer. The pump queues keep-alive comments every 15 seconds, including when no events are published. @@ -174,3 +175,26 @@ following; previously the gate held only network requests, and disk replay is chosen exactly when peers look comparable, so following could be announced during replay. Completion re-evaluates the target. Beacon-state's Status publications are unchanged. + +Amended 2026-09-15: delivery first attempts to drain the subscription's existing +output, then checks whether the new chunk, including its HTTP framing, fits +under the 512 KiB pending-output cap. This lets socket capacity that became +available since the last write attempt free room before rejecting the chunk, +without waiting for the next writable-readiness event. If the chunk still does +not fit, the connection closes. Otherwise, delivery queues it and attempts to +drain again. The pending response head also counts against the cap. + +Each drain writes until the buffer is empty or the socket returns `WouldBlock`; +unrecoverable write errors close the connection. Successful writes reset the +send deadline while output remains pending; draining it clears the deadline. +Progress means kernel acceptance, not peer consumption. + +Amended 2026-09-14: the send cap is 512 KiB. This accommodates one block's +128 `data_column_sidecar` events with 21 commitments each, estimated at +290 KiB, even when the socket accepts no bytes. Earlier pending events or +repeated publications can still exhaust the allowance. + +At 64 subscriptions, the pending-output allowance totals 32 MiB. Each +subscription reserves its buffer at construction and retains the allocation +until it closes. A full drain reuses the buffer from its beginning without +releasing it. Partial drains can advance through the entire allocation.