use std::collections::BTreeMap;
use marsdb_storage::{ReadableMultimapTable, ReadableTable, Txn};
use serde::{Deserialize, Serialize};
use crate::error::GraphError;
use crate::labels::lookup_label_id;
use crate::model::{NodeId, PropertyValue};
use crate::props::lookup_prop_id;
use crate::write_ctx::WriteCtx;
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub struct IndexDef {
pub unique: bool,
}
fn index_prefix(label_id: u32, prop_id: u32) -> [u8; 8] {
let mut out = [0u8; 8];
out[0..4].copy_from_slice(&label_id.to_be_bytes());
out[4..8].copy_from_slice(&prop_id.to_be_bytes());
out
}
pub(crate) fn encode_index_value(v: &PropertyValue) -> Vec<u8> {
match v {
PropertyValue::Null => vec![0x00],
PropertyValue::Bool(b) => vec![0x01, u8::from(*b)],
PropertyValue::Int(i) => {
let mut out = vec![0x02];
out.extend_from_slice(&((*i as u64) ^ 0x8000_0000_0000_0000).to_be_bytes());
out
}
PropertyValue::Float(f) => {
let bits = f.to_bits();
let sortable = if bits & 0x8000_0000_0000_0000 != 0 {
!bits
} else {
bits | 0x8000_0000_0000_0000
};
let mut out = vec![0x03];
out.extend_from_slice(&sortable.to_be_bytes());
out
}
PropertyValue::String(s) => {
let mut out = vec![0x04];
out.extend_from_slice(s.as_bytes());
out
}
PropertyValue::Date(days) => {
let mut out = vec![0x05];
out.extend_from_slice(&((*days as u64) ^ 0x8000_0000_0000_0000).to_be_bytes());
out
}
PropertyValue::Duration {
months,
days,
seconds,
nanos,
} => {
let mut out = vec![0x06];
out.extend_from_slice(&months.to_be_bytes());
out.extend_from_slice(&days.to_be_bytes());
out.extend_from_slice(&seconds.to_be_bytes());
out.extend_from_slice(&nanos.to_be_bytes());
out
}
PropertyValue::LocalTime(nanos_of_day) => {
let mut out = vec![0x07];
out.extend_from_slice(&nanos_of_day.to_be_bytes());
out
}
PropertyValue::Time {
nanos_of_day,
offset_seconds,
} => {
let instant = nanos_of_day - *offset_seconds as i64 * 1_000_000_000;
let mut out = vec![0x08];
out.extend_from_slice(&((instant as u64) ^ 0x8000_0000_0000_0000).to_be_bytes());
out
}
PropertyValue::LocalDateTime {
epoch_seconds,
nanos,
} => {
let mut out = vec![0x09];
out.extend_from_slice(&((*epoch_seconds as u64) ^ 0x8000_0000_0000_0000).to_be_bytes());
out.extend_from_slice(&nanos.to_be_bytes());
out
}
PropertyValue::DateTime {
epoch_seconds,
nanos,
..
} => {
let mut out = vec![0x0A];
out.extend_from_slice(&((*epoch_seconds as u64) ^ 0x8000_0000_0000_0000).to_be_bytes());
out.extend_from_slice(&nanos.to_be_bytes());
out
}
PropertyValue::List(items) => {
let mut out = vec![0x0B];
for item in items {
let encoded = encode_index_value(item);
out.extend_from_slice(&(encoded.len() as u32).to_be_bytes());
out.extend_from_slice(&encoded);
}
out
}
PropertyValue::Map(_) => {
unreachable!("PropertyValue::Map is never a real stored/indexed property value")
}
}
}
fn index_key(label_id: u32, prop_id: u32, value: &PropertyValue) -> Vec<u8> {
let mut out = index_prefix(label_id, prop_id).to_vec();
out.extend_from_slice(&encode_index_value(value));
out
}
pub(crate) fn create_index(
ctx: &mut WriteCtx,
label: &str,
prop: &str,
unique: bool,
) -> Result<(), GraphError> {
let label_id = crate::labels::intern_label(ctx, label)?;
let prop_id = crate::props::intern_prop(ctx, prop)?;
let prefix = index_prefix(label_id, prop_id);
if ctx.index_defs()?.get(prefix.as_slice())?.is_some() {
return Err(GraphError::CorruptData(format!(
"index on label {label:?} property {prop:?} already exists"
)));
}
let node_ids: Vec<u64> = ctx
.node_label_index()?
.get(label_id)?
.map(|entry| entry.map(|value| value.value()).map_err(GraphError::from))
.collect::<Result<Vec<_>, GraphError>>()?;
let mut entries: Vec<(Vec<u8>, u64)> = Vec::with_capacity(node_ids.len());
for node_id in &node_ids {
let Some(guard) = ctx.nodes()?.get(*node_id)? else {
continue;
};
if let Some(raw) = crate::encode::node_prop_raw(guard.value(), prop_id)? {
let value = crate::encode::decode_value(raw)?;
entries.push((index_key(label_id, prop_id, &value), *node_id));
}
}
if unique {
let mut seen = std::collections::HashSet::with_capacity(entries.len());
for (key, _) in &entries {
if !seen.insert(key.clone()) {
return Err(GraphError::UniqueConstraintViolation {
label: label.to_string(),
property: prop.to_string(),
});
}
}
}
let encoded = postcard::to_allocvec(&IndexDef { unique })?;
ctx.index_defs()?
.insert(prefix.as_slice(), encoded.as_slice())?;
for (key, node_id) in entries {
ctx.property_index()?.insert(key.as_slice(), node_id)?;
}
Ok(())
}
pub fn lookup_index_def(txn: Txn, label: &str, prop: &str) -> Result<Option<IndexDef>, GraphError> {
let Some(label_id) = lookup_label_id(txn, label)? else {
return Ok(None);
};
let Some(prop_id) = lookup_prop_id(txn, prop)? else {
return Ok(None);
};
let prefix = index_prefix(label_id, prop_id);
let defs = txn.open_table(marsdb_storage::tables::INDEX_DEFS)?;
let found = defs
.get(prefix.as_slice())?
.map(|guard| guard.value().to_vec());
drop(defs);
match found {
Some(bytes) => Ok(Some(postcard::from_bytes(&bytes)?)),
None => Ok(None),
}
}
pub fn lookup_range(
txn: Txn,
label: &str,
prop: &str,
lo: Option<(&PropertyValue, bool)>,
hi: Option<(&PropertyValue, bool)>,
limit: Option<usize>,
) -> Result<Vec<NodeId>, GraphError> {
let Some(mut cursor) = IndexRangeCursor::new(txn, label, prop, lo, hi)? else {
return Ok(Vec::new());
};
let mut out = Vec::new();
loop {
let want = match limit {
Some(l) => {
if out.len() >= l {
return Ok(out);
}
l - out.len()
}
None => usize::MAX,
};
let chunk = cursor.next_chunk(txn, want.min(4096))?;
if chunk.is_empty() {
return Ok(out);
}
out.extend(chunk);
}
}
type KeyRegion = (std::ops::Bound<Vec<u8>>, std::ops::Bound<Vec<u8>>);
pub struct IndexRangeCursor {
regions: Vec<KeyRegion>,
region_index: usize,
resume: Option<(Vec<u8>, u64)>,
}
impl IndexRangeCursor {
pub fn new(
txn: Txn,
label: &str,
prop: &str,
lo: Option<(&PropertyValue, bool)>,
hi: Option<(&PropertyValue, bool)>,
) -> Result<Option<Self>, GraphError> {
let Some(label_id) = lookup_label_id(txn, label)? else {
return Ok(None);
};
let Some(prop_id) = lookup_prop_id(txn, prop)? else {
return Ok(None);
};
let prefix = index_prefix(label_id, prop_id);
Ok(Some(Self {
regions: range_regions(&prefix, lo, hi),
region_index: 0,
resume: None,
}))
}
pub fn next_chunk(&mut self, txn: Txn, chunk_size: usize) -> Result<Vec<NodeId>, GraphError> {
use std::ops::Bound;
if chunk_size == 0 {
return Ok(Vec::new());
}
let index = txn.open_multimap_table(marsdb_storage::tables::PROPERTY_INDEX)?;
let mut out = Vec::new();
while self.region_index < self.regions.len() {
let (region_start, region_end) = &self.regions[self.region_index];
let start_owned;
let start_bound: Bound<&[u8]> = match &self.resume {
Some((key, _)) => {
start_owned = key.clone();
Bound::Included(start_owned.as_slice())
}
None => match region_start {
Bound::Included(k) => Bound::Included(k.as_slice()),
Bound::Excluded(k) => Bound::Excluded(k.as_slice()),
Bound::Unbounded => Bound::Unbounded,
},
};
let end_bound: Bound<&[u8]> = match region_end {
Bound::Included(k) => Bound::Included(k.as_slice()),
Bound::Excluded(k) => Bound::Excluded(k.as_slice()),
Bound::Unbounded => Bound::Unbounded,
};
for entry in index.range::<&[u8]>((start_bound, end_bound))? {
let (key, values) = entry?;
let key_bytes = key.value().to_vec();
let skip_through = match &self.resume {
Some((resume_key, resume_val)) if *resume_key == key_bytes => Some(*resume_val),
_ => None,
};
for value in values {
let node = value?.value();
if skip_through.is_some_and(|last| node <= last) {
continue;
}
out.push(NodeId(node));
self.resume = Some((key_bytes.clone(), node));
if out.len() >= chunk_size {
return Ok(out);
}
}
}
self.region_index += 1;
self.resume = None;
}
Ok(out)
}
}
fn range_regions(
prefix: &[u8],
lo: Option<(&PropertyValue, bool)>,
hi: Option<(&PropertyValue, bool)>,
) -> Vec<KeyRegion> {
use std::ops::Bound;
let key = |value: &PropertyValue| {
let mut k = prefix.to_vec();
k.extend_from_slice(&encode_index_value(value));
k
};
let tag_start = |tag: u8| {
let mut k = prefix.to_vec();
k.push(tag);
Bound::Included(k)
};
let tag_end = |tag: u8| {
let mut k = prefix.to_vec();
k.push(tag + 1);
Bound::Excluded(k)
};
let numeric = |v: &PropertyValue| matches!(v, PropertyValue::Int(_) | PropertyValue::Float(_));
let is_numeric = lo.map(|(v, _)| numeric(v)).unwrap_or(true)
&& hi.map(|(v, _)| numeric(v)).unwrap_or(true)
&& (lo.is_some() || hi.is_some())
&& (lo.is_some_and(|(v, _)| numeric(v)) || hi.is_some_and(|(v, _)| numeric(v)));
if is_numeric {
let int_bound = |side_lo: bool, bound: Option<(&PropertyValue, bool)>| match bound {
None => {
if side_lo {
tag_start(0x02)
} else {
tag_end(0x02)
}
}
Some((PropertyValue::Int(i), inclusive)) => {
let k = key(&PropertyValue::Int(*i));
if inclusive {
Bound::Included(k)
} else {
Bound::Excluded(k)
}
}
Some((PropertyValue::Float(f), _)) => {
let widened = if side_lo { f.floor() } else { f.ceil() };
let clamped = widened.clamp(i64::MIN as f64, i64::MAX as f64) as i64;
Bound::Included(key(&PropertyValue::Int(clamped)))
}
Some(_) => unreachable!("numeric region only built for numeric bounds"),
};
let float_bound = |side_lo: bool, bound: Option<(&PropertyValue, bool)>| match bound {
None => {
if side_lo {
tag_start(0x03)
} else {
tag_end(0x03)
}
}
Some((PropertyValue::Float(f), inclusive)) => {
let k = key(&PropertyValue::Float(*f));
if inclusive {
Bound::Included(k)
} else {
Bound::Excluded(k)
}
}
Some((PropertyValue::Int(i), _)) => {
let f = *i as f64;
let step = f.abs() * (2.0 * f64::EPSILON) + f64::MIN_POSITIVE;
let widened = if side_lo { f - step } else { f + step };
Bound::Included(key(&PropertyValue::Float(widened)))
}
Some(_) => unreachable!("numeric region only built for numeric bounds"),
};
return vec![
(int_bound(true, lo), int_bound(false, hi)),
(float_bound(true, lo), float_bound(false, hi)),
];
}
let tag = lo
.or(hi)
.map(|(v, _)| encode_index_value(v)[0])
.unwrap_or(0x00);
let start = match lo {
None => tag_start(tag),
Some((v, true)) => Bound::Included(key(v)),
Some((v, false)) => Bound::Excluded(key(v)),
};
let end = match hi {
None => tag_end(tag),
Some((v, true)) => Bound::Included(key(v)),
Some((v, false)) => Bound::Excluded(key(v)),
};
vec![(start, end)]
}
pub fn lookup_exact(
txn: Txn,
label: &str,
prop: &str,
value: &PropertyValue,
limit: Option<usize>,
) -> Result<Vec<NodeId>, GraphError> {
let Some(label_id) = lookup_label_id(txn, label)? else {
return Ok(Vec::new());
};
let Some(prop_id) = lookup_prop_id(txn, prop)? else {
return Ok(Vec::new());
};
let key = index_key(label_id, prop_id, value);
let index = txn.open_multimap_table(marsdb_storage::tables::PROPERTY_INDEX)?;
let iter = index.get(key.as_slice())?;
let ids: Vec<NodeId> = match limit {
Some(limit) => iter
.take(limit)
.map(|entry| {
entry
.map(|value| NodeId(value.value()))
.map_err(GraphError::from)
})
.collect::<Result<Vec<_>, GraphError>>()?,
None => iter
.map(|entry| {
entry
.map(|value| NodeId(value.value()))
.map_err(GraphError::from)
})
.collect::<Result<Vec<_>, GraphError>>()?,
};
drop(index);
Ok(ids)
}
pub fn match_count(
txn: Txn,
label: &str,
prop: &str,
value: &PropertyValue,
) -> Result<u64, GraphError> {
let Some(label_id) = lookup_label_id(txn, label)? else {
return Ok(0);
};
let Some(prop_id) = lookup_prop_id(txn, prop)? else {
return Ok(0);
};
let key = index_key(label_id, prop_id, value);
let index = txn.open_multimap_table(marsdb_storage::tables::PROPERTY_INDEX)?;
let count = index.get(key.as_slice())?.len();
Ok(count)
}
fn indexes_for_labels(
ctx: &mut WriteCtx,
label_ids: &[u32],
) -> Result<Vec<(u32, u32, String, IndexDef)>, GraphError> {
let raw: Vec<(u32, u32, IndexDef)> = {
let mut raw = Vec::new();
for entry in ctx.index_defs()?.iter()? {
let (key, value) = entry?;
let key_bytes = key.value();
let label_id = u32::from_be_bytes(
key_bytes[0..4]
.try_into()
.expect("index key prefix is 8 bytes"),
);
if !label_ids.contains(&label_id) {
continue;
}
let prop_id = u32::from_be_bytes(
key_bytes[4..8]
.try_into()
.expect("index key prefix is 8 bytes"),
);
let def: IndexDef = postcard::from_bytes(value.value())?;
raw.push((label_id, prop_id, def));
}
raw
};
raw.into_iter()
.map(|(label_id, prop_id, def)| {
let prop_name = resolve_prop_ctx(ctx, prop_id)?;
Ok((label_id, prop_id, prop_name, def))
})
.collect()
}
pub(crate) fn resolve_label_ctx(ctx: &mut WriteCtx, label_id: u32) -> Result<String, GraphError> {
let value = ctx.id_to_label()?.get(label_id)?.ok_or_else(|| {
GraphError::CorruptData(format!("label id {label_id} has no interned string"))
})?;
Ok(value.value().to_string())
}
pub(crate) fn resolve_prop_ctx(ctx: &mut WriteCtx, prop_id: u32) -> Result<String, GraphError> {
let value = ctx.id_to_prop()?.get(prop_id)?.ok_or_else(|| {
GraphError::CorruptData(format!("prop id {prop_id} has no interned string"))
})?;
Ok(value.value().to_string())
}
struct IndexTarget<'a> {
label_id: u32,
prop_id: u32,
label: &'a str,
prop: &'a str,
}
fn insert_entry(
ctx: &mut WriteCtx,
target: &IndexTarget<'_>,
value: &PropertyValue,
node_id: u64,
unique: bool,
) -> Result<(), GraphError> {
let key = index_key(target.label_id, target.prop_id, value);
if unique && ctx.property_index()?.get(key.as_slice())?.next().is_some() {
return Err(GraphError::UniqueConstraintViolation {
label: target.label.to_string(),
property: target.prop.to_string(),
});
}
ctx.property_index()?.insert(key.as_slice(), node_id)?;
Ok(())
}
fn remove_entry(
ctx: &mut WriteCtx,
label_id: u32,
prop_id: u32,
value: &PropertyValue,
node_id: u64,
) -> Result<(), GraphError> {
let key = index_key(label_id, prop_id, value);
ctx.property_index()?.remove(key.as_slice(), node_id)?;
Ok(())
}
pub(crate) fn on_node_created(
ctx: &mut WriteCtx,
node_id: u64,
label_ids: &[u32],
props: &BTreeMap<String, PropertyValue>,
) -> Result<(), GraphError> {
for (label_id, prop_id, prop_name, def) in indexes_for_labels(ctx, label_ids)? {
if let Some(value) = props.get(&prop_name) {
let label = resolve_label_ctx(ctx, label_id)?;
let target = IndexTarget {
label_id,
prop_id,
label: &label,
prop: &prop_name,
};
insert_entry(ctx, &target, value, node_id, def.unique)?;
}
}
Ok(())
}
pub(crate) fn on_node_deleted(
ctx: &mut WriteCtx,
node_id: u64,
label_ids: &[u32],
props: &BTreeMap<String, PropertyValue>,
) -> Result<(), GraphError> {
for (label_id, prop_id, prop_name, _def) in indexes_for_labels(ctx, label_ids)? {
if let Some(value) = props.get(&prop_name) {
remove_entry(ctx, label_id, prop_id, value, node_id)?;
}
}
Ok(())
}
pub(crate) fn on_node_prop_changed(
ctx: &mut WriteCtx,
node_id: u64,
label_ids: &[u32],
prop: &str,
old_value: Option<&PropertyValue>,
new_value: Option<&PropertyValue>,
) -> Result<(), GraphError> {
for (label_id, prop_id, prop_name, def) in indexes_for_labels(ctx, label_ids)? {
if prop_name != prop {
continue;
}
if let Some(old) = old_value {
remove_entry(ctx, label_id, prop_id, old, node_id)?;
}
if let Some(new) = new_value {
let label = resolve_label_ctx(ctx, label_id)?;
let target = IndexTarget {
label_id,
prop_id,
label: &label,
prop: &prop_name,
};
insert_entry(ctx, &target, new, node_id, def.unique)?;
}
}
Ok(())
}