use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, RwLock};
use anyhow::{Context, Result, anyhow, bail};
use redb::{Database, ReadableTable, ReadableTableMetadata, TableDefinition};
use crate::index_layout::{IndexEntry, IndexRow, ObjType, ObjectIndex, OneTableFourColumns};
const OBJECTS: TableDefinition<&[u8], &[u8]> = TableDefinition::new("objects");
const META: TableDefinition<&str, u64> = TableDefinition::new("meta");
const META_ARRIVAL_SEQ: &str = "arrival_seq";
const TAIL_ROW_BYTES: usize = 41;
pub const TAIL_ORDINAL_BIT: u32 = 0x8000_0000;
pub fn is_tail_row(row: &IndexRow) -> bool {
row.ordinal & TAIL_ORDINAL_BIT != 0
}
fn tail_ordinal(seq: u64) -> u32 {
TAIL_ORDINAL_BIT | (seq as u32 & !TAIL_ORDINAL_BIT)
}
fn encode_row(seq: u64, e: &IndexEntry) -> [u8; TAIL_ROW_BYTES] {
let mut b = [0u8; TAIL_ROW_BYTES];
b[0..8].copy_from_slice(&seq.to_le_bytes());
b[8..16].copy_from_slice(&e.offset.to_le_bytes());
b[16..24].copy_from_slice(&e.len.to_le_bytes());
b[24] = e.obj_type.code();
b[25..33].copy_from_slice(&e.uncompressed_size.to_le_bytes());
b[33..41].copy_from_slice(&e.delta_base.to_le_bytes());
b
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct TailRow {
seq: u64,
offset: u64,
len: u64,
obj_type: ObjType,
uncompressed_size: u64,
delta_base: u64,
}
fn decode_row(b: &[u8]) -> Result<TailRow> {
if b.len() != TAIL_ROW_BYTES {
bail!("tail row is {} bytes, expected {TAIL_ROW_BYTES}", b.len());
}
let u64_at = |i: usize| {
let mut w = [0u8; 8];
w.copy_from_slice(&b[i..i + 8]);
u64::from_le_bytes(w)
};
Ok(TailRow {
seq: u64_at(0),
offset: u64_at(8),
len: u64_at(16),
obj_type: ObjType::from_code(b[24])
.ok_or_else(|| anyhow!("tail row carries object type code {}", b[24]))?,
uncompressed_size: u64_at(25),
delta_base: u64_at(33),
})
}
impl TailRow {
fn as_index_row(&self) -> IndexRow {
IndexRow {
ordinal: tail_ordinal(self.seq),
offset: self.offset,
len: self.len,
obj_type: self.obj_type,
uncompressed_size: self.uncompressed_size,
delta_base: self.delta_base,
}
}
fn as_entry(&self, oid: &[u8]) -> IndexEntry {
IndexEntry {
oid: oid.to_vec(),
offset: self.offset,
len: self.len,
obj_type: self.obj_type,
uncompressed_size: self.uncompressed_size,
delta_base: self.delta_base,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct RebuildTriggers {
pub tail_hits_per_row: f64,
pub min_tail_hits: u64,
pub tail_bytes: u64,
}
impl RebuildTriggers {
pub const DEFAULT_TAIL_HITS_PER_ROW: f64 = 1.0;
pub const DEFAULT_MIN_TAIL_HITS: u64 = 4096;
pub const DEFAULT_TAIL_BYTES: u64 = 64 * 1024 * 1024;
pub fn miss_threshold(&self, projection_rows: u64) -> Option<u64> {
if self.min_tail_hits == 0 {
return None;
}
let scaled = (self.tail_hits_per_row * projection_rows as f64) as u64;
Some(scaled.max(self.min_tail_hits))
}
pub fn manual() -> Self {
Self {
tail_hits_per_row: 0.0,
min_tail_hits: 0,
tail_bytes: 0,
}
}
}
impl Default for RebuildTriggers {
fn default() -> Self {
Self {
tail_hits_per_row: Self::DEFAULT_TAIL_HITS_PER_ROW,
min_tail_hits: Self::DEFAULT_MIN_TAIL_HITS,
tail_bytes: Self::DEFAULT_TAIL_BYTES,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RebuildReason {
TailHits(u64),
TailBytes(u64),
Explicit,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct StackStats {
pub sealed_rows: u64,
pub total_rows: u64,
pub tail_hits: u64,
pub absent: u64,
pub unabsorbed_rows: u64,
pub tail_txns: u64,
pub tail_bytes: u64,
pub rebuilds: u64,
pub generation: u64,
}
pub struct ObjectReadStack<S: ObjectIndex = OneTableFourColumns> {
tail: Arc<Database>,
projection: RwLock<Arc<S>>,
triggers: RebuildTriggers,
total_rows: AtomicU64,
tail_hits: AtomicU64,
absent: AtomicU64,
tail_bytes: AtomicU64,
rebuilds: AtomicU64,
generation: AtomicU64,
unabsorbed_size: AtomicU64,
unabsorbed_types: [AtomicU64; 8],
unabsorbed_rows: AtomicU64,
tail_txns: AtomicU64,
}
impl<S: ObjectIndex> ObjectReadStack<S> {
pub fn open(tail_path: &Path, triggers: RebuildTriggers, cache_bytes: usize) -> Result<Self> {
let db = Database::builder()
.set_cache_size(cache_bytes)
.create(tail_path)
.with_context(|| format!("opening object tail at {}", tail_path.display()))?;
Self::from_db(db, triggers)
}
pub fn projection_name(&self) -> &'static str {
self.projection.read().expect("projection lock").name()
}
pub fn in_memory(triggers: RebuildTriggers) -> Result<Self> {
let db = Database::builder()
.create_with_backend(redb::backends::InMemoryBackend::new())
.context("creating an in-memory object tail")?;
Self::from_db(db, triggers)
}
fn from_db(db: Database, triggers: RebuildTriggers) -> Result<Self> {
let w = db.begin_write()?;
{
let _ = w.open_table(OBJECTS)?;
let _ = w.open_table(META)?;
}
w.commit()?;
let db = Arc::new(db);
let entries = scan(&db)?;
let total = entries.len() as u64;
let projection = Arc::new(S::build(&entries)?);
Ok(Self {
tail: db,
projection: RwLock::new(projection),
triggers,
total_rows: AtomicU64::new(total),
tail_hits: AtomicU64::new(0),
absent: AtomicU64::new(0),
tail_bytes: AtomicU64::new(0),
rebuilds: AtomicU64::new(1),
generation: AtomicU64::new(1),
unabsorbed_size: AtomicU64::new(0),
unabsorbed_types: Default::default(),
unabsorbed_rows: AtomicU64::new(0),
tail_txns: AtomicU64::new(0),
})
}
#[inline]
pub fn projection_is_complete(&self) -> bool {
self.unabsorbed_rows.load(Ordering::Acquire) == 0
}
pub fn append(&self, entries: &[IndexEntry]) -> Result<Option<RebuildReason>> {
if entries.is_empty() {
return Ok(None);
}
let width = entries[0].oid.len();
if width != 20 && width != 32 {
bail!("oid width {width} is neither sha1 (20) nor sha256 (32)");
}
let w = self.tail.begin_write()?;
let mut added = 0u64;
let mut bytes = 0u64;
let mut added_size = 0u64;
let mut added_types = [0u64; 8];
{
let mut objects = w.open_table(OBJECTS)?;
let mut meta = w.open_table(META)?;
let mut seq = meta.get(META_ARRIVAL_SEQ)?.map(|v| v.value()).unwrap_or(0);
for e in entries {
if e.oid.len() != width {
bail!(
"mixed oid widths in one append: {width} and {}",
e.oid.len()
);
}
if let Some(existing) = objects.get(e.oid.as_slice())? {
let old = decode_row(existing.value())?;
let encoded = |t: ObjType| matches!(t, ObjType::OfsDelta | ObjType::RefDelta);
let kind_disagrees = !encoded(old.obj_type)
&& !encoded(e.obj_type)
&& old.obj_type != e.obj_type;
if kind_disagrees || old.uncompressed_size != e.uncompressed_size {
bail!(
"identity violation: {} is already stored as \
(offset {}, len {}, {}, size {}, delta_base {}) and cannot be \
redeclared as (offset {}, len {}, {}, size {}, delta_base {}) — \
the TYPE or the inflated SIZE differs, so the same oid is \
describing a different object. Placement (offset, len, \
delta_base) may differ freely: a re-pushed pack lands its \
copies elsewhere and the first writer's row stands.",
hex::encode(&e.oid),
old.offset,
old.len,
old.obj_type.as_str(),
old.uncompressed_size,
old.delta_base,
e.offset,
e.len,
e.obj_type.as_str(),
e.uncompressed_size,
e.delta_base,
);
}
continue;
}
objects.insert(e.oid.as_slice(), encode_row(seq, e).as_slice())?;
seq += 1;
added += 1;
bytes += e.len;
added_size += e.uncompressed_size;
added_types[e.obj_type.code() as usize] += 1;
}
meta.insert(META_ARRIVAL_SEQ, seq)?;
}
w.commit()?;
self.total_rows.fetch_add(added, Ordering::AcqRel);
self.unabsorbed_rows.fetch_add(added, Ordering::AcqRel);
self.unabsorbed_size.fetch_add(added_size, Ordering::AcqRel);
for (slot, n) in self.unabsorbed_types.iter().zip(added_types) {
slot.fetch_add(n, Ordering::AcqRel);
}
let total_bytes = self.tail_bytes.fetch_add(bytes, Ordering::AcqRel) + bytes;
if self.triggers.tail_bytes != 0 && total_bytes >= self.triggers.tail_bytes {
self.rebuild()?;
return Ok(Some(RebuildReason::TailBytes(total_bytes)));
}
self.maybe_rebuild()
}
pub fn rebuild_due(&self) -> Option<RebuildReason> {
let hits = self.tail_hits.load(Ordering::Acquire);
if let Some(threshold) = self.triggers.miss_threshold(self.projection_len() as u64)
&& hits >= threshold
{
return Some(RebuildReason::TailHits(hits));
}
let bytes = self.tail_bytes.load(Ordering::Acquire);
if self.triggers.tail_bytes != 0 && bytes >= self.triggers.tail_bytes {
return Some(RebuildReason::TailBytes(bytes));
}
None
}
pub fn maybe_rebuild(&self) -> Result<Option<RebuildReason>> {
match self.rebuild_due() {
Some(reason) => {
self.rebuild()?;
Ok(Some(reason))
}
None => Ok(None),
}
}
pub fn rebuild(&self) -> Result<()> {
let entries = scan(&self.tail)?;
let total = entries.len() as u64;
let fresh = Arc::new(S::build(&entries)?);
*self
.projection
.write()
.map_err(|_| anyhow!("the projection lock is poisoned"))? = fresh;
self.total_rows.store(total, Ordering::Release);
self.tail_hits.store(0, Ordering::Release);
self.tail_bytes.store(0, Ordering::Release);
self.unabsorbed_size.store(0, Ordering::Release);
self.unabsorbed_rows.store(0, Ordering::Release);
self.tail_txns.store(0, Ordering::Release);
for slot in &self.unabsorbed_types {
slot.store(0, Ordering::Release);
}
self.rebuilds.fetch_add(1, Ordering::AcqRel);
self.generation.fetch_add(1, Ordering::AcqRel);
Ok(())
}
pub fn stats(&self) -> StackStats {
StackStats {
sealed_rows: self.projection_len() as u64,
total_rows: self.total_rows.load(Ordering::Acquire),
tail_hits: self.tail_hits.load(Ordering::Acquire),
absent: self.absent.load(Ordering::Acquire),
tail_bytes: self.tail_bytes.load(Ordering::Acquire),
rebuilds: self.rebuilds.load(Ordering::Acquire),
generation: self.generation.load(Ordering::Acquire),
unabsorbed_rows: self.unabsorbed_rows.load(Ordering::Acquire),
tail_txns: self.tail_txns.load(Ordering::Acquire),
}
}
pub fn projection_len(&self) -> usize {
self.projection.read().expect("projection lock").len()
}
pub fn oids_in_order(&self) -> Result<Vec<Vec<u8>>> {
let read = self.tail.begin_read()?;
let objects = read.open_table(OBJECTS)?;
let mut out = Vec::with_capacity(objects.len()? as usize);
for row in objects.iter()? {
let (k, _) = row?;
out.push(k.value().to_vec());
}
Ok(out)
}
pub fn extents_with_rows(&self, extents: &[(u64, u64)]) -> Result<Vec<bool>> {
let mut hit = vec![false; extents.len()];
let mut order: Vec<usize> = (0..extents.len()).filter(|&i| extents[i].1 > 0).collect();
order.sort_unstable_by_key(|&i| extents[i].0);
let starts: Vec<u64> = order.iter().map(|&i| extents[i].0).collect();
let mut wanted = order.len();
if wanted == 0 {
return Ok(hit);
}
let read = self.tail.begin_read()?;
let objects = read.open_table(OBJECTS)?;
for row in objects.iter()? {
let (_, v) = row?;
let offset = decode_row(v.value())?.offset;
let p = starts.partition_point(|&s| s <= offset);
if p == 0 {
continue;
}
let i = order[p - 1];
let (start, len) = extents[i];
if offset < start + len && !hit[i] {
hit[i] = true;
wanted -= 1;
if wanted == 0 {
break;
}
}
}
Ok(hit)
}
pub fn retain(&self, live: &dyn Fn(&[u8]) -> bool) -> Result<u64> {
let w = self.tail.begin_write()?;
let mut dropped = 0u64;
{
let mut objects = w.open_table(OBJECTS)?;
let dead: Vec<Vec<u8>> = objects
.iter()?
.filter_map(|row| row.ok())
.filter(|(k, _)| !live(k.value()))
.map(|(k, _)| k.value().to_vec())
.collect();
for oid in &dead {
objects.remove(oid.as_slice())?;
dropped += 1;
}
}
w.commit()?;
self.rebuild()?;
Ok(dropped)
}
fn fill_from_tail(&self, oids: &[&[u8]], out: &mut [Option<IndexRow>]) -> Result<()> {
self.tail_txns.fetch_add(1, Ordering::AcqRel);
let read = self.tail.begin_read()?;
let objects = read.open_table(OBJECTS)?;
let mut hits = 0u64;
let mut absent = 0u64;
for (slot, oid) in out.iter_mut().zip(oids) {
if slot.is_some() {
continue;
}
match objects.get(*oid)? {
Some(v) => {
*slot = Some(decode_row(v.value())?.as_index_row());
hits += 1;
}
None => absent += 1,
}
}
self.tail_hits.fetch_add(hits, Ordering::AcqRel);
self.absent.fetch_add(absent, Ordering::AcqRel);
Ok(())
}
fn absences_only(&self, out: &[Option<IndexRow>]) {
let n = out.iter().filter(|s| s.is_none()).count() as u64;
self.absent.fetch_add(n, Ordering::AcqRel);
}
}
fn scan(db: &Database) -> Result<Vec<IndexEntry>> {
let read = db.begin_read()?;
let objects = read.open_table(OBJECTS)?;
let mut out = Vec::with_capacity(objects.len()? as usize);
for row in objects.iter()? {
let (k, v) = row?;
out.push(decode_row(v.value())?.as_entry(k.value()));
}
Ok(out)
}
impl<S: ObjectIndex> ObjectIndex for ObjectReadStack<S> {
fn build(entries: &[IndexEntry]) -> Result<Self> {
let stack = Self::in_memory(RebuildTriggers::default())?;
stack.append(entries)?;
stack.rebuild()?;
Ok(stack)
}
fn lookup(&self, oid: &[u8]) -> Option<IndexRow> {
if let Some(row) = self.projection.read().expect("projection lock").lookup(oid) {
return Some(row);
}
if self.projection_is_complete() {
self.absent.fetch_add(1, Ordering::AcqRel);
return None;
}
let mut out = [None];
let _ = self.fill_from_tail(&[oid], &mut out);
out[0]
}
fn lookup_batch(&self, oids: &[&[u8]]) -> Vec<Option<IndexRow>> {
let mut out = self
.projection
.read()
.expect("projection lock")
.lookup_batch(oids);
if out.iter().any(Option::is_none) {
if self.projection_is_complete() {
self.absences_only(&out);
} else {
let _ = self.fill_from_tail(oids, &mut out);
}
}
out
}
fn ordinals_batch(&self, oids: &[&[u8]]) -> Vec<Option<u32>> {
let mut out = self
.projection
.read()
.expect("projection lock")
.ordinals_batch(oids);
if out.iter().any(Option::is_none) {
if self.projection_is_complete() {
let n = out.iter().filter(|s| s.is_none()).count() as u64;
self.absent.fetch_add(n, Ordering::AcqRel);
return out;
}
let mut rows: Vec<Option<IndexRow>> = out
.iter()
.map(|o| {
o.map(|ordinal| IndexRow {
ordinal,
offset: 0,
len: 0,
obj_type: ObjType::Blob,
uncompressed_size: 0,
delta_base: 0,
})
})
.collect();
let _ = self.fill_from_tail(oids, &mut rows);
for (slot, row) in out.iter_mut().zip(&rows) {
*slot = row.map(|r| r.ordinal);
}
}
out
}
fn extents_batch(&self, oids: &[&[u8]]) -> Vec<Option<(u64, u64)>> {
let mut out = self
.projection
.read()
.expect("projection lock")
.extents_batch(oids);
if out.iter().any(Option::is_none) {
if self.projection_is_complete() {
let n = out.iter().filter(|s| s.is_none()).count() as u64;
self.absent.fetch_add(n, Ordering::AcqRel);
return out;
}
let mut rows: Vec<Option<IndexRow>> = out
.iter()
.map(|e| {
e.map(|(offset, len)| IndexRow {
ordinal: 0,
offset,
len,
obj_type: ObjType::Blob,
uncompressed_size: 0,
delta_base: 0,
})
})
.collect();
let _ = self.fill_from_tail(oids, &mut rows);
for (slot, row) in out.iter_mut().zip(&rows) {
*slot = row.map(|r| (r.offset, r.len));
}
}
out
}
fn sum_uncompressed(&self) -> u64 {
self.projection
.read()
.expect("projection lock")
.sum_uncompressed()
+ self.unabsorbed_size.load(Ordering::Acquire)
}
fn count_type(&self, t: ObjType) -> usize {
self.projection
.read()
.expect("projection lock")
.count_type(t)
+ self.unabsorbed_types[t.code() as usize].load(Ordering::Acquire) as usize
}
fn name(&self) -> &'static str {
"ObjectReadStack"
}
fn len(&self) -> usize {
self.total_rows.load(Ordering::Acquire) as usize
}
fn ipc_bytes(&self) -> usize {
self.projection.read().expect("projection lock").ipc_bytes()
}
fn resident_bytes(&self) -> usize {
self.projection
.read()
.expect("projection lock")
.resident_bytes()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::arms::DEFAULT_REDB_CACHE_BYTES;
use crate::index_layout::{FourTables, synthetic_entries};
type Stack = ObjectReadStack<OneTableFourColumns>;
fn manual() -> Stack {
Stack::in_memory(RebuildTriggers::manual()).expect("in-memory stack")
}
fn split_stack(sealed: &[IndexEntry], tail: &[IndexEntry]) -> Stack {
let s = manual();
s.append(sealed).expect("seal half");
s.rebuild().expect("rebuild");
s.append(tail).expect("append tail");
s
}
#[test]
fn the_projection_is_incomplete_but_never_wrong() {
let all = synthetic_entries(500, 20, 0xC0FFEE);
let (sealed, tail) = all.split_at(300);
let stack = split_stack(sealed, tail);
assert_eq!(stack.projection_len(), 300, "the projection must be short");
assert_eq!(stack.len(), 500, "the stack answers for the repository");
assert!(
tail.iter().filter(|e| e.delta_base != 0).count() >= 10,
"the un-absorbed half carries no delta bases — the tail encoding would be untested"
);
let sealed_only = FourTables::build(sealed).expect("sealed-only index");
let complete = FourTables::build(&all).expect("complete index");
{
let proj = stack.projection.read().unwrap();
for e in tail {
assert!(
proj.lookup(&e.oid).is_none(),
"the projection already holds {} — this test would be vacuous",
hex::encode(&e.oid)
);
}
for e in sealed {
assert!(
proj.lookup(&e.oid).is_some(),
"the projection lost sealed object {}",
hex::encode(&e.oid)
);
}
}
let refs: Vec<&[u8]> = all.iter().map(|e| e.oid.as_slice()).collect();
let batched = stack.lookup_batch(&refs);
for (i, e) in all.iter().enumerate() {
let via_serial = stack.lookup(&e.oid).unwrap_or_else(|| {
panic!(
"the stack lost {} object {}, serial path",
if i < 300 { "sealed" } else { "tail" },
hex::encode(&e.oid)
)
});
let via_batch = batched[i].unwrap_or_else(|| {
panic!(
"the stack lost {} object {}, batch path",
if i < 300 { "sealed" } else { "tail" },
hex::encode(&e.oid)
)
});
assert_eq!(via_serial, via_batch, "serial and batch disagree");
assert_eq!(via_serial.offset, e.offset);
assert_eq!(via_serial.len, e.len);
assert_eq!(via_serial.obj_type, e.obj_type);
assert_eq!(via_serial.uncompressed_size, e.uncompressed_size);
assert_eq!(
via_serial.delta_base,
e.delta_base,
"the delta base of {} did not survive the {} path",
hex::encode(&e.oid),
if i < 300 { "projection" } else { "tail" }
);
let truth = complete.lookup(&e.oid).expect("complete index has it");
assert_eq!(via_serial.offset, truth.offset);
assert_eq!(via_serial.len, truth.len);
assert_eq!(via_serial.obj_type, truth.obj_type);
assert_eq!(via_serial.uncompressed_size, truth.uncompressed_size);
assert_eq!(via_serial.delta_base, truth.delta_base);
if i < 300 {
assert!(
!is_tail_row(&via_serial),
"sealed row wearing a tail ordinal"
);
assert_eq!(
via_serial.ordinal,
sealed_only.lookup(&e.oid).unwrap().ordinal,
"the projection's ordinal for {} is not its rank in the generation the \
projection covers",
hex::encode(&e.oid)
);
} else {
assert!(
is_tail_row(&via_serial),
"tail row {} has no tail ordinal",
hex::encode(&e.oid)
);
}
}
}
#[test]
fn only_a_tail_miss_is_absent_and_only_a_tail_hit_counts() {
let all = synthetic_entries(200, 20, 7);
let (sealed, tail) = all.split_at(120);
let stack = split_stack(sealed, tail);
let nowhere = synthetic_entries(64, 20, 0x00AB_5E47);
let mut refs: Vec<&[u8]> = Vec::new();
refs.extend(sealed.iter().map(|e| e.oid.as_slice()));
refs.extend(tail.iter().map(|e| e.oid.as_slice()));
refs.extend(nowhere.iter().map(|e| e.oid.as_slice()));
let rows = stack.lookup_batch(&refs);
assert_eq!(rows[..120].iter().filter(|r| r.is_some()).count(), 120);
assert_eq!(rows[120..200].iter().filter(|r| r.is_some()).count(), 80);
assert!(
rows[200..].iter().all(Option::is_none),
"an object in no repository resolved to a row"
);
let st = stack.stats();
assert_eq!(st.tail_hits, 80, "tail-served lookups miscounted");
assert_eq!(st.absent, 64, "genuine absences miscounted");
}
#[test]
fn an_append_only_violation_is_refused() {
let entries = synthetic_entries(16, 20, 11);
let stack = manual();
stack.append(&entries).expect("first append");
stack
.append(&entries)
.expect("identical re-append is idempotent");
assert_eq!(
stack.len(),
16,
"an idempotent re-append changed the row count"
);
let mut moved = entries[3].clone();
moved.offset += 1;
stack
.append(std::slice::from_ref(&moved))
.expect("a re-pushed pack at a new offset was refused");
let row = stack.lookup(&entries[3].oid).expect("still there");
assert_eq!(
row.offset, entries[3].offset,
"the second placement overwrote the first — the projection now names \
an offset that is not where the first writer put the bytes"
);
assert_eq!(stack.len(), 16, "a re-placement added a row");
let mut retyped = entries[4].clone();
retyped.uncompressed_size += 1;
let err = stack
.append(std::slice::from_ref(&retyped))
.expect_err("an oid that changed size was accepted");
assert!(
err.to_string().contains("identity violation"),
"wrong error: {err}"
);
let based = entries
.iter()
.find(|e| e.delta_base != 0)
.expect("the fixture must carry a delta");
let mut rebased = based.clone();
rebased.delta_base += 8;
stack
.append(std::slice::from_ref(&rebased))
.expect("a re-pushed pack that re-deltaed was refused");
assert_eq!(
stack.lookup(&based.oid).unwrap().delta_base,
based.delta_base,
"the second pack's base was adopted — the projection now names a base \
the FIRST writer never recorded, which is the thing the old guard \
was right to fear"
);
}
#[test]
fn each_trigger_fires_on_its_own_signal() {
let entries = synthetic_entries(64, 20, 21);
let bytes: u64 = entries.iter().map(|e| e.len).sum();
let stack = Stack::in_memory(RebuildTriggers {
tail_bytes: bytes,
..RebuildTriggers::manual()
})
.unwrap();
let g0 = stack.stats().generation;
let reason = stack.append(&entries).unwrap();
assert!(
matches!(reason, Some(RebuildReason::TailBytes(_))),
"the volume trigger did not fire at exactly its threshold: {reason:?}"
);
assert_eq!(stack.projection_len(), 64, "the rebuild absorbed nothing");
assert_eq!(stack.stats().generation, g0 + 1);
assert_eq!(
stack.stats().tail_bytes,
0,
"the byte counter was not reset"
);
let all = synthetic_entries(200, 20, 22);
let (sealed, tail) = all.split_at(100);
let stack = Stack::in_memory(RebuildTriggers {
tail_hits_per_row: 0.0,
min_tail_hits: 100,
tail_bytes: 0,
})
.unwrap();
stack.append(sealed).unwrap();
stack.rebuild().unwrap();
stack.append(tail).unwrap();
let g1 = stack.stats().generation;
assert_eq!(stack.projection_len(), 100);
let nowhere = synthetic_entries(1000, 20, 23);
let refs: Vec<&[u8]> = nowhere.iter().map(|e| e.oid.as_slice()).collect();
assert!(stack.lookup_batch(&refs).iter().all(Option::is_none));
assert_eq!(stack.stats().absent, 1000);
assert_eq!(
stack.rebuild_due(),
None,
"absences tripped the miss trigger"
);
stack.maybe_rebuild().unwrap();
assert_eq!(
stack.stats().generation,
g1,
"1000 absences must not rebuild"
);
assert_eq!(stack.projection_len(), 100);
let refs: Vec<&[u8]> = tail.iter().map(|e| e.oid.as_slice()).collect();
assert_eq!(stack.lookup_batch(&refs).iter().flatten().count(), 100);
assert_eq!(
stack.rebuild_due(),
Some(RebuildReason::TailHits(100)),
"100 tail-served lookups did not trip a threshold of 100"
);
assert!(stack.maybe_rebuild().unwrap().is_some());
assert_eq!(
stack.projection_len(),
200,
"the rebuild did not absorb the tail"
);
assert_eq!(stack.stats().generation, g1 + 1);
}
#[test]
fn the_column_scans_answer_for_the_repository_not_the_projection() {
let all = synthetic_entries(300, 20, 71);
let (sealed, tail) = all.split_at(180);
let stack = split_stack(sealed, tail);
let complete = OneTableFourColumns::build(&all).unwrap();
assert!(
stack.projection.read().unwrap().sum_uncompressed() < complete.sum_uncompressed(),
"the projection already sums to the whole repository — this test would be vacuous"
);
assert_eq!(
stack.sum_uncompressed(),
complete.sum_uncompressed(),
"sum_uncompressed answered for the projection, not the repository"
);
for t in ObjType::ALL {
assert_eq!(
stack.count_type(t),
complete.count_type(t),
"count_type({}) answered for the projection",
t.as_str()
);
}
let mut rewritten = all[0].clone();
rewritten.uncompressed_size += 1_000_000;
let batch = [all[7].clone(), rewritten];
assert!(stack.append(&batch).is_err());
assert_eq!(stack.sum_uncompressed(), complete.sum_uncompressed());
stack.rebuild().unwrap();
assert_eq!(stack.sum_uncompressed(), complete.sum_uncompressed());
for t in ObjType::ALL {
assert_eq!(stack.count_type(t), complete.count_type(t));
}
}
#[test]
fn the_narrow_batch_paths_also_fall_through() {
let all = synthetic_entries(300, 20, 81);
let (sealed, tail) = all.split_at(180);
let stack = split_stack(sealed, tail);
let refs: Vec<&[u8]> = all.iter().map(|e| e.oid.as_slice()).collect();
let ordinals = stack.ordinals_batch(&refs);
assert_eq!(
ordinals.iter().filter(|o| o.is_some()).count(),
300,
"the ordinal path lost {} of 300 objects",
300 - ordinals.iter().filter(|o| o.is_some()).count()
);
for (o, e) in ordinals[180..].iter().zip(tail) {
assert!(
o.unwrap() & TAIL_ORDINAL_BIT != 0,
"tail object {} got a projection ordinal",
hex::encode(&e.oid)
);
}
let extents = stack.extents_batch(&refs);
for (x, e) in extents.iter().zip(&all) {
assert_eq!(
*x,
Some((e.offset, e.len)),
"extent lost or wrong for {}",
hex::encode(&e.oid)
);
}
}
#[test]
fn a_rebuild_re_derives_ordinals_and_clears_the_tail_bit() {
let all = synthetic_entries(120, 32, 31);
let (sealed, tail) = all.split_at(60);
let stack = split_stack(sealed, tail);
let before = stack.lookup(&tail[0].oid).unwrap();
assert!(is_tail_row(&before));
stack.rebuild().unwrap();
let after = stack.lookup(&tail[0].oid).unwrap();
assert!(
!is_tail_row(&after),
"an absorbed row kept its tail ordinal"
);
assert_eq!(after.offset, before.offset, "absorption changed a fact");
assert_eq!(after.len, before.len);
assert_eq!(after.obj_type, before.obj_type);
assert_eq!(after.uncompressed_size, before.uncompressed_size);
assert_eq!(after.delta_base, before.delta_base);
let based = tail
.iter()
.find(|e| e.delta_base != 0)
.expect("the tail half must carry at least one delta");
let row = stack.lookup(&based.oid).unwrap();
assert_eq!(
row.delta_base,
based.delta_base,
"the rebuild moved {}'s delta base from {} to {}",
hex::encode(&based.oid),
based.delta_base,
row.delta_base
);
assert_ne!(
row.delta_base, row.ordinal as u64,
"a delta base that equals a row ordinal is the mistake §13 forbids"
);
let complete = OneTableFourColumns::build(&all).unwrap();
assert_eq!(
after.ordinal,
complete.lookup(&tail[0].oid).unwrap().ordinal,
"the rebuilt ordinal is not the lexicographic rank"
);
}
#[test]
fn the_crash_recovery_diff_is_exact_at_the_pack_boundary() {
let all = synthetic_entries(400, 20, 0xB17_5E7);
let bounds = [(0usize, 100usize), (100, 200), (200, 300), (300, 400)];
let extents: Vec<(u64, u64)> = bounds
.iter()
.map(|&(a, b)| {
let start = all[a].offset;
let end = all[b - 1].offset + all[b - 1].len;
(start, end - start)
})
.collect();
assert_eq!(
extents[1].0,
extents[0].0 + extents[0].1,
"the fixture's packs must be adjacent or the boundary is not under test"
);
let stack = manual();
for (i, &(a, b)) in bounds.iter().enumerate() {
if i != 1 {
stack.append(&all[a..b]).unwrap();
}
}
assert_eq!(stack.len(), 300, "the middle pack must be the missing one");
let hit = stack.extents_with_rows(&extents).unwrap();
assert_eq!(
hit,
vec![true, false, true, true],
"pack B has no rows in the tail and was reported absorbed"
);
assert_eq!(
stack.extents_with_rows(&extents[..2]).unwrap(),
vec![true, false],
"a row one byte past pack B's extent was counted as B's"
);
let past = all[399].offset + all[399].len;
assert_eq!(
stack.extents_with_rows(&[(0, 0), (past, 4096)]).unwrap(),
vec![false, false]
);
assert_eq!(
stack
.extents_with_rows(&[(all[7].offset, all[7].len)])
.unwrap(),
vec![true]
);
}
#[test]
fn a_file_backed_tail_reopens_warm() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("objects.tail.redb");
let all = synthetic_entries(90, 20, 41);
{
let stack =
ObjectReadStack::<OneTableFourColumns>::open(
&path,
RebuildTriggers::manual(),
DEFAULT_REDB_CACHE_BYTES,
)
.unwrap();
stack.append(&all[..40]).unwrap();
stack.rebuild().unwrap();
stack.append(&all[40..]).unwrap();
assert_eq!(stack.projection_len(), 40);
}
let stack =
ObjectReadStack::<OneTableFourColumns>::open(
&path,
RebuildTriggers::manual(),
DEFAULT_REDB_CACHE_BYTES,
).unwrap();
assert_eq!(
stack.projection_len(),
90,
"reopened projection covers {} of 90",
stack.projection_len()
);
for e in &all {
let row = stack
.lookup(&e.oid)
.unwrap_or_else(|| panic!("reopen lost {}", hex::encode(&e.oid)));
assert_eq!(row.offset, e.offset);
assert_eq!(row.uncompressed_size, e.uncompressed_size);
assert_eq!(row.delta_base, e.delta_base, "reopen lost a delta base");
}
assert!(
all.iter().filter(|e| e.delta_base != 0).count() >= 5,
"the fixture must carry delta bases across the reopen"
);
}
#[test]
fn a_packed_tail_row_round_trips_and_refuses_nonsense() {
let mut e = synthetic_entries(1, 20, 51)[0].clone();
e.delta_base = 4242;
let e = &e;
let packed = encode_row(9, e);
let back = decode_row(&packed).unwrap();
assert_eq!(back.seq, 9);
assert_eq!(back.offset, e.offset);
assert_eq!(back.len, e.len);
assert_eq!(back.obj_type, e.obj_type);
assert_eq!(back.uncompressed_size, e.uncompressed_size);
assert_eq!(
back.delta_base, e.delta_base,
"the delta base did not survive the tail row"
);
assert_ne!(
back.delta_base, back.uncompressed_size,
"the fixture must not let the two u64 fields alias"
);
assert!(decode_row(&packed[..TAIL_ROW_BYTES - 1]).is_err());
assert!(
decode_row(&packed[..33]).is_err(),
"a row in the old 33-byte encoding must be refused, not reinterpreted"
);
let mut bad = packed;
bad[24] = 5; assert!(
decode_row(&bad).is_err(),
"an unassigned object type code was accepted"
);
}
#[test]
fn a_fully_absorbed_stack_agrees_with_both_arrow_arms() {
let entries = synthetic_entries(400, 20, 61);
let absent = synthetic_entries(100, 20, 62);
let stack = Stack::build(&entries).unwrap();
let a = FourTables::build(&entries).unwrap();
let b = OneTableFourColumns::build(&entries).unwrap();
assert_eq!(stack.len(), 400);
assert_eq!(stack.projection_len(), 400, "build left rows in the tail");
let mut refs: Vec<&[u8]> = entries.iter().map(|e| e.oid.as_slice()).collect();
refs.extend(absent.iter().map(|e| e.oid.as_slice()));
let rs = stack.lookup_batch(&refs);
assert_eq!(
rs,
a.lookup_batch(&refs),
"the stack disagrees with FourTables"
);
assert_eq!(
rs,
b.lookup_batch(&refs),
"the stack disagrees with OneTableFourColumns"
);
assert_eq!(
stack.stats().tail_hits,
0,
"a fully absorbed stack still went to the tail for a row"
);
assert_eq!(stack.stats().absent, 100);
}
#[test]
fn a_complete_projection_answers_absent_without_asking_the_tail() {
let entries = synthetic_entries(200, 20, 91);
let nowhere = synthetic_entries(500, 20, 92);
let stack = Stack::build(&entries).unwrap();
assert!(stack.projection_is_complete());
let refs: Vec<&[u8]> = nowhere.iter().map(|e| e.oid.as_slice()).collect();
assert!(stack.lookup_batch(&refs).iter().all(Option::is_none));
assert!(stack.lookup(&nowhere[0].oid).is_none());
assert!(stack.ordinals_batch(&refs).iter().all(Option::is_none));
assert!(stack.extents_batch(&refs).iter().all(Option::is_none));
let st = stack.stats();
assert_eq!(st.unabsorbed_rows, 0);
assert_eq!(
st.tail_txns, 0,
"a complete projection opened {} redb transactions to say 'no'",
st.tail_txns
);
assert_eq!(st.absent, 500 + 1 + 500 + 500, "absences went uncounted");
let more = synthetic_entries(1, 20, 93);
stack.append(&more).unwrap();
assert!(!stack.projection_is_complete());
assert!(stack.lookup_batch(&refs).iter().all(Option::is_none));
assert!(
stack.stats().tail_txns > 0,
"an incomplete projection did not consult the tail"
);
assert!(
stack.lookup(&more[0].oid).is_some(),
"the appended row is not reachable"
);
}
}