use super::{codec_deser, codec_ser};
use crate::graph::algorithms::text_index::{Posting, TextIndex};
use crate::graph::index_freshness::IndexFreshness;
use crate::graph::schema::DirGraph;
use crate::graph::text_indexes::{self, TextIndexRead};
use crate::serde_codec;
use serde::{Deserialize, Serialize};
use std::io;
pub(super) const TEXT_INDEX_MAGIC: &[u8; 8] = b"KGLTIDX1";
const TEXT_INDEX_FORMAT_VERSION: u32 = 1;
struct HeldIndex<'a> {
node_type: &'a str,
property: &'a str,
resolved_field: &'a str,
skipped: usize,
watermark: u32,
limit: usize,
dirty: Vec<u32>,
guard: TextIndexRead<'a>,
}
#[derive(Serialize)]
struct PersistedTextIndexRef<'a> {
node_type: &'a str,
property: &'a str,
resolved_field: &'a str,
skipped: usize,
watermark: u32,
limit: usize,
dirty: Vec<u32>,
terms: Vec<(&'a str, &'a [Posting])>,
empty_docs: Vec<u32>,
}
#[derive(Serialize, Deserialize)]
struct PersistedTextIndex {
node_type: String,
property: String,
resolved_field: String,
skipped: usize,
watermark: u32,
limit: usize,
dirty: Vec<u32>,
terms: Vec<(String, Vec<Posting>)>,
empty_docs: Vec<u32>,
}
pub(super) fn encode_text_indexes(graph: &DirGraph) -> io::Result<Option<Vec<u8>>> {
let stores = text_indexes::list_text_indexes(graph);
if stores.is_empty() {
return Ok(None);
}
let held: Vec<HeldIndex<'_>> = stores
.into_iter()
.map(|(node_type, property, store)| {
let (watermark, limit, dirty) = store.freshness_state().persisted_parts();
HeldIndex {
node_type,
property,
resolved_field: store.resolved_field(),
skipped: store.skipped(),
watermark,
limit,
dirty,
guard: store.read(),
}
})
.collect();
let entries: Vec<PersistedTextIndexRef<'_>> = held
.iter()
.map(|held| {
let index = held.guard.index();
let mut terms: Vec<(&str, &[Posting])> = index.iter_terms().collect();
terms.sort_unstable_by_key(|(term, _)| *term);
let mut empty_docs: Vec<u32> = index
.doc_slots()
.filter(|slot| index.doc_len(*slot) == Some(0))
.collect();
empty_docs.sort_unstable();
PersistedTextIndexRef {
node_type: held.node_type,
property: held.property,
resolved_field: held.resolved_field,
skipped: held.skipped,
watermark: held.watermark,
limit: held.limit,
dirty: held.dirty.clone(),
terms,
empty_docs,
}
})
.collect();
let body = codec_ser(serde_codec::CodecVersion::PostcardV1, &entries)?;
let mut payload = Vec::with_capacity(12 + body.len());
payload.extend_from_slice(TEXT_INDEX_MAGIC);
payload.extend_from_slice(&TEXT_INDEX_FORMAT_VERSION.to_le_bytes());
payload.extend_from_slice(&body);
Ok(Some(payload))
}
pub(super) fn decode_text_indexes(payload: &[u8], graph: &mut DirGraph) {
if payload.len() < 12 || &payload[..8] != TEXT_INDEX_MAGIC {
return;
}
let ver = u32::from_le_bytes([payload[8], payload[9], payload[10], payload[11]]);
if ver != TEXT_INDEX_FORMAT_VERSION {
return; }
let codec = serde_codec::CodecVersion::PostcardV1;
let entries: Vec<PersistedTextIndex> =
match codec_deser(codec, &payload[12..], (payload.len() - 12) as u64) {
Ok(entries) => entries,
Err(_) => return,
};
for entry in entries {
if entry.dirty.iter().any(|slot| *slot >= entry.watermark) {
continue;
}
if !graph.has_node_type(&entry.node_type) {
continue;
}
let index = TextIndex::from_terms(entry.terms, &entry.empty_docs);
if index.validate().is_err() {
continue;
}
let freshness = IndexFreshness::restored(entry.watermark, entry.limit, &entry.dirty);
text_indexes::attach_persisted_text_index(
graph,
&entry.node_type,
&entry.property,
index,
freshness,
entry.resolved_field,
entry.skipped,
);
}
}
#[cfg(test)]
#[path = "text_index_persistence_tests.rs"]
mod tests;