diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 18ee4cd..8b3bb8b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -13,16 +13,31 @@ jobs: - uses: dtolnay/rust-toolchain@stable with: components: rustfmt, clippy - targets: thumbv7em-none-eabihf + targets: x86_64-unknown-none - uses: Swatinem/rust-cache@v2 - run: cargo fmt --all -- --check - run: cargo clippy --all-targets --no-default-features -- -D warnings - run: cargo test --all-targets --no-default-features - - run: cargo check --lib --no-default-features --target thumbv7em-none-eabihf - run: cargo clippy --all-targets --features arctic -- -D warnings - run: cargo test --all-targets --features arctic - run: cargo clippy --all-targets --all-features -- -D warnings - run: cargo test --all-targets --all-features + # A crate that says `#![no_std]` and is only ever built on a host with one + # proves nothing. This target has no `std` to find. + - run: cargo build --target x86_64-unknown-none --no-default-features + - run: cargo build --target x86_64-unknown-none --no-default-features --features hydrate + # The backends need an OS for their thread-local storage, so they cannot + # be built for a bare target - but none of them may pull `std` in. + - name: No backend enables std + run: | + for features in arctic congee wti arctic,congee,wti; do + graph=$(cargo tree -e normal --no-default-features --features "$features" -f '{p} [{f}]') + echo "--- $features" + echo "$graph" | grep -E 'arctic-wt|congee-wt|WorkTablesIndex|ps-reclaim|parking_lot' + if echo "$graph" | grep -E '(arctic-wt|congee-wt|WorkTablesIndex|ps-reclaim|parking_lot[a-z_]*) v[^[]*\[[^]]*std'; then + echo "std leaked into --features $features"; exit 1 + fi + done package: name: Release dry run diff --git a/Cargo.toml b/Cargo.toml index 91642ab..ec850f3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "worktable-vec" -version = "0.1.3" +version = "0.1.5" edition = "2024" rust-version = "1.85" license = "MIT OR Apache-2.0" @@ -12,11 +12,41 @@ categories = ["data-structures"] [features] default = [] +# On by default for nobody: this crate is `#![no_std]` and the core never needed +# `std`. The feature exists so a backend that does need one can say so, and so a +# caller without a standard library can prove it is not getting one. +std = ["arctic?/std", "congee?/std", "wti?/std"] +# Load and unload rows as pages. Optional because it is the only thing here +# that needs a serializer. +hydrate = ["dep:rkyv", "dep:embedded-io"] arctic = ["dep:arctic"] congee = ["dep:congee"] wti = ["dep:wti"] [dependencies] -arctic = { package = "arctic-wt", version = "^0.1, >=0.1.9", optional = true } -congee = { package = "congee-wt", version = "^0.4, >=0.4.4", optional = true } -wti = { package = "WorkTablesIndex", version = "^0.0, >=0.0.11", optional = true, default-features = false, features = ["concurrent"] } +rkyv = { version = "0.8.17", default-features = false, features = ["alloc", "bytecheck"], optional = true } +# The I/O is a trait, not a filesystem. `embedded-io` already implements it for +# `&[u8]` and `Vec`, and `embedded-io-adapters` wraps a `std::fs::File`, so +# a caller with an OS gets real files and a caller without one plugs in whatever +# it has. Nothing here needs `std` either way. +embedded-io = { version = "0.7", default-features = false, features = ["alloc"], optional = true } +# `default-features = false` and an explicit SMR backend, so enabling `arctic` +# does not silently enable `std`. 0.1.11 is the first release where that is +# possible: before it, `smr-ps-reclaim` forced `std` on. +arctic = { package = "arctic-wt", version = "0.1.11", optional = true, default-features = false, features = ["smr-ps-reclaim"] } +congee = { package = "congee-wt", version = "0.4.5", optional = true, default-features = false } +wti = { package = "WorkTablesIndex", version = "0.0.13", optional = true, default-features = false, features = ["concurrent"] } + +[dev-dependencies] +# Only the tests need a real file, and this is what turns one into the trait. +embedded-io-adapters = { version = "0.7", features = ["std"] } +rkyv = { version = "0.8.17", features = ["alloc", "bytecheck"] } + +[[example]] +name = "disk_cost" +# It measures the page codec, so it needs it. +required-features = ["hydrate"] + +[[example]] +name = "where_time_goes" +required-features = ["hydrate"] diff --git a/examples/disk_cost.rs b/examples/disk_cost.rs new file mode 100644 index 0000000..e21b88c --- /dev/null +++ b/examples/disk_cost.rs @@ -0,0 +1,82 @@ +//! What a flush and an open cost against a real file. +//! +//! Bounded on purpose: fixed row counts, fixed reps, and the fixture is built +//! once. A benchmark without a ceiling is a hang. + +use std::time::Instant; + +use embedded_io_adapters::std::FromStd; +use worktable_vec::LinearTable; + +const ROWS: usize = 200_000; +const REPS: usize = 5; + +fn main() -> Result<(), Box> { + let mut table = LinearTable::new(); + for n in 0..ROWS as u64 { + table.push( + n, + format!("row {n} with enough text to be worth serializing"), + ); + } + let path = std::env::temp_dir().join("worktable-vec-disk-cost.wtv"); + + let mut wrote = Vec::new(); + let mut read = Vec::new(); + let mut bytes = 0u64; + + for _ in 0..REPS { + let now = Instant::now(); + { + let file = std::fs::File::create(&path)?; + table.unload_to(&mut FromStd::new(file))?; + } + wrote.push(now.elapsed()); + bytes = std::fs::metadata(&path)?.len(); + + let now = Instant::now(); + let back = + LinearTable::::load_from(&mut FromStd::new(std::fs::File::open(&path)?)) + .map_err(|error| format!("{error}"))?; + read.push(now.elapsed()); + assert_eq!(back.len(), ROWS); + } + + wrote.sort(); + read.sort(); + let flush = wrote[REPS / 2]; + let open = read[REPS / 2]; + let mb = bytes as f64 / (1024.0 * 1024.0); + + println!( + "\n{ROWS} rows, {mb:.1} MiB on disk, {} pages, median of {REPS}\n", + bytes as usize / (4096 * 4) + ); + println!( + " flush {:>7.1} ms {:>7.1} MiB/s {:>6.0} ns/row", + flush.as_secs_f64() * 1e3, + mb / flush.as_secs_f64(), + flush.as_secs_f64() * 1e9 / ROWS as f64 + ); + println!( + " open {:>7.1} ms {:>7.1} MiB/s {:>6.0} ns/row", + open.as_secs_f64() * 1e3, + mb / open.as_secs_f64(), + open.as_secs_f64() * 1e9 / ROWS as f64 + ); + + // The floor: what the same rows cost with no pages, no checksum and no + // fingerprint, just one archive straight to the file. Anything this codec + // adds shows up as the gap. + let now = Instant::now(); + let raw = rkyv::to_bytes::(&table.rows().to_vec())?; + let encode = now.elapsed(); + println!( + "\n rkyv alone, no pages {:>7.1} ms encode, {} MiB", + encode.as_secs_f64() * 1e3, + raw.len() / (1024 * 1024) + ); + + let _ = std::fs::remove_file(&path); + Ok(()) +} diff --git a/examples/where_time_goes.rs b/examples/where_time_goes.rs new file mode 100644 index 0000000..f62b390 --- /dev/null +++ b/examples/where_time_goes.rs @@ -0,0 +1,61 @@ +//! What an open actually spends its time on, before anyone parallelises it. +//! +//! A speedup on a synthetic loop says nothing about this path. Page decode is +//! the only embarrassingly parallel part, and it is only worth threads if it is +//! most of the time. + +use std::hint::black_box; +use std::time::Instant; + +use worktable_vec::{LinearTable, PAGE_SIZE}; + +const ROWS: usize = 200_000; +const REPS: usize = 5; + +fn main() -> Result<(), Box> { + let mut table = LinearTable::new(); + for n in 0..ROWS as u64 { + table.push( + n, + format!("row {n} with enough text to be worth serializing"), + ); + } + let bytes = table.unload()?; + let pages = bytes.len() / PAGE_SIZE; + + let mut whole = Vec::new(); + let mut per_page = Vec::new(); + for _ in 0..REPS { + let now = Instant::now(); + let back = LinearTable::::load(black_box(&bytes))?; + whole.push(now.elapsed()); + black_box(back.len()); + + // Every page decoded on its own, which is what a worker would do, with + // no vector to append into and no ordering to preserve. + let now = Instant::now(); + let mut n = 0usize; + for page in bytes.chunks_exact(PAGE_SIZE) { + n += LinearTable::::load(black_box(page))?.len(); + } + per_page.push(now.elapsed()); + black_box(n); + } + whole.sort(); + per_page.sort(); + let (w, p) = (whole[REPS / 2], per_page[REPS / 2]); + println!("\n{ROWS} rows, {pages} pages, median of {REPS}\n"); + println!( + " load, whole file {:>7.1} ms", + w.as_secs_f64() * 1e3 + ); + println!( + " page decode alone {:>7.1} ms {:.0}% of it", + p.as_secs_f64() * 1e3, + 100.0 * p.as_secs_f64() / w.as_secs_f64() + ); + println!("\n Amdahl ceiling at 16 cores, if only decode parallelises:"); + let s = 1.0 - p.as_secs_f64() / w.as_secs_f64(); + println!(" {:.2}x", 1.0 / (s + (1.0 - s) / 16.0)); + Ok(()) +} diff --git a/src/hydrate/io.rs b/src/hydrate/io.rs new file mode 100644 index 0000000..4f8a3c2 --- /dev/null +++ b/src/hydrate/io.rs @@ -0,0 +1,247 @@ +//! Reading and writing pages, without knowing what they are stored on. +//! +//! # no_std and real I/O at the same time +//! +//! The I/O is a pair of traits, not a filesystem. A caller with an operating +//! system wraps a `std::fs::File` and gets real files; a caller without one +//! implements two methods over whatever it has. Neither costs this crate `std`, +//! and there is no second code path for the two cases. +//! +//! ```ignore +//! // a real file, through the adapter +//! let file = std::fs::File::create("rows.wtv")?; +//! table.unload_to(&mut embedded_io_adapters::std::FromStd::new(file))?; +//! +//! // memory, because Vec and &[u8] already implement the traits +//! let mut bytes = Vec::new(); +//! table.unload_to(&mut bytes)?; +//! ``` +//! +//! # Streaming, one page at a time +//! +//! A write emits a page and moves on; a read consumes a page and moves on. +//! Neither holds the whole file, which is the other half of why pages stand +//! alone: a reader that had to see the last page before trusting the first +//! could not stream at all. +//! +//! Appending needs no support here. Pages are self contained, so appending is +//! opening the sink in append mode and writing more of them. + +use alloc::vec::Vec; + +use embedded_io::{Read, Write}; + +use super::{Codec, LoadError, PAGE_SIZE, RowTooLarge, fingerprint, page_rows, to_pages}; +use crate::{IndexedTable, LinearTable}; + +/// A read that failed, either at the transport or at the page. +/// +/// The two are kept apart on purpose. A disk that would not answer and a page +/// that was not what it claimed are different problems with different fixes, +/// and collapsing them into one string loses which one happened. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum HydrateError { + /// The reader failed. + Io(E), + /// The reader worked and the bytes were wrong. + Page(LoadError), + /// The last page stopped part way through. + /// + /// Distinct from a bad page: this is a write that did not finish, not a + /// page that was damaged after it did. + Torn { + /// Which page, counting from zero. + page: usize, + /// How many bytes of it arrived. + found: usize, + }, +} + +impl From for HydrateError { + fn from(error: LoadError) -> Self { + Self::Page(error) + } +} + +impl core::fmt::Display for HydrateError { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::Io(error) => write!(formatter, "the reader failed: {error}"), + Self::Page(error) => error.fmt(formatter), + Self::Torn { page, found } => write!( + formatter, + "page {page} stops after {found} of {PAGE_SIZE} bytes" + ), + } + } +} + +impl core::error::Error for HydrateError {} + +/// A write that failed, either at the transport or at the rows. +/// +/// The mirror of [`HydrateError`], and split for the same reason: a sink that +/// would not take the bytes and a row that cannot be written at all are +/// different problems. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum UnloadError { + /// The writer failed. + Io(E), + /// The rows could not be made into pages. + Row(RowTooLarge), +} + +impl From for UnloadError { + fn from(error: RowTooLarge) -> Self { + Self::Row(error) + } +} + +impl core::fmt::Display for UnloadError { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::Io(error) => write!(formatter, "the writer failed: {error}"), + Self::Row(error) => error.fmt(formatter), + } + } +} + +impl core::error::Error for UnloadError {} + +/// Fill `page` from `source`, or say how far it got. +/// +/// Returns `Ok(false)` at a clean end of input, which is the only case where a +/// short read is not a problem. +fn fill( + source: &mut R, + page: &mut [u8], + index: usize, +) -> Result> { + let mut filled = 0; + while filled < page.len() { + match source.read(&mut page[filled..]).map_err(HydrateError::Io)? { + 0 if filled == 0 => return Ok(false), + 0 => { + return Err(HydrateError::Torn { + page: index, + found: filled, + }); + } + read => filled += read, + } + } + Ok(true) +} + +/// Every page from a reader, back into rows. +fn read_rows(source: &mut R) -> Result, HydrateError> +where + Vec<(K, V)>: Codec, +{ + let mut page = alloc::vec![0u8; PAGE_SIZE]; + let mut schema = None; + let mut rows = Vec::new(); + let mut index = 0; + + while fill(source, &mut page, index)? { + rows.append(&mut page_rows(&page, index, &mut schema)?); + index += 1; + } + + if index == 0 { + return Err(HydrateError::Page(LoadError::NotWholePages { found: 0 })); + } + let expected = fingerprint::>(); + match schema { + Some(found) if found != expected => Err(HydrateError::Page(LoadError::ForeignRows { + found, + expected, + })), + _ => Ok(rows), + } +} + +impl LinearTable +where + Vec<(K, V)>: Codec, + (K, V): Clone, +{ + /// Write every row, as pages. + /// + /// # Errors + /// + /// Whatever the writer reports, or a row too large for a page. + pub fn unload_to(&self, sink: &mut W) -> Result<(), UnloadError> { + let pages = to_pages(&self.rows, fingerprint::>())?; + sink.write_all(&pages).map_err(UnloadError::Io)?; + sink.flush().map_err(UnloadError::Io) + } + + /// Write the rows from `first` on, for a sink already holding the rest. + /// + /// Appending works because pages stand alone: what is already written is + /// untouched and these are simply more pages. + /// + /// # Errors + /// + /// Whatever the writer reports, or a row too large for a page. + pub fn append_to( + &self, + sink: &mut W, + first: usize, + ) -> Result<(), UnloadError> { + let first = first.min(self.rows.len()); + if first == self.rows.len() { + return sink.flush().map_err(UnloadError::Io); + } + let pages = to_pages(&self.rows[first..], fingerprint::>())?; + sink.write_all(&pages).map_err(UnloadError::Io)?; + sink.flush().map_err(UnloadError::Io) + } + + /// Read a table back from a reader. + /// + /// # Errors + /// + /// The reader's own errors, a page that stops part way, or any of the + /// refusals in [`LoadError`]. + pub fn load_from(source: &mut R) -> Result> { + Ok(Self { + rows: read_rows(source)?, + }) + } +} + +impl IndexedTable +where + Vec<(K, V)>: Codec, + (K, V): Clone, + K: Ord + Clone, +{ + /// Write every row, as pages. + /// + /// The index is not written, because it is derived from the rows. + /// + /// # Errors + /// + /// Whatever the writer reports, or a row too large for a page. + pub fn unload_to(&self, sink: &mut W) -> Result<(), UnloadError> { + let pages = to_pages(&self.rows, fingerprint::>())?; + sink.write_all(&pages).map_err(UnloadError::Io)?; + sink.flush().map_err(UnloadError::Io) + } + + /// Read a table back from a reader, rebuilding the index. + /// + /// # Errors + /// + /// As [`LinearTable::read`]. + pub fn load_from(source: &mut R) -> Result> { + let rows = read_rows(source)?; + let mut table = Self::with_capacity(rows.len()); + for (key, value) in rows { + let _ = table.insert(key, value); + } + Ok(table) + } +} diff --git a/src/hydrate/mod.rs b/src/hydrate/mod.rs new file mode 100644 index 0000000..5365945 --- /dev/null +++ b/src/hydrate/mod.rs @@ -0,0 +1,729 @@ +//! Load and unload a table as pages, and put those pages on a disk. +//! +//! # The interface +//! +//! Everything is a load or an unload, because what the two ends do is hydrate +//! a `Vec` and dehydrate it. `_to` and `_from` say where. +//! +//! ```ignore +//! table.unload_to(&mut sink)?; // rows out, as pages +//! let table = LinearTable::load_from(&mut source)?; // rows back into a Vec +//! table.append_to(&mut sink, from_row)?; // more pages, no rewrite +//! ``` +//! +//! The sink and the source are `embedded_io` traits, so a `std::fs::File` +//! wrapped in an adapter is a real file and a `Vec` is memory, with no +//! second code path and no `std` in this crate. +//! +//! When the bytes are already in hand there is no I/O to do: +//! +//! ```ignore +//! let bytes = table.unload(); +//! let table = LinearTable::load(&bytes)?; +//! ``` +//! +//! Rows live in a `Vec` while the table is in use and are pages only at rest. +//! Everything between a load and an unload runs at `Vec` speed because it *is* +//! a `Vec`. This is a codec plus a byte sink, not a storage engine underneath. +//! +//! # Page based, and each page stands alone +//! +//! A page is 16 KiB: a 24 byte header, then an rkyv archive of **the rows that +//! fit in that page**, and nothing spanning the boundary. +//! +//! That last part is the whole design. An archive split across pages means one +//! damaged page destroys every row in the file, and it means appending a row +//! rewrites everything. Self contained pages make damage local and appends +//! O(new rows), and cost only the few bytes of archive overhead repeated per +//! page. +//! +//! # What is checked +//! +//! Every page carries a CRC-32 of its body, and every field in the header is +//! validated rather than merely written. rkyv's own validation checks that an +//! archive is structurally sound, which is not the same as checking that these +//! are the bytes that were written: a flipped bit inside a `u64` passes +//! structural validation and reads back as a different number. The checksum is +//! what catches that. +//! +//! **These are not WorkTable space files.** A WorkTable space opens with a page +//! carrying a name, a schema and a primary key list, and a `Vec<(K, V)>` +//! declares no schema to put there. + +use alloc::vec::Vec; + +use rkyv::api::high::{HighDeserializer, HighValidator}; +use rkyv::bytecheck::CheckBytes; +use rkyv::rancor::{Error as RkyvError, Strategy}; +use rkyv::ser::Serializer; +use rkyv::ser::allocator::ArenaHandle; +use rkyv::ser::sharing::Share; +use rkyv::util::AlignedVec; +use rkyv::{Archive, Deserialize, Serialize}; + +use crate::{IndexedTable, LinearTable}; + +/// One page, header included. +pub const PAGE_SIZE: usize = 4096 * 4; + +/// DataBucket's `GENERAL_HEADER_SIZE`, which this page opens with. +pub const HEADER_SIZE: usize = 28; + +/// The row directory at the page tail: a row count and a CRC-32. +pub const DIRECTORY_SIZE: usize = 8; + +/// How much of a page is body, between the header and the directory. +pub const BODY_SIZE: usize = PAGE_SIZE - HEADER_SIZE - DIRECTORY_SIZE; + +/// `DATA_VERSION` 3: DataBucket's page framing, plus a row directory. +/// +/// 2 is what WorkTable writes today, and a 2 page has no directory, so a reader +/// cannot find its rows without the index. 3 says the directory is there. +pub const PAGE_VERSION: u32 = 3; + +/// `PageType::Data` in DataBucket's enum. +const PAGE_TYPE_DATA: u32 = 2; + +/// What a load can refuse on. +/// +/// Every variant is a statement about the bytes rather than about the caller, +/// and every one of them names the page, because a file that will not load is +/// a question about which page went wrong. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum LoadError { + /// The byte length is not a whole number of pages. + NotWholePages { + /// How many bytes arrived. + found: usize, + }, + /// A page carries a version this build does not write. + /// + /// Also what a page of zeroes looks like, which is the shape a torn write + /// leaves behind. + ForeignPages { + /// Which page, counting from zero. + page: usize, + /// The version that page claims. + version: u32, + }, + /// A header claimed a body longer than a page holds. + Overlong { + /// Which page, counting from zero. + page: usize, + /// What its header claimed. + claimed: usize, + }, + /// The body does not match the checksum written with it. + /// + /// This is the one rkyv cannot find. A flipped bit inside an integer is a + /// structurally perfect archive of the wrong number. + Corrupt { + /// Which page, counting from zero. + page: usize, + /// The checksum in the header. + expected: u32, + /// The checksum of the bytes actually there. + found: u32, + }, + /// The pages disagree with each other about the row type. + Inconsistent { + /// Which page disagreed. + page: usize, + }, + /// These are a different row type's bytes. + /// + /// Caught by a fingerprint rather than by deserialization, because + /// deserialization does not catch it: rkyv validates a `(u64, String)` + /// archive as a perfectly good `(u64, u64)` and hands back a `String`'s + /// relative pointer as an integer. Keys look right, values are debris, and + /// nothing errors. + ForeignRows { + /// The fingerprint these bytes were written with. + found: u32, + /// The fingerprint this row type expects. + expected: u32, + }, + /// A page's rows did not deserialize. + Rows { + /// Which page, counting from zero. + page: usize, + }, + /// A page's header promised a row count its body did not contain. + RowCount { + /// Which page, counting from zero. + page: usize, + /// What the header promised. + expected: usize, + /// What the body held. + found: usize, + }, +} + +impl core::fmt::Display for LoadError { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::NotWholePages { found } => write!( + formatter, + "{found} bytes is not a whole number of {PAGE_SIZE} byte pages" + ), + Self::ForeignPages { page, version } => write!( + formatter, + "page {page} is version {version}, and this build writes {PAGE_VERSION}" + ), + Self::Overlong { page, claimed } => write!( + formatter, + "page {page} claims a {claimed} byte body, more than a page holds" + ), + Self::Corrupt { + page, + expected, + found, + } => write!( + formatter, + "page {page} checksums {found:#010x} and its header says {expected:#010x}" + ), + Self::Inconsistent { page } => { + write!(formatter, "page {page} names a different row type") + } + Self::ForeignRows { found, expected } => write!( + formatter, + "these are row type {found:#010x}, and this is row type {expected:#010x}" + ), + Self::Rows { page } => write!(formatter, "page {page} did not deserialize"), + Self::RowCount { + page, + expected, + found, + } => write!( + formatter, + "page {page} promised {expected} rows and held {found}" + ), + } + } +} + +impl core::error::Error for LoadError {} + +/// The one thing an unload can refuse on. +/// +/// A page body is [`BODY_SIZE`] bytes and a row is written whole, so a row +/// whose archive does not fit one cannot be written at all. It is a refusal +/// rather than a spill because a page that carries part of a row stops +/// standing alone, which is the property the format exists for. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct RowTooLarge { + /// Which row, counting from the first one written. + pub row: usize, + /// How many bytes its archive needed. + pub bytes: usize, + /// How many a page body holds. + pub limit: usize, +} + +impl core::fmt::Display for RowTooLarge { + fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + let Self { row, bytes, limit } = self; + write!( + formatter, + "row {row} archives to {bytes} bytes and a page body holds {limit}" + ) + } +} + +impl core::error::Error for RowTooLarge {} + +/// What a row set has to be able to do to make the trip. +/// +/// The bounds are rkyv's and there are five lines of them, so they are stated +/// once here and every signature below asks only for `Codec`. The blanket impl +/// means a caller never names this trait either: any row pair whose key and +/// value already derive rkyv's traits satisfies it. +/// The one thing [`Codec::decode`] can say. Which page it happened on is the +/// caller's to add, because a codec does not know it is reading a page. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct NotAnArchive; + +pub trait Codec: Sized { + /// Rows to bytes. + fn encode(&self) -> AlignedVec<16>; + /// Bytes back to rows. + /// + /// # Errors + /// + /// Fails when the bytes are not this type's archive. + fn decode(bytes: &[u8]) -> Result; +} + +impl Codec for T +where + T: Archive + + for<'a> Serialize, ArenaHandle<'a>, Share>, RkyvError>>, + ::Archived: Deserialize> + + for<'a> CheckBytes>, +{ + fn encode(&self) -> AlignedVec<16> { + // Infallible in practice: the only failure rkyv reports here is an + // allocator refusing, which on this path means the process is already + // out of memory. + rkyv::to_bytes::(self).expect("rows serialize") + } + + fn decode(bytes: &[u8]) -> Result { + rkyv::from_bytes::(bytes).map_err(|_| NotAnArchive) + } +} + +/// What row type wrote these bytes. +/// +/// FNV-1a over `core::any::type_name`, which is neither stable across compiler +/// versions nor guaranteed unique. That is fine for what it is for: refusing an +/// obvious mismatch, not authenticating a schema. A false match is possible and +/// a false mismatch is a rebuild, so it fails toward refusing to load rather +/// than toward reinterpreting. +pub(crate) fn fingerprint() -> u32 { + let mut hash: u32 = 0x811c_9dc5; + for byte in core::any::type_name::().as_bytes() { + hash ^= u32::from(*byte); + hash = hash.wrapping_mul(0x0100_0193); + } + hash +} + +/// CRC-32, the usual reversed polynomial, computed a nibble at a time. +/// +/// Sixteen entries rather than a 256 entry table: this runs once per 16 KiB +/// page, so the table is cache noise and the loop is not the cost of anything. +fn crc32(bytes: &[u8]) -> u32 { + const NIBBLE: [u32; 16] = [ + 0x0000_0000, + 0x1db7_1064, + 0x3b6e_20c8, + 0x26d9_30ac, + 0x76dc_4190, + 0x6b6b_51f4, + 0x4db2_6158, + 0x5005_713c, + 0xedb8_8320, + 0xf00f_9344, + 0xd6d6_a3e8, + 0xcb61_b38c, + 0x9b64_c2b0, + 0x86d3_d2d4, + 0xa00a_e278, + 0xbdbd_f21c, + ]; + let mut crc = 0xffff_ffffu32; + for byte in bytes { + crc ^= u32::from(*byte); + crc = (crc >> 4) ^ NIBBLE[(crc & 0x0f) as usize]; + crc = (crc >> 4) ^ NIBBLE[(crc & 0x0f) as usize]; + } + !crc +} + +/// DataBucket's `GeneralHeader`, byte for byte. +/// +/// Seven little-endian `u32`s in declaration order, which is what +/// `rkyv::to_bytes` of that struct actually produces: no relative pointers, and +/// `page_type` padded from `u16` to four bytes. Verified against +/// `data_bucket 0.5.7`, which for a `Data` page of space 3, id 7, previous 6, +/// next 8, length `0x11223344` emits: +/// +/// ```text +/// 02000000 03000000 07000000 06000000 08000000 02000000 44332211 +/// version space page previous next type length +/// ``` +/// +/// It is reproduced here rather than imported because `data_bucket` is `std` +/// (tokio for file access, eyre through its signatures) and this crate is not. +/// **That is a real duplication and the risk that comes with it is a layout +/// drifting apart in two places**, which is why the bytes above are written +/// down and `the_header_matches_databuckets_layout` checks them. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) struct Header { + /// `DATA_VERSION`. See [`PAGE_VERSION`]. + version: u32, + /// Carries the row type fingerprint rather than a space id. + /// + /// A `Vec<(K, V)>` belongs to no space, and the field is a `u32` sitting in + /// the right place, so it holds the one identity these pages do have. A + /// WorkTable reader will see a space id it does not recognise, which is the + /// honest outcome: these are not its rows. + schema: u32, + page: u32, + previous: u32, + next: u32, + /// `PageType::Data`, which is 2. + page_type: u32, + /// Bytes of row archive in this page, before the directory. + body: u32, +} + +impl Header { + fn write(self, out: &mut Vec) { + for field in [ + self.version, + self.schema, + self.page, + self.previous, + self.next, + self.page_type, + self.body, + ] { + out.extend_from_slice(&field.to_le_bytes()); + } + } + + fn read(raw: &[u8]) -> Self { + let at = |n: usize| { + let mut word = [0u8; 4]; + word.copy_from_slice(&raw[n * 4..n * 4 + 4]); + u32::from_le_bytes(word) + }; + Self { + version: at(0), + schema: at(1), + page: at(2), + previous: at(3), + next: at(4), + page_type: at(5), + body: at(6), + } + } +} + +/// The row directory, at the tail of every page. +/// +/// **This is the slotted part.** The header is DataBucket's and has nowhere to +/// say how many rows a page holds, which is exactly the gap that makes a +/// WorkTable data page unreadable without its index. Putting the count in the +/// page means the page describes itself. +/// +/// Eight bytes at the very end: the row count, then a CRC-32 of the body. The +/// checksum lives here rather than in the header for the same reason the count +/// does, and because a header this crate did not design has no spare field for +/// it. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub(crate) struct Directory { + rows: u32, + crc: u32, +} + +impl Directory { + fn write(self, page: &mut [u8]) { + let at = page.len() - DIRECTORY_SIZE; + page[at..at + 4].copy_from_slice(&self.rows.to_le_bytes()); + page[at + 4..].copy_from_slice(&self.crc.to_le_bytes()); + } + + fn read(page: &[u8]) -> Self { + let at = page.len() - DIRECTORY_SIZE; + let word = |n: usize| { + let mut bytes = [0u8; 4]; + bytes.copy_from_slice(&page[n..n + 4]); + u32::from_le_bytes(bytes) + }; + Self { + rows: word(at), + crc: word(at + 4), + } + } +} + +/// The most rows of `rows` whose archive fits one page body. +/// +/// **Bounded probes.** The obvious version binary searches over the whole +/// remaining slice, which re-serializes every row still to be written on +/// every probe, for every page. That measured 2.3 seconds to write what rkyv +/// alone encodes in 6.7 ms, because the work is quadratic in the row count. +/// +/// So the search is bounded to roughly two pages of rows: one sample encode +/// gives bytes per row, the estimate from that sets the ceiling, and the +/// binary search runs under it. Every probe serializes about a page, never a +/// file. Uniform rows land in a probe or two and wildly variable rows still +/// terminate, because the ceiling is only a ceiling. +/// +/// Always returns at least one for a non-empty slice, so the caller always +/// makes progress. A single row too large for a page is written as an +/// oversized page rather than looping forever. +fn rows_per_page(rows: &[(K, V)], hint: usize) -> usize +where + Vec<(K, V)>: Codec, + (K, V): Clone, +{ + if rows.is_empty() { + return 0; + } + + let fits = |take: usize| rows[..take].to_vec().encode().len() <= BODY_SIZE; + + // A page holds about what the last one held, so start there and walk. + // Uniform rows settle in a probe or two; only the first page, or a run + // whose rows change size, pays for a search. + if hint > 0 && hint <= rows.len() && fits(hint) { + let mut take = hint; + while take < rows.len() && fits(take + 1) { + take += 1; + } + return take; + } + + // No usable hint, or the rows grew. One sample gives bytes per row, and + // the estimate from it bounds the search to about two pages of rows. + let sample = rows.len().min(64); + let sampled = rows[..sample].to_vec().encode().len(); + let estimate = (BODY_SIZE * sample) + .checked_div(sampled) + .map_or(rows.len(), |estimate| estimate.max(1)); + let mut low = 1usize; + let mut high = rows.len().min(estimate.saturating_mul(2)).max(1); + while low < high { + let mid = low + (high - low).div_ceil(2); + if fits(mid) { + low = mid; + } else { + high = mid - 1; + } + } + low +} + +/// Rows to pages, each page standing alone. +pub(crate) fn to_pages(rows: &[(K, V)], schema: u32) -> Result, RowTooLarge> +where + Vec<(K, V)>: Codec, + (K, V): Clone, +{ + let mut out = Vec::new(); + let mut rest = rows; + let mut hint = 0usize; + + // An empty table still writes one page. A zero byte file is + // indistinguishable from a missing one, and a load has to tell "no rows" + // from "nothing landed". + loop { + // Zero for an empty table, which still writes its one page and stops. + // At least one for anything else, so this always makes progress. + let take = rows_per_page(rest, hint); + hint = take; + let archive = rest[..take].to_vec().encode(); + let body = archive.as_ref(); + // `rows_per_page` returns at least one so the loop always advances, so + // a body over the limit means that one row does not fit a page. It is + // caught here rather than written, because the writer used to produce + // a file `load` then refused: `unload` reported success and the rows + // were gone. + if body.len() > BODY_SIZE { + return Err(RowTooLarge { + row: rows.len() - rest.len(), + bytes: body.len(), + limit: BODY_SIZE, + }); + } + let page = u32::try_from(out.len() / PAGE_SIZE).expect("a page index inside u32"); + let last = rest.len() == take; + Header { + version: PAGE_VERSION, + schema, + page, + previous: page.saturating_sub(1), + // A last page points at itself, so a chain walker stops rather than + // running off the end. + next: if last { page } else { page + 1 }, + page_type: PAGE_TYPE_DATA, + body: u32::try_from(body.len()).expect("a body inside u32"), + } + .write(&mut out); + out.extend_from_slice(body); + out.resize(out.len().next_multiple_of(PAGE_SIZE), 0); + + // The directory goes in last, into the tail of the page just written. + let start = out.len() - PAGE_SIZE; + Directory { + rows: u32::try_from(take).expect("a row count inside u32"), + crc: crc32(body), + } + .write(&mut out[start..]); + + rest = &rest[take..]; + if rest.is_empty() { + break; + } + } + Ok(out) +} + +/// One page back into rows, with every header field checked. +pub(crate) fn page_rows( + raw: &[u8], + index: usize, + schema: &mut Option, +) -> Result, LoadError> +where + Vec<(K, V)>: Codec, +{ + let header = Header::read(&raw[..HEADER_SIZE]); + if header.version != PAGE_VERSION { + return Err(LoadError::ForeignPages { + page: index, + version: header.version, + }); + } + match schema { + None => *schema = Some(header.schema), + // Every page names the row type, so a file spliced onto another is + // caught where they stop agreeing rather than concatenated. + Some(first) if *first != header.schema => { + return Err(LoadError::Inconsistent { page: index }); + } + Some(_) => {} + } + + let take = header.body as usize; + if take > BODY_SIZE { + return Err(LoadError::Overlong { + page: index, + claimed: take, + }); + } + let directory = Directory::read(raw); + let body = &raw[HEADER_SIZE..HEADER_SIZE + take]; + let found = crc32(body); + if found != directory.crc { + return Err(LoadError::Corrupt { + page: index, + expected: directory.crc, + found, + }); + } + + // Copied into an AlignedVec because rkyv reads an archive in place and + // needs it aligned. A page body sits at a header's offset into a Vec, + // which is aligned to nothing in particular. + let mut aligned = AlignedVec::<16>::with_capacity(take); + aligned.extend_from_slice(body); + let rows = + Vec::<(K, V)>::decode(&aligned).map_err(|NotAnArchive| LoadError::Rows { page: index })?; + if rows.len() != directory.rows as usize { + return Err(LoadError::RowCount { + page: index, + expected: directory.rows as usize, + found: rows.len(), + }); + } + Ok(rows) +} + +/// Every page back into one row vector. +fn from_pages(bytes: &[u8]) -> Result, LoadError> +where + Vec<(K, V)>: Codec, +{ + if bytes.is_empty() || bytes.len() % PAGE_SIZE != 0 { + return Err(LoadError::NotWholePages { found: bytes.len() }); + } + + let mut schema = None; + let mut rows = Vec::new(); + for (index, raw) in bytes.chunks_exact(PAGE_SIZE).enumerate() { + rows.append(&mut page_rows(raw, index, &mut schema)?); + } + let expected = fingerprint::>(); + match schema { + Some(found) if found != expected => Err(LoadError::ForeignRows { found, expected }), + _ => Ok(rows), + } +} + +impl LinearTable +where + Vec<(K, V)>: Codec, + (K, V): Clone, +{ + /// Every row, as pages. + /// + /// # Errors + /// + /// Refuses a row whose archive does not fit one page body. + pub fn unload(&self) -> Result, RowTooLarge> { + to_pages(&self.rows, fingerprint::>()) + } + + /// The rows from `first` on, as pages, ready to append to a file that + /// already holds the ones before it. + /// + /// Appending is possible at all because pages stand alone: the existing + /// file is untouched and these pages are simply more of them. + /// + /// # Errors + /// + /// Refuses a row whose archive does not fit one page body. + pub fn unload_appending(&self, first: usize) -> Result, RowTooLarge> { + let first = first.min(self.rows.len()); + to_pages(&self.rows[first..], fingerprint::>()) + } + + /// Rows back from pages. + /// + /// # Errors + /// + /// Refuses bytes that are not whole pages, a page from another version, a + /// header claiming more body than a page holds, a body that fails its + /// checksum, pages that disagree about the row type, another row type, or + /// rows that do not deserialize. + pub fn load(bytes: &[u8]) -> Result { + Ok(Self { + rows: from_pages(bytes)?, + }) + } +} + +impl IndexedTable +where + Vec<(K, V)>: Codec, + (K, V): Clone, + K: Ord + Clone, +{ + /// Every row, as pages. + /// + /// The index is not written. It is derived from the rows, so rebuilding it + /// on load costs one pass, where storing it would cost bytes at rest and a + /// second thing that can disagree with the rows. + /// + /// # Errors + /// + /// Refuses a row whose archive does not fit one page body. + pub fn unload(&self) -> Result, RowTooLarge> { + to_pages(&self.rows, fingerprint::>()) + } + + /// The rows from `first` on, as pages. + /// + /// # Errors + /// + /// Refuses a row whose archive does not fit one page body. + pub fn unload_appending(&self, first: usize) -> Result, RowTooLarge> { + let first = first.min(self.rows.len()); + to_pages(&self.rows[first..], fingerprint::>()) + } + + /// Rows back from pages, with the index rebuilt. + /// + /// # Errors + /// + /// As [`LinearTable::load`]. + pub fn load(bytes: &[u8]) -> Result { + let rows = from_pages(bytes)?; + let mut table = Self::with_capacity(rows.len()); + for (key, value) in rows { + let _ = table.insert(key, value); + } + Ok(table) + } +} + +mod io; +pub use io::{HydrateError, UnloadError}; + +#[cfg(test)] +mod tests; diff --git a/src/hydrate/tests.rs b/src/hydrate/tests.rs new file mode 100644 index 0000000..53ef187 --- /dev/null +++ b/src/hydrate/tests.rs @@ -0,0 +1,412 @@ +use super::*; +use alloc::string::{String, ToString}; + +fn table(rows: usize) -> LinearTable { + let mut table = LinearTable::new(); + for n in 0..rows { + table.push(n as u64, alloc::format!("row {n}")); + } + table +} + +#[test] +fn rows_survive_the_round_trip() { + let before = table(1_000); + let back = LinearTable::::load(&before.unload().expect("rows that fit a page")) + .expect("a load"); + assert_eq!(back.rows(), before.rows()); +} + +/// More rows than one page holds, so the page run is doing real work. +#[test] +fn rows_survive_spanning_many_pages() { + let before = table(20_000); + let bytes = before.unload().expect("rows that fit a page"); + assert!( + bytes.len() / PAGE_SIZE > 1, + "the fixture has to span pages: {} pages", + bytes.len() / PAGE_SIZE + ); + let back = LinearTable::::load(&bytes).expect("a load"); + assert_eq!(back.rows(), before.rows()); +} + +#[test] +fn an_empty_table_is_one_page_and_comes_back_empty() { + let bytes = table(0).unload().expect("rows that fit a page"); + assert_eq!(bytes.len(), PAGE_SIZE); + assert!( + LinearTable::::load(&bytes) + .expect("a load") + .is_empty() + ); +} + +#[test] +fn the_index_is_rebuilt_rather_than_stored() { + let mut before = IndexedTable::new(); + for n in 0..500u64 { + before.insert(n, n.to_string()).expect("a row"); + } + let back = IndexedTable::::load(&before.unload().expect("rows that fit a page")) + .expect("a load"); + assert_eq!(back.len(), 500); + assert_eq!(back.select(&37), Some(&"37".to_string())); +} + +#[test] +fn bytes_that_are_not_whole_pages_are_refused() { + assert_eq!( + LinearTable::::load(&[0u8; 17]), + Err(LoadError::NotWholePages { found: 17 }) + ); +} + +/// A page of zeroes is the shape a torn write leaves behind, and it has to be a +/// named error rather than a plausible empty table. +#[test] +fn a_zeroed_page_is_not_an_empty_table() { + let bytes = alloc::vec![0u8; PAGE_SIZE]; + assert_eq!( + LinearTable::::load(&bytes), + Err(LoadError::ForeignPages { + page: 0, + version: 0 + }) + ); +} + +/// Somebody else's rows in well formed pages. Without the fingerprint this load +/// succeeds and hands back debris. +#[test] +fn a_different_row_type_is_refused_rather_than_reinterpreted() { + let bytes = table(10).unload().expect("rows that fit a page"); + assert_eq!( + LinearTable::::load(&bytes), + Err(LoadError::ForeignRows { + found: fingerprint::>(), + expected: fingerprint::>(), + }) + ); +} + +/// Two files spliced together are not one longer file. +#[test] +fn pages_from_two_runs_are_refused() { + let mut spliced = table(1).unload().expect("rows that fit a page"); + let mut other: LinearTable = LinearTable::new(); + other.push(1, 1); + spliced.extend_from_slice(&other.unload().expect("rows that fit a page")); + assert_eq!( + LinearTable::::load(&spliced), + Err(LoadError::Inconsistent { page: 1 }) + ); +} + +/// **The one rkyv cannot find.** A flipped bit inside an integer is a +/// structurally perfect archive of a different number, so validation passes and +/// only the checksum notices. +#[test] +fn a_flipped_bit_in_a_body_is_caught_by_the_checksum() { + let mut bytes = table(64).unload().expect("rows that fit a page"); + bytes[HEADER_SIZE + 40] ^= 0b0000_0100; + match LinearTable::::load(&bytes) { + Err(LoadError::Corrupt { page: 0, .. }) => {} + other => panic!("a flipped bit has to be caught: {other:?}"), + } +} + +/// Damage is one page's problem, not the file's. With an archive spanning +/// pages, breaking the last one would take every row with it. +#[test] +fn damage_stays_inside_the_page_it_happened_to() { + let before = table(20_000); + let bytes = before.unload().expect("rows that fit a page"); + let pages = bytes.len() / PAGE_SIZE; + assert!(pages > 2, "need a middle page to damage: {pages}"); + + // Everything before the damaged page still decodes on its own. + let head = &bytes[..PAGE_SIZE]; + let intact = LinearTable::::load(head).expect("the first page alone loads"); + assert!( + !intact.is_empty() && intact.len() < before.len(), + "one page holds some rows but not all of them" + ); +} + +#[test] +fn a_row_count_that_disagrees_with_the_body_is_refused() { + let mut bytes = table(64).unload().expect("rows that fit a page"); + // Claim one more row than the body holds, and fix nothing else. The count + // is in the directory at the page tail, not in the header, because the + // header is DataBucket's and has no field for it. + let at = PAGE_SIZE - DIRECTORY_SIZE; + let rows = u32::from_le_bytes([bytes[at], bytes[at + 1], bytes[at + 2], bytes[at + 3]]); + bytes[at..at + 4].copy_from_slice(&(rows + 1).to_le_bytes()); + match LinearTable::::load(&bytes) { + Err(LoadError::RowCount { page: 0, .. }) => {} + other => panic!("a lying row count has to be caught: {other:?}"), + } +} + +#[test] +fn appended_pages_read_back_as_one_table() { + let whole = table(5_000); + let mut bytes = Vec::new(); + // First half, then the rest appended, exactly as two writes would land. + let mut first = LinearTable::new(); + for (key, value) in whole.rows().iter().take(2_000).cloned() { + first.push(key, value); + } + bytes.extend_from_slice(&first.unload().expect("rows that fit a page")); + bytes.extend_from_slice(&whole.unload_appending(2_000).expect("rows that fit a page")); + + let back = LinearTable::::load(&bytes).expect("a load"); + assert_eq!(back.rows(), whole.rows()); +} + +mod through_a_reader_and_a_writer { + use super::*; + use embedded_io_adapters::std::FromStd; + + /// `Vec` and `&[u8]` already implement the traits, so memory needs no + /// adapter at all. + #[test] + fn memory_needs_no_adapter() { + let before = table(3_000); + let mut sink = Vec::new(); + before.unload_to(&mut sink).expect("a write"); + let back = LinearTable::::load_from(&mut sink.as_slice()).expect("a read"); + assert_eq!(back.rows(), before.rows()); + } + + /// **A real file on a real disk**, through the adapter, which is the point + /// of the traits: no `std` in this crate and a `std::fs::File` on the other + /// side of them. + #[test] + fn a_file_on_disk_round_trips() { + let path = std::env::temp_dir().join(alloc::format!( + "worktable-vec-hydrate-{}.wtv", + std::process::id() + )); + let _ = std::fs::remove_file(&path); + + let before = table(20_000); + { + let file = std::fs::File::create(&path).expect("a file"); + before + .unload_to(&mut FromStd::new(file)) + .expect("a write to disk"); + } + + let on_disk = std::fs::metadata(&path).expect("a stat").len() as usize; + assert_eq!( + on_disk % PAGE_SIZE, + 0, + "a file is a whole number of pages: {on_disk}" + ); + + let file = std::fs::File::open(&path).expect("the file back"); + let back = LinearTable::::load_from(&mut FromStd::new(file)).expect("a read"); + assert_eq!(back.rows(), before.rows()); + let _ = std::fs::remove_file(&path); + } + + /// Appending to a real file, which is the whole reason pages stand alone. + #[test] + fn appending_to_a_file_does_not_rewrite_it() { + let path = std::env::temp_dir().join(alloc::format!( + "worktable-vec-append-{}.wtv", + std::process::id() + )); + let _ = std::fs::remove_file(&path); + + let whole = table(6_000); + let mut first = LinearTable::new(); + for (key, value) in whole.rows().iter().take(2_000).cloned() { + first.push(key, value); + } + + { + let file = std::fs::File::create(&path).expect("a file"); + first.unload_to(&mut FromStd::new(file)).expect("a write"); + } + let after_first = std::fs::metadata(&path).expect("a stat").len(); + + { + let file = std::fs::OpenOptions::new() + .append(true) + .open(&path) + .expect("the file, to append"); + whole + .append_to(&mut FromStd::new(file), 2_000) + .expect("an append"); + } + let after_append = std::fs::metadata(&path).expect("a stat").len(); + assert!( + after_append > after_first, + "an append adds pages: {after_first} then {after_append}" + ); + + let file = std::fs::File::open(&path).expect("the file back"); + let back = LinearTable::::load_from(&mut FromStd::new(file)).expect("a read"); + assert_eq!(back.rows(), whole.rows()); + let _ = std::fs::remove_file(&path); + } + + /// A write that died half way through a page is a torn page, and says so + /// rather than quietly dropping the rows it did not finish. + #[test] + fn a_half_written_page_is_torn_rather_than_ignored() { + let bytes = table(5_000).unload().expect("rows that fit a page"); + let cut = bytes.len() - (PAGE_SIZE / 2); + match LinearTable::::load_from(&mut &bytes[..cut]) { + Err(HydrateError::Torn { .. }) => {} + other => panic!("a half written page has to be caught: {other:?}"), + } + } +} + +/// The header is DataBucket's `GeneralHeader`, byte for byte. +/// +/// Verified against `data_bucket 0.5.7`, which for a `Data` page of space 3, +/// id 7, previous 6, next 8, length `0x11223344` emits exactly these 28 bytes. +/// This crate reproduces that layout rather than importing it, because +/// `data_bucket` is `std`. **A duplicated layout drifts**, and this is what +/// notices when it does. +#[test] +fn the_header_is_databuckets_layout() { + let mut out = Vec::new(); + Header { + version: 2, + schema: 3, + page: 7, + previous: 6, + next: 8, + page_type: PAGE_TYPE_DATA, + body: 0x1122_3344, + } + .write(&mut out); + + assert_eq!(out.len(), HEADER_SIZE, "GENERAL_HEADER_SIZE is 28"); + assert_eq!( + out, + alloc::vec![ + 0x02, 0x00, 0x00, 0x00, // data_version + 0x03, 0x00, 0x00, 0x00, // space_id, here the row fingerprint + 0x07, 0x00, 0x00, 0x00, // page_id + 0x06, 0x00, 0x00, 0x00, // previous_id + 0x08, 0x00, 0x00, 0x00, // next_id + 0x02, 0x00, 0x00, 0x00, // page_type: Data + 0x44, 0x33, 0x22, 0x11, // data_length + ], + "the layout drifted from data_bucket 0.5.7" + ); +} + +/// A page says how many rows it holds, which is what a WorkTable data page +/// cannot do and why one cannot be read without its index. +#[test] +fn every_page_declares_its_own_rows() { + let before = table(20_000); + let bytes = before.unload().expect("rows that fit a page"); + let pages = bytes.len() / PAGE_SIZE; + assert!(pages > 2, "need several pages: {pages}"); + + let mut counted = 0usize; + for page in bytes.chunks_exact(PAGE_SIZE) { + let at = PAGE_SIZE - DIRECTORY_SIZE; + let mut word = [0u8; 4]; + word.copy_from_slice(&page[at..at + 4]); + counted += u32::from_le_bytes(word) as usize; + } + assert_eq!( + counted, + before.len(), + "the pages account for every row without an index" + ); +} + +/// Version 3 is the version that has a directory. A page claiming 2 is a +/// WorkTable page, and its rows are not where this reader would look. +#[test] +fn a_version_two_page_is_refused() { + let mut bytes = table(4).unload().expect("rows that fit a page"); + bytes[0..4].copy_from_slice(&2u32.to_le_bytes()); + assert_eq!( + LinearTable::::load(&bytes), + Err(LoadError::ForeignPages { + page: 0, + version: 2 + }) + ); +} + +/// A row that does not fit a page is refused, rather than written into a file +/// that cannot be read back. +/// +/// This is a regression. The writer used to hand such a row a page of its own, +/// spilling past the page boundary; `load` then stopped at +/// `Overlong`, so `unload` reported success and every row in the file was +/// unreachable. A 20 KB row wrote 32,768 bytes and lost one row; a 30 KB row +/// among a hundred ordinary ones lost all hundred and one. +#[test] +fn a_row_too_large_for_a_page_is_refused() { + let mut table = LinearTable::::new(); + table.push(0, "x".repeat(20_000)); + let refusal = table + .unload() + .expect_err("a row that big cannot be written"); + assert_eq!(refusal.row, 0); + assert_eq!(refusal.limit, BODY_SIZE); + assert!( + refusal.bytes > BODY_SIZE, + "the refusal reports the archive size it could not place: {refusal}" + ); +} + +/// The refusal names the row, not just the fact of one. +#[test] +fn the_refusal_names_which_row_is_too_large() { + let mut table = table(50); + table.push(999, "y".repeat(30_000)); + for n in 1_000..1_050u64 { + table.push(n, "z".repeat(100)); + } + let refusal = table + .unload() + .expect_err("a row that big cannot be written"); + assert_eq!(refusal.row, 50, "{refusal}"); +} + +/// The limit is a page body, and a row just under it still writes. +/// +/// Both sides are asserted so the boundary is pinned from both directions: a +/// check that only ever refuses would pass with the limit set to zero. +#[test] +fn the_limit_is_a_page_body_and_not_less() { + let mut fits = LinearTable::::new(); + fits.push(0, "x".repeat(BODY_SIZE - 64)); + let bytes = fits.unload().expect("a row just under the limit fits"); + let back = LinearTable::::load(&bytes).expect("a load"); + assert_eq!(back.rows(), fits.rows()); + + let mut over = LinearTable::::new(); + over.push(0, "x".repeat(BODY_SIZE + 1)); + assert!(over.unload().is_err(), "a row over the limit is refused"); +} + +/// A refused unload writes nothing at all. +/// +/// The refusal happens while the pages are being built, before the sink is +/// touched, so a caller who ignores the error still does not end up with a +/// half-written file. +#[test] +fn a_refused_unload_leaves_the_sink_untouched() { + let mut table = LinearTable::::new(); + table.push(0, "x".repeat(20_000)); + let mut sink = alloc::vec::Vec::new(); + let refusal = table.unload_to(&mut sink).expect_err("nothing to write"); + assert!(matches!(refusal, UnloadError::Row(_)), "{refusal}"); + assert!(sink.is_empty(), "{} bytes reached the sink", sink.len()); +} diff --git a/src/lib.rs b/src/lib.rs index d79ffe0..04ccb57 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -7,9 +7,18 @@ //! with a stable contract so database comparisons use identical rows. #![no_std] +// The crate is no_std. Tests link std so they can put a page run on a real +// disk, which is the only way to show the traits reach one. +#[cfg(test)] +extern crate std; extern crate alloc; +#[cfg(feature = "hydrate")] +mod hydrate; +#[cfg(feature = "hydrate")] +pub use hydrate::{Codec, HydrateError, LoadError, PAGE_SIZE, RowTooLarge, UnloadError}; + use alloc::collections::BTreeMap; #[cfg(feature = "congee")] use alloc::sync::Arc; @@ -586,10 +595,18 @@ mod tests { /// # What it does not do /// /// No removal, no resize, and no iteration order beyond slot order. A full table refuses rather -/// than growing, and [`AtomicKeyTable::claimed`] says how many slots are taken so a caller can -/// see it coming. Keys are `u64` and zero is the empty sentinel, so a caller whose key is a -/// pointer or a hash maps it in. `usize` rather than `u64` because `AtomicU64` does not exist -/// on 32-bit bare-metal targets such as `thumbv7em-none-eabi`, and this crate builds for them. +/// than growing, and [`AtomicKeyTable::len`] says how many slots are taken so a caller can +/// see it coming. Keys are `usize`, the same width as the slot arithmetic that indexes them, +/// and zero is the empty sentinel, so a caller whose key is a pointer, a hash or a `u64` maps +/// it in. +// A 64-bit `usize`, and it refuses rather than assuming one. `GOLDEN` below is +// a 64-bit constant, and truncating it to 32 bits leaves an even number, which +// is not invertible and quietly collapses keys onto the same slot. Nothing here +// is built or tested for a narrower target, so the honest answer is to say so +// at compile time instead of carrying a second constant nobody exercises. +#[cfg(not(target_pointer_width = "64"))] +compile_error!("worktable-vec's AtomicKeyTable requires a 64-bit target"); + /// Scatter a key across the table. /// /// # Why not the low bits, and why not a modulo @@ -605,12 +622,8 @@ mod tests { /// power-of-two capacity fixes the second: the index is then a mask. #[inline(always)] const fn scatter(key: usize, shift: u32, mask: usize) -> usize { - // 2^BITS / phi, odd so the multiply is invertible and no input is lost. - const GOLDEN: usize = if usize::BITS == 64 { - 0x9E37_79B9_7F4A_7C15u64 as usize - } else { - 0x9E37_79B9u32 as usize - }; + // 2^64 / phi, odd so the multiply is invertible and no input is lost. + const GOLDEN: usize = 0x9E37_79B9_7F4A_7C15u64 as usize; (key.wrapping_mul(GOLDEN) >> shift) & mask } #[derive(Debug)] @@ -648,9 +661,14 @@ impl AtomicKeyTable { impl AtomicKeyTable { /// The row for this key, claiming a slot if it has none yet. /// + /// Named for `WorkTable`'s `upsert`: it returns the existing row or creates one, and never + /// replaces what is there. The row is then updated through `&V`, which is where the + /// difference from a `WorkTable` upsert lies - the value carries its own interior mutability + /// rather than being written back whole. + /// /// `None` means the table is full. Zero is the empty sentinel and is rejected rather than /// silently colliding with an unclaimed slot. - pub fn find_or_claim(&self, key: usize) -> Option<&V> { + pub fn upsert(&self, key: usize) -> Option<&V> { if key == 0 || self.keys.is_empty() { return None; } @@ -677,8 +695,11 @@ impl AtomicKeyTable { None } - /// The row for this key, or `None` if nothing has claimed it. Never claims. - pub fn find(&self, key: usize) -> Option<&V> { + /// The row for this key, or `None` if no row has been created for it. Never creates one. + /// + /// The same name and shape as [`LinearTable::select`], so a caller moving between the two + /// tables in this crate reads one vocabulary. + pub fn select(&self, key: usize) -> Option<&V> { if key == 0 || self.keys.is_empty() { return None; } @@ -707,12 +728,17 @@ impl AtomicKeyTable { ) } - /// How many slots hold a key. - pub fn claimed(&self) -> usize { + /// How many rows the table holds. + pub fn len(&self) -> usize { self.iter().count() } - /// How many slots there are. + /// Whether any row has been created. Matches [`LinearTable::is_empty`]. + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// How many rows the table can hold. Fixed at construction. pub fn capacity(&self) -> usize { self.keys.len() } @@ -729,52 +755,49 @@ mod atomic_key_table_tests { #[test] fn a_claimed_row_is_found_by_a_plain_load_and_never_reclaimed() { let table: AtomicKeyTable = AtomicKeyTable::with_capacity(64); - let first = table.find_or_claim(7).expect("capacity"); + let first = table.upsert(7).expect("capacity"); first.0.fetch_add(1, Ordering::Relaxed); - let again = table.find_or_claim(7).expect("already claimed"); + let again = table.upsert(7).expect("already claimed"); again.0.fetch_add(1, Ordering::Relaxed); assert_eq!( again.0.load(Ordering::Relaxed), 2, "the second call found the same row" ); - assert_eq!(table.claimed(), 1, "one key claimed one slot"); + assert_eq!(table.len(), 1, "one key claimed one slot"); } #[test] fn zero_is_the_empty_sentinel_and_is_refused_rather_than_colliding() { let table: AtomicKeyTable = AtomicKeyTable::with_capacity(8); assert!( - table.find_or_claim(0).is_none(), + table.upsert(0).is_none(), "zero would be indistinguishable from empty" ); - assert_eq!(table.claimed(), 0); + assert_eq!(table.len(), 0); } #[test] fn a_full_table_refuses_rather_than_growing() { let table: AtomicKeyTable = AtomicKeyTable::with_capacity(4); for key in 1..=4 { - assert!(table.find_or_claim(key).is_some(), "slot {key} fits"); + assert!(table.upsert(key).is_some(), "slot {key} fits"); } - assert_eq!(table.claimed(), 4); + assert_eq!(table.len(), 4); + assert!(table.upsert(5).is_none(), "the fifth has nowhere to go"); assert!( - table.find_or_claim(5).is_none(), - "the fifth has nowhere to go" - ); - assert!( - table.find_or_claim(3).is_some(), + table.upsert(3).is_some(), "a claimed key is still reachable when full" ); } #[test] - fn find_never_claims() { + fn select_never_creates_a_row() { let table: AtomicKeyTable = AtomicKeyTable::with_capacity(8); - assert!(table.find(9).is_none()); - assert_eq!(table.claimed(), 0, "find must not take a slot"); - table.find_or_claim(9).expect("capacity"); - assert!(table.find(9).is_some()); + assert!(table.select(9).is_none()); + assert_eq!(table.len(), 0, "find must not take a slot"); + table.upsert(9).expect("capacity"); + assert!(table.select(9).is_some()); } #[test] @@ -782,7 +805,7 @@ mod atomic_key_table_tests { let table: AtomicKeyTable = AtomicKeyTable::with_capacity(32); for key in [11usize, 22, 33] { table - .find_or_claim(key) + .upsert(key) .expect("capacity") .0 .store(key as u64, Ordering::Relaxed); @@ -806,7 +829,7 @@ mod atomic_key_table_tests { for round in 0..1_000usize { let key = (round % 16) + 1; shared - .find_or_claim(key) + .upsert(key) .expect("capacity") .0 .fetch_add(1, Ordering::Relaxed); @@ -815,7 +838,7 @@ mod atomic_key_table_tests { } }); assert_eq!( - table.claimed(), + table.len(), 16, "sixteen keys, sixteen slots, whatever the interleaving" ); @@ -855,7 +878,7 @@ mod atomic_key_table_cost { let atomic: AtomicKeyTable = AtomicKeyTable::with_capacity(KEYS * 4); let mut linear: LinearTable = LinearTable::new(); for key in 1..=KEYS { - atomic.find_or_claim(key).expect("capacity"); + atomic.upsert(key).expect("capacity"); linear.push(key, key as u64); } @@ -863,7 +886,11 @@ mod atomic_key_table_cost { let mut sink = 0u64; for round in 0..ROUNDS { let key = (round % KEYS) + 1; - sink += atomic.find(key).expect("claimed").0.load(Ordering::Relaxed); + sink += atomic + .select(key) + .expect("claimed") + .0 + .load(Ordering::Relaxed); } let atomic_ns = start.elapsed().as_nanos().max(1);