pub struct GraphWriter { /* private fields */ }Expand description
Buffered Parquet writer for graph topology and properties.
See the module docs for routing rules and limitations.
Implementations§
Source§impl GraphWriter
impl GraphWriter
Sourcepub fn open(dir: &Path, mode: OntologyMode) -> Result<Self, GfError>
pub fn open(dir: &Path, mode: OntologyMode) -> Result<Self, GfError>
Open (creating if necessary) a project directory for writing.
mode controls edge / property routing. The current wall-clock time is
captured once and reused for all row timestamps.
§Errors
Returns GfError::Storage if the directory cannot be created.
Sourcepub fn open_at(
dir: &Path,
mode: OntologyMode,
now_micros: i64,
) -> Result<Self, GfError>
pub fn open_at( dir: &Path, mode: OntologyMode, now_micros: i64, ) -> Result<Self, GfError>
Like open but with an injected timestamp (microseconds
since the Unix epoch). Used by tests for deterministic output.
§Errors
Returns GfError::Storage if the directory cannot be created.
Sourcepub fn create_node(
&mut self,
node_uuid: Uuid,
type_id: TypeId,
) -> Result<u64, GfError>
pub fn create_node( &mut self, node_uuid: Uuid, type_id: TypeId, ) -> Result<u64, GfError>
Buffer a new node and return its assigned node_id surrogate.
§Errors
Currently infallible; returns Result for forward compatibility.
Sourcepub fn create_node_with_labels(
&mut self,
node_uuid: Uuid,
type_ids: &[TypeId],
) -> Result<u64, GfError>
pub fn create_node_with_labels( &mut self, node_uuid: Uuid, type_ids: &[TypeId], ) -> Result<u64, GfError>
Buffer a node with its complete label set.
The first label is the immutable primary label used for legacy property
file routing. Label membership and labels() semantics use the complete
set. An empty slice creates an unlabelled node.
Sourcepub fn register_existing_node(&mut self, node_uuid: Uuid, node_id: u64)
pub fn register_existing_node(&mut self, node_uuid: Uuid, node_id: u64)
Register an already-persisted node’s identity so a subsequent
create_edge can resolve it as an endpoint — without
writing a new node row or minting a fresh surrogate.
Used by mixed MATCH … CREATE … execution (#703): a node bound by the
preceding MATCH is referenced (its node_uuid/node_id come from the
matched row), not created. Unlike create_node, this
does not push a [NodeRow] or advance next_node_id; it only teaches
the UUID→surrogate map.
Sourcepub fn node_id_for_uuid(&self, node_uuid: &Uuid) -> Option<u64>
pub fn node_id_for_uuid(&self, node_uuid: &Uuid) -> Option<u64>
Return the surrogate ID for a node already known to this write session. This includes both nodes buffered earlier in the statement and persisted nodes registered from a matched input row.
Sourcepub fn create_edge(
&mut self,
edge_uuid: Uuid,
rel_type: &str,
src_uuid: &Uuid,
dst_uuid: &Uuid,
) -> Result<u64, GfError>
pub fn create_edge( &mut self, edge_uuid: Uuid, rel_type: &str, src_uuid: &Uuid, dst_uuid: &Uuid, ) -> Result<u64, GfError>
Buffer a new edge and return its assigned edge_id surrogate.
Both endpoints must have been registered via
create_node first so their node_id surrogates
can be resolved.
§Errors
Returns GfError::Storage if either endpoint UUID is unknown.
Sourcepub fn set_properties(
&mut self,
node_uuid: &Uuid,
entity_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) -> Result<(), GfError>
pub fn set_properties( &mut self, node_uuid: &Uuid, entity_type: Option<&str>, props: HashMap<String, IrLiteral>, ) -> Result<(), GfError>
Buffer a property row for a node.
In Strict / Advisory mode with a known entity_type, properties route to
properties/TYPENAME.parquet; otherwise (exploratory, or no entity type)
they route to properties/_untyped.parquet.
§Errors
Currently infallible; returns Result for forward compatibility.
Sourcepub fn set_edge_properties(
&mut self,
edge_uuid: &Uuid,
rel_type: Option<&str>,
props: HashMap<String, IrLiteral>,
) -> Result<(), GfError>
pub fn set_edge_properties( &mut self, edge_uuid: &Uuid, rel_type: Option<&str>, props: HashMap<String, IrLiteral>, ) -> Result<(), GfError>
Buffer a property row for an edge, keyed by edge_uuid.
Edge properties route to edge_properties/<REL_TYPE>.parquet by relation
name in every mode (unlike node properties, which fall back to
_untyped in exploratory mode). The read side resolves the file stem from
the relation name, so a single namespace keyed by rel type keeps write and
read in lock-step and avoids colliding with the node properties/
directory. A None rel_type (an edge created without a known relation
name) routes to the _untyped catch-all.
§Errors
Currently infallible; returns Result for forward compatibility.
Sourcepub fn contains_pending_node(&self, node_uuid: &[u8; 16]) -> bool
pub fn contains_pending_node(&self, node_uuid: &[u8; 16]) -> bool
Whether a node with this uuid is buffered (created in this statement and not yet flushed or cancelled).
Sourcepub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32>
pub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32>
Return distinct label tokens on buffered nodes selected by UUID.
Sourcepub fn pending_nodes_batch(&self) -> Result<RecordBatch, GfError>
pub fn pending_nodes_batch(&self) -> Result<RecordBatch, GfError>
Materialize the currently buffered node topology without consuming it. Statement-local reads use this as an in-memory overlay before commit.
§Errors
Returns [GfError::Parquet] if the buffered values cannot form the
canonical topology batch.
Sourcepub fn find_pending_node(
&self,
labels: &[u32],
properties: &[(String, IrLiteral)],
) -> Option<PendingNodeMatch>
pub fn find_pending_node( &self, labels: &[u32], properties: &[(String, IrLiteral)], ) -> Option<PendingNodeMatch>
Find a buffered node whose labels and properties satisfy a MERGE pattern.
Sourcepub fn find_pending_nodes(
&self,
labels: &[u32],
properties: &[(String, IrLiteral)],
) -> Vec<PendingNodeMatch> ⓘ
pub fn find_pending_nodes( &self, labels: &[u32], properties: &[(String, IrLiteral)], ) -> Vec<PendingNodeMatch> ⓘ
Return every buffered node matching all requested labels and properties.
Sourcepub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool
pub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool
Whether an edge with this uuid is buffered.
Sourcepub fn find_pending_edge(
&self,
rel_type: &str,
src: &[u8; 16],
dst: &[u8; 16],
undirected: bool,
properties: &[(String, IrLiteral)],
) -> Option<([u8; 16], [u8; 16], [u8; 16], HashMap<String, IrLiteral>)>
pub fn find_pending_edge( &self, rel_type: &str, src: &[u8; 16], dst: &[u8; 16], undirected: bool, properties: &[(String, IrLiteral)], ) -> Option<([u8; 16], [u8; 16], [u8; 16], HashMap<String, IrLiteral>)>
Find a buffered edge matching type, endpoints, direction, and properties.
Sourcepub fn pending_incident_edge_uuids<S: BuildHasher>(
&self,
nodes: &HashSet<[u8; 16], S>,
) -> Vec<[u8; 16]>
pub fn pending_incident_edge_uuids<S: BuildHasher>( &self, nodes: &HashSet<[u8; 16], S>, ) -> Vec<[u8; 16]>
The uuids of buffered edges incident (as src or dst) to any of nodes.
The pending complement of
incident_edge_uuids, which only sees
committed files: openCypher’s “cannot delete a node that still has
relationships” must also count edges created earlier in the same
statement.
Sourcepub fn cancel_nodes<S: BuildHasher>(
&mut self,
targets: &HashSet<[u8; 16], S>,
) -> u64
pub fn cancel_nodes<S: BuildHasher>( &mut self, targets: &HashSet<[u8; 16], S>, ) -> u64
Drop buffered nodes (and their buffered property rows) whose uuid is in
targets, so a created-then-deleted node never hits disk. Forgets the
uuid→surrogate mapping too: the entity no longer exists, so a later
create_edge referencing it must fail. Returns the node rows dropped.
Sourcepub fn cancel_edges<S: BuildHasher>(
&mut self,
targets: &HashSet<[u8; 16], S>,
) -> u64
pub fn cancel_edges<S: BuildHasher>( &mut self, targets: &HashSet<[u8; 16], S>, ) -> u64
Drop buffered edges (and their buffered property rows) whose uuid is in
targets. Returns the edge rows dropped.
Sourcepub fn merge_pending_node_props(
&mut self,
node_uuid: &[u8; 16],
entity_type: Option<&str>,
props: HashMap<String, IrLiteral>,
)
pub fn merge_pending_node_props( &mut self, node_uuid: &[u8; 16], entity_type: Option<&str>, props: HashMap<String, IrLiteral>, )
Merge props into the buffered property row of a pending node
(SET on an entity created earlier in this statement), inserting a row
if it has none yet. Same stem routing as
set_properties.
Sourcepub fn add_pending_node_labels(
&mut self,
node_uuid: &[u8; 16],
labels: &[u32],
) -> u64
pub fn add_pending_node_labels( &mut self, node_uuid: &[u8; 16], labels: &[u32], ) -> u64
Add labels to a node buffered by this writer, preserving its primary label.
Sourcepub fn remove_pending_node_labels(
&mut self,
node_uuid: &[u8; 16],
labels: &[u32],
) -> u64
pub fn remove_pending_node_labels( &mut self, node_uuid: &[u8; 16], labels: &[u32], ) -> u64
Remove labels from a node buffered by this writer. The immutable scalar
type_id remains only as the property-file routing key; type_ids is
the authoritative membership set.
Sourcepub fn merge_pending_edge_props(
&mut self,
edge_uuid: &[u8; 16],
rel_type: Option<&str>,
props: HashMap<String, IrLiteral>,
)
pub fn merge_pending_edge_props( &mut self, edge_uuid: &[u8; 16], rel_type: Option<&str>, props: HashMap<String, IrLiteral>, )
Edge analogue of
merge_pending_node_props; same stem
routing as set_edge_properties.
Sourcepub fn remove_pending_node_props(
&mut self,
node_uuid: &[u8; 16],
keys: &HashSet<String>,
)
pub fn remove_pending_node_props( &mut self, node_uuid: &[u8; 16], keys: &HashSet<String>, )
Remove keys from a pending node’s buffered property rows (REMOVE on
an entity created earlier in this statement). Absent keys/rows are
no-ops (openCypher). Scans every stem — a REMOVE clause does not know
the routing the CREATE used.
Sourcepub fn remove_pending_edge_props(
&mut self,
edge_uuid: &[u8; 16],
keys: &HashSet<String>,
)
pub fn remove_pending_edge_props( &mut self, edge_uuid: &[u8; 16], keys: &HashSet<String>, )
Edge analogue of
remove_pending_node_props.
Sourcepub fn flush(&mut self) -> Result<(), GfError>
pub fn flush(&mut self) -> Result<(), GfError>
Merge all buffered rows with any existing on-disk data and write the result, then clear the row buffers.
Only creates a subdirectory when there are rows to write into it. Each target file is read, concatenated with the new rows (property files are decoded and re-inferred so the dynamic schema evolves), and rewritten — so separate write sessions accumulate (#733). All files stage and commit as one batch (#790), nodes first: a failure while building any file leaves the prior state fully intact, and a (rare) rename-phase failure can commit a node without its edges, never the reverse.
A batch that stages topology files bumps the project
topology_generation counter before committing (#759); property-only
flushes do not bump.
§Errors
Returns GfError::Storage on any I/O, Arrow, or Parquet failure.
Sourcepub fn take_pending_delta(&mut self) -> Vec<DeltaEdge>
pub fn take_pending_delta(&mut self) -> Vec<DeltaEdge>
Drain the edges captured for the next adjacency delta segment. The
statement driver (#792) calls this after flush_into to write or
discard the segment around its own commit (#765).
Sourcepub fn write_segment_best_effort(&self, generation: u64, edges: &[DeltaEdge])
pub fn write_segment_best_effort(&self, generation: u64, edges: &[DeltaEdge])
Best-effort write of the delta segment for generation — only when the
adjacency capability directory exists (never grow deltas/ for a project
that has no index). A failed write costs at most one future rebuild and
must never fail a already-committed flush, so the error is swallowed.
Sourcepub fn flush_into(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError>
pub fn flush_into(&mut self, staged: &mut RewriteBatch) -> Result<(), GfError>
Stage all buffered rows into staged (committed by the caller) and
clear the row buffers.
Reads through staged and restages: a file this statement already
staged (e.g. a DELETE rewrite of the same property file, #792) is the
merge base and its entry is replaced in place with the net content —
files new to the batch append after it, so created edges commit after
topology/nodes.parquet whether or not a delete staged it earlier.
On success the buffers are cleared even though nothing is committed yet; the writer is not reusable if the caller’s commit fails.
§Errors
Returns GfError::Storage on any I/O, Arrow, or Parquet failure.
Auto Trait Implementations§
impl Freeze for GraphWriter
impl RefUnwindSafe for GraphWriter
impl Send for GraphWriter
impl Sync for GraphWriter
impl Unpin for GraphWriter
impl UnsafeUnpin for GraphWriter
impl UnwindSafe for GraphWriter
Blanket Implementations§
impl<T> Allocation for T
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read more