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        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
266/// A temporary path for a target, unique per attempt.
267///
268/// Uniqueness combines a monotonic counter with a timestamp so two
269/// concurrent Create-mode uploads to the same target never collide on the
270/// same temp file (which would otherwise interleave/truncate each other).
271static 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
287/// Streaming write handle for the filesystem backend.
288pub struct FsSink {
289    file: tokio::fs::File,
290    tmp: Option<PathBuf>,
291    target: PathBuf,
292    /// The full virtual path (used for the committed `FileMeta`).
293    rel: String,
294    /// The mode the sink was opened with (used for the Create atomicity
295    /// check at commit time).
296    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            // Create mode must never clobber a target that appeared while
327            // we were streaming (closes the check-then-rename TOCTOU for
328            // the common non-concurrent-writer case).
329            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            // Durability: fsync the parent directory so the rename itself
341            // survives a crash.
342            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        // A symlink planted inside the root pointing outside it must not let
507        // reads/writes escape the mount root.
508        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}