global: snap

This commit is contained in:
nym21
2026-08-05 16:28:47 +02:00
parent c4decdef77
commit 781b0f0b09
46 changed files with 290 additions and 238 deletions
Generated
+56 -56
View File
@@ -491,6 +491,14 @@ dependencies = [
"serde_json",
]
[[package]]
name = "brk_byteview"
version = "0.3.6"
dependencies = [
"serde",
"serde_json",
]
[[package]]
name = "brk_cli"
version = "0.3.6"
@@ -574,7 +582,7 @@ name = "brk_error"
version = "0.3.6"
dependencies = [
"bitcoin",
"fjall",
"brk_fjall",
"jiff",
"jsonrpc",
"pco",
@@ -597,6 +605,23 @@ dependencies = [
"ureq",
]
[[package]]
name = "brk_fjall"
version = "0.3.6"
dependencies = [
"brk_byteview",
"brk_lsm_tree",
"byteorder-lite",
"dashmap",
"flume",
"log",
"lz4_flex",
"nanoid",
"tempfile",
"test-log",
"xxhash-rust",
]
[[package]]
name = "brk_indexer"
version = "0.3.6"
@@ -606,6 +631,7 @@ dependencies = [
"brk_bencher",
"brk_cohort",
"brk_error",
"brk_fjall",
"brk_logger",
"brk_reader",
"brk_rpc",
@@ -613,7 +639,6 @@ dependencies = [
"brk_traversable",
"brk_types",
"color-eyre",
"fjall",
"parking_lot",
"rayon",
"rustc-hash",
@@ -646,6 +671,32 @@ dependencies = [
"tracing-subscriber",
]
[[package]]
name = "brk_lsm_tree"
version = "0.3.6"
dependencies = [
"arc-swap",
"brk_byteview",
"byteorder-lite",
"criterion",
"crossbeam-skiplist",
"fs_extra",
"interval-heap",
"log",
"lz4_flex",
"nanoid",
"quick_cache",
"rand 0.10.2",
"rustc-hash",
"self_cell",
"sfa",
"strum",
"tempfile",
"test-log",
"varint-rs",
"xxhash-rust",
]
[[package]]
name = "brk_mcp"
version = "0.3.6"
@@ -783,10 +834,10 @@ dependencies = [
name = "brk_store"
version = "0.3.6"
dependencies = [
"brk_byteview",
"brk_error",
"brk_fjall",
"brk_types",
"byteview",
"fjall",
"rustc-hash",
"tempfile",
]
@@ -818,8 +869,8 @@ name = "brk_types"
version = "0.3.6"
dependencies = [
"bitcoin",
"brk_byteview",
"brk_error",
"byteview",
"derive_more",
"indexmap",
"itoa",
@@ -902,14 +953,6 @@ version = "1.12.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04"
[[package]]
name = "byteview"
version = "0.3.6"
dependencies = [
"serde",
"serde_json",
]
[[package]]
name = "cast"
version = "0.3.0"
@@ -1589,23 +1632,6 @@ version = "0.1.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5baebc0774151f905a1a2cc41989300b1e6fbb29aff0ceffa1064fdd3088d582"
[[package]]
name = "fjall"
version = "0.3.6"
dependencies = [
"byteorder-lite",
"byteview",
"dashmap",
"flume",
"log",
"lsm-tree",
"lz4_flex",
"nanoid",
"tempfile",
"test-log",
"xxhash-rust",
]
[[package]]
name = "flate2"
version = "1.1.9"
@@ -2404,32 +2430,6 @@ version = "0.4.33"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ceec5bc11778974d1bcb055b18002eba7f4b3518b6a0081b3af5f21666da9ad"
[[package]]
name = "lsm-tree"
version = "0.3.6"
dependencies = [
"arc-swap",
"byteorder-lite",
"byteview",
"criterion",
"crossbeam-skiplist",
"fs_extra",
"interval-heap",
"log",
"lz4_flex",
"nanoid",
"quick_cache",
"rand 0.10.2",
"rustc-hash",
"self_cell",
"sfa",
"strum",
"tempfile",
"test-log",
"varint-rs",
"xxhash-rust",
]
[[package]]
name = "lz4_flex"
version = "0.13.1"
+3 -3
View File
@@ -59,17 +59,17 @@ brk_traversable = { version = "0.3.6", path = "crates/brk_traversable", features
brk_traversable_derive = { version = "0.3.6", path = "crates/brk_traversable_derive" }
brk_types = { version = "0.3.6", path = "crates/brk_types" }
brk_website = { version = "0.3.6", path = "crates/brk_website" }
byteview = { version = "0.3.6", path = "crates/byteview" }
byteview = { package = "brk_byteview", version = "0.3.6", path = "crates/byteview" }
color-eyre = "0.6.5"
corepc-jsonrpc = { package = "jsonrpc", version = "0.19.0", features = ["simple_http"], default-features = false }
corepc-types = { version = "0.15.0", features = ["std"], default-features = false }
derive_more = { version = "2.1.1", features = ["deref", "deref_mut"] }
fjall = { path = "crates/fjall", version = "0.3.6" }
fjall = { package = "brk_fjall", path = "crates/fjall", version = "0.3.6" }
indexmap = { version = "2.14.0", features = ["serde"] }
jiff = { version = "0.2.35", features = ["perf-inline", "tz-system"], default-features = false }
libc = "0.2"
log = "0.4.33"
lsm-tree = { path = "crates/lsm-tree", version = "0.3.6", default-features = false }
lsm-tree = { package = "brk_lsm_tree", path = "crates/lsm-tree", version = "0.3.6", default-features = false }
lz4_flex = { version = "=0.13.1", default-features = false }
owo-colors = "4.3.0"
parking_lot = "0.12.5"
+10 -5
View File
@@ -3,7 +3,7 @@ mod manifest;
mod server;
mod upstream;
use std::{env, error::Error, io};
use std::{env, error::Error, io, process};
use config::api_bases;
use manifest::Catalog;
@@ -19,11 +19,11 @@ async fn main() -> Result<(), Box<dyn Error>> {
let mut arguments = env::args();
let _program = arguments.next();
let api_base = arguments
.next()
.ok_or_else(|| io::Error::other("usage: brk_mcp <REST_API_URL_OR_HOST>"))?;
let Some(api_base) = arguments.next() else {
usage();
};
if arguments.next().is_some() {
return Err(io::Error::other("usage: brk_mcp <REST_API_URL_OR_HOST>").into());
usage();
}
let api_bases = api_bases(&api_base).map_err(io::Error::other)?;
let catalog = Catalog::embedded().map_err(io::Error::other)?;
@@ -36,6 +36,11 @@ async fn main() -> Result<(), Box<dyn Error>> {
Ok(())
}
fn usage() -> ! {
eprintln!("Usage: brk_mcp <REST_API_URL_OR_HOST>");
process::exit(2);
}
async fn bind_available(
start: std::net::SocketAddr,
) -> io::Result<(TcpListener, std::net::SocketAddr)> {
+6 -2
View File
@@ -1,6 +1,6 @@
[package]
name = "byteview"
description = "Thin, immutable zero-copy slice type"
name = "brk_byteview"
description = "BRK-maintained fork of byteview, a thin immutable zero-copy byte slice"
license = "MIT OR Apache-2.0"
version.workspace = true
edition.workspace = true
@@ -10,6 +10,10 @@ homepage.workspace = true
categories = ["data-structures"]
keywords = ["german-string", "string-view", "byte-slice"]
[lib]
name = "byteview"
path = "src/lib.rs"
[features]
default = []
serde = ["dep:serde"]
+5 -3
View File
@@ -1,11 +1,13 @@
# byteview
# brk_byteview
[![CI](https://github.com/fjall-rs/byteview/actions/workflows/test.yml/badge.svg)](https://github.com/fjall-rs/byteview/actions/workflows/test.yml)
[![CI](https://github.com/fjall-rs/byteview/actions/workflows/miri.yml/badge.svg)](https://github.com/fjall-rs/byteview/actions/workflows/miri.yml)
[![docs.rs](https://img.shields.io/docsrs/byteview?color=green)](https://docs.rs/byteview)
[![Crates.io](https://img.shields.io/crates/v/byteview?color=blue)](https://crates.io/crates/byteview)
[![docs.rs](https://img.shields.io/docsrs/brk_byteview?color=green)](https://docs.rs/brk_byteview)
[![Crates.io](https://img.shields.io/crates/v/brk_byteview?color=blue)](https://crates.io/crates/brk_byteview)
![MSRV](https://img.shields.io/badge/MSRV-1.87-blue)
BRK-maintained fork of [`byteview`](https://github.com/fjall-rs/byteview), published separately for use by the [Bitcoin Research Kit](https://bitcoinresearchkit.org). Its Rust library name remains `byteview`.
An immutable byte slice that may be inlined, and can be partially cloned without heap allocation.
Think of it as a specialized `Arc<[u8]>` that can be inlined (skip allocation for small values) and no weak count.
+5 -5
View File
@@ -147,14 +147,13 @@ mod serde {
use serde::de::{self, Visitor};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::fmt;
use std::ops::Deref;
impl Serialize for StrView {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.deref())
serializer.serialize_str(self.as_ref())
}
}
@@ -165,7 +164,7 @@ mod serde {
{
struct StrViewVisitor;
impl<'de> Visitor<'de> for StrViewVisitor {
impl Visitor<'_> for StrViewVisitor {
type Value = StrView;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
@@ -192,10 +191,11 @@ mod tests {
#[cfg(feature = "serde")]
#[test]
fn serde_roundtrip() {
fn serde_roundtrip() -> serde_json::Result<()> {
let a = StrView::from("abcdef");
let b: StrView = serde_json::from_slice(&serde_json::to_vec(&a).unwrap()).unwrap();
let b: StrView = serde_json::from_slice(&serde_json::to_vec(&a)?)?;
assert_eq!(a, b);
Ok(())
}
#[test]
+3 -3
View File
@@ -1,10 +1,10 @@
[package]
name = "fjall"
description = "Log-structured, embeddable key-value storage engine"
name = "brk_fjall"
description = "BRK-maintained fork of Fjall, an embeddable log-structured key-value storage engine"
license.workspace = true
version.workspace = true
edition.workspace = true
readme.workspace = true
readme = "README.md"
repository.workspace = true
homepage.workspace = true
+5
View File
@@ -0,0 +1,5 @@
# brk_fjall
BRK-maintained fork of [Fjall](https://github.com/fjall-rs/fjall), published separately for use by the [Bitcoin Research Kit](https://bitcoinresearchkit.org). Its Rust library name remains `fjall`.
Fjall is an embeddable, log-structured key-value storage engine written in Rust.
+2 -3
View File
@@ -114,8 +114,8 @@ impl WriteBatch {
journal_writer.write_batch(self.data.iter(), self.data.len(), batch_seqno)?;
if let Some(mode) = self.durability {
if let Err(e) = journal_writer.persist(mode) {
if let Some(mode) = self.durability
&& let Err(e) = journal_writer.persist(mode) {
self.db.is_poisoned.poison();
log::error!(
@@ -124,7 +124,6 @@ impl WriteBatch {
return Err(crate::Error::Poisoned);
}
}
// TODO: maybe we can use a stack alloc hashset/vec here, such as smallset
#[expect(clippy::mutable_key_type)]
+1
View File
@@ -181,6 +181,7 @@ impl Builder {
/// #
/// # Ok::<_, fjall::Error>(())
/// ```
#[must_use]
pub fn with_compaction_filter_factories(mut self, f: CompactionFilterAssigner) -> Self {
self.inner.compaction_filter_factory_assigner = Some(f);
self
+4
View File
@@ -295,6 +295,10 @@ impl Database {
/// #
/// # Ok::<(), fjall::Error>(())
/// ```
///
/// # Errors
///
/// Returns an error if journal disk usage cannot be read.
pub fn disk_space(&self) -> crate::Result<u64> {
let journal_size = self.journal_disk_space()?;
+1 -3
View File
@@ -34,9 +34,7 @@ impl std::fmt::Debug for Journal {
write!(
f,
"{}",
self.path()
.map(|p| p.display().to_string())
.unwrap_or_else(|_| String::from("<failed to read path>"))
self.path().map_or_else(|_| String::from("<failed to read path>"), |p| p.display().to_string())
)
}
}
+3 -2
View File
@@ -771,10 +771,11 @@ impl Keyspace {
.expect("lock is poisoned")
.values()
{
if let Err(e) = keyspace.tree.get_version_history_lock().maintenance(
let maintenance_result = keyspace.tree.get_version_history_lock().maintenance(
keyspace.path(),
self.supervisor.snapshot_tracker.get_seqno_safe_to_gc(),
) {
);
if let Err(e) = maintenance_result {
log::warn!(
"Version history GC failed for keyspace {:?}: {e:?}",
keyspace.name,
+1 -3
View File
@@ -134,7 +134,7 @@ impl CreateOptions {
self
}
#[expect(clippy::expect_used, clippy::too_many_lines)]
#[expect(clippy::expect_used)]
pub(crate) fn from_kvs(
keyspace_id: InternalKeyspaceId,
meta_keyspace: &MetaKeyspace,
@@ -245,7 +245,6 @@ impl CreateOptions {
})
}
#[expect(clippy::too_many_lines)]
pub(crate) fn encode_kvs(&self, keyspace_id: InternalKeyspaceId) -> Vec<KvPair> {
use crate::keyspace::config::EncodeConfig;
@@ -507,7 +506,6 @@ mod tests {
}
#[test]
#[expect(clippy::unwrap_used)]
#[cfg(feature = "lz4")]
fn keyspace_opts_compression_default() {
use CompressionType::{Lz4, None as Uncompressed};
+17 -2
View File
@@ -79,9 +79,24 @@
#![deny(clippy::unwrap_used)]
#![deny(clippy::indexing_slicing)]
#![warn(clippy::pedantic, clippy::nursery)]
#![warn(clippy::expect_used)]
#![allow(
clippy::expect_used,
clippy::missing_panics_doc,
reason = "poisoned locks and violated persisted-data invariants are unrecoverable"
)]
#![allow(clippy::missing_const_for_fn, clippy::significant_drop_tightening)]
#![warn(clippy::multiple_crate_versions)]
#![allow(
clippy::multiple_crate_versions,
reason = "transitive dependencies currently require distinct hashbrown versions"
)]
#![cfg_attr(
test,
allow(
clippy::items_after_statements,
clippy::unwrap_used,
reason = "test fixtures favor direct assertions and local helper imports"
)
)]
#![cfg_attr(docsrs, feature(doc_cfg))]
macro_rules! fail_iter {
+9 -9
View File
@@ -243,7 +243,7 @@ mod tests {
None,
db.meta_keyspace
.inner
.get(&['n' as u8, 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.get([b'n', 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.as_deref(),
);
@@ -253,7 +253,7 @@ mod tests {
Some(b"items".as_slice()),
db.meta_keyspace
.inner
.get(&['n' as u8, 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.get([b'n', 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.as_deref(),
);
assert_eq!(
@@ -261,9 +261,9 @@ mod tests {
db.meta_keyspace
.inner
.get(
&[
'c' as u8, 0, 0, 0, 0, 0, 0, 0, 1, 'v' as u8, 'e' as u8, 'r' as u8,
's' as u8, 'i' as u8, 'o' as u8, 'n' as u8
[
b'c', 0, 0, 0, 0, 0, 0, 0, 1, b'v', b'e', b'r',
b's', b'i', b'o', b'n'
],
SeqNo::MAX
)?
@@ -279,7 +279,7 @@ mod tests {
Some(b"items".as_slice()),
db.meta_keyspace
.inner
.get(&['n' as u8, 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.get([b'n', 0, 0, 0, 0, 0, 0, 0, 1], SeqNo::MAX)?
.as_deref(),
);
assert_eq!(
@@ -287,9 +287,9 @@ mod tests {
db.meta_keyspace
.inner
.get(
&[
'c' as u8, 0, 0, 0, 0, 0, 0, 0, 1, 'v' as u8, 'e' as u8, 'r' as u8,
's' as u8, 'i' as u8, 'o' as u8, 'n' as u8
[
b'c', 0, 0, 0, 0, 0, 0, 0, 1, b'v', b'e', b'r',
b's', b'i', b'o', b'n'
],
SeqNo::MAX
)?
+1 -1
View File
@@ -376,7 +376,7 @@ mod tests {
let big = 100u64;
map.publish(big);
assert!(map.get() == (big + 1));
assert_eq!(map.get(), (big + 1));
let before = map.get();
map.publish(1);
+3 -3
View File
@@ -1,10 +1,10 @@
[package]
name = "lsm-tree"
description = "A K.I.S.S. implementation of log-structured merge trees (LSM-trees/LSMTs)"
name = "brk_lsm_tree"
description = "BRK-maintained fork of lsm-tree, a minimal log-structured merge tree implementation"
license.workspace = true
version.workspace = true
edition.workspace = true
readme.workspace = true
readme = "README.md"
repository.workspace = true
homepage.workspace = true
+5
View File
@@ -0,0 +1,5 @@
# brk_lsm_tree
BRK-maintained fork of [`lsm-tree`](https://github.com/fjall-rs/lsm-tree), published separately for use by the [Bitcoin Research Kit](https://bitcoinresearchkit.org). Its Rust library name remains `lsm_tree`.
`lsm-tree` is a minimal implementation of log-structured merge trees written in Rust.
+1 -1
View File
@@ -43,7 +43,7 @@ fn mvcc_stream(c: &mut Criterion) {
c.bench_function(&format!("MVCC stream {num} versions"), |b| {
let memtables = (0..num)
.map(|id| {
let table = Memtable::new(id as u64);
let table = Memtable::new(id);
for key in 'a'..='z' {
table.insert(InternalValue::from_components(
+4 -4
View File
@@ -210,7 +210,7 @@ fn tree_get_pairs(c: &mut Criterion) {
}
group.bench_function(
&format!("Tree::first_key_value (disjoint), {segment_count} segments"),
format!("Tree::first_key_value (disjoint), {segment_count} segments"),
|b| {
b.iter(|| {
assert!(tree.first_key_value(SeqNo::MAX, None).is_some());
@@ -219,7 +219,7 @@ fn tree_get_pairs(c: &mut Criterion) {
);
group.bench_function(
&format!("Tree::last_key_value (disjoint), {segment_count} segments"),
format!("Tree::last_key_value (disjoint), {segment_count} segments"),
|b| {
b.iter(|| {
assert!(tree.last_key_value(SeqNo::MAX, None).is_some());
@@ -250,7 +250,7 @@ fn tree_get_pairs(c: &mut Criterion) {
}
group.bench_function(
&format!("Tree::first_key_value (non-disjoint), {segment_count} segments"),
format!("Tree::first_key_value (non-disjoint), {segment_count} segments"),
|b| {
b.iter(|| {
assert!(tree.first_key_value(SeqNo::MAX, None).is_some());
@@ -259,7 +259,7 @@ fn tree_get_pairs(c: &mut Criterion) {
);
group.bench_function(
&format!("Tree::last_key_value (non-disjoint), {segment_count} segments"),
format!("Tree::last_key_value (non-disjoint), {segment_count} segments"),
|b| {
b.iter(|| {
assert!(tree.last_key_value(SeqNo::MAX, None).is_some());
@@ -237,7 +237,12 @@ impl CompactionStrategy for Strategy {
NAME
}
#[expect(clippy::too_many_lines)]
#[expect(
clippy::cast_possible_truncation,
clippy::expect_used,
clippy::too_many_lines,
reason = "the asserted seven-level invariant guarantees every level exists and its index fits in u8"
)]
fn choose(&self, version: &Version, _: &Config, state: &CompactionState) -> Choice {
assert!(version.level_count() == 7, "should have exactly 7 levels");
+18 -23
View File
@@ -614,7 +614,7 @@ mod tests {
!seqno,
)
.unwrap();
cursor.write_u8(if tomb { 1 } else { 0 }).unwrap();
cursor.write_u8(u8::from(tomb)).unwrap();
debug_assert_eq!(len, cursor.position() as usize);
@@ -630,7 +630,7 @@ mod tests {
/// The previous user key
///
/// Note that the user key is NOT the full KV key
/// because we embed MVCC information into the key (user_key#seqno#type).
/// because we embed MVCC information into the key (`user_key#seqno#type`).
prev_user_key: Option<UserKey>,
/// MVCC watermark we can safely delete if an item < watermark
@@ -645,30 +645,27 @@ mod tests {
// User key len
let ukl = l - TRAILER_SIZE;
match &self.prev_user_key {
Some(prev) => {
let user_key = &value.key.user_key[..ukl];
if let Some(prev) = &self.prev_user_key {
let user_key = &value.key.user_key[..ukl];
if prev == &user_key {
// We found another, older version of the previous key
let mut seqno = &value.key.user_key[(ukl + 1)..l - 1];
debug_assert_eq!(8, seqno.len());
if prev == &user_key {
// We found another, older version of the previous key
let mut seqno = &value.key.user_key[(ukl + 1)..l - 1];
debug_assert_eq!(8, seqno.len());
// IMPORTANT: Invert the seqno back to normal value
let seqno = !seqno.read_u64::<BE>().unwrap();
// IMPORTANT: Invert the seqno back to normal value
let seqno = !seqno.read_u64::<BE>().unwrap();
if seqno < self.mvcc_watermark {
return Ok(StreamFilterVerdict::Drop);
}
} else {
let user_key = &value.key.user_key.slice(..ukl);
self.prev_user_key = Some(user_key.clone());
if seqno < self.mvcc_watermark {
return Ok(StreamFilterVerdict::Drop);
}
}
None => {
} else {
let user_key = &value.key.user_key.slice(..ukl);
self.prev_user_key = Some(user_key.clone());
}
} else {
let user_key = &value.key.user_key.slice(..ukl);
self.prev_user_key = Some(user_key.clone());
}
Ok(StreamFilterVerdict::Keep)
@@ -677,11 +674,9 @@ mod tests {
#[test]
fn compaction_filter_custom_mvcc() {
let vec = vec![
kv(b"abc", 4, b"c", false),
let vec = [kv(b"abc", 4, b"c", false),
kv(b"abc", 3, b"b", false),
kv(b"abc", 2, b"a", false),
];
kv(b"abc", 2, b"a", false)];
let iter = vec.iter().cloned().map(Ok);
let iter = CompactionStream::new(iter, 995).with_filter(Filter {
+1
View File
@@ -198,6 +198,7 @@ fn move_tables(
log::error!("Manifest maintenance failed: {e:?}");
return Err(e);
}
drop(version_history_lock);
Ok(())
}
+14
View File
@@ -47,6 +47,20 @@
#![warn(clippy::multiple_crate_versions)]
#![allow(clippy::option_if_let_else)]
#![warn(clippy::redundant_feature_names)]
#![cfg_attr(
test,
allow(
clippy::cast_possible_truncation,
clippy::cast_precision_loss,
clippy::indexing_slicing,
clippy::items_after_statements,
clippy::too_many_lines,
clippy::unwrap_used,
clippy::useless_vec,
clippy::from_iter_instead_of_collect,
reason = "test fixtures favor direct assertions and intentionally bounded inputs"
)
)]
#![cfg_attr(coverage_nightly, feature(coverage_attribute))]
#[doc(hidden)]
@@ -66,16 +66,14 @@ impl Iter {
fn init_tli(&mut self) -> bool {
let mut iter = OwnedIndexBlockIter::new(self.tli_block.clone(), IndexBlock::iter);
if let Some((lo_key, lo_seqno)) = &self.lo {
if !iter.seek_lower(lo_key, *lo_seqno) {
if let Some((lo_key, lo_seqno)) = &self.lo
&& !iter.seek_lower(lo_key, *lo_seqno) {
return false;
}
}
if let Some((hi_key, hi_seqno)) = &self.hi {
if !iter.seek_upper(hi_key, *hi_seqno) {
if let Some((hi_key, hi_seqno)) = &self.hi
&& !iter.seek_upper(hi_key, *hi_seqno) {
return false;
}
}
self.tli = Some(iter);
@@ -99,11 +97,10 @@ impl Iterator for Iter {
type Item = crate::Result<KeyedBlockHandle>;
fn next(&mut self) -> Option<Self::Item> {
if let Some(lo_block) = &mut self.lo_consumer {
if let Some(item) = lo_block.next() {
if let Some(lo_block) = &mut self.lo_consumer
&& let Some(item) = lo_block.next() {
return Some(Ok(item));
}
}
if self.tli.is_none() && !self.init_tli() {
return None;
@@ -126,16 +123,14 @@ impl Iterator for Iter {
let mut iter = OwnedIndexBlockIter::new(index_block, IndexBlock::iter);
if let Some((lo_key, lo_seqno)) = &self.lo {
if !iter.seek_lower(lo_key, *lo_seqno) {
if let Some((lo_key, lo_seqno)) = &self.lo
&& !iter.seek_lower(lo_key, *lo_seqno) {
return None;
}
}
if let Some((hi_key, hi_seqno)) = &self.hi {
if !iter.seek_upper(hi_key, *hi_seqno) {
if let Some((hi_key, hi_seqno)) = &self.hi
&& !iter.seek_upper(hi_key, *hi_seqno) {
return None;
}
}
let next_item = iter.next().map(Ok);
@@ -148,11 +143,10 @@ impl Iterator for Iter {
}
// Nothing more found, consume from hi consumer
if let Some(hi_block) = &mut self.hi_consumer {
if let Some(item) = hi_block.next() {
if let Some(hi_block) = &mut self.hi_consumer
&& let Some(item) = hi_block.next() {
return Some(Ok(item));
}
}
None
}
@@ -160,11 +154,10 @@ impl Iterator for Iter {
impl DoubleEndedIterator for Iter {
fn next_back(&mut self) -> Option<Self::Item> {
if let Some(hi_block) = &mut self.hi_consumer {
if let Some(item) = hi_block.next_back() {
if let Some(hi_block) = &mut self.hi_consumer
&& let Some(item) = hi_block.next_back() {
return Some(Ok(item));
}
}
if self.tli.is_none() && !self.init_tli() {
return None;
@@ -187,16 +180,14 @@ impl DoubleEndedIterator for Iter {
let mut iter = OwnedIndexBlockIter::new(index_block, IndexBlock::iter);
if let Some((lo_key, lo_seqno)) = &self.lo {
if !iter.seek_lower(lo_key, *lo_seqno) {
if let Some((lo_key, lo_seqno)) = &self.lo
&& !iter.seek_lower(lo_key, *lo_seqno) {
return None;
}
}
if let Some((hi_key, hi_seqno)) = &self.hi {
if !iter.seek_upper(hi_key, *hi_seqno) {
if let Some((hi_key, hi_seqno)) = &self.hi
&& !iter.seek_upper(hi_key, *hi_seqno) {
return None;
}
}
let next_item = iter.next_back().map(Ok);
@@ -209,11 +200,10 @@ impl DoubleEndedIterator for Iter {
}
// Nothing more found, consume from lo consumer
if let Some(lo_block) = &mut self.lo_consumer {
if let Some(item) = lo_block.next_back() {
if let Some(lo_block) = &mut self.lo_consumer
&& let Some(item) = lo_block.next_back() {
return Some(Ok(item));
}
}
None
}
@@ -101,16 +101,14 @@ impl Iterator for Iter {
let mut iter = OwnedIndexBlockIter::new(index_block, IndexBlock::iter);
if let Some((lo_key, lo_seqno)) = &self.lo {
if !iter.seek_lower(lo_key, *lo_seqno) {
if let Some((lo_key, lo_seqno)) = &self.lo
&& !iter.seek_lower(lo_key, *lo_seqno) {
return None;
}
}
if let Some((hi_key, hi_seqno)) = &self.hi {
if !iter.seek_upper(hi_key, *hi_seqno) {
if let Some((hi_key, hi_seqno)) = &self.hi
&& !iter.seek_upper(hi_key, *hi_seqno) {
return None;
}
}
let next_item = iter.next().map(Ok);
@@ -139,16 +137,14 @@ impl DoubleEndedIterator for Iter {
let mut iter = OwnedIndexBlockIter::new(index_block, IndexBlock::iter);
if let Some((lo_key, lo_seqno)) = &self.lo {
if !iter.seek_lower(lo_key, *lo_seqno) {
if let Some((lo_key, lo_seqno)) = &self.lo
&& !iter.seek_lower(lo_key, *lo_seqno) {
return None;
}
}
if let Some((hi_key, hi_seqno)) = &self.hi {
if !iter.seek_upper(hi_key, *hi_seqno) {
if let Some((hi_key, hi_seqno)) = &self.hi
&& !iter.seek_upper(hi_key, *hi_seqno) {
return None;
}
}
let next_item = iter.next_back().map(Ok);
+3 -1
View File
@@ -48,7 +48,9 @@ impl BloomConstructionPolicy {
#[expect(
clippy::cast_precision_loss,
reason = "truncation is fine because this is an estimation"
clippy::cast_possible_truncation,
clippy::cast_sign_loss,
reason = "this positive estimate intentionally floors to a whole number of bytes"
)]
match self {
Self::BitsPerKey(bpk) => (*bpk * (n as f32)) as usize / 8,
+5 -1
View File
@@ -27,7 +27,11 @@ impl From<(TreeId, TableId)> for GlobalTableId {
}
}
pub(crate) fn next_table_id(counter: &SequenceNumberCounter) -> TableId {
#[expect(
clippy::expect_used,
reason = "exhausting the complete u32 table ID space is unrecoverable"
)]
pub fn next_table_id(counter: &SequenceNumberCounter) -> TableId {
counter.next().try_into().expect("ran out of table IDs")
}
+10 -15
View File
@@ -160,8 +160,8 @@ impl Iterator for Iter {
fn next(&mut self) -> Option<Self::Item> {
// Always try to keep iterating inside the already-materialized low data block first; this
// lets callers consume multiple entries without touching the index or cache again.
if let Some(block) = &mut self.lo_data_block {
if let Some(item) = block
if let Some(block) = &mut self.lo_data_block
&& let Some(item) = block
.next()
.map(|mut v| {
v.key.seqno += self.global_seqno;
@@ -171,7 +171,6 @@ impl Iterator for Iter {
{
return Some(item);
}
}
if !self.index_initialized {
// Lazily initialize the index iterator here (not in `new`) so callers can set bounds
@@ -214,8 +213,8 @@ impl Iterator for Iter {
let Some(handle) = self.index_iter.next() else {
// No more block handles coming from the index. Flush any pending items buffered on
// the high side (used by reverse iteration) before signalling completion.
if let Some(block) = &mut self.hi_data_block {
if let Some(item) = block
if let Some(block) = &mut self.hi_data_block
&& let Some(item) = block
.next()
.map(|mut v| {
v.key.seqno += self.global_seqno;
@@ -225,7 +224,6 @@ impl Iterator for Iter {
{
return Some(item);
}
}
// Nothing left to serve; drop both buffers so the iterator can be reused safely.
self.lo_data_block = None;
@@ -284,8 +282,8 @@ impl DoubleEndedIterator for Iter {
fn next_back(&mut self) -> Option<Self::Item> {
// Mirror the forward iterator: prefer consuming buffered items from the high data block to
// avoid touching the index once a block has been materialized.
if let Some(block) = &mut self.hi_data_block {
if let Some(item) = block
if let Some(block) = &mut self.hi_data_block
&& let Some(item) = block
.next_back()
.map(|mut v| {
v.key.seqno += self.global_seqno;
@@ -295,7 +293,6 @@ impl DoubleEndedIterator for Iter {
{
return Some(item);
}
}
if !self.index_initialized {
// Mirror forward iteration: initialize lazily so bounds can be applied up-front. The
@@ -310,14 +307,13 @@ impl DoubleEndedIterator for Iter {
true
};
if ok {
if let Some(bound) = &self.range.1 {
if ok
&& let Some(bound) = &self.range.1 {
let key = match bound {
Bound::Included(k) | Bound::Excluded(k) => k,
};
ok = self.index_iter.seek_upper(key, u64::MAX);
}
}
self.index_initialized = true;
@@ -333,8 +329,8 @@ impl DoubleEndedIterator for Iter {
let Some(handle) = self.index_iter.next_back() else {
// Once we exhaust the index in reverse order, flush any items that were buffered on
// the low side (set when iterating forward first) before signalling completion.
if let Some(block) = &mut self.lo_data_block {
if let Some(item) = block
if let Some(block) = &mut self.lo_data_block
&& let Some(item) = block
.next_back()
.map(|mut v| {
v.key.seqno += self.global_seqno;
@@ -344,7 +340,6 @@ impl DoubleEndedIterator for Iter {
{
return Some(item);
}
}
// Nothing left to produce; reset both buffers to keep the iterator reusable.
self.lo_data_block = None;
+5 -8
View File
@@ -99,15 +99,12 @@ impl ParsedMeta {
let block = DataBlock::new(block);
#[expect(clippy::indexing_slicing)]
{
let table_version = block
.point_read(b"table_version", SeqNo::MAX)
.expect("Table version should exist")
.value;
let table_version = block
.point_read(b"table_version", SeqNo::MAX)
.expect("Table version should exist")
.value;
assert_eq!(&*table_version, [5], "unsupported table version");
}
assert_eq!(&*table_version, [5], "unsupported table version");
{
let hash_type = block
+2 -2
View File
@@ -386,10 +386,10 @@ impl Table {
}
/// Tries to recover a table from a file.
#[warn(
#[expect(
clippy::too_many_arguments,
clippy::too_many_lines,
reason = "TODO: refactor"
reason = "table recovery mirrors the complete persisted table configuration"
)]
pub fn recover(
file_path: PathBuf,
+9 -11
View File
@@ -32,11 +32,10 @@ fn test_with_table(
}
for (idx, item) in items.iter().enumerate() {
if let Some(rotate) = rotate_every {
if idx % rotate == 0 {
if let Some(rotate) = rotate_every
&& idx % rotate == 0 {
writer.spill_block()?;
}
}
writer.write(item.clone())?;
}
let (_, checksum) = writer.finish()?.unwrap();
@@ -175,11 +174,10 @@ fn test_with_table(
}
for (idx, item) in items.iter().enumerate() {
if let Some(rotate) = rotate_every {
if idx % rotate == 0 {
if let Some(rotate) = rotate_every
&& idx % rotate == 0 {
writer.spill_block()?;
}
}
writer.write(item.clone())?;
}
let (_, checksum) = writer.finish()?.unwrap();
@@ -198,7 +196,7 @@ fn test_with_table(
assert_eq!(0, table.id());
assert_eq!(items.len(), table.metadata.item_count as usize);
assert!(table.regions.index.is_some(), "should use two-level index",);
assert!(table.regions.index.is_some(), "should use two-level index");
assert_eq!(0, table.pinned_filter_size(), "should not pin filter");
assert!(matches!(
table.file_accessor,
@@ -222,7 +220,7 @@ fn test_with_table(
assert_eq!(0, table.id());
assert_eq!(items.len(), table.metadata.item_count as usize);
assert!(table.regions.index.is_some(), "should use two-level index",);
assert!(table.regions.index.is_some(), "should use two-level index");
// assert!(table.pinned_filter_size() > 0, "should pin filter");
assert!(matches!(
table.file_accessor,
@@ -246,7 +244,7 @@ fn test_with_table(
assert_eq!(0, table.id());
assert_eq!(items.len(), table.metadata.item_count as usize);
assert!(table.regions.index.is_some(), "should use two-level index",);
assert!(table.regions.index.is_some(), "should use two-level index");
assert!(table.pinned_block_index_size() > 0, "should pin index");
// assert_eq!(0, table.pinned_filter_size(), "should not pin filter");
assert!(matches!(
@@ -271,7 +269,7 @@ fn test_with_table(
assert_eq!(0, table.id());
assert_eq!(items.len(), table.metadata.item_count as usize);
assert!(table.regions.index.is_some(), "should use two-level index",);
assert!(table.regions.index.is_some(), "should use two-level index");
assert!(table.pinned_block_index_size() > 0, "should pin index");
// assert!(table.pinned_filter_size() > 0, "should pin filter");
assert!(matches!(
@@ -296,7 +294,7 @@ fn test_with_table(
assert_eq!(0, table.id());
assert_eq!(items.len(), table.metadata.item_count as usize);
assert!(table.regions.index.is_some(), "should use two-level index",);
assert!(table.regions.index.is_some(), "should use two-level index");
assert!(table.pinned_block_index_size() > 0, "should pin index");
// assert!(table.pinned_filter_size() > 0, "should pin filter");
assert!(matches!(table.file_accessor, FileAccessor::File(..)));
+3 -5
View File
@@ -254,14 +254,12 @@ impl Writer {
self.meta.weak_tombstone_count += 1;
}
if value_type == ValueType::Value {
if let Some((prev_key, prev_type)) = &self.previous_item {
if prev_type == &ValueType::WeakTombstone && prev_key.as_ref() == user_key.as_ref()
if value_type == ValueType::Value
&& let Some((prev_key, prev_type)) = &self.previous_item
&& prev_type == &ValueType::WeakTombstone && prev_key.as_ref() == user_key.as_ref()
{
self.meta.weak_tombstone_reclaimable_count += 1;
}
}
}
// NOTE: Check if we visit a new key
if Some(&user_key) != self.current_key.as_ref() {
+7 -1
View File
@@ -194,6 +194,10 @@ impl<'a> Ingestion<'a> {
/// # Errors
///
/// Will return `Err` if an IO error occurs.
///
/// # Panics
///
/// In debug builds, panics if the version-history lock is poisoned.
#[allow(clippy::significant_drop_tightening)]
pub fn finish_exclusive(self) -> crate::Result<()> {
#[cfg(debug_assertions)]
@@ -265,7 +269,7 @@ impl<'a> Ingestion<'a> {
// compaction state lock and version history lock to safely modify
// the tree's version.
#[expect(clippy::expect_used, reason = "lock is expected to not be poisoned")]
let mut _compaction_state = self.tree.compaction_state.lock().expect("lock is poisoned");
let compaction_state = self.tree.compaction_state.lock().expect("lock is poisoned");
#[expect(clippy::expect_used, reason = "lock is expected to not be poisoned")]
let mut version_lock = self.tree.version_history.write().expect("lock is poisoned");
@@ -321,6 +325,8 @@ impl<'a> Ingestion<'a> {
if let Err(e) = version_lock.maintenance(&self.tree.config.path, 0) {
log::warn!("Version GC failed: {e:?}");
}
drop(version_lock);
drop(compaction_state);
Ok(())
}
+4
View File
@@ -24,6 +24,10 @@ pub type TreeId = u32;
pub type MemtableId = u64;
/// Hands out a unique (monotonically increasing) tree ID.
#[expect(
clippy::expect_used,
reason = "exhausting the complete u32 tree ID space is unrecoverable"
)]
pub fn get_next_tree_id() -> TreeId {
static TREE_ID_COUNTER: AtomicU64 = AtomicU64::new(0);
TREE_ID_COUNTER
+5 -1
View File
@@ -99,6 +99,10 @@ impl AbstractTree for Tree {
}
fn next_table_id(&self) -> TableId {
#[expect(
clippy::expect_used,
reason = "exhausting the complete u32 table ID space is unrecoverable"
)]
self.0
.table_id_counter
.get()
@@ -1132,7 +1136,7 @@ impl Tree {
log::debug!("Successfully recovered {} tables", tables.len());
let version = Version::from_recovery(recovery, &tables)?;
let version = Version::from_recovery(&recovery, &tables)?;
// NOTE: Cleanup old versions
// But only after we definitely recovered the latest version
+9 -7
View File
@@ -185,7 +185,7 @@ impl Version {
}
}
pub(crate) fn from_recovery(recovery: Recovery, tables: &[Table]) -> crate::Result<Self> {
pub(crate) fn from_recovery(recovery: &Recovery, tables: &[Table]) -> crate::Result<Self> {
let version_levels = recovery
.table_ids
.iter()
@@ -306,6 +306,10 @@ impl Version {
/// Returns a new version with a list of tables removed.
///
/// The table files are not immediately deleted, this is handled by the version system's free list.
#[expect(
clippy::unnecessary_wraps,
reason = "preserves the fallible version-transformation API"
)]
pub fn with_dropped(&self, ids: &[TableId]) -> crate::Result<Self> {
let id = self.id + 1;
@@ -353,11 +357,10 @@ impl Version {
.filter(|x| !x.is_empty())
.collect::<Vec<_>>();
if level_idx == dest_level {
if let Some(run) = Run::new(new_tables.to_vec()) {
if level_idx == dest_level
&& let Some(run) = Run::new(new_tables.to_vec()) {
runs.insert(0, run);
}
}
let runs = optimize_runs(runs);
@@ -395,11 +398,10 @@ impl Version {
.filter(|x| !x.is_empty())
.collect::<Vec<_>>();
if level_idx == dest_level {
if let Some(run) = Run::new(affected_tables.clone()) {
if level_idx == dest_level
&& let Some(run) = Run::new(affected_tables.clone()) {
runs.insert(0, run);
}
}
let runs = optimize_runs(runs);
+4 -4
View File
@@ -89,16 +89,16 @@ pub fn recover(folder: &Path) -> crate::Result<Recovery> {
}
let tree_type = {
let byte = toc
toc
.section(b"tree_type")
.ok_or(crate::Error::Unrecoverable)
.inspect_err(|_| {
log::error!("tree_type section not found in version #{curr_version_id} - maybe the file is corrupted?");
})?
.buf_reader(&version_file_path)?
.read_u8()?;
byte
.read_u8()?
};
if tree_type != 0 {
+4
View File
@@ -99,6 +99,10 @@ impl<T: Ranged> Run<T> {
// find last index where pred holds
let end = s.iter().rposition(&pred).map_or(start, |i| i + 1);
#[expect(
clippy::expect_used,
reason = "start and end are derived from positions in the same slice"
)]
s.get(start..end).expect("should be in range")
}
+1 -1
View File
@@ -2,7 +2,7 @@
name = "quickmatch"
description = "Lightning-fast fuzzy string matching with hybrid word and trigram indexing"
version.workspace = true
readme.workspace = true
readme = "docs/README.md"
license.workspace = true
edition.workspace = true
repository.workspace = true
+1 -1
View File
@@ -8,7 +8,7 @@ edition.workspace = true
license.workspace = true
homepage.workspace = true
repository.workspace = true
readme.workspace = true
readme = "README.md"
[dependencies]
libc = { workspace = true }
+1 -1
View File
@@ -8,7 +8,7 @@ edition.workspace = true
license.workspace = true
homepage.workspace = true
repository.workspace = true
readme.workspace = true
readme = "README.md"
[features]
derive = ["vecdb_derive"]
+4 -4
View File
@@ -8,7 +8,7 @@ High-performance mutable persistent vectors built on [`rawdb`](../rawdb/README.m
- **Multiple storage formats**:
- **Raw**: `BytesVec`, `ZeroCopyVec` (uncompressed)
- **Compressed**: `PcoVec`, `LZ4Vec`, `ZstdVec`
- **Computed vectors**: `EagerVec` (stored computations), `LazyVec` (single-source on-the-fly computation)
- **Computed vectors**: `EagerVec` (stored computations), `LazyVecFrom1/2/3` (on-the-fly computation)
- **Rollback support**: Time-travel via stamped change deltas without full snapshots
- **Sparse deletions**: Delete elements leaving holes, no reindexing required
- **Thread-safe**: Concurrent reads with exclusive writes
@@ -144,14 +144,14 @@ let mut derived: EagerVec<BytesVec<usize, f64>> =
// derived.compute_sma(&source, 20)?;
```
**`LazyVec<...>`** - Lazily computed vector from one source vector
**`LazyVecFrom1/2/3<...>`** - Lazily computed vectors from 1-3 source vectors
Values computed on-the-fly during iteration, nothing stored on disk. Use for temporary views or simple transformations.
```rust,ignore
use vecdb::LazyVec;
use vecdb::LazyVecFrom1;
let lazy = LazyVec::init(
let lazy = LazyVecFrom1::init(
"computed",
Version::TWO,
Box::new(source.clone()), // ScannableBoxedVec
+1 -1
View File
@@ -8,7 +8,7 @@ edition.workspace = true
license.workspace = true
homepage.workspace = true
repository.workspace = true
readme.workspace = true
readme = "README.md"
[lib]
proc-macro = true