use byteorder::{BE, ReadBytesExt, WriteBytesExt};
use std::{
io,
num::{NonZeroU32, NonZeroUsize},
time::Duration,
};
use crate::protocol_err;
#[cfg_attr(feature = "dump", derive(serde::Serialize, serde::Deserialize))]
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[non_exhaustive]
pub enum LinkPing {
Periodic(Duration),
WhenIdle(Duration),
WhenTimedOut,
}
#[cfg_attr(feature = "dump", derive(serde::Serialize, serde::Deserialize))]
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct UnackedIncrease {
pub consecutive: u32,
pub percent: u32,
}
#[cfg_attr(feature = "dump", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "dump", serde(default))]
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[allow(clippy::manual_non_exhaustive)]
pub struct Cfg {
pub io_write_size: NonZeroUsize,
pub ignore_flush: bool,
pub send_buffer: NonZeroU32,
pub send_queue: NonZeroUsize,
pub recv_buffer: NonZeroU32,
pub recv_queue: NonZeroUsize,
pub link_max_ping_spread: Option<NonZeroU32>,
pub no_link_timeout: Duration,
pub termination_timeout: Duration,
pub connect_queue: NonZeroUsize,
pub disconnect_on_server_id_mismatch: bool,
pub stats_intervals: Vec<Duration>,
pub link: LinkCfg,
#[doc(hidden)]
pub _non_exhaustive: (),
}
impl Default for Cfg {
fn default() -> Self {
Self {
io_write_size: NonZeroUsize::new(8_192).unwrap(),
ignore_flush: false,
send_buffer: NonZeroU32::new(134_217_728).unwrap(),
send_queue: NonZeroUsize::new(16).unwrap(),
recv_buffer: NonZeroU32::new(134_217_728).unwrap(),
recv_queue: NonZeroUsize::new(16).unwrap(),
link_max_ping_spread: Some(NonZeroU32::new(5).unwrap()),
no_link_timeout: Duration::from_secs(120),
termination_timeout: Duration::from_secs(300),
connect_queue: NonZeroUsize::new(32).unwrap(),
disconnect_on_server_id_mismatch: true,
stats_intervals: vec![
Duration::from_millis(100),
Duration::from_secs(1),
Duration::from_secs(5),
Duration::from_secs(10),
],
link: LinkCfg::default(),
_non_exhaustive: (),
}
}
}
#[cfg_attr(feature = "dump", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "dump", serde(default))]
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
#[allow(clippy::manual_non_exhaustive)]
pub struct LinkCfg {
pub io_packet_size: NonZeroUsize,
pub ack_timeout_min: Duration,
pub ack_timeout_roundtrip_factor: NonZeroU32,
pub ack_timeout_resent_factor: NonZeroU32,
pub ack_timeout_unreliable_factor: NonZeroU32,
pub ack_timeout_max: Duration,
pub unacked_init: NonZeroUsize,
pub unacked_limit: NonZeroUsize,
pub unacked_increase: Vec<UnackedIncrease>,
pub unacked_increase_single_link: u32,
pub ping: LinkPing,
pub ping_timeout: Duration,
pub max_ping: Option<Duration>,
pub test_data_limit: usize,
pub test_after_ack_timeout: bool,
pub retest_interval: Duration,
pub non_working_timeout: Duration,
pub flush_delay: Duration,
pub flush_interval: Option<Duration>,
pub ack_flush_interval: Option<Duration>,
pub unflushed_limit: Option<NonZeroUsize>,
#[doc(hidden)]
pub _non_exhaustive: (),
}
impl Default for LinkCfg {
fn default() -> Self {
Self {
io_packet_size: NonZeroUsize::new(65_536).unwrap(),
ack_timeout_min: Duration::from_secs(1),
ack_timeout_roundtrip_factor: NonZeroU32::new(3).unwrap(),
ack_timeout_resent_factor: NonZeroU32::new(3).unwrap(),
ack_timeout_unreliable_factor: NonZeroU32::new(3).unwrap(),
ack_timeout_max: Duration::from_secs(30),
unacked_init: NonZeroUsize::new(8192).unwrap(),
unacked_limit: NonZeroUsize::new(134_217_728).unwrap(),
unacked_increase: vec![
UnackedIncrease { consecutive: 0, percent: 101 },
UnackedIncrease { consecutive: 10, percent: 102 },
UnackedIncrease { consecutive: 25, percent: 105 },
UnackedIncrease { consecutive: 50, percent: 110 },
UnackedIncrease { consecutive: 100, percent: 120 },
],
unacked_increase_single_link: 200,
ping: LinkPing::WhenIdle(Duration::from_secs(15)),
ping_timeout: Duration::from_secs(40),
max_ping: None,
test_data_limit: 65_536,
test_after_ack_timeout: false,
retest_interval: Duration::from_secs(3),
non_working_timeout: Duration::from_secs(20),
flush_delay: Duration::from_millis(50),
flush_interval: None,
ack_flush_interval: Some(Duration::from_millis(50)),
unflushed_limit: Some(NonZeroUsize::new(131_072).unwrap()),
_non_exhaustive: (),
}
}
}
#[derive(Clone, Debug)]
pub(crate) struct ExchangedCfg {
pub recv_buffer: NonZeroU32,
}
impl ExchangedCfg {
pub fn write(&self, mut writer: impl io::Write) -> Result<(), io::Error> {
writer.write_u32::<BE>(self.recv_buffer.get())?;
Ok(())
}
pub fn read(mut reader: impl io::Read) -> Result<Self, io::Error> {
let this = Self {
recv_buffer: NonZeroU32::new(reader.read_u32::<BE>()?)
.ok_or_else(|| protocol_err!("recv_buffer must not be zero"))?,
};
Ok(this)
}
}
impl From<&Cfg> for ExchangedCfg {
fn from(cfg: &Cfg) -> Self {
Self { recv_buffer: cfg.recv_buffer }
}
}