use std::borrow::Cow;
use std::sync::Arc;
use common::future::stream::{self, Yielder};
use futures::StreamExt;
use surrealdb_datastore::values::inline_cache::{CacheValue, EdgeCacheEntry};
use super::common::{extract_record_ids_into, resolve_record_batch, resolve_version_stamp};
pub use super::graph_keys::EdgeTableSpec;
use super::graph_keys::{compute_graph_ranges, observe_fold_candidate};
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::exec::parts::LookupDirection;
use crate::exec::permission::{
PhysicalPermission, should_check_perms, validate_record_user_access,
};
use crate::exec::{
AccessMode, ContextLevel, ControlFlowExt, EvalContext, ExecOperator, ExecutionContext,
FlowResult, OperatorMetrics, PhysicalExpr, ValueBatch, ValueBatchStream, buffer_stream,
monitor_stream,
};
use crate::expr::{ControlFlow, Dir};
use crate::iam::Action;
use crate::idx::adjacency::{
AdjacencyScope, MergedAdjacencyCursor, VertexAdjacency, VisitFlow, vertex_adjacency_of,
};
use crate::key::Resumable;
use crate::key::schema::EdgeCacheKey;
use crate::kvs::{CachePolicy, Direction, Transaction};
use crate::val::{RecordId, RecordIdKey, TableName, Value};
type SourceTableInfo =
std::sync::Mutex<std::collections::HashMap<TableName, (VertexAdjacency, Option<u32>)>>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PayloadPredicateSpec {
pub(crate) tables: std::collections::HashMap<TableName, PayloadTableSpec>,
pub(crate) needs_id: bool,
pub(crate) needs_in: bool,
pub(crate) needs_out: bool,
pub(crate) key_only: bool,
pub(crate) matcher: Option<Arc<super::key_matcher::KeyPredicateMatcher>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PayloadTableSpec {
pub(crate) generation: u32,
pub(crate) fields: Vec<String>,
pub(crate) referenced: Vec<bool>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum GraphScanOutput {
#[default]
TargetId,
FullEdge,
TargetVertex,
}
#[derive(Debug, Clone)]
pub struct GraphEdgeScan {
pub(crate) input: Arc<dyn ExecOperator>,
pub(crate) direction: LookupDirection,
pub(crate) edge_tables: Vec<EdgeTableSpec>,
pub(crate) output_mode: GraphScanOutput,
pub(crate) target_tables: Vec<TableName>,
pub(crate) version: Option<Arc<dyn PhysicalExpr>>,
pub(crate) limit: Option<usize>,
pub(crate) limit_needs_id_order: bool,
pub(crate) predicate: Option<Arc<dyn PhysicalExpr>>,
pub(crate) predicate_key_resident: bool,
pub(crate) payload_spec: Option<Arc<PayloadPredicateSpec>>,
pub(crate) metrics: Arc<OperatorMetrics>,
pub(crate) table_info: Arc<SourceTableInfo>,
}
impl GraphEdgeScan {
pub(crate) fn new(
input: Arc<dyn ExecOperator>,
direction: LookupDirection,
edge_tables: Vec<EdgeTableSpec>,
output_mode: GraphScanOutput,
version: Option<Arc<dyn PhysicalExpr>>,
) -> Self {
Self {
input,
direction,
edge_tables,
output_mode,
target_tables: Vec::new(),
version,
limit: None,
limit_needs_id_order: false,
predicate: None,
predicate_key_resident: false,
payload_spec: None,
metrics: Arc::new(OperatorMetrics::new()),
table_info: Arc::new(SourceTableInfo::default()),
}
}
pub(crate) fn with_limit(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self
}
pub(crate) fn with_ordered_limit(mut self, limit: usize) -> Self {
self.limit = Some(limit);
self.limit_needs_id_order = true;
self
}
pub(crate) fn with_target_tables(mut self, tables: Vec<TableName>) -> Self {
self.target_tables = tables;
self
}
pub(crate) fn with_predicate(
mut self,
predicate: Arc<dyn PhysicalExpr>,
key_resident: bool,
) -> Self {
self.predicate = Some(predicate);
self.predicate_key_resident = key_resident;
self
}
pub(crate) fn with_payload_predicate(
mut self,
predicate: Arc<dyn PhysicalExpr>,
spec: PayloadPredicateSpec,
) -> Self {
self.predicate = Some(predicate);
self.predicate_key_resident = false;
self.payload_spec = Some(Arc::new(spec));
self
}
}
impl ExecOperator for GraphEdgeScan {
fn name(&self) -> &'static str {
"GraphEdgeScan"
}
fn attrs(&self) -> Vec<(String, String)> {
let dir = match self.direction {
LookupDirection::Out => "->",
LookupDirection::In => "<-",
LookupDirection::Both => "<->",
LookupDirection::Reference => "<~",
};
let tables = if self.edge_tables.is_empty() {
"*".to_string()
} else {
self.edge_tables.iter().map(|t| t.table.as_str()).collect::<Vec<_>>().join(", ")
};
let mut attrs = vec![
("direction".to_string(), dir.to_string()),
("tables".to_string(), tables),
("output".to_string(), format!("{:?}", self.output_mode)),
];
if let Some(ref version) = self.version {
attrs.push(("version".to_string(), version.to_sql()));
}
if let Some(limit) = self.limit {
attrs.push(("limit".to_string(), limit.to_string()));
if self.limit_needs_id_order {
attrs.push(("order_pushdown".to_string(), "id".to_string()));
}
}
if let Some(ref predicate) = self.predicate {
attrs.push(("predicate".to_string(), predicate.to_sql()));
attrs.push((
"predicate_scope".to_string(),
if self.predicate_key_resident {
"key".to_string()
} else if let Some(spec) = self.payload_spec.as_deref() {
if spec.key_only {
"key".to_string()
} else {
"payload".to_string()
}
} else {
"record".to_string()
},
));
}
attrs
}
fn required_context(&self) -> ContextLevel {
let mut level = self.input.required_context().max(ContextLevel::Database);
if let Some(ref predicate) = self.predicate {
level = level.max(predicate.required_context());
}
level
}
fn access_mode(&self) -> AccessMode {
let mut mode = self.input.access_mode();
if let Some(ref version) = self.version {
mode = mode.combine(version.access_mode());
}
if let Some(ref predicate) = self.predicate {
mode = mode.combine(predicate.access_mode());
}
mode
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
vec![&self.input]
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let db_ctx = ctx.database()?.clone();
validate_record_user_access(&db_ctx)?;
let check_perms = should_check_perms(&db_ctx, Action::View)?;
let mut input_stream = buffer_stream(
self.input.execute(ctx)?,
self.input.access_mode(),
self.input.cardinality_hint(),
ctx.root().ctx.config.exec.operator_buffer_size,
);
let direction = self.direction;
let edge_tables = Arc::new(self.edge_tables.clone());
let output_mode = self.output_mode;
let target_tables = Arc::new(self.target_tables.clone());
let edge_limit = self.limit;
let limit_needs_id_order = self.limit_needs_id_order;
let version_expr = self.version.clone();
let scan_batch_size = ctx.root().ctx.config.exec.scan_batch_size;
let fold_threshold = ctx.root().ctx.config.idx.graph_fold_threshold;
let ctx = ctx.clone();
let fetch_full = output_mode == GraphScanOutput::FullEdge;
let has_predicate = self.predicate.is_some();
let key_predicate: Option<Arc<dyn PhysicalExpr>> =
self.predicate.clone().filter(|_| self.predicate_key_resident);
let record_predicate: Option<Arc<dyn PhysicalExpr>> =
self.predicate.clone().filter(|_| !self.predicate_key_resident);
let payload_spec = self.payload_spec.clone();
let metrics = Arc::clone(&self.metrics);
let record_metrics = metrics.is_enabled();
let table_info = Arc::clone(&self.table_info);
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
let txn = ctx.txn();
let ns_id = db_ctx.ns_ctx.ns.namespace_id;
let db_id = db_ctx.db.database_id;
let mut perm_cache = PermCache::new();
let version: Option<u64> = resolve_version_stamp(&ctx, version_expr.as_ref()).await?;
let directions: Arc<Vec<Dir>> = Arc::new(match direction {
LookupDirection::Out => vec![Dir::Out],
LookupDirection::In => vec![Dir::In],
LookupDirection::Both => vec![Dir::In, Dir::Out],
LookupDirection::Reference => Err(ControlFlow::Err(anyhow::anyhow!(
"Reference lookups should use ReferenceScan, not GraphEdgeScan"
)))?,
});
let cache_eligible = payload_spec.as_deref().is_none_or(|s| s.key_only)
&& version.is_none()
&& edge_tables.iter().all(|s| {
matches!(s.range_start, std::ops::Bound::Unbounded)
&& matches!(s.range_end, std::ops::Bound::Unbounded)
});
let spec_prefixes: Vec<Vec<u8>> = if cache_eligible {
edge_tables
.iter()
.map(|spec| {
storekey::encode_vec(&spec.table).map_err(anyhow::Error::from_boxed)
})
.collect::<Result<_, anyhow::Error>>()?
} else {
Vec::new()
};
let mut rid_batch: Vec<RecordId> = Vec::with_capacity(scan_batch_size);
#[cfg(not(target_family = "wasm"))]
let pipelined_resolve = check_perms || fetch_full;
#[cfg(not(target_family = "wasm"))]
let mut inflight: Option<InflightResolve> = None;
#[cfg(not(target_family = "wasm"))]
let mut spare: Vec<RecordId> = Vec::new();
let mut payload_pending: Vec<PayloadPending> = Vec::new();
while let Some(batch_result) = input_stream.next().await {
let batch = batch_result?;
let source_rids: Vec<RecordId> = batch
.into_iter()
.flat_map(|v| {
let mut rids = Vec::new();
extract_record_ids_into(v, &mut rids);
rids
})
.collect();
let mut batch_table_info: std::collections::HashMap<
TableName,
(VertexAdjacency, Option<u32>),
> = std::collections::HashMap::new();
{
let memo = table_info.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
for rid in &source_rids {
if let Some(info) = memo.get(&rid.table) {
batch_table_info.insert(rid.table.clone(), *info);
}
}
}
let mut resolved: Vec<(TableName, (VertexAdjacency, Option<u32>))> = Vec::new();
for rid in &source_rids {
if batch_table_info.contains_key(&rid.table) {
continue;
}
let tb = txn
.get_tb(ns_id, db_id, &rid.table, None)
.await
.context("Failed to resolve the source vertex's table")?;
let info = (
vertex_adjacency_of(tb.as_deref()),
crate::idx::inline_cache::effective_edges_cap(tb.as_deref()),
);
batch_table_info.insert(rid.table.clone(), info);
resolved.push((rid.table.clone(), info));
}
if !resolved.is_empty() {
table_info
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.extend(resolved);
}
let mut cache_map: std::collections::HashMap<(usize, Dir), Vec<EdgeCacheEntry>> =
std::collections::HashMap::new();
if cache_eligible {
let mut slots: Vec<(usize, Dir)> = Vec::new();
let mut keys: Vec<EdgeCacheKey> = Vec::new();
for (src_idx, rid) in source_rids.iter().enumerate() {
let (_, edges_cap) = batch_table_info[&rid.table];
if edges_cap.is_none() {
continue;
}
for &dir in directions.iter() {
slots.push((src_idx, dir));
keys.push(EdgeCacheKey {
ns: ns_id,
db: db_id,
tb: Cow::Borrowed(&rid.table),
id: Cow::Borrowed(&rid.key),
dir,
});
}
}
if !keys.is_empty() {
let values = txn
.get_many_key(keys, None)
.await
.context("Failed to read the inline adjacency caches")?;
let mut hits: u64 = 0;
let mut misses: u64 = 0;
for (slot, value) in slots.into_iter().zip(values) {
match value {
Some(CacheValue::Live(entries)) => {
hits += 1;
cache_map.insert(slot, entries);
}
_ => misses += 1,
}
}
if record_metrics {
metrics.add_cache_hits(hits);
metrics.add_cache_misses(misses);
}
}
}
let mut prescanned: Vec<Option<Vec<RecordId>>> = Vec::new();
prescanned.resize_with(source_rids.len(), || None);
#[cfg(not(target_family = "wasm"))]
let plain_shape = !has_predicate && payload_spec.is_none();
#[cfg(not(target_family = "wasm"))]
let mut fanout_engaged: Option<bool> = None;
let mut win_start = 0usize;
while win_start < source_rids.len() {
let win_end = (win_start + GRAPH_FANOUT_CONCURRENCY).min(source_rids.len());
#[cfg(not(target_family = "wasm"))]
if plain_shape && fanout_engaged == Some(true) {
let missed: Vec<usize> = (win_start..win_end)
.filter(|src_idx| {
!directions
.iter()
.any(|dir| cache_map.contains_key(&(*src_idx, *dir)))
})
.collect();
if missed.len() > 1 {
let mut pending: std::collections::VecDeque<_> = missed
.into_iter()
.map(|src_idx| {
let ctx = ctx.clone();
let txn = Arc::clone(&txn);
let rid = source_rids[src_idx].clone();
let (adjacency, _) = batch_table_info[&rid.table];
let directions = Arc::clone(&directions);
let edge_tables = Arc::clone(&edge_tables);
let target_tables = Arc::clone(&target_tables);
let metrics = Arc::clone(&metrics);
async move {
let out = scan_source_plain(
&ctx,
&txn,
adjacency,
ns_id,
db_id,
&rid,
&directions,
&edge_tables,
output_mode,
&target_tables,
version,
edge_limit,
fold_threshold,
scan_batch_size,
&metrics,
record_metrics,
)
.await;
(src_idx, out)
}
})
.collect();
let mut set = tokio::task::JoinSet::new();
let mut first_err: Option<(usize, ControlFlow)> = None;
let mut join_err: Option<ControlFlow> = None;
loop {
if first_err.is_none() && join_err.is_none() {
while set.len() < GRAPH_FANOUT_CONCURRENCY {
let Some(task) = pending.pop_front() else {
break;
};
set.spawn(task);
}
}
let Some(joined) = set.join_next().await else {
break;
};
match joined {
Ok((src_idx, Ok(out))) => prescanned[src_idx] = out,
Ok((src_idx, Err(err))) => {
if first_err
.as_ref()
.is_none_or(|(lowest, _)| src_idx < *lowest)
{
first_err = Some((src_idx, err));
}
}
Err(e) => {
join_err.get_or_insert_with(|| {
ControlFlow::Err(anyhow::anyhow!(
"graph fan-out task failed: {e}"
))
});
}
}
}
if let Some((_, err)) = first_err {
return Err(err);
}
if let Some(err) = join_err {
return Err(err);
}
}
}
#[cfg(not(target_family = "wasm"))]
let win_timer = (plain_shape && fanout_engaged.is_none()).then(web_time::Instant::now);
for (offset, rid) in source_rids[win_start..win_end].iter().enumerate() {
let src_idx = win_start + offset;
if let Some(found) = prescanned[src_idx].take() {
for target in found {
rid_batch.push(target);
if rid_batch.len() >= scan_batch_size {
let values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
None,
None,
)
.await?;
rid_batch.clear();
yielder.emit(ValueBatch::new(values)).await;
}
}
continue;
}
let mut edges_yielded: usize = 0;
let (adjacency, _) = batch_table_info[&rid.table];
if adjacency == VertexAdjacency::Lightweight
&& !txn
.record_exists(ns_id, db_id, &rid.table, &rid.key, version)
.await
.context("Failed to check a lightweight source's existence")?
{
continue;
}
'dir_loop: for &dir in directions.iter() {
if let Some(entries) = cache_map.remove(&(src_idx, dir)) {
let selected: Box<
dyn Iterator<Item = &EdgeCacheEntry> + Send + '_,
> = if spec_prefixes.is_empty() {
Box::new(entries.iter())
} else {
Box::new(spec_prefixes.iter().flat_map(|prefix| {
entries
.iter()
.filter(move |e| e.edge.starts_with(prefix.as_slice()))
}))
};
let selected: Box<
dyn Iterator<Item = &EdgeCacheEntry> + Send + '_,
> = if limit_needs_id_order && edge_limit.is_some() {
let mut ordered: Vec<&EdgeCacheEntry> = selected.collect();
ordered.sort_unstable_by(|a, b| a.edge.cmp(&b.edge));
Box::new(ordered.into_iter())
} else {
selected
};
let mut scanned: u64 = 0;
let mut limit_hit = false;
for entry in selected {
scanned += 1;
let (edge, target) = decode_edge_cache_entry(entry)?;
match output_mode {
GraphScanOutput::TargetVertex => {
if target_tables.is_empty()
|| target_tables.contains(&target.table)
{
rid_batch.push(target);
edges_yielded += 1;
}
}
_ if payload_spec.is_some() => {
let spec = payload_spec
.as_deref()
.expect("guarded by the match arm");
let pred = record_predicate
.as_ref()
.expect("payload mode implies a predicate");
let item = crate::idx::adjacency::MergedAdjacencyEdge {
key: Vec::new(),
edge,
target: Some(target),
props: None,
};
let (edge, candidate) =
synthesize_payload_candidate(spec, rid, dir, item);
let Some(candidate) = candidate else {
return Err(ControlFlow::Err(anyhow::anyhow!(
"a key-only predicate cannot fall back on a cached entry"
)));
};
if record_metrics {
metrics.add_props_evals(1);
}
if pred
.evaluate(
EvalContext::from_exec_ctx(&ctx)
.with_value(&candidate),
)
.await?
.is_truthy()
{
rid_batch.push(edge);
}
}
_ => {
rid_batch.push(edge);
if !has_predicate {
edges_yielded += 1;
}
}
}
if !has_predicate {
if edge_limit.is_some_and(|l| edges_yielded >= l) {
limit_hit = true;
break;
}
if rid_batch.len() >= scan_batch_size {
let values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
key_predicate.as_ref(),
record_predicate.as_ref(),
)
.await?;
rid_batch.clear();
yielder.emit(ValueBatch::new(values)).await;
}
}
}
if record_metrics && scanned > 0 {
metrics.add_edges_scanned(scanned);
}
if has_predicate && !rid_batch.is_empty() {
let flush_record_predicate = if payload_spec.is_some() {
None
} else {
record_predicate.as_ref()
};
let mut values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
key_predicate.as_ref(),
flush_record_predicate,
)
.await?;
rid_batch.clear();
if let Some(l) = edge_limit {
let room = l.saturating_sub(edges_yielded);
if values.len() >= room {
values.truncate(room);
limit_hit = true;
}
}
edges_yielded += values.len();
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
}
if limit_hit {
break 'dir_loop;
}
continue;
}
let scopes =
compute_graph_ranges(ns_id, db_id, rid, dir, &edge_tables, &ctx)
.await?;
let mut dir_delta_hits: u64 = 0;
for scope in scopes {
let mut chunk = scope.range.clone();
let mut limit_hit = false;
'range_chunks: loop {
let mut legacy_edges: Vec<RecordId> = Vec::new();
let mut chunk_bound_hit = false;
let mut last_processed_key: Option<Vec<u8>> = None;
{
let adjacency_scope = AdjacencyScope {
ns: ns_id,
db: db_id,
vertex: rid,
dir: Some(dir),
edge_table: scope.edge_table.as_ref(),
fk_lower: scope.fk_lower.as_ref(),
delta_range: chunk.clone(),
};
let mut cursor = match adjacency {
VertexAdjacency::Lightweight => {
MergedAdjacencyCursor::open_lightweight(
&txn,
&adjacency_scope,
)
.context("Failed to open graph cursor")?
}
VertexAdjacency::Normal {
folded,
} if payload_spec
.as_deref()
.is_some_and(|s| !s.key_only) =>
{
MergedAdjacencyCursor::open_with_values(
&txn,
adjacency_scope,
folded,
version,
Some(
ctx.root()
.ctx
.get_index_stores()
.adjacency_resolve(),
),
)
.await
.context("Failed to open graph cursor")?
}
VertexAdjacency::Normal {
folded,
} => MergedAdjacencyCursor::open(
&txn,
adjacency_scope,
folded,
version,
Some(
ctx.root()
.ctx
.get_index_stores()
.adjacency_resolve(),
),
)
.await
.context("Failed to open graph cursor")?,
};
'cursor_loop: loop {
let batch_size = if has_predicate {
crate::kvs::NORMAL_BATCH_SIZE
} else {
let remaining = edge_limit.map(|l| {
l.saturating_sub(edges_yielded)
.min(crate::kvs::NORMAL_BATCH_SIZE as usize)
});
match remaining {
Some(0) => {
limit_hit = true;
break;
}
Some(r) => r as u32,
None => crate::kvs::NORMAL_BATCH_SIZE,
}
};
let visited = if payload_spec.is_none() {
cursor
.visit_unfolded_batch(batch_size, |key| {
let decoded =
crate::key::schema::DecodedGraph::decode(
key,
)?;
if output_mode
== GraphScanOutput::TargetVertex
{
match decoded.target {
Some(target)
if target_tables.is_empty()
|| target_tables
.contains(
&target.table,
) =>
{
rid_batch.push(target);
edges_yielded += 1;
}
Some(_) => {}
None => {
legacy_edges.push(decoded.edge);
if legacy_edges.len()
>= scan_batch_size
{
chunk_bound_hit = true;
last_processed_key =
Some(key.to_vec());
return Ok(VisitFlow::Stop);
}
}
}
} else {
rid_batch.push(decoded.edge);
if !has_predicate {
edges_yielded += 1;
}
}
if !has_predicate
&& edge_limit
.is_some_and(|l| edges_yielded >= l)
{
limit_hit = true;
return Ok(VisitFlow::Stop);
}
Ok(VisitFlow::Continue)
})
.await
.context("Failed to scan graph edge")?
} else {
None
};
match visited {
Some(0) => break,
Some(n) => {
if record_metrics {
metrics.add_edges_scanned(n as u64);
}
}
None => {
let batch = cursor
.next_batch_scan(batch_size)
.await
.context("Failed to scan graph edge")?;
if batch.is_empty() {
break;
}
let mut scanned_in_batch: u64 = 0;
for item in batch {
scanned_in_batch += 1;
if output_mode
== GraphScanOutput::TargetVertex
{
match item.target {
Some(target)
if target_tables.is_empty()
|| target_tables
.contains(
&target.table,
) =>
{
rid_batch.push(target);
edges_yielded += 1;
}
Some(_) => {
}
None => {
legacy_edges.push(item.edge);
if legacy_edges.len()
>= scan_batch_size
{
chunk_bound_hit = true;
last_processed_key =
Some(item.key);
break;
}
}
}
} else if let Some(spec) =
payload_spec.as_deref()
{
let matched =
match (&spec.matcher, &item.target)
{
(
Some(matcher),
Some(target),
) if dir != Dir::Both => Some(matcher.matches(
rid,
&item.edge,
target,
dir == Dir::Out,
)),
_ => None,
};
if let Some(keep) = matched {
if record_metrics {
metrics.add_props_evals(1);
}
if keep {
rid_batch.push(item.edge);
}
} else {
let entry =
synthesize_payload_candidate(
spec, rid, dir, item,
);
if record_metrics {
if entry.1.is_some() {
metrics.add_props_evals(1);
} else {
metrics
.add_props_fallbacks(1);
}
}
payload_pending.push(entry);
}
} else {
rid_batch.push(item.edge);
if !has_predicate {
edges_yielded += 1;
}
}
if !has_predicate
&& edge_limit
.is_some_and(|l| edges_yielded >= l)
{
limit_hit = true;
break;
}
}
if record_metrics && scanned_in_batch > 0 {
metrics.add_edges_scanned(scanned_in_batch);
}
}
}
let should_flush = if has_predicate {
!rid_batch.is_empty() || !payload_pending.is_empty()
} else {
rid_batch.len() >= scan_batch_size
};
if should_flush && !has_predicate {
#[cfg(not(target_family = "wasm"))]
{
if pipelined_resolve {
if let Some(pending) = inflight.take() {
let (values, cache, drained) =
join_inflight_resolve(pending)
.await?;
perm_cache = cache;
spare = drained;
yielder
.emit(ValueBatch::new(values))
.await;
}
debug_assert!(
!has_predicate,
"inflight stash is only valid without a pushed-down predicate"
);
let task_ctx = ctx.clone();
let task_txn = Arc::clone(&txn);
let task_cache =
std::mem::take(&mut perm_cache);
let rids = std::mem::replace(
&mut rid_batch,
std::mem::take(&mut spare),
);
if rid_batch.capacity() == 0 {
rid_batch.reserve(scan_batch_size);
}
let mut set = tokio::task::JoinSet::new();
set.spawn(async move {
let mut rids = rids;
let mut cache = task_cache;
let values = resolve_and_filter_batch(
&task_ctx,
&task_txn,
ns_id,
db_id,
&rids,
fetch_full,
check_perms,
version,
&mut cache,
None,
None,
)
.await?;
rids.clear();
Ok::<_, ControlFlow>((
values, cache, rids,
))
});
inflight = Some(set);
} else {
let values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
None,
None,
)
.await?;
rid_batch.clear();
yielder.emit(ValueBatch::new(values)).await;
}
}
#[cfg(target_family = "wasm")]
{
let values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
None,
None,
)
.await?;
rid_batch.clear();
yielder.emit(ValueBatch::new(values)).await;
}
} else if should_flush {
if !payload_pending.is_empty() {
let pred = record_predicate
.as_ref()
.expect("payload mode implies a predicate");
filter_payload_pending(
&ctx,
&txn,
ns_id,
db_id,
version,
pred,
std::mem::take(&mut payload_pending),
&mut rid_batch,
)
.await?;
}
let flush_record_predicate =
if payload_spec.is_some() {
None
} else {
record_predicate.as_ref()
};
let mut values = if rid_batch.is_empty() {
Vec::new()
} else {
resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
key_predicate.as_ref(),
flush_record_predicate,
)
.await?
};
rid_batch.clear();
if has_predicate {
if let Some(l) = edge_limit {
let room = l.saturating_sub(edges_yielded);
if values.len() >= room {
values.truncate(room);
limit_hit = true;
}
}
edges_yielded += values.len();
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
} else {
yielder.emit(ValueBatch::new(values)).await;
}
}
if limit_hit || chunk_bound_hit {
break 'cursor_loop;
}
}
let merge_stats = cursor.stats();
dir_delta_hits += merge_stats.delta_hits;
if record_metrics {
metrics.add_block_hits(merge_stats.block_hits);
metrics.add_delta_hits(merge_stats.delta_hits);
}
drop(cursor);
}
#[cfg(not(target_family = "wasm"))]
if let Some(pending) = inflight.take() {
let (values, cache, drained) =
join_inflight_resolve(pending).await?;
perm_cache = cache;
spare = drained;
yielder.emit(ValueBatch::new(values)).await;
}
if !legacy_edges.is_empty() && !limit_hit {
let inner_specs: Vec<EdgeTableSpec> =
if target_tables.is_empty() {
Vec::new()
} else {
target_tables
.iter()
.cloned()
.map(|t| EdgeTableSpec {
table: t,
range_start: std::ops::Bound::Unbounded,
range_end: std::ops::Bound::Unbounded,
})
.collect()
};
'legacy_loop: for edge_rid in legacy_edges {
let inner_scopes = compute_graph_ranges(
ns_id,
db_id,
&edge_rid,
dir,
&inner_specs,
&ctx,
)
.await?;
for inner in inner_scopes {
let mut inner_cursor = txn
.open_keys_cursor_raw(
inner.range,
Direction::Forward,
0,
version,
)
.await
.context(
"Failed to open legacy-fallback graph cursor",
)?;
loop {
let inner_batch = inner_cursor
.next_batch(crate::kvs::NORMAL_BATCH_SIZE)
.await
.context(
"Failed to scan edge adjacency for legacy graph fallback",
)?;
if inner_batch.is_empty() {
break;
}
for ik in &inner_batch {
let endpoint =
inner.decoder.decode_edge(ik).context(
"Failed to decode graph key",
)?;
rid_batch.push(endpoint);
edges_yielded += 1;
if edge_limit
.is_some_and(|l| edges_yielded >= l)
{
limit_hit = true;
break;
}
}
if rid_batch.len() >= scan_batch_size {
let values = resolve_record_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
CachePolicy::ReadWrite,
&mut perm_cache,
)
.await?;
yielder.emit(ValueBatch::new(values)).await;
rid_batch.clear();
}
if limit_hit {
break;
}
}
drop(inner_cursor);
if limit_hit {
break 'legacy_loop;
}
}
}
}
if !chunk_bound_hit || limit_hit {
break 'range_chunks;
}
let resume_from = last_processed_key
.expect("chunk_bound_hit implies a key was processed");
chunk = chunk.resume_after(&resume_from, Direction::Forward);
}
if limit_hit {
break 'dir_loop;
}
}
observe_fold_candidate(
&txn,
ns_id,
db_id,
rid,
dir,
dir_delta_hits,
fold_threshold,
);
}
}
#[cfg(not(target_family = "wasm"))]
if let Some(started) = win_timer {
let per_source =
started.elapsed().as_nanos() as u64 / (win_end - win_start) as u64;
fanout_engaged = Some(per_source >= GRAPH_SPAWN_PER_SOURCE_NANOS);
}
win_start = win_end;
}
}
if !rid_batch.is_empty() {
let values = resolve_and_filter_batch(
&ctx,
&txn,
ns_id,
db_id,
&rid_batch,
fetch_full,
check_perms,
version,
&mut perm_cache,
key_predicate.as_ref(),
record_predicate.as_ref(),
)
.await?;
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
}
Ok(())
});
Ok(monitor_stream(Box::pin(stream), "GraphEdgeScan", &self.metrics))
}
}
type PayloadPending = (RecordId, Option<Value>);
fn synthesize_payload_candidate(
spec: &PayloadPredicateSpec,
source: &RecordId,
dir: Dir,
item: crate::idx::adjacency::MergedAdjacencyEdge,
) -> PayloadPending {
use surrealdb_datastore::values::graph::InlineProps;
let edge = item.edge;
if item.target.is_none() && (!spec.key_only || spec.needs_in || spec.needs_out) {
return (edge, None);
}
if matches!(dir, Dir::Both) {
return (edge, None);
}
let fields: Vec<(String, Value)> = if spec.key_only {
Vec::new()
} else {
let Some(table_spec) = spec.tables.get(&edge.table) else {
return (edge, None);
};
let Some(body) = item.props.as_deref() else {
return (edge, None);
};
let Ok(props) = InlineProps::decode(body) else {
return (edge, None);
};
if props.generation != table_spec.generation {
return (edge, None);
}
let Some(values) = props.values else {
return (edge, None);
};
if values.len() != table_spec.fields.len() {
return (edge, None);
}
table_spec
.fields
.iter()
.zip(&table_spec.referenced)
.zip(values)
.filter(|((_, referenced), _)| **referenced)
.map(|((name, _), value)| (name.clone(), value))
.collect()
};
let mut object = crate::val::Object::default();
if spec.needs_id {
object.insert("id".to_owned(), Value::RecordId(edge.clone()));
}
if spec.needs_in || spec.needs_out {
let target = item
.target
.as_ref()
.expect("an in/out-reading predicate falls back on target-less entries");
let (in_rid, out_rid) = match dir {
Dir::Out => (source, target),
Dir::In => (target, source),
Dir::Both => unreachable!("handled above"),
};
if spec.needs_in {
object.insert("in".to_owned(), Value::RecordId(in_rid.clone()));
}
if spec.needs_out {
object.insert("out".to_owned(), Value::RecordId(out_rid.clone()));
}
}
for (name, value) in fields {
object.insert(name, value);
}
(edge, Some(Value::Object(object)))
}
#[allow(clippy::too_many_arguments)]
async fn filter_payload_pending(
ctx: &ExecutionContext,
txn: &Transaction,
ns_id: NamespaceId,
db_id: DatabaseId,
version: Option<u64>,
predicate: &Arc<dyn PhysicalExpr>,
pending: Vec<PayloadPending>,
out: &mut Vec<RecordId>,
) -> Result<(), ControlFlow> {
let mut candidates: Vec<Value> = Vec::new();
let mut fallback_rids: Vec<RecordId> = Vec::new();
let mut order: Vec<(RecordId, bool)> = Vec::with_capacity(pending.len());
for (rid, candidate) in pending {
match candidate {
Some(value) => {
candidates.push(value);
order.push((rid, true));
}
None => {
fallback_rids.push(rid.clone());
order.push((rid, false));
}
}
}
let candidate_results = if candidates.is_empty() {
Vec::new()
} else {
predicate.evaluate_batch(EvalContext::from_exec_ctx(ctx), &candidates).await?
};
let fallback_matches = if fallback_rids.is_empty() {
Vec::new()
} else {
let records = txn
.get_records(ns_id, db_id, &fallback_rids, version, CachePolicy::ReadWrite)
.await
.context("Failed to fetch the payload-fallback edge records")?;
let mut matches = vec![false; records.len()];
let mut present_slots: Vec<usize> = Vec::new();
let mut present_values: Vec<Value> = Vec::new();
for (slot, record) in records.iter().enumerate() {
if record.metadata.is_some() || !record.data.is_none() {
present_slots.push(slot);
present_values.push(record.data.clone());
}
}
if !present_values.is_empty() {
let results =
predicate.evaluate_batch(EvalContext::from_exec_ctx(ctx), &present_values).await?;
for (slot, result) in present_slots.into_iter().zip(results) {
matches[slot] = result.is_truthy();
}
}
matches
};
let mut candidate_results = candidate_results.into_iter();
let mut fallback_matches = fallback_matches.into_iter();
for (rid, evaluated) in order {
let keep = if evaluated {
candidate_results.next().is_some_and(|v| v.is_truthy())
} else {
fallback_matches.next().unwrap_or(false)
};
if keep {
out.push(rid);
}
}
Ok(())
}
fn decode_edge_cache_entry(entry: &EdgeCacheEntry) -> Result<(RecordId, RecordId), anyhow::Error> {
let (table, key): (TableName, RecordIdKey) = storekey::decode_borrow(&entry.edge)
.map_err(|e| anyhow::anyhow!("Failed to decode a cached edge identity: {e}"))?;
let (t_table, t_key): (TableName, RecordIdKey) = storekey::decode_borrow(&entry.target)
.map_err(|e| anyhow::anyhow!("Failed to decode a cached edge target: {e}"))?;
Ok((RecordId::new(table, key), RecordId::new(t_table, t_key)))
}
#[allow(clippy::too_many_arguments)]
async fn resolve_and_filter_batch(
ctx: &ExecutionContext,
txn: &Transaction,
ns_id: NamespaceId,
db_id: DatabaseId,
rids: &[RecordId],
fetch_full: bool,
check_perms: bool,
version: Option<u64>,
perm_cache: &mut PermCache,
key_predicate: Option<&Arc<dyn PhysicalExpr>>,
record_predicate: Option<&Arc<dyn PhysicalExpr>>,
) -> Result<Vec<Value>, ControlFlow> {
let values = if let Some(pred) = key_predicate {
let mut survivors: Vec<RecordId> = Vec::with_capacity(rids.len());
for rid in rids {
let candidate = Value::RecordId(rid.clone());
let keep = pred
.evaluate(EvalContext::from_exec_ctx(ctx).with_value(&candidate))
.await?
.is_truthy();
if keep {
survivors.push(rid.clone());
}
}
if survivors.is_empty() {
return Ok(Vec::new());
}
resolve_record_batch(
ctx,
txn,
ns_id,
db_id,
&survivors,
fetch_full,
check_perms,
version,
CachePolicy::ReadWrite,
perm_cache,
)
.await?
} else {
resolve_record_batch(
ctx,
txn,
ns_id,
db_id,
rids,
fetch_full,
check_perms,
version,
CachePolicy::ReadWrite,
perm_cache,
)
.await?
};
match record_predicate {
Some(pred) => apply_record_predicate(pred, ctx, values).await,
None => Ok(values),
}
}
type PermCache = std::collections::HashMap<TableName, PhysicalPermission>;
#[cfg(not(target_family = "wasm"))]
type InflightResolve =
tokio::task::JoinSet<Result<(Vec<Value>, PermCache, Vec<RecordId>), ControlFlow>>;
#[cfg(not(target_family = "wasm"))]
async fn join_inflight_resolve(
mut set: InflightResolve,
) -> Result<(Vec<Value>, PermCache, Vec<RecordId>), ControlFlow> {
let joined = set.join_next().await.expect("the inflight resolve set holds exactly one task");
joined.map_err(|e| ControlFlow::Err(anyhow::anyhow!("graph resolve task failed: {e}")))?
}
async fn apply_record_predicate(
predicate: &Arc<dyn PhysicalExpr>,
ctx: &ExecutionContext,
mut values: Vec<Value>,
) -> Result<Vec<Value>, ControlFlow> {
let results = {
let eval_ctx = EvalContext::from_exec_ctx(ctx);
predicate.evaluate_batch(eval_ctx, &values).await?
};
let mut write = 0;
for (read, result) in results.into_iter().enumerate() {
if result.is_truthy() {
if write != read {
values.swap(write, read);
}
write += 1;
}
}
values.truncate(write);
Ok(values)
}
const GRAPH_FANOUT_CONCURRENCY: usize = 16;
#[cfg(not(target_family = "wasm"))]
const GRAPH_SPAWN_PER_SOURCE_NANOS: u64 = 5_000;
#[cfg(not(target_family = "wasm"))]
#[allow(clippy::too_many_arguments)]
async fn scan_source_plain(
ctx: &ExecutionContext,
txn: &Transaction,
adjacency: VertexAdjacency,
ns_id: NamespaceId,
db_id: DatabaseId,
rid: &RecordId,
directions: &[Dir],
edge_tables: &[EdgeTableSpec],
output_mode: GraphScanOutput,
target_tables: &[TableName],
version: Option<u64>,
edge_limit: Option<usize>,
fold_threshold: usize,
cap: usize,
metrics: &OperatorMetrics,
record_metrics: bool,
) -> Result<Option<Vec<RecordId>>, ControlFlow> {
crate::exec::operators::check_cancelled(ctx)?;
if adjacency == VertexAdjacency::Lightweight
&& !txn
.record_exists(ns_id, db_id, &rid.table, &rid.key, version)
.await
.context("Failed to check a lightweight source's existence")?
{
return Ok(Some(Vec::new()));
}
let mut out: Vec<RecordId> = Vec::new();
let mut scanned: u64 = 0;
let mut block_hits: u64 = 0;
let mut delta_hits: u64 = 0;
let mut capped = false;
let mut observations: Vec<(Dir, u64)> = Vec::new();
for &dir in directions {
let scopes = compute_graph_ranges(ns_id, db_id, rid, dir, edge_tables, ctx).await?;
let mut dir_delta_hits: u64 = 0;
for scope in scopes {
let adjacency_scope = AdjacencyScope {
ns: ns_id,
db: db_id,
vertex: rid,
dir: Some(dir),
edge_table: scope.edge_table.as_ref(),
fk_lower: scope.fk_lower.as_ref(),
delta_range: scope.range.clone(),
};
let mut cursor = match adjacency {
VertexAdjacency::Lightweight => {
MergedAdjacencyCursor::open_lightweight(txn, &adjacency_scope)
.context("Failed to open graph cursor")?
}
VertexAdjacency::Normal {
folded,
} => MergedAdjacencyCursor::open(
txn,
adjacency_scope,
folded,
version,
Some(ctx.root().ctx.get_index_stores().adjacency_resolve()),
)
.await
.context("Failed to open graph cursor")?,
};
let mut bail = false;
loop {
let batch_size = match edge_limit {
Some(l) => {
let remaining = l.saturating_sub(out.len());
if remaining == 0 {
capped = true;
break;
}
remaining.min(crate::kvs::NORMAL_BATCH_SIZE as usize) as u32
}
None => crate::kvs::NORMAL_BATCH_SIZE,
};
let visited = cursor
.visit_unfolded_batch(batch_size, |key| {
let decoded = crate::key::schema::DecodedGraph::decode(key)?;
if output_mode == GraphScanOutput::TargetVertex {
match decoded.target {
Some(target) => {
if target_tables.is_empty()
|| target_tables.contains(&target.table)
{
out.push(target);
}
}
None => {
bail = true;
return Ok(VisitFlow::Stop);
}
}
} else {
out.push(decoded.edge);
}
if edge_limit.is_some_and(|l| out.len() >= l) {
capped = true;
return Ok(VisitFlow::Stop);
}
if out.len() >= cap {
bail = true;
return Ok(VisitFlow::Stop);
}
Ok(VisitFlow::Continue)
})
.await
.context("Failed to scan graph edge")?;
match visited {
Some(0) => break,
Some(n) => scanned += n as u64,
None => {
let batch = cursor
.next_batch_scan(batch_size)
.await
.context("Failed to scan graph edge")?;
if batch.is_empty() {
break;
}
scanned += batch.len() as u64;
for item in batch {
if output_mode == GraphScanOutput::TargetVertex {
match item.target {
Some(target) => {
if target_tables.is_empty()
|| target_tables.contains(&target.table)
{
out.push(target);
}
}
None => {
bail = true;
break;
}
}
} else {
out.push(item.edge);
}
if edge_limit.is_some_and(|l| out.len() >= l) {
capped = true;
break;
}
if out.len() >= cap {
bail = true;
break;
}
}
}
}
if bail || capped {
break;
}
}
let stats = cursor.stats();
dir_delta_hits += stats.delta_hits;
block_hits += stats.block_hits;
delta_hits += stats.delta_hits;
drop(cursor);
if bail {
return Ok(None);
}
if capped {
break;
}
}
if capped {
break;
}
observations.push((dir, dir_delta_hits));
}
if record_metrics {
metrics.add_edges_scanned(scanned);
metrics.add_block_hits(block_hits);
metrics.add_delta_hits(delta_hits);
}
for (dir, dir_delta_hits) in observations {
observe_fold_candidate(txn, ns_id, db_id, rid, dir, dir_delta_hits, fold_threshold);
}
Ok(Some(out))
}
#[cfg(all(feature = "kv-mem", not(target_family = "wasm")))]
#[cfg(test)]
mod prescan_tests {
use super::*;
use crate::exec::operators::test_util::TestDb;
async fn traversal_db() -> TestDb {
let mut setup = String::from("CREATE person:big; CREATE person:small;");
for i in 0..10 {
setup.push_str(&format!("CREATE person:t{i}; RELATE person:big->knows->person:t{i};"));
}
for i in 0..3 {
setup.push_str(&format!("RELATE person:small->knows->person:t{i};"));
}
TestDb::new(&setup).await
}
async fn prescan(db: &TestDb, source: &str, cap: usize) -> Option<Vec<RecordId>> {
let ctx = db.exec_ctx().await;
let txn = ctx.txn();
let db_ctx = ctx.database().expect("database context").clone();
let tb = txn
.get_tb(db_ctx.ns_ctx.ns.namespace_id, db_ctx.db.database_id, &"person".into(), None)
.await
.expect("the source table resolves");
let metrics = OperatorMetrics::new();
scan_source_plain(
&ctx,
&txn,
vertex_adjacency_of(tb.as_deref()),
db_ctx.ns_ctx.ns.namespace_id,
db_ctx.db.database_id,
&RecordId::new("person".into(), source.to_string()),
&[Dir::Out],
&[],
GraphScanOutput::TargetVertex,
&[],
None,
None,
usize::MAX,
cap,
&metrics,
false,
)
.await
.expect("prescan should succeed")
}
#[tokio::test]
async fn a_source_over_the_cap_defers_to_the_sequential_machinery() {
let db = traversal_db().await;
assert!(
prescan(&db, "big", 5).await.is_none(),
"a source with more edges than the cap must not buffer them all"
);
}
#[tokio::test]
async fn a_source_under_the_cap_returns_its_complete_adjacency() {
let db = traversal_db().await;
let mut out = prescan(&db, "small", 5).await.expect("three edges fit a cap of five");
out.sort();
let expected: Vec<RecordId> =
(0..3).map(|i| RecordId::new("person".into(), format!("t{i}"))).collect();
assert_eq!(out, expected);
}
}
#[cfg(all(feature = "kv-mem", not(target_family = "wasm")))]
#[cfg(test)]
mod window_tests {
use std::ops::Bound;
use futures::StreamExt;
use surrealdb_cnf::ConfigMap;
use super::*;
use crate::exec::operators::test_util::{TestDb, ValuesOperator, collect};
fn knows_scan(sources: Vec<Value>) -> (Arc<dyn ExecOperator>, Arc<OperatorMetrics>) {
let scan = GraphEdgeScan::new(
ValuesOperator::new(sources),
LookupDirection::Out,
vec![EdgeTableSpec {
table: "knows".into(),
range_start: Bound::Unbounded,
range_end: Bound::Unbounded,
}],
GraphScanOutput::TargetVertex,
None,
);
scan.metrics.enable();
let metrics = Arc::clone(&scan.metrics);
(Arc::new(scan), metrics)
}
fn person(key: impl std::fmt::Display) -> Value {
Value::RecordId(RecordId::new("person".into(), key.to_string()))
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn dropping_the_stream_stops_the_remaining_windows() {
const SOURCES: usize = 64;
const EDGES_PER_SOURCE: usize = 8;
let mut setup = String::new();
for j in 0..EDGES_PER_SOURCE {
setup.push_str(&format!("CREATE person:t{j};"));
}
for i in 0..SOURCES {
setup.push_str(&format!("CREATE person:s{i};"));
for j in 0..EDGES_PER_SOURCE {
setup.push_str(&format!("RELATE person:s{i}->knows->person:t{j};"));
}
}
let db = TestDb::new_with_config(
&setup,
ConfigMap::empty().with_key_value("scan_batch_size", "8"),
)
.await;
let ctx = db.exec_ctx().await;
let (op, metrics) = knows_scan((0..SOURCES).map(|i| person(format!("s{i}"))).collect());
let mut stream = op.execute(&ctx).expect("execute should succeed");
let first = stream.next().await.expect("one batch").expect("batch should be Ok");
assert!(!first.values().is_empty(), "the first batch should carry rows");
drop(stream);
let scanned = metrics.edges_scanned();
let window = (GRAPH_FANOUT_CONCURRENCY * EDGES_PER_SOURCE) as u64;
assert!(
scanned <= 2 * window,
"dropping after one batch should stop the later windows: \
scanned {scanned} edges, expected at most two windows ({})",
2 * window
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn results_are_exact_across_window_and_cap_boundaries() {
const SOURCES: usize = 40;
const HUB_EDGES: usize = 20;
let mut setup = String::from("CREATE person:s0;");
for j in 0..HUB_EDGES {
setup.push_str(&format!("CREATE person:u{j}; RELATE person:s0->knows->person:u{j};"));
}
for i in 1..SOURCES {
setup.push_str(&format!(
"CREATE person:s{i}; CREATE person:v{i}; \
RELATE person:s{i}->knows->person:v{i};"
));
}
let db = TestDb::new_with_config(
&setup,
ConfigMap::empty().with_key_value("scan_batch_size", "8"),
)
.await;
let ctx = db.exec_ctx().await;
let (op, _) = knows_scan((0..SOURCES).map(|i| person(format!("s{i}"))).collect());
let mut out = collect(&op, &ctx).await;
out.sort();
let mut expected: Vec<Value> = (0..HUB_EDGES)
.map(|j| person(format!("u{j}")))
.chain((1..SOURCES).map(|i| person(format!("v{i}"))))
.collect();
expected.sort();
assert_eq!(out, expected);
}
}
#[cfg(test)]
mod tests {
use std::ops::Bound;
use super::*;
use crate::exec::operators::CurrentValueSource;
#[test]
fn test_graph_edge_scan_attrs() {
let scan = GraphEdgeScan::new(
Arc::new(CurrentValueSource::new()),
LookupDirection::Out,
vec![
EdgeTableSpec {
table: "knows".into(),
range_start: Bound::Unbounded,
range_end: Bound::Unbounded,
},
EdgeTableSpec {
table: "follows".into(),
range_start: Bound::Unbounded,
range_end: Bound::Unbounded,
},
],
GraphScanOutput::TargetId,
None,
);
assert_eq!(scan.name(), "GraphEdgeScan");
let attrs = scan.attrs();
assert!(attrs.iter().any(|(k, v)| k == "direction" && v == "->"));
assert!(attrs.iter().any(|(k, v)| k == "tables" && v.contains("knows")));
}
}