use alloc::format;
use alloc::vec::Vec;
use plugmem_arena::{
Arena, ArenaCfg, BlobHeap, BlobHeapBuilder, BlobHeapCfg, BlobId, ChunkPool, ChunkPoolCfg,
ListHandle, ShardMode,
};
use crate::error::Error;
use crate::id::{FactId, NONE_U32};
use crate::index::IdListIndex;
use crate::index::bm25::Bm25Index;
use crate::index::hnsw::{HnswGraph, HnswScratch};
use crate::index::vecpool::VecPool;
use crate::journal::Op;
use crate::memory::persist::Sections;
use crate::memory::shards::ShardLayout;
use crate::model::{
EdgeHistorySlot, EdgeSlot, EntityByName, EntityRecord, FactAux, FactRecord, TemporalSlot,
};
use crate::snapshot::SnapshotSink;
use crate::storage::{Scratch, Storage};
use crate::tokenizer::Tokenizer;
use super::Memory;
pub(crate) const TOKENIZER_INDEX_VERSION: u32 = 2;
const AUTO_HNSW_INSERT_BUDGET: usize = 4096;
const NO_HNSW_INSERT_LIMIT: u32 = u32::MAX;
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub enum MaintenanceMode {
#[default]
Auto,
Compact,
ReindexText,
OptimizeVectors,
Full,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct MaintenanceOptions {
pub mode: MaintenanceMode,
pub max_hnsw_inserts: Option<usize>,
}
impl Default for MaintenanceOptions {
fn default() -> Self {
Self::auto()
}
}
impl MaintenanceOptions {
pub const fn auto() -> Self {
Self {
mode: MaintenanceMode::Auto,
max_hnsw_inserts: Some(AUTO_HNSW_INSERT_BUDGET),
}
}
pub const fn full() -> Self {
Self {
mode: MaintenanceMode::Full,
max_hnsw_inserts: None,
}
}
pub(crate) fn from_journal(mode: u8, max_hnsw_inserts: u32) -> Result<Self, Error> {
let mode = match mode {
0 => MaintenanceMode::Auto,
1 => MaintenanceMode::Compact,
2 => MaintenanceMode::ReindexText,
3 => MaintenanceMode::OptimizeVectors,
4 => MaintenanceMode::Full,
_ => return Err(Error::Corrupt("journal maintenance mode is invalid")),
};
let max_hnsw_inserts = if max_hnsw_inserts == NO_HNSW_INSERT_LIMIT {
None
} else {
Some(max_hnsw_inserts as usize)
};
Ok(Self {
mode,
max_hnsw_inserts,
})
}
pub(crate) fn journal_mode(self) -> u8 {
match self.mode {
MaintenanceMode::Auto => 0,
MaintenanceMode::Compact => 1,
MaintenanceMode::ReindexText => 2,
MaintenanceMode::OptimizeVectors => 3,
MaintenanceMode::Full => 4,
}
}
pub(crate) fn journal_max_hnsw_inserts(self) -> u32 {
self.max_hnsw_inserts
.map(|n| u32::try_from(n).unwrap_or(u32::MAX - 1))
.unwrap_or(NO_HNSW_INSERT_LIMIT)
}
}
const LIVE_DENSE_SLACK: usize = 8;
impl<'a> Memory<'a> {
pub(super) fn fact_ids_ascending(&self) -> Vec<u32> {
let mut ids: Vec<u32> = self.facts.iter().map(|rec| rec.id.0).collect();
ids.sort_unstable();
ids
}
fn dense_id_span(&self) -> usize {
let cap = self
.facts
.len()
.saturating_add(1)
.saturating_mul(LIVE_DENSE_SLACK);
(self.next_fact as usize).min(cap)
}
}
fn scratch_err<E: core::fmt::Debug>(e: E) -> Error {
Error::Storage(format!("{e:?}"))
}
trait PoolSink {
fn push_text(&mut self, bytes: &[u8]) -> Result<BlobId, Error>;
fn push_vector(&mut self, src: &VecPool<'_>, slot: u32) -> Result<u32, Error>;
}
struct OwnedPools {
texts: BlobHeap<'static>,
vecs: VecPool<'static>,
}
impl PoolSink for OwnedPools {
fn push_text(&mut self, bytes: &[u8]) -> Result<BlobId, Error> {
Ok(self.texts.push(bytes)?)
}
fn push_vector(&mut self, src: &VecPool<'_>, slot: u32) -> Result<u32, Error> {
Ok(self.vecs.copy_slot(src, slot))
}
}
struct StreamPools<'s, T: Scratch, V: Scratch> {
text_scratch: &'s mut T,
text_index: BlobHeapBuilder,
vec_scratch: &'s mut V,
vec_count: u32,
}
impl<T: Scratch, V: Scratch> PoolSink for StreamPools<'_, T, V> {
fn push_text(&mut self, bytes: &[u8]) -> Result<BlobId, Error> {
self.text_scratch.write(bytes).map_err(scratch_err)?;
Ok(self.text_index.push_len(bytes.len())?)
}
fn push_vector(&mut self, src: &VecPool<'_>, slot: u32) -> Result<u32, Error> {
self.vec_scratch
.write(src.slot_bytes(slot as usize))
.map_err(scratch_err)?;
let new = self.vec_count;
self.vec_count += 1;
Ok(new)
}
}
struct RebuildMeta {
facts: Arena<'static, FactRecord>,
fact_aux: Arena<'static, FactAux>,
entities: Arena<'static, EntityRecord>,
by_name: Arena<'static, EntityByName>,
temporal: Arena<'static, TemporalSlot>,
tag_lists: ChunkPool<'static>,
metas: BlobHeap<'static>,
bm25: Bm25Index<'static>,
tags_idx: IdListIndex<'static>,
entity_facts: IdListIndex<'static>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
pub struct MaintainReport {
pub purged: usize,
pub bytes_before: usize,
pub bytes_after: usize,
pub no_op: bool,
pub tombstones_before: usize,
pub facts_before: usize,
pub facts_after: usize,
pub vectors_before: usize,
pub vectors_after: usize,
pub hnsw_indexed_before: u32,
pub hnsw_indexed_after: u32,
pub shards_before: ShardLayout,
pub shards_after: ShardLayout,
pub structural_compacted: bool,
pub bm25_compacted: bool,
pub bm25_reindexed: bool,
pub hnsw_rebuilt: bool,
pub hnsw_remapped: bool,
pub hnsw_inserted: u32,
pub edges_compacted: bool,
pub edges_before: usize,
pub edge_versions_before: usize,
}
struct RepackedEdges {
out: Arena<'static, EdgeSlot>,
inn: Arena<'static, EdgeSlot>,
hist_out: Arena<'static, EdgeHistorySlot>,
hist_in: Arena<'static, EdgeHistorySlot>,
}
impl RepackedEdges {
fn pool_bytes(&self) -> usize {
self.out.pool_bytes()
+ self.inn.pool_bytes()
+ self.hist_out.pool_bytes()
+ self.hist_in.pool_bytes()
}
}
struct Rebuilt {
facts: Arena<'static, FactRecord>,
entities: Arena<'static, EntityRecord>,
by_name: Arena<'static, EntityByName>,
fact_aux: Arena<'static, FactAux>,
texts: BlobHeap<'static>,
metas: BlobHeap<'static>,
tag_lists: ChunkPool<'static>,
bm25: Bm25Index<'static>,
tags_idx: IdListIndex<'static>,
entity_facts: IdListIndex<'static>,
temporal: Arena<'static, TemporalSlot>,
vecs: VecPool<'static>,
hnsw: HnswGraph<'static>,
edges: Option<RepackedEdges>,
bm25_tokenizer_version: u32,
layout: ShardLayout,
report: MaintainReport,
}
#[derive(Clone, Copy)]
enum Bm25Policy {
Compact,
Reindex,
}
#[derive(Clone, Copy)]
struct WorkPlan {
compact: bool,
bm25_reindex: bool,
optimize_vectors: bool,
hnsw_full_rebuild: bool,
repack_edges: bool,
max_hnsw_inserts: Option<usize>,
layout: ShardLayout,
}
#[derive(Clone, Copy, Default)]
struct GraphWork {
rebuilt: bool,
remapped: bool,
inserted: u32,
}
impl WorkPlan {
fn needs_work(self) -> bool {
self.compact || self.bm25_reindex || self.optimize_vectors || self.repack_edges
}
}
fn hnsw_target(start: u32, total: u32, max_hnsw_inserts: Option<usize>) -> u32 {
let Some(max) = max_hnsw_inserts else {
return total;
};
let max = u32::try_from(max).unwrap_or(u32::MAX);
start.saturating_add(max).min(total)
}
impl Memory<'_> {
pub fn maintain<S: Storage>(
&mut self,
store: &mut S,
now: u64,
) -> Result<MaintainReport, Error> {
self.maintain_with_options(store, now, MaintenanceOptions::auto())
}
pub fn maintain_with_options<S: Storage>(
&mut self,
store: &mut S,
now: u64,
options: MaintenanceOptions,
) -> Result<MaintainReport, Error> {
let plan = self.work_plan(options);
let bytes_before = self.satellite_bytes(plan.repack_edges);
let mut report = self.report_skeleton(bytes_before);
if !plan.needs_work() {
report.no_op = true;
return Ok(report);
}
if !plan.compact {
let mut bm25 = None;
if plan.bm25_reindex {
bm25 = Some(self.reindex_bm25_from_text()?);
report.bm25_reindexed = true;
}
let mut hnsw = None;
if plan.optimize_vectors {
let (graph, work) =
self.optimize_graph(plan.hnsw_full_rebuild, plan.max_hnsw_inserts)?;
hnsw = Some(graph);
report.hnsw_rebuilt = work.rebuilt;
report.hnsw_remapped = work.remapped;
report.hnsw_inserted = work.inserted;
}
let mut entry = Vec::new();
Op::Maintain {
now,
mode: options.journal_mode(),
max_hnsw_inserts: options.journal_max_hnsw_inserts(),
}
.encode(&mut entry);
store
.append_journal(&entry)
.map_err(|e| Error::Storage(format!("{e:?}")))?;
if let Some(bm25) = bm25 {
self.bm25 = bm25;
self.bm25_tokenizer_version = TOKENIZER_INDEX_VERSION;
}
if let Some(hnsw) = hnsw {
self.hnsw = hnsw;
}
report.bytes_after = self.satellite_bytes(plan.repack_edges);
report.hnsw_indexed_after = self.hnsw.indexed();
return Ok(report);
}
let (rebuilt, _) = self.rebuild(plan)?;
report = rebuilt.report.clone();
let mut entry = Vec::new();
Op::Maintain {
now,
mode: options.journal_mode(),
max_hnsw_inserts: options.journal_max_hnsw_inserts(),
}
.encode(&mut entry);
store
.append_journal(&entry)
.map_err(|e| Error::Storage(format!("{e:?}")))?;
self.install(rebuilt);
Ok(report)
}
pub(super) fn replay_maintain_with_options(
&mut self,
options: MaintenanceOptions,
) -> Result<(), Error> {
let plan = self.work_plan(options);
if !plan.needs_work() {
return Ok(());
}
if !plan.compact {
if plan.bm25_reindex {
self.bm25 = self.reindex_bm25_from_text()?;
self.bm25_tokenizer_version = TOKENIZER_INDEX_VERSION;
}
if plan.optimize_vectors {
let (hnsw, _) =
self.optimize_graph(plan.hnsw_full_rebuild, plan.max_hnsw_inserts)?;
self.hnsw = hnsw;
}
return Ok(());
}
let (rebuilt, _) = self.rebuild(plan)?;
self.install(rebuilt);
Ok(())
}
fn satellite_bytes(&self, with_edges: bool) -> usize {
let edges = if with_edges { self.edge_bytes() } else { 0 };
edges
+ self.facts.pool_bytes()
+ self.fact_aux.pool_bytes()
+ self.entities.pool_bytes()
+ self.hnsw.pool_bytes()
+ self.texts.pool_bytes()
+ self.metas.pool_bytes()
+ self.tag_lists.pool_bytes()
+ self.bm25.pool_bytes()
+ self.tags_idx.pool_bytes()
+ self.entity_facts.pool_bytes()
+ self.temporal.pool_bytes()
+ self.vecs.pool_bytes()
}
fn edge_bytes(&self) -> usize {
self.edges_out.pool_bytes()
+ self.edges_in.pool_bytes()
+ self.edges_hist_out.pool_bytes()
+ self.edges_hist_in.pool_bytes()
}
pub fn maintenance_needed(&self, options: MaintenanceOptions) -> bool {
self.work_plan(options).needs_work()
}
pub fn maintenance_preview(
&self,
options: MaintenanceOptions,
bytes_before: usize,
) -> MaintainReport {
let mut report = self.report_skeleton(bytes_before);
if !self.maintenance_needed(options) {
report.no_op = true;
}
report
}
fn report_skeleton(&self, bytes_before: usize) -> MaintainReport {
MaintainReport {
purged: 0,
bytes_before,
bytes_after: bytes_before,
no_op: false,
tombstones_before: self.tombstones,
facts_before: self.facts.len(),
facts_after: self.facts.len(),
vectors_before: self.vecs.len(),
vectors_after: self.vecs.len(),
hnsw_indexed_before: self.hnsw.indexed(),
hnsw_indexed_after: self.hnsw.indexed(),
shards_before: ShardLayout::of_config(&self.cfg),
shards_after: ShardLayout::of_config(&self.cfg),
structural_compacted: false,
bm25_compacted: false,
bm25_reindexed: false,
hnsw_rebuilt: false,
hnsw_remapped: false,
hnsw_inserted: 0,
edges_compacted: false,
edges_before: self.edges_out.len(),
edge_versions_before: self.edges_hist_out.len(),
}
}
fn work_plan(&self, options: MaintenanceOptions) -> WorkPlan {
let tokenizer_stale = self.bm25_tokenizer_version != TOKENIZER_INDEX_VERSION;
let has_tombstones = self.tombstones != 0;
let compaction_due = has_tombstones || self.bm25.needs_resummarize();
let graph_tail = self.cfg.dim != 0
&& self.vecs.len() >= self.cfg.flat_to_hnsw
&& self.hnsw.indexed() < self.vecs.len() as u32;
let layout = self.target_layout();
let stored_layout = ShardLayout::of_config(&self.cfg);
let relayout_due = stored_layout.compacted_groups_earn_rebuild(&layout);
let edges_need_relayout = stored_layout.edges_earn_rebuild(&layout);
match options.mode {
MaintenanceMode::Auto => WorkPlan {
compact: compaction_due || relayout_due,
bm25_reindex: tokenizer_stale,
optimize_vectors: graph_tail,
hnsw_full_rebuild: false,
repack_edges: edges_need_relayout,
max_hnsw_inserts: options.max_hnsw_inserts,
layout,
},
MaintenanceMode::Compact => WorkPlan {
compact: compaction_due || relayout_due,
bm25_reindex: tokenizer_stale,
optimize_vectors: has_tombstones && self.hnsw.indexed() != 0,
hnsw_full_rebuild: false,
repack_edges: edges_need_relayout,
max_hnsw_inserts: options.max_hnsw_inserts,
layout,
},
MaintenanceMode::ReindexText => WorkPlan {
compact: compaction_due || relayout_due,
bm25_reindex: true,
optimize_vectors: has_tombstones && self.hnsw.indexed() != 0,
hnsw_full_rebuild: false,
repack_edges: edges_need_relayout,
max_hnsw_inserts: options.max_hnsw_inserts,
layout,
},
MaintenanceMode::OptimizeVectors => WorkPlan {
compact: false,
bm25_reindex: false,
optimize_vectors: self.cfg.dim != 0
&& self.vecs.len() >= self.cfg.flat_to_hnsw
&& (graph_tail || self.hnsw.indexed() == 0),
hnsw_full_rebuild: self.hnsw.indexed() == 0,
repack_edges: false,
max_hnsw_inserts: options.max_hnsw_inserts,
layout,
},
MaintenanceMode::Full => WorkPlan {
compact: true,
bm25_reindex: true,
optimize_vectors: self.cfg.dim != 0 && self.vecs.len() >= self.cfg.flat_to_hnsw,
hnsw_full_rebuild: true,
repack_edges: !self.edges_hist_out.is_empty() || edges_need_relayout,
max_hnsw_inserts: None,
layout,
},
}
}
fn repack_edges(&self, layout: &ShardLayout) -> Result<RepackedEdges, Error> {
let ord =
ArenaCfg::new(layout.edges, ShardMode::Ordered).with_max_bytes(self.cfg.max_bytes);
let mut out = Arena::new(ord)?;
let mut inn = Arena::new(ord)?;
let mut hist_out = Arena::new(ord)?;
let mut hist_in = Arena::new(ord)?;
for edge in self.edges_out.iter() {
out.insert(&edge)?;
}
for edge in self.edges_in.iter() {
inn.insert(&edge)?;
}
for version in self.edges_hist_out.iter() {
hist_out.insert(&version)?;
}
for version in self.edges_hist_in.iter() {
hist_in.insert(&version)?;
}
Ok(RepackedEdges {
out,
inn,
hist_out,
hist_in,
})
}
fn reindex_bm25_from_text(&self) -> Result<Bm25Index<'static>, Error> {
let cfg = &self.cfg;
let mut bm25 = Bm25Index::new(cfg.shards_postings, cfg.max_bytes)?;
let mut tokenizer = Tokenizer::new();
let mut tf: Vec<(u32, u8)> = Vec::new();
for fid in self.fact_ids_ascending() {
let id = FactId(fid);
let Some(rec) = self.facts.get(&fid.to_be_bytes()) else {
continue;
};
if rec.is_tombstone() {
continue;
}
let text = core::str::from_utf8(self.texts.get(rec.text))
.map_err(|_| Error::Corrupt("maintain: fact text is not UTF-8"))?;
tf.clear();
let terms = &self.terms;
let tf_ref = &mut tf;
tokenizer.tokenize(text, &mut |token| {
if let Some(term) = terms.lookup(token) {
match tf_ref.iter_mut().find(|(t, _)| *t == term.0) {
Some((_, c)) => *c = c.saturating_add(1),
None => tf_ref.push((term.0, 1)),
}
}
});
bm25.index_doc(id, &tf)?;
}
Ok(bm25)
}
fn optimize_graph(
&self,
full_rebuild: bool,
max_hnsw_inserts: Option<usize>,
) -> Result<(HnswGraph<'static>, GraphWork), Error> {
let total = self.vecs.len() as u32;
let mut graph = if full_rebuild || self.hnsw.indexed() == 0 {
HnswGraph::new(self.cfg.hnsw_m, self.cfg.hnsw_m0, self.cfg.max_bytes)?
} else {
self.hnsw.to_owned(self.cfg.max_bytes)?
};
let start = graph.indexed();
let target = hnsw_target(start, total, max_hnsw_inserts);
let mut scratch = HnswScratch::default();
graph.insert_bulk(
&self.vecs,
target,
self.cfg.hnsw_ef_construction,
&mut scratch,
)?;
Ok((
graph,
GraphWork {
rebuilt: full_rebuild || self.hnsw.indexed() == 0,
remapped: !full_rebuild && self.hnsw.indexed() != 0,
inserted: target.saturating_sub(start),
},
))
}
fn rebuild(&self, plan: WorkPlan) -> Result<(Rebuilt, usize), Error> {
let cfg = &self.cfg;
let blob = BlobHeapCfg::new()
.with_max_bytes(cfg.max_bytes)
.with_max_blob(cfg.max_blob);
let mut pools = OwnedPools {
texts: BlobHeap::new(blob),
vecs: VecPool::new(cfg.dim, cfg.max_bytes),
};
let bm25_policy = if plan.bm25_reindex {
Bm25Policy::Reindex
} else {
Bm25Policy::Compact
};
let (m, vec_map, purged) = self.rebuild_parts(&mut pools, bm25_policy, &plan.layout)?;
let (hnsw, graph_work) = self.rebuild_graph(
&vec_map,
&pools.vecs,
plan.hnsw_full_rebuild,
plan.max_hnsw_inserts,
)?;
let edges = plan
.repack_edges
.then(|| self.repack_edges(&plan.layout))
.transpose()?;
let mut report = self.report_skeleton(self.satellite_bytes(plan.repack_edges));
report.purged = purged;
report.bytes_after = edges.as_ref().map_or(0, RepackedEdges::pool_bytes)
+ m.facts.pool_bytes()
+ m.fact_aux.pool_bytes()
+ m.entities.pool_bytes()
+ hnsw.pool_bytes()
+ pools.texts.pool_bytes()
+ m.metas.pool_bytes()
+ m.tag_lists.pool_bytes()
+ m.bm25.pool_bytes()
+ m.tags_idx.pool_bytes()
+ m.entity_facts.pool_bytes()
+ m.temporal.pool_bytes()
+ pools.vecs.pool_bytes();
report.facts_after = m.facts.len();
report.vectors_after = pools.vecs.len();
report.hnsw_indexed_after = hnsw.indexed();
report.structural_compacted = true;
report.bm25_compacted = matches!(bm25_policy, Bm25Policy::Compact);
report.bm25_reindexed = matches!(bm25_policy, Bm25Policy::Reindex);
report.hnsw_rebuilt = graph_work.rebuilt;
report.hnsw_remapped = graph_work.remapped;
report.hnsw_inserted = graph_work.inserted;
report.edges_compacted = edges.is_some();
let layout = ShardLayout::of_config(&self.cfg).realized(&plan.layout, edges.is_some());
report.shards_before = ShardLayout::of_config(&self.cfg);
report.shards_after = layout;
Ok((
Rebuilt {
facts: m.facts,
entities: m.entities,
by_name: m.by_name,
fact_aux: m.fact_aux,
texts: pools.texts,
metas: m.metas,
tag_lists: m.tag_lists,
bm25: m.bm25,
tags_idx: m.tags_idx,
entity_facts: m.entity_facts,
temporal: m.temporal,
vecs: pools.vecs,
hnsw,
edges,
bm25_tokenizer_version: if plan.bm25_reindex {
TOKENIZER_INDEX_VERSION
} else {
self.bm25_tokenizer_version
},
layout,
report,
},
purged,
))
}
fn rebuild_parts<P: PoolSink>(
&self,
pools: &mut P,
bm25_policy: Bm25Policy,
layout: &ShardLayout,
) -> Result<(RebuildMeta, alloc::vec::Vec<u32>, usize), Error> {
let cfg = &self.cfg;
let uni =
|shards: usize| ArenaCfg::new(shards, ShardMode::Uniform).with_max_bytes(cfg.max_bytes);
let ord =
|shards: usize| ArenaCfg::new(shards, ShardMode::Ordered).with_max_bytes(cfg.max_bytes);
let mut entities = Arena::new(uni(layout.entities))?;
let mut by_name = Arena::new(ord(layout.entities))?;
for entry in self.by_name.iter() {
by_name.insert(&entry)?;
}
let mut facts = Arena::new(uni(layout.facts))?;
let mut fact_aux = Arena::new(uni(layout.facts))?;
let mut tag_lists = ChunkPool::new(ChunkPoolCfg::new().with_max_bytes(cfg.max_bytes));
let dense = self.dense_id_span();
let mut live = alloc::vec![false; dense];
for rec in self.facts.iter() {
let at = rec.id.0 as usize;
if at < dense && !rec.is_tombstone() {
live[at] = true;
}
}
let is_live = |id: FactId| match live.get(id.0 as usize) {
Some(&flag) => flag,
None => self
.facts
.get(&id.0.to_be_bytes())
.is_some_and(|rec| !rec.is_tombstone()),
};
let mut bm25 = match bm25_policy {
Bm25Policy::Compact => {
self.bm25
.compact_live(layout.postings, cfg.max_bytes, is_live)?
}
Bm25Policy::Reindex => Bm25Index::new(layout.postings, cfg.max_bytes)?,
};
let mut tags_idx = IdListIndex::new(layout.postings, cfg.max_bytes)?;
let mut entity_facts = IdListIndex::new(layout.entities, cfg.max_bytes)?;
let mut temporal = Arena::new(ord(layout.temporal))?;
let mut metas = BlobHeap::new(
BlobHeapCfg::new()
.with_max_bytes(cfg.max_bytes)
.with_max_blob(cfg.max_blob),
);
for eid in 0..self.next_entity {
let rec = self
.entities
.get(&eid.to_be_bytes())
.ok_or(Error::Corrupt("maintain: entity id gap"))?;
let name_id = pools.push_text(self.texts.get(rec.name))?;
entities.insert(&EntityRecord {
name: name_id,
..rec
})?;
}
let mut tokenizer = Tokenizer::new();
let mut tf: Vec<(u32, u8)> = Vec::new();
let mut vec_map = alloc::vec![NONE_U32; self.vecs.len()];
let mut purged = 0usize;
for fid in self.fact_ids_ascending() {
let id = FactId(fid);
let Some(rec) = self.facts.get(&fid.to_be_bytes()) else {
continue;
};
if rec.is_tombstone() {
purged += 1;
continue;
}
let text_bytes = self.texts.get(rec.text);
let text_id = pools.push_text(text_bytes)?;
if matches!(bm25_policy, Bm25Policy::Reindex) {
let text = core::str::from_utf8(text_bytes)
.map_err(|_| Error::Corrupt("maintain: fact text is not UTF-8"))?;
tf.clear();
let terms = &self.terms;
let tf_ref = &mut tf;
tokenizer.tokenize(text, &mut |token| {
if let Some(term) = terms.lookup(token) {
match tf_ref.iter_mut().find(|(t, _)| *t == term.0) {
Some((_, c)) => *c = c.saturating_add(1),
None => tf_ref.push((term.0, 1)),
}
}
});
bm25.index_doc(id, &tf)?;
}
let aux = self
.fact_aux
.get(&fid.to_be_bytes())
.ok_or(Error::Corrupt("maintain: fact aux gap"))?;
let mut tags = ListHandle::EMPTY;
for chunk in self.tag_lists.iter(&aux.tags) {
for raw in chunk.chunks_exact(4) {
let term = u32::from_be_bytes(raw.try_into().unwrap());
tag_lists.push(&mut tags, &term.to_be_bytes())?;
tags_idx.push(term, id, 0)?;
}
}
let meta = if aux.meta.0 == NONE_U32 {
BlobId(NONE_U32)
} else {
metas.push(self.metas.get(aux.meta))?
};
fact_aux.insert(&FactAux { id, tags, meta })?;
if let Some(entity) = rec.entity.some() {
entity_facts.push(entity.0, id, 0)?;
}
temporal.insert(&TemporalSlot {
recorded_at: rec.recorded_at,
fact: id,
})?;
let vector = if rec.has_vector() {
let new_slot = pools.push_vector(&self.vecs, rec.vector)?;
vec_map[rec.vector as usize] = new_slot;
new_slot
} else {
NONE_U32
};
facts.insert(&FactRecord {
text: text_id,
vector,
..rec
})?;
}
Ok((
RebuildMeta {
facts,
by_name,
fact_aux,
entities,
temporal,
tag_lists,
metas,
bm25,
tags_idx,
entity_facts,
},
vec_map,
purged,
))
}
pub fn snapshot_disk_first<T: Scratch, V: Scratch, Sk: SnapshotSink>(
&self,
created_at: u64,
text_scratch: &mut T,
vec_scratch: &mut V,
sink: Sk,
) -> Result<usize, Error> {
let options = MaintenanceOptions::auto();
if !self.maintenance_needed(options) {
self.write_snapshot_with(&self.sections(), created_at, sink)?;
return Ok(0);
}
Ok(self
.snapshot_disk_first_with_options(created_at, text_scratch, vec_scratch, sink, options)?
.purged)
}
pub fn snapshot_disk_first_with_options<T: Scratch, V: Scratch, Sk: SnapshotSink>(
&self,
created_at: u64,
text_scratch: &mut T,
vec_scratch: &mut V,
sink: Sk,
options: MaintenanceOptions,
) -> Result<MaintainReport, Error> {
let plan = self.work_plan(options);
let mut report = self.report_skeleton(self.satellite_bytes(plan.repack_edges));
if !plan.needs_work() {
report.no_op = true;
return Ok(report);
}
let cfg = &self.cfg;
let blob = BlobHeapCfg::new()
.with_max_bytes(cfg.max_bytes)
.with_max_blob(cfg.max_blob);
let mut pools = StreamPools {
text_scratch,
text_index: BlobHeapBuilder::new(blob),
vec_scratch,
vec_count: 0,
};
let bm25_policy = if plan.bm25_reindex {
Bm25Policy::Reindex
} else {
Bm25Policy::Compact
};
let (m, vec_map, purged) = self.rebuild_parts(&mut pools, bm25_policy, &plan.layout)?;
let StreamPools {
text_scratch,
text_index,
vec_scratch,
..
} = pools;
let mut text_index_bytes = Vec::new();
text_index.dump_index(&mut text_index_bytes);
let text_pool = text_scratch.freeze().map_err(scratch_err)?;
let vec_pool = vec_scratch.freeze().map_err(scratch_err)?;
let texts = BlobHeap::load_borrowed(blob, &text_index_bytes, text_pool)?;
let vecs = VecPool::from_parts_borrowed(cfg.dim, cfg.max_bytes, vec_pool)?;
let (hnsw, graph_work) = self.rebuild_graph(
&vec_map,
&vecs,
plan.hnsw_full_rebuild,
plan.max_hnsw_inserts,
)?;
let edges = plan
.repack_edges
.then(|| self.repack_edges(&plan.layout))
.transpose()?;
let sections = Sections {
facts: &m.facts,
fact_aux: &m.fact_aux,
entities: &m.entities,
by_name: &m.by_name,
temporal: &m.temporal,
texts: &texts,
metas: &m.metas,
tag_lists: &m.tag_lists,
bm25: &m.bm25,
tags_idx: &m.tags_idx,
entity_facts: &m.entity_facts,
vecs: &vecs,
hnsw: &hnsw,
edges_out: edges.as_ref().map_or(&self.edges_out, |e| &e.out),
edges_in: edges.as_ref().map_or(&self.edges_in, |e| &e.inn),
edges_hist_out: edges.as_ref().map_or(&self.edges_hist_out, |e| &e.hist_out),
edges_hist_in: edges.as_ref().map_or(&self.edges_hist_in, |e| &e.hist_in),
layout: ShardLayout::of_config(&self.cfg).realized(&plan.layout, edges.is_some()),
};
self.write_snapshot_with(§ions, created_at, sink)?;
report.purged = purged;
report.edges_compacted = edges.is_some();
report.bytes_after = edges.as_ref().map_or(0, RepackedEdges::pool_bytes)
+ m.facts.pool_bytes()
+ m.fact_aux.pool_bytes()
+ m.entities.pool_bytes()
+ hnsw.pool_bytes()
+ texts.pool_bytes()
+ m.metas.pool_bytes()
+ m.tag_lists.pool_bytes()
+ m.bm25.pool_bytes()
+ m.tags_idx.pool_bytes()
+ m.entity_facts.pool_bytes()
+ m.temporal.pool_bytes()
+ vecs.pool_bytes();
report.facts_after = m.facts.len();
report.vectors_after = vecs.len();
report.hnsw_indexed_after = hnsw.indexed();
report.structural_compacted = true;
report.bm25_compacted = matches!(bm25_policy, Bm25Policy::Compact);
report.bm25_reindexed = matches!(bm25_policy, Bm25Policy::Reindex);
report.hnsw_rebuilt = graph_work.rebuilt;
report.hnsw_remapped = graph_work.remapped;
report.hnsw_inserted = graph_work.inserted;
Ok(report)
}
fn rebuild_graph(
&self,
vec_map: &[u32],
pool: &VecPool<'_>,
full_rebuild: bool,
max_hnsw_inserts: Option<usize>,
) -> Result<(HnswGraph<'static>, GraphWork), Error> {
let cfg = &self.cfg;
let mut graph: HnswGraph<'static> = HnswGraph::new(cfg.hnsw_m, cfg.hnsw_m0, cfg.max_bytes)?;
let total = pool.len() as u32;
if cfg.dim == 0 || (total as usize) < cfg.flat_to_hnsw {
return Ok((graph, GraphWork::default()));
}
let old_indexed = self.hnsw.indexed() as usize;
let dead = vec_map[..old_indexed]
.iter()
.filter(|&&m| m == NONE_U32)
.count();
let mut scratch = HnswScratch::default();
let mut work = GraphWork::default();
if old_indexed > 0 && !full_rebuild {
graph = self.hnsw.remapped(vec_map, pool, cfg.max_bytes)?;
work.remapped = true;
} else if old_indexed > 0 || total > 0 {
work.rebuilt = true;
}
if full_rebuild && dead * 10 > old_indexed {
work.rebuilt = true;
}
let start = graph.indexed();
let target = hnsw_target(start, total, max_hnsw_inserts);
graph.insert_bulk(pool, target, cfg.hnsw_ef_construction, &mut scratch)?;
work.inserted = target.saturating_sub(start);
Ok((graph, work))
}
fn install(&mut self, r: Rebuilt) {
self.facts = r.facts;
self.entities = r.entities;
self.by_name = r.by_name;
self.fact_aux = r.fact_aux;
self.texts = r.texts;
self.metas = r.metas;
self.tag_lists = r.tag_lists;
self.bm25 = r.bm25;
self.tags_idx = r.tags_idx;
self.entity_facts = r.entity_facts;
self.temporal = r.temporal;
self.vecs = r.vecs;
self.hnsw = r.hnsw;
if let Some(edges) = r.edges {
self.edges_out = edges.out;
self.edges_in = edges.inn;
self.edges_hist_out = edges.hist_out;
self.edges_hist_in = edges.hist_in;
}
self.tombstones = 0;
self.bm25_tokenizer_version = r.bm25_tokenizer_version;
r.layout.apply(&mut self.cfg);
}
}