use std::cmp::Ordering;
use std::ops::Bound;
use std::sync::Arc;
use std::time::Duration;
use pigeonhole::{
Cell, Compaction, Condition, Durability, ErrorCode, Family, Options, Pigeonhole, Priority, Row,
Table, Value, ValueFilter, days,
};
use pigeonhole_io::sim::SimVfs;
fn db() -> Pigeonhole {
let vfs = SimVfs::new(7);
Pigeonhole::open("/db/api.phdb", sim_options(&vfs)).expect("open")
}
fn sim_options(vfs: &Arc<SimVfs>) -> Options {
Options::default()
.vfs(Arc::clone(vfs) as _)
.shards(2)
.memtable_budget(4 << 20)
.wal_segment_size(256 << 10)
}
fn table(db: &Pigeonhole) -> Table {
db.table("t")
.unwrap()
.family("a", Family::default())
.family("b", Family::default().max_versions(2))
.create_if_missing()
.unwrap()
}
fn cells(t: &Table, row: &[u8]) -> Vec<(String, Vec<u8>, Vec<u8>)> {
t.row(row)
.versions(0)
.read()
.unwrap()
.map(|r| {
r.iter()
.map(|e| {
(
e.family.to_owned(),
e.qualifier.to_vec(),
e.cell.value().to_vec(),
)
})
.collect()
})
.unwrap_or_default()
}
fn value(t: &Table, row: &[u8], family: &str, q: &[u8]) -> Option<Vec<u8>> {
t.get(row, family, q).unwrap().map(|c| c.value().to_vec())
}
struct TempDir(std::path::PathBuf);
impl TempDir {
fn new(name: &str) -> Self {
let dir =
std::env::temp_dir().join(format!("pigeonhole-api-{}-{name}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
TempDir(dir)
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = std::fs::remove_dir_all(&self.0);
}
}
#[test]
fn table_builder_create_open_and_families() {
let db = db();
assert_eq!(
db.table("t").unwrap().open().unwrap_err().code(),
ErrorCode::TableNotFound
);
let t = db
.table("t")
.unwrap()
.family("a", Family::default().max_versions(1))
.family("a", Family::default())
.create()
.unwrap();
assert_eq!(t.name(), "t");
assert_eq!(t.families(), ["a"]);
assert_eq!(
db.table("t").unwrap().create().unwrap_err().code(),
ErrorCode::TableExists
);
let t2 = db
.table("t")
.unwrap()
.family("a", Family::default().max_versions(0))
.family("b", Family::default())
.create_if_missing()
.unwrap();
assert_eq!(t2.families(), ["a", "b"]);
for ts in 1..=3 {
t2.mutate(b"r")
.put_at("a", b"q", ts, b"v")
.commit()
.unwrap();
}
assert_eq!(cells(&t2, b"r").len(), 1, "max_versions(1) kept");
t.mutate(b"r").put("b", b"q", b"x").commit().unwrap();
assert_eq!(value(&t, b"r", "b", b"q").as_deref(), Some(&b"x"[..]));
assert_eq!(db.tables(), ["t"]);
}
#[test]
fn drop_table_removes_it_and_its_data() {
let db = db();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
db.drop_table("t").unwrap();
assert!(db.tables().is_empty());
assert_eq!(
db.drop_table("t").unwrap_err().code(),
ErrorCode::TableNotFound
);
let t = table(&db);
assert_eq!(value(&t, b"r", "a", b"q"), None);
}
#[test]
fn family_settings_are_stored_and_unregistered_operators_refused() {
let db = db();
let err = db
.table("t")
.unwrap()
.family("f", Family::default().merge_operator("app.append"))
.create_if_missing()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::UnknownMergeOperator);
let t = db
.table("t")
.unwrap()
.family(
"f",
Family::default()
.max_versions(3)
.ttl(days(30))
.bloom_bits(0)
.blob_threshold(1 << 20)
.uncompressed()
.lz4()
.block_size(4096)
.cache_priority(Priority::High)
.compaction(Compaction::Leveled),
)
.create()
.unwrap();
assert_eq!(t.families(), ["f"]);
let u = db
.table("u")
.unwrap()
.family("tiered", Family::default().compaction(Compaction::Tiered))
.family(
"fifo",
Family::default()
.ttl(days(1))
.compaction(Compaction::FifoByTime),
)
.family("zstd", Family::default().zstd(19))
.create()
.unwrap();
assert_eq!(u.families(), ["tiered", "fifo", "zstd"]);
let value: Vec<u8> = (0..3000u32).map(|i| (i % 7) as u8).collect();
for row in 0..50u32 {
u.mutate(format!("r{row:03}").as_bytes())
.put("zstd", b"q", &value)
.commit()
.unwrap();
}
db.flush().unwrap();
db.compact().unwrap();
for row in 0..50u32 {
let cell = u
.get(format!("r{row:03}").as_bytes(), "zstd", b"q")
.unwrap();
assert_eq!(cell.map(|c| c.value().to_vec()), Some(value.clone()));
}
}
#[test]
fn builder_errors_surface_at_the_terminal_call() {
let db = db();
let t = table(&db);
let code = |r: pigeonhole::Result<_>| r.map(|_: ()| ()).unwrap_err().code();
let err = t
.mutate(b"r")
.put("nope", b"q", b"v")
.put("a", &[0u8; 70_000], b"v")
.commit()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::FamilyNotFound);
assert!(err.message().contains("nope"), "{err}");
assert_eq!(err.to_string(), err.message());
assert_eq!(
code(
t.mutate(b"r")
.put("a", &[0u8; 70_000], b"v")
.commit()
.map(|_| ())
),
ErrorCode::KeyTooLarge
);
assert_eq!(
code(t.mutate(&[0u8; 70_000]).delete_row().commit().map(|_| ())),
ErrorCode::KeyTooLarge
);
assert_eq!(
code(t.get(b"r", "nope", b"q").map(|_| ())),
ErrorCode::FamilyNotFound
);
assert_eq!(
code(t.row(b"r").family("a").family("nope").read().map(|_| ())),
ErrorCode::FamilyNotFound
);
assert_eq!(
code(t.scan_prefix(b"").families(["nope"]).iter().map(|_| ())),
ErrorCode::FamilyNotFound
);
let mut wb = db.write_batch();
wb.put(&t, b"r", "nope", b"q", b"v")
.put(&t, b"r", "a", b"q", b"v");
assert_eq!(code(wb.commit().map(|_| ())), ErrorCode::FamilyNotFound);
assert_eq!(cells(&t, b"r"), []);
}
#[test]
fn error_codes_are_stable_numbers() {
let codes = [
(ErrorCode::Io, 1),
(ErrorCode::Corruption, 2),
(ErrorCode::WriterLocked, 3),
(ErrorCode::ShmVersionMismatch, 4),
(ErrorCode::ShmUnavailable, 5),
(ErrorCode::UnsupportedFormat, 6),
(ErrorCode::NetworkFilesystem, 7),
(ErrorCode::TableNotFound, 8),
(ErrorCode::TableExists, 9),
(ErrorCode::FamilyNotFound, 10),
(ErrorCode::FamilyExists, 11),
(ErrorCode::UnknownMergeOperator, 12),
(ErrorCode::MergeFailed, 13),
(ErrorCode::Conflict, 14),
(ErrorCode::ReadOnly, 15),
(ErrorCode::KeyTooLarge, 16),
(ErrorCode::ValueTooLarge, 17),
(ErrorCode::NoSpace, 18),
(ErrorCode::InvalidArgument, 19),
(ErrorCode::Unsupported, 20),
(ErrorCode::Closed, 21),
(ErrorCode::NoReaderSlot, 22),
(ErrorCode::RecordTooLarge, 23),
(ErrorCode::Busy, 24),
(ErrorCode::SnapshotExpired, 25),
(ErrorCode::WouldDeadlock, 26),
(ErrorCode::BatchTooLarge, 27),
];
for (code, n) in codes {
assert_eq!(code as u32, n, "{code:?}");
}
}
#[test]
fn every_engine_error_maps_to_its_code() {
use pigeonhole_engine::Error as E;
let io = || pigeonhole_io::Error::new(pigeonhole_io::ErrorKind::Other, "x");
let merge = pigeonhole_engine::MergeError {
operator: "op".into(),
message: "bad".into(),
};
let cases: Vec<(E, ErrorCode)> = vec![
(E::Io(io()), ErrorCode::Io),
(E::Corruption("x".into()), ErrorCode::Corruption),
(E::WriterLocked, ErrorCode::WriterLocked),
(
E::ShmVersionMismatch {
found: 1,
expected: 2,
},
ErrorCode::ShmVersionMismatch,
),
(E::ShmUnavailable, ErrorCode::ShmUnavailable),
(E::UnsupportedFormat(9), ErrorCode::UnsupportedFormat),
(E::NetworkFilesystem, ErrorCode::NetworkFilesystem),
(E::TableNotFound("t".into()), ErrorCode::TableNotFound),
(E::TableExists("t".into()), ErrorCode::TableExists),
(E::FamilyNotFound("f".into()), ErrorCode::FamilyNotFound),
(E::FamilyExists("f".into()), ErrorCode::FamilyExists),
(
E::UnknownMergeOperator("op".into()),
ErrorCode::UnknownMergeOperator,
),
(E::Merge(merge), ErrorCode::MergeFailed),
(E::Conflict, ErrorCode::Conflict),
(E::ReadOnly, ErrorCode::ReadOnly),
(E::KeyTooLarge, ErrorCode::KeyTooLarge),
(E::ValueTooLarge, ErrorCode::ValueTooLarge),
(E::NoSpace, ErrorCode::NoSpace),
(E::InvalidArgument("x".into()), ErrorCode::InvalidArgument),
(E::Unsupported("x"), ErrorCode::Unsupported),
(E::Closed, ErrorCode::Closed),
(E::NoReaderSlot, ErrorCode::NoReaderSlot),
(E::RecordTooLarge, ErrorCode::RecordTooLarge),
(E::Busy, ErrorCode::Busy),
(E::BatchTooLarge, ErrorCode::BatchTooLarge),
(E::SnapshotExpired, ErrorCode::SnapshotExpired),
(E::WouldDeadlock, ErrorCode::WouldDeadlock),
];
for (e, code) in cases {
let what = format!("{e:?}");
let message = match &e {
E::Merge(_) => "merge operator \"op\" failed: bad".to_owned(),
E::Busy => "stalled: a flush or compaction did not free room within the write-stall \
timeout (30 s); back off and retry, and drop old snapshots"
.to_owned(),
E::BatchTooLarge => "the batch can never fit a shard's memtable arena (about half \
of Options::memtable_budget); split it or raise the budget"
.to_owned(),
_ => e.to_string(),
};
let public: pigeonhole::Error = e.into();
assert_eq!(public.code(), code, "{what}");
assert_eq!(public.message(), message);
}
}
#[test]
fn row_reads_order_families_by_creation_or_request() {
let db = db();
let t = table(&db);
t.mutate(b"r")
.put("b", b"y", b"1")
.put("a", b"z", b"2")
.put("a", b"x", b"3")
.commit()
.unwrap();
let order = |r: Option<pigeonhole::RowRef<'_>>| -> Vec<(String, Vec<u8>)> {
r.unwrap()
.iter()
.map(|e| (e.family.to_owned(), e.qualifier.to_vec()))
.collect()
};
let s = |f: &str, q: &[u8]| (f.to_owned(), q.to_vec());
assert_eq!(
order(t.row(b"r").read().unwrap()),
[s("a", b"x"), s("a", b"z"), s("b", b"y")]
);
assert_eq!(
order(t.row(b"r").families(["b", "a", "b"]).read().unwrap()),
[s("b", b"y"), s("a", b"x"), s("a", b"z")]
);
assert_eq!(
order(t.row(b"r").family("b").read().unwrap()),
[s("b", b"y")]
);
assert!(t.row(b"missing").read().unwrap().is_none());
assert!(t.row(b"r").qualifier_prefix(b"q").read().unwrap().is_none());
}
#[test]
fn qualifier_version_time_and_column_filters() {
let db = db();
let t = table(&db);
for (q, ts) in [
(&b"q1"[..], 10u64),
(b"q1", 20),
(b"q1", 30),
(b"q2", 40),
(b"q3", 50),
] {
t.mutate(b"r")
.put_at("a", q, ts, &ts.to_be_bytes())
.commit()
.unwrap();
}
let read = |r: pigeonhole::RowRead<'_>| -> Vec<(Vec<u8>, u64)> {
r.read()
.unwrap()
.map(|r| {
r.iter()
.map(|e| (e.qualifier.to_vec(), e.cell.timestamp()))
.collect()
})
.unwrap_or_default()
};
let v = |q: &[u8], ts: u64| (q.to_vec(), ts);
assert_eq!(
read(t.row(b"r")),
[v(b"q1", 30), v(b"q2", 40), v(b"q3", 50)]
);
assert_eq!(
read(t.row(b"r").versions(2)),
[v(b"q1", 30), v(b"q1", 20), v(b"q2", 40), v(b"q3", 50)]
);
assert_eq!(
read(t.row(b"r").versions(0).time_range(15..45)),
[v(b"q1", 30), v(b"q1", 20), v(b"q2", 40)]
);
assert_eq!(
read(t.row(b"r").versions(0).latest()),
[v(b"q1", 30), v(b"q2", 40), v(b"q3", 50)]
);
assert_eq!(read(t.row(b"r").qualifier_prefix(b"q2")), [v(b"q2", 40)]);
assert_eq!(
read(t.row(b"r").qualifier_range(&b"q2"[..]..)),
[v(b"q2", 40), v(b"q3", 50)]
);
assert_eq!(
read(
t.row(b"r")
.qualifier_bounds(Bound::Excluded(b"q1"), Bound::Included(b"q2"))
),
[v(b"q2", 40)]
);
assert_eq!(
read(t.row(b"r").column_limit(2)),
[v(b"q1", 30), v(b"q2", 40)]
);
assert_eq!(
read(
t.row(b"r")
.value_filter(ValueFilter::Equals(40u64.to_be_bytes().to_vec()))
),
[v(b"q2", 40)]
);
assert_eq!(
read(
t.row(b"r")
.value_filter(ValueFilter::Equals(20u64.to_be_bytes().to_vec()))
),
[]
);
assert_eq!(
read(t.row(b"r").value_filter(ValueFilter::Prefix(vec![0; 7]))),
[v(b"q1", 30), v(b"q2", 40), v(b"q3", 50)]
);
}
#[test]
fn scans_bounds_prefixes_limits_and_both_iterator_forms() {
let db = db();
let t = table(&db);
let keys: [&[u8]; 6] = [b"a", b"user:1", b"user:2", b"user:3", b"user;", b"\xff\xff"];
for k in keys {
t.mutate(k)
.put("a", b"q", k)
.put("b", b"n", b"x")
.commit()
.unwrap();
}
let owned = |s: pigeonhole::Scan<'_>| -> Vec<Vec<u8>> {
s.iter()
.unwrap()
.map(|r| r.unwrap().key().to_vec())
.collect()
};
let lent = |s: pigeonhole::Scan<'_>| -> Vec<Vec<u8>> {
let mut it = s.iter().unwrap();
let mut out = Vec::new();
while let Some(r) = it.next_ref().unwrap() {
out.push(r.key().to_vec());
}
assert!(it.next_ref().unwrap().is_none());
assert!(it.next().is_none());
out
};
let k = |ks: &[&[u8]]| ks.iter().map(|k| k.to_vec()).collect::<Vec<_>>();
assert_eq!(owned(t.scan_prefix(b"user:")), k(&keys[1..4]));
assert_eq!(lent(t.scan_prefix(b"user:")), k(&keys[1..4]));
assert_eq!(owned(t.scan_prefix(b"")), k(&keys));
assert_eq!(owned(t.scan_prefix(b"\xff")), k(&keys[5..]));
assert_eq!(owned(t.scan(&b"user:2"[..]..&b"user;"[..])), k(&keys[2..4]));
assert_eq!(
owned(t.scan(&b"user:2"[..]..=&b"user;"[..])),
k(&keys[2..5])
);
assert_eq!(owned(t.scan(b"user:2"..)), k(&keys[2..]));
assert_eq!(owned(t.scan::<[u8]>(..)), k(&keys));
assert_eq!(
lent(
t.scan_bounds(Bound::Excluded(b"user:1"), Bound::Unbounded)
.limit(2)
),
k(&keys[2..4])
);
assert_eq!(owned(t.scan_prefix(b"").limit(0)), k(&[]));
assert_eq!(owned(t.scan(&b"z"[..]..&b"a"[..])), k(&[]));
let row = t
.scan_prefix(b"a")
.family("b")
.iter()
.unwrap()
.next()
.unwrap()
.unwrap();
assert_eq!(row.len(), 1);
assert_eq!(row.entry(0).unwrap().0, "b");
let row = t
.scan_prefix(b"a")
.columns_per_row(1)
.iter()
.unwrap()
.next()
.unwrap()
.unwrap();
assert_eq!(row.len(), 2, "one column per family");
assert_eq!(
owned(
t.scan_prefix(b"")
.family("a")
.value_filter(ValueFilter::Prefix(b"user:".to_vec()))
),
k(&keys[1..4])
);
}
#[test]
fn row_and_rowref_agree_and_cells_outlive_their_handles() {
let db = db();
let t = table(&db);
let big = vec![9u8; 4000];
t.mutate(b"r")
.put("a", b"big", &big)
.put_i64("a", b"n", -5)
.put_f64("b", b"f", 2.5)
.commit()
.unwrap();
let r = t.row(b"r").read().unwrap().unwrap();
let row: Row = r.to_owned();
let view = row.view();
assert_eq!(r.len(), 3);
assert_eq!(row.len(), 3);
assert!(!row.is_empty() && !view.is_empty());
for i in 0..3 {
let (f, q, c) = row.entry(i).unwrap();
let e = view.entry(i).unwrap();
let e2 = r.entry(i).unwrap();
assert_eq!((f, q, c.value()), (e.family, e.qualifier, e.cell.value()));
assert_eq!((e.family, e.qualifier), (e2.family, e2.qualifier));
assert_eq!(c.timestamp(), e2.cell.timestamp());
}
assert!(row.entry(3).is_none() && r.entry(3).is_none());
assert_eq!(row.get("a", b"n").unwrap().as_i64(), Some(-5));
assert_eq!(view.get("b", b"f").unwrap().typed(), Value::F64(2.5));
assert_eq!(r.get("a", b"big").unwrap().value(), &big[..]);
assert!(row.get("a", b"none").is_none() && row.get("zz", b"n").is_none());
assert_eq!(r.get("b", b"f").unwrap().as_i64(), None);
let cell: Cell = t.get(b"r", "a", b"big").unwrap().unwrap().to_owned();
drop(r);
drop(t);
let t = table(&db);
t.mutate(b"r").delete_row().commit().unwrap();
drop(db);
assert_eq!(cell.value(), &big[..]);
assert_eq!(cell.typed(), Value::Bytes(&big));
assert_eq!(row.get("a", b"big").unwrap().value(), &big[..]);
let moved = std::thread::spawn(move || cell.value().len())
.join()
.unwrap();
assert_eq!(moved, 4000);
}
#[test]
fn snapshots_isolate_reads() {
let db = db();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"1").commit().unwrap();
let snap = db.snapshot().unwrap();
t.mutate(b"r").put("a", b"q", b"2").commit().unwrap();
t.mutate(b"s").put("a", b"q", b"3").commit().unwrap();
assert!(db.snapshot().unwrap().seqno() > snap.seqno());
assert_eq!(
t.get_at(&snap, b"r", "a", b"q").unwrap().unwrap().value(),
b"1"
);
let row = t.row(b"r").snapshot(&snap).read().unwrap().unwrap();
assert_eq!(row.get("a", b"q").unwrap().value(), b"1");
assert_eq!(
t.scan_prefix(b"").snapshot(&snap).iter().unwrap().count(),
1
);
assert_eq!(t.scan_prefix(b"").iter().unwrap().count(), 2);
let clone = snap.clone();
drop(snap);
assert!(t.get_at(&clone, b"s", "a", b"q").unwrap().is_none());
}
#[test]
fn deletes_follow_the_timestamp_rules() {
let db = db();
let t = table(&db);
for ts in [10u64, 20, 30] {
t.mutate(b"r")
.put_at("a", b"q", ts, &[ts as u8])
.commit()
.unwrap();
}
t.mutate(b"r").delete_cell("a", b"q", 30).commit().unwrap();
assert_eq!(value(&t, b"r", "a", b"q"), Some(vec![20]));
t.mutate(b"r")
.put_at("a", b"q", 30, b"again")
.commit()
.unwrap();
assert_eq!(value(&t, b"r", "a", b"q"), Some(vec![20]));
t.mutate(b"r").delete_column("a", b"q").commit().unwrap();
assert_eq!(value(&t, b"r", "a", b"q"), None);
t.mutate(b"r")
.put_at("a", b"q", 40, b"old")
.commit()
.unwrap();
assert_eq!(value(&t, b"r", "a", b"q"), None);
t.mutate(b"r").put("a", b"q", b"new").commit().unwrap();
assert_eq!(value(&t, b"r", "a", b"q").as_deref(), Some(&b"new"[..]));
t.mutate(b"r")
.put("b", b"x", b"1")
.put("a", b"y", b"2")
.commit()
.unwrap();
t.mutate(b"r").delete_family("a").commit().unwrap();
assert_eq!(
cells(&t, b"r"),
[("b".into(), b"x".to_vec(), b"1".to_vec())]
);
t.mutate(b"r").delete_row().commit().unwrap();
assert_eq!(cells(&t, b"r"), []);
t.mutate(b"s")
.put("a", b"q", b"1")
.put("a", b"q", b"2")
.commit()
.unwrap();
assert_eq!(value(&t, b"s", "a", b"q").as_deref(), Some(&b"2"[..]));
}
#[test]
fn a_compaction_purge_uncovers_later_puts_at_older_timestamps() {
let db = db();
let t = table(&db);
t.mutate(b"c").put("a", b"q", b"v").commit().unwrap();
t.mutate(b"c").delete_column("a", b"q").commit().unwrap();
t.mutate(b"f").put("a", b"q", b"v").commit().unwrap();
t.mutate(b"f").delete_family("a").commit().unwrap();
db.flush().unwrap();
for row in [&b"c"[..], b"f"] {
t.mutate(row)
.put_at("a", b"q", 1, b"before")
.commit()
.unwrap();
assert_eq!(value(&t, row, "a", b"q"), None, "hidden by the delete");
}
db.compact().unwrap();
for row in [&b"c"[..], b"f"] {
assert_eq!(value(&t, row, "a", b"q"), None);
t.mutate(row)
.put_at("a", b"q", 2, b"after")
.commit()
.unwrap();
assert_eq!(value(&t, row, "a", b"q").as_deref(), Some(&b"after"[..]));
}
}
fn counters(db: &Pigeonhole) -> Table {
db.table("counters")
.unwrap()
.family("c", Family::counter())
.family("a", Family::default())
.create_if_missing()
.unwrap()
}
fn count(t: &Table, row: &[u8], q: &[u8]) -> Option<i64> {
t.get(row, "c", q).unwrap().map(|c| c.as_i64().unwrap())
}
fn buckets(t: &Table, row: &[u8], q: &[u8]) -> Vec<(u64, i64)> {
t.row(row)
.family("c")
.qualifier_range(q..=q)
.versions(0)
.read()
.unwrap()
.map(|r| {
r.iter()
.map(|e| (e.cell.timestamp(), e.cell.as_i64().unwrap()))
.collect()
})
.unwrap_or_default()
}
#[test]
fn counters_and_typed_values() {
let db = db();
let t = counters(&db);
for d in [5, -2, 10] {
t.mutate(b"r").incr("c", b"hits", d).commit().unwrap();
}
assert_eq!(count(&t, b"r", b"hits"), Some(13));
assert_eq!(buckets(&t, b"r", b"hits"), [(0, 13)]);
t.mutate(b"r")
.put_i64("c", b"hits", 100)
.incr("c", b"hits", 7)
.commit()
.unwrap();
assert_eq!(count(&t, b"r", b"hits"), Some(107));
t.mutate(b"r")
.incr("c", b"hits", 1)
.incr("c", b"hits", 2)
.commit()
.unwrap();
assert_eq!(count(&t, b"r", b"hits"), Some(110));
t.mutate(b"r")
.incr("c", b"hits", 50)
.put_i64("c", b"hits", 100)
.commit()
.unwrap();
assert_eq!(count(&t, b"r", b"hits"), Some(100));
t.mutate(b"r").incr("c", b"hits", 1).commit().unwrap();
assert_eq!(count(&t, b"r", b"hits"), Some(101));
let mut wb = db.write_batch();
wb.incr(&t, b"r", "c", b"hits", 1)
.incr(&t, b"s", "c", b"hits", 1)
.incr(&t, b"s", "c", b"hits", 1)
.incr_at(&t, b"s", "c", b"hits", 9, 4)
.incr_at(&t, b"s", "c", b"hits", 9, 5);
wb.commit().unwrap();
assert_eq!(count(&t, b"r", b"hits"), Some(102));
assert_eq!(buckets(&t, b"s", b"hits"), [(9, 9), (0, 2)]);
t.mutate(b"w")
.put_i64("c", b"n", i64::MAX)
.commit()
.unwrap();
t.mutate(b"w").incr("c", b"n", 1).commit().unwrap();
assert_eq!(count(&t, b"w", b"n"), Some(i64::MIN));
}
#[test]
fn counter_buckets_are_versions() {
let db = db();
let t = counters(&db);
let hour = 3_600_000_000;
for (h, d) in [(7, 1), (8, 4), (8, 1), (9, 2)] {
t.mutate(b"r")
.incr_at("c", b"hourly", h * hour, d)
.commit()
.unwrap();
}
let mut wb = db.write_batch();
wb.incr_at(&t, b"r", "c", b"hourly", 9 * hour, 3);
wb.commit().unwrap();
assert_eq!(
buckets(&t, b"r", b"hourly"),
[(9 * hour, 5), (8 * hour, 5), (7 * hour, 1)]
);
t.mutate(b"r")
.put_i64_at("c", b"hourly", 8 * hour, 0)
.commit()
.unwrap();
t.mutate(b"r")
.incr_at("c", b"hourly", 8 * hour, 2)
.commit()
.unwrap();
let mut wb = db.write_batch();
wb.put_i64_at(&t, b"r", "c", b"hourly", 7 * hour, 50);
wb.commit().unwrap();
assert_eq!(
buckets(&t, b"r", b"hourly"),
[(9 * hour, 5), (8 * hour, 2), (7 * hour, 50)]
);
let window: Vec<i64> = t
.row(b"r")
.family("c")
.versions(0)
.time_range(8 * hour..10 * hour)
.read()
.unwrap()
.unwrap()
.iter()
.map(|e| e.cell.as_i64().unwrap())
.collect();
assert_eq!(window, [5, 2]);
t.mutate(b"r")
.delete_cell("c", b"hourly", 9 * hour)
.commit()
.unwrap();
db.flush().unwrap();
db.compact().unwrap();
assert_eq!(
buckets(&t, b"r", b"hourly"),
[(8 * hour, 2), (7 * hour, 50)]
);
let capped = db
.table("capped")
.unwrap()
.family("c", Family::counter().max_versions(2))
.create()
.unwrap();
for h in 1..=4 {
capped
.mutate(b"r")
.incr_at("c", b"n", h * hour, 1)
.commit()
.unwrap();
}
assert_eq!(buckets(&capped, b"r", b"n"), [(4 * hour, 1), (3 * hour, 1)]);
}
#[test]
fn counter_deletes_hide_only_what_came_before() {
let db = db();
let t = counters(&db);
t.mutate(b"r").incr("c", b"n", 5).commit().unwrap();
t.mutate(b"r").delete_column("c", b"n").commit().unwrap();
assert_eq!(count(&t, b"r", b"n"), None);
t.mutate(b"r").incr("c", b"n", 2).commit().unwrap();
assert_eq!(count(&t, b"r", b"n"), Some(2));
let before = db.snapshot().unwrap();
t.mutate(b"r").delete_row().commit().unwrap();
t.mutate(b"r").incr("c", b"n", 1).commit().unwrap();
assert_eq!(count(&t, b"r", b"n"), Some(1));
db.flush().unwrap();
db.compact().unwrap();
assert_eq!(count(&t, b"r", b"n"), Some(1));
assert_eq!(
t.get_at(&before, b"r", "c", b"n")
.unwrap()
.unwrap()
.as_i64(),
Some(2)
);
t.mutate(b"r").delete_family("c").commit().unwrap();
t.mutate(b"r").delete_cell("c", b"m", 0).commit().unwrap();
t.mutate(b"r").incr("c", b"m", 4).commit().unwrap();
assert_eq!(count(&t, b"r", b"n"), None);
assert_eq!(count(&t, b"r", b"m"), Some(4));
}
#[test]
fn counter_family_write_rules() {
let db = db();
let t = counters(&db);
let refused = |r: pigeonhole::Result<pigeonhole::CommitInfo>, what: &str| {
let err = r.unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidArgument, "{err}");
assert!(err.message().contains(what), "{err}");
};
refused(
t.mutate(b"r").incr("a", b"n", 1).commit(),
"family \"a\" of table \"counters\" has no merge operator",
);
refused(
t.mutate(b"r").incr_at("a", b"n", 5, 1).commit(),
"not a counter family",
);
for m in [
t.mutate(b"r").put("c", b"n", b"text"),
t.mutate(b"r").put_at("c", b"n", 5, b"text"),
t.mutate(b"r").put_f64("c", b"n", 1.5),
t.mutate(b"r").merge("c", b"n", &1i64.to_le_bytes()),
] {
refused(m.commit(), "only i64 values");
}
let mut wb = db.write_batch();
wb.put(&t, b"r", "c", b"n", b"text");
refused(wb.commit(), "only i64 values");
let ttl = db
.table("ttl")
.unwrap()
.family("c", Family::counter().ttl(days(1)))
.create()
.unwrap();
refused(ttl.mutate(b"r").incr("c", b"n", 1).commit(), "TTL");
refused(ttl.mutate(b"r").put_i64("c", b"n", 1).commit(), "TTL");
let bucket = 4_000_000_000_000_000; ttl.mutate(b"r")
.incr_at("c", b"n", bucket, 3)
.commit()
.unwrap();
assert_eq!(buckets(&ttl, b"r", b"n"), [(bucket, 3)]);
let err = db
.table("bad")
.unwrap()
.family("c", Family::counter().merge_operator(""))
.create()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidArgument, "{err}");
assert!(t.row(b"r").read().unwrap().is_none());
}
#[test]
fn families_with_the_i64_operator_keep_0_1_0_semantics() {
let db = db();
let t = db
.table("legacy")
.unwrap()
.family("a", Family::default().merge_operator("pigeonhole.i64_add"))
.create()
.unwrap();
for d in [5, -2, 10] {
t.mutate(b"r").incr("a", b"hits", d).commit().unwrap();
}
let cell = t.get(b"r", "a", b"hits").unwrap().unwrap();
assert_eq!(cell.as_i64(), Some(13));
assert!(cell.timestamp() > 0, "at the commit timestamp");
t.mutate(b"r")
.merge("a", b"hits", &3i64.to_le_bytes())
.commit()
.unwrap();
assert_eq!(
t.get(b"r", "a", b"hits").map(|_| ()).unwrap_err().code(),
ErrorCode::MergeFailed
);
t.mutate(b"m").put("a", b"c", b"text").commit().unwrap();
t.mutate(b"m").incr("a", b"c", 1).commit().unwrap();
assert_eq!(
t.get(b"m", "a", b"c").map(|_| ()).unwrap_err().code(),
ErrorCode::MergeFailed
);
let err = t
.mutate(b"r")
.incr_at("a", b"n", 5, 1)
.commit()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidArgument, "{err}");
}
#[test]
fn write_batches_span_tables_and_refuse_foreign_tables() {
let db = db();
let t = table(&db);
let u = db
.table("u")
.unwrap()
.family("f", Family::default())
.create()
.unwrap();
let mut wb = db.write_batch();
assert!(wb.is_empty());
wb.put(&t, b"r1", "a", b"q", b"1")
.put_at(&t, b"r2", "a", b"q", 77, b"2")
.put(&u, b"x", "f", b"q", b"3")
.delete_column(&t, b"r3", "a", b"q")
.delete_row(&u, b"gone");
assert_eq!(wb.len(), 5);
let info = wb.commit_with(Durability::Buffered).unwrap();
assert_eq!(info.durability, Durability::Buffered);
assert_eq!(value(&t, b"r1", "a", b"q").as_deref(), Some(&b"1"[..]));
assert_eq!(t.get(b"r2", "a", b"q").unwrap().unwrap().timestamp(), 77);
assert_eq!(u.get(b"x", "f", b"q").unwrap().unwrap().value(), b"3");
let other = self::db();
let foreign = table(&other);
let mut wb = db.write_batch();
wb.put(&foreign, b"r", "a", b"q", b"v");
assert_eq!(wb.commit().unwrap_err().code(), ErrorCode::InvalidArgument);
}
#[test]
fn durability_defaults_and_overrides() {
let vfs = SimVfs::new(3);
let db = Pigeonhole::open(
"/db/d.phdb",
sim_options(&vfs).durability(Durability::Buffered),
)
.unwrap();
let t = table(&db);
assert_eq!(db.default_durability(), Durability::Buffered);
let info = t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
assert_eq!(info.durability, Durability::Buffered);
db.set_default_durability(Durability::None);
assert_eq!(db.default_durability(), Durability::None);
let mut wb = db.write_batch();
wb.put(&t, b"s", "a", b"q", b"v");
assert_eq!(wb.commit().unwrap().durability, Durability::None);
for d in [
Durability::None,
Durability::Buffered,
Durability::GroupSync,
Durability::Sync,
] {
let info = t
.mutate(b"r")
.put("a", b"q", b"v")
.durability(d)
.commit()
.unwrap();
assert_eq!(info.durability, d);
let mut wb = db.write_batch();
wb.put(&t, b"s", "a", b"q", b"v");
assert_eq!(wb.commit_with(d).unwrap().durability, d);
}
}
#[test]
fn conditional_commits() {
let db = db();
let t = table(&db);
let absent = Condition::Absent {
family: "a".into(),
qualifier: b"lock".to_vec(),
};
let exists = Condition::Exists {
family: "a".into(),
qualifier: b"lock".to_vec(),
};
assert!(
t.mutate(b"r")
.put("a", b"x", b"1")
.commit_if(&exists)
.unwrap()
.is_none()
);
assert!(
t.mutate(b"r")
.put("a", b"lock", b"w1")
.commit_if(&absent)
.unwrap()
.is_some()
);
assert!(
t.mutate(b"r")
.put("a", b"lock", b"w2")
.commit_if(&absent)
.unwrap()
.is_none()
);
let owner_is = |who: &[u8]| Condition::Value {
family: "a".into(),
qualifier: b"lock".to_vec(),
filter: ValueFilter::Equals(who.to_vec()),
};
assert!(
t.mutate(b"r")
.delete_column("a", b"lock")
.commit_if(&owner_is(b"w2"))
.unwrap()
.is_none()
);
let info = t
.mutate(b"r")
.delete_column("a", b"lock")
.durability(Durability::Sync)
.commit_if(&owner_is(b"w1"))
.unwrap()
.unwrap();
assert_eq!(info.durability, Durability::Sync);
assert_eq!(value(&t, b"r", "a", b"lock"), None);
assert_eq!(value(&t, b"r", "a", b"x"), None);
t.mutate(b"n").put_i64("a", b"v", 7).commit().unwrap();
let gt = |n| Condition::Value {
family: "a".into(),
qualifier: b"v".to_vec(),
filter: ValueFilter::I64(Ordering::Greater, n),
};
assert!(
t.mutate(b"n")
.put_i64("a", b"v", 8)
.commit_if(>(7))
.unwrap()
.is_none()
);
assert!(
t.mutate(b"n")
.put_i64("a", b"v", 8)
.commit_if(>(6))
.unwrap()
.is_some()
);
let bad = Condition::Exists {
family: "nope".into(),
qualifier: vec![],
};
assert_eq!(
t.mutate(b"n")
.put("a", b"v", b"x")
.commit_if(&bad)
.unwrap_err()
.code(),
ErrorCode::FamilyNotFound
);
}
#[test]
fn transactions_commit_or_conflict() {
let db = db();
let t = table(&db);
t.mutate(b"a").put("a", b"bal", b"10").commit().unwrap();
let mut txn = db.transaction().unwrap();
assert_eq!(
txn.get(&t, b"a", "a", b"bal").unwrap().unwrap().value(),
b"10"
);
txn.put(&t, b"a", "a", b"bal", b"5")
.put(&t, b"b", "a", b"bal", b"5")
.delete_column(&t, b"c", "a", b"bal");
let info = txn.commit_with(Durability::GroupSync).unwrap();
assert_eq!(info.durability, Durability::GroupSync);
assert_eq!(value(&t, b"b", "a", b"bal").as_deref(), Some(&b"5"[..]));
let mut txn = db.transaction().unwrap();
let _ = txn.get(&t, b"a", "a", b"bal").unwrap();
txn.put(&t, b"a", "a", b"bal", b"0");
t.mutate(b"a").put("a", b"bal", b"99").commit().unwrap();
assert_eq!(txn.commit().unwrap_err().code(), ErrorCode::Conflict);
assert_eq!(value(&t, b"a", "a", b"bal").as_deref(), Some(&b"99"[..]));
let mut txn = db.transaction().unwrap();
txn.put(&t, b"a", "nope", b"q", b"v");
assert_eq!(txn.commit().unwrap_err().code(), ErrorCode::FamilyNotFound);
let mut txn = db.transaction().unwrap();
assert_eq!(
txn.get(&t, b"a", "nope", b"q")
.map(|_| ())
.unwrap_err()
.code(),
ErrorCode::FamilyNotFound
);
}
#[test]
fn closed_handles_fail_with_closed() {
let db = db();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
db.clone().close().unwrap();
assert_eq!(
t.mutate(b"r")
.put("a", b"q", b"v")
.commit()
.unwrap_err()
.code(),
ErrorCode::Closed
);
}
#[test]
fn application_owned_mode() {
let vfs = SimVfs::new(5);
let err =
Pigeonhole::open_application_owned("/db/app.phdb", sim_options(&vfs).compaction_cores(1))
.unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidArgument);
let (db, mut shards) =
Pigeonhole::open_application_owned("/db/app.phdb", sim_options(&vfs)).unwrap();
assert_eq!(shards.len(), 2);
assert_eq!(shards.iter().map(|s| s.index()).collect::<Vec<_>>(), [0, 1]);
let woken = Arc::new(std::sync::atomic::AtomicUsize::new(0));
for s in &mut shards {
let w = Arc::clone(&woken);
s.set_wakeup(Box::new(move || {
w.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}));
}
let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let threads: Vec<_> = shards
.into_iter()
.map(|mut s| {
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
while s.run_once(Duration::from_millis(1))
|| !stop.load(std::sync::atomic::Ordering::Acquire)
{
std::thread::yield_now();
}
})
})
.collect();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
assert_eq!(value(&t, b"r", "a", b"q").as_deref(), Some(&b"v"[..]));
assert!(woken.load(std::sync::atomic::Ordering::Relaxed) > 0);
drop(t);
db.close().unwrap();
stop.store(true, std::sync::atomic::Ordering::Release);
for th in threads {
th.join().unwrap();
}
}
#[test]
fn real_files_reopen_and_lock() {
let dir = TempDir::new("reopen");
let path = dir.0.join("db.phdb");
let opts = || Options::default().shards(2).memtable_budget(4 << 20);
{
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"durable").commit().unwrap();
let second = Pigeonhole::open(&path, opts()).unwrap_err();
assert_eq!(second.code(), ErrorCode::WriterLocked);
db.flush().unwrap();
db.close().unwrap();
}
let missing = Pigeonhole::open(dir.0.join("none.phdb"), opts().create_if_missing(false));
assert!(missing.is_err());
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = db.table("t").unwrap().open().unwrap();
assert_eq!(value(&t, b"r", "a", b"q").as_deref(), Some(&b"durable"[..]));
db.compact().unwrap();
db.backup(dir.0.join("copy.phdb")).unwrap();
db.close().unwrap();
let copy = Pigeonhole::open(dir.0.join("copy.phdb"), opts().create_if_missing(false)).unwrap();
let t = copy.table("t").unwrap().open().unwrap();
assert_eq!(value(&t, b"r", "a", b"q").as_deref(), Some(&b"durable"[..]));
drop(t);
copy.close().unwrap();
}
#[test]
fn reader_handles_see_the_writers_commits() {
let vfs = SimVfs::new(9);
let db = Pigeonhole::open("/db/r.phdb", sim_options(&vfs)).unwrap();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"1").commit().unwrap();
let reader = Pigeonhole::open_reader(
"/db/r.phdb",
pigeonhole::ReaderOptions::default()
.vfs(Arc::clone(&vfs) as _)
.block_cache(1 << 20),
)
.unwrap();
assert_eq!(reader.tables(), ["t"]);
let rt = reader.table("t").unwrap();
assert_eq!(rt.name(), "t");
assert_eq!(rt.get(b"r", "a", b"q").unwrap().unwrap().value(), b"1");
t.mutate(b"s").put("a", b"q", b"2").commit().unwrap();
let snap = reader.snapshot().unwrap();
assert_eq!(
rt.get_at(&snap, b"s", "a", b"q").unwrap().unwrap().value(),
b"2"
);
assert_eq!(rt.scan_prefix(b"").iter().unwrap().count(), 2);
assert_eq!(rt.scan(&b"s"[..]..).iter().unwrap().count(), 1);
assert_eq!(
rt.scan_bounds(Bound::Unbounded, Bound::Excluded(b"s"))
.iter()
.unwrap()
.count(),
1
);
assert!(rt.row(b"r").read().unwrap().is_some());
assert_eq!(
reader.table("nope").unwrap_err().code(),
ErrorCode::TableNotFound
);
drop(rt);
drop(reader);
drop(t);
db.close().unwrap();
}
#[test]
fn snapshots_from_another_database_are_refused() {
let db = db();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
let other = self::db();
let foreign = other.snapshot().unwrap();
let code = |r: pigeonhole::Result<()>| r.unwrap_err().code();
assert_eq!(
code(t.get_at(&foreign, b"r", "a", b"q").map(|_| ())),
ErrorCode::InvalidArgument
);
assert_eq!(
code(t.row(b"r").snapshot(&foreign).read().map(|_| ())),
ErrorCode::InvalidArgument
);
assert_eq!(
code(t.scan_prefix(b"").snapshot(&foreign).iter().map(|_| ())),
ErrorCode::InvalidArgument
);
let own = db.clone().snapshot().unwrap();
assert!(t.get_at(&own, b"r", "a", b"q").unwrap().is_some());
}
#[test]
fn reads_after_close_fail_with_closed() {
let db = db();
let t = table(&db);
t.mutate(b"r").put("a", b"q", b"v").commit().unwrap();
let snap = db.snapshot().unwrap();
let handle = db.clone();
db.close().unwrap();
let code = |r: pigeonhole::Result<()>| r.unwrap_err().code();
assert_eq!(code(t.get(b"r", "a", b"q").map(|_| ())), ErrorCode::Closed);
assert_eq!(
code(t.get_at(&snap, b"r", "a", b"q").map(|_| ())),
ErrorCode::Closed
);
assert_eq!(code(t.row(b"r").read().map(|_| ())), ErrorCode::Closed);
assert_eq!(
code(t.scan_prefix(b"").iter().map(|_| ())),
ErrorCode::Closed
);
assert_eq!(code(handle.snapshot().map(|_| ())), ErrorCode::Closed);
assert_eq!(code(handle.transaction().map(|_| ())), ErrorCode::Closed);
assert_eq!(
code(handle.table("t").unwrap().open().map(|_| ())),
ErrorCode::Closed
);
let mut wb = handle.write_batch();
wb.put(&t, b"r", "a", b"q", b"v");
assert_eq!(code(wb.commit().map(|_| ())), ErrorCode::Closed);
}
#[test]
fn a_full_memtable_arena_is_busy_with_a_precise_message() {
let vfs = SimVfs::new(11);
let opts = |budget| {
Options::default()
.vfs(Arc::clone(&vfs) as _)
.shards(1)
.memtable_budget(budget)
.wal_segment_size(256 << 10)
};
let db = Pigeonhole::open("/db/busy.phdb", opts(1 << 20)).unwrap();
let t = table(&db);
let cell = vec![1u8; 1000];
for i in 0..3_000u32 {
t.mutate(&i.to_be_bytes())
.put("a", b"q", &cell)
.commit()
.unwrap();
}
let mut wb = db.write_batch();
let big = vec![2u8; 4096];
for i in 0..400u32 {
wb.put(&t, &i.to_be_bytes(), "a", b"big", &big);
}
let err = wb
.commit()
.expect_err("a batch larger than the arena is refused");
assert_eq!(err.code(), ErrorCode::BatchTooLarge);
assert!(err.message().contains("memtable_budget"), "{err}");
drop(t);
db.close().unwrap();
let db = Pigeonhole::open("/db/busy.phdb", opts(256 << 10)).unwrap();
let t = db.table("t").unwrap().open().unwrap();
assert_eq!(
value(&t, &0u32.to_be_bytes(), "a", b"q").as_deref(),
Some(&cell[..])
);
drop(t);
db.close().unwrap();
}
#[test]
fn size_errors_name_the_size_and_the_limit() {
let db = db();
let t = table(&db);
let err = t
.mutate(&[0u8; 70_000])
.put("a", b"q", b"v")
.commit()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::KeyTooLarge);
assert!(err.message().contains("row key of 70000 bytes"), "{err}");
assert!(err.message().contains("65536"), "{err}");
let err = t
.mutate(b"r")
.put("a", &[0u8; 66_000], b"v")
.commit()
.unwrap_err();
assert!(err.message().contains("qualifier of 66000 bytes"), "{err}");
let mut wb = db.write_batch();
wb.delete_row(&t, &[0u8; 70_000]);
assert!(
wb.commit()
.unwrap_err()
.message()
.contains("row key of 70000 bytes")
);
let big = vec![7u8; 300_000];
t.mutate(b"r").put("a", b"q", &big).commit().unwrap();
assert_eq!(
t.get(b"r", "a", b"q").unwrap().map(|c| c.value().to_vec()),
Some(big.clone())
);
let mut txn = db.transaction().unwrap();
txn.put(&t, b"r2", "a", b"q", &big);
txn.commit().unwrap();
let m = db
.table("m")
.unwrap()
.family("a", Family::default().merge_operator("pigeonhole.i64_add"))
.create_if_missing()
.unwrap();
let err = m.mutate(b"r").merge("a", b"n", &big).commit().unwrap_err();
assert_eq!(err.code(), ErrorCode::ValueTooLarge);
assert!(err.message().contains("value of 300000 bytes"), "{err}");
assert!(err.message().contains(&(192 * 1024).to_string()), "{err}");
}
#[test]
fn write_batches_have_every_row_mutation() {
let db = db();
let t = table(&db);
t.mutate(b"r")
.put_at("a", b"x", 10, b"old")
.put("b", b"y", b"1")
.commit()
.unwrap();
let mut wb = db.write_batch();
wb.put_i64(&t, b"r", "a", b"n", 41)
.put_f64(&t, b"s", "a", b"f", 1.5)
.delete_cell(&t, b"r", "a", b"x", 10)
.delete_family(&t, b"r", "b");
assert_eq!(wb.len(), 4);
wb.commit().unwrap();
assert_eq!(t.get(b"r", "a", b"n").unwrap().unwrap().as_i64(), Some(41));
let mut wb = db.write_batch();
wb.merge(&t, b"m", "a", b"n", &1i64.to_le_bytes());
assert_eq!(wb.commit().unwrap_err().code(), ErrorCode::InvalidArgument);
assert_eq!(
t.get(b"s", "a", b"f").unwrap().unwrap().typed(),
Value::F64(1.5)
);
assert!(t.get(b"r", "a", b"x").unwrap().is_none());
assert!(t.get(b"r", "b", b"y").unwrap().is_none());
}
#[test]
fn durability_defaults_and_overrides_on_every_commit_path() {
let vfs = SimVfs::new(13);
let db = Pigeonhole::open("/db/dd.phdb", sim_options(&vfs)).unwrap();
assert_eq!(db.default_durability(), Durability::GroupSync);
let t = table(&db);
assert_eq!(
t.mutate(b"r")
.put("a", b"q", b"v")
.commit()
.unwrap()
.durability,
Durability::GroupSync
);
let mut txn = db.transaction().unwrap();
txn.put(&t, b"r", "a", b"q", b"t");
assert_eq!(
txn.commit_with(Durability::Buffered).unwrap().durability,
Durability::Buffered
);
let mut txn = db.transaction().unwrap();
txn.put(&t, b"r", "a", b"q", b"u");
assert_eq!(txn.commit().unwrap().durability, Durability::GroupSync);
let always = Condition::Exists {
family: "a".into(),
qualifier: b"q".to_vec(),
};
let info = t
.mutate(b"r")
.put("a", b"q", b"w")
.durability(Durability::None)
.commit_if(&always)
.unwrap()
.unwrap();
assert_eq!(info.durability, Durability::None);
}
#[test]
fn empty_table_names_and_tables_without_families_are_refused() {
let db = db();
let err = db
.table("")
.unwrap()
.family("f", Family::default())
.create()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidArgument);
for r in [
db.table("t").unwrap().create(),
db.table("t").unwrap().create_if_missing(),
] {
assert_eq!(r.unwrap_err().code(), ErrorCode::InvalidArgument);
}
assert!(db.tables().is_empty());
table(&db);
assert_eq!(
db.table("t").unwrap().open().unwrap().families(),
["a", "b"]
);
}
#[test]
fn days_saturates() {
assert_eq!(days(2), Duration::from_secs(2 * 86_400));
assert_eq!(days(u64::MAX), Duration::from_secs(u64::MAX));
}
#[test]
fn data_beyond_the_memtable_budget_lives_in_one_file() {
let dir = TempDir::new("beyond-budget");
let path = dir.0.join("big.phdb");
let opts = || Options::default().shards(1).memtable_budget(2 << 20);
let cell = vec![7u8; 1000];
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = table(&db);
for i in 0..8_000u32 {
t.mutate(&i.to_be_bytes())
.put("a", b"q", &cell)
.durability(Durability::Buffered)
.commit()
.unwrap();
}
db.compact().unwrap();
drop(t);
db.close().unwrap();
let names: Vec<_> = std::fs::read_dir(&dir.0)
.unwrap()
.map(|e| e.unwrap().file_name())
.collect();
assert_eq!(
names,
[std::ffi::OsString::from("big.phdb")],
"one file at rest"
);
assert!(std::fs::metadata(&path).unwrap().len() > 1 << 20);
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = db.table("t").unwrap().open().unwrap();
assert_eq!(t.scan_prefix(b"").iter().unwrap().count(), 8_000);
assert_eq!(
value(&t, &7_999u32.to_be_bytes(), "a", b"q").as_deref(),
Some(&cell[..])
);
drop(t);
db.close().unwrap();
}
#[test]
fn flush_compact_and_backup_through_the_public_api() {
let dir = TempDir::new("maintenance");
let path = dir.0.join("db.phdb");
let opts = || Options::default().shards(2).memtable_budget(4 << 20);
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = table(&db);
t.mutate(b"none")
.put("a", b"q", b"v1")
.durability(Durability::None)
.commit()
.unwrap();
for v in 0..5u8 {
t.mutate(b"ver").put("b", b"q", &[v]).commit().unwrap();
}
db.flush().unwrap();
db.compact().unwrap();
assert_eq!(cells(&t, b"ver").len(), 2);
let copy_path = dir.0.join("copy.phdb");
db.backup(©_path).unwrap();
t.mutate(b"none").put("a", b"q", b"v2").commit().unwrap();
t.mutate(b"later").put("a", b"q", b"x").commit().unwrap();
assert!(db.backup(©_path).is_err());
let copy = Pigeonhole::open(©_path, opts().create_if_missing(false)).unwrap();
let ct = copy.table("t").unwrap().open().unwrap();
assert_eq!(value(&ct, b"none", "a", b"q").as_deref(), Some(&b"v1"[..]));
assert_eq!(value(&ct, b"later", "a", b"q"), None);
assert_eq!(cells(&ct, b"ver").len(), 2);
drop(ct);
copy.close().unwrap();
let handle = db.clone();
drop(t);
db.close().unwrap();
assert_eq!(handle.flush().unwrap_err().code(), ErrorCode::Closed);
assert_eq!(handle.compact().unwrap_err().code(), ErrorCode::Closed);
assert_eq!(
handle.backup(dir.0.join("late.phdb")).unwrap_err().code(),
ErrorCode::Closed
);
}
#[test]
fn shrink_releases_file_space_after_deletes_and_keeps_data() {
let dir = TempDir::new("shrink");
let path = dir.0.join("db.phdb");
let opts = || Options::default().shards(1).memtable_budget(4 << 20);
let len = || std::fs::metadata(&path).unwrap().len();
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = table(&db);
t.mutate(b"keep").put("a", b"q", b"kept").commit().unwrap();
let junk = db
.table("junk")
.unwrap()
.family("f", Family::default())
.create_if_missing()
.unwrap();
let big = [9u8; 1024];
for round in 0..4u32 {
for i in 0..2_000u32 {
let row = (round * 2_000 + i).to_be_bytes();
junk.mutate(&row)
.put("f", b"q", &big)
.durability(Durability::None)
.commit()
.unwrap();
}
db.flush().unwrap();
}
let grown = len();
for i in 0..8_000u32 {
junk.mutate(&i.to_be_bytes())
.delete_row()
.durability(Durability::None)
.commit()
.unwrap();
}
drop(junk);
db.compact().unwrap();
let before = len();
assert!(before >= grown);
let released = db.shrink().unwrap();
let after = len();
assert!(released > 0, "no space released");
assert!(after < before, "{after} >= {before}");
assert_eq!(value(&t, b"keep", "a", b"q").as_deref(), Some(&b"kept"[..]));
t.mutate(b"late").put("a", b"q", b"x").commit().unwrap();
db.shrink().unwrap();
drop(t);
db.close().unwrap();
let db = Pigeonhole::open(&path, opts().create_if_missing(false)).unwrap();
let t = db.table("t").unwrap().open().unwrap();
assert_eq!(value(&t, b"keep", "a", b"q").as_deref(), Some(&b"kept"[..]));
assert_eq!(value(&t, b"late", "a", b"q").as_deref(), Some(&b"x"[..]));
let handle = db.clone();
drop(t);
db.close().unwrap();
assert_eq!(handle.shrink().unwrap_err().code(), ErrorCode::Closed);
}
#[test]
fn shrink_after_deleting_every_row_leaves_a_small_file() {
let dir = TempDir::new("shrink-all");
let path = dir.0.join("db.phdb");
let opts = || Options::default().shards(1);
let len = || std::fs::metadata(&path).unwrap().len();
let db = Pigeonhole::open(&path, opts()).unwrap();
let t = db
.table("t")
.unwrap()
.family("f", Family::default().max_versions(1))
.create_if_missing()
.unwrap();
let mut x = 88_172_645_463_325_252u64;
for i in 0..3_000u32 {
let mut v = [0u8; 1024];
for b in &mut v {
x ^= x << 13;
x ^= x >> 7;
x ^= x << 17;
*b = x as u8;
}
t.mutate(&i.to_be_bytes())
.put("f", b"q", &v)
.durability(Durability::Buffered)
.commit()
.unwrap();
}
db.flush().unwrap();
db.compact().unwrap();
for i in 0..3_000u32 {
t.mutate(&i.to_be_bytes())
.delete_row()
.durability(Durability::Buffered)
.commit()
.unwrap();
}
db.flush().unwrap();
db.compact().unwrap();
let before = len();
assert!(before >= 3 << 20, "the load grew the file: {before}");
let released = db.shrink().unwrap();
let after = len();
assert!(released > 0 && after < before, "{before} -> {after}");
assert!(
after <= 1 << 20,
"no row is left, yet the file is {after} bytes"
);
assert_eq!(t.scan_prefix(b"").iter().unwrap().count(), 0);
drop(t);
db.close().unwrap();
let db = Pigeonhole::open(&path, opts()).unwrap();
db.compact().unwrap();
db.shrink().unwrap();
db.close().unwrap();
assert!(len() <= 1 << 20, "after a reopen: {}", len());
}
#[test]
fn a_flush_makes_none_commits_survive_a_power_loss() {
let vfs = SimVfs::new(21);
let opts = || sim_options(&vfs).shards(1);
let db = Pigeonhole::open("/db/flush-crash.phdb", opts()).unwrap();
let t = table(&db);
t.mutate(b"flushed")
.put("a", b"q", b"v")
.durability(Durability::None)
.commit()
.unwrap();
db.flush().unwrap();
t.mutate(b"after")
.put("a", b"q", b"v")
.durability(Durability::None)
.commit()
.unwrap();
vfs.crash(pigeonhole_io::sim::CrashKind::Power);
drop(t);
drop(db);
let db = Pigeonhole::open("/db/flush-crash.phdb", opts()).unwrap();
let t = db.table("t").unwrap().open().unwrap();
assert_eq!(value(&t, b"flushed", "a", b"q").as_deref(), Some(&b"v"[..]));
assert_eq!(value(&t, b"after", "a", b"q"), None);
drop(t);
db.close().unwrap();
}
struct Driver {
ctl: Arc<(std::sync::Mutex<DriverCtl>, std::sync::Condvar)>,
thread: std::thread::JoinHandle<()>,
}
#[derive(Default)]
struct DriverCtl {
hold: bool,
held: bool,
stop: bool,
}
impl Driver {
fn start(mut shards: Vec<pigeonhole::Shard>) -> Self {
let ctl = Arc::new((
std::sync::Mutex::new(DriverCtl::default()),
std::sync::Condvar::new(),
));
let c = Arc::clone(&ctl);
let thread = std::thread::spawn(move || {
loop {
let mut busy = false;
for s in &mut shards {
busy |= s.run_once(Duration::from_millis(1));
}
let (m, cv) = &*c;
let mut st = m.lock().unwrap();
if st.stop && !busy {
return;
}
if st.hold && !busy {
st.held = true;
cv.notify_all();
while st.hold {
st = cv.wait(st).unwrap();
}
st.held = false;
}
drop(st);
if !busy {
std::thread::yield_now();
}
}
});
Self { ctl, thread }
}
fn quiesce(&self) {
let (m, cv) = &*self.ctl;
let mut st = m.lock().unwrap();
st.hold = true;
while !st.held {
st = cv.wait(st).unwrap();
}
}
fn resume(&self) {
let (m, cv) = &*self.ctl;
m.lock().unwrap().hold = false;
cv.notify_all();
}
fn stop(self) {
self.ctl.0.lock().unwrap().stop = true;
self.thread.join().unwrap();
}
}
#[test]
fn compact_after_drop_table_covers_the_live_tables() {
let vfs = SimVfs::new(83);
let path = "/db/drop-compact.phdb";
let file_len = || {
use pigeonhole_io::Vfs;
vfs.open(path.as_ref(), pigeonhole_io::OpenOptions::read())
.unwrap()
.len()
.unwrap()
};
let (db, shards) = Pigeonhole::open_application_owned(path, sim_options(&vfs)).unwrap();
let driver = Driver::start(shards);
let make = |name: &str| {
db.table(name)
.unwrap()
.family("a", Family::default())
.create_if_missing()
.unwrap()
};
let keep = make("keep");
let gone = make("gone");
for i in 0..400u32 {
let row = format!("row{i:05}");
keep.mutate(row.as_bytes())
.put("a", b"q", &[7; 300])
.commit()
.unwrap();
for j in 0..10u32 {
gone.mutate(format!("{row}-{j}").as_bytes())
.put("a", b"q", &[9; 1000])
.commit()
.unwrap();
}
}
db.flush().unwrap();
let before = file_len();
drop(gone);
db.drop_table("gone").unwrap();
db.compact().unwrap();
assert_eq!(keep.scan_prefix(b"").iter().unwrap().count(), 400);
db.compact().unwrap();
driver.quiesce();
db.shrink().unwrap();
driver.resume();
let after = file_len();
assert!(after < before / 2, "{after} vs {before}");
assert_eq!(keep.scan_prefix(b"").iter().unwrap().count(), 400);
drop(keep);
db.close().unwrap();
driver.stop();
let db = Pigeonhole::open(path, sim_options(&vfs)).unwrap();
assert!(db.table("gone").unwrap().open().is_err());
let keep = db.table("keep").unwrap().open().unwrap();
assert_eq!(keep.scan_prefix(b"").iter().unwrap().count(), 400);
assert_eq!(
value(&keep, b"row00399", "a", b"q").as_deref(),
Some(&[7u8; 300][..])
);
drop(keep);
db.close().unwrap();
}
#[derive(Debug)]
struct Append;
impl pigeonhole::MergeOperator for Append {
fn name(&self) -> &str {
"app.append"
}
fn merge(&self, acc: &mut Vec<u8>, older: &[u8]) -> Result<(), pigeonhole::MergeError> {
let newer = acc.split_off(1);
acc.extend_from_slice(&older[1..]);
acc.extend_from_slice(&newer);
Ok(())
}
fn finish(&self, base: Option<&[u8]>, acc: &mut Vec<u8>) -> Result<(), pigeonhole::MergeError> {
if let Some(base) = base {
let newer = acc.split_off(1);
acc.extend_from_slice(&base[1..]);
acc.extend_from_slice(&newer);
}
Ok(())
}
}
fn log_table(db: &Pigeonhole) -> pigeonhole::Result<Table> {
db.table("log")?
.family("l", Family::default().merge_operator("app.append"))
.family("plain", Family::default())
.create_if_missing()
}
fn log_value(t: &Table, row: &[u8]) -> pigeonhole::Result<Option<Vec<u8>>> {
Ok(t.get(row, "l", b"q")?.map(|c| c.value().to_vec()))
}
#[test]
fn registered_merge_operators_resolve_through_flush_compaction_and_reopen() {
let vfs = SimVfs::new(43);
let with_append = || sim_options(&vfs).merge_operator(Arc::new(Append));
let db = Pigeonhole::open("/db/merge.phdb", with_append()).unwrap();
let t = log_table(&db).unwrap();
t.mutate(b"r").put("l", b"q", b"a").commit().unwrap();
for part in [&b"b"[..], b"c"] {
t.mutate(b"r").merge("l", b"q", part).commit().unwrap();
}
t.mutate(b"r").put("plain", b"q", b"p").commit().unwrap();
assert_eq!(log_value(&t, b"r").unwrap().as_deref(), Some(&b"abc"[..]));
db.flush().unwrap();
db.compact().unwrap();
t.mutate(b"r").merge("l", b"q", b"d").commit().unwrap();
assert_eq!(log_value(&t, b"r").unwrap().as_deref(), Some(&b"abcd"[..]));
drop(t);
db.close().unwrap();
let db = Pigeonhole::open("/db/merge.phdb", with_append()).unwrap();
let t = log_table(&db).unwrap();
assert_eq!(log_value(&t, b"r").unwrap().as_deref(), Some(&b"abcd"[..]));
drop(t);
db.close().unwrap();
let err = Pigeonhole::open("/db/merge.phdb", sim_options(&vfs))
.expect_err("opened without the family's merge operator");
assert_eq!(err.code(), ErrorCode::UnknownMergeOperator, "{err}");
let db = Pigeonhole::open(
"/db/merge.phdb",
sim_options(&vfs).allow_unregistered_merge_operators(true),
)
.unwrap();
let t = db.table("log").unwrap().open().unwrap();
let err = log_value(&t, b"r").unwrap_err();
assert_eq!(err.code(), ErrorCode::UnknownMergeOperator, "{err}");
assert_eq!(
t.get(b"r", "plain", b"q")
.unwrap()
.map(|c| c.value().to_vec()),
Some(b"p".to_vec())
);
let err = t
.mutate(b"s")
.put("plain", b"q", b"x")
.commit()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::ReadOnly, "{err}");
drop(t);
db.close().unwrap();
}
#[test]
fn a_family_naming_an_unregistered_operator_is_refused_at_creation() {
let db = db();
let err = log_table(&db).unwrap_err();
assert_eq!(err.code(), ErrorCode::UnknownMergeOperator, "{err}");
assert!(db.tables().is_empty());
}
#[test]
fn the_write_stall_timeout_bounds_a_wait_for_room() {
let dir = TempDir::new("stall-timeout");
let timeout = Duration::from_millis(200);
let options = Options::default()
.shards(1)
.memtable_budget(4 << 20)
.tablet_changes(false)
.write_stall_timeout(timeout);
let db = Pigeonhole::open(dir.0.join("stall.phdb"), options).unwrap();
let t = table(&db);
let value = vec![7u8; 16 << 10];
let mut snaps = Vec::new();
let mut refused = None;
for i in 0..5_000u32 {
let started = std::time::Instant::now();
match t
.mutate(format!("r{i:05}").as_bytes())
.put("a", b"q", &value)
.commit()
{
Ok(_) => snaps.push(db.snapshot().unwrap()),
Err(e) => {
assert_eq!(e.code(), ErrorCode::Busy, "{e}");
refused = Some(started.elapsed());
break;
}
}
}
let waited = refused.expect("snapshots holding the arena never stalled a write");
assert!(
waited >= timeout && waited < Duration::from_secs(10),
"refused after {waited:?}, the timeout is {timeout:?}"
);
drop(snaps);
t.mutate(b"after").put("a", b"q", b"v").commit().unwrap();
}
#[test]
fn a_registered_operator_folds_onto_a_separated_base() {
let vfs = SimVfs::new(44);
let db = Pigeonhole::open(
"/db/merge-blob.phdb",
sim_options(&vfs).merge_operator(Arc::new(Append)),
)
.unwrap();
let t = db
.table("log")
.unwrap()
.family(
"l",
Family::default()
.merge_operator("app.append")
.blob_threshold(100),
)
.create_if_missing()
.unwrap();
let base = vec![b'b'; 500];
t.mutate(b"r").put("l", b"q", &base).commit().unwrap();
db.flush().unwrap();
t.mutate(b"r").merge("l", b"q", b"xy").commit().unwrap();
let mut want = base.clone();
want.extend_from_slice(b"xy");
for stage in ["operand in a memtable", "flushed", "compacted"] {
match stage {
"flushed" => db.flush().unwrap(),
"compacted" => db.compact().unwrap(),
_ => {}
}
assert_eq!(
log_value(&t, b"r").unwrap().as_deref(),
Some(&want[..]),
"{stage}: get"
);
let row = t.row(b"r").read().unwrap().unwrap();
assert_eq!(
row.get("l", b"q").map(|c| c.value().to_vec()).as_deref(),
Some(&want[..]),
"{stage}: row read"
);
}
drop(t);
db.close().unwrap();
}