use std::cell::UnsafeCell; use std::future::Future; use std::marker::PhantomPinned; use std::pin::Pin; use std::ptr::NonNull; use std::sync::atomic::Ordering::{AcqRel, Acquire}; use std::task::{Context, Poll, Waker};
/// Invoked when the IO driver is shut down; forces this `ScheduledIo` into a /// permanently shutdown state. pub(super) fn shutdown(&self) { let mask = SHUTDOWN.pack(1, 0); self.readiness.fetch_or(mask, AcqRel); self.wake(Ready::ALL);
}
/// Sets the readiness on this `ScheduledIo` by invoking the given closure on /// the current value, returning the previous readiness value. /// /// # Arguments /// - `tick`: whether setting the tick or trying to clear readiness for a /// specific tick. /// - `f`: a closure returning a new readiness value given the previous /// readiness. pub(super) fn set_readiness(&self, tick_op: Tick, f: implFn(Ready) -> Ready) { let _ = self.readiness.fetch_update(AcqRel, Acquire, |curr| { // If the io driver is shut down, then you are only allowed to clear readiness.
debug_assert!(SHUTDOWN.unpack(curr) == 0 || matches!(tick_op, Tick::Clear(_)));
let new_tick = match tick_op { // Trying to clear readiness with an old event!
Tick::Clear(t) if tick as u8 != t => return None,
Tick::Clear(t) => t as usize,
Tick::Set => tick.wrapping_add(1) % MAX_TICK,
}; let ready = Ready::from_usize(READINESS.unpack(curr));
Some(TICK.pack(new_tick, f(ready).as_usize()))
});
}
/// Notifies all pending waiters that have registered interest in `ready`. /// /// There may be many waiters to notify. Waking the pending task **must** be /// done from outside of the lock otherwise there is a potential for a /// deadlock. /// /// A stack array of wakers is created and filled with wakers to notify, the /// lock is released, and the wakers are notified. Because there may be more /// than 32 wakers to notify, if the stack array fills up, the lock is /// released, the array is cleared, and the iteration continues. pub(super) fn wake(&self, ready: Ready) { letmut wakers = WakeList::new();
letmut waiters = self.waiters.lock();
// check for AsyncRead slot if ready.is_readable() { iflet Some(waker) = waiters.reader.take() {
wakers.push(waker);
}
}
// check for AsyncWrite slot if ready.is_writable() { iflet Some(waker) = waiters.writer.take() {
wakers.push(waker);
}
}
'outer: loop { letmut iter = waiters.list.drain_filter(|w| ready.satisfies(w.interest));
while wakers.can_push() { match iter.next() {
Some(waiter) => { let waiter = unsafe { &mut *waiter.as_ptr() };
/// Polls for readiness events in a given direction. /// /// These are to support `AsyncRead` and `AsyncWrite` polling methods, /// which cannot use the `async fn` version. This uses reserved reader /// and writer slots. pub(super) fn poll_readiness(
&self,
cx: &mut Context<'_>,
direction: Direction,
) -> Poll<ReadyEvent> { let curr = self.readiness.load(Acquire);
let ready = direction.mask() & Ready::from_usize(READINESS.unpack(curr)); let is_shutdown = SHUTDOWN.unpack(curr) != 0;
if ready.is_empty() && !is_shutdown { // Update the task info letmut waiters = self.waiters.lock(); let waker = match direction {
Direction::Read => &mut waiters.reader,
Direction::Write => &mut waiters.writer,
};
// Avoid cloning the waker if one is already stored that matches the // current task. match waker {
Some(waker) => waker.clone_from(cx.waker()),
None => *waker = Some(cx.waker().clone()),
}
// Try again, in case the readiness was changed while we were // taking the waiters lock let curr = self.readiness.load(Acquire); let ready = direction.mask() & Ready::from_usize(READINESS.unpack(curr)); let is_shutdown = SHUTDOWN.unpack(curr) != 0; if is_shutdown {
Poll::Ready(ReadyEvent {
tick: TICK.unpack(curr) as u8,
ready: direction.mask(),
is_shutdown,
})
} elseif ready.is_empty() {
Poll::Pending
} else {
Poll::Ready(ReadyEvent {
tick: TICK.unpack(curr) as u8,
ready,
is_shutdown,
})
}
} else {
Poll::Ready(ReadyEvent {
tick: TICK.unpack(curr) as u8,
ready,
is_shutdown,
})
}
}
pub(crate) fn clear_readiness(&self, event: ReadyEvent) { // This consumes the current readiness state **except** for closed // states. Closed states are excluded because they are final states. let mask_no_closed = event.ready - Ready::READ_CLOSED - Ready::WRITE_CLOSED; self.set_readiness(Tick::Clear(event.tick), |curr| curr - mask_no_closed);
}
impl Drop for ScheduledIo { fn drop(&mutself) { self.wake(Ready::ALL);
}
}
unsafeimpl Send for ScheduledIo {} unsafeimpl Sync for ScheduledIo {}
impl ScheduledIo { /// An async version of `poll_readiness` which uses a linked list of wakers. pub(crate) asyncfn readiness(&self, interest: Interest) -> ReadyEvent { self.readiness_fut(interest).await
}
// This is in a separate function so that the borrow checker doesn't think // we are borrowing the `UnsafeCell` possibly over await boundaries. // // Go figure. fn readiness_fut(&self, interest: Interest) -> Readiness<'_> {
Readiness {
scheduled_io: self,
state: State::Init,
waiter: UnsafeCell::new(Waiter {
pointers: linked_list::Pointers::new(),
waker: None,
is_ready: false,
interest,
_p: PhantomPinned,
}),
}
}
}
unsafeimpl linked_list::Link for Waiter { type Handle = NonNull<Waiter>; type Target = Waiter;
let (scheduled_io, state, waiter) = { // Safety: `Self` is `!Unpin` // // While we could use `pin_project!` to remove // this unsafe block, there are already unsafe blocks here, // so it wouldn't significantly ease the mental burden // and would actually complicate the code. // That's why we didn't use it. let me = unsafe { self.get_unchecked_mut() };
(me.scheduled_io, &mut me.state, &me.waiter)
};
loop { match *state {
State::Init => { // Optimistically check existing readiness let curr = scheduled_io.readiness.load(SeqCst); let is_shutdown = SHUTDOWN.unpack(curr) != 0;
// Safety: `waiter.interest` never changes let interest = unsafe { (*waiter.get()).interest }; let ready = Ready::from_usize(READINESS.unpack(curr)).intersection(interest);
if !ready.is_empty() || is_shutdown { // Currently ready! let tick = TICK.unpack(curr) as u8;
*state = State::Done; return Poll::Ready(ReadyEvent {
tick,
ready,
is_shutdown,
});
}
// Wasn't ready, take the lock (and check again while locked). letmut waiters = scheduled_io.waiters.lock();
let curr = scheduled_io.readiness.load(SeqCst); letmut ready = Ready::from_usize(READINESS.unpack(curr)); let is_shutdown = SHUTDOWN.unpack(curr) != 0;
if is_shutdown {
ready = Ready::ALL;
}
let ready = ready.intersection(interest);
if !ready.is_empty() || is_shutdown { // Currently ready! let tick = TICK.unpack(curr) as u8;
*state = State::Done; return Poll::Ready(ReadyEvent {
tick,
ready,
is_shutdown,
});
}
// Not ready even after locked, insert into list...
// Safety: Since the `waiter` is not in the intrusive list yet, // we have exclusive access to it. The Mutex ensures // that this modification is visible to other threads that // acquire the same Mutex. let waker = unsafe { &mut (*waiter.get()).waker }; let old = waker.replace(cx.waker().clone());
debug_assert!(old.is_none(), "waker should be None at the first poll");
// Insert the waiter into the linked list // // safety: pointers from `UnsafeCell` are never null.
waiters
.list
.push_front(unsafe { NonNull::new_unchecked(waiter.get()) });
*state = State::Waiting;
}
State::Waiting => { // Currently in the "Waiting" state, implying the caller has // a waiter stored in the waiter list (guarded by // `notify.waiters`). In order to access the waker fields, // we must hold the lock.
let waiters = scheduled_io.waiters.lock();
// Safety: With the lock held, we have exclusive access to // the waiter. In other words, `ScheduledIo::wake()` // cannot access the waiter concurrently. let w = unsafe { &mut *waiter.get() };
if w.is_ready { // Our waker has been notified.
*state = State::Done;
} else { // Update the waker, if necessary.
w.waker.as_mut().unwrap().clone_from(cx.waker()); return Poll::Pending;
}
// Explicit drop of the lock to indicate the scope that the // lock is held. Because holding the lock is required to // ensure safe access to fields not held within the lock, it // is helpful to visualize the scope of the critical // section.
drop(waiters);
}
State::Done => { let curr = scheduled_io.readiness.load(Acquire); let is_shutdown = SHUTDOWN.unpack(curr) != 0;
// The returned tick might be newer than the event // which notified our waker. This is ok because the future // still didn't return `Poll::Ready`. let tick = TICK.unpack(curr) as u8;
// Safety: We don't need to acquire the lock here because // 1. `State::Done`` means `waiter` is no longer shared, // this means no concurrent access to `waiter` can happen // at this point. // 2. `waiter.interest` is never changed, this means // no side effects need to be synchronized by the lock. let interest = unsafe { (*waiter.get()).interest }; // The readiness state could have been cleared in the meantime, // but we allow the returned ready set to be empty. let ready = Ready::from_usize(READINESS.unpack(curr)).intersection(interest);
impl Drop for Readiness<'_> { fn drop(&mutself) { letmut waiters = self.scheduled_io.waiters.lock();
// Safety: `waiter` is only ever stored in `waiters` unsafe {
waiters
.list
.remove(NonNull::new_unchecked(self.waiter.get()))
};
}
}
unsafeimpl Send for Readiness<'_> {} unsafeimpl Sync for Readiness<'_> {}
Messung V0.5 in Prozent
¤ Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.0.8Bemerkung:
¤
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.