hermes_core/directories/
mmap.rs1use std::io;
6use std::ops::Range;
7use std::path::{Path, PathBuf};
8use std::sync::Arc;
9
10use async_trait::async_trait;
11use memmap2::Mmap;
12
13use super::{
14 Directory, DirectoryWriter, FileHandle, FileStreamingWriter, OwnedBytes, StreamingWriter,
15};
16
17pub struct MmapDirectory {
31 root: PathBuf,
32 label: super::IndexLabel,
33}
34
35impl MmapDirectory {
36 pub fn new(root: impl AsRef<Path>) -> Self {
38 Self {
39 root: root.as_ref().to_path_buf(),
40 label: super::IndexLabel::default(),
41 }
42 }
43
44 pub fn root(&self) -> &Path {
46 &self.root
47 }
48
49 fn resolve(&self, path: &Path) -> PathBuf {
50 self.root.join(path)
51 }
52
53 fn mmap_file(&self, path: &Path) -> io::Result<Arc<Mmap>> {
55 let full_path = self.resolve(path);
56 let file = std::fs::File::open(&full_path)?;
57 let mmap = unsafe { Mmap::map(&file)? };
58 Ok(Arc::new(mmap))
59 }
60}
61
62impl Clone for MmapDirectory {
63 fn clone(&self) -> Self {
64 Self {
65 root: self.root.clone(),
66 label: self.label.clone(),
68 }
69 }
70}
71
72#[async_trait]
73impl Directory for MmapDirectory {
74 async fn exists(&self, path: &Path) -> io::Result<bool> {
75 let full_path = self.resolve(path);
76 Ok(tokio::fs::try_exists(&full_path).await.unwrap_or(false))
77 }
78
79 async fn file_size(&self, path: &Path) -> io::Result<u64> {
80 let full_path = self.resolve(path);
81 let metadata = tokio::fs::metadata(&full_path).await?;
82 Ok(metadata.len())
83 }
84
85 async fn open_read(&self, path: &Path) -> io::Result<FileHandle> {
86 let mmap = self.mmap_file(path)?;
87 Ok(FileHandle::from_bytes(OwnedBytes::from_mmap(mmap)))
89 }
90
91 async fn read_range(&self, path: &Path, range: Range<u64>) -> io::Result<OwnedBytes> {
92 let mmap = self.mmap_file(path)?;
93 let start = range.start as usize;
94 let end = range.end as usize;
95
96 if end > mmap.len() {
97 return Err(io::Error::new(
98 io::ErrorKind::InvalidInput,
99 format!("Range {}..{} exceeds file size {}", start, end, mmap.len()),
100 ));
101 }
102
103 Ok(OwnedBytes::from_mmap_range(mmap, start..end))
105 }
106
107 async fn list_files(&self, prefix: &Path) -> io::Result<Vec<PathBuf>> {
108 let full_path = self.resolve(prefix);
109 let mut entries = tokio::fs::read_dir(&full_path).await?;
110 let mut files = Vec::new();
111
112 while let Some(entry) = entries.next_entry().await? {
113 if entry.file_type().await?.is_file() {
114 files.push(entry.path().strip_prefix(&self.root).unwrap().to_path_buf());
115 }
116 }
117
118 Ok(files)
119 }
120
121 async fn open_lazy(&self, path: &Path) -> io::Result<FileHandle> {
122 self.open_read(path).await
125 }
126
127 fn set_index_label(&self, label: &str) {
128 self.label.set(label);
129 }
130}
131
132#[async_trait]
133impl DirectoryWriter for MmapDirectory {
134 async fn write(&self, path: &Path, data: &[u8]) -> io::Result<()> {
135 let full_path = self.resolve(path);
136
137 if let Some(parent) = full_path.parent() {
139 tokio::fs::create_dir_all(parent).await?;
140 }
141
142 tokio::fs::write(&full_path, data).await
143 }
144
145 async fn delete(&self, path: &Path) -> io::Result<()> {
146 let full_path = self.resolve(path);
147 tokio::fs::remove_file(&full_path).await
148 }
149
150 async fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
151 let from_path = self.resolve(from);
152 let to_path = self.resolve(to);
153 std::fs::rename(&from_path, &to_path)
157 }
158
159 async fn sync(&self) -> io::Result<()> {
160 let dir = std::fs::File::open(&self.root)?;
162 dir.sync_all()?;
163 Ok(())
164 }
165
166 async fn streaming_writer(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
167 let full_path = self.resolve(path);
168 if let Some(parent) = full_path.parent() {
169 tokio::fs::create_dir_all(parent).await?;
170 }
171 let file = std::fs::File::create(&full_path)?;
172 Ok(Box::new(FileStreamingWriter::new(file)))
173 }
174
175 async fn streaming_writer_cold(&self, path: &Path) -> io::Result<Box<dyn StreamingWriter>> {
176 let full_path = self.resolve(path);
177 if let Some(parent) = full_path.parent() {
178 tokio::fs::create_dir_all(parent).await?;
179 }
180 let file = std::fs::File::create(&full_path)?;
181 Ok(Box::new(super::ColdStreamingWriter::new(
182 file,
183 self.label.get(),
184 )))
185 }
186}
187
188#[cfg(test)]
189mod tests {
190 use super::*;
191 use tempfile::TempDir;
192
193 #[tokio::test]
194 async fn test_mmap_directory_basic() {
195 let temp_dir = TempDir::new().unwrap();
196 let dir = MmapDirectory::new(temp_dir.path());
197
198 let test_data = b"Hello, mmap world!";
200 dir.write(Path::new("test.txt"), test_data).await.unwrap();
201
202 assert!(dir.exists(Path::new("test.txt")).await.unwrap());
204 assert!(!dir.exists(Path::new("nonexistent.txt")).await.unwrap());
205
206 assert_eq!(
208 dir.file_size(Path::new("test.txt")).await.unwrap(),
209 test_data.len() as u64
210 );
211
212 let slice = dir.open_read(Path::new("test.txt")).await.unwrap();
214 let bytes = slice.read_bytes().await.unwrap();
215 assert_eq!(bytes.as_slice(), test_data);
216
217 let range_bytes = dir.read_range(Path::new("test.txt"), 7..12).await.unwrap();
219 assert_eq!(range_bytes.as_slice(), b"mmap ");
220 }
221
222 #[tokio::test]
223 async fn test_mmap_directory_lazy_handle() {
224 let temp_dir = TempDir::new().unwrap();
225 let dir = MmapDirectory::new(temp_dir.path());
226
227 let data: Vec<u8> = (0..1000).map(|i| (i % 256) as u8).collect();
229 dir.write(Path::new("large.bin"), &data).await.unwrap();
230
231 let handle = dir.open_lazy(Path::new("large.bin")).await.unwrap();
233 assert_eq!(handle.len(), 1000);
234 assert!(handle.is_sync());
235
236 let range1 = handle.read_bytes_range(0..100).await.unwrap();
238 assert_eq!(range1.len(), 100);
239 assert_eq!(range1.as_slice(), &data[0..100]);
240
241 let range2 = handle.read_bytes_range_sync(500..600).unwrap();
243 assert_eq!(range2.as_slice(), &data[500..600]);
244 }
245}