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
22 changes: 15 additions & 7 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ jobs:
fail-fast: false
matrix:
# Native selectors and positioned file I/O are target-specific.
runner: [ubuntu-24.04, macos-15, windows-latest]
runner: [ubuntu-24.04, ubuntu-24.04-arm, macos-15-intel, macos-15, windows-latest]
steps:
- uses: actions/checkout@v7
- uses: actions-rust-lang/setup-rust-toolchain@v2
Expand Down Expand Up @@ -70,7 +70,7 @@ jobs:
io-uring:
name: io_uring ${{ matrix.runner }}
runs-on: ${{ matrix.runner }}
timeout-minutes: 10
timeout-minutes: 20
strategy:
fail-fast: false
matrix:
Expand All @@ -79,11 +79,15 @@ jobs:
- uses: actions/checkout@v7
- uses: dtolnay/rust-toolchain@stable
with:
components: clippy
components: clippy, llvm-tools-preview
- name: Install Bake launcher
run: cargo install socketry-cargo-bake --locked
- name: Install coverage tool
run: cargo install cargo-llvm-cov --locked
# Run on the host kernel: container policies commonly deny io_uring.
# Initialization errors fail these tests; they are not silently skipped.
- name: Exercise native completion and cancellation
run: cargo test --workspace --all-targets --all-features --locked
run: cargo bake --locked test:coverage --all-targets true --all-features true
- name: Check native selector code
run: cargo clippy --workspace --all-targets --all-features --locked -- -D warnings

Expand All @@ -96,10 +100,14 @@ jobs:
uses: vmactions/freebsd-vm@v1
with:
release: "15.1"
prepare: pkg install -y rust
prepare: pkg install -y curl ca_root_nss
run: |
cargo test --workspace --all-targets --features tokio --locked
cargo test --workspace --doc --features tokio --locked
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y --profile minimal
export PATH="$HOME/.cargo/bin:$PATH"
rustup component add llvm-tools-preview
cargo install socketry-cargo-bake --locked
cargo install cargo-llvm-cov --locked
cargo bake --locked test:coverage --all-targets true --features tokio

sanitizer:
name: Linux x86-64 ${{ matrix.name }}
Expand Down
35 changes: 23 additions & 12 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ manifest = "bake/Cargo.toml"
reviewers = ["socketry/managers"]

[workspace.package]
version = "0.1.4"
version = "0.1.5"
edition = "2024"
license = "MIT"
repository = "https://github.com/socketry/socketry-rust"
Expand All @@ -25,7 +25,7 @@ readme = "readme.md"
include = ["Cargo.toml", "readme.md", "license.md", "releases.md", "context/**", "src/**"]

[dependencies]
socketry-executor = { version = "0.1.4", path = "crates/executor", default-features = false }
socketry-executor = { version = "0.1.5", path = "crates/executor", default-features = false }

[features]
default = ["native"]
Expand Down
4 changes: 2 additions & 2 deletions context/design.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,7 +82,7 @@ Tokio and futures-io expose different AsyncRead and AsyncWrite traits. `tokio-ut

Entering a Tokio runtime context provides access to its services; entering alone does not drive the runtime. Some mixed execution is possible when the required services are running, but must be established for the concrete APIs being used. Do not advertise universal Tokio compatibility from a Waker or stream adapter alone.

The optional `scheduler::tokio` adapter implements Network, FileIo, Clock and Spawn against an existing runtime. It preserves direct-child barrier ownership and joins owned task destruction on asynchronous shutdown. The same generic TCP program runs on Socketry and Tokio. The adapter scopes runtime context to individual polls when registering resources; its futures can also be polled by Socketry workers while Tokio drives the underlying services. Socketry contextual lookups still identify Socketry execution; portable code passes handles explicitly.
The optional `scheduler::tokio` adapter implements Network, FileIO, Clock and Spawn against an existing runtime. It preserves direct-child barrier ownership and joins owned task destruction on asynchronous shutdown. The same generic TCP program runs on Socketry and Tokio. The adapter scopes runtime context to individual polls when registering resources; its futures can also be polled by Socketry workers while Tokio drives the underlying services. Socketry contextual lookups still identify Socketry execution; portable code passes handles explicitly.

For existing libraries tied to Tokio, either keep their work on Tokio and bridge owned messages/results, or provide the particular trait adapter they consume. Avoid a broad imitation of Tokio's API.

Expand Down Expand Up @@ -152,7 +152,7 @@ The algorithm's existing Ruby performance motivates the port. Rust performance c
1. Record boundaries and reusable Rust conventions (implemented).
2. Replace stackful execution with async-task and Crossbeam worker queues (implemented). Keep the coroutine prototype in its saved branch.
3. Implement explicit owners, barriers, cancellation and shutdown (implemented for direct children). Automatic descendant draining remains future work.
4. Establish minimal clock and I/O contracts with concrete consumers (implemented with Network, FileIo, Clock and a portable TCP example).
4. Establish minimal clock and I/O contracts with concrete consumers (implemented with Network, FileIO, Clock and a portable TCP example).
5. Implement native readiness and the Tokio adapter, running the same consumers with both (implemented). Native sleep uses async-io until the timer port.
6. Implement io\_uring's owned-buffer lifecycle, socket/file operations, cancellation and runtime probing (implemented). Improve operation reuse, buffer registration and submission backpressure in subsequent work.
7. Port the timer queue with upstream attribution and deterministic verification.
Expand Down
4 changes: 3 additions & 1 deletion context/implementation.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ This guide describes the current implementation and its boundaries. Read [the de

## I/O and runtime boundaries

- `scheduler/mod.rs` defines Network, FileIo and Clock. Operations return concrete Send futures; portable consumers receive the required capabilities.
- `scheduler.rs` re-exports Network, FileIO, Interest and Clock from `scheduler/network.rs`, `file_io.rs`, `interest.rs` and `clock.rs`. Operations return concrete Send futures; portable consumers receive the required capabilities.
- `scheduler/socketry.rs` owns the executor; `socketry/operations.rs` forwards capabilities to its lazily initialized, compile-time selected selector.
- `scheduler/selector/` contains readiness, epoll, kqueue, iocp and io\_uring. Platform readiness modules share async-io's persistent registrations and process-wide reactor. Registered sockets remain usable as tasks migrate.
- Default feature `native` provides TCP, positioned files and sleep. Feature `io-uring` selects Linux completion reads/writes; other supported platforms retain readiness. Feature `tokio` enables the separate runtime adapter. No default features builds the executor and contracts without native I/O.
Expand All @@ -53,3 +53,5 @@ This guide describes the current implementation and its boundaries. Read [the de
Public behavior is covered in crates/executor/tests. Deterministic channels force stealing, migration, concurrent wakeups and cancellation during polling. The parking test uses Loom to model the queue-publication/idle-registration handshake. Keep its atomics and fence order aligned with the implementation; this models the handshake, not Crossbeam or async-task internals. When verification is requested, run workspace tests and doctests with `tokio` enabled; on Linux also run all features to exercise io\_uring. Check executor-only and Tokio-only feature combinations. The same TCP/file consumers exercise both implementations; Linux tests cover cancellation batches and shutdown races. CI covers Linux, macOS, Windows and FreeBSD, with separate Linux io\_uring and sanitizer jobs; distinguish configured CI from executed results.

ThreadSanitizer loads `.github/tsan-suppressions.txt` to suppress Crossbeam's internal queue `Buffer::read` and `Buffer::write` race reports. Crossbeam reads slots speculatively and discards values when atomic validation fails. Its non-atomic volatile accesses remain a known Rust memory-model limitation; the suppression accepts that limitation rather than fixing it. See the [upstream discussion](https://github.com/crossbeam-rs/crossbeam/issues/589#issuecomment-720972996). Keep suppression patterns scoped to those buffer accesses and reassess them when updating Crossbeam. Other race reports continue to fail the sanitizer job.

Coverage measures source regions for native and Tokio schedulers on each supported OS and architecture, plus Linux io\_uring and FreeBSD. Executor-only and Tokio-only feature selections receive separate compilation checks.
1 change: 1 addition & 0 deletions crates/executor/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -34,4 +34,5 @@ io-uring = ["native", "dep:io-uring", "dep:futures-channel", "dep:libc", "dep:po
tokio = ["dep:tokio"]

[dev-dependencies]
tempfile = "3"
loom = "0.7.2"
4 changes: 2 additions & 2 deletions crates/executor/readme.md
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ 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`
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/`.

Expand Down Expand Up @@ -150,7 +150,7 @@ 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
`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
Expand Down
5 changes: 4 additions & 1 deletion crates/executor/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,5 +16,8 @@ mod worker;

pub use barrier::Barrier;
pub use owner::{Spawn, SpawnError};
pub use scheduler::{BufferResult, Clock, FileIo, Interest, Network, Scheduler, SchedulerHandle};
pub use scheduler::{BufferResult, Clock, FileIO, Interest, Network, Scheduler, SchedulerHandle};
pub use task::{Task, TaskError, TaskHandle, yield_now};

/// Compatibility spelling for [`FileIO`].
pub use scheduler::FileIo;
38 changes: 38 additions & 0 deletions crates/executor/src/scheduler.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
// Released under the MIT License.
// Copyright, 2026, by Samuel Williams.

//! Scheduler implementations and their I/O selectors.
use std::io;

pub mod selector;
pub mod socketry;

#[cfg(any(feature = "native", feature = "tokio"))]
mod file;

#[cfg(feature = "tokio")]
pub mod tokio;

pub use socketry::{Scheduler, SchedulerHandle};
pub(crate) use socketry::{Shared, enter};

/// An operation result together with its reusable, owned buffer.
///
/// Reads fill the existing buffer length and leave its length unchanged. Only
/// the first `result?` bytes contain newly read data. Writes may be partial.
/// The buffer is returned on both success and failure.
pub type BufferResult = (io::Result<usize>, Vec<u8>);

mod interest;
pub use interest::Interest;

mod network;
pub use network::Network;

mod file_io;
pub use file_io::FileIO;
/// Compatibility spelling for [`FileIO`].
pub use file_io::FileIO as FileIo;

mod clock;
pub use clock::Clock;
10 changes: 10 additions & 0 deletions crates/executor/src/scheduler/clock.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
// Released under the MIT License.
// Copyright, 2026, by Samuel Williams.

use std::future::Future;
use std::time::Duration;

/// A runtime's monotonic sleep facility.
pub trait Clock: Send + Sync {
fn sleep(&self, duration: Duration) -> impl Future<Output = ()> + Send;
}
31 changes: 31 additions & 0 deletions crates/executor/src/scheduler/file_io.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
// Released under the MIT License.
// Copyright, 2026, by Samuel Williams.

use super::BufferResult;
use std::fs::File;
use std::future::Future;
use std::sync::Arc;

/// Positioned file operations. A regular file does not support a universal
/// readiness fallback, so implementations use native completion or a blocking
/// pool. Use ordinary files opened without append mode, not pipes. Offsets
/// must fit in i64. The Unix implementation leaves the shared cursor unchanged;
/// the Windows blocking fallback updates it, as std's seek_read/seek_write do.
///
/// Buffers and the file remain owned by an in-flight operation even if the
/// waiting future is dropped. A write can still complete after cancellation.
pub trait FileIO: Send + Sync {
fn file_read_at(
&self,
file: Arc<File>,
buffer: Vec<u8>,
offset: u64,
) -> impl Future<Output = BufferResult> + Send;

fn file_write_at(
&self,
file: Arc<File>,
buffer: Vec<u8>,
offset: u64,
) -> impl Future<Output = BufferResult> + Send;
}
10 changes: 10 additions & 0 deletions crates/executor/src/scheduler/interest.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
// Released under the MIT License.
// Copyright, 2026, by Samuel Williams.

/// A socket readiness condition. Readiness can be spurious; retry nonblocking
/// operations and wait again when they return `WouldBlock`.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Interest {
Readable,
Writable,
}
Loading
Loading