Skip to main content

hermes_core/directories/
directory.rs

1//! Async Directory abstraction for IO operations
2//!
3//! Supports network, local filesystem, and in-memory storage.
4//! All reads are async to minimize blocking on network latency.
5
6use async_trait::async_trait;
7use parking_lot::RwLock;
8use std::collections::HashMap;
9use std::io;
10use std::ops::Range;
11use std::path::{Path, PathBuf};
12use std::sync::Arc;
13
14/// Callback type for lazy range reading
15#[cfg(not(target_arch = "wasm32"))]
16pub type RangeReadFn = Arc<
17    dyn Fn(
18            Range<u64>,
19        )
20            -> std::pin::Pin<Box<dyn std::future::Future<Output = io::Result<OwnedBytes>> + Send>>
21        + Send
22        + Sync,
23>;
24
25#[cfg(target_arch = "wasm32")]
26pub type RangeReadFn = Arc<
27    dyn Fn(
28        Range<u64>,
29    ) -> std::pin::Pin<Box<dyn std::future::Future<Output = io::Result<OwnedBytes>>>>,
30>;
31
32/// Unified file handle for both inline (mmap/RAM) and lazy (HTTP/filesystem) access.
33///
34/// Replaces the previous `FileSlice`, `LazyFileHandle`, and `LazyFileSlice` types.
35/// - **Inline**: data is available synchronously (mmap, RAM). Sync reads via `read_bytes_range_sync`.
36/// - **Lazy**: data is fetched on-demand via async callback (HTTP, filesystem).
37///
38/// Use `.slice()` to create sub-range views (zero-copy for Inline, offset-adjusted for Lazy).
39#[derive(Clone)]
40pub struct FileHandle {
41    inner: FileHandleInner,
42}
43
44#[derive(Clone)]
45enum FileHandleInner {
46    /// Data available inline — sync reads possible (mmap, RAM)
47    Inline {
48        data: OwnedBytes,
49        offset: u64,
50        len: u64,
51    },
52    /// Data fetched on-demand via async callback (HTTP, filesystem)
53    Lazy {
54        read_fn: RangeReadFn,
55        offset: u64,
56        len: u64,
57        /// Index name for the `hermes_directory_read_*` metric labels.
58        label: Arc<str>,
59    },
60}
61
62/// Late-bound index name for Directory-layer metric labels
63/// (`hermes_directory_read_*`, `hermes_cold_write_bytes_total`).
64///
65/// Directories are constructed before the schema is loaded, so the label is
66/// attached afterwards: `Index::open`/`create` call
67/// `Directory::set_index_label(schema.index_label())` on the index's
68/// directory instance. Reads happen at handle/writer creation, not per IO.
69#[derive(Clone, Debug)]
70pub struct IndexLabel(Arc<std::sync::RwLock<Arc<str>>>);
71
72impl Default for IndexLabel {
73    fn default() -> Self {
74        Self(Arc::new(std::sync::RwLock::new(Arc::from("unknown"))))
75    }
76}
77
78impl IndexLabel {
79    /// Current label ("unknown" until set).
80    pub fn get(&self) -> Arc<str> {
81        self.0.read().expect("IndexLabel lock poisoned").clone()
82    }
83
84    /// Set the label (idempotent; last write wins).
85    pub fn set(&self, label: &str) {
86        *self.0.write().expect("IndexLabel lock poisoned") = Arc::from(label);
87    }
88}
89
90impl std::fmt::Debug for FileHandle {
91    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
92        match &self.inner {
93            FileHandleInner::Inline { len, offset, .. } => f
94                .debug_struct("FileHandle::Inline")
95                .field("offset", offset)
96                .field("len", len)
97                .finish(),
98            FileHandleInner::Lazy { len, offset, .. } => f
99                .debug_struct("FileHandle::Lazy")
100                .field("offset", offset)
101                .field("len", len)
102                .finish(),
103        }
104    }
105}
106
107impl FileHandle {
108    /// Create an inline file handle from owned bytes (mmap, RAM).
109    /// Sync reads are available.
110    pub fn from_bytes(data: OwnedBytes) -> Self {
111        let len = data.len() as u64;
112        Self {
113            inner: FileHandleInner::Inline {
114                data,
115                offset: 0,
116                len,
117            },
118        }
119    }
120
121    /// Create an empty file handle.
122    pub fn empty() -> Self {
123        Self::from_bytes(OwnedBytes::empty())
124    }
125
126    /// Create a lazy file handle from an async range-read callback.
127    /// Only async reads are available. Reads emit `hermes_directory_read_*`
128    /// with `index="unknown"` — use [`FileHandle::lazy_labeled`] when the
129    /// owning index is known.
130    pub fn lazy(len: u64, read_fn: RangeReadFn) -> Self {
131        Self::lazy_labeled(len, read_fn, Arc::from("unknown"))
132    }
133
134    /// [`FileHandle::lazy`] with an index name for metric labels.
135    pub fn lazy_labeled(len: u64, read_fn: RangeReadFn, label: Arc<str>) -> Self {
136        Self {
137            inner: FileHandleInner::Lazy {
138                read_fn,
139                offset: 0,
140                len,
141                label,
142            },
143        }
144    }
145
146    /// Total length in bytes.
147    #[inline]
148    pub fn len(&self) -> u64 {
149        match &self.inner {
150            FileHandleInner::Inline { len, .. } => *len,
151            FileHandleInner::Lazy { len, .. } => *len,
152        }
153    }
154
155    /// Check if empty.
156    #[inline]
157    pub fn is_empty(&self) -> bool {
158        self.len() == 0
159    }
160
161    /// Whether synchronous reads are available (inline/mmap data).
162    #[inline]
163    pub fn is_sync(&self) -> bool {
164        matches!(&self.inner, FileHandleInner::Inline { .. })
165    }
166
167    /// Create a sub-range view. Zero-copy for Inline, offset-adjusted for Lazy.
168    pub fn slice(&self, range: Range<u64>) -> Self {
169        match &self.inner {
170            FileHandleInner::Inline { data, offset, len } => {
171                let new_offset = offset + range.start;
172                let new_len = range.end - range.start;
173                debug_assert!(
174                    new_offset + new_len <= offset + len,
175                    "slice out of bounds: {}+{} > {}+{}",
176                    new_offset,
177                    new_len,
178                    offset,
179                    len
180                );
181                Self {
182                    inner: FileHandleInner::Inline {
183                        data: data.clone(),
184                        offset: new_offset,
185                        len: new_len,
186                    },
187                }
188            }
189            FileHandleInner::Lazy {
190                read_fn,
191                offset,
192                len,
193                label,
194            } => {
195                let new_offset = offset + range.start;
196                let new_len = range.end - range.start;
197                debug_assert!(
198                    new_offset + new_len <= offset + len,
199                    "slice out of bounds: {}+{} > {}+{}",
200                    new_offset,
201                    new_len,
202                    offset,
203                    len
204                );
205                Self {
206                    inner: FileHandleInner::Lazy {
207                        read_fn: Arc::clone(read_fn),
208                        offset: new_offset,
209                        len: new_len,
210                        label: Arc::clone(label),
211                    },
212                }
213            }
214        }
215    }
216
217    /// Advise the kernel about the access pattern for a byte range of this handle.
218    ///
219    /// Only effective for Inline handles backed by mmap; no-op for Lazy
220    /// handles (HTTP, filesystem callbacks) and heap-backed data.
221    #[cfg(feature = "native")]
222    pub fn madvise_range(&self, range: Range<u64>, advice: libc::c_int) {
223        if let FileHandleInner::Inline { data, offset, len } = &self.inner {
224            let end = range.end.min(*len);
225            if range.start >= end {
226                return;
227            }
228            let start = (*offset + range.start) as usize;
229            let end = (*offset + end) as usize;
230            data.madvise_range(start..end, advice);
231        }
232    }
233
234    /// Async range read — works for both Inline and Lazy.
235    pub async fn read_bytes_range(&self, range: Range<u64>) -> io::Result<OwnedBytes> {
236        match &self.inner {
237            FileHandleInner::Inline { data, offset, len } => {
238                if range.end > *len {
239                    return Err(io::Error::new(
240                        io::ErrorKind::InvalidInput,
241                        format!("Range {:?} out of bounds (len: {})", range, len),
242                    ));
243                }
244                let start = (*offset + range.start) as usize;
245                let end = (*offset + range.end) as usize;
246                Ok(data.slice(start..end))
247            }
248            FileHandleInner::Lazy {
249                read_fn,
250                offset,
251                len,
252                label,
253            } => {
254                if range.end > *len {
255                    return Err(io::Error::new(
256                        io::ErrorKind::InvalidInput,
257                        format!("Range {:?} out of bounds (len: {})", range, len),
258                    ));
259                }
260                let abs_start = offset + range.start;
261                let abs_end = offset + range.end;
262                // Real IO (HTTP / custom read_fn) — mmap-backed Inline handles
263                // above are zero-copy slices whose latency materializes as
264                // page faults inside the query-phase histograms instead.
265                let t = crate::observe::Timer::start();
266                let result = (read_fn)(abs_start..abs_end).await;
267                if let Ok(bytes) = &result {
268                    crate::observe::directory_read(label, "lazy_range", t.secs(), bytes.len());
269                }
270                result
271            }
272        }
273    }
274
275    /// Read all bytes.
276    pub async fn read_bytes(&self) -> io::Result<OwnedBytes> {
277        self.read_bytes_range(0..self.len()).await
278    }
279
280    /// Synchronous range read — only works for Inline handles.
281    /// Returns `Err` if the handle is Lazy.
282    #[inline]
283    pub fn read_bytes_range_sync(&self, range: Range<u64>) -> io::Result<OwnedBytes> {
284        match &self.inner {
285            FileHandleInner::Inline { data, offset, len } => {
286                if range.end > *len {
287                    return Err(io::Error::new(
288                        io::ErrorKind::InvalidInput,
289                        format!("Range {:?} out of bounds (len: {})", range, len),
290                    ));
291                }
292                let start = (*offset + range.start) as usize;
293                let end = (*offset + range.end) as usize;
294                Ok(data.slice(start..end))
295            }
296            FileHandleInner::Lazy { .. } => Err(io::Error::new(
297                io::ErrorKind::Unsupported,
298                "Synchronous read not available on lazy file handle",
299            )),
300        }
301    }
302
303    /// Synchronous read of all bytes — only works for Inline handles.
304    #[inline]
305    pub fn read_bytes_sync(&self) -> io::Result<OwnedBytes> {
306        self.read_bytes_range_sync(0..self.len())
307    }
308}
309
310/// Backing store for OwnedBytes — supports both heap Vec and mmap.
311#[derive(Clone)]
312enum SharedBytes {
313    Vec(Arc<Vec<u8>>),
314    #[cfg(feature = "native")]
315    Mmap(Arc<memmap2::Mmap>),
316}
317
318impl SharedBytes {
319    #[inline]
320    fn as_bytes(&self) -> &[u8] {
321        match self {
322            SharedBytes::Vec(v) => v.as_slice(),
323            #[cfg(feature = "native")]
324            SharedBytes::Mmap(m) => m.as_ref(),
325        }
326    }
327}
328
329impl std::fmt::Debug for SharedBytes {
330    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
331        match self {
332            SharedBytes::Vec(v) => write!(f, "Vec(len={})", v.len()),
333            #[cfg(feature = "native")]
334            SharedBytes::Mmap(m) => write!(f, "Mmap(len={})", m.len()),
335        }
336    }
337}
338
339/// Owned bytes with cheap cloning (Arc-backed)
340///
341/// Supports two backing stores:
342/// - `Vec<u8>` for owned data (RamDirectory, FsDirectory, decompressed blocks)
343/// - `Mmap` for zero-copy memory-mapped files (MmapDirectory, native only)
344#[derive(Debug, Clone)]
345pub struct OwnedBytes {
346    data: SharedBytes,
347    range: Range<usize>,
348}
349
350impl OwnedBytes {
351    pub fn new(data: Vec<u8>) -> Self {
352        let len = data.len();
353        Self {
354            data: SharedBytes::Vec(Arc::new(data)),
355            range: 0..len,
356        }
357    }
358
359    pub fn empty() -> Self {
360        Self {
361            data: SharedBytes::Vec(Arc::new(Vec::new())),
362            range: 0..0,
363        }
364    }
365
366    /// Create from a pre-existing Arc<Vec<u8>> with a sub-range.
367    /// Used by RamDirectory and CachingDirectory to share data without copying.
368    pub(crate) fn from_arc_vec(data: Arc<Vec<u8>>, range: Range<usize>) -> Self {
369        Self {
370            data: SharedBytes::Vec(data),
371            range,
372        }
373    }
374
375    /// Create from a memory-mapped file (zero-copy).
376    #[cfg(feature = "native")]
377    pub(crate) fn from_mmap(mmap: Arc<memmap2::Mmap>) -> Self {
378        let len = mmap.len();
379        Self {
380            data: SharedBytes::Mmap(mmap),
381            range: 0..len,
382        }
383    }
384
385    /// Create from a memory-mapped file with a sub-range (zero-copy).
386    #[cfg(feature = "native")]
387    pub(crate) fn from_mmap_range(mmap: Arc<memmap2::Mmap>, range: Range<usize>) -> Self {
388        Self {
389            data: SharedBytes::Mmap(mmap),
390            range,
391        }
392    }
393
394    pub fn len(&self) -> usize {
395        self.range.len()
396    }
397
398    pub fn is_empty(&self) -> bool {
399        self.range.is_empty()
400    }
401
402    pub fn slice(&self, range: Range<usize>) -> Self {
403        let start = self.range.start + range.start;
404        let end = self.range.start + range.end;
405        Self {
406            data: self.data.clone(),
407            range: start..end,
408        }
409    }
410
411    pub fn as_slice(&self) -> &[u8] {
412        &self.data.as_bytes()[self.range.clone()]
413    }
414
415    /// Returns `true` if the backing store is a memory-mapped file.
416    ///
417    /// Used to guard `madvise` calls: `MADV_DONTNEED` on heap memory
418    /// zeroes pages on Linux and corrupts allocator metadata.
419    #[cfg(feature = "native")]
420    #[inline]
421    pub fn is_mmap(&self) -> bool {
422        matches!(self.data, SharedBytes::Mmap(_))
423    }
424
425    /// Advise the kernel about the access pattern for these bytes.
426    ///
427    /// No-op unless the backing store is mmap (heap memory must never be
428    /// madvised: `MADV_DONTNEED` on heap zeroes pages and corrupts allocator
429    /// metadata) or the range is empty.
430    #[cfg(feature = "native")]
431    pub fn madvise(&self, advice: libc::c_int) {
432        self.madvise_range(0..self.len(), advice);
433    }
434
435    /// Pin these bytes in physical memory (`mlock`). mmap-backed only —
436    /// heap memory is not evictable by the page cache. Returns whether the
437    /// lock succeeded; failure (e.g. RLIMIT_MEMLOCK) is not fatal.
438    /// Locks are released automatically when the mapping is unmapped.
439    #[cfg(feature = "native")]
440    pub fn mlock(&self) -> bool {
441        if !self.is_mmap() {
442            return false;
443        }
444        let slice = self.as_slice();
445        if slice.is_empty() {
446            return true;
447        }
448        let ptr = slice.as_ptr();
449        let len = slice.len();
450        let page_size = 4096usize;
451        let aligned_ptr = (ptr as usize) & !(page_size - 1);
452        let aligned_len = len + (ptr as usize - aligned_ptr);
453        unsafe { libc::mlock(aligned_ptr as *const libc::c_void, aligned_len) == 0 }
454    }
455
456    /// Advise the kernel about the access pattern for a sub-range.
457    ///
458    /// The range is relative to these bytes. Same mmap-only guard as
459    /// [`Self::madvise`]. The pointer is aligned down to a page boundary
460    /// as required by `madvise`.
461    #[cfg(feature = "native")]
462    pub fn madvise_range(&self, range: Range<usize>, advice: libc::c_int) {
463        if !self.is_mmap() {
464            return;
465        }
466        let slice = &self.as_slice()[range];
467        if slice.is_empty() {
468            return;
469        }
470        let ptr = slice.as_ptr();
471        let len = slice.len();
472        let page_size = 4096usize;
473        let aligned_ptr = (ptr as usize) & !(page_size - 1);
474        let aligned_len = len + (ptr as usize - aligned_ptr);
475        unsafe {
476            libc::madvise(aligned_ptr as *mut libc::c_void, aligned_len, advice);
477        }
478    }
479
480    pub fn to_vec(&self) -> Vec<u8> {
481        self.as_slice().to_vec()
482    }
483}
484
485impl AsRef<[u8]> for OwnedBytes {
486    fn as_ref(&self) -> &[u8] {
487        self.as_slice()
488    }
489}
490
491impl std::ops::Deref for OwnedBytes {
492    type Target = [u8];
493
494    fn deref(&self) -> &Self::Target {
495        self.as_slice()
496    }
497}
498
499/// Async directory trait for reading index files
500#[cfg(not(target_arch = "wasm32"))]
501#[async_trait]
502pub trait Directory: Send + Sync + 'static {
503    /// Check if a file exists
504    async fn exists(&self, path: &Path) -> io::Result<bool>;
505
506    /// Get file size
507    async fn file_size(&self, path: &Path) -> io::Result<u64>;
508
509    /// Open a file for reading (loads entire file into an inline FileHandle)
510    async fn open_read(&self, path: &Path) -> io::Result<FileHandle>;
511
512    /// Read a specific byte range from a file (optimized for network)
513    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes>;
514
515    /// List files in directory
516    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>>;
517
518    /// Open a file handle that fetches ranges on demand.
519    /// For mmap directories this returns an Inline handle (sync-capable).
520    /// For HTTP/filesystem directories this returns a Lazy handle.
521    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle>;
522
523    /// Attach the owning index's name for Directory-layer metric labels
524    /// (`hermes_directory_read_*`, `hermes_cold_write_bytes_total`).
525    /// Called by `Index::open`/`create` once the schema is loaded; wrappers
526    /// forward to their inner directory. Default: no-op (directories that
527    /// emit no Directory-layer metrics, e.g. RamDirectory).
528    fn set_index_label(&self, _label: &str) {}
529}
530
531/// Async directory trait for reading index files (WASM version - no Send requirement)
532#[cfg(target_arch = "wasm32")]
533#[async_trait(?Send)]
534pub trait Directory: 'static {
535    /// Check if a file exists
536    async fn exists(&self, path: &Path) -> io::Result<bool>;
537
538    /// Get file size
539    async fn file_size(&self, path: &Path) -> io::Result<u64>;
540
541    /// Open a file for reading (loads entire file into an inline FileHandle)
542    async fn open_read(&self, path: &Path) -> io::Result<FileHandle>;
543
544    /// Read a specific byte range from a file (optimized for network)
545    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes>;
546
547    /// List files in directory
548    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>>;
549
550    /// Open a file handle that fetches ranges on demand.
551    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle>;
552
553    /// Attach the owning index's name for Directory-layer metric labels.
554    /// No-op default; metrics are native-only but the label is harmless.
555    fn set_index_label(&self, _label: &str) {}
556}
557
558/// A writer for incrementally writing data to a directory file.
559///
560/// Avoids buffering entire files in memory during merge. File-backed
561/// directories write directly to disk; memory directories collect to Vec.
562pub trait StreamingWriter: io::Write + Send {
563    /// Finalize the write, making data available for reading.
564    fn finish(self: Box<Self>) -> io::Result<()>;
565
566    /// Bytes written so far.
567    fn bytes_written(&self) -> u64;
568}
569
570/// StreamingWriter backed by Vec<u8>, finalized via DirectoryWriter::write.
571/// Used as default/fallback and for RamDirectory.
572struct BufferedStreamingWriter {
573    path: PathBuf,
574    buffer: Vec<u8>,
575    /// Callback to write the buffer to the directory on finish.
576    /// We store the files Arc directly for RamDirectory.
577    files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
578}
579
580impl io::Write for BufferedStreamingWriter {
581    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
582        self.buffer.extend_from_slice(buf);
583        Ok(buf.len())
584    }
585
586    fn flush(&mut self) -> io::Result<()> {
587        Ok(())
588    }
589}
590
591impl StreamingWriter for BufferedStreamingWriter {
592    fn finish(self: Box<Self>) -> io::Result<()> {
593        self.files.write().insert(self.path, Arc::new(self.buffer));
594        Ok(())
595    }
596
597    fn bytes_written(&self) -> u64 {
598        self.buffer.len() as u64
599    }
600}
601
602/// Buffer size for FileStreamingWriter (8 MB).
603/// Large enough to coalesce millions of tiny writes (e.g. per-vector doc_id writes)
604/// into efficient sequential I/O.
605#[cfg(feature = "native")]
606const FILE_STREAMING_BUF_SIZE: usize = 8 * 1024 * 1024;
607
608/// StreamingWriter backed by a buffered std::fs::File for filesystem directories.
609#[cfg(feature = "native")]
610pub(crate) struct FileStreamingWriter {
611    pub(crate) file: io::BufWriter<std::fs::File>,
612    pub(crate) written: u64,
613}
614
615#[cfg(feature = "native")]
616impl FileStreamingWriter {
617    pub(crate) fn new(file: std::fs::File) -> Self {
618        Self {
619            file: io::BufWriter::with_capacity(FILE_STREAMING_BUF_SIZE, file),
620            written: 0,
621        }
622    }
623}
624
625#[cfg(feature = "native")]
626impl io::Write for FileStreamingWriter {
627    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
628        let n = self.file.write(buf)?;
629        self.written += n as u64;
630        Ok(n)
631    }
632
633    fn flush(&mut self) -> io::Result<()> {
634        self.file.flush()
635    }
636}
637
638#[cfg(feature = "native")]
639impl StreamingWriter for FileStreamingWriter {
640    fn finish(self: Box<Self>) -> io::Result<()> {
641        let file = self.file.into_inner().map_err(|e| e.into_error())?;
642        file.sync_all()?;
643        Ok(())
644    }
645
646    fn bytes_written(&self) -> u64 {
647        self.written
648    }
649}
650
651/// Async directory trait for writing index files
652#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
653#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
654pub trait DirectoryWriter: Directory {
655    /// Create/overwrite a file with data
656    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()>;
657
658    /// Create/overwrite a file with data, durably.
659    ///
660    /// [`Self::write`] does not guarantee the bytes reach stable storage
661    /// before returning (filesystem implementations leave them in the OS
662    /// page cache). Any file that durably-published metadata will reference
663    /// (e.g. segment `.meta`) must be written through this method instead:
664    /// it routes through [`Self::streaming_writer`], whose `finish()` fsyncs
665    /// on filesystem implementations.
666    async fn write_durable(&self, path: &Path, data: &[u8]) -> io::Result<()> {
667        use io::Write as _;
668        let mut writer = self.streaming_writer(path).await?;
669        writer.write_all(data)?;
670        writer.finish()
671    }
672
673    /// Delete a file
674    async fn delete(&self, path: &Path) -> io::Result<()>;
675
676    /// Atomic rename
677    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()>;
678
679    /// Create another immutable name for an existing file without copying its
680    /// contents when the backend supports it. Segment rewrites use this to
681    /// retain unchanged multi-gigabyte files while replacing only one index
682    /// payload. Backends without link semantics return `Unsupported`; callers
683    /// then fall back to a streaming copy.
684    async fn link(&self, _from: &Path, _to: &Path) -> io::Result<()> {
685        Err(io::Error::new(
686            io::ErrorKind::Unsupported,
687            "directory backend does not support immutable file links",
688        ))
689    }
690
691    /// Sync all pending writes
692    async fn sync(&self) -> io::Result<()>;
693
694    /// Create a streaming writer for incremental file writes.
695    /// Call finish() on the returned writer to finalize.
696    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>>;
697
698    /// Streaming writer for **bulk one-shot data** (merge/reorder outputs).
699    ///
700    /// Filesystem directories return a page-cache-dropping writer (see
701    /// `docs/cold-io.md`) so multi-GB merge writes cannot evict the serving
702    /// segments' warm pages. Output is byte-identical to the buffered
703    /// writer. Default impl delegates to [`Self::streaming_writer`].
704    async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
705        self.streaming_writer(path).await
706    }
707}
708
709/// In-memory directory for testing and small indexes
710#[derive(Debug, Default)]
711pub struct RamDirectory {
712    files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
713}
714
715impl Clone for RamDirectory {
716    fn clone(&self) -> Self {
717        Self {
718            files: Arc::clone(&self.files),
719        }
720    }
721}
722
723impl RamDirectory {
724    pub fn new() -> Self {
725        Self::default()
726    }
727
728    /// Synchronous file listing (for serialization).
729    pub fn list_files_sync(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
730        let files = self.files.read();
731        Ok(files
732            .keys()
733            .filter(|p| p.starts_with(prefix))
734            .cloned()
735            .collect())
736    }
737
738    /// Synchronous file read (for serialization).
739    pub fn read_file_sync(&self, path: &Path) -> io::Result<Vec<u8>> {
740        let files = self.files.read();
741        files
742            .get(path)
743            .map(|data| data.as_ref().clone())
744            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
745    }
746
747    /// Synchronous file write (for deserialization).
748    pub fn write_sync(&self, path: &Path, data: &[u8]) -> io::Result<()> {
749        self.files
750            .write()
751            .insert(path.to_path_buf(), Arc::new(data.to_vec()));
752        Ok(())
753    }
754}
755
756#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
757#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
758impl Directory for RamDirectory {
759    async fn exists(&self, path: &Path) -> io::Result<bool> {
760        Ok(self.files.read().contains_key(path))
761    }
762
763    async fn file_size(&self, path: &Path) -> io::Result<u64> {
764        self.files
765            .read()
766            .get(path)
767            .map(|data| data.len() as u64)
768            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
769    }
770
771    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
772        let files = self.files.read();
773        let data = files
774            .get(path)
775            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
776
777        Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
778            Arc::clone(data),
779            0..data.len(),
780        )))
781    }
782
783    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
784        let files = self.files.read();
785        let data = files
786            .get(path)
787            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
788
789        let start = range.start as usize;
790        let end = range.end as usize;
791
792        if end > data.len() {
793            return Err(io::Error::new(
794                io::ErrorKind::InvalidInput,
795                "Range out of bounds",
796            ));
797        }
798
799        Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end))
800    }
801
802    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
803        let files = self.files.read();
804        Ok(files
805            .keys()
806            .filter(|p| p.starts_with(prefix))
807            .cloned()
808            .collect())
809    }
810
811    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
812        // RAM data is always available synchronously — return Inline handle
813        self.open_read(path).await
814    }
815}
816
817#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
818#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
819impl DirectoryWriter for RamDirectory {
820    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
821        self.files
822            .write()
823            .insert(path.to_path_buf(), Arc::new(data.to_vec()));
824        Ok(())
825    }
826
827    async fn delete(&self, path: &Path) -> io::Result<()> {
828        self.files.write().remove(path);
829        Ok(())
830    }
831
832    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
833        let mut files = self.files.write();
834        if let Some(data) = files.remove(from) {
835            files.insert(to.to_path_buf(), data);
836        }
837        Ok(())
838    }
839
840    async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
841        let mut files = self.files.write();
842        let data = files.get(from).cloned().ok_or_else(|| {
843            io::Error::new(
844                io::ErrorKind::NotFound,
845                format!("source file {from:?} does not exist"),
846            )
847        })?;
848        files.insert(to.to_path_buf(), data);
849        Ok(())
850    }
851
852    async fn sync(&self) -> io::Result<()> {
853        Ok(())
854    }
855
856    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
857        Ok(Box::new(BufferedStreamingWriter {
858            path: path.to_path_buf(),
859            buffer: Vec::new(),
860            files: Arc::clone(&self.files),
861        }))
862    }
863}
864
865/// Local filesystem directory with async IO via tokio
866#[cfg(feature = "native")]
867#[derive(Debug, Clone)]
868pub struct FsDirectory {
869    root: PathBuf,
870    label: IndexLabel,
871}
872
873#[cfg(feature = "native")]
874impl FsDirectory {
875    pub fn new(root: impl AsRef<Path>) -> Self {
876        Self {
877            root: root.as_ref().to_path_buf(),
878            label: IndexLabel::default(),
879        }
880    }
881
882    fn resolve(&self, path: &Path) -> PathBuf {
883        self.root.join(path)
884    }
885}
886
887#[cfg(feature = "native")]
888#[async_trait]
889impl Directory for FsDirectory {
890    async fn exists(&self, path: &Path) -> io::Result<bool> {
891        let full_path = self.resolve(path);
892        // `try_exists` maps NotFound to Ok(false); any other stat failure
893        // (EACCES, EIO, ...) must propagate so callers can distinguish a
894        // genuinely missing file from a transient IO error — swallowing it
895        // as `false` quarantines a healthy segment as "missing mandatory
896        // files" instead of retrying.
897        tokio::fs::try_exists(&full_path).await
898    }
899
900    async fn file_size(&self, path: &Path) -> io::Result<u64> {
901        let full_path = self.resolve(path);
902        let metadata = tokio::fs::metadata(&full_path).await?;
903        Ok(metadata.len())
904    }
905
906    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
907        let full_path = self.resolve(path);
908        let data = tokio::fs::read(&full_path).await?;
909        Ok(FileHandle::from_bytes(OwnedBytes::new(data)))
910    }
911
912    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
913        use tokio::io::{AsyncReadExt, AsyncSeekExt};
914
915        let full_path = self.resolve(path);
916        let mut file = tokio::fs::File::open(&full_path).await?;
917
918        file.seek(std::io::SeekFrom::Start(range.start)).await?;
919
920        let len = (range.end - range.start) as usize;
921        let mut buffer = vec![0u8; len];
922        file.read_exact(&mut buffer).await?;
923
924        Ok(OwnedBytes::new(buffer))
925    }
926
927    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
928        let full_path = self.resolve(prefix);
929        let mut entries = tokio::fs::read_dir(&full_path).await?;
930        let mut files = Vec::new();
931
932        while let Some(entry) = entries.next_entry().await? {
933            if entry.file_type().await?.is_file() {
934                files.push(entry.path().strip_prefix(&self.root).unwrap().to_path_buf());
935            }
936        }
937
938        Ok(files)
939    }
940
941    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
942        let full_path = self.resolve(path);
943        let metadata = tokio::fs::metadata(&full_path).await?;
944        let file_size = metadata.len();
945
946        let read_fn: RangeReadFn = Arc::new(move |range: Range<u64>| {
947            let full_path = full_path.clone();
948            Box::pin(async move {
949                use tokio::io::{AsyncReadExt, AsyncSeekExt};
950
951                let mut file = tokio::fs::File::open(&full_path).await?;
952                file.seek(std::io::SeekFrom::Start(range.start)).await?;
953
954                let len = (range.end - range.start) as usize;
955                let mut buffer = vec![0u8; len];
956                file.read_exact(&mut buffer).await?;
957
958                Ok(OwnedBytes::new(buffer))
959            })
960        });
961
962        Ok(FileHandle::lazy_labeled(
963            file_size,
964            read_fn,
965            self.label.get(),
966        ))
967    }
968
969    fn set_index_label(&self, label: &str) {
970        self.label.set(label);
971    }
972}
973
974#[cfg(feature = "native")]
975#[async_trait]
976impl DirectoryWriter for FsDirectory {
977    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
978        let full_path = self.resolve(path);
979
980        // Ensure parent directory exists
981        if let Some(parent) = full_path.parent() {
982            tokio::fs::create_dir_all(parent).await?;
983        }
984
985        tokio::fs::write(&full_path, data).await
986    }
987
988    async fn delete(&self, path: &Path) -> io::Result<()> {
989        let full_path = self.resolve(path);
990        tokio::fs::remove_file(&full_path).await
991    }
992
993    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
994        let from_path = self.resolve(from);
995        let to_path = self.resolve(to);
996        // Metadata publication is the only rename user. Keep the atomic
997        // filesystem operation in a single future poll: tokio::fs::rename is
998        // backed by a cancellable await around spawn_blocking, so a dropped
999        // commit future could observe neither completion nor failure even
1000        // though the rename later succeeded. The caller must update its
1001        // in-memory metadata in the same poll after this returns.
1002        std::fs::rename(&from_path, &to_path)
1003    }
1004
1005    async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
1006        std::fs::hard_link(self.resolve(from), self.resolve(to))
1007    }
1008
1009    async fn sync(&self) -> io::Result<()> {
1010        // fsync the directory
1011        let dir = std::fs::File::open(&self.root)?;
1012        dir.sync_all()?;
1013        Ok(())
1014    }
1015
1016    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1017        let full_path = self.resolve(path);
1018        if let Some(parent) = full_path.parent() {
1019            tokio::fs::create_dir_all(parent).await?;
1020        }
1021        let file = std::fs::File::create(&full_path)?;
1022        Ok(Box::new(FileStreamingWriter::new(file)))
1023    }
1024
1025    async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1026        let full_path = self.resolve(path);
1027        if let Some(parent) = full_path.parent() {
1028            tokio::fs::create_dir_all(parent).await?;
1029        }
1030        let file = std::fs::File::create(&full_path)?;
1031        Ok(Box::new(super::ColdStreamingWriter::new(
1032            file,
1033            self.label.get(),
1034        )))
1035    }
1036}
1037
1038/// Caching wrapper for any Directory - caches file reads
1039pub struct CachingDirectory<D: Directory> {
1040    inner: D,
1041    cache: RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>,
1042    max_cached_bytes: usize,
1043    current_bytes: RwLock<usize>,
1044}
1045
1046impl<D: Directory> CachingDirectory<D> {
1047    pub fn new(inner: D, max_cached_bytes: usize) -> Self {
1048        Self {
1049            inner,
1050            cache: RwLock::new(HashMap::new()),
1051            max_cached_bytes,
1052            current_bytes: RwLock::new(0),
1053        }
1054    }
1055
1056    fn try_cache(&self, path: &Path, data: &[u8]) {
1057        let mut current = self.current_bytes.write();
1058        if *current + data.len() <= self.max_cached_bytes {
1059            self.cache
1060                .write()
1061                .insert(path.to_path_buf(), Arc::new(data.to_vec()));
1062            *current += data.len();
1063        }
1064    }
1065}
1066
1067#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1068#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1069impl<D: Directory> Directory for CachingDirectory<D> {
1070    async fn exists(&self, path: &Path) -> io::Result<bool> {
1071        if self.cache.read().contains_key(path) {
1072            return Ok(true);
1073        }
1074        self.inner.exists(path).await
1075    }
1076
1077    async fn file_size(&self, path: &Path) -> io::Result<u64> {
1078        if let Some(data) = self.cache.read().get(path) {
1079            return Ok(data.len() as u64);
1080        }
1081        self.inner.file_size(path).await
1082    }
1083
1084    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
1085        // Check cache first
1086        if let Some(data) = self.cache.read().get(path) {
1087            return Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
1088                Arc::clone(data),
1089                0..data.len(),
1090            )));
1091        }
1092
1093        // Read from inner and potentially cache
1094        let handle = self.inner.open_read(path).await?;
1095        let bytes = handle.read_bytes().await?;
1096
1097        self.try_cache(path, bytes.as_slice());
1098
1099        Ok(FileHandle::from_bytes(bytes))
1100    }
1101
1102    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
1103        // Check cache first
1104        if let Some(data) = self.cache.read().get(path) {
1105            let start = range.start as usize;
1106            let end = range.end as usize;
1107            return Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end));
1108        }
1109
1110        self.inner.read_range(path, range).await
1111    }
1112
1113    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
1114        self.inner.list_files(prefix).await
1115    }
1116
1117    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
1118        // For caching directory, delegate to inner - caching happens at read_range level
1119        self.inner.open_lazy(path).await
1120    }
1121
1122    fn set_index_label(&self, label: &str) {
1123        self.inner.set_index_label(label);
1124    }
1125}
1126
1127#[cfg(test)]
1128mod tests {
1129    use super::*;
1130
1131    #[tokio::test]
1132    async fn test_ram_directory() {
1133        let dir = RamDirectory::new();
1134
1135        // Write file
1136        dir.write(Path::new("test.txt"), b"hello world")
1137            .await
1138            .unwrap();
1139
1140        // Check exists
1141        assert!(dir.exists(Path::new("test.txt")).await.unwrap());
1142        assert!(!dir.exists(Path::new("nonexistent.txt")).await.unwrap());
1143
1144        // Read file
1145        let slice = dir.open_read(Path::new("test.txt")).await.unwrap();
1146        let data = slice.read_bytes().await.unwrap();
1147        assert_eq!(data.as_slice(), b"hello world");
1148
1149        // Read range
1150        let range_data = dir.read_range(Path::new("test.txt"), 0..5).await.unwrap();
1151        assert_eq!(range_data.as_slice(), b"hello");
1152
1153        // Delete
1154        dir.delete(Path::new("test.txt")).await.unwrap();
1155        assert!(!dir.exists(Path::new("test.txt")).await.unwrap());
1156    }
1157
1158    /// A transient stat failure (EACCES here, EIO on flaky storage) must
1159    /// surface as `Err`, not `Ok(false)`: callers classify a missing
1160    /// mandatory segment file as deterministic corruption and quarantine
1161    /// the segment until restart.
1162    #[cfg(all(unix, feature = "native"))]
1163    #[tokio::test]
1164    async fn test_fs_exists_propagates_stat_errors_instead_of_reporting_missing() {
1165        use std::os::unix::fs::PermissionsExt;
1166
1167        let temp_dir = tempfile::TempDir::new().unwrap();
1168        let dir = FsDirectory::new(temp_dir.path());
1169        dir.write(Path::new("locked/seg.meta"), b"data")
1170            .await
1171            .unwrap();
1172
1173        // Removing search permission from the parent makes stat on the child
1174        // fail with EACCES while the file itself still exists.
1175        let locked = temp_dir.path().join("locked");
1176        let original = std::fs::metadata(&locked).unwrap().permissions();
1177        std::fs::set_permissions(&locked, std::fs::Permissions::from_mode(0o000)).unwrap();
1178        if std::fs::metadata(locked.join("seg.meta")).is_ok() {
1179            // Running as root: directory permissions are not enforced, so the
1180            // stat failure cannot be provoked.
1181            std::fs::set_permissions(&locked, original).unwrap();
1182            return;
1183        }
1184        let result = dir.exists(Path::new("locked/seg.meta")).await;
1185        std::fs::set_permissions(&locked, original).unwrap();
1186
1187        let error =
1188            result.expect_err("stat failure must propagate as Err, not be misreported as missing");
1189        assert_ne!(error.kind(), io::ErrorKind::NotFound);
1190        // Once stat succeeds again the file is reported present.
1191        assert!(dir.exists(Path::new("locked/seg.meta")).await.unwrap());
1192    }
1193
1194    #[tokio::test]
1195    async fn test_file_handle() {
1196        let data = OwnedBytes::new(b"hello world".to_vec());
1197        let handle = FileHandle::from_bytes(data);
1198
1199        assert_eq!(handle.len(), 11);
1200        assert!(handle.is_sync());
1201
1202        let sub = handle.slice(0..5);
1203        let bytes = sub.read_bytes().await.unwrap();
1204        assert_eq!(bytes.as_slice(), b"hello");
1205
1206        let sub2 = handle.slice(6..11);
1207        let bytes2 = sub2.read_bytes().await.unwrap();
1208        assert_eq!(bytes2.as_slice(), b"world");
1209
1210        // Sync reads work on inline handles
1211        let sync_bytes = handle.read_bytes_range_sync(0..5).unwrap();
1212        assert_eq!(sync_bytes.as_slice(), b"hello");
1213    }
1214
1215    #[tokio::test]
1216    async fn test_owned_bytes() {
1217        let bytes = OwnedBytes::new(vec![1, 2, 3, 4, 5]);
1218
1219        assert_eq!(bytes.len(), 5);
1220        assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1221
1222        let sliced = bytes.slice(1..4);
1223        assert_eq!(sliced.as_slice(), &[2, 3, 4]);
1224
1225        // Original unchanged
1226        assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1227    }
1228}