use crate::columns::ColumnStore;
use crate::edge_props::EdgeProps;
use crate::idmap::IdMap;
use crate::interner::Interner;
use crate::snapshot::PerRuleIvfState;
use crate::topology::Topology;
use crate::types::{GraphError, Result, Value};
use crate::v8::layout::{
ColumnData, ColumnsData, CsrAdjMap, CsrData, CsrEtype, CsrRow, EdgePropEntry, EdgePropsData,
FieldEntry, HnswRuleEntry, HnswSectionData, IdMapData, InternerData, ProvenanceEntry,
ProvenanceSectionData, RuleFireEntry, RuleTripEntry, RulesMetaData, StringTableData, Triple,
ViewsSectionData,
};
use crate::v8::{
HEADER_SIZE, SECTION_COLUMNS, SECTION_EDGE_PROPS, SECTION_HNSW, SECTION_IDS, SECTION_IVF_STATE,
SECTION_LAST_CHANGE, SECTION_META, SECTION_PROVENANCE, SECTION_RULES_META, SECTION_STRINGS,
SECTION_SYMS, SECTION_TOPOLOGY, SECTION_VIEWS,
};
use serde::{Deserialize, Serialize};
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::io::Write;
#[derive(Serialize, Deserialize, Default)]
pub struct V8Meta {
pub labels: Vec<u32>,
pub wal_truncated: bool,
#[serde(skip)]
pub edge_props: EdgeProps,
#[serde(skip)]
pub rule_defs: Vec<Vec<u8>>,
#[serde(skip)]
pub provenance: BTreeMap<String, BTreeSet<(u32, u32, u32)>>,
#[serde(skip)]
pub rule_tripped: BTreeMap<String, bool>,
#[serde(skip)]
pub rule_fires: BTreeMap<String, u64>,
#[serde(skip)]
pub view_defs: Vec<Vec<u8>>,
#[serde(skip)]
pub hnsw: BTreeMap<String, (Vec<u8>, Vec<u8>)>,
#[serde(skip)]
pub ivf_bytes: Vec<u8>,
#[serde(skip)]
pub last_change: HashMap<u32, u64>,
}
#[allow(clippy::too_many_arguments)]
pub fn encode_v8<W: Write>(
base_topo: Option<&crate::v8::layout::ArchivedCsr>,
base_cols: Option<&crate::v8::layout::ArchivedColumns>,
base_strings: Option<&crate::v8::layout::ArchivedStringTable>,
base_edge_props: Option<(&crate::v8::layout::ArchivedEdgeProps, &[u8])>,
base_provenance_raw: Option<&[u8]>,
topo: &Topology,
cols: &ColumnStore,
ids: &IdMap,
syms: &Interner,
meta: &V8Meta,
out: &mut W,
) -> Result<()> {
let topo_bytes = match base_topo {
Some(archived_csr) => rkyv_encode(&topology_merge_to_csr(archived_csr, topo))?,
None => rkyv_encode(&topology_to_csr(topo))?,
};
let (cols_data, strings_data) = match base_cols {
Some(archived) => columns_merge_to_data(archived, base_strings, cols)?,
None => columnstore_to_data(cols)?,
};
let cols_bytes = rkyv_encode(&cols_data)?;
let strings_bytes = rkyv_encode(&strings_data)?;
drop(cols_data);
drop(strings_data);
let ids_bytes = rkyv_encode(&idmap_to_data(ids))?;
let syms_bytes = rkyv_encode(&interner_to_data(syms))?;
let meta_bytes = bincode::serialize(meta).map_err(|e| GraphError::Corrupt {
detail: format!("v8: meta bincode serialize: {e}"),
})?;
let edge_props_bytes_owned: Vec<u8>;
let edge_props_bytes: &[u8] = match base_edge_props {
Some((_archived, raw)) if meta.edge_props.is_clean() => {
raw
}
Some((archived, _)) => {
edge_props_bytes_owned =
rkyv_encode(&edge_props_merge_to_data(archived, &meta.edge_props))?;
&edge_props_bytes_owned
}
None => {
edge_props_bytes_owned = rkyv_encode(&edge_props_to_data(&meta.edge_props))?;
&edge_props_bytes_owned
}
};
let hnsw_bytes = rkyv_encode(&hnsw_to_data(&meta.hnsw))?;
let prov_bytes_owned: Vec<u8>;
let prov_bytes: &[u8] = match base_provenance_raw {
Some(raw) if meta.provenance.is_empty() => {
raw
}
_ => {
prov_bytes_owned = rkyv_encode(&provenance_to_data(&meta.provenance))?;
&prov_bytes_owned
}
};
let rules_meta_bytes = rkyv_encode(&rules_meta_to_data(
&meta.rule_defs,
&meta.rule_tripped,
&meta.rule_fires,
))?;
let views_bytes = rkyv_encode(&ViewsSectionData {
view_defs: meta.view_defs.clone(),
})?;
let ivf_bytes: &[u8] = &meta.ivf_bytes;
let last_change_bytes_owned: Vec<u8> =
bincode::serialize(&meta.last_change).map_err(|e| GraphError::Corrupt {
detail: format!("v8: last_change bincode serialize: {e}"),
})?;
let last_change_bytes: &[u8] = &last_change_bytes_owned;
let sections: &[(u8, &[u8])] = &[
(SECTION_TOPOLOGY, &topo_bytes),
(SECTION_COLUMNS, &cols_bytes),
(SECTION_STRINGS, &strings_bytes),
(SECTION_IDS, &ids_bytes),
(SECTION_SYMS, &syms_bytes),
(SECTION_META, &meta_bytes),
(SECTION_EDGE_PROPS, edge_props_bytes),
(SECTION_HNSW, &hnsw_bytes),
(SECTION_PROVENANCE, prov_bytes),
(SECTION_RULES_META, &rules_meta_bytes),
(SECTION_VIEWS, &views_bytes),
(SECTION_IVF_STATE, ivf_bytes),
(SECTION_LAST_CHANGE, last_change_bytes),
];
let n = sections.len();
let mut offsets = Vec::with_capacity(n);
let mut cur: u64 = HEADER_SIZE as u64;
for (_, bytes) in sections {
offsets.push(cur);
cur = align8(cur + bytes.len() as u64);
}
for (i, ((_, bytes), &offset)) in sections.iter().zip(offsets.iter()).enumerate() {
if offset > u32::MAX as u64 {
return Err(GraphError::Corrupt {
detail: format!("v8: section {i} offset {offset} exceeds u32"),
});
}
if bytes.len() > u32::MAX as usize {
return Err(GraphError::Corrupt {
detail: format!("v8: section {i} length {} exceeds u32", bytes.len()),
});
}
}
let mut header = vec![0u8; HEADER_SIZE];
header[0..4].copy_from_slice(b"GDB1");
header[4..6].copy_from_slice(&crate::snapshot::VERSION_9.to_le_bytes());
header[6..8].copy_from_slice(&(n as u16).to_le_bytes());
let mut pos = 8usize;
for ((section_id, bytes), &offset) in sections.iter().zip(offsets.iter()) {
let crc32 = crc32fast::hash(bytes);
header[pos] = *section_id;
header[pos + 1] = 0;
header[pos + 2] = 0;
header[pos + 3] = 0;
header[pos + 4..pos + 8].copy_from_slice(&(offset as u32).to_le_bytes());
header[pos + 8..pos + 12].copy_from_slice(&(bytes.len() as u32).to_le_bytes());
header[pos + 12..pos + 16].copy_from_slice(&crc32.to_le_bytes());
pos += 16;
}
let header_crc = crc32fast::hash(&header[0..pos]);
header[pos..pos + 4].copy_from_slice(&header_crc.to_le_bytes());
out.write_all(&header).map_err(GraphError::Io)?;
let zero_pad = [0u8; 8];
let last_idx = sections.len().saturating_sub(1);
for (i, ((_, bytes), &offset)) in sections.iter().zip(offsets.iter()).enumerate() {
out.write_all(bytes).map_err(GraphError::Io)?;
if i < last_idx {
let end = offset + bytes.len() as u64;
let next = align8(end);
let pad = (next - end) as usize;
if pad > 0 {
out.write_all(&zero_pad[..pad]).map_err(GraphError::Io)?;
}
}
}
Ok(())
}
fn align8(n: u64) -> u64 {
(n + 7) & !7
}
fn rkyv_encode<T>(value: &T) -> Result<Vec<u8>>
where
T: for<'a> rkyv::Serialize<
rkyv::api::high::HighSerializer<
rkyv::util::AlignedVec,
rkyv::ser::allocator::ArenaHandle<'a>,
rkyv::rancor::Error,
>,
>,
{
rkyv::api::high::to_bytes::<rkyv::rancor::Error>(value)
.map(|av| av.to_vec())
.map_err(|e| GraphError::Corrupt {
detail: format!("v8: rkyv encode: {e}"),
})
}
fn topology_to_csr(topo: &Topology) -> CsrData {
let mut etype_ids: Vec<u32> = topo.by_type.keys().copied().collect();
etype_ids.sort_unstable();
let etypes = etype_ids
.into_iter()
.map(|et| {
let adj = &topo.by_type[&et];
CsrEtype {
etype: et,
out_adj: adj_map_to_csr(&adj.out),
in_adj: adj_map_to_csr(&adj.inn),
}
})
.collect();
CsrData {
etypes,
edge_count: topo.edge_count(),
}
}
fn adj_map_to_csr(adj: &std::collections::HashMap<u32, crate::topology::AdjList>) -> CsrAdjMap {
let mut vertices: Vec<u32> = adj.keys().copied().collect();
vertices.sort_unstable();
let rows = vertices
.into_iter()
.filter_map(|v| {
let al = &adj[&v];
let neighbors = al.merged().into_owned();
if neighbors.is_empty() {
None
} else {
Some(CsrRow {
vertex: v,
neighbors,
})
}
})
.collect();
CsrAdjMap { rows }
}
fn topology_merge_to_csr(base: &crate::v8::layout::ArchivedCsr, overlay: &Topology) -> CsrData {
use std::collections::BTreeSet;
let mut etype_set: BTreeSet<u32> = overlay.by_type.keys().copied().collect();
for et in base.etypes.iter() {
etype_set.insert(u32::from(et.etype));
}
let etypes: Vec<CsrEtype> = etype_set
.into_iter()
.map(|et| {
let overlay_adj = overlay.by_type.get(&et);
let base_entry = base
.etypes
.binary_search_by_key(&et, |e| u32::from(e.etype))
.ok()
.map(|i| &base.etypes[i]);
let out_tombstones = overlay.out_tombstones.get(&et);
let in_tombstones = overlay.in_tombstones.get(&et);
let out_adj = merge_adj_map(
overlay_adj.map(|a| &a.out),
base_entry.map(|e| &e.out_adj),
out_tombstones,
);
let in_adj = merge_adj_map(
overlay_adj.map(|a| &a.inn),
base_entry.map(|e| &e.in_adj),
in_tombstones,
);
CsrEtype {
etype: et,
out_adj,
in_adj,
}
})
.collect();
let base_count = u64::from(base.edge_count);
let overlay_count = overlay.edge_count();
let tombstone_count: u64 = overlay
.out_tombstones
.values()
.flat_map(|m| m.values())
.map(|s| s.len() as u64)
.sum();
let edge_count = base_count.saturating_sub(tombstone_count) + overlay_count;
CsrData { etypes, edge_count }
}
fn merge_adj_map(
overlay_adj: Option<&std::collections::HashMap<u32, crate::topology::AdjList>>,
base_adj: Option<&crate::v8::layout::ArchivedCsrAdjMap>,
tombstones: Option<&std::collections::HashMap<u32, std::collections::BTreeSet<u32>>>,
) -> CsrAdjMap {
use std::collections::BTreeSet;
let mut vertex_set: BTreeSet<u32> = BTreeSet::new();
if let Some(o) = overlay_adj {
vertex_set.extend(o.keys().copied());
}
if let Some(b) = base_adj {
for row in b.rows.iter() {
vertex_set.insert(u32::from(row.vertex));
}
}
let rows: Vec<CsrRow> = vertex_set
.into_iter()
.filter_map(|v| {
let overlay_nbrs: Vec<u32> = overlay_adj
.and_then(|o| o.get(&v))
.map(|al| al.merged().into_owned())
.unwrap_or_default();
let base_nbrs: Vec<u32> = base_adj
.and_then(|b| {
b.rows
.binary_search_by_key(&v, |r| u32::from(r.vertex))
.ok()
.map(|i| &b.rows[i])
})
.map(|row| row.neighbors.iter().map(|n| u32::from(*n)).collect())
.unwrap_or_default();
let filtered_base_nbrs: Vec<u32> = match tombstones.and_then(|t| t.get(&v)) {
None => base_nbrs,
Some(t) => base_nbrs.into_iter().filter(|n| !t.contains(n)).collect(),
};
let merged = merge_sorted_unique_vecs(&overlay_nbrs, &filtered_base_nbrs);
if merged.is_empty() {
None
} else {
Some(CsrRow {
vertex: v,
neighbors: merged,
})
}
})
.collect();
CsrAdjMap { rows }
}
fn merge_sorted_unique_vecs(a: &[u32], b: &[u32]) -> Vec<u32> {
let mut out = Vec::with_capacity(a.len() + b.len());
let mut ai = 0;
let mut bi = 0;
while ai < a.len() && bi < b.len() {
match a[ai].cmp(&b[bi]) {
std::cmp::Ordering::Less => {
out.push(a[ai]);
ai += 1;
}
std::cmp::Ordering::Greater => {
out.push(b[bi]);
bi += 1;
}
std::cmp::Ordering::Equal => {
out.push(a[ai]);
ai += 1;
bi += 1;
}
}
}
out.extend_from_slice(&a[ai..]);
out.extend_from_slice(&b[bi..]);
out
}
fn columns_merge_to_data(
base: &crate::v8::layout::ArchivedColumns,
base_strings: Option<&crate::v8::layout::ArchivedStringTable>,
overlay: &ColumnStore,
) -> Result<(ColumnsData, StringTableData)> {
let mut merged = archived_to_columnstore(base, base_strings);
for (node, fields) in &overlay.prop_tombstones {
for field in fields {
merged.remove(*node, field);
}
}
for (field_name, node_map) in overlay.to_wire() {
for (node, value) in node_map {
merged.set(node, &field_name, value);
}
}
columnstore_to_data(&merged)
}
fn columnstore_to_data(store: &ColumnStore) -> Result<(ColumnsData, StringTableData)> {
let mut buf = Vec::new();
store.pack(&mut buf);
let (mut data, strings) = decode_all_columns(&buf)?;
for field in data.fields.iter_mut() {
if let ColumnData::Mixed(ref blob) = field.col {
let map: HashMap<u32, Value> =
bincode::deserialize(blob.as_slice()).unwrap_or_default();
if let Some(vec_col) = try_promote_to_vector(&map) {
field.col = vec_col;
}
}
}
Ok((data, strings))
}
fn try_promote_to_vector(map: &HashMap<u32, Value>) -> Option<ColumnData> {
if map.is_empty() {
return None;
}
let dim = map.values().next().and_then(|v| match v {
Value::List(items) => {
if items.iter().all(|i| matches!(i, Value::Float(_))) {
Some(items.len())
} else {
None
}
}
_ => None,
})?;
if dim == 0 {
return None;
}
let max_id = *map.keys().max().unwrap_or(&0) as usize;
for v in map.values() {
match v {
Value::List(items)
if items.len() == dim && items.iter().all(|i| matches!(i, Value::Float(_))) => {}
_ => return None,
}
}
let n_nodes = max_id + 1;
let mut data = vec![0.0f64; n_nodes * dim];
let mut present_words = vec![0u64; n_nodes.div_ceil(64)];
for (&id, v) in map {
let items = match v {
Value::List(items) => items,
_ => unreachable!(),
};
let start = id as usize * dim;
for (i, item) in items.iter().enumerate() {
if let Value::Float(f) = item {
data[start + i] = *f;
}
}
let word = id as usize / 64;
let bit = id as usize % 64;
present_words[word] |= 1u64 << bit;
}
Some(ColumnData::Vector {
dim: dim as u32,
data,
present: present_words,
})
}
fn decode_all_columns(buf: &[u8]) -> Result<(ColumnsData, StringTableData)> {
use crate::pack::{read_exact, read_f64s, read_i64s, read_str, read_u32, read_u32s};
let mut pos = 0usize;
let n_intern = read_u32(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: "v8: columns pack: truncated intern-string count".into(),
})? as usize;
let mut intern_strings: Vec<String> = Vec::with_capacity(n_intern);
for i in 0..n_intern {
intern_strings.push(read_str(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: columns pack: truncated intern-string[{i}]"),
})?);
}
let n_fields = read_u32(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: "v8: columns pack: truncated field count".into(),
})? as usize;
let mut fields = Vec::with_capacity(n_fields);
for field_idx in 0..n_fields {
let fname = read_str(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: columns pack: truncated field name at index {field_idx}"),
})?;
let tag = read_exact(buf, &mut pos, 1).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated tag byte"),
})?[0];
let col = match tag {
0 => {
let data = read_i64s(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Int data"),
})?;
let present = unpack_bitmap(buf, &mut pos).ok_or_else(|| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Int bitmap"),
})?;
ColumnData::Int { data, present }
}
1 => {
let data = read_f64s(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Float data"),
})?;
let present = unpack_bitmap(buf, &mut pos).ok_or_else(|| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Float bitmap"),
})?;
ColumnData::Float { data, present }
}
2 => {
let n = read_u32(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Bool length"),
})? as usize;
let raw = read_exact(buf, &mut pos, n)
.map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Bool data"),
})?
.to_vec();
let present = unpack_bitmap(buf, &mut pos).ok_or_else(|| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Bool bitmap"),
})?;
ColumnData::Bool { data: raw, present }
}
3 => {
let ids = read_u32s(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Str ids"),
})?;
let present = unpack_bitmap(buf, &mut pos).ok_or_else(|| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Str bitmap"),
})?;
ColumnData::Str {
ids,
present,
strings: Vec::new(),
}
}
4 => {
let blob_len = read_u32(buf, &mut pos).map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Mixed length"),
})? as usize;
let blob = read_exact(buf, &mut pos, blob_len)
.map_err(|_| GraphError::Corrupt {
detail: format!("v8: column '{fname}' pack decode: truncated Mixed data"),
})?
.to_vec();
ColumnData::Mixed(blob)
}
other => {
return Err(GraphError::Corrupt {
detail: format!(
"v8: column '{fname}' pack decode: unknown tag byte {other} at field \
index {field_idx}"
),
});
}
};
fields.push(FieldEntry { name: fname, col });
}
fields.sort_by(|a, b| a.name.cmp(&b.name));
Ok((
ColumnsData { fields },
StringTableData {
strings: intern_strings,
},
))
}
fn unpack_bitmap(buf: &[u8], pos: &mut usize) -> Option<Vec<u64>> {
use crate::pack::{read_exact, read_u32};
let n = read_u32(buf, pos).ok()? as usize;
let bytes = read_exact(buf, pos, n.saturating_mul(8)).ok()?;
Some(
bytes
.chunks_exact(8)
.map(|c| u64::from_le_bytes(c.try_into().unwrap())) .collect(),
)
}
fn idmap_to_data(ids: &IdMap) -> IdMapData {
let to_key: Vec<String> = ids.all_keys().to_vec();
let tombstones: Vec<u32> = ids
.all_keys()
.iter()
.enumerate()
.filter_map(|(i, _)| {
let id = i as u32;
if ids.is_tombstoned(id) {
Some(id)
} else {
None
}
})
.collect();
IdMapData { to_key, tombstones }
}
fn interner_to_data(syms: &Interner) -> InternerData {
let n = syms.len();
let mut to_str = Vec::with_capacity(n);
for i in 0u32..n as u32 {
to_str.push(syms.resolve(i).unwrap_or("").to_string());
}
InternerData { to_str }
}
pub fn csr_to_topology(archived: &crate::v8::layout::ArchivedCsr) -> Topology {
let mut topo = Topology::new();
for et_entry in archived.etypes.iter() {
let et = u32::from(et_entry.etype);
for row in et_entry.out_adj.rows.iter() {
let src = u32::from(row.vertex);
for &nbr in row.neighbors.iter() {
let dst = u32::from(nbr);
topo.add_edge(et, src, dst);
}
}
}
topo
}
pub fn archived_to_columnstore(
archived: &crate::v8::layout::ArchivedColumns,
shared: Option<&crate::v8::layout::ArchivedStringTable>,
) -> ColumnStore {
let mut store = ColumnStore::new();
for field in archived.fields.iter() {
let name: &str = field.name.as_str();
match &field.col {
crate::v8::layout::ArchivedColumnData::Int { data, present } => {
bitmap_for_each(present.as_slice(), |node| {
let idx = node as usize;
if idx < data.len() {
store.set(node, name, Value::Int(i64::from(data[idx])));
}
});
}
crate::v8::layout::ArchivedColumnData::Float { data, present } => {
bitmap_for_each(present.as_slice(), |node| {
let idx = node as usize;
if idx < data.len() {
store.set(node, name, Value::Float(f64::from(data[idx])));
}
});
}
crate::v8::layout::ArchivedColumnData::Bool { data, present } => {
bitmap_for_each(present.as_slice(), |node| {
let idx = node as usize;
if idx < data.len() {
store.set(node, name, Value::Bool(data[idx] != 0));
}
});
}
crate::v8::layout::ArchivedColumnData::Str {
ids,
present,
strings,
} => {
let table = match shared {
Some(t) => &t.strings,
None => strings,
};
let strings_vec: Vec<String> =
table.iter().map(|s| s.as_str().to_string()).collect();
bitmap_for_each(present.as_slice(), |node| {
let idx = node as usize;
if idx < ids.len() {
let sid = u32::from(ids[idx]) as usize;
if sid < strings_vec.len() {
store.set(node, name, Value::Str(strings_vec[sid].clone()));
}
}
});
}
crate::v8::layout::ArchivedColumnData::Mixed(blob) => {
let map: HashMap<u32, Value> =
bincode::deserialize(blob.as_slice()).unwrap_or_default();
for (node, v) in map {
store.set(node, name, v);
}
}
crate::v8::layout::ArchivedColumnData::Vector { dim, data, present } => {
let dim_val = u32::from(*dim) as usize;
bitmap_for_each(present.as_slice(), |node| {
let start = node as usize * dim_val;
let end = start + dim_val;
if end <= data.len() {
let floats: Vec<Value> = data[start..end]
.iter()
.map(|f| Value::Float(f64::from(*f)))
.collect();
store.set(node, name, Value::List(floats));
}
});
}
}
}
store
}
fn bitmap_for_each<F: FnMut(u32)>(words: &[rkyv::Archived<u64>], mut f: F) {
for (wi, word) in words.iter().enumerate() {
let mut w: u64 = u64::from(*word);
while w != 0 {
let bit = w.trailing_zeros();
f(wi as u32 * 64 + bit);
w &= w - 1;
}
}
}
pub fn archived_to_idmap(archived: &crate::v8::layout::ArchivedIdMap) -> IdMap {
let mut ids = IdMap::new();
let tombstoned: BTreeSet<u32> = archived.tombstones.iter().map(|t| u32::from(*t)).collect();
for (i, key) in archived.to_key.iter().enumerate() {
let s = key.as_str();
if tombstoned.contains(&(i as u32)) {
ids.get_or_insert(s);
ids.delete(s);
} else if !s.is_empty() {
ids.get_or_insert(s);
}
}
ids
}
pub fn archived_to_interner(archived: &crate::v8::layout::ArchivedInterner) -> Interner {
let mut syms = Interner::new();
for s in archived.to_str.iter() {
syms.intern(s.as_str());
}
syms
}
pub fn decode_meta(bytes: &[u8]) -> Result<V8Meta> {
bincode::deserialize(bytes).map_err(|e| GraphError::Corrupt {
detail: format!("v8: meta bincode deserialize: {e}"),
})
}
pub fn decode_last_change_bytes(bytes: &[u8]) -> HashMap<u32, u64> {
if bytes.is_empty() {
return HashMap::new();
}
bincode::deserialize(bytes).unwrap_or_default()
}
pub fn decode_ivf_bytes(bytes: &[u8]) -> BTreeMap<String, PerRuleIvfState> {
if bytes.is_empty() {
return BTreeMap::new();
}
bincode::deserialize(bytes).unwrap_or_default()
}
fn edge_props_to_data(ep: &EdgeProps) -> EdgePropsData {
let mut entries: Vec<EdgePropEntry> = ep
.sorted_entries()
.into_iter()
.map(|(etype, src, dst, props)| {
let props_blob = bincode::serialize(&props).unwrap_or_default();
EdgePropEntry {
etype,
src,
dst,
props_blob,
}
})
.collect();
entries.sort_by_key(|e| (e.etype, e.src, e.dst));
EdgePropsData { entries }
}
fn edge_props_merge_to_data(
base: &crate::v8::layout::ArchivedEdgeProps,
overlay: &EdgeProps,
) -> EdgePropsData {
use std::collections::BTreeSet as BSet;
let overlay_keys: BSet<(u32, u32, u32)> = overlay
.sorted_entries()
.iter()
.map(|&(et, s, d, _)| (et, s, d))
.collect();
let tombstoned_keys: BSet<(u32, u32, u32)> = overlay.tombstoned_keys().collect();
let mut entries: Vec<EdgePropEntry> = Vec::new();
for entry in base.entries.iter() {
let et = u32::from(entry.etype);
let s = u32::from(entry.src);
let d = u32::from(entry.dst);
let key = (et, s, d);
if tombstoned_keys.contains(&key) || overlay_keys.contains(&key) {
continue;
}
entries.push(EdgePropEntry {
etype: et,
src: s,
dst: d,
props_blob: entry.props_blob.as_slice().to_vec(),
});
}
for (et, s, d, props) in overlay.sorted_entries() {
let props_blob = bincode::serialize(props).unwrap_or_default();
entries.push(EdgePropEntry {
etype: et,
src: s,
dst: d,
props_blob,
});
}
entries.sort_by_key(|e| (e.etype, e.src, e.dst));
EdgePropsData { entries }
}
fn hnsw_to_data(hnsw: &BTreeMap<String, (Vec<u8>, Vec<u8>)>) -> HnswSectionData {
let mut rules: Vec<HnswRuleEntry> = hnsw
.iter()
.map(|(name, (src, dst))| HnswRuleEntry {
name: name.clone(),
src_blob: src.clone(),
dst_blob: dst.clone(),
})
.collect();
rules.sort_by(|a, b| a.name.cmp(&b.name));
HnswSectionData { rules }
}
fn provenance_to_data(prov: &BTreeMap<String, BTreeSet<(u32, u32, u32)>>) -> ProvenanceSectionData {
let mut entries: Vec<ProvenanceEntry> = prov
.iter()
.map(|(rule, triples)| {
let mut sorted: Vec<Triple> = triples
.iter()
.map(|&(etype, src, dst)| Triple { etype, src, dst })
.collect();
sorted.sort_by_key(|t| (t.etype, t.src, t.dst));
ProvenanceEntry {
rule: rule.clone(),
triples: sorted,
}
})
.collect();
entries.sort_by(|a, b| a.rule.cmp(&b.rule));
ProvenanceSectionData { entries }
}
fn rules_meta_to_data(
rule_defs: &[Vec<u8>],
tripped: &BTreeMap<String, bool>,
fires: &BTreeMap<String, u64>,
) -> RulesMetaData {
let tripped_vec: Vec<RuleTripEntry> = tripped
.iter()
.map(|(rule, &t)| RuleTripEntry {
rule: rule.clone(),
tripped: t,
})
.collect();
let fires_vec: Vec<RuleFireEntry> = fires
.iter()
.map(|(rule, &f)| RuleFireEntry {
rule: rule.clone(),
fires: f,
})
.collect();
RulesMetaData {
rule_defs: rule_defs.to_vec(),
tripped: tripped_vec,
fires: fires_vec,
}
}
pub fn archived_edge_props_to_owned(archived: &crate::v8::layout::ArchivedEdgeProps) -> EdgeProps {
let mut ep = EdgeProps::new();
for entry in archived.entries.iter() {
let etype = u32::from(entry.etype);
let src = u32::from(entry.src);
let dst = u32::from(entry.dst);
let props: std::collections::BTreeMap<String, Value> =
bincode::deserialize(entry.props_blob.as_slice()).unwrap_or_default();
for (field, value) in props {
ep.set(etype, src, dst, &field, value);
}
}
ep
}
pub fn archived_provenance_to_owned(
archived: &crate::v8::layout::ArchivedProvenance,
) -> BTreeMap<String, BTreeSet<(u32, u32, u32)>> {
let mut map = BTreeMap::new();
for entry in archived.entries.iter() {
let rule = entry.rule.as_str().to_string();
let triples: BTreeSet<(u32, u32, u32)> = entry
.triples
.iter()
.map(|t| (u32::from(t.etype), u32::from(t.src), u32::from(t.dst)))
.collect();
map.insert(rule, triples);
}
map
}
pub fn decode_provenance_bytes(
bytes: &[u8],
) -> BTreeMap<String, std::collections::BTreeSet<(u32, u32, u32)>> {
if bytes.is_empty() {
return BTreeMap::new();
}
match rkyv::access::<crate::v8::layout::ArchivedProvenanceSectionData, rkyv::rancor::Error>(
bytes,
) {
Ok(archived) => archived_provenance_to_owned(archived),
Err(_) => BTreeMap::new(),
}
}
pub type RulesMetaOwned = (Vec<Vec<u8>>, BTreeMap<String, bool>, BTreeMap<String, u64>);
pub fn archived_rules_meta_to_owned(
archived: &crate::v8::layout::ArchivedRulesMeta,
) -> RulesMetaOwned {
let rule_defs: Vec<Vec<u8>> = archived
.rule_defs
.iter()
.map(|b| b.as_slice().to_vec())
.collect();
let tripped: BTreeMap<String, bool> = archived
.tripped
.iter()
.map(|e| (e.rule.as_str().to_string(), e.tripped))
.collect();
let fires: BTreeMap<String, u64> = archived
.fires
.iter()
.map(|e| (e.rule.as_str().to_string(), u64::from(e.fires)))
.collect();
(rule_defs, tripped, fires)
}
pub fn archived_views_to_owned(archived: &crate::v8::layout::ArchivedViews) -> Vec<Vec<u8>> {
archived
.view_defs
.iter()
.map(|b| b.as_slice().to_vec())
.collect()
}
pub fn archived_provenance_contains(
archived: &crate::v8::layout::ArchivedProvenance,
rule: &str,
etype: u32,
src: u32,
dst: u32,
) -> bool {
let entry = archived.entries.iter().find(|e| e.rule.as_str() == rule);
let Some(entry) = entry else {
return false;
};
entry
.triples
.binary_search_by(|t| {
let te = u32::from(t.etype);
let ts = u32::from(t.src);
let td = u32::from(t.dst);
(te, ts, td).cmp(&(etype, src, dst))
})
.is_ok()
}
pub fn archived_hnsw_to_owned(
archived: &crate::v8::layout::ArchivedHnsw,
) -> BTreeMap<String, (Vec<u8>, Vec<u8>)> {
archived
.rules
.iter()
.map(|e| {
(
e.name.as_str().to_string(),
(
e.src_blob.as_slice().to_vec(),
e.dst_blob.as_slice().to_vec(),
),
)
})
.collect()
}