Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
183 changes: 183 additions & 0 deletions crates/turnloop-contract/src/executor_contract.rs
Original file line number Diff line number Diff line change
Expand Up @@ -394,3 +394,186 @@ pub fn pending_writes_may_replace_the_caller_slice<B: backend::Backend>() {
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<B: backend::Backend>() {
let mut ex = LocalExecutor::<B>::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<B: backend::Backend>() {
let mut ex = LocalExecutor::<B>::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");
}
}
Loading