aboutsummaryrefslogtreecommitdiffstats
path: root/log_db/src/log_reader_forward.rs
blob: a1979b3c999be64dbdf5690f57d298736c41d842 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
use super::*;

pub struct ForwardLogReader {
    metadata_reader: io::BufReader<fs::File>,
    data_reader: io::BufReader<fs::File>,
}

pub struct ForwardLogReaderItem {
    pub record: Record,
    pub index: u64,
}

impl ForwardLogReader {
    pub fn new(metadata_file: fs::File, data_file: fs::File) -> 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))
            .expect("Seek failed");

        ret
    }

    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 + METADATA_ROW_LENGTH as u64 * index,
            ))
            .expect("Seek failed");

        ret
    }

    fn read_record(&mut self) -> Result<Option<ForwardLogReaderItem>, io::Error> {
        loop {
            let pos = self.metadata_reader.stream_position()?;
            let index = (pos - METADATA_FILE_HEADER_SIZE as u64) / METADATA_ROW_LENGTH as u64;

            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 {
                    return Ok(None);
                } else {
                    return Err(e);
                }
            }

            // 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());

            if entry_offset == 0 && entry_length == 0 {
                // This is an unused entry in the metadata file, skip
                continue;
            }

            // Use .seek_relative instead of .seek to avoid dropping the BufReader internal buffer when
            // the seek distance is small
            let seek_distance = entry_offset as i64 - self.data_reader.stream_position()? as i64;
            self.data_reader.seek_relative(seek_distance)?;

            let mut result_buf = vec![0; entry_length as usize];
            self.data_reader.read_exact(&mut result_buf)?;

            let record = Record::deserialize(&result_buf);
            return Ok(Some(ForwardLogReaderItem { record, index }));
        }
    }
}

impl Iterator for ForwardLogReader {
    type Item = ForwardLogReaderItem;

    fn next(&mut self) -> Option<Self::Item> {
        self.read_record().unwrap_or_else(|err| {
            panic!("Error reading record: {:?}", err);
        })
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    const TEST_RESOURCES_DIR: &str = "tests/resources";

    #[test]
    fn test_forward_log_reader_fixture_db1() {
        let _ = env_logger::builder().is_test(true).try_init();
        let metadata_path = Path::new(TEST_RESOURCES_DIR).join("test_metadata_1");
        let data_path = Path::new(TEST_RESOURCES_DIR).join("test_data_1");
        let metadata_file = fs::OpenOptions::new()
            .read(true)
            .open(&metadata_path)
            .expect("Failed to open metadata file");
        let data_file = fs::OpenOptions::new()
            .read(true)
            .open(&data_path)
            .expect("Failed to open data file");

        let mut forward_log_reader = ForwardLogReader::new(metadata_file, data_file);

        // There are two records in the log with "schema" with one field: Bytes

        let first_record = forward_log_reader
            .next()
            .expect("Failed to read the first record");
        assert!(match &first_record.record.values[..] {
            [Value::Bytes(bytes)] => bytes.len() == 256,
            _ => false,
        });

        assert!(forward_log_reader.next().is_none());
    }
}