use anyhow::bail;
use bare_metrics_core::structures::{
get_supported_version, Frame, LogHeader, UnixTimestampMilliseconds,
};
use std::io::{Read, Seek, SeekFrom};
#[derive(Copy, Clone, Debug, Ord, PartialOrd, Eq, PartialEq, Hash)]
pub struct SeekToken {
pub start_ts: UnixTimestampMilliseconds,
pub offset: u64,
}
pub struct MetricsLogReader<R: Read> {
reader: R,
pub header: LogHeader,
last_read_ts: UnixTimestampMilliseconds,
}
impl<R: Read> MetricsLogReader<R> {
pub fn new(mut reader: R) -> anyhow::Result<Self> {
let header: LogHeader = serde_bare::from_reader(&mut reader)?;
if header.bare_metrics_version != get_supported_version() {
bail!("Wrong version. Expected {:?} got {:?}. Later versions of Bare Metrics may use a stable format.", get_supported_version(), header.bare_metrics_version);
}
let last_read_ts = header.start_time;
Ok(MetricsLogReader {
reader,
header,
last_read_ts,
})
}
pub fn read_frame(&mut self) -> anyhow::Result<Option<(UnixTimestampMilliseconds, Frame)>> {
let mut interceptor = EofTrackingReadInterceptor::new(&mut self.reader);
match serde_bare::from_reader::<_, Frame>(&mut interceptor) {
Ok(frame) => {
let start_ts = self.last_read_ts;
self.last_read_ts = frame.end_time;
Ok(Some((start_ts, frame)))
}
Err(err) if err.classify().is_eof() => Ok(None),
Err(other_err) => {
let eof_flag = interceptor.was_eof();
if eof_flag == Some(true) {
Ok(None)
} else {
bail!(
"Failed to read frame: {:?} class {:?}, intercepted eof flag {:?}",
other_err,
other_err.classify(),
eof_flag
);
}
}
}
}
}
struct EofTrackingReadInterceptor<R: Read> {
inner: R,
was_eof_flag: Option<bool>,
}
impl<R: Read> EofTrackingReadInterceptor<R> {
pub fn new(inner: R) -> EofTrackingReadInterceptor<R> {
EofTrackingReadInterceptor {
inner,
was_eof_flag: None,
}
}
pub fn was_eof(self) -> Option<bool> {
self.was_eof_flag
}
}
impl<R: Read> Read for EofTrackingReadInterceptor<R> {
fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
if self.was_eof_flag.is_none() {
let count = self.inner.read(buf)?;
if count == 0 {
self.was_eof_flag = Some(true);
} else {
self.was_eof_flag = Some(false);
}
Ok(count)
} else {
self.inner.read(buf)
}
}
}
impl<R: Read + Seek> MetricsLogReader<R> {
pub fn read_frame_rewindable(
&mut self,
) -> anyhow::Result<Option<(SeekToken, UnixTimestampMilliseconds, Frame)>> {
let current_pos_in_file = self.reader.stream_position()?;
if let Some((timestamp, frame)) = self.read_frame()? {
Ok(Some((
SeekToken {
start_ts: timestamp,
offset: current_pos_in_file,
},
timestamp,
frame,
)))
} else {
self.reader.seek(SeekFrom::Start(current_pos_in_file))?;
Ok(None)
}
}
pub fn seek(&mut self, seek_token: SeekToken) -> anyhow::Result<SeekToken> {
let SeekToken {
start_ts: seek_timestamp,
offset: seek_pos,
} = seek_token;
let old_pos_in_file = self.reader.stream_position()?;
let old_timestamp = self.last_read_ts;
self.reader.seek(SeekFrom::Start(seek_pos))?;
self.last_read_ts = seek_timestamp;
Ok(SeekToken {
start_ts: old_timestamp,
offset: old_pos_in_file,
})
}
}