Skip to main content

alopex_sql/executor/query/
scan.rs

1use alopex_core::KVTransaction;
2use alopex_core::kv::KVStore;
3
4use crate::catalog::TableMetadata;
5use crate::executor::Result;
6use crate::storage::{
7    KeyEncoder, RangeBoundedScanIterator, SqlTransaction, SqlTxn, StorageRangeConstraint,
8    TableScanIterator,
9};
10
11use super::Row;
12use super::iterator::ScanIterator;
13
14/// Execute a table scan and return rows with RowIDs.
15pub fn execute_scan<'txn, S: KVStore + 'txn>(
16    txn: &mut impl SqlTxn<'txn, S>,
17    table_meta: &crate::catalog::TableMetadata,
18) -> Result<Vec<Row>> {
19    Ok(txn.with_table(table_meta, |storage| {
20        let iter = storage.range_scan(0, u64::MAX)?;
21        let mut rows = Vec::new();
22        for entry in iter {
23            let (row_id, values) = entry?;
24            rows.push(Row::new(row_id, values));
25        }
26        Ok(rows)
27    })?)
28}
29
30/// Create a streaming scan iterator for FR-7 compliance.
31///
32/// This function creates a `ScanIterator` that streams rows directly from
33/// the underlying storage without materializing all rows upfront.
34///
35/// # Lifetime
36///
37/// The returned iterator borrows from the transaction (`'a`), so the
38/// transaction must remain valid while the iterator is in use.
39pub fn create_scan_iterator<'a, 'txn: 'a, S: KVStore + 'txn, T: SqlTxn<'txn, S>>(
40    txn: &'a mut T,
41    table_meta: &TableMetadata,
42) -> Result<ScanIterator<'a>> {
43    let table_id = table_meta.table_id;
44    let prefix = KeyEncoder::table_prefix(table_id);
45    let inner = txn.inner_mut().scan_prefix(&prefix)?;
46    let table_scan_iter = TableScanIterator::new(inner, table_id);
47    Ok(ScanIterator::new(table_scan_iter, table_meta))
48}
49
50/// Creates the only scan iterator intended for a fenced remote range worker.
51///
52/// Unlike [`create_scan_iterator`], this entry point rejects an ordinary local
53/// transaction and always uses concrete `scan_range` bounds.  The caller has
54/// already pinned the catalog, schema, and index identities in `constraint`.
55pub fn create_fenced_range_scan_iterator<'a, 'txn: 'a, S: KVStore + 'txn>(
56    txn: &'a mut SqlTransaction<'txn, S>,
57    table_meta: &TableMetadata,
58    constraint: &StorageRangeConstraint,
59) -> Result<RangeBoundedScanIterator<'a>> {
60    constraint.validate_table(table_meta)?;
61    constraint.validate_read_at(txn.read_at_point())?;
62    let (lower, upper) = constraint.encoded_bounds();
63    let inner = txn.inner_mut().scan_range(lower, upper)?;
64    Ok(RangeBoundedScanIterator::new(
65        TableScanIterator::new(inner, constraint.table_id()),
66        constraint.clone(),
67    ))
68}
69
70/// Executes a materialized fenced range scan for a remote worker.
71///
72/// This function is intentionally separate from the legacy local scan.  It
73/// cannot broaden a worker into a whole-table prefix scan and checks every
74/// returned primary row key before returning it.
75pub fn execute_fenced_range_scan<'txn, S: KVStore + 'txn>(
76    txn: &mut SqlTransaction<'txn, S>,
77    table_meta: &TableMetadata,
78    constraint: &StorageRangeConstraint,
79) -> Result<Vec<Row>> {
80    let iter = create_fenced_range_scan_iterator(txn, table_meta, constraint)?;
81    iter.map(|entry| entry.map(|(row_id, values)| Row::new(row_id, values)))
82        .collect::<std::result::Result<Vec<_>, _>>()
83        .map_err(Into::into)
84}
85
86#[cfg(test)]
87mod tests {
88    use std::sync::Arc;
89
90    use alopex_core::ReadAtPoint;
91    use alopex_core::kv::{KVStore, memory::MemoryKV};
92    use alopex_core::types::TxnMode;
93
94    use super::*;
95    use crate::catalog::ColumnMetadata;
96    use crate::planner::types::ResolvedType;
97    use crate::storage::{RangeReadSnapshot, SqlValue, TxnBridge};
98
99    fn table() -> TableMetadata {
100        TableMetadata::new(
101            "users",
102            vec![ColumnMetadata::new("id", ResolvedType::Integer)],
103        )
104        .with_table_id(7)
105    }
106
107    fn constraint(point: ReadAtPoint) -> StorageRangeConstraint {
108        StorageRangeConstraint::new(
109            "range-a",
110            3,
111            alopex_core::RowKeyRange::new(7, Some(2), Some(4)).unwrap(),
112            RangeReadSnapshot::new(point, "schema-13").unwrap(),
113        )
114        .unwrap()
115    }
116
117    #[test]
118    fn fenced_scan_uses_half_open_range_and_rechecks_each_row() {
119        let store = Arc::new(MemoryKV::new());
120        let bridge = TxnBridge::new(store.clone());
121        let table = table();
122        let mut write = bridge.begin_write().unwrap();
123        write
124            .with_table(&table, |storage| {
125                for row_id in 1..=4 {
126                    storage.insert(row_id, &[SqlValue::Integer(row_id as i32)])?;
127                }
128                Ok(())
129            })
130            .unwrap();
131        write.commit().unwrap();
132
133        let point = ReadAtPoint::new(7, 11, 13, 17);
134        // `from_read_at` models the transaction a capable remote backend has
135        // already opened.  MemoryKV itself is deliberately not read-at capable.
136        let inner = store.begin(TxnMode::ReadOnly).unwrap();
137        let mut read = TxnBridge::<MemoryKV>::from_read_at(inner, point);
138        let rows = execute_fenced_range_scan(&mut read, &table, &constraint(point)).unwrap();
139
140        assert_eq!(
141            rows.into_iter().map(|row| row.row_id).collect::<Vec<_>>(),
142            vec![2, 3]
143        );
144    }
145
146    #[test]
147    fn fenced_scan_rejects_an_unfenced_local_transaction_before_scanning() {
148        let store = Arc::new(MemoryKV::new());
149        let bridge = TxnBridge::new(store);
150        let table = table();
151        let point = ReadAtPoint::new(7, 11, 13, 17);
152        let mut local_read = bridge.begin_read().unwrap();
153
154        let error =
155            execute_fenced_range_scan(&mut local_read, &table, &constraint(point)).unwrap_err();
156        assert!(error.to_string().contains("requires a transaction opened"));
157    }
158
159    #[test]
160    fn fenced_scan_rejects_a_catalog_or_index_fence_mismatch_before_scanning() {
161        let store = Arc::new(MemoryKV::new());
162        let table = table();
163        let expected = ReadAtPoint::new(7, 11, 13, 17);
164        let inner = store.begin(TxnMode::ReadOnly).unwrap();
165        let mut read = TxnBridge::<MemoryKV>::from_read_at(inner, ReadAtPoint::new(7, 11, 13, 18));
166
167        let error =
168            execute_fenced_range_scan(&mut read, &table, &constraint(expected)).unwrap_err();
169        assert!(error.to_string().contains("read-at fence mismatch"));
170    }
171}