From 7bc186d7ac0380608f83efef606291360ae25648 Mon Sep 17 00:00:00 2001 From: pavanmanishd Date: Thu, 23 Apr 2026 03:14:19 +0530 Subject: [PATCH] fix: hold meta-page lock across reinit and writes in MetaPage::store Parallel index-build workers all call MetaPage::store(false) from finalize_index_build, which released the exclusive buffer lock between reinit and each subsequent write. Concurrent workers could interleave reinit/write and trigger the "offset N != 1" assertion failure on block 0. Restructure store() to hold one WritablePage across reinit, header write, and meta write, committing once at the end. Extract write_chain_item_to_page helper for the single-item, caller-holds-lock case so the chain-item serialization stays in util/chain.rs. Remove the now-unused ChainTapeWriter::reinit. Fixes #264 --- pgvectorscale/src/access_method/meta_page.rs | 34 ++++++++++++-------- pgvectorscale/src/util/chain.rs | 33 ++++++++----------- 2 files changed, 33 insertions(+), 34 deletions(-) diff --git a/pgvectorscale/src/access_method/meta_page.rs b/pgvectorscale/src/access_method/meta_page.rs index eecca125..d5534f3f 100644 --- a/pgvectorscale/src/access_method/meta_page.rs +++ b/pgvectorscale/src/access_method/meta_page.rs @@ -15,8 +15,8 @@ use crate::access_method::graph::start_nodes::StartNodes; use crate::access_method::node::{ReadableNode, WriteableNode}; use crate::access_method::options::TSVIndexOptions; use crate::access_method::stats::WriteStats; -use crate::util::chain::{ChainItemReader, ChainTapeWriter}; -use crate::util::page::{self, PageType}; +use crate::util::chain::{write_chain_item_to_page, ChainItemReader}; +use crate::util::page::{self, PageType, WritablePage}; use crate::util::*; const TSV_MAGIC_NUMBER: u32 = 768756476; //Magic number, random @@ -365,22 +365,28 @@ impl MetaPage { assert!(header.magic_number == TSV_MAGIC_NUMBER); assert!(header.version == TSV_VERSION); - let mut stats = WriteStats::default(); - let mut tape = if first_time { - ChainTapeWriter::new(index, PageType::Meta, &mut stats) + let header_bytes = header.serialize_to_vec(); + let meta_bytes = self.serialize_to_vec(); + + // Hold the exclusive page lock across reinit, header, and meta + // writes so concurrent callers (e.g. parallel index-build workers + // all finalizing at once) can't interleave reinit and write and + // race on block 0. See issue #264. + let mut page = if first_time { + WritablePage::new(index, PageType::Meta) } else { - ChainTapeWriter::reinit(index, PageType::Meta, &mut stats, META_BLOCK_NUMBER) + let mut p = WritablePage::modify(index, META_BLOCK_NUMBER); + p.reinit(PageType::Meta); + p }; - // Serialize the header - let bytes = header.serialize_to_vec(); - let off = tape.write(&bytes); - assert_eq!(off, ItemPointer::new(META_BLOCK_NUMBER, META_HEADER_OFFSET)); + let header_off = write_chain_item_to_page(&mut page, &header_bytes); + assert_eq!(header_off, META_HEADER_OFFSET); + + let meta_off = write_chain_item_to_page(&mut page, &meta_bytes); + assert_eq!(meta_off, META_OFFSET); - // Serialize the meta - let bytes = self.serialize_to_vec(); - let off = tape.write(&bytes); - assert_eq!(off, ItemPointer::new(META_BLOCK_NUMBER, META_OFFSET)); + page.commit(); } unsafe fn load(index: &PgRelation) -> MetaPage { diff --git a/pgvectorscale/src/util/chain.rs b/pgvectorscale/src/util/chain.rs index 4509cc6c..846bfb34 100644 --- a/pgvectorscale/src/util/chain.rs +++ b/pgvectorscale/src/util/chain.rs @@ -11,7 +11,7 @@ //! of the data. use pgrx::{ - pg_sys::{BlockNumber, InvalidBlockNumber}, + pg_sys::{BlockNumber, InvalidBlockNumber, OffsetNumber}, PgRelation, }; use rkyv::{Archive, Deserialize, Serialize}; @@ -31,6 +31,18 @@ struct ChainItemHeader { const CHAIN_ITEM_HEADER_SIZE: usize = std::mem::size_of::(); +/// Append a single chain-item to `page`. Caller retains the exclusive +/// lock; the item must fit on this page (no chaining). Returns the +/// offset number of the new item. +pub fn write_chain_item_to_page(page: &mut WritablePage, data: &[u8]) -> OffsetNumber { + let header = ChainItemHeader { + next: ItemPointer::new_invalid(), + }; + let header_bytes = rkyv::to_bytes::<_, 256>(&header).unwrap(); + let combined = [header_bytes.as_slice(), data].concat(); + page.add_item(&combined) +} + pub struct ChainTapeWriter<'a, S: StatsNodeWrite> { page_type: PageType, index: &'a PgRelation, @@ -53,25 +65,6 @@ impl<'a, S: StatsNodeWrite> ChainTapeWriter<'a, S> { } } - pub fn reinit( - index: &'a PgRelation, - page_type: PageType, - stats: &'a mut S, - block_number: BlockNumber, - ) -> Self { - assert!(page_type.is_chained()); - let mut page = WritablePage::modify(index, block_number); - page.reinit(page_type); - page.commit(); - - Self { - page_type, - index, - current: block_number, - stats, - } - } - /// Write chained data to the tape, returning an `ItemPointer` to the start of the data. pub fn write(&mut self, mut data: &[u8]) -> super::ItemPointer { let mut current_page = WritablePage::modify(self.index, self.current);