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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 15 additions & 3 deletions context/design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.

Expand Down Expand Up @@ -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.
Expand Down
31 changes: 31 additions & 0 deletions context/gaps.md
Original file line number Diff line number Diff line change
@@ -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.
13 changes: 12 additions & 1 deletion context/implementation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion crates/executor/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
Loading
Loading