diff --git a/codegen/src/persist_index/generator.rs b/codegen/src/persist_index/generator.rs index 7c7a864..e8419f8 100644 --- a/codegen/src/persist_index/generator.rs +++ b/codegen/src/persist_index/generator.rs @@ -356,24 +356,43 @@ impl Generator { .get(i) .expect("should be available as constructed from same values"); - if is_unsized(&ty.to_string()) { - let node = if is_unique { - quote! { - let node = page - .inner - .get_node() - .into_iter() - .map(|p| IndexPair { - key: p.key, - value: p.value.into(), - }) - .collect(); - let node = UnsizedNode::from_inner(node, #const_name); - #i.attach_node(node); - } - } else { - quote! { - let mut inner: Vec<_> = page + // Index pages carry entries in event-arrival order, not sorted + // order: only bootstrap-written pages happen to be sorted with + // the node maximum last, while pages that were incrementally + // updated through CDC events are arbitrary. Reconstruction + // therefore must sort every page and derive the node maximum + // from the sorted entries; assuming "last entry == node id" + // registers a wrong maximum in the in-memory node index and + // makes every entry above it unreachable. + let sort_page_entries = quote! { + inner.sort_by(|a, b| a.key.cmp(&b.key).then_with(|| a.value.cmp(&b.value))); + }; + // For non-unique indexes one key's duplicates can straddle + // nodes, and the B-tree invariant (every entry in a node + // compares <= that node's registered maximum) only survives + // reconstruction if a later node's segment of the same key + // compares greater than the previous node's maximum. So nodes + // are attached in ascending order of their true maximum and + // discriminators keep growing across node boundaries within + // one key. Discriminator 0 is reserved for the range infimum + // and u64::MAX for the supremum, so numbering starts at 1 and + // is capped below the supremum. + let multi_reconstruct = |attach: TokenStream| { + quote! { + let mut nodes: Vec<(IndexPair<#ty, OffsetEqLink>, Vec>)> = Vec::new(); + for page in persisted.#i.1 { + // The persisted node id is the node maximum the + // on-disk table of contents knows this node by. + // The reconstructed node must end with exactly + // this entry: CDC events emitted after the reload + // address nodes by their maximum, and a node whose + // maximum drifted from the TOC key can no longer + // be resolved by the space index. + let node_id = IndexPair { + key: page.inner.node_id.key.clone(), + value: OffsetEqLink(page.inner.node_id.link), + }; + let mut inner: Vec> = page .inner .get_node() .into_iter() @@ -382,41 +401,75 @@ impl Generator { value: OffsetEqLink(p.value), }) .collect(); - - let node_id = inner.pop(); - - let mut sorted: Vec<_> = inner.into_iter() - .map(|p| IndexMultiPair { - key: p.key, - value: p.value, - discriminator: 0, + inner.sort_by(|a, b| { + a.key.cmp(&b.key).then_with(|| { + let a_is_id = a.key == node_id.key && a.value == node_id.value; + let b_is_id = b.key == node_id.key && b.value == node_id.value; + a_is_id.cmp(&b_is_id).then_with(|| a.value.cmp(&b.value)) }) - .collect(); - sorted.sort(); - - let mut current_discriminator = 1u64; - for entry in sorted.iter_mut() { - entry.discriminator = current_discriminator.min(u64::MAX - 1); - current_discriminator += 1; + }); + if !inner.is_empty() { + nodes.push((node_id, inner)); } - if let Some(node_id) = node_id { - sorted.push(IndexMultiPair { key: node_id.key, value: node_id.value, discriminator: u64::MAX - 1 }); + } + nodes.sort_by(|(a, _), (b, _)| a.key.cmp(&b.key).then_with(|| a.value.cmp(&b.value))); + let mut prev_key = None; + let mut next_discriminator = 1u64; + for (_, inner) in nodes { + let mut sorted = Vec::with_capacity(inner.len()); + for p in inner { + if prev_key.as_ref() != Some(&p.key) { + prev_key = Some(p.key.clone()); + next_discriminator = 1; + } + sorted.push(IndexMultiPair { + key: p.key, + value: p.value, + discriminator: next_discriminator.min(u64::MAX - 1), + }); + next_discriminator = next_discriminator.saturating_add(1); } + #attach + } + } + }; + if is_unsized(&ty.to_string()) { + if is_unique { + quote! { + let #i: #t<_, OffsetEqLink, UnsizedNode<_>> = #t::with_maximum_node_size(#const_name); + for page in persisted.#i.1 { + let mut inner: Vec> = page + .inner + .get_node() + .into_iter() + .map(|p| IndexPair { + key: p.key, + value: p.value.into(), + }) + .collect(); + #sort_page_entries + let node = UnsizedNode::from_inner(inner, #const_name); + #i.attach_node(node); + } + } + } else { + let attach = quote! { let node = UnsizedNode::from_inner(sorted, #const_name); #i.attach_multi_node(node); + }; + let body = multi_reconstruct(attach); + quote! { + let #i: #t<_, OffsetEqLink, UnsizedNode<_>> = #t::with_maximum_node_size(#const_name); + #body } - }; + } + } else if is_unique { quote! { - let #i: #t<_, OffsetEqLink, UnsizedNode<_>> = #t::with_maximum_node_size(#const_name); + let size = get_index_page_size_from_data_length::<#ty>(#const_name); + let #i: #t<_, OffsetEqLink> = #t::with_maximum_node_size(size); for page in persisted.#i.1 { - #node - } - } - } else { - let node = if is_unique { - quote! { - let node = page + let mut inner: Vec> = page .inner .get_node() .into_iter() @@ -425,49 +478,19 @@ impl Generator { value: p.value.into(), }) .collect(); - #i.attach_node(node); - } - } else { - quote! { - let mut inner: Vec<_> = page - .inner - .get_node() - .into_iter() - .map(|p| IndexPair { - key: p.key, - value: OffsetEqLink(p.value), - }) - .collect(); - - let node_id = inner.pop(); - - let mut sorted: Vec<_> = inner.into_iter() - .map(|p| IndexMultiPair { - key: p.key, - value: p.value, - discriminator: 0, - }) - .collect(); - sorted.sort(); - - let mut current_discriminator = 1u64; - for entry in sorted.iter_mut() { - entry.discriminator = current_discriminator.min(u64::MAX - 1); - current_discriminator += 1; - } - if let Some(node_id) = node_id { - sorted.push(IndexMultiPair { key: node_id.key, value: node_id.value, discriminator: u64::MAX - 1 }); - } - - #i.attach_multi_node(sorted); + #sort_page_entries + #i.attach_node(inner); } + } + } else { + let attach = quote! { + #i.attach_multi_node(sorted); }; + let body = multi_reconstruct(attach); quote! { let size = get_index_page_size_from_data_length::<#ty>(#const_name); let #i: #t<_, OffsetEqLink> = #t::with_maximum_node_size(size); - for page in persisted.#i.1 { - #node - } + #body } } }) diff --git a/tests/persistence/duplicate_key_index_reload.rs b/tests/persistence/duplicate_key_index_reload.rs new file mode 100644 index 0000000..c25ff0c --- /dev/null +++ b/tests/persistence/duplicate_key_index_reload.rs @@ -0,0 +1,106 @@ +use std::time::Duration; + +use tokio::time::timeout; + +use crate::remove_dir_if_exists; +use worktable::prelude::PersistedWorkTable; +use worktable::prelude::*; +use worktable::worktable; + +worktable!( + name: DuplicateKeyReload, + persist: true, + columns: { + id: u64 primary_key autoincrement, + score: u64, + label: String, + }, + indexes: { + score_idx: score, + }, +); + +/// Regression test for lossy reconstruction of duplicate-key secondary +/// indexes. +/// +/// Index pages store entries in event-arrival order, but `from_persisted` +/// used to treat the last entry of every page as the node's maximum and +/// re-append it with discriminator `u64::MAX - 1`. For pages that were +/// incrementally updated through CDC events (any bulk load), that "maximum" +/// was an arbitrary entry, so the reconstructed in-memory node index +/// registered wrong node maxima and every entry sorting above one became +/// unreachable through `select_by_*` — around 12% of the entries in this +/// workload, although `select_all` still returned every row. Reconstruction +/// now sorts each page, orders nodes by their true maximum, and assigns +/// discriminators that keep growing across node boundaries within one key, +/// which restores the B-tree ordering invariant even when one key's +/// duplicates straddle nodes. +#[test] +fn test_duplicate_key_secondary_index_survives_reload() { + const ROWS: u64 = 10_000; + const KEYS: u64 = 97; + + let config = DiskConfig::new_with_table_name( + "tests/data/duplicate_key_index_reload/persisted", + DuplicateKeyReloadWorkTable::name_snake_case(), + DuplicateKeyReloadWorkTable::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("tests/data/duplicate_key_index_reload/persisted".to_string()).await; + + { + let engine = DuplicateKeyReloadPersistenceEngine::new(config.clone()).await.unwrap(); + let table = DuplicateKeyReloadWorkTable::load(engine).await.unwrap(); + + for i in 0..ROWS { + table + .insert(DuplicateKeyReloadRow { + id: i, + score: i % KEYS, + label: format!("row-{i}-{}", "x".repeat((i % 50) as usize)), + }) + .unwrap(); + } + + // The in-memory index is complete before shutdown. + for s in 0..KEYS { + let expected = ROWS / KEYS + u64::from(s < ROWS % KEYS); + assert_eq!( + table.select_by_score(s).execute().unwrap().len() as u64, + expected, + "in-memory index incomplete for key {s}" + ); + } + + timeout(Duration::from_secs(30), table.wait_for_ops()) + .await + .expect("persistence stalled on bulk insert"); + } + { + let engine = DuplicateKeyReloadPersistenceEngine::new(config.clone()).await.unwrap(); + let table = DuplicateKeyReloadWorkTable::load(engine).await.unwrap(); + + // Row data survives the round trip... + assert_eq!(table.select_all().execute().unwrap().len() as u64, ROWS); + + // ...and every row stays reachable through the secondary index, + // with per-key counts intact. + for s in 0..KEYS { + let expected = ROWS / KEYS + u64::from(s < ROWS % KEYS); + assert_eq!( + table.select_by_score(s).execute().unwrap().len() as u64, + expected, + "secondary index lost entries for key {s} across persist+reload" + ); + } + } + }) +} diff --git a/tests/persistence/mod.rs b/tests/persistence/mod.rs index bdcfa80..67e293b 100644 --- a/tests/persistence/mod.rs +++ b/tests/persistence/mod.rs @@ -4,6 +4,7 @@ use worktable::worktable; mod bulk_load_stall; mod concurrent; +mod duplicate_key_index_reload; mod failure; mod index_page; mod read;