From 64f583b9f2d16cba049effd04bf1574b0f6b35fa Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 30 Sep 2026 13:48:24 +0300 Subject: [PATCH] knot-lfs,knot-xrpc: split store into pools, gc-able Lewis: May this revision serve well! --- knot2/crates/knot-lfs/src/gc.rs | 57 ++++-- knot2/crates/knot-lfs/src/lib.rs | 8 +- knot2/crates/knot-lfs/src/store.rs | 228 ++++++++++++++++++----- knot2/crates/knot-lfs/src/transfer.rs | 54 ++++-- knot2/crates/knot-lfs/src/types.rs | 79 ++++++-- knot2/crates/knot-xrpc/src/lfs.rs | 29 ++- knot2/crates/knot-xrpc/tests/lfs.rs | 14 +- knot2/crates/knot-xrpc/tests/lfs_soak.rs | 10 +- 8 files changed, 348 insertions(+), 131 deletions(-) diff --git a/knot2/crates/knot-lfs/src/gc.rs b/knot2/crates/knot-lfs/src/gc.rs index a31b63bd1..d37336522 100644 --- a/knot2/crates/knot-lfs/src/gc.rs +++ b/knot2/crates/knot-lfs/src/gc.rs @@ -5,7 +5,7 @@ use knot_git::{GitError, Haves, RefClass, Repo, Wants}; use knot_types::{Oid, RepoDid, UnixSeconds}; use crate::store::{DiskStore, Reclaimed, expired}; -use crate::{LfsError, LfsOid, LfsSize, scan_pointers}; +use crate::{LfsError, LfsOid, LfsSize, Pool, scan_pointers}; #[derive(Debug, thiserror::Error)] pub enum GcError { @@ -14,6 +14,15 @@ pub enum GcError { #[error("store: {0}")] Store(#[from] LfsError), } +#[derive(Debug, thiserror::Error)] +pub enum RootsDenied { + #[error("The repo wouldn't open: {0}.")] + Unopened(GitError), + #[error("This repo's history is longer than the blob index window.")] + Windowed, + #[error("This repo's history stopped traversing or a stored frame wouldn't parse.")] + Untraversable, +} #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] pub struct GcReport { @@ -54,7 +63,7 @@ pub fn collect_repo( grace: Duration, now: SystemTime, ) -> Result { - let stored = store.enumerate(did)?; + let stored = store.enumerate(Pool::Lfs, did)?; if stored.is_empty() { return Ok(GcReport::default()); } @@ -64,7 +73,7 @@ pub fn collect_repo( .iter() .filter(|object| !reachable.contains(&object.oid)) .filter(|object| expired(now, object.mtime, grace)) - .map(|object| store.collect_expired(did, &object.oid, grace, now)) + .map(|object| store.collect_expired(Pool::Lfs, did, &object.oid, grace, now)) .try_fold( GcReport { scanned: stored.len(), @@ -85,15 +94,13 @@ pub fn collect_repo( mod tests { use std::time::Duration; + use super::*; + use crate::store::DiskStore; + use crate::{ClaimedSize, LfsOid, LfsSize, LfsStore, LfsStorePath}; use knot_git::{ EntryKind, Identity, Layout, NewCommit, RefUpdate, Repo, StagedAction, StagedChange, }; use knot_types::{AuthorName, BranchName, Email, Oid, RefName, RepoDid}; - use sha2::{Digest, Sha256}; - - use super::*; - use crate::store::DiskStore; - use crate::{ClaimedSize, LfsOid, LfsSize, LfsStore, LfsStorePath}; const EMPTY_TREE: &str = "4b825dc642cb6eb9a060e54bf8d69288fbee4904"; const DAY: Duration = Duration::from_secs(86_400); @@ -133,10 +140,16 @@ mod tests { } fn put_media(f: &Fixture, bytes: &[u8]) -> (LfsOid, LfsSize) { - let oid = LfsOid::from_digest(Sha256::digest(bytes).into()); + let oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(bytes)); let bytes_len = bytes.len() as u64; f.store - .put(&f.did, &oid, ClaimedSize::new(bytes_len), &mut &bytes[..]) + .put( + Pool::Lfs, + &f.did, + &oid, + ClaimedSize::new(bytes_len), + &mut &bytes[..], + ) .unwrap(); (oid, LfsSize::new(bytes_len)) } @@ -192,7 +205,12 @@ mod tests { } fn age(f: &Fixture, oid: &LfsOid, past: Duration) { - let path = f.store.object_file(&f.did, oid).unwrap().unwrap().1; + let path = f + .store + .object_file(Pool::Lfs, &f.did, oid) + .unwrap() + .unwrap() + .1; std::fs::OpenOptions::new() .write(true) .open(&path) @@ -232,7 +250,7 @@ mod tests { assert_eq!(report.swept, 0); [&live, &fresh, &forced].iter().for_each(|oid| { assert!( - f.store.probe(&f.did, oid).unwrap().is_some(), + f.store.probe(Pool::Lfs, &f.did, oid).unwrap().is_some(), "a live ref, the grace window, and the reflog old tip each keep their object" ); }); @@ -262,8 +280,8 @@ mod tests { "a plain orphan and a deleted-branch orphan are both swept once the mtime grace expires" ); assert_eq!(report.bytes, orphan_size.saturating_add(dropped_size)); - assert_eq!(f.store.probe(&f.did, &orphan).unwrap(), None); - assert_eq!(f.store.probe(&f.did, &dropped).unwrap(), None); + assert_eq!(f.store.probe(Pool::Lfs, &f.did, &orphan).unwrap(), None); + assert_eq!(f.store.probe(Pool::Lfs, &f.did, &dropped).unwrap(), None); } #[test] @@ -299,11 +317,11 @@ mod tests { .saturating_add(chain_size) .saturating_add(snapshot_size) ); - assert_eq!(f.store.probe(&f.did, &cob).unwrap(), None); - assert_eq!(f.store.probe(&f.did, &chain).unwrap(), None); - assert_eq!(f.store.probe(&f.did, &snapshot).unwrap(), None); + assert_eq!(f.store.probe(Pool::Lfs, &f.did, &cob).unwrap(), None); + assert_eq!(f.store.probe(Pool::Lfs, &f.did, &chain).unwrap(), None); + assert_eq!(f.store.probe(Pool::Lfs, &f.did, &snapshot).unwrap(), None); assert!( - f.store.probe(&f.did, &staged).unwrap().is_some(), + f.store.probe(Pool::Lfs, &f.did, &staged).unwrap().is_some(), "fork staging roots its bytes, so the three sweeps above are the class filter \ at work rather than every reserved name losing its objects" ); @@ -312,7 +330,8 @@ mod tests { #[test] fn a_repo_with_no_stored_objects_never_walks_git() { let f = fixture(); - let orphan_pointer = LfsOid::from_digest(Sha256::digest(b"pointer without bytes").into()); + let orphan_pointer = + LfsOid::from_digest(knot_types::Sha256Digest::hash(b"pointer without bytes")); commit_pointer(&f, "refs/heads/main", &orphan_pointer, LfsSize::new(21)); assert_eq!(collect(&f, DAY), GcReport::default()); } diff --git a/knot2/crates/knot-lfs/src/lib.rs b/knot2/crates/knot-lfs/src/lib.rs index e97968ecb..aa4ea3f7c 100644 --- a/knot2/crates/knot-lfs/src/lib.rs +++ b/knot2/crates/knot-lfs/src/lib.rs @@ -15,14 +15,14 @@ pub use batch::{ HashAlgo, MAX_BATCH_OBJECTS, TransferAdapter, }; pub use error::LfsError; -pub use gc::{GcError, GcReport, collect_repo}; +pub use gc::{GcError, GcReport, RootsDenied, collect_repo}; pub use pointer::{LfsPointer, POINTER_MAX_BYTES, parse_pointer}; pub use scan::scan_pointers; pub use store::{ DiskStore, LfsHandle, LfsStore, MemoryStore, OrphanSweep, Reclaimed, StoredObject, }; pub use transfer::{TransferOp, serve_transfer}; -pub use types::{ClaimedSize, FreeSpaceFloor, LfsOid, LfsSize, LfsStorePath, ObjectRelPath}; +pub use types::{ClaimedSize, FreeSpaceFloor, LfsOid, LfsSize, LfsStorePath, ObjectRelPath, Pool}; #[doc(hidden)] pub mod fuzz { @@ -30,7 +30,7 @@ pub mod fuzz { use crate::{ BatchRequest, FreeSpaceFloor, LfsSize, LfsStore, LfsStorePath, MAX_BATCH_OBJECTS, - MemoryStore, StoreAdmission, TransferOp, parse_pointer, serve_transfer, + MemoryStore, Pool, StoreAdmission, TransferOp, parse_pointer, serve_transfer, }; pub fn transfer(data: &[u8]) { @@ -67,7 +67,7 @@ pub mod fuzz { let repo = RepoDid::new("did:plc:squid").expect("static DID is valid"); let store = MemoryStore::new(); request.objects.iter().for_each(|object| { - let _ = store.probe(&repo, &object.oid); + let _ = store.probe(Pool::Lfs, &repo, &object.oid); }); } diff --git a/knot2/crates/knot-lfs/src/store.rs b/knot2/crates/knot-lfs/src/store.rs index fd2516cb9..d10838ccc 100644 --- a/knot2/crates/knot-lfs/src/store.rs +++ b/knot2/crates/knot-lfs/src/store.rs @@ -8,22 +8,30 @@ use knot_types::RepoDid; use sha2::{Digest, Sha256}; use crate::types::RepoPrefix; -use crate::{ClaimedSize, FreeSpaceFloor, LfsError, LfsOid, LfsSize, LfsStorePath, ObjectRelPath}; +use crate::{ + ClaimedSize, FreeSpaceFloor, LfsError, LfsOid, LfsSize, LfsStorePath, ObjectRelPath, Pool, +}; pub trait LfsStore: Send + Sync { fn put( &self, + pool: Pool, repo: &RepoDid, oid: &LfsOid, size: ClaimedSize, body: &mut dyn Read, ) -> Result<(), LfsError>; - fn read(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError>; + fn read( + &self, + pool: Pool, + repo: &RepoDid, + oid: &LfsOid, + ) -> Result, LfsError>; - fn probe(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError>; + fn probe(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError>; - fn touch(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError>; + fn touch(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError>; } #[derive(Debug, Clone, PartialEq, Eq)] @@ -98,7 +106,7 @@ fn write_verified( received: LfsSize::new(received), }); } - let computed = LfsOid::from_digest(hasher.finalize().into()); + let computed = LfsOid::from_digest(knot_types::Sha256Digest::new(hasher.finalize().into())); match computed == *declared { true => Ok(()), false => Err(LfsError::HashMismatch { @@ -258,14 +266,23 @@ impl DiskStore { Err(source) => Err(io_at("remove abandoned upload", &path)(source)), } })?; + migrate_legacy_prefixes(&root)?; let locks = std::iter::repeat_with(|| Mutex::new(())) .take(OID_LOCK_STRIPES) .collect(); Ok(Self { root, locks }) } - fn object_path(&self, repo: &RepoDid, oid: &LfsOid) -> Result { - Ok(self.root.object_path(&ObjectRelPath::new(repo, oid)?)) + fn object_path(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result { + Ok(self.root.object_path(&ObjectRelPath::new(pool, repo, oid)?)) + } + + fn prefix_path(&self, pool: Pool, repo: &RepoDid) -> Result { + Ok(self + .root + .as_path() + .join(pool.dir()) + .join(RepoPrefix::new(repo)?.as_path())) } fn oid_lock(&self, oid: &LfsOid) -> &Mutex<()> { @@ -274,15 +291,73 @@ impl DiskStore { } } +fn migrate_legacy_prefixes(root: &LfsStorePath) -> Result<(), LfsError> { + let lfs_pool = root.as_path().join(Pool::Lfs.dir()); + let legacy = subdirs(root.as_path())? + .into_iter() + .filter(|dir| { + dir.file_name().is_some_and(|name| { + !Pool::ALL.iter().any(|pool| name == pool.dir()) + && name != INCOMING_DIR + && name.to_str().is_some_and(is_did_method) + }) + }) + .collect::>(); + if legacy.is_empty() { + return Ok(()); + } + std::fs::create_dir_all(&lfs_pool).map_err(io_at("create dir", &lfs_pool))?; + legacy.iter().try_for_each(|dir| { + let into = lfs_pool.join( + dir.file_name() + .expect("a legacy prefix always has a final component"), + ); + merge_legacy_dir(dir, &into)?; + std::fs::remove_dir_all(dir).map_err(io_at("remove migrated prefix", dir)) + }) +} + +fn is_did_method(name: &str) -> bool { + name.chars() + .next() + .is_some_and(|first| first.is_ascii_lowercase()) + && name + .chars() + .all(|rest| rest.is_ascii_lowercase() || rest.is_ascii_digit()) +} + +fn merge_legacy_dir(from: &Path, into: &Path) -> Result<(), LfsError> { + std::fs::read_dir(from) + .map_err(io_at("read legacy prefix", from))? + .try_for_each(|entry| { + let entry = entry.map_err(io_at("read legacy entry", from))?; + let target = into.join(entry.file_name()); + let is_dir = entry + .file_type() + .map_err(io_at("stat legacy entry", &target))? + .is_dir(); + if is_dir { + std::fs::create_dir_all(&target).map_err(io_at("create dir", &target))?; + merge_legacy_dir(&entry.path(), &target) + } else if target.exists() { + Ok(()) + } else { + std::fs::rename(entry.path(), &target) + .map_err(io_at("migrate object into", &target)) + } + }) +} + impl LfsStore for DiskStore { fn put( &self, + pool: Pool, repo: &RepoDid, oid: &LfsOid, size: ClaimedSize, body: &mut dyn Read, ) -> Result<(), LfsError> { - let target = self.object_path(repo, oid)?; + let target = self.object_path(pool, repo, oid)?; let incoming = self.root.as_path().join(INCOMING_DIR); let mut temp = tempfile::Builder::new() .prefix("put-") @@ -313,8 +388,13 @@ impl LfsStore for DiskStore { fsync_chain(self.root.as_path(), parent) } - fn read(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - let path = self.object_path(repo, oid)?; + fn read( + &self, + pool: Pool, + repo: &RepoDid, + oid: &LfsOid, + ) -> Result, LfsError> { + let path = self.object_path(pool, repo, oid)?; match std::fs::File::open(&path) { Ok(file) => Ok(Box::new(file)), Err(source) if source.kind() == std::io::ErrorKind::NotFound => { @@ -324,12 +404,12 @@ impl LfsStore for DiskStore { } } - fn probe(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - Ok(self.object_file(repo, oid)?.map(|(size, _)| size)) + fn probe(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { + Ok(self.object_file(pool, repo, oid)?.map(|(size, _)| size)) } - fn touch(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - let path = self.object_path(repo, oid)?; + fn touch(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { + let path = self.object_path(pool, repo, oid)?; let _guard = self .oid_lock(oid) .lock() @@ -348,10 +428,11 @@ impl LfsStore for DiskStore { impl DiskStore { pub fn object_file( &self, + pool: Pool, repo: &RepoDid, oid: &LfsOid, ) -> Result, LfsError> { - let path = self.object_path(repo, oid)?; + let path = self.object_path(pool, repo, oid)?; match std::fs::metadata(&path) { Ok(meta) => Ok(Some((LfsSize::new(meta.len()), path))), Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(None), @@ -367,27 +448,29 @@ impl DiskStore { } pub fn remove_repo(&self, repo: &RepoDid) -> Result<(), LfsError> { - let prefix = self.root.as_path().join(RepoPrefix::new(repo)?.as_path()); - match std::fs::remove_dir_all(&prefix) { - Ok(()) => Ok(()), - Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(()), - Err(source) => Err(io_at("remove repo prefix", &prefix)(source)), - } + Pool::ALL.iter().try_for_each(|pool| { + let prefix = self.prefix_path(*pool, repo)?; + match std::fs::remove_dir_all(&prefix) { + Ok(()) => Ok(()), + Err(source) if source.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(source) => Err(io_at("remove repo prefix", &prefix)(source)), + } + }) } - pub fn enumerate(&self, repo: &RepoDid) -> Result, LfsError> { - let prefix = self.root.as_path().join(RepoPrefix::new(repo)?.as_path()); - enumerate_prefix(&prefix) + pub fn enumerate(&self, pool: Pool, repo: &RepoDid) -> Result, LfsError> { + enumerate_prefix(&self.prefix_path(pool, repo)?) } pub fn collect_expired( &self, + pool: Pool, repo: &RepoDid, oid: &LfsOid, grace: Duration, now: SystemTime, ) -> Result { - let path = self.object_path(repo, oid)?; + let path = self.object_path(pool, repo, oid)?; let _guard = self .oid_lock(oid) .lock() @@ -410,38 +493,75 @@ impl DiskStore { } } + pub fn sweep_unreferenced( + &self, + repo: &RepoDid, + referenced: &HashSet, + grace: Duration, + now: SystemTime, + ) -> Result { + let doomed: Vec = self + .enumerate(Pool::Attachments, repo)? + .into_iter() + .filter(|object| !referenced.contains(&object.oid) && expired(now, object.mtime, grace)) + .collect(); + doomed + .iter() + .map(|object| self.collect_expired(Pool::Attachments, repo, &object.oid, grace, now)) + .try_fold(OrphanSweep::default(), |acc, outcome| match outcome? { + Reclaimed::Swept(bytes) => Ok(OrphanSweep { + prefixes: acc.prefixes, + objects: acc.objects + 1, + bytes: acc.bytes.saturating_add(bytes), + }), + Reclaimed::Spared => Ok(acc), + }) + } + pub fn sweep_orphans( &self, hosted: &HashSet, grace: Duration, now: SystemTime, ) -> Result { - let root = self.root.as_path(); let expected: HashSet = hosted .iter() .filter_map(|repo| RepoPrefix::new(repo).ok()) .map(|prefix| prefix.as_path().to_path_buf()) .collect(); - let orphans: HashSet = discover_prefixes(root)? - .into_iter() - .filter(|prefix| { - prefix - .strip_prefix(root) - .ok() - .filter(|rel| matches!(rel.components().count(), 2 | 3)) - .map(|rel| !expected.contains(rel)) - .unwrap_or(false) - }) - .collect(); - orphans + Pool::ALL .iter() - .map(|prefix| self.reclaim_orphan(prefix, grace, now)) - .try_fold(OrphanSweep::default(), |acc, outcome| { - let outcome = outcome?; + .map(|pool| { + let pool_root = self.root.as_path().join(pool.dir()); + let orphans: HashSet = discover_prefixes(&pool_root)? + .into_iter() + .filter(|prefix| { + prefix + .strip_prefix(&pool_root) + .ok() + .filter(|rel| matches!(rel.components().count(), 2 | 3)) + .map(|rel| !expected.contains(rel)) + .unwrap_or(false) + }) + .collect(); + orphans + .iter() + .map(|prefix| self.reclaim_orphan(prefix, grace, now)) + .try_fold(OrphanSweep::default(), |acc, outcome| { + let outcome = outcome?; + Ok(OrphanSweep { + prefixes: acc.prefixes + outcome.prefixes, + objects: acc.objects + outcome.objects, + bytes: acc.bytes.saturating_add(outcome.bytes), + }) + }) + }) + .try_fold(OrphanSweep::default(), |acc, per_pool| { + let swept = per_pool?; Ok(OrphanSweep { - prefixes: acc.prefixes + outcome.prefixes, - objects: acc.objects + outcome.objects, - bytes: acc.bytes.saturating_add(outcome.bytes), + prefixes: acc.prefixes + swept.prefixes, + objects: acc.objects + swept.objects, + bytes: acc.bytes.saturating_add(swept.bytes), }) }) } @@ -514,12 +634,13 @@ impl MemoryStore { impl LfsStore for MemoryStore { fn put( &self, + pool: Pool, repo: &RepoDid, oid: &LfsOid, size: ClaimedSize, body: &mut dyn Read, ) -> Result<(), LfsError> { - let rel = ObjectRelPath::new(repo, oid)?; + let rel = ObjectRelPath::new(pool, repo, oid)?; let mut bytes = Vec::new(); write_verified(oid, size, body, |chunk| { bytes.extend_from_slice(chunk); @@ -529,8 +650,13 @@ impl LfsStore for MemoryStore { Ok(()) } - fn read(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - let rel = ObjectRelPath::new(repo, oid)?; + fn read( + &self, + pool: Pool, + repo: &RepoDid, + oid: &LfsOid, + ) -> Result, LfsError> { + let rel = ObjectRelPath::new(pool, repo, oid)?; self.locked() .get(&rel) .cloned() @@ -538,16 +664,16 @@ impl LfsStore for MemoryStore { .ok_or_else(|| LfsError::NotFound { oid: oid.clone() }) } - fn probe(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - let rel = ObjectRelPath::new(repo, oid)?; + fn probe(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { + let rel = ObjectRelPath::new(pool, repo, oid)?; Ok(self .locked() .get(&rel) .map(|bytes| LfsSize::new(bytes.len() as u64))) } - fn touch(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - self.probe(repo, oid) + fn touch(&self, pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { + self.probe(pool, repo, oid) } } diff --git a/knot2/crates/knot-lfs/src/transfer.rs b/knot2/crates/knot-lfs/src/transfer.rs index 0dadc5715..fd05aa576 100644 --- a/knot2/crates/knot-lfs/src/transfer.rs +++ b/knot2/crates/knot-lfs/src/transfer.rs @@ -10,7 +10,7 @@ use knot_types::{HttpStatus, RepoDid}; use crate::store::for_each_chunk; use crate::{ - BatchObject, ClaimedSize, LfsError, LfsOid, LfsSize, LfsStore, MAX_BATCH_OBJECTS, + BatchObject, ClaimedSize, LfsError, LfsOid, LfsSize, LfsStore, MAX_BATCH_OBJECTS, Pool, UploadAdmission, }; @@ -50,6 +50,7 @@ pub fn serve_transfer( store, admission, repo, + pool: Pool::Lfs, op, messages, lines: StreamingPeekableIter::new(input, &[], false), @@ -191,6 +192,7 @@ struct Session<'a, R: Read, W: Write> { store: &'a dyn LfsStore, admission: &'a dyn UploadAdmission, repo: &'a RepoDid, + pool: Pool, op: TransferOp, messages: &'a LfsMessages, lines: StreamingPeekableIter, @@ -407,8 +409,8 @@ impl Session<'_, R, W> { fn batch_line(&self, item: &BatchObject) -> Result { let stored = match self.op { - TransferOp::Upload => self.store.touch(self.repo, &item.oid)?, - TransferOp::Download => self.store.probe(self.repo, &item.oid)?, + TransferOp::Upload => self.store.touch(self.pool, self.repo, &item.oid)?, + TransferOp::Download => self.store.probe(self.pool, self.repo, &item.oid)?, }; let (size, action) = match (self.op, stored) { (TransferOp::Upload, Some(_)) => (item.size.get(), "noop"), @@ -456,7 +458,9 @@ impl Session<'_, R, W> { offset: 0, done: false, }; - let stored = self.store.put(self.repo, &oid, declared, &mut body); + let stored = self + .store + .put(self.pool, self.repo, &oid, declared, &mut body); drop(permit); let synced = body.done; if !synced { @@ -494,7 +498,7 @@ impl Session<'_, R, W> { .and_then(|value| value.parse().ok()) .map(ClaimedSize::new); let verdict = LfsOid::new(rest).and_then(|oid| { - match (self.store.touch(self.repo, &oid)?, declared) { + match (self.store.touch(self.pool, self.repo, &oid)?, declared) { (None, _) => Err(LfsError::NotFound { oid }), (Some(actual), Some(declared)) if !declared.matches(actual) => { Err(LfsError::SizeMismatch { @@ -522,9 +526,9 @@ impl Session<'_, R, W> { let opened = LfsOid::new(rest).and_then(|oid| { let size = self .store - .probe(self.repo, &oid)? + .probe(self.pool, self.repo, &oid)? .ok_or(LfsError::NotFound { oid: oid.clone() })?; - let body = self.store.read(self.repo, &oid)?; + let body = self.store.read(self.pool, self.repo, &oid)?; Ok((size, body)) }); let (size, mut body) = match opened { @@ -606,14 +610,12 @@ impl Read for PktBody<'_, R> { #[cfg(test)] mod tests { - use sha2::{Digest, Sha256}; - use super::*; use crate::admission::Unbounded; use crate::{FreeSpaceFloor, MemoryStore}; fn oid_of(bytes: &[u8]) -> LfsOid { - LfsOid::from_digest(Sha256::digest(bytes).into()) + LfsOid::from_digest(knot_types::Sha256Digest::hash(bytes)) } fn repo() -> RepoDid { @@ -712,6 +714,7 @@ mod tests { let seeded_oid = oid_of(seeded); store .put( + Pool::Lfs, &repo(), &seeded_oid, ClaimedSize::new(seeded.len() as u64), @@ -770,7 +773,7 @@ mod tests { assert_eq!(turns[4], vec![line("status 200"), Out::Flush]); assert_eq!(turns[5], vec![line("status 200"), Out::Flush]); assert_eq!( - store.probe(&repo(), &fresh_oid).unwrap(), + store.probe(Pool::Lfs, &repo(), &fresh_oid).unwrap(), Some(LfsSize::new(fresh_len as u64)) ); } @@ -783,6 +786,7 @@ mod tests { let absent = oid_of(b"never uploaded"); store .put( + Pool::Lfs, &repo(), &media_oid, ClaimedSize::new(media.len() as u64), @@ -849,7 +853,7 @@ mod tests { let turns = responses(&out); assert_eq!(turns[1][0], line("status 403")); assert_eq!(turns[2][0], line("status 403")); - assert_eq!(store.probe(&repo(), &oid).unwrap(), None); + assert_eq!(store.probe(Pool::Lfs, &repo(), &oid).unwrap(), None); let mut upload_script = Vec::new(); msg(&mut upload_script, &format!("get-object {oid}"), &[]); @@ -903,7 +907,7 @@ mod tests { let turns = responses(&parse_out(&output)); assert_eq!(turns[1][0], line("status 413")); assert_eq!(turns[2][0], line("status 429")); - assert_eq!(store.probe(&repo(), &oid).unwrap(), None); + assert_eq!(store.probe(Pool::Lfs, &repo(), &oid).unwrap(), None); } #[test] @@ -950,10 +954,10 @@ mod tests { "retry with matching bytes succeeds" ); assert_eq!( - store.probe(&repo(), &valid).unwrap(), + store.probe(Pool::Lfs, &repo(), &valid).unwrap(), Some(LfsSize::new(body.len() as u64)) ); - assert_eq!(store.probe(&repo(), &liar).unwrap(), None); + assert_eq!(store.probe(Pool::Lfs, &repo(), &liar).unwrap(), None); } #[test] @@ -1084,7 +1088,7 @@ mod tests { matches!(overshoot, Err(LfsError::Framing { .. })), "a body overshooting the object size limit ends the session instead of draining unbounded bytes" ); - assert_eq!(store.probe(&repo(), &oid).unwrap(), None); + assert_eq!(store.probe(Pool::Lfs, &repo(), &oid).unwrap(), None); } #[test] @@ -1093,6 +1097,7 @@ mod tests { impl LfsStore for BrokenStore { fn put( &self, + _pool: Pool, _repo: &RepoDid, _oid: &LfsOid, _size: ClaimedSize, @@ -1104,13 +1109,19 @@ mod tests { fn read( &self, + _pool: Pool, _repo: &RepoDid, _oid: &LfsOid, ) -> Result, LfsError> { unreachable!("this session never reads") } - fn probe(&self, _repo: &RepoDid, _oid: &LfsOid) -> Result, LfsError> { + fn probe( + &self, + _pool: Pool, + _repo: &RepoDid, + _oid: &LfsOid, + ) -> Result, LfsError> { Err(LfsError::Io { op: "stat", path: "/srv/secret-lfs-root/plc/sq/uid".into(), @@ -1118,8 +1129,13 @@ mod tests { }) } - fn touch(&self, repo: &RepoDid, oid: &LfsOid) -> Result, LfsError> { - self.probe(repo, oid) + fn touch( + &self, + pool: Pool, + repo: &RepoDid, + oid: &LfsOid, + ) -> Result, LfsError> { + self.probe(pool, repo, oid) } } diff --git a/knot2/crates/knot-lfs/src/types.rs b/knot2/crates/knot-lfs/src/types.rs index d91c1772f..91e1470c7 100644 --- a/knot2/crates/knot-lfs/src/types.rs +++ b/knot2/crates/knot-lfs/src/types.rs @@ -27,8 +27,25 @@ impl LfsOid { } } - pub fn from_digest(digest: [u8; 32]) -> Self { - Self(knot_types::lowercase_hex(&digest)) + pub fn from_digest(digest: knot_types::Sha256Digest) -> Self { + Self(knot_types::lowercase_hex(&digest.bytes())) + } + + pub fn digest(&self) -> knot_types::Sha256Digest { + knot_types::Sha256Digest::new( + self.0 + .as_bytes() + .as_chunks::<2>() + .0 + .iter() + .map(|pair| { + u8::from_str_radix(std::str::from_utf8(pair).expect("an oid is hex"), 16) + .expect("an oid is hex") + }) + .collect::>() + .try_into() + .expect("an oid digest is 32 bytes"), + ) } pub fn as_str(&self) -> &str { @@ -155,6 +172,23 @@ impl fmt::Display for FreeSpaceFloor { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +pub enum Pool { + Lfs, + Attachments, +} + +impl Pool { + pub(crate) const fn dir(self) -> &'static str { + match self { + Pool::Lfs => ".lfs", + Pool::Attachments => ".attachments", + } + } + + pub const ALL: [Pool; 2] = [Pool::Lfs, Pool::Attachments]; +} + #[derive(Debug, Clone, PartialEq, Eq, Hash)] pub struct RepoPrefix(PathBuf); @@ -176,10 +210,16 @@ impl RepoPrefix { pub struct ObjectRelPath(PathBuf); impl ObjectRelPath { - pub fn new(repo: &RepoDid, oid: &LfsOid) -> Result { + pub fn new(pool: Pool, repo: &RepoDid, oid: &LfsOid) -> Result { let prefix = RepoPrefix::new(repo)?; let hex = oid.as_str(); - Ok(Self(prefix.0.join(&hex[0..2]).join(&hex[2..4]).join(hex))) + Ok(Self( + Path::new(pool.dir()) + .join(prefix.0) + .join(&hex[0..2]) + .join(&hex[2..4]) + .join(hex), + )) } pub fn as_path(&self) -> &Path { @@ -217,8 +257,12 @@ mod tests { assert_eq!(oid.to_string(), SAMPLE_OID); assert_eq!(SAMPLE_OID.parse::().unwrap(), oid); - let from_digest = LfsOid::from_digest([0xab; 32]); + let from_digest = LfsOid::from_digest(knot_types::Sha256Digest::new([0xab; 32])); assert_eq!(from_digest.as_str().len(), 64); + assert_eq!( + from_digest.digest(), + knot_types::Sha256Digest::new([0xab; 32]) + ); assert_eq!(LfsOid::new(from_digest.as_str()).unwrap(), from_digest); } @@ -244,21 +288,28 @@ mod tests { } #[test] - fn object_paths_shard_by_repo_then_oid_and_stay_under_the_root() { + fn object_paths_shard_by_pool_repo_then_oid_and_stay_under_root() { let repo = RepoDid::new("did:plc:squid").unwrap(); let oid = LfsOid::new(SAMPLE_OID).unwrap(); - let rel = ObjectRelPath::new(&repo, &oid).unwrap(); + let rel = ObjectRelPath::new(Pool::Lfs, &repo, &oid).unwrap(); assert_eq!( rel.as_path(), - Path::new("plc/sq/uid/6c/17").join(SAMPLE_OID) + Path::new(".lfs/plc/sq/uid/6c/17").join(SAMPLE_OID) + ); + let attachments = ObjectRelPath::new(Pool::Attachments, &repo, &oid).unwrap(); + assert_eq!( + attachments.as_path(), + Path::new(".attachments/plc/sq/uid/6c/17").join(SAMPLE_OID) ); let root = LfsStorePath::new("/srv/lfs"); - let path = root.object_path(&rel); - assert!(path.starts_with(root.as_path())); - assert!( - path.components() - .all(|part| !matches!(part, std::path::Component::ParentDir)) - ); + Pool::ALL.iter().for_each(|pool| { + let path = root.object_path(&ObjectRelPath::new(*pool, &repo, &oid).unwrap()); + assert!(path.starts_with(root.as_path())); + assert!( + path.components() + .all(|part| !matches!(part, std::path::Component::ParentDir)) + ); + }); } } diff --git a/knot2/crates/knot-xrpc/src/lfs.rs b/knot2/crates/knot-xrpc/src/lfs.rs index fbfc4e3cf..13916b95e 100644 --- a/knot2/crates/knot-xrpc/src/lfs.rs +++ b/knot2/crates/knot-xrpc/src/lfs.rs @@ -12,7 +12,7 @@ use axum::routing::{get, post}; use knot_lfs::{ BATCH_MEDIA_TYPE, BatchAction, BatchActions, BatchObject, BatchObjectError, BatchOperation, BatchRequest, BatchResponse, BatchResponseObject, ClaimedSize, LfsError, LfsHandle, LfsOid, - LfsSize, LfsStore, MAX_BATCH_OBJECTS, UploadAdmission, + LfsSize, LfsStore, MAX_BATCH_OBJECTS, Pool, UploadAdmission, }; use knot_pack::{HaveOids, SocketPeer, WantOids}; use knot_runtime::{Clock, HttpTransport}; @@ -409,7 +409,7 @@ async fn probe_all( let oids: Vec = objects.iter().map(|object| object.oid.clone()).collect(); tokio::task::spawn_blocking(move || { oids.iter() - .map(|oid| store.probe(&target, oid)) + .map(|oid| store.probe(Pool::Lfs, &target, oid)) .collect::, _>>() }) .await @@ -517,7 +517,7 @@ async fn serve_object( let store = Arc::clone(&lfs.handle.store); let target = repo.clone(); let oid = oid.clone(); - tokio::task::spawn_blocking(move || store.object_file(&target, &oid)).await + tokio::task::spawn_blocking(move || store.object_file(Pool::Lfs, &target, &oid)).await }; let (size, path) = match located { Ok(Ok(Some((size, path)))) => (size, path), @@ -644,7 +644,7 @@ async fn serve_object_upload( ); tokio::task::spawn_blocking(move || { let mut body = tokio_util::io::SyncIoBridge::new(reader); - let outcome = store.put(&target, &object, size, &mut body); + let outcome = store.put(Pool::Lfs, &target, &object, size, &mut body); drop(permit); outcome }) @@ -769,7 +769,7 @@ async fn mirror_fork_objects_inner( knot_lfs::scan_pointers(&repo, wants.wants(), haves.haves()) .map_err(|error| error.to_string())? .into_iter() - .map(|(oid, size)| match store.probe(&fork, &oid) { + .map(|(oid, size)| match store.probe(Pool::Lfs, &fork, &oid) { Ok(None) => Ok(Some((oid, size))), Ok(Some(_)) => Ok(None), Err(fault) => Err(fault.to_string()), @@ -808,9 +808,9 @@ async fn mirror_fork_objects_inner( .into_iter() .filter_map(|(oid, size)| { let copied = admission.admit(size).and_then(|_permit| { - store - .read(&source, &oid) - .and_then(|mut body| store.put(&fork, &oid, size, &mut body)) + store.read(Pool::Lfs, &source, &oid).and_then(|mut body| { + store.put(Pool::Lfs, &fork, &oid, size, &mut body) + }) }); match copied { Ok(()) => None, @@ -1028,7 +1028,7 @@ async fn fetch_remote_object( let oid = object.oid.clone(); tokio::task::spawn_blocking(move || { let mut body = tokio_util::io::SyncIoBridge::new(reader); - let outcome = store.put(&fork, &oid, declared, &mut body); + let outcome = store.put(Pool::Lfs, &fork, &oid, declared, &mut body); drop(permit); outcome }) @@ -1098,18 +1098,17 @@ mod tests { } #[test] - fn oids_the_upstream_batch_never_answers_count_as_missing() { - use sha2::{Digest, Sha256}; - let held = LfsOid::from_digest(Sha256::digest(b"held").into()); - let ignored = LfsOid::from_digest(Sha256::digest(b"ignored").into()); + fn oids_absent_from_upstream_batch_count_as_missing() { + let stored = LfsOid::from_digest(knot_types::Sha256Digest::hash(b"stored")); + let ignored = LfsOid::from_digest(knot_types::Sha256Digest::hash(b"ignored")); let declared: BTreeMap = [ - (held.clone(), ClaimedSize::new(4)), + (stored.clone(), ClaimedSize::new(4)), (ignored.clone(), ClaimedSize::new(7)), ] .into_iter() .collect(); let answered = vec![BatchResponseObject { - oid: held, + oid: stored, size: ClaimedSize::new(4), authenticated: None, actions: None, diff --git a/knot2/crates/knot-xrpc/tests/lfs.rs b/knot2/crates/knot-xrpc/tests/lfs.rs index 7893ce30c..b61b47426 100644 --- a/knot2/crates/knot-xrpc/tests/lfs.rs +++ b/knot2/crates/knot-xrpc/tests/lfs.rs @@ -5,7 +5,7 @@ use http::{Request, StatusCode, header}; use sha2::{Digest, Sha256}; use tower::ServiceExt; -use knot_lfs::{LfsOid, LfsStore}; +use knot_lfs::{LfsOid, LfsStore, Pool}; use knot_types::RepoDid; use common::{OWNER, World, empty_repo, get}; @@ -13,7 +13,7 @@ use common::{OWNER, World, empty_repo, get}; const MEDIA: &[u8] = b"\xff\x00heavy media bytes that live outside the odb"; fn seeded_object(world: &World, repo: &RepoDid) -> (LfsOid, usize) { - let oid = LfsOid::from_digest(Sha256::digest(MEDIA).into()); + let oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(MEDIA)); world .state .lfs @@ -22,6 +22,7 @@ fn seeded_object(world: &World, repo: &RepoDid) -> (LfsOid, usize) { .handle .store .put( + Pool::Lfs, repo, &oid, knot_lfs::ClaimedSize::new(MEDIA.len() as u64), @@ -32,7 +33,7 @@ fn seeded_object(world: &World, repo: &RepoDid) -> (LfsOid, usize) { } fn absent_oid() -> LfsOid { - LfsOid::from_digest(Sha256::digest(b"never uploaded anywhere").into()) + LfsOid::from_digest(knot_types::Sha256Digest::hash(b"never uploaded anywhere")) } async fn post_batch(world: &World, path: &str, body: String) -> (StatusCode, serde_json::Value) { @@ -249,7 +250,12 @@ async fn hostile_batches_get_clean_typed_rejections() { let oversized_list: Vec<(String, u64)> = (0..1001) .map(|index| { let digest = Sha256::digest(index.to_string().as_bytes()); - (LfsOid::from_digest(digest.into()).as_str().to_string(), 1) + ( + LfsOid::from_digest(knot_types::Sha256Digest::new(digest.into())) + .as_str() + .to_string(), + 1, + ) }) .collect(); let (status, _) = post_batch(&world, &base, batch_body("download", &oversized_list)).await; diff --git a/knot2/crates/knot-xrpc/tests/lfs_soak.rs b/knot2/crates/knot-xrpc/tests/lfs_soak.rs index 3e83362b9..131271920 100644 --- a/knot2/crates/knot-xrpc/tests/lfs_soak.rs +++ b/knot2/crates/knot-xrpc/tests/lfs_soak.rs @@ -3,12 +3,11 @@ mod common; use axum::body::Body; use futures::StreamExt; use http::{Request, StatusCode, header}; +use knot_lfs::{LfsOid, LfsStore, Pool}; +use knot_types::RepoDid; use sha2::{Digest, Sha256}; use tower::ServiceExt; -use knot_lfs::{LfsOid, LfsStore}; -use knot_types::RepoDid; - use common::{World, empty_repo}; const OBJECT_BYTES: usize = 8 * 1024 * 1024; @@ -46,9 +45,10 @@ fn seed_objects(world: &World, repo: &RepoDid) -> Vec { (0..OBJECTS) .map(|index| { let body = incompressible(OBJECT_BYTES, 0x5eed_0000 + index as u64); - let oid = LfsOid::from_digest(Sha256::digest(&body).into()); + let oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&body)); store .put( + Pool::Lfs, repo, &oid, knot_lfs::ClaimedSize::new(body.len() as u64), @@ -80,7 +80,7 @@ async fn download(world: &World, did: &RepoDid, oid: &LfsOid, tag: &str) { .await; assert_eq!(streamed, OBJECT_BYTES as u64, "{tag}"); assert_eq!( - LfsOid::from_digest(hasher.finalize().into()), + LfsOid::from_digest(knot_types::Sha256Digest::new(hasher.finalize().into())), oid.clone(), "{tag}: downloaded bytes must hash to the requested oid" ); -- 2.51.2