use std::sync::Arc;
use common::future::stream::{self, Yielder};
use futures::StreamExt;
use crate::exec::{
AccessMode, CardinalityHint, ContextLevel, Error as ExecError, ExecOperator, ExecutionContext,
FlowResult, OperatorMetrics, OutputShape, ValueBatch, ValueBatchStream, buffer_stream,
monitor_stream,
};
use crate::expr::ControlFlow;
use crate::val::Value;
#[derive(Debug, Clone)]
pub struct UnwrapExactlyOne {
pub(crate) input: Arc<dyn ExecOperator>,
pub(crate) none_on_empty: bool,
pub(crate) metrics: Arc<OperatorMetrics>,
}
impl UnwrapExactlyOne {
pub(crate) fn new(input: Arc<dyn ExecOperator>, none_on_empty: bool) -> Self {
Self {
input,
none_on_empty,
metrics: Arc::new(OperatorMetrics::new()),
}
}
}
impl ExecOperator for UnwrapExactlyOne {
fn name(&self) -> &'static str {
"UnwrapExactlyOne"
}
fn required_context(&self) -> ContextLevel {
self.input.required_context()
}
fn access_mode(&self) -> AccessMode {
self.input.access_mode()
}
fn cardinality_hint(&self) -> CardinalityHint {
CardinalityHint::AtMostOne
}
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
vec![&self.input]
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn output_shape(&self) -> OutputShape {
OutputShape::Scalar
}
fn output_ordering(&self) -> crate::exec::OutputOrdering {
self.input.output_ordering()
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
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 none_on_empty = self.none_on_empty;
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
let mut collected: Vec<Value> = Vec::new();
while let Some(batch_result) = input_stream.next().await {
let batch = batch_result?;
collected.extend(batch.into_values());
if collected.len() > 1 {
Err(ControlFlow::Err(anyhow::anyhow!(ExecError::SingleOnlyOutput)))?;
}
}
if collected.is_empty() {
if none_on_empty {
yielder.emit(ValueBatch::new(vec![Value::None])).await;
} else {
Err(ControlFlow::Err(anyhow::anyhow!(ExecError::SingleOnlyOutput)))?;
}
} else {
let result = collected.pop().expect("collected has exactly one element");
yielder.emit(ValueBatch::new(vec![result])).await;
}
Ok(())
});
Ok(monitor_stream(Box::pin(stream), "UnwrapExactlyOne", &self.metrics))
}
}