Skip to main content

mkit_server/fs/
layout.rs

1//! `FsLayoutStore`: the ref class of one repo as files, delegating to
2//! `FileTransport`.
3
4use std::fmt;
5use std::fs;
6use std::io::ErrorKind;
7use std::path::{Path, PathBuf};
8use std::sync::Arc;
9
10use mkit_core::hash::{Hash, hash, to_hex, to_hex_bytes};
11use mkit_core::protocol::RefWriteCondition;
12use mkit_transport_file::{FileTransport, LockedRefs, RefFileError};
13
14use super::{io_error, ref_file_error, unavailable};
15use crate::refs;
16use crate::repo::{RepoId, RepoName};
17use crate::rt::{Clock, SystemClock};
18use crate::store::{
19    Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, PartitionStats, Precondition,
20    ScanPage, StoreCapabilities, StoreError, Value, Write, keys,
21};
22
23/// The marker a SQLite-metadata server deployment writes under the
24/// served root, binding it to one database: the root's refs live in `SQLite`, so
25/// [`FsLayoutStore::open`] refuses it (R-81), and so does every ref write
26/// through `FileTransport`. Never removed automatically.
27pub const META_MARKER: &str = mkit_transport_file::SERVER_META_MARKER;
28
29/// The directory `FileTransport` keeps refs in, as a ref-name prefix: the
30/// names the pipeline serves.
31const REFS_PREFIX: &str = refs::SERVED_REFS_PREFIX;
32
33/// Where the ref-class rows whose name is not a `refs/` ref name live,
34/// relative to the root: out of reach of every ref name (a ref name
35/// component never starts with `.`) and of `FileTransport::list_refs`.
36const ROWS_DIR: &str = ".mkit/server/rows";
37
38/// Where one key lives.
39enum Slot<'k> {
40    /// A valid ref name under `refs/`: `FileTransport`'s ref file
41    /// `<root>/<name>`, holding a 32-byte id as `<64-hex>\n`.
42    Ref(&'k str),
43    /// Any other name bytes the ref class allows (`store::keys`): a row
44    /// file in `ROWS_DIR` (see `row_file_name`), holding
45    /// `be16(len(name)) ‖ name ‖ value`.
46    Row(&'k [u8]),
47}
48
49/// A [`NamespaceStore`] holding the ref class (`r 00 <repo> 00 <name>`) of
50/// one repo, in one partition, as files under the served root. Its
51/// capabilities are [`StoreCapabilities::refs_only`]: one key per batch,
52/// no layout-version row (the `.mkit` on-disk format is layout version 1).
53///
54/// A ref (a `refs/` name) is exactly `FileTransport`'s ref file: reads are
55/// [`FileTransport`]'s strict reads, and every write is its CAS (`Missing`
56/// for an `Absent` guard, `Match` for an `Equals` guard, `Any` otherwise)
57/// and its atomic write, under its ref lock (`<root>/.mkit/refs/.lock`), so
58/// `mkit+file://` remotes and this store see the same files and serialize
59/// on the same lock. Local `mkit` commands use their own per-ref
60/// `refs-<digest>.lock` and are not coordinated with it
61/// (SPEC-CONCURRENCY §3.1). A ref's value is its 32-byte
62/// id. A ref file that does not decode is [`StoreError::Corrupt`] on a
63/// read or a precondition, never absent; a scan skips it with a warning,
64/// like a ref file whose name is over [`refs::MAX_REF_NAME_BYTES`] (written
65/// before SPEC-REFS §3 capped names), as `FileTransport::list_refs` skips
66/// both. A ref whose file would clash with
67/// another ref's directory, or the reverse, is [`StoreError::Invalid`]; a
68/// delete removes the directories it leaves empty.
69///
70/// The ref class also allows names that are not `refs/` ref names, with
71/// any value. The pipeline never writes one (it serves only `refs/` names,
72/// R-86); they live in row files under `.mkit/server/rows/`, written under
73/// the same lock and invisible to the CLI and `FileTransport` (so a name
74/// like `packs/<hex>` can never overwrite a pack). Ref files an older
75/// `mkit serve` wrote outside `refs/` (`<root>/main`) are not served.
76///
77/// `apply` takes the ref lock, reads the store clock (for a
78/// [`Precondition::NotAfter`]), checks every precondition and writes, all
79/// in one synchronous step (normative rules 4 and 8). Reads take no lock:
80/// every write is one atomic rename. A full disk or quota is
81/// [`StoreError::Full`], except for a delete-only batch (rule 7).
82///
83/// A process that crashes mid-write leaves its temp file
84/// (`.<file>.tmp.<pid>.<seq>`, next to the ref or row file) behind. It is
85/// at most a ref wire or a row long; scans skip it, and nothing sweeps it
86/// (the pack temp files a crashed upload leaves are swept, see
87/// `FsBlobStore::sweep_stale_uploads`).
88///
89/// It is the permanent metadata store of the server-free ssh path
90/// (reconciliation R-13), with `SinglePartition` routing.
91pub struct FsLayoutStore {
92    tx: FileTransport,
93    partition: Partition,
94    repo: RepoName,
95    clock: Arc<dyn Clock>,
96}
97
98impl fmt::Debug for FsLayoutStore {
99    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
100        f.debug_struct("FsLayoutStore")
101            .field("root", &self.tx.root())
102            .field("partition", &self.partition)
103            .field("repo", &self.repo)
104            .finish_non_exhaustive()
105    }
106}
107
108impl FsLayoutStore {
109    /// `repo`'s refs, served from `root`, in the partition `SinglePartition`
110    /// routes them to (`Partition::Namespace(repo.namespace)`).
111    #[must_use]
112    pub fn new(root: impl Into<PathBuf>, repo: &RepoId) -> Self {
113        let partition = Partition::Namespace(repo.namespace.clone());
114        Self::in_partition(root, partition, repo.name.clone())
115    }
116
117    /// [`Self::new`], refusing a root whose refs live in `SQLite` (R-81):
118    /// one carrying the [`META_MARKER`] a SQLite-metadata server
119    /// deployment writes. Serving its ref files too would keep a second,
120    /// diverging copy of the refs. Every server of a `.mkit` root opens
121    /// its ref store through here.
122    ///
123    /// # Errors
124    /// [`StoreError::Unsupported`] naming both ways out when the root is
125    /// marked; [`StoreError::Unavailable`] when the marker cannot be
126    /// checked.
127    pub fn open(root: impl Into<PathBuf>, repo: &RepoId) -> Result<Self, StoreError> {
128        let root = root.into();
129        let marker = root.join(META_MARKER);
130        match fs::symlink_metadata(&marker) {
131            Ok(_) => Err(StoreError::Unsupported(
132                format!(
133                    "repo root {} is served with --meta sqlite (marker {META_MARKER}): its refs \
134                     live in SQLite, and serving its file-based refs would keep a second, \
135                     diverging copy. Serve this root with --meta sqlite:<PATH>, or migrate the \
136                     refs back to files and remove {} by hand.",
137                    root.display(),
138                    marker.display()
139                )
140                .into(),
141            )),
142            Err(e) if e.kind() == ErrorKind::NotFound => Ok(Self::new(root, repo)),
143            Err(e) => Err(unavailable(e)),
144        }
145    }
146
147    /// `repo`'s refs in `partition`, served from `root`: for a deployment
148    /// that maps each partition to its own directory. Any other partition
149    /// or repo is [`StoreError::Unsupported`].
150    #[must_use]
151    pub fn in_partition(root: impl Into<PathBuf>, partition: Partition, repo: RepoName) -> Self {
152        Self {
153            tx: FileTransport::new(root),
154            partition,
155            repo,
156            clock: Arc::new(SystemClock),
157        }
158    }
159
160    /// Use `clock` for [`Precondition::NotAfter`] instead of the host
161    /// clock.
162    #[must_use]
163    pub fn with_clock(mut self, clock: Arc<dyn Clock>) -> Self {
164        self.clock = clock;
165        self
166    }
167
168    /// The served root.
169    #[must_use]
170    pub fn root(&self) -> &Path {
171        self.tx.root()
172    }
173
174    fn check_partition(&self, p: &Partition) -> Result<(), StoreError> {
175        if *p == self.partition {
176            Ok(())
177        } else {
178            Err(StoreError::Unsupported(
179                "this store holds a single partition".into(),
180            ))
181        }
182    }
183
184    /// Where `key` lives, if this store can hold it.
185    fn slot<'k>(&self, key: &'k Key) -> Result<Slot<'k>, StoreError> {
186        let Some(rest) = key.as_bytes().strip_prefix(b"r\0") else {
187            return Err(StoreError::Unsupported(
188                "this store holds only ref keys".into(),
189            ));
190        };
191        let name = rest
192            .strip_prefix(self.repo.as_str().as_bytes())
193            .and_then(|r| r.strip_prefix(b"\0"))
194            .ok_or_else(|| StoreError::Unsupported("this store holds one repo's refs".into()))?;
195        Ok(match core::str::from_utf8(name) {
196            Ok(name) if is_ref_name(name) => Slot::Ref(name),
197            _ => Slot::Row(name),
198        })
199    }
200
201    /// The value at `key`.
202    fn read(&self, key: &Key) -> Result<Option<Value>, StoreError> {
203        match self.slot(key)? {
204            Slot::Ref(name) => {
205                // Strict: a ref file that does not decode is `Corrupt`,
206                // never absent, so no precondition passes over it.
207                let id = self.tx.read_ref_strict(name).map_err(ref_file_error)?;
208                Ok(id.map(|id| Value::new(id.to_vec())))
209            }
210            Slot::Row(name) => {
211                let path = self.tx.server_path(&row_path(name));
212                match fs::read(path.map_err(ref_file_error)?) {
213                    Ok(bytes) => {
214                        let (stored, value) = decode_row(&bytes)?;
215                        if stored != name {
216                            return Err(StoreError::Corrupt("row file holds another key".into()));
217                        }
218                        Ok(Some(value))
219                    }
220                    Err(e) if e.kind() == ErrorKind::NotFound => Ok(None),
221                    Err(e) => Err(io_error(e)),
222                }
223            }
224        }
225    }
226
227    /// Every row whose key `keep` accepts, unordered.
228    fn rows(&self, keep: impl Fn(&Key) -> bool) -> Result<Vec<(Key, Value)>, StoreError> {
229        let mut rows = Vec::new();
230        // Listing `refs/` (not the whole root) never reads a pack. Unlike
231        // `read`, a listing skips a file it cannot serve, loudly, as
232        // `FileTransport::list_refs` (today's `mkit serve`) skips it: one
233        // stray file must not fail every listing of its directory.
234        let listed = self.tx.list_ref_files(REFS_PREFIX);
235        for (name, id) in listed.map_err(ref_file_error)? {
236            let key = keys::ref_key(&self.repo, &name);
237            if !keep(&key) {
238                continue;
239            }
240            if name.len() > refs::MAX_REF_NAME_BYTES {
241                // Written before SPEC-REFS §3 capped ref names.
242                tracing::warn!(
243                    len = name.len(),
244                    max = refs::MAX_REF_NAME_BYTES,
245                    "skipping a ref file whose name is over the ref-name limit"
246                );
247                continue;
248            }
249            match id {
250                Some(id) if is_ref_name(&name) => rows.push((key, Value::new(id.to_vec()))),
251                Some(_) => {}
252                None => tracing::warn!(file = %name, "skipping a ref file that holds no ref id"),
253            }
254        }
255        let rows_dir = self.tx.server_path(Path::new(ROWS_DIR));
256        let dir = match fs::read_dir(rows_dir.map_err(ref_file_error)?) {
257            Ok(dir) => dir,
258            Err(e) if e.kind() == ErrorKind::NotFound => return Ok(rows),
259            Err(e) => return Err(io_error(e)),
260        };
261        for entry in dir {
262            let entry = entry.map_err(io_error)?;
263            let file_name = entry.file_name();
264            let file_name = file_name.to_string_lossy();
265            if file_name.starts_with('.') {
266                continue; // an in-flight or abandoned temp file
267            }
268            // A short name is in the file name: skip rows out of range
269            // without reading them.
270            if let Some(hex) = file_name.strip_prefix(SHORT_ROW) {
271                let name = from_hex(hex).ok_or_else(|| corrupt_row_name(&file_name))?;
272                if !keep(&self.key(&name)) {
273                    continue;
274                }
275            }
276            let bytes = fs::read(entry.path()).map_err(io_error)?;
277            let (name, value) = decode_row(&bytes)?;
278            if row_file_name(name) != file_name {
279                return Err(corrupt_row_name(&file_name));
280            }
281            let key = self.key(name);
282            if keep(&key) {
283                rows.push((key, value));
284            }
285        }
286        Ok(rows)
287    }
288
289    /// `r 00 <repo> 00 <name>` for any name bytes.
290    fn key(&self, name: &[u8]) -> Key {
291        let repo = self.repo.as_str().as_bytes();
292        Key::new([b"r\0", repo, b"\0", name].concat())
293    }
294
295    /// The locked step of [`NamespaceStore::apply`]: read the clock, check
296    /// the preconditions in order, then write. [`Batch::validate`] allowed
297    /// at most one write and at most one key precondition, on its key.
298    fn check_and_write(
299        &self,
300        refs: &LockedRefs<'_>,
301        batch: &Batch,
302    ) -> Result<BatchOutcome, StoreError> {
303        // Rule 8: the store's clock, read once, under the lock. A reading
304        // before the epoch fails every deadline (fail closed).
305        let now = batch
306            .preconditions
307            .iter()
308            .any(|pre| matches!(pre, Precondition::NotAfter(_)))
309            .then(|| u64::try_from(self.clock.now_ms()).unwrap_or(u64::MAX));
310        // The guard FileTransport re-checks on a ref write, and the index
311        // of the key precondition it stands for.
312        let mut guard = (0, RefWriteCondition::Any);
313        for (index, pre) in batch.preconditions.iter().enumerate() {
314            let (holds, observed) = match pre {
315                Precondition::NotAfter(deadline) => {
316                    let backend_now = now.unwrap_or(u64::MAX);
317                    if backend_now > *deadline {
318                        return Ok(BatchOutcome::DeadlinePassed { backend_now });
319                    }
320                    continue;
321                }
322                Precondition::Absent(key) => {
323                    guard = (index, RefWriteCondition::Missing);
324                    let current = self.read(key)?;
325                    (current.is_none(), current)
326                }
327                Precondition::Present(key) => (self.read(key)?.is_some(), None),
328                Precondition::Equals(key, want) => {
329                    if let Ok(id) = Hash::try_from(want.as_bytes()) {
330                        guard = (index, RefWriteCondition::Match(id));
331                    }
332                    let current = self.read(key)?;
333                    (current.as_ref() == Some(want), current)
334                }
335            };
336            if !holds {
337                return Ok(BatchOutcome::PreconditionFailed { index, observed });
338            }
339        }
340        for write in &batch.writes {
341            match write {
342                Write::Put(key, value) => match self.slot(key)? {
343                    Slot::Ref(name) => {
344                        let id = ref_id(value)?;
345                        match refs.update_ref(name, guard.1, &id) {
346                            Ok(()) => {}
347                            // FileTransport's own re-check under the same
348                            // lock failed (unreachable unless a writer
349                            // bypassed the lock): report what a read sees,
350                            // as `mkit serve` does on a conflict.
351                            Err(RefFileError::Conflict) => {
352                                return Ok(BatchOutcome::PreconditionFailed {
353                                    index: guard.0,
354                                    observed: self.read(key)?,
355                                });
356                            }
357                            Err(e) => return Err(ref_file_error(e)),
358                        }
359                    }
360                    Slot::Row(name) => refs
361                        .write_file(&row_path(name), &encode_row(name, value))
362                        .map_err(ref_file_error)?,
363                },
364                Write::Delete(key) => {
365                    match self.slot(key)? {
366                        Slot::Ref(name) => refs.delete_ref(name),
367                        Slot::Row(name) => refs.remove_file(&row_path(name)),
368                    }
369                    .map_err(ref_file_error)?;
370                }
371            }
372        }
373        Ok(BatchOutcome::Committed)
374    }
375}
376
377/// Whether `name` is a ref this store keeps as a `FileTransport` ref file.
378fn is_ref_name(name: &str) -> bool {
379    refs::is_served_ref_name(name)
380}
381
382/// A ref's value: its 32-byte id.
383fn ref_id(value: &Value) -> Result<Hash, StoreError> {
384    Hash::try_from(value.as_bytes())
385        .map_err(|_| StoreError::Invalid("a ref's value is its 32-byte id".into()))
386}
387
388/// The row file of `name`, relative to the root.
389fn row_path(name: &[u8]) -> PathBuf {
390    Path::new(ROWS_DIR).join(row_file_name(name))
391}
392
393/// A row file named after its name: `n<hex(name)>`, for names up to
394/// [`MAX_SHORT_ROW`] bytes (a scan decodes the key from the file name and
395/// reads only the rows in range).
396const SHORT_ROW: &str = "n";
397/// A row file named after its name's hash: `h<hex(BLAKE3(name))>`, for
398/// longer names (a file name holds at most 255 bytes).
399const HASHED_ROW: &str = "h";
400/// The longest name kept in a file name (`1 + 2 × 100` bytes).
401pub(super) const MAX_SHORT_ROW: usize = 100;
402
403/// The longest temp file name `FileTransport` writes next to a row file
404/// `<file>`: `.<file>.tmp.<pid: u32>.<seq: u64>`.
405pub(super) const fn max_temp_name(file_name_len: usize) -> usize {
406    1 + file_name_len + ".tmp.".len() + 10 + 1 + 20
407}
408
409/// Every row file, and every temp file written to publish one, fits the
410/// 255-byte file-name limit of common filesystems (APFS, ext4, NTFS), for
411/// any pid and any value of the process-wide temp counter.
412pub(super) const NAME_MAX: usize = 255;
413const _: () = assert!(max_temp_name(SHORT_ROW.len() + 2 * MAX_SHORT_ROW) <= NAME_MAX);
414const _: () = assert!(max_temp_name(HASHED_ROW.len() + 64) <= NAME_MAX);
415
416fn row_file_name(name: &[u8]) -> String {
417    if name.len() <= MAX_SHORT_ROW {
418        format!("{SHORT_ROW}{}", to_hex_bytes(name))
419    } else {
420        format!("{HASHED_ROW}{}", to_hex(&hash(name)))
421    }
422}
423
424fn corrupt_row_name(file_name: &str) -> StoreError {
425    StoreError::Corrupt(format!("row file {file_name} does not hold its key").into())
426}
427
428/// Lowercase hex to bytes.
429fn from_hex(hex: &str) -> Option<Vec<u8>> {
430    let digit = |c: u8| match c {
431        b'0'..=b'9' => Some(c - b'0'),
432        b'a'..=b'f' => Some(c - b'a' + 10),
433        _ => None,
434    };
435    let (pairs, []) = hex.as_bytes().as_chunks::<2>() else {
436        return None;
437    };
438    pairs
439        .iter()
440        .map(|[hi, lo]| Some(digit(*hi)? << 4 | digit(*lo)?))
441        .collect()
442}
443
444fn encode_row(name: &[u8], value: &Value) -> Vec<u8> {
445    // A key is at most MAX_KEY_BYTES (1024), so its name fits a be16.
446    let len = u16::try_from(name.len()).unwrap_or(u16::MAX);
447    [&len.to_be_bytes()[..], name, value.as_bytes()].concat()
448}
449
450/// A row file's name and value.
451fn decode_row(bytes: &[u8]) -> Result<(&[u8], Value), StoreError> {
452    let corrupt = || StoreError::Corrupt("truncated row file".into());
453    let (len, rest) = bytes.split_first_chunk::<2>().ok_or_else(corrupt)?;
454    let len = usize::from(u16::from_be_bytes(*len));
455    let (name, value) = rest.split_at_checked(len).ok_or_else(corrupt)?;
456    Ok((name, Value::new(value.to_vec())))
457}
458
459impl NamespaceStore for FsLayoutStore {
460    fn capabilities(&self) -> StoreCapabilities {
461        StoreCapabilities::refs_only()
462    }
463
464    async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
465        self.check_partition(p)?;
466        self.read(key)
467    }
468
469    async fn scan(
470        &self,
471        p: &Partition,
472        start: &Key,
473        end: &Key,
474        after: Option<&Cursor>,
475        limit: u32,
476    ) -> Result<ScanPage, StoreError> {
477        if limit == 0 {
478            return Err(StoreError::Invalid("scan limit must be at least 1".into()));
479        }
480        self.check_partition(p)?;
481        // Every cursor this range returns is one of its keys.
482        let after = after.map(|c| Key::new(c.clone().into_bytes()));
483        if after.as_ref().is_some_and(|c| c < start || c >= end) {
484            return Err(StoreError::Invalid(
485                "scan cursor outside the scanned range".into(),
486            ));
487        }
488        let mut entries =
489            self.rows(|k| start <= k && k < end && after.as_ref().is_none_or(|c| k > c))?;
490        entries.sort_by(|a, b| a.0.cmp(&b.0));
491        let want = usize::try_from(limit).unwrap_or(usize::MAX);
492        let next = (entries.len() > want).then(|| {
493            entries.truncate(want);
494            Cursor::new(entries[want - 1].0.clone().into_bytes())
495        });
496        Ok(ScanPage { entries, next })
497    }
498
499    async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
500        batch.validate(&self.capabilities())?;
501        self.check_partition(p)?;
502        // Every key resolves, and a ref write carries an id, before the
503        // lock is taken: a batch this store cannot hold writes nothing.
504        for pre in &batch.preconditions {
505            match pre {
506                Precondition::Absent(key)
507                | Precondition::Present(key)
508                | Precondition::Equals(key, _) => {
509                    self.slot(key)?;
510                }
511                Precondition::NotAfter(_) => {}
512            }
513        }
514        for write in &batch.writes {
515            match write {
516                Write::Put(key, value) => {
517                    if let Slot::Ref(_) = self.slot(key)? {
518                        ref_id(value)?;
519                    }
520                }
521                Write::Delete(key) => {
522                    self.slot(key)?;
523                }
524            }
525        }
526        let result = self
527            .tx
528            .with_ref_lock(|refs| self.check_and_write(refs, &batch))
529            .map_err(ref_file_error)
530            .and_then(|outcome| outcome);
531        match result {
532            // Rule 7: a delete-only batch never reports `Full` (a delete
533            // frees space; a full disk while syncing it is an outage).
534            Err(StoreError::Full) if !batch.has_put() => Err(unavailable(std::io::Error::other(
535                "storage full while deleting",
536            ))),
537            other => other,
538        }
539    }
540
541    async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
542        self.check_partition(p)?;
543        let rows = self.rows(|_| true)?;
544        let bytes = rows
545            .iter()
546            .map(|(k, v)| (k.as_bytes().len() + v.as_bytes().len()) as u64)
547            .sum();
548        Ok(PartitionStats {
549            bytes,
550            keys: Some(rows.len() as u64),
551        })
552    }
553
554    async fn probe(&self) -> Result<(), StoreError> {
555        let meta = fs::metadata(self.root()).map_err(io_error)?;
556        if meta.is_dir() {
557            Ok(())
558        } else {
559            Err(unavailable(std::io::Error::other(
560                "ref root is not a directory",
561            )))
562        }
563    }
564}