use std::fmt;
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct GlobalCreditsReport {
pub current: u32,
pub min: u32,
pub(crate) seq: u8,
}
impl GlobalCreditsReport {
pub(crate) fn initial(credits: u32) -> Self {
Self { current: credits, min: credits, seq: 0 }
}
pub(crate) fn consume(&mut self, credits: u32) {
self.current = self.current.saturating_sub(credits);
self.min = self.min.min(self.current);
}
}
#[derive(Debug)]
#[non_exhaustive]
pub struct BufferSizeQuery<'a> {
pub current_size: u32,
pub used: u32,
pub returnable: u32,
pub seq: u8,
pub report: &'a GlobalCreditsReport,
pub report_is_current: bool,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct BufferSize {
pub size: u32,
pub return_threshold: u32,
pub force_return: bool,
}
impl BufferSize {
pub fn new(size: u32) -> Self {
Self { size, return_threshold: (size / 10).clamp(1, 1_048_576), force_return: false }
}
}
pub trait BufferSizer: fmt::Debug + Send + Sync + 'static {
fn duplicate(&self) -> Box<dyn BufferSizer>;
fn initial(&mut self) -> BufferSize;
fn size<'a>(&mut self, query: BufferSizeQuery<'a>) -> BufferSize;
}
impl Clone for Box<dyn BufferSizer> {
fn clone(&self) -> Self {
self.duplicate()
}
}
#[derive(Debug)]
pub(crate) struct DummySizer;
impl DummySizer {
#[allow(clippy::new_ret_no_self)]
pub fn new() -> Box<dyn BufferSizer> {
Box::new(Self)
}
}
impl BufferSizer for DummySizer {
fn duplicate(&self) -> Box<dyn BufferSizer> {
unreachable!()
}
fn initial(&mut self) -> BufferSize {
unreachable!()
}
fn size<'a>(&mut self, state: BufferSizeQuery<'a>) -> BufferSize {
let _ = state;
unreachable!()
}
}
#[derive(Debug, Clone)]
pub struct FixedBuffer(BufferSize);
impl FixedBuffer {
#[allow(clippy::new_ret_no_self)]
pub fn new(size: u32) -> Box<dyn BufferSizer> {
Box::new(Self(BufferSize::new(size)))
}
pub const fn size(&self) -> u32 {
self.0.size
}
}
impl BufferSizer for FixedBuffer {
fn duplicate(&self) -> Box<dyn BufferSizer> {
Box::new(self.clone())
}
fn initial(&mut self) -> BufferSize {
self.0.clone()
}
fn size<'a>(&mut self, query: BufferSizeQuery<'a>) -> BufferSize {
let _ = query;
self.0.clone()
}
}
#[derive(Debug, Clone)]
pub struct DynamicBuffer {
min: u32,
max: u32,
pub level_quot: u32,
current: BufferSize,
low_level: u32,
high_level: u32,
record_max: u32,
}
impl DynamicBuffer {
#[allow(clippy::new_ret_no_self)]
pub fn new(min: u32, max: u32) -> Box<dyn BufferSizer> {
assert!(min <= max);
let this = Self {
min,
max,
level_quot: 2,
current: BufferSize::new(0),
low_level: 0,
high_level: 0,
record_max: 0,
};
Box::new(this)
}
fn set_size(&mut self, mut size: u32) {
size = size.clamp(self.min, self.max);
if self.current.size == size {
return;
}
self.current = BufferSize::new(size);
self.low_level = (size / self.level_quot).clamp(1024, 1_048_576);
self.high_level = 4 * self.low_level;
const MB: f32 = 1_048_576.;
tracing::trace!("adjusting receive buffer size to {:.1} MB", size as f32 / MB);
if size > self.record_max {
self.record_max = size;
tracing::debug!("maximum receive buffer size increased to {:.1} MB", self.record_max as f32 / MB);
}
}
}
impl BufferSizer for DynamicBuffer {
fn duplicate(&self) -> Box<dyn BufferSizer> {
Box::new(self.clone())
}
fn initial(&mut self) -> BufferSize {
self.set_size(self.min);
self.current.clone()
}
fn size<'a>(&mut self, query: BufferSizeQuery<'a>) -> BufferSize {
if !query.report_is_current || query.current_size != self.current.size {
return self.current.clone();
}
if query.report.min < self.low_level {
tracing::trace!("computing receive buffer size: {query:?}");
self.set_size(self.current.size.saturating_mul(4));
} else if query.report.min > self.high_level {
tracing::trace!("computing receive buffer size: {query:?}");
let diff = ((query.report.min - self.high_level) / 2).min(65_536);
self.set_size(self.current.size.saturating_sub(diff));
}
self.current.clone()
}
}