use std::sync::Arc;
use surrealdb_datastore::values::inline_cache::CacheValue;
use crate::catalog::providers::TableProvider;
use crate::exec::operators::scan::EdgeTableSpec;
use crate::exec::operators::scan::graph_keys::compute_graph_ranges;
use crate::exec::parts::LookupDirection;
use crate::exec::physical_expr::{BoxFut, EvalContext, PhysicalExpr};
use crate::exec::{AccessMode, ContextLevel, FlowResult};
use crate::expr::Dir;
use crate::idx::adjacency::{VertexAdjacency, scope_degree, vertex_adjacency_of};
use crate::key::schema::EdgeCacheKey;
use crate::val::{RecordId, TableName, Value};
pub(crate) struct GraphDegreeExpr {
direction: LookupDirection,
edge_tables: Vec<TableName>,
fallback: Arc<dyn PhysicalExpr>,
prefixes: Vec<Vec<u8>>,
specs: Vec<EdgeTableSpec>,
table_info:
std::sync::Mutex<std::collections::HashMap<TableName, (VertexAdjacency, Option<u32>)>>,
}
impl std::fmt::Debug for GraphDegreeExpr {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("GraphDegreeExpr")
.field("direction", &self.direction)
.field("edge_tables", &self.edge_tables)
.finish_non_exhaustive()
}
}
impl surrealdb_types::ToSql for GraphDegreeExpr {
fn fmt_sql(&self, f: &mut String, fmt: surrealdb_types::SqlFormat) {
self.fallback.fmt_sql(f, fmt);
}
}
impl PhysicalExpr for GraphDegreeExpr {
fn name(&self) -> &'static str {
"GraphDegree"
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Database
}
fn access_mode(&self) -> AccessMode {
AccessMode::ReadOnly
}
fn evaluate<'a>(&'a self, ctx: EvalContext<'a>) -> BoxFut<'a, FlowResult<Value>> {
Box::pin(async move {
let rid = match ctx.current_value {
Some(Value::RecordId(rid)) => rid.clone(),
Some(Value::Object(obj)) => match obj.get("id") {
Some(Value::RecordId(rid)) => rid.clone(),
_ => return self.fallback.evaluate(ctx).await,
},
_ => return self.fallback.evaluate(ctx).await,
};
match self.degree(&ctx, &rid).await? {
Some(degree) => Ok(Value::Number(degree.into())),
None => self.fallback.evaluate(ctx).await,
}
})
}
}
struct DegreeCancel<'a> {
exec: &'a crate::exec::ExecutionContext,
}
impl crate::catalog::providers::CancellationProbe for DegreeCancel<'_> {
fn expect_not_timedout(
&self,
) -> crate::catalog::providers::BoxProviderFut<'_, anyhow::Result<()>> {
Box::pin(async move {
if self.exec.cancellation().is_cancelled() {
return Err(anyhow::anyhow!(crate::err::EngineError::QueryCancelled));
}
self.exec.ctx().expect_not_timedout().await
})
}
}
impl GraphDegreeExpr {
pub(crate) fn new(
direction: LookupDirection,
edge_tables: Vec<TableName>,
fallback: Arc<dyn PhysicalExpr>,
) -> anyhow::Result<Self> {
let prefixes = edge_tables
.iter()
.map(|t| storekey::encode_vec(t).map_err(anyhow::Error::from_boxed))
.collect::<anyhow::Result<_>>()?;
let specs = edge_tables
.iter()
.map(|t| EdgeTableSpec {
table: t.clone(),
range_start: std::ops::Bound::Unbounded,
range_end: std::ops::Bound::Unbounded,
})
.collect();
Ok(Self {
direction,
edge_tables,
fallback,
prefixes,
specs,
table_info: Default::default(),
})
}
async fn degree(&self, ctx: &EvalContext<'_>, rid: &RecordId) -> anyhow::Result<Option<u64>> {
let exec = ctx.exec_ctx;
let db_ctx = exec.database()?;
let ns_id = db_ctx.ns_ctx.ns.namespace_id;
let db_id = db_ctx.db.database_id;
let txn = exec.txn();
let info = self.table_info.lock().expect("degree memo poisoned").get(&rid.table).copied();
let (adjacency, edges_cap) = match info {
Some(info) => info,
None => {
let tb = txn.get_tb(ns_id, db_id, &rid.table, None).await?;
let info = (
vertex_adjacency_of(tb.as_deref()),
crate::idx::inline_cache::effective_edges_cap(tb.as_deref()),
);
self.table_info
.lock()
.expect("degree memo poisoned")
.insert(rid.table.clone(), info);
info
}
};
let VertexAdjacency::Normal {
folded,
} = adjacency
else {
return Ok(None);
};
let dirs: &[Dir] = match self.direction {
LookupDirection::Out => &[Dir::Out],
LookupDirection::In => &[Dir::In],
LookupDirection::Both => &[Dir::In, Dir::Out],
LookupDirection::Reference => return Ok(None),
};
let cancel = DegreeCancel {
exec,
};
let mut degree = 0u64;
for &dir in dirs {
if edges_cap.is_some_and(|cap| cap > 0) {
let key = EdgeCacheKey {
ns: ns_id,
db: db_id,
tb: std::borrow::Cow::Borrowed(&rid.table),
id: std::borrow::Cow::Borrowed(&rid.key),
dir,
};
if let Some(CacheValue::Live(entries)) = txn.get_key(&key, None).await? {
degree += entries
.iter()
.filter(|e| {
self.prefixes.is_empty()
|| self.prefixes.iter().any(|p| e.edge.starts_with(p))
})
.count() as u64;
continue;
}
}
let ranges = compute_graph_ranges(ns_id, db_id, rid, dir, &self.specs, exec)
.await
.map_err(|cf| match cf {
crate::expr::ControlFlow::Err(e) => e,
_ => anyhow::anyhow!("unexpected control flow while computing degree ranges"),
})?;
let resolve = exec.root().ctx.get_index_stores().adjacency_resolve();
let mut dir_degree = 0u64;
for range in ranges {
dir_degree += scope_degree(
&txn,
range.as_adjacency_scope(ns_id, db_id, rid, dir),
folded,
None,
Some(resolve),
Some(&cancel),
)
.await?;
}
if !folded {
crate::exec::operators::scan::graph_keys::observe_fold_candidate(
&txn,
ns_id,
db_id,
rid,
dir,
dir_degree,
exec.root().ctx.config.idx.graph_fold_threshold,
);
}
degree += dir_degree;
}
Ok(Some(degree))
}
}