1use 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#[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#[derive(Clone)]
40pub struct FileHandle {
41 inner: FileHandleInner,
42}
43
44#[derive(Clone)]
45enum FileHandleInner {
46 Inline {
48 data: OwnedBytes,
49 offset: u64,
50 len: u64,
51 },
52 Lazy {
54 read_fn: RangeReadFn,
55 offset: u64,
56 len: u64,
57 label: Arc<str>,
59 },
60}
61
62#[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 pub fn get(&self) -> Arc<str> {
81 self.0.read().expect("IndexLabel lock poisoned").clone()
82 }
83
84 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 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 pub fn empty() -> Self {
123 Self::from_bytes(OwnedBytes::empty())
124 }
125
126 pub fn lazy(len: u64, read_fn: RangeReadFn) -> Self {
131 Self::lazy_labeled(len, read_fn, Arc::from("unknown"))
132 }
133
134 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 #[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 #[inline]
157 pub fn is_empty(&self) -> bool {
158 self.len() == 0
159 }
160
161 #[inline]
163 pub fn is_sync(&self) -> bool {
164 matches!(&self.inner, FileHandleInner::Inline { .. })
165 }
166
167 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 #[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 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 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 pub async fn read_bytes(&self) -> io::Result<OwnedBytes> {
277 self.read_bytes_range(0..self.len()).await
278 }
279
280 #[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 #[inline]
305 pub fn read_bytes_sync(&self) -> io::Result<OwnedBytes> {
306 self.read_bytes_range_sync(0..self.len())
307 }
308}
309
310#[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#[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 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 #[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 #[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 #[cfg(feature = "native")]
420 #[inline]
421 pub fn is_mmap(&self) -> bool {
422 matches!(self.data, SharedBytes::Mmap(_))
423 }
424
425 #[cfg(feature = "native")]
431 pub fn madvise(&self, advice: libc::c_int) {
432 self.madvise_range(0..self.len(), advice);
433 }
434
435 #[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 #[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#[cfg(not(target_arch = "wasm32"))]
501#[async_trait]
502pub trait Directory: Send + Sync + 'static {
503 async fn exists(&self, path: &Path) -> io::Result<bool>;
505
506 async fn file_size(&self, path: &Path) -> io::Result<u64>;
508
509 async fn open_read(&self, path: &Path) -> io::Result<FileHandle>;
511
512 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes>;
514
515 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>>;
517
518 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle>;
522
523 fn set_index_label(&self, _label: &str) {}
529
530 fn local_path(&self, _path: &Path) -> Option<PathBuf> {
537 None
538 }
539}
540
541#[cfg(target_arch = "wasm32")]
543#[async_trait(?Send)]
544pub trait Directory: 'static {
545 async fn exists(&self, path: &Path) -> io::Result<bool>;
547
548 async fn file_size(&self, path: &Path) -> io::Result<u64>;
550
551 async fn open_read(&self, path: &Path) -> io::Result<FileHandle>;
553
554 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes>;
556
557 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>>;
559
560 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle>;
562
563 fn set_index_label(&self, _label: &str) {}
566
567 fn local_path(&self, _path: &Path) -> Option<PathBuf> {
569 None
570 }
571}
572
573pub trait StreamingWriter: io::Write + Send {
578 fn finish(self: Box<Self>) -> io::Result<()>;
580
581 fn bytes_written(&self) -> u64;
583}
584
585struct BufferedStreamingWriter {
588 path: PathBuf,
589 buffer: Vec<u8>,
590 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#[cfg(feature = "native")]
621const FILE_STREAMING_BUF_SIZE: usize = 8 * 1024 * 1024;
622
623#[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#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
668#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
669pub trait DirectoryWriter: Directory {
670 async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()>;
672
673 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 async fn delete(&self, path: &Path) -> io::Result<()>;
690
691 async fn rename(&self, from: &Path, to: &Path) -> io::Result<()>;
693
694 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 async fn sync(&self) -> io::Result<()>;
708
709 async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>>;
712
713 async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
720 self.streaming_writer(path).await
721 }
722}
723
724#[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 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 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 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 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#[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 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 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 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 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
1057pub 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 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 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 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 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 dir.write(Path::new("test.txt"), b"hello world")
1160 .await
1161 .unwrap();
1162
1163 assert!(dir.exists(Path::new("test.txt")).await.unwrap());
1165 assert!(!dir.exists(Path::new("nonexistent.txt")).await.unwrap());
1166
1167 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 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 dir.delete(Path::new("test.txt")).await.unwrap();
1178 assert!(!dir.exists(Path::new("test.txt")).await.unwrap());
1179 }
1180
1181 #[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 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 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 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 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 assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1250 }
1251}