alopex_sql/executor/query/
scan.rs1use 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
14pub 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
30pub 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
50pub 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
70pub 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 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}