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 8ea1068..885f04a 100644 --- a/crates/turnloop-contract/src/executor_contract.rs +++ b/crates/turnloop-contract/src/executor_contract.rs @@ -394,3 +394,186 @@ 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:?}"); +} + +/// 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 ebc2b11..e16bfbc 100644 --- a/crates/turnloop-contract/src/lib.rs +++ b/crates/turnloop-contract/src/lib.rs @@ -239,6 +239,14 @@ mod native { refused_connect_once::(); } #[test] + fn connect_timeout() { + tcp_connect_timeout::(); + } + #[test] + fn single_connection_preset() { + single_connection_config::(); + } + #[test] fn liveness() { ref_unref::(); } @@ -433,6 +441,226 @@ 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); +} +/// 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/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/executor.rs b/crates/turnloop-contract/tests/executor.rs index f519798..d6d4eca 100644 --- a/crates/turnloop-contract/tests/executor.rs +++ b/crates/turnloop-contract/tests/executor.rs @@ -44,3 +44,16 @@ 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, + >(); +} + +#[test] +fn single_connection_presets_carry_a_request() { + turnloop_contract::executor_contract::single_connection_presets::( + ); +} 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..1923c9a 100644 --- a/crates/turnloop-contract/tests/windows.rs +++ b/crates/turnloop-contract/tests/windows.rs @@ -15,6 +15,8 @@ contract!( sustained_posts_idle_io, cancel_close_ordering, refused_connect_once, + tcp_connect_timeout, + single_connection_config, ref_unref, no_spin, quiet_deadline_accounting, 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/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/executor.rs b/crates/turnloop/src/executor.rs index f83532d..f7f9dcc 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)] @@ -64,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, @@ -147,7 +173,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 +188,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 +204,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 +217,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 +240,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 +251,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 +309,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 +368,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 +411,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 +449,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 +974,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 +993,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 +1178,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 +1319,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 +1460,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); } } } diff --git a/crates/turnloop/src/types.rs b/crates/turnloop/src/types.rs index 26aa30b..7f9b532 100644 --- a/crates/turnloop/src/types.rs +++ b/crates/turnloop/src/types.rs @@ -204,13 +204,93 @@ 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)] -/// 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.