nodedb_columnar/mutation/
engine.rs1use std::collections::HashMap;
6
7use nodedb_types::columnar::ColumnarSchema;
8use nodedb_types::surrogate::Surrogate;
9use nodedb_types::value::Value;
10
11use crate::delete_bitmap::DeleteBitmap;
12use crate::error::ColumnarError;
13use crate::memtable::ColumnarMemtable;
14use crate::pk_index::{PkIndex, encode_pk};
15use crate::wal_record::ColumnarWalRecord;
16
17pub struct MutationEngine {
22 pub(super) collection: String,
23 pub(super) schema: ColumnarSchema,
24 pub(super) memtable: ColumnarMemtable,
25 pub(super) pk_index: PkIndex,
26 pub(super) delete_bitmaps: HashMap<u64, DeleteBitmap>,
28 pub(super) pk_col_indices: Vec<usize>,
30 pub(super) next_segment_id: u64,
32 pub(super) memtable_segment_id: u64,
35 pub(super) memtable_row_counter: u32,
37 pub(super) memtable_surrogates: Vec<Option<Surrogate>>,
43}
44
45#[derive(Debug)]
47pub struct MutationResult {
48 pub wal_records: Vec<ColumnarWalRecord>,
50}
51
52impl MutationEngine {
53 pub fn new(collection: String, schema: ColumnarSchema) -> Self {
55 Self::with_flush_threshold(collection, schema, crate::memtable::DEFAULT_FLUSH_THRESHOLD)
56 }
57
58 pub fn with_flush_threshold(
64 collection: String,
65 schema: ColumnarSchema,
66 flush_threshold: usize,
67 ) -> Self {
68 let pk_col_indices: Vec<usize> = schema
69 .columns
70 .iter()
71 .enumerate()
72 .filter(|(_, c)| c.primary_key)
73 .map(|(i, _)| i)
74 .collect();
75
76 let memtable = ColumnarMemtable::with_threshold(&schema, flush_threshold.max(1));
77 let memtable_segment_id = 0;
79
80 Self {
81 collection,
82 schema,
83 memtable,
84 pk_index: PkIndex::new(),
85 delete_bitmaps: HashMap::new(),
86 pk_col_indices,
87 next_segment_id: 1,
88 memtable_segment_id,
89 memtable_row_counter: 0,
90 memtable_surrogates: Vec::new(),
91 }
92 }
93
94 pub fn memtable(&self) -> &ColumnarMemtable {
98 &self.memtable
99 }
100
101 pub fn memtable_mut(&mut self) -> &mut ColumnarMemtable {
103 &mut self.memtable
104 }
105
106 pub fn pk_index(&self) -> &PkIndex {
108 &self.pk_index
109 }
110
111 pub fn pk_index_mut(&mut self) -> &mut PkIndex {
113 &mut self.pk_index
114 }
115
116 pub fn delete_bitmap(&self, segment_id: u64) -> Option<&DeleteBitmap> {
118 self.delete_bitmaps.get(&segment_id)
119 }
120
121 pub fn delete_bitmap_mut(&mut self, segment_id: u64) -> &mut DeleteBitmap {
126 self.delete_bitmaps.entry(segment_id).or_default()
127 }
128
129 pub fn memtable_segment_id(&self) -> u64 {
131 self.memtable_segment_id
132 }
133
134 pub fn pk_col_indices(&self) -> &[usize] {
136 &self.pk_col_indices
137 }
138
139 pub fn delete_bitmaps(&self) -> &HashMap<u64, DeleteBitmap> {
141 &self.delete_bitmaps
142 }
143
144 pub fn collection(&self) -> &str {
146 &self.collection
147 }
148
149 pub fn schema(&self) -> &ColumnarSchema {
151 &self.schema
152 }
153
154 pub fn should_flush(&self) -> bool {
156 self.memtable.should_flush()
157 }
158
159 pub fn memtable_surrogates(&self) -> &[Option<Surrogate>] {
164 &self.memtable_surrogates
165 }
166
167 pub fn scan_memtable_rows(&self) -> impl Iterator<Item = Vec<Value>> + '_ {
172 let deletes = self.delete_bitmaps.get(&self.memtable_segment_id);
173 self.memtable
174 .iter_rows()
175 .enumerate()
176 .filter_map(move |(row_idx, row)| {
177 if deletes.is_some_and(|bm| bm.is_deleted(row_idx as u32)) {
178 None
179 } else {
180 Some(row)
181 }
182 })
183 }
184
185 pub fn scan_memtable_rows_with_surrogates(
191 &self,
192 ) -> impl Iterator<Item = (Option<Surrogate>, Vec<Value>)> + '_ {
193 let deletes = self.delete_bitmaps.get(&self.memtable_segment_id);
194 let surrogates = &self.memtable_surrogates;
195 self.memtable
196 .iter_rows()
197 .enumerate()
198 .filter_map(move |(row_idx, row)| {
199 if deletes.is_some_and(|bm| bm.is_deleted(row_idx as u32)) {
200 return None;
201 }
202 let surrogate = surrogates.get(row_idx).copied().flatten();
203 Some((surrogate, row))
204 })
205 }
206
207 pub fn get_memtable_row(&self, row_idx: usize) -> Option<Vec<Value>> {
209 if self
210 .delete_bitmaps
211 .get(&self.memtable_segment_id)
212 .is_some_and(|bm| bm.is_deleted(row_idx as u32))
213 {
214 return None;
215 }
216 self.memtable.get_row(row_idx)
217 }
218
219 pub fn rollback_memtable_inserts(
230 &mut self,
231 row_count_before: usize,
232 inserted_pks: &[Vec<u8>],
233 displaced: &[(Vec<u8>, crate::pk_index::RowLocation)],
234 ) {
235 for pk in inserted_pks {
237 self.pk_index.remove(pk);
238 }
239 for (pk, prior_location) in displaced {
241 self.pk_index.upsert(pk.clone(), *prior_location);
242 if let Some(bm) = self.delete_bitmaps.get_mut(&prior_location.segment_id) {
244 bm.unmark_deleted(prior_location.row_index);
245 }
246 }
247 self.memtable.truncate_to(row_count_before);
249 self.memtable_surrogates.truncate(row_count_before);
250 self.memtable_row_counter = row_count_before as u32;
251 }
252
253 pub fn restore_deleted_rows(&mut self, rows: &[(Vec<u8>, crate::pk_index::RowLocation)]) {
263 for (pk_bytes, location) in rows {
264 if let Some(bm) = self.delete_bitmaps.get_mut(&location.segment_id) {
265 bm.unmark_deleted(location.row_index);
266 }
267 self.pk_index.upsert(pk_bytes.clone(), *location);
268 }
269 }
270
271 pub fn next_segment_id(&self) -> u64 {
275 self.next_segment_id
276 }
277
278 pub fn should_compact(&self, segment_id: u64, total_rows: u64) -> bool {
280 self.delete_bitmaps
281 .get(&segment_id)
282 .is_some_and(|bm| bm.should_compact(total_rows, 0.2))
283 }
284
285 pub fn encode_pk_from_row(&self, values: &[Value]) -> Result<Vec<u8>, ColumnarError> {
288 self.extract_pk_bytes(values)
289 }
290
291 pub(super) fn extract_pk_bytes(&self, values: &[Value]) -> Result<Vec<u8>, ColumnarError> {
295 if values.len() != self.schema.columns.len() {
296 return Err(ColumnarError::SchemaMismatch {
297 expected: self.schema.columns.len(),
298 got: values.len(),
299 });
300 }
301
302 if self.pk_col_indices.len() == 1 {
303 Ok(encode_pk(&values[self.pk_col_indices[0]]))
304 } else {
305 let pk_values: Vec<&Value> = self.pk_col_indices.iter().map(|&i| &values[i]).collect();
306 Ok(crate::pk_index::encode_composite_pk(&pk_values))
307 }
308 }
309}
310
311#[cfg(test)]
312mod tests {
313 use nodedb_types::columnar::{ColumnDef, ColumnType, ColumnarSchema};
314
315 use super::*;
316
317 fn minimal_schema() -> ColumnarSchema {
318 ColumnarSchema {
319 columns: vec![ColumnDef::required("id", ColumnType::Int64).with_primary_key()],
320 version: 1,
321 }
322 }
323
324 #[test]
325 fn segment_id_allocator_returns_err_at_u64_max() {
326 let mut engine = MutationEngine::new("test".to_string(), minimal_schema());
327 engine.next_segment_id = u64::MAX;
328 let result = engine.on_memtable_flushed(u64::MAX - 1);
329 assert!(
330 matches!(result, Err(ColumnarError::SegmentIdExhausted)),
331 "expected SegmentIdExhausted, got: {result:?}"
332 );
333 }
334
335 #[test]
336 fn with_flush_threshold_triggers_should_flush_at_threshold() {
337 use nodedb_types::value::Value;
338
339 let schema = minimal_schema();
340 let mut engine = MutationEngine::with_flush_threshold("col".to_string(), schema, 2);
342 assert!(!engine.should_flush(), "empty engine should not need flush");
343
344 engine.insert(&[Value::Integer(1)]).expect("insert row 1");
345 assert!(!engine.should_flush(), "1 row below threshold of 2");
346
347 engine.insert(&[Value::Integer(2)]).expect("insert row 2");
348 assert!(
349 engine.should_flush(),
350 "2 rows at threshold of 2 must trigger flush"
351 );
352 }
353
354 #[test]
355 fn with_flush_threshold_zero_is_clamped_to_one() {
356 use nodedb_types::value::Value;
357
358 let schema = minimal_schema();
359 let mut engine = MutationEngine::with_flush_threshold("col".to_string(), schema, 0);
360 engine.insert(&[Value::Integer(1)]).expect("insert");
361 assert!(
363 engine.should_flush(),
364 "clamped-to-1 threshold: 1 row must flush"
365 );
366 }
367}