use std::collections::BTreeSet;
use std::fs;
use std::io::{self, Read};
use std::path::{Path, PathBuf};
use super::{CommitHistory, HistoryError, InclusionProof, Position, verify_inclusion};
use crate::hash::{self, Hash};
use crate::layout::RepoLayout;
use crate::object::Object;
use crate::refs::ancestry_state::{self, Transaction};
use crate::refs::{self, RefMutation, RefWriteCondition};
use crate::store::ObjectStore;
const MAGIC: &[u8; 5] = b"MKHA\x01";
pub(crate) const MAX_ANCESTRY_LEAVES: usize = 1_000_000;
const MAX_SNAPSHOT_BYTES: u64 = (MAX_ANCESTRY_LEAVES as u64) * 32 + 8192;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AncestryDescriptor {
pub repository: Hash,
pub full_ref: String,
pub generation: Hash,
pub tip: Hash,
pub leaf_count: u64,
pub root: Hash,
}
#[derive(Debug, Clone)]
pub struct TrustedAncestryDescriptor(AncestryDescriptor);
impl TrustedAncestryDescriptor {
#[must_use]
pub fn descriptor(&self) -> &AncestryDescriptor {
&self.0
}
}
#[derive(Debug)]
pub struct AncestrySnapshot {
descriptor: AncestryDescriptor,
chain: Vec<Hash>,
mmb: CommitHistory,
}
impl AncestrySnapshot {
pub fn load(layout: &RepoLayout, branch: &str) -> Result<Self, HistoryError> {
let (_history_lock, mutation) = refs::acquire_history_mutation(layout, branch)?;
let dir = ancestry_state::branch_dir(layout.common_dir(), &format!("refs/heads/{branch}"));
if Transaction::read(&dir)?.is_some() {
return Err(HistoryError::Corrupted(
"history publication pending; retry the write to recover".into(),
));
}
let repository = read_repository_id(layout.common_dir())?.ok_or_else(|| {
HistoryError::Corrupted("no trusted local ancestry descriptor".into())
})?;
let snapshot = read_current(&dir)?
.ok_or_else(|| HistoryError::Corrupted("no trusted local ancestry snapshot".into()))?;
let current = mutation.current()?;
if snapshot.descriptor.repository != repository
|| snapshot.descriptor.full_ref != format!("refs/heads/{branch}")
|| Some(snapshot.descriptor.tip) != current
{
return Err(HistoryError::Corrupted(
"ancestry descriptor does not match the authoritative ref".into(),
));
}
let store = ObjectStore::open(layout)?;
if snapshot.chain != first_parent_chain(&store, snapshot.descriptor.tip)? {
return Err(HistoryError::Corrupted(
"snapshot is not the current first-parent chain".into(),
));
}
Ok(snapshot)
}
#[must_use]
pub fn descriptor(&self) -> &AncestryDescriptor {
&self.descriptor
}
#[must_use]
pub fn trusted_descriptor(&self) -> TrustedAncestryDescriptor {
TrustedAncestryDescriptor(self.descriptor.clone())
}
#[must_use]
pub fn len(&self) -> u64 {
self.descriptor.leaf_count
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.chain.is_empty()
}
#[must_use]
pub fn root(&self) -> Hash {
self.descriptor.root
}
pub fn prove(&self, position: Position) -> Result<InclusionProof, HistoryError> {
self.mmb.prove(position)
}
#[must_use]
pub fn position_of(&self, commit: &Hash) -> Option<Position> {
self.chain
.iter()
.position(|h| h == commit)
.map(|n| Position(n as u64))
}
fn build(
repository: Hash,
full_ref: String,
generation: Hash,
chain: Vec<Hash>,
) -> Result<Self, HistoryError> {
if chain.is_empty() || chain.len() > MAX_ANCESTRY_LEAVES {
return Err(HistoryError::Corrupted(
"invalid ancestry leaf count".into(),
));
}
let mut mmb = CommitHistory::open();
mmb.extend(&chain)?;
let descriptor = AncestryDescriptor {
repository,
full_ref,
generation,
tip: *chain.last().expect("nonempty chain"),
leaf_count: chain.len() as u64,
root: mmb.root(),
};
Ok(Self {
descriptor,
chain,
mmb,
})
}
fn encode(&self) -> Result<Vec<u8>, HistoryError> {
let d = &self.descriptor;
let name_len = u16::try_from(d.full_ref.len())
.map_err(|_| HistoryError::InvalidBranch(d.full_ref.clone()))?;
let mut bytes = Vec::with_capacity(self.chain.len() * 32 + 192 + d.full_ref.len());
bytes.extend_from_slice(MAGIC);
bytes.extend_from_slice(&d.repository);
bytes.extend_from_slice(&d.generation);
bytes.extend_from_slice(&d.tip);
bytes.extend_from_slice(&d.leaf_count.to_le_bytes());
bytes.extend_from_slice(&d.root);
bytes.extend_from_slice(&name_len.to_le_bytes());
bytes.extend_from_slice(d.full_ref.as_bytes());
for h in &self.chain {
bytes.extend_from_slice(h);
}
bytes.extend_from_slice(&hash::hash(&bytes));
Ok(bytes)
}
fn decode(bytes: &[u8]) -> Result<Self, HistoryError> {
let (claimed, chain) = decode_descriptor_and_chain(bytes)?;
let invalid = || HistoryError::Corrupted("malformed ancestry snapshot".into());
let snapshot = Self::build(
claimed.repository,
claimed.full_ref.clone(),
claimed.generation,
chain,
)?;
if snapshot.descriptor.tip != claimed.tip || snapshot.descriptor.root != claimed.root {
return Err(invalid());
}
Ok(snapshot)
}
}
fn take<'a>(input: &mut &'a [u8], n: usize) -> Option<&'a [u8]> {
if input.len() < n {
return None;
}
let (head, tail) = input.split_at(n);
*input = tail;
Some(head)
}
fn descriptor_header_invalid() -> HistoryError {
HistoryError::Corrupted("malformed ancestry snapshot".into())
}
fn parse_descriptor_header(input: &mut &[u8]) -> Result<AncestryDescriptor, HistoryError> {
fn digest(input: &mut &[u8]) -> Option<Hash> {
take(input, 32)?.try_into().ok()
}
let invalid = descriptor_header_invalid;
if take(input, 5) != Some(MAGIC.as_slice()) {
return Err(invalid());
}
let repository = digest(input).ok_or_else(invalid)?;
let generation = digest(input).ok_or_else(invalid)?;
let tip = digest(input).ok_or_else(invalid)?;
let count = u64::from_le_bytes(
take(input, 8)
.ok_or_else(invalid)?
.try_into()
.map_err(|_| invalid())?,
);
let root = digest(input).ok_or_else(invalid)?;
let name_len = u16::from_le_bytes(
take(input, 2)
.ok_or_else(invalid)?
.try_into()
.map_err(|_| invalid())?,
) as usize;
let full_ref = std::str::from_utf8(take(input, name_len).ok_or_else(invalid)?)
.map_err(|_| invalid())?
.to_owned();
if !full_ref.starts_with("refs/heads/")
|| !refs::validate_ref_name_grammar(&full_ref)
|| count == 0
|| count > MAX_ANCESTRY_LEAVES as u64
{
return Err(invalid());
}
Ok(AncestryDescriptor {
repository,
full_ref,
generation,
tip,
leaf_count: count,
root,
})
}
fn decode_descriptor_and_chain(
bytes: &[u8],
) -> Result<(AncestryDescriptor, Vec<Hash>), HistoryError> {
let invalid = descriptor_header_invalid;
let len = bytes.len() as u64;
if !(DESCRIPTOR_HEADER_LEN + 32..=MAX_SNAPSHOT_BYTES).contains(&len) {
return Err(invalid());
}
let (payload, checksum) = bytes.split_at(bytes.len() - 32);
if hash::hash(payload).as_slice() != checksum {
return Err(invalid());
}
let mut input = payload;
let descriptor = parse_descriptor_header(&mut input)?;
if input.len() as u64 != descriptor.leaf_count * 32 {
return Err(invalid());
}
let chain: Vec<Hash> = input
.chunks_exact(32)
.map(|c| c.try_into().expect("32-byte chunk"))
.collect();
Ok((descriptor, chain))
}
const DESCRIPTOR_HEADER_LEN: u64 = 143;
const DESCRIPTOR_HEADER_MAX_LEN: u64 = DESCRIPTOR_HEADER_LEN + u16::MAX as u64;
fn read_prefix(path: &Path, max_bytes: u64) -> Result<Option<Vec<u8>>, HistoryError> {
let file = match fs::File::open(path) {
Ok(file) => file,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(HistoryError::Io(e)),
};
let mut bytes = Vec::with_capacity(usize::try_from(max_bytes).unwrap_or(usize::MAX));
file.take(max_bytes)
.read_to_end(&mut bytes)
.map_err(HistoryError::Io)?;
Ok(Some(bytes))
}
#[must_use]
pub fn verify_ancestry(
commit: &Hash,
position: Position,
proof: &InclusionProof,
claimed: &AncestryDescriptor,
trusted: &TrustedAncestryDescriptor,
expected: &AncestryDescriptor,
) -> bool {
claimed == trusted.descriptor()
&& claimed == expected
&& position.0 < claimed.leaf_count
&& verify_inclusion(commit, position, proof, &claimed.root)
}
fn ancestry_cycle_or_limit() -> HistoryError {
HistoryError::Corrupted("ancestry cycle or traversal limit".into())
}
fn ancestry_not_commit_or_remix() -> HistoryError {
HistoryError::Corrupted("ancestry node is not a commit/remix".into())
}
fn first_parent_chain(store: &ObjectStore, tip: Hash) -> Result<Vec<Hash>, HistoryError> {
let mut chain = Vec::new();
let mut seen = BTreeSet::new();
let mut next = Some(tip);
while let Some(h) = next {
if chain.len() >= MAX_ANCESTRY_LEAVES || !seen.insert(h) {
return Err(ancestry_cycle_or_limit());
}
chain.push(h);
next = match store.read_object(&h)? {
Object::Commit(c) => c.parents.first().copied(),
Object::Remix(r) => r.parents.first().copied(),
_ => return Err(ancestry_not_commit_or_remix()),
};
}
chain.reverse();
Ok(chain)
}
enum SuffixWalk {
Reached(Vec<Hash>),
NotFound,
}
fn first_parent_suffix_to(
store: &ObjectStore,
target: Hash,
stop_at: Hash,
prefix_len: usize,
) -> Result<SuffixWalk, HistoryError> {
let mut suffix = Vec::new();
let mut seen = BTreeSet::new();
let mut next = Some(target);
while let Some(h) = next {
if h == stop_at {
suffix.reverse();
return Ok(SuffixWalk::Reached(suffix));
}
if suffix.len() + prefix_len >= MAX_ANCESTRY_LEAVES || !seen.insert(h) {
return Err(ancestry_cycle_or_limit());
}
suffix.push(h);
next = match store.read_object(&h)? {
Object::Commit(c) => c.parents.first().copied(),
Object::Remix(r) => r.parents.first().copied(),
_ => return Err(ancestry_not_commit_or_remix()),
};
}
Ok(SuffixWalk::NotFound)
}
fn fresh_id() -> Result<Hash, HistoryError> {
let mut id = [0; 32];
getrandom::fill(&mut id).map_err(|e| std::io::Error::other(e.to_string()))?;
Ok(id)
}
fn read_repository_id(common: &Path) -> Result<Option<Hash>, HistoryError> {
let Some(bytes) = ancestry_state::read_bounded(
&common.join(ancestry_state::DIRECTORY).join("repository-id"),
65,
)?
else {
return Ok(None);
};
refs::decode_ref_wire(&bytes)
.map(Some)
.ok_or_else(|| HistoryError::Corrupted("malformed history repository identity".into()))
}
fn repository_id(common: &Path) -> Result<Hash, HistoryError> {
if let Some(id) = read_repository_id(common)? {
return Ok(id);
}
let path = common.join(ancestry_state::DIRECTORY).join("repository-id");
let id = fresh_id()?;
crate::atomic::write_create_new(&path, &refs::encode_ref_wire(&id), true)?;
crate::atomic::sync_dir(common)?;
read_repository_id(common)?
.ok_or_else(|| HistoryError::Corrupted("history repository identity disappeared".into()))
}
fn snapshot_path(dir: &Path, generation: Hash) -> PathBuf {
dir.join("generations")
.join(format!("{}.snapshot", hash::to_hex(&generation)))
}
fn read_current(dir: &Path) -> Result<Option<AncestrySnapshot>, HistoryError> {
let Some(bytes) = ancestry_state::read_bounded(&dir.join("current"), 65)? else {
return Ok(None);
};
let generation = refs::decode_ref_wire(&bytes)
.ok_or_else(|| HistoryError::Corrupted("malformed history generation pointer".into()))?;
let raw = ancestry_state::read_bounded(&snapshot_path(dir, generation), MAX_SNAPSHOT_BYTES)?
.ok_or_else(|| HistoryError::Corrupted("missing ancestry generation snapshot".into()))?;
let snapshot = AncestrySnapshot::decode(&raw)?;
if snapshot.descriptor.generation != generation {
return Err(HistoryError::Corrupted(
"ancestry generation mismatch".into(),
));
}
Ok(Some(snapshot))
}
fn read_current_descriptor(dir: &Path) -> Result<Option<AncestryDescriptor>, HistoryError> {
let Some(bytes) = ancestry_state::read_bounded(&dir.join("current"), 65)? else {
return Ok(None);
};
let generation = refs::decode_ref_wire(&bytes)
.ok_or_else(|| HistoryError::Corrupted("malformed history generation pointer".into()))?;
let Some(prefix) = read_prefix(&snapshot_path(dir, generation), DESCRIPTOR_HEADER_MAX_LEN)?
else {
return Err(HistoryError::Corrupted(
"missing ancestry generation snapshot".into(),
));
};
let mut input = prefix.as_slice();
let descriptor = parse_descriptor_header(&mut input)?;
if descriptor.generation != generation {
return Err(HistoryError::Corrupted(
"ancestry generation mismatch".into(),
));
}
Ok(Some(descriptor))
}
fn read_snapshot_chain(dir: &Path, generation: Hash) -> Result<Vec<Hash>, HistoryError> {
let raw = ancestry_state::read_bounded(&snapshot_path(dir, generation), MAX_SNAPSHOT_BYTES)?
.ok_or_else(|| HistoryError::Corrupted("missing ancestry generation snapshot".into()))?;
let (descriptor, chain) = decode_descriptor_and_chain(&raw)?;
if descriptor.generation != generation {
return Err(HistoryError::Corrupted(
"ancestry generation mismatch".into(),
));
}
Ok(chain)
}
fn verify_scrub_window(
store: &ObjectStore,
prefix: &[Hash],
start: u64,
end: u64,
) -> Result<(), HistoryError> {
let start = usize::try_from(start).expect("bounded by MAX_ANCESTRY_LEAVES");
let end = usize::try_from(end).expect("bounded by MAX_ANCESTRY_LEAVES");
let window = prefix.get(start..end).ok_or_else(|| {
HistoryError::Corrupted("scrub window bounds exceed the on-disk ancestry prefix".into())
})?;
for h in window {
match store.read_object(h)? {
Object::Commit(_) | Object::Remix(_) => {}
_ => return Err(ancestry_not_commit_or_remix()),
}
}
Ok(())
}
const SCRUB_MIN_WINDOW: u64 = 512;
const SCRUB_LAP_FRACTION: u64 = 64;
const SCRUB_MAX_AGE_SECS: u64 = 7 * 24 * 60 * 60;
fn scrub_window(verified_through: u64) -> u64 {
(verified_through / SCRUB_LAP_FRACTION).max(SCRUB_MIN_WINDOW)
}
#[cfg(test)]
thread_local! { static NOW_OVERRIDE: std::cell::Cell<Option<u64>> = const { std::cell::Cell::new(None) }; }
fn now_unix() -> u64 {
#[cfg(test)]
if let Some(t) = NOW_OVERRIDE.with(std::cell::Cell::get) {
return t;
}
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ScrubState {
generation: Hash,
cursor: u64,
verified_through: u64,
last_full_verify_unix: u64,
}
const SCRUB_MAGIC: &[u8; 5] = b"MKSC\x02";
impl ScrubState {
fn fresh(chain_len: u64, now: u64) -> Self {
Self {
generation: [0; 32],
cursor: 0,
verified_through: chain_len,
last_full_verify_unix: now,
}
}
fn encode(self) -> [u8; 93] {
let mut bytes = [0u8; 93];
bytes[..5].copy_from_slice(SCRUB_MAGIC);
bytes[5..37].copy_from_slice(&self.generation);
bytes[37..45].copy_from_slice(&self.cursor.to_le_bytes());
bytes[45..53].copy_from_slice(&self.verified_through.to_le_bytes());
bytes[53..61].copy_from_slice(&self.last_full_verify_unix.to_le_bytes());
let checksum = hash::hash(&bytes[..61]);
bytes[61..93].copy_from_slice(&checksum);
bytes
}
fn decode(bytes: &[u8]) -> Option<Self> {
if bytes.len() != 93 || bytes[..5] != *SCRUB_MAGIC {
return None;
}
let (payload, checksum) = bytes.split_at(61);
if hash::hash(payload).as_slice() != checksum {
return None;
}
let field = |r: std::ops::Range<usize>| -> Option<u64> {
Some(u64::from_le_bytes(payload[r].try_into().ok()?))
};
Some(Self {
generation: payload[5..37].try_into().ok()?,
cursor: field(37..45)?,
verified_through: field(45..53)?,
last_full_verify_unix: field(53..61)?,
})
}
}
fn scrub_path(dir: &Path) -> PathBuf {
dir.join("scrub")
}
fn read_scrub_state(dir: &Path) -> Option<ScrubState> {
ancestry_state::read_bounded(&scrub_path(dir), 93)
.ok()
.flatten()
.and_then(|bytes| ScrubState::decode(&bytes))
}
fn write_scrub_state(dir: &Path, state: ScrubState) -> Result<(), HistoryError> {
crate::atomic::write_atomic(&scrub_path(dir), &state.encode(), false)?;
Ok(())
}
fn new_generation(
store: &ObjectStore,
target: Hash,
now: u64,
) -> Result<ChainDecision, HistoryError> {
let chain = first_parent_chain(store, target)?;
let scrub = ScrubState::fresh(chain.len() as u64, now);
Ok(ChainDecision::NewGeneration { chain, scrub })
}
fn decide_chain(
store: &ObjectStore,
dir: &Path,
compatible: Option<&AncestryDescriptor>,
target: Hash,
now: u64,
) -> Result<ChainDecision, HistoryError> {
let Some(d) = compatible else {
return new_generation(store, target, now);
};
let prefix_len = usize::try_from(d.leaf_count).expect("bounded by MAX_ANCESTRY_LEAVES");
let suffix = match first_parent_suffix_to(store, target, d.tip, prefix_len)? {
SuffixWalk::Reached(suffix) => suffix,
SuffixWalk::NotFound => return new_generation(store, target, now),
};
let scrub = read_scrub_state(dir).filter(|s| s.generation == d.generation);
let stale =
scrub.is_none_or(|s| now.saturating_sub(s.last_full_verify_unix) >= SCRUB_MAX_AGE_SECS);
if !stale {
let scrub = scrub.expect("`stale` is false only when `scrub` is Some");
let window = scrub_window(scrub.verified_through);
let end = scrub
.cursor
.saturating_add(window)
.min(scrub.verified_through);
if end < scrub.verified_through {
let prefix = read_snapshot_chain(dir, d.generation)?;
if end <= prefix.len() as u64 {
verify_scrub_window(store, &prefix, scrub.cursor, end)?;
let mut chain = prefix;
chain.extend(suffix);
return Ok(ChainDecision::SameGeneration {
chain,
scrub: ScrubState {
cursor: end,
..scrub
},
});
}
}
}
let mut chain = first_parent_chain(store, d.tip)?;
chain.extend(suffix);
let scrub = ScrubState::fresh(chain.len() as u64, now);
Ok(ChainDecision::SameGeneration { chain, scrub })
}
enum ChainDecision {
SameGeneration { chain: Vec<Hash>, scrub: ScrubState },
NewGeneration { chain: Vec<Hash>, scrub: ScrubState },
}
fn finish(
layout: &RepoLayout,
dir: &Path,
tx: &Transaction,
mutation: &RefMutation,
store: &ObjectStore,
prebuilt: Option<AncestrySnapshot>,
) -> Result<AncestrySnapshot, HistoryError> {
if tx.repository != repository_id(layout.common_dir())?
|| ancestry_state::branch_dir(layout.common_dir(), &tx.full_ref) != dir
{
return Err(HistoryError::Corrupted(
"history transaction context mismatch".into(),
));
}
let current = mutation.current()?;
if current != tx.previous && current != Some(tx.target) {
return Err(HistoryError::Corrupted(
"ref diverged from pending history transaction".into(),
));
}
let snapshot = match prebuilt {
Some(snapshot)
if snapshot.descriptor.repository == tx.repository
&& snapshot.descriptor.full_ref == tx.full_ref
&& snapshot.descriptor.generation == tx.generation
&& snapshot.descriptor.tip == tx.target =>
{
snapshot
}
_ => AncestrySnapshot::build(
tx.repository,
tx.full_ref.clone(),
tx.generation,
first_parent_chain(store, tx.target)?,
)?,
};
let encoded = snapshot.encode()?;
crate::atomic::write_atomic(&dir.join("pending-snapshot"), &encoded, true)?;
checkpoint(2)?;
mutation
.write_preserving_history(&refs::encode_ref_wire(&tx.target), RefWriteCondition::Any)?;
checkpoint(3)?;
let dest = snapshot_path(dir, tx.generation);
fs::create_dir_all(dest.parent().expect("snapshot parent"))?;
fs::rename(dir.join("pending-snapshot"), &dest)?;
crate::atomic::sync_dir(dest.parent().expect("snapshot parent"))?;
crate::atomic::sync_dir(dir)?;
checkpoint(4)?;
crate::atomic::write_atomic(
&dir.join("current"),
&refs::encode_ref_wire(&tx.generation),
true,
)?;
checkpoint(5)?;
ancestry_state::remove_synced(&dir.join("transaction"))?;
checkpoint(6)?;
Ok(snapshot)
}
pub(crate) fn recover(
layout: &RepoLayout,
branch: &str,
mutation: &RefMutation,
store: &ObjectStore,
) -> Result<(), HistoryError> {
let dir = ancestry_state::branch_dir(layout.common_dir(), &format!("refs/heads/{branch}"));
if let Some(tx) = Transaction::read(&dir)? {
finish(layout, &dir, &tx, mutation, store, None)?;
}
Ok(())
}
pub(crate) fn advance(
layout: &RepoLayout,
branch: &str,
mutation: &RefMutation,
condition: RefWriteCondition,
target: Hash,
store: &ObjectStore,
) -> Result<(), HistoryError> {
let full_ref = format!("refs/heads/{branch}");
let dir = ancestry_state::branch_dir(layout.common_dir(), &full_ref);
if let Some(tx) = Transaction::read(&dir)? {
finish(layout, &dir, &tx, mutation, store, None)?;
let retry_of_intent = target == tx.target
&& match condition {
RefWriteCondition::Any => true,
RefWriteCondition::Missing => tx.previous.is_none(),
RefWriteCondition::Match(expected) => tx.previous == Some(expected),
};
if retry_of_intent {
return Ok(());
}
}
mutation.check(condition)?;
let previous = mutation.current()?;
let repository = repository_id(layout.common_dir())?;
let old = read_current_descriptor(&dir)?;
let compatible = old.as_ref().filter(|d| {
d.repository == repository && d.full_ref == full_ref && Some(d.tip) == previous
});
if compatible.is_some_and(|d| d.tip == target) {
return Ok(());
}
let now = now_unix();
let decision = decide_chain(store, &dir, compatible, target, now)?;
let (chain, generation, scrub) = match decision {
ChainDecision::SameGeneration { chain, scrub } => (
chain,
compatible
.expect("SameGeneration is only returned when `compatible` is Some")
.generation,
scrub,
),
ChainDecision::NewGeneration { chain, scrub } => (chain, fresh_id()?, scrub),
};
let tx = Transaction {
repository,
full_ref,
previous,
target,
generation,
previous_generation: compatible.map(|d| d.generation),
};
let snapshot = AncestrySnapshot::build(repository, tx.full_ref.clone(), generation, chain)?;
crate::atomic::write_atomic(&dir.join("transaction"), &tx.encode(), true)?;
for parent in [
dir.parent(),
dir.parent().and_then(Path::parent),
Some(layout.common_dir()),
]
.into_iter()
.flatten()
{
crate::atomic::sync_dir(parent)?;
}
checkpoint(1)?;
finish(layout, &dir, &tx, mutation, store, Some(snapshot))?;
let _ = write_scrub_state(
&dir,
ScrubState {
generation,
..scrub
},
);
Ok(())
}
#[cfg(test)]
thread_local! { static FAIL_AFTER: std::cell::Cell<u8> = const { std::cell::Cell::new(0) }; }
#[cfg_attr(not(test), allow(clippy::unnecessary_wraps))]
fn checkpoint(stage: u8) -> Result<(), HistoryError> {
#[cfg(not(test))]
let _ = stage;
#[cfg(test)]
if FAIL_AFTER.with(|s| s.get() == stage) {
return Err(
std::io::Error::other(format!("injected history publication failure {stage}")).into(),
);
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::object::{Commit, Identity, Tree};
fn repo() -> (tempfile::TempDir, RepoLayout, ObjectStore) {
let dir = tempfile::tempdir().unwrap();
let layout = RepoLayout::single(dir.path());
let store = ObjectStore::init(&layout).unwrap();
refs::init(&layout).unwrap();
(dir, layout, store)
}
fn commit(store: &ObjectStore, parents: Vec<Hash>, message: &[u8]) -> Hash {
let tree = store
.write(&crate::serialize::serialize(&Object::Tree(Tree { entries: vec![] })).unwrap())
.unwrap();
let c = Commit::new_unannotated(
tree,
parents,
Identity::opaque(b"test".to_vec()),
[0; 32],
message.to_vec(),
0,
[0; 64],
);
store
.write(&crate::serialize::serialize(&Object::Commit(c)).unwrap())
.unwrap()
}
fn update(layout: &RepoLayout, store: &ObjectStore, branch: &str, target: Hash) {
refs::update_ref_with_ancestry(layout, branch, RefWriteCondition::Any, &target, store)
.unwrap();
}
#[test]
fn sequential_fast_forward_and_backfill_have_identical_roots_and_positions() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
let b = commit(&store, vec![a], b"b");
let c = commit(&store, vec![b], b"c");
for h in [a, b, c] {
update(&layout, &store, "sequential", h);
}
update(&layout, &store, "fast-forward", a);
update(&layout, &store, "fast-forward", c);
refs::write_ref(&layout, "backfill", &c).unwrap();
update(&layout, &store, "backfill", c);
let seq = AncestrySnapshot::load(&layout, "sequential").unwrap();
for branch in ["fast-forward", "backfill"] {
let snapshot = AncestrySnapshot::load(&layout, branch).unwrap();
assert_eq!(snapshot.root(), seq.root());
assert_eq!(snapshot.len(), 3);
for (i, h) in [a, b, c].iter().enumerate() {
assert_eq!(snapshot.position_of(h), Some(Position(i as u64)));
let proof = snapshot.prove(Position(i as u64)).unwrap();
assert!(verify_ancestry(
h,
Position(i as u64),
&proof,
snapshot.descriptor(),
&snapshot.trusted_descriptor(),
snapshot.descriptor()
));
}
}
}
#[test]
fn generation_changes_on_reset_and_recreation_but_not_noop_or_fast_forward() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
let b = commit(&store, vec![a], b"b");
update(&layout, &store, "main", a);
let original = AncestrySnapshot::load(&layout, "main").unwrap();
update(&layout, &store, "main", a);
assert_eq!(
AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor(),
original.descriptor()
);
update(&layout, &store, "main", b);
assert_eq!(
AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation,
original.descriptor().generation
);
update(&layout, &store, "main", a);
let reset = AncestrySnapshot::load(&layout, "main").unwrap();
assert_ne!(
reset.descriptor().generation,
original.descriptor().generation
);
assert_eq!(reset.root(), original.root());
refs::delete_ref_with_ancestry(&layout, "main", Some(a), &store).unwrap();
assert!(AncestrySnapshot::load(&layout, "main").is_err());
update(&layout, &store, "main", a);
let recreated = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(recreated.root(), original.root());
assert_ne!(
recreated.descriptor().generation,
reset.descriptor().generation
);
}
#[test]
fn merge_ancestry_uses_only_first_parent() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
let b = commit(&store, vec![a], b"b");
let side = commit(&store, vec![a], b"side");
let merge = commit(&store, vec![b, side], b"merge");
update(&layout, &store, "main", merge);
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.chain, vec![a, b, merge]);
assert_eq!(snapshot.position_of(&side), None);
}
#[test]
fn proof_cannot_substitute_repository_ref_generation_tip_count_or_root() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
update(&layout, &store, "main", a);
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
let proof = snapshot.prove(Position(0)).unwrap();
let trusted = snapshot.trusted_descriptor();
for field in 0..6 {
let mut wrong = snapshot.descriptor().clone();
match field {
0 => wrong.repository[0] ^= 1,
1 => wrong.full_ref = "refs/heads/other".into(),
2 => wrong.generation[0] ^= 1,
3 => wrong.tip[0] ^= 1,
4 => wrong.leaf_count += 1,
_ => wrong.root[0] ^= 1,
}
assert!(!verify_ancestry(
&a,
Position(0),
&proof,
&wrong,
&trusted,
&wrong
));
assert!(!verify_ancestry(
&a,
Position(0),
&proof,
snapshot.descriptor(),
&trusted,
&wrong
));
}
let (_foreign_dir, foreign, _) = repo();
assert!(AncestrySnapshot::load(&foreign, "main").is_err());
}
#[test]
fn every_publication_boundary_recovers_the_whole_fast_forward() {
for stage in 1..=6 {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
let b = commit(&store, vec![a], b"b");
let c = commit(&store, vec![b], b"c");
update(&layout, &store, "main", a);
FAIL_AFTER.with(|s| s.set(stage));
let failed = refs::update_ref_with_ancestry(
&layout,
"main",
RefWriteCondition::Match(a),
&c,
&store,
);
FAIL_AFTER.with(|s| s.set(0));
assert!(failed.is_err(), "stage {stage} must inject a failure");
if stage < 6 {
assert!(
AncestrySnapshot::load(&layout, "main").is_err(),
"pending proofs must be withheld"
);
let roots = refs::pending_history_roots(&layout).unwrap();
assert!(roots.contains(&a) && roots.contains(&c));
let live = crate::ops::gc::live_objects(&store, &layout).unwrap();
assert!(live.contains(&b) && live.contains(&c));
assert!(
refs::write_ref(&layout, "main", &a).is_err(),
"raw writer must not bypass recovery"
);
}
if stage < 6 {
refs::update_ref_with_ancestry(
&layout,
"main",
RefWriteCondition::Match(a),
&c,
&store,
)
.expect("retry of the interrupted CAS must finish its original intent");
} else {
update(&layout, &store, "main", c);
}
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.chain, vec![a, b, c], "stage {stage}");
assert!(refs::pending_history_roots(&layout).unwrap().is_empty());
}
}
#[test]
fn missing_ancestor_fails_before_publication() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
update(&layout, &store, "main", a);
let invalid = commit(&store, vec![[91; 32]], b"missing parent");
assert!(
refs::update_ref_with_ancestry(
&layout,
"main",
RefWriteCondition::Any,
&invalid,
&store
)
.is_err()
);
assert_eq!(refs::read_ref(&layout, "main").unwrap(), Some(a));
refs::delete_ref_with_ancestry(&layout, "main", None, &store).unwrap();
}
#[test]
fn raw_aba_mutation_invalidates_the_old_generation() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
let b = commit(&store, vec![a], b"b");
update(&layout, &store, "main", a);
let old = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
refs::write_ref(&layout, "main", &b).unwrap();
refs::write_ref(&layout, "main", &a).unwrap();
assert!(AncestrySnapshot::load(&layout, "main").is_err());
update(&layout, &store, "main", a);
assert_ne!(
AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation,
old
);
}
#[test]
fn tampered_snapshot_and_transaction_fail_closed() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
update(&layout, &store, "main", a);
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let path = snapshot_path(&dir, snapshot.descriptor().generation);
let mut bytes = fs::read(&path).unwrap();
bytes[7] ^= 1;
fs::write(&path, bytes).unwrap();
assert!(AncestrySnapshot::load(&layout, "main").is_err());
fs::write(dir.join("transaction"), b"broken").unwrap();
assert!(refs::pending_history_roots(&layout).is_err());
assert!(crate::ops::gc::live_objects(&store, &layout).is_err());
}
#[test]
fn a_write_scrub_state_failure_does_not_fail_the_publish() {
let (_dir, layout, store) = repo();
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
fs::create_dir_all(dir.join("scrub")).unwrap();
let a = commit(&store, vec![], b"a");
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &a, &store)
.expect("the publish itself must still succeed despite the scrub-state write failing");
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.descriptor().tip, a);
assert!(dir.join("scrub").is_dir());
}
#[test]
fn an_unreadable_scrub_file_does_not_brick_later_publishes() {
let (_dir, layout, store) = repo();
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
fs::create_dir_all(dir.join("scrub")).unwrap();
let a = commit(&store, vec![], b"a");
update(&layout, &store, "main", a);
assert!(dir.join("scrub").is_dir(), "still unreadable as a file");
let b = commit(&store, vec![a], b"b");
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &b, &store)
.expect("an unreadable scrub file must degrade the schedule, not fail the publish");
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.descriptor().tip, b);
}
#[test]
fn scrub_state_ahead_of_the_actual_prefix_falls_back_to_a_full_walk_instead_of_panicking() {
let (_dir, layout, store) = repo();
let a = commit(&store, vec![], b"a");
update(&layout, &store, "main", a);
let b = commit(&store, vec![a], b"b");
update(&layout, &store, "main", b);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let mut scrub = read_scrub_state(&dir).unwrap();
scrub.verified_through = 5_000;
scrub.cursor = 0;
write_scrub_state(&dir, scrub).unwrap();
let c = commit(&store, vec![b], b"c");
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &c, &store)
.expect("an out-of-range scrub window must fall back to a full walk, not panic");
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.descriptor().tip, c);
let after = read_scrub_state(&dir).unwrap();
assert_eq!(after.verified_through, snapshot.descriptor().leaf_count);
}
#[test]
fn read_current_descriptor_matches_full_snapshot_fields() {
for count in [1u64, 2, 50] {
let (_dir, layout, store) = repo();
let mut parents = vec![];
let mut tip = [0; 32];
for i in 0..count {
tip = commit(&store, parents, i.to_be_bytes().as_slice());
parents = vec![tip];
}
update(&layout, &store, "main", tip);
let branch_dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let full = read_current(&branch_dir).unwrap().unwrap();
let lite = read_current_descriptor(&branch_dir).unwrap().unwrap();
assert_eq!(lite, full.descriptor, "count={count}");
assert_eq!(lite.leaf_count, count);
assert_eq!(lite.tip, tip);
}
}
#[test]
fn read_current_descriptor_handles_a_long_branch_name() {
let (_dir, layout, store) = repo();
let long_branch = "b".repeat(200);
let a = commit(&store, vec![], b"a");
update(&layout, &store, &long_branch, a);
let branch_dir =
ancestry_state::branch_dir(layout.common_dir(), &format!("refs/heads/{long_branch}"));
let lite = read_current_descriptor(&branch_dir).unwrap().unwrap();
assert_eq!(lite.full_ref, format!("refs/heads/{long_branch}"));
assert_eq!(lite.tip, a);
assert_eq!(lite.leaf_count, 1);
}
#[test]
fn many_sequential_publishes_keep_one_generation_until_a_real_rewrite() {
let (_dir, layout, store) = repo();
let mut parents = vec![];
let mut tips = Vec::new();
for i in 0..40u64 {
let h = commit(&store, parents, i.to_be_bytes().as_slice());
parents = vec![h];
tips.push(h);
}
update(&layout, &store, "main", tips[0]);
let first_generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
for (i, &h) in tips.iter().enumerate().skip(1) {
update(&layout, &store, "main", h);
let generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
assert_eq!(
generation, first_generation,
"fast-forward to tips[{i}] must keep the same generation as the first publish"
);
}
let final_generation = first_generation;
update(&layout, &store, "main", tips[10]);
let reset_generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
assert_ne!(reset_generation, final_generation);
}
#[test]
#[ignore = "manual profiling tool, not a correctness assertion"]
fn profile_read_current_descriptor_vs_full_chain_read() {
const LEAVES: u64 = 50_000;
const ITERS: u32 = 200;
let (_dir, layout, store) = repo();
let tree = store
.write(&crate::serialize::serialize(&Object::Tree(Tree { entries: vec![] })).unwrap())
.unwrap();
let mut parents = vec![];
let mut tip = [0; 32];
for i in 0..LEAVES {
let c = Commit::new_unannotated(
tree,
parents,
Identity::opaque(b"profile".to_vec()),
[0; 32],
i.to_be_bytes().to_vec(),
0,
[0; 64],
);
tip = store
.write(&crate::serialize::serialize(&Object::Commit(c)).unwrap())
.unwrap();
parents = vec![tip];
}
update(&layout, &store, "main", tip);
let branch_dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
for _ in 0..5 {
std::hint::black_box(read_current_descriptor(&branch_dir).unwrap());
std::hint::black_box(read_current(&branch_dir).unwrap());
}
let old_style_read = || {
let bytes = ancestry_state::read_bounded(&branch_dir.join("current"), 65)
.unwrap()
.unwrap();
let generation = refs::decode_ref_wire(&bytes).unwrap();
let raw = ancestry_state::read_bounded(
&snapshot_path(&branch_dir, generation),
MAX_SNAPSHOT_BYTES,
)
.unwrap()
.unwrap();
decode_descriptor_and_chain(&raw).unwrap()
};
let start = std::time::Instant::now();
for _ in 0..ITERS {
std::hint::black_box(old_style_read());
}
let old_elapsed = start.elapsed();
let start = std::time::Instant::now();
for _ in 0..ITERS {
std::hint::black_box(read_current_descriptor(&branch_dir).unwrap());
}
let new_elapsed = start.elapsed();
eprintln!(
"read_current (old, full chain + checksum): {:?}/iter over {ITERS} iters, {LEAVES} leaves",
old_elapsed / ITERS
);
eprintln!(
"read_current_descriptor (new, header only): {:?}/iter over {ITERS} iters, {LEAVES} leaves",
new_elapsed / ITERS
);
eprintln!(
"speedup: {:.1}x",
old_elapsed.as_secs_f64() / new_elapsed.as_secs_f64()
);
}
fn set_now(t: u64) {
NOW_OVERRIDE.with(|c| c.set(Some(t)));
}
fn clear_now() {
NOW_OVERRIDE.with(|c| c.set(None));
}
fn corrupt_object(store: &ObjectStore, h: &Hash) {
let path = store.path_for(h);
let mut bytes = fs::read(&path).unwrap();
bytes[6] ^= 0xFF;
fs::write(&path, bytes).unwrap();
}
fn build_chain(store: &ObjectStore, count: usize) -> Vec<Hash> {
let mut parents = vec![];
let mut tips = Vec::with_capacity(count);
for i in 0..count {
let h = commit(store, parents, i.to_be_bytes().as_slice());
parents = vec![h];
tips.push(h);
}
tips
}
#[test]
fn fast_forward_scrubs_a_bounded_window_not_the_whole_prefix() {
let (_dir, layout, store) = repo();
let mut tips = build_chain(&store, 600);
update(&layout, &store, "main", tips[599]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(
scrub,
ScrubState {
generation,
cursor: 0,
verified_through: 600,
last_full_verify_unix: scrub.last_full_verify_unix,
}
);
let next = commit(&store, vec![tips[599]], b"600");
tips.push(next);
update(&layout, &store, "main", next);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(
scrub.cursor, SCRUB_MIN_WINDOW,
"window must be bounded, not full-prefix"
);
assert_eq!(scrub.verified_through, 600);
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.chain, tips);
}
#[test]
fn scrub_lap_completion_forces_a_full_walk_but_keeps_the_generation() {
let (_dir, layout, store) = repo();
let mut tips = build_chain(&store, 600);
update(&layout, &store, "main", tips[599]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let first_generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
let step1 = commit(&store, vec![tips[599]], b"600");
tips.push(step1);
update(&layout, &store, "main", step1);
assert_eq!(read_scrub_state(&dir).unwrap().cursor, 512);
let before = now_unix();
let step2 = commit(&store, vec![step1], b"601");
tips.push(step2);
update(&layout, &store, "main", step2);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(
scrub.cursor, 0,
"a completed lap resets to a fresh full verify"
);
assert_eq!(scrub.verified_through, 602);
assert!(scrub.last_full_verify_unix >= before);
let snapshot = AncestrySnapshot::load(&layout, "main").unwrap();
assert_eq!(snapshot.chain, tips);
assert_eq!(
snapshot.descriptor().generation,
first_generation,
"still a fast-forward on the same branch history — generation must not change"
);
}
#[test]
fn full_walk_fallback_still_verifies_the_spliced_prefix() {
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 1500);
set_now(1_000_000);
update(&layout, &store, "main", tips[1499]);
corrupt_object(&store, &tips[900]);
set_now(1_000_000 + SCRUB_MAX_AGE_SECS + 1);
let step = commit(&store, vec![tips[1499]], b"1500");
let result =
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &step, &store);
assert!(
result.is_err(),
"the spliced full-walk fallback must still verify the whole prefix, not just the suffix"
);
clear_now();
}
#[test]
fn corruption_outside_the_current_window_is_caught_within_a_bounded_number_of_publishes() {
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 1500);
update(&layout, &store, "main", tips[1499]);
corrupt_object(&store, &tips[600]);
let step1 = commit(&store, vec![tips[1499]], b"1500");
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &step1, &store)
.expect("window [0, 512) does not include index 600: must still succeed");
let step2 = commit(&store, vec![step1], b"1501");
let result =
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &step2, &store);
assert!(
result.is_err(),
"window [512, 1024) includes index 600: corruption must now be caught"
);
}
#[test]
fn stale_scrub_state_forces_a_full_walk_regardless_of_cursor_position() {
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 1500);
set_now(1_000_000);
update(&layout, &store, "main", tips[1499]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
assert_eq!(read_scrub_state(&dir).unwrap().cursor, 0);
set_now(1_000_000 + SCRUB_MAX_AGE_SECS + 1);
let step = commit(&store, vec![tips[1499]], b"1500");
update(&layout, &store, "main", step);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(scrub.cursor, 0, "a forced full walk resets the cursor");
assert_eq!(scrub.verified_through, 1501);
assert_eq!(
scrub.last_full_verify_unix,
1_000_000 + SCRUB_MAX_AGE_SECS + 1
);
clear_now();
}
#[test]
fn scrub_age_boundary_is_exclusive() {
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 600);
set_now(1_000_000);
update(&layout, &store, "main", tips[599]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
set_now(1_000_000 + SCRUB_MAX_AGE_SECS - 1);
let step = commit(&store, vec![tips[599]], b"600");
update(&layout, &store, "main", step);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(scrub.cursor, 512, "age < bound takes the window path");
assert_eq!(scrub.last_full_verify_unix, 1_000_000);
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 600);
set_now(1_000_000);
update(&layout, &store, "main", tips[599]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
set_now(1_000_000 + SCRUB_MAX_AGE_SECS);
let step = commit(&store, vec![tips[599]], b"600");
update(&layout, &store, "main", step);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(scrub.cursor, 0, "age == bound forces a full walk");
assert_eq!(scrub.last_full_verify_unix, 1_000_000 + SCRUB_MAX_AGE_SECS);
clear_now();
}
#[test]
fn missing_scrub_state_forces_a_full_walk_instead_of_trusting_nothing() {
let (_dir, layout, store) = repo();
let tips = build_chain(&store, 1500);
update(&layout, &store, "main", tips[1499]);
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
fs::remove_file(dir.join("scrub")).unwrap();
corrupt_object(&store, &tips[1000]);
let step = commit(&store, vec![tips[1499]], b"1500");
let result =
refs::update_ref_with_ancestry(&layout, "main", RefWriteCondition::Any, &step, &store);
assert!(
result.is_err(),
"missing scrub state must force a full walk, not skip verification"
);
}
#[test]
fn scrub_state_from_a_superseded_generation_is_discarded_not_misapplied() {
let (_dir, layout, store) = repo();
let dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
let long_tips = build_chain(&store, 600);
update(&layout, &store, "main", long_tips[599]);
let stale_scrub_bytes = fs::read(scrub_path(&dir)).unwrap();
assert_eq!(
ScrubState::decode(&stale_scrub_bytes)
.unwrap()
.verified_through,
600
);
let short_tips = build_chain(&store, 5);
update(&layout, &store, "main", short_tips[4]);
let new_generation = AncestrySnapshot::load(&layout, "main")
.unwrap()
.descriptor()
.generation;
assert_ne!(
ScrubState::decode(&stale_scrub_bytes).unwrap().generation,
new_generation,
"sanity: the two generations must actually differ"
);
fs::write(scrub_path(&dir), &stale_scrub_bytes).unwrap();
let next = commit(&store, vec![short_tips[4]], b"5");
update(&layout, &store, "main", next);
let scrub = read_scrub_state(&dir).unwrap();
assert_eq!(
scrub.generation, new_generation,
"the discarded stale state must be replaced with one bound to the current generation"
);
assert_eq!(scrub.verified_through, 6);
}
#[test]
#[ignore = "manual profiling tool, not a correctness assertion"]
fn profile_scrub_window_vs_full_walk_every_publish() {
const LEAVES: usize = 20_000;
const PUBLISHES: u32 = 20;
fn fixture() -> (
tempfile::TempDir,
RepoLayout,
ObjectStore,
PathBuf,
Vec<Hash>,
) {
let (dir, layout, store) = repo();
let base = build_chain(&store, LEAVES);
update(&layout, &store, "main", base[LEAVES - 1]);
let branch_dir = ancestry_state::branch_dir(layout.common_dir(), "refs/heads/main");
(dir, layout, store, branch_dir, base)
}
{
let (_dir, layout, store, _branch_dir, base) = fixture();
let extra = commit(&store, vec![base[LEAVES - 1]], b"warm");
update(&layout, &store, "main", extra);
}
let (_dir, layout, store, _branch_dir, base) = fixture();
let mut tip = base[LEAVES - 1];
let start = std::time::Instant::now();
for i in 0..PUBLISHES {
let next = commit(&store, vec![tip], i.to_be_bytes().as_slice());
update(&layout, &store, "main", next);
tip = next;
}
let scrubbed_elapsed = start.elapsed();
let (_dir, layout, store, branch_dir, base) = fixture();
let mut tip = base[LEAVES - 1];
let start = std::time::Instant::now();
for i in 0..PUBLISHES {
fs::remove_file(branch_dir.join("scrub")).unwrap();
let next = commit(&store, vec![tip], i.to_be_bytes().as_slice());
update(&layout, &store, "main", next);
tip = next;
}
let full_walk_elapsed = start.elapsed();
eprintln!(
"scrub window (new): {:?}/publish over {PUBLISHES} publishes, {LEAVES} leaves",
scrubbed_elapsed / PUBLISHES
);
eprintln!(
"full walk (old, forced every publish): {:?}/publish over {PUBLISHES} publishes, {LEAVES} leaves",
full_walk_elapsed / PUBLISHES
);
eprintln!(
"speedup: {:.1}x",
full_walk_elapsed.as_secs_f64() / scrubbed_elapsed.as_secs_f64()
);
}
}