use std::sync::Arc;
use lora_analyzer::symbols::VarId;
use lora_compiler::physical::{PhysicalNodeId, PhysicalPlan};
use lora_store::GraphStorage;
use crate::errors::ExecResult;
use crate::executor::merge_optional_rows;
use crate::value::{LoraValue, Row};
use super::traits::build_streaming_seeded;
use super::{drain, RowSource};
pub struct CallSubquerySource<'a, S: GraphStorage> {
upstream: Box<dyn RowSource + 'a>,
plan: &'a PhysicalPlan,
inner: PhysicalNodeId,
storage: &'a S,
params: Arc<std::collections::BTreeMap<String, LoraValue>>,
new_vars: &'a [VarId],
pending: std::vec::IntoIter<Row>,
pending_outer: Option<Row>,
}
impl<'a, S: GraphStorage> CallSubquerySource<'a, S> {
pub(super) fn new(
upstream: Box<dyn RowSource + 'a>,
plan: &'a PhysicalPlan,
inner: PhysicalNodeId,
storage: &'a S,
params: Arc<std::collections::BTreeMap<String, LoraValue>>,
new_vars: &'a [VarId],
) -> Self {
Self {
upstream,
plan,
inner,
storage,
params,
new_vars,
pending: Vec::new().into_iter(),
pending_outer: None,
}
}
}
impl<'a, S: GraphStorage> RowSource for CallSubquerySource<'a, S> {
fn next_row(&mut self) -> ExecResult<Option<Row>> {
loop {
if let Some(inner_row) = self.pending.next() {
let outer = self
.pending_outer
.as_ref()
.expect("pending_outer set when pending iter has rows");
return Ok(Some(merge_optional_rows(outer, &inner_row)));
}
let Some(outer_row) = self.upstream.next_row()? else {
return Ok(None);
};
let mut inner_source = build_streaming_seeded(
self.plan,
self.inner,
self.storage,
self.params.clone(),
outer_row.clone(),
)?;
let inner_rows = drain(inner_source.as_mut())?;
if inner_rows.is_empty() {
let _ = self.new_vars; continue;
}
self.pending_outer = Some(outer_row);
self.pending = inner_rows.into_iter();
}
}
}