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 #[cfg(feature = "native")]
591 fn copy_from_file_range(
592 &mut self,
593 _source: &std::fs::File,
594 _source_offset: &mut u64,
595 _len: usize,
596 ) -> io::Result<usize> {
597 Err(io::Error::new(
598 io::ErrorKind::Unsupported,
599 "streaming writer does not support kernel-assisted range copies",
600 ))
601 }
602}
603
604struct BufferedStreamingWriter {
607 path: PathBuf,
608 buffer: Vec<u8>,
609 files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
612}
613
614impl io::Write for BufferedStreamingWriter {
615 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
616 self.buffer.extend_from_slice(buf);
617 Ok(buf.len())
618 }
619
620 fn flush(&mut self) -> io::Result<()> {
621 Ok(())
622 }
623}
624
625impl StreamingWriter for BufferedStreamingWriter {
626 fn finish(self: Box<Self>) -> io::Result<()> {
627 self.files.write().insert(self.path, Arc::new(self.buffer));
628 Ok(())
629 }
630
631 fn bytes_written(&self) -> u64 {
632 self.buffer.len() as u64
633 }
634}
635
636#[cfg(feature = "native")]
640const FILE_STREAMING_BUF_SIZE: usize = 8 * 1024 * 1024;
641
642#[cfg(feature = "native")]
644pub(crate) struct FileStreamingWriter {
645 pub(crate) file: io::BufWriter<std::fs::File>,
646 pub(crate) written: u64,
647}
648
649#[cfg(feature = "native")]
650impl FileStreamingWriter {
651 pub(crate) fn new(file: std::fs::File) -> Self {
652 Self {
653 file: io::BufWriter::with_capacity(FILE_STREAMING_BUF_SIZE, file),
654 written: 0,
655 }
656 }
657}
658
659#[cfg(feature = "native")]
660impl io::Write for FileStreamingWriter {
661 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
662 let n = self.file.write(buf)?;
663 self.written += n as u64;
664 Ok(n)
665 }
666
667 fn flush(&mut self) -> io::Result<()> {
668 self.file.flush()
669 }
670}
671
672#[cfg(feature = "native")]
673impl StreamingWriter for FileStreamingWriter {
674 fn finish(self: Box<Self>) -> io::Result<()> {
675 let file = self.file.into_inner().map_err(|e| e.into_error())?;
676 file.sync_all()?;
677 Ok(())
678 }
679
680 fn bytes_written(&self) -> u64 {
681 self.written
682 }
683
684 fn copy_from_file_range(
685 &mut self,
686 source: &std::fs::File,
687 source_offset: &mut u64,
688 len: usize,
689 ) -> io::Result<usize> {
690 io::Write::flush(&mut self.file)?;
691 let copied = copy_file_range_once(source, source_offset, self.file.get_ref(), len)?;
692 self.written = self
693 .written
694 .checked_add(copied as u64)
695 .ok_or_else(|| io::Error::other("streaming-writer byte count overflow"))?;
696 Ok(copied)
697 }
698}
699
700#[cfg(feature = "native")]
701pub(crate) fn copy_file_range_once(
702 source: &std::fs::File,
703 source_offset: &mut u64,
704 destination: &std::fs::File,
705 len: usize,
706) -> io::Result<usize> {
707 #[cfg(target_os = "linux")]
708 {
709 use std::os::fd::AsRawFd;
710
711 let mut offset = libc::loff_t::try_from(*source_offset).map_err(|_| {
712 io::Error::new(io::ErrorKind::InvalidInput, "source offset exceeds i64")
713 })?;
714 let copied = unsafe {
715 libc::copy_file_range(
716 source.as_raw_fd(),
717 &mut offset,
718 destination.as_raw_fd(),
719 std::ptr::null_mut(),
720 len,
721 0,
722 )
723 };
724 if copied < 0 {
725 let error = io::Error::last_os_error();
726 let unsupported = error.raw_os_error().is_some_and(|code| {
727 code == libc::ENOSYS
728 || code == libc::EXDEV
729 || code == libc::EOPNOTSUPP
730 || code == libc::EINVAL
731 });
732 return if unsupported {
733 Err(io::Error::new(io::ErrorKind::Unsupported, error))
734 } else {
735 Err(error)
736 };
737 }
738 *source_offset = u64::try_from(offset)
739 .map_err(|_| io::Error::other("copy_file_range returned a negative source offset"))?;
740 Ok(copied as usize)
741 }
742 #[cfg(not(target_os = "linux"))]
743 {
744 let _ = (source, source_offset, destination, len);
745 Err(io::Error::new(
746 io::ErrorKind::Unsupported,
747 "kernel-assisted range copies are only available on Linux",
748 ))
749 }
750}
751
752#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
754#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
755pub trait DirectoryWriter: Directory {
756 async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()>;
758
759 async fn write_durable(&self, path: &Path, data: &[u8]) -> io::Result<()> {
768 use io::Write as _;
769 let mut writer = self.streaming_writer(path).await?;
770 writer.write_all(data)?;
771 writer.finish()
772 }
773
774 async fn delete(&self, path: &Path) -> io::Result<()>;
776
777 async fn rename(&self, from: &Path, to: &Path) -> io::Result<()>;
779
780 async fn link(&self, _from: &Path, _to: &Path) -> io::Result<()> {
786 Err(io::Error::new(
787 io::ErrorKind::Unsupported,
788 "directory backend does not support immutable file links",
789 ))
790 }
791
792 async fn sync(&self) -> io::Result<()>;
794
795 async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>>;
798
799 async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
806 self.streaming_writer(path).await
807 }
808}
809
810#[derive(Debug, Default)]
812pub struct RamDirectory {
813 files: Arc<RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>>,
814}
815
816impl Clone for RamDirectory {
817 fn clone(&self) -> Self {
818 Self {
819 files: Arc::clone(&self.files),
820 }
821 }
822}
823
824impl RamDirectory {
825 pub fn new() -> Self {
826 Self::default()
827 }
828
829 pub fn list_files_sync(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
831 let files = self.files.read();
832 Ok(files
833 .keys()
834 .filter(|p| p.starts_with(prefix))
835 .cloned()
836 .collect())
837 }
838
839 pub fn read_file_sync(&self, path: &Path) -> io::Result<Vec<u8>> {
841 let files = self.files.read();
842 files
843 .get(path)
844 .map(|data| data.as_ref().clone())
845 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
846 }
847
848 pub fn write_sync(&self, path: &Path, data: &[u8]) -> io::Result<()> {
850 self.files
851 .write()
852 .insert(path.to_path_buf(), Arc::new(data.to_vec()));
853 Ok(())
854 }
855}
856
857#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
858#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
859impl Directory for RamDirectory {
860 async fn exists(&self, path: &Path) -> io::Result<bool> {
861 Ok(self.files.read().contains_key(path))
862 }
863
864 async fn file_size(&self, path: &Path) -> io::Result<u64> {
865 self.files
866 .read()
867 .get(path)
868 .map(|data| data.len() as u64)
869 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))
870 }
871
872 async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
873 let files = self.files.read();
874 let data = files
875 .get(path)
876 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
877
878 Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
879 Arc::clone(data),
880 0..data.len(),
881 )))
882 }
883
884 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
885 let files = self.files.read();
886 let data = files
887 .get(path)
888 .ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "File not found"))?;
889
890 let start = range.start as usize;
891 let end = range.end as usize;
892
893 if end > data.len() {
894 return Err(io::Error::new(
895 io::ErrorKind::InvalidInput,
896 "Range out of bounds",
897 ));
898 }
899
900 Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end))
901 }
902
903 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
904 let files = self.files.read();
905 Ok(files
906 .keys()
907 .filter(|p| p.starts_with(prefix))
908 .cloned()
909 .collect())
910 }
911
912 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
913 self.open_read(path).await
915 }
916}
917
918#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
919#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
920impl DirectoryWriter for RamDirectory {
921 async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
922 self.files
923 .write()
924 .insert(path.to_path_buf(), Arc::new(data.to_vec()));
925 Ok(())
926 }
927
928 async fn delete(&self, path: &Path) -> io::Result<()> {
929 self.files.write().remove(path);
930 Ok(())
931 }
932
933 async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
934 let mut files = self.files.write();
935 if let Some(data) = files.remove(from) {
936 files.insert(to.to_path_buf(), data);
937 }
938 Ok(())
939 }
940
941 async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
942 let mut files = self.files.write();
943 let data = files.get(from).cloned().ok_or_else(|| {
944 io::Error::new(
945 io::ErrorKind::NotFound,
946 format!("source file {from:?} does not exist"),
947 )
948 })?;
949 files.insert(to.to_path_buf(), data);
950 Ok(())
951 }
952
953 async fn sync(&self) -> io::Result<()> {
954 Ok(())
955 }
956
957 async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
958 Ok(Box::new(BufferedStreamingWriter {
959 path: path.to_path_buf(),
960 buffer: Vec::new(),
961 files: Arc::clone(&self.files),
962 }))
963 }
964}
965
966#[cfg(feature = "native")]
968#[derive(Debug, Clone)]
969pub struct FsDirectory {
970 root: PathBuf,
971 label: IndexLabel,
972}
973
974#[cfg(feature = "native")]
975impl FsDirectory {
976 pub fn new(root: impl AsRef<Path>) -> Self {
977 Self {
978 root: root.as_ref().to_path_buf(),
979 label: IndexLabel::default(),
980 }
981 }
982
983 fn resolve(&self, path: &Path) -> PathBuf {
984 self.root.join(path)
985 }
986}
987
988#[cfg(feature = "native")]
989#[async_trait]
990impl Directory for FsDirectory {
991 async fn exists(&self, path: &Path) -> io::Result<bool> {
992 let full_path = self.resolve(path);
993 tokio::fs::try_exists(&full_path).await
999 }
1000
1001 async fn file_size(&self, path: &Path) -> io::Result<u64> {
1002 let full_path = self.resolve(path);
1003 let metadata = tokio::fs::metadata(&full_path).await?;
1004 Ok(metadata.len())
1005 }
1006
1007 async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
1008 let full_path = self.resolve(path);
1009 let data = tokio::fs::read(&full_path).await?;
1010 Ok(FileHandle::from_bytes(OwnedBytes::new(data)))
1011 }
1012
1013 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
1014 use tokio::io::{AsyncReadExt, AsyncSeekExt};
1015
1016 let full_path = self.resolve(path);
1017 let mut file = tokio::fs::File::open(&full_path).await?;
1018
1019 file.seek(std::io::SeekFrom::Start(range.start)).await?;
1020
1021 let len = (range.end - range.start) as usize;
1022 let mut buffer = vec![0u8; len];
1023 file.read_exact(&mut buffer).await?;
1024
1025 Ok(OwnedBytes::new(buffer))
1026 }
1027
1028 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
1029 let full_path = self.resolve(prefix);
1030 let mut entries = tokio::fs::read_dir(&full_path).await?;
1031 let mut files = Vec::new();
1032
1033 while let Some(entry) = entries.next_entry().await? {
1034 if entry.file_type().await?.is_file() {
1035 files.push(entry.path().strip_prefix(&self.root).unwrap().to_path_buf());
1036 }
1037 }
1038
1039 Ok(files)
1040 }
1041
1042 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
1043 let full_path = self.resolve(path);
1044 let metadata = tokio::fs::metadata(&full_path).await?;
1045 let file_size = metadata.len();
1046
1047 let read_fn: RangeReadFn = Arc::new(move |range: Range<u64>| {
1048 let full_path = full_path.clone();
1049 Box::pin(async move {
1050 use tokio::io::{AsyncReadExt, AsyncSeekExt};
1051
1052 let mut file = tokio::fs::File::open(&full_path).await?;
1053 file.seek(std::io::SeekFrom::Start(range.start)).await?;
1054
1055 let len = (range.end - range.start) as usize;
1056 let mut buffer = vec![0u8; len];
1057 file.read_exact(&mut buffer).await?;
1058
1059 Ok(OwnedBytes::new(buffer))
1060 })
1061 });
1062
1063 Ok(FileHandle::lazy_labeled(
1064 file_size,
1065 read_fn,
1066 self.label.get(),
1067 ))
1068 }
1069
1070 fn set_index_label(&self, label: &str) {
1071 self.label.set(label);
1072 }
1073
1074 fn local_path(&self, path: &Path) -> Option<PathBuf> {
1075 Some(self.resolve(path))
1076 }
1077}
1078
1079#[cfg(feature = "native")]
1080#[async_trait]
1081impl DirectoryWriter for FsDirectory {
1082 async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
1083 let full_path = self.resolve(path);
1084
1085 if let Some(parent) = full_path.parent() {
1087 tokio::fs::create_dir_all(parent).await?;
1088 }
1089
1090 tokio::fs::write(&full_path, data).await
1091 }
1092
1093 async fn delete(&self, path: &Path) -> io::Result<()> {
1094 let full_path = self.resolve(path);
1095 tokio::fs::remove_file(&full_path).await
1096 }
1097
1098 async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
1099 let from_path = self.resolve(from);
1100 let to_path = self.resolve(to);
1101 std::fs::rename(&from_path, &to_path)
1108 }
1109
1110 async fn link(&self, from: &Path, to: &Path) -> io::Result<()> {
1111 std::fs::hard_link(self.resolve(from), self.resolve(to))
1112 }
1113
1114 async fn sync(&self) -> io::Result<()> {
1115 let dir = std::fs::File::open(&self.root)?;
1117 dir.sync_all()?;
1118 Ok(())
1119 }
1120
1121 async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1122 let full_path = self.resolve(path);
1123 if let Some(parent) = full_path.parent() {
1124 tokio::fs::create_dir_all(parent).await?;
1125 }
1126 let file = std::fs::File::create(&full_path)?;
1127 Ok(Box::new(FileStreamingWriter::new(file)))
1128 }
1129
1130 async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
1131 let full_path = self.resolve(path);
1132 if let Some(parent) = full_path.parent() {
1133 tokio::fs::create_dir_all(parent).await?;
1134 }
1135 let file = std::fs::File::create(&full_path)?;
1136 Ok(Box::new(super::ColdStreamingWriter::new(
1137 file,
1138 self.label.get(),
1139 )))
1140 }
1141}
1142
1143pub struct CachingDirectory<D: Directory> {
1145 inner: D,
1146 cache: RwLock<HashMap<PathBuf, Arc<Vec<u8>>>>,
1147 max_cached_bytes: usize,
1148 current_bytes: RwLock<usize>,
1149}
1150
1151impl<D: Directory> CachingDirectory<D> {
1152 pub fn new(inner: D, max_cached_bytes: usize) -> Self {
1153 Self {
1154 inner,
1155 cache: RwLock::new(HashMap::new()),
1156 max_cached_bytes,
1157 current_bytes: RwLock::new(0),
1158 }
1159 }
1160
1161 fn try_cache(&self, path: &Path, data: &[u8]) {
1162 let mut current = self.current_bytes.write();
1163 if *current + data.len() <= self.max_cached_bytes {
1164 self.cache
1165 .write()
1166 .insert(path.to_path_buf(), Arc::new(data.to_vec()));
1167 *current += data.len();
1168 }
1169 }
1170}
1171
1172#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
1173#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
1174impl<D: Directory> Directory for CachingDirectory<D> {
1175 async fn exists(&self, path: &Path) -> io::Result<bool> {
1176 if self.cache.read().contains_key(path) {
1177 return Ok(true);
1178 }
1179 self.inner.exists(path).await
1180 }
1181
1182 async fn file_size(&self, path: &Path) -> io::Result<u64> {
1183 if let Some(data) = self.cache.read().get(path) {
1184 return Ok(data.len() as u64);
1185 }
1186 self.inner.file_size(path).await
1187 }
1188
1189 async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
1190 if let Some(data) = self.cache.read().get(path) {
1192 return Ok(FileHandle::from_bytes(OwnedBytes::from_arc_vec(
1193 Arc::clone(data),
1194 0..data.len(),
1195 )));
1196 }
1197
1198 let handle = self.inner.open_read(path).await?;
1200 let bytes = handle.read_bytes().await?;
1201
1202 self.try_cache(path, bytes.as_slice());
1203
1204 Ok(FileHandle::from_bytes(bytes))
1205 }
1206
1207 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
1208 if let Some(data) = self.cache.read().get(path) {
1210 let start = range.start as usize;
1211 let end = range.end as usize;
1212 return Ok(OwnedBytes::from_arc_vec(Arc::clone(data), start..end));
1213 }
1214
1215 self.inner.read_range(path, range).await
1216 }
1217
1218 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
1219 self.inner.list_files(prefix).await
1220 }
1221
1222 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
1223 self.inner.open_lazy(path).await
1225 }
1226
1227 fn set_index_label(&self, label: &str) {
1228 self.inner.set_index_label(label);
1229 }
1230
1231 fn local_path(&self, path: &Path) -> Option<PathBuf> {
1232 self.inner.local_path(path)
1233 }
1234}
1235
1236#[cfg(test)]
1237mod tests {
1238 use super::*;
1239
1240 #[tokio::test]
1241 async fn test_ram_directory() {
1242 let dir = RamDirectory::new();
1243
1244 dir.write(Path::new("test.txt"), b"hello world")
1246 .await
1247 .unwrap();
1248
1249 assert!(dir.exists(Path::new("test.txt")).await.unwrap());
1251 assert!(!dir.exists(Path::new("nonexistent.txt")).await.unwrap());
1252
1253 let slice = dir.open_read(Path::new("test.txt")).await.unwrap();
1255 let data = slice.read_bytes().await.unwrap();
1256 assert_eq!(data.as_slice(), b"hello world");
1257
1258 let range_data = dir.read_range(Path::new("test.txt"), 0..5).await.unwrap();
1260 assert_eq!(range_data.as_slice(), b"hello");
1261
1262 dir.delete(Path::new("test.txt")).await.unwrap();
1264 assert!(!dir.exists(Path::new("test.txt")).await.unwrap());
1265 }
1266
1267 #[cfg(all(unix, feature = "native"))]
1272 #[tokio::test]
1273 async fn test_fs_exists_propagates_stat_errors_instead_of_reporting_missing() {
1274 use std::os::unix::fs::PermissionsExt;
1275
1276 let temp_dir = tempfile::TempDir::new().unwrap();
1277 let dir = FsDirectory::new(temp_dir.path());
1278 dir.write(Path::new("locked/seg.meta"), b"data")
1279 .await
1280 .unwrap();
1281
1282 let locked = temp_dir.path().join("locked");
1285 let original = std::fs::metadata(&locked).unwrap().permissions();
1286 std::fs::set_permissions(&locked, std::fs::Permissions::from_mode(0o000)).unwrap();
1287 if std::fs::metadata(locked.join("seg.meta")).is_ok() {
1288 std::fs::set_permissions(&locked, original).unwrap();
1291 return;
1292 }
1293 let result = dir.exists(Path::new("locked/seg.meta")).await;
1294 std::fs::set_permissions(&locked, original).unwrap();
1295
1296 let error =
1297 result.expect_err("stat failure must propagate as Err, not be misreported as missing");
1298 assert_ne!(error.kind(), io::ErrorKind::NotFound);
1299 assert!(dir.exists(Path::new("locked/seg.meta")).await.unwrap());
1301 }
1302
1303 #[tokio::test]
1304 async fn test_file_handle() {
1305 let data = OwnedBytes::new(b"hello world".to_vec());
1306 let handle = FileHandle::from_bytes(data);
1307
1308 assert_eq!(handle.len(), 11);
1309 assert!(handle.is_sync());
1310
1311 let sub = handle.slice(0..5);
1312 let bytes = sub.read_bytes().await.unwrap();
1313 assert_eq!(bytes.as_slice(), b"hello");
1314
1315 let sub2 = handle.slice(6..11);
1316 let bytes2 = sub2.read_bytes().await.unwrap();
1317 assert_eq!(bytes2.as_slice(), b"world");
1318
1319 let sync_bytes = handle.read_bytes_range_sync(0..5).unwrap();
1321 assert_eq!(sync_bytes.as_slice(), b"hello");
1322 }
1323
1324 #[tokio::test]
1325 async fn test_owned_bytes() {
1326 let bytes = OwnedBytes::new(vec![1, 2, 3, 4, 5]);
1327
1328 assert_eq!(bytes.len(), 5);
1329 assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1330
1331 let sliced = bytes.slice(1..4);
1332 assert_eq!(sliced.as_slice(), &[2, 3, 4]);
1333
1334 assert_eq!(bytes.as_slice(), &[1, 2, 3, 4, 5]);
1336 }
1337}