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    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        // Restrict to exactly the requested range.
102        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                // Write to a temp file, rename on commit.
128                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
236/// A temporary path for a target, unique per attempt.
237fn 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
250/// Streaming write handle for the filesystem backend.
251pub 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}