use std::sync::Arc;
use futures::{StreamExt, TryStreamExt, stream};
use surrealdb_types::ToSql;
use crate::exec::parts::is_final;
use crate::exec::{BoxFut, ExecOperator, ExecutionContext, FlowResult};
use crate::val::Value;
pub(crate) const RECURSION_CONCURRENCY: usize = 16;
#[derive(Clone, Copy, Debug)]
pub(crate) struct RecursionBounds {
pub(crate) min: u32,
pub(crate) max: Option<u32>,
pub(crate) system_limit: u32,
}
impl RecursionBounds {
pub(crate) fn cap(&self) -> u32 {
self.max.unwrap_or(self.system_limit).min(self.system_limit)
}
pub(crate) fn errors_on_limit(&self) -> bool {
self.max.is_none()
}
}
pub(crate) fn is_recursion_target(value: &Value) -> bool {
match value {
Value::RecordId(_) => true,
Value::Array(arr) => arr.iter().any(is_recursion_target),
_ => false,
}
}
pub(crate) async fn eval_buffered<'a, T: 'a>(
futures: Vec<BoxFut<'a, FlowResult<T>>>,
) -> FlowResult<Vec<T>> {
if futures.len() < 2 {
let mut results = Vec::with_capacity(futures.len());
for fut in futures {
results.push(fut.await?);
}
Ok(results)
} else {
stream::iter(futures).buffered(RECURSION_CONCURRENCY).try_collect().await
}
}
pub(crate) async fn eval_buffered_all<'a>(
futures: Vec<BoxFut<'a, FlowResult<Value>>>,
) -> Vec<FlowResult<Value>> {
if futures.len() < 2 {
let mut results = Vec::with_capacity(futures.len());
for fut in futures {
results.push(fut.await);
}
results
} else {
stream::iter(futures).buffered(RECURSION_CONCURRENCY).collect().await
}
}
pub(crate) fn collect_discovery_targets(
v: Value,
out: &mut Vec<Value>,
) -> crate::exec::FlowResult<()> {
match v {
Value::Array(arr) => {
for inner in arr.0 {
if is_final(&inner) {
continue;
}
if !is_recursion_target(&inner) {
return Err(crate::err::Error::InvalidRecursionTarget {
value: inner.to_sql(),
}
.into());
}
out.push(inner);
}
}
v if is_final(&v) => {}
v if is_recursion_target(&v) => {
out.push(v);
}
v => {
return Err(crate::err::Error::InvalidRecursionTarget {
value: v.to_sql(),
}
.into());
}
}
Ok(())
}
pub(crate) fn discover_body_targets<'a>(
body_op: &'a Arc<dyn ExecOperator>,
body_ctx: ExecutionContext,
) -> BoxFut<'a, FlowResult<Vec<Value>>> {
Box::pin(async move {
let mut discovered = Vec::new();
let mut body_stream = body_op.execute(&body_ctx)?;
while let Some(batch_result) = body_stream.next().await {
let batch = batch_result?;
for v in batch.values {
collect_discovery_targets(v, &mut discovered)?;
}
}
Ok(discovered)
})
}