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: 3 additions & 1 deletion crates/storage/src/api/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,4 +17,6 @@ mod tables;
mod traits;

pub use tables::{ALL_TABLES, Table};
pub use traits::{Error, PrefixResult, StorageBackend, StorageReadView, StorageWriteBatch};
pub use traits::{
Error, PrefixResult, StorageBackend, StorageReadView, StorageReadViewExt, StorageWriteBatch,
};
61 changes: 59 additions & 2 deletions crates/storage/src/api/traits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,38 @@ pub trait StorageBackend: Send + Sync {

/// A read-only view of the storage.
pub trait StorageReadView {
/// Get a value by key from a table.
fn get(&self, table: Table, key: &[u8]) -> Result<Option<Vec<u8>>, Error>;
/// Calls `read_fn` with a borrow of the value stored under `key`, if any,
/// and returns whether the key was present.
///
/// The value is never copied: `read_fn` sees the backend's own buffer, so
/// a caller that only decodes or inspects the bytes avoids allocating a
/// value-sized `Vec` (full state snapshots are 100+ MB on mainnet-sized
/// beacon chains). `read_fn` runs at most once, and its error is returned
/// as-is. A `&mut dyn FnMut` rather than a generic closure, so the trait
/// stays usable as `dyn StorageReadView`.
fn read(
&self,
table: Table,
key: &[u8],
read_fn: &mut dyn FnMut(&[u8]) -> Result<(), Error>,
) -> Result<bool, Error>;

/// Get a value by key from a table, copied into an owned `Vec`.
///
/// Prefer [`StorageReadViewExt::read_with`] when the bytes are only
/// decoded: it decodes from the backend's buffer without the copy.
fn get(&self, table: Table, key: &[u8]) -> Result<Option<Vec<u8>>, Error> {
self.read_with(table, key, <[u8]>::to_vec)
}

/// Whether `key` is present in a table.
///
/// Same answer as `get(..)?.is_some()` but never materializes the value, so
/// the cost does not scale with its size. Prefer it for pure existence
/// checks on large values.
fn contains(&self, table: Table, key: &[u8]) -> Result<bool, Error> {
self.read(table, key, &mut |_| Ok(()))
}

/// Iterate over all entries with a given key prefix.
fn prefix_iterator(
Expand All @@ -35,6 +65,33 @@ pub trait StorageReadView {
) -> Result<Box<dyn Iterator<Item = PrefixResult> + '_>, Error>;
}

/// Generic conveniences over [`StorageReadView::read`].
///
/// A separate trait because its methods are generic, which would make
/// `StorageReadView` itself unusable as a trait object. The blanket impl
/// covers every view, `dyn StorageReadView` included.
pub trait StorageReadViewExt: StorageReadView {
/// Decodes the value stored under `key` straight from the backend's
/// buffer, without copying it first. `None` when the key is absent.
fn read_with<T>(
&self,
table: Table,
key: &[u8],
decode: impl FnOnce(&[u8]) -> T,
) -> Result<Option<T>, Error> {
let mut decode = Some(decode);
let mut value = None;
self.read(table, key, &mut |bytes| {
let decode = decode.take().expect("read calls read_fn at most once");
value = Some(decode(bytes));
Ok(())
})?;
Ok(value)
}
}

impl<V: StorageReadView + ?Sized> StorageReadViewExt for V {}

/// A write batch that can be committed atomically.
pub trait StorageWriteBatch: Send {
/// Put multiple key-value pairs into a table.
Expand Down
18 changes: 11 additions & 7 deletions crates/storage/src/backend/in_memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,13 +66,17 @@ struct InMemoryReadView<'a> {
}

impl StorageReadView for InMemoryReadView<'_> {
fn get(&self, table: Table, key: &[u8]) -> Result<Option<Vec<u8>>, Error> {
Ok(self
.guard
.get(&table)
.expect("table exists")
.get(key)
.cloned())
fn read(
&self,
table: Table,
key: &[u8],
read_fn: &mut dyn FnMut(&[u8]) -> Result<(), Error>,
) -> Result<bool, Error> {
let Some(value) = self.guard.get(&table).expect("table exists").get(key) else {
return Ok(false);
};
read_fn(value)?;
Ok(true)
}

fn prefix_iterator(
Expand Down
145 changes: 140 additions & 5 deletions crates/storage/src/backend/rocksdb.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,8 @@ use crate::api::{
ALL_TABLES, Error, PrefixResult, StorageBackend, StorageReadView, StorageWriteBatch, Table,
};
use rocksdb::{
BlockBasedOptions, Cache, ColumnFamilyDescriptor, DBWithThreadMode, MultiThreaded, Options,
WriteBatch, WriteOptions,
BlockBasedOptions, Cache, ColumnFamilyDescriptor, DBCompressionType, DBWithThreadMode,
MultiThreaded, Options, WriteBatch, WriteOptions,
};
use std::path::Path;
use std::sync::Arc;
Expand All @@ -18,6 +18,43 @@ fn cf_name(table: Table) -> &'static str {
table.name()
}

/// The smallest value a blob-file table stores out of line.
///
/// RocksDB's default data block size: a larger value would get an oversized
/// block of its own anyway, so it gains nothing from staying inline.
const MIN_BLOB_SIZE: u64 = 4 * 1024;

/// Whether a table keeps its values in blob files rather than inline in SSTs.
///
/// `States` holds full state snapshots (100+ MB on mainnet-sized beacon
/// chains) and `StateDiffs` the deltas between them (hundreds of KB). Inline,
/// every compaction that touches a file rewrites those values, and a lookup
/// that misses still reads the data block around the key, which here can be
/// a whole snapshot. In a blob file a value is written once, and the SSTs
/// hold only small references to it.
fn stores_values_in_blob_files(table: Table) -> bool {
matches!(table, Table::States | Table::StateDiffs)
}

/// Moves a column family's large values into blob files.
///
/// RocksDB applies the change to an existing database as it goes: new
/// writes land in blob files, and inline values move out as compaction
/// rewrites their SSTs. So a data directory written without it opens
/// unchanged, and no `DB_VERSION` bump is needed.
fn enable_blob_files(cf_opts: &mut Options) {
cf_opts.set_enable_blob_files(true);
cf_opts.set_min_blob_size(MIN_BLOB_SIZE);
// SST blocks get RocksDB's default Snappy, but blob files default to no
// compression, so moving the values out would otherwise grow the tables
// on disk. The values are raw SSZ.
cf_opts.set_blob_compression_type(DBCompressionType::Lz4);
// No blob garbage collection: nothing deletes or overwrites a state, so it
// would only relocate live blobs during compaction. Revisit if states are
// ever pruned. No blob cache either: the store caches decoded states
// itself, and a snapshot-sized entry would evict the whole block cache.
}

/// RocksDB storage backend.
#[derive(Clone)]
pub struct RocksDBBackend {
Expand Down Expand Up @@ -52,6 +89,9 @@ impl RocksDBBackend {
.map(|t| {
let mut cf_opts = Options::default();
cf_opts.set_block_based_table_factory(&block_opts);
if stores_values_in_blob_files(*t) {
enable_blob_files(&mut cf_opts);
}
ColumnFamilyDescriptor::new(cf_name(*t), cf_opts)
})
.collect();
Expand Down Expand Up @@ -93,7 +133,15 @@ impl StorageBackend for RocksDBBackend {
.ok()
.flatten()
.unwrap_or(0);
sst_bytes + memtable_bytes
// `estimate-live-data-size` counts SST files only, so a blob-file
// table's values would otherwise vanish from the estimate.
let blob_bytes = self
.db
.property_int_value_cf(&cf, "rocksdb.live-blob-file-size")
.ok()
.flatten()
.unwrap_or(0);
sst_bytes + memtable_bytes + blob_bytes
}
}

Expand All @@ -103,13 +151,24 @@ struct RocksDBReadView {
}

impl StorageReadView for RocksDBReadView {
fn get(&self, table: Table, key: &[u8]) -> Result<Option<Vec<u8>>, Error> {
fn read(
&self,
table: Table,
key: &[u8],
read_fn: &mut dyn FnMut(&[u8]) -> Result<(), Error>,
) -> Result<bool, Error> {
let cf = self
.db
.cf_handle(cf_name(table))
.ok_or_else(|| format!("Column family {} not found", cf_name(table)))?;

Ok(self.db.get_cf(&cf, key)?)
// Pinned: references RocksDB's own buffer instead of copying the
// value into a `Vec`.
let Some(value) = self.db.get_pinned_cf(&cf, key)? else {
return Ok(false);
};
read_fn(&value)?;
Ok(true)
}

fn prefix_iterator(
Expand Down Expand Up @@ -200,6 +259,82 @@ mod tests {
run_backend_tests(&backend);
}

/// A value big enough for a blob file, and the property counting them.
const BLOB_VALUE_LEN: usize = 64 * 1024;
const NUM_BLOB_FILES: &str = "rocksdb.num-blob-files";

fn num_blob_files(backend: &RocksDBBackend, table: Table) -> u64 {
let cf = backend.db.cf_handle(cf_name(table)).unwrap();
backend
.db
.property_int_value_cf(&cf, NUM_BLOB_FILES)
.unwrap()
.unwrap()
}

fn put_and_flush(backend: &RocksDBBackend, table: Table, key: &[u8], value: Vec<u8>) {
let mut batch = backend.begin_write().unwrap();
batch.put_batch(table, vec![(key.to_vec(), value)]).unwrap();
batch.commit().unwrap();
let cf = backend.db.cf_handle(cf_name(table)).unwrap();
backend.db.flush_cf(&cf).unwrap();
}

#[test]
fn large_state_values_live_in_blob_files() {
let dir = tempdir().unwrap();
let backend = RocksDBBackend::open(dir.path()).unwrap();
let value: Vec<u8> = (0..BLOB_VALUE_LEN).map(|i| (i % 251) as u8).collect();

for table in [Table::States, Table::StateDiffs] {
put_and_flush(&backend, table, b"big", value.clone());
assert_eq!(num_blob_files(&backend, table), 1, "{table:?}");

let view = backend.begin_read().unwrap();
assert_eq!(view.get(table, b"big").unwrap(), Some(value.clone()));
assert!(view.contains(table, b"big").unwrap());
assert!(!view.contains(table, b"absent").unwrap());
}
assert!(backend.estimate_table_bytes(Table::States) > 0);

// Small values, and every other table, stay inline.
put_and_flush(&backend, Table::States, b"small", vec![7; 16]);
assert_eq!(num_blob_files(&backend, Table::States), 1);
put_and_flush(&backend, Table::BlockHeaders, b"big", value);
assert_eq!(num_blob_files(&backend, Table::BlockHeaders), 0);
}

#[test]
fn a_directory_written_without_blob_files_still_reads() {
let dir = tempdir().unwrap();
let value: Vec<u8> = (0..BLOB_VALUE_LEN).map(|i| (i % 251) as u8).collect();

// The layout a data directory had before blob files: every table
// with plain options.
{
let mut opts = Options::default();
opts.create_if_missing(true);
opts.create_missing_column_families(true);
let cfs = ALL_TABLES.iter().map(|t| cf_name(*t));
let db = DBWithThreadMode::<MultiThreaded>::open_cf(&opts, dir.path(), cfs).unwrap();
let cf = db.cf_handle(cf_name(Table::States)).unwrap();
db.put_cf(&cf, b"old", &value).unwrap();
db.flush_cf(&cf).unwrap();
}

let backend = RocksDBBackend::open(dir.path()).unwrap();
assert_eq!(num_blob_files(&backend, Table::States), 0);
put_and_flush(&backend, Table::States, b"new", value.clone());
assert_eq!(num_blob_files(&backend, Table::States), 1);

let view = backend.begin_read().unwrap();
assert_eq!(
view.get(Table::States, b"old").unwrap(),
Some(value.clone())
);
assert_eq!(view.get(Table::States, b"new").unwrap(), Some(value));
}

#[test]
fn test_persistence() {
let dir = tempdir().unwrap();
Expand Down
68 changes: 68 additions & 0 deletions crates/storage/src/backend/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ pub fn run_backend_tests(backend: &dyn StorageBackend) {
test_delete(backend);
test_prefix_iterator(backend);
test_nonexistent_key(backend);
test_contains(backend);
test_read(backend);
test_delete_then_put(backend);
test_put_then_delete(backend);
test_delete_range(backend);
Expand Down Expand Up @@ -45,6 +47,72 @@ fn test_put_and_get(backend: &dyn StorageBackend) {
}
}

fn test_contains(backend: &dyn StorageBackend) {
{
let mut batch = backend.begin_write().unwrap();
batch
.put_batch(
Table::BlockHeaders,
vec![(b"test_contains_key".to_vec(), b"value1".to_vec())],
)
.unwrap();
batch.commit().unwrap();
}
let view = backend.begin_read().unwrap();
assert!(
view.contains(Table::BlockHeaders, b"test_contains_key")
.unwrap()
);
assert!(
!view
.contains(Table::BlockHeaders, b"test_contains_missing")
.unwrap()
);
}

fn test_read(backend: &dyn StorageBackend) {
{
let mut batch = backend.begin_write().unwrap();
batch
.put_batch(
Table::BlockHeaders,
vec![(b"test_read_key".to_vec(), b"value1".to_vec())],
)
.unwrap();
batch.commit().unwrap();
}
let view = backend.begin_read().unwrap();

let mut seen = Vec::new();
let found = view
.read(Table::BlockHeaders, b"test_read_key", &mut |bytes| {
seen.push(bytes.to_vec());
Ok(())
})
.unwrap();
assert!(found);
assert_eq!(seen, vec![b"value1".to_vec()]);

// A missing key never reaches the callback.
let mut calls = 0;
let found = view
.read(Table::BlockHeaders, b"test_read_missing", &mut |_| {
calls += 1;
Ok(())
})
.unwrap();
assert!(!found);
assert_eq!(calls, 0);

// The callback's own error is what `read` returns.
let err = view
.read(Table::BlockHeaders, b"test_read_key", &mut |_| {
Err("decode failed".into())
})
.unwrap_err();
assert_eq!(err.to_string(), "decode failed");
}

fn test_delete(backend: &dyn StorageBackend) {
// Write data
{
Expand Down
4 changes: 3 additions & 1 deletion crates/storage/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,9 @@ mod state_diff;
mod state_writer;
mod store;

pub use api::{ALL_TABLES, StorageBackend, StorageReadView, StorageWriteBatch, Table};
pub use api::{
ALL_TABLES, StorageBackend, StorageReadView, StorageReadViewExt, StorageWriteBatch, Table,
};
pub use committee_cache::{CommitteeCache, Lookup, ShufflingKey};
/// Error type returned by the fallible [`Store`] operations, exported so
/// callers can match on it (e.g. to distinguish [`Error::DbVersionMismatch`]).
Expand Down
Loading
Loading