use std::fmt::Debug;
use std::pin::Pin;
use std::sync::Arc;
use futures::Stream;
use crate::exe::FlowResultExt;
pub(crate) use crate::expr::{ControlFlowExt, FlowResult};
pub(crate) type BoxFut<'a, T> = Pin<Box<dyn std::future::Future<Output = T> + Send + 'a>>;
pub(crate) mod access_mode;
pub(crate) mod batch;
pub(crate) mod buffer;
pub(crate) mod cardinality;
pub(crate) mod config;
pub(crate) mod context;
pub(crate) mod error;
pub(crate) mod expression_registry;
pub(crate) mod fan_out;
pub(crate) use crate::val::field_path;
pub(crate) mod field_path_convert;
pub(crate) mod function;
pub(crate) mod index;
pub(crate) mod metrics;
pub(crate) mod object_extract;
pub(crate) mod operators;
pub(crate) mod ordering;
pub(crate) mod partitioning;
pub(crate) mod parts;
pub(crate) mod permission;
pub(crate) mod physical_expr;
pub(crate) mod plan_or_compute;
pub(crate) mod planner;
pub(crate) mod pre_decode_filter;
pub(crate) mod topk_pushdown;
pub(crate) use access_mode::{AccessMode, CombineAccessModes};
pub(crate) use batch::ValueBatch;
pub(crate) use buffer::buffer_stream;
pub(crate) use cardinality::CardinalityHint;
pub(crate) use context::{ContextLevel, DatabaseContext, ExecutionContext};
pub(crate) use error::Error;
pub(crate) use metrics::{OperatorMetrics, monitor_stream};
pub(crate) use ordering::OutputOrdering;
pub(crate) use partitioning::Partitioning;
pub(crate) use physical_expr::{EvalContext, PhysicalExpr};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum OutputShape {
Rows,
Scalar,
}
impl OutputShape {
pub(crate) fn is_scalar(self) -> bool {
matches!(self, OutputShape::Scalar)
}
}
pub(crate) type ValueBatchStream = Pin<Box<dyn Stream<Item = FlowResult<ValueBatch>> + Send>>;
pub(crate) trait ExecOperator: Debug + Send + Sync {
fn name(&self) -> &'static str;
fn attrs(&self) -> Vec<(String, String)> {
vec![]
}
fn required_context(&self) -> ContextLevel;
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream>;
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
vec![]
}
fn mutates_context(&self) -> bool {
false
}
fn output_context<'a>(
&'a self,
input: &'a ExecutionContext,
) -> BoxFut<'a, crate::expr::FlowResult<ExecutionContext>> {
Box::pin(async move { Ok(input.clone()) })
}
fn access_mode(&self) -> AccessMode;
fn cardinality_hint(&self) -> CardinalityHint {
CardinalityHint::Unbounded
}
fn output_partitioning(&self) -> Partitioning {
Partitioning::Single
}
fn output_shape(&self) -> OutputShape {
OutputShape::Rows
}
fn metrics(&self) -> Option<&OperatorMetrics> {
None
}
fn enable_metrics(&self) {
if let Some(m) = self.metrics() {
m.enable();
}
for (_, expr) in self.expressions() {
for (_, embedded) in expr.embedded_operators() {
embedded.enable_metrics();
}
}
for child in self.children() {
child.enable_metrics();
}
}
fn expressions(&self) -> Vec<(&str, &Arc<dyn PhysicalExpr>)> {
vec![]
}
fn output_ordering(&self) -> OutputOrdering {
OutputOrdering::Unordered
}
fn constant_output_fields(&self) -> Vec<crate::exec::field_path::FieldPath> {
vec![]
}
}