Expand description
Owned tasks and cooperative future polling on work-stealing workers.
Tasks require Send + 'static and use the ordinary worker thread stack.
The scheduler owns spawned tasks even when their join handles are dropped.
Use a Barrier to explicitly own a group of child tasks.
§socketry-executor
Owned future tasks and a work-stealing scheduler for Socketry’s Rust packages.
Use socketry as the common entry point, or depend on this package directly.
§Execution
Scheduler::new()starts workers using available hardware parallelism.Scheduler::with_workers(count)selects a nonzero worker count explicitly.scheduler.spawn(future)registers an owned task and schedules it immediately.handle.awaitreturnsResult<Output, TaskError>.scheduler.block_on(future)polls a root future on the calling thread. This root can borrow local data and need not be Send. It is not a spawned task.scheduler.run()waits until all owned tasks finish; it leaves admission open.scheduler.shutdown()closes admission, cancels tasks and joins workers.
use socketry_executor::{Scheduler, yield_now};
fn main() -> Result<(), Box<dyn std::error::Error>> {
let scheduler = Scheduler::with_workers(4)?;
let children = scheduler.barrier();
let task = children.spawn(async {
yield_now().await;
42
})?;
children.close();
scheduler.block_on(async {
assert_eq!(task.await?, 42);
children.wait().await;
Ok::<(), socketry_executor::TaskError>(())
})?;
Ok(())
}Workers poll pinned Send + 'static futures on ordinary thread stacks. A task
may migrate between polls; its pinned future stays at the same address. Only
one worker polls a given task at a time. Futures returning Pending must arrange
wakeups according to the standard Future contract.
Scheduler::current() returns a scheduler handle on workers and within a
block_on root. Task::current() returns the task being polled or destroyed;
it returns None at the root. Task references are thread-safe and support
identity and cancellation. They do not keep the task executing.
Blocking scheduler entry points panic on worker threads. Ordinary blocking system calls still block a worker, and CPU-intensive code must yield explicitly.
§Queues and affinity
Each worker owns a FIFO crossbeam-deque::Worker and has a concurrent incoming
queue for remote wakeups. New tasks submitted outside a worker enter the global
injector. New tasks spawned by a worker enter its local queue.
Once a task has run, wakeups target its last worker. That worker uses its local queue when waking itself; other threads use its incoming queue. A worker checks external work periodically even while its local queue stays busy. A worker without work steals batches from other workers’ local queues or incoming queues. This preserves affinity until another worker needs work; it does not pin tasks to threads.
Workers publish their sleeping state and recheck work before parking. Enqueuers wake a sleeping worker after publishing work. Crossbeam’s Retry result causes another search rather than parking. These transitions preserve racing wakeups. Paired sequentially consistent fences order publication against idle registration. A Loom model covers that handshake; integration tests exercise the actual queues and task implementation.
§Ownership and cancellation
Schedulers own top-level tasks. A Barrier explicitly owns its direct children
while the same scheduler executes them. Both implement Spawn; its associated
handle type leaves room for adapters using different runtimes.
- Dropping a join handle abandons its result; it does not cancel the owned task.
task.task().cancel()requests cancellation without waiting.task.cancel().awaitrequests cancellation and awaits future destruction.barrier.close()prevents new children while existing children continue.barrier.wait().awaitwaits until there are no direct children. Close first when the set of children must remain closed.barrier.stop().awaitcloses, cancels and waits.- Dropping a barrier closes it and requests cancellation without waiting.
- A scheduler’s surviving handles reject submissions after shutdown begins.
Cancellation is observed before the next poll. A running poll may return a successful result before cancellation is observed. Destruction waits until that poll returns. Dropping the future runs ordinary destructors, not the remainder of its async body. A task that never returns from poll prevents joined shutdown.
Future panics become TaskError::Panicked containing the original payload.
Await join handles to observe errors. Barrier waits and scheduler.run do not
aggregate task results or propagate unobserved failures.
Parent tasks must explicitly await their barriers for joined child cleanup.
Dropping a parent can drop its barriers and request cancellation, but parent
completion does not automatically wait for descendants. Owned tasks must be
'static; ownership does not permit borrowing a parent’s local variables.
Dropping a scheduler outside a worker cancels tasks and joins workers. Dropping one on a Socketry worker requests cancellation and lets workers finish without joining synchronously, avoiding a worker waiting for itself.
§Implementation costs
async-task supplies pinned task storage, wakers, runnable state and join
handles. Socketry adds a separately allocated, reference-counted task record
for identity, cancellation, ownership and affinity. Each barrier also has a
shared owner record. The task registry retains a waker while the task is alive.
Spawning and completion take the ownership registry mutex. Cancellation and owner closure also use it. Ordinary polling, wakeups and ready-queue operations do not take that mutex. Polling establishes task context with an Arc clone; explicit current-task/current-scheduler lookups clone shared references.
Queues allocate backing storage as needed. Rescheduling reuses the existing task; it does not allocate another future or a coroutine stack. Idle worker selection can scan worker flags, with a count allowing the scan to be skipped when all workers are busy. No throughput or allocation benchmark is claimed yet.
§I/O and selectors
The public scheduler module contains portable Network, FileIo, and Clock
traits, the Socketry implementation in socketry.rs, the optional Tokio adapter
in tokio.rs, and native implementations under selector/.
native(default) supplies TCP connect/accept/read/write/readiness through async-io’s process-wide reactor: epoll on Linux, kqueue on Apple/BSD, and IOCP/AFD socket readiness on Windows. The three platform modules expose the shared implementation; they do not duplicate its registration machinery.io-uringselects a dedicated Linux completion selector for socket and file reads/writes. Connection setup, accept, readiness and timers use async-io. On other supported platforms this feature leaves the platform default intact.tokioadds an adapter using a suppliedtokio::runtime::Handle. It does not construct or drive a runtime. Enable that runtime’s I/O and time facilities.- With default features disabled, the executor and public contracts still
build. Enabling only
tokioavoids the native I/O dependencies.
Scheduler and SchedulerHandle implement the traits. Register an owned TCP
socket/listener once, or use connect/accept, and retain the resulting resource
across operations. Returned operation futures are Send. Registrations remain
with their original reactor as tasks migrate; they are not re-created per poll.
Tokio resources passed to an adapter for another runtime return InvalidInput.
Read/write operations take ownership of a Vec<u8> and return the buffer with
the result, including ordinary errors. Read buffers must have an initialized
length (vec![0; capacity]); capacity alone supplies no writable bytes. Buffer
length is unchanged; only the first returned byte count contains new data.
Reads and writes can be partial. Each read/write call is one operation, not a
read-exact/write-all convenience method.
FileIo::file_read_at and file_write_at accept an Arc<std::fs::File> and an
explicit offset. The readiness and Tokio implementations use blocking pools.
Use ordinary files opened without append mode, and offsets fitting i64. Unix
positioned operations leave the shared cursor unchanged; the Windows blocking
fallback updates it according to std’s seek_read/seek_write semantics.
Clock::sleep uses async-io timers or Tokio timers. The io-event timer algorithm
has not yet been ported. There is no general public blocking-task API yet.
§Cancellation and shutdown
Dropping a read/write future abandons its result. A submitted operation may still consume or transmit bytes. Ownership of the buffer and file/socket continues until kernel access ends. Await a result when the byte count matters.
The io_uring selector probes required opcodes and completion-overflow support. It returns initialization errors rather than silently switching backend. Cancellation requests and original completions have distinct identifiers; only an original terminal completion releases the operation’s resources. Socketry shutdown outside a worker joins tasks and drains its ring. Dropping the scheduler on a worker requests cancellation without blocking that worker.
Readiness resources use a process-wide reactor, which is not shut down with an individual Socketry scheduler. An already-started blocking file operation can continue after its waiting task is cancelled. The Tokio adapter’s shutdown joins its owned tasks, not the external runtime or its blocking-operation pool.
An unexpected io_uring selector failure after submission cannot return memory whose kernel lifetime is unknown. That exceptional path retains the affected resources and panics the waiting operation; it does not free in-flight buffers.
§Costs and remaining work
Readiness registrations and buffers are reused. io_uring currently uses a dedicated selector thread, a command channel and one completion channel per operation; it does not pool operation records or register buffers. The Tokio file fallback retains an Arc/Mutex buffer owner to return the buffer even when a queued blocking job is cancelled during runtime shutdown. Neither fallback copies the buffer bytes just to transfer ownership.
Native overlapped Windows file operations, arbitrary descriptor/UDP APIs, operation and buffer pools, and the io-event timer port remain future work. There is no thread-local non-Send task facility. No I/O throughput claim has been established by benchmarks.
The portable example uses the same generic TCP exchange with either runtime:
cargo run -p socketry-executor --example portable_io
cargo run -p socketry-executor --example portable_io --no-default-features --features tokio
cargo run -p socketry-executor --example portable_io --features io-uring # Linux§Preserved prototype
The native coroutine prototype, verbatim CRuby vendor sources and its tests are
preserved in commit b520f3d on branch coroutine.
Run the example and tests with:
cargo run --package socketry-executor --example work_stealing
cargo test --workspace --all-targets --locked
cargo test --workspace --doc --lockedRe-exports§
pub use scheduler::BufferResult;pub use scheduler::Clock;pub use scheduler::FileIo;pub use scheduler::Interest;pub use scheduler::Network;pub use scheduler::Scheduler;pub use scheduler::SchedulerHandle;
Modules§
- scheduler
- Scheduler implementations and their I/O selectors.
Structs§
- Barrier
- Explicit ownership of child tasks running on a scheduler.
- Task
- A thread-safe task reference. Holding it does not prevent task cancellation.
- Task
Handle - An awaitable result of an owned task.
Enums§
- Spawn
Error - A task could not be registered with its owner.
- Task
Error - A task ended without producing its normal output.
Traits§
- Spawn
- The common spawning contract for a scheduler and an explicit child owner.
Functions§
- yield_
now - Cooperatively reschedule the current future once. Works with any executor implementing the standard future/waker contract.