use super::protocol::RecordType;
pub(crate) const DEFAULT_PEER_BYTE_BUDGET: u64 = 8 * 1024 * 1024;
pub(crate) const DEFAULT_SLOT_BYTE_BUDGET: u32 = 1024 * 1024;
pub(crate) const fn slot_buffer_depth(data_credit: u32) -> usize {
data_credit as usize + 1
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub(crate) enum NegotiationError {
#[error("peer advertised initial_credit = 0: legacy peer, the mux is unusable")]
LegacyPeer,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct NegotiatedLimits {
initial_credit: u32,
slot_byte_budget: u32,
}
impl NegotiatedLimits {
pub(crate) const fn from_wire(
initial_credit: u32,
slot_byte_budget: u32,
) -> Result<Self, NegotiationError> {
if initial_credit == 0 {
return Err(NegotiationError::LegacyPeer);
}
let slot_byte_budget = if slot_byte_budget == 0 {
DEFAULT_SLOT_BYTE_BUDGET
} else {
slot_byte_budget
};
Ok(Self {
initial_credit,
slot_byte_budget,
})
}
pub(crate) const fn initial_credit(&self) -> u32 {
self.initial_credit
}
pub(crate) const fn slot_byte_budget(&self) -> u32 {
self.slot_byte_budget
}
pub(crate) const fn slot_buffer_depth(&self) -> usize {
slot_buffer_depth(self.initial_credit)
}
pub(crate) const fn open_credit(&self) -> SlotCredit {
SlotCredit::new(self.initial_credit)
}
#[cfg(test)]
pub(crate) const fn open_account(&self) -> SlotCreditAccount {
SlotCreditAccount::new(self.initial_credit)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub(crate) enum CreditClass {
Data,
Terminal,
Control,
}
impl CreditClass {
pub(crate) const fn of(record_type: RecordType, is_terminal: bool) -> Self {
match record_type {
RecordType::OpenSlot | RecordType::CloseSlot | RecordType::CreditUpdate => {
Self::Control
}
RecordType::Data if is_terminal => Self::Terminal,
RecordType::Data | RecordType::SlotHeartbeat => Self::Data,
}
}
#[cfg(test)]
pub(crate) const fn occupies_buffer(self) -> bool {
!matches!(self, Self::Control)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub(crate) enum CreditError {
#[error("slot data credit exhausted")]
DataExhausted,
#[error("slot terminal reserve already spent")]
TerminalAlreadySpent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum TerminalReserve {
Unspent,
Spent,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SlotCredit {
data_available: u32,
terminal: TerminalReserve,
}
impl SlotCredit {
pub(crate) const fn new(initial: u32) -> Self {
Self {
data_available: initial,
terminal: TerminalReserve::Unspent,
}
}
pub(crate) const fn data_available(&self) -> u32 {
self.data_available
}
pub(crate) const fn terminal_available(&self) -> bool {
matches!(self.terminal, TerminalReserve::Unspent)
}
pub(crate) const fn can_spend(&self, class: CreditClass) -> bool {
match class {
CreditClass::Control => true,
CreditClass::Terminal => self.terminal_available(),
CreditClass::Data => self.data_available > 0,
}
}
pub(crate) fn try_spend(&mut self, class: CreditClass) -> Result<(), CreditError> {
match class {
CreditClass::Control => Ok(()),
CreditClass::Terminal => match self.terminal {
TerminalReserve::Unspent => {
self.terminal = TerminalReserve::Spent;
Ok(())
}
TerminalReserve::Spent => Err(CreditError::TerminalAlreadySpent),
},
CreditClass::Data => match self.data_available.checked_sub(1) {
Some(remaining) => {
self.data_available = remaining;
Ok(())
}
None => Err(CreditError::DataExhausted),
},
}
}
pub(crate) fn grant(&mut self, delta: u32) -> u32 {
self.data_available = self.data_available.saturating_add(delta);
self.data_available
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SlotCreditAccount {
limit: u32,
data_outstanding: u32,
buffered: u32,
ungranted: u32,
terminal: TerminalReserve,
}
impl SlotCreditAccount {
pub(crate) const fn new(limit: u32) -> Self {
Self {
limit,
data_outstanding: 0,
buffered: 0,
ungranted: 0,
terminal: TerminalReserve::Unspent,
}
}
#[cfg(test)]
pub(crate) const fn limit(&self) -> u32 {
self.limit
}
#[cfg(test)]
pub(crate) const fn buffer_depth(&self) -> usize {
slot_buffer_depth(self.limit)
}
pub(crate) const fn buffered(&self) -> u32 {
self.buffered
}
#[cfg(test)]
pub(crate) const fn data_outstanding(&self) -> u32 {
self.data_outstanding
}
#[cfg(test)]
pub(crate) const fn data_free(&self) -> u32 {
self.limit - self.data_outstanding
}
#[cfg(test)]
pub(crate) const fn pending_grant(&self) -> u32 {
self.ungranted
}
pub(crate) fn admit(&mut self, class: CreditClass) -> Result<(), CreditError> {
match class {
CreditClass::Control => return Ok(()),
CreditClass::Terminal => match self.terminal {
TerminalReserve::Unspent => self.terminal = TerminalReserve::Spent,
TerminalReserve::Spent => return Err(CreditError::TerminalAlreadySpent),
},
CreditClass::Data => {
if self.data_outstanding >= self.limit {
return Err(CreditError::DataExhausted);
}
self.data_outstanding += 1;
}
}
self.buffered += 1;
Ok(())
}
pub(crate) fn release(&mut self, drained: u32) -> u32 {
let leaving = drained.min(self.buffered);
self.buffered -= leaving;
let data = leaving.min(self.data_outstanding);
self.data_outstanding -= data;
self.ungranted = self.ungranted.saturating_add(data);
self.ungranted
}
pub(crate) fn take_pending_grant(&mut self) -> Option<u32> {
let delta = self.ungranted;
if delta == 0 {
return None;
}
self.ungranted = 0;
Some(delta)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
pub(crate) enum ByteBudgetError {
#[error("{requested} bytes exceeds the {limit}-byte budget outright")]
ExceedsBudget { requested: u64, limit: u64 },
#[error("{requested} bytes does not fit the {available} bytes left of {limit}")]
Exhausted {
requested: u64,
available: u64,
limit: u64,
},
}
impl ByteBudgetError {
#[cfg(test)]
pub(crate) const fn is_transient(&self) -> bool {
matches!(self, Self::Exhausted { .. })
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct ByteBudget {
limit: u64,
used: u64,
}
impl ByteBudget {
pub(crate) const fn new(limit: u64) -> Self {
Self { limit, used: 0 }
}
#[cfg(test)]
pub(crate) const fn per_peer() -> Self {
Self::new(DEFAULT_PEER_BYTE_BUDGET)
}
#[cfg(test)]
pub(crate) const fn per_slot(limits: &NegotiatedLimits) -> Self {
Self::new(limits.slot_byte_budget() as u64)
}
pub(crate) const fn limit(&self) -> u64 {
self.limit
}
pub(crate) const fn used(&self) -> u64 {
self.used
}
pub(crate) const fn available(&self) -> u64 {
self.limit - self.used
}
pub(crate) fn try_reserve(&mut self, bytes: usize) -> Result<(), ByteBudgetError> {
let requested = bytes as u64;
if requested > self.limit {
return Err(ByteBudgetError::ExceedsBudget {
requested,
limit: self.limit,
});
}
if requested > self.available() {
return Err(ByteBudgetError::Exhausted {
requested,
available: self.available(),
limit: self.limit,
});
}
self.used += requested;
Ok(())
}
pub(crate) fn release(&mut self, bytes: usize) {
self.used = self.used.saturating_sub(bytes as u64);
}
}
pub(crate) fn try_reserve_pair(
peer: &mut ByteBudget,
slot: &mut ByteBudget,
bytes: usize,
) -> Result<(), ByteBudgetError> {
slot.try_reserve(bytes)?;
match peer.try_reserve(bytes) {
Ok(()) => Ok(()),
Err(err) => {
slot.release(bytes);
Err(err)
}
}
}
pub(crate) fn release_pair(peer: &mut ByteBudget, slot: &mut ByteBudget, bytes: usize) {
slot.release(bytes);
peer.release(bytes);
}
#[cfg(test)]
mod tests;