use bytes::Bytes;
use serde::Deserialize;
use serde::Serialize;
use crate::consts::DEFAULT_TTL_MS;
use crate::consts::MAX_CHUNK_ENVELOPE_OVERHEAD;
use crate::consts::MIN_CHUNK_DATA;
use crate::consts::TRANSPORT_CUSTOM_OVERHEAD;
use crate::error::Error;
use crate::error::Result;
use crate::utils::get_epoch_ms;
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Chunk {
pub chunk: [usize; 2],
pub data: Bytes,
pub meta: ChunkMeta,
}
impl Chunk {
pub fn to_wire(&self) -> Result<Bytes> {
rings_codec::serialize(self)
.map(Bytes::from)
.map_err(Error::CodecSerialize)
}
pub fn from_wire(data: &[u8]) -> Result<Self> {
rings_codec::deserialize(data).map_err(Error::CodecDeserialize)
}
}
#[derive(Debug, Copy, Clone, Deserialize, Serialize)]
pub struct ChunkMeta {
pub id: uuid::Uuid,
pub ts_ms: u128,
pub ttl_ms: u64,
}
impl Default for ChunkMeta {
fn default() -> Self {
Self {
id: crate::utils::new_uuid(),
ts_ms: get_epoch_ms(),
ttl_ms: DEFAULT_TTL_MS,
}
}
}
#[derive(Debug, Clone, Default, Deserialize, Serialize)]
pub struct ChunkList(Vec<Chunk>);
impl ChunkList {
pub fn split(bytes: &Bytes, chunk_size: usize) -> Self {
let chunk_size = chunk_size.max(1);
let chunks: Vec<Bytes> = bytes
.chunks(chunk_size)
.map(|c| c.to_vec().into())
.collect();
let chunks_len: usize = chunks.len();
let meta = ChunkMeta::default();
Self(
chunks
.into_iter()
.enumerate()
.map(|(i, data)| Chunk {
meta,
chunk: [i, chunks_len],
data,
})
.collect::<Vec<Chunk>>(),
)
}
pub fn stream(bytes: Bytes, chunk_size: usize) -> impl Iterator<Item = Chunk> {
let chunk_size = chunk_size.max(1);
let total = bytes.len().div_ceil(chunk_size);
let meta = ChunkMeta::default();
(0..total).map(move |i| {
let start = i * chunk_size;
let end = start.saturating_add(chunk_size).min(bytes.len());
Chunk {
meta,
chunk: [i, total],
data: bytes.slice(start..end),
}
})
}
pub fn to_vec(&self) -> Vec<Chunk> {
self.0.clone()
}
pub fn as_vec(&self) -> &Vec<Chunk> {
&self.0
}
}
impl IntoIterator for &ChunkList {
type Item = Chunk;
type IntoIter = std::vec::IntoIter<Chunk>;
fn into_iter(self) -> Self::IntoIter {
self.to_vec().into_iter()
}
}
impl IntoIterator for ChunkList {
type Item = Chunk;
type IntoIter = std::vec::IntoIter<Chunk>;
fn into_iter(self) -> Self::IntoIter {
self.0.into_iter()
}
}
impl From<ChunkList> for Vec<Chunk> {
fn from(l: ChunkList) -> Self {
l.0
}
}
impl From<Vec<Chunk>> for ChunkList {
fn from(data: Vec<Chunk>) -> Self {
Self(data)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Framing {
Whole,
Chunked {
chunk_size: usize,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WireReserves {
pub whole: usize,
pub chunk: usize,
pub min_chunk_data: usize,
}
impl WireReserves {
pub const PRODUCTION: Self = Self {
whole: TRANSPORT_CUSTOM_OVERHEAD,
chunk: MAX_CHUNK_ENVELOPE_OVERHEAD + TRANSPORT_CUSTOM_OVERHEAD,
min_chunk_data: MIN_CHUNK_DATA,
};
pub fn plan(&self, payload_len: usize, max_message_size: usize) -> Option<Framing> {
let whole_fits = payload_len
.checked_add(self.whole)
.is_some_and(|wire| wire <= max_message_size);
if whole_fits {
return Some(Framing::Whole);
}
let min_viable = self.chunk.checked_add(self.min_chunk_data)?;
(max_message_size >= min_viable).then(|| Framing::Chunked {
chunk_size: max_message_size - self.chunk,
})
}
}