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::metadata(&full)
252 .await
253 .map_err(|e| StorageError::Other(e))?;
254 if meta.is_dir() {
255 tokio::fs::remove_dir_all(&full)
256 .await
257 .map_err(|e| StorageError::Other(e))
258 } else {
259 tokio::fs::remove_file(&full)
260 .await
261 .map_err(|e| StorageError::Other(e))
262 }
263 }
264}
265
266static TMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
272
273fn temp_path_for(target: &Path) -> PathBuf {
274 let nanos = SystemTime::now()
275 .duration_since(UNIX_EPOCH)
276 .map(|d| d.as_nanos())
277 .unwrap_or(0);
278 let counter = TMP_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
279 let file_name = target
280 .file_name()
281 .map(|n| n.to_string_lossy().to_string())
282 .unwrap_or_else(|| "upload".to_string());
283 let tmp_name = format!(".libfw-tmp-{file_name}-{nanos}-{counter}");
284 target.with_file_name(tmp_name)
285}
286
287pub struct FsSink {
289 file: tokio::fs::File,
290 tmp: Option<PathBuf>,
291 target: PathBuf,
292 rel: String,
294 mode: WriteMode,
297 written: u64,
298}
299
300#[async_trait]
301impl UploadSink for FsSink {
302 async fn write(&mut self, buf: &[u8]) -> Result<(), StorageError> {
303 use tokio::io::AsyncWriteExt;
304 self.file
305 .write_all(buf)
306 .await
307 .map_err(|e| StorageError::write_failed(self.written, e))?;
308 self.written += buf.len() as u64;
309 Ok(())
310 }
311
312 async fn commit(self: Box<Self>) -> Result<FileMeta, StorageError> {
313 use tokio::io::AsyncWriteExt;
314 let FsSink {
315 mut file,
316 tmp,
317 target,
318 rel,
319 mode,
320 ..
321 } = *self;
322 file.flush().await.map_err(|e| StorageError::Other(e))?;
323 file.sync_all().await.map_err(|e| StorageError::Other(e))?;
324 drop(file);
325 if let Some(tmp) = tmp {
326 if matches!(mode, WriteMode::Create)
330 && tokio::fs::try_exists(&target)
331 .await
332 .map_err(|e| StorageError::Other(e))?
333 {
334 let _ = tokio::fs::remove_file(&tmp).await;
335 return Err(StorageError::AlreadyExists(rel));
336 }
337 tokio::fs::rename(&tmp, &target)
338 .await
339 .map_err(|e| StorageError::Other(e))?;
340 if let Some(parent) = target.parent() {
343 if let Ok(d) = tokio::fs::File::open(parent).await {
344 let _ = d.sync_all().await;
345 }
346 }
347 }
348 let meta = tokio::fs::metadata(&target)
349 .await
350 .map_err(|e| StorageError::Other(e))?;
351 Ok(file_meta_at(&rel, &meta))
352 }
353
354 async fn abort(self: Box<Self>) -> Result<(), StorageError> {
355 let FsSink { file, tmp, .. } = *self;
356 drop(file);
357 if let Some(tmp) = tmp {
358 let _ = tokio::fs::remove_file(&tmp).await;
359 }
360 Ok(())
361 }
362}
363
364#[cfg(test)]
365mod tests {
366 use super::*;
367 use libfw_core::storage::StorageBackend;
368
369 #[tokio::test]
370 async fn write_read_roundtrip() {
371 let dir = tempfile::tempdir().unwrap();
372 let storage = FsStorage::new(dir.path());
373
374 let sink = storage
375 .write_stream("a/b.txt", WriteMode::Create)
376 .await
377 .unwrap();
378 let mut sink = sink;
379 sink.write(b"hello ").await.unwrap();
380 sink.write(b"world").await.unwrap();
381 let meta = sink.commit().await.unwrap();
382 assert_eq!(meta.size, 11);
383
384 let got = storage.file_meta("a/b.txt").await.unwrap().unwrap();
385 assert_eq!(got.size, 11);
386
387 let mut reader = storage
388 .read_stream("a/b.txt", RangeSpec::new(0, 11))
389 .await
390 .unwrap();
391 let mut buf = Vec::new();
392 reader.read_to_end(&mut buf).unwrap();
393 assert_eq!(buf, b"hello world");
394 }
395
396 #[tokio::test]
397 async fn create_rejects_existing() {
398 let dir = tempfile::tempdir().unwrap();
399 let storage = FsStorage::new(dir.path());
400 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
401 let mut sink = sink;
402 sink.write(b"x").await.unwrap();
403 sink.commit().await.unwrap();
404
405 let res = storage.write_stream("f.txt", WriteMode::Create).await;
406 assert!(matches!(res, Err(StorageError::AlreadyExists(_))));
407 }
408
409 #[tokio::test]
410 async fn resume_writes_at_offset() {
411 let dir = tempfile::tempdir().unwrap();
412 let storage = FsStorage::new(dir.path());
413
414 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
415 let mut sink = sink;
416 sink.write(b"ABCD").await.unwrap();
417 sink.commit().await.unwrap();
418
419 let mut sink = storage
420 .write_stream("f.txt", WriteMode::Resume { offset: 4 })
421 .await
422 .unwrap();
423 sink.write(b"EF").await.unwrap();
424 sink.commit().await.unwrap();
425
426 let mut reader = storage
427 .read_stream("f.txt", RangeSpec::full(6))
428 .await
429 .unwrap();
430 let mut buf = Vec::new();
431 reader.read_to_end(&mut buf).unwrap();
432 assert_eq!(buf, b"ABCDEF");
433 }
434
435 #[tokio::test]
436 async fn resume_offset_mismatch_fails() {
437 let dir = tempfile::tempdir().unwrap();
438 let storage = FsStorage::new(dir.path());
439 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
440 let mut sink = sink;
441 sink.write(b"AB").await.unwrap();
442 sink.commit().await.unwrap();
443
444 let res = storage.write_stream("f.txt", WriteMode::Resume { offset: 9 }).await;
445 assert!(res.is_err());
446 }
447
448 #[tokio::test]
449 async fn abort_leaves_no_partial_file() {
450 let dir = tempfile::tempdir().unwrap();
451 let storage = FsStorage::new(dir.path());
452 let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
453 let mut sink = sink;
454 sink.write(b"partial").await.unwrap();
455 sink.abort().await.unwrap();
456 assert!(!dir.path().join("f.txt").exists());
457 }
458
459 #[tokio::test]
460 async fn list_dir_and_remove() {
461 let dir = tempfile::tempdir().unwrap();
462 let storage = FsStorage::new(dir.path());
463 std::fs::create_dir_all(dir.path().join("sub")).unwrap();
464 std::fs::write(dir.path().join("sub/x.txt"), b"1").unwrap();
465 std::fs::write(dir.path().join("a.txt"), b"22").unwrap();
466
467 let entries = storage.list_dir("").await.unwrap();
468 assert_eq!(entries.len(), 2);
469 let names: Vec<_> = entries.iter().map(|e| e.path.as_str()).collect();
470 assert_eq!(names, vec!["a.txt", "sub"]);
471
472 storage.remove("sub").await.unwrap();
473 assert!(!dir.path().join("sub").exists());
474 }
475
476 #[tokio::test]
477 async fn rejects_path_traversal() {
478 let dir = tempfile::tempdir().unwrap();
479 let storage = FsStorage::new(dir.path());
480 assert!(storage.file_meta("../etc/passwd").await.is_err());
481 assert!(storage.file_meta("/etc/passwd").await.is_err());
482 }
483
484 #[tokio::test]
485 async fn committed_meta_has_full_relative_path() {
486 let dir = tempfile::tempdir().unwrap();
487 let storage = FsStorage::new(dir.path());
488 let sink = storage
489 .write_stream("deep/nested/f.txt", WriteMode::Create)
490 .await
491 .unwrap();
492 let mut sink = sink;
493 sink.write(b"hi").await.unwrap();
494 let meta = sink.commit().await.unwrap();
495 assert_eq!(meta.path, "deep/nested/f.txt");
496 }
497
498 #[cfg(unix)]
499 #[tokio::test]
500 async fn rejects_symlink_traversal() {
501 use std::os::unix::fs::symlink;
502
503 let dir = tempfile::tempdir().unwrap();
504 let storage = FsStorage::new(dir.path());
505
506 let outside = tempfile::tempdir().unwrap();
509 std::fs::write(outside.path().join("secret.txt"), b"top-secret").unwrap();
510 symlink(outside.path(), dir.path().join("evil")).unwrap();
511
512 assert!(storage.read_stream("evil/secret.txt", RangeSpec::full(10)).await.is_err());
513 assert!(storage.file_meta("evil/secret.txt").await.is_err());
514 assert!(storage.write_stream("evil/new.txt", WriteMode::Create).await.is_err());
515 }
516}