From 5529ae9593a3d10c8e602c56233be27eb1f9fcfa Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Fri, 17 Apr 2026 18:39:34 +0400 Subject: [PATCH 1/6] feat: add optimized accountsdb crate --- Cargo.toml | 11 +- README.md | 352 +++++++++++++++++++ accountsdb/Cargo.toml | 42 +++ accountsdb/README.md | 53 +++ accountsdb/src/lib.rs | 286 ++++++++++++++++ accountsdb/src/metrics.rs | 199 +++++++++++ accountsdb/src/snapshot.rs | 116 +++++++ accountsdb/src/store/defrag.rs | 136 ++++++++ accountsdb/src/store/index.rs | 207 +++++++++++ accountsdb/src/store/kv.rs | 123 +++++++ accountsdb/src/store/mmap.rs | 303 ++++++++++++++++ accountsdb/src/store/mod.rs | 301 ++++++++++++++++ accountsdb/src/tests.rs | 589 ++++++++++++++++++++++++++++++++ accountsdb/src/volatile.rs | 129 +++++++ solana/account/src/cow/tests.rs | 83 ----- 15 files changed, 2846 insertions(+), 84 deletions(-) create mode 100644 README.md create mode 100644 accountsdb/Cargo.toml create mode 100644 accountsdb/README.md create mode 100644 accountsdb/src/lib.rs create mode 100644 accountsdb/src/metrics.rs create mode 100644 accountsdb/src/snapshot.rs create mode 100644 accountsdb/src/store/defrag.rs create mode 100644 accountsdb/src/store/index.rs create mode 100644 accountsdb/src/store/kv.rs create mode 100644 accountsdb/src/store/mmap.rs create mode 100644 accountsdb/src/store/mod.rs create mode 100644 accountsdb/src/tests.rs create mode 100644 accountsdb/src/volatile.rs delete mode 100644 solana/account/src/cow/tests.rs diff --git a/Cargo.toml b/Cargo.toml index 18be369a..6b8f34c7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,6 @@ [workspace] members = [ + "accountsdb", "nucleus", "programs/magic-root-interface", "programs/v42-calculator-interface", @@ -21,6 +22,7 @@ rust-version = "1.94.1" version = "0.1.0" [workspace.dependencies] +accountsdb = { path = "accountsdb" } magic-root-interface = { path = "programs/magic-root-interface" } magic-root-program = { path = "programs/magic-root-program" } nucleus = { path = "nucleus", package = "magicblock-engine-nucleus" } @@ -35,16 +37,22 @@ arc-swap = "1.9.1" assert_matches = "1.5.0" base64 = "0.22.1" bincode = "1.3.3" +bitcode = "0.6.9" bitflags = "2.11.1" blake3 = "1.8.5" +bytemuck = "1.25" cfg-if = "1.0.4" +clonetree = "0.0.2" criterion = "0.8.2" derive_more = "2.1.1" env_logger = "0.11.8" futures = { version = "0.3.32", default-features = false } heed = { version = "0.22.1", default-features = false } itertools = "0.13.0" +memmap2 = "0.9.10" +num_cpus = "1.17.0" oneshot = "0.2.1" +parking_lot = "0.12.5" prometheus = { version = "0.14.0", default-features = false } qualifier_attr = "0.2.2" rand = "0.9.2" @@ -60,7 +68,8 @@ thiserror = "2.0.17" tokio = "1.52.1" tokio-util = "0.7.18" tracing = "0.1.44" -tracing-subscriber = { version = "0.3.23", features = ["env-filter", "fmt"] } +tracing-subscriber = "0.3.23" +twox-hash = { version = "2.1.2", default-features = false } wincode = "0.5.1" zstd = { version = "0.13.3", default-features = false } diff --git a/README.md b/README.md new file mode 100644 index 00000000..01c073d3 --- /dev/null +++ b/README.md @@ -0,0 +1,352 @@ +

MagicBlock Engine

+ +

+ Execution engine for ephemeral rollups — Solana transactions over durable, locally-owned state. +

+ +

+ License Apache-2.0 + Rust 1.96.1 + Edition 2024 + Solana SVM + Status experimental + Version 0.1.0 +

+ +--- + +MagicBlock Engine executes Solana transactions for ephemeral rollups. It owns +account state, records transaction and block history, and exposes asynchronous +APIs for execution, simulation, reads, and subscriptions. + +## ✨ Highlights + +| | | | +| :-- | :-- | :-- | +| ⚙️ **Runs Solana programs** — a real SVM, without the overhead of a validator | 🗃️ **Storage that fits the account** — engine-owned state on disk, chain-mirrored state in memory | 📚 **Retained history** — transactions and blocks kept in segments you can retain or drop wholesale | +| 🔁 **Replication** — mirror a live engine onto standby nodes over TCP | 🩹 **Recoverable startup** — restores snapshots and verifies replayed history after a crash | 📡 **Async APIs** — execute, simulate, read, and subscribe over live state | + +## 📖 Contents + +- [🚀 Starting the engine](#-starting-the-engine) +- [🛑 Shutdown](#-shutdown) +- [🔁 Replication](#-replication) +- [📦 Account state](#-account-state) +- [📨 Transactions](#-transactions) +- [📡 Subscriptions](#-subscriptions) +- [🩹 Startup and recovery](#-startup-and-recovery) +- [🧩 Workspace layout](#-workspace-layout) + +--- + +## 🚀 Starting the engine + +Bringing up an engine is mostly filling in one struct and awaiting one call — +everything underneath (storage, ledger, scheduler, background tasks) is wired up +for you. + +The embedding service must retain both the engine and its `ShutdownManager`. +The manager coordinates every background service started by `Engine::new`. + +```rust +use std::{num::NonZeroU64, path::PathBuf, time::Duration}; + +use engine::Engine; +use keeper::builder::KeeperBuilder; +use nucleus::{ + config::{AccountsDBParams, BlockstoreParams, LedgerParams}, + shutdown::ShutdownManager, +}; +use solana_keypair::Keypair; +use solana_sysvar::rent::Rent; + +async fn open_engine( + home: PathBuf, +) -> engine::Result<(Engine, ShutdownManager)> { + let mut shutdown = ShutdownManager::default(); + let builder = KeeperBuilder { + authority: Keypair::new().into(), + accountsdb: AccountsDBParams { + directory: home.join("accountsdb"), + lru_capacity: 10_000, + }, + ledger: LedgerParams { + directory: home.join("ledger"), + size_limit: 256 * 1024 * 1024 * 1024, + }, + blockstore: BlockstoreParams { + blocktime: Duration::from_millis(400), + superblock: NonZeroU64::new(16).unwrap(), + }, + builtins: Default::default(), + programs: Default::default(), + accounts: Default::default(), + rent: Rent::default(), + }; + + let engine = Engine::new(builder, None, &mut shutdown).await?; + Ok((engine, shutdown)) +} +``` + +The second argument chooses who advances blocks. `None` runs the built-in +pacer, which produces blocks on its own clock — the standalone case. Passing a +channel instead makes block boundaries caller-driven, as replication followers +do when they step in time with a leader. External producers supply the slot and +timestamp; the sequencer computes and overwrites the block hash and parent. + +The two modes also start differently: the built-in pacer wipes chain-mirrored +volatile accounts at startup (internal system accounts stay available), so a +standalone engine begins from clean external state. An external pacer keeps +whatever volatile state was restored, which replication depends on. + +--- + +## 🛑 Shutdown + +Shutdown isn't a hard stop — it unwinds in tiers, so in-flight work drains and +durable state lands on disk before the process goes away. + +The host waits for an OS signal or premature service termination with +`ShutdownManager::wait`. It should then stop external ingress and call +`ShutdownManager::terminate` while retaining the engine handle. + +```rust +let cause = shutdown.wait().await; + +// Stop accepting transactions and other external work here. +shutdown.terminate().await; +``` + +`wait` returns whether shutdown was requested by an OS signal or by a managed +service terminating early. Embedding processes can use the service reason to +distinguish recoverable lifecycle events, such as a replication snapshot that +requires reopening the engine, from fatal failures. + +Shutdown proceeds by service tier: + +1. A replication client stops consuming upstream state. +2. The pacemaker stops producing boundaries and calls `Engine::shutdown`. +3. The already-drained sequencer and terminally-synced ledger appender stop. +4. Ledger readers, simulation, subscriptions, and other backing services stop. + +Internal pacing publishes a final block and flushes durable state. External +pacing also writes volatile state to `CURRENT/volatile.db` after flushing the +corresponding ledger cursor. The final sync explicitly closes ledger workers, +so retained but inactive engine handles cannot hold shutdown open. Each tier +has a bounded termination window. + +--- + +## 🔁 Replication + +Point a follower at a leader and it keeps itself in sync — replaying the stream +when it can, and pulling a fresh snapshot when it has fallen too far behind. + +Replication keeps a standby engine in step with a live one: a **leader** serves +its history over TCP, and one or more **followers** replay that stream to stay +current. On the leader machine, bind a dispatcher to a reachable address and +serve the retained ledger: + +```rust +use std::sync::Arc; + +use replicator::ReplicationDispatcher; + +let allowed = Arc::from([follower_identity]); +ReplicationDispatcher::spawn(bind_addr, engine.clone(), allowed, &mut shutdown).await?; +``` + +On the follower machine, open its engine with an external pacer and connect the +client to the leader's address: + +```rust +use replicator::ReplicationClient; +use tokio::sync::mpsc; + +let (block_tx, block_rx) = mpsc::channel(16); +builder.authority.remote = Some(leader_identity); +let engine = Engine::new(builder, Some(block_rx), &mut shutdown).await?; +ReplicationClient::spawn(leader_addr, engine.clone(), block_tx, &mut shutdown)?; +``` + +Leader and follower local keypairs do not need to match. The server allowlist +contains follower local identities and denies all access when empty. The +follower's remote authority identifies its immediate upstream, whose signed +responses must arrive within 30 seconds of the follower's clock. + +The external pacer keeps replicated blocks ordered with transactions, resets, +and seals. If the leader's retained stream cannot satisfy the follower's cursor, +it sends the newest snapshot. The client stages it, reports `RestartRequired` +through the follower's shutdown manager, and the follower host reopens its +engine from the same directories. + +--- + +## 📦 Account state + +You never have to decide where an account lives — the engine watches what each +account *is* and keeps it in the right place on its own. + +The engine holds two kinds of accounts and stores each where it makes sense: + +- Accounts the engine controls — delegated, ephemeral, and transient — are + authoritative here and **persisted to disk**. +- Accounts that only mirror external chain or system state — read-only, + placeholders, and sysvars — are kept **in volatile memory**. + +An account's `AccountMode::authoritative()` classification decides which side it +belongs to. When that changes, accountsdb moves the account and drops the stale +copy from the other backend, so there is only ever one live copy. `Transient` +accounts remain authoritative and persisted even though runtime code cannot +mutate them. + +To change accounts directly, use `Engine::account(pubkey)`. `create`, `update`, +`patch`, and `delete` each run as one signed, committed transaction and require +the local signer to match the engine authority. + +```rust +use solana_account::{AccountBuilder, AccountFieldPatch, AccountMode}; +use solana_pubkey::Pubkey; + +let key = Pubkey::new_unique(); +let owner = Pubkey::new_unique(); +let account = AccountBuilder::default() + .lamports(2_000_000) + .owner(owner) + .mode(AccountMode::Delegated) + .data(vec![1, 2, 3, 4]) + .build(); + +engine.account(key).create(account, None).await?; +let current = engine.accounts().loader().load(&key)?; + +engine + .account(key) + .patch(vec![AccountFieldPatch::DataAt { + offset: 0, + data: vec![9; 4], + }]) + .await?; + +let replacement = AccountBuilder::default() + .lamports(2_000_000) + .owner(owner) + .mode(AccountMode::Delegated) + .data(vec![5; 4]) + .build(); +engine.account(key).update(replacement).await?; +engine.account(key).delete().await?; +``` + +Each mutation is one committed transaction. `create` can also run optional +post-finalize instructions in that transaction; if an instruction fails, the +creation does not commit. + +Missing external accounts can be coordinated with `Engine::accounts().ensure`. +The first caller receives `MissingAccount::Load`; concurrent callers receive a +wait handle for the same pubkey. After storing the account, the loader calls +`AccountLoad::complete(mode)` to publish success and update recency tracking for +non-authoritative accounts. Dropping the load guard instead wakes waiters with a +failed outcome. + +--- + +## 📨 Transactions + +Hand it whatever you've already got — a few instructions, a `Message`, or raw +encoded bytes — and pick how much you want to wait around for. + +`Engine::transaction` accepts an instruction slice, `Message`, sanitized +`TransactionView`, or encoded transaction bytes. Instruction slices and messages +use the effective authority as payer and the local signer with the latest +blockhash, so local composition requires those identities to match. + +```rust +use engine::Engine; +use solana_instruction::Instruction; + +async fn submit( + engine: &Engine, + instructions: &[Instruction], +) -> engine::Result<()> { + engine + .transaction(instructions)? + .execute() + .await? + .map_err(Into::into) +} +``` + +- `execute` waits for the committed transaction result. +- `schedule` queues execution without waiting for its result. +- `simulate` executes against owned account copies without committing state. + +--- + +## 📡 Subscriptions + +No polling loops — subscribe to what you care about and the engine pushes +updates as they happen. + +Keeper accessors expose Tokio broadcast receivers for live state: + +```rust +let mut account_updates = engine.accounts().subscribe(key).await; +let mut blocks = engine.blocks().subscribe(); + +let account = account_updates.recv().await?; +let block = blocks.recv().await?; +``` + +Related accessors subscribe to program-owned accounts, cache evictions, +snapshot completion, transaction status, logs, processed transactions, and +service messages. Broadcast consumers must handle `Lagged` when they fall +behind and `Closed` during shutdown; retained reads are available separately. + +--- + +## 🩹 Startup and recovery + +After an interrupted write, the next start checks local state against retained +history and restores a retained snapshot when necessary. + +Every startup reconciles the account store with the transaction history. A +crash, corruption, and a staged replication snapshot enter the same recovery +path, but recovery requires a valid retained snapshot when the current store +cannot be used. + +Concretely: keeper validates the account store against the retained ledger. A +corrupt store, or a valid one whose latest checkpoint trails the ledger, is +replaced with the newest retained snapshot. If that restored state still trails +the ledger tip, the engine replays the missing history to catch up, checking the +rebuilt state against each recorded checkpoint and refusing to continue +(`ReplayError::StateMismatch`) if they diverge. When the store is already +current, nothing runs. + +--- + +## 🧩 Workspace layout + +| Crate | Role | +| :-- | :-- | +| `nucleus` | Shared ledger, runtime, metrics, TLS, and shutdown types. | +| `solana/*` | The runtime forks required by the engine account model. | +| `accountsdb` | Owns persisted and volatile account storage and snapshots. | +| `ledger` | Stores transactions, execution records, blocks, and superblocks. | +| `keeper` | Opens both stores and provides caches, reads, and subscriptions. | +| `processor` | Schedules transactions across SVM executors and commits results. | +| `programs/*` | MagicRoot and the v42 test program and interfaces. | +| `engine` | Wires the execution engine and exposes the public handle. | +| `replicator` | Streams durable engine state between nodes. | + +Transactions are appended before execution, then paired with execution metadata. +Successful dirty accounts are written through accountsdb and live notifications +are published. Superblock boundaries quiesce execution while keeper snapshots +accountsdb and archives it beside the next retained ledger segment. + +--- + +

+ Built with 🦀 Rust · licensed under Apache-2.0 · © MagicBlock contributors +

diff --git a/accountsdb/Cargo.toml b/accountsdb/Cargo.toml new file mode 100644 index 00000000..59c77730 --- /dev/null +++ b/accountsdb/Cargo.toml @@ -0,0 +1,42 @@ +[package] +name = "magicblock-accountsdb" + +authors.workspace = true +edition.workspace = true +homepage.workspace = true +license.workspace = true +repository.workspace = true +rust-version.workspace = true +version.workspace = true + +[lib] +name = "accountsdb" + +[features] +testkit = [] + +[dependencies] +nucleus = { workspace = true, features = ["heed", "metrics"] } + +ahash = { workspace = true } +bincode = { workspace = true } +bytemuck = { workspace = true, features = ["derive", "extern_crate_std"] } +clonetree = { workspace = true } +derive_more = { workspace = true, features = ["from"] } +heed = { workspace = true } +memmap2 = { workspace = true } +parking_lot = { workspace = true } +scc = { workspace = true, features = ["serde"] } +thiserror = { workspace = true } +tracing = { workspace = true } +twox-hash = { workspace = true, features = ["alloc", "xxhash3_64"] } + +solana-account = { workspace = true, features = ["serde"] } +solana-pubkey = { workspace = true, features = ["bytemuck"] } + +[dev-dependencies] +assert_matches = { workspace = true } +nucleus = { workspace = true, features = ["testkit"] } + +[lints] +workspace = true diff --git a/accountsdb/README.md b/accountsdb/README.md new file mode 100644 index 00000000..aa9015f1 --- /dev/null +++ b/accountsdb/README.md @@ -0,0 +1,53 @@ +# `magicblock-accountsdb` + +Accountsdb routes account state between two backends according to +`AccountMode::authoritative()`: + +- `PersistedStore` is an mmap-backed account file with LMDB indexes. It holds + delegated, ephemeral, and transient accounts controlled by the engine. +- `VolatileStore` is an in-memory map for externally owned state that can be + fetched again. + +Every store operation touches the backend required by both the account's current +representation and authoritative classification. This commits borrowed images +in persistent storage, inserts owned images there, updates owned volatile +images, and removes stale copies after mode changes or closure. `Transient` +remains authoritative and runtime-immutable until its lifecycle state resolves. + +## Persisted layout + +`CURRENT/storage.db` contains a metadata header followed by account images in the +borrowed `solana-account` layout. Each image includes its full pubkey so scans can +recover keys without the index. Offsets are measured in 8-byte `StorageUnit`s. + +The LMDB index under `CURRENT/index` contains: + +- `accounts`: account key tag to storage offset and owner tag. +- `programs`: owner tag to account offsets. +- `freelist`: image size to reusable offsets. + +`PersistedProgramIter` retains its read transaction for the persisted portion of +iteration. The optional `testkit` feature uses smaller maps and growth blocks +without changing the on-disk format. + +## Writes and compaction + +A persisted batch commits its LMDB transaction once. If applying or committing +the batch fails, already committed borrowed images are rolled back so indexed +state remains authoritative. Freed image spans enter the freelist. + +Defragmentation requires exclusive access. It sorts free spans, compacts the +tail half, updates indexes before moving live bytes, and shrinks the mapped file. +Repeated passes reduce the remaining hole set geometrically. + +## Snapshots and volatile state + +`AccountsDB::snapshot` requires exclusive write access. It records the +superblock id, defragments and flushes persisted state, clones the active tree, +and serializes the current volatile map into the clone's `volatile.db`. + +`dump(None)` writes `CURRENT/volatile.db` for a clean externally paced shutdown. +The next open restores that file into memory and removes it. `reset` instead +removes chain-mirrored volatile accounts while preserving internal system +accounts and rebuilding their owner indexes. Persisted engine-authoritative +state is never reset. diff --git a/accountsdb/src/lib.rs b/accountsdb/src/lib.rs new file mode 100644 index 00000000..010d4cb6 --- /dev/null +++ b/accountsdb/src/lib.rs @@ -0,0 +1,286 @@ +#![doc = include_str!("../README.md")] + +use std::{ + cell::RefCell, + collections::BTreeSet, + path::{Path, PathBuf}, + sync::atomic::Ordering::*, +}; + +use derive_more::From; +use nucleus::Slot; +use nucleus::heed::RoTxnTls; +use solana_account::{AccountSharedData, CoWAccount}; +use solana_pubkey::Pubkey; +use tracing::{info, warn}; + +use crate::{ + store::{DatabaseVersion, PersistedProgramIter, PersistedStore}, + volatile::VolatileStore, +}; + +pub use snapshot::{BackupOp, SnapshotError, SnapshotResult}; +pub use store::mmap::STORAGE_FILE; + +mod metrics; +mod snapshot; +mod store; +mod volatile; + +#[cfg(test)] +mod tests; + +/// Active database subdirectory. +const ACTIVE_DIR: &str = "CURRENT"; + +/// Top-level account store backed by persisted and volatile backends. +pub struct AccountsDB { + /// On-disk store for engine-authoritative account modes. + persisted: PersistedStore, + /// Rebuildable in-memory store for non-authoritative account modes. + volatile: VolatileStore, + /// Database root directory. + root: PathBuf, +} + +impl AccountsDB { + /// Opens or creates the database at `root`. + pub fn new(root: impl AsRef) -> Result { + let root = root.as_ref().to_owned(); + let path = Self::directory(&root); + let persisted = PersistedStore::new(&path)?; + let volatile = VolatileStore::new(&path)?; + info!(?path, "opened accountsdb"); + let db = Self { persisted, volatile, root }; + metrics::init(&db); + Ok(db) + } + + /// Returns the active database directory under `root`. + pub fn directory(root: &Path) -> PathBuf { + root.join(ACTIVE_DIR) + } + + /// Stores accounts in the backend that matches their current form. + /// + /// Persistent modes are kept in persisted storage. Other modes are kept in + /// volatile storage. Each batch also touches the opposite backend so stale + /// copies are removed after mode changes. Persisted failures roll back + /// borrowed images before the caller sees the error. + pub fn store<'a, AC>(&self, accounts: AC) -> Result<()> + where + AC: IntoIterator + Clone, + ::IntoIter: Clone, + { + let iter = accounts.clone().into_iter().filter(persisted); + self.persisted.upsert(iter)?; + + let iter = accounts.into_iter().filter(volatile); + self.volatile.upsert(iter); + + Ok(()) + } + + /// Creates a loader that reuses a read transaction for persisted lookups. + pub fn loader(&self) -> AccountLoader<'_> { + AccountLoader::new(self) + } + + /// Iterates program-owned accounts across both backends. + pub fn program(&self, owner: &Pubkey) -> Result> { + let persisted = self.persisted.program(*owner)?; + let volatile = self.volatile.program(owner); + Ok(ProgramIter { persisted, volatile, db: self }) + } + + /// Returns the latest slot persisted in the database metadata. + pub fn slot(&self) -> Slot { + self.persisted.meta().slot.load(Acquire) + } + + /// Sets the database slot and flushes dirty pages asynchronously. + pub fn set_slot(&self, slot: Slot) -> Result<()> { + self.persisted.meta().slot.store(slot, Release); + self.flush(false) + } + + /// Returns the id of the last sealed superblock recorded in the database metadata. + pub fn superblock(&self) -> Slot { + self.persisted.meta().superblock.load(Acquire) + } + + /// Records the last sealed superblock id. Set on snapshot, and on replay + /// before recomputing the checksum to compare against a seal. + pub fn set_superblock(&self, superblock: u64) { + self.persisted.meta().superblock.store(superblock, Release); + } + + /// Flushes persisted account storage, forcing synchronous durability when requested. + pub fn flush(&self, force: bool) -> Result<()> { + self.persisted.flush(force).map_err(Into::into) + } + + /// Validates the persisted store checksum and on-disk format version. + pub fn validate(&self) -> Result<()> { + self.persisted.validate() + } + + /// Returns the last checksum published on superblock boundary. + pub fn checksum(&self) -> u64 { + self.persisted.meta().checksum.load(Acquire) + } + + /// Drops chain-mirrored volatile state while retaining system accounts; + /// persisted state is left untouched. + /// + /// Chain-owned accounts can be fetched again when synchronization resumes. + /// System accounts hold internal runtime state and survive the reset; their + /// volatile owner indexes are rebuilt. Persisted, engine-authoritative state + /// is never reset. + pub fn reset(&self) { + self.volatile.reset(); + } +} + +/// Loader that caches a read transaction for persisted account lookups. +pub struct AccountLoader<'a> { + /// Cached read transaction for the persisted index. + txn: RefCell>>, + /// Database handle used for volatile and persisted lookups. + db: &'a AccountsDB, +} + +impl<'a> AccountLoader<'a> { + /// Creates a new loader bound to `db`. + pub fn new(db: &'a AccountsDB) -> Self { + Self { txn: Default::default(), db } + } + + /// Loads one account, checking both backends in turn + pub fn load(&self, pubkey: &Pubkey) -> Result> { + let txn = &mut self.txn.borrow_mut(); + if let Some(acc) = self.db.persisted.load(txn, pubkey)? { + metrics::load(StoreKind::Persisted); + return Ok(Some(acc.into())); + } + let account = self.db.volatile.load(pubkey).map(Into::into); + if account.is_some() { + metrics::load(StoreKind::Volatile); + } else { + metrics::load(StoreKind::Absent); + } + Ok(account) + } + + /// Returns whether an account exists in either backend. + pub fn contains(&self, pubkey: &Pubkey) -> Result { + let txn = &mut self.txn.borrow_mut(); + if self.db.persisted.contains(txn, pubkey)? { + return Ok(true); + } + let contains = self.db.volatile.contains(pubkey); + Ok(contains) + } +} + +/// Iterates program-owned accounts across both backends. +pub struct ProgramIter<'a> { + /// Persisted program accounts. + persisted: Option>, + /// Volatile program pubkeys. + volatile: BTreeSet, + /// Database handle used to resolve volatile accounts. + db: &'a AccountsDB, +} + +impl<'a> Iterator for ProgramIter<'a> { + type Item = AccountEntry; + /// Yields authoritative accounts first, then volatile ones. + fn next(&mut self) -> Option { + if let Some(persisted) = &mut self.persisted { + // Yield authoritative entries first. + if let Some(item) = persisted.next() { + return Some(item); + } + } + // Release the persisted read txn before draining volatile entries. + let _ = self.persisted.take(); + // Then drain the in-memory set of non-authoritative accounts. + while let Some(pubkey) = self.volatile.pop_first() { + if let Some(account) = self.db.volatile.load(&pubkey) { + return Some((pubkey, account.into())); + } + warn!(%pubkey, "volatile program set references a missing account; skipping"); + } + None + } +} + +/// Errors returned by accountsdb. +#[derive(Debug, thiserror::Error, From)] +pub enum AccountsDBError { + /// LMDB key-value codec error. + #[error("LMDB key/value codec error: {0}")] + Codec(#[source] heed::BoxedError), + /// Filesystem error. + #[error("filesystem I/O error: {0}")] + IO(#[source] std::io::Error), + /// LMDB index access error. + #[error("LMDB index error: {0}")] + Index(#[source] heed::Error), + /// Storage allocation would exceed the maximum mapped size. + #[error("mapped storage exceeded the 32 GiB limit")] + Allocation, + /// Opened database version is not supported by current implementation. + #[error("unsupported database version: {0:?}")] + UnsupportedVersion(DatabaseVersion), + /// Database was corrupted during the shutdown/crash. + #[error("database integrity check failed")] + Corruption, + /// Volatile snapshot serialization error. + #[error("volatile snapshot serialization error: {0}")] + Serde(#[source] Box), +} + +/// Result type used by the accountsdb crate. +type Result = std::result::Result; +/// Account key plus shared account payload. +pub type AccountEntry = (Pubkey, AccountSharedData); + +/// Classification used by accountsdb metrics. +#[derive(Clone, Copy)] +pub(crate) enum StoreKind { + /// Mmap-backed persisted storage. + Persisted, + /// In-memory volatile storage. + Volatile, + /// Account was absent from both storage backends. + Absent, +} + +impl StoreKind { + /// Returns the Prometheus label value for this classification. + pub(crate) fn label(self) -> &'static str { + match self { + StoreKind::Persisted => "persisted", + StoreKind::Volatile => "volatile", + StoreKind::Absent => "absent", + } + } +} + +/// Returns `true` for entries that must touch persisted storage. +fn persisted(entry: &&AccountEntry) -> bool { + match entry.1.cow() { + CoWAccount::Borrowed(_) => true, + CoWAccount::Owned(_) => entry.1.mode().authoritative(), + } +} + +/// Returns `true` for entries that must touch volatile storage. +fn volatile(entry: &&AccountEntry) -> bool { + match entry.1.cow() { + CoWAccount::Borrowed(_) => !entry.1.mode().authoritative(), + CoWAccount::Owned(_) => true, + } +} diff --git a/accountsdb/src/metrics.rs b/accountsdb/src/metrics.rs new file mode 100644 index 00000000..c1255612 --- /dev/null +++ b/accountsdb/src/metrics.rs @@ -0,0 +1,199 @@ +//! Prometheus metrics for accountsdb. + +use std::sync::{OnceLock, atomic::Ordering::*}; + +use nucleus::metrics as metric; +use nucleus::metrics::{IntCounter, IntGaugeVec, MetricOperation, MetricSpec, OperationCounters}; + +use crate::{AccountsDB, StoreKind, store::Stats}; + +/// Process-wide accountsdb metrics registered in the default Prometheus registry. +static METRICS: OnceLock = OnceLock::new(); + +/// Persisted account image load counter. +const READS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_reads", + help: "Persisted account image loads.", +}; +/// Borrowed account commit counter. +const COMMITS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_commits", + help: "Borrowed account commits into persisted storage.", +}; +/// Fresh mapped-storage allocation counter. +const ALLOCS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_allocs", + help: "Fresh allocations from the mapped persisted storage file.", +}; +/// Persisted freelist reuse counter. +const REALLOCS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_reallocs", + help: "Allocations reused from the persisted freelist.", +}; +/// Defragmentation relocation counter. +const COMPACTIONS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_compactions", + help: "Persisted account relocations during defragmentation.", +}; +/// Persisted account removal counter. +const REMOVALS: MetricSpec = MetricSpec { + name: "accountsdb_persisted_removals", + help: "Persisted account removals.", +}; + +/// Persisted storage resize counter. +const RESIZES: MetricSpec = MetricSpec { + name: "accountsdb_persisted_resizes", + help: "Persisted storage file resizes.", +}; +/// Account load counter grouped by source or absence. +const LOADS: MetricSpec = MetricSpec { + name: "accountsdb_loads", + help: "Account loads by source or absence.", +}; +/// Operation latency histogram recorded in microseconds. +const OPERATION_TIME: MetricSpec = MetricSpec { + name: "accountsdb_operation_duration_micros", + help: "Accountsdb operation duration distribution in microseconds.", +}; +/// Account count gauge grouped by backend store. +const ACCOUNTS: MetricSpec = MetricSpec { + name: "accountsdb_accounts", + help: "Current accountsdb account count by backend store.", +}; + +/// Label used to separate persisted and volatile account counts. +const STORE_LABEL: &str = "store"; + +/// Accountsdb operation used as a low-cardinality operation label. +#[derive(Clone, Copy)] +pub(crate) enum Operation { + /// Persisted store flush path. + Flush, + /// Persisted checksum path. + Checksum, + /// Accountsdb snapshot path. + Snapshot, + /// Volatile-state dump path. + Dump, + /// Persisted store defragmentation path. + Defragmentation, +} + +impl MetricOperation for Operation { + /// Returns the Prometheus label value for this operation. + fn label(self) -> &'static str { + match self { + Operation::Flush => "flush", + Operation::Checksum => "checksum", + Operation::Snapshot => "snapshot", + Operation::Dump => "dump", + Operation::Defragmentation => "defragmentation", + } + } +} + +/// Registers accountsdb metrics once, seeding durable counters from persisted stats. +pub(crate) fn init(db: &AccountsDB) { + METRICS.get_or_init(|| Metrics::new(db.persisted.storage.stats())); +} + +/// Records one persisted account image load. +pub(crate) fn read() { + metric::with_metrics(&METRICS, |m| m.reads.inc()); +} + +/// Records one borrowed account commit into persisted storage. +pub(crate) fn commit() { + metric::with_metrics(&METRICS, |m| m.commits.inc()); +} + +/// Records one fresh allocation from the mapped persisted storage file. +pub(crate) fn alloc() { + metric::with_metrics(&METRICS, |m| m.allocs.inc()); +} + +/// Records one allocation reuse from the persisted freelist. +pub(crate) fn realloc() { + metric::with_metrics(&METRICS, |m| m.reallocs.inc()); +} + +/// Records one persisted account relocation during defragmentation. +pub(crate) fn compaction() { + metric::with_metrics(&METRICS, |m| m.compactions.inc()); +} + +/// Records one persisted account removal. +pub(crate) fn removal() { + metric::with_metrics(&METRICS, |m| m.removals.inc()); +} + +/// Records one persisted storage file resize. +pub(crate) fn resize() { + metric::with_metrics(&METRICS, |m| m.resizes.inc()); +} + +/// Refreshes the current account count for `store`. +pub(crate) fn accounts(store: StoreKind, count: u64) { + metric::with_metrics(&METRICS, |m| { + m.accounts.with_label_values(&[store.label()]).set(metric::gauge_value(count)); + }); +} + +/// Starts an operation timer that records latency when the returned guard drops. +pub(crate) fn time(op: Operation) -> metric::OperationTimer<'static> { + op.time(METRICS.get().map(|m| &m.operations)) +} + +/// Records one account load satisfied by `store`. +pub(crate) fn load(store: StoreKind) { + metric::with_metrics(&METRICS, |m| m.loads[store as usize].inc()); +} + +/// Owns all Prometheus collectors registered by accountsdb. +struct Metrics { + /// Durable persisted account image load counter. + reads: IntCounter, + /// Durable borrowed account commit counter. + commits: IntCounter, + /// Durable fresh allocation counter. + allocs: IntCounter, + /// Durable freelist reuse counter. + reallocs: IntCounter, + /// Durable defragmentation relocation counter. + compactions: IntCounter, + /// Durable persisted account removal counter. + removals: IntCounter, + /// Durable persisted storage resize counter. + resizes: IntCounter, + /// Per-`StoreKind` load counters pre-resolved from `loads_vec`. + loads: [IntCounter; 3], + /// Runtime operation duration and completion counters. + operations: OperationCounters, + /// Runtime account count gauge labeled by backend store. + accounts: IntGaugeVec, +} + +impl Metrics { + /// Builds collectors and seeds durable counters from persisted mmap stats. + fn new(stats: &Stats) -> Self { + let loads_vec = metric::counter_vec(LOADS, &[STORE_LABEL]); + let loads = [ + loads_vec.with_label_values(&[StoreKind::Persisted.label()]), + loads_vec.with_label_values(&[StoreKind::Volatile.label()]), + loads_vec.with_label_values(&[StoreKind::Absent.label()]), + ]; + Self { + reads: metric::counter(READS, stats.reads.load(Relaxed)), + commits: metric::counter(COMMITS, stats.commits.load(Relaxed)), + allocs: metric::counter(ALLOCS, stats.allocs.load(Relaxed)), + reallocs: metric::counter(REALLOCS, stats.reallocs.load(Relaxed)), + compactions: metric::counter(COMPACTIONS, stats.compactions.load(Relaxed)), + removals: metric::counter(REMOVALS, stats.removals.load(Relaxed)), + resizes: metric::counter(RESIZES, stats.resizes.load(Relaxed)), + loads, + operations: OperationCounters::new(OPERATION_TIME), + accounts: metric::gauge_vec(ACCOUNTS, &[STORE_LABEL]), + } + } +} diff --git a/accountsdb/src/snapshot.rs b/accountsdb/src/snapshot.rs new file mode 100644 index 00000000..7bc10fc4 --- /dev/null +++ b/accountsdb/src/snapshot.rs @@ -0,0 +1,116 @@ +//! Snapshot export helpers. + +use std::{ + fs::{self, File}, + io::{self, BufWriter}, + path::PathBuf, +}; + +use nucleus::MB; +use tracing::info; + +use crate::{ + ACTIVE_DIR, AccountsDB, + metrics::{self, Operation}, +}; + +/// Snapshot directory prefix. +const PREFIX: &str = "snapshot-"; +/// Snapshot payload filename for the volatile store. +pub(crate) const VOLATILE_DB_FILE: &str = "volatile.db"; + +/// Errors while writing a snapshot directory. +#[derive(thiserror::Error, Debug)] +pub enum SnapshotError { + /// I/O while writing the snapshot. + #[error("snapshot export I/O error")] + IO(#[from] io::Error), + /// Failed to flush the persisted store before copying the tree. + #[error("failed to flush persisted store")] + Flush(#[from] heed::Error), + /// Failed to serialize the volatile store into the snapshot. + #[error("failed to serialize volatile store")] + Serde(#[from] Box), + /// Failed to clone the active database tree into the snapshot slot. + #[error("failed to clone snapshot tree")] + FsClone(#[from] Box), + /// No archived snapshot could be restored. + #[error("no valid archived accountsdb snapshot found")] + Missing, +} + +/// Result type used by snapshot export and restore helpers. +pub type SnapshotResult = Result; + +/// Active database backup operation. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum BackupOp { + /// Move the active database tree to its backup path. + Save, + /// Move the saved backup tree back to the active database path. + Restore, +} + +impl AccountsDB { + /// Writes a superblock snapshot under `root`. + /// + /// # Safety + /// The caller must ensure exclusive write access while the snapshot is in + /// progress. The persisted backend is defragmented and flushed before the + /// active tree is cloned and the volatile store is rewritten in the clone. + /// That ordering keeps the exported state coherent only when no concurrent + /// writes can race with the export. + pub unsafe fn snapshot(&self, superblock: u64) -> SnapshotResult { + let _timer = metrics::time(Operation::Snapshot); + let src = self.root.join(ACTIVE_DIR); + let dst = self.root.join(format!("{PREFIX}{superblock:0>9}")); + self.set_superblock(superblock); + // SAFETY: snapshot owns exclusive write access, so defrag cannot race + // with concurrent mutation and can compact the persisted store first. + unsafe { self.persisted.defragment() }?; + // Persisted state must reach disk before we copy the active tree. + self.persisted.flush(true)?; + // Clone the whole active tree, then replace the volatile payload below. + clonetree::clone_tree(src, &dst, &Default::default()).map_err(Box::new)?; + self.dump(Some(&dst))?; + + Ok(dst) + } + + /// Serializes volatile accounts into `volatile.db` under `dst`. + /// + /// When `dst` is omitted, writes into the active database tree so the next + /// open restores the volatile store and consumes the file. Callers must + /// prevent concurrent account writes to obtain a coherent image. + pub fn dump(&self, dst: Option<&PathBuf>) -> SnapshotResult<()> { + let _timer = metrics::time(Operation::Dump); + let path = match dst { + Some(dst) => dst.join(VOLATILE_DB_FILE), + None => Self::directory(&self.root).join(VOLATILE_DB_FILE), + }; + let db = File::options().create(true).truncate(true).write(true).open(path)?; + let mut buffered = BufWriter::with_capacity(4 * MB, db); + bincode::serialize_into(&mut buffered, &self.volatile.accounts)?; + let db = buffered.into_inner().map_err(|e| e.into_error())?; + db.sync_data().map_err(Into::into) + } + + /// Saves or restores the active database tree and returns its destination. + /// + /// After restoring, callers must drop this instance and reopen the database: + /// its open handles still refer to the removed active tree. + pub fn backup(&self, op: BackupOp) -> SnapshotResult { + let active = self.root.join(ACTIVE_DIR); + let backup = self.root.join(format!("{ACTIVE_DIR}.bkp")); + let (from, to) = match op { + BackupOp::Save => (&active, &backup), + BackupOp::Restore => (&backup, &active), + }; + if to.exists() { + fs::remove_dir_all(to)?; + } + info!(?op, "accountsdb backup"); + fs::rename(from, to)?; + Ok(to.clone()) + } +} diff --git a/accountsdb/src/store/defrag.rs b/accountsdb/src/store/defrag.rs new file mode 100644 index 00000000..f7c593e4 --- /dev/null +++ b/accountsdb/src/store/defrag.rs @@ -0,0 +1,136 @@ +#![allow(unsafe_op_in_unsafe_fn)] + +use heed::Result; +use solana_account::BorrowedAccount; +use tracing::{debug, info}; + +use crate::{ + metrics::{self, Operation}, + store::kv::{Offset, OwnerAndOffset}, +}; + +use super::PersistedStore; + +/// Free span in the persisted image file, measured in storage units and +/// ordered by offset for the sweep. +#[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +struct Hole { + offset: Offset, + units: u32, +} + +impl Hole { + /// Builds one hole from the freelist entry tuple. + fn new((units, offset): (u32, Offset)) -> Self { + Self { offset, units } + } + + /// Returns the first offset past the hole. + fn end(&self) -> Offset { + self.offset + self.units + } +} + +impl PersistedStore { + /// Compacts the tail half of the freelist and repacks the storage. + /// + /// Caller must guarantee exclusive access while defrag runs. Holes are + /// swept in file order, contiguous holes are merged into one run, and the + /// following live span is reindexed before its bytes are slid left. Each + /// pass starts at the tail half of the sorted hole list so repeated runs + /// halve the remaining work. + /// + /// This operation is not crash-safe: interruption after reindexing can + /// leave the active tree inconsistent and require restoring a backup. + /// + /// # Safety + /// + /// No concurrent mutation may touch the persisted index or mapped + /// storage while offsets are rewritten and bytes are moved. + pub(crate) unsafe fn defragment(&self) -> Result<()> { + let _timer = metrics::time(Operation::Defragmentation); + let mut holes = { + let txn = self.index.env.read_txn()?; + self.index + .freelist + .iter(&txn)? + .map(|r| r.map(Hole::new)) + .collect::>>()? + }; + holes.sort_unstable(); + + let mut i = holes.len() / 2; + // `dst` is the next write head in storage units. + let Some(mut dst) = holes.get(i).map(|h| h.offset) else { + debug!("nothing to defragment"); + return Ok(()); + }; + // Tail is the live end; if the last hole reaches it, there is no + // trailing live span to move. + let tail = Offset(self.storage.cursor()); + while i < holes.len() { + let run = i; + let mut end = holes[i].end(); + i += 1; + + while let Some(&hole) = holes.get(i) { + if hole.offset != end { + break; + } + end = hole.end(); + i += 1; + } + + let next = holes.get(i).map(|h| h.offset).unwrap_or(tail); + // Defrag runs exclusively, so rewrite and commit the index before + // moving the live bytes that now point at the future offsets. + let shift = end - dst; + self.reindex(end, next, shift, &holes[run..i])?; + // Source and destination overlap by design during left-compaction. + dst = self.slide(end, dst, next - end); + } + let reclaimed = tail - dst; + info!(reclaimed, "defragmented persisted storage"); + self.storage.shrink(dst.0).map_err(Into::into) + } + + /// Rewrites the index for a live span before its bytes are moved. + /// + /// # Safety + /// + /// `src..end` must cover whole serialized accounts in mapped storage. + /// The caller must guarantee exclusive access while offsets are updated. + unsafe fn reindex(&self, src: Offset, end: Offset, shift: u32, run: &[Hole]) -> Result<()> { + let mut txn = self.index.env.write_txn()?; + let mut pos = src; + while pos < end { + let ptr = self.storage.at(pos); + let span = BorrowedAccount::span(ptr); + let next = pos + span; + // The bytes are still at `pos`; compute the future offset before copying. + let pubkey = BorrowedAccount::pubkey(ptr); + let account = BorrowedAccount::init(ptr); + let owner = account.owner().into(); + let data = OwnerAndOffset { owner, offset: pos - shift }; + self.index.relocate(&pubkey, pos, data, &mut txn)?; + self.storage.stats().compact(); + pos = next; + } + for hole in run { + self.index.freelist.delete_one_duplicate(&mut txn, &hole.units, &hole.offset)?; + } + txn.commit() + } + + /// Slides a live region left after the index has been committed. + /// + /// # Safety + /// + /// `src..src+len` and `dst..dst+len` must be valid mapped storage + /// ranges. The ranges may overlap. + unsafe fn slide(&self, src: Offset, dst: Offset, len: u32) -> Offset { + let src = self.storage.at(src); + src.copy_to(self.storage.at(dst), len as usize); + dst + len + } +} diff --git a/accountsdb/src/store/index.rs b/accountsdb/src/store/index.rs new file mode 100644 index 00000000..42d9449b --- /dev/null +++ b/accountsdb/src/store/index.rs @@ -0,0 +1,207 @@ +//! LMDB index for persisted accounts. +//! +//! The index maps compact pubkey tags to storage offsets and owner tags, +//! plus a freelist keyed by image size. + +use std::{fs, mem, path::Path}; + +use heed::{ + Database, DatabaseFlags, Env, EnvFlags, EnvOpenOptions, IntegerComparator, Result, RoIter, + RoTxn, RwTxn, iteration_method::MoveOnCurrentKeyDuplicates, +}; +use nucleus::heed::{DatabaseIndex, RoTxnTls}; +use solana_pubkey::Pubkey; + +use crate::store::kv::{KeyTail, Offset, OwnerAndOffset, PubkeyBytes, U32LE}; + +/// LMDB map size for the index database. +#[cfg(any(test, feature = "testkit"))] +const INDEX_MAP_SIZE: usize = nucleus::MB; +#[cfg(not(any(test, feature = "testkit")))] +const INDEX_MAP_SIZE: usize = nucleus::GB; +/// Subdirectory used for the LMDB index. +const INDEX_SUBDIR: &str = "index"; +/// Accounts table name. +const ACCOUNTS_INDEX: &str = "accounts"; +/// Program ownership table name. +const PROGRAMS_INDEX: &str = "programs"; +/// Freelist table name. +const FREELIST_INDEX: &str = "freelist"; + +/// Iterator over all persisted accounts in pubkey order. +type RoAccountIter<'a> = RoIter<'a, PubkeyBytes, OwnerAndOffset>; +/// Duplicate iterator over program-owned persisted accounts. +type RoProgramIter<'a> = RoIter<'a, KeyTail, Offset, MoveOnCurrentKeyDuplicates>; +/// Iterator over persisted accounts. +pub(crate) struct AccountIter<'a> { + /// Iterator over `pubkey -> account` entries. + pub(super) inner: RoAccountIter<'a>, + /// Keeps the read transaction alive for the iterator lifetime. + pub(super) _txn: RoTxnTls<'a>, +} +/// Duplicate iterator over persisted accounts for one owner. +pub(crate) struct OwnerIter<'a> { + /// Duplicates iterator over `owner -> account` entries. + pub(crate) inner: RoProgramIter<'a>, + /// Keeps the read transaction alive for the iterator lifetime. + pub(crate) _txn: RoTxnTls<'a>, +} + +/// LMDB index over persisted account offsets and owners. +pub(crate) struct Index { + /// LMDB environment for the on-disk index. + pub(super) env: Env, + /// Account pubkey -> offset + owner keytag. + pub(super) accounts: Database, + /// Owner keytag -> offset. + pub(super) programs: Database, + /// Image size -> offset. + pub(super) freelist: Database, +} + +impl Index { + /// Opens or creates the index directory and databases. + pub(crate) fn new(path: &Path) -> crate::Result { + let path = path.join(INDEX_SUBDIR); + fs::create_dir_all(&path)?; + // SAFETY: this process owns the index directory for the lifetime of + // the database, so the backing files are not mutated behind LMDB's back. + let env = unsafe { + EnvOpenOptions::new() + .max_dbs(3) + .map_size(INDEX_MAP_SIZE) + .flags(EnvFlags::WRITE_MAP) + .flags(EnvFlags::NO_READ_AHEAD) + .flags(EnvFlags::NO_SYNC) + .open(path)? + }; + + let mut txn = env.write_txn()?; + let accounts = env.database_options().name(ACCOUNTS_INDEX).types().create(&mut txn)?; + let programs = env + .database_options() + .name(PROGRAMS_INDEX) + .flags(DatabaseFlags::DUP_SORT | DatabaseFlags::DUP_FIXED) + .types() + .create(&mut txn)?; + let freelist = env + .database_options() + .name(FREELIST_INDEX) + .flags(DatabaseFlags::DUP_SORT | DatabaseFlags::DUP_FIXED) + .key_comparator() + .types() + .create(&mut txn)?; + txn.commit()?; + Ok(Self { + env, + accounts, + programs, + freelist, + }) + } + + /// Returns the persisted offset for `pubkey`. + pub(crate) fn offset(&self, key: &Pubkey, txn: &RoTxn<'_>) -> Result> { + let entry = self.accounts.get(txn, key)?; + Ok(entry.map(|e| e.offset)) + } + + /// Takes a freed span from the freelist when one matches `units`. + pub(crate) fn allocate(&self, units: u32, txn: &mut RwTxn<'_>) -> Result> { + let offset = self.freelist.get(txn, &units)?; + if let Some(offset) = offset { + self.freelist.delete_one_duplicate(txn, &units, &offset)?; + Ok(Some(offset)) + } else { + Ok(None) + } + } + + /// Inserts an account and its owner mapping. + pub(crate) fn insert( + &self, + key: &Pubkey, + data: OwnerAndOffset, + txn: &mut RwTxn<'_>, + ) -> Result<()> { + self.accounts.put(txn, key, &data)?; + let OwnerAndOffset { owner, offset } = data; + self.programs.put(txn, &owner, &offset) + } + + /// Removes an account and returns its persisted offset. + pub(crate) fn delete(&self, key: &Pubkey, txn: &mut RwTxn<'_>) -> Result> { + let Some(entry) = self.accounts.get(txn, key)? else { + return Ok(None); + }; + + let OwnerAndOffset { owner, offset } = entry; + self.accounts.delete(txn, key)?; + + self.programs.delete_one_duplicate(txn, &owner, &offset)?; + Ok(Some(offset)) + } + + /// Returns the duplicate iterator for accounts owned by `owner`. + pub(crate) fn program<'a>(&'a self, owner: Pubkey) -> Result>> { + let owner = owner.into(); + let txn = self.env.read_txn()?; + let Some(iter) = self.programs.get_duplicates(&txn, &owner)? else { + return Ok(None); + }; + // The duplicate iterator borrows `txn`; storing it in the wrapper keeps + // the borrow alive for the iterator lifetime. + // SAFETY: the wrapper owns `txn`, so the duplicate iterator cannot outlive it. + let iter = unsafe { mem::transmute::, RoProgramIter<'a>>(iter) }; + Ok(Some(OwnerIter { _txn: txn, inner: iter })) + } + + /// Returns an iterator over all accounts in pubkey order. + pub(crate) fn accounts<'a>(&'a self) -> Result> { + let txn = self.env.read_txn()?; + let iter = self.accounts.iter(&txn)?; + // The iterator borrows `txn`; storing it in the wrapper keeps the + // transaction alive for the iterator lifetime. + // SAFETY: the wrapper owns `txn`, so the iterator cannot outlive it. + let iter = unsafe { mem::transmute::, RoAccountIter<'a>>(iter) }; + Ok(AccountIter { _txn: txn, inner: iter }) + } + + /// Moves an account entry to a new owner while preserving its offset. + pub(crate) fn update_owner( + &self, + acc: &Pubkey, + new: KeyTail, + txn: &mut RwTxn<'_>, + ) -> Result<()> { + let Some(val) = self.accounts.get(txn, acc)? else { + return Ok(()); + }; + let OwnerAndOffset { owner: old, offset } = val; + self.programs.delete_one_duplicate(txn, &old, &offset)?; + + let data = OwnerAndOffset { owner: new, offset }; + self.accounts.put(txn, acc, &data)?; + self.programs.put(txn, &new, &offset) + } + + /// Moves an account entry to a new offset while preserving its owner. + pub(crate) fn relocate( + &self, + key: &Pubkey, + old: Offset, + new: OwnerAndOffset, + txn: &mut RwTxn<'_>, + ) -> Result<()> { + self.accounts.put(txn, key, &new)?; + let OwnerAndOffset { owner, offset } = new; + self.programs.delete_one_duplicate(txn, &owner, &old)?; + self.programs.put(txn, &owner, &offset) + } +} + +unsafe impl DatabaseIndex for Index { + fn env(&self) -> &Env { + &self.env + } +} diff --git a/accountsdb/src/store/kv.rs b/accountsdb/src/store/kv.rs new file mode 100644 index 00000000..75b26a1e --- /dev/null +++ b/accountsdb/src/store/kv.rs @@ -0,0 +1,123 @@ +use std::{array, borrow::Cow, ops}; + +use bytemuck::{Pod, Zeroable}; +use heed::{BoxedError, BytesDecode, BytesEncode, byteorder::LittleEndian, types::U32}; +use solana_pubkey::Pubkey; + +/// Result type used by LMDB byte codecs. +pub(crate) type CodecResult = Result; +/// Little-endian `u32` value stored in the freelist. +pub(super) type U32LE = U32; +/// Offset into mapped storage, measured in storage units. +#[derive(Clone, Copy, Zeroable, Pod, PartialEq, Eq, PartialOrd, Ord)] +#[repr(C)] +pub(crate) struct Offset(pub(super) u32); + +/// Compact 16-byte LMDB tag derived from the tail half of a pubkey. +#[derive(Clone, Copy, Pod, Zeroable)] +#[repr(C)] +pub(crate) struct KeyTail([u8; 16]); + +impl From for KeyTail { + fn from(v: Pubkey) -> Self { + Self(array::from_fn(|i| v.as_array()[i + size_of::()])) + } +} + +/// Full 32-byte pubkey codec for the accounts table. +pub(super) struct PubkeyBytes; + +/// LMDB value for the accounts table. +#[derive(Clone, Copy, Pod, Zeroable)] +#[repr(C)] +pub(crate) struct OwnerAndOffset { + /// Owner key tag for the stored account image. + pub(crate) owner: KeyTail, + /// Offset into mapped storage. + pub(crate) offset: Offset, +} + +impl<'a> BytesEncode<'a> for KeyTail { + type EItem = Self; + + fn bytes_encode(item: &'a Self::EItem) -> CodecResult> { + Ok(bytemuck::bytes_of(item).into()) + } +} + +impl<'a> BytesDecode<'a> for KeyTail { + type DItem = &'a Self; + + fn bytes_decode(bytes: &'a [u8]) -> CodecResult { + bytemuck::try_from_bytes(bytes).map_err(Into::into) + } +} + +impl<'a> BytesEncode<'a> for PubkeyBytes { + type EItem = Pubkey; + + fn bytes_encode(item: &'a Self::EItem) -> CodecResult> { + Ok(item.as_array().into()) + } +} + +impl<'a> BytesDecode<'a> for PubkeyBytes { + type DItem = &'a Pubkey; + + fn bytes_decode(bytes: &'a [u8]) -> CodecResult { + bytemuck::try_from_bytes(bytes).map_err(Into::into) + } +} + +impl<'a> BytesEncode<'a> for OwnerAndOffset { + type EItem = Self; + + fn bytes_encode(item: &'a Self::EItem) -> CodecResult> { + Ok(bytemuck::bytes_of(item).into()) + } +} + +impl<'a> BytesDecode<'a> for OwnerAndOffset { + type DItem = Self; + + fn bytes_decode(bytes: &'a [u8]) -> CodecResult { + bytemuck::try_pod_read_unaligned(bytes).map_err(Into::into) + } +} + +impl<'a> BytesEncode<'a> for Offset { + type EItem = Self; + + fn bytes_encode(item: &'a Self::EItem) -> CodecResult> { + U32LE::bytes_encode(&item.0) + } +} + +impl<'a> BytesDecode<'a> for Offset { + type DItem = Self; + + fn bytes_decode(bytes: &'a [u8]) -> CodecResult { + U32LE::bytes_decode(bytes).map(Self) + } +} + +impl ops::Add for Offset { + type Output = Self; + fn add(self, rhs: u32) -> Self::Output { + Self(self.0 + rhs) + } +} + +impl ops::Sub for Offset { + type Output = Self; + fn sub(self, rhs: u32) -> Self::Output { + Self(self.0 - rhs) + } +} + +impl ops::Sub for Offset { + type Output = u32; + fn sub(self, rhs: Self) -> Self::Output { + self.0 - rhs.0 + } +} diff --git a/accountsdb/src/store/mmap.rs b/accountsdb/src/store/mmap.rs new file mode 100644 index 00000000..e1685f5d --- /dev/null +++ b/accountsdb/src/store/mmap.rs @@ -0,0 +1,303 @@ +//! Mapped storage for persisted account images. +//! +//! The file reserves a small meta header at the front, followed by the raw +//! account images written in `solana-account`'s borrowed layout. + +use std::{ + fs::File, + io::{self, Write}, + ops::Range, + os::fd::AsRawFd, + path::Path, + ptr::NonNull, + sync::atomic::{AtomicU32, AtomicU64, Ordering::*}, +}; + +use memmap2::{MmapMut, MmapOptions}; +use nucleus::MB; +use parking_lot::Mutex; +use solana_account::{STORAGE_UNIT, StorageUnit}; +use tracing::{debug, error}; + +use crate::{ + AccountsDBError, Result, metrics, + store::{DatabaseVersion, VERSION, kv::Offset}, +}; + +/// Bytes reserved at the front of the mapped file for metadata. +const DATABASE_META_RESERVATION: usize = 256; +/// Filename used for the mapped storage file. +pub const STORAGE_FILE: &str = "storage.db"; +/// Growth block for the mapped storage file. +#[cfg(any(test, feature = "testkit"))] +pub(crate) const STORAGE_BLOCK: u64 = 16 * MB as u64; +#[cfg(not(any(test, feature = "testkit")))] +pub(crate) const STORAGE_BLOCK: u64 = 256 * MB as u64; +/// Initial file size: one storage block plus the metadata reservation. +const INIT_STORAGE_SIZE: u64 = STORAGE_BLOCK + DATABASE_META_RESERVATION as u64; +/// Maximum mapped storage size. +#[cfg(any(test, feature = "testkit"))] +const MMAP_SIZE: usize = 64 * MB; +#[cfg(not(any(test, feature = "testkit")))] +const MMAP_SIZE: usize = u32::MAX as usize * STORAGE_UNIT + DATABASE_META_RESERVATION; + +/// One allocation inside the mapped storage. +pub(crate) struct Allocation { + /// Offset from the start of the storage area, in storage units. + pub(crate) offset: Offset, + /// Pointer to the start of the allocated image. + pub(crate) ptr: NonNull, +} + +/// Mapped storage backing persisted account images. +pub(crate) struct MappedStorage { + /// Pointer to the reserved metadata header. + meta: NonNull, + /// Full file mapping. + mmap: MmapMut, + /// Start of the account image region. + head: NonNull, + /// File handle used for resizing. + file: Mutex, +} + +#[repr(C)] +#[derive(Default)] +/// Runtime counters for the persisted backend. +pub(crate) struct Stats { + /// Persisted image loads. + pub(crate) reads: AtomicU64, + /// `BorrowedAccount::commit` calls. + pub(crate) commits: AtomicU64, + /// Fresh allocations on backing storage. + pub(crate) allocs: AtomicU64, + /// Freelist allocation reuse. + pub(crate) reallocs: AtomicU64, + /// Account relocations during defrag. + pub(crate) compactions: AtomicU64, + /// Persisted deletes. + pub(crate) removals: AtomicU64, + /// File resizes. + pub(crate) resizes: AtomicU64, +} + +impl Stats { + /// Counts one persisted read. + pub(crate) fn read(&self) { + self.reads.fetch_add(1, Relaxed); + metrics::read(); + } + + /// Counts one borrowed account commit. + pub(crate) fn commit(&self) { + self.commits.fetch_add(1, Relaxed); + metrics::commit(); + } + + /// Counts one fresh allocation from the mapped file. + pub(crate) fn alloc(&self) { + self.allocs.fetch_add(1, Relaxed); + metrics::alloc(); + } + + /// Counts one freelist reuse. + pub(crate) fn realloc(&self) { + self.reallocs.fetch_add(1, Relaxed); + metrics::realloc(); + } + + /// Counts one relocation during defragmentation. + pub(crate) fn compact(&self) { + self.compactions.fetch_add(1, Relaxed); + metrics::compaction(); + } + + /// Counts one persisted removal. + pub(crate) fn remove(&self) { + self.removals.fetch_add(1, Relaxed); + metrics::removal(); + } + + /// Counts one file resize. + pub(crate) fn resize(&self) { + self.resizes.fetch_add(1, Relaxed); + metrics::resize(); + } +} + +#[repr(C)] +#[derive(Default)] +/// Metadata header stored at the front of the mapped file. +pub(crate) struct DatabaseMeta { + /// On-disk format version. + version: DatabaseVersion, + /// Last computed database checksum. + pub(crate) checksum: AtomicU64, + /// Current slot. + pub(crate) slot: AtomicU64, + /// Id of the last sealed superblock; folded into the checksum fingerprint. + pub(crate) superblock: AtomicU64, + /// Current backing file length in bytes. + len: AtomicU64, + /// Database statistics. + stats: Stats, + /// Next allocation cursor. + pub(super) cursor: AtomicU32, +} + +impl MappedStorage { + /// Opens or creates the mapped storage file. + pub(crate) fn new(path: &Path) -> Result { + let path = path.join(STORAGE_FILE); + let mut file = + File::options().create(true).truncate(false).read(true).write(true).open(path)?; + let fd = file.as_raw_fd(); + // SAFETY: the file is opened read/write and mapped for the full fixed size. + let mut mmap = unsafe { MmapOptions::new().len(MMAP_SIZE).map_mut(fd)? }; + if file.metadata()?.len() == 0 { + file.set_len(INIT_STORAGE_SIZE)?; + file.flush()?; + let meta = DatabaseMeta { + version: VERSION, + len: INIT_STORAGE_SIZE.into(), + slot: 1.into(), + ..Default::default() + }; + // SAFETY: the first bytes of the mapping are reserved for `DatabaseMeta`. + unsafe { mmap.as_mut_ptr().cast::().write(meta) }; + mmap.flush()?; + } + // SAFETY: the mapping is at least `DATABASE_META_RESERVATION` bytes long, + // so the meta header and account head pointers stay within the map. + let (meta, head) = unsafe { + let head = mmap.as_mut_ptr().add(DATABASE_META_RESERVATION); + let head = NonNull::new_unchecked(head.cast()); + let meta = NonNull::new_unchecked(mmap.as_mut_ptr().cast()); + (meta, head) + }; + let file = Mutex::new(file); + Ok(Self { meta, mmap, head, file }) + } + + /// Flushes dirty pages to durable storage. + pub(crate) fn flush(&self, sync: bool) -> io::Result<()> { + let range = self.active(); + if sync { + self.mmap.flush_range(range.start, range.len()) + } else { + self.mmap.flush_async_range(range.start, range.len()) + } + } + + /// Validates the opened storage format. + pub(crate) fn validate(&self) -> Result<()> { + let meta = self.meta(); + if meta.version != VERSION { + Err(AccountsDBError::UnsupportedVersion(meta.version)) + } else { + Ok(()) + } + } + + /// Returns a pointer inside the account image region. + pub(crate) fn at(&self, offset: Offset) -> NonNull { + // SAFETY: private call sites pass offsets from the index, allocator, or + // defrag cursor and uphold the mapped-region bounds. + unsafe { self.head.add(offset.0 as usize) } + } + + /// Returns the runtime counters. + pub(crate) fn stats(&self) -> &Stats { + &self.meta().stats + } + + /// Allocates a fresh span of `units` storage units. + pub(crate) fn allocate(&self, units: u32) -> Result { + let meta = self.meta(); + let mut offset = meta.cursor.load(Acquire); + loop { + let end = offset.checked_add(units).ok_or(AccountsDBError::Allocation)?; + let needed = Self::bytes(end as u64); + if needed > meta.len.load(Acquire) { + self.grow(needed)?; + } + if let Err(updated) = meta.cursor.compare_exchange(offset, end, AcqRel, Acquire) { + offset = updated; + } else { + break; + } + } + self.stats().alloc(); + let offset = Offset(offset); + let ptr = self.at(offset); + Ok(Allocation { offset, ptr }) + } + + /// Returns a shared reference to the metadata header. + pub(crate) fn meta(&self) -> &DatabaseMeta { + // SAFETY: `meta` points to the reserved header at the front of the map. + unsafe { &*self.meta.as_ptr() } + } + + /// Returns the current allocation cursor in storage units. + pub(crate) fn cursor(&self) -> u32 { + self.meta().cursor.load(Acquire) + } + + /// Shrinks the file to the current cursor. + pub(super) fn shrink(&self, units: u32) -> io::Result<()> { + self.resize(Self::bytes(units as u64), u64::le)?; + self.meta().cursor.store(units, Release); + Ok(()) + } + + /// Returns the active byte range, including the metadata reservation. + fn active(&self) -> Range { + 0..Self::bytes(self.cursor() as u64) as usize + } + + /// Converts storage units into file bytes, including the metadata reservation. + fn bytes(units: u64) -> u64 { + units * STORAGE_UNIT as u64 + DATABASE_META_RESERVATION as u64 + } + + /// Grows the file to at least `len` bytes. + /// Rounds up to a storage block before resizing. + fn grow(&self, mut len: u64) -> Result<()> { + len = len.div_ceil(STORAGE_BLOCK) * STORAGE_BLOCK; + if len > MMAP_SIZE as u64 { + error!( + requested = len, + limit = MMAP_SIZE, + "mapped storage limit exceeded" + ); + return Err(AccountsDBError::Allocation); + } + self.resize(len, u64::ge).map_err(Into::into) + } + + /// Resizes the file when the current size does not satisfy `cmp`. + /// + /// The file is updated before the new size is published into metadata so + /// readers never observe a larger size than the actual mapping. + fn resize(&self, len: u64, cmp: fn(&u64, &u64) -> bool) -> io::Result<()> { + let mut file = self.file.lock(); + if cmp(&file.metadata()?.len(), &len) { + return Ok(()); + } + // Resize the file first, then publish the new size into metadata. + file.set_len(len)?; + file.flush()?; + self.meta().len.store(len, Release); + self.stats().resize(); + self.mmap.flush()?; + debug!(len, "resized storage file"); + Ok(()) + } +} + +// SAFETY: the `NonNull` pointers point into the owned `mmap` and are never +// reseated; concurrent access is synchronized through atomics in the metadata +// header and the `Mutex`, so the storage is safe to send and share. +unsafe impl Send for MappedStorage {} +unsafe impl Sync for MappedStorage {} diff --git a/accountsdb/src/store/mod.rs b/accountsdb/src/store/mod.rs new file mode 100644 index 00000000..c7d23211 --- /dev/null +++ b/accountsdb/src/store/mod.rs @@ -0,0 +1,301 @@ +//! Persisted account load and write path. +//! +//! This module coordinates the mmap, LMDB index, and borrowed account layout. + +use core::{hash::Hasher, slice}; +use std::sync::atomic::Ordering::{Acquire, Release}; + +use solana_account::{ + AccountMode, AccountSharedData, BorrowedAccount, CoWAccount::*, DirtyMarkers, OwnedAccount, +}; +use solana_pubkey::Pubkey; +use tracing::{error, warn}; + +use nucleus::heed::{DatabaseIndex, OptRoTxn, OptRwTxn}; +use twox_hash::XxHash3_64; + +use crate::{ + AccountEntry, AccountsDBError, Result, StoreKind, + metrics::{self, Operation}, + store::{ + index::{Index, OwnerIter}, + kv::{Offset, OwnerAndOffset}, + mmap::{DatabaseMeta, MappedStorage}, + }, +}; + +mod defrag; +pub(crate) mod index; +mod kv; +pub(crate) mod mmap; + +pub(crate) use mmap::Stats; + +/// Current on-disk storage format version. +pub(crate) const VERSION: DatabaseVersion = 1; +/// Version tag stored in the metadata header. +pub(crate) type DatabaseVersion = u64; + +/// Persisted store backed by the mmap and LMDB index. +pub(crate) struct PersistedStore { + /// Mapped account storage. + pub(crate) storage: MappedStorage, + /// LMDB index over persisted accounts. + pub(crate) index: Index, +} + +/// Iterator over persisted program-owned accounts. +pub(crate) struct PersistedProgramIter<'a> { + /// Keeps the read transaction alive while iterating. + iter: OwnerIter<'a>, + /// Mapped storage backing the returned borrowed accounts. + mmap: &'a MappedStorage, +} + +impl PersistedStore { + /// Opens or creates the persisted store at `path`. + pub(crate) fn new(path: &std::path::Path) -> Result { + let index = Index::new(path)?; + let storage = MappedStorage::new(path)?; + Ok(Self { storage, index }) + } + + /// Loads the persisted image for `pubkey` from the mapped file. + pub(crate) fn load<'e>( + &'e self, + txn: OptRoTxn<'_, 'e>, + pubkey: &Pubkey, + ) -> Result> { + let txn = self.index.read_txn(txn)?; + let offset = self.index.offset(pubkey, txn)?; + offset.is_some().then(|| self.storage.stats().read()); + // SAFETY: offsets come from the persisted index and point into the map. + Ok(offset.map(|o| unsafe { BorrowedAccount::init(self.storage.at(o)) })) + } + + /// Returns whether a persisted account image exists for `pubkey`. + pub(crate) fn contains<'e>(&'e self, txn: OptRoTxn<'_, 'e>, pubkey: &Pubkey) -> Result { + let txn = self.index.read_txn(txn)?; + self.index.offset(pubkey, txn).map(|o| o.is_some()).map_err(Into::into) + } + + /// Applies a batch of account updates to the persisted store. + /// + /// Borrowed accounts in authoritative modes are committed in place. Owned + /// accounts in those modes are serialized into the mmap. Other modes delete + /// stale persisted entries. If the LMDB commit fails or database runs out of + /// space, the borrowed images are rolled back so in-memory state stays + /// aligned with the durable index. + pub(crate) fn upsert<'a, AC>(&self, accounts: AC) -> Result<()> + where + AC: IntoIterator + Clone, + { + let mut applied = 0; + let mut result = Ok(()); + let mut txn = None; + for entry in accounts.clone() { + result = self.apply(entry, &mut txn); + if result.is_err() { + break; + } + applied += 1; + } + // Commit once after the batch so the index and mmap stay in sync. + if let Some(txn) = txn + && result.is_ok() + { + metrics::accounts(StoreKind::Persisted, self.index.accounts.len(&txn)?); + result = txn.commit().map_err(Into::into); + } + if let Err(error) = &result { + warn!(applied, ?error, "accounts persistence failed; rolling back"); + // Only borrowed accounts need rollback here: owned inserts never + // mutate an existing borrowed image in place. + let processed = accounts.into_iter().take(applied).map(|(_, a)| a); + Self::rollback(processed); + } + + result + } + + /// Returns the persisted program iterator for `owner`. + pub(crate) fn program(&self, owner: Pubkey) -> Result>> { + let i = self.index.program(owner)?; + Ok(i.map(|iter| PersistedProgramIter { iter, mmap: &self.storage })) + } + + /// Flushes the mapped storage and LMDB index to durable storage. + pub(crate) fn flush(&self, sync: bool) -> heed::Result<()> { + let _timer = metrics::time(Operation::Flush); + self.index.flush()?; + if sync { + let checksum = self.checksum()?; + self.meta().checksum.store(checksum, Release); + } + self.storage.flush(sync)?; + Ok(()) + } + + /// Validates the persisted store checksum and on-disk format version. + pub(crate) fn validate(&self) -> Result<()> { + self.storage.validate()?; + if self.storage.cursor() == 0 { + return Ok(()); + } + let expected = self.meta().checksum.load(Acquire); + let actual = self.checksum()?; + if expected != actual { + error!(expected, actual, "state checksum mismatch"); + return Err(AccountsDBError::Corruption); + } + + Ok(()) + } + + /// Applies one account state transition to the persisted backend. + fn apply<'e>(&'e self, acc: &AccountEntry, txn: OptRwTxn<'_, 'e>) -> Result<()> { + let (pubkey, account) = acc; + // An account that has moved to a non-authoritative mode, or has been + // closed, no longer belongs here, so drop any stale persisted entry. + if !account.mode().authoritative() || account.is(AccountMode::Closed) { + self.delete(pubkey, txn)?; + if let Borrowed(acc) = account.cow() { + acc.commit(); + } + return Ok(()); + } + + let markers = account.markers(); + match account.cow() { + Borrowed(acc) => self.update(pubkey, acc, markers, txn), + Owned(acc) => self.insert(pubkey, acc, txn), + } + } + + /// Rolls back borrowed accounts that were already touched in the batch. + fn rollback<'a, AC>(accounts: AC) + where + AC: Iterator, + { + for acc in accounts { + if !acc.dirty() { + continue; + } + let Borrowed(acc) = acc.cow() else { continue }; + // SAFETY: only borrowed accounts were updated before the failed commit. + unsafe { acc.rollback() }; + } + } + + /// Commits a borrowed image after updating its owner mapping if needed. + fn update<'e>( + &'e self, + pubkey: &Pubkey, + acc: &BorrowedAccount, + markers: &DirtyMarkers, + txn: OptRwTxn<'_, 'e>, + ) -> Result<()> { + if markers.contains(DirtyMarkers::OWNER) { + let txn = self.index.write_txn(txn)?; + let owner = acc.owner().into(); + self.index.update_owner(pubkey, owner, txn)?; + } + if !markers.intersects(DirtyMarkers::all()) { + return Ok(()); + } + acc.commit(); + self.storage.stats().commit(); + Ok(()) + } + + /// Serializes an owned image into mapped storage and records its offset. + fn insert<'e>( + &'e self, + pubkey: &Pubkey, + acc: &OwnedAccount, + txn: OptRwTxn<'_, 'e>, + ) -> Result<()> { + let txn = self.index.write_txn(txn)?; + let units = acc.units(); + let owner = acc.owner().into(); + + let (ptr, offset) = if let Some(offset) = self.index.allocate(units, txn)? { + let ptr = self.storage.at(offset); + self.storage.stats().realloc(); + (ptr, offset) + } else { + let alloc = self.storage.allocate(units)?; + (alloc.ptr, alloc.offset) + }; + let data = OwnerAndOffset { owner, offset }; + if let Some(offset) = self.index.delete(pubkey, txn)? { + self.free(offset, txn)?; + } + self.index.insert(pubkey, data, txn)?; + // SAFETY: `ptr` points at a fresh span inside the mapped storage and + // `units` is the exact serialized size of this owned account. + unsafe { + let buffer = slice::from_raw_parts_mut(ptr.as_ptr(), units as usize); + acc.serialize(buffer, pubkey); + }; + Ok(()) + } + + /// Returns one persisted span to the freelist. + fn free(&self, offset: Offset, txn: &mut heed::RwTxn<'_>) -> Result<()> { + // SAFETY: `offset` was returned by the index and still points at a valid image. + let space = unsafe { BorrowedAccount::span(self.storage.at(offset)) }; + self.index.freelist.put(txn, &space, &offset)?; + Ok(()) + } + + /// Removes a persisted image and returns its storage span to the freelist. + fn delete<'e>(&'e self, pubkey: &Pubkey, txn: OptRwTxn<'_, 'e>) -> Result<()> { + let txn = self.index.write_txn(txn)?; + let Some(offset) = self.index.delete(pubkey, txn)? else { + return Ok(()); + }; + self.free(offset, txn)?; + self.storage.stats().remove(); + Ok(()) + } + + /// Returns the persisted storage metadata header. + pub(crate) fn meta(&self) -> &DatabaseMeta { + self.storage.meta() + } + + /// Computes a deterministic checksum over persisted accounts in pubkey order. + fn checksum(&self) -> heed::Result { + let _timer = metrics::time(Operation::Checksum); + let mut hasher = XxHash3_64::new(); + let mut iter = self.index.accounts()?; + hasher.write(&self.meta().slot.load(Acquire).to_le_bytes()); + hasher.write(&self.meta().superblock.load(Acquire).to_le_bytes()); + for entry in &mut iter.inner { + let (pubkey, data) = entry?; + hasher.write(pubkey.as_array()); + // SAFETY: offsets come from the persisted accounts index and point + // into the mapped storage for this store. + let account = unsafe { BorrowedAccount::init(self.storage.at(data.offset)) }; + hasher.write(account.storage()); + } + Ok(hasher.finish()) + } +} + +impl<'a> Iterator for PersistedProgramIter<'a> { + type Item = AccountEntry; + + fn next(&mut self) -> Option { + let (_, offset) = self.iter.inner.next()?.ok()?; + let ptr = self.mmap.at(offset); + // The image prefix stores the full pubkey, so iteration can recover it + // without consulting LMDB again. + // SAFETY: the iterator yields offsets stored in the same mapped database. + self.mmap.stats().read(); + let pubkey = unsafe { BorrowedAccount::pubkey(ptr) }; + let account = unsafe { BorrowedAccount::init(ptr).into() }; + Some((pubkey, account)) + } +} diff --git a/accountsdb/src/tests.rs b/accountsdb/src/tests.rs new file mode 100644 index 00000000..99834156 --- /dev/null +++ b/accountsdb/src/tests.rs @@ -0,0 +1,589 @@ +//! Integration-style unit tests for the two-backend account store. +//! +//! Each test drives a realistic multi-step flow through the public `AccountsDB` +//! surface and reaches into `pub(crate)` internals only to assert *which* +//! backend a given account landed in — the crate's central persisted/volatile +//! invariant that no public method exposes directly. + +use std::sync::atomic::Ordering::{Relaxed, Release}; + +use assert_matches::assert_matches; +use nucleus::testkit::{TempDir, init_tracing, tempdir}; +use solana_account::{ + AccountBuilder, AccountMode, AccountSharedData, ReadableAccount, WritableAccount, +}; +use solana_pubkey::Pubkey; + +use super::*; +use crate::snapshot::VOLATILE_DB_FILE; + +/// Fresh database on a throwaway directory; the `TempDir` must outlive the db. +fn db() -> (TempDir, AccountsDB) { + init_tracing(); + let dir = tempdir(); + let db = AccountsDB::new(dir.path()).unwrap(); + (dir, db) +} + +/// Owned mutable (persisted) account carrying `data`; its size follows the data. +fn mutable_data(lamports: u64, data: Vec, owner: &Pubkey) -> AccountSharedData { + let mut a = AccountSharedData::new(lamports, data.len(), owner); + a.set_data_from_slice(&data); + a.set_mode(AccountMode::Delegated).unwrap(); + a +} + +/// Empty mutable (persisted) account; `owner` defaults to the system program. +fn delegated(lamports: u64) -> AccountSharedData { + AccountBuilder::default() + .lamports(lamports) + .mode(AccountMode::Delegated) + .build() +} + +/// Stores one account, the shape every single-account write below takes. +fn store(db: &AccountsDB, pubkey: Pubkey, account: AccountSharedData) { + db.store(&[(pubkey, account)]).unwrap(); +} + +/// Whether a persisted image exists for `pubkey`. +fn in_persisted(db: &AccountsDB, pubkey: &Pubkey) -> bool { + let mut txn = None; + db.persisted.contains(&mut txn, pubkey).unwrap() +} + +/// Whether a volatile entry exists for `pubkey`. +fn in_volatile(db: &AccountsDB, pubkey: &Pubkey) -> bool { + db.volatile.contains(pubkey) +} + +/// Pubkeys `owner` owns, in iteration order (persisted first, then volatile). +fn program(db: &AccountsDB, owner: &Pubkey) -> Vec { + db.program(owner).unwrap().map(|(k, _)| k).collect() +} + +/// Balance of the account currently loaded for `pubkey`. +fn lamports(db: &AccountsDB, pubkey: &Pubkey) -> u64 { + db.loader().load(pubkey).unwrap().unwrap().lamports() +} + +/// Loads the account currently stored for `pubkey`. +/// +/// A persisted account comes back as a *borrowed* image and a volatile one as +/// *owned*; storing the loaded value back is how the engine drives mode changes +/// through the routing layer (a freshly built owned account with a +/// non-authoritative mode is filtered out of the persisted backend entirely). +fn reload(db: &AccountsDB, pubkey: &Pubkey) -> AccountSharedData { + db.loader().load(pubkey).unwrap().unwrap() +} + +/// Closes `pubkey`, deleting it from whichever backend currently holds it. +/// +/// Goes through the load→mutate→store path so the account is a *borrowed* image +/// the routing layer will actually evict (see [`reload`]). +fn close(db: &AccountsDB, pubkey: &Pubkey) { + let mut acc = reload(db, pubkey); + if acc.is(AccountMode::Delegated) { + acc.set_mode(AccountMode::Transient).unwrap(); + } + if acc.is(AccountMode::Transient) { + acc.set_mode(AccountMode::ReadOnly).unwrap(); + } + acc.set_mode(AccountMode::Closed).unwrap(); + store(db, *pubkey, acc); +} + +/// Allocation high-water mark of the persisted store, in storage units. +fn cursor(db: &AccountsDB) -> u32 { + db.persisted.storage.cursor() +} + +/// Defragments the persisted store until the cursor stops moving. +/// +/// A single pass sweeps only the tail half of the hole list, so full compaction +/// needs repeated passes. +fn defrag_to_stable(db: &AccountsDB) { + loop { + let before = cursor(db); + // SAFETY: the test is the sole owner of the store during defrag. + unsafe { db.persisted.defragment() }.unwrap(); + if cursor(db) == before { + break; + } + } +} + +// Routing, both eviction directions, owner remap and Closed/reset handling in +// one flow — the persisted-vs-volatile invariant is what this whole crate +// exists to enforce. +#[test] +fn test_routing_and_persistence_flips() { + let (_dir, db) = db(); + let (p, q) = (Pubkey::new_unique(), Pubkey::new_unique()); + let (a, b) = (Pubkey::new_unique(), Pubkey::new_unique()); + + let aacc = AccountBuilder::default().lamports(10).owner(p).mode(AccountMode::Delegated); + let bacc = AccountBuilder::default().lamports(20).owner(p); + // `a` is authoritative, `b` is non-authoritative; both are owned by `p`. + db.store(&[(a, aacc.build()), (b, bacc.build())]).unwrap(); + assert!(in_persisted(&db, &a) && !in_volatile(&db, &a)); + assert!(in_volatile(&db, &b) && !in_persisted(&db, &b)); + + // Loader reads across both backends; contains agrees. + let loader = db.loader(); + assert_eq!(loader.load(&a).unwrap().unwrap().lamports(), 10); + assert_eq!(loader.load(&b).unwrap().unwrap().lamports(), 20); + assert!(loader.contains(&a).unwrap() && loader.contains(&b).unwrap()); + assert!(!loader.contains(&Pubkey::new_unique()).unwrap()); + drop(loader); + + // Persisted account is yielded before the volatile one. + assert_eq!(program(&db, &p), vec![a, b]); + + // Loading a persisted account returns a borrowed image; mutating its owner + // and re-storing must commit in place and remap the program index. + let mut borrowed = reload(&db, &a); + borrowed.set_owner(q); + store(&db, a, borrowed); + assert_eq!(program(&db, &p), vec![b]); // `a` left p's set + assert_eq!(program(&db, &q), vec![a]); // and joined q's + + // Transient is immutable to programs but remains persistent while its + // lifecycle state is unresolved. + let mut flip = reload(&db, &a); + flip.set_mode(AccountMode::Transient).unwrap(); + flip.set_lamports(30); + store(&db, a, flip); + let transient = reload(&db, &a); + assert!(!transient.mutable()); + assert!(in_persisted(&db, &a) && !in_volatile(&db, &a)); + assert_eq!(lamports(&db, &a), 30); + + // Resolving to ReadOnly evicts the persisted copy into volatile. + let mut flip = reload(&db, &a); + flip.set_mode(AccountMode::ReadOnly).unwrap(); + store(&db, a, flip); + assert!(!in_persisted(&db, &a) && in_volatile(&db, &a)); + assert_eq!(lamports(&db, &a), 30); + + // ReadOnly → Delegated evicts it back into persisted. + let mut flip = reload(&db, &a); + flip.set_mode(AccountMode::Delegated).unwrap(); + store(&db, a, flip); + assert!(in_persisted(&db, &a) && !in_volatile(&db, &a)); + + // Closing removes it from both backends. + close(&db, &a); + assert!(!in_persisted(&db, &a) && !in_volatile(&db, &a)); + + // reset() drops volatile mirror only; persisted state is authoritative. + let c = Pubkey::new_unique(); + store(&db, c, mutable_data(50, vec![], &p)); + db.reset(); + assert!(!in_volatile(&db, &b)); + assert!(in_persisted(&db, &c)); +} + +// Both migration directions remove the source image and owner mapping, retain +// the account contents across reopen, and recycle persisted storage. +#[test] +fn test_store_kind_migration_invariants() { + let dir = tempdir(); + let (persisted_owner, volatile_owner, reuse_owner) = ( + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + ); + let (key, reuse) = (Pubkey::new_unique(), Pubkey::new_unique()); + let data = vec![1, 2, 3, 4]; + let assert_migrated = |db: &AccountsDB, persisted: bool, owner: Pubkey| { + assert_eq!(in_persisted(db, &key), persisted); + assert_eq!(in_volatile(db, &key), !persisted); + assert_eq!(program(db, &owner), vec![key]); + let account = reload(db, &key); + assert_eq!(account.owner(), &owner); + assert_eq!(account.lamports(), 20); + assert_eq!(account.data(), data); + }; + + { + let db = AccountsDB::new(dir.path()).unwrap(); + store(&db, key, mutable_data(10, data.clone(), &persisted_owner)); + let base = cursor(&db); + + let mut account = reload(&db, &key); + account.set_mode(AccountMode::Transient).unwrap(); + account.set_mode(AccountMode::ReadOnly).unwrap(); + account.set_owner(volatile_owner); + account.set_lamports(20); + store(&db, key, account); + + assert_migrated(&db, false, volatile_owner); + assert!(program(&db, &persisted_owner).is_empty()); + + // A same-sized persisted account must reuse the span released by the + // migration instead of extending the mmap. + store(&db, reuse, mutable_data(30, data.clone(), &reuse_owner)); + assert_eq!(cursor(&db), base); + + db.dump(None).unwrap(); + } + + { + let db = AccountsDB::new(dir.path()).unwrap(); + assert_migrated(&db, false, volatile_owner); + assert!(program(&db, &persisted_owner).is_empty()); + + let mut account = reload(&db, &key); + account.set_mode(AccountMode::Delegated).unwrap(); + account.set_owner(persisted_owner); + store(&db, key, account); + + assert_migrated(&db, true, persisted_owner); + assert!(program(&db, &volatile_owner).is_empty()); + + db.flush(true).unwrap(); + // Persist a stale volatile copy if cleanup regresses, so the final open + // can verify source-store cleanup rather than merely losing memory state. + db.dump(None).unwrap(); + } + + let db = AccountsDB::new(dir.path()).unwrap(); + assert_migrated(&db, true, persisted_owner); + assert!(program(&db, &volatile_owner).is_empty()); + assert_eq!(program(&db, &reuse_owner), vec![reuse]); +} + +// Persisted state and metadata survive a close/reopen, and validate() accepts +// the synced checksum. +#[test] +fn test_persistence_reopen_and_validate() { + let dir = tempdir(); + let keys: Vec = (0..8).map(|_| Pubkey::new_unique()).collect(); + + let checksum = { + let db = AccountsDB::new(dir.path()).unwrap(); + for (i, k) in keys.iter().enumerate() { + store(&db, *k, delegated(100 + i as u64)); + } + db.set_slot(42).unwrap(); + // Sync the checksum into the header so a reopen can validate against it. + db.persisted.flush(true).unwrap(); + assert!(db.validate().is_ok()); + db.checksum() + }; + + let db = AccountsDB::new(dir.path()).unwrap(); + for (i, k) in keys.iter().enumerate() { + assert_eq!(lamports(&db, k), 100 + i as u64); + } + assert_eq!(db.slot(), 42); + assert_eq!(db.checksum(), checksum); + assert!(db.validate().is_ok()); +} + +// A clean-shutdown dump lives in the active tree, is restored on the next open, +// and is then removed so volatile state returns to its in-memory-only form. +#[test] +fn test_dump_restores_volatile_on_reopen() { + let dir = tempdir(); + let key = Pubkey::new_unique(); + let active = AccountsDB::directory(dir.path()); + let dump = active.join(VOLATILE_DB_FILE); + + { + let db = AccountsDB::new(dir.path()).unwrap(); + let account = AccountBuilder::default().lamports(42).mode(AccountMode::ReadOnly).build(); + store(&db, key, account); + db.dump(None).unwrap(); + assert!(dump.exists(), "dump is written into the active tree"); + } + + let db = AccountsDB::new(dir.path()).unwrap(); + assert_eq!(lamports(&db, &key), 42); + assert!(in_volatile(&db, &key)); + assert!(!dump.exists(), "restored dump is consumed on open"); +} + +// A freed span is reused for a same-sized insert instead of growing the file; +// genuinely new accounts still extend it. +#[test] +fn test_freelist_reuse_and_growth() { + let (_dir, db) = db(); + let stats = || { + let s = db.persisted.storage.stats(); + (s.allocs.load(Relaxed), s.reallocs.load(Relaxed)) + }; + + let k1 = Pubkey::new_unique(); + store(&db, k1, delegated(1)); + let base = cursor(&db); + let (_, reallocs) = stats(); + + // Close k1 (returns its span to the freelist), then insert a same-sized + // account: it should land in the freed span without advancing the cursor. + close(&db, &k1); + store(&db, Pubkey::new_unique(), delegated(2)); + assert_eq!(cursor(&db), base); + assert_eq!(stats().1, reallocs + 1); + + // Fresh accounts have no reusable span, so the file grows. + let (allocs, _) = stats(); + for _ in 0..8 { + store(&db, Pubkey::new_unique(), delegated(1)); + } + assert!(cursor(&db) > base); + assert!(stats().0 > allocs); +} + +// Defragmentation reclaims interior holes while preserving every live account's +// content, ownership index, and checksum. +#[test] +fn test_defragment_preserves_live_accounts() { + let (_dir, db) = db(); + let owner = Pubkey::new_unique(); + let keys: Vec = (0..16).map(|_| Pubkey::new_unique()).collect(); + for (i, k) in keys.iter().enumerate() { + store(&db, *k, mutable_data(100 + i as u64, vec![], &owner)); + } + + // Punch alternating holes; keep the survivors for later comparison. + let mut live = Vec::new(); + for (i, k) in keys.iter().enumerate() { + if i % 2 == 0 { + close(&db, k); + } else { + live.push((*k, 100 + i as u64)); + } + } + db.persisted.flush(true).unwrap(); + let checksum = db.checksum(); + let before = cursor(&db); + + defrag_to_stable(&db); + assert!(cursor(&db) < before); + + // Every survivor still loads unchanged and remains program-indexed. + for (k, lam) in &live { + let acc = db.loader().load(k).unwrap().unwrap(); + assert_eq!(acc.lamports(), *lam); + assert_eq!(acc.owner(), &owner); + } + let mut owned = program(&db, &owner); + owned.sort(); + let mut expected: Vec = live.iter().map(|(k, _)| *k).collect(); + expected.sort(); + assert_eq!(owned, expected); + + // Relocating images must not change the content checksum. + db.persisted.flush(true).unwrap(); + assert_eq!(db.checksum(), checksum); +} + +// A snapshot is a self-contained tree: reopening it restores persisted accounts +// and bootstraps the volatile store from volatile.db, which is then consumed. +#[test] +fn test_snapshot_export_and_volatile_restore() { + let src = tempdir(); + let (a, b) = (Pubkey::new_unique(), Pubkey::new_unique()); + + let snapshot = { + let db = AccountsDB::new(src.path()).unwrap(); + let aacc = AccountBuilder::default().lamports(10).mode(AccountMode::Delegated); + let bacc = AccountBuilder::default().lamports(20).mode(AccountMode::ReadOnly); + db.store(&[(a, aacc.build()), (b, bacc.build())]).unwrap(); + // SAFETY: the test holds exclusive access to the store. + unsafe { db.snapshot(1) }.unwrap() + }; + + // Adopt the snapshot as a new database's active tree. + let dst = tempdir(); + let active = AccountsDB::directory(dst.path()); + std::fs::rename(&snapshot, &active).unwrap(); + let db = AccountsDB::new(dst.path()).unwrap(); + + assert_eq!(lamports(&db, &a), 10); + assert_eq!(lamports(&db, &b), 20); + assert!(in_persisted(&db, &a)); + assert!(in_volatile(&db, &b)); + // The volatile payload is single-sourced back into memory on open. + assert!(!active.join(VOLATILE_DB_FILE).exists()); + + // Backup renames the active tree out and back. + let saved = db.backup(BackupOp::Save).unwrap(); + assert!(saved.exists() && !active.exists()); + db.backup(BackupOp::Restore).unwrap(); + assert!(active.exists()); +} + +// validate() flags a persisted checksum that no longer matches the images. +#[test] +fn test_corruption_detection() { + let (_dir, db) = db(); + for _ in 0..4 { + store(&db, Pubkey::new_unique(), delegated(1)); + } + db.persisted.flush(true).unwrap(); + assert!(db.validate().is_ok()); + + // Corrupt the recorded checksum; recomputation must no longer agree. + db.persisted.meta().checksum.store(0xDEAD_BEEF, Release); + assert_matches!(db.validate(), Err(AccountsDBError::Corruption)); +} + +// Several freed spans of one size accumulate as duplicates under a single +// freelist key and are all reissued before the file grows — the N>1 duplicate +// case a broken DUP config silently loses. +#[test] +fn test_freelist_multi_duplicate_reuse() { + let (_dir, db) = db(); + let reallocs = || db.persisted.storage.stats().reallocs.load(Relaxed); + + const N: usize = 6; + let keys: Vec = (0..N).map(|_| Pubkey::new_unique()).collect(); + for k in &keys { + store(&db, *k, delegated(1)); + } + let base = cursor(&db); + // Close them all: N same-size spans return to the freelist as N duplicates. + for k in &keys { + close(&db, k); + } + let reused = reallocs(); + + // Each of N fresh same-size inserts must land in a freed span, so the cursor + // never advances and every insert is a reuse. + for _ in 0..N { + store(&db, Pubkey::new_unique(), delegated(2)); + } + assert_eq!(cursor(&db), base); + assert_eq!(reallocs(), reused + N as u64); +} + +// An immutable account changing owner is re-homed in the volatile program index +// and the now-empty old owner set is pruned. +#[test] +fn test_volatile_owner_remap() { + let (_dir, db) = db(); + let (x, y) = (Pubkey::new_unique(), Pubkey::new_unique()); + let k = Pubkey::new_unique(); + + store( + &db, + k, + AccountBuilder::default().lamports(10).owner(x).build(), + ); + assert_eq!(program(&db, &x), vec![k]); + + // Re-store the volatile account under a new owner. + let mut moved = reload(&db, &k); + moved.set_owner(y); + store(&db, k, moved); + + assert_eq!(program(&db, &x), Vec::::new()); // old set pruned + assert_eq!(program(&db, &y), vec![k]); + assert!(in_volatile(&db, &k)); +} + +// The freelist reuses a span only on an exact size match, and accounts of mixed +// sizes survive defragmentation with their data intact. +#[test] +fn test_variable_sizes_and_exact_freelist() { + let (_dir, db) = db(); + let owner = Pubkey::new_unique(); + let reallocs = || db.persisted.storage.stats().reallocs.load(Relaxed); + + // A freed large span cannot satisfy a smaller allocation: sizes differ, so + // the small insert allocates fresh rather than reusing the hole. + let big = Pubkey::new_unique(); + store(&db, big, mutable_data(1, vec![0; 4096], &owner)); + close(&db, &big); + let before = reallocs(); + let small = Pubkey::new_unique(); + store(&db, small, mutable_data(2, vec![0; 64], &owner)); + assert_eq!(reallocs(), before); // size mismatch -> no reuse + + // Store a spread of sizes with distinct data, punch an interior hole, then + // defragment and confirm every survivor keeps its exact bytes. + let sizes = [8usize, 512, 100, 4096, 1]; + let mut live = Vec::new(); + for (i, &space) in sizes.iter().enumerate() { + let k = Pubkey::new_unique(); + let data: Vec = (0..space).map(|b| (b as u8).wrapping_add(i as u8)).collect(); + store(&db, k, mutable_data(i as u64, data.clone(), &owner)); + live.push((k, data)); + } + close(&db, &small); + + defrag_to_stable(&db); + for (k, data) in &live { + assert_eq!( + db.loader().load(k).unwrap().unwrap().data(), + data.as_slice() + ); + } +} + +// The checksum hashes accounts in pubkey order, so it depends only on content — +// not on insertion order or the resulting on-disk offsets. +#[test] +fn test_checksum_order_independent() { + let keys: Vec = (0..8).map(|_| Pubkey::new_unique()).collect(); + + let checksum = |order: &[usize]| { + let (_dir, db) = db(); + for &i in order { + store(&db, keys[i], delegated(100 + i as u64)); + } + db.persisted.flush(true).unwrap(); + db.checksum() + }; + + let forward: Vec = (0..keys.len()).collect(); + let reversed: Vec = (0..keys.len()).rev().collect(); + assert_eq!(checksum(&forward), checksum(&reversed)); +} + +// 2 MiB accounts overflow the initial storage block, forcing the file to grow; +// removing half then defragmenting reclaims the large holes and shrinks the +// cursor back — growth and compaction over multi-megabyte images. +#[test] +fn test_large_accounts_growth_and_defrag() { + const SIZE: usize = 2 << 20; // 2 MiB of data per account + const COUNT: usize = 12; // ~24 MiB total, past the 16 MiB test block + + let (_dir, db) = db(); + let owner = Pubkey::new_unique(); + let resizes = || db.persisted.storage.stats().resizes.load(Relaxed); + let baseline = resizes(); + + // Distinct fill byte per account so content is verifiable without retaining + // the expected bytes. + let keys: Vec = (0..COUNT).map(|_| Pubkey::new_unique()).collect(); + for (i, k) in keys.iter().enumerate() { + store(&db, *k, mutable_data(i as u64, vec![i as u8; SIZE], &owner)); + } + // Crossing the initial block must have grown the file. + assert!(resizes() > baseline); + + // Close every other account to punch large interior holes. + let mut live = Vec::new(); + for (i, k) in keys.iter().enumerate() { + if i % 2 == 0 { + close(&db, k); + } else { + live.push((*k, i as u8)); + } + } + let before = cursor(&db); + + defrag_to_stable(&db); + assert!(cursor(&db) < before); + + // Every survivor keeps its full 2 MiB image byte-for-byte. + for (k, fill) in &live { + let acc = db.loader().load(k).unwrap().unwrap(); + assert_eq!(acc.data().len(), SIZE); + assert!(acc.data().iter().all(|&b| b == *fill)); + } +} diff --git a/accountsdb/src/volatile.rs b/accountsdb/src/volatile.rs new file mode 100644 index 00000000..8e296291 --- /dev/null +++ b/accountsdb/src/volatile.rs @@ -0,0 +1,129 @@ +//! In-memory account cache and program ownership sets. + +use std::{ + collections::BTreeSet, + fs::{self, File}, + io::BufReader, + path::Path, +}; + +use ahash::RandomState; +use scc::HashMap; +use solana_account::{AccountMode, AccountSharedData, OwnedAccount, ReadableAccount}; +use solana_pubkey::Pubkey; +use tracing::info; + +use crate::{Result, StoreKind, metrics, snapshot::VOLATILE_DB_FILE}; + +/// Owned accounts keyed by account pubkey. +type AccountsMap = HashMap; +/// Program ownership sets keyed by owner pubkey. +type ProgramsMap = HashMap, RandomState>; + +/// Volatile account store backed by concurrent hash maps. +pub(crate) struct VolatileStore { + /// Current owned accounts. + pub(crate) accounts: AccountsMap, + /// Program owner -> account pubkeys. + pub(crate) programs: ProgramsMap, +} + +impl VolatileStore { + /// Opens the volatile store, optionally bootstrapping from a snapshot file. + /// + /// If `volatile.db` exists, it is loaded into memory and then removed from + /// the snapshot directory so the active tree stays single-sourced. + pub(crate) fn new(path: &Path) -> Result { + const CAP: usize = 2048; + let snapshot = path.join(VOLATILE_DB_FILE); + let accounts: AccountsMap = if snapshot.exists() { + let mut r = BufReader::new(File::open(&snapshot)?); + let accs: AccountsMap = bincode::deserialize_from(&mut r)?; + fs::remove_file(snapshot)?; + info!( + count = accs.len(), + "restored volatile accounts from snapshot" + ); + accs + } else { + AccountsMap::with_capacity_and_hasher(CAP, Default::default()) + }; + + let programs = ProgramsMap::with_capacity_and_hasher(CAP, Default::default()); + accounts.iter_sync(|&pk, acc| { + BTreeSet::insert(&mut programs.entry_sync(acc.owner()).or_default(), pk) + }); + Ok(Self { accounts, programs }) + } + + /// Stores volatile accounts and keeps the program ownership sets in sync. + pub(crate) fn upsert<'a, AC>(&self, accounts: AC) + where + AC: IntoIterator, + { + for (pubkey, account) in accounts { + // An account that has moved to an authoritative mode, or has been + // closed, drops any stale volatile copy. + if account.mode().authoritative() || account.is(AccountMode::Closed) { + self.delete(pubkey); + continue; + } + // Non-authoritative accounts stay volatile and update the program + // mapping. + let owner = *account.owner(); + { + let mut set = self.programs.entry_sync(owner).or_default(); + BTreeSet::insert(&mut set, *pubkey); + } + let Some(prev) = self.accounts.upsert_sync(*pubkey, account.owned()) else { + continue; + }; + if prev.owner() == owner { + continue; + } + // Only the old owner set needs cleanup; the new owner was inserted above. + self.programs.remove_if_sync(&prev.owner(), |set| { + set.remove(pubkey); + set.is_empty() + }); + } + metrics::accounts(StoreKind::Volatile, self.accounts.len() as u64); + } + + /// Returns the owned account currently cached for `pubkey`. + pub(crate) fn load(&self, pubkey: &Pubkey) -> Option { + let entry = self.accounts.get_sync(pubkey)?; + Some(entry.get().clone()) + } + + /// Returns whether a volatile account exists for `pubkey`. + pub(crate) fn contains(&self, pubkey: &Pubkey) -> bool { + self.accounts.contains_sync(pubkey) + } + + /// Returns the owned accounts currently mapped to `owner`. + pub(crate) fn program(&self, owner: &Pubkey) -> BTreeSet { + self.programs.read_sync(owner, |_, s| s.clone()).unwrap_or_default() + } + + /// Drops chain-mirrored accounts while retaining internal system accounts. + pub(crate) fn reset(&self) { + self.programs.clear_sync(); + self.accounts.retain_sync(|k, a| { + if !a.is(AccountMode::System) { + return false; + } + let mut set = self.programs.entry_sync(a.owner()).or_default(); + BTreeSet::insert(&mut set, *k) + }); + } + + /// Removes the cached account and drops its owner mapping. + fn delete(&self, pubkey: &Pubkey) { + let Some(e) = self.accounts.remove_sync(pubkey) else { return }; + self.programs.remove_if_sync(&e.1.owner(), |set| { + set.remove(pubkey); + set.is_empty() + }); + } +} diff --git a/solana/account/src/cow/tests.rs b/solana/account/src/cow/tests.rs deleted file mode 100644 index 89ad2b65..00000000 --- a/solana/account/src/cow/tests.rs +++ /dev/null @@ -1,83 +0,0 @@ -use super::borrowed::BorrowedAccount; -use super::{StorageUnit, init, serialize_buf}; -use crate::AccountBuilder; -use solana_pubkey::Pubkey; -use std::sync::atomic::Ordering::Acquire; - -const BORROWED_LAMPORTS: u64 = 5; -const ACTIVE_DATA: &[u8] = &[1, 2, 3]; -const COMMIT_DATA: &[u8] = &[4, 5, 6]; -const COMMITTED_DATA: &[u8] = &[9, 2, 3]; -const ACTIVE_WRITE: u8 = 9; -const ROLLBACK_WRITE: u8 = 8; -const INITIAL_SEQUENCE: u32 = 0; -const COMMITTED_SEQUENCE: u32 = 1; - -// Serializes an owned account into a borrowed buffer image. -fn make_buf(data: &[u8]) -> Vec { - let owner = Pubkey::new_unique(); - let owned = AccountBuilder::default() - .lamports(BORROWED_LAMPORTS) - .data(data.to_vec()) - .owner(owner) - .build(); - serialize_buf(&owned) -} - -// Reads the active sequence counter. -fn seq(acc: &BorrowedAccount) -> u32 { - // SAFETY: test helpers only call this on a live borrowed buffer. - unsafe { acc.header.as_ref().sequence.load(Acquire) } -} - -// Returns the active image bytes for direct assertions. -fn data(acc: &BorrowedAccount) -> &[u8] { - &acc.data -} - -#[test] -// `init` should read the active image without changing the sequence. -fn test_init_reads_active_image() { - let mut buf = make_buf(ACTIVE_DATA); - let borrowed = init(&mut buf); - - assert_eq!(data(&borrowed), ACTIVE_DATA); - assert_eq!(seq(&borrowed), INITIAL_SEQUENCE); -} - -#[test] -// `translate` should copy the active image into the shadow view, and `commit` should publish it. -fn test_translate_commit_publishes_shadow_image() { - let mut buf = make_buf(ACTIVE_DATA); - let mut borrowed = init(&mut buf); - - // SAFETY: `borrowed` still points at the live borrowed image selected by `init`. - unsafe { borrowed.translate() }; - assert_eq!(seq(&borrowed), INITIAL_SEQUENCE); - - borrowed.data[0] = ACTIVE_WRITE; - borrowed.commit(); - assert_eq!(seq(&borrowed), COMMITTED_SEQUENCE); - - let borrowed = init(&mut buf); - assert_eq!(data(&borrowed), COMMITTED_DATA); -} - -#[test] -// `rollback` should discard shadow writes and restore the active view. -fn test_translate_rollback_discards_shadow_writes() { - let mut buf = make_buf(COMMIT_DATA); - let mut borrowed = init(&mut buf); - - // SAFETY: `borrowed` still points at the live borrowed image selected by `init`. - unsafe { borrowed.translate() }; - borrowed.data[0] = ROLLBACK_WRITE; - - // SAFETY: `reset` is paired with the preceding `translate`. - unsafe { borrowed.reset() }; - assert_eq!(seq(&borrowed), INITIAL_SEQUENCE); - assert_eq!(data(&borrowed), COMMIT_DATA); - - let borrowed = init(&mut buf); - assert_eq!(data(&borrowed), COMMIT_DATA); -} From 5d7bf1244f7ab53c8319da06babd122d7fd591cb Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Fri, 7 Aug 2026 19:33:55 +0400 Subject: [PATCH 2/6] feat: track committed transactions --- accountsdb/Cargo.toml | 1 + accountsdb/README.md | 6 +++++ accountsdb/src/lib.rs | 43 ++++++++++++++++++++++++++++++++++-- accountsdb/src/store/mmap.rs | 2 ++ 4 files changed, 50 insertions(+), 2 deletions(-) diff --git a/accountsdb/Cargo.toml b/accountsdb/Cargo.toml index 59c77730..d6cc7aeb 100644 --- a/accountsdb/Cargo.toml +++ b/accountsdb/Cargo.toml @@ -35,6 +35,7 @@ solana-account = { workspace = true, features = ["serde"] } solana-pubkey = { workspace = true, features = ["bytemuck"] } [dev-dependencies] +accountsdb = { workspace = true, features = ["testkit"] } assert_matches = { workspace = true } nucleus = { workspace = true, features = ["testkit"] } diff --git a/accountsdb/README.md b/accountsdb/README.md index aa9015f1..6d078560 100644 --- a/accountsdb/README.md +++ b/accountsdb/README.md @@ -14,11 +14,17 @@ in persistent storage, inserts owned images there, updates owned volatile images, and removes stale copies after mode changes or closure. `Transient` remains authoritative and runtime-immutable until its lifecycle state resolves. +`AccountsDB::commit` is the ledger-transaction boundary. It stores successful +account transitions and then advances a persistent transaction counter; empty +transitions from failed executions advance the counter as well. Direct `store` +operations used for initialization, sysvars, and administrative writes do not. + ## Persisted layout `CURRENT/storage.db` contains a metadata header followed by account images in the borrowed `solana-account` layout. Each image includes its full pubkey so scans can recover keys without the index. Offsets are measured in 8-byte `StorageUnit`s. +The transaction counter is metadata and is not part of the account checksum. The LMDB index under `CURRENT/index` contains: diff --git a/accountsdb/src/lib.rs b/accountsdb/src/lib.rs index 010d4cb6..9b08bed8 100644 --- a/accountsdb/src/lib.rs +++ b/accountsdb/src/lib.rs @@ -10,7 +10,7 @@ use std::{ use derive_more::From; use nucleus::Slot; use nucleus::heed::RoTxnTls; -use solana_account::{AccountSharedData, CoWAccount}; +use solana_account::{AccountSeqLock, AccountSharedData, CoWAccount}; use solana_pubkey::Pubkey; use tracing::{info, warn}; @@ -81,6 +81,21 @@ impl AccountsDB { Ok(()) } + /// Commits one ledger transaction's account transitions. + /// + /// The transaction count advances only after every supplied transition is + /// stored successfully. Empty transitions count, including failed SVM + /// executions that reached the commit path without account writes. + pub fn commit<'a, AC>(&self, accounts: AC) -> Result<()> + where + AC: IntoIterator + Clone, + ::IntoIter: Clone, + { + self.store(accounts)?; + self.persisted.meta().transactions.fetch_add(1, Release); + Ok(()) + } + /// Creates a loader that reuses a read transaction for persisted lookups. pub fn loader(&self) -> AccountLoader<'_> { AccountLoader::new(self) @@ -109,6 +124,11 @@ impl AccountsDB { self.persisted.meta().superblock.load(Acquire) } + /// Returns the number of successfully committed ledger transactions. + pub fn transactions(&self) -> u64 { + self.persisted.meta().transactions.load(Acquire) + } + /// Records the last sealed superblock id. Set on snapshot, and on replay /// before recomputing the checksum to compare against a seal. pub fn set_superblock(&self, superblock: u64) { @@ -156,7 +176,10 @@ impl<'a> AccountLoader<'a> { Self { txn: Default::default(), db } } - /// Loads one account, checking both backends in turn + /// Loads one account, reusing the persisted read transaction across calls. + /// + /// Reuse the loader for batch lookups to keep them on the same persisted + /// index snapshot. Persisted accounts take precedence over volatile ones. pub fn load(&self, pubkey: &Pubkey) -> Result> { let txn = &mut self.txn.borrow_mut(); if let Some(acc) = self.db.persisted.load(txn, pubkey)? { @@ -172,6 +195,22 @@ impl<'a> AccountLoader<'a> { Ok(account) } + /// Applies `reader` to an account image stable across a concurrent publish. + /// + /// Prefer this over [`Self::load`] when reading fields from persisted + /// accounts that may be updated concurrently. The reader may be called more + /// than once when the borrowed image changes, so it should have no side + /// effects. + pub fn read(&self, pubkey: &Pubkey, reader: F) -> Result> + where + F: Fn(&AccountSharedData) -> R, + { + let Some(account) = self.load(pubkey)? else { + return Ok(None); + }; + Ok(Some(AccountSeqLock::new(account).read(reader))) + } + /// Returns whether an account exists in either backend. pub fn contains(&self, pubkey: &Pubkey) -> Result { let txn = &mut self.txn.borrow_mut(); diff --git a/accountsdb/src/store/mmap.rs b/accountsdb/src/store/mmap.rs index e1685f5d..90f766bf 100644 --- a/accountsdb/src/store/mmap.rs +++ b/accountsdb/src/store/mmap.rs @@ -137,6 +137,8 @@ pub(crate) struct DatabaseMeta { pub(crate) slot: AtomicU64, /// Id of the last sealed superblock; folded into the checksum fingerprint. pub(crate) superblock: AtomicU64, + /// Transactions whose account-state commit completed successfully. + pub(crate) transactions: AtomicU64, /// Current backing file length in bytes. len: AtomicU64, /// Database statistics. From 48e0fb17c66a12938d8f7ee7a2aaa21115af0429 Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Fri, 7 Aug 2026 22:52:26 +0400 Subject: [PATCH 3/6] fix: safely pack fragmented account storage --- accountsdb/README.md | 17 +- accountsdb/src/lib.rs | 24 +++ accountsdb/src/metrics.rs | 6 +- accountsdb/src/snapshot.rs | 14 +- accountsdb/src/store/defrag.rs | 353 ++++++++++++++++++++++++++------- accountsdb/src/store/mmap.rs | 9 +- accountsdb/src/tests.rs | 194 ++++++++++++++++-- 7 files changed, 510 insertions(+), 107 deletions(-) diff --git a/accountsdb/README.md b/accountsdb/README.md index 6d078560..063c757b 100644 --- a/accountsdb/README.md +++ b/accountsdb/README.md @@ -42,15 +42,22 @@ A persisted batch commits its LMDB transaction once. If applying or committing the batch fails, already committed borrowed images are rolled back so indexed state remains authoritative. Freed image spans enter the freelist. -Defragmentation requires exclusive access. It sorts free spans, compacts the -tail half, updates indexes before moving live bytes, and shrinks the mapped file. -Repeated passes reduce the remaining hole set geometrically. +Defragmentation requires exclusive access. Snapshot export packs tail accounts +into exact holes or the smallest fitting holes that retain at least 33 storage +units. It copies only between non-overlapping spans and publishes all +relocations in one index transaction. Vacated source spans are deferred to the +next pass, so some fragmented layouts may stall. + +After validation, keeper startup repeats committed packing passes to a fixed +point before exposing the database to readers. Snapshot export runs one pass. +Both paths synchronously flush successful changes. ## Snapshots and volatile state `AccountsDB::snapshot` requires exclusive write access. It records the -superblock id, defragments and flushes persisted state, clones the active tree, -and serializes the current volatile map into the clone's `volatile.db`. +superblock id, runs one packing pass and flushes persisted state, clones the +active tree, and serializes the current volatile map into the clone's +`volatile.db`. `dump(None)` writes `CURRENT/volatile.db` for a clean externally paced shutdown. The next open restores that file into memory and removes it. `reset` instead diff --git a/accountsdb/src/lib.rs b/accountsdb/src/lib.rs index 9b08bed8..d1dc7ff0 100644 --- a/accountsdb/src/lib.rs +++ b/accountsdb/src/lib.rs @@ -145,6 +145,30 @@ impl AccountsDB { self.persisted.validate() } + /// Compacts persisted storage to a non-overlapping packing fixed point. + /// + /// This must run only after validation and before loaders or iterators are + /// created. Vacated sources become eligible on the following pass, and all + /// successful passes are flushed synchronously before returning. + pub fn compact(&mut self) -> Result { + let mut reclaimed = 0; + let mut changed = false; + loop { + // SAFETY: `&mut self` excludes readers and writers through this handle; + // the store owns its LMDB environment and mapped storage. + let pass = unsafe { self.persisted.defragment() }?; + reclaimed += pass.reclaimed; + changed |= pass.changed(); + if !pass.changed() { + break; + } + } + if changed { + self.flush(true)?; + } + Ok(reclaimed) + } + /// Returns the last checksum published on superblock boundary. pub fn checksum(&self) -> u64 { self.persisted.meta().checksum.load(Acquire) diff --git a/accountsdb/src/metrics.rs b/accountsdb/src/metrics.rs index c1255612..a1263c96 100644 --- a/accountsdb/src/metrics.rs +++ b/accountsdb/src/metrics.rs @@ -118,9 +118,9 @@ pub(crate) fn realloc() { metric::with_metrics(&METRICS, |m| m.reallocs.inc()); } -/// Records one persisted account relocation during defragmentation. -pub(crate) fn compaction() { - metric::with_metrics(&METRICS, |m| m.compactions.inc()); +/// Records persisted account relocations during defragmentation. +pub(crate) fn compaction(count: u64) { + metric::with_metrics(&METRICS, |m| m.compactions.inc_by(count)); } /// Records one persisted account removal. diff --git a/accountsdb/src/snapshot.rs b/accountsdb/src/snapshot.rs index 7bc10fc4..d5190008 100644 --- a/accountsdb/src/snapshot.rs +++ b/accountsdb/src/snapshot.rs @@ -55,18 +55,18 @@ impl AccountsDB { /// Writes a superblock snapshot under `root`. /// /// # Safety - /// The caller must ensure exclusive write access while the snapshot is in - /// progress. The persisted backend is defragmented and flushed before the - /// active tree is cloned and the volatile store is rewritten in the clone. - /// That ordering keeps the exported state coherent only when no concurrent - /// writes can race with the export. + /// The caller must ensure exclusive access while the snapshot is in + /// progress. The persisted backend runs one non-overlapping packing pass + /// and is flushed before the active tree is cloned and the volatile store + /// is rewritten in the clone. That ordering keeps the exported state + /// coherent only when no concurrent access can race with the export. pub unsafe fn snapshot(&self, superblock: u64) -> SnapshotResult { let _timer = metrics::time(Operation::Snapshot); let src = self.root.join(ACTIVE_DIR); let dst = self.root.join(format!("{PREFIX}{superblock:0>9}")); self.set_superblock(superblock); - // SAFETY: snapshot owns exclusive write access, so defrag cannot race - // with concurrent mutation and can compact the persisted store first. + // SAFETY: snapshot owns exclusive access, so defrag cannot race with + // readers or writers while compacting the persisted store. unsafe { self.persisted.defragment() }?; // Persisted state must reach disk before we copy the active tree. self.persisted.flush(true)?; diff --git a/accountsdb/src/store/defrag.rs b/accountsdb/src/store/defrag.rs index f7c593e4..f07c8dbd 100644 --- a/accountsdb/src/store/defrag.rs +++ b/accountsdb/src/store/defrag.rs @@ -1,8 +1,10 @@ #![allow(unsafe_op_in_unsafe_fn)] +use std::{collections::BTreeSet, ops::Range}; + use heed::Result; use solana_account::BorrowedAccount; -use tracing::{debug, info}; +use tracing::info; use crate::{ metrics::{self, Operation}, @@ -11,8 +13,23 @@ use crate::{ use super::PersistedStore; -/// Free span in the persisted image file, measured in storage units and -/// ordered by offset for the sweep. +/// Smallest useful destination remainder, in 8-byte storage units. +const MIN_REMAINDER: u32 = 33; +type Fit = (u32, Offset, usize); + +/// Result of one committed packing pass. +pub(crate) struct Defragged { + pub(crate) moved: usize, + pub(crate) reclaimed: u32, +} + +impl Defragged { + pub(crate) fn changed(&self) -> bool { + self.moved > 0 || self.reclaimed > 0 + } +} + +/// Free span in the persisted image file, measured in storage units. #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] struct Hole { offset: Offset, @@ -20,117 +37,307 @@ struct Hole { } impl Hole { - /// Builds one hole from the freelist entry tuple. fn new((units, offset): (u32, Offset)) -> Self { Self { offset, units } } - /// Returns the first offset past the hole. - fn end(&self) -> Offset { + fn end(self) -> Offset { self.offset + self.units } } +/// Adjacent entry-time holes treated as one packing destination. +struct Run { + parts: Range, + free: Hole, +} + +impl Run { + fn take(&mut self, units: u32) -> Offset { + debug_assert!(units <= self.free.units); + let dst = self.free.offset; + self.free.offset = self.free.offset + units; + self.free.units -= units; + dst + } +} + +/// One account relocation planned against entry-time free space. +#[derive(Clone, Copy)] +struct Move { + src: Offset, + dst: Offset, + units: u32, +} + +impl Move { + fn source(self) -> Hole { + Hole { + offset: self.src, + units: self.units, + } + } +} + +/// Temporary state for one non-overlapping packing pass. +struct Defrag<'a> { + store: &'a PersistedStore, + holes: Vec, + runs: Vec, + moves: Vec, + tail: Offset, +} + impl PersistedStore { - /// Compacts the tail half of the freelist and repacks the storage. + /// Packs tail accounts into holes that existed at the start of this pass. /// - /// Caller must guarantee exclusive access while defrag runs. Holes are - /// swept in file order, contiguous holes are merged into one run, and the - /// following live span is reindexed before its bytes are slid left. Each - /// pass starts at the tail half of the sorted hole list so repeated runs - /// halve the remaining work. + /// Adjacent freelist entries form logical runs. An account uses an exact + /// fit when available, otherwise the smallest run that leaves at least 33 + /// units. Destination remainders may accept more accounts in this pass; + /// vacated source spans are deferred until a later pass. Some fragmented + /// layouts therefore cannot progress. /// - /// This operation is not crash-safe: interruption after reindexing can - /// leave the active tree inconsistent and require restoring a backup. + /// This operation is not crash-safe: interruption after publishing moved + /// offsets can leave the active tree inconsistent and require a backup. /// /// # Safety /// - /// No concurrent mutation may touch the persisted index or mapped - /// storage while offsets are rewritten and bytes are moved. - pub(crate) unsafe fn defragment(&self) -> Result<()> { + /// No concurrent access may touch the persisted index or mapped storage + /// while offsets are rewritten and bytes are moved. + pub(crate) unsafe fn defragment(&self) -> Result { let _timer = metrics::time(Operation::Defragmentation); - let mut holes = { - let txn = self.index.env.read_txn()?; - self.index + Defrag::new(self)?.execute() + } +} + +impl<'a> Defrag<'a> { + /// Reads a consistent entry-time layout and plans tail-to-left moves. + /// + /// # Safety + /// + /// The store must be exclusively accessed, and indexed offsets must point + /// to valid serialized accounts in its mapped storage. + unsafe fn new(store: &'a PersistedStore) -> Result { + let (mut holes, mut accounts) = { + let txn = store.index.env.read_txn()?; + let holes = store + .index .freelist .iter(&txn)? .map(|r| r.map(Hole::new)) - .collect::>>()? + .collect::>>()?; + let accounts = if holes.is_empty() { + Vec::new() + } else { + store + .index + .accounts + .iter(&txn)? + .map(|r| r.map(|(_, data)| data.offset)) + .collect::>>()? + }; + (holes, accounts) }; holes.sort_unstable(); + accounts.sort_unstable(); - let mut i = holes.len() / 2; - // `dst` is the next write head in storage units. - let Some(mut dst) = holes.get(i).map(|h| h.offset) else { - debug!("nothing to defragment"); - return Ok(()); + let runs = Self::runs(&holes); + let mut defrag = Self { + store, + holes, + runs, + moves: Vec::new(), + tail: Offset(store.storage.cursor()), }; - // Tail is the live end; if the last hole reaches it, there is no - // trailing live span to move. - let tail = Offset(self.storage.cursor()); + defrag.pack(accounts.into_iter().rev()); + Ok(defrag) + } + + /// Groups physically adjacent holes without changing their freelist shape. + fn runs(holes: &[Hole]) -> Vec { + let mut runs = Vec::new(); + let mut i = 0; while i < holes.len() { - let run = i; + let first = i; + let offset = holes[i].offset; let mut end = holes[i].end(); i += 1; - - while let Some(&hole) = holes.get(i) { - if hole.offset != end { - break; - } + while let Some(hole) = holes.get(i) + && hole.offset == end + { end = hole.end(); i += 1; } + runs.push(Run { + parts: first..i, + free: Hole { offset, units: end - offset }, + }); + } + runs + } - let next = holes.get(i).map(|h| h.offset).unwrap_or(tail); - // Defrag runs exclusively, so rewrite and commit the index before - // moving the live bytes that now point at the future offsets. - let shift = end - dst; - self.reindex(end, next, shift, &holes[run..i])?; - // Source and destination overlap by design during left-compaction. - dst = self.slide(end, dst, next - end); + /// Selects the best exact fit or the best fit with a useful remainder. + fn fit(fit: &BTreeSet, units: u32) -> Option { + let &(largest, _, _) = fit.last()?; + if units > largest { + return None; } - let reclaimed = tail - dst; - info!(reclaimed, "defragmented persisted storage"); - self.storage.shrink(dst.0).map_err(Into::into) + + let low = (units, Offset(0), 0); + let high = (units, Offset(u32::MAX), usize::MAX); + if let Some(exact) = fit.range(low..=high).next() { + return Some(*exact); + } + + let minimum = units.checked_add(MIN_REMAINDER)?; + if minimum > largest { + return None; + } + fit.range((minimum, Offset(0), 0)..).next().copied() } - /// Rewrites the index for a live span before its bytes are moved. + /// Packs accounts in descending source order into eligible runs. /// /// # Safety /// - /// `src..end` must cover whole serialized accounts in mapped storage. - /// The caller must guarantee exclusive access while offsets are updated. - unsafe fn reindex(&self, src: Offset, end: Offset, shift: u32, run: &[Hole]) -> Result<()> { - let mut txn = self.index.env.write_txn()?; - let mut pos = src; - while pos < end { - let ptr = self.storage.at(pos); - let span = BorrowedAccount::span(ptr); - let next = pos + span; - // The bytes are still at `pos`; compute the future offset before copying. - let pubkey = BorrowedAccount::pubkey(ptr); - let account = BorrowedAccount::init(ptr); - let owner = account.owner().into(); - let data = OwnerAndOffset { owner, offset: pos - shift }; - self.index.relocate(&pubkey, pos, data, &mut txn)?; - self.storage.stats().compact(); - pos = next; + /// Every supplied offset must point to a valid serialized account, and no + /// concurrent access may modify the index, freelist, or mapped storage. + unsafe fn pack(&mut self, accounts: impl Iterator) { + // Best fit by remaining units, then by the lowest current offset. + let mut fit: BTreeSet = self + .runs + .iter() + .enumerate() + .map(|(i, run)| (run.free.units, run.free.offset, i)) + .collect(); + let mut eligible = self.runs.len(); + + for src in accounts { + // Runs are already ordered by their physical end. + while eligible > 0 && self.runs[eligible - 1].free.end() > src { + let i = eligible - 1; + fit.remove(&(self.runs[i].free.units, self.runs[i].free.offset, i)); + eligible -= 1; + } + if fit.is_empty() { + break; + } + + let units = BorrowedAccount::span(self.store.storage.at(src)); + let Some((remaining, start, i)) = Self::fit(&fit, units) else { + continue; + }; + fit.remove(&(remaining, start, i)); + let dst = self.runs[i].take(units); + self.moves.push(Move { src, dst, units }); + let free = self.runs[i].free; + if free.units > 0 { + fit.insert((free.units, free.offset, i)); + } } - for hole in run { - self.index.freelist.delete_one_duplicate(&mut txn, &hole.units, &hole.offset)?; + } + + /// Returns the first unit in the final free suffix without re-sorting it. + fn compacted_tail(&self) -> Offset { + let mut run = self.runs.len(); + let mut movement = 0; + let mut tail = self.tail; + + loop { + while run > 0 && self.runs[run - 1].free.units == 0 { + run -= 1; + } + let free = (run > 0).then(|| self.runs[run - 1].free); + let source = self.moves.get(movement).copied().map(Move::source); + let (hole, from_run) = match (free, source) { + (Some(free), Some(source)) => (free.max(source), free.offset >= source.offset), + (Some(free), None) => (free, true), + (None, Some(source)) => (source, false), + (None, None) => break, + }; + if hole.end() != tail { + break; + } + tail = hole.offset; + if from_run { + run -= 1; + } else { + movement += 1; + } } - txn.commit() + tail } - /// Slides a live region left after the index has been committed. + /// Copies the plan and publishes all index and freelist changes. /// /// # Safety /// - /// `src..src+len` and `dst..dst+len` must be valid mapped storage - /// ranges. The ranges may overlap. - unsafe fn slide(&self, src: Offset, dst: Offset, len: u32) -> Offset { - let src = self.storage.at(src); - src.copy_to(self.storage.at(dst), len as usize); - dst + len + /// The entry-time layout must remain unchanged since planning, and no + /// concurrent access may observe or modify storage while moves publish. + unsafe fn execute(self) -> Result { + let tail = self.compacted_tail(); + let outcome = Defragged { + moved: self.moves.len(), + reclaimed: self.tail - tail, + }; + if !outcome.changed() { + info!("nothing to defragment"); + return Ok(outcome); + } + + // Entry-time destinations are disjoint, so every source remains intact + // until the complete plan has been copied. + for movement in &self.moves { + self.store.storage.at(movement.src).copy_to_nonoverlapping( + self.store.storage.at(movement.dst), + movement.units as usize, + ); + } + + let mut txn = self.store.index.env.write_txn()?; + for movement in &self.moves { + let ptr = self.store.storage.at(movement.src); + let pubkey = BorrowedAccount::pubkey(ptr); + let owner = BorrowedAccount::init(ptr).owner().into(); + let data = OwnerAndOffset { owner, offset: movement.dst }; + self.store.index.relocate(&pubkey, movement.src, data, &mut txn)?; + } + self.publish(tail, &mut txn)?; + txn.commit()?; + + self.store.storage.stats().compact(outcome.moved); + if outcome.reclaimed > 0 { + self.store.storage.shrink(tail.0)?; + } + info!( + moved = outcome.moved, + reclaimed = outcome.reclaimed, + "defragmented persisted storage" + ); + Ok(outcome) + } + + /// Publishes final free spans while retaining untouched component sizes. + fn publish(&self, tail: Offset, txn: &mut heed::RwTxn<'_>) -> Result<()> { + for run in &self.runs { + for &hole in &self.holes[run.parts.clone()] { + let offset = hole.offset.max(run.free.offset); + let end = hole.end().min(tail); + if offset == hole.offset && end == hole.end() { + continue; + } + self.store.index.freelist.delete_one_duplicate(txn, &hole.units, &hole.offset)?; + if offset < end { + self.store.index.freelist.put(txn, &(end - offset), &offset)?; + } + } + } + for movement in &self.moves { + if movement.src < tail { + let end = movement.source().end().min(tail); + self.store.index.freelist.put(txn, &(end - movement.src), &movement.src)?; + } + } + Ok(()) } } diff --git a/accountsdb/src/store/mmap.rs b/accountsdb/src/store/mmap.rs index 90f766bf..e394016d 100644 --- a/accountsdb/src/store/mmap.rs +++ b/accountsdb/src/store/mmap.rs @@ -106,10 +106,11 @@ impl Stats { metrics::realloc(); } - /// Counts one relocation during defragmentation. - pub(crate) fn compact(&self) { - self.compactions.fetch_add(1, Relaxed); - metrics::compaction(); + /// Counts relocations during defragmentation. + pub(crate) fn compact(&self, count: usize) { + let count = count as u64; + self.compactions.fetch_add(count, Relaxed); + metrics::compaction(count); } /// Counts one persisted removal. diff --git a/accountsdb/src/tests.rs b/accountsdb/src/tests.rs index 99834156..a91c6637 100644 --- a/accountsdb/src/tests.rs +++ b/accountsdb/src/tests.rs @@ -8,7 +8,10 @@ use std::sync::atomic::Ordering::{Relaxed, Release}; use assert_matches::assert_matches; -use nucleus::testkit::{TempDir, init_tracing, tempdir}; +use nucleus::{ + heed::DatabaseIndex, + testkit::{TempDir, init_tracing, tempdir}, +}; use solana_account::{ AccountBuilder, AccountMode, AccountSharedData, ReadableAccount, WritableAccount, }; @@ -98,21 +101,34 @@ fn cursor(db: &AccountsDB) -> u32 { db.persisted.storage.cursor() } -/// Defragments the persisted store until the cursor stops moving. -/// -/// A single pass sweeps only the tail half of the hole list, so full compaction -/// needs repeated passes. -fn defrag_to_stable(db: &AccountsDB) { +/// Current persisted offset for one live account. +fn offset(db: &AccountsDB, pubkey: &Pubkey) -> impl Copy + PartialEq + use<> { + let mut txn = None; + let txn = db.persisted.index.read_txn(&mut txn).unwrap(); + db.persisted.index.offset(pubkey, txn).unwrap().unwrap() +} + +/// Defragments until a pass makes no change and returns the reclaimed total. +fn defrag_to_stable(db: &AccountsDB) -> u32 { + let mut total = 0; loop { - let before = cursor(db); // SAFETY: the test is the sole owner of the store during defrag. - unsafe { db.persisted.defragment() }.unwrap(); - if cursor(db) == before { - break; + let pass = unsafe { db.persisted.defragment() }.unwrap(); + total += pass.reclaimed; + if !pass.changed() { + return total; } } } +/// Builds a mutable account with an exact persisted span. +fn mutable_units(lamports: u64, units: u32, owner: &Pubkey) -> AccountSharedData { + (0..1024) + .map(|len| mutable_data(lamports, vec![0; len], owner)) + .find(|account| account.owned().units() == units) + .unwrap() +} + // Routing, both eviction directions, owner remap and Closed/reset handling in // one flow — the persisted-vs-volatile invariant is what this whole crate // exists to enforce. @@ -261,19 +277,26 @@ fn test_persistence_reopen_and_validate() { let dir = tempdir(); let keys: Vec = (0..8).map(|_| Pubkey::new_unique()).collect(); - let checksum = { + let (checksum, before) = { let db = AccountsDB::new(dir.path()).unwrap(); for (i, k) in keys.iter().enumerate() { store(&db, *k, delegated(100 + i as u64)); } + let discarded = Pubkey::new_unique(); + store(&db, discarded, delegated(0)); + close(&db, &discarded); db.set_slot(42).unwrap(); // Sync the checksum into the header so a reopen can validate against it. db.persisted.flush(true).unwrap(); assert!(db.validate().is_ok()); - db.checksum() + (db.checksum(), cursor(&db)) }; - let db = AccountsDB::new(dir.path()).unwrap(); + let mut db = AccountsDB::new(dir.path()).unwrap(); + assert!(db.validate().is_ok()); + let reclaimed = db.compact().unwrap(); + assert_eq!(reclaimed, before - cursor(&db)); + assert!(reclaimed > 0); for (i, k) in keys.iter().enumerate() { assert_eq!(lamports(&db, k), 100 + i as u64); } @@ -360,7 +383,8 @@ fn test_defragment_preserves_live_accounts() { let checksum = db.checksum(); let before = cursor(&db); - defrag_to_stable(&db); + let reclaimed = defrag_to_stable(&db); + assert_eq!(reclaimed, before - cursor(&db)); assert!(cursor(&db) < before); // Every survivor still loads unchanged and remains program-indexed. @@ -380,6 +404,145 @@ fn test_defragment_preserves_live_accounts() { assert_eq!(db.checksum(), checksum); } +/// Exact and thresholded best-fit packing updates component holes correctly, +/// while deferred source holes become usable only by a later committed pass. +#[test] +fn test_defragment_best_fit_and_deferred_holes() { + let owner = Pubkey::new_unique(); + let small = mutable_units(1, 21, &owner); + let medium = mutable_units(2, 33, &owner); + + // Two adjacent component holes form one run. The small tail account leaves + // an exact medium remainder, which the following account consumes. + { + let (_dir, db) = db(); + let (small_hole, medium_hole, anchor, medium_key, small_key) = ( + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + ); + store(&db, small_hole, small.clone()); + store(&db, medium_hole, medium.clone()); + store(&db, anchor, small.clone()); + store(&db, medium_key, medium.clone()); + store(&db, small_key, small.clone()); + let medium_dst = offset(&db, &medium_hole); + let small_dst = offset(&db, &small_hole); + close(&db, &small_hole); + close(&db, &medium_hole); + + let pass = unsafe { db.persisted.defragment() }.unwrap(); + assert_eq!(pass.moved, 2); + assert_eq!(pass.reclaimed, 54); + assert!(offset(&db, &medium_key) == medium_dst); + assert!(offset(&db, &small_key) == small_dst); + assert_eq!(lamports(&db, &medium_key), 2); + assert_eq!(lamports(&db, &small_key), 1); + assert_eq!(lamports(&db, &anchor), 1); + + let before = cursor(&db); + store(&db, Pubkey::new_unique(), small.clone()); + assert_eq!(cursor(&db), before + 21); + } + + // Exact fit wins first. The next account skips a 23-unit hole that would + // leave two units. Two accounts instead use a 75-unit hole and publish its + // reusable 33-unit suffix. + { + let (_dir, db) = db(); + let near = mutable_units(3, 23, &owner); + let wide = mutable_units(4, 75, &owner); + let ( + near_hole, + anchor_a, + wide_hole, + anchor_b, + exact_hole, + anchor_c, + wide_key_a, + wide_key_b, + exact_key, + ) = ( + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + ); + store(&db, near_hole, near); + store(&db, anchor_a, small.clone()); + store(&db, wide_hole, wide); + store(&db, anchor_b, small.clone()); + store(&db, exact_hole, small.clone()); + store(&db, anchor_c, small.clone()); + store(&db, wide_key_a, small.clone()); + store(&db, wide_key_b, small.clone()); + store(&db, exact_key, small.clone()); + let near_dst = offset(&db, &near_hole); + let wide_dst = offset(&db, &wide_hole); + let exact_dst = offset(&db, &exact_hole); + close(&db, &near_hole); + close(&db, &wide_hole); + close(&db, &exact_hole); + + let pass = unsafe { db.persisted.defragment() }.unwrap(); + assert_eq!(pass.moved, 3); + assert_eq!(pass.reclaimed, 63); + assert!(offset(&db, &wide_key_b) == wide_dst); + assert!(offset(&db, &exact_key) == exact_dst); + + let before = cursor(&db); + store(&db, Pubkey::new_unique(), medium.clone()); + assert_eq!(cursor(&db), before); + let near_key = Pubkey::new_unique(); + store(&db, near_key, mutable_units(5, 23, &owner)); + assert_eq!(cursor(&db), before); + assert!(offset(&db, &near_key) == near_dst); + } + + // The first pass moves only the middle account. Its source joins the next + // hole after commit, and public startup compaction exhausts later passes. + { + let (_dir, mut db) = db(); + let second = mutable_units(3, 45, &owner); + let (first_hole, middle, second_hole, tail) = ( + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + Pubkey::new_unique(), + ); + store(&db, first_hole, small.clone()); + store(&db, middle, small.clone()); + store(&db, second_hole, second); + store(&db, tail, medium.clone()); + let middle_dst = offset(&db, &first_hole); + let tail_dst = offset(&db, &middle); + close(&db, &first_hole); + close(&db, &second_hole); + let before = cursor(&db); + + let pass = unsafe { db.persisted.defragment() }.unwrap(); + assert_eq!((pass.moved, pass.reclaimed), (1, 0)); + assert_eq!(cursor(&db), before); + assert!(offset(&db, &middle) == middle_dst); + + assert_eq!(db.compact().unwrap(), 66); + assert!(offset(&db, &tail) == tail_dst); + assert_eq!(lamports(&db, &middle), 1); + assert_eq!(lamports(&db, &tail), 2); + + let before = cursor(&db); + store(&db, Pubkey::new_unique(), small.clone()); + assert_eq!(cursor(&db), before + 21); + } +} + // A snapshot is a self-contained tree: reopening it restores persisted accounts // and bootstraps the volatile store from volatile.db, which is then consumed. #[test] @@ -577,7 +740,8 @@ fn test_large_accounts_growth_and_defrag() { } let before = cursor(&db); - defrag_to_stable(&db); + let reclaimed = defrag_to_stable(&db); + assert_eq!(reclaimed, before - cursor(&db)); assert!(cursor(&db) < before); // Every survivor keeps its full 2 MiB image byte-for-byte. From f18be2f1684adb681dd6612602b47083bc5dc70f Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Sat, 8 Aug 2026 00:56:32 +0400 Subject: [PATCH 4/6] fix: increase the min slack space size to 43 units --- accountsdb/README.md | 4 ++-- accountsdb/src/store/defrag.rs | 10 ++++---- accountsdb/src/store/mod.rs | 2 ++ accountsdb/src/tests.rs | 43 +++++++++++++++++++++------------- 4 files changed, 36 insertions(+), 23 deletions(-) diff --git a/accountsdb/README.md b/accountsdb/README.md index 063c757b..2c25bfb5 100644 --- a/accountsdb/README.md +++ b/accountsdb/README.md @@ -43,8 +43,8 @@ the batch fails, already committed borrowed images are rolled back so indexed state remains authoritative. Freed image spans enter the freelist. Defragmentation requires exclusive access. Snapshot export packs tail accounts -into exact holes or the smallest fitting holes that retain at least 33 storage -units. It copies only between non-overlapping spans and publishes all +into exact holes or the smallest fitting holes that leave a minimum useful +remainder. It copies only between non-overlapping spans and publishes all relocations in one index transaction. Vacated source spans are deferred to the next pass, so some fragmented layouts may stall. diff --git a/accountsdb/src/store/defrag.rs b/accountsdb/src/store/defrag.rs index f07c8dbd..a4fecd81 100644 --- a/accountsdb/src/store/defrag.rs +++ b/accountsdb/src/store/defrag.rs @@ -14,7 +14,7 @@ use crate::{ use super::PersistedStore; /// Smallest useful destination remainder, in 8-byte storage units. -const MIN_REMAINDER: u32 = 33; +pub(crate) const MIN_REMAINDER: u32 = 43; type Fit = (u32, Offset, usize); /// Result of one committed packing pass. @@ -92,10 +92,10 @@ impl PersistedStore { /// Packs tail accounts into holes that existed at the start of this pass. /// /// Adjacent freelist entries form logical runs. An account uses an exact - /// fit when available, otherwise the smallest run that leaves at least 33 - /// units. Destination remainders may accept more accounts in this pass; - /// vacated source spans are deferred until a later pass. Some fragmented - /// layouts therefore cannot progress. + /// fit when available, otherwise the smallest run that leaves at least + /// [`MIN_REMAINDER`] units. Destination remainders may accept more accounts + /// in this pass; vacated source spans are deferred until a later pass. Some + /// fragmented layouts therefore cannot progress. /// /// This operation is not crash-safe: interruption after publishing moved /// offsets can leave the active tree inconsistent and require a backup. diff --git a/accountsdb/src/store/mod.rs b/accountsdb/src/store/mod.rs index c7d23211..0833f049 100644 --- a/accountsdb/src/store/mod.rs +++ b/accountsdb/src/store/mod.rs @@ -29,6 +29,8 @@ pub(crate) mod index; mod kv; pub(crate) mod mmap; +#[cfg(test)] +pub(crate) use defrag::MIN_REMAINDER; pub(crate) use mmap::Stats; /// Current on-disk storage format version. diff --git a/accountsdb/src/tests.rs b/accountsdb/src/tests.rs index a91c6637..083e12a1 100644 --- a/accountsdb/src/tests.rs +++ b/accountsdb/src/tests.rs @@ -18,7 +18,7 @@ use solana_account::{ use solana_pubkey::Pubkey; use super::*; -use crate::snapshot::VOLATILE_DB_FILE; +use crate::{snapshot::VOLATILE_DB_FILE, store::MIN_REMAINDER}; /// Fresh database on a throwaway directory; the `TempDir` must outlive the db. fn db() -> (TempDir, AccountsDB) { @@ -123,9 +123,16 @@ fn defrag_to_stable(db: &AccountsDB) -> u32 { /// Builds a mutable account with an exact persisted span. fn mutable_units(lamports: u64, units: u32, owner: &Pubkey) -> AccountSharedData { - (0..1024) + let account = mutable_at_least(lamports, units, owner); + assert_eq!(account.owned().units(), units); + account +} + +/// Builds the smallest mutable account spanning at least `units` storage units. +fn mutable_at_least(lamports: u64, units: u32, owner: &Pubkey) -> AccountSharedData { + (0..=units as usize * solana_account::STORAGE_UNIT) .map(|len| mutable_data(lamports, vec![0; len], owner)) - .find(|account| account.owned().units() == units) + .find(|account| account.owned().units() >= units) .unwrap() } @@ -410,10 +417,12 @@ fn test_defragment_preserves_live_accounts() { fn test_defragment_best_fit_and_deferred_holes() { let owner = Pubkey::new_unique(); let small = mutable_units(1, 21, &owner); - let medium = mutable_units(2, 33, &owner); + let small_units = small.owned().units(); + let medium = mutable_at_least(2, MIN_REMAINDER, &owner); + let medium_units = medium.owned().units(); // Two adjacent component holes form one run. The small tail account leaves - // an exact medium remainder, which the following account consumes. + // a useful remainder, which the following account consumes exactly. { let (_dir, db) = db(); let (small_hole, medium_hole, anchor, medium_key, small_key) = ( @@ -435,7 +444,7 @@ fn test_defragment_best_fit_and_deferred_holes() { let pass = unsafe { db.persisted.defragment() }.unwrap(); assert_eq!(pass.moved, 2); - assert_eq!(pass.reclaimed, 54); + assert_eq!(pass.reclaimed, small_units + medium_units); assert!(offset(&db, &medium_key) == medium_dst); assert!(offset(&db, &small_key) == small_dst); assert_eq!(lamports(&db, &medium_key), 2); @@ -444,16 +453,18 @@ fn test_defragment_best_fit_and_deferred_holes() { let before = cursor(&db); store(&db, Pubkey::new_unique(), small.clone()); - assert_eq!(cursor(&db), before + 21); + assert_eq!(cursor(&db), before + small_units); } - // Exact fit wins first. The next account skips a 23-unit hole that would - // leave two units. Two accounts instead use a 75-unit hole and publish its - // reusable 33-unit suffix. + // Exact fit wins first. The next account skips a hole whose remainder is + // just below the threshold. Two accounts instead use a wider hole and + // publish a useful suffix. { let (_dir, db) = db(); - let near = mutable_units(3, 23, &owner); - let wide = mutable_units(4, 75, &owner); + let short_remainder = (MIN_REMAINDER - 1) & !1; + let near_units = small_units + short_remainder; + let near = mutable_units(3, near_units, &owner); + let wide = mutable_units(4, 2 * small_units + medium_units, &owner); let ( near_hole, anchor_a, @@ -501,7 +512,7 @@ fn test_defragment_best_fit_and_deferred_holes() { store(&db, Pubkey::new_unique(), medium.clone()); assert_eq!(cursor(&db), before); let near_key = Pubkey::new_unique(); - store(&db, near_key, mutable_units(5, 23, &owner)); + store(&db, near_key, mutable_units(5, near_units, &owner)); assert_eq!(cursor(&db), before); assert!(offset(&db, &near_key) == near_dst); } @@ -510,7 +521,7 @@ fn test_defragment_best_fit_and_deferred_holes() { // hole after commit, and public startup compaction exhausts later passes. { let (_dir, mut db) = db(); - let second = mutable_units(3, 45, &owner); + let second = mutable_units(3, 2 * medium_units - small_units, &owner); let (first_hole, middle, second_hole, tail) = ( Pubkey::new_unique(), Pubkey::new_unique(), @@ -532,14 +543,14 @@ fn test_defragment_best_fit_and_deferred_holes() { assert_eq!(cursor(&db), before); assert!(offset(&db, &middle) == middle_dst); - assert_eq!(db.compact().unwrap(), 66); + assert_eq!(db.compact().unwrap(), 2 * medium_units); assert!(offset(&db, &tail) == tail_dst); assert_eq!(lamports(&db, &middle), 1); assert_eq!(lamports(&db, &tail), 2); let before = cursor(&db); store(&db, Pubkey::new_unique(), small.clone()); - assert_eq!(cursor(&db), before + 21); + assert_eq!(cursor(&db), before + small_units); } } From 1ea00ae80a0ee3ce30cf1bb767cea7a3b4b4e2c2 Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Mon, 10 Aug 2026 14:44:37 +0400 Subject: [PATCH 5/6] refactor: use safe heed transaction helpers --- accountsdb/src/store/index.rs | 2 +- accountsdb/src/store/mod.rs | 12 ++++++------ accountsdb/src/tests.rs | 4 ++-- 3 files changed, 9 insertions(+), 9 deletions(-) diff --git a/accountsdb/src/store/index.rs b/accountsdb/src/store/index.rs index 42d9449b..776be137 100644 --- a/accountsdb/src/store/index.rs +++ b/accountsdb/src/store/index.rs @@ -200,7 +200,7 @@ impl Index { } } -unsafe impl DatabaseIndex for Index { +impl DatabaseIndex for Index { fn env(&self) -> &Env { &self.env } diff --git a/accountsdb/src/store/mod.rs b/accountsdb/src/store/mod.rs index 0833f049..9214362c 100644 --- a/accountsdb/src/store/mod.rs +++ b/accountsdb/src/store/mod.rs @@ -11,7 +11,7 @@ use solana_account::{ use solana_pubkey::Pubkey; use tracing::{error, warn}; -use nucleus::heed::{DatabaseIndex, OptRoTxn, OptRwTxn}; +use nucleus::heed::{DatabaseIndex, OptRoTxn, OptRwTxn, read_txn, write_txn}; use twox_hash::XxHash3_64; use crate::{ @@ -68,7 +68,7 @@ impl PersistedStore { txn: OptRoTxn<'_, 'e>, pubkey: &Pubkey, ) -> Result> { - let txn = self.index.read_txn(txn)?; + let txn = read_txn(self.index.env(), txn)?; let offset = self.index.offset(pubkey, txn)?; offset.is_some().then(|| self.storage.stats().read()); // SAFETY: offsets come from the persisted index and point into the map. @@ -77,7 +77,7 @@ impl PersistedStore { /// Returns whether a persisted account image exists for `pubkey`. pub(crate) fn contains<'e>(&'e self, txn: OptRoTxn<'_, 'e>, pubkey: &Pubkey) -> Result { - let txn = self.index.read_txn(txn)?; + let txn = read_txn(self.index.env(), txn)?; self.index.offset(pubkey, txn).map(|o| o.is_some()).map_err(Into::into) } @@ -198,7 +198,7 @@ impl PersistedStore { txn: OptRwTxn<'_, 'e>, ) -> Result<()> { if markers.contains(DirtyMarkers::OWNER) { - let txn = self.index.write_txn(txn)?; + let txn = write_txn(self.index.env(), txn)?; let owner = acc.owner().into(); self.index.update_owner(pubkey, owner, txn)?; } @@ -217,7 +217,7 @@ impl PersistedStore { acc: &OwnedAccount, txn: OptRwTxn<'_, 'e>, ) -> Result<()> { - let txn = self.index.write_txn(txn)?; + let txn = write_txn(self.index.env(), txn)?; let units = acc.units(); let owner = acc.owner().into(); @@ -253,7 +253,7 @@ impl PersistedStore { /// Removes a persisted image and returns its storage span to the freelist. fn delete<'e>(&'e self, pubkey: &Pubkey, txn: OptRwTxn<'_, 'e>) -> Result<()> { - let txn = self.index.write_txn(txn)?; + let txn = write_txn(self.index.env(), txn)?; let Some(offset) = self.index.delete(pubkey, txn)? else { return Ok(()); }; diff --git a/accountsdb/src/tests.rs b/accountsdb/src/tests.rs index 083e12a1..ac2ff604 100644 --- a/accountsdb/src/tests.rs +++ b/accountsdb/src/tests.rs @@ -9,7 +9,7 @@ use std::sync::atomic::Ordering::{Relaxed, Release}; use assert_matches::assert_matches; use nucleus::{ - heed::DatabaseIndex, + heed::{DatabaseIndex, read_txn}, testkit::{TempDir, init_tracing, tempdir}, }; use solana_account::{ @@ -104,7 +104,7 @@ fn cursor(db: &AccountsDB) -> u32 { /// Current persisted offset for one live account. fn offset(db: &AccountsDB, pubkey: &Pubkey) -> impl Copy + PartialEq + use<> { let mut txn = None; - let txn = db.persisted.index.read_txn(&mut txn).unwrap(); + let txn = read_txn(db.persisted.index.env(), &mut txn).unwrap(); db.persisted.index.offset(pubkey, txn).unwrap().unwrap() } From 7beb4e637060794f79496de117049e18c511d7f8 Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Mon, 10 Aug 2026 14:44:37 +0400 Subject: [PATCH 6/6] refactor: use safe heed transaction helpers --- Cargo.toml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 6b8f34c7..bfd1735d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -22,15 +22,15 @@ rust-version = "1.94.1" version = "0.1.0" [workspace.dependencies] -accountsdb = { path = "accountsdb" } +accountsdb = { path = "accountsdb", package = "magicblock-accountsdb" } magic-root-interface = { path = "programs/magic-root-interface" } magic-root-program = { path = "programs/magic-root-program" } nucleus = { path = "nucleus", package = "magicblock-engine-nucleus" } -v42-calculator-interface = { path = "programs/v42-calculator-interface", default-features = false } solana-account = { path = "solana/account" } solana-program-runtime = { path = "solana/program-runtime" } solana-svm = { path = "solana/svm" } solana-transaction-context = { path = "solana/transaction-context" } +v42-calculator-interface = { path = "programs/v42-calculator-interface", default-features = false } ahash = "0.8.12" arc-swap = "1.9.1"