use super::*;
use crate::nsw::*;
use alloc::string::ToString;
use alloc::vec;
#[cfg(test)]
fn normalize_rowids_dense(c: &mut Catalog, tables: &[&str]) {
for name in tables {
c.get_mut(name).unwrap().assign_dense_rowids();
}
}
#[test]
fn redo_apply_matches_direct_position_ops() {
fn fresh() -> Catalog {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("v", DataType::Text, true),
],
))
.unwrap();
c
}
let rows = [
Row::new(alloc::vec![Value::BigInt(1), Value::text("a")]),
Row::new(alloc::vec![Value::BigInt(2), Value::text("b")]),
Row::new(alloc::vec![Value::BigInt(3), Value::text("c")]),
Row::new(alloc::vec![Value::BigInt(4), Value::text("d")]),
];
let upd = alloc::vec![Value::BigInt(2), Value::text("B")];
let mut c1 = fresh();
{
let t = c1.get_mut("t").unwrap();
for r in &rows {
t.insert(r.clone()).unwrap();
}
t.update_row(1, upd.clone()).unwrap(); t.delete_rows(&[0, 2]); }
let mut log = alloc::vec::Vec::new();
for r in &rows {
log.push(RowChange::Insert {
table: "t".to_string(),
row: r.clone(),
rowid: row_header::RowId::UNASSIGNED,
writer_version: 0,
});
}
log.push(RowChange::Update {
table: "t".to_string(),
pos: 1,
new_row: upd,
rowid: row_header::RowId::UNASSIGNED,
writer_version: 0,
});
log.push(RowChange::Delete {
table: "t".to_string(),
positions: vec![0, 2],
rowids: vec![row_header::RowId::UNASSIGNED; 2],
writer_version: 0,
});
let mut c2 = fresh();
c2.apply_redo(&log).unwrap();
normalize_rowids_dense(&mut c1, &["t"]);
normalize_rowids_dense(&mut c2, &["t"]);
assert_eq!(
c1.serialize(),
c2.serialize(),
"redo apply diverged from direct position ops"
);
let mut c3 = fresh();
assert!(
c3.apply_redo(&[RowChange::Insert {
table: "nope".to_string(),
row: rows[0].clone(),
rowid: row_header::RowId::UNASSIGNED,
writer_version: 0,
}])
.is_err()
);
}
#[test]
fn redo_capture_replays_to_identical_state() {
fn fresh() -> Catalog {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("v", DataType::Text, true),
],
))
.unwrap();
c
}
let mk = |id: i64, v: &str| Row::new(alloc::vec![Value::BigInt(id), Value::text(v)]);
let mut c1 = fresh();
{
let t = c1.get_mut("t").unwrap();
t.enable_redo();
t.insert(mk(1, "a")).unwrap();
t.insert(mk(2, "b")).unwrap();
t.insert(mk(3, "c")).unwrap();
t.update_row(1, alloc::vec![Value::BigInt(2), Value::text("B")])
.unwrap();
t.delete_rows(&[0]); t.delete_rows(&[99]); }
let log = c1.get_mut("t").unwrap().take_redo();
assert_eq!(log.len(), 5, "captured log: {log:?}");
let mut c2 = fresh();
c2.apply_redo(&log).unwrap();
normalize_rowids_dense(&mut c1, &["t"]);
normalize_rowids_dense(&mut c2, &["t"]);
assert_eq!(
c1.serialize(),
c2.serialize(),
"replayed capture diverged from execution"
);
assert!(c1.get_mut("t").unwrap().take_redo().is_empty());
}
#[test]
fn catalog_drain_redo_replays_multi_table() {
fn fresh() -> Catalog {
let mut c = Catalog::new();
for name in ["a", "b"] {
c.create_table(TableSchema::new(
name,
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("v", DataType::Text, true),
],
))
.unwrap();
}
c
}
let mk = |id: i64, v: &str| Row::new(alloc::vec![Value::BigInt(id), Value::text(v)]);
let mut c1 = fresh();
c1.enable_redo_all();
c1.get_mut("a").unwrap().insert(mk(1, "a1")).unwrap();
c1.get_mut("b").unwrap().insert(mk(2, "b1")).unwrap();
c1.get_mut("a").unwrap().insert(mk(3, "a2")).unwrap();
c1.get_mut("a")
.unwrap()
.update_row(0, alloc::vec![Value::BigInt(1), Value::text("A1")])
.unwrap();
c1.get_mut("b").unwrap().delete_rows(&[0]);
let log = c1.drain_redo();
let mut c2 = fresh();
c2.apply_redo(&log).unwrap();
normalize_rowids_dense(&mut c1, &["a", "b"]);
normalize_rowids_dense(&mut c2, &["a", "b"]);
assert_eq!(c1.serialize(), c2.serialize(), "multi-table redo diverged");
assert!(c1.drain_redo().is_empty());
}
#[test]
fn redo_log_codec_round_trips() {
use row_header::RowId;
let changes = vec![
RowChange::Insert {
table: "t".to_string(),
row: Row::new(alloc::vec![
Value::BigInt(1),
Value::text("a"),
Value::Null,
Value::Bool(true),
]),
rowid: RowId(11),
writer_version: 101,
},
RowChange::Update {
table: "users".to_string(),
pos: 42,
new_row: alloc::vec![Value::Int(7), Value::bytes(alloc::vec![1, 2, 3])],
rowid: RowId(22),
writer_version: 202,
},
RowChange::Delete {
table: "t".to_string(),
positions: alloc::vec![0, 5, 99],
rowids: alloc::vec![RowId(3), RowId(4), RowId(5)],
writer_version: 303,
},
RowChange::Delete {
table: "empty".to_string(),
positions: alloc::vec![],
rowids: alloc::vec![],
writer_version: 0,
},
RowChange::Tombstone {
table: "t".to_string(),
rowids: alloc::vec![RowId(7), RowId(8)],
xmax: 404,
},
RowChange::Tombstone {
table: "solo".to_string(),
rowids: alloc::vec![RowId(9)],
xmax: 505,
},
];
let bytes = encode_redo_log(&changes);
assert_eq!(decode_redo_log(&bytes).unwrap(), changes);
let empty = encode_redo_log(&[]);
assert_eq!(decode_redo_log(&empty).unwrap(), Vec::<RowChange>::new());
assert!(decode_redo_log(&bytes[..bytes.len() / 2]).is_err());
assert!(decode_redo_log(&[]).is_err());
}
#[test]
fn redo_log_old_format_decodes_and_replays_identically() {
use crate::codec;
use row_header::RowId;
fn encode_old(changes: &[RowChange]) -> Vec<u8> {
let mut out = Vec::new();
out.push(crate::FILE_VERSION);
codec::write_u32(&mut out, changes.len() as u32);
let write_values = |out: &mut Vec<u8>, vals: &[Value<'static>]| {
codec::write_u32(out, vals.len() as u32);
for v in vals {
codec::write_value(out, v);
}
};
for change in changes {
match change {
RowChange::Insert { table, row, .. } => {
out.push(0);
codec::write_str(&mut out, table);
write_values(&mut out, &row.values);
}
RowChange::Update {
table,
pos,
new_row,
..
} => {
out.push(1);
codec::write_str(&mut out, table);
codec::write_u32(&mut out, *pos as u32);
write_values(&mut out, new_row);
}
RowChange::Delete {
table, positions, ..
} => {
out.push(2);
codec::write_str(&mut out, table);
codec::write_u32(&mut out, positions.len() as u32);
for p in positions {
codec::write_u32(&mut out, *p as u32);
}
}
RowChange::Tombstone { .. } => {
unreachable!("old redo layout has no Tombstone op")
}
}
}
out
}
fn fresh() -> Catalog {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("v", DataType::Text, true),
],
))
.unwrap();
c
}
let mk = |id: i64, v: &str| Row::new(alloc::vec![Value::BigInt(id), Value::text(v)]);
let logical = vec![
RowChange::Insert {
table: "t".to_string(),
row: mk(1, "a"),
rowid: RowId(999), writer_version: 7, },
RowChange::Insert {
table: "t".to_string(),
row: mk(2, "b"),
rowid: RowId(999),
writer_version: 7,
},
RowChange::Update {
table: "t".to_string(),
pos: 0,
new_row: alloc::vec![Value::BigInt(1), Value::text("A")],
rowid: RowId(999),
writer_version: 7,
},
RowChange::Delete {
table: "t".to_string(),
positions: alloc::vec![1],
rowids: alloc::vec![RowId(999)],
writer_version: 7,
},
];
let old_bytes = encode_old(&logical);
assert_ne!(old_bytes[0], 0xFF, "old layout must not look like new");
let decoded = decode_redo_log(&old_bytes).unwrap();
let expected = vec![
RowChange::Insert {
table: "t".to_string(),
row: mk(1, "a"),
rowid: RowId::UNASSIGNED,
writer_version: 0,
},
RowChange::Insert {
table: "t".to_string(),
row: mk(2, "b"),
rowid: RowId::UNASSIGNED,
writer_version: 0,
},
RowChange::Update {
table: "t".to_string(),
pos: 0,
new_row: alloc::vec![Value::BigInt(1), Value::text("A")],
rowid: RowId::UNASSIGNED,
writer_version: 0,
},
RowChange::Delete {
table: "t".to_string(),
positions: alloc::vec![1],
rowids: alloc::vec![], writer_version: 0,
},
];
assert_eq!(decoded, expected, "old-format decode diverged");
let mut direct = fresh();
{
let t = direct.get_mut("t").unwrap();
t.insert(mk(1, "a")).unwrap();
t.insert(mk(2, "b")).unwrap();
t.update_row(0, alloc::vec![Value::BigInt(1), Value::text("A")])
.unwrap();
t.delete_rows(&[1]);
}
let mut replayed = fresh();
replayed.apply_redo(&decoded).unwrap();
normalize_rowids_dense(&mut direct, &["t"]);
normalize_rowids_dense(&mut replayed, &["t"]);
assert_eq!(
direct.serialize(),
replayed.serialize(),
"old-format redo replay diverged from direct ops"
);
}
#[test]
fn redo_tombstone_survives_replay_hidden_from_snapshot() {
use crate::row_header::XMAX_ALIVE;
use crate::snapshot::Snapshot;
fn fresh() -> Catalog {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
c
}
let mk = |id: i32| Row::new(alloc::vec![Value::Int(id)]);
let snap = Snapshot::unbounded();
const TOMB_XMAX: u64 = 42;
let mut cap = fresh();
cap.enable_redo_all();
{
let t = cap.get_mut("t").unwrap();
t.insert(mk(1)).unwrap();
t.insert(mk(2)).unwrap();
t.insert(mk(3)).unwrap();
t.mark_row_deleted(1, TOMB_XMAX).unwrap();
}
let log = cap.drain_redo();
let tomb_ids: Vec<_> = log
.iter()
.filter_map(|c| match c {
RowChange::Tombstone { rowids, xmax, .. } => Some((rowids.clone(), *xmax)),
_ => None,
})
.collect();
assert_eq!(
tomb_ids.len(),
1,
"one in-place delete → one Tombstone redo"
);
assert_eq!(
tomb_ids[0].1, TOMB_XMAX,
"tombstone carries the writer version"
);
assert_eq!(tomb_ids[0].0.len(), 1, "one row tombstoned");
let bytes = encode_redo_log(&log);
let decoded = decode_redo_log(&bytes).unwrap();
assert_eq!(decoded, log, "tombstone redo must round-trip the codec");
let unresolved_before = crate::unresolved_tombstone_count();
let mut rep = fresh();
rep.apply_redo(&decoded).unwrap();
assert_eq!(
crate::unresolved_tombstone_count(),
unresolved_before,
"the tombstone target must resolve by RowId within the replay run"
);
let t = rep.get("t").unwrap();
assert_eq!(
t.rows().len(),
3,
"tombstone must NOT physically remove the row"
);
let visible: Vec<i32> = t
.scan_visible(&snap)
.filter_map(|(_, r)| match r.values.first() {
Some(Value::Int(v)) => Some(*v),
_ => None,
})
.collect();
assert_eq!(
visible,
alloc::vec![1, 3],
"tombstoned row must be hidden after replay"
);
let hidden_idx = t
.rows()
.iter()
.position(|r| r.values.first() == Some(&Value::Int(2)))
.expect("id=2 physically present");
let h = t.headers().get(hidden_idx).expect("header lock-step");
assert_eq!(
h.xmax, TOMB_XMAX,
"recovered row must carry the tombstone xmax"
);
assert!(h.is_deleted(), "recovered row must read as deleted");
for (i, r) in t.rows().iter().enumerate() {
if r.values.first() != Some(&Value::Int(2)) {
assert_eq!(
t.headers().get(i).unwrap().xmax,
XMAX_ALIVE,
"survivor {i} must stay alive"
);
}
}
let mut capc = fresh();
capc.enable_redo_all();
{
let t = capc.get_mut("t").unwrap();
t.insert(mk(1)).unwrap();
t.insert(mk(2)).unwrap();
t.insert(mk(3)).unwrap();
t.delete_rows(&[1]); }
let logc = capc.drain_redo();
assert!(
logc.iter()
.all(|c| !matches!(c, RowChange::Tombstone { .. })),
"gate-off DELETE must NOT emit a Tombstone redo"
);
let mut repc = fresh();
repc.apply_redo(&decode_redo_log(&encode_redo_log(&logc)).unwrap())
.unwrap();
let tc = repc.get("t").unwrap();
assert_eq!(
tc.rows().len(),
2,
"gate-off replay physically removes the row"
);
let ids: Vec<i32> = tc
.rows()
.iter()
.filter_map(|r| match r.values.first() {
Some(Value::Int(v)) => Some(*v),
_ => None,
})
.collect();
assert_eq!(
ids,
alloc::vec![1, 3],
"gate-off replay keeps only survivors"
);
}
#[test]
fn redo_update_tombstone_plus_insert_survives_replay() {
use crate::row_header::XMAX_ALIVE;
use crate::snapshot::Snapshot;
fn fresh() -> Catalog {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
c
}
let mk = |id: i32| Row::new(alloc::vec![Value::Int(id)]);
let snap = Snapshot::unbounded();
const STMT_V: u64 = 42;
const OLD_VAL: i32 = 2;
const NEW_VAL: i32 = 20;
let mut cap = fresh();
cap.enable_redo_all();
{
let t = cap.get_mut("t").unwrap();
t.insert(mk(1)).unwrap();
t.insert(mk(OLD_VAL)).unwrap();
t.insert(mk(3)).unwrap();
t.mark_row_deleted(1, STMT_V).unwrap();
t.insert_with_xmin(mk(NEW_VAL), STMT_V).unwrap();
}
let log = cap.drain_redo();
let tomb_pos = log
.iter()
.position(|c| matches!(c, RowChange::Tombstone { .. }))
.expect("in-place UPDATE emits a Tombstone for the old version");
let (tomb_rowids, tomb_xmax) = match &log[tomb_pos] {
RowChange::Tombstone { rowids, xmax, .. } => (rowids.clone(), *xmax),
_ => unreachable!(),
};
assert_eq!(
tomb_rowids.len(),
1,
"one row superseded → one tombstone target"
);
assert_eq!(
tomb_xmax, STMT_V,
"tombstone carries the statement writer version"
);
assert_eq!(
log.iter()
.filter(|c| matches!(c, RowChange::Tombstone { .. }))
.count(),
1,
"an in-place UPDATE tombstones the old version exactly once"
);
let new_ins_pos = log
.iter()
.position(|c| {
matches!(
c,
RowChange::Insert { row, .. } if row.values.first() == Some(&Value::Int(NEW_VAL))
)
})
.expect("in-place UPDATE emits an Insert carrying the new values");
assert!(
tomb_pos < new_ins_pos,
"tombstone(old) must precede insert(new) in the redo run"
);
let bytes = encode_redo_log(&log);
let decoded = decode_redo_log(&bytes).unwrap();
assert_eq!(
decoded, log,
"UPDATE tombstone+insert redo must round-trip the codec"
);
let unresolved_before = crate::unresolved_tombstone_count();
let mut rep = fresh();
rep.apply_redo(&decoded).unwrap();
assert_eq!(
crate::unresolved_tombstone_count(),
unresolved_before,
"the UPDATE's tombstone must resolve by RowId within the replay run"
);
let t = rep.get("t").unwrap();
assert_eq!(
t.rows().len(),
4,
"in-place UPDATE keeps the old row + appends the new"
);
let mut visible: Vec<i32> = t
.scan_visible(&snap)
.filter_map(|(_, r)| match r.values.first() {
Some(Value::Int(v)) => Some(*v),
_ => None,
})
.collect();
visible.sort_unstable();
assert_eq!(
visible,
alloc::vec![1, 3, NEW_VAL],
"old version hidden, new version visible, survivors intact"
);
assert!(
!visible.contains(&OLD_VAL),
"the superseded (old) value must NOT be visible after replay"
);
let old_idx = t
.rows()
.iter()
.position(|r| r.values.first() == Some(&Value::Int(OLD_VAL)))
.expect("old version physically present (tombstone keeps the slot)");
let old_h = t.headers().get(old_idx).expect("header lock-step");
assert_eq!(
old_h.xmax, STMT_V,
"superseded row must carry the tombstone xmax"
);
assert!(old_h.is_deleted(), "superseded row must read as deleted");
let new_idx = t
.rows()
.iter()
.position(|r| r.values.first() == Some(&Value::Int(NEW_VAL)))
.expect("new version physically present");
let new_h = t.headers().get(new_idx).expect("header lock-step");
assert_eq!(new_h.xmax, XMAX_ALIVE, "new version must stay alive");
assert!(!new_h.is_deleted(), "new version must not read as deleted");
for (i, r) in t.rows().iter().enumerate() {
match r.values.first() {
Some(&Value::Int(1)) | Some(&Value::Int(3)) => assert_eq!(
t.headers().get(i).unwrap().xmax,
XMAX_ALIVE,
"survivor {i} must stay alive"
),
_ => {}
}
}
let mut capc = fresh();
capc.enable_redo_all();
{
let t = capc.get_mut("t").unwrap();
t.insert(mk(1)).unwrap();
t.insert(mk(OLD_VAL)).unwrap();
t.insert(mk(3)).unwrap();
t.update_row(1, alloc::vec![Value::Int(NEW_VAL)]).unwrap(); }
let logc = capc.drain_redo();
assert!(
logc.iter()
.all(|c| !matches!(c, RowChange::Tombstone { .. })),
"gate-off UPDATE must NOT emit a Tombstone redo"
);
assert!(
logc.iter().any(|c| matches!(c, RowChange::Update { .. })),
"gate-off UPDATE emits an in-place Update redo"
);
let mut repc = fresh();
repc.apply_redo(&decode_redo_log(&encode_redo_log(&logc)).unwrap())
.unwrap();
let tc = repc.get("t").unwrap();
assert_eq!(
tc.rows().len(),
3,
"gate-off UPDATE replays in place (no extra row)"
);
let mut ids: Vec<i32> = tc
.rows()
.iter()
.filter_map(|r| match r.values.first() {
Some(Value::Int(v)) => Some(*v),
_ => None,
})
.collect();
ids.sort_unstable();
assert_eq!(
ids,
alloc::vec![1, 3, NEW_VAL],
"gate-off replay updates the row in place to the new value"
);
}
#[cfg(test)]
fn mvcc_appendix_bytes(t: &Table) -> Vec<u8> {
let mut a = Vec::new();
a.extend_from_slice(&(t.rows().len() as u32).to_le_bytes());
for (h, rid) in t.headers().iter().zip(t.rowids().iter()) {
a.extend_from_slice(&h.xmin.to_le_bytes());
a.extend_from_slice(&h.xmax.to_le_bytes());
a.push(h.flags);
a.extend_from_slice(&rid.0.to_le_bytes());
}
a.extend_from_slice(&t.next_rowid_for_test().to_le_bytes());
a
}
#[test]
fn v52_snapshot_without_mvcc_appendix_loads_frozen_and_dense() {
use crate::row_header::{RowHeader, RowId, XMAX_ALIVE, XMIN_FROZEN};
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
{
let t = c.get_mut("t").unwrap();
t.insert(Row::new(alloc::vec![Value::Int(0x1111)])).unwrap();
t.insert(Row::new(alloc::vec![Value::Int(0x2222)])).unwrap();
t.insert(Row::new(alloc::vec![Value::Int(0x3333)])).unwrap();
}
const EMPTY_COMMENT_BLOCK: usize = 4;
let v53 = {
let mut full = c.serialize();
const EMPTY_NONTABLE_ACL_BLOCK: usize = 4 + 2 + 2 + 4;
const EMPTY_RULE_BLOCK: usize = 4;
const EMPTY_STATS_EXT_BLOCK: usize = 4;
const EMPTY_LARGE_OBJECT_BLOCK: usize = 4;
const EMPTY_FUNCTION_ATTR_BLOCK: usize = 4;
const EMPTY_DB_ROLE_SETTING_BLOCK: usize = 4;
const EMPTY_REPLICATION_SLOT_BLOCK: usize = 4;
full.truncate(
full.len()
- 4
- EMPTY_COMMENT_BLOCK
- EMPTY_NONTABLE_ACL_BLOCK
- EMPTY_RULE_BLOCK
- EMPTY_STATS_EXT_BLOCK
- EMPTY_LARGE_OBJECT_BLOCK
- EMPTY_FUNCTION_ATTR_BLOCK
- EMPTY_DB_ROLE_SETTING_BLOCK
- EMPTY_REPLICATION_SLOT_BLOCK,
);
full
};
let appendix = mvcc_appendix_bytes(c.get("t").unwrap());
let hits: Vec<usize> = v53
.windows(appendix.len())
.enumerate()
.filter_map(|(i, w)| {
if w == appendix.as_slice() {
Some(i)
} else {
None
}
})
.collect();
assert_eq!(
hits.len(),
1,
"MVCC appendix must appear exactly once in the image"
);
let start = hits[0];
const EMPTY_POLICY_APPENDIX: usize = 4;
const EMPTY_DEFAULT_TEXT_APPENDIX: usize = 2;
let trailing_v53plus = EMPTY_DEFAULT_TEXT_APPENDIX + EMPTY_POLICY_APPENDIX;
const EMPTY_CONSTRAINT_NAME_APPENDIX: usize = 4;
const EMPTY_COMPOSITE_APPENDIX: usize = 2;
const EMPTY_OWNER_ACL_APPENDIX: usize = 3;
const EMPTY_COLUMN_ACL_APPENDIX: usize = 2;
const EMPTY_EXCLUSION_APPENDIX: usize = 2;
const EMPTY_RESTART_APPENDIX: usize = 2;
const EMPTY_MYSQL_INT_WIDTH_APPENDIX: usize = 2;
const EMPTY_MYSQL_FSP_APPENDIX: usize = 2;
const EMPTY_CHECK_VALIDATED_APPENDIX: usize = 2;
const EMPTY_COLLATION_APPENDIX: usize = 2;
const EMPTY_UNIQUE_TIMING_APPENDIX: usize = 2;
let tail_v60plus = EMPTY_CONSTRAINT_NAME_APPENDIX
+ EMPTY_COMPOSITE_APPENDIX
+ EMPTY_OWNER_ACL_APPENDIX
+ EMPTY_COLUMN_ACL_APPENDIX
+ EMPTY_EXCLUSION_APPENDIX
+ EMPTY_RESTART_APPENDIX
+ EMPTY_MYSQL_INT_WIDTH_APPENDIX
+ EMPTY_MYSQL_FSP_APPENDIX
+ EMPTY_CHECK_VALIDATED_APPENDIX
+ EMPTY_COLLATION_APPENDIX
+ EMPTY_UNIQUE_TIMING_APPENDIX;
let mut v52 = Vec::with_capacity(v53.len() - appendix.len() - trailing_v53plus - tail_v60plus);
v52.extend_from_slice(&v53[..start - trailing_v53plus]);
v52.extend_from_slice(&v53[start + appendix.len() + tail_v60plus..]);
v52[FILE_MAGIC.len()] = 52;
let restored = Catalog::deserialize(&v52).expect("v52 image must still load");
let t = restored.get("t").unwrap();
assert_eq!(t.rows().len(), 3, "rows survive the v52 load");
for i in 0..t.rows().len() {
let h = *t.headers().get(i).unwrap();
assert_eq!(h, RowHeader::frozen(), "v52 row {i} must load frozen");
assert_eq!(h.xmin, XMIN_FROZEN);
assert_eq!(h.xmax, XMAX_ALIVE);
}
let ids: Vec<RowId> = t.rowids().iter().copied().collect();
assert_eq!(
ids,
alloc::vec![RowId(1), RowId(2), RowId(3)],
"v52 dense rowids"
);
assert_eq!(t.next_rowid_for_test(), 4, "v52 next_rowid = N + 1");
}
#[test]
fn truncated_mvcc_appendix_errors_cleanly() {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
{
let t = c.get_mut("t").unwrap();
t.insert(Row::new(alloc::vec![Value::Int(7)])).unwrap();
t.insert(Row::new(alloc::vec![Value::Int(8)])).unwrap();
}
let full = c.serialize();
let truncated = &full[..full.len() - 5];
let err = Catalog::deserialize(truncated);
assert!(err.is_err(), "a truncated snapshot must error, not panic");
}
#[test]
fn v68_overloads_are_separate_functions_with_separate_acls() {
use crate::{AclItem, FunctionDef, priv_bits};
let mk = |args: &str, body: &str, grantee: &str| FunctionDef {
name: "f".into(),
args_repr: args.into(),
returns: "TEXT".into(),
language: "sql".into(),
body: body.into(),
owner: Some("alice".into()),
volatility: crate::FN_VOLATILE,
strict: false,
security_definer: false,
leakproof: false,
parallel: crate::FN_PARALLEL_UNSAFE,
cost: None,
rows: None,
acl: alloc::vec![AclItem {
grantee: grantee.into(),
privs: priv_bits::EXECUTE,
grantable: 0,
grantor: "alice".into(),
}],
};
let mut c = Catalog::new();
c.create_function(mk("(x INT)", "SELECT 'int'", "bob"), false)
.unwrap();
c.create_function(mk("(x TEXT)", "SELECT 'text'", "eve"), false)
.unwrap();
assert_eq!(c.functions_named("f").len(), 2, "two overloads coexist");
let restored = Catalog::deserialize(&c.serialize()).expect("v68 image loads");
assert_eq!(restored.functions_named("f").len(), 2);
let by_alias = restored
.function_by_key(&crate::function_signature_key("f", "(x integer)"))
.expect("integer folds to int");
assert_eq!(by_alias.body, "SELECT 'int'");
assert_eq!(
by_alias.acl[0].grantee, "bob",
"each overload keeps its OWN acl"
);
let text_one = restored
.function_by_key(&crate::function_signature_key("f", "(TEXT)"))
.expect("bare type, no arg name");
assert_eq!(text_one.acl[0].grantee, "eve");
}
#[test]
fn v67_roundtrip_preserves_function_owner_and_acl() {
use crate::{AclItem, FunctionDef, priv_bits};
let mut c = Catalog::new();
c.create_function(
FunctionDef {
name: "f1".into(),
args_repr: "(x INT)".into(),
returns: "INT".into(),
language: "sql".into(),
body: "SELECT x + 1".into(),
owner: Some("alice".into()),
volatility: crate::FN_VOLATILE,
strict: false,
security_definer: false,
leakproof: false,
parallel: crate::FN_PARALLEL_UNSAFE,
cost: None,
rows: None,
acl: alloc::vec![AclItem {
grantee: "fred".into(),
privs: priv_bits::EXECUTE,
grantable: 0,
grantor: "alice".into(),
}],
},
false,
)
.unwrap();
let restored = Catalog::deserialize(&c.serialize()).expect("v67 image loads");
let f = restored
.function_by_key(&crate::function_signature_key("f1", "(x INT)"))
.unwrap();
assert_eq!(f.owner.as_deref(), Some("alice"));
assert_eq!(f.acl.len(), 1);
assert_eq!(f.acl[0].grantee, "fred");
assert_eq!(f.acl[0].privs, priv_bits::EXECUTE);
}
#[test]
fn v66_roundtrip_preserves_sequence_schema_and_database_acls() {
use crate::{AclItem, SequenceDataType, SequenceDef, priv_bits};
let mut c = Catalog::new();
c.create_sequence(
SequenceDef {
name: "sq".into(),
data_type: SequenceDataType::BigInt,
start: 1,
increment: 1,
min_value: 1,
max_value: i64::MAX,
cache: 1,
cycle: false,
owned_by: None,
last_value: 0,
is_called: false,
owner: Some("alice".into()),
acl: alloc::vec![AclItem {
grantee: "eve".into(),
privs: priv_bits::USAGE,
grantable: 0,
grantor: "alice".into(),
}],
},
false,
)
.unwrap();
c.schema_acl_mut().push(AclItem {
grantee: String::new(),
privs: priv_bits::USAGE,
grantable: 0,
grantor: "pg_database_owner".into(),
});
c.database_acl_mut().push(AclItem {
grantee: "eve".into(),
privs: priv_bits::CONNECT | priv_bits::TEMPORARY,
grantable: 0,
grantor: "admin".into(),
});
let restored = Catalog::deserialize(&c.serialize()).expect("v66 image loads");
let sq = restored.sequence("sq").unwrap();
assert_eq!(sq.owner.as_deref(), Some("alice"));
assert_eq!(sq.acl.len(), 1);
assert_eq!(sq.acl[0].privs, priv_bits::USAGE);
assert_eq!(restored.schema_acl().len(), 1);
assert_eq!(restored.schema_acl()[0].grantee, "", "PUBLIC");
assert_eq!(restored.database_acl()[0].grantee, "eve");
}
#[test]
fn v64_roundtrip_preserves_owner_and_acl() {
use crate::{AclItem, priv_bits};
let mut c = Catalog::new();
let mut schema = TableSchema::new("ac", vec![ColumnSchema::new("id", DataType::Int, false)]);
schema.owner = Some("alice".into());
schema.acl = alloc::vec![
AclItem {
grantee: "alice".into(),
privs: priv_bits::ALL,
grantable: 0,
grantor: "alice".into(),
},
AclItem {
grantee: "bob".into(),
privs: priv_bits::SELECT | priv_bits::UPDATE,
grantable: priv_bits::SELECT,
grantor: "alice".into(),
},
AclItem {
grantee: String::new(),
privs: priv_bits::SELECT,
grantable: 0,
grantor: "alice".into(),
},
];
c.create_table(schema).unwrap();
let restored = Catalog::deserialize(&c.serialize()).expect("v64 image loads");
let sc = restored.get("ac").unwrap().schema();
assert_eq!(sc.owner.as_deref(), Some("alice"));
assert_eq!(sc.acl.len(), 3);
assert_eq!(sc.acl[1].grantee, "bob");
assert_eq!(sc.acl[1].privs, priv_bits::SELECT | priv_bits::UPDATE);
assert_eq!(sc.acl[1].grantable, priv_bits::SELECT);
assert_eq!(sc.acl[1].grantor, "alice");
assert_eq!(sc.acl[2].grantee, "", "PUBLIC rides the empty grantee");
}
#[test]
fn v63_roundtrip_preserves_user_composite_type_binding() {
let mut c = Catalog::new();
let mut p = ColumnSchema::new("p", DataType::Jsonb, true);
p.user_composite_type = Some("pt".into());
c.create_table(TableSchema::new(
"cp",
vec![ColumnSchema::new("id", DataType::Int, false), p],
))
.unwrap();
let restored = Catalog::deserialize(&c.serialize()).expect("v63 image loads");
let cols = &restored.get("cp").unwrap().schema().columns;
assert_eq!(cols[1].user_composite_type.as_deref(), Some("pt"));
assert_eq!(cols[0].user_composite_type, None);
}
#[test]
fn v53_roundtrip_preserves_mixed_headers_and_rowids() {
use crate::row_header::{HEAP_XMIN_FROZEN, RowHeader, RowId, XMAX_ALIVE};
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
{
let t = c.get_mut("t").unwrap();
for v in 0..6 {
t.insert(Row::new(alloc::vec![Value::Int(100 + v)]))
.unwrap();
}
t.delete_rows(&[1, 3, 5]);
assert_eq!(t.rows().len(), 3);
assert_eq!(
t.rowids().iter().copied().collect::<Vec<_>>(),
alloc::vec![RowId(1), RowId(3), RowId(5)],
"survivors keep their real (non-dense) rowids"
);
assert_eq!(
t.next_rowid_for_test(),
7,
"next_rowid unaffected by delete"
);
let headers = t.headers_mut_for_test();
*headers.get_mut(0).unwrap() = RowHeader::frozen();
*headers.get_mut(1).unwrap() = RowHeader {
xmin: 5,
xmax: 99,
flags: 0,
};
*headers.get_mut(2).unwrap() = RowHeader {
xmin: 42,
xmax: XMAX_ALIVE,
flags: HEAP_XMIN_FROZEN,
};
}
let want_headers: Vec<RowHeader> = c.get("t").unwrap().headers().iter().copied().collect();
let bytes = c.serialize();
let restored = Catalog::deserialize(&bytes).expect("v53 image loads");
let t = restored.get("t").unwrap();
assert_eq!(
t.rowids().iter().copied().collect::<Vec<_>>(),
alloc::vec![RowId(1), RowId(3), RowId(5)],
"rowids must round-trip verbatim"
);
let got_headers: Vec<RowHeader> = t.headers().iter().copied().collect();
assert_eq!(
got_headers, want_headers,
"headers must round-trip verbatim"
);
assert!(
got_headers[1].is_deleted(),
"the tombstoned row stays tombstoned"
);
assert_eq!(got_headers[1].xmax, 99, "tombstone xmax preserved");
assert_eq!(t.next_rowid_for_test(), 7, "next_rowid restored verbatim");
assert!(
t.next_rowid_for_test() > 5,
"next_rowid must exceed every loaded id"
);
}
#[test]
fn rls_policies_round_trip() {
let mut c = Catalog::new();
let mut schema = TableSchema::new("d", vec![ColumnSchema::new("id", DataType::Int, false)]);
schema.row_security = true;
schema.force_row_security = true;
schema.policies.push(crate::PolicyDef {
name: "p_sel".into(),
cmd: crate::PolicyCmd::Select,
permissive: true,
roles: alloc::vec::Vec::new(),
using_expr: Some("(id > 5)".into()),
with_check_expr: None,
});
schema.policies.push(crate::PolicyDef {
name: "p_upd".into(),
cmd: crate::PolicyCmd::Update,
permissive: false,
roles: alloc::vec![String::from("alice"), String::from("bob")],
using_expr: Some("(id > 0)".into()),
with_check_expr: Some("(id > 10)".into()),
});
c.create_table(schema).unwrap();
let restored = Catalog::deserialize(&c.serialize()).expect("v59 image loads");
let s = restored.get("d").unwrap().schema();
assert!(s.row_security && s.force_row_security, "RLS flags survive");
assert_eq!(s.policies.len(), 2);
assert_eq!(s.policies[0].name, "p_sel");
assert_eq!(s.policies[0].cmd, crate::PolicyCmd::Select);
assert!(s.policies[0].permissive);
assert!(s.policies[0].roles.is_empty());
assert_eq!(s.policies[0].using_expr.as_deref(), Some("(id > 5)"));
assert_eq!(s.policies[0].with_check_expr, None);
assert_eq!(s.policies[1].name, "p_upd");
assert_eq!(s.policies[1].cmd, crate::PolicyCmd::Update);
assert!(!s.policies[1].permissive);
assert_eq!(s.policies[1].roles, alloc::vec!["alice", "bob"]);
assert_eq!(s.policies[1].with_check_expr.as_deref(), Some("(id > 10)"));
}
#[test]
fn restored_rows_are_visible_to_a_fresh_snapshot() {
use crate::row_header::{self, RowHeader, XMAX_ALIVE};
use crate::snapshot::{InProgressSet, Snapshot};
const RESTORED_XMIN: u64 = 9_000_000_001;
const RESTORED_XMAX: u64 = 9_000_000_500;
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
{
let t = c.get_mut("t").unwrap();
t.insert(Row::new(vec![Value::Int(1)])).unwrap();
t.insert(Row::new(vec![Value::Int(2)])).unwrap();
let headers = t.headers_mut_for_test();
*headers.get_mut(0).unwrap() = RowHeader {
xmin: RESTORED_XMIN,
xmax: XMAX_ALIVE,
flags: 0,
};
*headers.get_mut(1).unwrap() = RowHeader {
xmin: RESTORED_XMIN,
xmax: RESTORED_XMAX,
flags: 0,
};
}
let restored = Catalog::deserialize(&c.serialize()).expect("image loads");
let version = row_header::current_version();
assert!(
version > RESTORED_XMAX,
"cursor {version} must exceed every restored version ({RESTORED_XMAX})"
);
let snap = Snapshot::new(version, InProgressSet::empty(), version, 0);
let headers: Vec<RowHeader> = restored
.get("t")
.unwrap()
.headers()
.iter()
.copied()
.collect();
assert!(
snap.visible(&headers[0]),
"a committed restored row must not read as a future write"
);
assert!(
!snap.visible(&headers[1]),
"a restored tombstone must stay deleted, not resurrect"
);
}
#[test]
fn cross_checkpoint_tombstone_resolves_after_snapshot_restore() {
use crate::row_header::RowId;
use crate::snapshot::Snapshot;
const TOMB_XMAX: u64 = 77;
let snap = Snapshot::unbounded();
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
{
let t = c.get_mut("t").unwrap();
for v in [10, 20, 30, 40] {
t.insert(Row::new(alloc::vec![Value::Int(v)])).unwrap();
}
t.delete_rows(&[0, 1, 2]); assert_eq!(t.rows().len(), 1);
assert_eq!(
t.rowids().iter().copied().collect::<Vec<_>>(),
alloc::vec![RowId(4)],
"surviving row carries the pre-checkpoint RowId(4)"
);
}
let base = c.serialize();
c.enable_redo_all();
{
let t = c.get_mut("t").unwrap();
t.mark_row_deleted(0, TOMB_XMAX).unwrap();
}
let log = c.drain_redo();
let tomb_targets: Vec<_> = log
.iter()
.filter_map(|ch| match ch {
RowChange::Tombstone { rowids, xmax, .. } => Some((rowids.clone(), *xmax)),
_ => None,
})
.collect();
assert_eq!(tomb_targets.len(), 1, "one in-place tombstone captured");
assert_eq!(
tomb_targets[0].0,
alloc::vec![RowId(4)],
"tombstone names the pre-checkpoint id"
);
let decoded = decode_redo_log(&encode_redo_log(&log)).unwrap();
let mut restored = Catalog::deserialize(&base).expect("base snapshot loads");
assert_eq!(
restored
.get("t")
.unwrap()
.rowids()
.iter()
.copied()
.collect::<Vec<_>>(),
alloc::vec![RowId(4)],
"v53 restore preserves RowId(4); pre-v53 would have dense-assigned RowId(1)"
);
let unresolved_before = crate::unresolved_tombstone_count();
restored.apply_redo(&decoded).unwrap();
assert_eq!(
crate::unresolved_tombstone_count(),
unresolved_before,
"the cross-checkpoint tombstone must resolve by RowId (0 unresolved)"
);
let t = restored.get("t").unwrap();
assert_eq!(t.rows().len(), 1, "tombstone keeps the physical row");
let h = *t.headers().get(0).unwrap();
assert_eq!(
h.xmax, TOMB_XMAX,
"recovered row carries the tombstone xmax"
);
assert!(h.is_deleted(), "recovered row reads as deleted");
let visible: Vec<i32> = t
.scan_visible(&snap)
.filter_map(|(_, r)| match r.values.first() {
Some(Value::Int(v)) => Some(*v),
_ => None,
})
.collect();
assert!(
visible.is_empty(),
"the tombstoned survivor must be hidden after restore"
);
}
#[test]
fn snapshot_round_trips_large_bytea_and_text_array_element() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"q",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("data", DataType::Bytes, true),
ColumnSchema::new("uris", DataType::TextArray, true),
],
))
.unwrap();
let big_blob = alloc::vec![0xAB_u8; 200_000];
let big_elem = "u".repeat(100_000);
cat.get_mut("q")
.unwrap()
.insert(Row::new(alloc::vec![
Value::BigInt(1),
Value::bytes(big_blob.clone()),
Value::TextArray(alloc::vec![Some(big_elem.clone()), None, Some("s".into())]),
]))
.unwrap();
let bytes = cat.serialize();
let re = Catalog::deserialize(&bytes).unwrap();
let row = re.get("q").unwrap().rows.get(0).unwrap().clone();
match &row.values[1] {
Value::Bytes(b) => assert_eq!(b.len(), big_blob.len()),
other => panic!("expected Bytes, got {other:?}"),
}
match &row.values[2] {
Value::TextArray(items) => {
assert_eq!(items[0].as_ref().unwrap().len(), big_elem.len());
assert!(items[1].is_none());
}
other => panic!("expected TextArray, got {other:?}"),
}
}
#[test]
fn plain_u16_bytea_len_ffff_decodes_under_v46_rules() {
let payload = alloc::vec![7_u8; 65_535];
let mut buf = Vec::new();
write_u16(&mut buf, 65_535);
buf.extend_from_slice(&payload);
let mut cur = Cursor::new(&buf).with_codec_version(46);
let len = cur.read_len_escaped_v47().unwrap();
assert_eq!(len, 65_535);
assert_eq!(cur.take(len).unwrap().len(), 65_535);
}
#[test]
fn escaped_string_codec_round_trips_large_text() {
for len in [0usize, 1, 65_534, 65_535, 65_536, 1_048_576] {
let s: String = "x".repeat(len);
let mut buf = Vec::new();
write_str(&mut buf, &s);
let expected_header = if len >= STR_LEN_ESCAPE as usize { 6 } else { 2 };
assert_eq!(buf.len(), expected_header + len, "header width for {len}");
let mut cur = Cursor::new(&buf).with_codec_version(FILE_VERSION);
assert_eq!(cur.read_str().unwrap().len(), len, "round-trip {len}");
}
}
#[test]
fn plain_u16_len_ffff_decodes_under_old_rules() {
let s = "y".repeat(65_535);
let mut buf = Vec::new();
write_u16(&mut buf, 65_535);
buf.extend_from_slice(s.as_bytes());
let mut old = Cursor::new(&buf); assert_eq!(old.read_str().unwrap(), s);
}
#[test]
fn snapshot_round_trips_megabyte_text_row() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"mail",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("body", DataType::Text, false),
],
))
.unwrap();
let body = "m".repeat(1_048_576);
cat.get_mut("mail")
.unwrap()
.insert(Row::new(vec![Value::BigInt(1), Value::text(body.clone())]))
.unwrap();
let bytes = cat.serialize();
let re = Catalog::deserialize(&bytes).unwrap();
let t = re.get("mail").unwrap();
match &t.rows.get(0).unwrap().values[1] {
Value::Text(s) => assert_eq!(s.len(), body.len()),
other => panic!("expected Text, got {other:?}"),
}
}
#[test]
fn segment_v3_round_trips_large_text_rows() {
let schema = TableSchema::new(
"mail",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("body", DataType::Text, false),
],
);
let big = "b".repeat(200_000);
let rows: Vec<(u64, Vec<u8>)> = (0u64..3)
.map(|i| {
let row = Row::new(vec![
Value::BigInt(i.cast_signed()),
Value::text(big.clone()),
]);
(i, encode_row_body_dense(&row, &schema))
})
.collect();
let (bytes, _) = encode_segment(rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES).unwrap();
assert_eq!(&bytes[..8], b"SPGSEG\x04\x00", "new segments are V4");
let seg = OwnedSegment::from_bytes(bytes).unwrap();
assert!(seg.codec_version() >= 47);
let payload = seg.lookup(1).expect("pk 1 present");
let (row, _) = decode_row_body_dense(&payload, &schema, seg.codec_version()).unwrap();
match &row.values[1] {
Value::Text(s) => assert_eq!(s.len(), big.len()),
other => panic!("expected Text, got {other:?}"),
}
}
#[test]
fn index_key_round_trips_large_text() {
let key = IndexKey::Text("k".repeat(100_000));
let mut buf = Vec::new();
write_index_key(&mut buf, &key);
let mut cur = Cursor::new(&buf).with_codec_version(FILE_VERSION);
let back = cur.read_index_key().unwrap();
assert_eq!(back, key);
}
#[cfg(target_arch = "aarch64")]
#[test]
fn neon_l2_matches_scalar() {
let dims = [4usize, 8, 12, 16, 64, 128, 256, 384, 512, 768, 1024, 1536];
for &d in &dims {
let mut state: u64 = (d as u64).wrapping_mul(0x9E37_79B9_7F4A_7C15);
let mut a = Vec::with_capacity(d);
let mut b = Vec::with_capacity(d);
for _ in 0..d {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let x = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let y = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
a.push(x);
b.push(y);
}
let scalar = l2_distance_sq_scalar(&a, &b);
let neon = unsafe { l2_distance_sq_neon(&a, &b) };
let tol = (scalar.abs().max(1e-6)) * 1e-4;
assert!(
(scalar - neon).abs() <= tol,
"dim={d}: scalar={scalar} neon={neon} diff={}",
(scalar - neon).abs()
);
}
}
#[cfg(target_arch = "aarch64")]
#[test]
fn neon_inner_product_matches_scalar() {
let dims = [4usize, 8, 12, 16, 64, 128, 256, 512, 1024];
for &d in &dims {
let mut state: u64 = (d as u64).wrapping_mul(0x9E37_79B9_7F4A_7C15);
let mut a = Vec::with_capacity(d);
let mut b = Vec::with_capacity(d);
for _ in 0..d {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let x = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let y = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
a.push(x);
b.push(y);
}
let scalar = inner_product_scalar(&a, &b);
let neon = unsafe { inner_product_neon(&a, &b) };
#[allow(clippy::cast_precision_loss)]
let tol = (scalar.abs().max(1e-6)) * 1e-4 + (d as f32) * 1e-6;
assert!(
(scalar - neon).abs() <= tol,
"IP dim={d}: scalar={scalar} neon={neon} diff={}",
(scalar - neon).abs()
);
}
}
#[cfg(target_arch = "aarch64")]
#[allow(clippy::similar_names)]
#[test]
fn neon_cosine_dot_norms_matches_scalar() {
let dims = [4usize, 8, 12, 16, 64, 128, 256, 512, 1024];
for &d in &dims {
let mut state: u64 = (d as u64).wrapping_mul(0xBF58_476D_1CE4_E5B9);
let mut a = Vec::with_capacity(d);
let mut b = Vec::with_capacity(d);
for _ in 0..d {
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let x = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
state = state
.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1);
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
let y = (((state >> 32) & 0x00FF_FFFF) as f32) / (0x80_0000_u32 as f32) - 1.0;
a.push(x);
b.push(y);
}
let (dot_s, na_s, nb_s) = cosine_dot_norms_scalar(&a, &b);
let (dot_n, na_n, nb_n) = unsafe { cosine_dot_norms_neon(&a, &b) };
#[allow(clippy::cast_precision_loss)]
let tol_d = (dot_s.abs().max(1e-6)) * 1e-4 + (d as f32) * 1e-6;
#[allow(clippy::cast_precision_loss)]
let tol_n = (na_s.abs().max(1e-6)) * 1e-4 + (d as f32) * 1e-6;
assert!(
(dot_s - dot_n).abs() <= tol_d,
"cosine dot dim={d}: scalar={dot_s} neon={dot_n}"
);
assert!(
(na_s - na_n).abs() <= tol_n,
"cosine na dim={d}: scalar={na_s} neon={na_n}"
);
assert!(
(nb_s - nb_n).abs() <= tol_n,
"cosine nb dim={d}: scalar={nb_s} neon={nb_n}"
);
}
}
fn make_users_schema() -> TableSchema {
TableSchema::new(
"users",
vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("name", DataType::Text, false),
ColumnSchema::new("score", DataType::Float, true),
],
)
}
#[test]
fn value_type_tag_matches_variant() {
assert_eq!(Value::Int(1).data_type(), Some(DataType::Int));
assert_eq!(Value::BigInt(1).data_type(), Some(DataType::BigInt));
assert_eq!(Value::Float(1.0).data_type(), Some(DataType::Float));
assert_eq!(Value::text("x").data_type(), Some(DataType::Text));
assert_eq!(Value::Bool(true).data_type(), Some(DataType::Bool));
assert_eq!(Value::Null.data_type(), None);
assert!(Value::Null.is_null());
assert!(!Value::Int(0).is_null());
}
#[test]
fn sq8_value_reports_sq8_data_type() {
let q = crate::quantize::quantize(&[0.0, 0.25, 0.5, 0.75, 1.0]);
let v = Value::Sq8Vector(q);
assert_eq!(
v.data_type(),
Some(DataType::Vector {
dim: 5,
encoding: VecEncoding::Sq8,
}),
);
}
#[test]
fn datatype_display_matches_pg_keyword() {
assert_eq!(DataType::Int.to_string(), "INT");
assert_eq!(DataType::BigInt.to_string(), "BIGINT");
assert_eq!(DataType::Float.to_string(), "FLOAT");
assert_eq!(DataType::Text.to_string(), "TEXT");
assert_eq!(DataType::Bool.to_string(), "BOOL");
}
#[test]
fn row_len_and_emptiness() {
let r = Row::new(vec![Value::Int(1), Value::Null]);
assert_eq!(r.len(), 2);
assert!(!r.is_empty());
assert!(Row::new(Vec::new()).is_empty());
}
#[test]
fn table_schema_column_position() {
let s = make_users_schema();
assert_eq!(s.column_position("id"), Some(0));
assert_eq!(s.column_position("score"), Some(2));
assert_eq!(s.column_position("missing"), None);
}
#[test]
fn catalog_create_table_then_lookup() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
assert_eq!(cat.table_count(), 1);
assert!(cat.get("users").is_some());
assert!(cat.get("nope").is_none());
}
#[test]
fn catalog_duplicate_table_is_rejected() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let err = cat.create_table(make_users_schema()).unwrap_err();
assert!(matches!(err, StorageError::DuplicateTable { ref name } if name == "users"));
}
#[test]
fn table_insert_happy_path_appends_row() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(Row::new(vec![
Value::Int(1),
Value::text("alice"),
Value::Float(99.5),
]))
.unwrap();
assert_eq!(t.row_count(), 1);
assert_eq!(t.rows()[0].values[1], Value::text("alice"));
}
#[test]
fn rowid_monotonic_survives_delete_and_never_reused() {
use crate::row_header::RowId;
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for i in 0..3 {
t.insert(Row::new(vec![
Value::Int(i),
Value::text("x"),
Value::Float(0.0),
]))
.unwrap();
}
assert_eq!(t.rows().len(), t.rowids().len());
assert_eq!(
t.rowids().iter().copied().collect::<alloc::vec::Vec<_>>(),
alloc::vec![RowId(1), RowId(2), RowId(3)]
);
t.delete_rows(&[1]);
assert_eq!(t.rows().len(), 2);
assert_eq!(t.rows().len(), t.rowids().len());
assert_eq!(
t.rowids().iter().copied().collect::<alloc::vec::Vec<_>>(),
alloc::vec![RowId(1), RowId(3)]
);
t.insert(Row::new(vec![
Value::Int(9),
Value::text("y"),
Value::Float(0.0),
]))
.unwrap();
assert_eq!(
t.rowids().iter().copied().collect::<alloc::vec::Vec<_>>(),
alloc::vec![RowId(1), RowId(3), RowId(4)]
);
t.truncate();
assert_eq!(t.rowids().len(), 0);
t.insert(Row::new(vec![
Value::Int(0),
Value::text("z"),
Value::Float(0.0),
]))
.unwrap();
assert_eq!(
t.rowids().iter().copied().collect::<alloc::vec::Vec<_>>(),
alloc::vec![RowId(5)]
);
}
#[test]
fn relid_assigned_monotonic_by_create_table() {
use crate::row_header::RelId;
let bare = Table::new(make_users_schema());
assert_eq!(bare.rel_id(), RelId::UNASSIGNED);
let mut cat = Catalog::new();
for n in ["a", "b", "c"] {
let mut s = make_users_schema();
s.name = n.into();
cat.create_table(s).unwrap();
}
assert_eq!(cat.get("a").unwrap().rel_id(), RelId(1));
assert_eq!(cat.get("b").unwrap().rel_id(), RelId(2));
assert_eq!(cat.get("c").unwrap().rel_id(), RelId(3));
let bytes = cat.serialize();
let restored = Catalog::deserialize(&bytes).unwrap();
assert_eq!(restored.get("a").unwrap().rel_id(), RelId(1));
assert_eq!(restored.get("c").unwrap().rel_id(), RelId(3));
}
#[test]
fn table_insert_arity_mismatch() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
let err = t.insert(Row::new(vec![Value::Int(1)])).unwrap_err();
assert!(matches!(
err,
StorageError::ArityMismatch {
expected: 3,
actual: 1
}
));
assert_eq!(t.row_count(), 0);
}
#[test]
fn table_insert_type_mismatch_reports_column() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
let err = t
.insert(Row::new(vec![
Value::Int(1),
Value::Int(42), Value::Float(0.0),
]))
.unwrap_err();
match err {
StorageError::TypeMismatch {
ref column,
expected,
actual,
position,
} => {
assert_eq!(column, "name");
assert_eq!(expected, DataType::Text);
assert_eq!(actual, DataType::Int);
assert_eq!(position, 1);
}
other => panic!("unexpected: {other:?}"),
}
assert_eq!(t.row_count(), 0);
}
#[test]
fn table_insert_null_into_not_null_rejected() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
let err = t
.insert(Row::new(vec![
Value::Int(1),
Value::Null, Value::Float(1.0),
]))
.unwrap_err();
assert!(matches!(err, StorageError::NullInNotNull { ref column } if column == "name"));
}
#[test]
fn table_insert_null_into_nullable_ok() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(Row::new(vec![
Value::Int(1),
Value::text("bob"),
Value::Null,
]))
.unwrap();
assert_eq!(t.row_count(), 1);
}
#[test]
fn catalog_get_mut_independent_per_table() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"a",
vec![ColumnSchema::new("v", DataType::Int, false)],
))
.unwrap();
cat.create_table(TableSchema::new(
"b",
vec![ColumnSchema::new("v", DataType::Int, false)],
))
.unwrap();
cat.get_mut("a")
.unwrap()
.insert(Row::new(vec![Value::Int(1)]))
.unwrap();
assert_eq!(cat.get("a").unwrap().row_count(), 1);
assert_eq!(cat.get("b").unwrap().row_count(), 0);
}
fn assert_round_trip(cat: &Catalog) {
let bytes = cat.serialize();
let restored = Catalog::deserialize(&bytes).expect("deserialize");
assert_eq!(restored.table_count(), cat.table_count());
for (a, b) in cat.tables.iter().zip(restored.tables.iter()) {
assert_eq!(a.schema, b.schema);
assert_eq!(a.rows, b.rows);
}
}
#[test]
fn serialize_empty_catalog_round_trips() {
assert_round_trip(&Catalog::new());
}
#[test]
fn serialize_single_empty_table_round_trips() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
assert_round_trip(&cat);
}
#[test]
fn nsw_clone_is_o1() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"docs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32
},
true
),
],
))
.unwrap();
let t = cat.get_mut("docs").unwrap();
for i in 0..1500_i32 {
#[allow(clippy::cast_precision_loss)] let base = (i as f32) * 0.01;
t.insert(Row::new(alloc::vec![
Value::Int(i),
Value::vector(alloc::vec![base, base + 0.05, base + 0.1]),
]))
.unwrap();
}
t.add_nsw_index("docs_nsw".into(), "v", NSW_DEFAULT_M)
.unwrap();
let g = match &cat.get("docs").unwrap().indices()[0].kind {
IndexKind::Nsw(g) => g,
IndexKind::BTree(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => {
panic!("expected NSW")
}
};
assert_eq!(g.levels.len(), 1500, "one level slot per inserted row");
assert!(
g.layers.len() >= 2,
"1500 nodes should populate at least two HNSW layers, got {}",
g.layers.len()
);
let cloned = g.clone();
assert!(
g.levels.shares_storage_with(&cloned.levels),
"levels PV not shared after clone — clone copied elements (O(N))"
);
assert_eq!(g.layers.len(), cloned.layers.len());
for (l, (orig, cl)) in g.layers.iter().zip(cloned.layers.iter()).enumerate() {
assert!(
orig.shares_storage_with(cl),
"layer {l} PV not shared after clone — clone copied elements (O(N))"
);
}
}
#[test]
fn sq8_catalog_serialise_roundtrip_preserves_cells_and_index() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 8,
encoding: VecEncoding::Sq8,
},
false,
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
for i in 0..32_i32 {
#[allow(clippy::cast_precision_loss)]
let base = (i as f32) * 0.03;
let v: Vec<f32> = (0..8_i32)
.map(|j| {
#[allow(clippy::cast_precision_loss)]
let off = (j as f32) * 0.01;
base + off
})
.collect();
t.insert(Row::new(alloc::vec![
Value::Int(i),
Value::Sq8Vector(quantize::quantize(&v)),
]))
.unwrap();
}
t.add_nsw_index("v_idx".into(), "v", NSW_DEFAULT_M).unwrap();
let query = alloc::vec![0.15_f32, 0.16, 0.17, 0.18, 0.19, 0.20, 0.21, 0.22];
let (before_cell, before_ty, before_hits) = {
let t_ref = cat.get("vecs").unwrap();
(
t_ref.rows()[5].values[1].clone(),
t_ref.schema().columns[1].ty,
nsw_query(t_ref, "v_idx", &query, 5, NswMetric::L2),
)
};
let bytes = cat.serialize();
let restored = Catalog::deserialize(&bytes).expect("deserialize ok");
let rt = restored.get("vecs").unwrap();
assert_eq!(rt.schema().columns[1].ty, before_ty);
assert_eq!(rt.rows()[5].values[1], before_cell);
let after_hits = nsw_query(rt, "v_idx", &query, 5, NswMetric::L2);
assert_eq!(before_hits, after_hits);
}
#[test]
fn half_catalog_serialise_roundtrip_preserves_cells_and_index() {
use crate::halfvec;
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 8,
encoding: VecEncoding::F16,
},
false,
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
for i in 0..32_i32 {
#[allow(clippy::cast_precision_loss)]
let base = (i as f32) * 0.03;
let v: Vec<f32> = (0..8_i32)
.map(|j| {
#[allow(clippy::cast_precision_loss)]
let off = (j as f32) * 0.01;
base + off
})
.collect();
t.insert(Row::new(alloc::vec![
Value::Int(i),
Value::HalfVector(halfvec::HalfVector::from_f32_slice(&v)),
]))
.unwrap();
}
t.add_nsw_index("v_idx".into(), "v", NSW_DEFAULT_M).unwrap();
let query = alloc::vec![0.15_f32, 0.16, 0.17, 0.18, 0.19, 0.20, 0.21, 0.22];
let (before_cell, before_ty, before_hits) = {
let t_ref = cat.get("vecs").unwrap();
(
t_ref.rows()[5].values[1].clone(),
t_ref.schema().columns[1].ty,
nsw_query(t_ref, "v_idx", &query, 5, NswMetric::L2),
)
};
let bytes = cat.serialize();
let restored = Catalog::deserialize(&bytes).expect("deserialize ok");
let rt = restored.get("vecs").unwrap();
assert_eq!(rt.schema().columns[1].ty, before_ty);
assert_eq!(rt.rows()[5].values[1], before_cell);
let after_hits = nsw_query(rt, "v_idx", &query, 5, NswMetric::L2);
assert_eq!(before_hits, after_hits);
}
#[test]
#[allow(clippy::similar_names)]
fn hnsw_half_recall_at_10_matches_f32_groundtruth() {
use crate::halfvec;
fn next(state: &mut u64) -> f32 {
*state = state
.wrapping_add(0x9E37_79B9_7F4A_7C15)
.wrapping_mul(0xBF58_476D_1CE4_E5B9);
#[allow(clippy::cast_precision_loss)]
let u = ((*state >> 32) as u32 as f32) / (u32::MAX as f32);
2.0 * u - 1.0
}
let dim: u32 = 32;
let n: usize = 512;
let dim_us = dim as usize;
let mut seed: u64 = 0xF16_F16_F16_F16_u64;
let corpus: Vec<Vec<f32>> = (0..n)
.map(|_| (0..dim_us).map(|_| next(&mut seed)).collect())
.collect();
let queries: Vec<Vec<f32>> = (0..32)
.map(|_| (0..dim_us).map(|_| next(&mut seed)).collect())
.collect();
let exact_top10: Vec<Vec<usize>> = queries
.iter()
.map(|q| {
let mut scored: Vec<(f32, usize)> = corpus
.iter()
.enumerate()
.map(|(i, v)| (l2_distance_sq(v, q), i))
.collect();
scored.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(core::cmp::Ordering::Equal));
scored.into_iter().take(10).map(|(_, i)| i).collect()
})
.collect();
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim,
encoding: VecEncoding::F16,
},
false,
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
for (i, v) in corpus.iter().enumerate() {
t.insert(Row::new(alloc::vec![
Value::Int(i32::try_from(i).unwrap()),
Value::HalfVector(halfvec::HalfVector::from_f32_slice(v)),
]))
.unwrap();
}
t.add_nsw_index("v_idx".into(), "v", NSW_DEFAULT_M).unwrap();
let table = cat.get("vecs").unwrap();
let mut total_overlap = 0_usize;
for (q, exact) in queries.iter().zip(exact_top10.iter()) {
let hits = nsw_query(table, "v_idx", q, 10, NswMetric::L2);
for h in &hits {
if exact.contains(h) {
total_overlap += 1;
}
}
}
#[allow(clippy::cast_precision_loss)]
let recall = total_overlap as f32 / (10.0 * queries.len() as f32);
assert!(
recall >= 0.95,
"HALF HNSW recall@10 = {recall:.3}, below floor 0.95 — \
check halfvec dispatch in `cell_to_query_metric_distance`"
);
}
#[test]
fn hnsw_sq8_recall_at_10_above_0_95_vs_f32_groundtruth() {
use crate::quantize;
fn next(state: &mut u64) -> f32 {
*state = state
.wrapping_add(0x9E37_79B9_7F4A_7C15)
.wrapping_mul(0xBF58_476D_1CE4_E5B9);
#[allow(clippy::cast_precision_loss)]
let u = ((*state >> 32) as u32 as f32) / (u32::MAX as f32);
2.0 * u - 1.0
}
let dim: u32 = 32;
let n: usize = 512;
let dim_us = dim as usize;
let mut seed: u64 = 0xCAFE_BABE_DEAD_BEEFu64;
let corpus: Vec<Vec<f32>> = (0..n)
.map(|_| (0..dim_us).map(|_| next(&mut seed)).collect())
.collect();
let queries: Vec<Vec<f32>> = (0..32)
.map(|_| (0..dim_us).map(|_| next(&mut seed)).collect())
.collect();
let exact_top10: Vec<Vec<usize>> = queries
.iter()
.map(|q| {
let mut scored: Vec<(f32, usize)> = corpus
.iter()
.enumerate()
.map(|(i, v)| (l2_distance_sq(v, q), i))
.collect();
scored.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(core::cmp::Ordering::Equal));
scored.into_iter().take(10).map(|(_, i)| i).collect()
})
.collect();
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim,
encoding: VecEncoding::Sq8,
},
false,
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
for (i, v) in corpus.iter().enumerate() {
t.insert(Row::new(alloc::vec![
Value::Int(i32::try_from(i).unwrap()),
Value::Sq8Vector(quantize::quantize(v)),
]))
.unwrap();
}
t.add_nsw_index("v_idx".into(), "v", NSW_DEFAULT_M).unwrap();
let table = cat.get("vecs").unwrap();
let mut total_overlap = 0_usize;
for (q, exact) in queries.iter().zip(exact_top10.iter()) {
let hits = nsw_query(table, "v_idx", q, 10, NswMetric::L2);
for h in &hits {
if exact.contains(h) {
total_overlap += 1;
}
}
}
#[allow(clippy::cast_precision_loss)]
let recall = total_overlap as f32 / (10.0 * queries.len() as f32);
assert!(
recall >= 0.95,
"SQ8 HNSW recall@10 = {recall:.3}, below floor 0.95 — \
check `sq8_rerank` is wired in `nsw_search` for SQ8 columns"
);
}
#[test]
fn nsw_index_topology_persists_through_round_trip() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"docs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32
},
true
),
],
))
.unwrap();
let t = cat.get_mut("docs").unwrap();
for i in 0..6_i32 {
#[allow(clippy::cast_precision_loss)] let base = (i as f32) * 0.1;
let row = Row::new(alloc::vec![
Value::Int(i),
Value::vector(alloc::vec![base, base + 0.05, base + 0.1]),
]);
t.insert(row).unwrap();
}
t.add_nsw_index("docs_nsw".into(), "v", NSW_DEFAULT_M)
.unwrap();
let original = match &cat.get("docs").unwrap().indices()[0].kind {
IndexKind::Nsw(g) => g.clone(),
IndexKind::BTree(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => {
panic!("expected NSW")
}
};
let bytes = cat.serialize();
let restored = Catalog::deserialize(&bytes).expect("deserialize");
let restored_graph = match &restored.get("docs").unwrap().indices()[0].kind {
IndexKind::Nsw(g) => g.clone(),
IndexKind::BTree(_)
| IndexKind::Brin { .. }
| IndexKind::Gin(_)
| IndexKind::GinTrgm(_)
| IndexKind::GinFulltext(_)
| IndexKind::GinJsonb(_) => {
panic!("expected NSW")
}
};
assert_eq!(restored_graph.m, original.m);
assert_eq!(restored_graph.m_max_0, original.m_max_0);
assert_eq!(restored_graph.entry, original.entry);
assert_eq!(restored_graph.entry_level, original.entry_level);
assert_eq!(restored_graph.levels, original.levels);
assert_eq!(restored_graph.layers, original.layers);
}
#[test]
fn hnsw_level_assignment_is_deterministic() {
for i in 0..32usize {
assert_eq!(nsw_assign_level(i), nsw_assign_level(i));
}
}
#[test]
fn hnsw_layer_0_dominates_population() {
let on_zero = (0..200usize).filter(|&i| nsw_assign_level(i) == 0).count();
assert!(on_zero > 150, "level-0 nodes too few: {on_zero}");
}
#[test]
fn hnsw_search_matches_brute_force_for_l2_top1() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
alloc::vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32
},
true
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
let dataset: alloc::vec::Vec<(i32, [f32; 3])> = alloc::vec![
(1, [0.0, 0.0, 0.0]),
(2, [1.0, 0.0, 0.0]),
(3, [0.0, 1.0, 0.0]),
(4, [0.0, 0.0, 1.0]),
(5, [1.0, 1.0, 0.0]),
(6, [1.0, 0.0, 1.0]),
(7, [0.0, 1.0, 1.0]),
(8, [1.0, 1.0, 1.0]),
(9, [0.5, 0.5, 0.5]),
(10, [0.2, 0.8, 0.5]),
];
for &(id, v) in &dataset {
t.insert(Row::new(alloc::vec![
Value::Int(id),
Value::vector(alloc::vec![v[0], v[1], v[2]]),
]))
.unwrap();
}
t.add_nsw_index("v_idx".into(), "v", NSW_DEFAULT_M).unwrap();
let idx_pos = cat
.get("vecs")
.unwrap()
.indices()
.iter()
.position(|i| i.name == "v_idx")
.unwrap();
for query in [[0.4, 0.4, 0.4], [0.9, 0.1, 0.0], [0.0, 0.9, 0.9]] {
let table = cat.get("vecs").unwrap();
let hnsw_top = nsw_search(table, idx_pos, &query, 1, 16, NswMetric::L2);
let mut brute: alloc::vec::Vec<(f32, usize)> = (0..table.rows.len())
.map(|i| {
let Value::Vector(v) = &table.rows[i].values[1] else {
return (f32::INFINITY, i);
};
(l2_distance_sq(v, &query), i)
})
.collect();
brute.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(core::cmp::Ordering::Equal));
assert!(!hnsw_top.is_empty(), "HNSW returned no results");
assert_eq!(
hnsw_top[0].1, brute[0].1,
"HNSW top-1 != brute-force top-1 for {query:?}"
);
}
}
#[test]
fn serialize_table_with_rows_round_trips() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(Row::new(vec![
Value::Int(1),
Value::text("alice"),
Value::Float(95.5),
]))
.unwrap();
t.insert(Row::new(vec![
Value::Int(2),
Value::text("bob"),
Value::Null,
]))
.unwrap();
assert_round_trip(&cat);
}
#[test]
fn serialize_multiple_tables_round_trips() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
cat.create_table(TableSchema::new(
"flags",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("active", DataType::Bool, false),
],
))
.unwrap();
cat.get_mut("flags")
.unwrap()
.insert(Row::new(vec![Value::BigInt(7), Value::Bool(true)]))
.unwrap();
assert_round_trip(&cat);
}
#[test]
fn deserialize_rejects_bad_magic() {
let mut buf = b"BADMAGIC".to_vec();
buf.push(FILE_VERSION);
buf.extend_from_slice(&0u32.to_le_bytes());
let err = Catalog::deserialize(&buf).unwrap_err();
assert!(matches!(err, StorageError::Corrupt(_)));
}
#[test]
fn deserialize_rejects_unsupported_version() {
let mut buf = FILE_MAGIC.to_vec();
buf.push(99); buf.extend_from_slice(&0u32.to_le_bytes());
let err = Catalog::deserialize(&buf).unwrap_err();
assert!(matches!(err, StorageError::Corrupt(ref s) if s.contains("version")));
}
#[test]
fn deserialize_rejects_truncated_file() {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let bytes = cat.serialize();
let truncated = &bytes[..bytes.len() - 1];
assert!(matches!(
Catalog::deserialize(truncated),
Err(StorageError::Corrupt(_))
));
}
#[test]
fn deserialize_rejects_trailing_garbage() {
let cat = Catalog::new();
let mut bytes = cat.serialize();
bytes.push(0xFF);
assert!(matches!(
Catalog::deserialize(&bytes),
Err(StorageError::Corrupt(ref s)) if s.contains("trailing")
));
}
fn populated_users() -> Catalog {
let mut cat = Catalog::new();
cat.create_table(make_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name, score) in [
(1, "alice", Some(90.0)),
(2, "bob", None),
(3, "alice", Some(70.0)), ] {
t.insert(Row::new(vec![
Value::Int(id),
Value::text(name),
score.map_or(Value::Null, Value::Float),
]))
.unwrap();
}
cat
}
#[test]
fn add_index_builds_from_existing_rows() {
let mut cat = populated_users();
cat.get_mut("users")
.unwrap()
.add_index("by_id".into(), "id")
.unwrap();
let t = cat.get("users").unwrap();
let idx = t.index_on(0).expect("index_on(0)");
assert_eq!(idx.lookup_eq(&IndexKey::Int(2)), &[RowLocator::Hot(1)]);
assert_eq!(idx.lookup_eq(&IndexKey::Int(99)), &[] as &[RowLocator]);
}
#[test]
fn add_index_dup_name_rejected() {
let mut cat = populated_users();
let t = cat.get_mut("users").unwrap();
t.add_index("ix".into(), "id").unwrap();
let err = t.add_index("ix".into(), "name").unwrap_err();
assert!(matches!(err, StorageError::DuplicateIndex { ref name } if name == "ix"));
}
#[test]
fn add_index_unknown_column_rejected() {
let mut cat = populated_users();
let err = cat
.get_mut("users")
.unwrap()
.add_index("ix".into(), "ghost")
.unwrap_err();
assert!(matches!(err, StorageError::ColumnNotFound { ref column } if column == "ghost"));
}
#[test]
fn insert_after_create_index_updates_it() {
let mut cat = populated_users();
let t = cat.get_mut("users").unwrap();
t.add_index("by_name".into(), "name").unwrap();
t.insert(Row::new(vec![
Value::Int(4),
Value::text("dave"),
Value::Null,
]))
.unwrap();
let idx = t.index_on(1).unwrap();
assert_eq!(
idx.lookup_eq(&IndexKey::Text("dave".into())),
&[RowLocator::Hot(3)]
);
assert_eq!(
idx.lookup_eq(&IndexKey::Text("alice".into())),
&[RowLocator::Hot(0), RowLocator::Hot(2)]
);
}
#[test]
fn null_or_float_values_are_not_indexed() {
let mut cat = populated_users();
let t = cat.get_mut("users").unwrap();
t.add_index("by_score".into(), "score").unwrap();
let idx = t.index_on(2).unwrap();
assert_eq!(idx.lookup_eq(&IndexKey::Int(90)), &[] as &[RowLocator]);
}
#[test]
fn vector_value_data_type_carries_dim() {
let v = Value::vector(vec![1.0, 2.0, 3.0]);
assert_eq!(
v.data_type(),
Some(DataType::Vector {
dim: 3,
encoding: VecEncoding::F32
})
);
}
#[test]
fn vector_column_insert_matching_dim_ok() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"emb",
vec![ColumnSchema::new(
"v",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32,
},
false,
)],
))
.unwrap();
cat.get_mut("emb")
.unwrap()
.insert(Row::new(vec![Value::vector(vec![1.0, 2.0, 3.0])]))
.unwrap();
}
#[test]
fn vector_column_insert_dim_mismatch_rejected() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"emb",
vec![ColumnSchema::new(
"v",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32,
},
false,
)],
))
.unwrap();
let err = cat
.get_mut("emb")
.unwrap()
.insert(Row::new(vec![Value::vector(vec![1.0, 2.0])]))
.unwrap_err();
assert!(matches!(err, StorageError::TypeMismatch { .. }));
}
#[test]
fn vector_value_survives_catalog_round_trip() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"emb",
vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 4,
encoding: VecEncoding::F32,
},
false,
),
],
))
.unwrap();
cat.get_mut("emb")
.unwrap()
.insert(Row::new(vec![
Value::Int(1),
Value::vector(vec![0.5, -1.25, 3.0, 7.0]),
]))
.unwrap();
let restored = Catalog::deserialize(&cat.serialize()).expect("round-trip");
let table = restored.get("emb").unwrap();
assert_eq!(
table.schema().columns[1].ty,
DataType::Vector {
dim: 4,
encoding: VecEncoding::F32
}
);
assert_eq!(
table.rows()[0].values[1],
Value::vector(vec![0.5, -1.25, 3.0, 7.0])
);
}
#[test]
fn index_survives_serialize_deserialize_round_trip() {
let mut cat = populated_users();
cat.get_mut("users")
.unwrap()
.add_index("by_name".into(), "name")
.unwrap();
let restored = Catalog::deserialize(&cat.serialize()).unwrap();
let idx = restored
.get("users")
.unwrap()
.index_on(1)
.expect("index_on(1) after restore");
assert_eq!(idx.name, "by_name");
assert_eq!(
idx.lookup_eq(&IndexKey::Text("alice".into())),
&[RowLocator::Hot(0), RowLocator::Hot(2)]
);
}
fn bigint_pk_users_schema() -> TableSchema {
TableSchema::new(
"users",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("name", DataType::Text, false),
],
)
}
fn make_user_row(id: i64, name: &str) -> Row<'static> {
Row::new(vec![Value::BigInt(id), Value::text(name.to_string())])
}
#[test]
fn update_row_non_indexed_column_keeps_index_intact() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob"), (3, "carol")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
t.update_row(1, vec![Value::BigInt(2), Value::text("bobby")])
.unwrap();
let idx = t.index_on(0).unwrap();
assert_eq!(
idx.lookup_eq(&IndexKey::Int(2)),
&[RowLocator::Hot(1)],
"old key still resolves the in-place position"
);
assert_eq!(t.rows()[1].values[1], Value::text("bobby"));
}
#[test]
fn update_row_indexed_column_moves_entry() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob"), (3, "carol")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
t.update_row(1, vec![Value::BigInt(20), Value::text("bob")])
.unwrap();
let idx = t.index_on(0).unwrap();
assert!(
idx.lookup_eq(&IndexKey::Int(2)).is_empty(),
"old key entry removed"
);
assert_eq!(
idx.lookup_eq(&IndexKey::Int(20)),
&[RowLocator::Hot(1)],
"new key entry resolves the position"
);
assert_eq!(idx.lookup_eq(&IndexKey::Int(1)), &[RowLocator::Hot(0)]);
assert_eq!(idx.lookup_eq(&IndexKey::Int(3)), &[RowLocator::Hot(2)]);
}
#[test]
fn update_row_duplicate_key_moves_only_target_position() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(7i64, "a"), (7, "b"), (9, "c")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
t.update_row(1, vec![Value::BigInt(8), Value::text("b")])
.unwrap();
let idx = t.index_on(0).unwrap();
assert_eq!(idx.lookup_eq(&IndexKey::Int(7)), &[RowLocator::Hot(0)]);
assert_eq!(idx.lookup_eq(&IndexKey::Int(8)), &[RowLocator::Hot(1)]);
assert_eq!(idx.lookup_eq(&IndexKey::Int(9)), &[RowLocator::Hot(2)]);
}
#[test]
fn update_row_null_transition_on_indexed_nullable_column() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"n",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("tag", DataType::BigInt, true),
],
))
.unwrap();
let t = cat.get_mut("n").unwrap();
t.insert(Row::new(vec![Value::BigInt(1), Value::BigInt(5)]))
.unwrap();
t.add_index("by_tag".into(), "tag").unwrap();
t.update_row(0, vec![Value::BigInt(1), Value::Null])
.unwrap();
let idx = t.index_on(1).unwrap();
assert!(idx.lookup_eq(&IndexKey::Int(5)).is_empty());
t.update_row(0, vec![Value::BigInt(1), Value::BigInt(6)])
.unwrap();
let idx = t.index_on(1).unwrap();
assert_eq!(idx.lookup_eq(&IndexKey::Int(6)), &[RowLocator::Hot(0)]);
}
#[test]
fn lookup_by_pk_finds_row_via_hot_index() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob"), (3, "carol")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let got = cat
.lookup_by_pk("users", "by_id", &IndexKey::Int(2))
.unwrap();
assert_eq!(got, make_user_row(2, "bob"));
assert_eq!(cat.cold_segment_count(), 0);
}
#[test]
fn lookup_by_pk_returns_none_when_key_missing() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(1, "alice")).unwrap();
t.add_index("by_id".into(), "id").unwrap();
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(999))
.is_none()
);
assert!(
cat.lookup_by_pk("other_table", "by_id", &IndexKey::Int(1))
.is_none()
);
assert!(
cat.lookup_by_pk("users", "no_such_index", &IndexKey::Int(1))
.is_none()
);
}
#[test]
fn lookup_by_pk_resolves_cold_locator_via_loaded_segment() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.add_index("by_id".into(), "id").unwrap();
let schema = t.schema.clone();
let cold_rows: Vec<(i64, &str)> = vec![(100, "ivy"), (200, "joe"), (300, "kim"), (400, "lin")];
let seg_rows: Vec<(u64, Vec<u8>)> = cold_rows
.iter()
.map(|(id, name)| {
let row = make_user_row(*id, name);
((*id).cast_unsigned(), encode_row_body_dense(&row, &schema))
})
.collect();
let (seg_bytes, _meta) =
encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES).unwrap();
let seg_id = cat.load_segment_bytes(seg_bytes).unwrap();
assert_eq!(seg_id, 0);
assert_eq!(cat.cold_segment_count(), 1);
let pairs: Vec<(IndexKey, RowLocator)> = cold_rows
.iter()
.map(|(id, _)| {
(
IndexKey::Int(*id),
RowLocator::Cold {
segment_id: seg_id,
page_offset: 0,
},
)
})
.collect();
let registered = cat
.get_mut("users")
.unwrap()
.register_cold_locators("by_id", pairs)
.unwrap();
assert_eq!(registered, 4);
for (id, name) in &cold_rows {
let got = cat
.lookup_by_pk("users", "by_id", &IndexKey::Int(*id))
.unwrap_or_else(|| panic!("cold key {id} not found"));
assert_eq!(got, make_user_row(*id, name));
}
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(999))
.is_none()
);
}
#[test]
fn lookup_by_pk_mixes_hot_and_cold_tiers() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let schema = t.schema.clone();
let cold_rows: Vec<(i64, &str)> = vec![(100, "ivy"), (200, "joe")];
let seg_rows: Vec<(u64, Vec<u8>)> = cold_rows
.iter()
.map(|(id, name)| {
let row = make_user_row(*id, name);
((*id).cast_unsigned(), encode_row_body_dense(&row, &schema))
})
.collect();
let (seg_bytes, _) = encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES).unwrap();
let seg_id = cat.load_segment_bytes(seg_bytes).unwrap();
let pairs: Vec<(IndexKey, RowLocator)> = cold_rows
.iter()
.map(|(id, _)| {
(
IndexKey::Int(*id),
RowLocator::Cold {
segment_id: seg_id,
page_offset: 0,
},
)
})
.collect();
cat.get_mut("users")
.unwrap()
.register_cold_locators("by_id", pairs)
.unwrap();
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(1))
.unwrap(),
make_user_row(1, "alice")
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(2))
.unwrap(),
make_user_row(2, "bob")
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(100))
.unwrap(),
make_user_row(100, "ivy")
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(200))
.unwrap(),
make_user_row(200, "joe")
);
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(50))
.is_none()
);
}
#[test]
fn register_cold_locators_rejects_nsw_index() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"vecs",
vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new(
"v",
DataType::Vector {
dim: 4,
encoding: VecEncoding::F32,
},
false,
),
],
))
.unwrap();
let t = cat.get_mut("vecs").unwrap();
t.insert(Row::new(vec![
Value::Int(1),
Value::vector(vec![1.0, 0.0, 0.0, 0.0]),
]))
.unwrap();
t.add_nsw_index("by_v".into(), "v", NSW_DEFAULT_M).unwrap();
let err = t
.register_cold_locators(
"by_v",
vec![(
IndexKey::Int(1),
RowLocator::Cold {
segment_id: 0,
page_offset: 0,
},
)],
)
.unwrap_err();
assert!(matches!(err, StorageError::Corrupt(ref s) if s.contains("not BTree")));
}
#[test]
fn load_segment_bytes_rejects_garbage() {
let mut cat = Catalog::new();
let err = cat.load_segment_bytes(vec![0u8; 10]).unwrap_err();
assert!(matches!(err, StorageError::Corrupt(ref s) if s.contains("segment")));
assert_eq!(cat.cold_segment_count(), 0);
}
#[test]
fn load_segment_bytes_returns_sequential_ids() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let schema = cat.get("users").unwrap().schema.clone();
for batch in 0u32..3 {
let rows: Vec<(u64, Vec<u8>)> = (0u64..4)
.map(|i| {
let id = u64::from(batch) * 100 + i;
let row = make_user_row(id.cast_signed(), "x");
(id, encode_row_body_dense(&row, &schema))
})
.collect();
let (bytes, _) = encode_segment(rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES).unwrap();
assert_eq!(cat.load_segment_bytes(bytes).unwrap(), batch);
}
assert_eq!(cat.cold_segment_count(), 3);
}
#[test]
fn v8_catalog_decodes_as_all_hot_under_v9_reader() {
let mut cat = populated_users();
cat.get_mut("users")
.unwrap()
.add_index("by_name".into(), "name")
.unwrap();
let v8_bytes = encode_as_v8(&cat);
assert_eq!(v8_bytes[FILE_MAGIC.len()], 8, "version byte must be 8");
let restored = Catalog::deserialize(&v8_bytes).expect("v9 reader accepts v8 stream");
let idx = restored
.get("users")
.unwrap()
.index_on(1)
.expect("index_on(1) after restore");
assert_eq!(
idx.lookup_eq(&IndexKey::Text("alice".into())),
&[RowLocator::Hot(0), RowLocator::Hot(2)]
);
for entry in idx.lookup_eq(&IndexKey::Text("alice".into())) {
assert!(entry.is_hot(), "v8 → v9 read must yield Hot only");
}
}
fn encode_as_v8(cat: &Catalog) -> Vec<u8> {
let mut out = Vec::with_capacity(64);
out.extend_from_slice(FILE_MAGIC);
out.push(8u8);
write_u32(&mut out, u32::try_from(cat.tables.len()).unwrap());
for t in &cat.tables {
write_str(&mut out, &t.schema.name);
write_u16(&mut out, u16::try_from(t.schema.columns.len()).unwrap());
for c in &t.schema.columns {
write_str(&mut out, &c.name);
write_data_type(&mut out, c.ty);
out.push(u8::from(c.nullable));
match &c.default {
None => out.push(0),
Some(v) => {
out.push(1);
write_value(&mut out, v);
}
}
out.push(u8::from(c.auto_increment));
}
write_u32(&mut out, u32::try_from(t.rows.len()).unwrap());
for row in &t.rows {
out.extend_from_slice(&encode_row_body_dense(row, &t.schema));
}
write_u16(&mut out, u16::try_from(t.indices.len()).unwrap());
for idx in &t.indices {
write_str(&mut out, &idx.name);
write_u16(&mut out, u16::try_from(idx.column_position).unwrap());
match &idx.kind {
IndexKind::BTree(_) => out.push(0),
IndexKind::Nsw(g) => {
out.push(1);
write_u16(&mut out, u16::try_from(g.m).unwrap());
write_nsw_graph(&mut out, g);
}
IndexKind::Brin { .. } => panic!(
"v8 catalog writer cannot serialise BRIN — \
tests with BRIN indices must use the current writer"
),
IndexKind::Gin(_) => panic!(
"v8 catalog writer cannot serialise GIN — \
tests with GIN indices must use the current writer"
),
IndexKind::GinTrgm(_) => panic!(
"v8 catalog writer cannot serialise trigram-GIN — \
tests with trgm indices must use the current writer"
),
IndexKind::GinFulltext(_) => panic!(
"v8 catalog writer cannot serialise fulltext-GIN — \
tests with FULLTEXT KEY must use the current writer"
),
IndexKind::GinJsonb(_) => panic!(
"v8 catalog writer cannot serialise JSONB-GIN — \
tests with JSONB-GIN must use the current writer"
),
}
}
}
out
}
#[test]
fn v9_catalog_round_trip_preserves_cold_locators() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob")] {
t.insert(make_user_row(id, name)).unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let schema = t.schema.clone();
let cold_rows: Vec<(i64, &str)> = vec![(100, "ivy"), (200, "joe"), (300, "kim")];
let seg_rows: Vec<(u64, Vec<u8>)> = cold_rows
.iter()
.map(|(id, name)| {
let row = make_user_row(*id, name);
((*id).cast_unsigned(), encode_row_body_dense(&row, &schema))
})
.collect();
let (seg_bytes, _) = encode_segment(seg_rows.into_iter(), 0.01, SEGMENT_PAGE_BYTES).unwrap();
let seg_id = cat.load_segment_bytes(seg_bytes.clone()).unwrap();
let pairs: Vec<(IndexKey, RowLocator)> = cold_rows
.iter()
.map(|(id, _)| {
(
IndexKey::Int(*id),
RowLocator::Cold {
segment_id: seg_id,
page_offset: 0,
},
)
})
.collect();
cat.get_mut("users")
.unwrap()
.register_cold_locators("by_id", pairs)
.unwrap();
let bytes = cat.serialize();
assert_eq!(bytes[FILE_MAGIC.len()], FILE_VERSION);
let mut restored = Catalog::deserialize(&bytes).expect("v9 round-trip parses");
let restored_seg_id = restored.load_segment_bytes(seg_bytes).unwrap();
assert_eq!(restored_seg_id, seg_id);
let idx = restored.get("users").unwrap().index_on(0).unwrap();
assert_eq!(idx.lookup_eq(&IndexKey::Int(1)), &[RowLocator::Hot(0)]);
assert_eq!(idx.lookup_eq(&IndexKey::Int(2)), &[RowLocator::Hot(1)]);
for (id, _) in &cold_rows {
assert_eq!(
idx.lookup_eq(&IndexKey::Int(*id)),
&[RowLocator::Cold {
segment_id: seg_id,
page_offset: 0,
}]
);
}
assert_eq!(
restored
.lookup_by_pk("users", "by_id", &IndexKey::Int(2))
.unwrap(),
make_user_row(2, "bob")
);
for (id, name) in &cold_rows {
assert_eq!(
restored
.lookup_by_pk("users", "by_id", &IndexKey::Int(*id))
.unwrap(),
make_user_row(*id, name)
);
}
}
#[test]
fn row_body_encoded_len_matches_actual_encode_for_all_types() {
let schema = TableSchema::new(
"wide",
vec![
ColumnSchema::new("a", DataType::SmallInt, true),
ColumnSchema::new("b", DataType::Int, false),
ColumnSchema::new("c", DataType::BigInt, false),
ColumnSchema::new("d", DataType::Float, false),
ColumnSchema::new("e", DataType::Bool, false),
ColumnSchema::new("f", DataType::Text, false),
ColumnSchema::new(
"g",
DataType::Vector {
dim: 3,
encoding: VecEncoding::F32,
},
false,
),
ColumnSchema::new(
"h",
DataType::Numeric {
precision: 18,
scale: 2,
},
false,
),
ColumnSchema::new("i", DataType::Date, false),
ColumnSchema::new("j", DataType::Timestamp, false),
],
);
let cases: &[Row] = &[
Row::new(vec![
Value::SmallInt(7),
Value::Int(42),
Value::BigInt(1_000_000),
Value::Float(1.5),
Value::Bool(true),
Value::text("hello"),
Value::vector(vec![1.0, 2.0, 3.0]),
Value::Numeric {
scaled: 12345,
scale: 2,
kind: crate::NumericKind::Finite,
},
Value::Date(20_000),
Value::Timestamp(1_700_000_000_000_000),
]),
Row::new(vec![
Value::Null,
Value::Int(0),
Value::BigInt(0),
Value::Float(0.0),
Value::Bool(false),
Value::text(""),
Value::vector(vec![]),
Value::Numeric {
scaled: 0,
scale: 2,
kind: crate::NumericKind::Finite,
},
Value::Date(0),
Value::Timestamp(0),
]),
Row::new(vec![
Value::SmallInt(-1),
Value::Int(-1),
Value::BigInt(-1),
Value::Float(-0.5),
Value::Bool(true),
Value::text("a much longer payload here"),
Value::vector(vec![0.1, 0.2, 0.3]),
Value::Numeric {
scaled: -999_999_999,
scale: 2,
kind: crate::NumericKind::Finite,
},
Value::Date(-1),
Value::Timestamp(-1),
]),
];
for row in cases {
let actual = encode_row_body_dense(row, &schema).len();
let fast = row_body_encoded_len(row, &schema);
assert_eq!(actual, fast, "row {row:?}");
}
}
#[test]
fn hot_bytes_grows_on_insert_and_matches_encoded_sum() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
assert_eq!(t.hot_bytes(), 0);
let mut expected: u64 = 0;
for (id, name) in [(1i64, "alice"), (2, "bob"), (3, "carol")] {
let row = make_user_row(id, name);
expected += encode_row_body_dense(&row, &t.schema).len() as u64;
t.insert(row).unwrap();
}
assert_eq!(t.hot_bytes(), expected);
assert_eq!(cat.hot_tier_bytes(), expected);
}
#[test]
fn hot_bytes_shrinks_on_delete() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for (id, name) in [(1i64, "alice"), (2, "bob"), (3, "carol")] {
t.insert(make_user_row(id, name)).unwrap();
}
let before = t.hot_bytes();
let bob_row = make_user_row(2, "bob");
let bob_bytes = encode_row_body_dense(&bob_row, &t.schema).len() as u64;
let removed = t.delete_rows(&[1]);
assert_eq!(removed, 1);
assert_eq!(t.hot_bytes(), before - bob_bytes);
}
#[test]
fn hot_bytes_diffs_on_update_for_variable_width_columns() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(1, "alice")).unwrap();
let after_insert = t.hot_bytes();
let new_row = make_user_row(1, "alice-the-longer-name");
let old_len = encode_row_body_dense(&make_user_row(1, "alice"), &t.schema).len() as u64;
let new_len = encode_row_body_dense(&new_row, &t.schema).len() as u64;
t.update_row(0, new_row.values).unwrap();
assert_eq!(t.hot_bytes(), after_insert - old_len + new_len);
assert!(t.hot_bytes() > after_insert, "longer text grew the counter");
}
#[test]
fn hot_bytes_round_trips_through_serialize_deserialize() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for i in 0..10 {
t.insert(make_user_row(i, &alloc::format!("name-{i}")))
.unwrap();
}
let pre = cat.hot_tier_bytes();
let restored = Catalog::deserialize(&cat.serialize()).unwrap();
assert_eq!(restored.hot_tier_bytes(), pre);
assert_eq!(restored.get("users").unwrap().hot_bytes(), pre);
}
#[test]
fn freeze_oldest_to_cold_moves_rows_and_keeps_lookups_working() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..10i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let total_bytes_before = t.hot_bytes();
let report = cat
.freeze_oldest_to_cold("users", "by_id", 6)
.expect("freeze succeeds");
assert_eq!(report.frozen_rows, 6);
assert_eq!(report.segment_id, 0);
assert!(report.bytes_freed > 0);
assert!(!report.segment_bytes.is_empty());
let t = cat.get("users").unwrap();
assert_eq!(t.row_count(), 4, "4 hot rows remain (10 - 6 frozen)");
assert_eq!(cat.cold_segment_count(), 1);
assert_eq!(
t.hot_bytes(),
total_bytes_before - report.bytes_freed,
"hot_bytes accounting matches FreezeReport"
);
for id in 0..10i64 {
let got = cat
.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
.unwrap_or_else(|| panic!("PK {id} disappeared after freeze"));
assert_eq!(got, make_user_row(id, &alloc::format!("u-{id}")));
}
}
#[test]
fn freeze_twice_preserves_prior_cold_locators() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..12i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 4)
.expect("first freeze ok");
cat.freeze_oldest_to_cold("users", "by_id", 4)
.expect("second freeze ok");
assert_eq!(cat.get("users").unwrap().row_count(), 4);
assert_eq!(cat.cold_segment_count(), 2);
for id in 0..12i64 {
let got = cat
.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
.unwrap_or_else(|| panic!("PK {id} not resolvable after two freezes"));
assert_eq!(got, make_user_row(id, &alloc::format!("u-{id}")));
}
}
#[test]
fn v7_37_15_headers_stay_lock_step_with_rows() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..10i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
assert_eq!(t.row_count(), 10);
assert_eq!(
t.headers().len(),
t.rows().len(),
"headers in lock-step after 10 inserts"
);
t.delete_rows(&[0, 2, 4]);
assert_eq!(t.row_count(), 7);
assert_eq!(
t.headers().len(),
t.rows().len(),
"headers in lock-step after delete_rows"
);
t.insert_no_index(make_user_row(100, "replay-100")).unwrap();
assert_eq!(t.row_count(), 8);
assert_eq!(
t.headers().len(),
t.rows().len(),
"headers in lock-step after insert_no_index"
);
t.delete_rows_no_index(&[0, 1]);
assert_eq!(t.row_count(), 6);
assert_eq!(
t.headers().len(),
t.rows().len(),
"headers in lock-step after delete_rows_no_index"
);
t.truncate();
assert_eq!(t.row_count(), 0);
assert_eq!(
t.headers().len(),
t.rows().len(),
"headers in lock-step after truncate"
);
}
#[test]
fn v7_37_15_fresh_inserts_default_to_frozen_header() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(7, "alice")).unwrap();
let header = t.headers().get(0).expect("header for first row");
assert!(
header.is_all_visible_fast(),
"fresh insert must default to frozen + alive (got {header:?})"
);
}
#[test]
fn v7_37_15_phase_b_scan_visible_filters_correctly() {
use crate::snapshot::{InProgressSet, Snapshot};
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..5i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
let unbounded = Snapshot::unbounded();
let all: Vec<_> = t.scan_visible(&unbounded).collect();
assert_eq!(all.len(), 5, "unbounded snapshot must see every frozen row");
for i in 0..5 {
assert!(t.is_row_visible(i, &unbounded), "row {i} visible");
}
let snap = Snapshot::new(
100, InProgressSet::from_sorted(alloc::vec![50]), 50, 0, );
{
let header = t.headers().get(0).expect("row 0 header exists");
let in_flight = crate::row_header::RowHeader {
xmin: 50,
xmax: crate::row_header::XMAX_ALIVE,
flags: 0,
};
let _ = header; let new_headers = t
.headers_mut_for_test()
.set(0, in_flight)
.expect("row 0 exists");
*t.headers_mut_for_test() = new_headers;
}
let visible: Vec<_> = t.scan_visible(&snap).map(|(i, _)| i).collect();
assert!(
!visible.contains(&0),
"row 0 has xmin=50 (in-flight) — must be invisible to snap; \
visible set = {visible:?}"
);
assert_eq!(
visible.len(),
4,
"other 4 rows still frozen + alive; visible to snap"
);
}
#[test]
fn v7_37_15_phase_c_insert_with_xmin_stamps_header() {
use crate::snapshot::{InProgressSet, Snapshot};
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert_with_xmin(make_user_row(7, "alice"), 17).unwrap();
let header = t.headers().get(0).expect("header exists");
assert_eq!(header.xmin, 17, "header xmin = writer's tx version");
assert_eq!(header.xmax, crate::row_header::XMAX_ALIVE);
let before = Snapshot::new(
20, InProgressSet::from_sorted(alloc::vec![17]), 10, 0, );
assert!(
!t.is_row_visible(0, &before),
"row hidden while writer in-flight"
);
let after = Snapshot::new(20, InProgressSet::empty(), 20, 0);
assert!(
t.is_row_visible(0, &after),
"row visible after writer commits"
);
}
#[test]
fn v7_37_15_phase_c_mark_row_deleted_writes_xmax() {
use crate::snapshot::{InProgressSet, Snapshot};
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert_with_xmin(make_user_row(7, "alice"), 17).unwrap();
t.mark_row_deleted(0, 25).unwrap();
let h = t.headers().get(0).expect("header exists");
assert_eq!(h.xmin, 17);
assert_eq!(h.xmax, 25, "xmax = deleter's version");
assert_eq!(t.row_count(), 1);
let before = Snapshot::new(30, InProgressSet::from_sorted(alloc::vec![25]), 17, 0);
assert!(
t.is_row_visible(0, &before),
"row visible while deleter in-flight"
);
let after = Snapshot::new(30, InProgressSet::empty(), 30, 0);
assert!(
!t.is_row_visible(0, &after),
"row hidden after deleter commits"
);
}
#[test]
fn v7_37_15_phase_c_repeated_delete_preserves_original_xmax() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert_with_xmin(make_user_row(7, "alice"), 17).unwrap();
t.mark_row_deleted(0, 25).unwrap();
t.mark_row_deleted(0, 100).unwrap(); let h = t.headers().get(0).expect("header exists");
assert_eq!(h.xmax, 25, "original xmax preserved; first-deleter-wins");
}
#[test]
fn v7_37_15_phase_c_frozen_xmin_short_circuits_to_legacy_insert() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert_with_xmin(make_user_row(7, "alice"), crate::row_header::XMIN_FROZEN)
.unwrap();
let h = t.headers().get(0).expect("header exists");
assert!(
h.is_all_visible_fast(),
"frozen-xmin call must produce a frozen-fast header"
);
}
#[test]
fn v7_37_15_phase_d_vacuum_reclaims_only_safe_rows() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..5i64 {
t.insert_with_xmin(make_user_row(id, &alloc::format!("u-{id}")), 10 + id as u64)
.unwrap();
}
t.mark_row_deleted(1, 100).unwrap();
t.mark_row_deleted(3, 200).unwrap();
assert_eq!(t.row_count(), 5, "physical row count unchanged by delete");
let dry = t.vacuum(150, true);
assert_eq!(dry.rows_reclaimed, 1, "dry-run counts safe-to-reclaim only");
assert_eq!(t.row_count(), 5, "dry-run does not mutate");
let expected_survivor_rowids: alloc::vec::Vec<crate::row_header::RowId> = (0..t.row_count())
.filter(|&i| i != 1)
.filter_map(|i| t.rowids().get(i).copied())
.collect();
let real = t.vacuum(150, false);
assert_eq!(real.rows_reclaimed, 1);
assert_eq!(t.row_count(), 4);
let survivor_rowids: alloc::vec::Vec<crate::row_header::RowId> =
t.rowids().iter().copied().collect();
assert_eq!(
survivor_rowids, expected_survivor_rowids,
"survivors keep their stable RowIds across vacuum compaction"
);
assert_eq!(
t.rowids().len(),
t.rows().len(),
"rowids lock-step with rows"
);
}
#[test]
fn v7_37_15_phase_d_vacuum_all_aggregates_per_table() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let users_schema = TableSchema::new(
"logs",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("msg", DataType::Text, true),
],
);
cat.create_table(users_schema).unwrap();
let u = cat.get_mut("users").unwrap();
u.insert_with_xmin(make_user_row(0, "alice"), 10).unwrap();
u.mark_row_deleted(0, 50).unwrap();
let l = cat.get_mut("logs").unwrap();
l.insert_with_xmin(Row::new(vec![Value::BigInt(1), Value::text("hello")]), 20)
.unwrap();
let report = cat.vacuum_all(100, false);
assert_eq!(
report.rows_reclaimed, 1,
"one row reclaimed across both tables"
);
assert_eq!(
report.per_table.len(),
1,
"only tables with reclaimed rows appear; got {:?}",
report.per_table
);
assert_eq!(report.per_table[0].0, "users");
assert_eq!(report.per_table[0].1, 1);
}
#[test]
fn v7_37_15_phase_d_is_all_visible_tracks_mvcc_writers() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
assert!(t.is_all_visible());
for id in 0..3i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
assert!(t.is_all_visible(), "frozen-only table is all-visible");
t.insert_with_xmin(make_user_row(99, "mvcc"), 42).unwrap();
assert!(
!t.is_all_visible(),
"non-frozen xmin must clear all-visible"
);
t.mark_row_deleted(3, 50).unwrap();
let _ = t.vacuum(100, false);
assert!(
t.is_all_visible(),
"table is all-visible again after the MVCC row is reclaimed"
);
}
#[test]
fn v7_37_15_end_to_end_mvcc_lifecycle() {
use crate::snapshot::{InProgressSet, Snapshot};
use crate::vacuum::is_reclaimable;
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert_with_xmin(make_user_row(42, "alice"), 10).unwrap();
let mid_insert = Snapshot::new(15, InProgressSet::from_sorted(alloc::vec![10]), 10, 0);
assert!(
!t.is_row_visible(0, &mid_insert),
"writer in flight ⇒ hidden"
);
let post_insert = Snapshot::new(20, InProgressSet::empty(), 20, 0);
assert!(
t.is_row_visible(0, &post_insert),
"committed insert ⇒ visible"
);
t.mark_row_deleted(0, 30).unwrap();
let pre_delete = Snapshot::new(25, InProgressSet::empty(), 25, 0);
assert!(
t.is_row_visible(0, &pre_delete),
"pre-delete snapshot sees row"
);
let post_delete = Snapshot::new(50, InProgressSet::empty(), 30, 0);
assert!(
!t.is_row_visible(0, &post_delete),
"post-delete snapshot hides row"
);
assert!(!is_reclaimable(30, 25));
let dry_at_25 = t.vacuum(25, true);
assert_eq!(
dry_at_25.rows_reclaimed, 0,
"vacuum waits for oldest_active > xmax"
);
assert!(is_reclaimable(30, 40));
let real_at_40 = t.vacuum(40, false);
assert_eq!(real_at_40.rows_reclaimed, 1);
assert_eq!(t.row_count(), 0, "row physically reclaimed by vacuum");
}
#[test]
fn freeze_oldest_to_cold_rejects_invalid_input() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..3i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
assert!(matches!(
cat.freeze_oldest_to_cold("users", "by_id", 0),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.freeze_oldest_to_cold("missing", "by_id", 1),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.freeze_oldest_to_cold("users", "no_such_index", 1),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.freeze_oldest_to_cold("users", "by_id", 999),
Err(StorageError::Corrupt(_))
));
assert_eq!(cat.get("users").unwrap().row_count(), 3);
assert_eq!(cat.cold_segment_count(), 0);
}
#[test]
fn freeze_oldest_to_cold_rejects_non_integer_pk() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"by_name",
vec![
ColumnSchema::new("name", DataType::Text, false),
ColumnSchema::new("payload", DataType::BigInt, false),
],
))
.unwrap();
let t = cat.get_mut("by_name").unwrap();
t.insert(Row::new(vec![Value::text("a"), Value::BigInt(1)]))
.unwrap();
t.add_index("by_n".into(), "name").unwrap();
let err = cat
.freeze_oldest_to_cold("by_name", "by_n", 1)
.expect_err("non-integer PK rejected");
match err {
StorageError::Corrupt(s) => assert!(
s.contains("non-integer"),
"error message names the constraint: {s}"
),
other => panic!("expected Corrupt, got {other:?}"),
}
assert_eq!(cat.get("by_name").unwrap().row_count(), 1);
assert_eq!(cat.cold_segment_count(), 0);
}
#[test]
fn freeze_keeps_remaining_hot_rows_addressable_via_secondary_index() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
t.add_index("by_name".into(), "name").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
let idx = cat.get("users").unwrap().index_on(1).unwrap();
let got = idx.lookup_eq(&IndexKey::Text("u-4".into()));
assert_eq!(got.len(), 1);
assert!(got[0].is_hot(), "kept-hot rows still surface as Hot");
match got[0] {
RowLocator::Hot(i) => {
assert_eq!(i, 1);
}
RowLocator::Cold { .. } => unreachable!(),
}
}
#[test]
fn promote_cold_row_pulls_frozen_row_back_to_hot_tier() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 4).unwrap();
let hot_bytes_before = cat.get("users").unwrap().hot_bytes();
let new_idx = cat
.promote_cold_row("users", "by_id", &IndexKey::Int(2))
.expect("promote ok")
.expect("PK 2 was cold");
assert_eq!(
new_idx, 2,
"promoted row appended after the 2 surviving hot rows"
);
let t = cat.get("users").unwrap();
assert_eq!(t.row_count(), 3, "hot tier grew from 2 to 3");
let row = make_user_row(2, "u-2");
let row_len = encode_row_body_dense(&row, &t.schema).len() as u64;
assert_eq!(t.hot_bytes(), hot_bytes_before + row_len);
let entries = t.index_on(0).unwrap().lookup_eq(&IndexKey::Int(2));
assert_eq!(entries.len(), 1, "exactly one locator per key");
assert!(entries[0].is_hot(), "promote retired the Cold locator");
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(2))
.unwrap(),
row
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(0))
.unwrap(),
make_user_row(0, "u-0")
);
}
#[test]
fn promote_cold_row_returns_none_when_key_is_not_cold() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(7, "alice")).unwrap();
t.add_index("by_id".into(), "id").unwrap();
assert!(
cat.promote_cold_row("users", "by_id", &IndexKey::Int(7))
.unwrap()
.is_none()
);
assert!(
cat.promote_cold_row("users", "by_id", &IndexKey::Int(99))
.unwrap()
.is_none()
);
assert_eq!(cat.get("users").unwrap().row_count(), 1);
assert_eq!(cat.cold_segment_count(), 0);
}
#[test]
fn shadow_cold_row_removes_cold_locators_and_drops_lookup() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..5i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(1))
.is_some(),
"frozen PK resolves before shadow"
);
let removed = cat
.shadow_cold_row("users", "by_id", &IndexKey::Int(1))
.unwrap();
assert_eq!(removed, 1, "exactly one cold locator retired");
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(1))
.is_none(),
"shadowed key no longer resolves"
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(0))
.unwrap(),
make_user_row(0, "u-0")
);
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(2))
.unwrap(),
make_user_row(2, "u-2")
);
}
#[test]
fn shadow_cold_row_returns_zero_when_key_is_not_cold() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(1, "alice")).unwrap();
t.add_index("by_id".into(), "id").unwrap();
assert_eq!(
cat.shadow_cold_row("users", "by_id", &IndexKey::Int(1))
.unwrap(),
0,
"hot-only key drops no cold locators"
);
assert_eq!(
cat.shadow_cold_row("users", "by_id", &IndexKey::Int(999))
.unwrap(),
0,
"absent key drops no cold locators"
);
assert_eq!(cat.get("users").unwrap().row_count(), 1);
}
#[test]
fn promote_and_shadow_reject_invalid_inputs() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
t.insert(make_user_row(1, "alice")).unwrap();
t.add_index("by_id".into(), "id").unwrap();
assert!(matches!(
cat.promote_cold_row("missing", "by_id", &IndexKey::Int(1)),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.shadow_cold_row("missing", "by_id", &IndexKey::Int(1)),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.promote_cold_row("users", "no_such_index", &IndexKey::Int(1)),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.shadow_cold_row("users", "no_such_index", &IndexKey::Int(1)),
Err(StorageError::Corrupt(_))
));
}
#[test]
fn commit_freeze_slices_single_slice_matches_freeze_oldest() {
let mut a = Catalog::new();
let mut b = Catalog::new();
for cat in [&mut a, &mut b] {
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..10i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
}
let single = a.freeze_oldest_to_cold("users", "by_id", 6).unwrap();
let slice = b
.prepare_freeze_slice("users", "by_id", 0..6)
.expect("prepare");
let parallel = b
.commit_freeze_slices("users", "by_id", alloc::vec![slice])
.expect("commit");
assert_eq!(single.segment_id, parallel.segment_id);
assert_eq!(single.frozen_rows, parallel.frozen_rows);
assert_eq!(single.bytes_freed, parallel.bytes_freed);
assert_eq!(single.segment_bytes, parallel.segment_bytes);
for id in 0..10i64 {
assert_eq!(
a.lookup_by_pk("users", "by_id", &IndexKey::Int(id)),
b.lookup_by_pk("users", "by_id", &IndexKey::Int(id)),
"PK {id} differs after single vs slice freeze"
);
}
}
#[test]
fn commit_freeze_slices_two_slices_match_single_slice() {
let mut a = Catalog::new();
let mut b = Catalog::new();
for cat in [&mut a, &mut b] {
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in [3, 7, 1, 9, 5, 0, 8, 4, 2, 6].iter().copied() {
t.insert(make_user_row(id as i64, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
}
let single = a
.prepare_freeze_slice("users", "by_id", 0..8)
.expect("prepare");
let one = a
.commit_freeze_slices("users", "by_id", alloc::vec![single])
.expect("commit one");
let s1 = b
.prepare_freeze_slice("users", "by_id", 0..4)
.expect("prepare s1");
let s2 = b
.prepare_freeze_slice("users", "by_id", 4..8)
.expect("prepare s2");
let two = b
.commit_freeze_slices("users", "by_id", alloc::vec![s1, s2])
.expect("commit two");
assert_eq!(one.segment_bytes, two.segment_bytes);
assert_eq!(one.frozen_rows, two.frozen_rows);
for id in 0..10i64 {
assert_eq!(
a.lookup_by_pk("users", "by_id", &IndexKey::Int(id)),
b.lookup_by_pk("users", "by_id", &IndexKey::Int(id)),
"PK {id} differs after one-slice vs two-slice freeze"
);
}
}
#[test]
fn commit_freeze_slices_rejects_gap() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let s1 = cat.prepare_freeze_slice("users", "by_id", 0..2).unwrap();
let s2 = cat.prepare_freeze_slice("users", "by_id", 3..5).unwrap();
assert!(matches!(
cat.commit_freeze_slices("users", "by_id", alloc::vec![s1, s2]),
Err(StorageError::Corrupt(_))
));
assert_eq!(cat.cold_segment_count(), 0);
assert_eq!(cat.get("users").unwrap().row_count(), 6);
}
#[test]
fn commit_freeze_slices_empty_is_noop() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..3i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let report = cat
.commit_freeze_slices("users", "by_id", Vec::new())
.unwrap();
assert_eq!(report.frozen_rows, 0);
assert_eq!(cat.cold_segment_count(), 0);
assert_eq!(cat.get("users").unwrap().row_count(), 3);
}
#[test]
fn compact_merges_small_segments_storage_unit() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..8i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
assert_eq!(cat.cold_segment_count(), 2);
assert_eq!(cat.cold_segment_slot_count(), 2);
let max_seg_bytes = cat
.cold_segment_ids_global()
.iter()
.map(|id| cat.cold_segment(*id).unwrap().bytes().len() as u64)
.max()
.unwrap();
let target = max_seg_bytes + 1;
let report = cat
.compact_cold_segments("users", "by_id", target)
.expect("compact succeeds");
assert_eq!(report.sources.len(), 2);
let merged_id = report.merged_segment_id.expect("merge happened");
assert_eq!(report.merged_rows, 6);
assert_eq!(report.deleted_rows_pruned, 0);
assert!(!report.merged_segment_bytes.is_empty());
assert_eq!(cat.cold_segment_count(), 1);
assert_eq!(cat.cold_segment_slot_count(), 3);
assert_eq!(cat.cold_segment_ids_global(), alloc::vec![merged_id]);
for id in 0..8i64 {
let got = cat
.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
.unwrap_or_else(|| panic!("PK {id} lost after compaction"));
assert_eq!(got, make_user_row(id, &alloc::format!("u-{id}")));
}
}
#[test]
fn compact_drops_shadowed_cold_rows() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
assert_eq!(
cat.shadow_cold_row("users", "by_id", &IndexKey::Int(1))
.unwrap(),
1
);
assert_eq!(
cat.shadow_cold_row("users", "by_id", &IndexKey::Int(4))
.unwrap(),
1
);
let max_seg_bytes = cat
.cold_segment_ids_global()
.iter()
.map(|id| cat.cold_segment(*id).unwrap().bytes().len() as u64)
.max()
.unwrap();
let report = cat
.compact_cold_segments("users", "by_id", max_seg_bytes + 1)
.expect("compact succeeds");
assert_eq!(report.sources.len(), 2);
assert_eq!(report.merged_rows, 4, "6 frozen − 2 shadowed = 4 live");
assert_eq!(report.deleted_rows_pruned, 2);
for shadowed in [1i64, 4i64] {
assert!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(shadowed))
.is_none(),
"shadowed PK {shadowed} must remain invisible after compact"
);
}
for live in [0i64, 2, 3, 5] {
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(live))
.unwrap_or_else(|| panic!("live PK {live} lost after compact"));
}
}
#[test]
fn compact_is_noop_below_two_candidates() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let report = cat
.compact_cold_segments("users", "by_id", 1 << 30)
.expect("noop ok");
assert!(report.merged_segment_id.is_none());
assert!(report.sources.is_empty());
cat.freeze_oldest_to_cold("users", "by_id", 4).unwrap();
let report = cat
.compact_cold_segments("users", "by_id", 1 << 30)
.expect("noop ok");
assert!(report.merged_segment_id.is_none());
assert_eq!(cat.cold_segment_count(), 1);
let report = cat
.compact_cold_segments("users", "by_id", 1)
.expect("noop ok");
assert!(report.merged_segment_id.is_none());
assert_eq!(cat.cold_segment_count(), 1);
}
#[test]
fn compact_swap_survives_catalog_roundtrip_via_load_at() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..6i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 3).unwrap();
let max_seg_bytes = cat
.cold_segment_ids_global()
.iter()
.map(|id| cat.cold_segment(*id).unwrap().bytes().len() as u64)
.max()
.unwrap();
let report = cat
.compact_cold_segments("users", "by_id", max_seg_bytes + 1)
.expect("compact ok");
let merged_id = report.merged_segment_id.unwrap();
let cat_bytes = cat.serialize();
let merged_bytes = report.merged_segment_bytes.clone();
let mut restored = Catalog::deserialize(&cat_bytes).expect("deserialize ok");
restored
.load_segment_bytes_at(merged_id, merged_bytes)
.expect("reload merged ok");
for id in 0..6i64 {
let got = restored
.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
.unwrap_or_else(|| panic!("PK {id} lost across roundtrip"));
assert_eq!(got, make_user_row(id, &alloc::format!("u-{id}")));
}
assert_eq!(restored.cold_segment_count(), 1);
}
#[test]
fn load_segment_bytes_at_pads_and_rejects_collision() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..4i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
let report = cat.freeze_oldest_to_cold("users", "by_id", 2).unwrap();
let bytes_seg0 = report.segment_bytes.clone();
cat.load_segment_bytes_at(5, bytes_seg0.clone())
.expect("pad + load ok");
assert_eq!(cat.cold_segment_slot_count(), 6);
assert_eq!(cat.cold_segment_count(), 2);
assert!(matches!(
cat.load_segment_bytes_at(5, bytes_seg0.clone()),
Err(StorageError::Corrupt(_))
));
assert!(matches!(
cat.load_segment_bytes_at(0, bytes_seg0),
Err(StorageError::Corrupt(_))
));
}
#[test]
fn promote_then_refreeze_does_not_leave_orphan_locators() {
let mut cat = Catalog::new();
cat.create_table(bigint_pk_users_schema()).unwrap();
let t = cat.get_mut("users").unwrap();
for id in 0..4i64 {
t.insert(make_user_row(id, &alloc::format!("u-{id}")))
.unwrap();
}
t.add_index("by_id".into(), "id").unwrap();
cat.freeze_oldest_to_cold("users", "by_id", 2).unwrap();
let promoted = cat
.promote_cold_row("users", "by_id", &IndexKey::Int(0))
.unwrap();
assert!(promoted.is_some());
let entries_after_promote = cat
.get("users")
.unwrap()
.index_on(0)
.unwrap()
.lookup_eq(&IndexKey::Int(0))
.to_vec();
assert_eq!(entries_after_promote.len(), 1);
assert!(entries_after_promote[0].is_hot());
for id in [2i64, 3] {
assert_eq!(
cat.lookup_by_pk("users", "by_id", &IndexKey::Int(id))
.unwrap(),
make_user_row(id, &alloc::format!("u-{id}"))
);
}
}
fn partition_parent_schema() -> TableSchema {
let mut s = TableSchema::new(
"events_partitioned",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("ts", DataType::Timestamptz, false),
ColumnSchema::new("payload", DataType::Jsonb, true),
],
);
s.partition_role = Some(PartitionRole::Parent {
kind: PartitionKind::Range,
key_column_positions: vec![1],
index_template_sources: vec![
"CREATE INDEX events_partitioned_ts_idx ON events_partitioned (ts DESC)".to_string(),
"CREATE INDEX events_partitioned_pid_ts_idx ON events_partitioned (payload, ts)"
.to_string(),
],
});
s
}
fn partition_range_child_schema(name: &str, lower_micros: i64, upper_micros: i64) -> TableSchema {
let mut s = TableSchema::new(
name,
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("ts", DataType::Timestamptz, false),
ColumnSchema::new("payload", DataType::Jsonb, true),
],
);
s.partition_role = Some(PartitionRole::Range {
parent_name: "events_partitioned".to_string(),
lower: PartitionBound::TimestampTz(lower_micros),
upper: PartitionBound::TimestampTz(upper_micros),
});
s
}
fn partition_default_child_schema(name: &str) -> TableSchema {
let mut s = TableSchema::new(
name,
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("ts", DataType::Timestamptz, false),
ColumnSchema::new("payload", DataType::Jsonb, true),
],
);
s.partition_role = Some(PartitionRole::Default {
parent_name: "events_partitioned".to_string(),
});
s
}
#[test]
fn partition_role_none_round_trips() {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"plain",
vec![ColumnSchema::new("id", DataType::BigInt, false)],
))
.unwrap();
let bytes = c.serialize();
let back = Catalog::deserialize(&bytes).unwrap();
assert!(back.get("plain").unwrap().schema().partition_role.is_none());
assert_eq!(bytes, back.serialize());
}
#[test]
fn partition_role_all_three_variants_round_trip() {
let mut c = Catalog::new();
c.create_table(partition_parent_schema()).unwrap();
c.create_table(partition_range_child_schema(
"events_2026_06",
1_748_736_000_000_000, 1_751_328_000_000_000, ))
.unwrap();
c.create_table(partition_range_child_schema(
"events_2026_07",
1_751_328_000_000_000,
1_754_006_400_000_000,
))
.unwrap();
c.create_table(partition_default_child_schema("events_default"))
.unwrap();
let bytes = c.serialize();
let back = Catalog::deserialize(&bytes).unwrap();
match back
.get("events_partitioned")
.unwrap()
.schema()
.partition_role
.as_ref()
.unwrap()
{
PartitionRole::Parent {
kind,
key_column_positions,
index_template_sources,
} => {
assert_eq!(*kind, PartitionKind::Range);
assert_eq!(key_column_positions, &vec![1usize]);
assert_eq!(index_template_sources.len(), 2);
assert!(index_template_sources[0].contains("ts DESC"));
assert!(index_template_sources[1].contains("payload, ts"));
}
other => panic!("expected Parent, got {other:?}"),
}
match back
.get("events_2026_06")
.unwrap()
.schema()
.partition_role
.as_ref()
.unwrap()
{
PartitionRole::Range {
parent_name,
lower,
upper,
} => {
assert_eq!(parent_name, "events_partitioned");
assert_eq!(*lower, PartitionBound::TimestampTz(1_748_736_000_000_000));
assert_eq!(*upper, PartitionBound::TimestampTz(1_751_328_000_000_000));
}
other => panic!("expected Range, got {other:?}"),
}
match back
.get("events_default")
.unwrap()
.schema()
.partition_role
.as_ref()
.unwrap()
{
PartitionRole::Default { parent_name } => {
assert_eq!(parent_name, "events_partitioned");
}
other => panic!("expected Default, got {other:?}"),
}
assert_eq!(bytes, back.serialize());
}
#[test]
fn partition_bound_minvalue_maxvalue_round_trip() {
let mut c = Catalog::new();
c.create_table(partition_parent_schema()).unwrap();
let mut s = TableSchema::new(
"events_all_time",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("ts", DataType::Timestamptz, false),
ColumnSchema::new("payload", DataType::Jsonb, true),
],
);
s.partition_role = Some(PartitionRole::Range {
parent_name: "events_partitioned".to_string(),
lower: PartitionBound::MinValue,
upper: PartitionBound::MaxValue,
});
c.create_table(s).unwrap();
let bytes = c.serialize();
let back = Catalog::deserialize(&bytes).unwrap();
match back
.get("events_all_time")
.unwrap()
.schema()
.partition_role
.as_ref()
.unwrap()
{
PartitionRole::Range { lower, upper, .. } => {
assert_eq!(*lower, PartitionBound::MinValue);
assert_eq!(*upper, PartitionBound::MaxValue);
}
other => panic!("expected Range, got {other:?}"),
}
assert_eq!(bytes, back.serialize());
}
#[test]
fn arena_clone_into_round_trip_heap_variants() {
let arena = bumpalo::Bump::new();
let owned_text = Value::text("hello world");
let arena_text = owned_text.clone_into(&arena);
assert_eq!(arena_text, owned_text);
assert_eq!(arena_text.clone().into_owned(), owned_text);
let owned_json = Value::json(r#"{"k":1}"#);
let arena_json = owned_json.clone_into(&arena);
assert_eq!(arena_json, owned_json);
let owned_xml = Value::xml("<a/>");
let arena_xml = owned_xml.clone_into(&arena);
assert_eq!(arena_xml, owned_xml);
let owned_bytes = Value::bytes(alloc::vec![1u8, 2, 3, 4]);
let arena_bytes = owned_bytes.clone_into(&arena);
assert_eq!(arena_bytes, owned_bytes);
let owned_vec = Value::vector(alloc::vec![1.0f32, 2.0, 3.0]);
let arena_vec = owned_vec.clone_into(&arena);
assert_eq!(arena_vec, owned_vec);
let owned_bits = Value::bit_string(12, alloc::vec![0xAB, 0xC0]);
let arena_bits = owned_bits.clone_into(&arena);
assert_eq!(arena_bits, owned_bits);
}
#[test]
fn arena_clone_into_round_trip_copy_and_nested_variants() {
let arena = bumpalo::Bump::new();
for v in [
Value::SmallInt(7),
Value::Int(-3),
Value::BigInt(42),
Value::Float(core::f64::consts::PI),
Value::Bool(true),
Value::Date(20_000),
Value::Timestamp(1_700_000_000_000_000),
Value::Uuid([1u8; 16]),
Value::Null,
Value::TextArray(alloc::vec![
Some("a".to_string()),
None,
Some("b".to_string()),
]),
Value::IntArray(alloc::vec![Some(1), None, Some(3)]),
] {
let lifted = v.clone_into(&arena);
assert_eq!(lifted, v, "clone_into changed scalar/array Value: {v:?}");
}
}
#[test]
fn arena_row_clone_into_then_into_owned_round_trip() {
let arena = bumpalo::Bump::new();
let row = Row::new(alloc::vec![
Value::BigInt(1),
Value::text("ham"),
Value::Null,
Value::bytes(alloc::vec![0xDE, 0xAD]),
]);
let arena_row = row.clone_into(&arena);
assert_eq!(arena_row, row);
let lifted: Row<'static> = arena_row.into_owned();
assert_eq!(lifted, row);
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![
ColumnSchema::new("id", DataType::BigInt, false),
ColumnSchema::new("v", DataType::Text, true),
ColumnSchema::new("opt", DataType::SmallInt, true),
ColumnSchema::new("blob", DataType::Bytes, true),
],
))
.unwrap();
c.get_mut("t").unwrap().insert(lifted.clone()).unwrap();
let bytes = c.serialize();
let back = Catalog::deserialize(&bytes).unwrap();
let stored = back.get("t").unwrap().rows().get(0).cloned().unwrap();
assert_eq!(stored, lifted, "catalog round-trip diverged");
assert_eq!(bytes, back.serialize(), "catalog bytes round-trip diverged");
}
#[test]
fn v54_snapshot_crc_detects_corruption() {
let mut c = Catalog::new();
c.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("id", DataType::Int, false)],
))
.unwrap();
c.get_mut("t")
.unwrap()
.insert(Row::new(alloc::vec![Value::Int(42)]))
.unwrap();
let bytes = c.serialize();
assert_eq!(bytes[FILE_MAGIC.len()], FILE_VERSION);
let restored = Catalog::deserialize(&bytes).expect("round-trips");
assert_eq!(restored.get("t").unwrap().rows().len(), 1);
let mut corrupt = bytes.clone();
let mid = corrupt.len() / 2;
corrupt[mid] ^= 0x01;
let err = Catalog::deserialize(&corrupt).unwrap_err();
assert!(
format!("{err:?}").contains("CRC mismatch") || format!("{err:?}").contains("Corrupt"),
"expected a corruption error, got {err:?}"
);
}
mod tx_writeset {
use crate::row_header::{RowId, XMAX_ALIVE};
use crate::{ColumnSchema, DataType, Row, Table, TableSchema, Value};
fn table() -> Table {
Table::new(TableSchema::new(
alloc::string::String::from("t"),
alloc::vec![ColumnSchema::new(
alloc::string::String::from("x"),
DataType::Int,
false
)],
))
}
fn row(x: i32) -> Row<'static> {
Row::new(alloc::vec![Value::Int(x)])
}
#[test]
fn extract_and_replay_round_trip_onto_fresher_clone() {
let mut old = table();
old.insert(row(1)).unwrap();
old.insert(row(2)).unwrap();
old.insert_with_xmin(row(3), 77).unwrap();
let rid1 = old.rowids().get(0).copied().unwrap();
old.mark_rows_deleted(&[0], 77);
let ws = old.extract_tx_writeset(77);
assert_eq!(ws.inserted.len(), 1);
assert_eq!(ws.tombstoned, alloc::vec![rid1]);
let inserted_rid = ws.inserted[0].0;
let mut fresh = table();
fresh.insert(row(1)).unwrap();
fresh.insert(row(2)).unwrap();
fresh.insert(row(9)).unwrap();
let conflicts = fresh.replay_tx_writeset(&ws, 77);
assert!(conflicts.is_empty(), "clean replay: {conflicts:?}");
let pos3 = (0..fresh.rows().len())
.find(|&i| fresh.rows().get(i).unwrap().values[0] == Value::Int(3))
.expect("replayed insert present");
assert_eq!(fresh.headers().get(pos3).unwrap().xmin, 77);
assert_eq!(fresh.rowids().get(pos3).copied(), Some(inserted_rid));
assert_eq!(fresh.headers().get(0).unwrap().xmax, 77);
assert_eq!(fresh.headers().get(1).unwrap().xmax, XMAX_ALIVE);
assert_eq!(fresh.dead_rows(), 1);
}
#[test]
fn replay_tombstone_survives_slot_shift() {
let mut old = table();
old.insert(row(1)).unwrap();
old.insert(row(2)).unwrap();
old.insert(row(3)).unwrap();
let rid3 = old.rowids().get(2).copied().unwrap();
old.mark_rows_deleted(&[2], 55);
let ws = old.extract_tx_writeset(55);
let mut fresh = table();
fresh.insert(row(1)).unwrap();
fresh.insert(row(2)).unwrap();
fresh.insert(row(3)).unwrap();
assert_eq!(fresh.rowids().get(2).copied(), Some(rid3));
fresh.delete_rows_no_index(&[0]);
let conflicts = fresh.replay_tx_writeset(&ws, 55);
assert!(conflicts.is_empty());
let pos3 = (0..fresh.rowids().len())
.find(|&i| fresh.rowids().get(i) == Some(&rid3))
.expect("row 3 present");
assert_eq!(fresh.headers().get(pos3).unwrap().xmax, 55);
}
#[test]
fn replay_reports_conflicts_and_is_idempotent() {
let mut old = table();
old.insert(row(1)).unwrap();
old.insert(row(2)).unwrap();
let rid1 = old.rowids().get(0).copied().unwrap();
let rid2 = old.rowids().get(1).copied().unwrap();
old.mark_rows_deleted(&[0, 1], 88);
let ws = old.extract_tx_writeset(88);
let mut fresh = table();
fresh.insert(row(1)).unwrap();
fresh.mark_rows_deleted(&[0], 99);
let conflicts = fresh.replay_tx_writeset(&ws, 88);
assert_eq!(conflicts, alloc::vec![rid1, rid2], "both targets conflict");
let mut fresh2 = table();
fresh2.insert(row(1)).unwrap();
fresh2.mark_rows_deleted(&[0], 88);
let ws1 = crate::TxWriteSet {
inserted: alloc::vec::Vec::new(),
tombstoned: alloc::vec![RowId(1)],
};
assert!(fresh2.replay_tx_writeset(&ws1, 88).is_empty());
}
}
#[cfg(test)]
mod mysql_ci_fold_tests {
use crate::mysql_ci_fold;
fn eq(a: &str, b: &str) -> bool {
mysql_ci_fold(a) == mysql_ci_fold(b)
}
#[test]
fn case_folds() {
assert!(eq("Foo", "foo"));
assert!(eq("FOO", "foo"));
assert!(eq("a", "A"));
assert!(eq("z", "Z"));
assert_eq!(mysql_ci_fold("MixedCase"), "mixedcase");
}
#[test]
fn accents_strip_to_the_base() {
for accented in ["á", "à", "â", "ä", "ã", "å"] {
assert!(eq("a", accented), "{accented} should fold to a");
}
for accented in ["é", "è", "ê", "ë"] {
assert!(eq("e", accented), "{accented} should fold to e");
}
assert!(eq("i", "í"));
assert!(eq("o", "ó"));
assert!(eq("o", "ö"));
assert!(eq("u", "ü"));
assert!(eq("n", "ñ"));
assert!(eq("c", "ç"));
assert!(eq("y", "ý"));
assert!(eq("A", "Ä"));
assert!(eq("O", "Ö"));
assert!(eq("U", "Ü"));
}
#[test]
fn ligatures_expand() {
assert!(eq("ss", "ß"));
assert!(!eq("s", "ß"));
assert!(eq("ae", "æ"));
assert!(!eq("a", "æ"));
assert!(eq("oe", "œ"));
assert!(eq("o", "ø"));
assert!(eq("d", "ð"));
}
#[test]
fn words_measured_on_mariadb() {
assert!(eq("Bär", "bar"));
assert!(!eq("Bär", "baer"));
assert!(eq("straße", "strasse"));
assert!(!eq("straße", "strase"));
assert!(eq("café", "cafe"));
assert!(eq("naïve", "naive"));
assert!(eq("RÉSUMÉ", "resume"));
}
#[test]
fn ascii_and_unknown_scripts_pass_through() {
assert_eq!(mysql_ci_fold("hello123"), "hello123");
assert_eq!(mysql_ci_fold(""), "");
assert_eq!(mysql_ci_fold("日本語"), "日本語");
}
}
#[test]
fn check_validated_flag_round_trips() {
let mut cat = Catalog::new();
cat.create_table(TableSchema::new(
"t",
vec![ColumnSchema::new("a", DataType::Int, true)],
))
.unwrap();
{
let t = cat.get_mut("t").unwrap();
t.schema_mut().checks = alloc::vec![
crate::CheckConstraint {
name: Some("c_valid".into()),
expr: "a > 0".into(),
validated: true,
},
crate::CheckConstraint {
name: Some("c_unvalidated".into()),
expr: "a < 100".into(),
validated: false,
},
];
}
let bytes = cat.serialize();
let back = Catalog::deserialize(&bytes).expect("round-trip");
let checks = &back.get("t").unwrap().schema().checks;
assert_eq!(checks.len(), 2);
assert!(checks[0].validated, "validated one stays validated");
assert!(!checks[1].validated, "NOT VALID one stays unvalidated");
assert_eq!(checks[1].name.as_deref(), Some("c_unvalidated"));
}
#[test]
fn check_constraint_size_is_pinned() {
assert_eq!(
core::mem::size_of::<crate::CheckConstraint>(),
56,
"CheckConstraint changed shape; re-measure with the boot-level \
server-side EXPLAIN ANALYZE instrument before accepting"
);
}
#[test]
fn v88_collation_name_survives_serialize_deserialize() {
let mut cat = Catalog::new();
let mut b = ColumnSchema::new("b", DataType::Text, true);
b.collation_name = Some("C".into());
let mut c = ColumnSchema::new("c", DataType::Text, true);
c.collation_name = Some("POSIX".into());
cat.create_table(TableSchema::new(
"ct",
vec![ColumnSchema::new("a", DataType::Text, true), b, c],
))
.unwrap();
let bytes = cat.serialize();
let back = Catalog::deserialize(&bytes).expect("round trip");
let cols = &back.get("ct").unwrap().schema().columns;
assert_eq!(cols[0].collation_name, None, "no COLLATE stays None");
assert_eq!(cols[1].collation_name.as_deref(), Some("C"));
assert_eq!(cols[2].collation_name.as_deref(), Some("POSIX"));
}
#[test]
fn v88_a_table_with_no_collation_pays_two_bytes() {
let mut plain = Catalog::new();
plain
.create_table(TableSchema::new(
"p",
vec![ColumnSchema::new("a", DataType::Text, true)],
))
.unwrap();
let mut collated = Catalog::new();
let mut a = ColumnSchema::new("a", DataType::Text, true);
a.collation_name = Some("C".into());
collated
.create_table(TableSchema::new("p", vec![a]))
.unwrap();
assert!(
collated.serialize().len() > plain.serialize().len(),
"a declared collation has to cost something"
);
let back = Catalog::deserialize(&plain.serialize()).expect("round trip");
assert_eq!(
back.get("p").unwrap().schema().columns[0].collation_name,
None
);
}
#[test]
fn pruned_decode_agrees_with_the_full_decode_and_ends_in_the_same_place() {
let schema = TableSchema::new(
"rec",
vec![
ColumnSchema::new("id", DataType::Int, false),
ColumnSchema::new("a", DataType::Text, true),
ColumnSchema::new("n", DataType::BigInt, true),
ColumnSchema::new("b", DataType::Text, true),
],
);
let rows = vec![
Row::new(vec![
Value::Int(1),
Value::text("short".to_string()),
Value::BigInt(7),
Value::text("x".repeat(300)),
]),
Row::new(vec![Value::Int(2), Value::Null, Value::Null, Value::Null]),
Row::new(vec![
Value::Int(3),
Value::text(String::new()),
Value::BigInt(-9),
Value::text("tail".to_string()),
]),
];
for mask in [
vec![], vec![true, false, true, true], vec![true, true, true, false], vec![true, false, true, false], vec![false, false, false, false], ] {
let mut arena = Vec::new();
for r in &rows {
encode_row_body_dense_into(r, &schema, &mut arena);
}
let mut at_full = 0usize;
let mut at_pruned = 0usize;
for (i, _) in rows.iter().enumerate() {
let (full, used_full) =
decode_row_body_dense(&arena[at_full..], &schema, CURRENT_ROW_CODEC_VERSION)
.unwrap();
let (pruned, used_pruned) = decode_row_body_dense_pruned(
&arena[at_pruned..],
&schema,
CURRENT_ROW_CODEC_VERSION,
&mask,
)
.unwrap();
assert_eq!(
used_full, used_pruned,
"row {i} mask {mask:?}: pruning changed how many bytes the row occupies"
);
assert_eq!(
full.values.len(),
pruned.values.len(),
"row {i} mask {mask:?}: arity has to survive pruning, positions are indexes"
);
for (c, (f, p)) in full.values.iter().zip(pruned.values.iter()).enumerate() {
let kept = mask.get(c).copied().unwrap_or(true);
if kept {
assert_eq!(f, p, "row {i} col {c} mask {mask:?}: kept column differs");
}
}
at_full += used_full;
at_pruned += used_pruned;
}
assert_eq!(
at_full,
arena.len(),
"mask {mask:?}: full decode consumed the arena"
);
assert_eq!(
at_pruned,
arena.len(),
"mask {mask:?}: pruned decode consumed the arena"
);
}
}