use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use crate::distribution::etf::MAX_DIST_FRAME_BYTES;
pub(super) const INBOUND_RESIDENCY_ENVELOPE_BYTES: u64 = 4 * 1024 * 1024 * 1024;
pub(super) const INBOUND_RESIDENCY_PER_PEER_BYTES: u64 = MAX_DIST_FRAME_BYTES as u64;
pub(super) struct InboundResidency {
pub(super) charged: AtomicU64,
envelope: u64,
per_peer: u64,
pub(super) refused: AtomicU64,
}
impl InboundResidency {
pub(super) fn new() -> Self {
Self {
charged: AtomicU64::new(0),
envelope: INBOUND_RESIDENCY_ENVELOPE_BYTES,
per_peer: INBOUND_RESIDENCY_PER_PEER_BYTES,
refused: AtomicU64::new(0),
}
}
pub(super) fn try_admit(self: &Arc<Self>) -> Option<InboundAdmissionPermit> {
let mut current = self.charged.load(Ordering::Relaxed);
loop {
let next = current.saturating_add(self.per_peer);
if next > self.envelope {
self.refused.fetch_add(1, Ordering::Relaxed);
return None;
}
match self.charged.compare_exchange_weak(
current,
next,
Ordering::AcqRel,
Ordering::Relaxed,
) {
Ok(_) => {
return Some(InboundAdmissionPermit {
residency: Arc::clone(self),
bytes: self.per_peer,
});
}
Err(observed) => current = observed,
}
}
}
}
pub(super) struct InboundAdmissionPermit {
residency: Arc<InboundResidency>,
bytes: u64,
}
impl Drop for InboundAdmissionPermit {
fn drop(&mut self) {
self.residency
.charged
.fetch_sub(self.bytes, Ordering::AcqRel);
}
}