/home/runner/work/feoxdb/feoxdb/src/storage/format.rs
Line | Count | Source |
1 | | use crate::constants::*; |
2 | | use crate::core::record::Record; |
3 | | use crate::storage::seq_token::seq_token; |
4 | | use std::sync::atomic::Ordering; |
5 | | |
6 | | /// Trait for handling different record format versions |
7 | | pub trait RecordFormat: Send + Sync { |
8 | | /// Calculate the size of a record on disk (excluding value) |
9 | | fn record_header_size(&self, key_len: usize) -> usize; |
10 | | |
11 | | /// Calculate total size including value |
12 | | fn total_size(&self, key_len: usize, value_len: usize) -> usize; |
13 | | |
14 | | /// Serialize a record to bytes for disk storage |
15 | | fn serialize_record(&self, record: &Record, include_value: bool) -> Vec<u8>; |
16 | | |
17 | | /// Append a serialized record without its sector header. |
18 | 0 | fn serialize_record_into(&self, record: &Record, include_value: bool, data: &mut Vec<u8>) { |
19 | 0 | data.extend_from_slice(&self.serialize_record(record, include_value)); |
20 | 0 | } |
21 | | |
22 | | /// Parse a record from disk bytes (returns key, value_len, timestamp, ttl_expiry) |
23 | | fn parse_record(&self, data: &[u8]) -> Option<(Vec<u8>, usize, u64, u64)>; |
24 | | |
25 | | /// Get the offset where value data starts in the serialized format |
26 | | fn value_offset(&self, key_len: usize) -> usize; |
27 | | } |
28 | | |
29 | 346 | fn serialization_buffer( |
30 | 346 | format: &dyn RecordFormat, |
31 | 346 | record: &Record, |
32 | 346 | include_value: bool, |
33 | 346 | ) -> Vec<u8> { |
34 | 346 | let value_len = if include_value { record.value_len } else { 00 }; |
35 | 346 | Vec::with_capacity(format.total_size(record.key.len(), value_len) - SECTOR_HEADER_SIZE) |
36 | 346 | } |
37 | | |
38 | | /// Version 1 format (no TTL support) |
39 | | pub struct FormatV1; |
40 | | static FORMAT_V1: FormatV1 = FormatV1; |
41 | | |
42 | | impl RecordFormat for FormatV1 { |
43 | 48 | fn record_header_size(&self, key_len: usize) -> usize { |
44 | 48 | SECTOR_HEADER_SIZE + 2 + key_len + 8 + 8 // header + key_len(2) + key + value_len(8) + timestamp(8) |
45 | 48 | } |
46 | | |
47 | 37 | fn total_size(&self, key_len: usize, value_len: usize) -> usize { |
48 | 37 | self.record_header_size(key_len) + value_len |
49 | 37 | } |
50 | | |
51 | 8 | fn serialize_record(&self, record: &Record, include_value: bool) -> Vec<u8> { |
52 | 8 | let mut data = serialization_buffer(self, record, include_value); |
53 | 8 | self.serialize_record_into(record, include_value, &mut data); |
54 | 8 | data |
55 | 8 | } |
56 | | |
57 | 15 | fn serialize_record_into(&self, record: &Record, include_value: bool, data: &mut Vec<u8>) { |
58 | 15 | data.extend_from_slice(&(record.key.len() as u16).to_le_bytes()); |
59 | 15 | data.extend_from_slice(&record.key); |
60 | 15 | data.extend_from_slice(&(record.value_len as u64).to_le_bytes()); |
61 | 15 | data.extend_from_slice(&record.timestamp.to_le_bytes()); |
62 | | |
63 | 15 | if include_value { |
64 | 9 | if let Some(value) = record.value.read().as_ref() { |
65 | 9 | data.extend_from_slice(value); |
66 | 9 | }0 |
67 | 6 | } |
68 | 15 | } |
69 | | |
70 | 11 | fn parse_record(&self, data: &[u8]) -> Option<(Vec<u8>, usize, u64, u64)> { |
71 | 11 | if data.len() < SECTOR_HEADER_SIZE + 2 { |
72 | 0 | return None; |
73 | 11 | } |
74 | | |
75 | 11 | let mut offset = SECTOR_HEADER_SIZE + 2; |
76 | 11 | let key_len = u16::from_le_bytes( |
77 | 11 | data[SECTOR_HEADER_SIZE..SECTOR_HEADER_SIZE + 2] |
78 | 11 | .try_into() |
79 | 11 | .ok()?0 , |
80 | | ) as usize; |
81 | | |
82 | 11 | if offset + key_len + 16 > data.len() { |
83 | 0 | return None; |
84 | 11 | } |
85 | | |
86 | 11 | let key = data[offset..offset + key_len].to_vec(); |
87 | 11 | offset += key_len; |
88 | | |
89 | 11 | let value_len = u64::from_le_bytes(data[offset..offset + 8].try_into().ok()?0 ) as usize; |
90 | 11 | offset += 8; |
91 | | |
92 | 11 | let timestamp = u64::from_le_bytes(data[offset..offset + 8].try_into().ok()?0 ); |
93 | | |
94 | 11 | Some((key, value_len, timestamp, 0)) // No TTL in v1 |
95 | 11 | } |
96 | | |
97 | 11 | fn value_offset(&self, key_len: usize) -> usize { |
98 | 11 | SECTOR_HEADER_SIZE + 2 + key_len + 8 + 8 |
99 | 11 | } |
100 | | } |
101 | | |
102 | | /// Version 2 format (with TTL support) |
103 | | pub struct FormatV2; |
104 | | static FORMAT_V2: FormatV2 = FormatV2; |
105 | | |
106 | | impl RecordFormat for FormatV2 { |
107 | 23.9k | fn record_header_size(&self, key_len: usize) -> usize { |
108 | 23.9k | SECTOR_HEADER_SIZE + 2 + key_len + 8 + 8 + 8 // header + key_len(2) + key + value_len(8) + timestamp(8) + ttl(8) |
109 | 23.9k | } |
110 | | |
111 | 22.9k | fn total_size(&self, key_len: usize, value_len: usize) -> usize { |
112 | 22.9k | self.record_header_size(key_len) + value_len |
113 | 22.9k | } |
114 | | |
115 | 338 | fn serialize_record(&self, record: &Record, include_value: bool) -> Vec<u8> { |
116 | 338 | let mut data = serialization_buffer(self, record, include_value); |
117 | 338 | self.serialize_record_into(record, include_value, &mut data); |
118 | 338 | data |
119 | 338 | } |
120 | | |
121 | 21.7k | fn serialize_record_into(&self, record: &Record, include_value: bool, data: &mut Vec<u8>) { |
122 | 21.7k | data.extend_from_slice(&(record.key.len() as u16).to_le_bytes()); |
123 | 21.7k | data.extend_from_slice(&record.key); |
124 | 21.7k | data.extend_from_slice(&(record.value_len as u64).to_le_bytes()); |
125 | 21.7k | data.extend_from_slice(&record.timestamp.to_le_bytes()); |
126 | 21.7k | data.extend_from_slice(&record.ttl_expiry.load(Ordering::Acquire).to_le_bytes()); |
127 | | |
128 | 21.7k | if include_value { |
129 | 339 | if let Some(value) = record.value.read().as_ref() { |
130 | 339 | data.extend_from_slice(value); |
131 | 339 | }0 |
132 | 21.4k | } |
133 | 21.7k | } |
134 | | |
135 | 588 | fn parse_record(&self, data: &[u8]) -> Option<(Vec<u8>, usize, u64, u64)> { |
136 | 588 | if data.len() < SECTOR_HEADER_SIZE + 2 { |
137 | 0 | return None; |
138 | 588 | } |
139 | | |
140 | 588 | let mut offset = SECTOR_HEADER_SIZE + 2; |
141 | 588 | let key_len = u16::from_le_bytes( |
142 | 588 | data[SECTOR_HEADER_SIZE..SECTOR_HEADER_SIZE + 2] |
143 | 588 | .try_into() |
144 | 588 | .ok()?0 , |
145 | | ) as usize; |
146 | | |
147 | 588 | if offset + key_len + 24 > data.len() { |
148 | | // 24 = value_len(8) + timestamp(8) + ttl(8) |
149 | 0 | return None; |
150 | 588 | } |
151 | | |
152 | 588 | let key = data[offset..offset + key_len].to_vec(); |
153 | 588 | offset += key_len; |
154 | | |
155 | 588 | let value_len = u64::from_le_bytes(data[offset..offset + 8].try_into().ok()?0 ) as usize; |
156 | 588 | offset += 8; |
157 | | |
158 | 588 | let timestamp = u64::from_le_bytes(data[offset..offset + 8].try_into().ok()?0 ); |
159 | 588 | offset += 8; |
160 | | |
161 | 588 | let ttl_expiry = u64::from_le_bytes(data[offset..offset + 8].try_into().ok()?0 ); |
162 | | |
163 | 588 | Some((key, value_len, timestamp, ttl_expiry)) |
164 | 588 | } |
165 | | |
166 | 108 | fn value_offset(&self, key_len: usize) -> usize { |
167 | 108 | SECTOR_HEADER_SIZE + 2 + key_len + 8 + 8 + 8 |
168 | 108 | } |
169 | | } |
170 | | |
171 | | /// Check that a head sector still stores the record it was read for. |
172 | | /// |
173 | | /// A retired extent can be released and handed to another record while a reader |
174 | | /// still holds the old `Record` and its stale sector number, so the bytes coming |
175 | | /// back from disk must be proven to belong to this record before they are used. |
176 | | /// Everything needed is already in the buffer the reader just filled, so this |
177 | | /// costs no extra I/O. The layout prefix is identical in every format version: |
178 | | /// marker(2) seq(2) key_len(2) key value_len(8) timestamp(8). |
179 | 104 | pub fn sector_holds_record(data: &[u8], record: &Record) -> bool { |
180 | 104 | if data.len() < SECTOR_HEADER_SIZE + 2 { |
181 | 0 | return false; |
182 | 104 | } |
183 | 104 | if u16::from_le_bytes([data[0], data[1]]) != SECTOR_MARKER { |
184 | 0 | return false; |
185 | 104 | } |
186 | 104 | let key_len = |
187 | 104 | u16::from_le_bytes([data[SECTOR_HEADER_SIZE], data[SECTOR_HEADER_SIZE + 1]]) as usize; |
188 | 104 | if key_len != record.key.len() { |
189 | 0 | return false; |
190 | 104 | } |
191 | 104 | let key_at = SECTOR_HEADER_SIZE + 2; |
192 | 104 | let value_len_at = key_at + key_len; |
193 | 104 | if value_len_at + 16 > data.len() { |
194 | 0 | return false; |
195 | 104 | } |
196 | 104 | if data[key_at..value_len_at] != record.key[..] { |
197 | 0 | return false; |
198 | 104 | } |
199 | 104 | let Ok(value_len) = data[value_len_at..value_len_at + 8].try_into() else { |
200 | 0 | return false; |
201 | | }; |
202 | 104 | if u64::from_le_bytes(value_len) as usize != record.value_len { |
203 | 0 | return false; |
204 | 104 | } |
205 | 104 | let Ok(timestamp) = data[value_len_at + 8..value_len_at + 16].try_into() else { |
206 | 0 | return false; |
207 | | }; |
208 | 104 | u64::from_le_bytes(timestamp) == record.timestamp |
209 | 104 | } |
210 | | |
211 | | #[cfg(test)] |
212 | 1 | pub(crate) fn fill_retirement_extent(retired: &mut [u8], sector: u64, sectors: usize) { |
213 | 1 | debug_assert_eq!(retired.len(), sectors * FEOX_BLOCK_SIZE); |
214 | 1 | retired.fill(0); |
215 | 1 | fill_retirement_markers(retired, sector, sectors); |
216 | 1 | } |
217 | | |
218 | 70 | pub(crate) fn fill_retirement_markers(retired: &mut [u8], sector: u64, remaining: usize) { |
219 | 70 | debug_assert_eq!(retired.len() % FEOX_BLOCK_SIZE, 0); |
220 | 70 | let blocks = retired.len() / FEOX_BLOCK_SIZE; |
221 | 70 | debug_assert!(blocks <= remaining); |
222 | 1.46k | for offset in 0..blocks70 { |
223 | 1.46k | let start = offset * FEOX_BLOCK_SIZE; |
224 | 1.46k | fill_retirement_marker( |
225 | 1.46k | &mut retired[start..start + DELETION_MARKER_SIZE], |
226 | 1.46k | sector + offset as u64, |
227 | 1.46k | remaining - offset, |
228 | 1.46k | ); |
229 | 1.46k | } |
230 | 70 | } |
231 | | |
232 | 1.46k | pub(crate) fn fill_retirement_marker(marker: &mut [u8], sector: u64, remaining: usize) { |
233 | 1.46k | write_retirement_marker(marker, sector, remaining, RETIREMENT_COMPLETE); |
234 | 1.46k | } |
235 | | |
236 | | #[cfg(test)] |
237 | 1 | pub(crate) fn retirement_block(sector: u64, remaining: usize) -> Vec<u8> { |
238 | 1 | let mut retired = vec![0u8; FEOX_BLOCK_SIZE]; |
239 | 1 | fill_retirement_marker(&mut retired[..DELETION_MARKER_SIZE], sector, remaining); |
240 | 1 | retired |
241 | 1 | } |
242 | | |
243 | | #[cfg(test)] |
244 | 4 | pub(crate) fn pending_retirement_block(sector: u64, remaining: usize) -> Vec<u8> { |
245 | 4 | let mut retired = vec![0u8; FEOX_BLOCK_SIZE]; |
246 | 4 | write_retirement_marker( |
247 | 4 | &mut retired[..DELETION_MARKER_SIZE], |
248 | 4 | sector, |
249 | 4 | remaining, |
250 | | RETIREMENT_PENDING, |
251 | | ); |
252 | 4 | retired |
253 | 4 | } |
254 | | |
255 | 2.10k | pub(crate) fn retirement_marker_token(sector: u64, marker: &[u8]) -> u16 { |
256 | 2.10k | let mut protected = [0; 17]; |
257 | 2.10k | protected[..16].copy_from_slice(&marker[..16]); |
258 | 2.10k | protected[16] = marker[18]; |
259 | 2.10k | seq_token(sector, &protected) |
260 | 2.10k | } |
261 | | |
262 | 1.47k | fn write_retirement_marker(marker: &mut [u8], sector: u64, remaining: usize, state: u8) { |
263 | 1.47k | marker[..8].copy_from_slice(DELETION_MARKER); |
264 | 1.47k | marker[8..16].copy_from_slice(&(remaining as u64).to_le_bytes()); |
265 | 1.47k | marker[18] = state; |
266 | 1.47k | let token = retirement_marker_token(sector, marker); |
267 | 1.47k | marker[16..18].copy_from_slice(&token.to_le_bytes()); |
268 | 1.47k | } |
269 | | |
270 | 0 | pub fn get_format(version: u32) -> Box<dyn RecordFormat> { |
271 | 0 | match version { |
272 | 0 | 1 => Box::new(FormatV1), |
273 | 0 | 2 | 3 => Box::new(FormatV2), |
274 | 0 | _ => Box::new(FormatV2), |
275 | | } |
276 | 0 | } |
277 | | |
278 | 927 | pub(crate) fn get_format_ref(version: u32) -> &'static dyn RecordFormat { |
279 | 927 | match version { |
280 | 58 | 1 => &FORMAT_V1, |
281 | 869 | 2 | 3 => &FORMAT_V2, |
282 | 0 | _ => &FORMAT_V2, |
283 | | } |
284 | 927 | } |