use crate::Position;
use crate::event::Event;
use crate::log::set::{BATCH_OVERHEAD, PositionRange, RECORD_OVERHEAD};
use super::tips::{StagedTips, TagTips};
use super::{AppendError, Reply};
pub(super) fn measure(events: &[Event]) -> (usize, usize) {
let mut total = 0;
let mut largest = 0;
for event in events {
let framed = RECORD_OVERHEAD + event.as_bytes().len();
total += framed;
if framed > largest {
largest = framed;
}
}
(total, largest)
}
pub(super) const MARKER_BYTES: usize = BATCH_OVERHEAD;
struct Staged<'a> {
reply: &'a Reply,
range: PositionRange,
token: u64,
}
pub(super) struct Batch<'a> {
records: Vec<&'a [u8]>,
staged: Vec<Staged<'a>>,
tips: StagedTips,
first: Position,
next: Position,
}
impl<'a> Batch<'a> {
pub(super) fn new(next: Position) -> Self {
Batch {
records: Vec::new(),
staged: Vec::new(),
tips: StagedTips::new(),
first: next,
next,
}
}
pub(super) fn staged_tips(&self) -> &StagedTips {
&self.tips
}
pub(super) fn is_empty(&self) -> bool {
self.records.is_empty()
}
pub(super) fn records(&self) -> &[&'a [u8]] {
&self.records
}
pub(super) fn committed_records(&self) -> impl Iterator<Item = (Position, &'a [u8])> + '_ {
let first = self.first.get();
self.records
.iter()
.enumerate()
.map(move |(i, &bytes)| (Position::new(first + i as u64), bytes))
}
pub(super) fn stage(&mut self, events: &'a [Event], token: u64, reply: &'a Reply) {
let first = self.next;
let mut p = first;
for event in events {
self.records.push(event.as_bytes());
for tag in event.as_ref().tags() {
self.tips.record(tag, p);
}
p = p.next();
}
let range = PositionRange {
first,
last: Position::new(p.get() - 1),
};
self.next = p;
self.staged.push(Staged {
reply,
range,
token,
});
}
pub(super) fn commit_ok(
self,
committed: PositionRange,
main: &mut TagTips,
next_position: Position,
) {
debug_assert_eq!(committed.first, self.first, "batch first position mismatch");
debug_assert_eq!(
committed.last,
Position::new(self.next.get() - 1),
"batch last position mismatch",
);
main.absorb(self.tips, next_position);
for s in self.staged {
seglog::crash_point!("partial_ack");
let _ = s.reply.send((s.token, Ok(s.range)));
}
}
pub(super) fn commit_err(self, err: AppendError) {
for s in self.staged {
let _ = s.reply.send((s.token, Err(err.clone())));
}
}
}
#[cfg(test)]
mod tests {
use flume::{self as channel, Receiver};
use crate::event::{Event, EventType, Tags};
use crate::writer::WriterConfig;
use super::super::AppendError;
use super::*;
fn event(payload: &[u8]) -> Event {
Event::new(&EventType::new("E").unwrap(), &Tags::empty(), payload).unwrap()
}
type ReplyRx = Receiver<(u64, Result<PositionRange, AppendError>)>;
fn reply() -> (Reply, ReplyRx) {
channel::unbounded()
}
#[test]
fn commit_err_replies_every_staged_request() {
let e1 = vec![event(b"a")];
let e2 = vec![event(b"b"), event(b"c")];
let (tx1, rx1) = reply();
let (tx2, rx2) = reply();
let mut batch = Batch::new(Position::new(1));
batch.stage(&e1, 10, &tx1);
batch.stage(&e2, 20, &tx2);
batch.commit_err(AppendError::TooLarge { size: 99 });
assert!(matches!(
rx1.try_recv(),
Ok((10, Err(AppendError::TooLarge { size: 99 })))
));
assert!(matches!(
rx2.try_recv(),
Ok((20, Err(AppendError::TooLarge { size: 99 })))
));
}
#[test]
fn commit_ok_assigns_dense_ranges_per_request() {
let e1 = vec![event(b"a"), event(b"b")]; let e2 = vec![event(b"c")]; let (tx1, rx1) = reply();
let (tx2, rx2) = reply();
let mut batch = Batch::new(Position::new(1));
batch.stage(&e1, 10, &tx1);
batch.stage(&e2, 20, &tx2);
let cfg = WriterConfig::default();
let mut main = TagTips::new(Position::new(1), cfg.tips_window);
batch.commit_ok(
PositionRange {
first: Position::new(1),
last: Position::new(3),
},
&mut main,
Position::new(4),
);
let (token1, res1) = rx1.try_recv().unwrap();
assert_eq!(token1, 10);
assert_eq!(
res1.unwrap(),
PositionRange {
first: Position::new(1),
last: Position::new(2)
}
);
let (token2, res2) = rx2.try_recv().unwrap();
assert_eq!(token2, 20);
assert_eq!(
res2.unwrap(),
PositionRange {
first: Position::new(3),
last: Position::new(3)
}
);
}
}