use std::sync::Arc;
use crate::persistent_artrie::block_storage::BlockStorage;
use crate::persistent_artrie::char::arena_manager::ArenaSlot;
use crate::persistent_artrie::char::nodes::persistent_node::PersistentCharNode;
use crate::persistent_artrie::char::nodes::CharNode;
use crate::persistent_artrie::char::persist::overlay_inner_single_node_with_prefix;
use crate::persistent_artrie::char::relative_encoding::SerializationContext;
use crate::persistent_artrie::char::serialization_char::{
deserialize_char_node_v2, serialize_char_node_v2, DeserializationContext,
};
use crate::persistent_artrie::char::types::CharTrieNodeInner;
use crate::persistent_artrie::core::eviction::DiskLocationRegistry;
use crate::persistent_artrie::core::key_encoding::CharKey;
use crate::persistent_artrie::core::overlay::compressed_serialize::OverlayCompressedSerialize;
use crate::persistent_artrie::error::{PersistentARTrieError, Result};
use crate::persistent_artrie::swizzled_ptr::SwizzledPtr;
use crate::persistent_artrie::NodeType;
type VocabOverlayNode = PersistentCharNode<u64>;
type VocabScan = (
std::collections::BTreeMap<Vec<u32>, Option<u64>>,
Option<Option<u64>>,
);
type DecodedVocabNode = (bool, Option<u64>, Vec<u32>, Vec<(u32, SwizzledPtr)>);
impl<S: BlockStorage> super::dict_impl::PersistentVocabARTrie<S> {
pub(super) fn build_disk_char_node_static(
original: &CharNode,
disk_children: &[(u32, SwizzledPtr)],
) -> CharNode {
let mut new_node = match original {
CharNode::N4(_) => CharNode::N4(Box::default()),
CharNode::N16(_) => CharNode::N16(Box::default()),
CharNode::N48(_) => CharNode::N48(Box::default()),
CharNode::Bucket(_) => CharNode::Bucket(Box::default()),
};
{
let new_header = new_node.header_mut();
let orig_header = original.header();
new_header.prefix_len = orig_header.prefix_len;
new_header.flags = orig_header.flags;
}
*new_node.prefix_mut() = *original.prefix();
for &(key, ref ptr) in disk_children {
match new_node.add_child_growing(key, ptr.clone()) {
Ok(Some(grown)) => new_node = grown,
Ok(None) => {}
Err(e) => {
eprintln!(
"Warning: failed to add child in build_disk_char_node_static: {:?}",
e
);
}
}
}
new_node
}
#[cfg(test)]
pub(super) fn serialize_overlay_to_disk(
&self,
root: &Arc<VocabOverlayNode>,
) -> Result<SwizzledPtr> {
struct PendingChild {
key: u32,
ptr: Option<SwizzledPtr>,
}
struct Frame {
node: Arc<VocabOverlayNode>,
parent_key: Option<u32>,
parent_pushed_path: bool,
pending_in_mem: Vec<(u32, Arc<VocabOverlayNode>)>,
slots: Vec<PendingChild>,
}
fn make_frame(
node: Arc<VocabOverlayNode>,
parent_key: Option<u32>,
parent_pushed_path: bool,
) -> Frame {
let n = node.num_children();
let mut slots: Vec<PendingChild> = Vec::with_capacity(n);
let mut pending_in_mem: Vec<(u32, Arc<VocabOverlayNode>)> = Vec::with_capacity(n);
for (&key, child) in node.iter_children() {
if let Some(child_arc) = child.as_in_mem() {
slots.push(PendingChild { key, ptr: None });
pending_in_mem.push((key, Arc::clone(child_arc)));
} else if let Some(on_disk) = child.as_on_disk() {
if !on_disk.is_null() {
slots.push(PendingChild {
key,
ptr: Some(on_disk.clone()),
});
}
}
}
pending_in_mem.reverse();
Frame {
node,
parent_key,
parent_pushed_path,
pending_in_mem,
slots,
}
}
let mut path: Vec<char> = Vec::new();
let mut stack: Vec<Frame> = Vec::new();
stack.push(make_frame(Arc::clone(root), None, false));
let mut completed: Option<(u32, SwizzledPtr)> = None;
loop {
let frame = stack
.last_mut()
.expect("serialize_overlay_to_disk: non-empty work-stack");
if let Some((key, ptr)) = completed.take() {
let slot = frame
.slots
.iter_mut()
.find(|s| s.key == key && s.ptr.is_none())
.expect("completed child key has a matching unfilled parent slot");
slot.ptr = Some(ptr);
}
if let Some((key, child_arc)) = frame.pending_in_mem.pop() {
let pushed = char::from_u32(key).map(|ch| path.push(ch)).is_some();
stack.push(make_frame(child_arc, Some(key), pushed));
continue;
}
let frame = stack
.pop()
.expect("serialize_overlay_to_disk: frame to finalize");
let child_disk_ptrs: Vec<(u32, SwizzledPtr)> = frame
.slots
.into_iter()
.map(|s| {
(
s.key,
s.ptr.expect(
"every in-mem child slot is filled before its parent node is \
serialized (post-order invariant)",
),
)
})
.collect();
let inner = overlay_inner_single_node(frame.node.as_ref(), &child_disk_ptrs);
let node_ptr = self.serialize_one_overlay_node(&inner, &child_disk_ptrs)?;
if frame.parent_pushed_path {
path.pop();
}
match frame.parent_key {
Some(key) => completed = Some((key, node_ptr)),
None => return Ok(node_ptr),
}
}
}
pub(super) fn serialize_overlay_snapshot_compressed(
&self,
root: &Arc<VocabOverlayNode>,
) -> Result<SwizzledPtr> {
OverlayCompressedSerialize::<CharKey, u64>::serialize_compressed_loop(self, root, None)
}
fn serialize_one_overlay_node(
&self,
node: &CharTrieNodeInner<u64>,
child_disk_ptrs: &[(u32, SwizzledPtr)],
) -> Result<SwizzledPtr> {
let arena_manager = self.arena_manager.as_ref().ok_or_else(|| {
PersistentARTrieError::internal("No arena manager for overlay serialization")
})?;
let parent_arena_id = arena_manager.read().next_slot().arena_id;
let (parent_slot, arena_node_count) = {
let mgr = arena_manager.read();
let slot = mgr.next_slot();
let node_count = mgr
.get_arena(parent_arena_id)
.map(|a| a.node_count())
.unwrap_or(0);
(slot, node_count)
};
let ctx = if parent_slot.slot_id < child_disk_ptrs.len() as u32 {
SerializationContext::full_encoding(parent_slot)
} else if let Some(first_child) =
check_sequential_char_children(child_disk_ptrs, parent_arena_id, arena_node_count)
{
SerializationContext::sequential(parent_slot, first_child)
} else {
SerializationContext::new(parent_slot)
};
let disk_node = Self::build_disk_char_node_static(&node.node, child_disk_ptrs);
let value_bytes: Vec<u8> = if let Some(ref value) = node.value {
crate::serialization::bincode_compat::serialize(value).map_err(|e| {
PersistentARTrieError::internal(format!("Failed to serialize value: {}", e))
})?
} else {
Vec::new()
};
let mut node_buffer = Vec::new();
serialize_char_node_v2(&disk_node, &mut node_buffer, &ctx)?;
let build_data = |node_buf: &[u8], value_buf: &[u8]| -> Vec<u8> {
let total_size = node_buf.len() + 4 + value_buf.len();
let mut data = Vec::with_capacity(total_size);
data.extend_from_slice(node_buf);
data.extend_from_slice(&(value_buf.len() as u32).to_le_bytes());
data.extend_from_slice(value_buf);
data
};
let data = build_data(&node_buffer, &value_bytes);
let slot = arena_manager.write().allocate(&data)?;
let final_slot = if slot != ctx.parent_slot {
let corrected_ctx = SerializationContext::new(slot);
let mut corrected_buffer = Vec::new();
serialize_char_node_v2(&disk_node, &mut corrected_buffer, &corrected_ctx)?;
let corrected_data = build_data(&corrected_buffer, &value_bytes);
if corrected_data.len() == data.len() {
arena_manager.write().update(slot, &corrected_data)?;
slot
} else {
arena_manager.write().allocate(&corrected_data)?
}
} else {
slot
};
Ok(SwizzledPtr::on_disk(
final_slot.arena_id + 1,
final_slot.slot_id,
NodeType::CharNode4,
))
}
fn read_overlay_record_fields(&self, node_ptr: &SwizzledPtr) -> Result<DecodedVocabNode> {
use std::io::Cursor;
let arena_manager = self.arena_manager.as_ref().ok_or_else(|| {
PersistentARTrieError::internal("vocab overlay enumerate: no arena manager")
})?;
let disk_loc = node_ptr.disk_location().ok_or_else(|| {
PersistentARTrieError::corrupted("vocab overlay enumerate: swizzled/null node ptr")
})?;
let arena_id = disk_loc.block_id.checked_sub(1).ok_or_else(|| {
PersistentARTrieError::corrupted("vocab overlay enumerate: invalid block_id 0")
})?;
let slot = ArenaSlot::new(arena_id, disk_loc.offset);
let am = arena_manager.read();
let node_data = am.read(slot)?;
let deser_ctx = DeserializationContext::new(slot);
let mut cursor = Cursor::new(node_data);
let char_node = deserialize_char_node_v2(&mut cursor, &deser_ctx)?;
let offset = cursor.position() as usize;
if offset + 4 > node_data.len() {
return Err(PersistentARTrieError::corrupted(
"vocab overlay enumerate: value_len extends past node record",
));
}
let value_len = u32::from_le_bytes([
node_data[offset],
node_data[offset + 1],
node_data[offset + 2],
node_data[offset + 3],
]) as usize;
let value: Option<u64> = if value_len > 0 {
let value_start = offset + 4;
let value_end = value_start.checked_add(value_len).ok_or_else(|| {
PersistentARTrieError::corrupted("vocab overlay enumerate: value length overflow")
})?;
if value_end > node_data.len() {
return Err(PersistentARTrieError::corrupted(
"vocab overlay enumerate: value bytes extend past node record",
));
}
Some(
crate::serialization::bincode_compat::deserialize(
&node_data[value_start..value_end],
)
.map_err(|e| {
PersistentARTrieError::corrupted(format!(
"vocab overlay enumerate: value deserialize failed: {}",
e
))
})?,
)
} else {
None
};
let is_final = char_node.is_final();
let plen = char_node.header().prefix_len as usize;
let prefix_units: Vec<u32> = char_node.prefix().chars[..plen].to_vec();
let children: Vec<(u32, SwizzledPtr)> = char_node
.iter_children()
.filter(|(_, ptr)| !ptr.is_null())
.map(|(key, ptr)| (key, ptr.clone()))
.collect();
drop(am);
Ok((is_final, value, prefix_units, children))
}
pub(super) fn enumerate_overlay_terms_from_disk(&self, root_ptr: u64) -> Result<VocabScan> {
use std::collections::BTreeMap;
let mut all: BTreeMap<Vec<u32>, Option<u64>> = BTreeMap::new();
if root_ptr == 0 {
return Ok((all, None)); }
let mut stack: Vec<(SwizzledPtr, Vec<u32>)> =
vec![(SwizzledPtr::from_raw(root_ptr), Vec::new())];
while let Some((ptr, parent_path)) = stack.pop() {
let (is_final, value, prefix_units, children) =
self.read_overlay_record_fields(&ptr)?;
let mut here = parent_path;
here.extend_from_slice(&prefix_units);
if is_final {
all.insert(here.clone(), value);
}
for (edge, child_ptr) in children.into_iter().rev() {
let mut p = here.clone();
p.push(edge);
stack.push((child_ptr, p));
}
}
let empty_term: Option<Option<u64>> = all.remove(&Vec::<u32>::new());
Ok((all, empty_term))
}
pub(super) fn reestablish_overlay_from_image(&mut self, root_ptr: u64) -> Result<()> {
use crate::persistent_artrie::core::key_encoding::CharKey;
use crate::persistent_artrie::core::overlay::f5_build::build_overlay_root_from_terms;
let (terms, empty) = self.enumerate_overlay_terms_from_disk(root_ptr)?;
let rev = dashmap::DashMap::new();
let mut entries: Vec<(Vec<u32>, Option<u64>)> = Vec::with_capacity(terms.len());
for (units, id_opt) in terms {
if let Some(id) = id_opt {
let term: String = units.iter().filter_map(|&u| char::from_u32(u)).collect();
rev.insert(id, term);
}
entries.push((units, id_opt));
}
if let Some(Some(id)) = empty {
rev.insert(id, String::new());
}
let overlay_root = build_overlay_root_from_terms::<CharKey, u64, _>(entries, empty);
self.install_prebuilt_overlay_root_inherent(overlay_root);
self.reverse_term_map = Some(rev);
if let Some(ref wal) = self.wal_writer {
if let Err(e) = wal.set_overlay_regime() {
log::warn!(
"vocab reestablish_overlay: could not stamp Overlay regime: {:?}",
e
);
}
}
Ok(())
}
pub(super) fn replay_wal_into_overlay_rank_aware(
&mut self,
wal_path: &std::path::Path,
checkpoint_lsn: u64,
) -> Result<(usize, usize)> {
use crate::persistent_artrie::core::overlay::flip::LockFreeOverlay;
use crate::persistent_artrie::wal::{WalReader, WalRecord};
use std::collections::HashSet;
use std::sync::atomic::Ordering;
let segments = self
.wal_writer
.as_ref()
.and_then(|wal| wal.collect_wal_segments(&self.wal_config).ok())
.filter(|segments| !segments.is_empty())
.unwrap_or_else(|| vec![wal_path.to_path_buf()]);
let mut pending: Vec<(u64, String, u64)> = Vec::new();
let mut ranked: HashSet<u64> = HashSet::new();
let mut max_generation: u64 = 0;
let mut records_seen = 0usize;
for segment in segments {
let reader = WalReader::new(&segment)?;
for rr in reader.iter() {
let (lsn, record) = rr?;
if lsn <= checkpoint_lsn {
continue;
}
match record {
WalRecord::Insert { term, value } => {
records_seen += 1;
if let Some(vb) = value {
if vb.len() >= 8 {
let id = u64::from_le_bytes(vb[..8].try_into().expect("len>=8"));
let t = String::from_utf8(term).map_err(|e| {
PersistentARTrieError::corrupted(format!(
"vocab overlay replay: invalid UTF-8 term: {e}"
))
})?;
pending.push((lsn, t, id));
}
}
}
WalRecord::CommitRank {
data_lsn,
generation,
..
} => {
ranked.insert(data_lsn);
max_generation = max_generation.max(generation);
}
WalRecord::Checkpoint {
checkpoint_lsn: new_lsn,
..
} => {
self.synced_lsn.fetch_max(new_lsn, Ordering::AcqRel);
}
_ => {}
}
self.next_lsn.fetch_max(lsn + 1, Ordering::AcqRel);
}
}
let mut applied = 0usize;
let mut max_id: u64 = 0;
for (lsn, term, id) in pending {
if !ranked.contains(&lsn) {
continue; }
max_id = max_id.max(id);
if self.get_index_lockfree(&term).is_none() {
let units: Vec<u32> = term.chars().map(|c| c as u32).collect();
self.overlay_publish_value(&units, id);
if let Some(ref rev) = self.reverse_term_map {
rev.insert(id, term);
}
self.entry_count.fetch_add(1, Ordering::AcqRel);
applied += 1;
}
}
if max_id > 0 {
self.next_index.fetch_max(max_id + 1, Ordering::AcqRel);
}
self.commit_seq.fetch_max(max_generation, Ordering::AcqRel);
Ok((records_seen, applied))
}
}
impl<S: BlockStorage> OverlayCompressedSerialize<CharKey, u64>
for super::dict_impl::PersistentVocabARTrie<S>
{
type Projected = CharTrieNodeInner<u64>;
fn project_node(
node: &VocabOverlayNode,
child_disk_ptrs: &[(u32, SwizzledPtr)],
) -> Result<Self::Projected> {
Ok(overlay_inner_single_node(node, child_disk_ptrs))
}
fn project_chunk(
synth: &VocabOverlayNode,
child_disk_ptrs: &[(u32, SwizzledPtr)],
prefix: &[u32],
) -> Result<Self::Projected> {
Ok(overlay_inner_single_node_with_prefix::<u64>(
synth,
child_disk_ptrs,
prefix,
))
}
fn serialize_projected_node(
&self,
projected: &Self::Projected,
child_disk_ptrs: &[(u32, SwizzledPtr)],
_path: &[u32],
_registry: Option<&mut DiskLocationRegistry>,
) -> Result<SwizzledPtr> {
self.serialize_one_overlay_node(projected, child_disk_ptrs)
}
fn new_synth_node() -> VocabOverlayNode {
VocabOverlayNode::new()
}
fn stamp_durable(_live: &VocabOverlayNode, _raw: u64) {
}
}
fn overlay_inner_single_node(
node: &VocabOverlayNode,
child_disk_ptrs: &[(u32, SwizzledPtr)],
) -> CharTrieNodeInner<u64> {
let mut inner = CharTrieNodeInner::<u64>::default();
inner.node.header_mut().set_final(node.is_final());
inner.value = node.get_value();
for &(key, ref ptr) in child_disk_ptrs {
if let Some(grown) = inner
.node
.add_child_growing(key, ptr.clone())
.expect("overlay_inner_single_node: add on-disk child within capacity")
{
inner.node = grown;
}
}
inner
}
fn check_sequential_char_children(
child_ptrs: &[(u32, SwizzledPtr)],
parent_arena_id: u32,
arena_node_count: u32,
) -> Option<ArenaSlot> {
if child_ptrs.len() < 2 {
return None;
}
let mut slots: Vec<ArenaSlot> = Vec::with_capacity(child_ptrs.len());
for (_, ptr) in child_ptrs {
let loc = ptr.disk_location()?;
let arena_id = loc.block_id.checked_sub(1)?;
if arena_id != parent_arena_id {
return None;
}
slots.push(ArenaSlot::new(arena_id, loc.offset));
}
let first = slots[0];
for (i, slot) in slots.iter().enumerate() {
if slot.slot_id != first.slot_id.checked_add(i as u32)? {
return None;
}
}
let count = slots.len() as u32;
let last_slot = first.slot_id + count - 1;
if last_slot >= arena_node_count {
return None;
}
Some(first)
}
#[cfg(test)]
mod tests {
use crate::persistent_artrie::vocab::PersistentVocabARTrie;
#[test]
fn overlay_serialize_enumerate_roundtrip() {
let scratch_root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("target/test-tmp");
std::fs::create_dir_all(&scratch_root).expect("scratch dir");
let dir = tempfile::Builder::new()
.prefix("vocab-overlay-rt-")
.tempdir_in(&scratch_root)
.expect("scratch tempdir");
let path = dir.path().join("fixture.vocab");
let mut vocab = PersistentVocabARTrie::create(&path).expect("create vocab");
vocab.install_overlay();
let terms = ["apple", "app", "applet", "banana", "band", "can", "candy"];
let mut expected: Vec<(Vec<u32>, u64)> = Vec::with_capacity(terms.len());
for t in &terms {
let id = vocab.insert(t).expect("overlay insert");
expected.push((t.chars().map(|c| c as u32).collect(), id));
}
let root = vocab
.lockfree_root
.as_ref()
.and_then(|r| r.load())
.expect("overlay root present");
let root_ptr = vocab
.serialize_overlay_to_disk(&root)
.expect("serialize overlay");
let (enumerated, empty) = vocab
.enumerate_overlay_terms_from_disk(root_ptr.to_raw())
.expect("enumerate overlay");
assert!(empty.is_none(), "no empty term in this fixture");
assert_eq!(enumerated.len(), terms.len(), "every term round-trips");
for (units, id) in &expected {
assert_eq!(
enumerated.get(units).copied(),
Some(Some(*id)),
"term {:?} preserved its id {}",
units,
id
);
}
drop(vocab);
}
#[test]
fn flip_checkpoint_reopen_roundtrip() {
let dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("target/test-scratch");
std::fs::create_dir_all(&dir).expect("scratch dir");
let path = dir.join(format!("vocab_flip_rt_{}.vocab", std::process::id()));
let cleanup = |p: &std::path::Path| {
let _ = std::fs::remove_file(p);
let _ = std::fs::remove_file(p.with_extension("vocab.wal"));
let _ = std::fs::remove_file(p.with_extension("vocab.idx"));
};
cleanup(&path);
let terms = ["alpha", "beta", "al", "alp", "gamma", "be", "alphabet"];
let mut expected: Vec<(String, u64)> = Vec::with_capacity(terms.len());
{
let vocab = PersistentVocabARTrie::create(&path).expect("create");
assert!(vocab.route_overlay(), "overlay is the live rep after flip");
for t in &terms {
let id = vocab.insert(t).expect("overlay insert"); expected.push((t.to_string(), id));
}
for (t, id) in &expected {
assert_eq!(vocab.get_index(t), Some(*id), "in-mem forward");
assert_eq!(
vocab.get_term(*id).as_deref(),
Some(t.as_str()),
"in-mem reverse"
);
}
vocab.checkpoint().expect("checkpoint_overlay"); }
let (vocab, _report) =
PersistentVocabARTrie::open_with_recovery(&path).expect("reopen v2 image");
assert!(
vocab.route_overlay(),
"overlay is the live rep after reopen"
);
for (t, id) in &expected {
assert_eq!(vocab.get_index(t), Some(*id), "reopened forward {t:?}");
assert_eq!(
vocab.get_term(*id).as_deref(),
Some(t.as_str()),
"reopened reverse id {id}"
);
}
assert_eq!(vocab.len(), terms.len(), "reopened entry_count");
drop(vocab);
cleanup(&path);
}
#[test]
fn cx_compressed_long_chain_reopen_roundtrip() {
let dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("target/test-scratch");
std::fs::create_dir_all(&dir).expect("scratch dir");
let path = dir.join(format!("vocab_cx_longchain_{}.vocab", std::process::id()));
let cleanup = |p: &std::path::Path| {
let _ = std::fs::remove_file(p);
let _ = std::fs::remove_file(p.with_extension("vocab.wal"));
let _ = std::fs::remove_file(p.with_extension("vocab.idx"));
};
cleanup(&path);
let alphabet = "abcdefghijklmnopqrstuvwxyz0123456789ABCDEFGHIJKLMN"; let mut terms: Vec<String> = Vec::new();
for (i, &len) in [1usize, 7, 13, 21, 40].iter().enumerate() {
let first = (b'A' + i as u8) as char;
let body: String = alphabet.chars().take(len.saturating_sub(1)).collect();
terms.push(format!("{first}{body}"));
}
terms.push("sharedprefixzzzzzzzz".to_string());
terms.push("sharedprefixqqqqqqqq".to_string());
terms.push("\u{1F3AF}".repeat(10));
let mut expected: Vec<(String, u64)> = Vec::with_capacity(terms.len());
{
let vocab = PersistentVocabARTrie::create(&path).expect("create");
assert!(vocab.route_overlay(), "overlay is the live rep");
for t in &terms {
let id = vocab.insert(t).expect("overlay insert");
expected.push((t.clone(), id));
}
vocab.checkpoint().expect("compressed checkpoint"); }
let (vocab, _report) =
PersistentVocabARTrie::open_with_recovery(&path).expect("reopen compressed v2 image");
assert!(
vocab.route_overlay(),
"overlay is the live rep after reopen"
);
for (t, id) in &expected {
assert_eq!(
vocab.get_index(t),
Some(*id),
"reopened forward {t:?} (no truncation)"
);
assert_eq!(
vocab.get_term(*id).as_deref(),
Some(t.as_str()),
"reopened reverse id {id}"
);
}
assert_eq!(vocab.len(), terms.len(), "reopened entry_count");
drop(vocab);
cleanup(&path);
}
#[test]
fn flip_crash_recovery_replays_wal_tail() {
let dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("target/test-scratch");
std::fs::create_dir_all(&dir).expect("scratch dir");
let pid = std::process::id();
let path_a = dir.join(format!("vocab_crash_a_{}.vocab", pid));
let path_b = dir.join(format!("vocab_crash_b_{}.vocab", pid));
let record_segment_dir = |p: &std::path::Path| {
let wal = p.with_extension("vocab.wal");
let mut dir = wal;
dir.set_extension("wal.d");
dir
};
let cleanup = |p: &std::path::Path| {
let _ = std::fs::remove_file(p);
let _ = std::fs::remove_file(p.with_extension("vocab.wal"));
let _ = std::fs::remove_file(p.with_extension("vocab.idx"));
let _ = std::fs::remove_dir_all(record_segment_dir(p));
};
cleanup(&path_a);
cleanup(&path_b);
let first = ["aa", "ab", "ba", "bb"];
let second = ["car", "cat", "dog", "do"]; let mut expected: Vec<(String, u64)> = Vec::new();
{
let vocab = PersistentVocabARTrie::create(&path_a).expect("create");
for t in &first {
expected.push((t.to_string(), vocab.insert(t).expect("insert")));
}
vocab.checkpoint().expect("checkpoint"); for t in &second {
expected.push((t.to_string(), vocab.insert(t).expect("insert")));
}
vocab.sync().expect("sync WAL"); std::fs::copy(&path_a, &path_b).expect("copy image");
std::fs::copy(
path_a.with_extension("vocab.wal"),
path_b.with_extension("vocab.wal"),
)
.expect("copy wal");
let src_segments = record_segment_dir(&path_a);
let dst_segments = record_segment_dir(&path_b);
if src_segments.exists() {
std::fs::create_dir_all(&dst_segments).expect("copy WAL segment dir");
for entry in std::fs::read_dir(&src_segments).expect("read WAL segment dir") {
let entry = entry.expect("WAL segment entry");
std::fs::copy(entry.path(), dst_segments.join(entry.file_name()))
.expect("copy WAL record segment");
}
}
}
let (vocab, _report) =
PersistentVocabARTrie::open_with_recovery(&path_b).expect("crash recovery reopen");
assert!(vocab.route_overlay(), "overlay live after crash recovery");
for (t, id) in &expected {
assert_eq!(vocab.get_index(t), Some(*id), "recovered forward {t:?}");
assert_eq!(
vocab.get_term(*id).as_deref(),
Some(t.as_str()),
"recovered reverse {id}"
);
}
assert_eq!(
vocab.len(),
first.len() + second.len(),
"recovered entry_count"
);
let new_id = vocab.insert("zzz").expect("post-recovery insert");
assert!(
expected.iter().all(|(_, id)| *id != new_id),
"post-recovery id {new_id} collided with a recovered id"
);
drop(vocab);
cleanup(&path_a);
cleanup(&path_b);
}
}