From c633345e6eac000db28aab09210e25f7b79c0249 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 05:59:01 +0700 Subject: [PATCH 01/14] fix: preserve unsized page identity across batch splits --- src/persistence/space/index/unsized_.rs | 61 ++++++++++---- tests/persistence/sync/mod.rs | 1 + .../sync/repeated_string_upsert.rs | 80 +++++++++++++++++++ 3 files changed, 125 insertions(+), 17 deletions(-) create mode 100644 tests/persistence/sync/repeated_string_upsert.rs diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 42070a4..601df04 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -390,15 +390,27 @@ where async fn process_change_event_batch(&mut self, events: BatchChangeEvent) -> eyre::Result<()> { let mut pages: HashMap = HashMap::new(); + // A split can change a page's maximum and therefore its table-of- + // contents key while later events in the same CDC batch still refer + // to the pre-split maximum. Keep those historical identities scoped + // to this batch so the event reaches the page it was generated from. + let mut page_aliases: HashMap<(T, Link), PageId> = HashMap::new(); for ev in events { match &ev { ChangeEvent::InsertAt { max_value, .. } | ChangeEvent::RemoveAt { max_value, .. } => { - let page_id = &(max_value.key.clone(), max_value.value); + let event_page_key = (max_value.key.clone(), max_value.value); let page_index = self .table_of_contents - .get(page_id) - .expect("page should be available in table of contents"); + .get(&event_page_key) + .or_else(|| page_aliases.get(&event_page_key).copied()) + .ok_or_else(|| { + eyre!( + "unsized index event references missing page {event_page_key:?}: event={ev:?}; table_of_contents={:?}; buffered_pages={:?}", + self.table_of_contents.iter().collect::>(), + pages.keys().collect::>() + ) + })?; let page = pages.get_mut(&page_index); let page_to_update = if let Some(page) = page { page @@ -413,19 +425,19 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; + let current_page_key = ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, + ); page_to_update.inner.apply_change_event(ev.clone())?; - if &( + let updated_page_key = ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, - ) != page_id - { - self.table_of_contents.update_key( - page_id, - ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ), - ); + ); + page_aliases.insert(event_page_key, page_index); + page_aliases.insert(current_page_key.clone(), page_index); + if updated_page_key != current_page_key { + self.table_of_contents.update_key(¤t_page_key, updated_page_key); } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -452,12 +464,19 @@ where max_value, split_index, } => { - let page_id = &(max_value.key.clone(), max_value.value); + let event_page_key = (max_value.key.clone(), max_value.value); let page_index = self .table_of_contents - .get(page_id) - .expect("page should be available in table of contents"); + .get(&event_page_key) + .or_else(|| page_aliases.get(&event_page_key).copied()) + .ok_or_else(|| { + eyre!( + "unsized index split references missing page {event_page_key:?}: event={ev:?}; table_of_contents={:?}; buffered_pages={:?}", + self.table_of_contents.iter().collect::>(), + pages.keys().collect::>() + ) + })?; let page = pages.get_mut(&page_index); let page_to_update = if let Some(page) = page { page @@ -472,6 +491,10 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; + let current_page_key = ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, + ); let splitted_page = page_to_update.inner.split(*split_index); let new_page_id = if let Some(id) = self.table_of_contents.pop_empty_page_id() { @@ -481,7 +504,7 @@ where }; self.table_of_contents.update_key( - page_id, + ¤t_page_key, ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, @@ -497,6 +520,10 @@ where header, }; pages.insert(new_page_id, general_page); + // The pre-split maximum remains the right page's identity. + // A following remove/insert pair can still name it even + // after the remove temporarily lowers that maximum. + page_aliases.insert(event_page_key, new_page_id); } } } diff --git a/tests/persistence/sync/mod.rs b/tests/persistence/sync/mod.rs index 1306a34..23bbb71 100644 --- a/tests/persistence/sync/mod.rs +++ b/tests/persistence/sync/mod.rs @@ -9,6 +9,7 @@ mod failure; mod failure_multi_index; mod many_strings; mod option; +mod repeated_string_upsert; mod string_primary_index; mod string_re_read; mod string_secondary_index; diff --git a/tests/persistence/sync/repeated_string_upsert.rs b/tests/persistence/sync/repeated_string_upsert.rs new file mode 100644 index 0000000..de9e97a --- /dev/null +++ b/tests/persistence/sync/repeated_string_upsert.rs @@ -0,0 +1,80 @@ +use std::time::Duration; + +use tokio::time::timeout; +use worktable::prelude::*; +use worktable::worktable; + +use crate::remove_dir_if_exists; + +worktable!( + name: StringBlob, + persist: true, + columns: { + key: String primary_key, + value: String, + updated_at: String, + }, +); + +#[test] +fn repeated_varying_string_upserts_keep_the_worker_healthy() { + let path = "tests/data/sync/repeated_string_upsert"; + let config = DiskConfig::new_with_table_name( + path, + StringBlobWorkTable::name_snake_case(), + StringBlobWorkTable::version(), + ); + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_io() + .enable_time() + .build() + .unwrap(); + + runtime.block_on(async { + remove_dir_if_exists(path.to_string()).await; + { + let engine = StringBlobPersistenceEngine::new(config.clone()).await.unwrap(); + let table = StringBlobWorkTable::load(engine).await.unwrap(); + + for index in 0..256 { + table + .insert(StringBlobRow { + key: format!("marker:{index:04}"), + value: format!("initial-{index}"), + updated_at: format!("2026-08-08T00:{:02}:00Z", index % 60), + }) + .unwrap(); + } + table + .insert(StringBlobRow { + key: "settings".into(), + value: "{}".into(), + updated_at: "2026-08-08T00:00:00Z".into(), + }) + .unwrap(); + table.wait_for_ops().await.unwrap(); + + for revision in 0..1_000 { + table + .upsert(StringBlobRow { + key: "settings".into(), + value: "x".repeat(64 + revision % 4_096), + updated_at: format!("2026-08-08T01:{:02}:{:02}Z", revision % 60, revision % 60), + }) + .await + .unwrap(); + } + + timeout(Duration::from_secs(15), table.wait_for_ops()) + .await + .expect("persistence stalled after repeated string upserts") + .expect("persistence worker failed after repeated string upserts"); + } + { + let engine = StringBlobPersistenceEngine::new(config).await.unwrap(); + let table = StringBlobWorkTable::load(engine).await.unwrap(); + assert_eq!(table.select("settings".to_string()).unwrap().value.len(), 1_063); + } + }); +} From ee65f653a495ecae77b907cf0643ee4b0e5d3cc1 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 05:59:08 +0700 Subject: [PATCH 02/14] fix: report persistence worker panics immediately --- .../generator/space_file/worktable_impls.rs | 20 ++++++++ src/persistence/task.rs | 50 ++++++++++++++++++- 2 files changed, 69 insertions(+), 1 deletion(-) diff --git a/codegen/src/persist_table/generator/space_file/worktable_impls.rs b/codegen/src/persist_table/generator/space_file/worktable_impls.rs index 20745e6..735e7cd 100644 --- a/codegen/src/persist_table/generator/space_file/worktable_impls.rs +++ b/codegen/src/persist_table/generator/space_file/worktable_impls.rs @@ -10,6 +10,7 @@ impl Generator { let space_info_fn = self.gen_worktable_space_info_fn(); let persisted_pk_fn = self.gen_worktable_persisted_primary_key_fn(); let wait_for_ops_fn = self.gen_worktable_wait_for_ops_fn(); + let wait_for_failure_fn = self.gen_worktable_wait_for_failure_fn(); let close_fn = self.gen_worktable_close_fn(); let persisted_data_file_size_fn = self.gen_persisted_data_file_size_fn(); @@ -18,6 +19,7 @@ impl Generator { #space_info_fn #persisted_pk_fn #wait_for_ops_fn + #wait_for_failure_fn #close_fn #persisted_data_file_size_fn } @@ -55,6 +57,24 @@ impl Generator { } } + fn gen_worktable_wait_for_failure_fn(&self) -> TokenStream { + if self.attributes.read_only { + quote! { + pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { + std::future::pending().await + } + } + } else { + quote! { + /// Waits for this table's persistence worker to fail. + /// An idle healthy worker does not complete this future. + pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { + self.1.wait_for_failure().await + } + } + } + } + fn gen_worktable_close_fn(&self) -> TokenStream { if self.attributes.read_only { quote! { diff --git a/src/persistence/task.rs b/src/persistence/task.rs index f26a5c2..a7c2d9c 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -2,11 +2,13 @@ use std::collections::{HashMap, HashSet, VecDeque}; use std::fmt::Debug; use std::hash::Hash; use std::marker::PhantomData; +use std::panic::AssertUnwindSafe; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::time::Duration; use data_bucket::page::PageId; +use futures::FutureExt; use parking_lot::Mutex as ParkingMutex; use tokio::sync::Notify; use tokio::task::JoinHandle; @@ -417,6 +419,7 @@ mod lifecycle_tests { None, Engine, IndexCorruption, + Panic, } impl PersistenceEngine<(), u64, TestEvents, TestIndex> for TestEngine { @@ -447,6 +450,7 @@ mod lifecycle_tests { PersistenceIndexCorruption::new("table/primary.wt.idx", "injected shadow divergence").into(), ); } + TestFailure::Panic => panic!("injected persistence worker panic"), } self.batches.fetch_add(1, Ordering::Relaxed); self.events.lock().push("batch"); @@ -567,6 +571,23 @@ mod lifecycle_tests { assert!(Arc::ptr_eq(&wait_error, &close_error)); } + #[tokio::test] + async fn engine_panic_is_terminal_and_reported_to_waiters() { + let task = PersistenceTask::run_engine(TestEngine { + batches: Arc::new(AtomicUsize::new(0)), + events: Arc::new(ParkingMutex::new(Vec::new())), + config: TestConfig, + failure: TestFailure::Panic, + }); + + task.apply_operation(insert_operation(1)).unwrap(); + let wait_error = task.wait_for_failure().await.unwrap_err(); + assert!(wait_error.to_string().contains("injected persistence worker panic")); + + let intake_error = task.apply_operation(insert_operation(2)).unwrap_err(); + assert!(Arc::ptr_eq(&wait_error, &intake_error)); + } + #[tokio::test] async fn vacuum_reclamation_waits_for_preceding_row_moves() { let events = Arc::new(ParkingMutex::new(Vec::new())); @@ -864,7 +885,7 @@ impl let analyzer_in_progress = Arc::new(AtomicBool::new(true)); let task_analyzer_in_progress = analyzer_in_progress.clone(); - let task = async move { + let worker = async move { let mut pending_reclaim: Option> = None; loop { let message = if pending_reclaim.is_none() { @@ -951,6 +972,17 @@ impl } } }; + let supervisor_lifecycle = lifecycle.clone(); + let task = async move { + if let Err(payload) = AssertUnwindSafe(worker).catch_unwind().await { + let reason = payload + .downcast_ref::<&str>() + .copied() + .or_else(|| payload.downcast_ref::().map(String::as_str)) + .unwrap_or("persistence worker panicked without a message"); + supervisor_lifecycle.fail(eyre::eyre!("persistence worker panicked: {reason}")); + } + }; let engine_task_handle = tokio::spawn(task); Self { queue, @@ -1008,6 +1040,22 @@ impl } } + /// Waits until the worker fails or closes. + /// + /// Unlike [`Self::wait_for_ops`], an idle healthy worker does not satisfy + /// this future. Applications can keep it alive as a failure notification + /// without forcing persistence queues to drain or polling their state. + pub async fn wait_for_failure(&self) -> PersistenceResult { + loop { + let notified = self.lifecycle.notify.notified(); + match self.lifecycle.state() { + PersistenceState::Failed(error) => return Err(error), + PersistenceState::Closed => return Ok(()), + PersistenceState::Running | PersistenceState::Closing => notified.await, + } + } + } + pub async fn close(mut self) -> PersistenceResult { let begin_result = self.lifecycle.begin_close(); self.queue.wake(); From c6e85ce3d6a98f3586c296decb0fc107fd9f9cc1 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 06:13:22 +0700 Subject: [PATCH 03/14] release: prepare 1.0.0-beta.9 --- Cargo.toml | 4 ++-- codegen/Cargo.toml | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index f92e8cc..8dbe242 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.8" +version = "1.0.0-beta.9" 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.8" } +worktable_codegen = { path = "codegen", version = "=1.0.0-beta.9" } [dev-dependencies] chrono = "0.4.43" diff --git a/codegen/Cargo.toml b/codegen/Cargo.toml index 45016bf..a5daf24 100644 --- a/codegen/Cargo.toml +++ b/codegen/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "worktable_codegen" -version = "1.0.0-beta.8" +version = "1.0.0-beta.9" edition = "2024" license = "MIT" description = "Proc-macro companion crate for worktable: the worktable! macro and its derives." From f820a6d2780a943fdee70dbf4dda71bd881e28cd Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 07:58:37 +0700 Subject: [PATCH 04/14] fix: harden persistence failure reporting --- .../generator/space_file/worktable_impls.rs | 3 +- src/persistence/space/index/unsized_.rs | 30 +++++++++++-------- src/persistence/task.rs | 22 ++++++++------ 3 files changed, 32 insertions(+), 23 deletions(-) diff --git a/codegen/src/persist_table/generator/space_file/worktable_impls.rs b/codegen/src/persist_table/generator/space_file/worktable_impls.rs index 735e7cd..64950f0 100644 --- a/codegen/src/persist_table/generator/space_file/worktable_impls.rs +++ b/codegen/src/persist_table/generator/space_file/worktable_impls.rs @@ -66,8 +66,9 @@ impl Generator { } } else { quote! { - /// Waits for this table's persistence worker to fail. + /// Waits for this table's persistence worker to fail or close. /// An idle healthy worker does not complete this future. + /// A graceful close returns `Ok(())`. pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { self.1.wait_for_failure().await } diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 601df04..63a1da0 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -406,9 +406,10 @@ where .or_else(|| page_aliases.get(&event_page_key).copied()) .ok_or_else(|| { eyre!( - "unsized index event references missing page {event_page_key:?}: event={ev:?}; table_of_contents={:?}; buffered_pages={:?}", - self.table_of_contents.iter().collect::>(), - pages.keys().collect::>() + "unsized index event references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + self.table_of_contents.pages.len(), + pages.len(), + page_aliases.len() ) })?; let page = pages.get_mut(&page_index); @@ -430,13 +431,14 @@ where page_to_update.inner.node_id.link, ); page_to_update.inner.apply_change_event(ev.clone())?; - let updated_page_key = ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ); - page_aliases.insert(event_page_key, page_index); - page_aliases.insert(current_page_key.clone(), page_index); - if updated_page_key != current_page_key { + if page_to_update.inner.node_id.key != current_page_key.0 + || page_to_update.inner.node_id.link != current_page_key.1 + { + let updated_page_key = ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, + ); + page_aliases.insert(current_page_key.clone(), page_index); self.table_of_contents.update_key(¤t_page_key, updated_page_key); } } @@ -472,9 +474,10 @@ where .or_else(|| page_aliases.get(&event_page_key).copied()) .ok_or_else(|| { eyre!( - "unsized index split references missing page {event_page_key:?}: event={ev:?}; table_of_contents={:?}; buffered_pages={:?}", - self.table_of_contents.iter().collect::>(), - pages.keys().collect::>() + "unsized index split references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + self.table_of_contents.pages.len(), + pages.len(), + page_aliases.len() ) })?; let page = pages.get_mut(&page_index); @@ -524,6 +527,7 @@ where // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. page_aliases.insert(event_page_key, new_page_id); + page_aliases.insert(current_page_key, new_page_id); } } } diff --git a/src/persistence/task.rs b/src/persistence/task.rs index a7c2d9c..6e7b302 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -582,7 +582,10 @@ mod lifecycle_tests { task.apply_operation(insert_operation(1)).unwrap(); let wait_error = task.wait_for_failure().await.unwrap_err(); - assert!(wait_error.to_string().contains("injected persistence worker panic")); + assert_eq!( + wait_error.to_string(), + "persistence engine failed: persistence worker panicked" + ); let intake_error = task.apply_operation(insert_operation(2)).unwrap_err(); assert!(Arc::ptr_eq(&wait_error, &intake_error)); @@ -975,12 +978,11 @@ impl let supervisor_lifecycle = lifecycle.clone(); let task = async move { if let Err(payload) = AssertUnwindSafe(worker).catch_unwind().await { - let reason = payload - .downcast_ref::<&str>() - .copied() - .or_else(|| payload.downcast_ref::().map(String::as_str)) - .unwrap_or("persistence worker panicked without a message"); - supervisor_lifecycle.fail(eyre::eyre!("persistence worker panicked: {reason}")); + // The process panic hook already records the local diagnostic. + // Do not propagate arbitrary panic payloads to API consumers: + // engines can panic with paths or row-derived strings. + drop(payload); + supervisor_lifecycle.fail(eyre::eyre!("persistence worker panicked")); } }; let engine_task_handle = tokio::spawn(task); @@ -1043,8 +1045,10 @@ impl /// Waits until the worker fails or closes. /// /// Unlike [`Self::wait_for_ops`], an idle healthy worker does not satisfy - /// this future. Applications can keep it alive as a failure notification - /// without forcing persistence queues to drain or polling their state. + /// this future. Applications can keep it alive as a terminal-state + /// notification without forcing persistence queues to drain or polling + /// their state. A graceful close returns `Ok(())`; a failure returns the + /// shared terminal error. pub async fn wait_for_failure(&self) -> PersistenceResult { loop { let notified = self.lifecycle.notify.notified(); From 4651d752cb32371025b9bfd5d39e1e638b59d262 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 10:26:58 +0700 Subject: [PATCH 05/14] fix: register persistence terminal waiters eagerly --- src/persistence/task.rs | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/persistence/task.rs b/src/persistence/task.rs index 6e7b302..f307a94 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -1052,6 +1052,12 @@ impl pub async fn wait_for_failure(&self) -> PersistenceResult { loop { let notified = self.lifecycle.notify.notified(); + tokio::pin!(notified); + // `notify_waiters` does not retain a permit. Register this waiter + // before reading the lifecycle state so a terminal transition + // cannot land between the state read and the first poll of + // `notified` and be lost forever. + notified.as_mut().enable(); match self.lifecycle.state() { PersistenceState::Failed(error) => return Err(error), PersistenceState::Closed => return Ok(()), From 9c7b1c27b24c2ebe260190a069ea9572c5704a10 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 11:22:34 +0700 Subject: [PATCH 06/14] fix: handle worker cancellation and bound aliases --- src/persistence/space/index/unsized_.rs | 62 +++++++++++++++++++++-- src/persistence/task.rs | 67 +++++++++++++++++++++++++ 2 files changed, 126 insertions(+), 3 deletions(-) diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 63a1da0..3e006e2 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -394,6 +394,10 @@ where // contents key while later events in the same CDC batch still refer // to the pre-split maximum. Keep those historical identities scoped // to this batch so the event reaches the page it was generated from. + // At most two transitional identities per buffered page. Replacing a + // page maximum consumes its prior aliases; a split may retain both the + // event's identity and the page's actual pre-split identity for the + // new right page. Memory is therefore bounded by pages, not operations. let mut page_aliases: HashMap<(T, Link), PageId> = HashMap::new(); for ev in events { match &ev { @@ -438,7 +442,7 @@ where page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, ); - page_aliases.insert(current_page_key.clone(), page_index); + replace_page_aliases(&mut page_aliases, page_index, [current_page_key.clone()]); self.table_of_contents.update_key(¤t_page_key, updated_page_key); } } @@ -526,8 +530,8 @@ where // The pre-split maximum remains the right page's identity. // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. - page_aliases.insert(event_page_key, new_page_id); - page_aliases.insert(current_page_key, new_page_id); + page_aliases.retain(|_, target| *target != page_index); + replace_page_aliases(&mut page_aliases, new_page_id, [event_page_key, current_page_key]); } } } @@ -544,3 +548,55 @@ where Ok(()) } } + +fn replace_page_aliases( + aliases: &mut HashMap<(T, Link), PageId>, + page_id: PageId, + keys: impl IntoIterator, +) { + aliases.retain(|_, target| *target != page_id); + for key in keys { + aliases.insert(key, page_id); + } +} + +#[cfg(test)] +mod alias_tests { + use super::*; + + fn link(offset: u32) -> Link { + Link { + page_id: 1.into(), + offset, + length: 8, + } + } + + #[test] + fn repeated_maximum_changes_keep_one_alias_per_page() { + let page_id = PageId::from(7); + let mut aliases = HashMap::new(); + for revision in 0..1_000 { + replace_page_aliases(&mut aliases, page_id, [(format!("key-{revision}"), link(revision))]); + } + + assert_eq!(aliases.len(), 1); + assert_eq!(aliases.get(&("key-999".into(), link(999))).copied(), Some(page_id)); + } + + #[test] + fn split_keeps_only_the_two_live_transitional_identities() { + let old_page = PageId::from(3); + let right_page = PageId::from(4); + let mut aliases = HashMap::from([(("older".to_string(), link(1)), old_page)]); + aliases.retain(|_, target| *target != old_page); + replace_page_aliases( + &mut aliases, + right_page, + [("event".to_string(), link(2)), ("current".to_string(), link(3))], + ); + + assert_eq!(aliases.len(), 2); + assert!(aliases.values().all(|target| *target == right_page)); + } +} diff --git a/src/persistence/task.rs b/src/persistence/task.rs index f307a94..47af72b 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -107,6 +107,32 @@ impl PersistenceLifecycle { } } +/// Marks a worker cancellation as terminal even when its future never runs to +/// completion. Tokio cancellation drops the future; it is not a panic and +/// therefore cannot be observed by `catch_unwind`. +struct WorkerCompletionGuard { + lifecycle: Arc, + armed: bool, +} + +impl WorkerCompletionGuard { + fn new(lifecycle: Arc) -> Self { + Self { lifecycle, armed: true } + } + + fn disarm(mut self) { + self.armed = false; + } +} + +impl Drop for WorkerCompletionGuard { + fn drop(&mut self) { + if self.armed { + self.lifecycle.fail(eyre::eyre!("persistence worker was cancelled")); + } + } +} + pub struct QueueAnalyzer { operations: OptimizedVec>, queue_inner_wt: Arc, @@ -591,6 +617,43 @@ mod lifecycle_tests { assert!(Arc::ptr_eq(&wait_error, &intake_error)); } + #[test] + fn runtime_shutdown_is_terminal_and_rejects_later_operations() { + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .enable_all() + .build() + .unwrap(); + let task = runtime.block_on(async { + let task = PersistenceTask::run_engine(TestEngine { + batches: Arc::new(AtomicUsize::new(0)), + events: Arc::new(ParkingMutex::new(Vec::new())), + config: TestConfig, + failure: TestFailure::None, + }); + tokio::task::yield_now().await; + task + }); + + drop(runtime); + + let verifier = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .unwrap(); + let wait_error = verifier + .block_on(async { tokio::time::timeout(Duration::from_secs(1), task.wait_for_failure()).await }) + .expect("cancelled worker must notify terminal waiters") + .unwrap_err(); + assert_eq!( + wait_error.to_string(), + "persistence engine failed: persistence worker was cancelled" + ); + + let intake_error = task.apply_operation(insert_operation(1)).unwrap_err(); + assert!(Arc::ptr_eq(&wait_error, &intake_error)); + } + #[tokio::test] async fn vacuum_reclamation_waits_for_preceding_row_moves() { let events = Arc::new(ParkingMutex::new(Vec::new())); @@ -976,6 +1039,9 @@ impl } }; let supervisor_lifecycle = lifecycle.clone(); + // Constructed outside the async block so cancellation before its first + // poll still drops the guard and publishes terminal failure. + let completion_guard = WorkerCompletionGuard::new(lifecycle.clone()); let task = async move { if let Err(payload) = AssertUnwindSafe(worker).catch_unwind().await { // The process panic hook already records the local diagnostic. @@ -984,6 +1050,7 @@ impl drop(payload); supervisor_lifecycle.fail(eyre::eyre!("persistence worker panicked")); } + completion_guard.disarm(); }; let engine_task_handle = tokio::spawn(task); Self { From 496e7f9ba43a12318315d36459354a94bbf26524 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 11:32:25 +0700 Subject: [PATCH 07/14] fix: resolve read-only failure monitoring --- .../src/persist_table/generator/space_file/worktable_impls.rs | 4 +++- codegen/src/persist_table/mod.rs | 4 ++++ 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/codegen/src/persist_table/generator/space_file/worktable_impls.rs b/codegen/src/persist_table/generator/space_file/worktable_impls.rs index 64950f0..b94f595 100644 --- a/codegen/src/persist_table/generator/space_file/worktable_impls.rs +++ b/codegen/src/persist_table/generator/space_file/worktable_impls.rs @@ -60,8 +60,10 @@ impl Generator { fn gen_worktable_wait_for_failure_fn(&self) -> TokenStream { if self.attributes.read_only { quote! { + /// Read-only tables have no persistence worker to monitor, so + /// they are already terminal from this API's perspective. pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { - std::future::pending().await + Ok(()) } } } else { diff --git a/codegen/src/persist_table/mod.rs b/codegen/src/persist_table/mod.rs index e0d12cf..e1989a0 100644 --- a/codegen/src/persist_table/mod.rs +++ b/codegen/src/persist_table/mod.rs @@ -78,6 +78,10 @@ mod tests { output.contains("fn into_worktable_with_mode"), "read_only should generate explicit recovery-mode conversion" ); + assert!( + output.contains("async fn wait_for_persistence_failure (& self) -> PersistenceResult { Ok (()) }"), + "read_only failure monitoring should resolve immediately because it has no worker" + ); assert!(output.contains("LoadMode :: Strict")); } From 107eb5bb52aadd922f9bbde9434a73338693b930 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 11:39:53 +0700 Subject: [PATCH 08/14] perf: index page aliases by page --- src/persistence/space/index/unsized_.rs | 136 ++++++++++++++++++++---- 1 file changed, 114 insertions(+), 22 deletions(-) diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 3e006e2..1d00a8f 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -398,7 +398,7 @@ where // page maximum consumes its prior aliases; a split may retain both the // event's identity and the page's actual pre-split identity for the // new right page. Memory is therefore bounded by pages, not operations. - let mut page_aliases: HashMap<(T, Link), PageId> = HashMap::new(); + let mut page_aliases = PageAliases::default(); for ev in events { match &ev { ChangeEvent::InsertAt { max_value, .. } | ChangeEvent::RemoveAt { max_value, .. } => { @@ -407,7 +407,7 @@ where let page_index = self .table_of_contents .get(&event_page_key) - .or_else(|| page_aliases.get(&event_page_key).copied()) + .or_else(|| page_aliases.get(&event_page_key)) .ok_or_else(|| { eyre!( "unsized index event references a missing page (toc_segments={}, buffered_pages={}, aliases={})", @@ -442,8 +442,8 @@ where page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, ); - replace_page_aliases(&mut page_aliases, page_index, [current_page_key.clone()]); self.table_of_contents.update_key(¤t_page_key, updated_page_key); + page_aliases.replace(page_index, [current_page_key]); } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -475,7 +475,7 @@ where let page_index = self .table_of_contents .get(&event_page_key) - .or_else(|| page_aliases.get(&event_page_key).copied()) + .or_else(|| page_aliases.get(&event_page_key)) .ok_or_else(|| { eyre!( "unsized index split references a missing page (toc_segments={}, buffered_pages={}, aliases={})", @@ -530,8 +530,8 @@ where // The pre-split maximum remains the right page's identity. // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. - page_aliases.retain(|_, target| *target != page_index); - replace_page_aliases(&mut page_aliases, new_page_id, [event_page_key, current_page_key]); + page_aliases.remove_page(page_index); + page_aliases.replace(new_page_id, [event_page_key, current_page_key]); } } } @@ -549,14 +549,86 @@ where } } -fn replace_page_aliases( - aliases: &mut HashMap<(T, Link), PageId>, - page_id: PageId, - keys: impl IntoIterator, -) { - aliases.retain(|_, target| *target != page_id); - for key in keys { - aliases.insert(key, page_id); +struct PageAliases { + by_key: HashMap, PageId>, + by_page: HashMap>>, +} + +impl Default for PageAliases { + fn default() -> Self { + Self { + by_key: HashMap::new(), + by_page: HashMap::new(), + } + } +} + +impl PageAliases { + fn get(&self, key: &(T, Link)) -> Option { + self.by_key.get(key).copied() + } + + fn len(&self) -> usize { + self.by_key.len() + } + + fn remove_page(&mut self, page_id: PageId) { + let Some(keys) = self.by_page.remove(&page_id) else { + return; + }; + for key in keys { + if self.by_key.get(key.as_ref()) == Some(&page_id) { + self.by_key.remove(key.as_ref()); + } + } + } + + fn replace(&mut self, page_id: PageId, keys: impl IntoIterator) { + self.remove_page(page_id); + let mut page_keys = Vec::with_capacity(2); + for key in keys { + if page_keys + .iter() + .any(|existing: &Arc<(T, Link)>| existing.as_ref() == &key) + { + continue; + } + + let key = Arc::new(key); + let existing = self + .by_key + .get_key_value(key.as_ref()) + .map(|(stored, previous_page)| (stored.clone(), *previous_page)); + let shared_key = if let Some((stored, previous_page)) = existing { + *self + .by_key + .get_mut(stored.as_ref()) + .expect("alias found immediately before update") = page_id; + if previous_page != page_id { + self.remove_key_from_page(previous_page, stored.as_ref()); + } + stored + } else { + self.by_key.insert(key.clone(), page_id); + key + }; + page_keys.push(shared_key); + } + if !page_keys.is_empty() { + self.by_page.insert(page_id, page_keys); + } + } + + fn remove_key_from_page(&mut self, page_id: PageId, key: &(T, Link)) { + let remove_page = if let Some(keys) = self.by_page.get_mut(&page_id) { + keys.retain(|existing| existing.as_ref() != key); + keys.is_empty() + } else { + false + }; + if remove_page { + self.by_page.remove(&page_id); + } } } @@ -575,28 +647,48 @@ mod alias_tests { #[test] fn repeated_maximum_changes_keep_one_alias_per_page() { let page_id = PageId::from(7); - let mut aliases = HashMap::new(); + let mut aliases = PageAliases::default(); for revision in 0..1_000 { - replace_page_aliases(&mut aliases, page_id, [(format!("key-{revision}"), link(revision))]); + aliases.replace(page_id, [(format!("key-{revision}"), link(revision))]); } assert_eq!(aliases.len(), 1); - assert_eq!(aliases.get(&("key-999".into(), link(999))).copied(), Some(page_id)); + assert_eq!(aliases.get(&("key-999".into(), link(999))), Some(page_id)); } #[test] fn split_keeps_only_the_two_live_transitional_identities() { let old_page = PageId::from(3); let right_page = PageId::from(4); - let mut aliases = HashMap::from([(("older".to_string(), link(1)), old_page)]); - aliases.retain(|_, target| *target != old_page); - replace_page_aliases( - &mut aliases, + let mut aliases = PageAliases::default(); + aliases.replace(old_page, [("older".to_string(), link(1))]); + aliases.remove_page(old_page); + aliases.replace( right_page, [("event".to_string(), link(2)), ("current".to_string(), link(3))], ); assert_eq!(aliases.len(), 2); - assert!(aliases.values().all(|target| *target == right_page)); + assert_eq!(aliases.by_page.get(&right_page).map(Vec::len), Some(2)); + } + + #[test] + fn replacing_many_pages_preserves_constant_work_per_page() { + let mut aliases = PageAliases::default(); + for page in 1..=1_000 { + aliases.replace(PageId::from(page), [(format!("old-{page}"), link(page))]); + } + for page in 1..=1_000 { + aliases.replace(PageId::from(page), [(format!("new-{page}"), link(page + 1_000))]); + } + + assert_eq!(aliases.len(), 1_000); + assert_eq!(aliases.by_page.len(), 1_000); + for page in 1..=1_000 { + assert_eq!( + aliases.get(&(format!("new-{page}"), link(page + 1_000))), + Some(PageId::from(page)) + ); + } } } From 267c68e85d0a9d8a0817c5a2be638f247451f22e Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 11:56:58 +0700 Subject: [PATCH 09/14] refactor: simplify page alias ownership --- src/persistence/space/index/unsized_.rs | 54 ++++++------------------- 1 file changed, 12 insertions(+), 42 deletions(-) diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 1d00a8f..1930025 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -550,8 +550,8 @@ where } struct PageAliases { - by_key: HashMap, PageId>, - by_page: HashMap>>, + by_key: HashMap<(T, Link), PageId>, + by_page: HashMap>, } impl Default for PageAliases { @@ -563,7 +563,7 @@ impl Default for PageAliases { } } -impl PageAliases { +impl PageAliases { fn get(&self, key: &(T, Link)) -> Option { self.by_key.get(key).copied() } @@ -577,9 +577,7 @@ impl PageAliases { return; }; for key in keys { - if self.by_key.get(key.as_ref()) == Some(&page_id) { - self.by_key.remove(key.as_ref()); - } + self.by_key.remove(&key); } } @@ -587,49 +585,21 @@ impl PageAliases { self.remove_page(page_id); let mut page_keys = Vec::with_capacity(2); for key in keys { - if page_keys - .iter() - .any(|existing: &Arc<(T, Link)>| existing.as_ref() == &key) - { + if page_keys.contains(&key) { continue; } - - let key = Arc::new(key); - let existing = self - .by_key - .get_key_value(key.as_ref()) - .map(|(stored, previous_page)| (stored.clone(), *previous_page)); - let shared_key = if let Some((stored, previous_page)) = existing { - *self - .by_key - .get_mut(stored.as_ref()) - .expect("alias found immediately before update") = page_id; - if previous_page != page_id { - self.remove_key_from_page(previous_page, stored.as_ref()); - } - stored - } else { - self.by_key.insert(key.clone(), page_id); - key - }; - page_keys.push(shared_key); + // `(value, Link)` identifies one physical index entry, so it + // cannot be owned by another page at the same time. Keep one copy + // in each direction instead of interning keys behind `Arc` for an + // impossible cross-page migration path. + let previous = self.by_key.insert(key.clone(), page_id); + debug_assert!(previous.is_none(), "an alias cannot belong to two pages"); + page_keys.push(key); } if !page_keys.is_empty() { self.by_page.insert(page_id, page_keys); } } - - fn remove_key_from_page(&mut self, page_id: PageId, key: &(T, Link)) { - let remove_page = if let Some(keys) = self.by_page.get_mut(&page_id) { - keys.retain(|existing| existing.as_ref() != key); - keys.is_empty() - } else { - false - }; - if remove_page { - self.by_page.remove(&page_id); - } - } } #[cfg(test)] From 718052bb2d400c1960f0c3358b5ab62c48340e4b Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 12:07:31 +0700 Subject: [PATCH 10/14] fix: preserve transitional page identities --- .../generator/space_file/worktable_impls.rs | 8 +-- codegen/src/persist_table/mod.rs | 4 +- src/persistence/space/index/unsized_.rs | 64 ++++++++++++++++--- 3 files changed, 57 insertions(+), 19 deletions(-) diff --git a/codegen/src/persist_table/generator/space_file/worktable_impls.rs b/codegen/src/persist_table/generator/space_file/worktable_impls.rs index b94f595..3bf3e79 100644 --- a/codegen/src/persist_table/generator/space_file/worktable_impls.rs +++ b/codegen/src/persist_table/generator/space_file/worktable_impls.rs @@ -59,13 +59,7 @@ impl Generator { fn gen_worktable_wait_for_failure_fn(&self) -> TokenStream { if self.attributes.read_only { - quote! { - /// Read-only tables have no persistence worker to monitor, so - /// they are already terminal from this API's perspective. - pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { - Ok(()) - } - } + quote! {} } else { quote! { /// Waits for this table's persistence worker to fail or close. diff --git a/codegen/src/persist_table/mod.rs b/codegen/src/persist_table/mod.rs index e1989a0..d54455b 100644 --- a/codegen/src/persist_table/mod.rs +++ b/codegen/src/persist_table/mod.rs @@ -79,8 +79,8 @@ mod tests { "read_only should generate explicit recovery-mode conversion" ); assert!( - output.contains("async fn wait_for_persistence_failure (& self) -> PersistenceResult { Ok (()) }"), - "read_only failure monitoring should resolve immediately because it has no worker" + !output.contains("wait_for_persistence_failure"), + "read_only should not expose monitoring for a worker it does not have" ); assert!(output.contains("LoadMode :: Strict")); } diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 1930025..fa8cc09 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -394,10 +394,10 @@ where // contents key while later events in the same CDC batch still refer // to the pre-split maximum. Keep those historical identities scoped // to this batch so the event reaches the page it was generated from. - // At most two transitional identities per buffered page. Replacing a - // page maximum consumes its prior aliases; a split may retain both the - // event's identity and the page's actual pre-split identity for the - // new right page. Memory is therefore bounded by pages, not operations. + // At most two transitional identities per buffered page: the event's + // identity and the page's actual pre-apply identity. Keeping both is + // required when a split is followed by a remove/insert pair that still + // names the pre-split maximum. Memory is bounded by pages, not events. let mut page_aliases = PageAliases::default(); for ev in events { match &ev { @@ -443,7 +443,7 @@ where page_to_update.inner.node_id.link, ); self.table_of_contents.update_key(¤t_page_key, updated_page_key); - page_aliases.replace(page_index, [current_page_key]); + page_aliases.replace(page_index, [event_page_key, current_page_key]); } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -551,7 +551,36 @@ where struct PageAliases { by_key: HashMap<(T, Link), PageId>, - by_page: HashMap>, + by_page: HashMap>, +} + +struct PageAliasKeys { + slots: [Option<(T, Link)>; 2], +} + +impl Default for PageAliasKeys { + fn default() -> Self { + Self { slots: [None, None] } + } +} + +impl PageAliasKeys { + fn contains(&self, key: &(T, Link)) -> bool { + self.slots.iter().flatten().any(|stored| stored == key) + } + + fn push(&mut self, key: (T, Link)) { + let slot = self + .slots + .iter_mut() + .find(|slot| slot.is_none()) + .expect("a page has at most two transitional aliases"); + *slot = Some(key); + } + + fn len(&self) -> usize { + self.slots.iter().flatten().count() + } } impl Default for PageAliases { @@ -576,14 +605,14 @@ impl PageAliases { let Some(keys) = self.by_page.remove(&page_id) else { return; }; - for key in keys { + for key in keys.slots.into_iter().flatten() { self.by_key.remove(&key); } } fn replace(&mut self, page_id: PageId, keys: impl IntoIterator) { self.remove_page(page_id); - let mut page_keys = Vec::with_capacity(2); + let mut page_keys = PageAliasKeys::default(); for key in keys { if page_keys.contains(&key) { continue; @@ -596,7 +625,7 @@ impl PageAliases { debug_assert!(previous.is_none(), "an alias cannot belong to two pages"); page_keys.push(key); } - if !page_keys.is_empty() { + if page_keys.len() != 0 { self.by_page.insert(page_id, page_keys); } } @@ -639,7 +668,22 @@ mod alias_tests { ); assert_eq!(aliases.len(), 2); - assert_eq!(aliases.by_page.get(&right_page).map(Vec::len), Some(2)); + assert_eq!(aliases.by_page.get(&right_page).map(PageAliasKeys::len), Some(2)); + } + + #[test] + fn split_remove_insert_preserves_the_event_identity() { + let right_page = PageId::from(4); + let pre_split = ("pre-split".to_string(), link(1)); + let post_split = ("post-split".to_string(), link(2)); + let after_remove = ("after-remove".to_string(), link(3)); + let mut aliases = PageAliases::default(); + + aliases.replace(right_page, [pre_split.clone(), post_split.clone()]); + aliases.replace(right_page, [pre_split.clone(), after_remove]); + + assert_eq!(aliases.get(&pre_split), Some(right_page)); + assert_eq!(aliases.get(&post_split), None); } #[test] From 0b86385f887b9b39dfbc72c2f942cdb3b2d9fced Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 12:25:28 +0700 Subject: [PATCH 11/14] fix: isolate terminal persistence notifications --- src/persistence/space/index/unsized_.rs | 148 ++++++++++++++++++------ src/persistence/task.rs | 22 ++-- 2 files changed, 129 insertions(+), 41 deletions(-) diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index fa8cc09..d4ca88a 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -430,20 +430,48 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; - let current_page_key = ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ); + // `apply_change_event` can only rewrite `node_id` in these + // exact cases. Avoid cloning an unsized maximum for the + // overwhelmingly common interior insert/remove path. + let can_change_max = match &ev { + ChangeEvent::InsertAt { value, index, .. } => { + value.key > page_to_update.inner.node_id.key + || *index == page_to_update.inner.slots_size as usize + } + ChangeEvent::RemoveAt { + max_value, + value, + index, + .. + } => { + value == max_value + && *index != 0 + && page_to_update.inner.slots_size != 0 + && *index == page_to_update.inner.slots_size as usize - 1 + } + _ => false, + }; + let current_page_key = can_change_max.then(|| { + ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, + ) + }); page_to_update.inner.apply_change_event(ev.clone())?; - if page_to_update.inner.node_id.key != current_page_key.0 - || page_to_update.inner.node_id.link != current_page_key.1 + if let Some(current_page_key) = current_page_key + && (page_to_update.inner.node_id.key != current_page_key.0 + || page_to_update.inner.node_id.link != current_page_key.1) { let updated_page_key = ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, ); + // The TOC owns the buffered page's actual identity. + // `event_page_key` may be a historical alias, so using + // it as the canonical update target can remove or + // rewrite the wrong segment. self.table_of_contents.update_key(¤t_page_key, updated_page_key); - page_aliases.replace(page_index, [event_page_key, current_page_key]); + page_aliases.replace(page_index, [event_page_key, current_page_key])?; } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -531,7 +559,7 @@ where // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. page_aliases.remove_page(page_index); - page_aliases.replace(new_page_id, [event_page_key, current_page_key]); + page_aliases.replace(new_page_id, [event_page_key, current_page_key])?; } } } @@ -569,18 +597,21 @@ impl PageAliasKeys { self.slots.iter().flatten().any(|stored| stored == key) } - fn push(&mut self, key: (T, Link)) { - let slot = self - .slots - .iter_mut() - .find(|slot| slot.is_none()) - .expect("a page has at most two transitional aliases"); + fn push(&mut self, key: (T, Link)) -> bool { + let Some(slot) = self.slots.iter_mut().find(|slot| slot.is_none()) else { + return false; + }; *slot = Some(key); + true } fn len(&self) -> usize { self.slots.iter().flatten().count() } + + fn iter(&self) -> impl Iterator { + self.slots.iter().flatten() + } } impl Default for PageAliases { @@ -610,24 +641,39 @@ impl PageAliases { } } - fn replace(&mut self, page_id: PageId, keys: impl IntoIterator) { - self.remove_page(page_id); + fn replace(&mut self, page_id: PageId, keys: impl IntoIterator) -> eyre::Result<()> { let mut page_keys = PageAliasKeys::default(); for key in keys { if page_keys.contains(&key) { continue; } - // `(value, Link)` identifies one physical index entry, so it - // cannot be owned by another page at the same time. Keep one copy - // in each direction instead of interning keys behind `Arc` for an - // impossible cross-page migration path. - let previous = self.by_key.insert(key.clone(), page_id); - debug_assert!(previous.is_none(), "an alias cannot belong to two pages"); - page_keys.push(key); + if !page_keys.push(key) { + return Err(eyre!("page {page_id:?} exceeded the two transitional alias invariant")); + } + } + + // `(value, Link)` identifies one physical index entry, so it cannot be + // owned by another page. Enforce that invariant in release builds + // before changing either direction, rather than leaving the maps + // inconsistent if a future event-shape regression violates it. + for key in page_keys.iter() { + if let Some(previous_page) = self.by_key.get(key) + && *previous_page != page_id + { + return Err(eyre!( + "page alias ownership collision between {previous_page:?} and {page_id:?}" + )); + } + } + + self.remove_page(page_id); + for key in page_keys.iter() { + self.by_key.insert(key.clone(), page_id); } if page_keys.len() != 0 { self.by_page.insert(page_id, page_keys); } + Ok(()) } } @@ -648,7 +694,9 @@ mod alias_tests { let page_id = PageId::from(7); let mut aliases = PageAliases::default(); for revision in 0..1_000 { - aliases.replace(page_id, [(format!("key-{revision}"), link(revision))]); + aliases + .replace(page_id, [(format!("key-{revision}"), link(revision))]) + .unwrap(); } assert_eq!(aliases.len(), 1); @@ -660,12 +708,14 @@ mod alias_tests { let old_page = PageId::from(3); let right_page = PageId::from(4); let mut aliases = PageAliases::default(); - aliases.replace(old_page, [("older".to_string(), link(1))]); + aliases.replace(old_page, [("older".to_string(), link(1))]).unwrap(); aliases.remove_page(old_page); - aliases.replace( - right_page, - [("event".to_string(), link(2)), ("current".to_string(), link(3))], - ); + aliases + .replace( + right_page, + [("event".to_string(), link(2)), ("current".to_string(), link(3))], + ) + .unwrap(); assert_eq!(aliases.len(), 2); assert_eq!(aliases.by_page.get(&right_page).map(PageAliasKeys::len), Some(2)); @@ -679,8 +729,10 @@ mod alias_tests { let after_remove = ("after-remove".to_string(), link(3)); let mut aliases = PageAliases::default(); - aliases.replace(right_page, [pre_split.clone(), post_split.clone()]); - aliases.replace(right_page, [pre_split.clone(), after_remove]); + aliases + .replace(right_page, [pre_split.clone(), post_split.clone()]) + .unwrap(); + aliases.replace(right_page, [pre_split.clone(), after_remove]).unwrap(); assert_eq!(aliases.get(&pre_split), Some(right_page)); assert_eq!(aliases.get(&post_split), None); @@ -690,10 +742,14 @@ mod alias_tests { fn replacing_many_pages_preserves_constant_work_per_page() { let mut aliases = PageAliases::default(); for page in 1..=1_000 { - aliases.replace(PageId::from(page), [(format!("old-{page}"), link(page))]); + aliases + .replace(PageId::from(page), [(format!("old-{page}"), link(page))]) + .unwrap(); } for page in 1..=1_000 { - aliases.replace(PageId::from(page), [(format!("new-{page}"), link(page + 1_000))]); + aliases + .replace(PageId::from(page), [(format!("new-{page}"), link(page + 1_000))]) + .unwrap(); } assert_eq!(aliases.len(), 1_000); @@ -705,4 +761,30 @@ mod alias_tests { ); } } + + #[test] + fn alias_invariants_fail_without_corrupting_existing_ownership() { + let first_page = PageId::from(1); + let second_page = PageId::from(2); + let shared = ("shared".to_string(), link(1)); + let mut aliases = PageAliases::default(); + aliases.replace(first_page, [shared.clone()]).unwrap(); + + assert!(aliases.replace(second_page, [shared.clone()]).is_err()); + assert_eq!(aliases.get(&shared), Some(first_page)); + assert!( + aliases + .replace( + second_page, + [ + ("one".to_string(), link(2)), + ("two".to_string(), link(3)), + ("three".to_string(), link(4)), + ], + ) + .is_err() + ); + assert_eq!(aliases.get(&shared), Some(first_page)); + assert!(!aliases.by_page.contains_key(&second_page)); + } } diff --git a/src/persistence/task.rs b/src/persistence/task.rs index 47af72b..d3e32b4 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -43,14 +43,18 @@ const MAX_PAGE_AMOUNT: usize = 16; #[derive(Debug)] struct PersistenceLifecycle { state: ParkingMutex, - notify: Notify, + /// Terminal transitions only: `Failed` or `Closed`. + terminal_notify: Notify, + /// Queue-drain and lifecycle progress observed by `wait_for_ops`. + progress_notify: Notify, } impl PersistenceLifecycle { fn new() -> Self { Self { state: ParkingMutex::new(PersistenceState::Running), - notify: Notify::new(), + terminal_notify: Notify::new(), + progress_notify: Notify::new(), } } @@ -63,7 +67,7 @@ impl PersistenceLifecycle { match &*state { PersistenceState::Running => { *state = PersistenceState::Closing; - self.notify.notify_waiters(); + self.progress_notify.notify_waiters(); Ok(()) } PersistenceState::Closing => Ok(()), @@ -77,7 +81,8 @@ impl PersistenceLifecycle { if matches!(*state, PersistenceState::Closing) { *state = PersistenceState::Closed; } - self.notify.notify_waiters(); + self.terminal_notify.notify_waiters(); + self.progress_notify.notify_waiters(); } fn fail(&self, report: eyre::Report) -> Arc { @@ -93,7 +98,8 @@ impl PersistenceLifecycle { error } }; - self.notify.notify_waiters(); + self.terminal_notify.notify_waiters(); + self.progress_notify.notify_waiters(); error } @@ -963,7 +969,7 @@ impl message } else if analyzer.len() == 0 && pending_reclaim.is_none() { task_analyzer_in_progress.store(false, Ordering::Release); - engine_lifecycle.notify.notify_waiters(); + engine_lifecycle.progress_notify.notify_waiters(); if matches!(engine_lifecycle.state(), PersistenceState::Closing) { engine_lifecycle.finish_close(); return; @@ -1103,7 +1109,7 @@ impl } tokio::select! { - _ = self.lifecycle.notify.notified() => {}, + _ = self.lifecycle.progress_notify.notified() => {}, _ = tokio::time::sleep(Duration::from_secs(1)) => {} } } @@ -1118,7 +1124,7 @@ impl /// shared terminal error. pub async fn wait_for_failure(&self) -> PersistenceResult { loop { - let notified = self.lifecycle.notify.notified(); + let notified = self.lifecycle.terminal_notify.notified(); tokio::pin!(notified); // `notify_waiters` does not retain a permit. Register this waiter // before reading the lifecycle state so a terminal transition From 000ec7ddf145bdc928fc703fb85c7175d29a3fbb Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 12:45:04 +0700 Subject: [PATCH 12/14] fix: decouple persistence monitoring --- .../generator/space_file/worktable_impls.rs | 15 ++- codegen/src/persist_table/mod.rs | 3 +- src/lib.rs | 11 ++- src/persistence/mod.rs | 2 +- src/persistence/space/index/unsized_.rs | 34 ++++--- src/persistence/task.rs | 98 ++++++++++++------- 6 files changed, 104 insertions(+), 59 deletions(-) diff --git a/codegen/src/persist_table/generator/space_file/worktable_impls.rs b/codegen/src/persist_table/generator/space_file/worktable_impls.rs index 3bf3e79..bd3dac7 100644 --- a/codegen/src/persist_table/generator/space_file/worktable_impls.rs +++ b/codegen/src/persist_table/generator/space_file/worktable_impls.rs @@ -10,7 +10,7 @@ impl Generator { let space_info_fn = self.gen_worktable_space_info_fn(); let persisted_pk_fn = self.gen_worktable_persisted_primary_key_fn(); let wait_for_ops_fn = self.gen_worktable_wait_for_ops_fn(); - let wait_for_failure_fn = self.gen_worktable_wait_for_failure_fn(); + let persistence_monitor_fn = self.gen_worktable_persistence_monitor_fn(); let close_fn = self.gen_worktable_close_fn(); let persisted_data_file_size_fn = self.gen_persisted_data_file_size_fn(); @@ -19,7 +19,7 @@ impl Generator { #space_info_fn #persisted_pk_fn #wait_for_ops_fn - #wait_for_failure_fn + #persistence_monitor_fn #close_fn #persisted_data_file_size_fn } @@ -57,16 +57,15 @@ impl Generator { } } - fn gen_worktable_wait_for_failure_fn(&self) -> TokenStream { + fn gen_worktable_persistence_monitor_fn(&self) -> TokenStream { if self.attributes.read_only { quote! {} } else { quote! { - /// Waits for this table's persistence worker to fail or close. - /// An idle healthy worker does not complete this future. - /// A graceful close returns `Ok(())`. - pub async fn wait_for_persistence_failure(&self) -> PersistenceResult { - self.1.wait_for_failure().await + /// Returns a cloneable terminal-state monitor that does not + /// borrow the table and can therefore observe `close()`. + pub fn persistence_monitor(&self) -> PersistenceMonitor { + self.1.monitor() } } } diff --git a/codegen/src/persist_table/mod.rs b/codegen/src/persist_table/mod.rs index d54455b..c57e3ad 100644 --- a/codegen/src/persist_table/mod.rs +++ b/codegen/src/persist_table/mod.rs @@ -79,7 +79,7 @@ mod tests { "read_only should generate explicit recovery-mode conversion" ); assert!( - !output.contains("wait_for_persistence_failure"), + !output.contains("persistence_monitor"), "read_only should not expose monitoring for a worker it does not have" ); assert!(output.contains("LoadMode :: Strict")); @@ -111,6 +111,7 @@ mod tests { output.contains("async fn into_worktable_with_mode"), "normal should generate explicit recovery-mode conversion" ); + assert!(output.contains("fn persistence_monitor")); assert!(output.contains("LoadMode :: Strict")); } diff --git a/src/lib.rs b/src/lib.rs index 9e28c54..d39af05 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -36,11 +36,12 @@ pub mod prelude { pub use crate::persistence::{ AcknowledgeOperation, ArtPersistenceKey, DeleteOperation, DiskConfig, DiskPersistenceEngine, IndexTableOfContents, InsertOperation, LoadMode, Operation, OperationId, PersistedWorkTable, PersistenceConfig, - PersistenceEngine, PersistenceError, PersistenceIndexCorruption, PersistenceLoadError, PersistenceResult, - PersistenceState, PersistenceTask, ReadOnlyPersistenceEngine, SpaceArcticIndex, SpaceCongeeIndex, SpaceData, - SpaceDataOps, SpaceIndex, SpaceIndexOps, SpaceIndexUnsized, SpaceLogicalIndex, SpaceLogicalIndexUnsized, - SpaceSecondaryIndexOps, UpdateOperation, load_persisted_state, map_index_pages_to_toc_and_general, - map_unsized_index_pages_to_toc_and_general, reconstruct_multi_index_nodes, validate_events, + PersistenceEngine, PersistenceError, PersistenceIndexCorruption, PersistenceLoadError, PersistenceMonitor, + PersistenceResult, PersistenceState, PersistenceTask, ReadOnlyPersistenceEngine, SpaceArcticIndex, + SpaceCongeeIndex, SpaceData, SpaceDataOps, SpaceIndex, SpaceIndexOps, SpaceIndexUnsized, SpaceLogicalIndex, + SpaceLogicalIndexUnsized, SpaceSecondaryIndexOps, UpdateOperation, load_persisted_state, + map_index_pages_to_toc_and_general, map_unsized_index_pages_to_toc_and_general, reconstruct_multi_index_nodes, + validate_events, }; pub use crate::primary_key::{PrimaryKeyGenerator, PrimaryKeyGeneratorState, TablePrimaryKey}; pub use crate::table::select::{Order, QueryParams, SelectQueryBuilder, SelectQueryExecutor}; diff --git a/src/persistence/mod.rs b/src/persistence/mod.rs index 121e27a..c875061 100644 --- a/src/persistence/mod.rs +++ b/src/persistence/mod.rs @@ -20,7 +20,7 @@ pub use space::{ SpaceIndexOps, SpaceIndexUnsized, SpaceLogicalIndex, SpaceLogicalIndexUnsized, SpaceSecondaryIndexOps, map_index_pages_to_toc_and_general, map_unsized_index_pages_to_toc_and_general, reconstruct_multi_index_nodes, }; -pub use task::PersistenceTask; +pub use task::{PersistenceMonitor, PersistenceTask}; mod engine; mod error; diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index d4ca88a..7bd2c4b 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -410,7 +410,7 @@ where .or_else(|| page_aliases.get(&event_page_key)) .ok_or_else(|| { eyre!( - "unsized index event references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + "unsized index event for {event_page_key:?} references a missing page (toc_segments={}, buffered_pages={}, aliases={})", self.table_of_contents.pages.len(), pages.len(), page_aliases.len() @@ -451,16 +451,28 @@ where } _ => false, }; - let current_page_key = can_change_max.then(|| { + #[cfg(debug_assertions)] + let debug_pre_event_page_key = ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, + ); + let pre_event_page_key = can_change_max.then(|| { ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, ) }); page_to_update.inner.apply_change_event(ev.clone())?; - if let Some(current_page_key) = current_page_key - && (page_to_update.inner.node_id.key != current_page_key.0 - || page_to_update.inner.node_id.link != current_page_key.1) + #[cfg(debug_assertions)] + debug_assert!( + can_change_max + || (page_to_update.inner.node_id.key == debug_pre_event_page_key.0 + && page_to_update.inner.node_id.link == debug_pre_event_page_key.1), + "data_bucket changed an unsized page identity outside WorkTable's can_change_max predicate" + ); + if let Some(pre_event_page_key) = pre_event_page_key + && (page_to_update.inner.node_id.key != pre_event_page_key.0 + || page_to_update.inner.node_id.link != pre_event_page_key.1) { let updated_page_key = ( page_to_update.inner.node_id.key.clone(), @@ -470,8 +482,8 @@ where // `event_page_key` may be a historical alias, so using // it as the canonical update target can remove or // rewrite the wrong segment. - self.table_of_contents.update_key(¤t_page_key, updated_page_key); - page_aliases.replace(page_index, [event_page_key, current_page_key])?; + self.table_of_contents.update_key(&pre_event_page_key, updated_page_key); + page_aliases.replace(page_index, [event_page_key, pre_event_page_key])?; } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -506,7 +518,7 @@ where .or_else(|| page_aliases.get(&event_page_key)) .ok_or_else(|| { eyre!( - "unsized index split references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + "unsized index split for {event_page_key:?} references a missing page (toc_segments={}, buffered_pages={}, aliases={})", self.table_of_contents.pages.len(), pages.len(), page_aliases.len() @@ -526,7 +538,7 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; - let current_page_key = ( + let pre_split_page_key = ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, ); @@ -539,7 +551,7 @@ where }; self.table_of_contents.update_key( - ¤t_page_key, + &pre_split_page_key, ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, @@ -559,7 +571,7 @@ where // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. page_aliases.remove_page(page_index); - page_aliases.replace(new_page_id, [event_page_key, current_page_key])?; + page_aliases.replace(new_page_id, [event_page_key, pre_split_page_key])?; } } } diff --git a/src/persistence/task.rs b/src/persistence/task.rs index d3e32b4..88ca551 100644 --- a/src/persistence/task.rs +++ b/src/persistence/task.rs @@ -2,13 +2,11 @@ use std::collections::{HashMap, HashSet, VecDeque}; use std::fmt::Debug; use std::hash::Hash; use std::marker::PhantomData; -use std::panic::AssertUnwindSafe; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::time::Duration; use data_bucket::page::PageId; -use futures::FutureExt; use parking_lot::Mutex as ParkingMutex; use tokio::sync::Notify; use tokio::task::JoinHandle; @@ -113,9 +111,8 @@ impl PersistenceLifecycle { } } -/// Marks a worker cancellation as terminal even when its future never runs to -/// completion. Tokio cancellation drops the future; it is not a panic and -/// therefore cannot be observed by `catch_unwind`. +/// Marks every non-graceful worker exit as terminal, including cancellation +/// before the worker future's first poll and panic unwinding during a poll. struct WorkerCompletionGuard { lifecycle: Arc, armed: bool, @@ -134,7 +131,42 @@ impl WorkerCompletionGuard { impl Drop for WorkerCompletionGuard { fn drop(&mut self) { if self.armed { - self.lifecycle.fail(eyre::eyre!("persistence worker was cancelled")); + let reason = if std::thread::panicking() { + "persistence worker panicked" + } else { + "persistence worker was cancelled" + }; + self.lifecycle.fail(eyre::eyre!(reason)); + } + } +} + +/// Cloneable terminal-state handle independent of table ownership. +/// +/// Create this handle before spawning a supervisor. The table can then still +/// be moved into `close()`, while the supervisor observes either graceful +/// closure or a terminal persistence failure. +#[derive(Clone, Debug)] +pub struct PersistenceMonitor { + lifecycle: Arc, +} + +impl PersistenceMonitor { + /// Waits until the worker fails or closes. + pub async fn wait_for_failure(self) -> PersistenceResult { + loop { + let notified = self.lifecycle.terminal_notify.notified(); + tokio::pin!(notified); + // `notify_waiters` does not retain a permit. Register this waiter + // before reading the lifecycle state so a terminal transition + // cannot land between the state read and the first poll of + // `notified` and be lost forever. + notified.as_mut().enable(); + match self.lifecycle.state() { + PersistenceState::Failed(error) => return Err(error), + PersistenceState::Closed => return Ok(()), + PersistenceState::Running | PersistenceState::Closing => notified.await, + } } } } @@ -583,6 +615,22 @@ mod lifecycle_tests { assert_eq!(batches.load(Ordering::Relaxed), 1); } + #[tokio::test] + async fn detached_monitor_observes_graceful_close() { + let task = PersistenceTask::run_engine(TestEngine { + batches: Arc::new(AtomicUsize::new(0)), + events: Arc::new(ParkingMutex::new(Vec::new())), + config: TestConfig, + failure: TestFailure::None, + }); + let monitor = task.monitor(); + let waiter = tokio::spawn(monitor.wait_for_failure()); + + task.close().await.unwrap(); + + waiter.await.unwrap().unwrap(); + } + #[tokio::test] async fn engine_failure_is_terminal_and_reused_for_later_callers() { let task = PersistenceTask::run_engine(TestEngine { @@ -1044,18 +1092,11 @@ impl } } }; - let supervisor_lifecycle = lifecycle.clone(); // Constructed outside the async block so cancellation before its first // poll still drops the guard and publishes terminal failure. let completion_guard = WorkerCompletionGuard::new(lifecycle.clone()); let task = async move { - if let Err(payload) = AssertUnwindSafe(worker).catch_unwind().await { - // The process panic hook already records the local diagnostic. - // Do not propagate arbitrary panic payloads to API consumers: - // engines can panic with paths or row-derived strings. - drop(payload); - supervisor_lifecycle.fail(eyre::eyre!("persistence worker panicked")); - } + worker.await; completion_guard.disarm(); }; let engine_task_handle = tokio::spawn(task); @@ -1115,28 +1156,19 @@ impl } } + /// Returns a cloneable monitor independent of this task's ownership. + pub fn monitor(&self) -> PersistenceMonitor { + PersistenceMonitor { + lifecycle: self.lifecycle.clone(), + } + } + /// Waits until the worker fails or closes. /// - /// Unlike [`Self::wait_for_ops`], an idle healthy worker does not satisfy - /// this future. Applications can keep it alive as a terminal-state - /// notification without forcing persistence queues to drain or polling - /// their state. A graceful close returns `Ok(())`; a failure returns the - /// shared terminal error. + /// Prefer [`Self::monitor`] when another task must keep waiting while this + /// task is moved into [`Self::close`]. pub async fn wait_for_failure(&self) -> PersistenceResult { - loop { - let notified = self.lifecycle.terminal_notify.notified(); - tokio::pin!(notified); - // `notify_waiters` does not retain a permit. Register this waiter - // before reading the lifecycle state so a terminal transition - // cannot land between the state read and the first poll of - // `notified` and be lost forever. - notified.as_mut().enable(); - match self.lifecycle.state() { - PersistenceState::Failed(error) => return Err(error), - PersistenceState::Closed => return Ok(()), - PersistenceState::Running | PersistenceState::Closing => notified.await, - } - } + self.monitor().wait_for_failure().await } pub async fn close(mut self) -> PersistenceResult { From a69512c674e9869ec2cf39fd2e8603d6761ed614 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 13:05:33 +0700 Subject: [PATCH 13/14] fix: derive persisted page identities --- .../space/index/table_of_contents.rs | 29 +- src/persistence/space/index/unsized_.rs | 383 +++++++++--------- 2 files changed, 226 insertions(+), 186 deletions(-) diff --git a/src/persistence/space/index/table_of_contents.rs b/src/persistence/space/index/table_of_contents.rs index 5703c34..181ef08 100644 --- a/src/persistence/space/index/table_of_contents.rs +++ b/src/persistence/space/index/table_of_contents.rs @@ -113,15 +113,30 @@ where pub fn update_key(&mut self, old_key: &T, new_key: T) where T: Clone + Debug, + { + assert!( + self.try_update_key(old_key, new_key), + "Page with key {old_key:?} not found" + ); + } + + /// Updates a page identity without panicking when the old identity is + /// absent. Batch replay uses this checked form because an absent key is a + /// persistence invariant failure that must surface through `Result`. + pub fn try_update_key(&mut self, old_key: &T, new_key: T) -> bool + where + T: Clone, { let page = self.get_current_page_mut(); if page.inner.update_key(old_key, new_key.clone()).is_none() { for page in self.pages.iter_mut() { if page.inner.update_key(old_key, new_key.clone()).is_some() { - return; + return true; } } - panic!("Page with key {old_key:?} not found"); + false + } else { + true } } @@ -229,6 +244,16 @@ mod tests { ); } + #[test] + fn checked_update_reports_a_missing_identity_without_mutating_the_toc() { + let mut toc = IndexTableOfContents::::new(0.into(), Arc::new(AtomicU32::new(1))); + toc.insert(7, 2.into()); + + assert!(!toc.try_update_key(&8, 9)); + assert_eq!(toc.get(&7), Some(2.into())); + assert_eq!(toc.get(&9), None); + } + #[test] fn insert_more_than_one_page() { let mut toc = IndexTableOfContents::::new(0.into(), Arc::new(AtomicU32::new(0))); diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index 7bd2c4b..cb4ab5a 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -28,6 +28,10 @@ use crate::persistence::space::BatchChangeEvent; use crate::persistence::{IndexTableOfContents, SpaceIndex, SpaceIndexOps}; use crate::prelude::WT_INDEX_EXTENSION; +// Persistence batches are limited to 16 source pages. Allow one additional +// split page per source page while keeping alias bookkeeping entirely inline. +const MAX_BATCH_ALIASED_PAGES: usize = 32; + #[derive(Debug)] pub struct SpaceIndexUnsized { space_id: SpaceId, @@ -403,19 +407,25 @@ where match &ev { ChangeEvent::InsertAt { max_value, .. } | ChangeEvent::RemoveAt { max_value, .. } => { let event_page_key = (max_value.key.clone(), max_value.value); - - let page_index = self - .table_of_contents - .get(&event_page_key) - .or_else(|| page_aliases.get(&event_page_key)) - .ok_or_else(|| { - eyre!( - "unsized index event for {event_page_key:?} references a missing page (toc_segments={}, buffered_pages={}, aliases={})", - self.table_of_contents.pages.len(), - pages.len(), - page_aliases.len() - ) - })?; + // A direct TOC hit means the event key is the page's + // canonical pre-event identity. An alias hit carries the + // current canonical identity captured when the alias was + // installed. This lets us compare the actual post-apply + // identity without predicting DataBucket's mutation rules. + let (page_index, aliased_page_key) = if let Some(page_index) = + self.table_of_contents.get(&event_page_key) + { + (page_index, None) + } else if let Some((page_index, current_key)) = page_aliases.resolve(&event_page_key) { + (page_index, Some(current_key.clone())) + } else { + return Err(eyre!( + "unsized index event references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + self.table_of_contents.pages.len(), + pages.len(), + page_aliases.len() + )); + }; let page = pages.get_mut(&page_index); let page_to_update = if let Some(page) = page { page @@ -430,50 +440,12 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; - // `apply_change_event` can only rewrite `node_id` in these - // exact cases. Avoid cloning an unsized maximum for the - // overwhelmingly common interior insert/remove path. - let can_change_max = match &ev { - ChangeEvent::InsertAt { value, index, .. } => { - value.key > page_to_update.inner.node_id.key - || *index == page_to_update.inner.slots_size as usize - } - ChangeEvent::RemoveAt { - max_value, - value, - index, - .. - } => { - value == max_value - && *index != 0 - && page_to_update.inner.slots_size != 0 - && *index == page_to_update.inner.slots_size as usize - 1 - } - _ => false, - }; - #[cfg(debug_assertions)] - let debug_pre_event_page_key = ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ); - let pre_event_page_key = can_change_max.then(|| { - ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ) - }); + let canonical_page_key = aliased_page_key.as_ref().unwrap_or(&event_page_key); page_to_update.inner.apply_change_event(ev.clone())?; - #[cfg(debug_assertions)] - debug_assert!( - can_change_max - || (page_to_update.inner.node_id.key == debug_pre_event_page_key.0 - && page_to_update.inner.node_id.link == debug_pre_event_page_key.1), - "data_bucket changed an unsized page identity outside WorkTable's can_change_max predicate" - ); - if let Some(pre_event_page_key) = pre_event_page_key - && (page_to_update.inner.node_id.key != pre_event_page_key.0 - || page_to_update.inner.node_id.link != pre_event_page_key.1) + if page_to_update.inner.node_id.key != canonical_page_key.0 + || page_to_update.inner.node_id.link != canonical_page_key.1 { + let pre_event_page_key = aliased_page_key.unwrap_or_else(|| event_page_key.clone()); let updated_page_key = ( page_to_update.inner.node_id.key.clone(), page_to_update.inner.node_id.link, @@ -482,8 +454,16 @@ where // `event_page_key` may be a historical alias, so using // it as the canonical update target can remove or // rewrite the wrong segment. - self.table_of_contents.update_key(&pre_event_page_key, updated_page_key); - page_aliases.replace(page_index, [event_page_key, pre_event_page_key])?; + if !self + .table_of_contents + .try_update_key(&pre_event_page_key, updated_page_key.clone()) + { + return Err(eyre!( + "unsized index page identity is absent from the table of contents (page={page_index:?}, toc_segments={})", + self.table_of_contents.pages.len() + )); + } + page_aliases.replace(page_index, updated_page_key, event_page_key, pre_event_page_key)?; } } ChangeEvent::CreateNode { event_id: _, max_value } => { @@ -512,18 +492,20 @@ where } => { let event_page_key = (max_value.key.clone(), max_value.value); - let page_index = self - .table_of_contents - .get(&event_page_key) - .or_else(|| page_aliases.get(&event_page_key)) - .ok_or_else(|| { - eyre!( - "unsized index split for {event_page_key:?} references a missing page (toc_segments={}, buffered_pages={}, aliases={})", - self.table_of_contents.pages.len(), - pages.len(), - page_aliases.len() - ) - })?; + let (page_index, aliased_page_key) = if let Some(page_index) = + self.table_of_contents.get(&event_page_key) + { + (page_index, None) + } else if let Some((page_index, current_key)) = page_aliases.resolve(&event_page_key) { + (page_index, Some(current_key.clone())) + } else { + return Err(eyre!( + "unsized index split references a missing page (toc_segments={}, buffered_pages={}, aliases={})", + self.table_of_contents.pages.len(), + pages.len(), + page_aliases.len() + )); + }; let page = pages.get_mut(&page_index); let page_to_update = if let Some(page) = page { page @@ -538,10 +520,15 @@ where .get_mut(&page_index) .expect("should be available as was just inserted before") }; - let pre_split_page_key = ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ); + let canonical_page_key = aliased_page_key.as_ref().unwrap_or(&event_page_key); + if page_to_update.inner.node_id.key != canonical_page_key.0 + || page_to_update.inner.node_id.link != canonical_page_key.1 + { + return Err(eyre!( + "unsized index split found a buffered page with a mismatched identity (page={page_index:?})" + )); + } + let pre_split_page_key = aliased_page_key.unwrap_or_else(|| event_page_key.clone()); let splitted_page = page_to_update.inner.split(*split_index); let new_page_id = if let Some(id) = self.table_of_contents.pop_empty_page_id() { @@ -550,17 +537,20 @@ where self.next_page_id.fetch_add(1, Ordering::Relaxed).into() }; - self.table_of_contents.update_key( - &pre_split_page_key, - ( - page_to_update.inner.node_id.key.clone(), - page_to_update.inner.node_id.link, - ), - ); - self.table_of_contents.insert( - (splitted_page.node_id.key.clone(), splitted_page.node_id.link), - new_page_id, + let left_page_key = ( + page_to_update.inner.node_id.key.clone(), + page_to_update.inner.node_id.link, ); + if !self + .table_of_contents + .try_update_key(&pre_split_page_key, left_page_key) + { + return Err(eyre!( + "unsized index split identity is absent from the table of contents (page={page_index:?})" + )); + } + let right_page_key = (splitted_page.node_id.key.clone(), splitted_page.node_id.link); + self.table_of_contents.insert(right_page_key.clone(), new_page_id); let header = GeneralHeader::new(new_page_id, PageType::Index, self.space_id); let general_page = GeneralPage { inner: splitted_page, @@ -571,7 +561,7 @@ where // A following remove/insert pair can still name it even // after the remove temporarily lowers that maximum. page_aliases.remove_page(page_index); - page_aliases.replace(new_page_id, [event_page_key, pre_split_page_key])?; + page_aliases.replace(new_page_id, right_page_key, event_page_key, pre_split_page_key)?; } } } @@ -590,101 +580,98 @@ where } struct PageAliases { - by_key: HashMap<(T, Link), PageId>, - by_page: HashMap>, + entries: [Option>; MAX_BATCH_ALIASED_PAGES], } -struct PageAliasKeys { - slots: [Option<(T, Link)>; 2], -} - -impl Default for PageAliasKeys { - fn default() -> Self { - Self { slots: [None, None] } - } -} - -impl PageAliasKeys { - fn contains(&self, key: &(T, Link)) -> bool { - self.slots.iter().flatten().any(|stored| stored == key) - } - - fn push(&mut self, key: (T, Link)) -> bool { - let Some(slot) = self.slots.iter_mut().find(|slot| slot.is_none()) else { - return false; - }; - *slot = Some(key); - true - } - - fn len(&self) -> usize { - self.slots.iter().flatten().count() - } - - fn iter(&self) -> impl Iterator { - self.slots.iter().flatten() - } +struct PageAliasEntry { + page_id: PageId, + current_key: (T, Link), + aliases: [Option<(T, Link)>; 2], } impl Default for PageAliases { fn default() -> Self { Self { - by_key: HashMap::new(), - by_page: HashMap::new(), + entries: std::array::from_fn(|_| None), } } } -impl PageAliases { +impl PageAliases { + fn resolve(&self, key: &(T, Link)) -> Option<(PageId, &(T, Link))> { + self.entries.iter().flatten().find_map(|entry| { + entry + .aliases + .iter() + .flatten() + .any(|alias| alias == key) + .then_some((entry.page_id, &entry.current_key)) + }) + } + + #[cfg(test)] fn get(&self, key: &(T, Link)) -> Option { - self.by_key.get(key).copied() + self.resolve(key).map(|(page_id, _)| page_id) } fn len(&self) -> usize { - self.by_key.len() + self.entries + .iter() + .flatten() + .map(|entry| entry.aliases.iter().flatten().count()) + .sum() } - fn remove_page(&mut self, page_id: PageId) { - let Some(keys) = self.by_page.remove(&page_id) else { - return; - }; - for key in keys.slots.into_iter().flatten() { - self.by_key.remove(&key); - } + #[cfg(test)] + fn page_len(&self) -> usize { + self.entries.iter().flatten().count() } - fn replace(&mut self, page_id: PageId, keys: impl IntoIterator) -> eyre::Result<()> { - let mut page_keys = PageAliasKeys::default(); - for key in keys { - if page_keys.contains(&key) { - continue; - } - if !page_keys.push(key) { - return Err(eyre!("page {page_id:?} exceeded the two transitional alias invariant")); - } + fn remove_page(&mut self, page_id: PageId) { + if let Some(slot) = self + .entries + .iter_mut() + .find(|slot| slot.as_ref().is_some_and(|entry| entry.page_id == page_id)) + { + *slot = None; } + } - // `(value, Link)` identifies one physical index entry, so it cannot be - // owned by another page. Enforce that invariant in release builds - // before changing either direction, rather than leaving the maps - // inconsistent if a future event-shape regression violates it. - for key in page_keys.iter() { - if let Some(previous_page) = self.by_key.get(key) - && *previous_page != page_id + fn replace( + &mut self, + page_id: PageId, + current_key: (T, Link), + event_key: (T, Link), + pre_event_key: (T, Link), + ) -> eyre::Result<()> { + let first_alias = (event_key != current_key).then_some(event_key); + let second_alias = + (pre_event_key != current_key && first_alias.as_ref() != Some(&pre_event_key)).then_some(pre_event_key); + + for alias in [first_alias.as_ref(), second_alias.as_ref()].into_iter().flatten() { + if let Some(owner) = + self.entries.iter().flatten().find(|entry| { + entry.page_id != page_id && entry.aliases.iter().flatten().any(|stored| stored == alias) + }) { return Err(eyre!( - "page alias ownership collision between {previous_page:?} and {page_id:?}" + "page alias ownership collision between {:?} and {page_id:?}", + owner.page_id )); } } - self.remove_page(page_id); - for key in page_keys.iter() { - self.by_key.insert(key.clone(), page_id); - } - if page_keys.len() != 0 { - self.by_page.insert(page_id, page_keys); - } + let slot_index = self + .entries + .iter() + .position(|slot| slot.as_ref().is_some_and(|entry| entry.page_id == page_id)) + .or_else(|| self.entries.iter().position(Option::is_none)) + .ok_or_else(|| eyre!("unsized index batch exceeded its inline alias-page capacity"))?; + self.entries[slot_index] = Some(PageAliasEntry { + page_id, + current_key, + aliases: [first_alias, second_alias], + }); Ok(()) } } @@ -706,13 +693,20 @@ mod alias_tests { let page_id = PageId::from(7); let mut aliases = PageAliases::default(); for revision in 0..1_000 { + let current = (format!("current-{revision}"), link(revision + 1_000)); aliases - .replace(page_id, [(format!("key-{revision}"), link(revision))]) + .replace( + page_id, + current, + (format!("key-{revision}"), link(revision)), + (format!("key-{revision}"), link(revision)), + ) .unwrap(); } assert_eq!(aliases.len(), 1); assert_eq!(aliases.get(&("key-999".into(), link(999))), Some(page_id)); + assert_eq!(aliases.page_len(), 1); } #[test] @@ -720,17 +714,26 @@ mod alias_tests { let old_page = PageId::from(3); let right_page = PageId::from(4); let mut aliases = PageAliases::default(); - aliases.replace(old_page, [("older".to_string(), link(1))]).unwrap(); + aliases + .replace( + old_page, + ("old-current".to_string(), link(9)), + ("older".to_string(), link(1)), + ("older".to_string(), link(1)), + ) + .unwrap(); aliases.remove_page(old_page); aliases .replace( right_page, - [("event".to_string(), link(2)), ("current".to_string(), link(3))], + ("right-current".to_string(), link(4)), + ("event".to_string(), link(2)), + ("pre-split".to_string(), link(3)), ) .unwrap(); assert_eq!(aliases.len(), 2); - assert_eq!(aliases.by_page.get(&right_page).map(PageAliasKeys::len), Some(2)); + assert_eq!(aliases.page_len(), 1); } #[test] @@ -742,36 +745,45 @@ mod alias_tests { let mut aliases = PageAliases::default(); aliases - .replace(right_page, [pre_split.clone(), post_split.clone()]) + .replace(right_page, post_split.clone(), pre_split.clone(), pre_split.clone()) + .unwrap(); + aliases + .replace(right_page, after_remove.clone(), pre_split.clone(), post_split.clone()) .unwrap(); - aliases.replace(right_page, [pre_split.clone(), after_remove]).unwrap(); assert_eq!(aliases.get(&pre_split), Some(right_page)); - assert_eq!(aliases.get(&post_split), None); + assert_eq!(aliases.get(&post_split), Some(right_page)); + assert_eq!( + aliases.resolve(&pre_split).map(|(_, current)| current), + Some(&after_remove) + ); } #[test] - fn replacing_many_pages_preserves_constant_work_per_page() { + fn inline_capacity_fails_without_overwriting_existing_pages() { let mut aliases = PageAliases::default(); - for page in 1..=1_000 { + for page in 1..=MAX_BATCH_ALIASED_PAGES as u32 { aliases - .replace(PageId::from(page), [(format!("old-{page}"), link(page))]) + .replace( + PageId::from(page), + (format!("current-{page}"), link(page + 1_000)), + (format!("old-{page}"), link(page)), + (format!("old-{page}"), link(page)), + ) .unwrap(); } - for page in 1..=1_000 { + assert_eq!(aliases.page_len(), MAX_BATCH_ALIASED_PAGES); + assert!( aliases - .replace(PageId::from(page), [(format!("new-{page}"), link(page + 1_000))]) - .unwrap(); - } - - assert_eq!(aliases.len(), 1_000); - assert_eq!(aliases.by_page.len(), 1_000); - for page in 1..=1_000 { - assert_eq!( - aliases.get(&(format!("new-{page}"), link(page + 1_000))), - Some(PageId::from(page)) - ); - } + .replace( + PageId::from(MAX_BATCH_ALIASED_PAGES as u32 + 1), + ("overflow-current".to_string(), link(10_000)), + ("overflow-old".to_string(), link(10_001)), + ("overflow-old".to_string(), link(10_001)), + ) + .is_err() + ); + assert_eq!(aliases.get(&("old-1".into(), link(1))), Some(PageId::from(1))); } #[test] @@ -780,23 +792,26 @@ mod alias_tests { let second_page = PageId::from(2); let shared = ("shared".to_string(), link(1)); let mut aliases = PageAliases::default(); - aliases.replace(first_page, [shared.clone()]).unwrap(); + aliases + .replace( + first_page, + ("first-current".to_string(), link(9)), + shared.clone(), + shared.clone(), + ) + .unwrap(); - assert!(aliases.replace(second_page, [shared.clone()]).is_err()); - assert_eq!(aliases.get(&shared), Some(first_page)); assert!( aliases .replace( second_page, - [ - ("one".to_string(), link(2)), - ("two".to_string(), link(3)), - ("three".to_string(), link(4)), - ], + ("second-current".to_string(), link(10)), + shared.clone(), + shared.clone(), ) .is_err() ); assert_eq!(aliases.get(&shared), Some(first_page)); - assert!(!aliases.by_page.contains_key(&second_page)); + assert_eq!(aliases.page_len(), 1); } } From ff47c43396e66ddbd7093115aa60e11c6eae9770 Mon Sep 17 00:00:00 2001 From: meh Date: Sat, 8 Aug 2026 13:22:25 +0700 Subject: [PATCH 14/14] fix: preserve oversized alias batches --- src/persistence/space/index/unsized_.rs | 152 ++++++++++++++++-------- 1 file changed, 100 insertions(+), 52 deletions(-) diff --git a/src/persistence/space/index/unsized_.rs b/src/persistence/space/index/unsized_.rs index cb4ab5a..b7efef1 100644 --- a/src/persistence/space/index/unsized_.rs +++ b/src/persistence/space/index/unsized_.rs @@ -28,9 +28,10 @@ use crate::persistence::space::BatchChangeEvent; use crate::persistence::{IndexTableOfContents, SpaceIndex, SpaceIndexOps}; use crate::prelude::WT_INDEX_EXTENSION; -// Persistence batches are limited to 16 source pages. Allow one additional -// split page per source page while keeping alias bookkeeping entirely inline. -const MAX_BATCH_ALIASED_PAGES: usize = 32; +// Normal persistence batches begin at 16 source pages. Keep that common case +// inline while allowing analyzer retries to grow beyond it without turning a +// recovery batch into a terminal capacity error. +const INLINE_BATCH_ALIASED_PAGES: usize = 16; #[derive(Debug)] pub struct SpaceIndexUnsized { @@ -63,6 +64,21 @@ where + Debug + for<'a> rkyv::bytecheck::CheckBytes>, { + fn resolve_batch_page( + &self, + aliases: &PageAliases, + event_page_key: &(T, Link), + ) -> Option<(PageId, Option<(T, Link)>)> { + self.table_of_contents + .get(event_page_key) + .map(|page_id| (page_id, None)) + .or_else(|| { + aliases + .resolve(event_page_key) + .map(|(page_id, current_key)| (page_id, Some(current_key.clone()))) + }) + } + fn compact_page_if_needed(page: &mut UnsizedIndexPage) -> eyre::Result<()> { let persisted_size = UnsizedIndexPageUtility::::persisted_size(page.slots_size as usize, page.node_id_size as usize) @@ -412,13 +428,8 @@ where // current canonical identity captured when the alias was // installed. This lets us compare the actual post-apply // identity without predicting DataBucket's mutation rules. - let (page_index, aliased_page_key) = if let Some(page_index) = - self.table_of_contents.get(&event_page_key) - { - (page_index, None) - } else if let Some((page_index, current_key)) = page_aliases.resolve(&event_page_key) { - (page_index, Some(current_key.clone())) - } else { + let Some((page_index, aliased_page_key)) = self.resolve_batch_page(&page_aliases, &event_page_key) + else { return Err(eyre!( "unsized index event references a missing page (toc_segments={}, buffered_pages={}, aliases={})", self.table_of_contents.pages.len(), @@ -463,6 +474,11 @@ where self.table_of_contents.pages.len() )); } + if self.table_of_contents.get(&updated_page_key) != Some(page_index) { + return Err(eyre!( + "unsized index page identity update did not become canonical (page={page_index:?})" + )); + } page_aliases.replace(page_index, updated_page_key, event_page_key, pre_event_page_key)?; } } @@ -492,13 +508,8 @@ where } => { let event_page_key = (max_value.key.clone(), max_value.value); - let (page_index, aliased_page_key) = if let Some(page_index) = - self.table_of_contents.get(&event_page_key) - { - (page_index, None) - } else if let Some((page_index, current_key)) = page_aliases.resolve(&event_page_key) { - (page_index, Some(current_key.clone())) - } else { + let Some((page_index, aliased_page_key)) = self.resolve_batch_page(&page_aliases, &event_page_key) + else { return Err(eyre!( "unsized index split references a missing page (toc_segments={}, buffered_pages={}, aliases={})", self.table_of_contents.pages.len(), @@ -551,6 +562,11 @@ where } let right_page_key = (splitted_page.node_id.key.clone(), splitted_page.node_id.link); self.table_of_contents.insert(right_page_key.clone(), new_page_id); + if self.table_of_contents.get(&right_page_key) != Some(new_page_id) { + return Err(eyre!( + "unsized index split identity did not become canonical (page={new_page_id:?})" + )); + } let header = GeneralHeader::new(new_page_id, PageType::Index, self.space_id); let general_page = GeneralPage { inner: splitted_page, @@ -579,8 +595,16 @@ where } } +/// Transitional event identities for pages whose canonical TOC key changed +/// earlier in the same CDC batch. +/// +/// Each page owns at most the event identity and its actual pre-event identity. +/// The current canonical identity is retained alongside them so alias lookup +/// never has to predict page mutation semantics. Normal batches use the inline +/// slots; analyzer retries may spill into `overflow` without losing events. struct PageAliases { - entries: [Option>; MAX_BATCH_ALIASED_PAGES], + inline: [Option>; INLINE_BATCH_ALIASED_PAGES], + overflow: Vec>, } struct PageAliasEntry { @@ -592,14 +616,19 @@ struct PageAliasEntry { impl Default for PageAliases { fn default() -> Self { Self { - entries: std::array::from_fn(|_| None), + inline: std::array::from_fn(|_| None), + overflow: Vec::new(), } } } impl PageAliases { + fn entries(&self) -> impl Iterator> { + self.inline.iter().flatten().chain(self.overflow.iter()) + } + fn resolve(&self, key: &(T, Link)) -> Option<(PageId, &(T, Link))> { - self.entries.iter().flatten().find_map(|entry| { + self.entries().find_map(|entry| { entry .aliases .iter() @@ -615,25 +644,23 @@ impl PageAliases { } fn len(&self) -> usize { - self.entries - .iter() - .flatten() - .map(|entry| entry.aliases.iter().flatten().count()) - .sum() + self.entries().map(|entry| entry.aliases.iter().flatten().count()).sum() } #[cfg(test)] fn page_len(&self) -> usize { - self.entries.iter().flatten().count() + self.entries().count() } fn remove_page(&mut self, page_id: PageId) { if let Some(slot) = self - .entries + .inline .iter_mut() .find(|slot| slot.as_ref().is_some_and(|entry| entry.page_id == page_id)) { *slot = None; + } else if let Some(index) = self.overflow.iter().position(|entry| entry.page_id == page_id) { + self.overflow.swap_remove(index); } } @@ -648,11 +675,15 @@ impl PageAliases { let second_alias = (pre_event_key != current_key && first_alias.as_ref() != Some(&pre_event_key)).then_some(pre_event_key); + if first_alias.is_none() && second_alias.is_none() { + self.remove_page(page_id); + return Ok(()); + } + for alias in [first_alias.as_ref(), second_alias.as_ref()].into_iter().flatten() { - if let Some(owner) = - self.entries.iter().flatten().find(|entry| { - entry.page_id != page_id && entry.aliases.iter().flatten().any(|stored| stored == alias) - }) + if let Some(owner) = self + .entries() + .find(|entry| entry.page_id != page_id && entry.aliases.iter().flatten().any(|stored| stored == alias)) { return Err(eyre!( "page alias ownership collision between {:?} and {page_id:?}", @@ -661,17 +692,24 @@ impl PageAliases { } } - let slot_index = self - .entries - .iter() - .position(|slot| slot.as_ref().is_some_and(|entry| entry.page_id == page_id)) - .or_else(|| self.entries.iter().position(Option::is_none)) - .ok_or_else(|| eyre!("unsized index batch exceeded its inline alias-page capacity"))?; - self.entries[slot_index] = Some(PageAliasEntry { + let entry = PageAliasEntry { page_id, current_key, aliases: [first_alias, second_alias], - }); + }; + if let Some(slot) = self + .inline + .iter_mut() + .find(|slot| slot.as_ref().is_some_and(|stored| stored.page_id == page_id)) + { + *slot = Some(entry); + } else if let Some(slot) = self.overflow.iter_mut().find(|stored| stored.page_id == page_id) { + *slot = entry; + } else if let Some(slot) = self.inline.iter_mut().find(|slot| slot.is_none()) { + *slot = Some(entry); + } else { + self.overflow.push(entry); + } Ok(()) } } @@ -709,6 +747,20 @@ mod alias_tests { assert_eq!(aliases.page_len(), 1); } + #[test] + fn canonical_only_transition_stores_no_alias_entry() { + let page_id = PageId::from(7); + let canonical = ("current".to_string(), link(1)); + let mut aliases = PageAliases::default(); + + aliases + .replace(page_id, canonical.clone(), canonical.clone(), canonical) + .unwrap(); + + assert_eq!(aliases.page_len(), 0); + assert_eq!(aliases.len(), 0); + } + #[test] fn split_keeps_only_the_two_live_transitional_identities() { let old_page = PageId::from(3); @@ -760,9 +812,10 @@ mod alias_tests { } #[test] - fn inline_capacity_fails_without_overwriting_existing_pages() { + fn batches_beyond_inline_capacity_preserve_every_alias() { let mut aliases = PageAliases::default(); - for page in 1..=MAX_BATCH_ALIASED_PAGES as u32 { + let page_count = INLINE_BATCH_ALIASED_PAGES as u32 + 8; + for page in 1..=page_count { aliases .replace( PageId::from(page), @@ -772,18 +825,13 @@ mod alias_tests { ) .unwrap(); } - assert_eq!(aliases.page_len(), MAX_BATCH_ALIASED_PAGES); - assert!( - aliases - .replace( - PageId::from(MAX_BATCH_ALIASED_PAGES as u32 + 1), - ("overflow-current".to_string(), link(10_000)), - ("overflow-old".to_string(), link(10_001)), - ("overflow-old".to_string(), link(10_001)), - ) - .is_err() - ); + assert_eq!(aliases.page_len(), page_count as usize); + assert_eq!(aliases.overflow.len(), 8); assert_eq!(aliases.get(&("old-1".into(), link(1))), Some(PageId::from(1))); + assert_eq!( + aliases.get(&(format!("old-{page_count}"), link(page_count))), + Some(PageId::from(page_count)) + ); } #[test]