/home/runner/work/feoxdb/feoxdb/src/core/store/persistence.rs
Line | Count | Source |
1 | | use std::fs::{File, OpenOptions}; |
2 | | use std::io::{self, Read}; |
3 | | use std::sync::atomic::Ordering; |
4 | | use std::sync::Arc; |
5 | | |
6 | | use bytes::Bytes; |
7 | | |
8 | | use crate::constants::*; |
9 | | use crate::core::record::Record; |
10 | | use crate::error::{FeoxError, Result}; |
11 | | use crate::storage::format::{get_format_ref, sector_holds_record}; |
12 | | |
13 | | use super::FeoxStore; |
14 | | |
15 | | const ZERO_SCAN_BLOCKS: usize = 256; |
16 | | |
17 | 154 | fn validate_device_size(size: u64) -> Result<()> { |
18 | 154 | let reserved_size = FEOX_DATA_START_BLOCK * FEOX_BLOCK_SIZE as u64; |
19 | 154 | if size <= reserved_size |
20 | 152 | || size > MAX_DEVICE_SIZE |
21 | 151 | || !size.is_multiple_of(FEOX_BLOCK_SIZE as u64) |
22 | | { |
23 | 5 | return Err(FeoxError::InvalidDevice); |
24 | 149 | } |
25 | 149 | Ok(()) |
26 | 154 | } |
27 | | |
28 | 76 | fn file_is_all_zero(_file: &std::fs::File, path: &str, size: u64) -> Result<bool> { |
29 | | #[cfg(target_os = "linux")] |
30 | 76 | if sparse_file_has_no_data(_file)?0 { |
31 | 1 | return Ok(true); |
32 | 75 | } |
33 | | |
34 | 75 | let mut contents = std::fs::OpenOptions::new() |
35 | 75 | .read(true) |
36 | 75 | .open(path) |
37 | 75 | .map_err(FeoxError::IoError)?0 ; |
38 | 75 | let mut buffer = vec![0; ZERO_SCAN_BLOCKS * FEOX_BLOCK_SIZE]; |
39 | 75 | let mut remaining = size; |
40 | 75 | while remaining > 0 { |
41 | 75 | let read_len = remaining.min(buffer.len() as u64) as usize; |
42 | 75 | contents |
43 | 75 | .read_exact(&mut buffer[..read_len]) |
44 | 75 | .map_err(FeoxError::IoError)?0 ; |
45 | 69.7k | if buffer[..read_len].iter()75 .any75 (|byte| *byte != 0) { |
46 | 75 | return Ok(false); |
47 | 0 | } |
48 | 0 | remaining -= read_len as u64; |
49 | | } |
50 | 0 | Ok(true) |
51 | 76 | } |
52 | | |
53 | | #[cfg(target_os = "linux")] |
54 | 76 | fn sparse_file_has_no_data(file: &std::fs::File) -> Result<bool> { |
55 | | use std::os::fd::AsRawFd; |
56 | | |
57 | 76 | let offset = unsafe { libc::lseek(file.as_raw_fd(), 0, libc::SEEK_DATA) }; |
58 | 76 | if offset >= 0 { |
59 | 75 | return Ok(false); |
60 | 1 | } |
61 | | |
62 | 1 | let error = std::io::Error::last_os_error(); |
63 | 1 | match error.raw_os_error() { |
64 | 1 | Some(libc::ENXIO) => Ok(true), |
65 | 0 | Some(libc::EINVAL) => Ok(false), |
66 | 0 | _ => Err(FeoxError::IoError(error)), |
67 | | } |
68 | 76 | } |
69 | | |
70 | | impl FeoxStore { |
71 | | /// Force flush all pending writes to disk. |
72 | | /// |
73 | | /// In persistent mode, ensures all buffered writes are flushed to disk. |
74 | | /// In memory-only mode, this is a no-op. |
75 | | /// |
76 | | /// # Example |
77 | | /// |
78 | | /// ```no_run |
79 | | /// # use feoxdb::FeoxStore; |
80 | | /// # fn main() -> feoxdb::Result<()> { |
81 | | /// let store = FeoxStore::new(Some("/path/to/data.feox".to_string()))?; |
82 | | /// store.insert(b"important", b"data")?; |
83 | | /// store.flush_all()?; // Ensure data is persisted |
84 | | /// # Ok(()) |
85 | | /// # } |
86 | | /// ``` |
87 | 109 | pub fn flush_all(&self) -> Result<()> { |
88 | 109 | if self.initialized && !self.memory_only { |
89 | | // First flush the write buffer to ensure all data is written |
90 | 109 | if let Some(ref wb) = self.write_buffer { |
91 | 109 | wb.force_flush()?4 ; |
92 | 0 | } |
93 | | |
94 | 105 | if let Some(ref disk_io) = self.disk_io { |
95 | | // Update metadata with current stats |
96 | 105 | let mut metadata = self._metadata.write(); |
97 | 105 | metadata.total_records = self.stats.record_count.load(Ordering::Relaxed) as u64; |
98 | 105 | metadata.total_size = self.stats.disk_usage.load(Ordering::Relaxed); |
99 | 105 | metadata.fragmentation = self.free_space.read().get_fragmentation(); |
100 | 105 | metadata.update(); |
101 | | |
102 | | // Write metadata |
103 | 105 | disk_io.write().write_store_metadata(&mut metadata)?0 ; |
104 | 0 | } |
105 | 0 | } |
106 | 105 | Ok(()) |
107 | 109 | } |
108 | | |
109 | 99 | pub(super) fn load_value_from_disk(&self, record: &Arc<Record>) -> Result<Bytes> { |
110 | 99 | let mut source = Arc::clone(record); |
111 | | loop { |
112 | 99 | if let Some(value0 ) = source.get_value() { |
113 | 0 | return Ok(value); |
114 | 99 | } |
115 | 99 | if source.sector.load(Ordering::Acquire) != 0 { |
116 | 99 | break; |
117 | 0 | } |
118 | 0 | source = source.value_source().ok_or(FeoxError::StaleExtent)?; |
119 | | } |
120 | 99 | let extent = source.acquire_extent().ok_or(FeoxError::StaleExtent)?0 ; |
121 | 99 | let sector = source.sector.load(Ordering::Acquire); |
122 | 99 | if self.memory_only || sector == 0 { |
123 | 0 | return Err(FeoxError::StaleExtent); |
124 | 99 | } |
125 | 99 | crate::test_hooks::pause_at(crate::test_hooks::AFTER_SECTOR_LOAD); |
126 | | |
127 | | // Get the appropriate format handler |
128 | 99 | let format = get_format_ref(self.format_version); |
129 | | |
130 | | // Calculate how many sectors we need to read |
131 | 99 | let total_size = format.total_size(source.key.len(), source.value_len); |
132 | 99 | let sectors_needed = total_size.div_ceil(FEOX_BLOCK_SIZE); |
133 | | |
134 | | // Read the sectors |
135 | 99 | let disk_io = self |
136 | 99 | .disk_io |
137 | 99 | .as_ref() |
138 | 99 | .ok_or_else(|| {0 |
139 | 0 | FeoxError::IoError(io::Error::new( |
140 | 0 | io::ErrorKind::NotFound, |
141 | 0 | "No disk IO available", |
142 | 0 | )) |
143 | 0 | })? |
144 | 99 | .read(); |
145 | | |
146 | 99 | let data = disk_io.read_sectors_sync(sector, sectors_needed as u64)?0 ; |
147 | 99 | drop(extent); |
148 | | |
149 | 99 | if !sector_holds_record(&data, &source) { |
150 | 0 | return Err(FeoxError::StaleExtent); |
151 | 99 | } |
152 | | |
153 | 99 | let offset = format.value_offset(source.key.len()); |
154 | 99 | let end = offset |
155 | 99 | .checked_add(source.value_len) |
156 | 99 | .filter(|end| *end <= data.len()) |
157 | 99 | .ok_or(FeoxError::InvalidRecord)?0 ; |
158 | 99 | if source.value_len <= data.len() / 2 { |
159 | 83 | Ok(Bytes::copy_from_slice(&data[offset..end])) |
160 | | } else { |
161 | 16 | Ok(Bytes::from(data).slice(offset..end)) |
162 | | } |
163 | 99 | } |
164 | | |
165 | 118 | pub(super) fn open_device( |
166 | 118 | &mut self, |
167 | 118 | device_path: &Option<String>, |
168 | 118 | file_size: Option<u64>, |
169 | 118 | ) -> Result<()> { |
170 | 118 | if let Some(path) = device_path { |
171 | | // Open the device/file |
172 | | #[cfg(target_os = "linux")] |
173 | | use std::os::unix::fs::OpenOptionsExt; |
174 | | |
175 | | #[cfg(unix)] |
176 | 118 | let (file, use_direct_io) = if std::path::Path::new("/.dockerenv").exists() { |
177 | 0 | let file = OpenOptions::new() |
178 | 0 | .read(true) |
179 | 0 | .write(true) |
180 | 0 | .create(true) |
181 | 0 | .truncate(false) |
182 | 0 | .open(path) |
183 | 0 | .map_err(FeoxError::IoError)?; |
184 | 0 | (file, false) // Don't use O_DIRECT in Docker |
185 | | } else { |
186 | | // Try with O_DIRECT on Linux, fall back without it on other Unix systems |
187 | | #[cfg(target_os = "linux")] |
188 | | { |
189 | | // Try to open with O_DIRECT first |
190 | 118 | match OpenOptions::new() |
191 | 118 | .read(true) |
192 | 118 | .write(true) |
193 | 118 | .create(true) |
194 | 118 | .truncate(false) |
195 | 118 | .custom_flags(libc::O_DIRECT) |
196 | 118 | .open(path) |
197 | | { |
198 | 118 | Ok(file) => (file, true), // Successfully opened with O_DIRECT |
199 | | Err(_) => { |
200 | | // Fallback to regular open |
201 | 0 | let file = OpenOptions::new() |
202 | 0 | .read(true) |
203 | 0 | .write(true) |
204 | 0 | .create(true) |
205 | 0 | .truncate(false) |
206 | 0 | .open(path) |
207 | 0 | .map_err(FeoxError::IoError)?; |
208 | 0 | (file, false) |
209 | | } |
210 | | } |
211 | | } |
212 | | #[cfg(not(target_os = "linux"))] |
213 | | { |
214 | | let file = OpenOptions::new() |
215 | | .read(true) |
216 | | .write(true) |
217 | | .create(true) |
218 | | .truncate(false) |
219 | | .open(path) |
220 | | .map_err(FeoxError::IoError)?; |
221 | | (file, false) // O_DIRECT not supported on this platform |
222 | | } |
223 | | }; |
224 | | |
225 | | #[cfg(not(unix))] |
226 | | let file = OpenOptions::new() |
227 | | .read(true) |
228 | | .write(true) |
229 | | .create(true) |
230 | | .truncate(false) |
231 | | .open(path) |
232 | | .map_err(FeoxError::IoError)?; |
233 | | |
234 | | // Get file size |
235 | 118 | let metadata = file.metadata().map_err(FeoxError::IoError)?0 ; |
236 | 118 | self.device_size = metadata.len(); |
237 | | |
238 | | // Track whether this is a newly created file |
239 | 118 | self.fresh_device = self.device_size == 0; |
240 | | |
241 | 118 | if self.fresh_device { |
242 | 40 | self.initialize_fresh_device(&file, file_size)?3 ; |
243 | | } else { |
244 | 78 | validate_device_size(self.device_size)?2 ; |
245 | 76 | let is_empty_file = file_is_all_zero(&file, path, self.device_size)?0 ; |
246 | | |
247 | 76 | if is_empty_file { |
248 | 1 | self.free_space.write().initialize(self.device_size)?0 ; |
249 | 1 | self.fresh_device = true; |
250 | 1 | let mut metadata = self._metadata.write(); |
251 | 1 | metadata.device_size = self.device_size; |
252 | 1 | metadata.update(); |
253 | 75 | } else { |
254 | 75 | self.free_space.write().set_device_size(self.device_size); |
255 | 75 | } |
256 | | } |
257 | | |
258 | | #[cfg(not(unix))] |
259 | | let use_direct_io = false; |
260 | 113 | self.attach_device_file(file, use_direct_io)?0 ; |
261 | 0 | } |
262 | 113 | Ok(()) |
263 | 118 | } |
264 | | |
265 | 26 | pub(super) fn open_device_read_only(&mut self, file: File) -> Result<()> { |
266 | 26 | self.device_size = file.metadata().map_err(FeoxError::IoError)?0 .len(); |
267 | 26 | validate_device_size(self.device_size)?0 ; |
268 | 26 | self.free_space.write().set_device_size(self.device_size); |
269 | 26 | self.attach_device_file(file, false) |
270 | 26 | } |
271 | | |
272 | 10 | pub(super) fn open_fresh_device(&mut self, file: File, file_size: Option<u64>) -> Result<()> { |
273 | 10 | if file.metadata().map_err(FeoxError::IoError)?0 .len() != 0 { |
274 | 0 | return Err(FeoxError::InvalidDevice); |
275 | 10 | } |
276 | 10 | self.fresh_device = true; |
277 | 10 | self.initialize_fresh_device(&file, file_size)?0 ; |
278 | 10 | self.attach_device_file(file, false) |
279 | 10 | } |
280 | | |
281 | 50 | fn initialize_fresh_device(&mut self, file: &File, file_size: Option<u64>) -> Result<()> { |
282 | 50 | let target_size = file_size.unwrap_or(DEFAULT_DEVICE_SIZE); |
283 | 50 | validate_device_size(target_size)?3 ; |
284 | 47 | file.set_len(target_size).map_err(FeoxError::IoError)?0 ; |
285 | 47 | self.device_size = target_size; |
286 | 47 | self.free_space.write().initialize(self.device_size)?0 ; |
287 | | |
288 | 47 | let mut metadata = self._metadata.write(); |
289 | 47 | metadata.device_size = self.device_size; |
290 | 47 | metadata.update(); |
291 | 47 | Ok(()) |
292 | 50 | } |
293 | | |
294 | 149 | fn attach_device_file(&mut self, file: File, use_direct_io: bool) -> Result<()> { |
295 | | #[cfg(unix)] |
296 | | { |
297 | | use std::os::unix::io::AsRawFd; |
298 | 149 | let file = Arc::new(file); |
299 | 149 | self.device_fd = Some(file.as_raw_fd()); |
300 | 149 | self.device_file = Some(file.as_ref().try_clone().map_err(FeoxError::IoError)?0 ); |
301 | 149 | let disk_io = crate::storage::io::DiskIO::new(file, use_direct_io)?0 ; |
302 | 149 | self.disk_io = Some(Arc::new(parking_lot::RwLock::new(disk_io))); |
303 | | } |
304 | | |
305 | | #[cfg(not(unix))] |
306 | | { |
307 | | let _ = use_direct_io; |
308 | | self.device_file = Some(file.try_clone().map_err(FeoxError::IoError)?); |
309 | | let disk_io = crate::storage::io::DiskIO::new_from_file(file)?; |
310 | | self.disk_io = Some(Arc::new(parking_lot::RwLock::new(disk_io))); |
311 | | } |
312 | | |
313 | 149 | Ok(()) |
314 | 149 | } |
315 | | } |
316 | | |
317 | | impl Drop for FeoxStore { |
318 | 247 | fn drop(&mut self) { |
319 | | // Stop TTL sweeper if running |
320 | 247 | if let Some(mut sweeper0 ) = self.ttl_sweeper.write().take() { |
321 | 0 | sweeper.stop(); |
322 | 247 | } |
323 | | |
324 | | // Signal shutdown to write buffer workers |
325 | 247 | if let Some(ref wb106 ) = self.write_buffer { |
326 | 106 | wb.initiate_shutdown(); |
327 | 141 | } |
328 | | |
329 | 247 | if let Some(write_buffer106 ) = self.write_buffer.take() { |
330 | 106 | write_buffer.finish_shutdown(); |
331 | 141 | } |
332 | | |
333 | | // Write metadata directly without using the write buffer |
334 | 247 | if self.initialized && !self.memory_only229 && !self.read_only130 { |
335 | 106 | if let Some(ref disk_io) = self.disk_io { |
336 | 106 | // Update metadata with current stats |
337 | 106 | let mut metadata = self._metadata.write(); |
338 | 106 | metadata.total_records = self.stats.record_count.load(Ordering::Relaxed) as u64; |
339 | 106 | metadata.total_size = self.stats.disk_usage.load(Ordering::Relaxed); |
340 | 106 | metadata.fragmentation = self.free_space.read().get_fragmentation(); |
341 | 106 | metadata.update(); |
342 | 106 | |
343 | 106 | // Write metadata |
344 | 106 | let _ = disk_io.write().write_store_metadata(&mut metadata); |
345 | 106 | }0 |
346 | 141 | } |
347 | | |
348 | | // Now it's safe to shutdown disk I/O since workers have exited |
349 | 247 | if let Some(ref disk_io143 ) = self.disk_io { |
350 | 143 | disk_io.write().shutdown(); |
351 | 143 | }104 |
352 | 247 | } |
353 | | } |