use std::collections::{BTreeMap, HashMap};
use std::fs::File;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, RwLock};
use anyhow::{anyhow, bail, Context, Result};
use znippy_common::ReservedSection;
use crate::arms::StoreConfig;
use crate::exploded::ExplodedStats;
use crate::exploded_arrow::ExplodedArchive;
use crate::gc::{Gc, GcReport};
use crate::graph::{assign_generations, CommitNode};
use crate::index_layout::{ObjectIndex, OneTableFourColumns};
use crate::indexer::{IndexJob, ObjectAbsorb, PushPath};
use crate::object::{GitHashKind, GitObjectKind};
use crate::pack_walk::PackWalk;
use crate::reach::ReachEntry;
use crate::read_stack::{ObjectReadStack, RebuildTriggers};
use crate::refs::{RefLog, RefUpdate};
use crate::resolve::BaseSource;
pub use git_storage_trait::{Extent, GitOps, Oid, RefRow, Stored, TxId};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LookupPath {
Serial,
Batch,
}
pub const BATCH_SATURATES_AT: usize = 100;
pub fn lookup_path(n: usize) -> LookupPath {
if n == 1 {
LookupPath::Serial
} else {
LookupPath::Batch
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct PendingPack {
pack_id: u64,
offset: u64,
len: u64,
}
impl PendingPack {
fn from_job(j: IndexJob) -> Self {
Self {
pack_id: j.pack_id,
offset: j.offset,
len: j.len,
}
}
}
#[derive(Default)]
struct PackState {
words: Vec<u64>,
extent: Vec<Option<Extent>>,
unabsorbed: usize,
}
impl PackState {
fn is_absorbed(&self, pack_id: u64) -> bool {
let i = pack_id as usize;
self.words
.get(i / 64)
.is_some_and(|w| w >> (i % 64) & 1 == 1)
}
fn note(&mut self, pack_id: u64, extent: Extent) {
let i = pack_id as usize;
self.grow_to(i);
if self.extent[i].is_none() {
self.extent[i] = Some(extent);
self.unabsorbed += 1;
}
}
fn grow_to(&mut self, i: usize) {
if self.words.len() <= i / 64 {
self.words.resize(i / 64 + 1, 0);
}
if self.extent.len() <= i {
self.extent.resize(i + 1, None);
}
}
fn mark_absorbed(&mut self, pack_id: u64, extent: Extent) {
let i = pack_id as usize;
self.grow_to(i);
if self.words[i / 64] >> (i % 64) & 1 == 0 {
self.words[i / 64] |= 1 << (i % 64);
if self.extent[i].is_some() {
self.unabsorbed -= 1;
}
}
self.extent[i] = Some(extent);
}
fn pending(&self) -> Vec<PendingPack> {
self.extent
.iter()
.enumerate()
.filter(|(i, e)| e.is_some() && !self.is_absorbed(*i as u64))
.map(|(i, e)| {
let (offset, len) = e.expect("filtered on Some");
PendingPack {
pack_id: i as u64,
offset,
len,
}
})
.collect()
}
}
#[derive(Default)]
struct Derived {
graph: Vec<CommitNode>,
trees: HashMap<String, Vec<u8>>,
ordinal: HashMap<String, u32>,
oids: Vec<String>,
oids_raw: Vec<u8>,
reach: Arc<Vec<ReachEntry>>,
commit_raw: Option<Arc<std::collections::HashSet<Vec<u8>>>>,
packs: Vec<(u64, u64)>,
}
const FAN_OUT_AT: usize = 4096;
const LIVE_REACH_COMMITS: usize = 512;
fn live_reach_policy() -> crate::reach::ReachPolicy {
static MAX: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
let max_commits = *MAX.get_or_init(|| {
crate::arms::read_env(crate::arms::ENV_REACH_COMMITS)
.and_then(|v| v.parse::<usize>().ok())
.unwrap_or(LIVE_REACH_COMMITS)
});
crate::reach::ReachPolicy { max_commits }
}
const MAX_RESOLVE_WORKERS: usize = 4;
fn emit_workers() -> usize {
static WORKERS: std::sync::OnceLock<usize> = std::sync::OnceLock::new();
*WORKERS.get_or_init(|| {
crate::arms::read_env(crate::arms::ENV_EMIT_WORKERS)
.and_then(|v| v.parse::<usize>().ok())
.unwrap_or(MAX_RESOLVE_WORKERS)
})
}
pub(crate) struct Absorber<S: ObjectIndex = OneTableFourColumns> {
blobs: PathBuf,
reader: File,
map: crate::archive_map::ArchiveMap,
hash: GitHashKind,
objects: ObjectReadStack<S>,
exploded: ExplodedArchive,
derived: RwLock<Derived>,
packs: Mutex<PackState>,
gate: Mutex<()>,
}
pub struct GitStore<S: ObjectIndex = OneTableFourColumns> {
root: PathBuf,
blobs: PathBuf,
archive: PathBuf,
push: PushPath,
account: String,
hash: GitHashKind,
absorber: Arc<Absorber<S>>,
refs: RefLog,
ref_gate: Mutex<()>,
gc: Box<dyn Gc + Send + Sync>,
arms: StoreConfig,
}
impl GitStore<OneTableFourColumns> {
pub fn open(root: &Path, account: &str) -> Result<Self> {
Self::open_with(root, account, GitHashKind::Sha1)
}
pub fn open_with(root: &Path, account: &str, hash: GitHashKind) -> Result<Self> {
Self::open_with_arms(
root,
account,
hash,
StoreConfig::DEFAULT.with_redb_cache_bytes(crate::arms::redb_cache_bytes()?),
)
}
}
enum BaseOf {
Whole,
Inside { at: u64 },
Outside {
offset: u64,
oid: Option<Vec<u8>>,
},
}
struct ExternalBases {
by_offset: HashMap<u64, OidKey>,
held: std::collections::HashSet<OidKey>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
struct OidKey {
bytes: [u8; 32],
len: u8,
}
impl OidKey {
fn new(oid: &[u8]) -> Result<Self> {
let len = oid.len();
if len == 0 || len > 32 {
bail!("{len} is not the width of any git object id (1..=32 bytes)");
}
let mut bytes = [0u8; 32];
bytes[..len].copy_from_slice(oid);
Ok(OidKey {
bytes,
len: len as u8,
})
}
fn as_slice(&self) -> &[u8] {
&self.bytes[..self.len as usize]
}
}
impl<S: ObjectIndex + 'static> GitStore<S> {
pub fn open_with_arms(
root: &Path,
account: &str,
hash: GitHashKind,
arms: StoreConfig,
) -> Result<Self> {
std::fs::create_dir_all(root)
.with_context(|| format!("creating the store root {}", root.display()))?;
let blobs = root.join("objects.pack");
let writer = arms.writer.create(&blobs)?;
let journal = arms.writer.journal(&blobs);
let cache_bytes = arms.redb_cache_bytes;
let objects = ObjectReadStack::<S>::open(
&root.join("objects.tail"),
RebuildTriggers::default(),
cache_bytes,
)?;
let _ = cache_bytes; let exploded = ExplodedArchive::open(&root.join("objects.exploded"))?;
let absorber = Arc::new(Absorber::<S>::open(&blobs, hash, objects, exploded)?);
let acked = match journal.as_deref() {
Some(p) if p.exists() => crate::archive_write::read_journal(p).with_context(|| {
format!(
"reading {} — it is the durable half of the indexed bit and a store \
cannot be opened without knowing which packs it owes index work for",
p.display()
)
})?,
_ => Vec::new(),
};
let requeue = absorber.adopt_journal(&acked)?;
let push = PushPath::with_absorber(
writer,
&blobs,
journal,
absorber.clone() as Arc<dyn ObjectAbsorb>,
)?;
let refs = RefLog::new(root.join("refs.log"));
let archive = root.join("repository.znippy");
let store = Self {
root: root.to_path_buf(),
blobs,
archive,
push,
account: account.to_string(),
hash,
absorber,
refs,
ref_gate: Mutex::new(()),
gc: arms.gc.create(),
arms,
};
let indexer = store.push.indexer(&store.account);
for p in &requeue {
indexer.submit(IndexJob {
pack_id: p.pack_id,
offset: p.offset,
len: p.len,
})?;
}
store.refold()?;
Ok(store)
}
pub fn arms(&self) -> StoreConfig {
self.arms
}
pub fn writer_name(&self) -> &'static str {
self.push.name()
}
pub fn writer_durability(&self) -> &'static str {
self.push.durability()
}
pub(crate) fn gc_arm(&self) -> &dyn Gc {
self.gc.as_ref()
}
pub fn blobs_path(&self) -> &Path {
&self.blobs
}
pub fn archive_path(&self) -> &Path {
&self.archive
}
pub fn root(&self) -> &Path {
&self.root
}
pub fn hash_kind(&self) -> GitHashKind {
self.hash
}
pub fn index(&self) -> &ObjectReadStack<S> {
&self.absorber.objects
}
pub fn object_count(&self) -> usize {
self.absorber.objects.len()
}
pub fn commit_count(&self) -> usize {
self.absorber
.derived
.read()
.expect("derived lock")
.graph
.len()
}
pub fn content(&self, oid: Oid<'_>) -> Result<Option<(GitObjectKind, Vec<u8>)>> {
if self.unindexed_packs() > 0 {
self.absorb_pending()?;
}
self.absorber.resolved(oid)
}
pub fn exploded_stats(&self) -> ExplodedStats {
self.absorber.exploded.engine_stats()
}
pub fn exploded_path(&self) -> &Path {
self.absorber.exploded.path()
}
pub fn exploded_of_kind(&self, kind: GitObjectKind) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.absorber.exploded.of_kind(kind)
}
pub fn graph_snapshot(&self) -> Vec<CommitNode> {
self.absorber
.derived
.read()
.expect("derived lock")
.graph
.clone()
}
pub fn unindexed_packs(&self) -> usize {
self.absorber.unindexed_packs()
}
pub fn indexer(&self) -> std::sync::Arc<crate::indexer::AccountIndexer> {
self.push.indexer(&self.account)
}
pub fn wait_indexed(&self) {
self.push.pool().wait_caught_up();
}
pub fn absorb_pending(&self) -> Result<usize> {
self.absorber.absorb_pending()
}
pub fn oids_in_extent(&self, extent: Extent) -> Result<Vec<Vec<u8>>> {
self.absorb_pending()
.context("absorbing the pending index work before enumerating a pack's objects")?;
let (start, len) = extent;
let end = start.saturating_add(len);
let all = self.index().oids_in_order()?;
let refs: Vec<&[u8]> = all.iter().map(Vec::as_slice).collect();
let extents = self.index().extents_batch(&refs);
Ok(all
.into_iter()
.zip(extents)
.filter_map(|(oid, at)| match at {
Some((offset, _)) if offset >= start && offset < end => Some(oid),
_ => None,
})
.collect())
}
pub fn emit_set(
&self,
oids: &[Oid<'_>],
ofs_delta_ok: bool,
thin_haves: Option<&[Oid<'_>]>,
) -> Result<Vec<crate::pack_walk::EmitEntry>> {
use crate::index_layout::{ObjType, ObjectIndex as _};
use std::collections::hash_map::Entry;
self.absorb_pending()
.context("absorbing pending index work before emitting a pack")?;
let mut asked: Vec<&[u8]> = Vec::with_capacity(oids.len());
let mut seen: HashMap<&[u8], usize> = HashMap::with_capacity(oids.len());
for oid in oids {
if let Entry::Vacant(slot) = seen.entry(*oid) {
slot.insert(asked.len());
asked.push(*oid);
}
}
let rows = self.index().lookup_batch(&asked);
let snap = self.archive_snapshot()?;
let in_set_offsets: HashMap<u64, usize> = rows
.iter()
.enumerate()
.filter_map(|(i, r)| r.as_ref().map(|r| (r.offset, i)))
.collect();
let oid_len = self.hash_kind().oid_len();
let mut base_content: Option<(u64, Vec<u8>)> = None;
let mut client_bases: Option<ExternalBases> = None;
let mut out = Vec::with_capacity(asked.len());
for (oid, row) in asked.iter().zip(rows.iter()) {
let Some(row) = row.as_ref() else {
bail!(
"this pack was asked for {}, which this repository does not hold; emitting \
the rest would be a short pack reported as a success, so this refuses \
instead",
hex::encode(oid)
);
};
let base = match row.obj_type {
ObjType::OfsDelta => match in_set_offsets.get(&row.delta_base) {
Some(_) => BaseOf::Inside {
at: row.delta_base,
},
None => BaseOf::Outside {
offset: row.delta_base,
oid: None,
},
},
ObjType::RefDelta => {
let head = self.extent(&snap, row.offset, row.len.min(64))?;
let (_, _, n) = crate::pack_walk::type_and_size_of(&head)?;
match head.get(n..n + oid_len) {
Some(b) => match seen
.get(b)
.and_then(|i| rows[*i].as_ref())
.map(|r| r.offset)
{
Some(at) => BaseOf::Inside { at },
None => BaseOf::Outside {
offset: 0,
oid: Some(b.to_vec()),
},
},
None => BaseOf::Outside {
offset: 0,
oid: None,
},
}
}
_ => BaseOf::Whole,
};
if let BaseOf::Whole | BaseOf::Inside { .. } = base {
let inside_at = match base {
BaseOf::Inside { at } => Some(at),
_ => None,
};
let extent = crate::pack_walk::EntryBytes::Extent {
offset: row.offset,
len: row.len,
};
let (stored, obj_type, delta_base) = match (row.obj_type, inside_at) {
(ObjType::OfsDelta, Some(_)) if !ofs_delta_ok => {
let base_oid = self.oid_at_offset(row.delta_base, &asked, &rows)?;
let bytes = self.extent(&snap, row.offset, row.len)?;
(
crate::pack_walk::EntryBytes::Owned(Self::as_ref_delta(
&bytes, &base_oid,
)?),
ObjType::RefDelta,
0,
)
}
(ObjType::RefDelta, Some(at)) if ofs_delta_ok => {
let bytes = self.extent(&snap, row.offset, row.len)?;
(
crate::pack_walk::EntryBytes::Owned(Self::as_ofs_delta(
&bytes, oid_len, row.offset, at,
)?),
ObjType::OfsDelta,
at,
)
}
_ => (extent, row.obj_type, row.delta_base),
};
out.push(crate::pack_walk::EmitEntry {
oid: oid.to_vec(),
stored,
obj_type,
uncompressed_size: row.uncompressed_size,
delta_base,
offset: row.offset,
recompressed: false,
deltified: false,
});
continue;
}
if let Some(haves) = thin_haves {
if client_bases.is_none() {
client_bases = Some(self.client_bases(haves)?);
}
let ext = client_bases
.as_ref()
.expect("the client's bases were just built");
let base_oid: Option<Vec<u8>> = match &base {
BaseOf::Whole | BaseOf::Inside { .. } => None,
BaseOf::Outside {
oid: Some(named), ..
} => ext
.held
.contains(&OidKey::new(named)?)
.then(|| named.clone()),
BaseOf::Outside {
offset,
oid: None,
} => ext.by_offset.get(offset).map(|k| k.as_slice().to_vec()),
};
if let Some(base_oid) = base_oid {
let stored = if row.obj_type == ObjType::OfsDelta {
let bytes = self.extent(&snap, row.offset, row.len)?;
crate::pack_walk::EntryBytes::Owned(Self::as_ref_delta(&bytes, &base_oid)?)
} else {
crate::pack_walk::EntryBytes::Extent {
offset: row.offset,
len: row.len,
}
};
out.push(crate::pack_walk::EmitEntry {
oid: oid.to_vec(),
stored,
obj_type: ObjType::RefDelta,
uncompressed_size: row.uncompressed_size,
delta_base: 0,
offset: row.offset,
recompressed: false,
deltified: false,
});
continue;
}
}
let (kind, body) = self.content(oid)?.ok_or_else(|| {
anyhow!(
"{} is a delta whose base is outside this request, so it must be sent whole, \
and this repository cannot produce its content — neither §14's exploded \
table nor a re-derivation from the verbatim packs answered",
hex::encode(oid)
)
})?;
let candidate = if crate::delta::enabled() {
self.in_request_chain_ancestor(row.delta_base, row.offset, &in_set_offsets)?
} else {
None
};
if let Some(base_at) = candidate {
let at = *in_set_offsets
.get(&base_at)
.expect("the ancestor was found by lookup in this very map");
if base_content.as_ref().is_none_or(|(o, _)| *o != base_at) {
let base_oid = asked[at];
let (_, base_body) = self.content(base_oid)?.ok_or_else(|| {
anyhow!(
"{} is in this request and is the delta base chosen for {}, and this \
repository cannot produce its content",
hex::encode(base_oid),
hex::encode(oid)
)
})?;
base_content = Some((base_at, base_body));
}
let (_, base_body) = base_content.as_ref().expect("just filled");
if let Some(d) = crate::delta::delta(base_body, &body)? {
let stored = Self::ofs_delta_entry(&d, row.offset, base_at)?;
let (stored, obj_type, delta_base) = if ofs_delta_ok {
(
crate::pack_walk::EntryBytes::Owned(stored),
ObjType::OfsDelta,
base_at,
)
} else {
(
crate::pack_walk::EntryBytes::Owned(Self::as_ref_delta(
&stored, asked[at],
)?),
ObjType::RefDelta,
0,
)
};
out.push(crate::pack_walk::EmitEntry {
oid: oid.to_vec(),
stored,
obj_type,
uncompressed_size: row.uncompressed_size,
delta_base,
offset: row.offset,
recompressed: true,
deltified: true,
});
continue;
}
}
out.push(crate::pack_walk::EmitEntry {
oid: oid.to_vec(),
stored: crate::pack_walk::EntryBytes::Owned(Self::whole_entry(kind, &body)?),
obj_type: crate::serve::resolved_type(kind),
uncompressed_size: body.len() as u64,
delta_base: 0,
offset: row.offset,
recompressed: true,
deltified: false,
});
}
Ok(out)
}
pub(crate) fn resolve_emit_payloads<'a>(
&self,
entries: &'a [crate::pack_walk::EmitEntry],
snap: &'a crate::archive_map::Mapped,
) -> Result<Vec<std::borrow::Cow<'a, [u8]>>>
where
Self: Sync,
{
let workers = if entries.len() < FAN_OUT_AT {
1
} else {
emit_workers()
};
self.resolve_with(entries, snap, workers)
}
pub(crate) fn resolve_with<'a>(
&self,
entries: &'a [crate::pack_walk::EmitEntry],
snap: &'a crate::archive_map::Mapped,
workers: usize,
) -> Result<Vec<std::borrow::Cow<'a, [u8]>>>
where
Self: Sync,
{
use znippy_zoomies::gatling_forkjoin::gatling_for_each;
let one = |i: usize| -> Result<std::borrow::Cow<'a, [u8]>> {
let e = &entries[i];
let bytes = match &e.stored {
crate::pack_walk::EntryBytes::Owned(v) => std::borrow::Cow::Borrowed(v.as_slice()),
crate::pack_walk::EntryBytes::Extent { offset, len } => {
self.extent(snap, *offset, *len)?
}
};
crate::pack_walk::type_and_size_of(&bytes).with_context(|| {
format!(
"resolving the stored bytes of {} for emission",
hex::encode(&e.oid)
)
})?;
Ok(bytes)
};
if workers <= 1 {
return (0..entries.len()).map(one).collect();
}
gatling_for_each(entries.len(), workers, one)
.into_iter()
.collect()
}
pub(crate) fn emit_ordered(
&self,
ordered: &[crate::pack_walk::EmitEntry],
out: &mut dyn std::io::Write,
) -> Result<crate::pack_walk::EmitReport>
where
Self: Sync,
{
let snap = self
.archive_snapshot()
.context("mapping the archive to emit")?;
let payloads = self
.resolve_emit_payloads(ordered, &snap)
.context("resolving the emit set's stored bytes")?;
crate::pack_walk::emit_pack(
ordered,
self.hash_kind(),
out,
&|i| Ok(payloads[i].as_ref()),
)
}
fn in_request_chain_ancestor(
&self,
from: u64,
entry_at: u64,
in_request: &HashMap<u64, usize>,
) -> Result<Option<u64>> {
use crate::index_layout::{ObjType, ObjectIndex as _};
const MAX_LINKS: usize = 100;
const HEADER_PROBE: u64 = 64;
let oid_len = self.hash_kind().oid_len();
let mut at = from;
for _ in 0..MAX_LINKS {
if at != 0 && at < entry_at && in_request.contains_key(&at) {
return Ok(Some(at));
}
if at == 0 {
return Ok(None);
}
let head = self.read_extent(at, HEADER_PROBE)?;
let (t, _, n) = crate::pack_walk::type_and_size_of(&head)?;
at = match t {
ObjType::Commit | ObjType::Tree | ObjType::Blob | ObjType::Tag => {
return Ok(None)
}
ObjType::OfsDelta => {
let (distance, _) = crate::pack_walk::ofs_distance_of(&head[n..])?;
match at.checked_sub(distance) {
Some(next) if next != 0 => next,
_ => return Ok(None),
}
}
ObjType::RefDelta => {
let Some(base_oid) = head.get(n..n + oid_len) else {
return Ok(None);
};
match self.index().lookup(base_oid) {
Some(row) => row.offset,
None => return Ok(None),
}
}
};
}
Ok(None)
}
fn ofs_delta_entry(d: &[u8], entry_at: u64, base_at: u64) -> Result<Vec<u8>> {
use crate::index_layout::ObjType;
use std::io::Write as _;
let distance = entry_at.checked_sub(base_at).ok_or_else(|| {
anyhow!(
"a computed delta at archive offset {entry_at} names a base at {base_at}, which \
is after it — an ofs-delta distance is backwards and this would not encode"
)
})?;
let mut out = Vec::with_capacity(d.len() / 2 + 32);
crate::pack_walk::encode_type_and_size(&mut out, ObjType::OfsDelta, d.len() as u64);
crate::pack_walk::encode_ofs_distance(&mut out, distance);
let mut z = flate2::write::ZlibEncoder::new(out, flate2::Compression::fast());
z.write_all(d)?;
Ok(z.finish()?)
}
fn client_bases(&self, haves: &[Oid<'_>]) -> Result<ExternalBases> {
use crate::index_layout::ObjectIndex as _;
let held = self.reachable_raw(haves, &[]).context(
"closing over the negotiated common tips, to learn which bases the client can supply \
for itself",
)?;
let refs: Vec<&[u8]> = held.iter().collect();
let rows = self.index().lookup_batch(&refs);
let mut by_offset = HashMap::with_capacity(refs.len());
let mut keys = std::collections::HashSet::with_capacity(refs.len());
for (oid, row) in refs.iter().zip(rows.iter()) {
let key = OidKey::new(oid)?;
if let Some(row) = row {
by_offset.insert(row.offset, key);
}
keys.insert(key);
}
Ok(ExternalBases {
by_offset,
held: keys,
})
}
fn oid_at_offset(
&self,
offset: u64,
asked: &[&[u8]],
rows: &[Option<crate::index_layout::IndexRow>],
) -> Result<Vec<u8>> {
asked
.iter()
.zip(rows)
.find(|(_, r)| r.as_ref().is_some_and(|r| r.offset == offset))
.map(|(oid, _)| oid.to_vec())
.ok_or_else(|| {
anyhow!(
"the entry at archive offset {offset} was checked to be in this request and \
then could not be found in it"
)
})
}
fn as_ref_delta(stored: &[u8], base_oid: &[u8]) -> Result<Vec<u8>> {
use crate::index_layout::ObjType;
let (t, stated_size, n) = crate::pack_walk::type_and_size_of(stored)?;
if t != ObjType::OfsDelta {
bail!("only an ofs-delta can be re-headed as a ref-delta, this entry is {t:?}");
}
let (_, d) = crate::pack_walk::ofs_distance_of(stored.get(n..).unwrap_or(&[]))?;
let mut out = Vec::with_capacity(stored.len() + base_oid.len());
crate::pack_walk::encode_type_and_size(&mut out, ObjType::RefDelta, stated_size);
out.extend_from_slice(base_oid);
out.extend_from_slice(stored.get(n + d..).unwrap_or(&[]));
Ok(out)
}
fn as_ofs_delta(
stored: &[u8],
oid_len: usize,
entry_at: u64,
base_at: u64,
) -> Result<Vec<u8>> {
use crate::index_layout::ObjType;
let (t, stated_size, n) = crate::pack_walk::type_and_size_of(stored)?;
if t != ObjType::RefDelta {
bail!("only a ref-delta can be re-headed as an ofs-delta, this entry is {t:?}");
}
let payload = stored.get(n + oid_len..).ok_or_else(|| {
anyhow!(
"a ref-delta entry at archive offset {entry_at} is {} bytes, which is not even its \
header plus a {oid_len}-byte base oid",
stored.len()
)
})?;
let mut out = Vec::with_capacity(payload.len() + 32);
crate::pack_walk::encode_type_and_size(&mut out, ObjType::OfsDelta, stated_size);
crate::pack_walk::encode_ofs_distance(
&mut out,
entry_at.checked_sub(base_at).filter(|d| *d > 0).unwrap_or(1),
);
out.extend_from_slice(payload);
Ok(out)
}
fn whole_entry(kind: crate::object::GitObjectKind, body: &[u8]) -> Result<Vec<u8>> {
use std::io::Write as _;
let mut out = Vec::with_capacity(body.len() / 2 + 32);
crate::pack_walk::encode_type_and_size(
&mut out,
crate::serve::resolved_type(kind),
body.len() as u64,
);
let mut z = flate2::write::ZlibEncoder::new(out, flate2::Compression::fast());
z.write_all(body)?;
Ok(z.finish()?)
}
#[cfg(test)]
pub(crate) fn hold_absorb_gate(&self) -> std::sync::MutexGuard<'_, ()> {
self.absorber.gate.lock().expect("absorb gate poisoned")
}
pub(crate) fn read_extent(&self, offset: u64, len: u64) -> Result<Vec<u8>> {
self.absorber.read_extent(offset, len)
}
pub(crate) fn archive_snapshot(&self) -> Result<Arc<crate::archive_map::Mapped>> {
self.absorber.map.snapshot()
}
pub(crate) fn extent<'m>(
&self,
snap: &'m crate::archive_map::Mapped,
offset: u64,
len: u64,
) -> Result<std::borrow::Cow<'m, [u8]>> {
match snap.get(offset, len) {
Some(b) => Ok(std::borrow::Cow::Borrowed(b)),
None => Ok(std::borrow::Cow::Owned(self.read_extent(offset, len)?)),
}
}
pub(crate) fn queue(&self, pack_id: u64, extent: Extent) -> Result<()> {
self.absorber.queue(pack_id, extent)
}
fn reach_bitmaps(&self) -> Result<Arc<Vec<ReachEntry>>> {
self.absorber.reach_bitmaps()
}
pub(crate) fn reach_bitmaps_with(
&self,
policy: crate::reach::ReachPolicy,
cache: bool,
) -> Result<Arc<Vec<ReachEntry>>> {
self.absorber.reach_bitmaps_with(policy, cache)
}
pub(crate) fn commit_oids_raw(
&self,
) -> Result<Option<Arc<std::collections::HashSet<Vec<u8>>>>> {
self.absorber.commit_oids_raw()
}
fn refold(&self) -> Result<()> {
self.absorber.refold()
}
}
impl<S: ObjectIndex> Absorber<S> {
fn open(
blobs: &Path,
hash: GitHashKind,
objects: ObjectReadStack<S>,
exploded: ExplodedArchive,
) -> Result<Self> {
let reader =
File::open(blobs).with_context(|| format!("opening {} for reads", blobs.display()))?;
Ok(Self {
blobs: blobs.to_path_buf(),
reader,
map: crate::archive_map::ArchiveMap::new(blobs),
hash,
objects,
exploded,
derived: RwLock::new(Derived::default()),
packs: Mutex::new(PackState::default()),
gate: Mutex::new(()),
})
}
fn unindexed_packs(&self) -> usize {
self.packs.lock().expect("pack state lock").unabsorbed
}
fn queue(&self, pack_id: u64, extent: Extent) -> Result<()> {
let mut s = self
.packs
.lock()
.map_err(|_| anyhow!("pack state poisoned"))?;
s.note(pack_id, extent);
Ok(())
}
fn adopt_journal(&self, journal: &[Extent]) -> Result<Vec<PendingPack>> {
if journal.is_empty() {
return Ok(Vec::new());
}
let packs = crate::archive_write::acked_packs(journal);
let retired = crate::archive_write::retired_offsets(journal);
if packs.is_empty() {
return Ok(Vec::new());
}
let exploded_is_whole = match self.exploded.policy() {
crate::exploded_arrow::ExplodePolicy::Graph => {
self.objects.len() == 0 || self.exploded.rows()? > 0
}
_ => self.exploded.rows()? >= self.objects.len() as u64,
};
let has_rows = if exploded_is_whole {
self.objects.extents_with_rows(&packs)?
} else {
vec![false; packs.len()]
};
{
let mut d = self
.derived
.write()
.map_err(|_| anyhow!("derived poisoned"))?;
d.packs = packs.clone();
}
let mut s = self
.packs
.lock()
.map_err(|_| anyhow!("pack state poisoned"))?;
let mut requeue = Vec::new();
for (i, (&extent, &absorbed)) in packs.iter().zip(&has_rows).enumerate() {
let pack_id = i as u64;
if absorbed {
s.mark_absorbed(pack_id, extent);
} else if retired.contains(&extent.0) {
s.mark_absorbed(pack_id, extent);
} else {
s.note(pack_id, extent);
requeue.push(PendingPack {
pack_id,
offset: extent.0,
len: extent.1,
});
}
}
Ok(requeue)
}
fn absorb_pending(&self) -> Result<usize> {
if self.unindexed_packs() == 0 {
return Ok(0);
}
let _gate = self
.gate
.lock()
.map_err(|_| anyhow!("absorb gate poisoned"))?;
let jobs: Vec<PendingPack> = {
let s = self
.packs
.lock()
.map_err(|_| anyhow!("pack state poisoned"))?;
s.pending()
};
let mut absorbed = 0usize;
for job in &jobs {
self.absorb_gated(job)?;
absorbed += 1;
}
Ok(absorbed)
}
fn absorb_gated(&self, job: &PendingPack) -> Result<()> {
{
let s = self
.packs
.lock()
.map_err(|_| anyhow!("pack state poisoned"))?;
if s.is_absorbed(job.pack_id) {
return Ok(());
}
}
self.absorb_one(job).map_err(|e| {
e.context(format!(
"absorbing pack {} at ({}, {}) — its bytes are durable and it stays \
un-indexed; reads will keep falling back rather than answer absent",
job.pack_id, job.offset, job.len
))
})?;
let mut s = self
.packs
.lock()
.map_err(|_| anyhow!("pack state poisoned"))?;
s.mark_absorbed(job.pack_id, (job.offset, job.len));
Ok(())
}
fn absorb_one(&self, job: &PendingPack) -> Result<()> {
let bytes = self.read_extent(job.offset, job.len)?;
let walked = crate::pack_walk::walk(&bytes, self.hash.oid_len())?;
let rows = crate::resolve::resolve_walked(
&bytes,
&walked,
self.hash,
job.offset,
self,
&self.exploded,
)?;
self.exploded.flush()?;
let entries: Vec<crate::index_layout::IndexEntry> =
rows.iter().map(|r| r.index_entry()).collect();
self.objects.append(&entries)?;
{
let mut d = self
.derived
.write()
.map_err(|_| anyhow!("derived poisoned"))?;
if !d.packs.contains(&(job.offset, job.len)) {
d.packs.push((job.offset, job.len));
}
}
self.refold()?;
let _ = self.objects.maybe_rebuild()?;
Ok(())
}
fn refold(&self) -> Result<()> {
let commits = self.exploded.of_kind(GitObjectKind::Commit)?;
let trees = self.exploded.of_kind(GitObjectKind::Tree)?;
let mut d = self
.derived
.write()
.map_err(|_| anyhow!("derived poisoned"))?;
d.trees = trees
.into_iter()
.map(|(oid, payload)| (hex::encode(oid), payload))
.collect();
let graph: Vec<CommitNode> = commits
.iter()
.map(|(oid, payload)| {
let h = crate::object::parse_commit(payload);
CommitNode {
oid: hex::encode(oid),
parents: h.parents,
tree: h.tree,
committer_time: h.committer_time,
generation: 0,
}
})
.collect();
d.graph = assign_generations(graph);
d.commit_raw = d
.graph
.iter()
.map(|n| hex::decode(&n.oid).ok())
.collect::<Option<std::collections::HashSet<Vec<u8>>>>()
.map(Arc::new);
let mut raw = self.tail_oids()?;
raw.sort_unstable();
let mut oids_raw = Vec::with_capacity(raw.iter().map(Vec::len).sum());
for oid in &raw {
oids_raw.extend_from_slice(oid);
}
let oids: Vec<String> = raw.iter().map(hex::encode).collect();
d.ordinal = oids
.iter()
.enumerate()
.map(|(i, o)| (o.clone(), i as u32))
.collect();
d.oids = oids;
d.oids_raw = oids_raw;
d.reach = Arc::default();
Ok(())
}
fn tail_oids(&self) -> Result<Vec<Vec<u8>>> {
self.objects.oids_in_order()
}
fn reach_bitmaps(&self) -> Result<Arc<Vec<ReachEntry>>> {
self.reach_bitmaps_with(live_reach_policy(), true)
}
fn reach_bitmaps_with(
&self,
policy: crate::reach::ReachPolicy,
cache: bool,
) -> Result<Arc<Vec<ReachEntry>>> {
{
let d = self
.derived
.read()
.map_err(|_| anyhow!("derived poisoned"))?;
if cache && !d.reach.is_empty() {
return Ok(Arc::clone(&d.reach));
}
if d.graph.is_empty() {
return Ok(Arc::default());
}
}
let mut d = self
.derived
.write()
.map_err(|_| anyhow!("derived poisoned"))?;
let facts = crate::reach::ObjectFacts {
ordinal: &d.ordinal,
trees: &d.trees,
oid_len: self.hash.oid_len(),
};
let built = Arc::new(crate::reach::build_reach(&d.graph, &facts, policy));
if cache {
d.reach = Arc::clone(&built);
}
Ok(built)
}
fn commit_oids_raw(&self) -> Result<Option<Arc<std::collections::HashSet<Vec<u8>>>>> {
Ok(self
.derived
.read()
.map_err(|_| anyhow!("derived poisoned"))?
.commit_raw
.clone())
}
fn read_extent(&self, offset: u64, len: u64) -> Result<Vec<u8>> {
use std::os::fd::AsRawFd;
let want = usize::try_from(len).context("extent length overflows usize")?;
let mut buf: Vec<u8> = Vec::with_capacity(want);
let fd = self.reader.as_raw_fd();
let mut filled = 0usize;
while filled < want {
let n = unsafe {
libc::pread(
fd,
buf.as_mut_ptr().add(filled).cast(),
want - filled,
i64::try_from(offset + filled as u64).context("extent offset overflows off_t")?,
)
};
if n < 0 {
let err = std::io::Error::last_os_error();
if err.kind() == std::io::ErrorKind::Interrupted {
continue;
}
return Err(err).with_context(|| {
format!("reading ({offset}, {len}) out of {}", self.blobs.display())
});
}
if n == 0 {
bail!(
"short read at ({offset}, {len}) out of {}: {filled} of {want} bytes before \
EOF — the archive is truncated relative to its index",
self.blobs.display()
);
}
filled += n as usize;
}
unsafe { buf.set_len(want) };
Ok(buf)
}
}
impl<S: ObjectIndex> ObjectAbsorb for Absorber<S> {
fn absorb(&self, job: IndexJob) -> Result<()> {
let _gate = self
.gate
.lock()
.map_err(|_| anyhow!("absorb gate poisoned"))?;
self.absorb_gated(&PendingPack::from_job(job))
}
}
impl<S: ObjectIndex + 'static> GitStore<S> {
pub(crate) fn push_path(&self) -> &PushPath {
&self.push
}
pub(crate) fn account(&self) -> &str {
&self.account
}
pub(crate) fn ref_log(&self) -> &RefLog {
&self.refs
}
pub(crate) fn ref_gate(&self) -> &Mutex<()> {
&self.ref_gate
}
pub(crate) fn external_bases_exist(&self, bytes: &[u8], w: &PackWalk) -> Result<()> {
let c = w.closure();
if !c.broken_offsets.is_empty() {
bail!(
"this pack is corrupt: {} delta base offset(s) do not land on an entry — first {}",
c.broken_offsets.len(),
c.broken_offsets[0]
);
}
let mut unheld: Vec<&[u8]> = Vec::new();
for oid in &c.external_refs {
if self.lookup_one(oid)?.is_none() {
unheld.push(oid.as_slice());
}
}
if unheld.is_empty() {
return Ok(());
}
let resolved = crate::resolve::resolve_walked(
bytes,
w,
self.hash_kind(),
0,
&*self.absorber,
&crate::exploded::NoSink,
)
.with_context(|| {
format!(
"this pack deltas against {} object(s) this repository does not have, so it was \
resolved to find out whether the pack carries them itself — first {}",
unheld.len(),
hex::encode(unheld[0])
)
})?;
let in_pack: std::collections::HashSet<&[u8]> =
resolved.iter().map(|r| r.oid.as_slice()).collect();
for oid in unheld {
if !in_pack.contains(oid) {
bail!(
"this pack deltas against {}, which this repository does not have and which \
the pack does not carry either — the push is refused rather than stored with \
a dangling base",
hex::encode(oid)
);
}
}
Ok(())
}
pub(crate) fn lookup_one(&self, oid: Oid<'_>) -> Result<Option<crate::index_layout::IndexRow>> {
if self.unindexed_packs() > 0 {
self.absorb_pending()?;
}
Ok(self.absorber.objects.lookup(oid))
}
pub(crate) fn ref_state(&self) -> Result<BTreeMap<String, crate::refs::RefState>> {
self.refs.current()
}
pub(crate) fn reachable_oids(&self, want: &[Oid<'_>], have: &[Oid<'_>]) -> Result<Vec<String>> {
Ok(self
.reachable_raw_with(want, have, live_reach_policy(), true)?
.iter()
.map(hex::encode)
.collect())
}
#[cfg(test)]
pub(crate) fn reachable_oids_with(
&self,
want: &[Oid<'_>],
have: &[Oid<'_>],
policy: crate::reach::ReachPolicy,
cache: bool,
) -> Result<Vec<String>> {
Ok(self
.reachable_raw_with(want, have, policy, cache)?
.iter()
.map(hex::encode)
.collect())
}
pub(crate) fn reachable_raw(
&self,
want: &[Oid<'_>],
have: &[Oid<'_>],
) -> Result<git_storage_trait::OidList> {
self.reachable_raw_with(want, have, live_reach_policy(), true)
}
pub(crate) fn reachable_raw_with(
&self,
want: &[Oid<'_>],
have: &[Oid<'_>],
policy: crate::reach::ReachPolicy,
cache: bool,
) -> Result<git_storage_trait::OidList> {
if self.unindexed_packs() > 0 {
self.absorb_pending()?;
}
let bitmaps = self.reach_bitmaps_with(policy, cache)?;
let by_commit: HashMap<&str, &roaring::RoaringBitmap> = bitmaps
.iter()
.map(|e| (e.commit.as_str(), &e.bitmap))
.collect();
let d = self
.absorber
.derived
.read()
.map_err(|_| anyhow!("derived poisoned"))?;
let space = &d.oids_raw;
let ordinal = &d.ordinal;
let oid_len = self.hash.oid_len();
let graph: HashMap<&str, &CommitNode> =
d.graph.iter().map(|n| (n.oid.as_str(), n)).collect();
let facts = crate::reach::ObjectFacts {
ordinal,
trees: &d.trees,
oid_len,
};
let mut union = roaring::RoaringBitmap::new();
for oid in want {
let hex = hex::encode(oid);
match by_commit.get(hex.as_str()) {
Some(bm) => union |= *bm,
None => {
if graph.contains_key(hex.as_str()) {
crate::reach::accumulate(
hex.as_str(),
&by_commit,
&graph,
&facts,
&mut union,
);
} else if let Some(&o) = ordinal.get(hex.as_str()) {
union.insert(o);
} else {
bail!(
"want {hex} is not in this repository — a negotiation must not be \
answered from a partial set"
);
}
}
}
}
let mut had = roaring::RoaringBitmap::new();
for oid in have {
let hex = hex::encode(oid);
if let Some(bm) = by_commit.get(hex.as_str()) {
had |= *bm;
} else if graph.contains_key(hex.as_str()) {
crate::reach::accumulate(hex.as_str(), &by_commit, &graph, &facts, &mut had);
} else if let Some(&o) = ordinal.get(hex.as_str()) {
had.insert(o);
}
}
let delta = union - had;
let mut out = git_storage_trait::OidList::with_capacity(delta.len() as usize, oid_len);
for o in delta {
let at = o as usize * oid_len;
let raw = space
.get(at..at + oid_len)
.ok_or_else(|| anyhow!("ordinal {o} is outside this store's ordinal space"))?;
out.push(raw)?;
}
Ok(out)
}
pub(crate) fn live_set(&self) -> Result<std::collections::HashSet<String>> {
let refs = self.ref_state()?;
let tips: Vec<Vec<u8>> = refs
.values()
.filter_map(|s| s.target.as_deref())
.filter_map(|t| hex::decode(t).ok())
.collect();
if tips.is_empty() {
bail!(
"this repository has no ref pointing at an object, so every object in it would be \
dead. A GC that would delete everything is refused: name a ref first"
);
}
let borrowed: Vec<Oid<'_>> = tips.iter().map(|t| t.as_slice()).collect();
let mut live: std::collections::HashSet<String> =
self.reachable_oids(&borrowed, &[])?.into_iter().collect();
for s in refs.values() {
if let Some(p) = &s.peeled {
live.insert(p.clone());
}
if let Some(t) = &s.target {
live.insert(t.clone());
}
}
Ok(live)
}
pub(crate) fn drop_dead_rows(&self, live: &std::collections::HashSet<String>) -> Result<u64> {
let dropped = self
.absorber
.objects
.retain(&|oid: &[u8]| live.contains(&hex::encode(oid)))?;
self.absorber
.exploded
.retain(&|oid: &[u8]| live.contains(&hex::encode(oid)))?;
self.absorber.refold()?;
Ok(dropped)
}
pub(crate) fn retire_dead_packs(
&self,
live: &std::collections::HashSet<String>,
) -> Result<Vec<u64>> {
let Some(journal) = self.arms.writer.journal(&self.blobs) else {
return Ok(Vec::new());
};
if !journal.exists() {
return Ok(Vec::new());
}
let rows = crate::archive_write::read_journal(&journal)?;
let packs = crate::archive_write::acked_packs(&rows);
let already = crate::archive_write::retired_offsets(&rows);
if packs.is_empty() {
return Ok(Vec::new());
}
let occupied = self.absorber.objects.extents_with_rows(&packs)?;
let mut order: Vec<usize> = (0..packs.len()).filter(|&i| packs[i].1 > 0).collect();
order.sort_unstable_by_key(|&i| packs[i].0);
let starts: Vec<u64> = order.iter().map(|&i| packs[i].0).collect();
let live_oids: Vec<Vec<u8>> = live.iter().filter_map(|h| hex::decode(h).ok()).collect();
let borrowed: Vec<&[u8]> = live_oids.iter().map(|o| o.as_slice()).collect();
let mut has_live = vec![false; packs.len()];
for (offset, _) in self
.absorber
.objects
.extents_batch(&borrowed)
.into_iter()
.flatten()
{
let p = starts.partition_point(|&s| s <= offset);
if p == 0 {
continue;
}
let i = order[p - 1];
if offset < packs[i].0 + packs[i].1 {
has_live[i] = true;
}
}
let dead: Vec<u64> = (0..packs.len())
.filter(|&i| occupied[i] && !has_live[i] && !already.contains(&packs[i].0))
.map(|i| packs[i].0)
.collect();
crate::archive_write::retire_packs(&journal, &dead)?;
Ok(dead)
}
pub(crate) fn seal_archive(
&self,
reserved: Vec<ReservedSection>,
) -> Result<crate::archive_write::SealReport> {
crate::archive_write::seal_generation_zero(
&self.blobs,
self.arms.writer.journal(&self.blobs).as_deref(),
&self.archive,
reserved,
)
}
pub(crate) fn reserved_sections(&self) -> Result<Vec<ReservedSection>> {
use znippy_common::{GUNNAR_GRAPH_MODULE, GUNNAR_REACH_MODULE};
let mut out = vec![self.refs.seal_section()?];
let d = self
.absorber
.derived
.read()
.map_err(|_| anyhow!("derived poisoned"))?;
if !d.graph.is_empty() {
out.push(ReservedSection::arrow(
GUNNAR_GRAPH_MODULE,
crate::graph::graph_schema(),
vec![crate::graph::build_graph_batch(&d.graph)?],
));
}
drop(d);
let reach = self.reach_bitmaps()?;
if !reach.is_empty() {
out.push(ReservedSection::arrow(
GUNNAR_REACH_MODULE,
crate::reach::reach_schema(),
vec![crate::reach::build_reach_batch(&reach)?],
));
}
Ok(out)
}
}
pub enum SelectedStore {
OneTableFourColumns(GitStore<OneTableFourColumns>),
FourTables(GitStore<crate::index_layout::FourTables>),
PackedPayload(GitStore<crate::index_layout::PackedPayload>),
}
macro_rules! on_arm {
($self:ident, $s:ident => $body:expr) => {
match $self {
SelectedStore::OneTableFourColumns($s) => $body,
SelectedStore::FourTables($s) => $body,
SelectedStore::PackedPayload($s) => $body,
}
};
}
impl SelectedStore {
pub fn arms(&self) -> StoreConfig {
on_arm!(self, s => s.arms())
}
pub fn writer_name(&self) -> &'static str {
on_arm!(self, s => s.writer_name())
}
pub fn writer_durability(&self) -> &'static str {
on_arm!(self, s => s.writer_durability())
}
pub fn index_ipc_bytes(&self) -> usize {
on_arm!(self, s => s.index().ipc_bytes())
}
pub fn index_name(&self) -> &'static str {
on_arm!(self, s => s.index().projection_name())
}
pub fn oids_in_extent(&self, extent: Extent) -> Result<Vec<Vec<u8>>> {
on_arm!(self, s => s.oids_in_extent(extent))
}
pub fn object_count(&self) -> usize {
on_arm!(self, s => s.object_count())
}
pub fn wait_indexed(&self) {
on_arm!(self, s => s.wait_indexed())
}
pub fn rebuild_projection(&self) -> Result<()> {
on_arm!(self, s => s.index().rebuild())
}
pub fn seal(&self) -> Result<Vec<ReservedSection>> {
on_arm!(self, s => s.seal())
}
}
impl GitOps for SelectedStore {
fn put(&self, pack: &[u8], refs: &[RefUpdate]) -> Result<TxId> {
on_arm!(self, s => s.put(pack, refs))
}
fn put_pack(&self, bytes: &[u8]) -> Result<TxId> {
on_arm!(self, s => s.put_pack(bytes))
}
fn put_refs(&self, updates: &[RefUpdate]) -> Result<TxId> {
on_arm!(self, s => s.put_refs(updates))
}
fn get(&self, oid: Oid<'_>) -> Result<Option<Stored>> {
on_arm!(self, s => s.get(oid))
}
fn has(&self, oid: Oid<'_>) -> Result<bool> {
on_arm!(self, s => s.has(oid))
}
fn size(&self, oid: Oid<'_>) -> Result<Option<u64>> {
on_arm!(self, s => s.size(oid))
}
fn extents(&self, oids: &[Oid<'_>]) -> Result<Vec<Option<Extent>>> {
on_arm!(self, s => s.extents(oids))
}
fn refs(&self) -> Result<Vec<RefRow>> {
on_arm!(self, s => s.refs())
}
fn update_ref(&self, name: &str, old: Option<Oid<'_>>, new: Option<Oid<'_>>) -> Result<TxId> {
on_arm!(self, s => s.update_ref(name, old, new))
}
fn put_refs_cas(&self, edits: &[git_storage_trait::RefCas<'_>]) -> Result<TxId> {
on_arm!(self, s => s.put_refs_cas(edits))
}
fn reachable(&self, want: &[Oid<'_>], have: &[Oid<'_>]) -> Result<Vec<Vec<u8>>> {
on_arm!(self, s => s.reachable(want, have))
}
fn gc(&self) -> Result<GcReport> {
on_arm!(self, s => s.gc())
}
}
impl crate::serve::GitServe for SelectedStore {
fn read(&self, oid: Oid<'_>) -> Result<Option<(crate::index_layout::ObjType, Vec<u8>)>> {
on_arm!(self, s => crate::serve::GitServe::read(s, oid))
}
fn header(&self, oid: Oid<'_>) -> Result<Option<(crate::index_layout::ObjType, u64)>> {
on_arm!(self, s => crate::serve::GitServe::header(s, oid))
}
fn sizes(&self, oids: &[Oid<'_>]) -> Result<Vec<Option<u64>>> {
on_arm!(self, s => crate::serve::GitServe::sizes(s, oids))
}
fn head(&self) -> Result<Option<RefRow>> {
on_arm!(self, s => crate::serve::GitServe::head(s))
}
fn set_head(&self, target: &str) -> Result<TxId> {
on_arm!(self, s => crate::serve::GitServe::set_head(s, target))
}
fn emit_pack(
&self,
objects: &[Oid<'_>],
have: &[Oid<'_>],
caps: &git_storage_trait::Caps,
out: &mut dyn std::io::Write,
) -> Result<git_storage_trait::PackStats> {
on_arm!(self, s => crate::serve::GitServe::emit_pack(s, objects, have, caps, out))
}
fn select(
&self,
want: &[Oid<'_>],
have: &[Oid<'_>],
) -> Result<Option<git_storage_trait::ReachSet>> {
on_arm!(self, s => crate::serve::GitServe::select(s, want, have))
}
}
pub fn open_selected(
root: &Path,
account: &str,
hash: GitHashKind,
arms: StoreConfig,
) -> Result<SelectedStore> {
use crate::arms::IndexArm;
use crate::index_layout::{FourTables, PackedPayload};
Ok(match arms.index {
IndexArm::OneTableFourColumns => {
SelectedStore::OneTableFourColumns(GitStore::<OneTableFourColumns>::open_with_arms(
root, account, hash, arms,
)?)
}
IndexArm::FourTables => SelectedStore::FourTables(GitStore::<FourTables>::open_with_arms(
root, account, hash, arms,
)?),
IndexArm::PackedPayload => SelectedStore::PackedPayload(
GitStore::<PackedPayload>::open_with_arms(root, account, hash, arms)?,
),
})
}
pub fn open_from_env(root: &Path, account: &str, hash: GitHashKind) -> Result<SelectedStore> {
let arms = StoreConfig::from_env()?;
open_selected(root, account, hash, arms)
}
impl<S: ObjectIndex> Absorber<S> {
fn resolved(&self, oid: &[u8]) -> Result<Option<(GitObjectKind, Vec<u8>)>> {
if let Some(hit) = self.exploded.content(oid)? {
return Ok(Some(hit));
}
let Some(row) = self.objects.lookup(oid) else {
return Ok(None);
};
let found = {
let d = self
.derived
.read()
.map_err(|_| anyhow!("derived poisoned"))?;
d.packs
.iter()
.find(|(o, l)| row.offset >= *o && row.offset < o + l)
.copied()
};
let Some((pack_offset, pack_len)) = found else {
return Ok(None);
};
self.exploded.note_rederived();
let bytes = self.read_extent(pack_offset, pack_len)?;
let walked = crate::pack_walk::walk(&bytes, self.hash.oid_len())?;
let capture = crate::exploded::CaptureOne::new(oid);
crate::resolve::resolve_walked(
&bytes,
&walked,
self.hash,
pack_offset,
&crate::resolve::NoBases,
&capture,
)?;
Ok(capture.take())
}
}
impl<S: ObjectIndex> BaseSource for Absorber<S> {
fn content(&self, oid: &[u8]) -> Option<(GitObjectKind, Vec<u8>)> {
self.resolved(oid).ok().flatten()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::archive_write::{read_journal, SafeWriter};
use crate::object::GitObjectKind;
use crate::store::tests::{one_blob_pack, real_pack, tmpdir};
use std::collections::HashSet;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering as AtomicOrdering};
use std::time::{Duration, Instant};
use znippy_zoomies::background::Job;
use znippy_zoomies::gatling_forkjoin::gatling_for_each;
fn loadavg() -> String {
std::fs::read_to_string("/proc/loadavg")
.unwrap_or_default()
.split_whitespace()
.take(3)
.collect::<Vec<_>>()
.join(" ")
}
#[test]
fn the_raw_ordinal_space_and_the_hex_one_are_the_same_sequence() {
let dir = tmpdir("ordinal-space-two-spellings");
let (pack, rows) = real_pack();
let store = GitStore::open(&dir, "rickard").unwrap();
store.put_pack(&pack).unwrap();
store.wait_indexed();
let oid_len = store.hash_kind().oid_len();
let (hex, raw) = {
let d = store.absorber.derived.read().unwrap();
(d.oids.clone(), d.oids_raw.clone())
};
assert!(
!hex.is_empty(),
"the fixture pack ({} entries) folded to an EMPTY ordinal space, so every \
comparison below would hold vacuously",
rows.len()
);
assert_eq!(
raw.len(),
hex.len() * oid_len,
"the raw ordinal space is {} bytes for {} oids of {oid_len} bytes",
raw.len(),
hex.len()
);
let mut checked = 0usize;
for (o, h) in hex.iter().enumerate() {
assert_eq!(
&raw[o * oid_len..(o + 1) * oid_len],
hex::decode(h).unwrap().as_slice(),
"ordinal {o} is {h} in the hex space and something else in the raw one"
);
checked += 1;
}
assert_eq!(checked, hex.len());
let tip = store
.graph_snapshot()
.into_iter()
.max_by_key(|c| c.generation)
.expect("a graph");
let tip_raw = hex::decode(&tip.oid).unwrap();
let flat = store.reachable_raw(&[&tip_raw], &[]).unwrap();
let hexed = store.reachable_oids(&[&tip_raw], &[]).unwrap();
assert!(
flat.len() > 1,
"the tip's closure is {} object(s) — too small to tell an ordering bug from a \
coincidence",
flat.len()
);
assert_eq!(
flat.iter().map(hex::encode).collect::<Vec<_>>(),
hexed,
"the serving answer and the maintenance answer name different objects"
);
assert!(flat.contains(&tip_raw), "a closure without its own tip");
eprintln!(
"load {}; {} ordinals agree in both spellings, tip closure {} objects",
loadavg(),
checked,
flat.len()
);
}
fn commits_in(rows: &[crate::resolve::Resolved]) -> usize {
rows.iter()
.filter(|r| r.kind == GitObjectKind::Commit)
.map(|r| r.oid.clone())
.collect::<HashSet<Vec<u8>>>()
.len()
}
#[test]
fn a_push_ends_in_object_rows_and_no_read_asked_for_them() {
let dir = tmpdir("drain-rows");
let store = GitStore::open(&dir, "rickard").unwrap();
let (pack, rows) = real_pack();
let t = Instant::now();
let tx = store.put_pack(&pack).unwrap();
let ack = t.elapsed();
let pack_id = tx.pack_id.expect("a pack push assigns an id");
store.wait_indexed();
let drained = t.elapsed();
assert_eq!(
store.object_count(),
rows.len(),
"the drain built no object rows: {} of {}",
store.object_count(),
rows.len()
);
assert_eq!(store.unindexed_packs(), 0, "the bit is still clear");
assert!(
store.indexer().is_indexed(pack_id),
"the pack's indexed bit is not up after the drain"
);
assert_eq!(
store.indexer().absorb_failures(),
0,
"the drain recorded an absorb failure: {:?}",
store.indexer().last_absorb_error()
);
for (i, r) in rows.iter().enumerate() {
let row = store
.index()
.lookup(&r.oid)
.unwrap_or_else(|| panic!("object {i} {} has no row", hex::encode(&r.oid)));
assert_eq!(
(row.offset, row.len),
(tx.extent.unwrap().0 + r.offset, r.len)
);
assert_eq!(row.uncompressed_size, r.uncompressed_size);
}
assert_eq!(
store.absorb_pending().unwrap(),
0,
"the drain left work for a read to do"
);
assert_eq!(
store.commit_count(),
commits_in(&rows),
"the commit graph does not match the pack"
);
let secs = (drained - ack).as_secs_f64();
eprintln!(
"load {}; {} objects in {} bytes: ack {:.0} µs, drain {:.1} ms, {:.0} object rows/s",
loadavg(),
rows.len(),
pack.len(),
ack.as_secs_f64() * 1e6,
(drained - ack).as_secs_f64() * 1e3,
rows.len() as f64 / secs.max(1e-9),
);
}
const KILL_DIR: &str = "GUNNAR_KILL_MID_ABSORB_DIR";
#[test]
#[ignore = "spawned by a_store_reopened_after_a_kill_mid_absorb_requeues_the_pack; it aborts on purpose"]
fn the_child_that_dies_between_durable_bytes_and_indexed_rows() {
let Ok(dir) = std::env::var(KILL_DIR) else {
return;
};
let dir = PathBuf::from(dir);
let pack = std::fs::read(dir.join("fixture.pack")).expect("the parent's fixture");
let store = GitStore::open(&dir, "rickard").expect("open");
let _gate = store.hold_absorb_gate();
let tx = store.put_pack(&pack).expect("the ack path");
assert_eq!(tx.pack_id, Some(0), "a fresh archive starts at ordinal 0");
assert_eq!(store.unindexed_packs(), 1, "the pack is not queued");
assert_eq!(store.index().len(), 0, "a row landed before the kill");
std::fs::write(dir.join("ready"), b"durable, not indexed").expect("marker");
std::process::abort();
}
#[test]
fn a_store_reopened_after_a_kill_mid_absorb_requeues_the_pack() {
use std::os::unix::process::ExitStatusExt;
let dir = tmpdir("kill-mid-absorb");
let (pack, rows) = real_pack();
std::fs::write(dir.join("fixture.pack"), &pack).unwrap();
let status = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"--ignored",
"--nocapture",
"git_ops::tests::the_child_that_dies_between_durable_bytes_and_indexed_rows",
])
.env(KILL_DIR, &dir)
.status()
.expect("spawn the child that dies");
assert_eq!(
status.signal(),
Some(6),
"the child did not die by SIGABRT — it {status:?}, so nothing was interrupted"
);
assert!(
dir.join("ready").exists(),
"the child never reached the kill point"
);
let blobs = dir.join("objects.pack");
let acked = read_journal(&SafeWriter::journal_path(&blobs)).unwrap();
assert_eq!(
acked.len(),
1,
"the pack's extent is not durable: {acked:?}"
);
assert_eq!(acked[0].1, pack.len() as u64);
let on_disk = std::fs::read(&blobs).unwrap();
assert_eq!(
&on_disk[acked[0].0 as usize..(acked[0].0 + acked[0].1) as usize],
&pack[..],
"the verbatim bytes did not survive the kill"
);
{
let tail = ObjectReadStack::<OneTableFourColumns>::open(
&dir.join("objects.tail"),
RebuildTriggers::default(),
crate::arms::DEFAULT_REDB_CACHE_BYTES,
)
.unwrap();
assert_eq!(
tail.len(),
0,
"the kill landed after the rows were indexed, not between the bytes and the rows \
— this guard would then be testing nothing"
);
}
let store = GitStore::open(&dir, "rickard").unwrap();
assert!(
store.has(&rows[0].oid).unwrap(),
"a durable object read `absent` from a store reopened over the pack that holds it — \
the crash-recovery bit did not re-queue anything"
);
store.wait_indexed();
assert_eq!(
store.object_count(),
rows.len(),
"the re-queued pack's objects never reached the index — left: {}, right: {}",
store.object_count(),
rows.len()
);
assert!(
store.indexer().is_indexed(0),
"the re-queued pack's indexed bit never went up in the new process"
);
assert_eq!(store.unindexed_packs(), 0);
assert_eq!(store.absorb_pending().unwrap(), 0);
for r in rows.iter().take(512) {
let row = store
.index()
.lookup(&r.oid)
.unwrap_or_else(|| panic!("{} has no row after recovery", hex::encode(&r.oid)));
assert_eq!((row.offset, row.len), (acked[0].0 + r.offset, r.len));
}
assert_eq!(
store.commit_count(),
commits_in(&rows),
"the commit graph did not come back: {} of {}",
store.commit_count(),
commits_in(&rows)
);
let tip = store
.graph_snapshot()
.into_iter()
.max_by_key(|c| c.generation)
.expect("a graph");
let tip_raw = hex::decode(&tip.oid).unwrap();
let closure = store.reachable(&[&tip_raw], &[]).unwrap();
assert!(
!closure.is_empty(),
"reachable() came back empty after the reopen"
);
assert!(closure.contains(&tip_raw));
eprintln!(
"load {}; killed mid-absorb: {} journal extent(s) recovered {} object rows, {} \
commits, closure of the tip {} objects",
loadavg(),
acked.len(),
store.object_count(),
store.commit_count(),
closure.len(),
);
}
#[test]
fn a_clean_reopen_requeues_nothing_and_the_next_push_gets_a_fresh_ordinal() {
let dir = tmpdir("clean-reopen");
let (pack, rows) = real_pack();
let (second, blob_oid) = one_blob_pack(b"a push that arrives after the restart");
{
let store = GitStore::open(&dir, "rickard").unwrap();
store.put_pack(&pack).unwrap();
store.wait_indexed();
assert_eq!(
store.object_count(),
rows.len(),
"the first run did not index"
);
}
let store = GitStore::open(&dir, "rickard").unwrap();
store.wait_indexed();
assert!(
!store.indexer().is_indexed(0),
"a fully absorbed pack was re-queued on open — the diff is not exact"
);
assert_eq!(store.unindexed_packs(), 0);
assert_eq!(
store.absorb_pending().unwrap(),
0,
"the reopen left index work"
);
assert_eq!(
store.object_count(),
rows.len(),
"the reopened store lost rows"
);
assert!(store.has(&rows[0].oid).unwrap());
let tx = store.put_pack(&second).unwrap();
assert_eq!(
tx.pack_id,
Some(1),
"a reopened store handed a fresh pack the ordinal of an absorbed one"
);
store.wait_indexed();
assert!(
store.has(&blob_oid).unwrap(),
"the second push's object never reached the index — its ordinal collided with an \
absorbed pack's, so every writer skipped it"
);
assert_eq!(store.object_count(), rows.len() + 1);
assert!(store.indexer().is_indexed(1));
let acked = read_journal(&SafeWriter::journal_path(store.blobs_path())).unwrap();
assert_eq!(
acked.len(),
2,
"the reopen truncated the journal: {acked:?}"
);
assert_eq!(acked[1], tx.extent.unwrap());
}
#[test]
fn the_derive_on_open_diff_costs_milliseconds_at_a_realistic_pack_count() {
const PACKS: usize = 1000;
let dir = tmpdir("diff-cost");
let (real, rows) = real_pack();
{
let store = GitStore::open(&dir, "rickard").unwrap();
store.put_pack(&real).unwrap();
for i in 0..PACKS - 1 {
let (p, _) = one_blob_pack(format!("pack number {i}").as_bytes());
store.put_pack(&p).unwrap();
}
store.wait_indexed();
assert_eq!(
store.object_count(),
rows.len() + PACKS - 1,
"the fixture did not absorb"
);
}
let acked = read_journal(&SafeWriter::journal_path(&dir.join("objects.pack"))).unwrap();
assert_eq!(acked.len(), PACKS, "one journal row per pack");
let tail = ObjectReadStack::<OneTableFourColumns>::open(
&dir.join("objects.tail"),
RebuildTriggers::default(),
crate::arms::DEFAULT_REDB_CACHE_BYTES,
)
.unwrap();
let t = Instant::now();
let hit = tail.extents_with_rows(&acked).unwrap();
let clean = t.elapsed();
assert!(
hit.iter().all(|&h| h),
"an absorbed pack was reported unabsorbed"
);
let mut in_flight = acked.clone();
let end = acked[PACKS - 1].0 + acked[PACKS - 1].1;
in_flight.push((end, 4096));
let t = Instant::now();
let hit = tail.extents_with_rows(&in_flight).unwrap();
let crashed = t.elapsed();
assert_eq!(
hit.iter().filter(|h| !**h).count(),
1,
"the pack with no rows is not the only one reported unabsorbed"
);
eprintln!(
"load {}; {PACKS} packs / {} object rows: derive-on-open diff {:.2} ms clean \
(early exit), {:.2} ms with one pack un-absorbed (full tail scan)",
loadavg(),
tail.len(),
clean.as_secs_f64() * 1e3,
crashed.as_secs_f64() * 1e3,
);
}
#[test]
fn the_indexed_bit_is_a_bitset_over_dense_pack_ordinals() {
let ids = [0u64, 1, 63, 64, 65, 4095];
let mut s = PackState::default();
for id in ids {
s.note(id, (id * 4096 + 12, 4096));
}
assert_eq!(s.unabsorbed, ids.len());
assert!(
ids.iter().all(|&i| !s.is_absorbed(i)),
"a bit is up already"
);
s.mark_absorbed(64, (64 * 4096 + 12, 4096));
assert!(s.is_absorbed(64), "the bit for ordinal 64 did not go up");
for other in [0u64, 1, 63, 65, 4095] {
assert!(
!s.is_absorbed(other),
"setting ordinal 64 also set ordinal {other} — one bit spilled onto a neighbour"
);
}
assert_eq!(s.unabsorbed, ids.len() - 1);
assert_eq!(
s.pending().iter().map(|p| p.pack_id).collect::<Vec<_>>(),
vec![0, 1, 63, 65, 4095],
"the pending list is not the clear bits"
);
assert_eq!(
s.pending()[0],
PendingPack {
pack_id: 0,
offset: 12,
len: 4096
},
"a pending pack lost the extent it must be absorbed from"
);
s.mark_absorbed(64, (64 * 4096 + 12, 4096));
s.note(64, (64 * 4096 + 12, 4096));
assert!(s.is_absorbed(64));
assert_eq!(
s.unabsorbed,
ids.len() - 1,
"an absorbed pack was queued again"
);
assert_eq!(s.words.len(), 64);
assert_eq!(s.words.len() * 8, 512, "4096 packs is 512 bytes of bits");
}
#[test]
fn a_read_before_the_drain_waits_and_never_answers_absent() {
let dir = tmpdir("drain-before");
let store = Arc::new(GitStore::open(&dir, "rickard").unwrap());
let (pack, rows) = real_pack();
let gate = store.hold_absorb_gate();
let tx = store.put_pack(&pack).unwrap();
let pack_id = tx.pack_id.unwrap();
assert_eq!(store.unindexed_packs(), 1, "the pack is not queued");
assert_eq!(
store.index().len(),
0,
"the index already holds rows — the 'before the drain' state is not the one under test"
);
assert!(
!store.indexer().is_indexed(pack_id),
"the bit is already up"
);
let oid = rows[0].oid.clone();
let answered = Arc::new(AtomicBool::new(false));
let (s, a) = (store.clone(), answered.clone());
let reader = Job::spawn(move || {
let got = s.has(&oid).unwrap();
a.store(true, AtomicOrdering::Release);
got
});
std::thread::sleep(Duration::from_millis(150));
assert!(
!answered.load(AtomicOrdering::Acquire),
"a read that arrived before the drain answered in 150 ms with an index holding {} \
rows — the only answer it can have given is `absent`",
store.index().len()
);
assert_eq!(
store.index().len(),
0,
"the index gained rows while the absorb gate was held"
);
assert!(
!store.indexer().is_indexed(pack_id),
"the pack's indexed bit went up while the drain is still parked in front of its \
absorb — the bit does not gate the object rows"
);
drop(gate);
assert!(
reader.join().unwrap(),
"a read that waited for the drain still answered absent"
);
store.wait_indexed();
assert_eq!(store.unindexed_packs(), 0);
assert!(store.indexer().is_indexed(pack_id));
assert_eq!(store.index().len(), rows.len());
assert_eq!(store.absorb_pending().unwrap(), 0);
for r in rows.iter().take(256) {
assert!(
store.index().lookup(&r.oid).is_some(),
"{} is not in the rows after the drain",
hex::encode(&r.oid)
);
}
}
#[test]
fn no_read_is_told_absent_while_the_drain_is_running() {
let dir = tmpdir("drain-during");
let store = GitStore::open(&dir, "rickard").unwrap();
let (pack, rows) = real_pack();
let oids: Vec<&[u8]> = rows.iter().map(|r| r.oid.as_slice()).collect();
let tx = store.put_pack(&pack).unwrap();
let pack_id = tx.pack_id.unwrap();
let early = AtomicU64::new(0);
let absent = AtomicU64::new(0);
let reads = AtomicU64::new(0);
let workers = 4usize;
gatling_for_each(workers, workers, |w| loop {
let still_pending = store.unindexed_packs() > 0;
for oid in oids.iter().skip(w).step_by(workers * 8) {
if still_pending {
early.fetch_add(1, AtomicOrdering::Relaxed);
}
reads.fetch_add(1, AtomicOrdering::Relaxed);
if !store.has(oid).unwrap() {
absent.fetch_add(1, AtomicOrdering::Relaxed);
}
}
if store.indexer().is_indexed(pack_id) {
break;
}
});
let absent = absent.load(AtomicOrdering::Acquire);
assert_eq!(
absent,
0,
"{workers} reader(s) were told a durable object was absent — {absent} absent answers \
over {} objects",
rows.len()
);
assert_eq!(store.object_count(), rows.len());
eprintln!(
"load {}; {} reads across {workers} gatling workers, {} of them while the pack was \
still un-indexed",
loadavg(),
reads.load(AtomicOrdering::Acquire),
early.load(AtomicOrdering::Acquire),
);
}
#[test]
fn a_pack_is_absorbed_once_even_when_both_callers_race_for_it() {
let dir = tmpdir("drain-once");
let store = Arc::new(GitStore::open(&dir, "rickard").unwrap());
let (pack, rows) = real_pack();
let gate = store.hold_absorb_gate();
store.put_pack(&pack).unwrap();
let s = store.clone();
let reader = Job::spawn(move || s.absorb_pending().unwrap());
std::thread::sleep(Duration::from_millis(100));
drop(gate);
let by_read = reader.join().unwrap();
store.wait_indexed();
assert_eq!(
store.commit_count(),
commits_in(&rows),
"the commit graph carries {} commits, the pack has {}",
store.commit_count(),
commits_in(&rows)
);
assert_eq!(store.object_count(), rows.len());
assert_eq!(store.unindexed_packs(), 0);
eprintln!(
"load {}; the read absorbed {by_read} pack(s), the drain absorbed the rest",
loadavg()
);
}
#[test]
fn a_push_produces_an_exploded_row_for_every_object_in_the_pack() {
let dir = tmpdir("exploded-every-object");
let store = GitStore::open(&dir, "rickard").unwrap();
let (pack, rows) = real_pack();
let t = Instant::now();
store.put_pack(&pack).unwrap();
let ack = t.elapsed();
store.wait_indexed();
let drain = t.elapsed();
let distinct: HashSet<Vec<u8>> = rows.iter().map(|r| r.oid.clone()).collect();
let stats = store.exploded_stats();
assert_eq!(
stats.rows,
distinct.len() as u64,
"the drain built no exploded rows: {} of {}",
stats.rows,
distinct.len()
);
assert_eq!(
stats.written, stats.rows,
"rows exist that this process never committed"
);
let hash = store.hash_kind();
let mut checked = 0usize;
for oid in &distinct {
let (kind, payload) = store
.content(oid)
.unwrap()
.unwrap_or_else(|| panic!("{} has no exploded row", hex::encode(oid)));
let rehashed = hash.oid_of(&crate::object::canonical(kind, &payload));
assert_eq!(
&rehashed,
oid,
"the exploded row filed under {} contains an object that hashes to {}",
hex::encode(oid),
hex::encode(&rehashed)
);
checked += 1;
}
assert_eq!(checked, distinct.len());
assert_eq!(
store.exploded_stats().rederived,
0,
"a content read re-resolved a pack even though the table was whole"
);
let inflated: u64 = rows.iter().map(|r| r.uncompressed_size).sum();
let table = std::fs::metadata(store.exploded_path())
.map(|m| m.len())
.unwrap_or(0);
eprintln!(
"load {}; {} pack entries → {} exploded rows, all {} re-hashed to their own oid; \
ack {:.1} ms, drain {:.1} ms, {:.0} rows/s; disk: pack {:.1} MiB verbatim, objects \
inflate to {:.1} MiB, table file {:.1} MiB ({:.0}x the pack, {:.3}x the payload)",
loadavg(),
rows.len(),
stats.rows,
checked,
ack.as_secs_f64() * 1e3,
drain.as_secs_f64() * 1e3,
stats.rows as f64 / drain.as_secs_f64(),
pack.len() as f64 / (1 << 20) as f64,
inflated as f64 / (1 << 20) as f64,
table as f64 / (1 << 20) as f64,
table as f64 / pack.len() as f64,
table as f64 / inflated.max(1) as f64,
);
}
#[test]
#[ignore = "needs a pack: ZNIPPY_EXPLODE_BENCH_PACK=<path> cargo test … -- --ignored --nocapture"]
fn the_table_costs_what_it_holds() {
let Ok(path) = std::env::var("ZNIPPY_EXPLODE_BENCH_PACK") else {
eprintln!("no ZNIPPY_EXPLODE_BENCH_PACK named; nothing measured");
return;
};
let pack = std::fs::read(&path).unwrap_or_else(|e| panic!("reading {path}: {e}"));
let dir = tmpdir("exploded-cost");
let store = GitStore::open(&dir, "rickard").unwrap();
let t = Instant::now();
store.put_pack(&pack).unwrap();
let ack = t.elapsed();
store.wait_indexed();
let drain = t.elapsed();
let stats = store.exploded_stats();
let table = std::fs::metadata(store.exploded_path())
.map(|m| m.len())
.unwrap_or(0);
let mut held = 0u64;
for kind in [
GitObjectKind::Commit,
GitObjectKind::Tree,
GitObjectKind::Blob,
GitObjectKind::Tag,
] {
held += store
.exploded_of_kind(kind)
.unwrap()
.iter()
.map(|(_, p)| p.len() as u64)
.sum::<u64>();
}
let mib = |n: u64| n as f64 / (1 << 20) as f64;
eprintln!(
"load {}; pack {} — {:.1} MiB verbatim → {} rows holding {:.1} MiB, table file \
{:.1} MiB = {:.3}x the payload and {:.1}x the pack; ack {:.0} ms, drain {:.0} ms",
loadavg(),
path,
mib(pack.len() as u64),
stats.rows,
mib(held),
mib(table),
table as f64 / held.max(1) as f64,
table as f64 / pack.len() as f64,
ack.as_secs_f64() * 1e3,
drain.as_secs_f64() * 1e3,
);
}
const CLEAN_DIR: &str = "GUNNAR_CLEAN_SHUTDOWN_DIR";
#[test]
#[ignore = "spawned by a_clean_shutdown_and_reopen_keeps_the_graph_and_reachable"]
fn the_child_that_pushes_and_shuts_down_cleanly() {
let Ok(dir) = std::env::var(CLEAN_DIR) else {
return;
};
let dir = PathBuf::from(dir);
let pack = std::fs::read(dir.join("fixture.pack")).expect("the parent's fixture");
let store = GitStore::open(&dir, "rickard").expect("open");
store.put_pack(&pack).expect("the ack path");
store.wait_indexed();
assert_eq!(store.unindexed_packs(), 0, "the child shut down mid-index");
let tip = store
.graph_snapshot()
.into_iter()
.max_by_key(|c| c.generation)
.expect("the child built no graph at all");
let tip_raw = hex::decode(&tip.oid).unwrap();
let closure = store.reachable(&[&tip_raw], &[]).unwrap();
assert!(!closure.is_empty(), "the child's own reachable() was empty");
std::fs::write(dir.join("tip"), &tip.oid).expect("marker");
std::fs::write(dir.join("closure"), closure.len().to_string()).expect("marker");
std::fs::write(dir.join("commits"), store.commit_count().to_string()).expect("marker");
std::fs::write(dir.join("objects"), store.object_count().to_string()).expect("marker");
drop(store);
std::fs::write(dir.join("clean"), b"closed").expect("marker");
}
#[test]
fn a_clean_shutdown_and_reopen_keeps_the_graph_and_reachable() {
let dir = tmpdir("clean-shutdown-graph");
let (pack, rows) = real_pack();
std::fs::write(dir.join("fixture.pack"), &pack).unwrap();
let status = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"--ignored",
"--nocapture",
"git_ops::tests::the_child_that_pushes_and_shuts_down_cleanly",
])
.env(CLEAN_DIR, &dir)
.status()
.expect("spawn the child that shuts down cleanly");
assert!(
status.success(),
"the child did not exit cleanly — it {status:?}, so this is a crash guard and not a \
clean-shutdown one"
);
assert!(
dir.join("clean").exists(),
"the child never reached its clean shutdown"
);
let tip_hex = std::fs::read_to_string(dir.join("tip")).unwrap();
let closure_before: usize = std::fs::read_to_string(dir.join("closure"))
.unwrap()
.parse()
.unwrap();
let commits_before: usize = std::fs::read_to_string(dir.join("commits"))
.unwrap()
.parse()
.unwrap();
assert_eq!(
commits_before,
commits_in(&rows),
"the child's own graph was already wrong, so the reopen proves nothing"
);
let store = GitStore::open(&dir, "rickard").unwrap();
assert_eq!(
store.commit_count(),
commits_in(&rows),
"the commit graph did not survive a clean shutdown: {} of {} commits — the exploded \
table is a warm cache, not a durable table",
store.commit_count(),
commits_in(&rows)
);
let tip_raw = hex::decode(&tip_hex).unwrap();
let closure = store.reachable(&[&tip_raw], &[]).unwrap();
assert_eq!(
closure.len(),
closure_before,
"reachable() answers differently after a clean reopen: {} objects against the {} the \
same tip closed over before the shutdown",
closure.len(),
closure_before
);
assert!(closure.contains(&tip_raw));
store.wait_indexed();
assert!(
!store.indexer().is_indexed(0),
"the reopen re-queued the pack, so the graph above is a recovery artefact rather \
than the table's"
);
assert_eq!(store.unindexed_packs(), 0);
assert_eq!(
store.absorb_pending().unwrap(),
0,
"the reopen left index work"
);
assert_eq!(store.object_count(), rows.len());
assert_eq!(
store.exploded_stats().written,
0,
"this process wrote exploded rows, so the table was rebuilt rather than read"
);
eprintln!(
"load {}; clean shutdown → reopen: {} rows, {} commits, tip closure {} objects \
(rebuilt nothing)",
loadavg(),
store.object_count(),
store.commit_count(),
closure.len(),
);
}
#[test]
fn a_content_read_is_served_from_the_table_and_not_re_derived() {
const READS: usize = 512;
let dir = tmpdir("exploded-served");
let (pack, rows) = real_pack();
let oids: Vec<Vec<u8>> = rows
.iter()
.map(|r| r.oid.clone())
.collect::<HashSet<Vec<u8>>>()
.into_iter()
.take(READS)
.collect();
{
let store = GitStore::open(&dir, "rickard").unwrap();
store.put_pack(&pack).unwrap();
store.wait_indexed();
let before = store.exploded_stats();
let t = Instant::now();
for oid in &oids {
assert!(
store.content(oid).unwrap().is_some(),
"the table lost an object"
);
}
let served_in = t.elapsed();
let after = store.exploded_stats();
assert_eq!(
(
after.served - before.served,
after.rederived - before.rederived
),
(oids.len() as u64, 0),
"{} content reads re-resolved a whole pack instead of hitting the table — \
served {}, rederived {}",
oids.len(),
after.served - before.served,
after.rederived - before.rederived,
);
eprintln!(
"load {}; {} content reads from §14's table in {:.1} ms ({:.0} ns/read)",
loadavg(),
oids.len(),
served_in.as_secs_f64() * 1e3,
served_in.as_secs_f64() * 1e9 / oids.len() as f64,
);
}
std::fs::remove_file(dir.join("objects.exploded")).unwrap();
let store = GitStore::open(&dir, "rickard").unwrap();
let gate = store.hold_absorb_gate();
let before = store.exploded_stats();
let mut rederived = 0u64;
for oid in oids.iter().take(8) {
assert!(
store.absorber.resolved(oid).unwrap().is_some(),
"the verbatim truth could not re-derive {}",
hex::encode(oid)
);
rederived += 1;
}
let after = store.exploded_stats();
drop(gate);
assert_eq!(
after.served - before.served,
0,
"a dropped table still reported {} reads served",
after.served - before.served
);
assert_eq!(
after.rederived - before.rederived,
rederived,
"the fallback ran but was not counted"
);
}
#[test]
fn dropping_the_exploded_table_and_reopening_still_answers() {
let dir = tmpdir("exploded-dropped");
let (pack, rows) = real_pack();
let tip_hex;
let closure_before;
{
let store = GitStore::open(&dir, "rickard").unwrap();
store.put_pack(&pack).unwrap();
store.wait_indexed();
let tip = store
.graph_snapshot()
.into_iter()
.max_by_key(|c| c.generation)
.expect("a graph");
let tip_raw = hex::decode(&tip.oid).unwrap();
closure_before = store.reachable(&[&tip_raw], &[]).unwrap().len();
tip_hex = tip.oid;
}
let table = dir.join("objects.exploded");
assert!(table.exists(), "the store never built an exploded table");
std::fs::remove_file(&table).unwrap();
{
let fresh = crate::exploded_arrow::ExplodedArchive::open(&table).unwrap();
assert_eq!(fresh.rows().unwrap(), 0, "the drop did not drop anything");
}
let _ = std::fs::remove_file(&table);
let store = GitStore::open(&dir, "rickard").unwrap();
let tip_raw = hex::decode(&tip_hex).unwrap();
let closure = store.reachable(&[&tip_raw], &[]).unwrap();
assert_eq!(
closure.len(),
closure_before,
"reachable() after dropping the table: {} objects, against {} before",
closure.len(),
closure_before
);
assert!(store.has(&rows[0].oid).unwrap());
store.wait_indexed();
let stats = store.exploded_stats();
let distinct: HashSet<Vec<u8>> = rows.iter().map(|r| r.oid.clone()).collect();
assert_eq!(
stats.rows,
distinct.len() as u64,
"the dropped table was never rebuilt: {} of {} rows",
stats.rows,
distinct.len()
);
assert!(
stats.written > 0,
"the table came back without this process writing a row — it was never dropped"
);
assert_eq!(
store.commit_count(),
commits_in(&rows),
"the graph did not come back with the table: {} of {}",
store.commit_count(),
commits_in(&rows)
);
assert_eq!(
store.object_count(),
rows.len(),
"the rebuild lost index rows"
);
eprintln!(
"load {}; dropped and rebuilt: {} rows written, {} commits, tip closure {} objects",
loadavg(),
stats.written,
store.commit_count(),
closure.len(),
);
}
}