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}