Skip to main content

mkit_server/fs/
blob.rs

1//! `FsBlobStore`: content-addressed blobs as files, `<root>/packs/<64-hex>`
2//! by default, the layout `FileTransport::upload_pack` writes.
3
4use std::collections::HashMap;
5use std::fmt;
6use std::fs::{self, File, OpenOptions};
7use std::io::{self, ErrorKind, Read as _, Seek as _, SeekFrom, Write as _};
8use std::path::{Path, PathBuf};
9use std::pin::Pin;
10use std::sync::{Arc, Mutex, Weak};
11use std::task::{Context, Poll};
12use std::time::{Duration, SystemTime};
13
14use bytes::{Bytes, BytesMut};
15use futures_core::Stream;
16use mkit_core::hash::{Hash, Hasher};
17use mkit_transport_file::{create_dir_all_durably, sync_dir, temp_path};
18
19use super::{io_error, unavailable};
20use crate::store::{
21    BlobBody, BlobKey, BlobMeta, BlobStore, ByteRange, CommitOutcome, MAX_BLOB_PIECE_BYTES,
22    PackSink, StoreError,
23};
24
25/// The size of each piece of a streamed body.
26pub(super) const READ_BLOCK: usize = 64 * 1024;
27
28type MultipartLocks = Arc<Mutex<HashMap<[u8; 32], Weak<tokio::sync::Mutex<()>>>>>;
29
30/// A [`BlobStore`] over `<root>/<keyspace>/<64-hex>` for packs (`packs` by
31/// default), `<root>/upload-markers/v1/<64-hex>` for upload markers,
32/// `<root>/objects/<64-hex>` for extracted objects and
33/// `<root>/object-offsets/v1/<64-hex>` for their offset sidecars.
34/// An upload streams into a temp file in its destination
35/// directory (named like `FileTransport`'s own, `.<hex>.tmp.<pid>.<seq>`)
36/// while hashing it, and becomes visible only once its BLAKE3 and length
37/// verify: fsync, rename over the destination, fsync the directory. A
38/// failed, aborted or dropped upload removes its temp file and never
39/// touches an existing blob. A new directory's entry is fsynced into its
40/// parent before anything is published in it. A full disk or quota is
41/// [`StoreError::Full`].
42///
43/// A process that crashes mid-upload leaves its temp file,
44/// `<keyspace>/.<64-hex>.tmp.<pid>.<seq>` or a corresponding marker temp
45/// file, behind. Neither is visible as a blob;
46/// [`FsBlobStore::sweep_stale_uploads`] removes old ones.
47#[derive(Debug, Clone)]
48pub struct FsBlobStore {
49    pub(super) root: PathBuf,
50    keyspace: &'static str,
51    pub(super) multipart_locks: MultipartLocks,
52}
53
54impl FsBlobStore {
55    /// The `packs` keyspace under `root`, the directory `FileTransport`
56    /// uploads to.
57    #[must_use]
58    pub fn new(root: impl Into<PathBuf>) -> Self {
59        Self::with_keyspace(root, "packs")
60    }
61
62    /// The `keyspace` directory under `root`. `keyspace` is one plain path
63    /// component; `objects`, `object-offsets` and `upload-markers` are the
64    /// sibling namespaces' own directories, so a pack keyspace cannot use
65    /// them (its keys are refused).
66    ///
67    /// # Panics
68    /// If `keyspace` is empty, starts with `.` or holds a path separator.
69    #[must_use]
70    pub fn with_keyspace(root: impl Into<PathBuf>, keyspace: &'static str) -> Self {
71        assert!(
72            !keyspace.is_empty() && !keyspace.starts_with('.') && !keyspace.contains(['/', '\\']),
73            "a keyspace is one plain path component: {keyspace:?}"
74        );
75        assert!(
76            !crate::store::is_reserved_pack_keyspace(keyspace),
77            "a keyspace must not alias a sibling namespace: {keyspace:?}"
78        );
79        Self {
80            root: root.into(),
81            keyspace,
82            multipart_locks: Arc::new(Mutex::new(HashMap::new())),
83        }
84    }
85
86    /// The root directory.
87    #[must_use]
88    pub fn root(&self) -> &Path {
89        &self.root
90    }
91
92    /// The keyspace this store serves.
93    #[must_use]
94    pub fn keyspace(&self) -> &'static str {
95        self.keyspace
96    }
97
98    fn dir(&self) -> PathBuf {
99        self.root.join(self.keyspace)
100    }
101
102    fn path(&self, key: &BlobKey) -> Result<PathBuf, StoreError> {
103        Ok(self.root.join(key.relative_path(self.keyspace)?))
104    }
105
106    /// Remove temp files crashed uploads left in the pack and marker
107    /// directories: regular files named exactly `.<64-hex>.tmp.<pid>.<seq>`
108    /// (the names [`temp_path`] gives an upload, from this store or
109    /// `FileTransport::upload_pack`) last modified at least `min_age` ago.
110    /// Nothing else is touched: no blob, no symlink, no other name, no file
111    /// modified in the future. It also removes `server-uploads` session
112    /// directories whose immutable `meta` file's mtime is at least seven days
113    /// plus one hour old, independently of `min_age`. An incomplete session
114    /// without `meta` uses the directory mtime. Returns how many entries
115    /// were removed; an entry that cannot be inspected or removed is skipped.
116    ///
117    /// A live upload keeps its temp file's modification time fresh as it
118    /// writes, so a `min_age` well above any pause between two writes of
119    /// one upload is safe against `FileTransport` writers, which take no
120    /// lock. The caller must also rule out a live writer that can pause for
121    /// longer, a stalled streaming upload: `mkit serve` sweeps only while
122    /// it holds `serve.lock` exclusively, which no other `mkit serve` or
123    /// `mkit-server` (each holds it shared) can then hold.
124    ///
125    /// # Errors
126    /// I/O listing the directory; a missing directory sweeps nothing.
127    pub fn sweep_stale_uploads(&self, min_age: Duration) -> io::Result<usize> {
128        let now = SystemTime::now();
129        let mut removed = 0;
130        for dir in [
131            self.dir(),
132            self.root.join("upload-markers/v1"),
133            self.root.join("objects"),
134            self.root.join("object-offsets/v1"),
135        ] {
136            let entries = match fs::read_dir(dir) {
137                Ok(entries) => entries,
138                Err(e) if e.kind() == ErrorKind::NotFound => continue,
139                Err(e) => return Err(e),
140            };
141            for entry in entries {
142                let Ok(entry) = entry else { continue };
143                let name = entry.file_name();
144                if !name.to_str().is_some_and(is_upload_temp_name) {
145                    continue;
146                }
147                // `DirEntry::metadata` does not follow a symlink.
148                let Ok(meta) = entry.metadata() else { continue };
149                let stale = meta.is_file()
150                    && meta
151                        .modified()
152                        .ok()
153                        .and_then(|m| now.duration_since(m).ok())
154                        .is_some_and(|age| age >= min_age);
155                if stale && fs::remove_file(entry.path()).is_ok() {
156                    removed += 1;
157                }
158            }
159        }
160        Ok(removed + super::multipart::sweep_sessions(&self.root, now)?)
161    }
162}
163
164/// Whether `name` is an upload's temp file name, `.<64 lowercase
165/// hex>.tmp.<pid>.<seq>`, with a decimal `u32` pid and `u64` sequence.
166pub(super) fn is_upload_temp_name(name: &str) -> bool {
167    let decimal = |s: &str, max: usize| {
168        !s.is_empty() && s.len() <= max && s.bytes().all(|b| b.is_ascii_digit())
169    };
170    let Some(rest) = name.strip_prefix('.') else {
171        return false;
172    };
173    let Some((hex, rest)) = rest.split_at_checked(64) else {
174        return false;
175    };
176    let Some((pid, seq)) = rest.strip_prefix(".tmp.").and_then(|r| r.split_once('.')) else {
177        return false;
178    };
179    hex.bytes().all(|b| matches!(b, b'0'..=b'9' | b'a'..=b'f'))
180        && decimal(pid, 10)
181        && decimal(seq, 20)
182}
183
184/// The upload handle of [`FsBlobStore`]: a temp file and a running hash.
185/// Memory is one chunk, never the blob.
186pub struct FsPackSink {
187    /// The open temp file; `None` once closed for the rename.
188    file: Option<File>,
189    /// The temp file's path; `None` once it was renamed or removed.
190    tmp: Option<PathBuf>,
191    dest: PathBuf,
192    key: BlobKey,
193    declared: u64,
194    written: u64,
195    hasher: Hasher,
196}
197
198impl fmt::Debug for FsPackSink {
199    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
200        f.debug_struct("FsPackSink")
201            .field("tmp", &self.tmp)
202            .field("dest", &self.dest)
203            .field("declared", &self.declared)
204            .field("written", &self.written)
205            .finish_non_exhaustive()
206    }
207}
208
209impl FsPackSink {
210    /// Verify the upload and move it into place; `Ok(true)` if the
211    /// destination already existed.
212    fn publish(&mut self, root: Option<Hash>) -> Result<bool, StoreError> {
213        let expected = self.key.expected_root(root)?;
214        if self.written != self.declared {
215            return Err(StoreError::Invalid("blob length does not match".into()));
216        }
217        if self.hasher.finalize() != expected {
218            return Err(StoreError::Invalid(
219                "blob hash does not match its key".into(),
220            ));
221        }
222        let file = self.file.take().ok_or_else(closed)?;
223        file.sync_all().map_err(io_error)?;
224        drop(file);
225        let tmp = self.tmp.as_ref().ok_or_else(closed)?;
226        let existed = self.dest.exists();
227        // Identical bytes by construction, so replacing a present blob is
228        // harmless, and it repairs one a pre-atomic writer left short.
229        fs::rename(tmp, &self.dest).map_err(io_error)?;
230        self.tmp = None;
231        if let Some(dir) = self.dest.parent() {
232            sync_dir(dir).map_err(io_error)?;
233        }
234        Ok(existed)
235    }
236}
237
238fn outcome(existed: bool) -> CommitOutcome {
239    if existed {
240        CommitOutcome::AlreadyPresent
241    } else {
242        CommitOutcome::Created
243    }
244}
245
246/// The error for a sink used after it was closed (unreachable: `commit`
247/// and `abort` consume it).
248fn closed() -> StoreError {
249    StoreError::unavailable(io::Error::other("blob upload already closed"))
250}
251
252impl Drop for FsPackSink {
253    fn drop(&mut self) {
254        self.file = None;
255        if let Some(tmp) = self.tmp.take() {
256            let _ = fs::remove_file(tmp);
257        }
258    }
259}
260
261impl PackSink for FsPackSink {
262    async fn write(&mut self, chunk: Bytes) -> Result<(), StoreError> {
263        let total = self.written.checked_add(chunk.len() as u64);
264        let Some(total) = total.filter(|t| *t <= self.declared) else {
265            return Err(StoreError::Invalid("blob is longer than declared".into()));
266        };
267        let file = self.file.as_mut().ok_or_else(closed)?;
268        file.write_all(&chunk).map_err(io_error)?;
269        self.hasher.update(&chunk);
270        self.written = total;
271        Ok(())
272    }
273
274    async fn commit(mut self) -> Result<CommitOutcome, StoreError> {
275        // On any error, dropping `self` removes the temp file.
276        Ok(outcome(self.publish(None)?))
277    }
278
279    async fn commit_with_root(mut self, content_root: Hash) -> Result<CommitOutcome, StoreError> {
280        Ok(outcome(self.publish(Some(content_root))?))
281    }
282
283    async fn abort(self) {}
284}
285
286/// The remaining bytes of a streamed body, read [`READ_BLOCK`] at a time.
287struct Blocks {
288    file: File,
289    remaining: u64,
290}
291
292impl Stream for Blocks {
293    type Item = Result<Bytes, StoreError>;
294
295    fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<Option<Self::Item>> {
296        let this = self.get_mut();
297        if this.remaining == 0 {
298            return Poll::Ready(None);
299        }
300        let n = usize::try_from(this.remaining).map_or(READ_BLOCK, |r| r.min(READ_BLOCK));
301        let mut piece = BytesMut::zeroed(n);
302        match this.file.read_exact(&mut piece) {
303            Ok(()) => {
304                this.remaining -= n as u64;
305                Poll::Ready(Some(Ok(piece.freeze())))
306            }
307            Err(e) => {
308                this.remaining = 0;
309                Poll::Ready(Some(Err(io_error(e))))
310            }
311        }
312    }
313}
314
315// The piece size a streamed body must respect.
316const _: () = assert!(READ_BLOCK <= MAX_BLOB_PIECE_BYTES);
317
318impl BlobStore for FsBlobStore {
319    type Sink = FsPackSink;
320
321    async fn begin(&self, key: BlobKey, len: u64) -> Result<FsPackSink, StoreError> {
322        let dest = self.path(&key)?;
323        let dir = dest
324            .parent()
325            .ok_or_else(|| StoreError::Invalid("blob path has no directory".into()))?;
326        create_dir_all_durably(dir).map_err(io_error)?;
327        let tmp = temp_path(&dest).map_err(io_error)?;
328        let file = OpenOptions::new()
329            .write(true)
330            .create_new(true)
331            .open(&tmp)
332            .map_err(io_error)?;
333        Ok(FsPackSink {
334            file: Some(file),
335            tmp: Some(tmp),
336            dest,
337            key,
338            declared: len,
339            written: 0,
340            hasher: Hasher::new(),
341        })
342    }
343
344    async fn get(
345        &self,
346        key: &BlobKey,
347        range: Option<ByteRange>,
348    ) -> Result<Option<BlobBody>, StoreError> {
349        let mut file = match File::open(self.path(key)?) {
350            Ok(file) => file,
351            Err(e) if e.kind() == ErrorKind::NotFound => return Ok(None),
352            Err(e) => return Err(io_error(e)),
353        };
354        let len = file.metadata().map_err(io_error)?.len();
355        let span = match range {
356            Some(range) => range.resolve(len)?,
357            None => 0..len,
358        };
359        if span.start > 0 {
360            file.seek(SeekFrom::Start(span.start)).map_err(io_error)?;
361        }
362        let n = span.end - span.start;
363        if let Ok(whole) = usize::try_from(n)
364            && whole <= MAX_BLOB_PIECE_BYTES
365        {
366            let mut buf = vec![0; whole];
367            file.read_exact(&mut buf).map_err(io_error)?;
368            return Ok(Some(BlobBody::Bytes(Bytes::from(buf))));
369        }
370        Ok(Some(BlobBody::Stream {
371            len: n,
372            stream: Box::pin(Blocks { file, remaining: n }),
373        }))
374    }
375
376    async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
377        match fs::metadata(self.path(key)?) {
378            Ok(meta) => Ok(Some(BlobMeta { len: meta.len() })),
379            Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
380            Err(e) => Err(io_error(e)),
381        }
382    }
383
384    async fn probe(&self) -> Result<(), StoreError> {
385        let meta = fs::metadata(&self.root).map_err(io_error)?;
386        if meta.is_dir() {
387            Ok(())
388        } else {
389            Err(unavailable(io::Error::other(
390                "blob root is not a directory",
391            )))
392        }
393    }
394
395    async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
396        let path = self.path(key)?;
397        match fs::remove_file(&path) {
398            Ok(()) => {}
399            Err(e) if e.kind() == ErrorKind::NotFound => return Ok(false),
400            Err(e) => return Err(io_error(e)),
401        }
402        if let Some(dir) = path.parent() {
403            sync_dir(dir).map_err(io_error)?;
404        }
405        Ok(true)
406    }
407}