use std::sync::atomic::{
AtomicI64,
Ordering,
};
use futures::task::AtomicWaker;
pub(crate) struct SendWindow {
remaining: AtomicI64,
waker: AtomicWaker,
}
impl SendWindow {
pub(crate) const fn new(initial: u32) -> Self {
Self {
remaining: AtomicI64::new(initial as i64),
waker: AtomicWaker::new(),
}
}
pub(crate) fn available(&self) -> i64 {
self.remaining.load(Ordering::Acquire)
}
pub(crate) fn consume(&self, n: usize) -> bool {
let n = n as i64;
loop {
let cur = self.remaining.load(Ordering::Acquire);
if cur == i64::MIN {
return false; }
if self
.remaining
.compare_exchange(cur, cur - n, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return true;
}
}
}
pub(crate) fn replenish(&self, delta: u32) {
self.remaining.fetch_add(delta as i64, Ordering::Release);
self.waker.wake();
}
pub(crate) fn apply_delta(&self, delta: i64) {
self.remaining.fetch_add(delta, Ordering::Release);
if delta > 0 {
self.waker.wake();
}
}
pub(crate) fn register_waker(&self, waker: &std::task::Waker) {
self.waker.register(waker);
}
pub(crate) fn close(&self) {
self.remaining.store(i64::MIN, Ordering::Release);
self.waker.wake();
}
pub(crate) fn is_closed(&self) -> bool {
self.remaining.load(Ordering::Acquire) == i64::MIN
}
}