From 28bc32ddcd8c115bca56d3abc19c4d2cc701c7d6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:32:38 +0200 Subject: [PATCH 1/3] Hand a host's own completions back instead of dropping them LocalExecutor::turn drained the loop's output and passed every completion to Shared::dispatch, which returned early for any token without the executor's tag bit. A host that owns the loop and submits its own operations through LocalExecutor::driver therefore never saw their completions: no error, no counter, a socket that simply stopped delivering. That made the asynchronous protocol layers unusable from a host-driven loop. The issue offers two fixes; this takes the first, handing unrouted completions back, because it is the smaller change and composes with a host that already has a dispatch table keyed on the token, with no callback to re-enter the executor from inside its own turn. LocalExecutor::turn_into mirrors Driver::turn: it fills the host's buffer, claims the executor's completions out of it in place (Vec::extract_if, no allocation), and leaves every other completion there in delivery order while the host is free to submit more work. Plain turn keeps those completions in the executor's own buffer, readable through LocalExecutor::unclaimed until the next turn, so it no longer discards anything either. Executor-owned closes used Token(0), which would now surface as foreign completions. They use an internal token carrying the tag with generation zero, which reserve never issues, so it is claimed and matches no slot. Close futures still observe every Closed for their handle, whoever closed it. The contract test drives two real TCP pairs on one loop: the host reads and writes one with its own tokens while an executor task writes and reads the other. Both complete; then the host closes its handles and gets exactly its four Closed back and nothing else. With the old drop-everything dispatch the host read never completes (fails at the 5 s bound with nothing read), and with the internal token reverted to Token(0) the adapters' two closes leak to the host. Closes #45 --- .../src/executor_contract.rs | 103 ++++++++++++++++++ crates/turnloop-contract/tests/executor.rs | 7 ++ crates/turnloop/README.md | 5 +- crates/turnloop/src/executor.rs | 93 +++++++++++++--- 4 files changed, 189 insertions(+), 19 deletions(-) diff --git a/crates/turnloop-contract/src/executor_contract.rs b/crates/turnloop-contract/src/executor_contract.rs index 8ea1068..efe4c33 100644 --- a/crates/turnloop-contract/src/executor_contract.rs +++ b/crates/turnloop-contract/src/executor_contract.rs @@ -394,3 +394,106 @@ pub fn pending_writes_may_replace_the_caller_slice() { assert_eq!(&actual[..32], [7; 32]); assert_eq!(&actual[32..], b"new"); } + +/// A host that owns the loop submits its own read alongside an executor-driven +/// one, and both complete: the executor wakes its future and hands the host's +/// completions back through `unclaimed` instead of dropping them. Executor-owned +/// closes are claimed, so the host sees only the tokens it issued. +pub fn host_operations_share_the_loop() { + let mut ex = LocalExecutor::::new(Config::default()).expect("executor"); + let (host_server, host_client, host_conn) = crate::pair(&mut ex.driver()); + let (exec_server, exec_client, exec_conn) = crate::pair(&mut ex.driver()); + let h = ex.handle(); + let task = ex + .spawn_local(async move { + let mut writer = h.io(exec_client); + let mut reader = h.io(exec_conn); + let sent = b"executor bytes"; + let mut wrote = 0; + while wrote < sent.len() { + wrote += poll_fn(|cx| Pin::new(&mut writer).poll_write(cx, &sent[wrote..])) + .await + .expect("executor write"); + } + poll_fn(|cx| Pin::new(&mut writer).poll_flush(cx)) + .await + .expect("executor flush"); + let mut bytes = [0; 32]; + let mut read = 0; + while read < sent.len() { + let n = poll_fn(|cx| Pin::new(&mut reader).poll_read(cx, &mut bytes[read..])) + .await + .expect("executor read"); + assert!(n > 0, "early EOF"); + read += n; + } + bytes[..read].to_vec() + // Both adapters drop here and close their handles internally. + }) + .expect("spawn executor task"); + ex.driver() + .read(host_conn, ReadBuf::Pooled, Token(10)) + .expect("host read"); + ex.driver() + .write( + host_client, + WriteBuf::Owned(b"host bytes".to_vec()), + Token(11), + ) + .expect("host write"); + let mut host_read = Vec::new(); + let mut host_wrote = None; + let mut unexpected = Vec::new(); + let mut out = Completions::default(); + let until = ex.driver().now() + Duration::from_secs(5); + while !task.is_finished() || host_wrote.is_none() || host_read.len() < 10 { + assert!( + ex.driver().now() < until, + "host read {:?}, host write {host_wrote:?}, task finished {}", + host_read, + task.is_finished() + ); + ex.turn_into(Timeout::Until(until), &mut out) + .expect("executor turn"); + for c in out.drain() { + match (c.token, c.result) { + (Token(10), OpResult::Read { lease: Some(b), .. }) => { + host_read.extend_from_slice(b.as_slice()); + if host_read.len() < 10 { + ex.driver() + .read(host_conn, ReadBuf::Pooled, Token(10)) + .expect("continue host read"); + } + } + (Token(11), OpResult::Wrote(n)) => host_wrote = Some(n), + (token, result) => unexpected.push(format!("{token:?} {result:?}")), + } + } + } + assert_eq!(host_read, b"host bytes"); + assert_eq!(host_wrote, Some(10)); + let mut task = task; + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + let std::task::Poll::Ready(value) = Pin::new(&mut task).poll(&mut cx) else { + panic!("finished task must be ready") + }; + assert_eq!(value.expect("executor task"), b"executor bytes"); + // The adapters' own closes are the executor's; the host's closes, collected + // here through plain `turn`, are the only ones handed back. + for handle in [host_client, host_conn, host_server, exec_server] { + ex.driver().close(handle, Token(20)).expect("host close"); + } + let mut closed = 0; + while ex.alive() { + assert!(ex.driver().now() < until, "drain closes"); + ex.turn(Timeout::Until(until)).expect("close turn"); + for c in ex.unclaimed().drain() { + match (c.token, c.result) { + (Token(20), OpResult::Closed) => closed += 1, + (token, result) => unexpected.push(format!("{token:?} {result:?}")), + } + } + } + assert_eq!(closed, 4, "every host close is handed back exactly once"); + assert!(unexpected.is_empty(), "foreign completions: {unexpected:?}"); +} diff --git a/crates/turnloop-contract/tests/executor.rs b/crates/turnloop-contract/tests/executor.rs index f519798..bad647e 100644 --- a/crates/turnloop-contract/tests/executor.rs +++ b/crates/turnloop-contract/tests/executor.rs @@ -44,3 +44,10 @@ fn pending_write_can_replace_its_buffer() { turnloop::backend::Platform, >(); } + +#[test] +fn host_operations_complete_beside_executor_futures() { + turnloop_contract::executor_contract::host_operations_share_the_loop::< + turnloop::backend::Platform, + >(); +} diff --git a/crates/turnloop/README.md b/crates/turnloop/README.md index 8aee0a5..9dc1a8e 100644 --- a/crates/turnloop/README.md +++ b/crates/turnloop/README.md @@ -40,7 +40,10 @@ driver.turn(Timeout::Until(deadline), &mut completions)?; Enable `executor` for `LocalExecutor` and its futures-io adapters. Read/write staging and task tables are reserved at construction. Writes are buffered; flush or close before dropping an adapter to confirm underlying completion. -`AsyncIo::poll_shutdown` half-closes a stream without releasing its handle. The crate's +`AsyncIo::poll_shutdown` half-closes a stream without releasing its handle. A host +that also submits its own operations on the executor's loop turns it with +`LocalExecutor::turn_into`, which hands back every completion the executor did not +issue. The crate's rustdoc includes runnable loop and executor examples. See the [revision 2 handoff](https://github.com/PerryTS/turnloop/blob/main/docs/BACKEND_REVISION_2.md) for native ownership, process teardown, platform integration and contract details. diff --git a/crates/turnloop/src/executor.rs b/crates/turnloop/src/executor.rs index f83532d..c8b0b4f 100644 --- a/crates/turnloop/src/executor.rs +++ b/crates/turnloop/src/executor.rs @@ -44,6 +44,10 @@ use std::{ time::Duration, }; const TAG: u64 = 1 << 63; +/// Token for executor-owned closes that no future awaits by token. It carries the +/// executor's tag, so its completion is claimed rather than handed to the host, +/// and generation zero is never issued, so it can never match a live slot. +const INTERNAL: Token = Token(TAG); /// Fixed executor capacities, allocated at construction. #[derive(Clone, Copy, Debug)] @@ -147,7 +151,7 @@ impl Shared { | OpResult::HandleReceived { handle: conn }, ) = result { - let _ = self.driver.borrow_mut().close(conn, Token(0)); + let _ = self.driver.borrow_mut().close(conn, INTERNAL); } } fn abandon(&self, key: Key) { @@ -162,7 +166,7 @@ impl Shared { (slot.op, slot.timer.take(), slot.result.is_some()) }; if let Some(timer) = timer { - let _ = self.driver.borrow_mut().close(timer, Token(0)); + let _ = self.driver.borrow_mut().close(timer, INTERNAL); } else if let Some(op) = op { self.driver.borrow_mut().cancel(op); } @@ -178,7 +182,9 @@ impl Shared { } slot.result.take() } - fn dispatch(&self, completion: Completion) { + /// Whether this completion belongs to the executor. Every `Closed` is also + /// shown to the [`Close`] futures watching its handle, whoever closed it. + fn claims(&self, completion: &Completion) -> bool { if matches!(completion.result, OpResult::Closed) { for slot in self.slots.borrow_mut().iter_mut() { if slot.used && slot.closing.is_some() && slot.closing == completion.handle { @@ -189,9 +195,12 @@ impl Shared { } } } - if completion.token.0 & TAG == 0 { - return; - } + completion.token.0 & TAG != 0 + } + /// Route a claimed completion to its slot. A stale generation (an abandoned + /// operation's late result, or an [`INTERNAL`] close) is dropped here. + fn dispatch(&self, completion: Completion) { + debug_assert!(completion.token.0 & TAG != 0, "unclaimed completion"); let key = Key { index: completion.token.0 as u32 as usize, generation: ((completion.token.0 & !TAG) >> 32) as u32, @@ -209,7 +218,7 @@ impl Shared { (slot.waker.take(), slot.abandoned, timer) }; if let Some(timer) = timer { - let _ = self.driver.borrow_mut().close(timer, Token(0)); + let _ = self.driver.borrow_mut().close(timer, INTERNAL); } if abandoned { self.free(key); @@ -220,6 +229,13 @@ impl Shared { } /// A `!Send` executor whose owning host explicitly drives every turn. +/// +/// The host may share the loop with the executor: operations it submits itself +/// through [`LocalExecutor::driver`], with a token whose top bit is clear, are +/// never consumed by the executor. [`LocalExecutor::turn_into`] hands their +/// completions back in the host's own buffer, so a host with a dispatch table +/// keyed on the token routes them exactly as it would from [`Driver::turn`]; +/// [`LocalExecutor::turn`] keeps them in [`LocalExecutor::unclaimed`]. pub struct LocalExecutor { shared: Rc>, out: Completions, @@ -271,11 +287,23 @@ impl LocalExecutor { shared: self.shared.clone(), } } - /// Borrow the loop for synchronous resource creation/configuration. Tokens - /// with their top bit set are reserved; executor turn consumes completions. + /// Borrow the loop for synchronous resource creation/configuration, and for + /// the host's own operations. Tokens with their top bit set are reserved for + /// the executor; every other token's completion is handed back by + /// [`LocalExecutor::turn_into`] or kept in [`LocalExecutor::unclaimed`]. pub fn driver(&self) -> RefMut<'_, Driver> { self.shared.driver.borrow_mut() } + /// Completions from the most recent [`LocalExecutor::turn`] that the executor + /// did not issue: the host's own operations (any token with its top bit + /// clear) and posts, in delivery order. + /// + /// Like the output of [`Driver::turn`], the next turn clears this buffer, so + /// drain it after every turn; whatever is left, pooled-read leases included, + /// is dropped then. Executor-owned closes are claimed and never listed here. + pub fn unclaimed(&mut self) -> &mut Completions { + &mut self.out + } /// Spawn a local future. Dropping its JoinHandle requests cancellation. pub fn spawn_local(&self, future: F) -> Result> { self.handle().spawn_local(future) @@ -318,7 +346,34 @@ impl LocalExecutor { /// Drive one bounded loop turn and then one pass of runnable futures. Pending /// tasks that woke themselves keep this turn nonblocking, without being repolled /// repeatedly inside the same call. + /// + /// Completions for executor operations wake their futures; every other + /// completion stays in [`LocalExecutor::unclaimed`] for the host. The returned + /// [`TurnInfo::completions`] counts both. A host with its own dispatch table + /// uses [`LocalExecutor::turn_into`] instead. pub fn turn(&mut self, timeout: crate::Timeout) -> Result { + // An empty Vec does not allocate; the retained buffer is put back below. + let mut out = std::mem::replace( + &mut self.out, + Completions { + entries: Vec::new(), + }, + ); + let result = self.turn_into(timeout, &mut out); + self.out = out; + result + } + /// [`LocalExecutor::turn`] for a host that also submits its own operations: + /// like [`Driver::turn`], it clears `out` and fills it up to its capacity, and + /// then claims the executor's completions from it. What is left in `out`, in + /// delivery order, is every completion the executor did not issue (the + /// host's own tokens and posts), for the host to route while it is free to + /// submit more work through [`LocalExecutor::driver`]. + pub fn turn_into( + &mut self, + timeout: crate::Timeout, + out: &mut Completions, + ) -> Result { self.run_ready(); let ready = self .shared @@ -334,10 +389,11 @@ impl LocalExecutor { } else { timeout }, - &mut self.out, + out, )?; - for completion in self.out.drain() { - self.shared.dispatch(completion); + let shared = &self.shared; + for completion in out.entries.extract_if(.., |c| shared.claims(c)) { + shared.dispatch(completion); } self.run_ready(); Ok(info) @@ -371,7 +427,8 @@ impl Clone for ExecutorHandle { } impl ExecutorHandle { /// Borrow the driver for resource configuration. Release before awaiting. - /// Tokens with the top bit set belong to the executor. + /// Tokens with the top bit set belong to the executor; completions for any + /// other token are handed back to the host by [`LocalExecutor::turn_into`]. pub fn driver(&self) -> RefMut<'_, Driver> { self.shared.driver.borrow_mut() } @@ -895,7 +952,7 @@ impl AsyncWrite for AsyncIo { this.shared .driver .borrow_mut() - .close(this.handle, Token(0)) + .close(this.handle, INTERNAL) .map_err(io_error)?; this.closed = true; } @@ -914,7 +971,7 @@ impl Drop for AsyncIo { self.shared.abandon(key); } if !self.closed { - let _ = self.shared.driver.borrow_mut().close(self.handle, Token(0)); + let _ = self.shared.driver.borrow_mut().close(self.handle, INTERNAL); } } } @@ -1099,7 +1156,7 @@ impl Future for Close { if driver.is_closing(handle) { Ok(()) } else { - driver.close(handle, Token(0)) + driver.close(handle, INTERNAL) } }; if let Err(e) = result { @@ -1240,7 +1297,7 @@ impl Drop for Connect { .shared .driver .borrow_mut() - .close(handle, Token(0)); + .close(handle, INTERNAL); } if let Some(key) = self.key.take() { // Connect has no borrowed buffer. Generation tags reject its later @@ -1381,7 +1438,7 @@ impl Drop for Fetch { self.executor.shared.abandon(key); } if let Some(handle) = self.handle.take() { - let _ = self.executor.driver().close(handle, Token(0)); + let _ = self.executor.driver().close(handle, INTERNAL); } } } From daa48db202e287ce9591ecf779f7b4cc8581ddd8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:34:06 +0200 Subject: [PATCH 2/3] Give TcpOpts a connect timeout and document its boundary TcpOpts had only nodelay, so every HTTP client had to build its own connect deadline out of a timer plus close, and nothing said whether the single field was an oversight or a deliberate limit next to Loop::set_option. The issue allows either extending TcpOpts or documenting the boundary; this does both, adding the one field the backends can honour portably. The loop already enforces pipe_connect_until's deadline in the driver, on its own clock: expiry cancels the pending connect and reports TimedOut only after the backend acknowledges the cancellation. TcpOpts::connect_timeout arms that same queue from tcp_connect, so it behaves identically on kqueue, epoll, IOCP and WASI with no backend change. The deadline is retired with the operation, covers the attempt only, and a zero budget is InvalidInput. A local bind address is not added: it would need a bind step in every backend's socket creation, which is a larger change than this issue. The TcpOpts docs now say what is here, what goes through set_option on the handle tcp_connect returns (and that a receive buffer set that way cannot change the negotiated window scale), and that a pre-connect local bind is not offered. Breaking: TcpOpts gains a public field, so struct literals need ..TcpOpts::default(). The two contract uses are updated. The contract test, wired for the Unix, Windows and WASI lanes, checks that a zero budget is refused with nothing retained, that a 5 s budget to a live listener connects and leaves no deadline behind, and that a 1 ms budget that has passed before the first turn completes exactly once with TimedOut and then closes cleanly. Without the deadline insertion the budget is never armed and the late attempt reports Connected. Closes #82 --- crates/turnloop-contract/src/lib.rs | 93 +++++++++++++++++++++++ crates/turnloop-contract/src/sockopts.rs | 10 ++- crates/turnloop-contract/tests/wasi.rs | 4 + crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/driver.rs | 27 +++++-- crates/turnloop/src/types.rs | 31 +++++++- 6 files changed, 154 insertions(+), 12 deletions(-) diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index ebc2b11..05e960d 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -239,6 +239,10 @@ mod native { refused_connect_once::(); } #[test] + fn connect_timeout() { + tcp_connect_timeout::(); + } + #[test] fn liveness() { ref_unref::(); } @@ -433,6 +437,95 @@ pub fn refused_connect_once() { l.close(conn, Token(8)).expect("close"); assert_eq!(count, 1); } +/// `TcpOpts::connect_timeout` is enforced by the loop: an attempt still pending +/// at its deadline completes once with `TimedOut`, an attempt that connects in +/// time is unaffected and leaves no deadline behind, and a zero budget is +/// refused before any socket exists. +pub fn tcp_connect_timeout() { + let mut l = Driver::::new(Config::default()).expect("loop"); + let zero = TcpOpts { + connect_timeout: Some(Duration::ZERO), + ..TcpOpts::default() + }; + assert_eq!( + l.tcp_connect(localhost(), &zero, Token(1)) + .expect_err("zero budget") + .kind, + ErrorKind::InvalidInput + ); + assert!(!l.alive(), "a refused connect retains nothing"); + + let server = l + .tcp_listen(localhost(), &ListenOpts::default()) + .expect("listen"); + let addr = l.local_addr(server).expect("addr"); + let mut out = Completions::default(); + + // In time: Connected, and the deadline is retired with the operation. + let generous = TcpOpts { + connect_timeout: Some(Duration::from_secs(5)), + ..TcpOpts::default() + }; + let on_time = l.tcp_connect(addr, &generous, Token(2)).expect("connect"); + assert!(l.next_deadline().is_some(), "the budget is armed"); + let until = l.now() + Duration::from_secs(3); + let mut results = Vec::new(); + while results.is_empty() { + assert!(l.now() < until, "connect in time"); + l.turn(Timeout::Until(until), &mut out).expect("turn"); + results.extend(out.drain().map(|c| (c.token, c.terminal, c.result))); + } + assert!( + matches!(results[..], [(Token(2), true, OpResult::Connected)]), + "{results:?}" + ); + assert_eq!(l.next_deadline(), None, "no deadline outlives its connect"); + + // Expired before the loop ever looked: exactly one TimedOut, then a clean close. + let tight = TcpOpts { + connect_timeout: Some(Duration::from_millis(1)), + ..TcpOpts::default() + }; + let late = l.tcp_connect(addr, &tight, Token(3)).expect("connect"); + thread::sleep(Duration::from_millis(20)); + results.clear(); + while results.is_empty() { + assert!(l.now() < until, "timeout delivered"); + l.turn(Timeout::Until(until), &mut out).expect("turn"); + results.extend(out.drain().map(|c| (c.token, c.terminal, c.result))); + } + assert!( + matches!( + results[..], + [( + Token(3), + true, + OpResult::Err(Error { + kind: ErrorKind::TimedOut, + .. + }) + )] + ), + "{results:?}" + ); + for _ in 0..3 { + l.turn(Timeout::Now, &mut out).expect("no duplicate"); + assert!(out.is_empty(), "the timeout is reported once"); + } + for (h, token) in [(late, 4), (on_time, 5), (server, 6)] { + l.close(h, Token(token)).expect("close"); + } + let mut closed = 0; + while l.alive() { + assert!(l.now() < until, "closes"); + l.turn(Timeout::Until(until), &mut out).expect("turn"); + for c in out.drain() { + assert!(matches!(c.result, OpResult::Closed), "{c:?}"); + closed += 1; + } + } + assert_eq!(closed, 3); +} pub fn ref_unref() { let mut l = Driver::::new(Config::default()).expect("loop"); assert!(!l.alive()); diff --git a/crates/turnloop-contract/src/sockopts.rs b/crates/turnloop-contract/src/sockopts.rs index 7b353fd..8587900 100644 --- a/crates/turnloop-contract/src/sockopts.rs +++ b/crates/turnloop-contract/src/sockopts.rs @@ -426,7 +426,10 @@ pub fn nodelay_round_trip_and_accept_default() { let eager = l .tcp_connect( (std::net::Ipv4Addr::LOCALHOST, 9).into(), - &TcpOpts { nodelay: true }, + &TcpOpts { + nodelay: true, + ..TcpOpts::default() + }, Token(30), ) .expect("connecting socket"); @@ -634,7 +637,10 @@ pub fn unsupported_connect_nodelay_is_rejected() { assert_eq!( l.tcp_connect( (std::net::Ipv4Addr::LOCALHOST, 9).into(), - &TcpOpts { nodelay: true }, + &TcpOpts { + nodelay: true, + ..TcpOpts::default() + }, Token(1), ) .expect_err("no Nagle control here") diff --git a/crates/turnloop-contract/tests/wasi.rs b/crates/turnloop-contract/tests/wasi.rs index 851b0de..0f91c77 100644 --- a/crates/turnloop-contract/tests/wasi.rs +++ b/crates/turnloop-contract/tests/wasi.rs @@ -106,6 +106,10 @@ fn refused_connect() { contract::refused_connect_once::(); } #[test] +fn connect_timeout() { + contract::tcp_connect_timeout::(); +} +#[test] fn ref_unref() { contract::ref_unref::(); } diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index 7529ab3..91fa4d0 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -15,6 +15,7 @@ contract!( sustained_posts_idle_io, cancel_close_ordering, refused_connect_once, + tcp_connect_timeout, ref_unref, no_spin, quiet_deadline_accounting, diff --git a/crates/turnloop/src/driver.rs b/crates/turnloop/src/driver.rs index 45303f4..7290b9d 100644 --- a/crates/turnloop/src/driver.rs +++ b/crates/turnloop/src/driver.rs @@ -750,21 +750,34 @@ impl Driver { pub fn udp_bind(&mut self, addr: SocketAddr, opts: &UdpOpts) -> Result { self.open(Open::Udp { addr, opts: *opts }) } - /// Create a TCP socket and submit its connection operation with the supplied token. + /// Create a TCP socket and submit its connection operation with the supplied + /// token. [`TcpOpts::connect_timeout`] is enforced exactly like the deadline + /// of [`Driver::pipe_connect_until`]. pub fn tcp_connect( &mut self, addr: SocketAddr, opts: &TcpOpts, token: Token, ) -> Result { + if opts.connect_timeout == Some(Duration::ZERO) { + return Err(Error::new(ErrorKind::InvalidInput)); + } let h = self.open(Open::Tcp { addr, opts: *opts })?; - if let Err(e) = self.submit(h, Operation::Connect, token) { - self.backend.release(h); - self.handles.remove(h.key); - self.refs -= 1; - return Err(e); + match self.submit(h, Operation::Connect, token) { + Ok(op) => { + if let Some(at) = opts.connect_timeout.and_then(|t| self.now().checked_add(t)) { + self.connect_deadlines.insert(op.key, at); + self.backend.deadline_changed(self.next_deadline()); + } + Ok(h) + } + Err(e) => { + self.backend.release(h); + self.handles.remove(h.key); + self.refs -= 1; + Err(e) + } } - Ok(h) } /// Return the local IP endpoint of a bound socket. pub fn local_addr(&self, h: Handle) -> Result { diff --git a/crates/turnloop/src/types.rs b/crates/turnloop/src/types.rs index 26aa30b..428e4a9 100644 --- a/crates/turnloop/src/types.rs +++ b/crates/turnloop/src/types.rs @@ -205,12 +205,37 @@ impl Default for Config { } } #[derive(Clone, Copy, Debug, Default)] -/// TCP connection options applied when creating the socket. Anything a host -/// needs to change later goes through `Loop::set_option` and -/// [`SocketOption`] instead. +/// Connect-time options for [`Loop::tcp_connect`]. +/// +/// This struct holds what must be decided before the connection attempt starts +/// and what the loop itself enforces. Every socket option that can be changed +/// on a live socket goes through [`Loop::set_option`] and [`SocketOption`] +/// instead, and can be applied to the handle `tcp_connect` returns, before the +/// `Connected` completion arrives, wherever the backend supports that option: +/// keep-alive, linger, TTL and buffer sizes. A +/// receive buffer set that way is applied after the handshake has started, so it +/// cannot change the window scale the peers negotiate. +/// +/// Not offered at all, on any backend: binding the connecting socket to a +/// chosen local address or port before it connects. `IPV6_V6ONLY` is not needed +/// here: a connecting socket's family is the destination's. +/// +/// [`Loop::tcp_connect`]: crate::Driver::tcp_connect +/// [`Loop::set_option`]: crate::Driver::set_option pub struct TcpOpts { /// Disable the TCP Nagle algorithm for latency-sensitive small writes. pub nodelay: bool, + /// Give up on the connection attempt after this long. + /// + /// Enforced by the loop on its own clock, identically on every backend: when + /// the deadline passes, the pending connect is cancelled and completes once, + /// with `Err` of kind `TimedOut`, only after the backend acknowledges the + /// cancellation. It covers the attempt only, never later stream I/O, and it + /// cannot outlast the operating system's own connect timeout, which still + /// applies. The handle stays open after a timeout; close it as after any + /// other failed connect. `None` (the default) leaves the attempt to the OS. + /// A zero duration is `InvalidInput`. + pub connect_timeout: Option, } #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] /// What a second bind of the same address is allowed to do. From 716360a1fa3b3173d9d5500461298a8fe8854c75 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ralph=20K=C3=BCpper?= Date: Tue, 22 Sep 2026 19:41:16 +0200 Subject: [PATCH 3/3] Give a one-connection caller a documented Config preset Config's defaults size every table for a server. A client that builds one loop, drives a single outbound connection one request at a time and drops the loop had to trim max_handles and max_operations by hand with no way to know which values were safe, and every such caller guessed independently. Config::single_connection() is that preset: 8 handles, 16 operations, 16 events per turn, two pooled buffers of the default 16 KiB and a post queue of 16. The blocking pool keeps its default because it is process-wide and must match every other loop's. The docs give the floor it covers as a table: the worst moment for this caller is a reconnect (redirect or retry) started before the previous attempt's completions are turned out, which needs 4 handles and 10 operations; the preset leaves headroom above that. They also say which fields fail loudly with ResourceLimit when undersized and which only pace the loop. ExecutorConfig::single_connection() is the executor half: 4 tasks and 8 operation slots, which fits inside the loop preset's 16 operations. Two contract tests exercise the presets. single_connection_config drives the loop directly through that worst moment: a lookup, a connect with a deadline, a read, write and shutdown, then a redirect that closes the socket and its timer, arms a replacement timer, resolves again and reconnects before any of the first attempt's completions are turned out. Nothing is refused with ResourceLimit and the second exchange echoes back byte for byte. It is wired for the Unix and Windows lanes; it needs a threaded std listener as its peer. single_connection_presets carries one request end to end through LocalExecutor under both presets, after abandoning a first adapter with a read pending. The README points one-connection callers at the preset. Closes #81 --- README.md | 4 + .../src/executor_contract.rs | 80 +++++++++++ crates/turnloop-contract/src/lib.rs | 135 ++++++++++++++++++ crates/turnloop-contract/tests/executor.rs | 6 + crates/turnloop-contract/tests/windows.rs | 1 + crates/turnloop/src/executor.rs | 22 +++ crates/turnloop/src/types.rs | 55 +++++++ 7 files changed, 303 insertions(+) diff --git a/README.md b/README.md index 73e95ce..38cd5cd 100644 --- a/README.md +++ b/README.md @@ -40,6 +40,10 @@ driver.turn(Timeout::Until(deadline), &mut completions)?; # Ok::<(), turnloop::Error>(()) ``` +A short-lived loop that drives one outbound connection at a time can start +from `Config::single_connection()` (and `ExecutorConfig::single_connection()`) +instead of trimming the server-sized defaults by hand. + Use the pinned nightly for dependency resolution under the seven-day publication soak. The workspace also builds with stable Rust 1.97.1. diff --git a/crates/turnloop-contract/src/executor_contract.rs b/crates/turnloop-contract/src/executor_contract.rs index efe4c33..885f04a 100644 --- a/crates/turnloop-contract/src/executor_contract.rs +++ b/crates/turnloop-contract/src/executor_contract.rs @@ -497,3 +497,83 @@ pub fn host_operations_share_the_loop() { assert_eq!(closed, 4, "every host close is handed back exactly once"); assert!(unexpected.is_empty(), "foreign completions: {unexpected:?}"); } + +/// The executor presets carry one request end to end: a lookup, a connect with +/// its deadline, a write, a half-close and a read to EOF, all under a request +/// timeout, then a second request on a fresh connection after the first +/// adapter was dropped mid-read. +pub fn single_connection_presets() { + let mut ex = LocalExecutor::::with_config( + Config::single_connection(), + turnloop::ExecutorConfig::single_connection(), + ) + .expect("executor"); + let (addr, peer) = crate::echo_peer(2); + let h = ex.handle(); + let task = ex + .spawn_local(async move { + let opts = TcpOpts { + connect_timeout: Some(Duration::from_secs(5)), + ..TcpOpts::default() + }; + let lookup = DnsRequest { + host: "localhost".into(), + port: addr.port(), + }; + assert!(!h.resolve(lookup).await.expect("resolve").is_empty()); + // Abandoned with a read pending: its slots return on cancellation. + let mut first = h.connect(addr, opts).await.expect("connect"); + let mut bytes = [0; 16]; + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + assert!( + Pin::new(&mut first) + .poll_read(&mut cx, &mut bytes) + .is_pending() + ); + drop(first); + let request = async { + let mut stream = h.connect(addr, opts).await.expect("reconnect"); + let sent = b"one request"; + let mut wrote = 0; + while wrote < sent.len() { + wrote += poll_fn(|cx| Pin::new(&mut stream).poll_write(cx, &sent[wrote..])) + .await + .expect("write"); + } + poll_fn(|cx| stream.poll_shutdown(cx)) + .await + .expect("shutdown"); + let mut reply = Vec::new(); + loop { + let n = poll_fn(|cx| Pin::new(&mut stream).poll_read(cx, &mut bytes)) + .await + .expect("read"); + if n == 0 { + break reply; + } + reply.extend_from_slice(&bytes[..n]); + } + }; + h.timeout(Duration::from_secs(5), request) + .await + .expect("request deadline") + }) + .expect("spawn"); + let until = ex.driver().now() + Duration::from_secs(10); + while !task.is_finished() { + assert!(ex.driver().now() < until, "request stalled"); + ex.turn(Timeout::Until(until)).expect("turn"); + } + let mut task = task; + let mut cx = std::task::Context::from_waker(std::task::Waker::noop()); + let std::task::Poll::Ready(reply) = Pin::new(&mut task).poll(&mut cx) else { + panic!("finished task must be ready") + }; + assert_eq!(reply.expect("task"), b"one request"); + let requests = peer.join().expect("peer"); + assert_eq!(requests[1], b"one request"); + while ex.alive() { + assert!(ex.driver().now() < until, "teardown"); + ex.turn(Timeout::Until(until)).expect("turn"); + } +} diff --git a/crates/turnloop-contract/src/lib.rs b/crates/turnloop-contract/src/lib.rs index 05e960d..e16bfbc 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -243,6 +243,10 @@ mod native { tcp_connect_timeout::(); } #[test] + fn single_connection_preset() { + single_connection_config::(); + } + #[test] fn liveness() { ref_unref::(); } @@ -526,6 +530,137 @@ pub fn tcp_connect_timeout() { } assert_eq!(closed, 3); } +/// A peer that reads each of `connections` requests to EOF and echoes it back. +/// Echo failures are ignored: a client may close without reading its reply. +pub(crate) fn echo_peer(connections: usize) -> (SocketAddr, thread::JoinHandle>>) { + use std::io::{Read, Write}; + let listener = std::net::TcpListener::bind(localhost()).expect("peer listen"); + let addr = listener.local_addr().expect("peer addr"); + let peer = thread::spawn(move || { + (0..connections) + .map(|_| { + let (mut stream, _) = listener.accept().expect("peer accept"); + stream + .set_read_timeout(Some(Duration::from_secs(5))) + .expect("peer timeout"); + let mut request = Vec::new(); + stream.read_to_end(&mut request).expect("peer read"); + let _ = stream.write_all(&request); + request + }) + .collect() + }); + (addr, peer) +} +/// Turn until every token in `wanted` has a completion, recording each result. +fn turn_for( + l: &mut Driver, + out: &mut Completions, + until: turnloop::Instant, + wanted: &[u64], + seen: &mut Vec<(u64, String)>, +) { + while !wanted.iter().all(|w| seen.iter().any(|(t, _)| t == w)) { + assert!(l.now() < until, "waiting for {wanted:?}, saw {seen:?}"); + l.turn(Timeout::Until(until), out).expect("turn"); + seen.extend(out.drain().map(|c| (c.token.0, format!("{:?}", c.result)))); + } +} +/// `Config::single_connection` carries a client through its worst moment: a +/// redirect started while the first attempt's socket, its cancelled read, write +/// and shutdown and its request timer are still closing, with a second DNS +/// lookup and connect deadlines armed. Every creation and submission succeeds and the second +/// exchange completes byte for byte. +pub fn single_connection_config() { + let mut l = Driver::::new(Config::single_connection()).expect("loop"); + let (peer_addr, peer) = echo_peer(2); + let mut out = Completions::default(); + let mut seen = Vec::new(); + let until = l.now() + Duration::from_secs(10); + let opts = TcpOpts { + connect_timeout: Some(Duration::from_secs(5)), + ..TcpOpts::default() + }; + + let first_timer = l + .timer(l.now() + Duration::from_secs(30), None, Token(1)) + .expect("request timer"); + l.resolve( + DnsRequest { + host: "localhost".into(), + port: peer_addr.port(), + }, + Token(2), + ) + .expect("resolve"); + turn_for(&mut l, &mut out, until, &[2], &mut seen); + assert!(seen[0].1.starts_with("Resolved("), "{seen:?}"); + + let first = l.tcp_connect(peer_addr, &opts, Token(3)).expect("connect"); + turn_for(&mut l, &mut out, until, &[3], &mut seen); + l.read(first, ReadBuf::Pooled, Token(4)).expect("read"); + l.write(first, WriteBuf::Owned(b"first".to_vec()), Token(5)) + .expect("write"); + l.shutdown(first, Token(6)).expect("shutdown"); + + // A redirect, before the first attempt's read, write and shutdown have been + // turned out: abandon it, replace the timer, look up the new name and connect. + l.close(first, Token(7)).expect("close first socket"); + l.close(first_timer, Token(8)).expect("close first timer"); + let second_timer = l + .timer(l.now() + Duration::from_secs(30), None, Token(9)) + .expect("replacement timer"); + l.resolve( + DnsRequest { + host: "localhost".into(), + port: peer_addr.port(), + }, + Token(16), + ) + .expect("second lookup"); + let second = l + .tcp_connect(peer_addr, &opts, Token(10)) + .expect("reconnect"); + turn_for(&mut l, &mut out, until, &[4, 5, 6, 7, 8, 10, 16], &mut seen); + l.read(second, ReadBuf::Pooled, Token(11)).expect("read"); + l.write(second, WriteBuf::Owned(b"second".to_vec()), Token(12)) + .expect("write"); + l.shutdown(second, Token(13)).expect("shutdown"); + let mut reply = Vec::new(); + let mut eof = false; + while !eof { + assert!(l.now() < until, "reply stalled at {reply:?}"); + l.turn(Timeout::Until(until), &mut out).expect("turn"); + for c in out.drain() { + match c.result { + OpResult::Read { lease: Some(b), .. } if c.token == Token(11) => { + reply.extend_from_slice(b.as_slice()); + drop(b); + l.read(second, ReadBuf::Pooled, Token(11)).expect("read on"); + } + OpResult::Eof => eof = true, + ref result => seen.push((c.token.0, format!("{result:?}"))), + } + } + } + assert_eq!(reply, b"second"); + l.close(second, Token(14)).expect("close"); + l.close(second_timer, Token(15)).expect("close"); + while l.alive() { + assert!(l.now() < until, "teardown"); + l.turn(Timeout::Until(until), &mut out).expect("turn"); + seen.extend(out.drain().map(|c| (c.token.0, format!("{:?}", c.result)))); + } + assert!( + !seen.iter().any(|(_, r)| r.contains("ResourceLimit")), + "{seen:?}" + ); + // The abandoned request may have reached the peer in full, in part or not + // at all; only the second one's bytes are specified. + let requests = peer.join().expect("peer"); + assert_eq!(requests.len(), 2); + assert_eq!(requests[1], b"second"); +} pub fn ref_unref() { let mut l = Driver::::new(Config::default()).expect("loop"); assert!(!l.alive()); diff --git a/crates/turnloop-contract/tests/executor.rs b/crates/turnloop-contract/tests/executor.rs index bad647e..d6d4eca 100644 --- a/crates/turnloop-contract/tests/executor.rs +++ b/crates/turnloop-contract/tests/executor.rs @@ -51,3 +51,9 @@ fn host_operations_complete_beside_executor_futures() { turnloop::backend::Platform, >(); } + +#[test] +fn single_connection_presets_carry_a_request() { + turnloop_contract::executor_contract::single_connection_presets::( + ); +} diff --git a/crates/turnloop-contract/tests/windows.rs b/crates/turnloop-contract/tests/windows.rs index 91fa4d0..1923c9a 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -16,6 +16,7 @@ contract!( cancel_close_ordering, refused_connect_once, tcp_connect_timeout, + single_connection_config, ref_unref, no_spin, quiet_deadline_accounting, diff --git a/crates/turnloop/src/executor.rs b/crates/turnloop/src/executor.rs index c8b0b4f..f7f9dcc 100644 --- a/crates/turnloop/src/executor.rs +++ b/crates/turnloop/src/executor.rs @@ -68,6 +68,28 @@ impl Default for ExecutorConfig { } } } +impl ExecutorConfig { + /// The executor half of [`Config::single_connection`]: one connection's + /// futures, one request at a time. + /// + /// Each pending executor future holds one slot, and an abandoned one keeps + /// it until its cancellation is acknowledged. The connection's resolve or + /// connect, then its read, write and shutdown, a request deadline's + /// [`Sleep`] and one [`Close`] need 6 at most, and the preset keeps 8, which + /// is 128 KiB of staging at the default 16 KiB `buffer_size`. Every slot + /// that submits I/O also holds one of the loop's operations, so 8 stays + /// within `Config::single_connection`'s 16 with room for the host's own. + /// Four tasks cover the request, a watchdog beside it and a background task + /// a protocol layer may spawn for its connection. A future or task past + /// either bound fails with `ResourceLimit`; nothing is dropped. + pub fn single_connection() -> Self { + Self { + tasks: 4, + operations: 8, + buffer_size: 16 * 1024, + } + } +} struct IoSlot { generation: u32, used: bool, diff --git a/crates/turnloop/src/types.rs b/crates/turnloop/src/types.rs index 428e4a9..7f9b532 100644 --- a/crates/turnloop/src/types.rs +++ b/crates/turnloop/src/types.rs @@ -204,6 +204,61 @@ impl Default for Config { } } } +impl Config { + /// Capacities for a short-lived loop that drives one outbound connection, + /// one request at a time: a client that resolves a name, connects, writes a + /// request, reads the response and drops the loop. + /// + /// [`Config::default`] sizes every table for a server. This preset keeps the + /// defaults only where they are not per-loop: the pooled buffer size (one + /// read still takes a full TLS record) and [`Config::blocking_pool`], which is + /// process-wide and must match every other loop's, so shrinking it here would + /// make the first DNS lookup fail with `InvalidInput` in a process that also + /// runs default loops. + /// + /// # The floor, and what it covers + /// + /// A handle's slot is held from creation until its `Closed` completion is + /// delivered, and an operation's until its terminal completion is. The worst + /// moment for this caller is a reconnect (a redirect or a retry) that starts + /// before the previous attempt's completions have been turned out: + /// + /// | Holding a slot at that moment | handles | operations | + /// |---|---|---| + /// | the old socket, closing, and its cancelled read, write and shutdown | 1 | 3 | + /// | the new socket: its connect, then read, write and shutdown | 1 | 3 | + /// | a request deadline timer, and the one it replaces | 2 | 2 | + /// | a DNS lookup, and one other blocking-pool job | 0 | 2 | + /// | **needed** | **4** | **10** | + /// | **this preset** | **8** | **16** | + /// + /// `max_handles` and `max_operations` are the only fields a caller can size + /// too small by guessing, and they fail loudly: a creation or submission past + /// either ceiling is refused with `ResourceLimit`, nothing is dropped. The + /// other fields only pace the loop: `events_per_turn` bounds how many native + /// events one turn collects, `pooled_buffers` how many + /// [`ReadBuf::Pooled`](crate::ReadBuf::Pooled) reads can hold data at once + /// (a read waits for a free buffer rather than failing; two let the next read + /// fill while the host still holds the last one), and `post_capacity` how + /// many cross-thread posts can queue before `Poster::post` refuses one. + /// + /// Anything beyond the table, such as a second concurrent connection, a + /// listener, child processes, signal subscriptions or filesystem requests, + /// needs its own handles and operations on top. With the `executor` + /// feature's `LocalExecutor`, pair this with + /// `ExecutorConfig::single_connection`, whose operations fit inside these. + pub fn single_connection() -> Self { + Self { + max_handles: 8, + max_operations: 16, + events_per_turn: 16, + pooled_buffers: 2, + pooled_buffer_size: 16 * 1024, + post_capacity: 16, + blocking_pool: crate::PoolConfig::default(), + } + } +} #[derive(Clone, Copy, Debug, Default)] /// Connect-time options for [`Loop::tcp_connect`]. ///