Skip to main content

dope_core/driver/
ready.rs

1use std::cell::Cell;
2use std::io::{Error, ErrorKind, Result};
3use std::marker::PhantomData;
4use std::ptr::NonNull;
5
6use o3::collections::BatchSet;
7use o3::marker::ThreadBound;
8
9use crate::io::fd::FdSlot;
10
11use super::DriverRef;
12use super::token::{Epoch, ROUTE_FRAMEWORK, SlotIndex, Token};
13
14const NIL: u32 = u32::MAX;
15
16/// A generation-checked address in a driver's local ready queue.
17#[derive(Clone, Copy)]
18pub struct ReadyKey<'d> {
19    index: u32,
20    epoch: u32,
21    _arena: PhantomData<&'d Arena>,
22    _thread: ThreadBound,
23}
24
25impl ReadyKey<'static> {
26    pub const NONE: Self = Self {
27        index: NIL,
28        epoch: 0,
29        _arena: PhantomData,
30        _thread: ThreadBound::NEW,
31    };
32}
33
34impl PartialEq for ReadyKey<'_> {
35    fn eq(&self, other: &Self) -> bool {
36        self.index == other.index && self.epoch == other.epoch
37    }
38}
39
40impl Eq for ReadyKey<'_> {}
41
42/// A driver-scoped, type-erased wake target for an in-flight operation.
43///
44/// Unlike [`ReadyKey`], this preserves hierarchical task wakeups. A completion
45/// can therefore wake the exact child task that registered the operation,
46/// rather than only waking the root runtime task.
47#[derive(Clone, Copy)]
48pub struct CompletionWaker<'d> {
49    target: CompletionTarget<'d>,
50    _thread: ThreadBound,
51}
52
53#[derive(Clone, Copy)]
54enum CompletionTarget<'d> {
55    Ready(DriverRef<'d>, ReadyKey<'d>),
56    Callback(NonNull<()>, unsafe fn(NonNull<()>)),
57}
58
59impl<'d> CompletionWaker<'d> {
60    pub fn from_ready(driver: DriverRef<'d>, key: ReadyKey<'d>) -> Self {
61        Self {
62            target: CompletionTarget::Ready(driver, key),
63            _thread: ThreadBound::NEW,
64        }
65    }
66
67    /// # Safety
68    ///
69    /// `target` must remain valid for every call to `wake` while this handle is
70    /// live, and `callback` must accept that exact target.
71    pub unsafe fn from_callback(target: NonNull<()>, callback: unsafe fn(NonNull<()>)) -> Self {
72        Self {
73            target: CompletionTarget::Callback(target, callback),
74            _thread: ThreadBound::NEW,
75        }
76    }
77
78    #[inline]
79    pub fn wake(self) {
80        match self.target {
81            CompletionTarget::Ready(driver, key) => driver.activate_ready(key),
82            CompletionTarget::Callback(target, callback) => unsafe { callback(target) },
83        }
84    }
85}
86
87#[derive(Clone, Copy)]
88pub struct ReadyHandle<'d> {
89    arena: &'d Arena,
90    key: ReadyKey<'d>,
91}
92
93impl<'d> ReadyHandle<'d> {
94    pub fn set_target(self, target: Token) {
95        self.arena.set_target(self.key, target);
96    }
97
98    #[inline]
99    pub fn activate(self) {
100        self.arena.activate(self.key);
101    }
102
103    pub fn key(self) -> ReadyKey<'d> {
104        self.key
105    }
106}
107
108pub struct ReadySlot<'d> {
109    arena: &'d Arena,
110    key: ReadyKey<'d>,
111}
112
113impl<'d> ReadySlot<'d> {
114    fn new(arena: &'d Arena, index: u32) -> Self {
115        arena.live[index as usize].set(true);
116        Self {
117            arena,
118            key: ReadyKey {
119                index,
120                epoch: arena.epochs[index as usize].get(),
121                _arena: PhantomData,
122                _thread: ThreadBound::NEW,
123            },
124        }
125    }
126
127    pub fn set_target(&self, target: Token) {
128        self.arena.set_target(self.key, target);
129    }
130
131    #[inline]
132    pub fn activate(&self) {
133        self.arena.activate(self.key);
134    }
135
136    pub fn key(&self) -> ReadyKey<'d> {
137        self.key
138    }
139}
140
141impl Drop for ReadySlot<'_> {
142    fn drop(&mut self) {
143        self.arena.release(self.key);
144    }
145}
146
147pub(super) struct Arena {
148    fixed: usize,
149    ready: BatchSet,
150    targets: Box<[Cell<Token>]>,
151    epochs: Box<[Cell<u32>]>,
152    live: Box<[Cell<bool>]>,
153    next_free: Box<[Cell<u32>]>,
154    free: Cell<u32>,
155    free_len: Cell<usize>,
156}
157
158impl Arena {
159    pub(super) fn new(fixed: usize, dynamic: usize) -> Result<Box<Self>> {
160        let capacity = fixed
161            .checked_add(dynamic)
162            .ok_or_else(|| Error::new(ErrorKind::InvalidInput, "dope: ready capacity overflow"))?;
163        if capacity > u32::MAX as usize {
164            return Err(Error::new(
165                ErrorKind::InvalidInput,
166                "dope: ready capacity exceeds u32",
167            ));
168        }
169
170        let dummy = Token::new(ROUTE_FRAMEWORK, SlotIndex::new(0), Epoch::INITIAL);
171        Ok(Box::new(Self {
172            fixed,
173            ready: BatchSet::with_capacity(capacity),
174            targets: (0..capacity).map(|_| Cell::new(dummy)).collect(),
175            epochs: (0..capacity).map(|_| Cell::new(0)).collect(),
176            live: (0..capacity)
177                .map(|index| Cell::new(index < fixed))
178                .collect(),
179            next_free: (0..capacity)
180                .map(|index| {
181                    let next = if index + 1 < capacity {
182                        index as u32 + 1
183                    } else {
184                        NIL
185                    };
186                    Cell::new(next)
187                })
188                .collect(),
189            free: Cell::new(if fixed < capacity { fixed as u32 } else { NIL }),
190            free_len: Cell::new(capacity - fixed),
191        }))
192    }
193
194    fn valid(&self, key: ReadyKey<'_>) -> bool {
195        let index = key.index as usize;
196        self.live.get(index).is_some_and(Cell::get) && self.epochs[index].get() == key.epoch
197    }
198
199    pub(crate) fn fixed_slot(&self, slot: FdSlot) -> ReadyHandle<'_> {
200        let index = slot.raw();
201        debug_assert!((index as usize) < self.fixed);
202        ReadyHandle {
203            arena: self,
204            key: ReadyKey {
205                index,
206                epoch: 0,
207                _arena: PhantomData,
208                _thread: ThreadBound::NEW,
209            },
210        }
211    }
212
213    pub(crate) fn make_slot(&self, target: Token) -> Result<ReadySlot<'_>> {
214        self.make_slot_reserving(target, 0)
215    }
216
217    pub(crate) fn make_slots<I>(&self, targets: I) -> Result<Box<[ReadySlot<'_>]>>
218    where
219        I: IntoIterator<Item = Token>,
220        I::IntoIter: ExactSizeIterator,
221    {
222        let targets = targets.into_iter();
223        let requested = targets.len();
224        let available = self.free_len.get();
225        if requested > available {
226            return Err(Self::capacity_error(requested, available));
227        }
228        targets
229            .enumerate()
230            .map(|(index, target)| self.make_slot_reserving(target, requested - index - 1))
231            .collect()
232    }
233
234    pub(crate) fn make_slot_reserving(
235        &self,
236        target: Token,
237        reserve: usize,
238    ) -> Result<ReadySlot<'_>> {
239        if self.free_len.get() <= reserve {
240            return Err(Self::capacity_error(1, self.free_len.get()));
241        }
242        let index = self.free.get();
243        if index == NIL {
244            return Err(Self::capacity_error(1, 0));
245        }
246        self.free.set(self.next_free[index as usize].get());
247        self.free_len.set(self.free_len.get() - 1);
248        self.targets[index as usize].set(target);
249        Ok(ReadySlot::new(self, index))
250    }
251
252    fn capacity_error(requested: usize, available: usize) -> Error {
253        Error::new(
254            ErrorKind::WouldBlock,
255            format!(
256                "dope: dynamic ready capacity exhausted: requested {requested}, available {available}"
257            ),
258        )
259    }
260
261    fn set_target(&self, key: ReadyKey<'_>, target: Token) {
262        if self.valid(key) {
263            self.targets[key.index as usize].set(target);
264        }
265    }
266
267    pub(crate) fn activate(&self, key: ReadyKey<'_>) {
268        if self.valid(key) {
269            self.ready.insert(key.index as usize);
270        }
271    }
272
273    fn release(&self, key: ReadyKey<'_>) {
274        if !self.valid(key) {
275            return;
276        }
277        let index = key.index as usize;
278        self.live[index].set(false);
279        self.ready.remove(index);
280        let Some(epoch) = key.epoch.checked_add(1) else {
281            self.next_free[index].set(NIL);
282            return;
283        };
284        self.epochs[index].set(epoch);
285        self.next_free[index].set(self.free.get());
286        self.free.set(key.index);
287        self.free_len.set(self.free_len.get() + 1);
288    }
289
290    pub(crate) fn drain(&self, mut activate: impl FnMut(Token)) {
291        let Some(mut ready) = self.ready.drain_batch() else {
292            return;
293        };
294        for index in &mut ready {
295            if self.live[index].get() {
296                activate(self.targets[index].get());
297            }
298        }
299    }
300
301    pub(crate) fn has_ready(&self) -> bool {
302        !self.ready.is_empty()
303    }
304}