Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 20 additions & 14 deletions pgvectorscale/src/access_method/meta_page.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
33 changes: 13 additions & 20 deletions pgvectorscale/src/util/chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@
//! of the data.

use pgrx::{
pg_sys::{BlockNumber, InvalidBlockNumber},
pg_sys::{BlockNumber, InvalidBlockNumber, OffsetNumber},
PgRelation,
};
use rkyv::{Archive, Deserialize, Serialize};
Expand All @@ -31,6 +31,18 @@ struct ChainItemHeader {

const CHAIN_ITEM_HEADER_SIZE: usize = std::mem::size_of::<ArchivedChainItemHeader>();

/// 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,
Expand All @@ -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);
Expand Down