use std::sync::Arc;
use common::future::stream::{self, Yielder};
use surrealdb_types::{SqlFormat, ToSql};
use crate::exec::context::{ContextLevel, ExecutionContext};
use crate::exec::physical_expr::{EvalContext, PhysicalExpr};
use crate::exec::{
AccessMode, ExecOperator, FlowResult, OperatorMetrics, OutputShape, ValueBatch,
ValueBatchStream, monitor_stream,
};
use crate::val::Value;
#[derive(Debug, Clone)]
pub struct SourceExpr {
pub expr: Arc<dyn PhysicalExpr>,
pub(crate) metrics: Arc<OperatorMetrics>,
}
impl SourceExpr {
pub(crate) fn new(expr: Arc<dyn PhysicalExpr>) -> Self {
Self {
expr,
metrics: Arc::new(OperatorMetrics::new()),
}
}
}
impl ExecOperator for SourceExpr {
fn name(&self) -> &'static str {
"SourceExpr"
}
fn attrs(&self) -> Vec<(String, String)> {
vec![("expr".to_string(), self.expr.to_sql())]
}
fn required_context(&self) -> ContextLevel {
self.expr.required_context().max(ContextLevel::Database)
}
fn access_mode(&self) -> AccessMode {
self.expr.access_mode()
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn expressions(&self) -> Vec<(&str, &Arc<dyn PhysicalExpr>)> {
vec![("expr", &self.expr)]
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let expr = Arc::clone(&self.expr);
let ctx = ctx.clone();
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
let this_value = ctx.value("this").cloned();
let eval_ctx = EvalContext::from_exec_ctx(&ctx);
let eval_ctx = match this_value {
Some(ref v) => eval_ctx.with_value(v),
None => eval_ctx,
};
let value = expr.evaluate(eval_ctx).await?;
match value {
Value::Array(arr) => {
let mut values: Vec<Value> = arr
.into_iter()
.filter(|v| !matches!(v, Value::None | Value::Null))
.collect();
if !values.is_empty() {
super::fetch::batch_fetch_in_place(&ctx, &mut values).await?;
values.retain(|v| !matches!(v, Value::None | Value::Null));
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
}
}
Value::None | Value::Null => {}
Value::RecordId(ref rid) => {
let fetched = super::fetch::fetch_record(&ctx, rid).await?;
if !matches!(fetched, Value::None) {
yielder.emit(ValueBatch::new(vec![fetched])).await;
}
}
other => {
yielder.emit(ValueBatch::new(vec![other])).await;
}
}
Ok(())
});
Ok(monitor_stream(Box::pin(stream), "SourceExpr", &self.metrics))
}
fn output_shape(&self) -> OutputShape {
OutputShape::Rows
}
}
impl ToSql for SourceExpr {
fn fmt_sql(&self, f: &mut String, fmt: SqlFormat) {
self.expr.fmt_sql(f, fmt);
}
}