Skip to main content

GraphWriter

Struct GraphWriter 

Source
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

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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).

Source

pub fn pending_node_labels(&self, targets: &HashSet<[u8; 16]>) -> HashSet<u32>

Return distinct label tokens on buffered nodes selected by UUID.

Source

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.

Source

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.

Source

pub fn find_pending_nodes( &self, labels: &[u32], properties: &[(String, IrLiteral)], ) -> Vec<PendingNodeMatch> ⓘ

Return every buffered node matching all requested labels and properties.

Source

pub fn contains_pending_edge(&self, edge_uuid: &[u8; 16]) -> bool

Whether an edge with this uuid is buffered.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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.

Source

pub fn remove_pending_edge_props( &mut self, edge_uuid: &[u8; 16], keys: &HashSet<String>, )

Edge analogue of remove_pending_node_props.

Source

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.

Source

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).

Source

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.

Source

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§

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts 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
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V