From 087816802b551e2b97b89ea07038b2c85494750d Mon Sep 17 00:00:00 2001 From: meh Date: Fri, 7 Aug 2026 12:51:36 +0700 Subject: [PATCH] release: prepare 1.0.0-beta.8 --- Cargo.toml | 4 +- codegen/Cargo.toml | 2 +- src/persistence/operation/batch.rs | 135 +++++++++++-------- tests/persistence/mod.rs | 1 + tests/persistence/multi_row_backend_order.rs | 108 +++++++++++++++ 5 files changed, 192 insertions(+), 58 deletions(-) create mode 100644 tests/persistence/multi_row_backend_order.rs diff --git a/Cargo.toml b/Cargo.toml index 7fa0570..f92e8cc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,7 +3,7 @@ members = ["codegen", "examples", "performance_measurement", "performance_measur [package] name = "worktable" -version = "1.0.0-beta.7" +version = "1.0.0-beta.8" edition = "2024" authors = ["Handy-caT"] license = "MIT" @@ -66,7 +66,7 @@ tracing = "0.1" url = { version = "2", optional = true } uuid = { version = "1.24.0", features = ["v4", "v7"] } walkdir = { version = "2", optional = true } -worktable_codegen = { path = "codegen", version = "=1.0.0-beta.7" } +worktable_codegen = { path = "codegen", version = "=1.0.0-beta.8" } [dev-dependencies] chrono = "0.4.43" diff --git a/codegen/Cargo.toml b/codegen/Cargo.toml index 759dcf3..45016bf 100644 --- a/codegen/Cargo.toml +++ b/codegen/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "worktable_codegen" -version = "1.0.0-beta.7" +version = "1.0.0-beta.8" edition = "2024" license = "MIT" description = "Proc-macro companion crate for worktable: the worktable! macro and its derives." diff --git a/src/persistence/operation/batch.rs b/src/persistence/operation/batch.rs index 5cbbfc0..1fb7d7e 100644 --- a/src/persistence/operation/batch.rs +++ b/src/persistence/operation/batch.rs @@ -15,6 +15,9 @@ use crate::persistence::task::{LastEventIds, QueueInnerRow}; use crate::prelude::*; use crate::prelude::{From, Order, SelectQueryExecutor}; +// Ephemeral metadata rebuilt for every persistence batch, not a persisted +// schema. One Multi operation deliberately owns several rows, so operation_id +// is non-unique while pos is the unique association back to the ops vector. worktable! ( name: BatchInner, columns: { @@ -80,43 +83,51 @@ fn latest_data_writes( ops: &[Operation], ) -> BatchData { 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; - }; - let link = op.link(); - let operation_id = op.operation_id(); - let key = (link.page_id, link.offset); - match latest.entry(key) { - std::collections::hash_map::Entry::Occupied(mut entry) => { - if (operation_id, sequence) > (entry.get().0, entry.get().1) { - entry.insert((operation_id, sequence, link, bytes.to_vec())); - } + fn collect_in_order( + ops: &[Operation], + order: impl Iterator + Clone, + ) -> BatchData { + let mut latest: HashMap = HashMap::with_capacity(ops.len()); + for sequence in order.clone() { + let op = &ops[sequence]; + if op.bytes().is_some() { + let link = op.link(); + latest.insert((link.page_id, link.offset), sequence); } - std::collections::hash_map::Entry::Vacant(entry) => { - entry.insert((operation_id, sequence, link, bytes.to_vec())); + } + + let mut ordered = HashMap::new(); + for sequence in order { + let op = &ops[sequence]; + let Some(bytes) = op.bytes() else { + continue; + }; + let link = op.link(); + if latest.get(&(link.page_id, link.offset)) != Some(&sequence) { + continue; } + ordered + .entry(link.page_id) + .or_insert_with(Vec::new) + .push((link, bytes.to_vec())); } + ordered } - let mut ordered = HashMap::new(); - for (_, (operation_id, sequence, link, bytes)) in latest { - ordered - .entry(link.page_id) - .or_insert_with(Vec::new) - .push((operation_id, sequence, link, bytes)); + // The analyzer already establishes this order. Keep that production path + // linear; only defensive callers that construct an unsorted BatchOperation + // pay for an index sort. + if ops + .windows(2) + .all(|pair| pair[0].operation_id() <= pair[1].operation_id()) + { + collect_in_order(ops, 0..ops.len()) + } else { + let mut order = (0..ops.len()).collect::>(); + order.sort_unstable_by_key(|sequence| (ops[*sequence].operation_id(), *sequence)); + collect_in_order(ops, order.into_iter()) } - ordered - .into_iter() - .map(|(page_id, mut writes)| { - 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() } #[derive(Debug)] @@ -151,6 +162,30 @@ where phantom_data: PhantomData, } } + + /// Remove metadata immediately after `self.ops.remove(removed_pos)`. + /// + /// At entry, `self.ops.len()` is already one shorter while `info_wt` still + /// has the old positions. Shifting upward positions in ascending order + /// keeps every unique destination vacant as it is filled. + async fn remove_info_at_pos(&self, removed_pos: usize) -> eyre::Result<()> { + let row = self + .info_wt + .select_by_pos(removed_pos) + .ok_or_else(|| eyre::eyre!("batch metadata position {removed_pos} is missing"))?; + self.info_wt.delete_without_lock::<_>(row.id).await?; + + for old_pos in (removed_pos + 1)..=self.ops.len() { + let row = self + .info_wt + .select_by_pos(old_pos) + .ok_or_else(|| eyre::eyre!("batch metadata position {old_pos} is missing during reindex"))?; + self.info_wt + .update_pos_by_id(PosByIdQuery { pos: old_pos - 1 }, row.id) + .await?; + } + Ok(()) + } } impl @@ -168,7 +203,7 @@ where async fn remove_operations_from_events( &mut self, invalid_events: PreparedIndexEvents, - ) -> Vec> { + ) -> eyre::Result>> { let mut removed_ops = Vec::new(); for ev in &invalid_events.primary_evs { @@ -186,7 +221,7 @@ where }) { let pos = self.ops.len() - (operation_pos_rev + 1); let op = self.ops.remove(pos); - self.remove_info_at_pos(pos).await; + self.remove_info_at_pos(pos).await?; removed_ops.push(op); } } @@ -198,7 +233,7 @@ where }) { let pos = self.ops.len() - (operation_pos_rev + 1); let op = self.ops.remove(pos); - self.remove_info_at_pos(pos).await; + self.remove_info_at_pos(pos).await?; removed_ops.push(op); }; // else it was already removed with primary @@ -222,26 +257,7 @@ where prepared_evs.secondary_evs.remove(op_secondary); } - 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(); - } + Ok(removed_ops) } pub fn get_last_event_ids(&self) -> LastEventIds { @@ -304,7 +320,7 @@ where primary_evs: primary_invalid_events, secondary_evs: secondary_invalid_events, }; - let ops = self.remove_operations_from_events(events_to_remove).await; + let ops = self.remove_operations_from_events(events_to_remove).await?; ops_to_remove.extend(ops); } @@ -446,7 +462,7 @@ mod tests { use data_bucket::Link; use uuid::Uuid; - use super::latest_data_writes; + use super::{BatchOperation, latest_data_writes}; use crate::persistence::operation::{InsertOperation, Operation, OperationId}; fn insert(id: u128, link: Link, bytes: Vec) -> Operation<(), u64, ()> { @@ -557,4 +573,13 @@ mod tests { assert_eq!(batch.get(&1.into()).unwrap(), &vec![(new_link, vec![2; 6])]); } + + #[tokio::test] + async fn missing_batch_metadata_returns_an_error_instead_of_panicking() { + let batch: BatchOperation<(), u64, (), ()> = BatchOperation::new(vec![], Default::default()); + + let error = batch.remove_info_at_pos(0).await.unwrap_err(); + + assert!(error.to_string().contains("batch metadata position 0 is missing")); + } } diff --git a/tests/persistence/mod.rs b/tests/persistence/mod.rs index 7dc62be..ceb4ef2 100644 --- a/tests/persistence/mod.rs +++ b/tests/persistence/mod.rs @@ -8,6 +8,7 @@ mod duplicate_key_index_reload; mod failure; mod index_page; mod loaded_index_growth; +mod multi_row_backend_order; mod read; mod recovery_load; mod schema; diff --git a/tests/persistence/multi_row_backend_order.rs b/tests/persistence/multi_row_backend_order.rs new file mode 100644 index 0000000..657404e --- /dev/null +++ b/tests/persistence/multi_row_backend_order.rs @@ -0,0 +1,108 @@ +use worktable::prelude::*; +use worktable::worktable; + +use crate::remove_dir_if_exists; + +macro_rules! persisted_multi_row_backend_case { + ( + $module:ident, + $name:ident, + $table:ident, + $row:ident, + $engine:ident, + $backend:ident, + $root:literal + ) => { + mod $module { + use super::*; + + worktable! { + name: $name, + persist: true, + columns: { + id: u64 primary_key autoincrement using $backend, + group_id: u64, + payload: String, + }, + indexes: { + group_idx: group_id, + }, + queries: { + update: { + PayloadByGroup(payload) by group_id, + } + } + } + + #[tokio::test] + async fn multi_row_relocation_survives_persisted_reload() { + const ROOT: &str = $root; + remove_dir_if_exists(ROOT.to_string()).await; + let config = DiskConfig::new_with_table_name(ROOT, $table::name_snake_case(), $table::version()); + + let engine = $engine::new(config.clone()).await.unwrap(); + let table = $table::load(engine).await.unwrap(); + for length in 1..=128 { + table + .insert($row { + id: table.get_next_pk().into(), + group_id: 7, + payload: "x".repeat(length), + }) + .unwrap(); + } + table.wait_for_ops().await.unwrap(); + + let replacement = "new-payload".repeat(64); + table + .update_payload_by_group( + PayloadByGroupQuery { + payload: replacement.clone(), + }, + 7, + ) + .await + .unwrap(); + table.wait_for_ops().await.unwrap(); + drop(table); + + let engine = $engine::new(config).await.unwrap(); + let table = $table::load(engine).await.unwrap(); + let rows = table.select_by_group_id(7).execute().unwrap(); + assert_eq!(rows.len(), 128); + assert!(rows.iter().all(|row| row.payload == replacement)); + table.wait_for_ops().await.unwrap(); + drop(table); + remove_dir_if_exists(ROOT.to_string()).await; + } + } + }; +} + +persisted_multi_row_backend_case!( + wti, + MultiRowWti, + MultiRowWtiWorkTable, + MultiRowWtiRow, + MultiRowWtiPersistenceEngine, + worktables_index, + "tests/data/multi_row_backend_wti" +); +persisted_multi_row_backend_case!( + congee, + MultiRowCongee, + MultiRowCongeeWorkTable, + MultiRowCongeeRow, + MultiRowCongeePersistenceEngine, + congee, + "tests/data/multi_row_backend_congee" +); +persisted_multi_row_backend_case!( + arctic, + MultiRowArctic, + MultiRowArcticWorkTable, + MultiRowArcticRow, + MultiRowArcticPersistenceEngine, + arctic, + "tests/data/multi_row_backend_arctic" +);