use std::sync::Arc;
use futures::stream;
use crate::exec::context::{ContextLevel, ExecutionContext};
use crate::exec::{
AccessMode, CardinalityHint, ExecOperator, FlowResult, OperatorMetrics, ValueBatch,
ValueBatchStream, monitor_stream,
};
use crate::val::Value;
#[derive(Debug, Clone)]
pub struct CurrentValueSource {
pub(crate) metrics: Arc<OperatorMetrics>,
}
impl CurrentValueSource {
pub(crate) fn new() -> Self {
Self {
metrics: Arc::new(OperatorMetrics::new()),
}
}
}
impl ExecOperator for CurrentValueSource {
fn name(&self) -> &'static str {
"CurrentValueSource"
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Root
}
fn access_mode(&self) -> AccessMode {
AccessMode::ReadOnly
}
fn cardinality_hint(&self) -> CardinalityHint {
CardinalityHint::AtMostOne
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let value = ctx.current_value().cloned().unwrap_or(Value::None);
Ok(monitor_stream(
Box::pin(stream::once(async move {
Ok(ValueBatch {
values: vec![value],
})
})),
"CurrentValueSource",
&self.metrics,
))
}
}