1use 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
23pub const META_MARKER: &str = mkit_transport_file::SERVER_META_MARKER;
28
29const REFS_PREFIX: &str = refs::SERVED_REFS_PREFIX;
32
33const ROWS_DIR: &str = ".mkit/server/rows";
37
38enum Slot<'k> {
40 Ref(&'k str),
43 Row(&'k [u8]),
47}
48
49pub 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 #[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 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 #[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 #[must_use]
163 pub fn with_clock(mut self, clock: Arc<dyn Clock>) -> Self {
164 self.clock = clock;
165 self
166 }
167
168 #[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 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 fn read(&self, key: &Key) -> Result<Option<Value>, StoreError> {
203 match self.slot(key)? {
204 Slot::Ref(name) => {
205 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 fn rows(&self, keep: impl Fn(&Key) -> bool) -> Result<Vec<(Key, Value)>, StoreError> {
229 let mut rows = Vec::new();
230 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 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; }
268 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 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 fn check_and_write(
299 &self,
300 refs: &LockedRefs<'_>,
301 batch: &Batch,
302 ) -> Result<BatchOutcome, StoreError> {
303 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 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 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
377fn is_ref_name(name: &str) -> bool {
379 refs::is_served_ref_name(name)
380}
381
382fn 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
388fn row_path(name: &[u8]) -> PathBuf {
390 Path::new(ROWS_DIR).join(row_file_name(name))
391}
392
393const SHORT_ROW: &str = "n";
397const HASHED_ROW: &str = "h";
400pub(super) const MAX_SHORT_ROW: usize = 100;
402
403pub(super) const fn max_temp_name(file_name_len: usize) -> usize {
406 1 + file_name_len + ".tmp.".len() + 10 + 1 + 20
407}
408
409pub(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
428fn 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 let len = u16::try_from(name.len()).unwrap_or(u16::MAX);
447 [&len.to_be_bytes()[..], name, value.as_bytes()].concat()
448}
449
450fn 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 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 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 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}