diff --git a/context/design.md b/context/design.md index e620289..d4cb4dc 100644 --- a/context/design.md +++ b/context/design.md @@ -60,6 +60,18 @@ Distinguish cooperative cancellation, dropping a future, and waiting for termina Contextual APIs such as `Scheduler::current()` can remain conveniences. Set and restore their context around every poll, including when polling panics; do not assume a future always runs on the thread where it was created. +### Cooperative cancellation + +The shutdown direction is cooperative: request graceful termination and let the application continue polling until its work completes. There is no default forced abort, second-request escalation, or grace-period deadline. If a process refuses to stop, its supervisor can follow SIGTERM with SIGKILL. + +`Cancellation` is implemented in `socketry-executor` and re-exported by `socketry`. It works independently of the executor. Clones share a persistent signal; child signals receive ancestor cancellation without propagating it to parents or siblings. This explicit cancellation tree expresses service boundaries without imposing an implicit task hierarchy. Dropping a signal does not cancel it. Rust destructors release retained resources, but asynchronous draining still requires live futures and explicit completion tracking through task handles or barriers. + +`cancel()` signals intent; `check()` returns `Cancelled` at cooperative cancellation points; `cancelled().await` suspends using wakers. `never()` represents work without cooperative cancellation. Passing signals explicitly to future I/O operations remains planned; current Socket, File, and Clock signatures are unchanged. + +`defer_cancel(&signal, future, on_cancel)` observes the signal, invokes a synchronous shutdown callback once, and keeps awaiting the future's normal output. The callback tells the protected work to stop accepting new work or begin draining. Async cleanup belongs in that future, with cancellation inputs that allow it to finish. The wrapper uses standard future polling, handles pinned futures, and neither busy-waits nor prevents its own destruction. + +The intended application sequence is to request application or service cancellation, await task/group completion while runtime services remain available, then release the owner and scheduler. Process signal installation and cooperative scheduler shutdown are future integration work. Existing Task cancellation, Barrier stopping and scheduler shutdown still destroy futures and bypass deferred async cleanup; they are distinct from cooperative cancellation. + ### Current executor The implementation currently lives in socketry-executor. async-task owns pinned future storage and its runnable/waker state. Socketry supplies thread management, ownership, cancellation flags, task identity and queue selection. Each worker has a FIFO Crossbeam queue plus an incoming injector for remote wakeups. External submissions enter a global injector; worker submissions go to that worker's queue. Wakeups target the last worker. Only workers without their own work steal from other local queues or incoming queues. @@ -82,7 +94,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 Socket, File, 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. @@ -151,8 +163,8 @@ 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). +3. Implement explicit owners, barriers, cancellation and shutdown (implemented for direct children). Runtime-independent cooperative signals and deferred cancellation are implemented; cooperative scheduler shutdown, cancellation-aware I/O, and automatic descendant draining remain future work. +4. Establish minimal clock and I/O contracts with concrete consumers (implemented with Socket, File, 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. diff --git a/context/gaps.md b/context/gaps.md new file mode 100644 index 0000000..850f554 --- /dev/null +++ b/context/gaps.md @@ -0,0 +1,31 @@ +# Scheduler Gaps + +Socketry has a working future executor, explicit task ownership, portable I/O capabilities, and a Tokio adapter. These four gaps describe the remaining architectural work toward a scheduler inspired by [Zig's `std.Io` model](https://ziglang.org/download/0.16.0/release-notes.html#I-O-as-an-Interface), while retaining ordinary Rust futures and `.await`. + +## 1. Graceful cancellation needs runtime integration + +Task cancellation currently drops the future before its next poll. An active poll must return first. Synchronous destructors run, but the remainder of the async body, including asynchronous cleanup, does not execute. + +`Cancellation` and `defer_cancel` now provide runtime-independent cooperative requests and deferred cleanup. Signals can be shared or linked through child signals to shut down the application or an individual service. They signal intent; task handles and barriers still confirm completion. + +Integrating these primitives into scheduler shutdown and cancellation-aware I/O remains pending. Graceful shutdown must keep workers, I/O services, and task owners available until draining completes. The intended model has no default forced abort or internal escalation; a process supervisor handles SIGTERM followed by SIGKILL. Existing task cancellation and scheduler shutdown still destroy futures, so callers must request cooperative cancellation and await work first. + +## 2. Ownership currently joins direct children + +Barriers own and join their direct children. Cancelling a parent can drop its barriers and request child cancellation, but parent completion does not automatically wait for descendants to finish. Callers must explicitly await their barriers when joined child cleanup is required. + +`Spawn` makes task submission portable, while creating and controlling barriers still uses runtime-specific APIs. A generic library can spawn children, but cannot express the full group lifecycle through a shared contract. Define portable group creation, admission closure, waiting, and stopping, and decide whether parent completion should include descendant draining. + +## 3. The I/O boundary is still narrow + +`Socket` covers TCP registration, connect, accept, reads, writes, and readiness. Listener binding happens through `std::net`. `File` accepts an already-open `std::fs::File`, and `Clock` only supplies relative sleep. + +Opening resources, DNS, deadlines, synchronization, randomness, and processes remain outside the portable boundary. Expand capabilities as real consumers need them, keeping the application's implementation choice explicit. A broader boundary would also allow more operating-system behavior to be substituted in deterministic tests. + +## 4. The backends establish behavior before performance + +The default selector delegates readiness and timers to async-io's shared reactor. Linux io\_uring uses a dedicated selector thread and retains buffers until original operations complete; cancellation completions alone do not release them. + +The io\_uring implementation currently uses an unbounded command channel and a completion channel per operation. Operation records are not pooled, and buffers are not registered with the kernel. Add submission backpressure and measure allocation, contention, throughput, and retained memory before choosing optimizations or making performance claims. + +See the [design guide](design.md) and [implementation guide](implementation.md) for the current contracts and implementation sequence. diff --git a/context/implementation.md b/context/implementation.md index cabadbf..bc983de 100644 --- a/context/implementation.md +++ b/context/implementation.md @@ -26,9 +26,20 @@ This guide describes the current implementation and its boundaries. Read [the de - Parents must explicitly await barriers for joined cleanup. Automatic waiting for descendants after parent destruction is not implemented. - Scheduler Drop cancels all tasks, joining threads outside a worker. On a worker it requests shutdown without joining. Surviving handles reject spawn. +## Cooperative cancellation + +- `cancellation.rs`, `cancelled.rs`, and `defer_cancel.rs` implement runtime-independent primitives, exported by both `socketry-executor` and the `socketry` facade. +- Cancellation clones share state. Children receive ancestor cancellation, retain ancestors after intermediate handles are dropped, and never cancel parents or siblings. Dropping signals does not cancel work. `Cancellation::never()` allocates nothing; its children are independent cancellable signals. +- Each family serializes child registration and cancellation with one mutex. Atomic flags support cheap checks. Nodes hold weak child registrations; child destruction unregisters them. Cancellation and exclusively owned ancestor destruction are iterative, avoiding stack exhaustion on deep chains. +- Register event listeners before checking cancellation to avoid lost notifications. Dropping a wait unregisters its listener. Release all signal locks before invoking wakers so wakeups can reenter cancellation and child registration. +- defer\_cancel pins the work and its signal wait. Observed cancellation consumes the synchronous callback once, before polling the work again, and the wrapper preserves the work's normal output. Callback panics propagate; dropping the wrapper destroys its future. It does not mask operation signals or intercept task destruction. +- For graceful shutdown, request cancellation and await task handles or barriers before releasing the owner or shutting down runtime services. Signals express intent rather than completion. There is no forced abort or escalation in this cooperative contract; a supervisor can enforce SIGTERM followed by SIGKILL. +- Existing Task cancellation, Barrier stopping, and scheduler shutdown retain their destruction semantics. Signal handlers, cooperative runtime shutdown and cancellation inputs on I/O methods are not yet integrated. + ## I/O and runtime boundaries -- `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.rs` re-exports Socket, File, Interest and Clock from `scheduler/socket.rs`, `file.rs`, `interest.rs` and `clock.rs`. Operations return concrete Send futures; portable consumers receive the required capabilities. +- `File` names the scheduler capability, while `std::fs::File` remains the resource type. Use the `StdFile` alias where both names are imported. Blocking positioned file helpers live in `scheduler/positioned_file.rs`. - `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. diff --git a/crates/executor/Cargo.toml b/crates/executor/Cargo.toml index 2541178..d0f7845 100644 --- a/crates/executor/Cargo.toml +++ b/crates/executor/Cargo.toml @@ -6,7 +6,7 @@ description = "Owned asynchronous tasks and a work-stealing future scheduler" license.workspace = true repository.workspace = true readme = "readme.md" -include = ["Cargo.toml", "license.md", "readme.md", "src/**", "examples/**"] +include = ["Cargo.toml", "license.md", "readme.md", "design.md", "src/**", "examples/**"] [lib] name = "socketry_executor" diff --git a/crates/executor/design.md b/crates/executor/design.md new file mode 100644 index 0000000..4afba79 --- /dev/null +++ b/crates/executor/design.md @@ -0,0 +1,594 @@ +# Executor Design + +This document records the agreed direction for `socketry-executor`: an explicit scheduler and I/O interface inspired by [Zig 0.16's I/O model](https://ziglang.org/download/0.16.0/release-notes.html#I-O-as-an-Interface), expressed through ordinary Rust futures, capability traits, and cooperative cancellation. It includes proposed interfaces; the implementation status below distinguishes them from APIs that exist today. + +## Execution model + +Use Rust's existing `Future`, `Context`, and `Waker` contracts. A separate computation trait is unnecessary for the current design. Calling an operation constructs a future; awaiting it allows the executor to poll it and suspend until a wakeup permits progress. Polling does not mean busy waiting. + +The scheduler executes pinned futures on ordinary worker stacks. Independently spawned tasks are `Send + 'static` and can migrate between workers, while their pinned storage remains stationary. Only one worker polls a task at a time. Blocking system calls must run outside the workers that poll futures. + +Pass the implementation into libraries explicitly. A library should require the capabilities it actually uses rather than implicitly selecting a global runtime. Spawning additionally requires an explicit task owner. Existing resources can be passed directly to code that does not need to create them. + +## Capability traits and static dispatch + +`File`, `Socket`, and `Clock` name scheduler capabilities. `Spawn` expresses task submission and ownership separately. `File` is the capability trait; `std::fs::File` is currently the resource type, conventionally imported as `StdFile` when both names are needed. `Socket` replaces the earlier capability name `Network`. + +Prefer generic trait bounds, concrete resource types, and concrete future types. These allow compile-time capability checking, monomorphization, and inlining without requiring a boxed future for each operation. The compiler decides which calls to inline; the interface should make that optimization possible. + +Keep core traits small and dependable. Implementing a core trait requires working implementations of its operations for the documented resource types and inputs. Avoid default methods that simply return `Unsupported`. + +Additional capabilities belong in extension traits where their semantics or platform availability differ. A consumer requiring an extension declares that requirement through its trait bounds. An incompatible scheduler then fails at compile time. For example, `S: File + FileAllocate` illustrates how a consumer might require allocation; `FileAllocate` is a proposed name, not an existing trait or a settled trait boundary. + +Do not create one trait per syscall automatically. Group operations around capabilities that real consumers need, and keep inherently platform-specific types and flags in platform extensions. + +## Platform support and fallbacks + +Backend selection belongs inside scheduler construction or resource initialization, rather than in application error handling around every operation. + +| Situation | Intended contract | +| --- | --- | +| Native asynchronous operation exists | Use the native implementation. | +| Equivalent blocking operation exists | Run it on a blocking pool and preserve the same future and result semantics. | +| Entire capability has no valid implementation | Do not implement the corresponding trait for that scheduler/platform. | +| Capability depends on the running kernel or a registration facility | Probe when initializing the backend or creating the capability/resource. | +| A particular filesystem, resource, or flag combination rejects an operation | Return an ordinary operation error, preserving the OS error where available. | +| Interface is inherently platform-specific | Expose it through an extension trait or platform module. | + +A trait bound guarantees an implementation exists; it cannot guarantee that every filesystem supports every allocation mode or that every operation succeeds. Runtime failures remain part of the result contract. `io::ErrorKind::Unsupported` can represent a genuine resource or feature limitation, but should not be the routine result of calling a core operation on a supported platform. + +Use a blocking pool when native asynchronous support is unavailable. This is acceptable on Windows and other platforms. A fallback may change performance, but must preserve the documented semantics. Allocation cannot silently become truncation, message I/O cannot discard ancillary data, and a positioned operation cannot alter the shared cursor if its contract promises to preserve it. + +Bound blocking concurrency and provide asynchronous admission/backpressure. Cancellation can prevent queued work from starting. Once a blocking syscall has started, the operation may have to finish; retain its resources and buffer throughout. A cancellation result must not imply that completed side effects were undone. + +An application explicitly selecting a specialized backend may receive an initialization failure instead of automatic fallback. The policy must be clear at that boundary. The current Linux io\_uring backend follows this approach: it probes its requirements and fails initialization when they are unavailable. Automatic backend fallback remains a design option, not implemented behaviour. + +## Cancellation and task ownership + +Cancellation signals intent. Task handles and barriers establish completion. Resource ownership determines what remains alive. Keep these responsibilities distinct. + +Rust's destruction semantics help release owned resources and make an implicit cancellation hierarchy unnecessary. They do not replace task ownership or asynchronous cleanup: ordinary `Drop` cannot await, and dropping a future does not execute the rest of its async body. Independently spawned tasks still need owners and explicit completion tracking. + +Use explicit cancellation scopes to describe shutdown relationships. A process can have a root signal and each service can have a child signal. This cancellation tree need not mirror the task ownership graph. A task tree is not imposed on every future. + +### Cancellation interface + +`Cancellation` lives in `socketry-executor` and is re-exported by the `socketry` facade. It works independently of a scheduler. + +| Operation | Behaviour | +| --- | --- | +| `Cancellation::new()` / `Default` | Create an independent, uncancelled signal. | +| `clone()` | Share the same signal. | +| `child()` | Create a distinct signal that also receives ancestor cancellation. | +| `cancel() -> bool` | Request cancellation; return true only for the first request on that signal. | +| `is_cancelled()` | Observe the persistent state. | +| `check() -> Result<(), Cancelled>` | Check at an explicit cooperative cancellation point. | +| `cancelled().await` | Wait using normal future polling and wakeups. | +| `Cancellation::never()` | Supply an allocation-free signal that cannot be cancelled. | + +Cancelling a child leaves its parent and siblings running. A child created after its parent is cancelled starts cancelled. Dropping a signal does not cancel it. A child of `never()` is an independent cancellable signal. Repeated cancellation requests are idempotent and never escalate. + +### System and service shutdown + +The agreed shutdown contract is cooperative. Process signal integration should translate SIGINT/SIGTERM into a request on the application's root signal. Cancelling a service's child signal requests shutdown of that service without shutting down unrelated work. + +Graceful shutdown proceeds as follows: + +1. Request cancellation for the application or service scope. +2. Stop accepting new work into that scope. +3. Keep polling existing work so it can drain and perform asynchronous cleanup. +4. Await task handles or barriers to confirm completion. +5. Release resources and shut down runtime services after draining finishes. + +Workers, selectors, timers, blocking-operation services, and task owners must remain available throughout draining. A process-level cancellation request must not immediately disable the services that cleanup needs. + +There is no default forced abort, grace-period deadline, or escalation on a second request. Work may choose to ignore cancellation. A process supervisor can enforce termination by following SIGTERM with SIGKILL. An application can also abandon futures and destroy its runtime, giving up asynchronous cleanup; this does not remove the runtime's obligation to retain memory still accessible by the kernel. + +### Deferred cancellation + +The implemented wrapper is: + +```text +defer_cancel(&cancellation, future, on_cancel) +``` + +It observes cancellation while polling the protected future. On the first observation, it invokes the synchronous `FnOnce()` callback and continues polling the future to its normal output. The callback can stop an accept loop, request draining, or wake the protected work. Asynchronous cleanup stays inside the future. + +The callback runs before the next protected poll, including the first poll if the signal is already cancelled. If cancellation and completion are observed in the same wrapper poll, the callback runs and the normal output is returned. Callback panics propagate. + +The wrapper does not mask cancellation inputs passed to operations inside the future. Cleanup must use inputs that permit it to complete, such as a separate scope or `Cancellation::never()`. Dropping the wrapper still drops its future; it does not intercept task destruction. + +### Transitional runtime behaviour + +The cooperative primitives are implemented, but existing `Task::cancel`, `Barrier::stop`, scheduler shutdown, and scheduler destruction still request destruction of task futures. They can bypass deferred asynchronous cleanup. Callers must currently request cooperative cancellation and await draining before invoking those lifecycle operations. + +Barriers own and join their direct children. Parent completion does not automatically await all descendants; owners must explicitly await the groups whose completion matters. Signal integration, cooperative scheduler shutdown, and a portable group lifecycle remain future work. + +## Scheduler operation interface + +The intended low-level shape is a method on the scheduler capability that returns a future and accepts an explicit cancellation input. Keep operations recognizable from their system interfaces, for example: + +```text +scheduler.file_open(path, options, cancellation) +scheduler.file_read(file, buffer, cancellation) +scheduler.file_read_at(file, buffer, offset, cancellation) +scheduler.socket_accept(listener, cancellation) +scheduler.socket_sendmsg(socket, message, flags, cancellation) +scheduler.address_resolve(name, options, cancellation) +scheduler.sleep(duration, cancellation) +``` + +These are illustrative proposed signatures, not current APIs. Exact resource, options, cancellation borrowing, and error types still need to be chosen. Operations that return owned buffers must preserve that ownership contract for errors and cancellation as well as success. + +An optional context pairing a scheduler reference with a cancellation signal could reduce repetitive arguments. Such a context should preserve explicit dependency and cancellation boundaries; it is a convenience, not a new computation abstraction. + +Sleep uses monotonic time and completes when the requested interval has elapsed. It need not return the elapsed time; callers can measure that separately. Cancellation-aware sleep is expected to return completion or `Cancelled`. Current `Clock::sleep` returns `()` and has no cancellation argument. + +Blocking and unblocking a suspended computation should use registration and wakeups with protection against lost notifications. This synchronization interface is separate from running a blocking syscall on a thread pool; its precise API remains open. + +## High-level interface sketch + +The following is a proposed API for discussion, not a description of implemented interfaces or a compilable example. It specifies caller-visible behaviour without choosing backend machinery. Resource, options, buffer, and error type names are provisional. `IoBuf` and `IoBufMut` stand for the owned-buffer contracts described below; their definitions remain open. + +### Results and ownership + +Use a common operation error that distinguishes cooperative cancellation from an OS error, while preserving the latter for inspection: + +```rust +pub enum OperationError { + Cancelled(Cancelled), + Io(std::io::Error), +} + +pub type OperationResult = Result; +pub type BufferResult = (OperationResult, B); +``` + +Operations borrowing a resource return `OperationResult`. Operations taking a caller-owned buffer return `BufferResult`, with the buffer outside the `Result` so it can be recovered on failure. Reads and writes return partial transfer counts. A read's count identifies the valid received prefix; a write sends only initialized bytes. Read-exact, write-all, and send-all loops belong above this interface. + +An already-cancelled input prevents new ordinary work from starting. For work already in flight, cancellation requests a cooperative stop; awaiting the operation still establishes its completion and resource ownership. Report a successful transfer that races cancellation as progress rather than replacing its count with `Cancelled`. Cancellation never promises to undo side effects. + +The sketches borrow resource handles and cancellation signals. In-flight operations must retain the underlying resources for as long as necessary, even if the waiting future is abandoned. Resource handles belong to their creating scheduler/backend; cross-instance compatibility is checked at the appropriate boundary. + +### Core capability shapes + +Keep the receiver, resource, operation arguments, and cancellation input explicit: + +```rust +pub trait File: Send + Sync { + type File: Send + Sync; + + fn file_open( + &self, + path: &Path, + options: FileOpenOptions, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn file_read( + &self, + file: &Self::File, + buffer: B, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn file_write( + &self, + file: &Self::File, + buffer: B, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn file_read_at( + &self, + file: &Self::File, + buffer: B, + offset: u64, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn file_write_at( + &self, + file: &Self::File, + buffer: B, + offset: u64, + cancellation: &Cancellation, + ) -> impl Future> + Send; +} + +pub trait Socket: Send + Sync { + type Socket: Send + Sync; + type Listener: Send + Sync; + + fn socket_accept( + &self, + listener: &Self::Listener, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn socket_connect( + &self, + address: SocketAddr, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn socket_recv( + &self, + socket: &Self::Socket, + buffer: B, + flags: RecvFlags, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn socket_send( + &self, + socket: &Self::Socket, + buffer: B, + flags: SendFlags, + cancellation: &Cancellation, + ) -> impl Future> + Send; +} + +pub trait Clock: Send + Sync { + fn sleep( + &self, + duration: Duration, + cancellation: &Cancellation, + ) -> impl Future> + Send; +} +``` + +The proposed associated file resource allows a backend-owned handle rather than fixing future implementations to `std::fs::File`. This is a change from the current `File` trait. The connect convenience above creates and connects a socket; a lower-level variant taking an existing socket may be needed for callers that configure it before connecting. Socket creation, binding, listening, registration, explicit close, message ownership, and portable flags still need their own signatures. Consuming resource operations such as explicit close must define what ownership is returned if cancellation prevents them from starting. + +Additional capability shapes follow the same convention: + +| Capability | Proposed operation shape | Output | +| --- | --- | --- | +| File allocation | `file_allocate(file, offset, length, mode, cancellation)` | `OperationResult<()>` | +| File synchronization | `file_sync(file, mode, cancellation)` | `OperationResult<()>` | +| Message send | `socket_sendmsg(socket, owned_message, flags, cancellation)` | Transfer result together with the owned message | +| Message receive | `socket_recvmsg(socket, owned_message, flags, cancellation)` | Receive result including address/control/status, together with the owned message | +| Address resolution | `address_resolve(name, options, cancellation)` | `OperationResult` | + +These rows identify extension capabilities; they do not require every scheduler to implement every operation. Synchronization wait/wake remains a separate interface to specify, including lost-wakeup protection. `Spawn` continues to require an explicit task owner. + +### Pool receive capability + +A pool receive returns a selected filled lease, rather than accepting an already-selected buffer: + +```rust +pub trait PoolReceive: Socket { + type Pool: Send + Sync; + type Lease: IoBufMut; + + fn buffer_pool( + &self, + options: BufferPoolOptions, + cancellation: &Cancellation, + ) -> impl Future> + Send; + + fn socket_recv_with_pool( + &self, + socket: &Self::Socket, + pool: &Self::Pool, + flags: RecvFlags, + cancellation: &Cancellation, + ) -> impl Future> + Send; +} +``` + +The returned lease exposes only the initialized received payload for sending, and can therefore be passed directly to `socket_send`. Its readable length is the receive count, not the capacity of its backing allocation. For stream sockets, an empty successful lease represents EOF. The pool owns storage that does not escape as a lease: on error or cancellation, it is safely reclaimed only after backend access ends. This differs from an operation that must return a buffer supplied by the caller. If a receive has produced usable data, report that lease as progress rather than discarding it solely because cancellation raced completion. + +`BufferPoolOptions` would describe a bounded number of buffers and their capacities. Native selection or fixed-registration requirements belong in separate extensions, rather than changing the meaning of this portable capability. Pool message receives need a corresponding result carrying the filled lease and message metadata. + +### Example: positioned file read + +```rust +async fn read_header( + scheduler: &S, + path: &Path, + cancellation: &Cancellation, +) -> OperationResult> { + let file = scheduler + .file_open(path, FileOpenOptions::read_only(), cancellation) + .await?; + + let (result, mut buffer) = scheduler + .file_read_at(&file, vec![0; 4096], 0, cancellation) + .await; + + let count = result?; + buffer.truncate(count); + Ok(buffer) +} +``` + +This deliberately performs one read and permits a short result. A caller needing a complete header would use a read-exact helper. Keeping `result` and `buffer` separate also lets a caller reuse the allocation after an error instead of propagating it immediately. + +### Example: forwarding pooled socket data + +```rust +async fn forward( + scheduler: &S, + source: &S::Socket, + destination: &S::Socket, + pool: &S::Pool, + cancellation: &Cancellation, +) -> OperationResult<()> { + loop { + let lease = scheduler + .socket_recv_with_pool(source, pool, RecvFlags::empty(), cancellation) + .await?; + + if lease.is_empty() { + return Ok(()); + } + + let (result, lease) = send_all( + scheduler, destination, lease, SendFlags::empty(), cancellation, + ).await; + + // Sending is complete; the lease can return to its pool. + drop(lease); + result?; + } +} +``` + +`send_all` here is a proposed higher-level helper, not a scheduler hook. It repeatedly calls `socket_send` with owned views of the unsent range, handles zero progress, and returns the lease on success or failure after all backend access ends. The view interface also remains to be specified. This example uses explicit sequential awaits and makes no assumption about `with_previous_buffer`. + +Create the pool once and share it across connections, for example: + +```rust +let pool = scheduler.buffer_pool( + BufferPoolOptions { buffers: 128, buffer_capacity: 16 * 1024 }, + &cancellation, +).await?; + +forward(&scheduler, &source, &destination, &pool, &cancellation).await?; +``` + +### Example: prepared accept followed by receive + +```rust +let accept = scheduler.prepare_socket_accept(&listener); + +let read = scheduler.prepare_socket_recv( + accept.socket(), + vec![0; 4096], + RecvFlags::empty(), +); + +let outcomes = scheduler + .submit(accept.link(read), &cancellation) + .await; +``` + +`accept.socket()` is proposed syntax for an owned, typed reference to the socket that accept will produce. It is not an already-accepted socket or a borrow that would prevent moving the prepared accept into the chain. The receive starts only after accept succeeds; a failed accept skips it. Outcomes must preserve the accepted socket even if the receive fails or is skipped, and return the receive buffer with its partial byte count or error. + +Unlike forwarding a selected receive buffer, this dependency has a native io\_uring implementation target: direct accept into a reserved, known descriptor slot, followed by a linked receive using that slot. The Rust output-reference API and its type/lifetime rules are still proposed and unimplemented. They must prevent a consumer from running before its producing accept, or without that dependency. See the direct-descriptor linking section below. + +### Example: prepared write followed by synchronization + +```rust +let write = scheduler.prepare_file_write_at(&file, buffer, offset); +let sync = scheduler.prepare_file_sync(&file, SyncMode::Data); + +let outcomes = scheduler.submit(write.link(sync), &cancellation).await; +``` + +This requires the proposed file-synchronization and `Linked` extensions. Both operations have known inputs. The intended example policy runs synchronization only after the complete requested write; errors or short writes skip it. `outcomes` must retain the write buffer and distinguish a completed successor from a skipped one. Its exact type and policy-selection syntax remain open, as detailed in the linked-operation section. This does not make a partial write transactional. + +### Example: cooperative service shutdown + +```rust +let system = Cancellation::new(); +let service = system.child(); +let drain = Cancellation::never(); + +defer_cancel( + &service, + server.run(&scheduler, &drain), + || server.request_shutdown(), +).await?; +``` + +`server` is an application-defined service. `request_shutdown` synchronously signals and wakes its accept loop; `run` stops accepting and awaits its owned connection tasks before returning. Process signal integration requests `system.cancel()`, while a local service stop requests `service.cancel()`. The service remains polled during draining, and operations needed to finish accepted work use the non-cancelled drain input. Cancellation-aware accept uses the service's stopping signal or equivalent notification. This example adds no forced abort or shutdown deadline. + +## Linked operations + +Linking is a portable dependency capability. A proposed `Linked` extension trait would express ordered execution of prepared I/O operations, allowing consumers to require it through a bound such as `S: Linked`. A scheduler can implement this capability through kernel linking, a blocking pool, or a readiness state machine. Implement the trait when the backend can preserve the dependency contract; an unavailable capability is a compile-time failure. + +### Preparation and submission + +Prepare the whole chain before submitting it. Preparation retains arguments, buffers, and resources without starting I/O. Submission returns an ordinary future. A chain can use concrete generic types to preserve each operation's result type and allow static dispatch and compiler specialization. + +An illustrative proposed API is: + +```text +let write = scheduler.prepare_file_write_at(&file, buffer, offset); +let sync = scheduler.prepare_file_sync(&file, SyncMode::Data); + +let results = scheduler.submit(write.link(sync), &cancellation).await; +``` + +The names and exact type signatures remain open. Existing single-operation methods can remain convenient wrappers around preparation and submission. This introduces a representation for scheduler I/O that can be composed before execution, while retaining `Future` as the execution interface. + +Ordinary sequential `.await` calls do not expose the whole dependency chain to the backend in advance. An executor cannot inspect an arbitrary future to recover that sequence and turn it into a kernel chain. Explicit prepared operations give the backend the information it needs. + +### Dependency contract + +A chain must define when each successor may start, which preceding outcomes stop execution, and how errors, partial transfers, skipped operations, and cancellation are reported. Return individual outcomes and owned buffers, including those belonging to operations that never started. The exact result representation and continuation policies remain to be chosen. + +Linking supplies ordering without transactional rollback. Kernel links do not automatically substitute an earlier operation's result into a later operation's arguments. A write followed by synchronization can be prepared in advance; a write whose length depends on the actual byte count of a preceding read generally needs a userspace continuation. + +io\_uring's `IOSQE_IO_LINK` starts a successor after its predecessor completes, and breaks the chain on errors or unexpected results, including short reads. Unstarted successors then complete with `-ECANCELED`. `IOSQE_IO_HARDLINK` permits continuation despite completion errors, although submission failures can still break the chain. These dependency policies are distinct from task shutdown. See [kernel link semantics](https://man7.org/linux/man-pages/man2/io_uring_enter.2.html). + +Use native linking only when its behaviour matches the defined chain contract. A fallback must apply the same outcome and continuation rules; submitting independent, unlinked operations concurrently would lose the ordering guarantee. + +### Direct-descriptor linking for accept and receive + +io\_uring supports a native accept-to-receive chain when the accepted socket is installed into a known slot in the ring's registered file table. This uses [io\_uring\_prep\_accept\_direct](https://man7.org/linux/man-pages/man3/io_uring_prep_accept_direct.3.html), rather than expecting an ordinary accept's returned OS descriptor to be substituted into the next SQE. + +The backend can reserve an unused slot before submission, prepare accept with that explicit `file_index` and `IOSQE_IO_LINK`, then prepare receive with its `fd` set to the same slot and `IOSQE_FIXED_FILE` set. Both requests are submitted as one chain. The receive resolves its socket after accept succeeds, allowing execution without userspace processing the accept completion first. + +This requires `IORING_FEAT_LINKED_FILE`, available since Linux 5.17, which defers descriptor lookup for dependent requests until they execute. The backend must also establish direct-accept support and initialize a suitable registered file table. See [dependent descriptor lookup](https://man7.org/linux/man-pages/man2/io_uring_setup.2.html). + +Use an explicitly reserved slot for this native mapping. `IORING_FILE_INDEX_ALLOC` chooses an index dynamically and reports it in the accept completion, so it does not provide the known index needed by a prebuilt successor. An explicit slot must be empty and exclusively reserved: direct accept can replace an existing entry. Direct descriptors are private to their ring, not ordinary OS descriptor numbers. + +The prepared chain must retain the slot through all dependent access and safely release it when the chain is abandoned. If accept succeeds but receive fails or is cancelled, ownership of the accepted socket must still be resolved; skipping the receive does not undo accept. A portable fallback can keep the accepted resource in chain state and pass it to the receive after completion. + +This is a specific resource dependency with a known native target, not general substitution of arbitrary completion results into later SQEs. It does not implement the hypothetical `with_previous_buffer` handoff of a selected buffer and actual received length. + +### Backend implementations + +| Backend | Linking strategy | +| --- | --- | +| io_uring | Submit the prepared operations together as linked SQEs. | +| Blocking pool | Submit the whole chain as one worker job that executes its operations sequentially. | +| Readiness | Advance the chain through successive asynchronous operations as their dependencies complete. | + +A blocking worker can retain the whole chain and return its results after execution, avoiding executor wakeups and job submission between operations. It can check cooperative cancellation between operations and before starting queued work. An already-running blocking syscall may need to finish. Keep resources alive through completion and allow cleanup chains to use cancellation inputs that permit draining. + +For io\_uring, send the chain to the selector as one command, reserve sufficient SQ capacity, and publish its entries contiguously in the same submission. Links cannot cross submission boundaries. Track every member's completion and retain shared resources until all relevant kernel access ends. The current per-operation command model must be extended to support this. See [linked request submission](https://man7.org/linux/man-pages/man7/io_uring_linked_requests.7.html). + +The portable `Linked` capability guarantees dependency behaviour. If a caller specifically requires kernel-linked submission, a separate native extension can express that requirement at compile time. SQE LINK is an implementation strategy for portable linking, rather than a restriction of linking to Linux. + +## Candidate I/O hooks + +Use the [liburing preparation helpers](https://github.com/axboe/liburing/blob/master/src/include/liburing.h) as an inventory of operations to consider. This is a candidate surface, not a requirement that every opcode become part of every scheduler's core traits. + +| Capability area | Candidate operations | +| --- | --- | +| File lifecycle | `file_open`, `file_open_at`, `file_close` | +| File transfer | Stateful `file_read`/`file_write`, positioned `file_read_at`/`file_write_at`, vectored equivalents | +| File management | `file_seek`, `file_allocate`, `file_truncate`, file/data synchronization, range synchronization, access advice | +| Metadata and paths | Stat, extended attributes, mkdir, unlink, rename, hard links, symbolic links | +| Socket lifecycle | Create, bind, listen, accept, connect, shutdown, close | +| Socket transfer | Send/receive, sendto/recvfrom, `socket_sendmsg`/`socket_recvmsg` | +| Socket control | Options, local/peer addresses, readiness | +| Address resolution | `address_resolve` through an appropriate resolver implementation | +| Time and synchronization | Monotonic clock access, sleep, deadlines, wait/wake operations | +| Additional extensions | Process waits, pipes, splice/tee and other descriptor operations | + +The canonical message syscall names are `sendmsg` and `recvmsg`. `writemsg` is not the intended name. + +Not every useful scheduler operation has an io\_uring opcode; seek and address resolution can use other implementation strategies. Multishot operations, fixed descriptors, ring messaging, device commands, and specialized zero-copy operations should remain optional extensions with their own lifecycle contracts. + +### File semantics + +Stateful reads and writes use and advance the file's shared cursor. Positioned operations correspond to `pread`/`pwrite` and should leave that cursor unchanged. Expose both explicitly rather than encoding stateful access through a magic public offset. io\_uring internally supports current-position reads using offset `-1`; its documentation warns that asynchronous shared-cursor access requires serialization for predictable behaviour. See [io\_uring\_prep\_read](https://man7.org/linux/man-pages/man3/io_uring_prep_read.3.html). + +Allocation is distinct from truncation. An allocation interface needs to account for an offset, length, and supported modes. Platform-specific modes may require an extension rather than a universal flags type. See [io\_uring\_prep\_fallocate](https://man7.org/linux/man-pages/man3/io_uring_prep_fallocate.3.html). + +Current `File` only provides positioned reads and writes of already-open, non-append regular files, using offsets that fit in `i64`. Its Windows fallback uses `seek_read`/`seek_write`, which update the shared cursor according to [Windows FileExt](https://doc.rust-lang.org/std/os/windows/fs/trait.FileExt.html). Resolving this difference is required before promising uniform cursor-preserving semantics. + +### Syscall proximity and messages + +Preserve partial transfer counts, relevant flags, addresses, ancillary data, receive status, and OS errors. Keep read-exact/write-all loops and protocol buffering as higher-level conveniences. + +A safe owned message representation should retain its payload buffers and all address, iovec, control, and header storage required by the backend for the duration of kernel access. It can construct native metadata without inherently copying the payload. Specific raw layouts and ancillary features may be platform extensions. See [io\_uring\_prep\_sendmsg](https://man7.org/linux/man-pages/man3/io_uring_prep_sendmsg.3.html). + +Cancellation is a request, not a transactional rollback: a cancelled read may consume bytes and a cancelled write may transmit bytes. Define how callers observe completion and recover buffer ownership before committing to the cancellation-aware result types. + +## Buffer ownership + +Completion-based I/O requires a buffer whose address remains valid while the kernel uses it. The operation must retain buffers and descriptors even if its waiting future is dropped. A cancellation completion alone does not establish that the original operation has finished accessing memory. + +The current API transfers an owned `Vec` and returns `(io::Result, Vec)`. Reads operate on initialized length rather than spare capacity, leave that length unchanged, and report the number of valid bytes separately. Capacity alone is not writable input under this contract. + +There is no standard owned completion-buffer trait to adopt directly. `Vec` and `Box<[u8]>` provide useful owned storage. Standard [`IoSlice`](https://doc.rust-lang.org/std/io/struct.IoSlice.html) and `IoSliceMut` provide borrowed vectored views rather than owning the underlying memory. [`bytes::Bytes` and `BytesMut`](https://docs.rs/bytes/latest/bytes/) provide useful sharing and slicing; `Buf`/`BufMut` alone do not establish completion-I/O lifetime safety. + +Consider small owned buffer traits modeled on the contracts described by Tokio-uring's [`IoBuf`](https://docs.rs/tokio-uring/latest/tokio_uring/buf/trait.IoBuf.html) and [`IoBufMut`](https://docs.rs/tokio-uring/latest/tokio_uring/buf/trait.IoBufMut.html). They would express: + +- A stable backing address while the owned buffer value moves. +- Initialized length distinct from writable capacity. +- Exclusive writable access and a sound way to mark newly initialized bytes. +- Ownership sufficient to keep storage alive through completion. +- `Send` where operations cross Socketry worker or selector threads. + +Support ordinary owned buffers first, with optional implementations for bytes types and later registered-buffer leases. Any unsafe buffer contract must be explicit and enforceable. Supporting uninitialized capacity could avoid initialization work, but is a deliberate change from the current read contract, not an automatic property of passing a `Vec`. + +## Buffer pools, registered buffers, and specialized I/O + +Buffer pools should be optional capabilities that compose with owned-buffer operations. Keep a portable pool contract separate from requirements for native registration or kernel buffer selection. A software pool is an acceptable fallback where it preserves the contract; a consumer specifically requiring a native facility should declare that through an extension trait. Probe kernel-dependent facilities when constructing the pool or registration. + +### Fixed registered buffers + +Fixed registered buffers identify pre-registered memory by index and can reduce repeated memory-mapping work for I/O. A pool can own registrations and issue leases retaining both the memory and its registration lifetime. Validate compatible ring identity before submission and prevent reuse while an operation still accesses the lease. See [buffer registration](https://man7.org/linux/man-pages/man3/io_uring_register_buffers.3.html). + +An operation using a fixed buffer normally receives a lease selected before submission. Registration optimizes access to that memory; it does not defer the assignment of a buffer to a waiting operation. + +### Provided-buffer pools + +Provided-buffer rings solve a different problem: many pending receives can share fewer payload buffers. The kernel selects an available buffer when a receive can make progress, rather than requiring one reserved buffer per idle socket. Registering a provided-buffer ring is distinct from registering fixed payload buffers. See [provided-buffer rings](https://man7.org/linux/man-pages/man3/io_uring_setup_buf_ring.3.html). + +The proposed portable receive interface accepts the pool itself and returns a filled buffer lease, including the valid received length. Acquiring a lease before waiting for data would lose the memory advantage. A lease must distinguish readable initialized bytes from writable capacity and retain its pool and storage until all consumers finish. + +For ordinary single-buffer selection, the lifecycle is: + +1. The pool offers available buffers to the backend. +2. A receive selects a buffer when data is available. +3. Completion identifies the selected buffer and the received byte count. +4. The caller owns a filled lease and may process it or transfer it into a send/write operation. +5. Once no consumer or kernel operation can access the buffer, it can be returned to the pool and offered again. + +Do not offer a leased buffer for another receive while it is still being read or sent. Dropping a waiting future or requesting cancellation must not recycle storage still accessible by the kernel. Pool shutdown must retain outstanding leases, registrations, and offered storage until their respective backend lifetimes have ended. + +Pool exhaustion requires an explicit backpressure and replenishment policy. io\_uring can return `-ENOBUFS` rather than waiting for a buffer; exhaustion can also terminate a multishot receive. The backend must replenish and rearm as appropriate, or expose the error according to the operation contract. See [provided-buffer ownership and exhaustion](https://man7.org/linux/man-pages/man7/io_uring_provided_buffers.7.html). + +A readiness fallback can wait for readability, acquire a lease, and attempt a nonblocking receive. If the receive returns `WouldBlock`, return the lease before waiting again. Waiting for pool availability must also suspend through wakeups. A blocking fallback may retain a lease throughout a running syscall; its bounded concurrency limits how many such leases are occupied. + +Multishot receives could expose a stream of filled leases as a separate extension. The implementation must distinguish individual buffer lifetimes from the lifetime of the receive request, and detect when it needs rearming. See [multishot receive completions](https://man7.org/linux/man-pages/man3/io_uring_prep_recv_multishot.3.html). Streams, bundles, and incremental buffer consumption need their own contracts before implementation. + +### Operation names and preparation + +Use `with_pool` for operations that receive a pool argument. The proposed naming pattern is: + +| Operation | Direct future | Prepared operation | +| --- | --- | --- | +| Receive into an owned buffer | `socket_recv` | `prepare_socket_recv` | +| Receive using a pool | `socket_recv_with_pool` | `prepare_socket_recv_with_pool` | +| Receive a message into owned storage | `socket_recvmsg` | `prepare_socket_recvmsg` | +| Receive a message using a pool | `socket_recvmsg_with_pool` | `prepare_socket_recvmsg_with_pool` | + +These names are proposed, not implemented APIs. Exact arguments, lease types, message metadata, and error/cancellation results remain to be chosen. Direct operations accept explicit cancellation; prepared operations retain the pool and resource inputs, with cancellation supplied at submission as described under linked operations. Preparation must not reserve a payload buffer for each waiting receive. + +`with_pool` describes how an operation obtains storage. The [Rust API Guidelines](https://rust-lang.github.io/api-guidelines/naming.html) use `from_*` for conversions; they do not prescribe `from_pool` for this operation. + +Pool receives can participate in prepared chains whose successors have known inputs and preserve the agreed dependency rules. A successor that needs the selected buffer and actual received length has an additional data dependency that ordinary SQE linking does not resolve. + +### Hypothetical forwarding from a preceding operation + +`with_previous_buffer` is only a hypothetical interface. There is no selected implementation target, agreed signature, or commitment to add it to the proposed `Linked` capability. It could express that a successor consumes the compatible filled lease produced by its predecessor, using the initialized length rather than the original capacity. Any eventual design would need typed compatibility, error and cancellation behaviour, partial-transfer results, and ownership for skipped successors. + +Native SQE links provide ordering, not automatic forwarding of a receive's selected buffer and byte count. Send-side provided-buffer selection is supported on newer kernels, but it selects from an offered group rather than inheriting a preceding receive's result. See [send buffer selection](https://man7.org/linux/man-pages/man3/io_uring_prep_send.3.html). + +Relevant liburing discussions include [using a preceding receive's result as the send length (#58)](https://github.com/axboe/liburing/issues/58) and [moving selected receive buffers into a send ring (#1126)](https://github.com/axboe/liburing/issues/1126). Userspace can reuse the received storage by preparing a send after completion, or offering the filled range to a send buffer ring. This avoids an additional userspace payload copy but still requires userspace to arrange the handoff. + +A selector-managed continuation or a whole-chain blocking job might eventually implement forwarding without waking the application task between operations. These are possible strategies to investigate, not an implementation plan for `with_previous_buffer` or a promise of one kernel-linked submission. + +### Zero-copy lifetime + +Registration does not itself guarantee zero-copy transport. Zero-copy sends can report a transfer result before a subsequent notification allows buffer reuse, so one result is not necessarily the end of kernel access. See [io\_uring\_prep\_send\_zc](https://man7.org/linux/man-pages/man3/io_uring_prep_send_zc.3.html). The current selector's single-terminal-completion request model must be extended before exposing those operations safely. + +## Current implementation and remaining work + +| Area | Current status | +| --- | --- | +| Executor | Future workers, work stealing, explicit owners, task handles and direct-child barriers implemented. | +| Cancellation | `Cancellation`, `Cancelled`, and `defer_cancel` implemented and re-exported. | +| Graceful runtime shutdown | Signal installation, cooperative scheduler draining, and cancellation-aware I/O remain pending; existing shutdown destroys task futures. | +| Capabilities | `File`, `Socket`, `Clock`, and `Spawn` implemented; broader hooks and optional extension boundaries remain proposed. | +| Linked operations | Prepared chains and a portable `Linked` capability remain proposed; native SQE linking, whole-chain blocking jobs, and readiness sequencing are not implemented. | +| File operations | Owned-buffer positioned read/write implemented; stateful I/O, lifecycle, allocation and metadata hooks pending. | +| Socket operations | TCP registration, connect, accept, read/write and readiness implemented; general sockets and message I/O pending. | +| Backends | Shared native readiness, blocking file fallbacks, a Linux io_uring selector, and an adapter to an existing Tokio runtime implemented. | +| Buffers | Owned `Vec` implemented; generic owned-buffer traits, fixed registered leases, and provided/software pools remain proposed. | +| Buffer forwarding | `with_previous_buffer` is hypothetical, with no selected implementation target. | +| Optimization | Dedicated io_uring selector uses an unbounded command channel and a completion channel per operation; pooling and submission backpressure pending. | +| Timers | Existing backend timers provide sleep; the io-event timer algorithm port remains planned. | + +Next design work is to choose the dependable core and extension boundaries, cancellation-aware error and buffer-return contracts, linked-chain continuation rules, resource lifecycle semantics, and owned-buffer trait safety requirements. Then integrate cooperative cancellation into operations and runtime shutdown before expanding optional optimizations. + +The workspace's [architecture guide](https://github.com/socketry/socketry-rust/blob/main/context/design.md) and [implementation guide](https://github.com/socketry/socketry-rust/blob/main/context/implementation.md) provide further details. In the source checkout, `context/gaps.md` records the four architectural gaps. The executor [README](readme.md) documents existing APIs. diff --git a/crates/executor/examples/portable_io.rs b/crates/executor/examples/portable_io.rs index d48aa0e..61faa2b 100644 --- a/crates/executor/examples/portable_io.rs +++ b/crates/executor/examples/portable_io.rs @@ -3,14 +3,14 @@ //! The same TCP exchange, compiled against Socketry or Tokio. #[cfg(any(feature = "native", feature = "tokio"))] -use socketry_executor::{Network, Spawn}; +use socketry_executor::{Socket, Spawn}; #[cfg(any(feature = "native", feature = "tokio"))] use std::{io, net::TcpListener}; #[cfg(any(feature = "native", feature = "tokio"))] async fn exchange(scheduler: SchedulerType) -> io::Result where - SchedulerType: Network + Spawn + Clone + 'static, + SchedulerType: Socket + Spawn + Clone + 'static, { let listener = TcpListener::bind("127.0.0.1:0")?; let address = listener.local_addr()?; diff --git a/crates/executor/readme.md b/crates/executor/readme.md index 8d65fcc..79b8283 100644 --- a/crates/executor/readme.md +++ b/crates/executor/readme.md @@ -1,7 +1,8 @@ # 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. +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. + +The [executor design](design.md) records the agreed direction for cooperative cancellation, scheduler capabilities, portable I/O, and owned buffers, distinguishing implemented APIs from proposals. ## Execution @@ -9,8 +10,7 @@ Use `socketry` as the common entry point, or depend on this package directly. - `Scheduler::with_workers(count)` selects a nonzero worker count explicitly. - `scheduler.spawn(future)` registers an owned task and schedules it immediately. - `handle.await` returns `Result`. -- `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.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. @@ -35,165 +35,111 @@ fn main() -> Result<(), Box> { } ``` -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. +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. +`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. +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. +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. +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. +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. +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().await` requests cancellation and awaits future destruction. - `barrier.close()` prevents new children while existing children continue. -- `barrier.wait().await` waits until there are no direct children. Close first - when the set of children must remain closed. +- `barrier.wait().await` waits until there are no direct children. Close first when the set of children must remain closed. - `barrier.stop().await` closes, 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. +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. + +## Cooperative cancellation + +`Cancellation` is a runtime-independent shutdown signal. It does not own tasks or destroy futures. `cancel()` requests cancellation and returns true for the first request on that signal; subsequent requests return false and never escalate. `is_cancelled()` observes the persistent state, `check()` returns `Err(Cancelled)` when cancelled, and `cancelled().await` waits using ordinary polling and wakeups. + +Clones share state. `child()` creates an independent cancellation boundary that also receives ancestor cancellation. Cancelling a child leaves its parent and siblings running. Descendants retain this relationship even after intermediate handles are dropped. Dropping signals never requests cancellation. A child created after its parent is cancelled starts cancelled. `Cancellation::never()` requires no allocation and remains uncancelled; its children are independent cancellable signals. + +`defer_cancel(&signal, future, on_cancel)` invokes a synchronous `FnOnce()` callback when cancellation is observed, then continues polling the protected future to completion. Keep asynchronous cleanup in that future; the callback should request shutdown or wake it. The callback runs before the next protected poll, including the first poll when the signal is already cancelled. If both cancellation and completion are observed in the same poll, the callback runs and the wrapper returns the future's normal output. Cancellation arriving during the protected poll can lose that race to completion. -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. +```rust +use socketry_executor::{Cancellation, Scheduler, defer_cancel, yield_now}; + +let scheduler = Scheduler::with_workers(1)?; +let shutdown = Cancellation::new(); +let drain = Cancellation::new(); +shutdown.cancel(); +let output = scheduler.block_on(defer_cancel(&shutdown, async { + drain.cancelled().await; + yield_now().await; // Asynchronous cleanup remains executable. + 42 +}, || { drain.cancel(); })); +assert_eq!(output, 42); +# Ok::<(), std::io::Error>(()) +``` -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. +Use a root signal for process shutdown and child signals for individual services. Keep task owners and the runtime alive until task handles or barriers confirm draining is complete. The protected work must use cancellation inputs that allow cleanup to continue; the wrapper does not mask signals passed into its operations. Callback panics propagate, and dropping the wrapper still drops its work. -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. +These primitives do not change `Task::cancel`, `Barrier::stop`, or scheduler shutdown, which still destroy futures. Signal handlers and cancellation-aware I/O signatures remain separate integration work. Cooperative requests have no forced abort or deadline escalation; a process supervisor can enforce SIGTERM followed by SIGKILL. ## 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. +`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. +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. +Each cancellable signal allocates a reference-counted node. A signal family shares a mutex for child registration and cancellation propagation; parent nodes hold weak child registrations, and children retain ancestors. Cancellation walks descendants iteratively and releases locks before waking listeners. Waiting uses the existing event-listener dependency. Dropped children remove their registrations, and deep ancestor chains are released iteratively. Ordinary signal checks use an atomic flag; `never()` needs no node or listener. + +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-uring` selects 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. -- `tokio` adds an adapter using a supplied `tokio::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 `tokio` avoids 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` 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` 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. +The public `scheduler` module contains portable `Socket`, `File`, 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-uring` selects 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. +- `tokio` adds an adapter using a supplied `tokio::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 `tokio` avoids 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` 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. + +`File::file_read_at` and `file_write_at` accept an `Arc` 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. +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. +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. +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. +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. +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. +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: @@ -205,8 +151,7 @@ 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`. +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: diff --git a/crates/executor/src/cancellation.rs b/crates/executor/src/cancellation.rs new file mode 100644 index 0000000..ada1c51 --- /dev/null +++ b/crates/executor/src/cancellation.rs @@ -0,0 +1,185 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +use crate::Cancelled; +use event_listener::Event; +use std::collections::HashMap; +use std::future::pending; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{Arc, Mutex, PoisonError, Weak}; + +struct Node { + cancelled: AtomicBool, + event: Event, + // Serialize child registration and cancellation across this signal family. + tree: Arc>, + // Children retain ancestors, but ancestors only hold weak child references. + parent: Option>, + children: Mutex>>, +} + +impl Node { + fn new(tree: Arc>, parent: Option>) -> Self { + Self { + cancelled: AtomicBool::new(false), + event: Event::new(), + tree, + parent, + children: Mutex::new(HashMap::new()), + } + } +} + +impl Drop for Node { + fn drop(&mut self) { + let mut identifier = self as *const Self as usize; + let mut parent = self.parent.take(); + // Release an exclusively owned ancestor chain without recursive drops. + while let Some(node) = parent { + // The allocation's address identifies its weak registration until + // destruction removes it, so recycled addresses cannot alias it. + node.children + .lock() + .unwrap_or_else(PoisonError::into_inner) + .remove(&identifier); + identifier = Arc::as_ptr(&node) as usize; + match Arc::try_unwrap(node) { + Ok(mut node) => { + parent = node.parent.take(); + } + Err(_) => break, + } + } + } +} + +/// A runtime-independent, cooperative cancellation signal. +/// +/// Clones share a persistent cancellation state. Child signals receive ancestor +/// cancellation, but cancelling a child does not cancel its parent or siblings. +/// Dropping a signal does not request cancellation. Signals neither own tasks nor +/// abort futures; callers decide how to respond and separately await completion. +/// +/// ``` +/// use socketry_executor::{Cancellation, Cancelled}; +/// +/// let system = Cancellation::new(); +/// let server = system.child(); +/// assert!(server.cancel()); +/// assert_eq!(server.check(), Err(Cancelled)); +/// assert!(!system.is_cancelled()); +/// assert!(system.cancel()); +/// ``` +#[derive(Clone)] +pub struct Cancellation { + node: Option>, +} + +impl Cancellation { + /// Create an independent, uncancelled signal. + pub fn new() -> Self { + Self { + node: Some(Arc::new(Node::new(Arc::new(Mutex::new(())), None))), + } + } + + /// Create a signal which cannot be cancelled and requires no allocation. + /// `cancel` always returns false and `cancelled` remains pending forever. + pub const fn never() -> Self { + Self { node: None } + } + + /// Create a signal cancelled by this signal or any of its ancestors. + /// A child of an already-cancelled signal starts cancelled. A child of a + /// never-cancelled signal is an independent, cancellable signal. + pub fn child(&self) -> Self { + let Some(parent) = &self.node else { + return Self::new(); + }; + let _tree = parent.tree.lock().unwrap_or_else(PoisonError::into_inner); + let child = Arc::new(Node::new( + Arc::clone(&parent.tree), + Some(Arc::clone(parent)), + )); + if parent.cancelled.load(Ordering::Acquire) { + child.cancelled.store(true, Ordering::Release); + } else { + parent + .children + .lock() + .unwrap_or_else(PoisonError::into_inner) + .insert(Arc::as_ptr(&child) as usize, Arc::downgrade(&child)); + } + Self { node: Some(child) } + } + + /// Request cancellation of this signal and all its descendants. + /// Returns true only for the first request on this signal. Repeated requests + /// never escalate. Waiters are woken after the entire family lock is released. + /// Concurrent observers may see propagation in progress, but all descendants + /// are cancelled when this call returns. This does not wait for their work. + pub fn cancel(&self) -> bool { + let Some(root) = &self.node else { + return false; + }; + let tree = root.tree.lock().unwrap_or_else(PoisonError::into_inner); + if root.cancelled.load(Ordering::Acquire) { + return false; + } + let mut pending = vec![Arc::clone(root)]; + let mut cancelled = Vec::new(); + while let Some(node) = pending.pop() { + node.cancelled.store(true, Ordering::Release); + pending.extend( + std::mem::take(&mut *node.children.lock().unwrap_or_else(PoisonError::into_inner)) + .into_values() + .filter_map(|child| child.upgrade()), + ); + cancelled.push(node); + } + drop(tree); + for node in cancelled { + node.event.notify(usize::MAX); + } + true + } + + /// Whether cancellation has been requested. + pub fn is_cancelled(&self) -> bool { + self.node + .as_ref() + .is_some_and(|node| node.cancelled.load(Ordering::Acquire)) + } + + /// Check for cancellation at an explicit cooperative cancellation point. + pub fn check(&self) -> Result<(), Cancelled> { + if self.is_cancelled() { + Err(Cancelled) + } else { + Ok(()) + } + } + + /// Wait for cancellation using standard future polling and wakeups. + /// This future is immediately ready if already cancelled. Dropping the wait + /// unregisters its listener without cancelling the signal or other waiters. + pub async fn cancelled(&self) { + let Some(node) = &self.node else { + return pending().await; + }; + let listener = node.event.listen(); + // Register before checking to preserve a racing cancellation notification. + if !node.cancelled.load(Ordering::Acquire) { + listener.await; + } + } +} + +impl Default for Cancellation { + fn default() -> Self { + Self::new() + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/executor/src/cancellation/tests.rs b/crates/executor/src/cancellation/tests.rs new file mode 100644 index 0000000..1750a82 --- /dev/null +++ b/crates/executor/src/cancellation/tests.rs @@ -0,0 +1,39 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +use super::*; + +#[test] +fn dropped_children_release_their_registrations_and_ancestors() { + let root = Cancellation::new(); + let node = root.node.as_ref().unwrap(); + for _ in 0..100 { + let child = root.child(); + let grandchild = child.child(); + drop(child); + drop(grandchild); + assert!(node.children.lock().unwrap().is_empty()); + assert_eq!(Arc::strong_count(node), 1); + } +} + +#[test] +fn an_exclusively_owned_deep_chain_is_released_without_recursive_destruction() { + std::thread::Builder::new() + .stack_size(64 * 1024) + .spawn(|| { + let root = Cancellation::new(); + let weak = Arc::downgrade(root.node.as_ref().unwrap()); + let mut leaf = root.clone(); + for _ in 0..10_000 { + leaf = leaf.child(); + } + drop(root); + assert!(weak.upgrade().is_some()); + drop(leaf); + assert!(weak.upgrade().is_none()); + }) + .unwrap() + .join() + .unwrap(); +} diff --git a/crates/executor/src/cancelled.rs b/crates/executor/src/cancelled.rs new file mode 100644 index 0000000..88f8762 --- /dev/null +++ b/crates/executor/src/cancelled.rs @@ -0,0 +1,16 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +use std::fmt; + +/// An operation observed a cooperative cancellation request. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct Cancelled; + +impl fmt::Display for Cancelled { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("operation cancelled") + } +} + +impl std::error::Error for Cancelled {} diff --git a/crates/executor/src/defer_cancel.rs b/crates/executor/src/defer_cancel.rs new file mode 100644 index 0000000..66b9ecd --- /dev/null +++ b/crates/executor/src/defer_cancel.rs @@ -0,0 +1,56 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +use crate::Cancellation; +use std::future::{Future, poll_fn}; +use std::pin::pin; + +/// Notify protected work of cancellation and continue awaiting its completion. +/// +/// When cancellation is observed, call `on_cancel` exactly once before +/// polling the protected future again. This includes an already-cancelled signal +/// on the first poll. If completion and cancellation are both ready in that poll, +/// the callback runs before returning the future's normal output. A callback +/// panic propagates to the caller. +/// +/// The callback is synchronous: request graceful shutdown and wake the work; +/// keep asynchronous draining in the protected future. Its operations must use +/// cancellation inputs which allow draining. This wrapper does not mask signals, +/// prevent its own destruction, or intercept scheduler task cancellation. +/// +/// ``` +/// use socketry_executor::{Cancellation, Scheduler, defer_cancel, yield_now}; +/// +/// let scheduler = Scheduler::with_workers(1)?; +/// let shutdown = Cancellation::new(); +/// let drain = Cancellation::new(); +/// shutdown.cancel(); +/// let output = scheduler.block_on(defer_cancel(&shutdown, async { +/// drain.cancelled().await; +/// yield_now().await; +/// 42 +/// }, || { drain.cancel(); })); +/// assert_eq!(output, 42); +/// # Ok::<(), std::io::Error>(()) +/// ``` +pub async fn defer_cancel( + cancellation: &Cancellation, + future: FutureType, + on_cancel: OnCancel, +) -> FutureType::Output +where + FutureType: Future, + OnCancel: FnOnce(), +{ + let mut notification = pin!(cancellation.cancelled()); + let mut future = pin!(future); + let mut on_cancel = Some(on_cancel); + poll_fn(move |context| { + if on_cancel.is_some() && notification.as_mut().poll(context).is_ready() { + // Remove the callback before invoking user code, including on panic. + on_cancel.take().expect("cancellation callback is present")(); + } + future.as_mut().poll(context) + }) + .await +} diff --git a/crates/executor/src/lib.rs b/crates/executor/src/lib.rs index 3efe239..062bf28 100644 --- a/crates/executor/src/lib.rs +++ b/crates/executor/src/lib.rs @@ -9,12 +9,18 @@ #![doc = include_str!("../readme.md")] mod barrier; +mod cancellation; +mod cancelled; +mod defer_cancel; mod owner; pub mod scheduler; mod task; mod worker; pub use barrier::Barrier; +pub use cancellation::Cancellation; +pub use cancelled::Cancelled; +pub use defer_cancel::defer_cancel; pub use owner::{Spawn, SpawnError}; -pub use scheduler::{BufferResult, Clock, FileIo, Interest, Network, Scheduler, SchedulerHandle}; +pub use scheduler::{BufferResult, Clock, File, Interest, Scheduler, SchedulerHandle, Socket}; pub use task::{Task, TaskError, TaskHandle, yield_now}; diff --git a/crates/executor/src/scheduler.rs b/crates/executor/src/scheduler.rs index 7c75073..503a8e6 100644 --- a/crates/executor/src/scheduler.rs +++ b/crates/executor/src/scheduler.rs @@ -8,7 +8,7 @@ pub mod selector; pub mod socketry; #[cfg(any(feature = "native", feature = "tokio"))] -mod file; +mod positioned_file; #[cfg(feature = "tokio")] pub mod tokio; @@ -26,11 +26,11 @@ pub type BufferResult = (io::Result, Vec); mod interest; pub use interest::Interest; -mod network; -pub use network::Network; +mod socket; +pub use socket::Socket; -mod file_io; -pub use file_io::FileIo; +mod file; +pub use file::File; mod clock; pub use clock::Clock; diff --git a/crates/executor/src/scheduler/file.rs b/crates/executor/src/scheduler/file.rs index c8c27e0..acc4d22 100644 --- a/crates/executor/src/scheduler/file.rs +++ b/crates/executor/src/scheduler/file.rs @@ -1,62 +1,36 @@ // Released under the MIT License. // Copyright, 2026, by Samuel Williams. -//! Blocking positioned file operations shared by selector adapters. -use std::fs::File; -use std::io; +use super::BufferResult; +use std::fs::File as StdFile; +use std::future::Future; +use std::sync::Arc; -pub(crate) fn read_at(file: &File, buffer: &mut [u8], offset: u64) -> io::Result { - if offset > i64::MAX as u64 { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - "file offset exceeds i64::MAX", - )); - } - let operation = || { - #[cfg(unix)] - { - use std::os::unix::fs::FileExt; - file.read_at(buffer, offset) - } - #[cfg(windows)] - { - use std::os::windows::fs::FileExt; - file.seek_read(buffer, offset) - } - }; - retry_interrupted(operation) -} - -pub(crate) fn write_at(file: &File, buffer: &[u8], offset: u64) -> io::Result { - if offset > i64::MAX as u64 { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - "file offset exceeds i64::MAX", - )); - } - let operation = || { - #[cfg(unix)] - { - use std::os::unix::fs::FileExt; - file.write_at(buffer, offset) - } - #[cfg(windows)] - { - use std::os::windows::fs::FileExt; - file.seek_write(buffer, offset) - } - }; - retry_interrupted(operation) -} +/// A scheduler's positioned file operations. +/// +/// This trait is implemented by schedulers. File resources are [`std::fs::File`]; +/// alias that type as `StdFile` when importing both names. +/// +/// 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 File: Send + Sync { + fn file_read_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> impl Future + Send; -fn retry_interrupted(mut operation: impl FnMut() -> io::Result) -> io::Result { - loop { - match operation() { - Err(error) if error.kind() == io::ErrorKind::Interrupted => continue, - result => return result, - } - } + fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> impl Future + Send; } - -#[cfg(test)] -mod tests; diff --git a/crates/executor/src/scheduler/file_io.rs b/crates/executor/src/scheduler/file_io.rs deleted file mode 100644 index 9d20faa..0000000 --- a/crates/executor/src/scheduler/file_io.rs +++ /dev/null @@ -1,31 +0,0 @@ -// 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, - buffer: Vec, - offset: u64, - ) -> impl Future + Send; - - fn file_write_at( - &self, - file: Arc, - buffer: Vec, - offset: u64, - ) -> impl Future + Send; -} diff --git a/crates/executor/src/scheduler/positioned_file.rs b/crates/executor/src/scheduler/positioned_file.rs new file mode 100644 index 0000000..c8c27e0 --- /dev/null +++ b/crates/executor/src/scheduler/positioned_file.rs @@ -0,0 +1,62 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +//! Blocking positioned file operations shared by selector adapters. +use std::fs::File; +use std::io; + +pub(crate) fn read_at(file: &File, buffer: &mut [u8], offset: u64) -> io::Result { + if offset > i64::MAX as u64 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "file offset exceeds i64::MAX", + )); + } + let operation = || { + #[cfg(unix)] + { + use std::os::unix::fs::FileExt; + file.read_at(buffer, offset) + } + #[cfg(windows)] + { + use std::os::windows::fs::FileExt; + file.seek_read(buffer, offset) + } + }; + retry_interrupted(operation) +} + +pub(crate) fn write_at(file: &File, buffer: &[u8], offset: u64) -> io::Result { + if offset > i64::MAX as u64 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "file offset exceeds i64::MAX", + )); + } + let operation = || { + #[cfg(unix)] + { + use std::os::unix::fs::FileExt; + file.write_at(buffer, offset) + } + #[cfg(windows)] + { + use std::os::windows::fs::FileExt; + file.seek_write(buffer, offset) + } + }; + retry_interrupted(operation) +} + +fn retry_interrupted(mut operation: impl FnMut() -> io::Result) -> io::Result { + loop { + match operation() { + Err(error) if error.kind() == io::ErrorKind::Interrupted => continue, + result => return result, + } + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/executor/src/scheduler/file/tests.rs b/crates/executor/src/scheduler/positioned_file/tests.rs similarity index 100% rename from crates/executor/src/scheduler/file/tests.rs rename to crates/executor/src/scheduler/positioned_file/tests.rs diff --git a/crates/executor/src/scheduler/selector/io_uring.rs b/crates/executor/src/scheduler/selector/io_uring.rs index 58c2065..429a17a 100644 --- a/crates/executor/src/scheduler/selector/io_uring.rs +++ b/crates/executor/src/scheduler/selector/io_uring.rs @@ -15,13 +15,13 @@ //! channel per operation. It does not yet pool operation records or register //! buffers with the kernel. use super::readiness::{self, Listener, Socket}; -use crate::scheduler::{BufferResult, Clock, FileIo, Interest, Network}; +use crate::scheduler::{BufferResult, Clock, File, Interest, Socket as SocketOperations}; use event_listener::{Event as CompletionEvent, Listener as _}; use futures_channel::oneshot; use io_uring::{IoUring, opcode, squeue, types}; use polling::{Event, Events, Poller}; use std::collections::{HashMap, VecDeque}; -use std::fs::File; +use std::fs::File as StdFile; use std::io::{self, Read, Write}; use std::net::{SocketAddr, TcpListener, TcpStream}; use std::os::fd::AsRawFd; @@ -105,7 +105,7 @@ where enum Resource { Socket(Socket), - File(Arc, u64), + File(Arc, u64), } struct Request { @@ -430,7 +430,7 @@ impl Selector { } } -impl Network for Selector { +impl SocketOperations for Selector { type Socket = Socket; type Listener = Listener; @@ -472,12 +472,17 @@ impl Network for Selector { } } -impl FileIo for Selector { - async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { +impl File for Selector { + async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { self.operation(Resource::File(file, offset), buffer, false) .await } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { self.operation(Resource::File(file, offset), buffer, true) .await } diff --git a/crates/executor/src/scheduler/selector/io_uring/tests.rs b/crates/executor/src/scheduler/selector/io_uring/tests.rs index fe3f2ec..9b296b5 100644 --- a/crates/executor/src/scheduler/selector/io_uring/tests.rs +++ b/crates/executor/src/scheduler/selector/io_uring/tests.rs @@ -11,7 +11,7 @@ fn queued_request_returns_its_buffer_if_the_selector_exits() { let request = Request { identifier: 0, resource: Resource::File( - Arc::new(File::open(std::env::current_exe().unwrap()).unwrap()), + Arc::new(StdFile::open(std::env::current_exe().unwrap()).unwrap()), 0, ), buffer, diff --git a/crates/executor/src/scheduler/selector/readiness.rs b/crates/executor/src/scheduler/selector/readiness.rs index 3484159..68eaef8 100644 --- a/crates/executor/src/scheduler/selector/readiness.rs +++ b/crates/executor/src/scheduler/selector/readiness.rs @@ -7,10 +7,10 @@ //! It remains available while registered resources are alive; Socketry task //! shutdown does not shut down that shared reactor. Futures contain no private //! coroutine stacks and may be polled on different workers. -use crate::scheduler::file::{read_at, write_at}; -use crate::scheduler::{BufferResult, Clock, FileIo, Interest, Network}; +use crate::scheduler::positioned_file::{read_at, write_at}; +use crate::scheduler::{BufferResult, Clock, File, Interest, Socket as SocketOperations}; use async_io::Async; -use std::fs::File; +use std::fs::File as StdFile; use std::io::{self, Read, Write}; use std::net::{SocketAddr, TcpListener, TcpStream}; use std::sync::Arc; @@ -64,7 +64,7 @@ fn wrap_accepted( accepted.map(|(socket, address)| (Socket(Arc::new(socket)), address)) } -impl Network for Selector { +impl SocketOperations for Selector { type Socket = Socket; type Listener = Listener; @@ -108,10 +108,10 @@ impl Network for Selector { } } -impl FileIo for Selector { +impl File for Selector { async fn file_read_at( &self, - file: Arc, + file: Arc, mut buffer: Vec, offset: u64, ) -> BufferResult { @@ -122,7 +122,12 @@ impl FileIo for Selector { .await } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { blocking::unblock(move || { let result = write_at(&file, &buffer, offset); (result, buffer) diff --git a/crates/executor/src/scheduler/network.rs b/crates/executor/src/scheduler/socket.rs similarity index 91% rename from crates/executor/src/scheduler/network.rs rename to crates/executor/src/scheduler/socket.rs index 605acbd..0a186f2 100644 --- a/crates/executor/src/scheduler/network.rs +++ b/crates/executor/src/scheduler/socket.rs @@ -8,6 +8,9 @@ use std::net::{SocketAddr, TcpListener, TcpStream}; /// Portable socket operations, selected through the concrete implementation. /// +/// This trait is implemented by schedulers. Its associated [`Socket::Socket`] +/// and [`Socket::Listener`] types represent their registered resources. +/// /// Registrations belong to their creating implementation. Keep a socket's /// registration across operations and worker migration. Implementations must /// return an error when a resource belongs to an incompatible runtime instance. @@ -17,7 +20,7 @@ use std::net::{SocketAddr, TcpListener, TcpStream}; /// bytes, and a cancelled write can transmit bytes. Implementations retain any /// kernel-accessible memory until the operation finishes. Await completion when /// the amount transferred matters. No asynchronous cleanup is promised by Drop. -pub trait Network: Send + Sync { +pub trait Socket: Send + Sync { type Socket: Send + Sync; type Listener: Send + Sync; diff --git a/crates/executor/src/scheduler/socketry/operations.rs b/crates/executor/src/scheduler/socketry/operations.rs index 2d61e6a..110d511 100644 --- a/crates/executor/src/scheduler/socketry/operations.rs +++ b/crates/executor/src/scheduler/socketry/operations.rs @@ -4,8 +4,8 @@ //! Native selector operations exposed by Socketry schedulers and handles. use super::{Scheduler, SchedulerHandle}; use crate::scheduler::selector::DefaultSelector; -use crate::scheduler::{BufferResult, Clock, FileIo, Interest, Network}; -use std::fs::File; +use crate::scheduler::{BufferResult, Clock, File, Interest, Socket}; +use std::fs::File as StdFile; use std::io; use std::net::{SocketAddr, TcpListener, TcpStream}; use std::sync::Arc; @@ -75,9 +75,9 @@ impl SchedulerHandle { } } -impl Network for SchedulerHandle { - type Socket = ::Socket; - type Listener = ::Listener; +impl Socket for SchedulerHandle { + type Socket = ::Socket; + type Listener = ::Listener; fn register_socket(&self, socket: TcpStream) -> io::Result { self.selector()?.register_socket(socket) @@ -114,15 +114,20 @@ impl Network for SchedulerHandle { } } -impl FileIo for SchedulerHandle { - async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { +impl File for SchedulerHandle { + async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { match self.selector() { Ok(selector) => selector.file_read_at(file, buffer, offset).await, Err(error) => (Err(error), buffer), } } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { match self.selector() { Ok(selector) => selector.file_write_at(file, buffer, offset).await, Err(error) => (Err(error), buffer), @@ -137,9 +142,9 @@ impl Clock for SchedulerHandle { } } -impl Network for Scheduler { - type Socket = ::Socket; - type Listener = ::Listener; +impl Socket for Scheduler { + type Socket = ::Socket; + type Listener = ::Listener; fn register_socket(&self, socket: TcpStream) -> io::Result { self.handle.register_socket(socket) @@ -170,12 +175,17 @@ impl Network for Scheduler { } } -impl FileIo for Scheduler { - async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { +impl File for Scheduler { + async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { self.handle.file_read_at(file, buffer, offset).await } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { self.handle.file_write_at(file, buffer, offset).await } } diff --git a/crates/executor/src/scheduler/socketry/operations/tests.rs b/crates/executor/src/scheduler/socketry/operations/tests.rs index 28af597..747b7e3 100644 --- a/crates/executor/src/scheduler/socketry/operations/tests.rs +++ b/crates/executor/src/scheduler/socketry/operations/tests.rs @@ -2,7 +2,7 @@ // Copyright, 2026, by Samuel Williams. use super::{DefaultSelector, Scheduler, check_initialized_open, check_open}; -use crate::scheduler::{Clock, FileIo, Interest, Network}; +use crate::scheduler::{Clock, File, Interest, Socket}; use std::fs::OpenOptions; use std::io; use std::net::{SocketAddr, TcpListener, TcpStream}; diff --git a/crates/executor/src/scheduler/tokio.rs b/crates/executor/src/scheduler/tokio.rs index 7f1d128..69dc59b 100644 --- a/crates/executor/src/scheduler/tokio.rs +++ b/crates/executor/src/scheduler/tokio.rs @@ -7,7 +7,7 @@ //! runtime. Dropping it closes admission and requests cancellation. Await //! shutdown to join task destruction. Socketry's Task::current and //! Scheduler::current describe Socketry execution, not Tokio tasks. -use super::{BufferResult, Clock, FileIo, Interest, Network}; +use super::{BufferResult, Clock, File, Interest, Socket as SocketOperations}; use crate::owner::Owner; use crate::{Spawn, SpawnError, TaskError}; use ::tokio::runtime::Handle; @@ -15,7 +15,7 @@ use ::tokio::task::{AbortHandle, JoinHandle}; use pin_project_lite::pin_project; use std::cell::RefCell; use std::collections::HashMap; -use std::fs::File; +use std::fs::File as StdFile; use std::future::Future; use std::io; use std::net::{SocketAddr, TcpListener, TcpStream}; @@ -619,7 +619,7 @@ impl Future for InRuntime { } } -impl Network for SchedulerHandle { +impl SocketOperations for SchedulerHandle { type Socket = Socket; type Listener = Listener; @@ -706,12 +706,17 @@ impl Network for SchedulerHandle { } } -impl FileIo for SchedulerHandle { - async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { +impl File for SchedulerHandle { + async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { self.file_operation(file, buffer, offset, false).await } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { self.file_operation(file, buffer, offset, true).await } } @@ -719,7 +724,7 @@ impl FileIo for SchedulerHandle { impl SchedulerHandle { async fn file_operation( &self, - file: Arc, + file: Arc, buffer: Vec, offset: u64, write: bool, @@ -737,9 +742,9 @@ impl SchedulerHandle { .spawn_blocking(move || { with_file_buffer(&operation_buffer, |buffer| { if write { - super::file::write_at(&file, buffer, offset) + super::positioned_file::write_at(&file, buffer, offset) } else { - super::file::read_at(&file, buffer, offset) + super::positioned_file::read_at(&file, buffer, offset) } }) }) @@ -811,7 +816,7 @@ impl Spawn for Barrier { } } -impl Network for Scheduler { +impl SocketOperations for Scheduler { type Socket = Socket; type Listener = Listener; fn register_socket(&self, socket: TcpStream) -> io::Result { @@ -836,11 +841,16 @@ impl Network for Scheduler { self.handle.io_wait(socket, interest).await } } -impl FileIo for Scheduler { - async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { +impl File for Scheduler { + async fn file_read_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { self.handle.file_read_at(file, buffer, offset).await } - async fn file_write_at(&self, file: Arc, buffer: Vec, offset: u64) -> BufferResult { + async fn file_write_at( + &self, + file: Arc, + buffer: Vec, + offset: u64, + ) -> BufferResult { self.handle.file_write_at(file, buffer, offset).await } } diff --git a/crates/executor/tests/cancellation.rs b/crates/executor/tests/cancellation.rs new file mode 100644 index 0000000..4d7c13f --- /dev/null +++ b/crates/executor/tests/cancellation.rs @@ -0,0 +1,553 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +mod support; + +use socketry_executor::{Cancellation, Cancelled, Scheduler, defer_cancel, yield_now}; +use std::cell::Cell; +use std::future::{Future, pending, poll_fn}; +use std::marker::PhantomPinned; +use std::pin::{Pin, pin}; +use std::rc::Rc; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Barrier, mpsc}; +use std::task::{Context, Poll, Wake, Waker}; +use support::{CountDrop, receive}; + +#[derive(Default)] +struct WakeCounter(AtomicUsize); + +impl Wake for WakeCounter { + fn wake(self: Arc) { + self.wake_by_ref(); + } + + fn wake_by_ref(self: &Arc) { + self.0.fetch_add(1, Ordering::SeqCst); + } +} + +fn poll( + future: Pin<&mut FutureType>, + waker: &Waker, +) -> Poll { + future.poll(&mut Context::from_waker(waker)) +} + +#[test] +fn clones_share_persistent_cancellation_and_a_descriptive_error() { + fn thread_safe() {} + fn send(_: Value) {} + thread_safe::(); + thread_safe::(); + + let token = Cancellation::default(); + send(token.cancelled()); + let clone = token.clone(); + assert_eq!(clone.check(), Ok(())); + assert!(token.cancel()); + assert!(clone.is_cancelled()); + assert!(!clone.cancel()); + assert_eq!(clone.check(), Err(Cancelled)); + assert_eq!(Cancelled.to_string(), "operation cancelled"); + let error: &dyn std::error::Error = &Cancelled; + assert!(error.source().is_none()); + + let mut notification = pin!(clone.cancelled()); + assert!(poll(notification.as_mut(), Waker::noop()).is_ready()); +} + +#[test] +fn independent_and_child_signals_have_separate_cancellation_boundaries() { + let root = Cancellation::new(); + let child = root.child(); + let grandchild = child.child(); + let sibling = root.child(); + let independent = Cancellation::new(); + assert!(child.cancel()); + assert!(grandchild.is_cancelled()); + assert!(!root.is_cancelled()); + assert!(!sibling.is_cancelled()); + assert!(root.cancel()); + assert!(sibling.is_cancelled()); + assert!(!independent.is_cancelled()); + assert!(!child.cancel()); + assert!(!grandchild.cancel()); + assert!(child.child().child().is_cancelled()); +} + +#[test] +fn ancestor_cancellation_reaches_children_after_intermediate_handles_are_dropped() { + let root = Cancellation::new(); + let intermediate = root.child(); + let child = intermediate.child(); + drop(intermediate); + assert!(!child.is_cancelled()); + root.cancel(); + assert!(child.is_cancelled()); +} + +#[test] +fn dropping_all_parent_handles_does_not_cancel_surviving_children() { + let root = Cancellation::new(); + let child = root.child(); + drop(root.clone()); + drop(root); + assert!(!child.is_cancelled()); + assert!(child.cancel()); +} + +#[test] +fn never_cancelled_tokens_remain_pending_and_can_create_independent_children() { + let token = Cancellation::never(); + let counter = Arc::new(WakeCounter::default()); + let waker = Waker::from(Arc::clone(&counter)); + let mut notification = pin!(token.cancelled()); + assert!(poll(notification.as_mut(), &waker).is_pending()); + assert!(!token.cancel()); + assert!(!token.clone().cancel()); + assert!(!token.is_cancelled()); + assert_eq!(token.check(), Ok(())); + assert_eq!(counter.0.load(Ordering::SeqCst), 0); + let child = token.child(); + assert!(child.cancel()); + assert!(!token.is_cancelled()); + assert!(poll(notification.as_mut(), &waker).is_pending()); +} + +#[test] +fn cancellation_wakes_every_registered_waiter_including_descendants() { + let root = Cancellation::new(); + let tokens = [root.clone(), root.clone(), root.child()]; + let counters: Vec<_> = tokens + .iter() + .map(|_| Arc::new(WakeCounter::default())) + .collect(); + let wakers: Vec<_> = counters + .iter() + .map(|counter| Waker::from(Arc::clone(counter))) + .collect(); + let mut notifications: Vec<_> = tokens + .iter() + .map(|token| Box::pin(token.cancelled())) + .collect(); + for (notification, waker) in notifications.iter_mut().zip(&wakers) { + assert!(poll(notification.as_mut(), waker).is_pending()); + } + assert!(root.cancel()); + for ((notification, waker), counter) in notifications.iter_mut().zip(&wakers).zip(&counters) { + assert_eq!(counter.0.load(Ordering::SeqCst), 1); + assert!(poll(notification.as_mut(), waker).is_ready()); + } + assert!(!root.cancel()); + assert!( + counters + .iter() + .all(|counter| counter.0.load(Ordering::SeqCst) == 1) + ); +} + +#[test] +fn dropping_a_wait_unregisters_it_without_affecting_other_waiters() { + let token = Cancellation::new(); + let abandoned = Arc::new(WakeCounter::default()); + let retained = Arc::new(WakeCounter::default()); + let mut first = Box::pin(token.cancelled()); + let mut second = pin!(token.cancelled()); + let retained_waker = Waker::from(Arc::clone(&retained)); + assert!(poll(first.as_mut(), &Waker::from(Arc::clone(&abandoned))).is_pending()); + assert!(poll(second.as_mut(), &retained_waker).is_pending()); + drop(first); + assert!(!token.is_cancelled()); + token.cancel(); + assert_eq!(abandoned.0.load(Ordering::SeqCst), 0); + assert_eq!(retained.0.load(Ordering::SeqCst), 1); + assert!(poll(second.as_mut(), &retained_waker).is_ready()); +} + +#[test] +fn a_wait_uses_the_waker_from_its_latest_poll() { + let token = Cancellation::new(); + let old = Arc::new(WakeCounter::default()); + let new = Arc::new(WakeCounter::default()); + let mut notification = pin!(token.cancelled()); + let new_waker = Waker::from(Arc::clone(&new)); + assert!(poll(notification.as_mut(), &Waker::from(Arc::clone(&old))).is_pending()); + assert!(poll(notification.as_mut(), &new_waker).is_pending()); + token.cancel(); + assert_eq!(old.0.load(Ordering::SeqCst), 0); + assert_eq!(new.0.load(Ordering::SeqCst), 1); + assert!(poll(notification.as_mut(), &new_waker).is_ready()); +} + +#[test] +fn cancellation_wakers_can_reenter_token_operations() { + struct ReentrantWake { + token: Cancellation, + count: AtomicUsize, + } + impl Wake for ReentrantWake { + fn wake(self: Arc) { + self.wake_by_ref(); + } + fn wake_by_ref(self: &Arc) { + assert!(!self.token.cancel()); + assert!(self.token.child().is_cancelled()); + self.count.fetch_add(1, Ordering::SeqCst); + } + } + + let token = Cancellation::new(); + let counter = Arc::new(ReentrantWake { + token: token.clone(), + count: AtomicUsize::new(0), + }); + let waker = Waker::from(Arc::clone(&counter)); + let mut notification = pin!(token.cancelled()); + assert!(poll(notification.as_mut(), &waker).is_pending()); + let (sender, receiver) = mpsc::channel(); + let cancellation = token.clone(); + let thread = std::thread::spawn(move || sender.send(cancellation.cancel()).unwrap()); + assert!(receive(&receiver)); + thread.join().unwrap(); + assert_eq!(counter.count.load(Ordering::SeqCst), 1); + assert!(poll(notification.as_mut(), &waker).is_ready()); +} + +#[test] +fn cancellation_racing_child_registration_cannot_escape_propagation() { + for _ in 0..64 { + let root = Cancellation::new(); + let start = Arc::new(Barrier::new(2)); + let parent = root.clone(); + let child_start = Arc::clone(&start); + let child = std::thread::spawn(move || { + child_start.wait(); + parent.child().child() + }); + start.wait(); + root.cancel(); + assert!(child.join().unwrap().is_cancelled()); + } +} + +#[test] +fn cancellation_racing_waiter_registration_preserves_wakeups() { + for _ in 0..64 { + let token = Cancellation::new(); + let cancellation = token.clone(); + let start = Arc::new(Barrier::new(2)); + let cancellation_start = Arc::clone(&start); + let thread = std::thread::spawn(move || { + cancellation_start.wait(); + cancellation.cancel(); + }); + let counter = Arc::new(WakeCounter::default()); + let waker = Waker::from(Arc::clone(&counter)); + let mut notification = pin!(token.cancelled()); + start.wait(); + let first = poll(notification.as_mut(), &waker); + thread.join().unwrap(); + if first.is_pending() { + assert_eq!(counter.0.load(Ordering::SeqCst), 1); + assert!(poll(notification.as_mut(), &waker).is_ready()); + } + } +} + +#[test] +fn concurrent_cancellation_requests_have_one_winner() { + let root = Cancellation::new(); + let child = root.child(); + let start = Arc::new(Barrier::new(8)); + let threads: Vec<_> = (0..8) + .map(|_| { + let token = root.clone(); + let start = Arc::clone(&start); + std::thread::spawn(move || { + start.wait(); + token.cancel() + }) + }) + .collect(); + assert_eq!( + threads + .into_iter() + .map(|thread| thread.join().unwrap()) + .filter(|won| *won) + .count(), + 1 + ); + assert!(child.is_cancelled()); +} + +#[test] +fn cancellation_traverses_a_deep_chain_without_recursion() { + let root = Cancellation::new(); + let mut leaf = root.clone(); + for _ in 0..10_000 { + leaf = leaf.child(); + } + assert!(root.cancel()); + assert!(leaf.is_cancelled()); +} + +#[test] +fn deferred_work_completes_normally_without_a_cancellation_callback() { + let scheduler = Scheduler::with_workers(1).unwrap(); + let token = Cancellation::new(); + let called = Rc::new(Cell::new(false)); + let callback_called = Rc::clone(&called); + let callback_dropped = Arc::new(AtomicUsize::new(0)); + let guard = CountDrop(Arc::clone(&callback_dropped)); + let value = scheduler.block_on(defer_cancel( + &token, + async { + yield_now().await; + 42 + }, + move || { + let _guard = guard; + callback_called.set(true); + }, + )); + assert_eq!(value, 42); + assert!(!called.get()); + assert_eq!(callback_dropped.load(Ordering::SeqCst), 1); +} + +#[test] +fn already_cancelled_tokens_notify_before_the_first_protected_poll() { + let scheduler = Scheduler::with_workers(1).unwrap(); + let token = Cancellation::new(); + token.cancel(); + let called = Cell::new(false); + let value = scheduler.block_on(defer_cancel( + &token, + async { + assert!(called.get()); + 42 + }, + || called.set(true), + )); + assert_eq!(value, 42); +} + +#[test] +fn deferred_cancellation_keeps_pending_work_alive_and_calls_a_consuming_callback_once() { + let token = Cancellation::new(); + let callback_count = Cell::new(0); + let released = Cell::new(false); + let dropped = Arc::new(AtomicUsize::new(0)); + let guard = CountDrop(Arc::clone(&dropped)); + let future = poll_fn(|_| { + let _ = &guard; + if released.get() { + Poll::Ready(42) + } else { + Poll::Pending + } + }); + let payload = String::from("shutdown"); + let counter = Arc::new(WakeCounter::default()); + let waker = Waker::from(Arc::clone(&counter)); + let mut wrapper = pin!(defer_cancel(&token, future, || { + assert_eq!(payload, "shutdown"); + drop(payload); + callback_count.set(callback_count.get() + 1); + })); + assert!(poll(wrapper.as_mut(), &waker).is_pending()); + assert!(token.cancel()); + assert_eq!(counter.0.load(Ordering::SeqCst), 1); + assert!(poll(wrapper.as_mut(), &waker).is_pending()); + assert_eq!(callback_count.get(), 1); + assert_eq!(dropped.load(Ordering::SeqCst), 0); + assert!(!token.cancel()); + assert!(poll(wrapper.as_mut(), &waker).is_pending()); + assert_eq!(callback_count.get(), 1); + released.set(true); + assert_eq!(poll(wrapper.as_mut(), &waker), Poll::Ready(42)); +} + +#[test] +fn the_callback_can_make_protected_work_ready_in_the_same_poll() { + let token = Cancellation::new(); + let ready = Cell::new(false); + let future = poll_fn(|_| { + if ready.get() { + Poll::Ready(42) + } else { + Poll::Pending + } + }); + let mut wrapper = pin!(defer_cancel(&token, future, || ready.set(true))); + assert!(poll(wrapper.as_mut(), Waker::noop()).is_pending()); + token.cancel(); + assert_eq!(poll(wrapper.as_mut(), Waker::noop()), Poll::Ready(42)); +} + +#[test] +fn nested_wrappers_observe_their_own_tokens_and_keep_draining() { + let outer = Cancellation::new(); + let inner = Cancellation::new(); + let outer_count = Cell::new(0); + let inner_count = Cell::new(0); + let ready = Cell::new(false); + let future = poll_fn(|_| { + if ready.get() { + Poll::Ready(42) + } else { + Poll::Pending + } + }); + let mut wrapper = pin!(defer_cancel( + &outer, + defer_cancel(&inner, future, || inner_count.set(1)), + || outer_count.set(1) + )); + assert!(poll(wrapper.as_mut(), Waker::noop()).is_pending()); + outer.cancel(); + assert!(poll(wrapper.as_mut(), Waker::noop()).is_pending()); + assert_eq!(outer_count.get(), 1); + assert_eq!(inner_count.get(), 0); + inner.cancel(); + assert!(poll(wrapper.as_mut(), Waker::noop()).is_pending()); + assert_eq!(inner_count.get(), 1); + ready.set(true); + assert_eq!(poll(wrapper.as_mut(), Waker::noop()), Poll::Ready(42)); +} + +#[test] +fn dropping_a_wrapper_drops_its_work_and_does_not_call_on_cancel() { + let token = Cancellation::new(); + let count = Cell::new(0); + let dropped = Arc::new(AtomicUsize::new(0)); + let guard = CountDrop(Arc::clone(&dropped)); + let future = async move { + let _guard = guard; + pending::<()>().await; + }; + let mut wrapper = Box::pin(defer_cancel(&token, future, || count.set(1))); + assert!(poll(wrapper.as_mut(), Waker::noop()).is_pending()); + drop(wrapper); + token.cancel(); + assert_eq!(count.get(), 0); + assert_eq!(dropped.load(Ordering::SeqCst), 1); +} + +#[test] +fn callback_panics_propagate_and_the_wrapper_can_be_destroyed() { + let token = Cancellation::new(); + token.cancel(); + let dropped = Arc::new(AtomicUsize::new(0)); + let guard = CountDrop(Arc::clone(&dropped)); + let mut wrapper = Box::pin(defer_cancel( + &token, + async move { + let _guard = guard; + pending::<()>().await; + }, + || panic!("shutdown callback failed"), + )); + let panic = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + poll(wrapper.as_mut(), Waker::noop()) + })); + assert!(panic.is_err()); + drop(wrapper); + assert_eq!(dropped.load(Ordering::SeqCst), 1); +} + +#[test] +fn protected_futures_need_not_be_unpin() { + struct PinnedFuture { + address: Cell, + _pin: PhantomPinned, + } + impl Future for PinnedFuture { + type Output = usize; + fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll { + let this = self.as_ref().get_ref(); + let address = this as *const Self as usize; + if this.address.get() == 0 { + this.address.set(address); + context.waker().wake_by_ref(); + Poll::Pending + } else { + assert_eq!(this.address.get(), address); + Poll::Ready(address) + } + } + } + let scheduler = Scheduler::with_workers(1).unwrap(); + let address = scheduler.block_on(defer_cancel( + &Cancellation::never(), + PinnedFuture { + address: Cell::new(0), + _pin: PhantomPinned, + }, + || panic!("never-cancelled token invoked callback"), + )); + assert_ne!(address, 0); +} + +#[test] +fn socketry_tasks_can_finish_async_cleanup_after_token_cancellation() { + let scheduler = Scheduler::with_workers(2).unwrap(); + let shutdown = Cancellation::new(); + let token = shutdown.clone(); + let (started, waiting) = mpsc::channel(); + let task = scheduler + .spawn(async move { + let drain = Cancellation::new(); + let callback_drain = drain.clone(); + defer_cancel( + &token, + async { + started.send(()).unwrap(); + drain.cancelled().await; + yield_now().await; + 42 + }, + move || { + callback_drain.cancel(); + }, + ) + .await + }) + .unwrap(); + receive(&waiting); + shutdown.cancel(); + assert_eq!(scheduler.block_on(task).unwrap(), 42); +} + +#[cfg(feature = "tokio")] +#[test] +fn tokio_can_drive_the_same_token_and_deferred_cleanup() { + let runtime = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let shutdown = Cancellation::new(); + let token = shutdown.clone(); + let task = runtime.spawn(async move { + let drain = Cancellation::new(); + let callback_drain = drain.clone(); + defer_cancel( + &token, + async { + drain.cancelled().await; + tokio::task::yield_now().await; + 42 + }, + move || { + callback_drain.cancel(); + }, + ) + .await + }); + let value = runtime.block_on(async { + tokio::task::yield_now().await; + shutdown.cancel(); + task.await.unwrap() + }); + assert_eq!(value, 42); +} diff --git a/crates/executor/tests/io.rs b/crates/executor/tests/io.rs index 2c0575f..4b03a21 100644 --- a/crates/executor/tests/io.rs +++ b/crates/executor/tests/io.rs @@ -5,8 +5,8 @@ mod support; -use socketry_executor::{Clock, FileIo, Interest, Network, Spawn}; -use std::fs::{File, OpenOptions}; +use socketry_executor::{Clock, File, Interest, Socket, Spawn}; +use std::fs::{File as StdFile, OpenOptions}; use std::future::{Future, poll_fn}; use std::io; use std::net::{SocketAddr, TcpListener, TcpStream}; @@ -57,7 +57,7 @@ fn error_kind(result: io::Result) -> io::ErrorKind { async fn round_trip(scheduler: SchedulerType) where - SchedulerType: Network + Spawn + Clock + Clone + 'static, + SchedulerType: Socket + Spawn + Clock + Clone + 'static, { let (listener, address) = listening_socket(); let listener = scheduler.register_listener(listener).unwrap(); @@ -106,7 +106,7 @@ where async fn cancel_pending_read(scheduler: &SchedulerType) where - SchedulerType: Network + Clock, + SchedulerType: Socket + Clock, { let (client, server) = socket_pair(); let socket = scheduler.register_socket(client).unwrap(); @@ -128,7 +128,7 @@ where async fn concurrent_reads(scheduler: SchedulerType) where - SchedulerType: Network + Spawn + Clone + 'static, + SchedulerType: Socket + Spawn + Clone + 'static, { let (client, mut server) = socket_pair(); let socket = Arc::new(scheduler.register_socket(client).unwrap()); @@ -167,7 +167,7 @@ impl Drop for TemporaryFile { } } -async fn positioned_files(scheduler: &SchedulerType) { +async fn positioned_files(scheduler: &SchedulerType) { static NEXT_FILE: AtomicU64 = AtomicU64::new(0); let path = std::env::temp_dir().join(format!( "socketry-io-{}-{}", @@ -199,7 +199,7 @@ async fn positioned_files(scheduler: &SchedulerType) { .await; assert_eq!(result.unwrap_err().kind(), io::ErrorKind::InvalidInput); assert_eq!(buffer.as_ptr(), allocation); - let read_only = Arc::new(File::open(&path).unwrap()); + let read_only = Arc::new(StdFile::open(&path).unwrap()); let (result, returned) = scheduler.file_write_at(read_only, buffer, 0).await; assert!(result.is_err()); assert_eq!(returned.as_ptr(), allocation); @@ -346,7 +346,7 @@ mod native { let socket = handle.register_socket(client).unwrap(); let (raw_listener, address) = listening_socket(); let registered_listener = handle.register_listener(raw_listener).unwrap(); - let file = Arc::new(File::open(std::env::current_exe().unwrap()).unwrap()); + let file = Arc::new(StdFile::open(std::env::current_exe().unwrap()).unwrap()); scheduler.shutdown(); let (client, _peer) = socket_pair(); @@ -637,7 +637,7 @@ mod tokio_adapter { let socket = handle.register_socket(client).unwrap(); let (listener, address) = listening_socket(); let listener = handle.register_listener(listener).unwrap(); - let file = Arc::new(File::open(std::env::current_exe().unwrap()).unwrap()); + let file = Arc::new(StdFile::open(std::env::current_exe().unwrap()).unwrap()); runtime.block_on(scheduler.shutdown()); let (new_socket, _) = socket_pair(); diff --git a/readme.md b/readme.md index 8fc665d..5f26293 100644 --- a/readme.md +++ b/readme.md @@ -38,9 +38,56 @@ See the [executor package](crates/executor/readme.md) for scheduling, ownership, cargo run --package socketry-executor --example work_stealing ``` +## Cooperative cancellation + +`Cancellation` signals a shutdown request without destroying futures. Clones share the request; `child()` creates a boundary that receives parent cancellation without cancelling its parent or siblings. Use one root for application shutdown and children for independently stoppable services. Dropping a signal does not cancel it. + +Work can await `signal.cancelled()` or use `signal.check()` at explicit cancellation points, returning `Cancelled`. Waiting uses wakers and does not busy-wait. `Cancellation::never()` supplies an allocation-free input for work that cannot be cancelled through its signal. + +`defer_cancel(&signal, future, on_cancel)` calls the synchronous callback once when cancellation is observed and continues awaiting the future's normal output. The callback requests graceful shutdown; the future performs asynchronous draining. These primitives work with standard futures on Socketry or Tokio, without a runtime dependency. + +```rust +use socketry::{Cancellation, Scheduler, defer_cancel, yield_now}; + +fn main() -> Result<(), Box> { + let scheduler = Scheduler::with_workers(2)?; + let shutdown = Cancellation::new(); + let server_shutdown = shutdown.child(); + let server = scheduler.spawn(async move { + let drain = Cancellation::new(); + defer_cancel(&server_shutdown, async { + drain.cancelled().await; + // Finish accepted work while the scheduler remains available. + yield_now().await; + 42 + }, || { drain.cancel(); }).await + })?; + + shutdown.cancel(); // Request shutdown after receiving a process signal. + assert_eq!(scheduler.block_on(server)?, 42); + scheduler.shutdown(); + Ok(()) +} +``` + +Cancellation expresses intent; task handles and barriers confirm completion. Keep the scheduler, I/O services, and owners alive while draining. Existing task cancellation, barrier stopping, and scheduler shutdown still destroy futures; `defer_cancel` cannot protect a future from destruction. Signal handling and cancellation inputs on I/O operations are not yet integrated. There is no signal escalation policy: an external supervisor can enforce SIGTERM followed by SIGKILL. + ## Portable I/O and runtime selection -Generic code can accept `Network`, `FileIo`, `Clock`, and `Spawn` capabilities. Socketry and the optional Tokio adapter implement these contracts with concrete future and resource types. Import the traits to call their methods. +Generic code can accept `Socket`, `File`, `Clock`, and `Spawn` capabilities. Socketry and the optional Tokio adapter implement these contracts with concrete future and resource types. Import the traits to call their methods. + +`File` is the scheduler capability trait; `std::fs::File` is the file resource. Alias the resource when using both names: + +```rust +use socketry::File; +use std::{fs::File as StdFile, io, sync::Arc}; + +async fn read_prefix(scheduler: &S, file: Arc) -> io::Result> { + let (result, mut buffer) = scheduler.file_read_at(file, vec![0; 4096], 0).await; + buffer.truncate(result?); + Ok(buffer) +} +``` | Cargo configuration | Implementation | | --- | --- | diff --git a/releases.md b/releases.md index ddd0234..46dd78e 100644 --- a/releases.md +++ b/releases.md @@ -2,6 +2,9 @@ ## Unreleased +- Rename the portable socket operations trait from `Network` to `Socket` in `socketry`, `socketry-executor`, and their public `scheduler` module. Update imports and trait bounds. +- Rename the portable file operations trait from `FileIo` to `File` in `socketry`, `socketry-executor`, and their public `scheduler` module. Update imports and trait bounds; alias `std::fs::File` when both names are used. +- Add runtime-independent `Cancellation`, `Cancelled`, and `defer_cancel` for cooperative service shutdown and asynchronous draining. Export them through `socketry-executor` and `socketry`. - Let Cargo select the dependency version in the installation example. ## v0.2.0 diff --git a/src/lib.rs b/src/lib.rs index 2963d6e..6d34a3c 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,8 +6,9 @@ pub use socketry_executor as executor; pub use socketry_executor::{ - Barrier, BufferResult, Clock, FileIo, Interest, Network, Scheduler, SchedulerHandle, Spawn, - SpawnError, Task, TaskError, TaskHandle, yield_now, + Barrier, BufferResult, Cancellation, Cancelled, Clock, File, Interest, Scheduler, + SchedulerHandle, Socket, Spawn, SpawnError, Task, TaskError, TaskHandle, defer_cancel, + yield_now, }; pub use socketry_executor::scheduler;