Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions db4-graph/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -248,7 +248,7 @@ where
.iter_entries()
.filter(|entry| {
!entry.node_additions(STATIC_GRAPH_LAYER_ID).is_empty()
|| entry.has_layer_inner(*id)
|| entry.has_layer(*id)
})
.count()
})
Expand All @@ -264,7 +264,7 @@ where
.iter_entries()
.filter(|entry| {
!entry.node_additions(STATIC_GRAPH_LAYER_ID).is_empty()
|| ids.iter().any(|layer| entry.has_layer_inner(layer))
|| ids.iter().any(|layer| entry.has_layer(layer))
})
.count()
})
Expand Down
67 changes: 67 additions & 0 deletions db4-graph/src/replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -541,6 +541,73 @@ where
Ok(())
}

fn replay_delete_node(
&mut self,
lsn: LSN,
_transaction_id: TransactionID,
t: EventTime,
node_name: Option<GID>,
node_id: VID,
layer_name: Option<String>,
layer_id: LayerId,
) -> Result<(), StorageError> {
// Insert node id into resolver.
if let Some(ref name) = node_name {
self.graph()
.logical_to_physical
.set(name.as_ref(), node_id)?;
}

// Make layer name -> id mapping available to both edge and node meta.
if let Some(name) = layer_name.as_deref() {
self.graph()
.edge_meta()
.layer_meta()
.set_id(name, layer_id.0);

self.graph()
.node_meta()
.layer_meta()
.set_id(name, layer_id.0);
}

// Resolve segment and check LSN.
let (segment_id, pos) = self.graph().storage().nodes().resolve_pos(node_id);
self.resize_segments_to_vid(node_id);

let segment = self
.graph()
.storage()
.nodes()
.get_or_create_segment(segment_id);

let immut_lsn = segment.immut_lsn();

// Replay this entry only if it doesn't exist in immut.
if immut_lsn < lsn {
let node_segment = self.nodes.get_mut(segment_id).ok_or_else(|| {
StorageError::GenericFailure(format!(
"Node segment {segment_id} not found during replay_add_node"
))
})?;

let mut node_writer = node_segment.writer();

if !node_writer.has_node(pos, STATIC_GRAPH_LAYER_ID) {
node_writer.increment_seg_num_nodes();
}

if let Some(name) = node_name {
node_writer.store_node_id(pos, name);
}

node_writer.delete(t, pos, layer_id);
node_writer.set_lsn(lsn);
}

Ok(())
}

fn replay_add_node_metadata(
&mut self,
lsn: LSN,
Expand Down
115 changes: 80 additions & 35 deletions db4-storage/src/api/nodes.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,25 @@
use crate::{
LocalPOS,
error::StorageError,
generic_time_ops::LayerIter,
pages::node_store::increment_and_clamp,
persist::strategy::PersistenceStrategy,
segments::node::segment::MemNodeSegment,
utils::{Iter2, Iter3, Iter4},
wal::LSN,
};
use itertools::Itertools;
use parking_lot::{RwLockReadGuard, RwLockWriteGuard, lock_api::ArcRwLockReadGuard};
use raphtory_api::{
core::{
Direction,
entities::properties::{
meta::{Meta, NODE_ID_PROP_ID, NODE_TYPE_PROP_ID},
prop::{AsPropRef, Prop, PropUnwrap},
tprop::TPropOps,
entities::{
LayerId,
properties::{
meta::{Meta, NODE_ID_PROP_ID, NODE_TYPE_PROP_ID, STATIC_GRAPH_LAYER_ID},
prop::{AsPropRef, Prop, PropUnwrap, prop_hashable::HashableProp},
tprop::TPropOps,
},
},
},
iter::IntoDynBoxed,
Expand All @@ -17,6 +30,8 @@ use raphtory_core::{
storage::timeindex::{EventTime, TimeIndexOps},
utils::iter::GenLockedIter,
};
use raphtory_itertools::FastMergeExt;
use rayon::prelude::*;
use std::{
borrow::Cow,
collections::HashSet,
Expand All @@ -29,23 +44,6 @@ use std::{
},
};

use crate::{
LocalPOS,
error::StorageError,
generic_time_ops::LayerIter,
pages::node_store::increment_and_clamp,
persist::strategy::PersistenceStrategy,
segments::node::segment::MemNodeSegment,
utils::{Iter2, Iter3, Iter4},
wal::LSN,
};
use raphtory_api::core::entities::{
LayerId,
properties::{meta::STATIC_GRAPH_LAYER_ID, prop::prop_hashable::HashableProp},
};
use raphtory_itertools::FastMergeExt;
use rayon::prelude::*;

/// A property predicate a storage backend may resolve to candidate rows via a
/// secondary index. String operators use the semantics of the corresponding
/// `str` methods; comparisons use the property value's natural order.
Expand Down Expand Up @@ -185,7 +183,7 @@ pub type NodeTypeIndexOf<NS> = <<NS as NodeSegmentOps>::Extension as Persistence
pub trait NodeSegmentOps: Send + Sync + Debug + 'static {
type Extension;

type Entry<'a>: NodeEntryOps<'a>
type Entry<'a>: NodeEntryOps + 'a
where
Self: 'a;

Expand Down Expand Up @@ -314,16 +312,15 @@ pub trait LockedNSSegment: Debug + Send + Sync {
}
}

pub trait NodeEntryOps<'a>: Send + Sync + 'a {
pub trait NodeEntryOps: Send + Sync {
type Ref<'b>: NodeRefOps<'b>
where
'a: 'b,
Self: 'b;

fn as_ref<'b>(&'b self) -> Self::Ref<'b>
where
'a: 'b;
fn as_ref<'b>(&'b self) -> Self::Ref<'b>;
}

pub trait IntoEdges<'a>: NodeEntryOps + Send + Sync + 'a {
fn into_edges<'b: 'a>(
self,
layers: &'b LayerIds,
Expand All @@ -338,9 +335,13 @@ pub trait NodeEntryOps<'a>: Send + Sync + 'a {
}
}

impl<'a, T: NodeEntryOps + Send + Sync + 'a> IntoEdges<'a> for T {}

pub trait NodeRefOps<'a>: Copy + Clone + Send + Sync + 'a {
type Additions: TimeIndexOps<'a, IndexType = EventTime>;
type EdgeAdditions: TimeIndexOps<'a, IndexType = EventTime>;

type Deletions: TimeIndexOps<'a, IndexType = EventTime>;
type TProps: TPropOps<'a>;

fn out_edges(self, layer_id: LayerId) -> impl Iterator<Item = (VID, EID)> + Send + Sync + 'a;
Expand Down Expand Up @@ -446,14 +447,15 @@ pub trait NodeRefOps<'a>: Copy + Clone + Send + Sync + 'a {
}
}

fn node_meta(&self) -> &Arc<Meta>;
fn node_meta(self) -> &'a Arc<Meta>;

fn temp_prop_rows(
fn t_prop_rows<L: Into<LayerIter<'a>>>(
self,
w: Option<Range<EventTime>>,
prop_ids: Arc<[usize]>,
) -> impl Iterator<Item = (EventTime, usize, Vec<(usize, Prop)>)> + 'a {
(0..self.internal_num_layers()).flat_map(move |layer_id| {
layers: L,
) -> impl Iterator<Item = (EventTime, LayerId, Vec<(usize, Prop)>)> + 'a {
self.layer_ids_iter(layers).flat_map(move |layer_id| {
let w = w.clone();
let prop_ids = Arc::clone(&prop_ids);
let additions = self.node_additions(layer_id);
Expand All @@ -466,7 +468,7 @@ pub trait NodeRefOps<'a>: Copy + Clone + Send + Sync + 'a {
.iter()
.copied()
.map(move |prop_id| {
self.temporal_prop_layer(LayerId(layer_id), prop_id)
self.t_prop(layer_id, prop_id)
.iter_inner(w.clone())
.map(move |(t, prop)| (t, (prop_id, prop)))
})
Expand Down Expand Up @@ -538,11 +540,39 @@ pub trait NodeRefOps<'a>: Copy + Clone + Send + Sync + 'a {

fn node_additions<L: Into<LayerIter<'a>>>(self, layer_id: L) -> Self::Additions;

fn node_deletions<L: Into<LayerIter<'a>>>(self, layer_id: L) -> Self::Deletions;

fn node_updates_iter<L: Into<LayerIter<'a>>>(
self,
layer_ids: L,
) -> impl Iterator<Item = (LayerId, Self::Additions, Self::Deletions)> + 'a {
self.layer_ids_iter(layer_ids).map(move |layer_id| {
(
layer_id,
self.node_additions(layer_id),
self.node_deletions(layer_id),
)
})
}

fn c_prop(self, layer_id: LayerId, prop_id: usize) -> Option<Prop>;

fn c_prop_str(self, layer_id: LayerId, prop_id: usize) -> Option<&'a str>;

fn temporal_prop_layer(self, layer_id: LayerId, prop_id: usize) -> Self::TProps;
fn t_prop<L: Into<LayerIter<'a>>>(self, layer_ids: L, prop_id: usize) -> Self::TProps;

/// Iterate over `NodeTProps` for each layer specified by `layer_ids`, always
/// including `STATIC_GRAPH_LAYER_ID` (the layer for nodes added without an
/// explicit layer name). This mirrors the behaviour of `layer_ids_with_static`
/// used for node additions: unlayered nodes must be visible in every view.
fn t_prop_iter_layers<L: Into<LayerIter<'a>>>(
self,
layer_ids: L,
prop_id: usize,
) -> impl Iterator<Item = (LayerId, Self::TProps)> + Send + Sync + 'a {
self.layer_ids_iter(layer_ids)
.map(move |id| (id, self.t_prop(id, prop_id)))
}

fn degree(self, layers: &LayerIds, dir: Direction) -> usize;

Expand All @@ -568,7 +598,22 @@ pub trait NodeRefOps<'a>: Copy + Clone + Send + Sync + 'a {
.map_or(0, |id| id as usize)
}

fn internal_num_layers(&self) -> usize;
fn num_layers(&self) -> usize;

fn has_layer_inner(self, layer_id: LayerId) -> bool;
fn has_layer(self, layer_id: LayerId) -> bool;

fn layer_ids_iter<L: Into<LayerIter<'a>>>(
self,
layer_ids: L,
) -> impl Iterator<Item = LayerId> + Send + Sync + 'a {
layer_ids
.into()
.into_iter(self.num_layers())
.filter(move |layer| self.has_layer(*layer))
}

fn has_layer_additions<L: Into<LayerIter<'a>>>(self, layer_ids: L) -> bool {
let layers = layer_ids.into();
!self.node_additions(layers).is_empty() || !self.edge_additions(layers).is_empty()
}
}
30 changes: 13 additions & 17 deletions db4-storage/src/generic_t_props.rs
Original file line number Diff line number Diff line change
@@ -1,15 +1,12 @@
use std::{borrow::Borrow, ops::Range};

use either::Either;
use crate::{generic_time_ops::LayerIter, utils::Iter4};
use itertools::Itertools;
use raphtory_api::core::entities::{
LayerId,
properties::{prop::Prop, tprop::TPropOps},
};
use raphtory_api_macros::box_on_debug_lifetime;
use raphtory_core::{entities::LayerIds, storage::timeindex::EventTime};

use crate::utils::Iter4;
use std::{borrow::Borrow, ops::Range};

/// `WithTProps` defines behavior for types that store multiple temporal
/// properties either in memory or on disk.
Expand Down Expand Up @@ -62,23 +59,23 @@ where
#[derive(Clone, Copy)]
pub struct GenericTProps<'a, Ref: WithTProps<'a>> {
reference: Ref,
layer_id: Either<&'a LayerIds, LayerId>,
layer_id: LayerIter<'a>,
prop_id: usize,
}

impl<'a, Ref: WithTProps<'a>> GenericTProps<'a, Ref> {
pub fn new(reference: Ref, layer_id: &'a LayerIds, prop_id: usize) -> Self {
pub fn new(reference: Ref, layer_id: LayerIter<'a>, prop_id: usize) -> Self {
Self {
reference,
layer_id: Either::Left(layer_id),
layer_id,
prop_id,
}
}

pub fn new_with_layer(reference: Ref, layer_id: LayerId, prop_id: usize) -> Self {
Self {
reference,
layer_id: Either::Right(layer_id),
layer_id: layer_id.into(),
prop_id,
}
}
Expand All @@ -87,18 +84,17 @@ impl<'a, Ref: WithTProps<'a>> GenericTProps<'a, Ref> {
impl<'a, Ref: WithTProps<'a>> GenericTProps<'a, Ref> {
#[inline]
fn tprops(self, prop_id: usize) -> impl Iterator<Item = Ref::TProp> + Send + Sync + 'a {
match self.layer_id {
Either::Left(layer_ids) => {
Either::Left(self.reference.into_t_props_layers(layer_ids, prop_id))
}
Either::Right(layer_id) => {
Either::Right(self.reference.into_t_props(layer_id, prop_id))
}
}
self.layer_id
.into_iter(self.reference.num_layers())
.flat_map(move |layer_id| self.reference.into_t_props(layer_id, prop_id))
}
}

impl<'a, Ref: WithTProps<'a>> TPropOps<'a> for GenericTProps<'a, Ref> {
fn active(&self, w: Range<EventTime>) -> bool {
self.tprops(self.prop_id)
.any(|t_props| t_props.active(w.clone()))
}
fn last_before(&self, t: EventTime) -> Option<(EventTime, Prop)> {
self.tprops(self.prop_id)
.filter_map(|t_props| t_props.last_before(t))
Expand Down
Loading
Loading