From 1ffc186f59f25a5cb5f2fe3f716fce24cf32dbfc Mon Sep 17 00:00:00 2001 From: meh Date: Fri, 7 Aug 2026 11:49:49 +0700 Subject: [PATCH 1/2] fix: order overlapping durable row writes --- src/persistence/operation/batch.rs | 55 ++++++++++++++++++++++++------ 1 file changed, 44 insertions(+), 11 deletions(-) diff --git a/src/persistence/operation/batch.rs b/src/persistence/operation/batch.rs index 02d4dc8..abdf447 100644 --- a/src/persistence/operation/batch.rs +++ b/src/persistence/operation/batch.rs @@ -66,16 +66,18 @@ impl From for BatchInnerRow { } } -/// Coalesces durable row writes by physical storage slot. +/// Coalesces durable row writes by physical storage slot and preserves their +/// creation order. /// /// `Link::length` can change when an unsized row is reinserted into a reused /// `(page_id, offset)`. Treating the two lengths as different keys leaves -/// overlapping writes in the same batch, whose eventual application order is -/// derived from a hash map. The newest operation must be the only write for a -/// physical slot. WorkTable-generated operation IDs use `Uuid::now_v7`, whose -/// shared process context guarantees creation-order sorting even within one -/// millisecond; callers constructing `Operation` values manually must preserve -/// that ordering contract. +/// overlapping writes in the same batch. The newest operation must be the only +/// write for an identical physical start, and writes at different starts must +/// still be applied oldest-to-newest: range splitting can make them overlap. +/// WorkTable-generated operation IDs use `Uuid::now_v7`, whose shared process +/// context guarantees creation-order sorting even within one millisecond; +/// callers constructing `Operation` values manually must preserve that +/// ordering contract. fn latest_data_writes( ops: &[Operation], ) -> BatchData { @@ -99,11 +101,21 @@ fn latest_data_writes( } } - let mut data = HashMap::new(); - for (_, (_, link, bytes)) in latest { - data.entry(link.page_id).or_insert_with(Vec::new).push((link, bytes)); + let mut ordered = HashMap::new(); + for (_, (operation_id, link, bytes)) in latest { + ordered + .entry(link.page_id) + .or_insert_with(Vec::new) + .push((operation_id, link, bytes)); } - data + ordered + .into_iter() + .map(|(page_id, mut writes)| { + writes.sort_unstable_by_key(|(operation_id, _, _)| *operation_id); + let writes = writes.into_iter().map(|(_, link, bytes)| (link, bytes)).collect(); + (page_id, writes) + }) + .collect() } #[derive(Debug)] @@ -455,4 +467,25 @@ mod tests { assert_eq!(writes, &vec![(new_link, vec![2; 6])]); } + + #[test] + fn overlapping_reused_ranges_remain_in_creation_order() { + let older_link = Link { + page_id: 1.into(), + offset: 128, + length: 8, + }; + let newer_link = Link { + page_id: 1.into(), + offset: 132, + length: 8, + }; + + for _ in 0..128 { + let batch = latest_data_writes(&[insert(2, newer_link, vec![2; 8]), insert(1, older_link, vec![1; 8])]); + let writes = batch.get(&1.into()).unwrap(); + + assert_eq!(writes, &vec![(older_link, vec![1; 8]), (newer_link, vec![2; 8])]); + } + } } From 9aee1a92f15a1df468ffd592cf906e3f86367f3b Mon Sep 17 00:00:00 2001 From: meh Date: Fri, 7 Aug 2026 12:23:34 +0700 Subject: [PATCH 2/2] fix: preserve multi-row persistence order --- src/persistence/operation/batch.rs | 137 ++++++++++++++++++++++------- src/persistence/operation/mod.rs | 2 +- src/persistence/task.rs | 87 +++++++++++++++--- 3 files changed, 177 insertions(+), 49 deletions(-) diff --git a/src/persistence/operation/batch.rs b/src/persistence/operation/batch.rs index abdf447..5cbbfc0 100644 --- a/src/persistence/operation/batch.rs +++ b/src/persistence/operation/batch.rs @@ -1,4 +1,4 @@ -use std::collections::{HashMap, HashSet}; +use std::collections::HashMap; use std::fmt::Debug; use std::hash::Hash; use std::marker::PhantomData; @@ -26,17 +26,15 @@ worktable! ( pos: usize, }, indexes: { - operation_id_idx: operation_id unique, + operation_id_idx: operation_id, page_id_idx: page_id, link_idx: link, op_type_idx: op_type, + pos_idx: pos unique, }, queries: { update: { - PosByOpId(pos) by operation_id, - }, - delete: { - ByOpId() by operation_id, + PosById(pos) by id, } } ); @@ -81,8 +79,11 @@ impl From for BatchInnerRow { fn latest_data_writes( ops: &[Operation], ) -> BatchData { - let mut latest: HashMap<(PageId, u32), (OperationId, Link, Vec)> = HashMap::new(); - for op in ops { + type PhysicalSlot = (PageId, u32); + type SequencedWrite = (OperationId, usize, Link, Vec); + + let mut latest: HashMap = HashMap::new(); + for (sequence, op) in ops.iter().enumerate() { let Some(bytes) = op.bytes() else { continue; }; @@ -91,28 +92,28 @@ fn latest_data_writes( let key = (link.page_id, link.offset); match latest.entry(key) { std::collections::hash_map::Entry::Occupied(mut entry) => { - if operation_id > entry.get().0 { - entry.insert((operation_id, link, bytes.to_vec())); + if (operation_id, sequence) > (entry.get().0, entry.get().1) { + entry.insert((operation_id, sequence, link, bytes.to_vec())); } } std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert((operation_id, link, bytes.to_vec())); + entry.insert((operation_id, sequence, link, bytes.to_vec())); } } } let mut ordered = HashMap::new(); - for (_, (operation_id, link, bytes)) in latest { + for (_, (operation_id, sequence, link, bytes)) in latest { ordered .entry(link.page_id) .or_insert_with(Vec::new) - .push((operation_id, link, bytes)); + .push((operation_id, sequence, link, bytes)); } ordered .into_iter() .map(|(page_id, mut writes)| { - writes.sort_unstable_by_key(|(operation_id, _, _)| *operation_id); - let writes = writes.into_iter().map(|(_, link, bytes)| (link, bytes)).collect(); + writes.sort_unstable_by_key(|(operation_id, sequence, _, _)| (*operation_id, *sequence)); + let writes = writes.into_iter().map(|(_, _, link, bytes)| (link, bytes)).collect(); (page_id, writes) }) .collect() @@ -167,8 +168,8 @@ where async fn remove_operations_from_events( &mut self, invalid_events: PreparedIndexEvents, - ) -> HashSet> { - let mut removed_ops = HashSet::new(); + ) -> Vec> { + let mut removed_ops = Vec::new(); for ev in &invalid_events.primary_evs { if let Some(operation_pos_rev) = self.ops.iter().rev().position(|op| { @@ -183,27 +184,26 @@ where false } }) { - let op = self.ops.remove(self.ops.len() - (operation_pos_rev + 1)); - removed_ops.insert(op); + let pos = self.ops.len() - (operation_pos_rev + 1); + let op = self.ops.remove(pos); + self.remove_info_at_pos(pos).await; + removed_ops.push(op); } } - for (index, id) in invalid_events.secondary_evs.iter_event_ids() { + let secondary_event_ids = invalid_events.secondary_evs.iter_event_ids().collect::>(); + for (index, id) in secondary_event_ids { if let Some(operation_pos_rev) = self.ops.iter().rev().position(|op| { let evs = op.secondary_key_events(); evs.contains_event(index, id) }) { - let op = self.ops.remove(self.ops.len() - (operation_pos_rev + 1)); - removed_ops.insert(op); + let pos = self.ops.len() - (operation_pos_rev + 1); + let op = self.ops.remove(pos); + self.remove_info_at_pos(pos).await; + removed_ops.push(op); }; // else it was already removed with primary } for op in &removed_ops { - let pk = self - .info_wt - .select_by_operation_id(op.operation_id()) - .expect("exists as all should be inserted on prepare step") - .id; - self.info_wt.delete_without_lock::<_>(pk).await.unwrap(); let prepared_evs = self .prepared_index_evs .as_mut() @@ -225,6 +225,25 @@ where removed_ops } + async fn remove_info_at_pos(&self, removed_pos: usize) { + let row = self + .info_wt + .select_by_pos(removed_pos) + .expect("batch metadata position must exist"); + self.info_wt.delete_without_lock::<_>(row.id).await.unwrap(); + + for old_pos in (removed_pos + 1)..=self.ops.len() { + let row = self + .info_wt + .select_by_pos(old_pos) + .expect("later batch metadata position must exist"); + self.info_wt + .update_pos_by_id(PosByIdQuery { pos: old_pos - 1 }, row.id) + .await + .unwrap(); + } + } + pub fn get_last_event_ids(&self) -> LastEventIds { let prepared_evs = self .prepared_index_evs @@ -351,12 +370,6 @@ where } } - for (pos, op) in self.ops.iter().enumerate() { - let op_id = op.operation_id(); - let q = PosByOpIdQuery { pos }; - self.info_wt.update_pos_by_op_id(q, op_id).await? - } - Ok(Some(ops_to_remove)) } @@ -447,6 +460,17 @@ mod tests { }) } + fn multi_insert(id: u128, link: Link, bytes: Vec) -> Operation<(), u64, ()> { + Operation::Insert(InsertOperation { + id: OperationId::Multi(Uuid::from_u128(id)), + primary_key_events: vec![], + secondary_keys_events: (), + pk_gen_state: (), + bytes, + link, + }) + } + #[test] fn variable_length_link_reuse_keeps_only_the_newest_physical_write() { let old_link = Link { @@ -488,4 +512,49 @@ mod tests { assert_eq!(writes, &vec![(older_link, vec![1; 8]), (newer_link, vec![2; 8])]); } } + + #[test] + fn equal_id_overlapping_writes_preserve_batch_order() { + let older_link = Link { + page_id: 1.into(), + offset: 128, + length: 8, + }; + let newer_link = Link { + page_id: 1.into(), + offset: 132, + length: 8, + }; + + let batch = latest_data_writes(&[ + multi_insert(1, older_link, vec![1; 8]), + multi_insert(1, newer_link, vec![2; 8]), + ]); + + assert_eq!( + batch.get(&1.into()).unwrap(), + &vec![(older_link, vec![1; 8]), (newer_link, vec![2; 8])] + ); + } + + #[test] + fn equal_id_physical_slot_reuse_keeps_later_batch_write() { + let old_link = Link { + page_id: 1.into(), + offset: 128, + length: 4, + }; + let new_link = Link { + page_id: 1.into(), + offset: 128, + length: 6, + }; + + let batch = latest_data_writes(&[ + multi_insert(1, old_link, vec![1; 4]), + multi_insert(1, new_link, vec![2; 6]), + ]); + + assert_eq!(batch.get(&1.into()).unwrap(), &vec![(new_link, vec![2; 6])]); + } } diff --git a/src/persistence/operation/mod.rs b/src/persistence/operation/mod.rs index c933948..5f7954f 100644 --- a/src/persistence/operation/mod.rs +++ b/src/persistence/operation/mod.rs @@ -14,7 +14,7 @@ use uuid::Uuid; use crate::prelude::From; -pub use batch::{BatchInnerRow, BatchInnerWorkTable, BatchOperation, PosByOpIdQuery}; +pub use batch::{BatchInnerRow, BatchInnerWorkTable, BatchOperation}; pub use operation::{AcknowledgeOperation, DeleteOperation, InsertOperation, Operation, UpdateOperation}; pub use util::validate_events; diff --git a/src/persistence/task.rs b/src/persistence/task.rs index a80a3a7..f26a5c2 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -12,7 +12,7 @@ use tokio::sync::Notify; use tokio::task::JoinHandle; use worktable_codegen::worktable; -use crate::persistence::operation::{BatchInnerRow, BatchInnerWorkTable, BatchOperation, OperationId, PosByOpIdQuery}; +use crate::persistence::operation::{BatchInnerRow, BatchInnerWorkTable, BatchOperation, OperationId}; use crate::persistence::{ PersistenceEngine, PersistenceError, PersistenceIndexCorruption, PersistenceResult, PersistenceState, }; @@ -288,27 +288,33 @@ where ops_pos_set.extend(rows.into_iter().map(|r| (r.pos, r.id))) } - let mut ops = Vec::with_capacity(ops_pos_set.len()); + // Queue row IDs are monotonic insertion sequence numbers. Restore that + // sequence after HashSet collection so operations sharing one Multi ID + // remain in their original order through the stable OperationId sort. + let mut ops_positions = ops_pos_set.into_iter().collect::>(); + ops_positions.sort_unstable_by_key(|(_, id)| *id); + + let mut queued_ops = Vec::with_capacity(ops_positions.len()); let info_wt = BatchInnerWorkTable::default(); - for (pos, id) in ops_pos_set { - let mut row: BatchInnerRow = self.queue_inner_wt.select(id).expect("exists as Id exists").into(); + for (pos, id) in ops_positions { + let row: BatchInnerRow = self.queue_inner_wt.select(id).expect("exists as Id exists").into(); let op = self .operations .remove(pos) .expect("should be available as presented in table"); - row.pos = ops.len(); - row.op_type = op.operation_type(); - ops.push(op); - info_wt.insert(row)?; + queued_ops.push((op, row)); self.queue_inner_wt.delete_without_lock::<_>(id).await? } // println!("New wt generated {:?}", start.elapsed()); - // return ops sorted by `OperationId` - ops.sort_by_key(|k| k.operation_id()); - for (pos, op) in ops.iter().enumerate() { - let op_id = op.operation_id(); - let q = PosByOpIdQuery { pos }; - info_wt.update_pos_by_op_id(q, op_id).await?; + // The sort is stable, so queue creation order breaks ties for rows + // produced by one multi-row operation. + queued_ops.sort_by_key(|(op, _)| op.operation_id()); + let mut ops = Vec::with_capacity(queued_ops.len()); + for (pos, (op, mut row)) in queued_ops.into_iter().enumerate() { + row.pos = pos; + row.op_type = op.operation_type(); + info_wt.insert(row)?; + ops.push(op); } let mut op = BatchOperation::new(ops, info_wt); @@ -472,6 +478,59 @@ mod lifecycle_tests { }) } + fn multi_insert_operation(id: u128, offset: u32, byte: u8) -> Operation<(), u64, TestEvents> { + Operation::Insert(InsertOperation { + id: OperationId::Multi(uuid::Uuid::from_u128(id)), + pk_gen_state: (), + primary_key_events: vec![], + secondary_keys_events: TestEvents, + bytes: vec![byte; 8], + link: Link { + page_id: 1.into(), + offset, + length: 8, + }, + }) + } + + #[tokio::test] + async fn analyzer_preserves_queue_order_for_operations_sharing_a_multi_id() { + let queue_inner_wt = Arc::new(QueueInnerWorkTable::default()); + let mut analyzer: QueueAnalyzer<(), u64, TestEvents, TestIndex> = QueueAnalyzer::new(queue_inner_wt); + analyzer.push(multi_insert_operation(1, 128, 1)).unwrap(); + analyzer.push(multi_insert_operation(1, 132, 2)).unwrap(); + + let batch = analyzer + .collect_batch_from_op_id(OperationId::Multi(uuid::Uuid::from_u128(1))) + .await + .unwrap() + .unwrap() + .get_batch_data_op() + .unwrap(); + + assert_eq!( + batch.get(&1.into()).unwrap(), + &vec![ + ( + Link { + page_id: 1.into(), + offset: 128, + length: 8, + }, + vec![1; 8], + ), + ( + Link { + page_id: 1.into(), + offset: 132, + length: 8, + }, + vec![2; 8], + ), + ] + ); + } + #[tokio::test] async fn close_drains_and_joins_the_engine() { let batches = Arc::new(AtomicUsize::new(0));