use core_storage::v8::layout::ArchivedColumns;
use core_storage::v8::seam::{ColumnsView, TopologyView};
use core_storage::{ColumnStore, Direction, IdMap, Interner, Value};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum AggFn {
Sum,
Avg,
Min,
Max,
Count,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum ViewSource {
Degree {
edge_type: String,
direction: Direction,
},
NeighborAgg {
edge_type: String,
direction: Direction,
agg: AggFn,
prop: String,
},
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ViewDef {
pub name: String,
pub label: String,
pub view_prop: String,
pub source: ViewSource,
}
impl ViewDef {
pub fn validate(&self) -> Result<(), String> {
if self.name.is_empty() {
return Err("view name must not be empty".into());
}
if self.label.is_empty() {
return Err("view label must not be empty".into());
}
if self.view_prop.is_empty() {
return Err("view_prop must not be empty".into());
}
match &self.source {
ViewSource::Degree { edge_type, .. } => {
if edge_type.is_empty() {
return Err("Degree edge_type must not be empty".into());
}
}
ViewSource::NeighborAgg {
edge_type, prop, ..
} => {
if edge_type.is_empty() {
return Err("NeighborAgg edge_type must not be empty".into());
}
if prop.is_empty() {
return Err("NeighborAgg prop must not be empty".into());
}
}
}
Ok(())
}
fn edge_type(&self) -> &str {
match &self.source {
ViewSource::Degree { edge_type, .. } => edge_type,
ViewSource::NeighborAgg { edge_type, .. } => edge_type,
}
}
fn direction(&self) -> Direction {
match &self.source {
ViewSource::Degree { direction, .. } => *direction,
ViewSource::NeighborAgg { direction, .. } => *direction,
}
}
}
#[derive(Debug, Default, Clone)]
pub struct ViewStore {
views: BTreeMap<String, ViewDef>,
}
impl ViewStore {
pub fn new() -> Self {
Self::default()
}
pub fn views(&self) -> impl Iterator<Item = &ViewDef> {
self.views.values()
}
pub fn is_empty(&self) -> bool {
self.views.is_empty()
}
pub fn has_view(&self, name: &str) -> bool {
self.views.contains_key(name)
}
pub fn view_for_prop(&self, prop_name: &str) -> Option<&str> {
self.views
.values()
.find(|v| v.view_prop == prop_name)
.map(|v| v.name.as_str())
}
pub fn create_view(
&mut self,
def: ViewDef,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
) -> Result<(), String> {
def.validate()?;
if self.views.contains_key(&def.name) {
return Err(format!("view {:?} already exists", def.name));
}
if let Some(existing) = self.views.values().find(|v| v.view_prop == def.view_prop) {
return Err(format!(
"view_prop {:?} is already used by view {:?}",
def.view_prop, existing.name
));
}
let view_props: std::collections::HashSet<&str> =
self.views.values().map(|v| v.view_prop.as_str()).collect();
if props
.fields()
.any(|f| f == def.view_prop && !view_props.contains(f))
{
return Err(format!(
"view_prop {:?} conflicts with an existing node property",
def.view_prop
));
}
backfill_view(&def, props, topo, ids, syms, labels);
self.views.insert(def.name.clone(), def);
Ok(())
}
pub fn restore_view(&mut self, def: ViewDef) -> Result<(), String> {
def.validate()?;
if self.views.contains_key(&def.name) {
return Ok(()); }
self.views.insert(def.name.clone(), def);
Ok(())
}
pub fn delete_view(
&mut self,
name: &str,
props: &mut ColumnStore,
ids: &IdMap,
labels: &[u32],
syms: &Interner,
) -> Result<(), String> {
let def = self
.views
.remove(name)
.ok_or_else(|| format!("view {:?} not found", name))?;
if let Some(label_sym) = syms.get(&def.label) {
for id in 0..ids.len() as u32 {
if labels.get(id as usize).copied() == Some(label_sym) {
props.remove(id, &def.view_prop);
}
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub fn on_edge_changed(
&self,
etype: u32,
src: u32,
dst: u32,
inserted: bool,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
base_cols: Option<&ArchivedColumns>,
) {
for def in self.views.values() {
let Some(et_sym) = syms.get(def.edge_type()) else {
continue;
};
if et_sym != etype {
continue;
}
let direction = def.direction();
let subject = match direction {
Direction::Out => src,
Direction::In => dst,
};
let neighbor = match direction {
Direction::Out => dst,
Direction::In => src,
};
update_node_view(
def, subject, neighbor, inserted, props, topo, ids, syms, labels, base_cols,
);
}
}
#[allow(clippy::too_many_arguments)]
pub fn on_prop_changed(
&self,
changed_node: u32,
field: &str,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
base_cols: Option<&ArchivedColumns>,
) {
for def in self.views.values() {
let ViewSource::NeighborAgg {
edge_type,
direction,
prop,
..
} = &def.source
else {
continue;
};
if prop != field {
continue;
}
let Some(et_sym) = syms.get(edge_type) else {
continue;
};
let reverse_dir = match direction {
Direction::Out => Direction::In,
Direction::In => Direction::Out,
};
let subjects: Vec<u32> = topo.neighbors(et_sym, reverse_dir, changed_node).to_vec();
for subject in subjects {
let val = {
let pv = build_cols_view(props, base_cols);
compute_view_value(def, subject, pv, topo, ids, syms, labels)
};
match val {
Some(v) => props.set(subject, &def.view_prop, v),
None => {
props.remove(subject, &def.view_prop);
}
}
}
}
}
pub fn init_node_views(
&self,
node: u32,
props: &mut ColumnStore,
syms: &Interner,
labels: &[u32],
) {
for def in self.views.values() {
let Some(label_sym) = syms.get(&def.label) else {
continue;
};
if labels.get(node as usize).copied() != Some(label_sym) {
continue;
}
match &def.source {
ViewSource::Degree { .. } => {
if props.get(node, &def.view_prop).is_none() {
props.set(node, &def.view_prop, Value::Int(0));
}
}
ViewSource::NeighborAgg {
agg: AggFn::Count, ..
} => {
if props.get(node, &def.view_prop).is_none() {
props.set(node, &def.view_prop, Value::Int(0));
}
}
ViewSource::NeighborAgg {
agg: AggFn::Sum, ..
} => {
props.set(node, &def.view_prop, Value::Float(0.0));
}
_ => {} }
}
}
pub fn rebuild_all(
&self,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
) {
for def in self.views.values() {
backfill_view(def, props, topo, ids, syms, labels);
}
}
}
fn build_cols_view<'a>(
overlay: &'a ColumnStore,
base_cols: Option<&'a ArchivedColumns>,
) -> ColumnsView<'a> {
match base_cols {
None => ColumnsView::owned(overlay),
Some(b) => ColumnsView::with_base(overlay, b),
}
}
fn backfill_view(
def: &ViewDef,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
) {
let Some(label_sym) = syms.get(&def.label) else {
return;
};
let Some(et_sym) = syms.get(def.edge_type()) else {
for id in 0..ids.len() as u32 {
if labels.get(id as usize).copied() == Some(label_sym) {
match &def.source {
ViewSource::Degree { .. } => {
props.set(id, &def.view_prop, Value::Int(0));
}
ViewSource::NeighborAgg {
agg: AggFn::Count, ..
} => {
props.set(id, &def.view_prop, Value::Int(0));
}
ViewSource::NeighborAgg {
agg: AggFn::Sum, ..
} => {
props.set(id, &def.view_prop, Value::Float(0.0));
}
_ => {}
}
}
}
return;
};
for id in 0..ids.len() as u32 {
if labels.get(id as usize).copied() != Some(label_sym) {
continue;
}
match compute_view_value(
def,
id,
ColumnsView::owned(&*props),
topo,
ids,
syms,
labels,
) {
Some(val) => props.set(id, &def.view_prop, val),
None => {
props.remove(id, &def.view_prop);
}
}
let _ = et_sym;
}
}
pub fn compute_view_value(
def: &ViewDef,
node: u32,
props: ColumnsView<'_>,
topo: &TopologyView<'_>,
_ids: &IdMap,
syms: &Interner,
labels: &[u32],
) -> Option<Value> {
let label_sym = syms.get(&def.label)?;
if labels.get(node as usize).copied() != Some(label_sym) {
return None;
}
let et_sym = match syms.get(def.edge_type()) {
Some(s) => s,
None => {
return match &def.source {
ViewSource::Degree { .. } => Some(Value::Int(0)),
ViewSource::NeighborAgg {
agg: AggFn::Count, ..
} => Some(Value::Int(0)),
ViewSource::NeighborAgg {
agg: AggFn::Sum, ..
} => Some(Value::Float(0.0)),
ViewSource::NeighborAgg { .. } => None,
};
}
};
let direction = def.direction();
let neighbors = topo.neighbors(et_sym, direction, node);
match &def.source {
ViewSource::Degree { .. } => Some(Value::Int(neighbors.len() as i64)),
ViewSource::NeighborAgg { agg, prop, .. } => match agg {
AggFn::Count => Some(Value::Int(
neighbors
.iter()
.filter(|&&n| props.get(n, prop).is_some())
.count() as i64,
)),
AggFn::Sum => {
let mut sum = 0.0f64;
for &nbr in neighbors.as_ref() {
if let Some(vr) = props.get(nbr, prop) {
if let Some(n) = as_float(vr.as_value()) {
sum += n;
}
}
}
Some(Value::Float(sum))
}
AggFn::Avg => {
let mut sum = 0.0f64;
let mut count = 0usize;
for &nbr in neighbors.as_ref() {
if let Some(vr) = props.get(nbr, prop) {
if let Some(n) = as_float(vr.as_value()) {
sum += n;
count += 1;
}
}
}
if count == 0 {
None
} else {
Some(Value::Float(sum / count as f64))
}
}
AggFn::Min => {
let mut best: Option<f64> = None;
for &nbr in neighbors.as_ref() {
if let Some(vr) = props.get(nbr, prop) {
if let Some(n) = as_float(vr.as_value()) {
best = Some(best.map_or(n, |m: f64| m.min(n)));
}
}
}
best.map(Value::Float)
}
AggFn::Max => {
let mut best: Option<f64> = None;
for &nbr in neighbors.as_ref() {
if let Some(vr) = props.get(nbr, prop) {
if let Some(n) = as_float(vr.as_value()) {
best = Some(best.map_or(n, |m: f64| m.max(n)));
}
}
}
best.map(Value::Float)
}
},
}
}
#[allow(clippy::too_many_arguments)]
fn update_node_view(
def: &ViewDef,
subject: u32,
neighbor: u32,
inserted: bool,
props: &mut ColumnStore,
topo: &TopologyView<'_>,
ids: &IdMap,
syms: &Interner,
labels: &[u32],
base_cols: Option<&ArchivedColumns>,
) {
let Some(label_sym) = syms.get(&def.label) else {
return;
};
if labels.get(subject as usize).copied() != Some(label_sym) {
return;
}
match &def.source {
ViewSource::Degree { direction, .. } => {
let et_sym = match syms.get(def.edge_type()) {
Some(s) => s,
None => {
props.set(subject, &def.view_prop, Value::Int(0));
return;
}
};
let count = topo.neighbors(et_sym, *direction, subject).len() as i64;
props.set(subject, &def.view_prop, Value::Int(count));
}
ViewSource::NeighborAgg { .. } => {
let val = {
let pv = build_cols_view(props, base_cols);
compute_view_value(def, subject, pv, topo, ids, syms, labels)
};
match val {
Some(v) => props.set(subject, &def.view_prop, v),
None => {
props.remove(subject, &def.view_prop);
}
}
}
}
let _ = (neighbor, inserted); }
fn as_float(v: &Value) -> Option<f64> {
match v {
Value::Int(i) => Some(*i as f64),
Value::Float(f) if f.is_finite() => Some(*f),
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
use core_storage::{IdMap, Interner, Topology};
fn make_setup() -> (ViewStore, ColumnStore, Topology, IdMap, Interner, Vec<u32>) {
let mut ids = IdMap::new();
let mut syms = Interner::new();
let mut labels = Vec::new();
let mut topo = Topology::new();
let mut props = ColumnStore::new();
let person_sym = syms.intern("Person");
let city_sym = syms.intern("City");
let edge_sym = syms.intern("LIVES_IN");
for key in &["p0", "p1", "p2"] {
let id = ids.get_or_insert(key);
if labels.len() <= id as usize {
labels.resize(id as usize + 1, u32::MAX);
}
labels[id as usize] = person_sym;
}
let c0 = ids.get_or_insert("c0");
if labels.len() <= c0 as usize {
labels.resize(c0 as usize + 1, u32::MAX);
}
labels[c0 as usize] = city_sym;
let p0 = ids.get("p0").unwrap();
let p1 = ids.get("p1").unwrap();
topo.add_edge(edge_sym, p0, c0);
topo.add_edge(edge_sym, p1, c0);
props.set(p0, "score", Value::Float(3.0));
props.set(p1, "score", Value::Float(7.0));
(ViewStore::new(), props, topo, ids, syms, labels)
}
#[test]
fn degree_view_basic() {
let (mut vs, mut props, topo, ids, syms, labels) = make_setup();
let def = ViewDef {
name: "city_pop".into(),
label: "City".into(),
view_prop: "in_deg".into(),
source: ViewSource::Degree {
edge_type: "LIVES_IN".into(),
direction: Direction::In,
},
};
vs.create_view(
def,
&mut props,
&TopologyView::owned(&topo),
&ids,
&syms,
&labels,
)
.unwrap();
let c0 = ids.get("c0").unwrap();
assert_eq!(props.get(c0, "in_deg"), Some(&Value::Int(2)));
}
#[test]
fn neighbor_agg_sum() {
let (mut vs, mut props, topo, ids, syms, labels) = make_setup();
let def = ViewDef {
name: "city_score_sum".into(),
label: "City".into(),
view_prop: "score_sum".into(),
source: ViewSource::NeighborAgg {
edge_type: "LIVES_IN".into(),
direction: Direction::In,
agg: AggFn::Sum,
prop: "score".into(),
},
};
vs.create_view(
def,
&mut props,
&TopologyView::owned(&topo),
&ids,
&syms,
&labels,
)
.unwrap();
let c0 = ids.get("c0").unwrap();
assert_eq!(props.get(c0, "score_sum"), Some(&Value::Float(10.0)));
}
#[test]
fn neighbor_agg_count_skips_missing_prop() {
let (mut vs, mut props, topo, ids, syms, labels) = make_setup();
let p1 = ids.get("p1").unwrap();
props.remove(p1, "score");
let def = ViewDef {
name: "city_score_n".into(),
label: "City".into(),
view_prop: "score_n".into(),
source: ViewSource::NeighborAgg {
edge_type: "LIVES_IN".into(),
direction: Direction::In,
agg: AggFn::Count,
prop: "score".into(),
},
};
vs.create_view(
def,
&mut props,
&TopologyView::owned(&topo),
&ids,
&syms,
&labels,
)
.unwrap();
let c0 = ids.get("c0").unwrap();
assert_eq!(props.get(c0, "score_n"), Some(&Value::Int(1)));
}
#[test]
fn view_prop_collision_rejected() {
let (mut vs, mut props, topo, ids, syms, labels) = make_setup();
let p0 = ids.get("p0").unwrap();
props.set(p0, "collision_prop", Value::Int(1));
let def = ViewDef {
name: "test_view".into(),
label: "Person".into(),
view_prop: "collision_prop".into(),
source: ViewSource::Degree {
edge_type: "LIVES_IN".into(),
direction: Direction::Out,
},
};
let err = vs
.create_view(
def,
&mut props,
&TopologyView::owned(&topo),
&ids,
&syms,
&labels,
)
.unwrap_err();
assert!(
err.contains("conflicts with an existing node property"),
"{err}"
);
}
#[test]
fn delete_view_removes_values() {
let (mut vs, mut props, topo, ids, syms, labels) = make_setup();
let def = ViewDef {
name: "city_pop".into(),
label: "City".into(),
view_prop: "in_deg".into(),
source: ViewSource::Degree {
edge_type: "LIVES_IN".into(),
direction: Direction::In,
},
};
vs.create_view(
def,
&mut props,
&TopologyView::owned(&topo),
&ids,
&syms,
&labels,
)
.unwrap();
let c0 = ids.get("c0").unwrap();
assert!(props.get(c0, "in_deg").is_some());
vs.delete_view("city_pop", &mut props, &ids, &labels, &syms)
.unwrap();
assert!(props.get(c0, "in_deg").is_none());
}
}