Skip to content
Merged
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
4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ members = ["codegen", "examples", "performance_measurement", "performance_measur

[package]
name = "worktable"
version = "1.0.0-beta.7"
version = "1.0.0-beta.8"
edition = "2024"
authors = ["Handy-caT"]
license = "MIT"
Expand Down Expand Up @@ -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.7" }
worktable_codegen = { path = "codegen", version = "=1.0.0-beta.8" }

[dev-dependencies]
chrono = "0.4.43"
Expand Down
2 changes: 1 addition & 1 deletion codegen/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "worktable_codegen"
version = "1.0.0-beta.7"
version = "1.0.0-beta.8"
edition = "2024"
license = "MIT"
description = "Proc-macro companion crate for worktable: the worktable! macro and its derives."
Expand Down
135 changes: 80 additions & 55 deletions src/persistence/operation/batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@ use crate::persistence::task::{LastEventIds, QueueInnerRow};
use crate::prelude::*;
use crate::prelude::{From, Order, SelectQueryExecutor};

// Ephemeral metadata rebuilt for every persistence batch, not a persisted
// schema. One Multi operation deliberately owns several rows, so operation_id
// is non-unique while pos is the unique association back to the ops vector.
worktable! (
name: BatchInner,
columns: {
Expand Down Expand Up @@ -80,43 +83,51 @@ fn latest_data_writes<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>(
ops: &[Operation<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>],
) -> BatchData {
type PhysicalSlot = (PageId, u32);
type SequencedWrite = (OperationId, usize, Link, Vec<u8>);

let mut latest: HashMap<PhysicalSlot, SequencedWrite> = HashMap::new();
for (sequence, op) in ops.iter().enumerate() {
let Some(bytes) = op.bytes() else {
continue;
};
let link = op.link();
let operation_id = op.operation_id();
let key = (link.page_id, link.offset);
match latest.entry(key) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
if (operation_id, sequence) > (entry.get().0, entry.get().1) {
entry.insert((operation_id, sequence, link, bytes.to_vec()));
}
fn collect_in_order<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>(
ops: &[Operation<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>],
order: impl Iterator<Item = usize> + Clone,
) -> BatchData {
let mut latest: HashMap<PhysicalSlot, usize> = HashMap::with_capacity(ops.len());
for sequence in order.clone() {
let op = &ops[sequence];
if op.bytes().is_some() {
let link = op.link();
latest.insert((link.page_id, link.offset), sequence);
}
std::collections::hash_map::Entry::Vacant(entry) => {
entry.insert((operation_id, sequence, link, bytes.to_vec()));
}

let mut ordered = HashMap::new();
for sequence in order {
let op = &ops[sequence];
let Some(bytes) = op.bytes() else {
continue;
};
let link = op.link();
if latest.get(&(link.page_id, link.offset)) != Some(&sequence) {
continue;
}
ordered
.entry(link.page_id)
.or_insert_with(Vec::new)
.push((link, bytes.to_vec()));
}
ordered
}

let mut ordered = HashMap::new();
for (_, (operation_id, sequence, link, bytes)) in latest {
ordered
.entry(link.page_id)
.or_insert_with(Vec::new)
.push((operation_id, sequence, link, bytes));
// The analyzer already establishes this order. Keep that production path
// linear; only defensive callers that construct an unsorted BatchOperation
// pay for an index sort.
if ops
.windows(2)
.all(|pair| pair[0].operation_id() <= pair[1].operation_id())
{
collect_in_order(ops, 0..ops.len())
} else {
let mut order = (0..ops.len()).collect::<Vec<_>>();
order.sort_unstable_by_key(|sequence| (ops[*sequence].operation_id(), *sequence));
collect_in_order(ops, order.into_iter())
}
ordered
.into_iter()
.map(|(page_id, mut writes)| {
writes.sort_unstable_by_key(|(operation_id, sequence, _, _)| (*operation_id, *sequence));
let writes = writes.into_iter().map(|(_, _, link, bytes)| (link, bytes)).collect();
(page_id, writes)
})
.collect()
}

#[derive(Debug)]
Expand Down Expand Up @@ -151,6 +162,30 @@ where
phantom_data: PhantomData,
}
}

/// Remove metadata immediately after `self.ops.remove(removed_pos)`.
///
/// At entry, `self.ops.len()` is already one shorter while `info_wt` still
/// has the old positions. Shifting upward positions in ascending order
/// keeps every unique destination vacant as it is filled.
async fn remove_info_at_pos(&self, removed_pos: usize) -> eyre::Result<()> {
let row = self
.info_wt
.select_by_pos(removed_pos)
.ok_or_else(|| eyre::eyre!("batch metadata position {removed_pos} is missing"))?;
self.info_wt.delete_without_lock::<_>(row.id).await?;

for old_pos in (removed_pos + 1)..=self.ops.len() {
let row = self
.info_wt
.select_by_pos(old_pos)
.ok_or_else(|| eyre::eyre!("batch metadata position {old_pos} is missing during reindex"))?;
self.info_wt
.update_pos_by_id(PosByIdQuery { pos: old_pos - 1 }, row.id)
.await?;
}
Ok(())
}
}

impl<PrimaryKeyGenState, PrimaryKey, SecondaryEvents, AvailableIndexes>
Expand All @@ -168,7 +203,7 @@ where
async fn remove_operations_from_events(
&mut self,
invalid_events: PreparedIndexEvents<PrimaryKey, SecondaryEvents>,
) -> Vec<Operation<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>> {
) -> eyre::Result<Vec<Operation<PrimaryKeyGenState, PrimaryKey, SecondaryEvents>>> {
let mut removed_ops = Vec::new();

for ev in &invalid_events.primary_evs {
Expand All @@ -186,7 +221,7 @@ where
}) {
let pos = self.ops.len() - (operation_pos_rev + 1);
let op = self.ops.remove(pos);
self.remove_info_at_pos(pos).await;
self.remove_info_at_pos(pos).await?;
removed_ops.push(op);
}
}
Expand All @@ -198,7 +233,7 @@ where
}) {
let pos = self.ops.len() - (operation_pos_rev + 1);
let op = self.ops.remove(pos);
self.remove_info_at_pos(pos).await;
self.remove_info_at_pos(pos).await?;
removed_ops.push(op);
};
// else it was already removed with primary
Expand All @@ -222,26 +257,7 @@ where
prepared_evs.secondary_evs.remove(op_secondary);
}

removed_ops
}

async fn remove_info_at_pos(&self, removed_pos: usize) {
let row = self
.info_wt
.select_by_pos(removed_pos)
.expect("batch metadata position must exist");
self.info_wt.delete_without_lock::<_>(row.id).await.unwrap();

for old_pos in (removed_pos + 1)..=self.ops.len() {
let row = self
.info_wt
.select_by_pos(old_pos)
.expect("later batch metadata position must exist");
self.info_wt
.update_pos_by_id(PosByIdQuery { pos: old_pos - 1 }, row.id)
.await
.unwrap();
}
Ok(removed_ops)
}

pub fn get_last_event_ids(&self) -> LastEventIds<AvailableIndexes> {
Expand Down Expand Up @@ -304,7 +320,7 @@ where
primary_evs: primary_invalid_events,
secondary_evs: secondary_invalid_events,
};
let ops = self.remove_operations_from_events(events_to_remove).await;
let ops = self.remove_operations_from_events(events_to_remove).await?;
ops_to_remove.extend(ops);
}

Expand Down Expand Up @@ -446,7 +462,7 @@ mod tests {
use data_bucket::Link;
use uuid::Uuid;

use super::latest_data_writes;
use super::{BatchOperation, latest_data_writes};
use crate::persistence::operation::{InsertOperation, Operation, OperationId};

fn insert(id: u128, link: Link, bytes: Vec<u8>) -> Operation<(), u64, ()> {
Expand Down Expand Up @@ -557,4 +573,13 @@ mod tests {

assert_eq!(batch.get(&1.into()).unwrap(), &vec![(new_link, vec![2; 6])]);
}

#[tokio::test]
async fn missing_batch_metadata_returns_an_error_instead_of_panicking() {
let batch: BatchOperation<(), u64, (), ()> = BatchOperation::new(vec![], Default::default());

let error = batch.remove_info_at_pos(0).await.unwrap_err();

assert!(error.to_string().contains("batch metadata position 0 is missing"));
}
}
1 change: 1 addition & 0 deletions tests/persistence/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ mod duplicate_key_index_reload;
mod failure;
mod index_page;
mod loaded_index_growth;
mod multi_row_backend_order;
mod read;
mod recovery_load;
mod schema;
Expand Down
108 changes: 108 additions & 0 deletions tests/persistence/multi_row_backend_order.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
use worktable::prelude::*;
use worktable::worktable;

use crate::remove_dir_if_exists;

macro_rules! persisted_multi_row_backend_case {
(
$module:ident,
$name:ident,
$table:ident,
$row:ident,
$engine:ident,
$backend:ident,
$root:literal
) => {
mod $module {
use super::*;

worktable! {
name: $name,
persist: true,
columns: {
id: u64 primary_key autoincrement using $backend,
group_id: u64,
payload: String,
},
indexes: {
group_idx: group_id,
},
queries: {
update: {
PayloadByGroup(payload) by group_id,
}
}
}

#[tokio::test]
async fn multi_row_relocation_survives_persisted_reload() {
const ROOT: &str = $root;
remove_dir_if_exists(ROOT.to_string()).await;
let config = DiskConfig::new_with_table_name(ROOT, $table::name_snake_case(), $table::version());

let engine = $engine::new(config.clone()).await.unwrap();
let table = $table::load(engine).await.unwrap();
for length in 1..=128 {
table
.insert($row {
id: table.get_next_pk().into(),
group_id: 7,
payload: "x".repeat(length),
})
.unwrap();
}
table.wait_for_ops().await.unwrap();

let replacement = "new-payload".repeat(64);
table
.update_payload_by_group(
PayloadByGroupQuery {
payload: replacement.clone(),
},
7,
)
.await
.unwrap();
table.wait_for_ops().await.unwrap();
drop(table);

let engine = $engine::new(config).await.unwrap();
let table = $table::load(engine).await.unwrap();
let rows = table.select_by_group_id(7).execute().unwrap();
assert_eq!(rows.len(), 128);
assert!(rows.iter().all(|row| row.payload == replacement));
table.wait_for_ops().await.unwrap();
drop(table);
remove_dir_if_exists(ROOT.to_string()).await;
}
}
};
}

persisted_multi_row_backend_case!(
wti,
MultiRowWti,
MultiRowWtiWorkTable,
MultiRowWtiRow,
MultiRowWtiPersistenceEngine,
worktables_index,
"tests/data/multi_row_backend_wti"
);
persisted_multi_row_backend_case!(
congee,
MultiRowCongee,
MultiRowCongeeWorkTable,
MultiRowCongeeRow,
MultiRowCongeePersistenceEngine,
congee,
"tests/data/multi_row_backend_congee"
);
persisted_multi_row_backend_case!(
arctic,
MultiRowArctic,
MultiRowArcticWorkTable,
MultiRowArcticRow,
MultiRowArcticPersistenceEngine,
arctic,
"tests/data/multi_row_backend_arctic"
);
Loading