Skip to main content

libfw_server/
storage.rs

1//! Local filesystem storage backend for `libfw-server`.
2//!
3//! Writes go to a temporary file next to the destination and are
4//! atomically renamed into place on [`UploadSink::commit`], so a failed or
5//! aborted upload never leaves a partial target behind.
6
7use 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/// A [`StorageBackend`] rooted at a local directory.
18///
19/// Paths passed to the backend are treated as relative to `root`; any path
20/// escaping the root (absolute, `..`, symlink-traversing) is rejected.
21#[derive(Debug, Clone)]
22pub struct FsStorage {
23    root: PathBuf,
24}
25
26impl FsStorage {
27    /// Create a backend serving files under `root`.
28    pub fn new(root: impl Into<PathBuf>) -> Self {
29        FsStorage { root: root.into() }
30    }
31
32    /// Resolve a virtual path against the root, rejecting traversal.
33    ///
34    /// Both textual `..` segments and symlinked path components are
35    /// rejected so a read/write can never escape the mount root through a
36    /// symlink planted inside it.
37    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                    // Reject any component that is itself a symlink.
48                    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                        // A not-yet-existing component is fine (e.g. a new
56                        // upload target); it will be validated as it is
57                        // created component by component.
58                        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        // Restrict to exactly the requested range.
122        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                // Write to a temp file, rename on commit.
148                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        // Use `symlink_metadata` (not `metadata`) so a symlink that appeared
252        // since `resolve` checked is never followed — we refuse to remove
253        // *through* it. (The check-then-use race cannot be fully closed on
254        // all platforms without openat/O_NOFOLLOW; this narrows it for the
255        // destructive operation.)
256        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
276/// A temporary path for a target, unique per attempt.
277///
278/// Uniqueness combines a monotonic counter with a timestamp so two
279/// concurrent Create-mode uploads to the same target never collide on the
280/// same temp file (which would otherwise interleave/truncate each other).
281static 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
297/// Streaming write handle for the filesystem backend.
298pub struct FsSink {
299    file: tokio::fs::File,
300    tmp: Option<PathBuf>,
301    target: PathBuf,
302    /// The full virtual path (used for the committed `FileMeta`).
303    rel: String,
304    /// The mode the sink was opened with (used for the Create atomicity
305    /// check at commit time).
306    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            // Create mode must never clobber a target that appeared while
337            // we were streaming (closes the check-then-rename TOCTOU for
338            // the common non-concurrent-writer case).
339            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            // Durability: fsync the parent directory so the rename itself
351            // survives a crash.
352            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        // A symlink planted inside the root pointing outside it must not let
517        // reads/writes escape the mount root.
518        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}