Skip to main content

socketry_executor/scheduler/
mod.rs

1// Released under the MIT License.
2// Copyright, 2026, by Samuel Williams.
3
4//! Scheduler implementations and their I/O selectors.
5use std::fs::File;
6use std::future::Future;
7use std::io;
8use std::net::{SocketAddr, TcpListener, TcpStream};
9use std::sync::Arc;
10use std::time::Duration;
11
12pub mod selector;
13pub mod socketry;
14
15#[cfg(any(feature = "native", feature = "tokio"))]
16mod file;
17
18#[cfg(feature = "tokio")]
19pub mod tokio;
20
21pub use socketry::{Scheduler, SchedulerHandle};
22pub(crate) use socketry::{Shared, enter};
23
24/// An operation result together with its reusable, owned buffer.
25///
26/// Reads fill the existing buffer length and leave its length unchanged. Only
27/// the first `result?` bytes contain newly read data. Writes may be partial.
28/// The buffer is returned on both success and failure.
29pub type BufferResult = (io::Result<usize>, Vec<u8>);
30
31/// A socket readiness condition. Readiness can be spurious; retry nonblocking
32/// operations and wait again when they return `WouldBlock`.
33#[derive(Clone, Copy, Debug, Eq, PartialEq)]
34pub enum Interest {
35    Readable,
36    Writable,
37}
38
39/// Portable socket operations, selected through the concrete implementation.
40///
41/// Registrations belong to their creating implementation. Keep a socket's
42/// registration across operations and worker migration. Implementations must
43/// return an error when a resource belongs to an incompatible runtime instance.
44///
45/// Dropping an operation future abandons its result. It need not undo an I/O
46/// operation already submitted to the kernel: a cancelled read can consume
47/// bytes, and a cancelled write can transmit bytes. Implementations retain any
48/// kernel-accessible memory until the operation finishes. Await completion when
49/// the amount transferred matters. No asynchronous cleanup is promised by Drop.
50pub trait Network: Send + Sync {
51    type Socket: Send + Sync;
52    type Listener: Send + Sync;
53
54    /// Register an owned socket once. The implementation sets nonblocking mode
55    /// when required. Do not change its mode through another OS handle.
56    fn register_socket(&self, socket: TcpStream) -> io::Result<Self::Socket>;
57
58    /// Register an owned listener once, with the same mode requirements.
59    fn register_listener(&self, listener: TcpListener) -> io::Result<Self::Listener>;
60
61    fn connect(&self, address: SocketAddr)
62    -> impl Future<Output = io::Result<Self::Socket>> + Send;
63
64    fn accept(
65        &self,
66        listener: &Self::Listener,
67    ) -> impl Future<Output = io::Result<(Self::Socket, SocketAddr)>> + Send;
68
69    fn io_read(
70        &self,
71        socket: &Self::Socket,
72        buffer: Vec<u8>,
73    ) -> impl Future<Output = BufferResult> + Send;
74
75    fn io_write(
76        &self,
77        socket: &Self::Socket,
78        buffer: Vec<u8>,
79    ) -> impl Future<Output = BufferResult> + Send;
80
81    fn io_wait(
82        &self,
83        socket: &Self::Socket,
84        interest: Interest,
85    ) -> impl Future<Output = io::Result<()>> + Send;
86}
87
88/// Positioned file operations. A regular file does not support a universal
89/// readiness fallback, so implementations use native completion or a blocking
90/// pool. Use ordinary files opened without append mode, not pipes. Offsets
91/// must fit in i64. The Unix implementation leaves the shared cursor unchanged;
92/// the Windows blocking fallback updates it, as std's seek_read/seek_write do.
93///
94/// Buffers and the file remain owned by an in-flight operation even if the
95/// waiting future is dropped. A write can still complete after cancellation.
96pub trait FileIo: Send + Sync {
97    fn file_read_at(
98        &self,
99        file: Arc<File>,
100        buffer: Vec<u8>,
101        offset: u64,
102    ) -> impl Future<Output = BufferResult> + Send;
103
104    fn file_write_at(
105        &self,
106        file: Arc<File>,
107        buffer: Vec<u8>,
108        offset: u64,
109    ) -> impl Future<Output = BufferResult> + Send;
110}
111
112/// A runtime's monotonic sleep facility.
113pub trait Clock: Send + Sync {
114    fn sleep(&self, duration: Duration) -> impl Future<Output = ()> + Send;
115}