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> {
38 let rel = Path::new(path);
39 if rel.is_absolute() {
40 return Err(StorageError::Unsupported("absolute paths are not allowed"));
41 }
42 let mut joined = self.root.clone();
43 for component in rel.components() {
44 match component {
45 Component::Normal(seg) => {
46 joined.push(seg);
47 match std::fs::symlink_metadata(&joined) {
49 Ok(m) if m.file_type().is_symlink() => {
50 return Err(StorageError::Unsupported(
51 "path must not traverse a symlink",
52 ))
53 }
54 Ok(_) => {}
55 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
59 Err(e) => return Err(StorageError::Other(e)),
60 }
61 }
62 Component::CurDir => {}
63 _ => {
64 return Err(StorageError::Unsupported(
65 "path must not contain '..' or special components",
66 ))
67 }
68 }
69 }
70 Ok(joined)
71 }
72}
73
74fn file_meta_at(rel: &str, meta: &std::fs::Metadata) -> FileMeta {
75 let mtime = meta
76 .modified()
77 .ok()
78 .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
79 .map(|d| d.as_secs())
80 .unwrap_or(0);
81 FileMeta {
82 path: rel.to_string(),
83 size: meta.len(),
84 mtime,
85 etag: etag_from_size_mtime(meta.len(), mtime),
86 }
87}
88
89#[async_trait]
90impl StorageBackend for FsStorage {
91 async fn file_meta(&self, path: &str) -> Result<Option<FileMeta>, StorageError> {
92 let full = self.resolve(path)?;
93 let meta = match tokio::fs::metadata(&full).await {
94 Ok(m) => m,
95 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
96 Err(e) => return Err(StorageError::Other(e)),
97 };
98 if meta.is_dir() {
99 return Err(StorageError::Unsupported("path is a directory"));
100 }
101 Ok(Some(file_meta_at(path, &meta)))
102 }
103
104 async fn read_stream(
105 &self,
106 path: &str,
107 range: RangeSpec,
108 ) -> Result<Box<dyn Read + Send>, StorageError> {
109 let full = self.resolve(path)?;
110 let mut file = tokio::fs::File::open(&full)
111 .await
112 .map_err(|e| StorageError::Other(e))?;
113 if range.start > 0 {
114 tokio::io::AsyncSeekExt::seek(&mut file, SeekFrom::Start(range.start))
115 .await
116 .map_err(|e| StorageError::Other(e))?;
117 }
118 let std_file = file
119 .try_into_std()
120 .map_err(|e| StorageError::Other(std::io::Error::other(format!("{e:?}"))))?;
121 let limited = std_file.take(range.len());
123 Ok(Box::new(limited))
124 }
125
126 async fn write_stream(
127 &self,
128 path: &str,
129 mode: WriteMode,
130 ) -> Result<Box<dyn UploadSink>, StorageError> {
131 let full = self.resolve(path)?;
132 if let Some(parent) = full.parent() {
133 tokio::fs::create_dir_all(parent)
134 .await
135 .map_err(|e| StorageError::Other(e))?;
136 }
137
138 match mode {
139 WriteMode::Create | WriteMode::Overwrite => {
140 if mode == WriteMode::Create
141 && tokio::fs::try_exists(&full)
142 .await
143 .map_err(|e| StorageError::Other(e))?
144 {
145 return Err(StorageError::AlreadyExists(path.to_string()));
146 }
147 let tmp = temp_path_for(&full);
149 let file = tokio::fs::File::create(&tmp)
150 .await
151 .map_err(|e| StorageError::Other(e))?;
152 Ok(Box::new(FsSink {
153 file,
154 tmp: Some(tmp),
155 target: full,
156 rel: path.to_string(),
157 mode,
158 written: 0,
159 }))
160 }
161 WriteMode::Resume { offset } => {
162 let file = tokio::fs::OpenOptions::new()
163 .write(true)
164 .append(true)
165 .open(&full)
166 .await
167 .map_err(|e| StorageError::Other(e))?;
168 let current = file
169 .metadata()
170 .await
171 .map_err(|e| StorageError::Other(e))?
172 .len();
173 if current != offset {
174 return Err(StorageError::write_failed(
175 offset,
176 std::io::Error::other(format!(
177 "existing file is {current} bytes, expected {offset}"
178 )),
179 ));
180 }
181 Ok(Box::new(FsSink {
182 file,
183 tmp: None,
184 target: full,
185 rel: path.to_string(),
186 mode,
187 written: offset,
188 }))
189 }
190 }
191 }
192
193 async fn list_dir(&self, path: &str) -> Result<Vec<DirEntry>, StorageError> {
194 let full = if path.is_empty() {
195 self.root.clone()
196 } else {
197 self.resolve(path)?
198 };
199 let mut entries = Vec::new();
200 let mut read_dir = tokio::fs::read_dir(&full)
201 .await
202 .map_err(|e| {
203 if e.kind() == std::io::ErrorKind::NotFound {
204 StorageError::NotFound(path.to_string())
205 } else {
206 StorageError::Other(e)
207 }
208 })?;
209 while let Some(entry) = read_dir
210 .next_entry()
211 .await
212 .map_err(|e| StorageError::Other(e))?
213 {
214 let name = entry.file_name().to_string_lossy().to_string();
215 let rel = if path.is_empty() {
216 name.clone()
217 } else {
218 format!("{path}/{name}")
219 };
220 let meta = entry
221 .metadata()
222 .await
223 .map_err(|e| StorageError::Other(e))?;
224 let mtime = meta
225 .modified()
226 .ok()
227 .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
228 .map(|d| d.as_secs())
229 .unwrap_or(0);
230 entries.push(DirEntry {
231 path: rel,
232 is_dir: meta.is_dir(),
233 size: if meta.is_dir() { 0 } else { meta.len() },
234 mtime,
235 });
236 }
237 entries.sort_by(|a, b| a.path.cmp(&b.path));
238 Ok(entries)
239 }
240
241 async fn mkdir_all(&self, path: &str) -> Result<(), StorageError> {
242 let full = self.resolve(path)?;
243 tokio::fs::create_dir_all(&full)
244 .await
245 .map_err(|e| StorageError::Other(e))?;
246 Ok(())
247 }
248
249 async fn remove(&self, path: &str) -> Result<(), StorageError> {
250 let full = self.resolve(path)?;
251 let meta = tokio::fs::symlink_metadata(&full)
257 .await
258 .map_err(|e| StorageError::Other(e))?;
259 if meta.file_type().is_symlink() {
260 return Err(StorageError::Unsupported(
261 "refusing to remove through a symlink",
262 ));
263 }
264 if meta.is_dir() {
265 tokio::fs::remove_dir_all(&full)
266 .await
267 .map_err(|e| StorageError::Other(e))
268 } else {
269 tokio::fs::remove_file(&full)
270 .await
271 .map_err(|e| StorageError::Other(e))
272 }
273 }
274}
275
276static TMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
282
283fn temp_path_for(target: &Path) -> PathBuf {
284 let nanos = SystemTime::now()
285 .duration_since(UNIX_EPOCH)
286 .map(|d| d.as_nanos())
287 .unwrap_or(0);
288 let counter = TMP_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
289 let file_name = target
290 .file_name()
291 .map(|n| n.to_string_lossy().to_string())
292 .unwrap_or_else(|| "upload".to_string());
293 let tmp_name = format!(".libfw-tmp-{file_name}-{nanos}-{counter}");
294 target.with_file_name(tmp_name)
295}
296
297pub struct FsSink {
299 file: tokio::fs::File,
300 tmp: Option<PathBuf>,
301 target: PathBuf,
302 rel: String,
304 mode: WriteMode,
307 written: u64,
308}
309
310#[async_trait]
311impl UploadSink for FsSink {
312 async fn write(&mut self, buf: &[u8]) -> Result<(), StorageError> {
313 use tokio::io::AsyncWriteExt;
314 self.file
315 .write_all(buf)
316 .await
317 .map_err(|e| StorageError::write_failed(self.written, e))?;
318 self.written += buf.len() as u64;
319 Ok(())
320 }
321
322 async fn commit(self: Box<Self>) -> Result<FileMeta, StorageError> {
323 use tokio::io::AsyncWriteExt;
324 let FsSink {
325 mut file,
326 tmp,
327 target,
328 rel,
329 mode,
330 ..
331 } = *self;
332 file.flush().await.map_err(|e| StorageError::Other(e))?;
333 file.sync_all().await.map_err(|e| StorageError::Other(e))?;
334 drop(file);
335 if let Some(tmp) = tmp {
336 if matches!(mode, WriteMode::Create)
340 && tokio::fs::try_exists(&target)
341 .await
342 .map_err(|e| StorageError::Other(e))?
343 {
344 let _ = tokio::fs::remove_file(&tmp).await;
345 return Err(StorageError::AlreadyExists(rel));
346 }
347 tokio::fs::rename(&tmp, &target)
348 .await
349 .map_err(|e| StorageError::Other(e))?;
350 if let Some(parent) = target.parent() {
353 if let Ok(d) = tokio::fs::File::open(parent).await {
354 let _ = d.sync_all().await;
355 }
356 }
357 }
358 let meta = tokio::fs::metadata(&target)
359 .await
360 .map_err(|e| StorageError::Other(e))?;
361 Ok(file_meta_at(&rel, &meta))
362 }
363
364 async fn abort(self: Box<Self>) -> Result<(), StorageError> {
365 let FsSink { file, tmp, .. } = *self;
366 drop(file);
367 if let Some(tmp) = tmp {
368 let _ = tokio::fs::remove_file(&tmp).await;
369 }
370 Ok(())
371 }
372}
373
374#[cfg(test)]
375mod tests {
376 use super::*;
377 use libfw_core::storage::StorageBackend;
378
379 #[tokio::test]
380 async fn write_read_roundtrip() {
381 let dir = tempfile::tempdir().unwrap();
382 let storage = FsStorage::new(dir.path());
383
384 let sink = storage
385 .write_stream("a/b.txt", WriteMode::Create)
386 .await
387 .unwrap();
388 let mut sink = sink;
389 sink.write(b"hello ").await.unwrap();
390 sink.write(b"world").await.unwrap();
391 let meta = sink.commit().await.unwrap();
392 assert_eq!(meta.size, 11);
393
394 let got = storage.file_meta("a/b.txt").await.unwrap().unwrap();
395 assert_eq!(got.size, 11);
396
397 let mut reader = storage
398 .read_stream("a/b.txt", RangeSpec::new(0, 11))
399 .await
400 .unwrap();
401 let mut buf = Vec::new();
402 reader.read_to_end(&mut buf).unwrap();
403 assert_eq!(buf, b"hello world");
404 }
405
406 #[tokio::test]
407 async fn create_rejects_existing() {
408 let dir = tempfile::tempdir().unwrap();
409 let storage = FsStorage::new(dir.path());
410 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
411 let mut sink = sink;
412 sink.write(b"x").await.unwrap();
413 sink.commit().await.unwrap();
414
415 let res = storage.write_stream("f.txt", WriteMode::Create).await;
416 assert!(matches!(res, Err(StorageError::AlreadyExists(_))));
417 }
418
419 #[tokio::test]
420 async fn resume_writes_at_offset() {
421 let dir = tempfile::tempdir().unwrap();
422 let storage = FsStorage::new(dir.path());
423
424 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
425 let mut sink = sink;
426 sink.write(b"ABCD").await.unwrap();
427 sink.commit().await.unwrap();
428
429 let mut sink = storage
430 .write_stream("f.txt", WriteMode::Resume { offset: 4 })
431 .await
432 .unwrap();
433 sink.write(b"EF").await.unwrap();
434 sink.commit().await.unwrap();
435
436 let mut reader = storage
437 .read_stream("f.txt", RangeSpec::full(6))
438 .await
439 .unwrap();
440 let mut buf = Vec::new();
441 reader.read_to_end(&mut buf).unwrap();
442 assert_eq!(buf, b"ABCDEF");
443 }
444
445 #[tokio::test]
446 async fn resume_offset_mismatch_fails() {
447 let dir = tempfile::tempdir().unwrap();
448 let storage = FsStorage::new(dir.path());
449 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
450 let mut sink = sink;
451 sink.write(b"AB").await.unwrap();
452 sink.commit().await.unwrap();
453
454 let res = storage.write_stream("f.txt", WriteMode::Resume { offset: 9 }).await;
455 assert!(res.is_err());
456 }
457
458 #[tokio::test]
459 async fn abort_leaves_no_partial_file() {
460 let dir = tempfile::tempdir().unwrap();
461 let storage = FsStorage::new(dir.path());
462 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
463 let mut sink = sink;
464 sink.write(b"partial").await.unwrap();
465 sink.abort().await.unwrap();
466 assert!(!dir.path().join("f.txt").exists());
467 }
468
469 #[tokio::test]
470 async fn list_dir_and_remove() {
471 let dir = tempfile::tempdir().unwrap();
472 let storage = FsStorage::new(dir.path());
473 std::fs::create_dir_all(dir.path().join("sub")).unwrap();
474 std::fs::write(dir.path().join("sub/x.txt"), b"1").unwrap();
475 std::fs::write(dir.path().join("a.txt"), b"22").unwrap();
476
477 let entries = storage.list_dir("").await.unwrap();
478 assert_eq!(entries.len(), 2);
479 let names: Vec<_> = entries.iter().map(|e| e.path.as_str()).collect();
480 assert_eq!(names, vec!["a.txt", "sub"]);
481
482 storage.remove("sub").await.unwrap();
483 assert!(!dir.path().join("sub").exists());
484 }
485
486 #[tokio::test]
487 async fn rejects_path_traversal() {
488 let dir = tempfile::tempdir().unwrap();
489 let storage = FsStorage::new(dir.path());
490 assert!(storage.file_meta("../etc/passwd").await.is_err());
491 assert!(storage.file_meta("/etc/passwd").await.is_err());
492 }
493
494 #[tokio::test]
495 async fn committed_meta_has_full_relative_path() {
496 let dir = tempfile::tempdir().unwrap();
497 let storage = FsStorage::new(dir.path());
498 let sink = storage
499 .write_stream("deep/nested/f.txt", WriteMode::Create)
500 .await
501 .unwrap();
502 let mut sink = sink;
503 sink.write(b"hi").await.unwrap();
504 let meta = sink.commit().await.unwrap();
505 assert_eq!(meta.path, "deep/nested/f.txt");
506 }
507
508 #[cfg(unix)]
509 #[tokio::test]
510 async fn rejects_symlink_traversal() {
511 use std::os::unix::fs::symlink;
512
513 let dir = tempfile::tempdir().unwrap();
514 let storage = FsStorage::new(dir.path());
515
516 let outside = tempfile::tempdir().unwrap();
519 std::fs::write(outside.path().join("secret.txt"), b"top-secret").unwrap();
520 symlink(outside.path(), dir.path().join("evil")).unwrap();
521
522 assert!(storage.read_stream("evil/secret.txt", RangeSpec::full(10)).await.is_err());
523 assert!(storage.file_meta("evil/secret.txt").await.is_err());
524 assert!(storage.write_stream("evil/new.txt", WriteMode::Create).await.is_err());
525 }
526}