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, ChunkRange, 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                    blocks_path: None,
160                    ranges: Vec::new(),
161                }))
162            }
163            WriteMode::Resume { offset } => {
164                let file = tokio::fs::OpenOptions::new()
165                    .write(true)
166                    .append(true)
167                    .open(&full)
168                    .await
169                    .map_err(|e| StorageError::Other(e))?;
170                let current = file
171                    .metadata()
172                    .await
173                    .map_err(|e| StorageError::Other(e))?
174                    .len();
175                if current != offset {
176                    return Err(StorageError::write_failed(
177                        offset,
178                        std::io::Error::other(format!(
179                            "existing file is {current} bytes, expected {offset}"
180                        )),
181                    ));
182                }
183                Ok(Box::new(FsSink {
184                    file,
185                    tmp: None,
186                    target: full,
187                    rel: path.to_string(),
188                    mode,
189                    written: offset,
190                    blocks_path: None,
191                    ranges: Vec::new(),
192                }))
193            }
194        }
195    }
196
197    async fn write_stream_session(
198        &self,
199        path: &str,
200        session: &str,
201        mode: WriteMode,
202    ) -> Result<Box<dyn UploadSink>, StorageError> {
203        let full = self.resolve(path)?;
204        if let Some(parent) = full.parent() {
205            tokio::fs::create_dir_all(parent)
206                .await
207                .map_err(|e| StorageError::Other(e))?;
208        }
209        // The session string is embedded in a temp filename, so it must never
210        // be able to inject path separators or `..` (a malicious client could
211        // otherwise write outside the mount root). Restrict to safe chars.
212        let safe: String = session
213            .chars()
214            .map(|c| {
215                if c.is_ascii_alphanumeric() || c == '-' || c == '_' {
216                    c
217                } else {
218                    '_'
219                }
220            })
221            .collect();
222        let name = full
223            .file_name()
224            .map(|n| n.to_string_lossy().to_string())
225            .unwrap_or_else(|| "upload".to_string());
226        let tmp = full.with_file_name(format!(".libfw-sess-{safe}-{name}"));
227
228        let exists = tokio::fs::try_exists(&tmp)
229            .await
230            .map_err(|e| StorageError::Other(e))?;
231        let file = if exists {
232            // Subsequent chunk / resume of an in-flight session: open the
233            // shared temp for positional (seek + write) access. `mode` is
234            // ignored; the already-received ranges are reloaded from the
235            // sidecar so the client can resume only the missing parts.
236            tokio::fs::OpenOptions::new()
237                .read(true)
238                .write(true)
239                .open(&tmp)
240                .await
241                .map_err(|e| StorageError::Other(e))?
242        } else {
243            // First request for this session: create the shared temp.
244            match mode {
245                WriteMode::Create | WriteMode::Overwrite => {
246                    if mode == WriteMode::Create
247                        && tokio::fs::try_exists(&full)
248                            .await
249                            .map_err(|e| StorageError::Other(e))?
250                    {
251                        return Err(StorageError::AlreadyExists(path.to_string()));
252                    }
253                    tokio::fs::OpenOptions::new()
254                        .write(true)
255                        .create_new(true)
256                        .open(&tmp)
257                        .await
258                        .map_err(|e| StorageError::Other(e))?
259                }
260                // Resumable sessions are driven by the per-block probe (the
261                // client asks "what ranges do you have?" and sends only the
262                // gaps), not by a contiguous offset append, so a legacy
263                // `Resume` mode is not applicable here.
264                WriteMode::Resume { .. } => {
265                    return Err(StorageError::Unsupported(
266                        "session upload does not support contiguous resume; use the block probe".into(),
267                    ))
268                }
269            }
270        };
271        // Load any already-received byte ranges from the sidecar so a pause /
272        // resume only re-sends the missing blocks.
273        let blocks_path = blocks_path_for(&tmp);
274        let ranges = read_ranges(&blocks_path).await;
275        Ok(Box::new(FsSink {
276            file,
277            tmp: Some(tmp),
278            target: full,
279            rel: path.to_string(),
280            mode,
281            written: 0,
282            blocks_path: Some(blocks_path),
283            ranges,
284        }))
285    }
286
287    async fn list_dir(&self, path: &str) -> Result<Vec<DirEntry>, StorageError> {
288        let full = if path.is_empty() {
289            self.root.clone()
290        } else {
291            self.resolve(path)?
292        };
293        let mut entries = Vec::new();
294        let mut read_dir = tokio::fs::read_dir(&full)
295            .await
296            .map_err(|e| {
297                if e.kind() == std::io::ErrorKind::NotFound {
298                    StorageError::NotFound(path.to_string())
299                } else {
300                    StorageError::Other(e)
301                }
302            })?;
303        while let Some(entry) = read_dir
304            .next_entry()
305            .await
306            .map_err(|e| StorageError::Other(e))?
307        {
308            let name = entry.file_name().to_string_lossy().to_string();
309            let rel = if path.is_empty() {
310                name.clone()
311            } else {
312                format!("{path}/{name}")
313            };
314            let meta = entry
315                .metadata()
316                .await
317                .map_err(|e| StorageError::Other(e))?;
318            let mtime = meta
319                .modified()
320                .ok()
321                .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
322                .map(|d| d.as_secs())
323                .unwrap_or(0);
324            entries.push(DirEntry {
325                path: rel,
326                is_dir: meta.is_dir(),
327                size: if meta.is_dir() { 0 } else { meta.len() },
328                mtime,
329            });
330        }
331        entries.sort_by(|a, b| a.path.cmp(&b.path));
332        Ok(entries)
333    }
334
335    async fn mkdir_all(&self, path: &str) -> Result<(), StorageError> {
336        let full = self.resolve(path)?;
337        tokio::fs::create_dir_all(&full)
338            .await
339            .map_err(|e| StorageError::Other(e))?;
340        Ok(())
341    }
342
343    async fn remove(&self, path: &str) -> Result<(), StorageError> {
344        let full = self.resolve(path)?;
345        // Use `symlink_metadata` (not `metadata`) so a symlink that appeared
346        // since `resolve` checked is never followed — we refuse to remove
347        // *through* it. (The check-then-use race cannot be fully closed on
348        // all platforms without openat/O_NOFOLLOW; this narrows it for the
349        // destructive operation.)
350        let meta = tokio::fs::symlink_metadata(&full)
351            .await
352            .map_err(|e| StorageError::Other(e))?;
353        if meta.file_type().is_symlink() {
354            return Err(StorageError::Unsupported(
355                "refusing to remove through a symlink",
356            ));
357        }
358        if meta.is_dir() {
359            tokio::fs::remove_dir_all(&full)
360                .await
361                .map_err(|e| StorageError::Other(e))
362        } else {
363            tokio::fs::remove_file(&full)
364                .await
365                .map_err(|e| StorageError::Other(e))
366        }
367    }
368
369    async fn cleanup_stale_sessions(
370        &self,
371        max_age: std::time::Duration,
372    ) -> Result<usize, StorageError> {
373        // tus-style expiry: a client that vanishes mid-upload leaves its
374        // `.libfw-sess-<id>-<name>` temp (and `.blocks` sidecar) behind. Walk
375        // the root, removing the ones whose last write is older than
376        // `max_age`. Only session temps are touched — never committed user
377        // files — and symlinked directories are never followed.
378        let deadline = SystemTime::now()
379            .checked_sub(max_age)
380            .unwrap_or(UNIX_EPOCH);
381        let mut removed = 0usize;
382        let mut stack = vec![self.root.clone()];
383        while let Some(dir) = stack.pop() {
384            let mut rd = match tokio::fs::read_dir(&dir).await {
385                Ok(rd) => rd,
386                Err(_) => continue,
387            };
388            while let Ok(Some(entry)) = rd.next_entry().await {
389                let ft = match entry.file_type().await {
390                    Ok(ft) => ft,
391                    Err(_) => continue,
392                };
393                if ft.is_dir() {
394                    if !ft.is_symlink() {
395                        stack.push(entry.path());
396                    }
397                    continue;
398                }
399                if ft.is_symlink() {
400                    continue;
401                }
402                let name = entry.file_name().to_string_lossy().to_string();
403                if !name.starts_with(".libfw-sess-") {
404                    continue;
405                }
406                let modified = entry
407                    .metadata()
408                    .await
409                    .ok()
410                    .and_then(|m| m.modified().ok())
411                    .and_then(|t| t.duration_since(UNIX_EPOCH).ok())
412                    .map(|d| UNIX_EPOCH + d)
413                    .unwrap_or(UNIX_EPOCH);
414                if modified < deadline {
415                    let _ = tokio::fs::remove_file(entry.path()).await;
416                    // Remove the parallel range sidecar (`<temp>.blocks`).
417                    let mut sidecar = entry.path().as_os_str().to_owned();
418                    sidecar.push(".blocks");
419                    let _ = tokio::fs::remove_file(PathBuf::from(sidecar)).await;
420                    removed += 1;
421                }
422            }
423        }
424        Ok(removed)
425    }
426}
427
428/// A temporary path for a target, unique per attempt.
429///
430/// Uniqueness combines a monotonic counter with a timestamp so two
431/// concurrent Create-mode uploads to the same target never collide on the
432/// same temp file (which would otherwise interleave/truncate each other).
433static TMP_COUNTER: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
434
435fn temp_path_for(target: &Path) -> PathBuf {
436    let nanos = SystemTime::now()
437        .duration_since(UNIX_EPOCH)
438        .map(|d| d.as_nanos())
439        .unwrap_or(0);
440    let counter = TMP_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
441    let file_name = target
442        .file_name()
443        .map(|n| n.to_string_lossy().to_string())
444        .unwrap_or_else(|| "upload".to_string());
445    let tmp_name = format!(".libfw-tmp-{file_name}-{nanos}-{counter}");
446    target.with_file_name(tmp_name)
447}
448
449/// Streaming write handle for the filesystem backend.
450pub struct FsSink {
451    file: tokio::fs::File,
452    tmp: Option<PathBuf>,
453    target: PathBuf,
454    /// The full virtual path (used for the committed `FileMeta`).
455    rel: String,
456    /// The mode the sink was opened with (used for the Create atomicity
457    /// check at commit time).
458    mode: WriteMode,
459    written: u64,
460    /// Optional sidecar path tracking received byte ranges for a resumable
461    /// "session" upload (parallel to the session temp file). `None` for
462    /// ordinary (non-session) sinks.
463    blocks_path: Option<PathBuf>,
464    /// In-memory copy of the received ranges (kept in sync with
465    /// `blocks_path`).
466    ranges: Vec<ChunkRange>,
467}
468
469/// Merge `new` (a `[start, end)` half-open range) into a sorted, disjoint
470/// list of ranges, coalescing overlaps/adjacencies. Returns the updated list.
471fn merge_range(ranges: &mut Vec<ChunkRange>, new: ChunkRange) {
472    if new.is_empty() {
473        return;
474    }
475    ranges.push(new);
476    ranges.sort_by_key(|r| r.start);
477    let mut merged: Vec<ChunkRange> = Vec::with_capacity(ranges.len());
478    for r in ranges.drain(..) {
479        if let Some(last) = merged.last_mut() {
480            // Overlapping or adjacent ranges coalesce.
481            if r.start <= last.end {
482                if r.end > last.end {
483                    last.end = r.end;
484                }
485                continue;
486            }
487        }
488        merged.push(r);
489    }
490    *ranges = merged;
491}
492
493#[async_trait]
494impl UploadSink for FsSink {
495    async fn write(&mut self, buf: &[u8]) -> Result<(), StorageError> {
496        use tokio::io::AsyncWriteExt;
497        self.file
498            .write_all(buf)
499            .await
500            .map_err(|e| StorageError::write_failed(self.written, e))?;
501        self.written += buf.len() as u64;
502        Ok(())
503    }
504
505    async fn write_at(&mut self, offset: u64, buf: &[u8]) -> Result<(), StorageError> {
506        use tokio::io::{AsyncSeekExt, AsyncWriteExt};
507        self.file
508            .seek(SeekFrom::Start(offset))
509            .await
510            .map_err(|e| StorageError::write_failed(offset, e))?;
511        self.file
512            .write_all(buf)
513            .await
514            .map_err(|e| StorageError::write_failed(offset, e))?;
515        // Track the highest extent written (used by `len`-style bookkeeping;
516        // the real size for commit validation comes from file metadata).
517        self.written = self.written.max(offset.saturating_add(buf.len() as u64));
518        // For a resumable session sink, record the received byte range in the
519        // sidecar so a later probe / resume knows this part is already on
520        // disk and only missing gaps need to be re-sent.
521        if let Some(blocks) = self.blocks_path.clone() {
522            merge_range(&mut self.ranges, ChunkRange {
523                start: offset,
524                end: offset.saturating_add(buf.len() as u64),
525            });
526            persist_ranges(&blocks, &self.ranges).await?;
527        }
528        Ok(())
529    }
530
531    async fn received_ranges(&mut self) -> Result<Vec<ChunkRange>, StorageError> {
532        Ok(self.ranges.clone())
533    }
534
535    async fn len(&self) -> Result<u64, StorageError> {
536        self.file
537            .metadata()
538            .await
539            .map(|m| m.len())
540            .map_err(|e| StorageError::Other(e))
541    }
542
543    async fn commit(self: Box<Self>) -> Result<FileMeta, StorageError> {
544        use tokio::io::AsyncWriteExt;
545        let FsSink {
546            mut file,
547            tmp,
548            target,
549            rel,
550            mode,
551            blocks_path,
552            ..
553        } = *self;
554        let _ = remove_blocks_sidecar(blocks_path.as_deref()).await;
555        file.flush().await.map_err(|e| StorageError::Other(e))?;
556        file.sync_all().await.map_err(|e| StorageError::Other(e))?;
557        drop(file);
558        if let Some(tmp) = tmp {
559            // Create mode must never clobber a target that appeared while
560            // we were streaming (closes the check-then-rename TOCTOU for
561            // the common non-concurrent-writer case).
562            if matches!(mode, WriteMode::Create)
563                && tokio::fs::try_exists(&target)
564                    .await
565                    .map_err(|e| StorageError::Other(e))?
566            {
567                let _ = tokio::fs::remove_file(&tmp).await;
568                return Err(StorageError::AlreadyExists(rel));
569            }
570            tokio::fs::rename(&tmp, &target)
571                .await
572                .map_err(|e| StorageError::Other(e))?;
573            // Durability: fsync the parent directory so the rename itself
574            // survives a crash.
575            if let Some(parent) = target.parent() {
576                if let Ok(d) = tokio::fs::File::open(parent).await {
577                    let _ = d.sync_all().await;
578                }
579            }
580        }
581        let meta = tokio::fs::metadata(&target)
582            .await
583            .map_err(|e| StorageError::Other(e))?;
584        Ok(file_meta_at(&rel, &meta))
585    }
586
587    async fn abort(self: Box<Self>) -> Result<(), StorageError> {
588        let FsSink {
589            file,
590            tmp,
591            blocks_path,
592            ..
593        } = *self;
594        drop(file);
595        if let Some(tmp) = tmp {
596            let _ = tokio::fs::remove_file(&tmp).await;
597        }
598        let _ = remove_blocks_sidecar(blocks_path.as_deref()).await;
599        Ok(())
600    }
601}
602
603/// Sidecar filename for a session temp (e.g. `<temp>.blocks`).
604fn blocks_path_for(tmp: &Path) -> PathBuf {
605    let mut name = tmp.as_os_str().to_owned();
606    name.push(".blocks");
607    PathBuf::from(name)
608}
609
610/// Read the persisted received byte ranges for a session temp sidecar.
611async fn read_ranges(blocks: &Path) -> Vec<ChunkRange> {
612    match tokio::fs::read_to_string(blocks).await {
613        Ok(text) => serde_json::from_str::<Vec<ChunkRange>>(&text).unwrap_or_default(),
614        Err(_) => Vec::new(),
615    }
616}
617
618/// Persist the received byte ranges to a session temp sidecar (best-effort).
619async fn persist_ranges(blocks: &Path, ranges: &[ChunkRange]) -> Result<(), StorageError> {
620    let text = serde_json::to_string(ranges).unwrap_or_else(|_| "[]".to_string());
621    tokio::fs::write(blocks, text)
622        .await
623        .map_err(|e| StorageError::Other(e))
624}
625
626/// Remove a session temp sidecar (best-effort; missing file is fine).
627async fn remove_blocks_sidecar(blocks: Option<&Path>) -> Result<(), StorageError> {
628    if let Some(blocks) = blocks {
629        let _ = tokio::fs::remove_file(blocks).await;
630    }
631    Ok(())
632}
633
634#[cfg(test)]
635mod tests {
636    use super::*;
637    use libfw_core::storage::StorageBackend;
638
639    /// Force a file's mtime (used to age a fake session temp).
640    fn filetime_set(path: &Path, time: std::time::SystemTime) -> std::io::Result<()> {
641        // Open writable: on Windows, setting file times through a read-only
642        // handle is refused.
643        let file = std::fs::OpenOptions::new().write(true).open(path)?;
644        file.set_times(std::fs::FileTimes::new().set_modified(time))
645    }
646
647    #[tokio::test]
648    async fn write_read_roundtrip() {
649        let dir = tempfile::tempdir().unwrap();
650        let storage = FsStorage::new(dir.path());
651
652        let sink = storage
653            .write_stream("a/b.txt", WriteMode::Create)
654            .await
655            .unwrap();
656        let mut sink = sink;
657        sink.write(b"hello ").await.unwrap();
658        sink.write(b"world").await.unwrap();
659        let meta = sink.commit().await.unwrap();
660        assert_eq!(meta.size, 11);
661
662        let got = storage.file_meta("a/b.txt").await.unwrap().unwrap();
663        assert_eq!(got.size, 11);
664
665        let mut reader = storage
666            .read_stream("a/b.txt", RangeSpec::new(0, 11))
667            .await
668            .unwrap();
669        let mut buf = Vec::new();
670        reader.read_to_end(&mut buf).unwrap();
671        assert_eq!(buf, b"hello world");
672    }
673
674    #[tokio::test]
675    async fn create_rejects_existing() {
676        let dir = tempfile::tempdir().unwrap();
677        let storage = FsStorage::new(dir.path());
678        let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
679        let mut sink = sink;
680        sink.write(b"x").await.unwrap();
681        sink.commit().await.unwrap();
682
683        let res = storage.write_stream("f.txt", WriteMode::Create).await;
684        assert!(matches!(res, Err(StorageError::AlreadyExists(_))));
685    }
686
687    #[tokio::test]
688    async fn resume_writes_at_offset() {
689        let dir = tempfile::tempdir().unwrap();
690        let storage = FsStorage::new(dir.path());
691
692        let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
693        let mut sink = sink;
694        sink.write(b"ABCD").await.unwrap();
695        sink.commit().await.unwrap();
696
697        let mut sink = storage
698            .write_stream("f.txt", WriteMode::Resume { offset: 4 })
699            .await
700            .unwrap();
701        sink.write(b"EF").await.unwrap();
702        sink.commit().await.unwrap();
703
704        let mut reader = storage
705            .read_stream("f.txt", RangeSpec::full(6))
706            .await
707            .unwrap();
708        let mut buf = Vec::new();
709        reader.read_to_end(&mut buf).unwrap();
710        assert_eq!(buf, b"ABCDEF");
711    }
712
713    #[tokio::test]
714    async fn resume_offset_mismatch_fails() {
715        let dir = tempfile::tempdir().unwrap();
716        let storage = FsStorage::new(dir.path());
717        let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
718        let mut sink = sink;
719        sink.write(b"AB").await.unwrap();
720        sink.commit().await.unwrap();
721
722        let res = storage.write_stream("f.txt", WriteMode::Resume { offset: 9 }).await;
723        assert!(res.is_err());
724    }
725
726    #[tokio::test]
727    async fn abort_leaves_no_partial_file() {
728        let dir = tempfile::tempdir().unwrap();
729        let storage = FsStorage::new(dir.path());
730        let sink = storage.write_stream("f.txt", WriteMode::Create).await.unwrap();
731        let mut sink = sink;
732        sink.write(b"partial").await.unwrap();
733        sink.abort().await.unwrap();
734        assert!(!dir.path().join("f.txt").exists());
735    }
736
737    #[tokio::test]
738    async fn list_dir_and_remove() {
739        let dir = tempfile::tempdir().unwrap();
740        let storage = FsStorage::new(dir.path());
741        std::fs::create_dir_all(dir.path().join("sub")).unwrap();
742        std::fs::write(dir.path().join("sub/x.txt"), b"1").unwrap();
743        std::fs::write(dir.path().join("a.txt"), b"22").unwrap();
744
745        let entries = storage.list_dir("").await.unwrap();
746        assert_eq!(entries.len(), 2);
747        let names: Vec<_> = entries.iter().map(|e| e.path.as_str()).collect();
748        assert_eq!(names, vec!["a.txt", "sub"]);
749
750        storage.remove("sub").await.unwrap();
751        assert!(!dir.path().join("sub").exists());
752    }
753
754    #[tokio::test]
755    async fn rejects_path_traversal() {
756        let dir = tempfile::tempdir().unwrap();
757        let storage = FsStorage::new(dir.path());
758        assert!(storage.file_meta("../etc/passwd").await.is_err());
759        assert!(storage.file_meta("/etc/passwd").await.is_err());
760    }
761
762    #[tokio::test]
763    async fn cleanup_stale_sessions_removes_old_temps_only() {
764        let dir = tempfile::tempdir().unwrap();
765        let storage = FsStorage::new(dir.path());
766
767        // A stale session temp (very old mtime) + its range sidecar.
768        let stale = dir.path().join(".libfw-sess-oldid-a.bin");
769        std::fs::write(&stale, b"partial").unwrap();
770        std::fs::write(dir.path().join(".libfw-sess-oldid-a.bin.blocks"), b"[[0,7]]").unwrap();
771        // Age it well beyond the 1-hour TTL (7 days old). Windows/FAT can
772        // clamp very old timestamps, so use a recent-but-stale date.
773        let old = std::time::SystemTime::now()
774            .checked_sub(std::time::Duration::from_secs(7 * 24 * 3600))
775            .unwrap();
776        assert!(
777            filetime_set(&stale, old).is_ok(),
778            "failed to age the stale session temp"
779        );
780
781        // A fresh session temp (now) must be kept.
782        let fresh = dir.path().join(".libfw-sess-newid-a.bin");
783        std::fs::write(&fresh, b"partial").unwrap();
784
785        // A committed user file must never be touched.
786        std::fs::write(dir.path().join("real.txt"), b"real").unwrap();
787
788        let removed = storage
789            .cleanup_stale_sessions(std::time::Duration::from_secs(3600))
790            .await
791            .unwrap();
792        assert_eq!(removed, 1, "only the stale temp is removed");
793        assert!(!stale.exists());
794        assert!(!dir.path().join(".libfw-sess-oldid-a.bin.blocks").exists());
795        assert!(fresh.exists(), "fresh temp survives");
796        assert!(dir.path().join("real.txt").exists(), "user file survives");
797    }
798
799    #[tokio::test]
800    async fn committed_meta_has_full_relative_path() {
801        let dir = tempfile::tempdir().unwrap();
802        let storage = FsStorage::new(dir.path());
803        let sink = storage
804            .write_stream("deep/nested/f.txt", WriteMode::Create)
805            .await
806            .unwrap();
807        let mut sink = sink;
808        sink.write(b"hi").await.unwrap();
809        let meta = sink.commit().await.unwrap();
810        assert_eq!(meta.path, "deep/nested/f.txt");
811    }
812
813    #[cfg(unix)]
814    #[tokio::test]
815    async fn rejects_symlink_traversal() {
816        use std::os::unix::fs::symlink;
817
818        let dir = tempfile::tempdir().unwrap();
819        let storage = FsStorage::new(dir.path());
820
821        // A symlink planted inside the root pointing outside it must not let
822        // reads/writes escape the mount root.
823        let outside = tempfile::tempdir().unwrap();
824        std::fs::write(outside.path().join("secret.txt"), b"top-secret").unwrap();
825        symlink(outside.path(), dir.path().join("evil")).unwrap();
826
827        assert!(storage.read_stream("evil/secret.txt", RangeSpec::full(10)).await.is_err());
828        assert!(storage.file_meta("evil/secret.txt").await.is_err());
829        assert!(storage.write_stream("evil/new.txt", WriteMode::Create).await.is_err());
830    }
831}