Skip to content

Commit b520f3d

Browse files
Implement coroutine tasks and nested future waits
1 parent 64651b0 commit b520f3d

16 files changed

Lines changed: 1812 additions & 416 deletions

File tree

‎agent-context.md‎

Lines changed: 11 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,9 @@ layout.
77
## Workspace
88

99
- `socketry` is the public facade and re-exports the common concurrency API.
10-
- `socketry-concurrent` owns the stackful fiber implementation and the
11-
cooperative scheduler.
10+
- `socketry-concurrent` owns the coroutine implementation and the cooperative
11+
future scheduler. `Scheduler::spawn` accepts futures; `wait` lets ordinary
12+
synchronous functions suspend on futures using the current task's stack.
1213
- Each workspace package is published independently to crates.io.
1314

1415
## Implementation boundaries
@@ -17,7 +18,13 @@ layout.
1718
revision and attribution are recorded in
1819
`crates/concurrent/vendor/cruby/upstream.md` and
1920
`crates/concurrent/vendor/cruby/license.md`.
20-
- Suspended fibers are thread-affine. Do not migrate their stacks between
21+
- Suspended tasks are thread-affine. Do not migrate their stacks between
2122
worker threads.
2223
- Future wakers may enqueue work from another thread; the owning scheduler
23-
resumes its fibers.
24+
resumes its tasks.
25+
- A nested `wait` preserves an in-progress `Future::poll` call. Resume that
26+
stack before polling the outer future again. Preserve outer wakeups consumed
27+
while nested waits are active.
28+
- Both nested waits and ordinary future suspension currently stay on the
29+
owning thread. Migration between completed polls is a possible future
30+
extension and is not implemented.

‎crates/concurrent/Cargo.toml‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,11 @@
22
name = "socketry-concurrent"
33
version.workspace = true
44
edition.workspace = true
5-
description = "Stackful fibers and a small cooperative scheduler"
5+
description = "A cooperative future scheduler with nested stackful waits"
66
license.workspace = true
77
repository.workspace = true
88
readme = "readme.md"
9-
include = ["Cargo.toml", "build.rs", "license.md", "readme.md", "src/**", "vendor/**"]
9+
include = ["Cargo.toml", "build.rs", "license.md", "readme.md", "src/**", "examples/**", "vendor/**"]
1010

1111
[lib]
1212
name = "socketry_concurrent"
Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,55 @@
1+
use socketry_concurrent::{Scheduler, wait};
2+
use std::cell::RefCell;
3+
use std::future::poll_fn;
4+
use std::rc::Rc;
5+
use std::task::{Poll, Waker};
6+
7+
#[derive(Default)]
8+
struct Message {
9+
value: Option<String>,
10+
reader: Option<Waker>,
11+
}
12+
13+
// An ordinary synchronous function can wait for a future at any call depth.
14+
fn read_message(message: &RefCell<Message>) -> String {
15+
wait(poll_fn(|context| {
16+
let mut message = message.borrow_mut();
17+
if let Some(value) = message.value.take() {
18+
Poll::Ready(value)
19+
} else {
20+
message.reader = Some(context.waker().clone());
21+
Poll::Pending
22+
}
23+
}))
24+
}
25+
26+
fn process_message(message: &RefCell<Message>) {
27+
println!("consumer: waiting");
28+
let value = read_message(message);
29+
println!("consumer: received {value}");
30+
}
31+
32+
fn main() -> std::io::Result<()> {
33+
let mut scheduler = Scheduler::new(256 * 1024);
34+
let message = Rc::new(RefCell::new(Message::default()));
35+
36+
let consumer_message = Rc::clone(&message);
37+
scheduler.spawn(async move {
38+
process_message(&consumer_message);
39+
})?;
40+
41+
scheduler.spawn(async move {
42+
println!("producer: sending");
43+
let reader = {
44+
let mut message = message.borrow_mut();
45+
message.value = Some(String::from("hello"));
46+
message.reader.take()
47+
};
48+
if let Some(reader) = reader {
49+
reader.wake();
50+
}
51+
})?;
52+
53+
scheduler.run();
54+
Ok(())
55+
}

‎crates/concurrent/readme.md‎

Lines changed: 62 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -1,41 +1,81 @@
11
# socketry-concurrent
22

3-
Stackful fibers and cooperative scheduling for Socketry's Rust packages. Use
3+
Futures with nested stackful waits for Socketry's Rust packages. Use
44
the `socketry` package for the one-stop public entry point, or depend on this
55
package directly to use the concurrency implementation on its own.
66

7-
The API starts with three pieces:
7+
The library provides three building blocks:
88

9-
- Stack reserves memory with inaccessible guard pages at both ends.
10-
- Fiber owns a stack and switches between its caller and a closure.
11-
- Pool keeps multiple stacks available for reuse.
9+
- `Stack` reserves memory with inaccessible guard pages at both ends.
10+
- `Pool` keeps multiple stacks available for reuse.
11+
- `Scheduler` accepts futures and polls each task on its own coroutine stack.
1212

13-
Scheduler is an optional, single-threaded executor. It can run Rust futures and
14-
also exposes block_current and task unblock operations for stackful code.
15-
Future wakers may run on any thread; the fiber itself always resumes on the
16-
thread that owns its stack. Arbitrary Rust locals on a suspended stack are not
17-
tracked by the type system, so moving a suspended fiber between workers is not
18-
safe. The scheduler does not migrate fibers.
13+
`wait(future)` lets an ordinary function wait for an asynchronous result. When
14+
the future is pending, the task's stack is suspended and the executor can run
15+
other tasks. The function returns the future's output after it completes. Calls
16+
can be nested, including inside another future's `poll`, without making the
17+
calling functions async.
18+
19+
`Task::current` provides a thread-affine reference for blocking and task-to-task
20+
transfer. `Scheduler::spawn` returns a separate, thread-safe `TaskHandle` for
21+
waking the task. Future wakers may run on any thread; each task always resumes
22+
on the thread that owns its stack.
1923

2024
## Example
2125

22-
use socketry_concurrent::Scheduler;
26+
use socketry_concurrent::{Scheduler, wait};
27+
use std::future::poll_fn;
28+
use std::task::Poll;
29+
30+
fn answer() -> usize {
31+
let mut first_poll = true;
32+
wait(poll_fn(|context| {
33+
if first_poll {
34+
first_poll = false;
35+
context.waker().wake_by_ref();
36+
Poll::Pending
37+
} else {
38+
Poll::Ready(42)
39+
}
40+
}))
41+
}
2342

2443
fn main() -> std::io::Result<()> {
2544
let mut scheduler = Scheduler::new(256 * 1024);
2645
scheduler.spawn(async {
27-
// Await ordinary Rust futures here.
46+
assert_eq!(answer(), 42);
2847
})?;
2948
scheduler.run();
3049
Ok(())
3150
}
3251

33-
Fiber::yield_now returns control to the caller. Resuming that fiber continues
34-
after the yield. Dropping a suspended fiber resumes it with a private
35-
cancellation panic so Rust unwinds the stack and runs local destructors.
52+
`Scheduler::current().unwrap().wait(future)` is also available. Both forms
53+
require a current scheduler task. A future can borrow local data and does not
54+
need to implement Send or Unpin. Each wait pins the future on the task stack
55+
and reuses the task's wake signal, without allocating a separate future or
56+
waker. A future from another runtime still needs that runtime's I/O, timer,
57+
or other services to be running.
58+
59+
Run the producer/consumer example with:
60+
61+
cargo run --package socketry-concurrent --example nested_wait
62+
63+
Use `Task::current().unwrap().block()` or `Scheduler::block_current()` to park a
64+
task. Another thread can make it runnable through its `TaskHandle`. Dropping a
65+
scheduler unwinds suspended task stacks and runs their local destructors.
66+
67+
## Polling and nested waits
68+
69+
An ordinary `.await` can return `Poll::Pending` from the task's future. A nested
70+
`wait`, however, can suspend while that same `poll` call is still executing.
71+
The scheduler resumes the saved stack at that wait; it does not call the outer
72+
future's `poll` again until the previous invocation has returned. Wakeups for
73+
outer futures are preserved while inner waits run.
3674

37-
Use Scheduler::spawn_fiber for synchronous stackful tasks. Such a task can call
38-
Scheduler::block_current and be made runnable again through its TaskHandle.
75+
This implementation uses one scheduler thread for both cases. A future
76+
executor could move a Send future between completed polls, but a stack
77+
suspended inside a poll must remain on its original worker. The current
78+
scheduler implements neither work stealing nor cross-thread migration.
3979

4080
## Native context switches
4181

@@ -53,14 +93,14 @@ vendored tree for reference.
5393

5494
The x86-64 Linux build includes CRuby's CET shadow-stack switch path. It checks
5595
whether shadow stacks are enabled at runtime and allocates a shadow stack for
56-
each fiber only when needed.
96+
each task only when needed.
5797

5898
## Sanitizers
5999

60100
The `address-sanitizer` and `thread-sanitizer` Cargo features enable the
61101
compiler runtime's fiber-switch hooks. Pair one feature with the matching Rust
62102
sanitizer flag on nightly; the hooks let each runtime follow the custom stacks
63-
used by Fiber.
103+
used by scheduled tasks.
64104

65105
RUSTFLAGS="-Zsanitizer=address" cargo +nightly test -Zbuild-std --target aarch64-apple-darwin --package socketry-concurrent --features address-sanitizer
66106

@@ -78,10 +118,10 @@ embedded header.
78118

79119
## Current boundaries
80120

81-
- Fiber and Scheduler are thread-affine and are not Send.
121+
- Task contexts and Scheduler are thread-affine and are not Send.
82122
- Scheduler tasks may contain non-Send futures because they are polled only on
83123
the scheduler's owning thread.
84124
- TaskHandle::unblock and future wakers are thread-safe and only enqueue work;
85-
they never resume a fiber on the waking thread.
125+
they never resume a task on the waking thread.
86126
- This first version does not implement work stealing or cross-thread stack
87127
migration.

‎crates/concurrent/src/context.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -294,7 +294,7 @@ impl Drop for Context {
294294
#[cfg(coroutine_thread_sanitizer)]
295295
if self.thread_sanitizer_fiber_owned != 0 && !self.thread_sanitizer_fiber.is_null() {
296296
// SAFETY: this context owns the fiber handle and is no longer
297-
// running when its containing Fiber is dropped.
297+
// running when its containing task coroutine is dropped.
298298
unsafe {
299299
thread_sanitizer_destroy_fiber(self.thread_sanitizer_fiber);
300300
}

0 commit comments

Comments
 (0)