use std::collections::{BTreeMap, BTreeSet};
use std::sync::Mutex;
use async_trait::async_trait;
use graph_storage_sdk::models::{
AdjacencyEntry, AdjacencySide, AdmissionBasis, ComponentReadiness, DeleteOutcome,
DeleteRequest, EdgeKey, EdgeView, ElementEnvelope, GraphRevision, GtsTypeId, IngestOutcome,
IngestRequest, ItemError, ItemFamily, ItemOutcome, NodeId, NodeKey, NodeRow, NodeView,
OnExisting, Page, ProjectionRequest, ReadSnapshot, ReadinessState, RegisteredType,
SchemaDiagnostic, SearchMode, SearchRequest, SearchResponse, SourceNamespaceOwner,
StoreCapabilities, Subject, TopologyPage, TopologyRequest, TypeChange, TypeChangeState,
TypeIdSet, TypeOutcome, TypeQuery, TypeRecord, TypeRegistration, TypeRegistrationOptions,
};
use graph_storage_sdk::plugin_api::{
EmbeddingPlan, EmbeddingState, GraphStoreError, GraphStoreV1, StoreCtx, VectorArm,
};
use time::OffsetDateTime;
use uuid::Uuid;
use crate::domain::embedding::{PlannedVector, StoredVector, VectorOutcome, decide_vector};
use crate::domain::tally::IngestTally;
use crate::domain::{evolution, identity, ontology, ownership, projection};
#[derive(Clone)]
struct FakeNode {
id: i64,
key: String,
type_id: String,
name: Option<String>,
search_text: String,
payload: Option<serde_json::Value>,
embedding: Option<Vec<f32>>,
embedding_epoch: Option<i64>,
embedding_input_hash: Option<String>,
version: i64,
deleted: bool,
audit: FakeAudit,
}
#[derive(Clone)]
struct FakeEdge {
key: String,
type_id: String,
src: i64,
dst: i64,
discriminator: Option<String>,
payload: Option<serde_json::Value>,
deleted: bool,
scope: Option<(String, String)>,
audit: FakeAudit,
}
#[derive(Clone)]
struct FakeAudit {
created_at: OffsetDateTime,
created_by: Subject,
updated_at: OffsetDateTime,
updated_by: Subject,
deleted_at: Option<OffsetDateTime>,
deleted_by: Option<Subject>,
}
impl FakeAudit {
fn created(subject: &Subject) -> Self {
let now = OffsetDateTime::now_utc();
Self {
created_at: now,
created_by: subject.clone(),
updated_at: now,
updated_by: subject.clone(),
deleted_at: None,
deleted_by: None,
}
}
fn updated(&mut self, subject: &Subject) {
self.updated_at = OffsetDateTime::now_utc();
self.updated_by = subject.clone();
self.deleted_at = None;
self.deleted_by = None;
}
fn tombstoned(&mut self, subject: &Subject) {
self.deleted_at = Some(OffsetDateTime::now_utc());
self.deleted_by = Some(subject.clone());
}
fn envelope(&self, key: String, tenant_id: Uuid, revision: GraphRevision) -> ElementEnvelope {
ElementEnvelope {
tenant_id,
key,
created_at: self.created_at,
created_by: self.created_by.clone(),
updated_at: self.updated_at,
updated_by: self.updated_by.clone(),
deleted_at: self.deleted_at,
deleted_by: self.deleted_by.clone(),
graph_revision: revision,
}
}
}
#[derive(Clone)]
struct Receipt {
request_hash: String,
epoch: i64,
outcome: IngestOutcome,
}
#[derive(Clone)]
struct FakeScope {
owner_producer: String,
generation: i64,
request_hash: String,
}
#[derive(Default)]
struct Tenant {
types: BTreeMap<String, TypeRecord>,
namespaces: BTreeMap<String, SourceNamespaceOwner>,
index_kinds: BTreeMap<String, BTreeMap<String, ontology::ScalarKind>>,
nodes: Vec<FakeNode>,
edges: Vec<FakeEdge>,
revision: i64,
receipts: BTreeMap<(String, String), Receipt>,
scopes: BTreeMap<(String, String), FakeScope>,
snapshots: BTreeMap<Uuid, Box<TenantData>>,
}
#[derive(Clone, Default)]
struct TenantData {
nodes: Vec<FakeNode>,
edges: Vec<FakeEdge>,
revision: i64,
}
pub struct FakeGraphStore {
epoch: i64,
tenants: Mutex<BTreeMap<Uuid, Tenant>>,
next_id: std::sync::atomic::AtomicI64,
max_chain_depth: usize,
declines_snapshots: bool,
drift_per_revision_read: i64,
revision_reads: std::sync::atomic::AtomicI64,
rows_hydrated: std::sync::atomic::AtomicU64,
hydrate_calls: std::sync::atomic::AtomicU64,
node_types_answer: NodeTypesAnswer,
}
#[derive(Clone, Copy, Debug)]
enum NodeTypesAnswer {
Typed,
Unsupported,
Failing,
}
impl Default for FakeGraphStore {
fn default() -> Self {
Self::new()
}
}
impl FakeGraphStore {
#[must_use]
pub fn new() -> Self {
Self {
epoch: 1,
tenants: Mutex::new(BTreeMap::new()),
next_id: std::sync::atomic::AtomicI64::new(0),
max_chain_depth: 3,
declines_snapshots: false,
drift_per_revision_read: 0,
revision_reads: std::sync::atomic::AtomicI64::new(0),
rows_hydrated: std::sync::atomic::AtomicU64::new(0),
hydrate_calls: std::sync::atomic::AtomicU64::new(0),
node_types_answer: NodeTypesAnswer::Typed,
}
}
#[must_use]
pub fn without_node_types() -> Self {
Self {
node_types_answer: NodeTypesAnswer::Unsupported,
..Self::new()
}
}
#[must_use]
pub fn failing_node_types() -> Self {
Self {
node_types_answer: NodeTypesAnswer::Failing,
..Self::new()
}
}
#[must_use]
pub fn hydrate_calls(&self) -> u64 {
self.hydrate_calls
.load(std::sync::atomic::Ordering::Relaxed)
}
#[must_use]
pub fn rows_hydrated(&self) -> u64 {
self.rows_hydrated
.load(std::sync::atomic::Ordering::Relaxed)
}
#[must_use]
pub fn declining_snapshots(drift: i64) -> Self {
Self {
declines_snapshots: true,
drift_per_revision_read: drift,
..Self::new()
}
}
#[must_use]
pub fn with_max_chain_depth(mut self, depth: usize) -> Self {
self.max_chain_depth = depth;
self
}
fn revision_of(&self, tenant: &Tenant) -> GraphRevision {
GraphRevision {
source_epoch: self.epoch,
revision: tenant.revision,
}
}
}
fn snapshot_data(tenant: &Tenant) -> TenantData {
TenantData {
nodes: tenant.nodes.clone(),
edges: tenant.edges.clone(),
revision: tenant.revision,
}
}
fn scope_admits(scope: &toolkit_security::AccessScope, tenant: Uuid) -> bool {
if scope.is_deny_all() {
return false;
}
if scope.is_unconstrained() {
return true;
}
let tenants = scope.all_uuid_values_for(toolkit_security::pep_properties::OWNER_TENANT_ID);
tenants.is_empty() || tenants.contains(&tenant)
}
fn visible<'a>(tenant: &'a Tenant, ctx: &StoreCtx<'_>) -> (&'a [FakeNode], &'a [FakeEdge], i64) {
if let Some(snapshot) = ctx.snapshot
&& let Some(data) = tenant.snapshots.get(&snapshot.id)
{
return (&data.nodes, &data.edges, data.revision);
}
(&tenant.nodes, &tenant.edges, tenant.revision)
}
#[async_trait]
impl GraphStoreV1 for FakeGraphStore {
fn capabilities(&self) -> StoreCapabilities {
StoreCapabilities {
scope_replace: true,
snapshots: !self.declines_snapshots,
vector_search: true,
labels: false,
chunks: false,
topology: false,
}
}
async fn register_types_with(
&self,
ctx: &StoreCtx<'_>,
batch: Vec<TypeRegistration>,
options: TypeRegistrationOptions,
) -> Result<Vec<RegisteredType>, GraphStoreError> {
if let Some(duplicate) =
ontology::duplicate_type_id(batch.iter().map(|r| r.type_id.as_str()))
{
return Err(GraphStoreError::InvalidQuery {
what: format!(
"`{duplicate}` is registered twice in one batch; a batch is one act and \
cannot name a type twice"
),
});
}
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.entry(ctx.tenant).or_default();
let mut prepared = Vec::new();
for registration in &batch {
let chain = ontology::ancestors(®istration.type_id);
let mut ancestors: Vec<(String, serde_json::Value)> = Vec::new();
for ancestor in &chain[..chain.len().saturating_sub(1)] {
let schema = batch
.iter()
.find(|r| &r.type_id == ancestor)
.map(|r| r.schema.clone())
.or_else(|| tenant.types.get(ancestor).map(|t| t.schema.clone()))
.ok_or_else(|| {
validation(
0,
ItemFamily::Node,
®istration.type_id,
&format!("ancestor `{ancestor}` is not registered"),
)
})?;
ancestors.push((ancestor.clone(), schema));
}
let refs: Vec<&serde_json::Value> =
ancestors.iter().map(|(_, schema)| schema).collect();
let descriptor = ontology::analyze(
®istration.type_id,
®istration.schema,
&refs,
self.max_chain_depth,
)
.map_err(|error| {
validation(
0,
ItemFamily::Node,
®istration.type_id,
&error.to_string(),
)
})?;
{
let mut chain: Vec<(String, serde_json::Value)> = ancestors.clone();
chain.push((descriptor.type_id.clone(), descriptor.schema.clone()));
ontology::ChainValidator::compile(&descriptor.schema, chain).map_err(|error| {
validation(0, ItemFamily::Node, &descriptor.type_id, &error.to_string())
})?;
}
if let Some(spec) = options.migration_for(&descriptor.type_id)
&& !tenant.types.contains_key(&descriptor.type_id)
{
return Err(GraphStoreError::InvalidQuery {
what: format!(
"the migration for `{}` has nothing to migrate: the type is not \
registered yet, so it holds no rows",
spec.type_id
),
});
}
let decided = match tenant.types.get(&descriptor.type_id).cloned() {
None => Decided {
outcome: TypeOutcome::Created,
basis: None,
change: fresh_change(&descriptor.type_id),
revision: 1,
created_at: OffsetDateTime::now_utc(),
rewrites: Vec::new(),
},
Some(existing) => {
decide_existing(tenant, &existing, &descriptor, &ancestors, &options)?
}
};
let kinds: BTreeMap<String, ontology::ScalarKind> = descriptor
.index_paths
.iter()
.map(|p| (p.pointer.clone(), p.kind))
.collect();
prepared.push((
TypeRecord {
type_id: descriptor.type_id,
type_uuid: descriptor.type_uuid,
kind: descriptor.kind,
is_abstract: descriptor.is_abstract,
schema: descriptor.schema,
effective_traits: descriptor.effective_traits,
created_at: decided.created_at,
revision: decided.revision,
},
kinds,
decided,
));
}
let mut out = Vec::with_capacity(prepared.len());
for (record, kinds, decided) in prepared {
if !options.dry_run && decided.outcome != TypeOutcome::Unchanged {
tenant.types.insert(record.type_id.clone(), record.clone());
tenant.index_kinds.insert(record.type_id.clone(), kinds);
for rewrite in &decided.rewrites {
match rewrite.family {
ItemFamily::Node => {
if let Some(node) =
tenant.nodes.iter_mut().find(|n| n.key == rewrite.key)
{
node.payload = Some(rewrite.payload.clone());
node.version += 1;
if !record.effective_traits.vector_search.is_empty() {
node.embedding_epoch = None;
}
node.audit.updated(&ctx.subject);
}
}
ItemFamily::Edge => {
if let Some(edge) =
tenant.edges.iter_mut().find(|e| e.key == rewrite.key)
{
edge.payload = Some(rewrite.payload.clone());
edge.audit.updated(&ctx.subject);
}
}
}
}
if decided.outcome == TypeOutcome::Updated {
tenant.revision += 1;
}
}
out.push(RegisteredType {
record,
outcome: decided.outcome,
basis: decided.basis,
change: Some(decided.change),
});
}
Ok(out)
}
async fn probe_readiness(&self) -> Vec<ComponentReadiness> {
vec![
ComponentReadiness::healthy(graph_storage_sdk::models::DATABASE),
ComponentReadiness::new(
graph_storage_sdk::models::SQLPGQ,
ReadinessState::Degraded,
"this store has no property graph",
"nothing: every traversal is served by the two-query hop",
"not applicable to an in-memory store",
),
]
}
async fn list_source_namespaces(
&self,
ctx: &StoreCtx<'_>,
) -> Result<Vec<SourceNamespaceOwner>, GraphStoreError> {
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
Ok(tenants
.get(&ctx.tenant)
.map(|tenant| tenant.namespaces.values().cloned().collect())
.unwrap_or_default())
}
async fn transfer_source_namespace(
&self,
ctx: &StoreCtx<'_>,
namespace: &str,
owner_principal: &str,
) -> Result<SourceNamespaceOwner, GraphStoreError> {
if owner_principal.trim().is_empty() {
return Err(GraphStoreError::InvalidQuery {
what: "a transfer needs the principal to transfer to".to_owned(),
});
}
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.entry(ctx.tenant).or_default();
let now = OffsetDateTime::now_utc();
let previous = tenant
.namespaces
.get(namespace)
.map(|row| row.owner_principal.clone());
let row = SourceNamespaceOwner {
namespace: namespace.to_owned(),
owner_principal: owner_principal.to_owned(),
claimed_at: tenant
.namespaces
.get(namespace)
.map_or(now, |row| row.claimed_at),
previous_owner: previous.filter(|owner| owner != owner_principal),
transferred_at: Some(now),
transferred_by: Some(ctx.subject.clone()),
};
tenant.namespaces.insert(namespace.to_owned(), row.clone());
Ok(row)
}
async fn get_type(
&self,
ctx: &StoreCtx<'_>,
id: &GtsTypeId,
) -> Result<TypeRecord, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Err(GraphStoreError::NotFound);
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
tenants
.get(&ctx.tenant)
.and_then(|t| t.types.get(id).cloned())
.ok_or(GraphStoreError::NotFound)
}
async fn list_types(
&self,
ctx: &StoreCtx<'_>,
query: TypeQuery,
) -> Result<Page<TypeRecord>, GraphStoreError> {
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(Page {
items: Vec::new(),
next_cursor: None,
revision: GraphRevision {
source_epoch: self.epoch,
revision: 0,
},
});
};
let (items, next_cursor) = catalogue_page(tenant, &query);
Ok(Page {
items,
next_cursor,
revision: self.revision_of(tenant),
})
}
async fn resolve_type_set(
&self,
ctx: &StoreCtx<'_>,
patterns: &[String],
) -> Result<TypeIdSet, GraphStoreError> {
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(TypeIdSet::default());
};
Ok(TypeIdSet(
tenant
.types
.keys()
.filter(|id| ontology::matches_any_pattern(id, patterns).unwrap_or(false))
.cloned()
.collect(),
))
}
async fn ingest(
&self,
ctx: &StoreCtx<'_>,
req: IngestRequest,
embedding: EmbeddingPlan,
) -> Result<IngestOutcome, GraphStoreError> {
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.entry(ctx.tenant).or_default();
let request_hash = identity::ingest_request_hash(&req);
let producer = ctx.subject.principal();
if let Some(key) = &req.idempotency_key
&& let Some(receipt) = tenant.receipts.get(&(producer.clone(), key.clone()))
{
if receipt.request_hash != request_hash {
return Err(GraphStoreError::IdempotencyMismatch);
}
if receipt.epoch != self.epoch {
return Err(GraphStoreError::IdempotencyExpired);
}
let mut outcome = receipt.outcome.clone();
outcome.replayed = true;
return Ok(outcome);
}
if let Some(replace) = &req.replace_scope {
fence(tenant, &producer, replace, &request_hash)?;
}
let writer = ctx.subject.principal();
let mut claims: Vec<(String, SourceNamespaceOwner)> = Vec::new();
for spec in &req.nodes {
let family = tenant
.types
.get(&spec.type_id)
.and_then(|record| record.effective_traits.family.clone());
let namespace = match ownership::namespace_of(family.as_deref(), spec.payload.as_ref())
.map_err(|error| {
validation(0, ItemFamily::Node, &spec.type_id, &error.to_string())
})? {
ownership::Namespaced::None => continue,
ownership::Namespaced::Under(namespace) => namespace.to_owned(),
};
let held = tenant
.namespaces
.get(&namespace)
.map(|row| row.owner_principal.clone())
.or_else(|| {
claims
.iter()
.find(|(name, _)| name == &namespace)
.map(|(_, row)| row.owner_principal.clone())
});
match ownership::decide(held.as_deref(), &writer) {
ownership::Claim::Allowed => {}
ownership::Claim::Forbidden => {
return Err(GraphStoreError::SourceNamespaceForbidden { namespace });
}
ownership::Claim::Take => claims.push((
namespace.clone(),
SourceNamespaceOwner {
namespace,
owner_principal: writer.clone(),
claimed_at: OffsetDateTime::now_utc(),
previous_owner: None,
transferred_at: None,
transferred_by: None,
},
)),
}
}
let managed_node_types: BTreeSet<String> = tenant
.types
.values()
.filter(|record| {
record.kind == graph_storage_sdk::models::TypeKind::Node
&& record.effective_traits.scope_managed
})
.map(|record| record.type_id.clone())
.collect();
let static_edge_types: BTreeSet<String> = tenant
.types
.values()
.filter(|record| record.effective_traits.family.as_deref() == Some("static"))
.map(|record| record.type_id.clone())
.collect();
let mut nodes = tenant.nodes.clone();
let mut edges = tenant.edges.clone();
let mut next_id = self.next_id.load(std::sync::atomic::Ordering::SeqCst);
let mut tally = IngestTally::new(req.options.report_per_item);
let mut changed = false;
let mut state = BatchState {
next_id: &mut next_id,
tally: &mut tally,
subject: &ctx.subject,
};
for (index, spec) in req.nodes.iter().enumerate() {
let decided = embedding.nodes.get(index).ok_or_else(|| {
GraphStoreError::Internal(format!(
"embedding plan covers {} nodes; the batch has {}",
embedding.nodes.len(),
req.nodes.len()
))
})?;
let vector = plan_vector(
nodes.iter().find(|n| n.key == spec.node_key),
PlannedVector {
decided,
active_epoch: embedding.epoch,
},
);
changed |= apply_node(tenant, &mut nodes, &edges, &mut state, index, spec, vector)?;
}
for (index, spec) in req.edges.iter().enumerate() {
changed |= apply_edge(
tenant,
&mut nodes,
&mut edges,
&mut state,
index,
spec,
&EdgeDeclaration {
create_phantoms: req.options.create_phantoms.unwrap_or(true),
scope: req
.replace_scope
.as_ref()
.map(|replace| (replace.attribute.clone(), replace.value.clone())),
},
)?;
}
if let Some(replace) = &req.replace_scope {
changed |= replace_scope(
ScopeReplacement {
nodes: &mut nodes,
edges: &mut edges,
managed_node_types: &managed_node_types,
static_edge_types: &static_edge_types,
types: &tenant.types,
},
&req,
replace,
&mut tally,
);
}
tenant.nodes = nodes;
tenant.edges = edges;
self.next_id
.store(next_id, std::sync::atomic::Ordering::SeqCst);
for (namespace, row) in claims {
tenant.namespaces.insert(namespace, row);
}
if changed {
tenant.revision += 1;
}
if let Some(replace) = &req.replace_scope {
let key = (replace.attribute.clone(), replace.value.clone());
let owner_producer = tenant
.scopes
.get(&key)
.map(|scope| scope.owner_producer.clone())
.filter(|owner| !owner.is_empty())
.unwrap_or_else(|| producer.clone());
tenant.scopes.insert(
key,
FakeScope {
owner_producer,
generation: replace.generation,
request_hash: request_hash.clone(),
},
);
}
let (counts, per_item_nodes, per_item_edges) = tally.into_parts();
let outcome = IngestOutcome {
revision: self.revision_of(tenant),
replayed: false,
counts,
per_item_nodes,
per_item_edges,
};
if let Some(key) = &req.idempotency_key {
tenant.receipts.insert(
(producer, key.clone()),
Receipt {
request_hash,
epoch: self.epoch,
outcome: outcome.clone(),
},
);
}
Ok(outcome)
}
async fn soft_delete(
&self,
ctx: &StoreCtx<'_>,
req: DeleteRequest,
) -> Result<DeleteOutcome, GraphStoreError> {
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.entry(ctx.tenant).or_default();
let (nodes, edges) = match req {
DeleteRequest::Node(key) => {
let Some(index) = tenant.nodes.iter().position(|n| n.key == key && !n.deleted)
else {
let tombstoned = tenant.nodes.iter().any(|n| n.key == key);
return settle_no_op(self, tenant, tombstoned);
};
let id = tenant.nodes[index].id;
let mut tombstoned = 0u64;
for edge in &mut tenant.edges {
if !edge.deleted && (edge.src == id || edge.dst == id) {
edge.deleted = true;
edge.audit.tombstoned(&ctx.subject);
tombstoned += 1;
}
}
tenant.nodes[index].deleted = true;
tenant.nodes[index].audit.tombstoned(&ctx.subject);
(1u64, tombstoned)
}
DeleteRequest::Edge(key) => {
let Some(position) = tenant.edges.iter().position(|e| e.key == key && !e.deleted)
else {
let tombstoned = tenant.edges.iter().any(|e| e.key == key);
return settle_no_op(self, tenant, tombstoned);
};
let (src, dst) = (tenant.edges[position].src, tenant.edges[position].dst);
let visible = |id: i64| tenant.nodes.iter().any(|n| n.id == id && !n.deleted);
if !(visible(src) && visible(dst)) {
return Err(GraphStoreError::NotFound);
}
let edge = &mut tenant.edges[position];
edge.deleted = true;
edge.audit.tombstoned(&ctx.subject);
(0u64, 1u64)
}
};
tenant.revision += 1;
Ok(DeleteOutcome {
revision: self.revision_of(tenant),
tombstoned_nodes: nodes,
tombstoned_edges: edges,
})
}
async fn begin_read(&self, ctx: &StoreCtx<'_>) -> Result<ReadSnapshot, GraphStoreError> {
assert!(
!self.declines_snapshots,
"a caller asked for a snapshot from a store that declares it has none; \
the declaration is the contract, not a hint"
);
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.entry(ctx.tenant).or_default();
let id = Uuid::now_v7();
let data = snapshot_data(tenant);
let revision = self.revision_of(tenant);
tenant.snapshots.insert(id, Box::new(data));
Ok(ReadSnapshot { id, revision })
}
async fn end_read(&self, snapshot: ReadSnapshot) -> Result<(), GraphStoreError> {
let mut tenants = self.tenants.lock().map_err(|_| poisoned())?;
for tenant in tenants.values_mut() {
tenant.snapshots.remove(&snapshot.id);
}
Ok(())
}
async fn revision(&self, ctx: &StoreCtx<'_>) -> Result<GraphRevision, GraphStoreError> {
let drift = self.drift_per_revision_read
* self
.revision_reads
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let observed = tenants.get(&ctx.tenant).map_or(
GraphRevision {
source_epoch: self.epoch,
revision: 0,
},
|t| self.revision_of(t),
);
Ok(GraphRevision {
revision: observed.revision + drift,
..observed
})
}
async fn get_node(
&self,
ctx: &StoreCtx<'_>,
key: &NodeKey,
adjacency_limit: u32,
) -> Result<NodeView, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Err(GraphStoreError::NotFound);
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.get(&ctx.tenant).ok_or(GraphStoreError::NotFound)?;
let (nodes, edges, revision) = visible(tenant, ctx);
let node = nodes
.iter()
.find(|n| &n.key == key && !n.deleted)
.ok_or(GraphStoreError::NotFound)?;
let by_id: BTreeMap<i64, &FakeNode> = nodes
.iter()
.filter(|n| !n.deleted)
.map(|n| (n.id, n))
.collect();
let mut adjacency = Vec::new();
let mut truncated = false;
let mut taken = [0u32; 2];
for edge in edges.iter().filter(|e| !e.deleted) {
let (side, other, slot) = if edge.src == node.id {
(AdjacencySide::Outgoing, edge.dst, 0)
} else if edge.dst == node.id {
(AdjacencySide::Incoming, edge.src, 1)
} else {
continue;
};
if taken[slot] >= adjacency_limit {
truncated = true;
continue;
}
taken[slot] += 1;
let Some(neighbour) = by_id.get(&other) else {
continue;
};
adjacency.push(AdjacencyEntry {
edge_key: edge.key.clone(),
edge_type_id: edge.type_id.clone(),
side,
neighbor_key: neighbour.key.clone(),
neighbor_type_id: neighbour.type_id.clone(),
});
}
Ok(view_of(
node,
ctx.tenant,
GraphRevision {
source_epoch: self.epoch,
revision,
},
adjacency,
truncated,
))
}
async fn get_edge(
&self,
ctx: &StoreCtx<'_>,
key: &EdgeKey,
) -> Result<EdgeView, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Err(GraphStoreError::NotFound);
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let tenant = tenants.get(&ctx.tenant).ok_or(GraphStoreError::NotFound)?;
let (nodes, edges, revision) = visible(tenant, ctx);
let edge = edges
.iter()
.find(|e| &e.key == key && !e.deleted)
.ok_or(GraphStoreError::NotFound)?;
let endpoint = |id: i64| nodes.iter().find(|n| n.id == id && !n.deleted);
let (Some(src), Some(dst)) = (endpoint(edge.src), endpoint(edge.dst)) else {
return Err(GraphStoreError::NotFound);
};
Ok(EdgeView {
edge_key: edge.key.clone(),
edge_type_id: edge.type_id.clone(),
src: src.key.clone(),
dst: dst.key.clone(),
discriminator: edge.discriminator.clone(),
payload: edge.payload.clone(),
envelope: edge.audit.envelope(
edge.key.clone(),
ctx.tenant,
GraphRevision {
source_epoch: self.epoch,
revision,
},
),
})
}
async fn hydrate_nodes(
&self,
ctx: &StoreCtx<'_>,
ids: &[NodeId],
) -> Result<Vec<NodeView>, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(Vec::new());
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(Vec::new());
};
let (nodes, _, revision) = visible(tenant, ctx);
let revision = GraphRevision {
source_epoch: self.epoch,
revision,
};
let views: Vec<NodeView> = ids
.iter()
.filter_map(|id| nodes.iter().find(|n| n.id == *id && !n.deleted))
.map(|n| view_of(n, ctx.tenant, revision, Vec::new(), false))
.collect();
self.rows_hydrated
.fetch_add(views.len() as u64, std::sync::atomic::Ordering::Relaxed);
self.hydrate_calls
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
Ok(views)
}
async fn node_types(
&self,
ctx: &StoreCtx<'_>,
ids: &[NodeId],
) -> Result<Vec<(NodeId, GtsTypeId)>, GraphStoreError> {
match self.node_types_answer {
NodeTypesAnswer::Typed => {}
NodeTypesAnswer::Unsupported => {
return Err(GraphStoreError::Unsupported { what: "node_types" });
}
NodeTypesAnswer::Failing => {
return Err(GraphStoreError::Unavailable {
reason: "node_types is failing, as this test asked".to_owned(),
});
}
}
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(Vec::new());
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(Vec::new());
};
let (nodes, _, _) = visible(tenant, ctx);
Ok(ids
.iter()
.filter_map(|id| nodes.iter().find(|n| n.id == *id && !n.deleted))
.map(|n| (n.id, n.type_id.clone()))
.collect())
}
async fn search(
&self,
ctx: &StoreCtx<'_>,
req: SearchRequest,
vector: Option<VectorArm>,
) -> Result<SearchResponse, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(SearchResponse {
truncated: None,
hits: Vec::new(),
revision: GraphRevision {
source_epoch: self.epoch,
revision: 0,
},
});
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Err(GraphStoreError::NotFound);
};
let (nodes, _, revision) = visible(tenant, ctx);
let admitted: Option<BTreeSet<String>> = if req.type_patterns.is_empty() {
None
} else {
Some(
tenant
.types
.values()
.filter(|record| {
ontology::matches_any_pattern(&record.type_id, &req.type_patterns)
.unwrap_or(false)
})
.map(|record| record.type_id.clone())
.collect(),
)
};
let live: Vec<&FakeNode> = nodes
.iter()
.filter(|n| !n.deleted)
.filter(|n| {
admitted
.as_ref()
.is_none_or(|types| types.contains(&n.type_id))
})
.collect();
let lexical: Vec<&FakeNode> =
if matches!(req.mode, SearchMode::Lexical | SearchMode::Hybrid) {
let needle = req.query.clone().unwrap_or_default().to_lowercase();
live.iter()
.copied()
.filter(|n| needle.is_empty() || n.search_text.to_lowercase().contains(&needle))
.take(req.arm_limit as usize)
.collect()
} else {
Vec::new()
};
let ranked: Vec<&FakeNode> = match &vector {
Some(arm) if matches!(req.mode, SearchMode::Vector | SearchMode::Hybrid) => {
let mut scored: Vec<(f64, &FakeNode)> = live
.iter()
.copied()
.filter(|n| n.embedding_epoch == Some(arm.epoch))
.filter_map(|n| {
n.embedding
.as_ref()
.map(|stored| (cosine_distance(stored, &arm.query_vector), n))
})
.collect();
scored.sort_by(|a, b| {
a.0.partial_cmp(&b.0)
.unwrap_or(std::cmp::Ordering::Equal)
.then_with(|| a.1.id.cmp(&b.1.id))
});
scored
.into_iter()
.map(|(_, n)| n)
.take(req.arm_limit as usize)
.collect()
}
_ => Vec::new(),
};
Ok(SearchResponse {
truncated: None,
hits: fuse_arms(&lexical, &ranked, req.limit as usize),
revision: GraphRevision {
source_epoch: self.epoch,
revision,
},
})
}
async fn project_table(
&self,
ctx: &StoreCtx<'_>,
req: ProjectionRequest,
) -> Result<toolkit_odata::Page<NodeRow>, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(toolkit_odata::Page {
items: Vec::new(),
page_info: empty_page_info(),
});
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Err(GraphStoreError::NotFound);
};
let (nodes, _, revision) = visible(tenant, ctx);
let limit = req.query.limit.unwrap_or(200);
let selected: Vec<&FakeNode> = nodes
.iter()
.filter(|n| !n.deleted)
.filter(|n| {
req.type_set
.as_ref()
.is_none_or(|set| set.contains(&n.type_id))
})
.collect();
let selected: Vec<&FakeNode> = if projection::mentions_payload(&req.query) {
let admitted = req.type_set.as_ref().map(|set| {
let kinds: Vec<BTreeMap<String, ontology::ScalarKind>> = set
.0
.iter()
.map(|type_id| tenant.index_kinds.get(type_id).cloned().unwrap_or_default())
.collect();
projection::admitted_paths(&kinds)
});
let plan = projection::plan(&req.query, admitted.as_ref())
.map_err(|error| GraphStoreError::InvalidQuery { what: error.0 })?;
if req.query.cursor.is_some() {
return Err(GraphStoreError::InvalidQuery {
what: "the in-memory store does not page a payload projection".to_owned(),
});
}
let mut rows: Vec<&FakeNode> = selected
.into_iter()
.filter(|n| plan.filter.as_ref().is_none_or(|p| eval::holds(p, n)))
.collect();
rows.sort_by(|a, b| eval::order(&plan, a, b));
rows
} else {
selected
};
let items = selected
.into_iter()
.take(usize::try_from(limit).unwrap_or(usize::MAX))
.map(|n| NodeRow {
envelope: n.audit.envelope(
n.key.clone(),
ctx.tenant,
GraphRevision {
source_epoch: self.epoch,
revision,
},
),
node_key: n.key.clone(),
type_id: n.type_id.clone(),
name: n.name.clone(),
payload: n.payload.clone(),
})
.collect();
Ok(toolkit_odata::Page {
items,
page_info: toolkit_odata::page::PageInfo {
next_cursor: None,
prev_cursor: None,
limit,
},
})
}
async fn load_topology(
&self,
_ctx: &StoreCtx<'_>,
_req: TopologyRequest,
) -> Result<TopologyPage, GraphStoreError> {
Err(GraphStoreError::Unsupported { what: "topology" })
}
async fn resolve_node_ids(
&self,
ctx: &StoreCtx<'_>,
keys: &[NodeKey],
) -> Result<Vec<(NodeKey, NodeId)>, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(Vec::new());
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(Vec::new());
};
let (nodes, _, _) = visible(tenant, ctx);
Ok(keys
.iter()
.filter_map(|key| {
nodes
.iter()
.find(|n| &n.key == key && !n.deleted)
.map(|n| (key.clone(), n.id))
})
.collect())
}
async fn embedding_state(
&self,
ctx: &StoreCtx<'_>,
keys: &[NodeKey],
) -> Result<Vec<Option<EmbeddingState>>, GraphStoreError> {
if !scope_admits(ctx.scope, ctx.tenant) {
return Ok(keys.iter().map(|_| None).collect());
}
let tenants = self.tenants.lock().map_err(|_| poisoned())?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(keys.iter().map(|_| None).collect());
};
let (nodes, _, _) = visible(tenant, ctx);
Ok(keys
.iter()
.map(|key| {
nodes
.iter()
.find(|n| &n.key == key && !n.deleted)
.map(|n| EmbeddingState {
input_hash: n.embedding_input_hash.clone(),
vector_epoch: if n.embedding.is_some() {
n.embedding_epoch
} else {
None
},
})
})
.collect())
}
}
pub struct FakeGraphEngine {
store: std::sync::Arc<FakeGraphStore>,
}
impl FakeGraphEngine {
pub fn new(store: std::sync::Arc<FakeGraphStore>) -> Self {
Self { store }
}
}
#[async_trait]
impl graph_storage_sdk::plugin_api::GraphEngineV1 for FakeGraphEngine {
fn capabilities(&self) -> graph_storage_sdk::models::EngineCapabilities {
graph_storage_sdk::models::EngineCapabilities::default()
}
async fn cursor(
&self,
ctx: &StoreCtx<'_>,
) -> Result<
graph_storage_sdk::plugin_api::EngineCursor,
graph_storage_sdk::plugin_api::GraphEngineError,
> {
let revision = GraphStoreV1::revision(self.store.as_ref(), ctx)
.await
.map_err(|error| {
graph_storage_sdk::plugin_api::GraphEngineError::Internal(error.to_string())
})?;
Ok(graph_storage_sdk::plugin_api::EngineCursor { revision })
}
async fn expand(
&self,
ctx: &StoreCtx<'_>,
req: graph_storage_sdk::plugin_api::ExpandRequest,
) -> Result<
graph_storage_sdk::plugin_api::ExpandResponse,
graph_storage_sdk::plugin_api::GraphEngineError,
> {
use graph_storage_sdk::models::Direction;
let tenants = self.store.tenants.lock().map_err(|_| {
graph_storage_sdk::plugin_api::GraphEngineError::Internal("poisoned".into())
})?;
let Some(tenant) = tenants.get(&ctx.tenant) else {
return Ok(graph_storage_sdk::plugin_api::ExpandResponse {
reached: Vec::new(),
degrees: Vec::new(),
edges: Vec::new(),
truncated: None,
served_by: graph_storage_sdk::plugin_api::HopBackend::TwoQuery,
});
};
let (nodes, edges, _) = visible(tenant, ctx);
let live: BTreeMap<i64, &FakeNode> = nodes
.iter()
.filter(|n| !n.deleted)
.map(|n| (n.id, n))
.collect();
let frontier: std::collections::BTreeSet<i64> = req.frontier.iter().copied().collect();
let mut reached = Vec::new();
let mut out_edges = Vec::new();
for edge in edges.iter().filter(|e| !e.deleted) {
if let Some(set) = &req.edge_types
&& !set.contains(&edge.type_id)
{
continue;
}
let (Some(src), Some(dst)) = (live.get(&edge.src), live.get(&edge.dst)) else {
continue;
};
let forward = frontier.contains(&edge.src)
&& matches!(req.direction, Direction::Outgoing | Direction::Either);
let backward = frontier.contains(&edge.dst)
&& matches!(req.direction, Direction::Incoming | Direction::Either);
if !forward && !backward {
continue;
}
if forward {
reached.push(edge.dst);
}
if backward {
reached.push(edge.src);
}
out_edges.push(graph_storage_sdk::models::EdgeRef {
edge_key: edge.key.clone(),
edge_type_id: edge.type_id.clone(),
src: src.key.clone(),
dst: dst.key.clone(),
});
}
reached.sort_unstable();
reached.dedup();
let degrees = if req.with_degrees {
let mut incident: BTreeMap<i64, u32> = BTreeMap::new();
for edge in edges.iter().filter(|e| !e.deleted) {
if live.contains_key(&edge.src) && live.contains_key(&edge.dst) {
for id in [edge.src, edge.dst] {
*incident.entry(id).or_default() += 1;
}
}
}
reached
.iter()
.map(|id| incident.get(id).copied().unwrap_or(0))
.collect()
} else {
Vec::new()
};
Ok(graph_storage_sdk::plugin_api::ExpandResponse {
reached,
degrees,
edges: out_edges,
truncated: None,
served_by: graph_storage_sdk::plugin_api::HopBackend::TwoQuery,
})
}
async fn shortest_path(
&self,
_ctx: &StoreCtx<'_>,
_req: graph_storage_sdk::plugin_api::ShortestPathRequest,
) -> Result<
graph_storage_sdk::plugin_api::PathResponse,
graph_storage_sdk::plugin_api::GraphEngineError,
> {
Err(
graph_storage_sdk::plugin_api::GraphEngineError::Unsupported {
what: "shortest_path",
},
)
}
async fn match_pattern(
&self,
_ctx: &StoreCtx<'_>,
_req: graph_storage_sdk::plugin_api::PatternRequest,
) -> Result<
graph_storage_sdk::plugin_api::PatternResponse,
graph_storage_sdk::plugin_api::GraphEngineError,
> {
Err(
graph_storage_sdk::plugin_api::GraphEngineError::Unsupported {
what: "match_pattern",
},
)
}
}
fn poisoned() -> GraphStoreError {
GraphStoreError::Internal("fake store lock is poisoned".into())
}
struct Decided {
outcome: TypeOutcome,
basis: Option<AdmissionBasis>,
change: TypeChange,
revision: i32,
created_at: OffsetDateTime,
rewrites: Vec<Rewrite>,
}
struct Rewrite {
family: ItemFamily,
key: String,
payload: serde_json::Value,
}
fn fresh_change(type_id: &str) -> TypeChange {
TypeChange {
type_id: type_id.to_owned(),
state: TypeChangeState::New,
backward: "compatible".to_owned(),
forward: "compatible".to_owned(),
diagnostics: Vec::new(),
traits_changed: Vec::new(),
rows: None,
rows_rewritten: None,
levels_not_evolvable_in_place: Vec::new(),
migration_required: false,
admissible: true,
}
}
fn decide_existing(
tenant: &Tenant,
existing: &TypeRecord,
descriptor: &ontology::TypeDescriptor,
ancestors: &[(String, serde_json::Value)],
options: &TypeRegistrationOptions,
) -> Result<Decided, GraphStoreError> {
let update = options.on_existing == OnExisting::Update;
let traits_changed =
evolution::traits_diff(&existing.effective_traits, &descriptor.effective_traits);
let unchanged = |traits_changed: Vec<graph_storage_sdk::models::TraitChange>| Decided {
outcome: TypeOutcome::Unchanged,
basis: None,
change: TypeChange {
type_id: descriptor.type_id.clone(),
state: TypeChangeState::Unchanged,
backward: "compatible".to_owned(),
forward: "compatible".to_owned(),
diagnostics: Vec::new(),
traits_changed,
rows: None,
rows_rewritten: None,
levels_not_evolvable_in_place: Vec::new(),
migration_required: false,
admissible: true,
},
revision: existing.revision,
created_at: existing.created_at,
rewrites: Vec::new(),
};
if existing.schema == descriptor.schema {
if options.migration_for(&descriptor.type_id).is_some() {
return Err(GraphStoreError::InvalidQuery {
what: format!(
"the migration for `{}` has nothing to migrate towards: the candidate \
schema is byte-identical to the registered one",
descriptor.type_id
),
});
}
return Ok(unchanged(traits_changed));
}
let comparison = evolution::compare(
&existing.schema,
&descriptor.schema,
ancestors.iter().cloned(),
)
.map_err(|error| validation(0, ItemFamily::Node, &descriptor.type_id, &error.to_string()))?;
let state = comparison.state();
let mut change = TypeChange {
type_id: descriptor.type_id.clone(),
state,
backward: comparison.backward.as_str().to_owned(),
forward: comparison.forward.as_str().to_owned(),
diagnostics: comparison.diagnostics,
traits_changed,
rows: None,
rows_rewritten: None,
levels_not_evolvable_in_place: comparison.levels_not_evolvable_in_place,
migration_required: !matches!(state, TypeChangeState::Compatible),
admissible: false,
};
let accepted = |change: TypeChange, basis: AdmissionBasis, rewrites: Vec<Rewrite>| Decided {
outcome: TypeOutcome::Updated,
basis: Some(basis),
change,
revision: existing.revision.saturating_add(1),
created_at: existing.created_at,
rewrites,
};
let migration = options.migration_for(&descriptor.type_id);
match evolution::decide(
state,
evolution::Asked {
update,
offered: evolution::offered(migration.is_some(), options.revalidate),
},
) {
evolution::Decision::Refuse => {
if options.dry_run {
return Ok(Decided {
outcome: TypeOutcome::Unchanged,
basis: None,
change,
revision: existing.revision,
created_at: existing.created_at,
rewrites: Vec::new(),
});
}
Err(GraphStoreError::Conflict {
reason: if update {
evolution::refusal_reason(&descriptor.type_id, state, &change.diagnostics, 5)
} else {
format!(
"type `{}` is already registered with a different schema",
descriptor.type_id
)
},
})
}
evolution::Decision::Accept => {
change.admissible = true;
Ok(accepted(change, AdmissionBasis::SchemaProved, Vec::new()))
}
evolution::Decision::Migrate => {
let Some(spec) = migration else {
return Err(GraphStoreError::Internal(
"the rule asked for a migration where none was declared".to_owned(),
));
};
let plan = crate::domain::migration::compile(spec).map_err(|error| {
GraphStoreError::InvalidQuery {
what: error.to_string(),
}
})?;
let mut chain: Vec<(String, serde_json::Value)> = ancestors.to_vec();
chain.push((descriptor.type_id.clone(), descriptor.schema.clone()));
let validator =
ontology::ChainValidator::compile(&descriptor.schema, chain).map_err(|error| {
validation(0, ItemFamily::Node, &descriptor.type_id, &error.to_string())
})?;
let (scanned, rewrites, failures) =
migrate_fake(tenant, &descriptor.type_id, &plan, &validator);
change.rows = Some(scanned);
change.rows_rewritten = Some(rewrites.len() as u64);
if failures.is_empty() {
change.admissible = true;
let basis = AdmissionBasis::Migrated {
rows_scanned: scanned,
rows_rewritten: rewrites.len() as u64,
};
return Ok(accepted(
change,
basis,
if options.dry_run {
Vec::new()
} else {
rewrites
},
));
}
if !options.dry_run {
return Err(GraphStoreError::Validation { items: failures });
}
for failure in &failures {
change.diagnostics.push(SchemaDiagnostic {
location: failure.pointer.clone().unwrap_or_default(),
finding: "row_invalid_after_migration".to_owned(),
message: failure.message.clone(),
});
}
Ok(Decided {
outcome: TypeOutcome::Unchanged,
basis: None,
change,
revision: existing.revision,
created_at: existing.created_at,
rewrites: Vec::new(),
})
}
evolution::Decision::Revalidate => {
let mut chain: Vec<(String, serde_json::Value)> = ancestors.to_vec();
chain.push((descriptor.type_id.clone(), descriptor.schema.clone()));
let validator =
ontology::ChainValidator::compile(&descriptor.schema, chain).map_err(|error| {
validation(0, ItemFamily::Node, &descriptor.type_id, &error.to_string())
})?;
let rows = rows_of_type(tenant, &descriptor.type_id);
change.rows = Some(rows);
let failures = revalidate_fake(tenant, &descriptor.type_id, &validator);
if failures.is_empty() {
change.admissible = true;
return Ok(accepted(
change,
AdmissionBasis::DataBacked {
rows_validated: rows,
},
Vec::new(),
));
}
if !options.dry_run {
return Err(GraphStoreError::Validation { items: failures });
}
for failure in &failures {
change.diagnostics.push(SchemaDiagnostic {
location: failure.pointer.clone().unwrap_or_default(),
finding: "stored_row_invalid".to_owned(),
message: failure.message.clone(),
});
}
Ok(Decided {
outcome: TypeOutcome::Unchanged,
basis: None,
change,
revision: existing.revision,
created_at: existing.created_at,
rewrites: Vec::new(),
})
}
}
}
fn rows_of_type(tenant: &Tenant, type_id: &str) -> u64 {
let nodes = tenant
.nodes
.iter()
.filter(|node| !node.deleted && node.type_id == type_id)
.count();
let edges = tenant
.edges
.iter()
.filter(|edge| !edge.deleted && edge.type_id == type_id)
.count();
(nodes + edges) as u64
}
fn migrate_fake(
tenant: &Tenant,
type_id: &str,
plan: &crate::domain::migration::Plan,
validator: &ontology::ChainValidator,
) -> (u64, Vec<Rewrite>, Vec<ItemError>) {
let mut scanned = 0u64;
let mut rewrites = Vec::new();
let mut failures = Vec::new();
for node in tenant
.nodes
.iter()
.filter(|node| !node.deleted && node.type_id == type_id)
{
scanned += 1;
let mut payload = node.payload.clone().unwrap_or(serde_json::Value::Null);
if payload.is_null() {
payload = serde_json::json!({});
}
let changed = plan.apply(&mut payload);
let mut instance = serde_json::json!({ "node_key": node.key, "type": type_id });
if let Some(name) = &node.name {
instance["name"] = serde_json::json!(name);
}
instance["payload"] = payload.clone();
let violations = validator.validate(&instance);
if !violations.is_empty() {
for (pointer, message) in violations {
failures.push(ItemError {
index: usize::try_from(scanned - 1).unwrap_or(usize::MAX),
family: ItemFamily::Node,
gts_type: Some(type_id.to_owned()),
pointer: Some(pointer),
message: format!(
"node `{}` does not satisfy the candidate after the migration: {message}",
node.key
),
});
}
continue;
}
if changed {
rewrites.push(Rewrite {
family: ItemFamily::Node,
key: node.key.clone(),
payload,
});
}
}
for edge in tenant
.edges
.iter()
.filter(|edge| !edge.deleted && edge.type_id == type_id)
{
scanned += 1;
let mut payload = edge.payload.clone().unwrap_or(serde_json::json!({}));
let changed = plan.apply(&mut payload);
let key_of = |id: i64| {
tenant
.nodes
.iter()
.find(|node| node.id == id)
.map_or_else(String::new, |node| node.key.clone())
};
let mut instance = serde_json::json!({
"type": type_id,
"src_node_key": key_of(edge.src),
"dst_node_key": key_of(edge.dst),
});
instance["payload"] = payload.clone();
let violations = validator.validate(&instance);
if !violations.is_empty() {
for (pointer, message) in violations {
failures.push(ItemError {
index: usize::try_from(scanned - 1).unwrap_or(usize::MAX),
family: ItemFamily::Edge,
gts_type: Some(type_id.to_owned()),
pointer: Some(pointer),
message: format!(
"edge `{}` does not satisfy the candidate after the migration: {message}",
edge.key
),
});
}
continue;
}
if changed {
rewrites.push(Rewrite {
family: ItemFamily::Edge,
key: edge.key.clone(),
payload,
});
}
}
(scanned, rewrites, failures)
}
fn revalidate_fake(
tenant: &Tenant,
type_id: &str,
validator: &ontology::ChainValidator,
) -> Vec<ItemError> {
let mut errors = Vec::new();
for (index, node) in tenant
.nodes
.iter()
.filter(|node| !node.deleted && node.type_id == type_id)
.enumerate()
{
let mut instance = serde_json::json!({ "node_key": node.key, "type": type_id });
if let Some(name) = &node.name {
instance["name"] = serde_json::json!(name);
}
if let Some(payload) = &node.payload {
instance["payload"] = payload.clone();
}
for (pointer, message) in validator.validate(&instance) {
errors.push(ItemError {
index,
family: ItemFamily::Node,
gts_type: Some(type_id.to_owned()),
pointer: Some(pointer),
message: format!("node `{}`: {message}", node.key),
});
}
}
for (index, edge) in tenant
.edges
.iter()
.filter(|edge| !edge.deleted && edge.type_id == type_id)
.enumerate()
{
let key_of = |id: i64| {
tenant
.nodes
.iter()
.find(|node| node.id == id)
.map_or_else(String::new, |node| node.key.clone())
};
let mut instance = serde_json::json!({
"type": type_id,
"src_node_key": key_of(edge.src),
"dst_node_key": key_of(edge.dst),
});
if let Some(payload) = &edge.payload {
instance["payload"] = payload.clone();
}
for (pointer, message) in validator.validate(&instance) {
errors.push(ItemError {
index,
family: ItemFamily::Edge,
gts_type: Some(type_id.to_owned()),
pointer: Some(pointer),
message: format!("edge `{}`: {message}", edge.key),
});
}
}
errors
}
fn validation(index: usize, family: ItemFamily, type_id: &str, message: &str) -> GraphStoreError {
GraphStoreError::Validation {
items: vec![ItemError {
index,
family,
gts_type: Some(type_id.to_owned()),
pointer: None,
message: message.to_owned(),
}],
}
}
fn view_of(
node: &FakeNode,
tenant_id: Uuid,
revision: GraphRevision,
adjacency: Vec<AdjacencyEntry>,
truncated: bool,
) -> NodeView {
NodeView {
node_key: node.key.clone(),
type_id: node.type_id.clone(),
name: node.name.clone(),
payload: node.payload.clone(),
has_embedding: node.embedding.is_some(),
labels: Vec::new(),
adjacency,
adjacency_truncated: truncated,
envelope: node.audit.envelope(node.key.clone(), tenant_id, revision),
}
}
fn settle_no_op(
store: &FakeGraphStore,
tenant: &Tenant,
tombstoned: bool,
) -> Result<DeleteOutcome, GraphStoreError> {
if !tombstoned {
return Err(GraphStoreError::NotFound);
}
Ok(DeleteOutcome {
revision: store.revision_of(tenant),
tombstoned_nodes: 0,
tombstoned_edges: 0,
})
}
fn fence(
tenant: &Tenant,
producer: &str,
replace: &graph_storage_sdk::models::ReplaceScope,
request_hash: &str,
) -> Result<(), GraphStoreError> {
let key = (replace.attribute.clone(), replace.value.clone());
let Some(scope) = tenant.scopes.get(&key) else {
return Ok(());
};
if !scope.owner_producer.is_empty() && scope.owner_producer != producer {
return Err(GraphStoreError::Conflict {
reason: format!(
"scope `{}={}` is owned by another producer; a replacement may only be \
submitted by its owner",
replace.attribute, replace.value
),
});
}
if replace.generation < scope.generation {
return Err(GraphStoreError::StaleGeneration {
recorded: scope.generation,
offered: replace.generation,
});
}
if replace.generation == scope.generation && scope.request_hash != request_hash {
return Err(GraphStoreError::Conflict {
reason: "same source generation with different content".into(),
});
}
Ok(())
}
struct BatchState<'a> {
next_id: &'a mut i64,
tally: &'a mut IngestTally,
subject: &'a Subject,
}
struct VectorWrite {
embedding: Option<Vec<f32>>,
epoch: Option<i64>,
input_hash: Option<String>,
}
fn plan_vector(current: Option<&FakeNode>, planned: PlannedVector<'_>) -> VectorWrite {
let stored = current.map(|row| StoredVector {
has_vector: row.embedding.is_some(),
input_hash: row.embedding_input_hash.as_deref(),
});
match decide_vector(stored, planned) {
VectorOutcome::Store {
vector,
epoch,
input_hash,
} => VectorWrite {
embedding: Some(vector),
epoch,
input_hash: Some(input_hash),
},
VectorOutcome::Absent { input_hash } => VectorWrite {
embedding: None,
epoch: None,
input_hash: Some(input_hash),
},
VectorOutcome::Preserve => VectorWrite {
embedding: current.and_then(|row| row.embedding.clone()),
epoch: current.and_then(|row| row.embedding_epoch),
input_hash: current.and_then(|row| row.embedding_input_hash.clone()),
},
VectorOutcome::Stale => VectorWrite {
embedding: current.and_then(|row| row.embedding.clone()),
epoch: None,
input_hash: current.and_then(|row| row.embedding_input_hash.clone()),
},
}
}
fn apply_node(
tenant: &Tenant,
nodes: &mut Vec<FakeNode>,
edges: &[FakeEdge],
state: &mut BatchState<'_>,
index: usize,
spec: &graph_storage_sdk::models::NodeSpec,
vector: VectorWrite,
) -> Result<bool, GraphStoreError> {
let record = tenant.types.get(&spec.type_id).ok_or_else(|| {
validation(
index,
ItemFamily::Node,
&spec.type_id,
"type is not registered",
)
})?;
if record.is_abstract {
return Err(validation(
index,
ItemFamily::Node,
&spec.type_id,
"abstract types cannot be instantiated",
));
}
let Some(existing) = nodes.iter_mut().find(|n| n.key == spec.node_key) else {
if let Some(expected) = spec.expected_version
&& expected != 0
{
return Err(GraphStoreError::Conflict {
reason: format!(
"expected version {expected}, but no node is stored under key `{}`",
spec.node_key
),
});
}
*state.next_id += 1;
nodes.push(FakeNode {
id: *state.next_id,
key: spec.node_key.clone(),
type_id: spec.type_id.clone(),
name: spec.name.clone(),
search_text: crate::infra::store::ingest::compose_search_text(
spec.name.as_deref(),
spec.payload.as_ref(),
&record.effective_traits.full_text_search,
),
payload: spec.payload.clone(),
embedding: vector.embedding,
embedding_epoch: vector.epoch,
embedding_input_hash: vector.input_hash,
version: 1,
deleted: false,
audit: FakeAudit::created(state.subject),
});
return Ok(state.tally.node(&ItemOutcome::Inserted));
};
if existing.deleted {
return Err(GraphStoreError::Conflict {
reason: format!(
"node key `{}` is tombstoned and cannot be re-ingested",
spec.node_key
),
});
}
if let Some(expected) = spec.expected_version
&& expected != existing.version
{
return Err(GraphStoreError::Conflict {
reason: format!(
"expected version {expected}, stored version is {}",
existing.version
),
});
}
let same_type = existing.type_id == spec.type_id;
if !same_type {
let was_phantom = tenant
.types
.get(&existing.type_id)
.and_then(|t| t.effective_traits.family.clone())
.as_deref()
== Some("phantom");
if !was_phantom {
return Err(GraphStoreError::Conflict {
reason: format!(
"node `{}` is already registered under type `{}`; a same-key ingest may \
not change it",
spec.node_key, existing.type_id
),
});
}
}
let vector_unchanged = (
&existing.embedding,
existing.embedding_epoch,
&existing.embedding_input_hash,
) == (&vector.embedding, vector.epoch, &vector.input_hash);
let unchanged = same_type
&& existing.name == spec.name
&& existing.payload == spec.payload
&& vector_unchanged;
if unchanged {
return Ok(state.tally.node(&ItemOutcome::Unchanged));
}
if !same_type {
revalidate_incident_edges(tenant, edges, existing.id, &spec.type_id, index)?;
}
existing.type_id.clone_from(&spec.type_id);
existing.name.clone_from(&spec.name);
existing.search_text = crate::infra::store::ingest::compose_search_text(
spec.name.as_deref(),
spec.payload.as_ref(),
&record.effective_traits.full_text_search,
);
existing.payload.clone_from(&spec.payload);
existing.embedding = vector.embedding;
existing.embedding_epoch = vector.epoch;
existing.embedding_input_hash = vector.input_hash;
existing.version += 1;
existing.audit.updated(state.subject);
Ok(state.tally.node(&if same_type {
ItemOutcome::Updated
} else {
ItemOutcome::Materialized
}))
}
fn revalidate_incident_edges(
tenant: &Tenant,
edges: &[FakeEdge],
node_id: i64,
concrete_type: &str,
index: usize,
) -> Result<(), GraphStoreError> {
for edge in edges.iter().filter(|e| !e.deleted) {
let Some(record) = tenant.types.get(&edge.type_id) else {
continue;
};
for (is_end, patterns, which) in [
(
edge.src == node_id,
&record.effective_traits.src_types,
"source",
),
(
edge.dst == node_id,
&record.effective_traits.dst_types,
"destination",
),
] {
if !is_end || patterns.is_empty() {
continue;
}
if !ontology::matches_any_pattern(concrete_type, patterns).unwrap_or(false) {
return Err(GraphStoreError::Validation {
items: vec![ItemError {
index,
family: ItemFamily::Node,
gts_type: Some(concrete_type.to_owned()),
pointer: Some("/type".to_owned()),
message: format!(
"materializing this node as `{concrete_type}` would leave edge \
`{}` invalid: `{}` does not admit it as a {which} (accepts {})",
edge.key,
edge.type_id,
patterns.join(", ")
),
}],
});
}
}
}
Ok(())
}
fn apply_endpoint(
tenant: &Tenant,
nodes: &mut Vec<FakeNode>,
state: &mut BatchState<'_>,
index: usize,
type_id: &str,
key: &str,
create_phantoms: bool,
) -> Result<i64, GraphStoreError> {
if let Some(existing) = nodes.iter().find(|n| n.key == key && !n.deleted) {
return Ok(existing.id);
}
if nodes.iter().any(|n| n.key == key && n.deleted) {
return Err(GraphStoreError::Conflict {
reason: format!(
"edge[{index}] names endpoint `{key}`, which is tombstoned; the key cannot be \
linked to or re-ingested before purge"
),
});
}
if !create_phantoms {
return Err(validation(
index,
ItemFamily::Edge,
type_id,
"an endpoint does not exist and phantom creation is disabled",
));
}
let phantom_type = tenant
.types
.values()
.find(|t| t.effective_traits.family.as_deref() == Some("phantom"))
.ok_or_else(|| {
validation(
index,
ItemFamily::Edge,
type_id,
"an endpoint does not exist and no phantom type is registered",
)
})?;
*state.next_id += 1;
nodes.push(FakeNode {
id: *state.next_id,
key: key.to_owned(),
type_id: phantom_type.type_id.clone(),
name: None,
search_text: String::new(),
payload: None,
embedding: None,
embedding_epoch: None,
embedding_input_hash: None,
version: 1,
deleted: false,
audit: FakeAudit::created(state.subject),
});
state.tally.phantom_created();
Ok(*state.next_id)
}
struct EdgeDeclaration {
create_phantoms: bool,
scope: Option<(String, String)>,
}
fn apply_edge(
tenant: &Tenant,
nodes: &mut Vec<FakeNode>,
edges: &mut Vec<FakeEdge>,
state: &mut BatchState<'_>,
index: usize,
spec: &graph_storage_sdk::models::EdgeSpec,
declaration: &EdgeDeclaration,
) -> Result<bool, GraphStoreError> {
let EdgeDeclaration {
create_phantoms,
scope,
} = declaration;
let create_phantoms = *create_phantoms;
let scope = scope.clone();
let record = tenant.types.get(&spec.type_id).cloned().ok_or_else(|| {
validation(
index,
ItemFamily::Edge,
&spec.type_id,
"type is not registered",
)
})?;
let before = state.tally.counts.phantoms_created;
let src = apply_endpoint(
tenant,
nodes,
state,
index,
&spec.type_id,
&spec.src_node_key,
create_phantoms,
)?;
let dst = apply_endpoint(
tenant,
nodes,
state,
index,
&spec.type_id,
&spec.dst_node_key,
create_phantoms,
)?;
let mut changed = state.tally.counts.phantoms_created > before;
for (id, patterns, key, pointer) in [
(
src,
&record.effective_traits.src_types,
&spec.src_node_key,
"/src_node_key",
),
(
dst,
&record.effective_traits.dst_types,
&spec.dst_node_key,
"/dst_node_key",
),
] {
let Some(endpoint) = nodes.iter().find(|n| n.id == id) else {
continue;
};
let is_phantom = tenant
.types
.get(&endpoint.type_id)
.and_then(|t| t.effective_traits.family.as_deref())
== Some("phantom");
if is_phantom || patterns.is_empty() {
continue;
}
let admitted = ontology::matches_any_pattern(&endpoint.type_id, patterns).unwrap_or(false);
if !admitted {
return Err(GraphStoreError::Validation {
items: vec![ItemError {
index,
family: ItemFamily::Edge,
gts_type: Some(spec.type_id.clone()),
pointer: Some(pointer.to_owned()),
message: format!(
"endpoint `{key}` is a `{}`, which `{}` does not admit; this edge \
type accepts {}",
endpoint.type_id,
spec.type_id,
patterns.join(", ")
),
}],
});
}
}
let edge_key = identity::derive_edge_key(record.type_uuid, spec);
if let Some(owner) = edges
.iter()
.find(|e| e.key == edge_key)
.and_then(|e| e.scope.as_ref())
&& let Some(declaring) = scope.as_ref()
&& owner != declaring
{
return Err(GraphStoreError::Conflict {
reason: format!(
"edge `{edge_key}` was declared by scope `{}={}` and may not be \
re-declared under `{}={}`; a move between scopes is a \
deletion and a re-declaration, not a write",
owner.0, owner.1, declaring.0, declaring.1
),
});
}
let outcome = match edges.iter_mut().find(|e| e.key == edge_key) {
Some(existing) if existing.payload == spec.payload && !existing.deleted => {
if existing.scope.is_none() && scope.is_some() {
existing.scope.clone_from(&scope);
}
ItemOutcome::Unchanged
}
Some(existing) => {
existing.payload.clone_from(&spec.payload);
existing.deleted = false;
if scope.is_some() {
existing.scope.clone_from(&scope);
}
existing.audit.updated(state.subject);
ItemOutcome::Updated
}
None => {
edges.push(FakeEdge {
key: edge_key,
type_id: spec.type_id.clone(),
src,
dst,
discriminator: spec.discriminator.clone(),
payload: spec.payload.clone(),
deleted: false,
scope,
audit: FakeAudit::created(state.subject),
});
ItemOutcome::Inserted
}
};
changed |= state.tally.edge(&outcome);
Ok(changed)
}
struct ScopeReplacement<'a> {
nodes: &'a mut Vec<FakeNode>,
edges: &'a mut Vec<FakeEdge>,
managed_node_types: &'a BTreeSet<String>,
static_edge_types: &'a BTreeSet<String>,
types: &'a BTreeMap<String, TypeRecord>,
}
fn replace_scope(
working: ScopeReplacement<'_>,
req: &IngestRequest,
replace: &graph_storage_sdk::models::ReplaceScope,
tally: &mut IngestTally,
) -> bool {
let ScopeReplacement {
nodes,
edges,
managed_node_types,
static_edge_types,
types,
} = working;
let written: BTreeSet<&str> = req
.nodes
.iter()
.map(|spec| spec.node_key.as_str())
.collect();
let declared: BTreeSet<String> = req
.edges
.iter()
.filter_map(|spec| {
let record = types.get(&spec.type_id)?;
Some(identity::derive_edge_key(record.type_uuid, spec))
})
.collect();
let owner = (replace.attribute.clone(), replace.value.clone());
let dropped: Vec<i64> = nodes
.iter()
.filter(|node| {
!node.deleted
&& managed_node_types.contains(&node.type_id)
&& !written.contains(node.key.as_str())
&& node
.payload
.as_ref()
.and_then(|payload| payload.get(&replace.attribute))
.and_then(serde_json::Value::as_str)
== Some(replace.value.as_str())
})
.map(|node| node.id)
.collect();
let before = edges.len();
edges.retain(|edge| {
if !static_edge_types.contains(&edge.type_id) {
return true;
}
let abandoned = edge.scope.as_ref() == Some(&owner) && !declared.contains(&edge.key);
let orphaned = dropped.contains(&edge.src) || dropped.contains(&edge.dst);
!(abandoned || orphaned)
});
tally.counts.scope_removed_edges = (before - edges.len()) as u64;
if dropped.is_empty() {
return tally.counts.scope_removed_edges > 0;
}
edges.retain(|edge| {
!(edge.deleted && (dropped.contains(&edge.src) || dropped.contains(&edge.dst)))
});
let referenced: BTreeSet<i64> = edges.iter().flat_map(|edge| [edge.src, edge.dst]).collect();
let before = nodes.len();
nodes.retain(|node| !dropped.contains(&node.id) || referenced.contains(&node.id));
tally.counts.scope_removed_nodes = (before - nodes.len()) as u64;
tally.counts.scope_removed_nodes > 0 || tally.counts.scope_removed_edges > 0
}
fn catalogue_page(tenant: &Tenant, query: &TypeQuery) -> (Vec<TypeRecord>, Option<String>) {
let limit = query.top.map_or(usize::MAX, |top| top as usize);
let mut items: Vec<TypeRecord> = Vec::new();
let mut next_cursor = None;
let mut walked = tenant
.types
.values()
.filter(|record| {
query
.cursor
.as_ref()
.is_none_or(|cursor| &record.type_id > cursor)
})
.filter(|record| query.kind.is_none_or(|kind| kind == record.kind))
.peekable();
while let Some(record) = walked.next() {
let admitted = query.pattern.as_ref().is_none_or(|pattern| {
ontology::matches_any_pattern(&record.type_id, std::slice::from_ref(pattern))
.unwrap_or(false)
});
if admitted {
items.push(record.clone());
if items.len() >= limit {
next_cursor = walked.peek().map(|_| record.type_id.clone());
break;
}
}
}
(items, next_cursor)
}
fn empty_page_info() -> toolkit_odata::page::PageInfo {
toolkit_odata::page::PageInfo {
next_cursor: None,
prev_cursor: None,
limit: 0,
}
}
fn cosine_distance(one: &[f32], other: &[f32]) -> f64 {
let dot: f64 = one
.iter()
.zip(other)
.map(|(a, b)| f64::from(*a) * f64::from(*b))
.sum();
let norm = |v: &[f32]| -> f64 {
v.iter()
.map(|x| f64::from(*x) * f64::from(*x))
.sum::<f64>()
.sqrt()
};
let (left, right) = (norm(one), norm(other));
if left == 0.0 || right == 0.0 {
return 1.0;
}
1.0 - dot / (left * right)
}
fn fuse_arms(
lexical: &[&FakeNode],
vector: &[&FakeNode],
limit: usize,
) -> Vec<graph_storage_sdk::models::SearchHit> {
use graph_storage_sdk::models::{ArmHit, SearchArm, SearchHit};
const K: f64 = 60.0;
let mut order: Vec<i64> = Vec::new();
let mut fused: BTreeMap<i64, (f64, Vec<ArmHit>, &FakeNode)> = BTreeMap::new();
for (arm, rows) in [(SearchArm::Lexical, lexical), (SearchArm::Vector, vector)] {
for (position, node) in rows.iter().enumerate() {
let rank = u32::try_from(position)
.unwrap_or(u32::MAX)
.saturating_add(1);
let contribution = 1.0 / (K + f64::from(rank));
let entry = fused
.entry(node.id)
.or_insert_with(|| (0.0, Vec::new(), node));
if entry.1.is_empty() {
order.push(node.id);
}
entry.0 += contribution;
entry.1.push(ArmHit {
arm,
rank,
score: contribution,
});
}
}
let mut hits: Vec<SearchHit> = order
.into_iter()
.filter_map(|id| fused.get(&id))
.map(|(score, arms, node)| SearchHit {
node_key: node.key.clone(),
type_id: node.type_id.clone(),
name: node.name.clone(),
score: *score,
arms: arms.clone(),
snippet: None,
})
.collect();
hits.sort_by(|a, b| {
b.score
.partial_cmp(&a.score)
.unwrap_or(std::cmp::Ordering::Equal)
});
hits.truncate(limit);
hits
}
mod eval {
use std::cmp::Ordering;
use toolkit_odata::SortDir;
use super::FakeNode;
use crate::domain::ontology::ScalarKind;
use crate::domain::projection::{CmpOp, FieldRef, Plan, Predicate, Scalar, TextOp};
#[derive(Clone, Debug, PartialEq)]
enum Cell {
Null,
Text(String),
Num(f64),
Bool(bool),
}
fn cell(field: &FieldRef, node: &FakeNode) -> Cell {
match field {
FieldRef::NodeKey => Cell::Text(node.key.clone()),
FieldRef::Name => Cell::Text(node.name.clone().unwrap_or_default()),
FieldRef::CreatedAt => rfc3339(node.audit.created_at),
FieldRef::UpdatedAt => rfc3339(node.audit.updated_at),
FieldRef::Payload { pointer, kind } => {
let inner = pointer.strip_prefix("/payload").unwrap_or(pointer);
let Some(value) = node.payload.as_ref().and_then(|p| p.pointer(inner)) else {
return Cell::Null;
};
match (kind, value) {
(ScalarKind::String | ScalarKind::DateTime, serde_json::Value::String(s)) => {
Cell::Text(s.clone())
}
(ScalarKind::Number | ScalarKind::Integer, serde_json::Value::Number(n)) => {
n.as_f64().map_or(Cell::Null, Cell::Num)
}
(ScalarKind::Boolean, serde_json::Value::Bool(b)) => Cell::Bool(*b),
_ => Cell::Null,
}
}
}
}
fn rfc3339(at: time::OffsetDateTime) -> Cell {
at.format(&time::format_description::well_known::Rfc3339)
.map_or(Cell::Null, Cell::Text)
}
fn literal(value: &Scalar) -> Cell {
match value {
Scalar::Str(s) | Scalar::DateTime(s) => Cell::Text(s.clone()),
Scalar::Num(n) => n.parse::<f64>().map_or(Cell::Null, Cell::Num),
Scalar::Bool(b) => Cell::Bool(*b),
}
}
fn compare(a: &Cell, b: &Cell) -> Option<Ordering> {
match (a, b) {
(Cell::Text(x), Cell::Text(y)) => Some(x.cmp(y)),
(Cell::Num(x), Cell::Num(y)) => x.partial_cmp(y),
(Cell::Bool(x), Cell::Bool(y)) => Some(x.cmp(y)),
_ => None,
}
}
pub(super) fn holds(predicate: &Predicate, node: &FakeNode) -> bool {
match predicate {
Predicate::Compare { field, op, value } => {
let Some(ordering) = compare(&cell(field, node), &literal(value)) else {
return false;
};
match op {
CmpOp::Eq => ordering == Ordering::Equal,
CmpOp::Ne => ordering != Ordering::Equal,
CmpOp::Gt => ordering == Ordering::Greater,
CmpOp::Ge => ordering != Ordering::Less,
CmpOp::Lt => ordering == Ordering::Less,
CmpOp::Le => ordering != Ordering::Greater,
}
}
Predicate::In { field, values } => {
let actual = cell(field, node);
values
.iter()
.any(|v| compare(&actual, &literal(v)) == Some(Ordering::Equal))
}
Predicate::Text { field, op, needle } => match cell(field, node) {
Cell::Text(text) => match op {
TextOp::Contains => text.contains(needle.as_str()),
TextOp::StartsWith => text.starts_with(needle.as_str()),
TextOp::EndsWith => text.ends_with(needle.as_str()),
},
_ => false,
},
Predicate::And(children) => children.iter().all(|c| holds(c, node)),
Predicate::Or(children) => children.iter().any(|c| holds(c, node)),
Predicate::Not(inner) => !holds(inner, node),
}
}
pub(super) fn order(plan: &Plan, a: &FakeNode, b: &FakeNode) -> Ordering {
for term in &plan.order {
let (x, y) = (cell(&term.field, a), cell(&term.field, b));
let ordering = match (&x, &y) {
(Cell::Null, Cell::Null) => Ordering::Equal,
(Cell::Null, _) => Ordering::Greater,
(_, Cell::Null) => Ordering::Less,
_ => {
let natural = compare(&x, &y).unwrap_or(Ordering::Equal);
match term.dir {
SortDir::Asc => natural,
SortDir::Desc => natural.reverse(),
}
}
};
if ordering != Ordering::Equal {
return ordering;
}
}
Ordering::Equal
}
}