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    /// Resolve a directory-relative file to a native filesystem path.
531    ///
532    /// Local backends expose this so large, short-lived merge scratch files
533    /// can live beside the index instead of silently spilling to the
534    /// container's root filesystem. Remote and in-memory backends return
535    /// `None`.
536    fn local_path(&self, _path: &Path) -> Option<PathBuf> {
537        None
538    }
539}
540
541/// Async directory trait for reading index files (WASM version - no Send requirement)
542#[cfg(target_arch = "wasm32")]
543#[async_trait(?Send)]
544pub trait Directory: 'static {
545    /// Check if a file exists
546    async fn exists(&self, path: &Path) -> io::Result<bool>;
547
548    /// Get file size
549    async fn file_size(&self, path: &Path) -> io::Result<u64>;
550
551    /// Open a file for reading (loads entire file into an inline FileHandle)
552    async fn open_read(&self, path: &Path) -> io::Result<FileHandle>;
553
554    /// Read a specific byte range from a file (optimized for network)
555    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes>;
556
557    /// List files in directory
558    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>>;
559
560    /// Open a file handle that fetches ranges on demand.
561    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle>;
562
563    /// Attach the owning index's name for Directory-layer metric labels.
564    /// No-op default; metrics are native-only but the label is harmless.
565    fn set_index_label(&self, _label: &str) {}
566
567    /// WASM backends do not expose a native filesystem path.
568    fn local_path(&self, _path: &Path) -> Option<PathBuf> {
569        None
570    }
571}
572
573/// A writer for incrementally writing data to a directory file.
574///
575/// Avoids buffering entire files in memory during merge. File-backed
576/// directories write directly to disk; memory directories collect to Vec.
577pub trait StreamingWriter: io::Write + Send {
578    /// Finalize the write, making data available for reading.
579    fn finish(self: Box<Self>) -> io::Result<()>;
580
581    /// Bytes written so far.
582    fn bytes_written(&self) -> u64;
583}
584
585/// StreamingWriter backed by Vec<u8>, finalized via DirectoryWriter::write.
586/// Used as default/fallback and for RamDirectory.
587struct BufferedStreamingWriter {
588    path: PathBuf,
589    buffer: Vec<u8>,
590    /// Callback to write the buffer to the directory on finish.
591    /// We store the files Arc directly for RamDirectory.
592    files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
593}
594
595impl io::Write for BufferedStreamingWriter {
596    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
597        self.buffer.extend_from_slice(buf);
598        Ok(buf.len())
599    }
600
601    fn flush(&mut self) -> io::Result<()> {
602        Ok(())
603    }
604}
605
606impl StreamingWriter for BufferedStreamingWriter {
607    fn finish(self: Box<Self>) -> io::Result<()> {
608        self.files.write().insert(self.path, Arc::new(self.buffer));
609        Ok(())
610    }
611
612    fn bytes_written(&self) -> u64 {
613        self.buffer.len() as u64
614    }
615}
616
617/// Buffer size for FileStreamingWriter (8 MB).
618/// Large enough to coalesce millions of tiny writes (e.g. per-vector doc_id writes)
619/// into efficient sequential I/O.
620#[cfg(feature = "native")]
621const FILE_STREAMING_BUF_SIZE: usize = 8 * 1024 * 1024;
622
623/// StreamingWriter backed by a buffered std::fs::File for filesystem directories.
624#[cfg(feature = "native")]
625pub(crate) struct FileStreamingWriter {
626    pub(crate) file: io::BufWriter<std::fs::File>,
627    pub(crate) written: u64,
628}
629
630#[cfg(feature = "native")]
631impl FileStreamingWriter {
632    pub(crate) fn new(file: std::fs::File) -> Self {
633        Self {
634            file: io::BufWriter::with_capacity(FILE_STREAMING_BUF_SIZE, file),
635            written: 0,
636        }
637    }
638}
639
640#[cfg(feature = "native")]
641impl io::Write for FileStreamingWriter {
642    fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
643        let n = self.file.write(buf)?;
644        self.written += n as u64;
645        Ok(n)
646    }
647
648    fn flush(&mut self) -> io::Result<()> {
649        self.file.flush()
650    }
651}
652
653#[cfg(feature = "native")]
654impl StreamingWriter for FileStreamingWriter {
655    fn finish(self: Box<Self>) -> io::Result<()> {
656        let file = self.file.into_inner().map_err(|e| e.into_error())?;
657        file.sync_all()?;
658        Ok(())
659    }
660
661    fn bytes_written(&self) -> u64 {
662        self.written
663    }
664}
665
666/// Async directory trait for writing index files
667#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
668#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
669pub trait DirectoryWriter: Directory {
670    /// Create/overwrite a file with data
671    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()>;
672
673    /// Create/overwrite a file with data, durably.
674    ///
675    /// [`Self::write`] does not guarantee the bytes reach stable storage
676    /// before returning (filesystem implementations leave them in the OS
677    /// page cache). Any file that durably-published metadata will reference
678    /// (e.g. segment `.meta`) must be written through this method instead:
679    /// it routes through [`Self::streaming_writer`], whose `finish()` fsyncs
680    /// on filesystem implementations.
681    async fn write_durable(&self, path: &Path, data: &[u8]) -> io::Result<()> {
682        use io::Write as _;
683        let mut writer = self.streaming_writer(path).await?;
684        writer.write_all(data)?;
685        writer.finish()
686    }
687
688    /// Delete a file
689    async fn delete(&self, path: &Path) -> io::Result<()>;
690
691    /// Atomic rename
692    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()>;
693
694    /// Create another immutable name for an existing file without copying its
695    /// contents when the backend supports it. Segment rewrites use this to
696    /// retain unchanged multi-gigabyte files while replacing only one index
697    /// payload. Backends without link semantics return `Unsupported`; callers
698    /// then fall back to a streaming copy.
699    async fn link(&self, _from: &Path, _to: &Path) -> io::Result<()> {
700        Err(io::Error::new(
701            io::ErrorKind::Unsupported,
702            "directory backend does not support immutable file links",
703        ))
704    }
705
706    /// Sync all pending writes
707    async fn sync(&self) -> io::Result<()>;
708
709    /// Create a streaming writer for incremental file writes.
710    /// Call finish() on the returned writer to finalize.
711    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>>;
712
713    /// Streaming writer for **bulk one-shot data** (merge/reorder outputs).
714    ///
715    /// Filesystem directories return a page-cache-dropping writer (see
716    /// `docs/cold-io.md`) so multi-GB merge writes cannot evict the serving
717    /// segments' warm pages. Output is byte-identical to the buffered
718    /// writer. Default impl delegates to [`Self::streaming_writer`].
719    async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
720        self.streaming_writer(path).await
721    }
722}
723
724/// In-memory directory for testing and small indexes
725#[derive(Debug, Default)]
726pub struct RamDirectory {
727    files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
728}
729
730impl Clone for RamDirectory {
731    fn clone(&self) -> Self {
732        Self {
733            files: Arc::clone(&self.files),
734        }
735    }
736}
737
738impl RamDirectory {
739    pub fn new() -> Self {
740        Self::default()
741    }
742
743    /// Synchronous file listing (for serialization).
744    pub fn list_files_sync(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
745        let files = self.files.read();
746        Ok(files
747            .keys()
748            .filter(|p| p.starts_with(prefix))
749            .cloned()
750            .collect())
751    }
752
753    /// Synchronous file read (for serialization).
754    pub fn read_file_sync(&self, path: &Path) -> io::Result<Vec<u8>> {
755        let files = self.files.read();
756        files
757            .get(path)
758            .map(|data| data.as_ref().clone())
759            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
760    }
761
762    /// Synchronous file write (for deserialization).
763    pub fn write_sync(&self, path: &Path, data: &[u8]) -> io::Result<()> {
764        self.files
765            .write()
766            .insert(path.to_path_buf(), Arc::new(data.to_vec()));
767        Ok(())
768    }
769}
770
771#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
772#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
773impl Directory for RamDirectory {
774    async fn exists(&self, path: &Path) -> io::Result<bool> {
775        Ok(self.files.read().contains_key(path))
776    }
777
778    async fn file_size(&self, path: &Path) -> io::Result<u64> {
779        self.files
780            .read()
781            .get(path)
782            .map(|data| data.len() as u64)
783            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
784    }
785
786    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
787        let files = self.files.read();
788        let data = files
789            .get(path)
790            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
791
792        Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
793            Arc::clone(data),
794            0..data.len(),
795        )))
796    }
797
798    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
799        let files = self.files.read();
800        let data = files
801            .get(path)
802            .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
803
804        let start = range.start as usize;
805        let end = range.end as usize;
806
807        if end > data.len() {
808            return Err(io::Error::new(
809                io::ErrorKind::InvalidInput,
810                "Range out of bounds",
811            ));
812        }
813
814        Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end))
815    }
816
817    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
818        let files = self.files.read();
819        Ok(files
820            .keys()
821            .filter(|p| p.starts_with(prefix))
822            .cloned()
823            .collect())
824    }
825
826    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
827        // RAM data is always available synchronously — return Inline handle
828        self.open_read(path).await
829    }
830}
831
832#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
833#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
834impl DirectoryWriter for RamDirectory {
835    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
836        self.files
837            .write()
838            .insert(path.to_path_buf(), Arc::new(data.to_vec()));
839        Ok(())
840    }
841
842    async fn delete(&self, path: &Path) -> io::Result<()> {
843        self.files.write().remove(path);
844        Ok(())
845    }
846
847    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
848        let mut files = self.files.write();
849        if let Some(data) = files.remove(from) {
850            files.insert(to.to_path_buf(), data);
851        }
852        Ok(())
853    }
854
855    async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
856        let mut files = self.files.write();
857        let data = files.get(from).cloned().ok_or_else(|| {
858            io::Error::new(
859                io::ErrorKind::NotFound,
860                format!("source file {from:?} does not exist"),
861            )
862        })?;
863        files.insert(to.to_path_buf(), data);
864        Ok(())
865    }
866
867    async fn sync(&self) -> io::Result<()> {
868        Ok(())
869    }
870
871    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
872        Ok(Box::new(BufferedStreamingWriter {
873            path: path.to_path_buf(),
874            buffer: Vec::new(),
875            files: Arc::clone(&self.files),
876        }))
877    }
878}
879
880/// Local filesystem directory with async IO via tokio
881#[cfg(feature = "native")]
882#[derive(Debug, Clone)]
883pub struct FsDirectory {
884    root: PathBuf,
885    label: IndexLabel,
886}
887
888#[cfg(feature = "native")]
889impl FsDirectory {
890    pub fn new(root: impl AsRef<Path>) -> Self {
891        Self {
892            root: root.as_ref().to_path_buf(),
893            label: IndexLabel::default(),
894        }
895    }
896
897    fn resolve(&self, path: &Path) -> PathBuf {
898        self.root.join(path)
899    }
900}
901
902#[cfg(feature = "native")]
903#[async_trait]
904impl Directory for FsDirectory {
905    async fn exists(&self, path: &Path) -> io::Result<bool> {
906        let full_path = self.resolve(path);
907        // `try_exists` maps NotFound to Ok(false); any other stat failure
908        // (EACCES, EIO, ...) must propagate so callers can distinguish a
909        // genuinely missing file from a transient IO error — swallowing it
910        // as `false` quarantines a healthy segment as "missing mandatory
911        // files" instead of retrying.
912        tokio::fs::try_exists(&full_path).await
913    }
914
915    async fn file_size(&self, path: &Path) -> io::Result<u64> {
916        let full_path = self.resolve(path);
917        let metadata = tokio::fs::metadata(&full_path).await?;
918        Ok(metadata.len())
919    }
920
921    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
922        let full_path = self.resolve(path);
923        let data = tokio::fs::read(&full_path).await?;
924        Ok(FileHandle::from_bytes(OwnedBytes::new(data)))
925    }
926
927    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
928        use tokio::io::{AsyncReadExt, AsyncSeekExt};
929
930        let full_path = self.resolve(path);
931        let mut file = tokio::fs::File::open(&full_path).await?;
932
933        file.seek(std::io::SeekFrom::Start(range.start)).await?;
934
935        let len = (range.end - range.start) as usize;
936        let mut buffer = vec![0u8; len];
937        file.read_exact(&mut buffer).await?;
938
939        Ok(OwnedBytes::new(buffer))
940    }
941
942    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
943        let full_path = self.resolve(prefix);
944        let mut entries = tokio::fs::read_dir(&full_path).await?;
945        let mut files = Vec::new();
946
947        while let Some(entry) = entries.next_entry().await? {
948            if entry.file_type().await?.is_file() {
949                files.push(entry.path().strip_prefix(&self.root).unwrap().to_path_buf());
950            }
951        }
952
953        Ok(files)
954    }
955
956    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
957        let full_path = self.resolve(path);
958        let metadata = tokio::fs::metadata(&full_path).await?;
959        let file_size = metadata.len();
960
961        let read_fn: RangeReadFn = Arc::new(move |range: Range<u64>| {
962            let full_path = full_path.clone();
963            Box::pin(async move {
964                use tokio::io::{AsyncReadExt, AsyncSeekExt};
965
966                let mut file = tokio::fs::File::open(&full_path).await?;
967                file.seek(std::io::SeekFrom::Start(range.start)).await?;
968
969                let len = (range.end - range.start) as usize;
970                let mut buffer = vec![0u8; len];
971                file.read_exact(&mut buffer).await?;
972
973                Ok(OwnedBytes::new(buffer))
974            })
975        });
976
977        Ok(FileHandle::lazy_labeled(
978            file_size,
979            read_fn,
980            self.label.get(),
981        ))
982    }
983
984    fn set_index_label(&self, label: &str) {
985        self.label.set(label);
986    }
987
988    fn local_path(&self, path: &Path) -> Option<PathBuf> {
989        Some(self.resolve(path))
990    }
991}
992
993#[cfg(feature = "native")]
994#[async_trait]
995impl DirectoryWriter for FsDirectory {
996    async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
997        let full_path = self.resolve(path);
998
999        // Ensure parent directory exists
1000        if let Some(parent) = full_path.parent() {
1001            tokio::fs::create_dir_all(parent).await?;
1002        }
1003
1004        tokio::fs::write(&full_path, data).await
1005    }
1006
1007    async fn delete(&self, path: &Path) -> io::Result<()> {
1008        let full_path = self.resolve(path);
1009        tokio::fs::remove_file(&full_path).await
1010    }
1011
1012    async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
1013        let from_path = self.resolve(from);
1014        let to_path = self.resolve(to);
1015        // Metadata publication is the only rename user. Keep the atomic
1016        // filesystem operation in a single future poll: tokio::fs::rename is
1017        // backed by a cancellable await around spawn_blocking, so a dropped
1018        // commit future could observe neither completion nor failure even
1019        // though the rename later succeeded. The caller must update its
1020        // in-memory metadata in the same poll after this returns.
1021        std::fs::rename(&from_path, &to_path)
1022    }
1023
1024    async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
1025        std::fs::hard_link(self.resolve(from), self.resolve(to))
1026    }
1027
1028    async fn sync(&self) -> io::Result<()> {
1029        // fsync the directory
1030        let dir = std::fs::File::open(&self.root)?;
1031        dir.sync_all()?;
1032        Ok(())
1033    }
1034
1035    async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1036        let full_path = self.resolve(path);
1037        if let Some(parent) = full_path.parent() {
1038            tokio::fs::create_dir_all(parent).await?;
1039        }
1040        let file = std::fs::File::create(&full_path)?;
1041        Ok(Box::new(FileStreamingWriter::new(file)))
1042    }
1043
1044    async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1045        let full_path = self.resolve(path);
1046        if let Some(parent) = full_path.parent() {
1047            tokio::fs::create_dir_all(parent).await?;
1048        }
1049        let file = std::fs::File::create(&full_path)?;
1050        Ok(Box::new(super::ColdStreamingWriter::new(
1051            file,
1052            self.label.get(),
1053        )))
1054    }
1055}
1056
1057/// Caching wrapper for any Directory - caches file reads
1058pub struct CachingDirectory<D: Directory> {
1059    inner: D,
1060    cache: RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>,
1061    max_cached_bytes: usize,
1062    current_bytes: RwLock<usize>,
1063}
1064
1065impl<D: Directory> CachingDirectory<D> {
1066    pub fn new(inner: D, max_cached_bytes: usize) -> Self {
1067        Self {
1068            inner,
1069            cache: RwLock::new(HashMap::new()),
1070            max_cached_bytes,
1071            current_bytes: RwLock::new(0),
1072        }
1073    }
1074
1075    fn try_cache(&self, path: &Path, data: &[u8]) {
1076        let mut current = self.current_bytes.write();
1077        if *current + data.len() <= self.max_cached_bytes {
1078            self.cache
1079                .write()
1080                .insert(path.to_path_buf(), Arc::new(data.to_vec()));
1081            *current += data.len();
1082        }
1083    }
1084}
1085
1086#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1087#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1088impl<D: Directory> Directory for CachingDirectory<D> {
1089    async fn exists(&self, path: &Path) -> io::Result<bool> {
1090        if self.cache.read().contains_key(path) {
1091            return Ok(true);
1092        }
1093        self.inner.exists(path).await
1094    }
1095
1096    async fn file_size(&self, path: &Path) -> io::Result<u64> {
1097        if let Some(data) = self.cache.read().get(path) {
1098            return Ok(data.len() as u64);
1099        }
1100        self.inner.file_size(path).await
1101    }
1102
1103    async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
1104        // Check cache first
1105        if let Some(data) = self.cache.read().get(path) {
1106            return Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
1107                Arc::clone(data),
1108                0..data.len(),
1109            )));
1110        }
1111
1112        // Read from inner and potentially cache
1113        let handle = self.inner.open_read(path).await?;
1114        let bytes = handle.read_bytes().await?;
1115
1116        self.try_cache(path, bytes.as_slice());
1117
1118        Ok(FileHandle::from_bytes(bytes))
1119    }
1120
1121    async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
1122        // Check cache first
1123        if let Some(data) = self.cache.read().get(path) {
1124            let start = range.start as usize;
1125            let end = range.end as usize;
1126            return Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end));
1127        }
1128
1129        self.inner.read_range(path, range).await
1130    }
1131
1132    async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
1133        self.inner.list_files(prefix).await
1134    }
1135
1136    async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
1137        // For caching directory, delegate to inner - caching happens at read_range level
1138        self.inner.open_lazy(path).await
1139    }
1140
1141    fn set_index_label(&self, label: &str) {
1142        self.inner.set_index_label(label);
1143    }
1144
1145    fn local_path(&self, path: &Path) -> Option<PathBuf> {
1146        self.inner.local_path(path)
1147    }
1148}
1149
1150#[cfg(test)]
1151mod tests {
1152    use super::*;
1153
1154    #[tokio::test]
1155    async fn test_ram_directory() {
1156        let dir = RamDirectory::new();
1157
1158        // Write file
1159        dir.write(Path::new("test.txt"), b"hello world")
1160            .await
1161            .unwrap();
1162
1163        // Check exists
1164        assert!(dir.exists(Path::new("test.txt")).await.unwrap());
1165        assert!(!dir.exists(Path::new("nonexistent.txt")).await.unwrap());
1166
1167        // Read file
1168        let slice = dir.open_read(Path::new("test.txt")).await.unwrap();
1169        let data = slice.read_bytes().await.unwrap();
1170        assert_eq!(data.as_slice(), b"hello world");
1171
1172        // Read range
1173        let range_data = dir.read_range(Path::new("test.txt"), 0..5).await.unwrap();
1174        assert_eq!(range_data.as_slice(), b"hello");
1175
1176        // Delete
1177        dir.delete(Path::new("test.txt")).await.unwrap();
1178        assert!(!dir.exists(Path::new("test.txt")).await.unwrap());
1179    }
1180
1181    /// A transient stat failure (EACCES here, EIO on flaky storage) must
1182    /// surface as `Err`, not `Ok(false)`: callers classify a missing
1183    /// mandatory segment file as deterministic corruption and quarantine
1184    /// the segment until restart.
1185    #[cfg(all(unix, feature = "native"))]
1186    #[tokio::test]
1187    async fn test_fs_exists_propagates_stat_errors_instead_of_reporting_missing() {
1188        use std::os::unix::fs::PermissionsExt;
1189
1190        let temp_dir = tempfile::TempDir::new().unwrap();
1191        let dir = FsDirectory::new(temp_dir.path());
1192        dir.write(Path::new("locked/seg.meta"), b"data")
1193            .await
1194            .unwrap();
1195
1196        // Removing search permission from the parent makes stat on the child
1197        // fail with EACCES while the file itself still exists.
1198        let locked = temp_dir.path().join("locked");
1199        let original = std::fs::metadata(&locked).unwrap().permissions();
1200        std::fs::set_permissions(&locked, std::fs::Permissions::from_mode(0o000)).unwrap();
1201        if std::fs::metadata(locked.join("seg.meta")).is_ok() {
1202            // Running as root: directory permissions are not enforced, so the
1203            // stat failure cannot be provoked.
1204            std::fs::set_permissions(&locked, original).unwrap();
1205            return;
1206        }
1207        let result = dir.exists(Path::new("locked/seg.meta")).await;
1208        std::fs::set_permissions(&locked, original).unwrap();
1209
1210        let error =
1211            result.expect_err("stat failure must propagate as Err, not be misreported as missing");
1212        assert_ne!(error.kind(), io::ErrorKind::NotFound);
1213        // Once stat succeeds again the file is reported present.
1214        assert!(dir.exists(Path::new("locked/seg.meta")).await.unwrap());
1215    }
1216
1217    #[tokio::test]
1218    async fn test_file_handle() {
1219        let data = OwnedBytes::new(b"hello world".to_vec());
1220        let handle = FileHandle::from_bytes(data);
1221
1222        assert_eq!(handle.len(), 11);
1223        assert!(handle.is_sync());
1224
1225        let sub = handle.slice(0..5);
1226        let bytes = sub.read_bytes().await.unwrap();
1227        assert_eq!(bytes.as_slice(), b"hello");
1228
1229        let sub2 = handle.slice(6..11);
1230        let bytes2 = sub2.read_bytes().await.unwrap();
1231        assert_eq!(bytes2.as_slice(), b"world");
1232
1233        // Sync reads work on inline handles
1234        let sync_bytes = handle.read_bytes_range_sync(0..5).unwrap();
1235        assert_eq!(sync_bytes.as_slice(), b"hello");
1236    }
1237
1238    #[tokio::test]
1239    async fn test_owned_bytes() {
1240        let bytes = OwnedBytes::new(vec![1, 2, 3, 4, 5]);
1241
1242        assert_eq!(bytes.len(), 5);
1243        assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1244
1245        let sliced = bytes.slice(1..4);
1246        assert_eq!(sliced.as_slice(), &[2, 3, 4]);
1247
1248        // Original unchanged
1249        assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1250    }
1251}