1use std::io::{Read, SeekFrom};
8use std::path::{Component, Path, PathBuf};
9use std::time::{SystemTime, UNIX_EPOCH};
10
11use async_trait::async_trait;
12use libfw_core::metadata::{etag_from_size_mtime, FileMeta};
13use libfw_core::range::RangeSpec;
14use libfw_core::storage::{DirEntry, StorageBackend, UploadSink, WriteMode};
15use libfw_core::StorageError;
16
17#[derive(Debug, Clone)]
22pub struct FsStorage {
23 root: PathBuf,
24}
25
26impl FsStorage {
27 pub fn new(root: impl Into<PathBuf>) -> Self {
29 FsStorage { root: root.into() }
30 }
31
32 fn resolve(&self, path: &str) -> Result<PathBuf, StorageError> {
34 let rel = Path::new(path);
35 if rel.is_absolute() {
36 return Err(StorageError::Unsupported("absolute paths are not allowed"));
37 }
38 let mut joined = self.root.clone();
39 for component in rel.components() {
40 match component {
41 Component::Normal(seg) => joined.push(seg),
42 Component::CurDir => {}
43 _ => {
44 return Err(StorageError::Unsupported(
45 "path must not contain '..' or special components",
46 ))
47 }
48 }
49 }
50 Ok(joined)
51 }
52}
53
54fn file_meta_at(rel: &str, meta: &std::fs::Metadata) -> FileMeta {
55 let mtime = meta
56 .modified()
57 .ok()
58 .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
59 .map(|d| d.as_secs())
60 .unwrap_or(0);
61 FileMeta {
62 path: rel.to_string(),
63 size: meta.len(),
64 mtime,
65 etag: etag_from_size_mtime(meta.len(), mtime),
66 }
67}
68
69#[async_trait]
70impl StorageBackend for FsStorage {
71 async fn file_meta(&self, path: &str) -> Result<Option<FileMeta>, StorageError> {
72 let full = self.resolve(path)?;
73 let meta = match tokio::fs::metadata(&full).await {
74 Ok(m) => m,
75 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
76 Err(e) => return Err(StorageError::Other(e)),
77 };
78 if meta.is_dir() {
79 return Err(StorageError::Unsupported("path is a directory"));
80 }
81 Ok(Some(file_meta_at(path, &meta)))
82 }
83
84 async fn read_stream(
85 &self,
86 path: &str,
87 range: RangeSpec,
88 ) -> Result<Box<dyn Read + Send>, StorageError> {
89 let full = self.resolve(path)?;
90 let mut file = tokio::fs::File::open(&full)
91 .await
92 .map_err(|e| StorageError::Other(e))?;
93 if range.start > 0 {
94 tokio::io::AsyncSeekExt::seek(&mut file, SeekFrom::Start(range.start))
95 .await
96 .map_err(|e| StorageError::Other(e))?;
97 }
98 let std_file = file
99 .try_into_std()
100 .map_err(|e| StorageError::Other(std::io::Error::other(format!("{e:?}"))))?;
101 let limited = std_file.take(range.len());
103 Ok(Box::new(limited))
104 }
105
106 async fn write_stream(
107 &self,
108 path: &str,
109 mode: WriteMode,
110 ) -> Result<Box<dyn UploadSink>, StorageError> {
111 let full = self.resolve(path)?;
112 if let Some(parent) = full.parent() {
113 tokio::fs::create_dir_all(parent)
114 .await
115 .map_err(|e| StorageError::Other(e))?;
116 }
117
118 match mode {
119 WriteMode::Create | WriteMode::Overwrite => {
120 if mode == WriteMode::Create
121 && tokio::fs::try_exists(&full)
122 .await
123 .map_err(|e| StorageError::Other(e))?
124 {
125 return Err(StorageError::AlreadyExists(path.to_string()));
126 }
127 let tmp = temp_path_for(&full);
129 let file = tokio::fs::File::create(&tmp)
130 .await
131 .map_err(|e| StorageError::Other(e))?;
132 Ok(Box::new(FsSink {
133 file,
134 tmp: Some(tmp),
135 target: full,
136 written: 0,
137 }))
138 }
139 WriteMode::Resume { offset } => {
140 let file = tokio::fs::OpenOptions::new()
141 .write(true)
142 .append(true)
143 .open(&full)
144 .await
145 .map_err(|e| StorageError::Other(e))?;
146 let current = file
147 .metadata()
148 .await
149 .map_err(|e| StorageError::Other(e))?
150 .len();
151 if current != offset {
152 return Err(StorageError::write_failed(
153 offset,
154 std::io::Error::other(format!(
155 "existing file is {current} bytes, expected {offset}"
156 )),
157 ));
158 }
159 Ok(Box::new(FsSink {
160 file,
161 tmp: None,
162 target: full,
163 written: offset,
164 }))
165 }
166 }
167 }
168
169 async fn list_dir(&self, path: &str) -> Result<Vec<DirEntry>, StorageError> {
170 let full = if path.is_empty() {
171 self.root.clone()
172 } else {
173 self.resolve(path)?
174 };
175 let mut entries = Vec::new();
176 let mut read_dir = tokio::fs::read_dir(&full)
177 .await
178 .map_err(|e| StorageError::Other(e))?;
179 while let Some(entry) = read_dir
180 .next_entry()
181 .await
182 .map_err(|e| StorageError::Other(e))?
183 {
184 let name = entry.file_name().to_string_lossy().to_string();
185 let rel = if path.is_empty() {
186 name.clone()
187 } else {
188 format!("{path}/{name}")
189 };
190 let meta = entry
191 .metadata()
192 .await
193 .map_err(|e| StorageError::Other(e))?;
194 let mtime = meta
195 .modified()
196 .ok()
197 .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
198 .map(|d| d.as_secs())
199 .unwrap_or(0);
200 entries.push(DirEntry {
201 path: rel,
202 is_dir: meta.is_dir(),
203 size: if meta.is_dir() { 0 } else { meta.len() },
204 mtime,
205 });
206 }
207 entries.sort_by(|a, b| a.path.cmp(&b.path));
208 Ok(entries)
209 }
210
211 async fn mkdir_all(&self, path: &str) -> Result<(), StorageError> {
212 let full = self.resolve(path)?;
213 tokio::fs::create_dir_all(&full)
214 .await
215 .map_err(|e| StorageError::Other(e))?;
216 Ok(())
217 }
218
219 async fn remove(&self, path: &str) -> Result<(), StorageError> {
220 let full = self.resolve(path)?;
221 let meta = tokio::fs::metadata(&full)
222 .await
223 .map_err(|e| StorageError::Other(e))?;
224 if meta.is_dir() {
225 tokio::fs::remove_dir_all(&full)
226 .await
227 .map_err(|e| StorageError::Other(e))
228 } else {
229 tokio::fs::remove_file(&full)
230 .await
231 .map_err(|e| StorageError::Other(e))
232 }
233 }
234}
235
236fn temp_path_for(target: &Path) -> PathBuf {
238 let nanos = SystemTime::now()
239 .duration_since(UNIX_EPOCH)
240 .map(|d| d.as_nanos())
241 .unwrap_or(0);
242 let file_name = target
243 .file_name()
244 .map(|n| n.to_string_lossy().to_string())
245 .unwrap_or_else(|| "upload".to_string());
246 let tmp_name = format!(".libfw-tmp-{file_name}-{nanos}");
247 target.with_file_name(tmp_name)
248}
249
250pub struct FsSink {
252 file: tokio::fs::File,
253 tmp: Option<PathBuf>,
254 target: PathBuf,
255 written: u64,
256}
257
258#[async_trait]
259impl UploadSink for FsSink {
260 async fn write(&mut self, buf: &[u8]) -> Result<(), StorageError> {
261 use tokio::io::AsyncWriteExt;
262 self.file
263 .write_all(buf)
264 .await
265 .map_err(|e| StorageError::write_failed(self.written, e))?;
266 self.written += buf.len() as u64;
267 Ok(())
268 }
269
270 async fn commit(self: Box<Self>) -> Result<FileMeta, StorageError> {
271 use tokio::io::AsyncWriteExt;
272 let FsSink {
273 mut file,
274 tmp,
275 target,
276 ..
277 } = *self;
278 file.flush().await.map_err(|e| StorageError::Other(e))?;
279 file.sync_all().await.map_err(|e| StorageError::Other(e))?;
280 drop(file);
281 if let Some(tmp) = tmp {
282 tokio::fs::rename(&tmp, &target)
283 .await
284 .map_err(|e| StorageError::Other(e))?;
285 }
286 let meta = tokio::fs::metadata(&target)
287 .await
288 .map_err(|e| StorageError::Other(e))?;
289 let rel = target
290 .strip_prefix(&target.parent().map(Path::to_path_buf).unwrap_or_default())
291 .map(|p| p.to_string_lossy().to_string())
292 .unwrap_or_default();
293 Ok(file_meta_at(&rel, &meta))
294 }
295
296 async fn abort(self: Box<Self>) -> Result<(), StorageError> {
297 let FsSink { file, tmp, .. } = *self;
298 drop(file);
299 if let Some(tmp) = tmp {
300 let _ = tokio::fs::remove_file(&tmp).await;
301 }
302 Ok(())
303 }
304}
305
306#[cfg(test)]
307mod tests {
308 use super::*;
309 use libfw_core::storage::StorageBackend;
310
311 #[tokio::test]
312 async fn write_read_roundtrip() {
313 let dir = tempfile::tempdir().unwrap();
314 let storage = FsStorage::new(dir.path());
315
316 let sink = storage
317 .write_stream("a/b.txt", WriteMode::Create)
318 .await
319 .unwrap();
320 let mut sink = sink;
321 sink.write(b"hello ").await.unwrap();
322 sink.write(b"world").await.unwrap();
323 let meta = sink.commit().await.unwrap();
324 assert_eq!(meta.size, 11);
325
326 let got = storage.file_meta("a/b.txt").await.unwrap().unwrap();
327 assert_eq!(got.size, 11);
328
329 let mut reader = storage
330 .read_stream("a/b.txt", RangeSpec::new(0, 11))
331 .await
332 .unwrap();
333 let mut buf = Vec::new();
334 reader.read_to_end(&mut buf).unwrap();
335 assert_eq!(buf, b"hello world");
336 }
337
338 #[tokio::test]
339 async fn create_rejects_existing() {
340 let dir = tempfile::tempdir().unwrap();
341 let storage = FsStorage::new(dir.path());
342 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
343 let mut sink = sink;
344 sink.write(b"x").await.unwrap();
345 sink.commit().await.unwrap();
346
347 let res = storage.write_stream("f.txt", WriteMode::Create).await;
348 assert!(matches!(res, Err(StorageError::AlreadyExists(_))));
349 }
350
351 #[tokio::test]
352 async fn resume_writes_at_offset() {
353 let dir = tempfile::tempdir().unwrap();
354 let storage = FsStorage::new(dir.path());
355
356 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
357 let mut sink = sink;
358 sink.write(b"ABCD").await.unwrap();
359 sink.commit().await.unwrap();
360
361 let mut sink = storage
362 .write_stream("f.txt", WriteMode::Resume { offset: 4 })
363 .await
364 .unwrap();
365 sink.write(b"EF").await.unwrap();
366 sink.commit().await.unwrap();
367
368 let mut reader = storage
369 .read_stream("f.txt", RangeSpec::full(6))
370 .await
371 .unwrap();
372 let mut buf = Vec::new();
373 reader.read_to_end(&mut buf).unwrap();
374 assert_eq!(buf, b"ABCDEF");
375 }
376
377 #[tokio::test]
378 async fn resume_offset_mismatch_fails() {
379 let dir = tempfile::tempdir().unwrap();
380 let storage = FsStorage::new(dir.path());
381 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
382 let mut sink = sink;
383 sink.write(b"AB").await.unwrap();
384 sink.commit().await.unwrap();
385
386 let res = storage.write_stream("f.txt", WriteMode::Resume { offset: 9 }).await;
387 assert!(res.is_err());
388 }
389
390 #[tokio::test]
391 async fn abort_leaves_no_partial_file() {
392 let dir = tempfile::tempdir().unwrap();
393 let storage = FsStorage::new(dir.path());
394 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
395 let mut sink = sink;
396 sink.write(b"partial").await.unwrap();
397 sink.abort().await.unwrap();
398 assert!(!dir.path().join("f.txt").exists());
399 }
400
401 #[tokio::test]
402 async fn list_dir_and_remove() {
403 let dir = tempfile::tempdir().unwrap();
404 let storage = FsStorage::new(dir.path());
405 std::fs::create_dir_all(dir.path().join("sub")).unwrap();
406 std::fs::write(dir.path().join("sub/x.txt"), b"1").unwrap();
407 std::fs::write(dir.path().join("a.txt"), b"22").unwrap();
408
409 let entries = storage.list_dir("").await.unwrap();
410 assert_eq!(entries.len(), 2);
411 let names: Vec<_> = entries.iter().map(|e| e.path.as_str()).collect();
412 assert_eq!(names, vec!["a.txt", "sub"]);
413
414 storage.remove("sub").await.unwrap();
415 assert!(!dir.path().join("sub").exists());
416 }
417
418 #[tokio::test]
419 async fn rejects_path_traversal() {
420 let dir = tempfile::tempdir().unwrap();
421 let storage = FsStorage::new(dir.path());
422 assert!(storage.file_meta("../etc/passwd").await.is_err());
423 assert!(storage.file_meta("/etc/passwd").await.is_err());
424 }
425}