aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src/log_reader_forward.rs
diff options
context:
space:
mode:
authorJan Tuomi <jan@jantuomi.fi>2024-11-09 22:30:51 +0200
committerJan Tuomi <jan@jantuomi.fi>2024-11-10 00:02:32 +0200
commitbfa8fab553eab1b0daa062010fb2ccf5d365d21f (patch)
treeebbb5575dcf361c176dc721a4d1c2b4183de90fa /log_db/src/log_reader_forward.rs
parent2975980717626b12eb959c1ecd93038b5be1d7c6 (diff)
Refresh indexes on read or at maintenance
Diffstat (limited to 'log_db/src/log_reader_forward.rs')
-rw-r--r--log_db/src/log_reader_forward.rs53
1 files changed, 44 insertions, 9 deletions
diff --git a/log_db/src/log_reader_forward.rs b/log_db/src/log_reader_forward.rs
index 343d7ff..bfd329b 100644
--- a/log_db/src/log_reader_forward.rs
+++ b/log_db/src/log_reader_forward.rs
@@ -21,7 +21,29 @@ impl<'a> ForwardLogReader {
ret
}
- fn read_record(&mut self) -> Result<Option<Record>, io::Error> {
+ pub fn new_with_index(
+ metadata_file: fs::File,
+ data_file: fs::File,
+ index: u64,
+ ) -> ForwardLogReader {
+ let mut ret = ForwardLogReader {
+ metadata_reader: io::BufReader::new(metadata_file),
+ data_reader: io::BufReader::new(data_file),
+ };
+
+ ret.metadata_reader
+ .seek(io::SeekFrom::Start(
+ METADATA_FILE_HEADER_SIZE as u64 + 16 * index,
+ ))
+ .expect("Seek failed");
+
+ ret
+ }
+
+ fn read_record(&mut self) -> Result<Option<ForwardLogReaderItem>, io::Error> {
+ let metadata_index =
+ (self.metadata_reader.stream_position()? - METADATA_FILE_HEADER_SIZE as u64) / 16;
+
let mut metadata_entry_buf = vec![0; 16]; // 2x u64
if let Err(e) = self.metadata_reader.read_exact(&mut metadata_entry_buf) {
if e.kind() == io::ErrorKind::UnexpectedEof {
@@ -32,24 +54,37 @@ impl<'a> ForwardLogReader {
}
// First u64 is the offset of the record in the data file, second is the length of the record
- let entry_offset = u64::from_be_bytes(metadata_entry_buf[0..8].try_into().unwrap());
- let entry_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap());
+ let data_offset = u64::from_be_bytes(metadata_entry_buf[0..8].try_into().unwrap());
+ let data_length = u64::from_be_bytes(metadata_entry_buf[8..16].try_into().unwrap());
// Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when
// the seek distance is small
- let seek_distance = entry_offset - self.data_reader.stream_position()?;
+ let seek_distance = data_offset - self.data_reader.stream_position()?;
self.data_reader.seek_relative(seek_distance as i64)?;
- let mut result_buf = vec![0; entry_length as usize];
+ let mut result_buf = vec![0; data_length as usize];
self.data_reader.read_exact(&mut result_buf)?;
let record = Record::deserialize(&result_buf);
- Ok(Some(record))
+ let item = ForwardLogReaderItem {
+ metadata_index,
+ data_offset,
+ data_length,
+ record,
+ };
+ Ok(Some(item))
}
}
+pub struct ForwardLogReaderItem {
+ pub metadata_index: u64,
+ pub data_offset: u64,
+ pub data_length: u64,
+ pub record: Record,
+}
+
impl Iterator for ForwardLogReader {
- type Item = Record;
+ type Item = ForwardLogReaderItem;
fn next(&mut self) -> Option<Self::Item> {
match self.read_record() {
@@ -83,10 +118,10 @@ mod tests {
// There are two records in the log with "schema" with one field: Bytes
- let first_record = forward_log_reader
+ let first_item = forward_log_reader
.next()
.expect("Failed to read the first record");
- assert!(match first_record.values() {
+ assert!(match first_item.record.values() {
[Value::Bytes(bytes)] => bytes.len() == 256,
_ => false,
});