Skip to main content

lora_executor/pull/
source.rs

1use crate::errors::ExecResult;
2use crate::value::Row;
3
4/// Fallible pull-based row cursor.
5///
6/// Each call to [`RowSource::next_row`] returns the next row,
7/// `Ok(None)` when the cursor is exhausted, or an error if execution
8/// fails. The cursor stays in a valid state after an error — callers
9/// may drop it without observing additional side effects.
10pub trait RowSource {
11    /// Pull the next row.
12    fn next_row(&mut self) -> ExecResult<Option<Row>>;
13}
14
15/// Drain a row source into a `Vec<Row>`, propagating the first error.
16pub fn drain<S: RowSource + ?Sized>(source: &mut S) -> ExecResult<Vec<Row>> {
17    let mut out = Vec::new();
18    while let Some(row) = source.next_row()? {
19        out.push(row);
20    }
21    Ok(out)
22}
23
24/// Buffered cursor backed by a pre-computed `Vec<Row>`. Used both as
25/// a simple "rows already collected" adapter and as the leaf fallback
26/// for operators whose internals still require full materialization.
27pub struct BufferedRowSource {
28    iter: std::vec::IntoIter<Row>,
29}
30
31impl BufferedRowSource {
32    pub fn new(rows: Vec<Row>) -> Self {
33        Self {
34            iter: rows.into_iter(),
35        }
36    }
37}
38
39impl RowSource for BufferedRowSource {
40    fn next_row(&mut self) -> ExecResult<Option<Row>> {
41        Ok(self.iter.next())
42    }
43}
44
45/// Yields a single empty row exactly once. The bottom of every plan
46/// chain that doesn't start with an explicit input.
47pub struct ArgumentSource {
48    yielded: bool,
49}
50
51impl ArgumentSource {
52    pub fn new() -> Self {
53        Self { yielded: false }
54    }
55}
56
57impl Default for ArgumentSource {
58    fn default() -> Self {
59        Self::new()
60    }
61}
62
63impl RowSource for ArgumentSource {
64    fn next_row(&mut self) -> ExecResult<Option<Row>> {
65        if self.yielded {
66            Ok(None)
67        } else {
68            self.yielded = true;
69            Ok(Some(Row::new()))
70        }
71    }
72}
73
74/// Enforces a query deadline inside a pull pipeline. Wrapped around every
75/// source when the query has a deadline, so even a loop that never yields
76/// a row (a filter rejecting a cartesian product, say) keeps checking:
77/// each pull from upstream passes through a wrapper. The clock is read
78/// every 64 pulls to keep the overhead negligible.
79pub(crate) struct DeadlineSource<'a> {
80    inner: Box<dyn RowSource + 'a>,
81    deadline: web_time::Instant,
82    tick: u32,
83}
84
85impl<'a> DeadlineSource<'a> {
86    pub(crate) fn wrap(
87        inner: Box<dyn RowSource + 'a>,
88        deadline: Option<web_time::Instant>,
89    ) -> Box<dyn RowSource + 'a> {
90        match deadline {
91            Some(deadline) => Box::new(Self {
92                inner,
93                deadline,
94                tick: 0,
95            }),
96            None => inner,
97        }
98    }
99}
100
101impl RowSource for DeadlineSource<'_> {
102    fn next_row(&mut self) -> ExecResult<Option<Row>> {
103        self.tick = self.tick.wrapping_add(1);
104        if self.tick.is_multiple_of(64) && crate::cancel::deadline_reached(self.deadline) {
105            return Err(crate::errors::ExecutorError::QueryTimeout);
106        }
107        self.inner.next_row()
108    }
109}