Skip to main content

feoxdb/core/store/
persistence.rs

1use std::fs::{File, OpenOptions};
2use std::io::{self, Read};
3use std::sync::atomic::Ordering;
4use std::sync::Arc;
5
6use bytes::Bytes;
7
8use crate::constants::*;
9use crate::core::record::Record;
10use crate::error::{FeoxError, Result};
11use crate::storage::format::{get_format_ref, sector_holds_record};
12
13use super::FeoxStore;
14
15const ZERO_SCAN_BLOCKS: usize = 256;
16
17fn validate_device_size(size: u64) -> Result<()> {
18    let reserved_size = FEOX_DATA_START_BLOCK * FEOX_BLOCK_SIZE as u64;
19    if size <= reserved_size
20        || size > MAX_DEVICE_SIZE
21        || !size.is_multiple_of(FEOX_BLOCK_SIZE as u64)
22    {
23        return Err(FeoxError::InvalidDevice);
24    }
25    Ok(())
26}
27
28fn file_is_all_zero(_file: &std::fs::File, path: &str, size: u64) -> Result<bool> {
29    #[cfg(target_os = "linux")]
30    if sparse_file_has_no_data(_file)? {
31        return Ok(true);
32    }
33
34    let mut contents = std::fs::OpenOptions::new()
35        .read(true)
36        .open(path)
37        .map_err(FeoxError::IoError)?;
38    let mut buffer = vec![0; ZERO_SCAN_BLOCKS * FEOX_BLOCK_SIZE];
39    let mut remaining = size;
40    while remaining > 0 {
41        let read_len = remaining.min(buffer.len() as u64) as usize;
42        contents
43            .read_exact(&mut buffer[..read_len])
44            .map_err(FeoxError::IoError)?;
45        if buffer[..read_len].iter().any(|byte| *byte != 0) {
46            return Ok(false);
47        }
48        remaining -= read_len as u64;
49    }
50    Ok(true)
51}
52
53#[cfg(target_os = "linux")]
54fn sparse_file_has_no_data(file: &std::fs::File) -> Result<bool> {
55    use std::os::fd::AsRawFd;
56
57    let offset = unsafe { libc::lseek(file.as_raw_fd(), 0, libc::SEEK_DATA) };
58    if offset >= 0 {
59        return Ok(false);
60    }
61
62    let error = std::io::Error::last_os_error();
63    match error.raw_os_error() {
64        Some(libc::ENXIO) => Ok(true),
65        Some(libc::EINVAL) => Ok(false),
66        _ => Err(FeoxError::IoError(error)),
67    }
68}
69
70impl 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    pub fn flush_all(&self) -> Result<()> {
88        if self.initialized && !self.memory_only {
89            // First flush the write buffer to ensure all data is written
90            if let Some(ref wb) = self.write_buffer {
91                wb.force_flush()?;
92            }
93
94            if let Some(ref disk_io) = self.disk_io {
95                // Update metadata with current stats
96                let mut metadata = self._metadata.write();
97                metadata.total_records = self.stats.record_count.load(Ordering::Relaxed) as u64;
98                metadata.total_size = self.stats.disk_usage.load(Ordering::Relaxed);
99                metadata.fragmentation = self.free_space.read().get_fragmentation();
100                metadata.update();
101
102                // Write metadata
103                disk_io.write().write_store_metadata(&mut metadata)?;
104            }
105        }
106        Ok(())
107    }
108
109    pub(super) fn load_value_from_disk(&self, record: &Arc<Record>) -> Result<Bytes> {
110        let mut source = Arc::clone(record);
111        loop {
112            if let Some(value) = source.get_value() {
113                return Ok(value);
114            }
115            if source.sector.load(Ordering::Acquire) != 0 {
116                break;
117            }
118            source = source.value_source().ok_or(FeoxError::StaleExtent)?;
119        }
120        let extent = source.acquire_extent().ok_or(FeoxError::StaleExtent)?;
121        let sector = source.sector.load(Ordering::Acquire);
122        if self.memory_only || sector == 0 {
123            return Err(FeoxError::StaleExtent);
124        }
125        crate::test_hooks::pause_at(crate::test_hooks::AFTER_SECTOR_LOAD);
126
127        // Get the appropriate format handler
128        let format = get_format_ref(self.format_version);
129
130        // Calculate how many sectors we need to read
131        let total_size = format.total_size(source.key.len(), source.value_len);
132        let sectors_needed = total_size.div_ceil(FEOX_BLOCK_SIZE);
133
134        // Read the sectors
135        let disk_io = self
136            .disk_io
137            .as_ref()
138            .ok_or_else(|| {
139                FeoxError::IoError(io::Error::new(
140                    io::ErrorKind::NotFound,
141                    "No disk IO available",
142                ))
143            })?
144            .read();
145
146        let data = disk_io.read_sectors_sync(sector, sectors_needed as u64)?;
147        drop(extent);
148
149        if !sector_holds_record(&data, &source) {
150            return Err(FeoxError::StaleExtent);
151        }
152
153        let offset = format.value_offset(source.key.len());
154        let end = offset
155            .checked_add(source.value_len)
156            .filter(|end| *end <= data.len())
157            .ok_or(FeoxError::InvalidRecord)?;
158        if source.value_len <= data.len() / 2 {
159            Ok(Bytes::copy_from_slice(&data[offset..end]))
160        } else {
161            Ok(Bytes::from(data).slice(offset..end))
162        }
163    }
164
165    pub(super) fn open_device(
166        &mut self,
167        device_path: &Option<String>,
168        file_size: Option<u64>,
169    ) -> Result<()> {
170        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            let (file, use_direct_io) = if std::path::Path::new("/.dockerenv").exists() {
177                let file = OpenOptions::new()
178                    .read(true)
179                    .write(true)
180                    .create(true)
181                    .truncate(false)
182                    .open(path)
183                    .map_err(FeoxError::IoError)?;
184                (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                    match OpenOptions::new()
191                        .read(true)
192                        .write(true)
193                        .create(true)
194                        .truncate(false)
195                        .custom_flags(libc::O_DIRECT)
196                        .open(path)
197                    {
198                        Ok(file) => (file, true), // Successfully opened with O_DIRECT
199                        Err(_) => {
200                            // Fallback to regular open
201                            let file = OpenOptions::new()
202                                .read(true)
203                                .write(true)
204                                .create(true)
205                                .truncate(false)
206                                .open(path)
207                                .map_err(FeoxError::IoError)?;
208                            (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            let metadata = file.metadata().map_err(FeoxError::IoError)?;
236            self.device_size = metadata.len();
237
238            // Track whether this is a newly created file
239            self.fresh_device = self.device_size == 0;
240
241            if self.fresh_device {
242                self.initialize_fresh_device(&file, file_size)?;
243            } else {
244                validate_device_size(self.device_size)?;
245                let is_empty_file = file_is_all_zero(&file, path, self.device_size)?;
246
247                if is_empty_file {
248                    self.free_space.write().initialize(self.device_size)?;
249                    self.fresh_device = true;
250                    let mut metadata = self._metadata.write();
251                    metadata.device_size = self.device_size;
252                    metadata.update();
253                } else {
254                    self.free_space.write().set_device_size(self.device_size);
255                }
256            }
257
258            #[cfg(not(unix))]
259            let use_direct_io = false;
260            self.attach_device_file(file, use_direct_io)?;
261        }
262        Ok(())
263    }
264
265    pub(super) fn open_device_read_only(&mut self, file: File) -> Result<()> {
266        self.device_size = file.metadata().map_err(FeoxError::IoError)?.len();
267        validate_device_size(self.device_size)?;
268        self.free_space.write().set_device_size(self.device_size);
269        self.attach_device_file(file, false)
270    }
271
272    pub(super) fn open_fresh_device(&mut self, file: File, file_size: Option<u64>) -> Result<()> {
273        if file.metadata().map_err(FeoxError::IoError)?.len() != 0 {
274            return Err(FeoxError::InvalidDevice);
275        }
276        self.fresh_device = true;
277        self.initialize_fresh_device(&file, file_size)?;
278        self.attach_device_file(file, false)
279    }
280
281    fn initialize_fresh_device(&mut self, file: &File, file_size: Option<u64>) -> Result<()> {
282        let target_size = file_size.unwrap_or(DEFAULT_DEVICE_SIZE);
283        validate_device_size(target_size)?;
284        file.set_len(target_size).map_err(FeoxError::IoError)?;
285        self.device_size = target_size;
286        self.free_space.write().initialize(self.device_size)?;
287
288        let mut metadata = self._metadata.write();
289        metadata.device_size = self.device_size;
290        metadata.update();
291        Ok(())
292    }
293
294    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            let file = Arc::new(file);
299            self.device_fd = Some(file.as_raw_fd());
300            self.device_file = Some(file.as_ref().try_clone().map_err(FeoxError::IoError)?);
301            let disk_io = crate::storage::io::DiskIO::new(file, use_direct_io)?;
302            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        Ok(())
314    }
315}
316
317impl Drop for FeoxStore {
318    fn drop(&mut self) {
319        // Stop TTL sweeper if running
320        if let Some(mut sweeper) = self.ttl_sweeper.write().take() {
321            sweeper.stop();
322        }
323
324        // Signal shutdown to write buffer workers
325        if let Some(ref wb) = self.write_buffer {
326            wb.initiate_shutdown();
327        }
328
329        if let Some(write_buffer) = self.write_buffer.take() {
330            write_buffer.finish_shutdown();
331        }
332
333        // Write metadata directly without using the write buffer
334        if self.initialized && !self.memory_only && !self.read_only {
335            if let Some(ref disk_io) = self.disk_io {
336                // Update metadata with current stats
337                let mut metadata = self._metadata.write();
338                metadata.total_records = self.stats.record_count.load(Ordering::Relaxed) as u64;
339                metadata.total_size = self.stats.disk_usage.load(Ordering::Relaxed);
340                metadata.fragmentation = self.free_space.read().get_fragmentation();
341                metadata.update();
342
343                // Write metadata
344                let _ = disk_io.write().write_store_metadata(&mut metadata);
345            }
346        }
347
348        // Now it's safe to shutdown disk I/O since workers have exited
349        if let Some(ref disk_io) = self.disk_io {
350            disk_io.write().shutdown();
351        }
352    }
353}