From 1a54518148c38e09a789511de2b5ba29d938a334 Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Wed, 16 Sep 2026 14:12:07 +0700 Subject: [PATCH 1/2] fix: drain account readers before snapshot compaction --- Cargo.lock | 2 + Cargo.toml | 1 + accountsdb/Cargo.toml | 4 + accountsdb/README.md | 36 ++++- accountsdb/src/lib.rs | 117 ++++++++++++--- accountsdb/src/readers.rs | 234 ++++++++++++++++++++++++++++++ accountsdb/src/snapshot.rs | 18 ++- accountsdb/src/tests.rs | 166 +++++++++++++++++++-- engine/src/testkit.rs | 2 +- keeper/README.md | 14 +- keeper/src/accessor.rs | 20 ++- keeper/src/builder.rs | 12 +- keeper/src/lib.rs | 12 +- keeper/src/subscriptions.rs | 13 +- keeper/src/tests/recovery.rs | 15 +- keeper/src/tests/subscriptions.rs | 2 +- processor/README.md | 6 + processor/src/callback.rs | 18 ++- processor/src/executor.rs | 7 +- processor/src/tests.rs | 7 +- 20 files changed, 616 insertions(+), 90 deletions(-) create mode 100644 accountsdb/src/readers.rs diff --git a/Cargo.lock b/Cargo.lock index d8c0156f..42455f06 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2314,10 +2314,12 @@ dependencies = [ "magicblock-engine-nucleus", "memmap2", "parking_lot", + "rustix", "scc", "solana-account", "solana-pubkey", "thiserror 2.0.20", + "thread_local", "tracing", "twox-hash", ] diff --git a/Cargo.toml b/Cargo.toml index 87e52bd8..cc6863f4 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -79,6 +79,7 @@ snedfile = "0.1.0" tar = "0.4.46" tempfile = "3.27.0" thiserror = "2.0.20" +thread_local = "1.1.10" tokio = "1.53.1" tokio-util = "0.7.19" tracing = "0.1.44" diff --git a/accountsdb/Cargo.toml b/accountsdb/Cargo.toml index 3f8c92e8..75808542 100644 --- a/accountsdb/Cargo.toml +++ b/accountsdb/Cargo.toml @@ -28,12 +28,16 @@ memmap2 = { workspace = true } parking_lot = { workspace = true } scc = { workspace = true, features = ["serde"] } thiserror = { workspace = true } +thread_local = { 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"] } +[target.'cfg(target_os = "linux")'.dependencies] +rustix = { workspace = true, features = ["thread"] } + [dev-dependencies] accountsdb = { workspace = true, features = ["testkit"] } assert_matches = { workspace = true } diff --git a/accountsdb/README.md b/accountsdb/README.md index 88c7f238..de7a6cd0 100644 --- a/accountsdb/README.md +++ b/accountsdb/README.md @@ -43,17 +43,43 @@ limits old-page retention without giving up reader-slot reuse. iteration. The optional `testkit` feature uses smaller maps and growth blocks without changing the on-disk format. +## Reader scopes + +Use `AccountLoader::read` and `AccountsDB::program` for scoped zero-copy +reads; callbacks can return encoded results or owned snapshots. Program iteration +only exposes callback results. `unsafe load` is reserved for +zero-copy transaction execution: borrowed results must remain protected through +commit, with concurrent account writes, deletion, and storage reuse excluded. +Retaining a loader or iterator excludes compaction, not ordinary account writes. +`AccountLoader::unguarded` is the exception: callers must exclude compaction +themselves for the loader and all borrowed results. It skips reader admission +without introducing a separate account lookup path. + +Guarded loaders and program iterators enter the scope before opening their LMDB +transactions. Registered readers update only their own cache-line-separated +slot; first registration and reads meeting compaction use the maintenance mutex. +Nested synchronous scopes share the outer admission. Do not hold a reader scope +across asynchronous work or invoke compaction from inside one. + +Linux requires expedited private `membarrier` support, registered through rustix +when opening the database. Readers use compiler fences; compaction issues the +process-wide barrier. Registration and barrier errors propagate, with admission +restored on maintenance failure. macOS uses full memory fences on both sides of +the same protocol. + ## 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. Snapshot export packs tail accounts -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. +Defragmentation requires exclusive access. Snapshot export drains registered +readers before packing and blocks new readers until relocation and truncation +finish. Account writes must still be quiesced by the caller. Snapshot export +packs tail accounts 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. After validation, keeper startup repeats committed packing passes to a fixed point before exposing the database to readers. Snapshot export runs one pass. diff --git a/accountsdb/src/lib.rs b/accountsdb/src/lib.rs index 2af9cf5d..3892e699 100644 --- a/accountsdb/src/lib.rs +++ b/accountsdb/src/lib.rs @@ -14,6 +14,7 @@ use solana_pubkey::Pubkey; use tracing::{info, warn}; use crate::{ + readers::{ReadGuard, Readers}, store::{DatabaseVersion, PersistedProgramIter, PersistedStore, index::RoTxnTls}, volatile::VolatileStore, }; @@ -22,6 +23,7 @@ pub use snapshot::{BackupOp, SnapshotError, SnapshotResult}; pub use store::mmap::STORAGE_FILE; mod metrics; +mod readers; mod snapshot; mod store; mod volatile; @@ -40,17 +42,25 @@ pub struct AccountsDB { volatile: VolatileStore, /// Database root directory. root: PathBuf, + /// Reader admission while the persisted layout is being relocated. + readers: Readers, } impl AccountsDB { /// Opens or creates the database at `root`. pub fn new(root: impl AsRef) -> Result { + let readers = Readers::new()?; 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 }; + let db = Self { + persisted, + volatile, + root, + readers, + }; metrics::init(&db); Ok(db) } @@ -100,11 +110,31 @@ impl AccountsDB { 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 }) + /// Reads each program-owned account without letting borrowed images escape. + /// + /// The iterator retains reader admission and the persisted index snapshot. + /// `reader` may run again after a concurrent image publish and must have no + /// side effects. Only its result escapes; copying account data is optional. + pub fn program<'a, F, R>( + &'a self, + owner: &Pubkey, + reader: F, + ) -> Result + 'a> + where + F: Fn(&Pubkey, &AccountSharedData) -> R + 'a, + R: 'a, + { + let guard = self.readers.enter(); + let iter = ProgramIter { + persisted: self.persisted.program(*owner)?, + volatile: self.volatile.program(owner), + db: self, + _reader: guard, + }; + Ok(iter.map(move |(pubkey, account)| { + let result = AccountSeqLock::new(account).read(|account| reader(&pubkey, account)); + (pubkey, result) + })) } /// Returns the latest slot persisted in the database metadata. @@ -136,11 +166,14 @@ impl AccountsDB { /// Flushes persisted account storage, forcing synchronous durability when requested. pub fn flush(&self, force: bool) -> Result<()> { + // A synchronous flush walks indexed account images for the checksum. + let _reader = force.then(|| self.readers.enter()); self.persisted.flush(force).map_err(Into::into) } /// Validates the persisted store checksum and on-disk format version. pub fn validate(&self) -> Result<()> { + let _reader = self.readers.enter(); self.persisted.validate() } @@ -186,25 +219,64 @@ impl AccountsDB { } } -/// Loader that caches a read transaction for persisted account lookups. +/// Synchronous reader scope caching a persisted index transaction. +/// +/// Holding a guarded loader delays compaction. Finish the batch and drop it before +/// awaiting unrelated work. [`Self::read`] keeps account access within the scope; +/// raw execution views must obey [`Self::load`]'s safety contract. +/// Execution with its own compaction barrier can use [`Self::unguarded`]. pub struct AccountLoader<'a> { /// Cached read transaction for the persisted index. txn: RefCell>>, /// Database handle used for volatile and persisted lookups. db: &'a AccountsDB, + /// Present for ordinary readers, absent when the caller excludes compaction. + /// Drops after the cached transaction so old offsets cannot escape admission. + _reader: Option>, } impl<'a> AccountLoader<'a> { /// Creates a new loader bound to `db`. + #[inline] pub fn new(db: &'a AccountsDB) -> Self { - Self { txn: Default::default(), db } + Self { + txn: RefCell::new(None), + db, + _reader: Some(db.readers.enter()), + } } - /// Loads one account, reusing the persisted read transaction across calls. + /// Creates a loader without reader-admission bookkeeping. + /// + /// Intended for execution already covered by the sequencer's barrier. It + /// uses the same account lookup and cached index transaction as [`Self::new`]. /// - /// 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> { + /// # Safety + /// The caller must independently exclude compaction for this loader's + /// entire lifetime, including its cached index transaction, and until all + /// returned borrowed accounts have had their final access. + #[inline] + pub unsafe fn unguarded(db: &'a AccountsDB) -> Self { + Self { + txn: RefCell::new(None), + db, + _reader: None, + } + } + + /// Loads a zero-copy account view for externally synchronized execution. + /// + /// The SVM retains these views through transaction commit. Ordinary readers + /// must use [`Self::read`] instead. Both paths reuse the same cached index + /// transaction and give volatile accounts precedence over persisted ones. + /// + /// # Safety + /// The database must outlive borrowed results. Retain a guarded loader until + /// their last access, or independently exclude relocation for their use. + /// Concurrent image updates require `AccountSeqLock`; account deletion and + /// storage reuse must also be excluded while borrowed results are used. + /// Mutating borrowed results additionally requires exclusive account access. + pub unsafe fn load(&self, pubkey: &Pubkey) -> Result> { if let Some(account) = self.db.volatile.load(pubkey) { metrics::load(StoreKind::Volatile); return Ok(Some(account.into())); @@ -221,18 +293,17 @@ impl<'a> AccountLoader<'a> { /// 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. + /// The reader may be called more than once when the borrowed image changes, + /// so it should have no side effects. Only its result escapes the scope; + /// clone the account inside the callback when an owned snapshot is needed. 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))) + // SAFETY: the callback and sequence checks finish within this loader's + // admission scope, or the caller's unguarded-construction contract. + let account = unsafe { self.load(pubkey) }?; + Ok(account.map(|account| AccountSeqLock::new(account).read(reader))) } /// Returns whether an account exists in either backend. @@ -246,14 +317,16 @@ impl<'a> AccountLoader<'a> { } } -/// Iterates program-owned accounts across both backends. -pub struct ProgramIter<'a> { +/// Internal images consumed only by scoped program reads. +struct ProgramIter<'a> { /// Persisted program accounts. persisted: Option>, /// Volatile program pubkeys. volatile: BTreeSet, /// Database handle used to resolve volatile accounts. db: &'a AccountsDB, + /// Outlives the persisted iterator and its LMDB transaction. + _reader: ReadGuard<'a>, } impl<'a> Iterator for ProgramIter<'a> { diff --git a/accountsdb/src/readers.rs b/accountsdb/src/readers.rs new file mode 100644 index 00000000..3f7d7474 --- /dev/null +++ b/accountsdb/src/readers.rs @@ -0,0 +1,234 @@ +//! Reader admission for relocation of the live mapped store. +//! +//! Each thread writes only its own slot. Linux pairs compiler fences on entry/exit +//! with a process-wide membarrier during maintenance; other platforms use full +//! fences on both sides. The gate/slot handshake prevents a reader and the +//! compactor from both overlooking the other's announcement. +//! +//! Admission protects mapped images and cached index offsets from relocation, +//! not concurrent account writes. Callers must enter before opening an index +//! transaction and drop that transaction and all borrowed views before leaving. +//! Maintenance must independently exclude writers while holding its pause. + +use std::{ + io, + marker::PhantomData, + sync::atomic::{AtomicBool, AtomicUsize, Ordering::*, fence}, +}; + +use parking_lot::{Condvar, Mutex, MutexGuard}; +#[cfg(target_os = "linux")] +use rustix::thread::{MembarrierCommand, membarrier}; +use thread_local::ThreadLocal; + +/// Coordinates reader scopes with exclusive relocation of one database. +pub(crate) struct Readers { + /// Stable nesting counters, written only by their owning threads and scanned + /// by maintenance. First registration is serialized with that scan. + slots: ThreadLocal, + /// Admission gate: true while maintenance drains readers or relocates data. + closed: AtomicBool, + /// Serializes pauses and first registration; retained throughout relocation, + /// including while the compactor releases `idle` to wait for readers. + maintenance: Mutex<()>, + /// Couples the active-slot scan with exit notification to prevent lost + /// wakeups. Separate from `maintenance` so waiting never admits a new pause. + idle: Mutex<()>, + /// Wakes the sole compactor when an outer reader exits or withdraws. + drained: Condvar, + /// Test hook before parking for active readers; may fire more than once. + #[cfg(test)] + pub(crate) on_drain: Observer, + /// Test hook after a registered reader withdraws from a closed gate. + #[cfg(test)] + pub(crate) on_block: Observer, +} + +/// Per-database scheduling observations, absent from production builds. +/// Callbacks run under internal locks and must not reenter this protocol. +#[cfg(test)] +type Observer = Mutex>>; + +/// Isolates one thread's writes, including paired 64-byte cache-line prefetches +/// on x86-64. Padding is per thread, not per account. +#[derive(Default)] +#[repr(align(128))] +struct Slot { + /// Number of live guards on the owning thread; zero means inactive. + /// Single-writer ownership permits loads/stores rather than atomic RMWs. + depth: AtomicUsize, +} + +/// Keeps this thread admitted until its outermost synchronous read ends. +pub(crate) struct ReadGuard<'a> { + /// This thread's stable slot, shared by its nested reader scopes. + slot: &'a Slot, + /// Admission state to notify when the final nested scope exits. + readers: &'a Readers, + /// Makes the guard neither Send nor Sync, preserving single-threaded nesting + /// and preventing concurrent non-RMW updates to its slot. + _local: PhantomData<*mut ()>, +} + +/// Reopens admission on success, error, or unwind, before releasing the mutex. +pub(crate) struct Pause<'a> { + /// Gate reopened by Drop, including on barrier failure or unwinding. + readers: &'a Readers, + /// Exclusive maintenance ownership, released only after Drop reopens admission. + _maintenance: MutexGuard<'a, ()>, +} + +impl Readers { + /// Opens admission and registers Linux's process-wide expedited barrier. + /// Registration failure aborts construction; the fence protocol never changes + /// after readers begin using the database. + pub(crate) fn new() -> io::Result { + // Registration is idempotent across databases in the same process. + #[cfg(target_os = "linux")] + membarrier(MembarrierCommand::RegisterPrivateExpedited)?; + Ok(Self { + slots: ThreadLocal::new(), + closed: AtomicBool::new(false), + maintenance: Mutex::new(()), + idle: Mutex::new(()), + drained: Condvar::new(), + #[cfg(test)] + on_drain: Mutex::new(None), + #[cfg(test)] + on_block: Mutex::new(None), + }) + } + + /// Enters a synchronous read scope, waiting if maintenance has closed admission. + /// Nested scopes inherit admission so an existing reader can finish while + /// maintenance waits. Registered, uncontended readers take no mutex. + #[inline] + pub(crate) fn enter(&self) -> ReadGuard<'_> { + let slot = self.slots.get().unwrap_or_else(|| self.register()); + loop { + let depth = slot.depth.load(Relaxed); + slot.depth.store(depth + 1, Relaxed); + let guard = ReadGuard { + slot, + readers: self, + _local: PhantomData, + }; + // Nested scopes inherit outer admission even while maintenance + // is waiting for that outer scope to finish. + if depth != 0 { + return guard; + } + + light_fence(); + if !self.closed.load(Acquire) { + return guard; + } + drop(guard); + #[cfg(test)] + if let Some(observe) = self.on_block.lock().as_ref() { + observe(); + } + // Only readers that meet maintenance touch the mutex. Never wait + // while marked active: the compactor is waiting for these slots. + drop(self.maintenance.lock()); + } + } + + /// Closes admission and drains existing scopes before granting relocation. + /// Returns WouldBlock if this thread holds a reader, avoiding self-deadlock. + /// Barrier failure reopens admission through the temporary pause guard. + pub(crate) fn pause(&self) -> io::Result> { + if self.slots.get().is_some_and(|slot| slot.depth.load(Relaxed) != 0) { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + "cannot compact accountsdb from an active reader scope", + )); + } + let maintenance = self.maintenance.lock(); + let pause = Pause { + readers: self, + _maintenance: maintenance, + }; + self.closed.store(true, Relaxed); + + #[cfg(target_os = "linux")] + membarrier(MembarrierCommand::PrivateExpedited)?; + #[cfg(not(target_os = "linux"))] + fence(SeqCst); + + // Paired entry barriers guarantee that a racing reader either appears + // active here or observes the closed gate before touching the index. + // Registration stays excluded for the full pause, so this set is fixed. + // The separate wait mutex lets exiting readers notify without releasing + // exclusive maintenance ownership or allowing a second compactor in. + let mut idle = self.idle.lock(); + while self.slots.iter().any(|slot| slot.depth.load(Acquire) != 0) { + #[cfg(test)] + if let Some(observe) = self.on_drain.lock().as_ref() { + observe(); + } + self.drained.wait(&mut idle); + } + // Keep subsequent relocation after the complete slot scan, including + // on weakly ordered platforms when a slot was already inactive. + fence(SeqCst); + Ok(pause) + } + + /// Registers a thread once, synchronized with the closed gate and slot scan. + /// Waiting on maintenance prevents adding a slot to an in-progress pause. + #[cold] + fn register(&self) -> &Slot { + let _registration = self.maintenance.lock(); + self.slots.get_or_default() + } + + /// Wakes the sole compactor after an outer reader exits or withdraws. + /// Holding `idle` pairs notification with its scan-and-wait transition. + #[cold] + fn notify(&self) { + let _idle = self.idle.lock(); + self.drained.notify_one(); + } +} + +impl Drop for ReadGuard<'_> { + /// Releases one nesting level; the last exit publishes completed accesses + /// and notifies a waiting compactor after the gate/slot handshake. + #[inline] + fn drop(&mut self) { + // Publish the last mmap access before allowing reclamation. Only this + // thread changes its nesting depth, so no atomic RMW is needed. + let depth = self.slot.depth.load(Relaxed); + self.slot.depth.store(depth - 1, Release); + if depth != 1 { + return; + } + // Pair with maintenance's gate publication and heavy fence: either + // its scan sees our exit or we see the gate and notify. The mutex + // then closes the race between the scan and entering the wait. + light_fence(); + if self.readers.closed.load(Acquire) { + self.readers.notify(); + } + } +} + +impl Drop for Pause<'_> { + /// Publishes completed relocation and reopens admission before unlocking + /// maintenance, so blocked readers can retry against the new layout. + fn drop(&mut self) { + self.readers.closed.store(false, Release); + } +} + +/// Orders the slot announcement before inspecting the admission gate. +/// Linux pairs a compiler fence with maintenance's process-wide membarrier; +/// macOS and other platforms need a full hardware fence on the reader side too. +#[inline] +fn light_fence() { + #[cfg(target_os = "linux")] + std::sync::atomic::compiler_fence(SeqCst); + #[cfg(not(target_os = "linux"))] + fence(SeqCst); +} diff --git a/accountsdb/src/snapshot.rs b/accountsdb/src/snapshot.rs index d5190008..a8471975 100644 --- a/accountsdb/src/snapshot.rs +++ b/accountsdb/src/snapshot.rs @@ -55,19 +55,21 @@ impl AccountsDB { /// Writes a superblock snapshot under `root`. /// /// # Safety - /// 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. + /// The caller must exclude account writes for the duration of export and + /// quiesce any raw borrowed views not covered by reader scopes. + /// Reader admission is paused internally during the packing pass; + /// ordinary reads may run again while the flushed tree is cloned. 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 access, so defrag cannot race with - // readers or writers while compacting the persisted store. - unsafe { self.persisted.defragment() }?; + { + let _pause = self.readers.pause()?; + // SAFETY: the caller excludes writers; the pause drains readers + // and prevents new index snapshots until relocation is complete. + 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. diff --git a/accountsdb/src/tests.rs b/accountsdb/src/tests.rs index 167789fb..a006ba09 100644 --- a/accountsdb/src/tests.rs +++ b/accountsdb/src/tests.rs @@ -53,12 +53,12 @@ fn in_volatile(db: &AccountsDB, pubkey: &Pubkey) -> bool { /// 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() + 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() + db.loader().read(pubkey, |account| account.lamports()).unwrap().unwrap() } /// Loads the account currently stored for `pubkey`. @@ -68,7 +68,9 @@ fn lamports(db: &AccountsDB, pubkey: &Pubkey) -> u64 { /// 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() + // SAFETY: these synchronous tests exclusively own the database and finish + // using borrowed results before relocation, deletion, or storage reuse. + unsafe { db.loader().load(pubkey) }.unwrap().unwrap() } /// Closes `pubkey`, deleting it from whichever backend currently holds it. @@ -145,8 +147,8 @@ fn test_routing_and_persistence_flips() { // 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_eq!(loader.read(&a, |acc| acc.lamports()).unwrap().unwrap(), 10); + assert_eq!(loader.read(&b, |acc| acc.lamports()).unwrap().unwrap(), 20); assert!(loader.contains(&a).unwrap() && loader.contains(&b).unwrap()); assert!(!loader.contains(&Pubkey::new_unique()).unwrap()); drop(loader); @@ -399,7 +401,7 @@ fn test_defragment_preserves_live_accounts() { // Every survivor still loads unchanged and remains program-indexed. for (k, lam) in &live { - let acc = db.loader().load(k).unwrap().unwrap(); + let acc = reload(&db, k); assert_eq!(acc.lamports(), *lam); assert_eq!(acc.owner(), &owner); } @@ -698,10 +700,7 @@ fn test_variable_sizes_and_exact_freelist() { defrag_to_stable(&db); for (k, data) in &live { - assert_eq!( - db.loader().load(k).unwrap().unwrap().data(), - data.as_slice() - ); + assert_eq!(reload(&db, k).data(), data.as_slice()); } } @@ -768,8 +767,153 @@ fn test_large_accounts_growth_and_defrag() { // Every survivor keeps its full 2 MiB image byte-for-byte. for (k, fill) in &live { - let acc = db.loader().load(k).unwrap().unwrap(); + let acc = reload(&db, k); assert_eq!(acc.data().len(), SIZE); assert!(acc.data().iter().all(|&b| b == *fill)); } } + +/// Proves snapshot relocation waits for a cached old index view, rejects a new +/// reader during that wait, and only then moves the account below a real EOF +/// truncation. Channel events force the schedule; timeouts only bound failures. +#[test] +fn test_snapshot_drains_readers_before_truncation() { + use std::{sync::mpsc, thread, time::Duration}; + + #[derive(Debug)] + enum Event { + Ready(u64), + Registered, + Draining, + Blocked, + Finished, + } + + const SIZE: usize = 2 << 20; + const WATCHDOG: Duration = Duration::from_secs(10); + let (dir, db) = db(); + let keys = [Pubkey::new_unique(), Pubkey::new_unique(), Pubkey::new_unique()]; + for key in keys { + store( + &db, + key, + delegated_account(1, vec![0x5a; SIZE], Pubkey::default()).build(), + ); + } + close(&db, &keys[0]); + close(&db, &keys[1]); + let key = keys[2]; + let original = offset(&db, &key); + let file = AccountsDB::directory(dir.path()).join(STORAGE_FILE); + let length = || std::fs::metadata(&file).unwrap().len(); + let before = length(); + let valid = |account: &AccountSharedData| { + account.data().len() == SIZE && account.data().iter().all(|&byte| byte == 0x5a) + }; + + let (events, observed) = mpsc::channel(); + let sender = events.clone(); + *db.readers.on_drain.lock() = Some(Box::new(move || { + sender.send(Event::Draining).unwrap(); + })); + let sender = events.clone(); + *db.readers.on_block.lock() = Some(Box::new(move || { + sender.send(Event::Blocked).unwrap(); + })); + + thread::scope(|scope| { + let db = &db; + let (release, resume) = mpsc::channel(); + let sender = events.clone(); + let old_reader = scope.spawn(move || { + let loader = db.loader(); + assert!(loader.contains(&key).unwrap()); + let old = db + .persisted + .index + .offset(&key, loader.txn.borrow().as_ref().unwrap()) + .unwrap() + .unwrap(); + // This is a lower bound on the file offset: it excludes the metadata + // header, so exceeding the new EOF proves the old image was removed. + let old_bytes = bytemuck::cast::<_, u64>(old) * solana_account::STORAGE_UNIT as u64; + sender.send(Event::Ready(old_bytes)).unwrap(); + // On a broken-barrier negative control, discard the stale transaction + // without dereferencing its truncated image. Failure is an assertion, + // not an intentional SIGBUS or undefined mapped-memory access. + resume + .recv_timeout(WATCHDOG) + .unwrap_or(false) + .then(|| loader.read(&key, valid).unwrap().unwrap()) + }); + + let (start, proceed) = mpsc::channel(); + let sender = events.clone(); + let new_reader = scope.spawn(move || { + // Exercise the registered-slot admission path, not first registration. + drop(db.loader()); + sender.send(Event::Registered).unwrap(); + proceed + .recv_timeout(WATCHDOG) + .unwrap_or(false) + .then(|| db.loader().read(&key, valid).unwrap().unwrap()) + }); + let mut old_bytes = None; + for _ in 0..2 { + match observed.recv_timeout(WATCHDOG).unwrap() { + Event::Ready(bytes) => old_bytes = Some(bytes), + Event::Registered => {} + event => panic!("unexpected reader setup event: {event:?}"), + } + } + let snapshot = scope.spawn(|| { + // SAFETY: setup writes are finished; all readers use guarded loaders. + let result = unsafe { db.snapshot(1) }; + events.send(Event::Finished).unwrap(); + result + }); + + let draining = matches!(observed.recv_timeout(WATCHDOG).unwrap(), Event::Draining); + start.send(draining).unwrap(); + let blocked = draining + && loop { + match observed.recv_timeout(WATCHDOG).unwrap() { + Event::Draining => continue, + Event::Blocked => break true, + Event::Finished => break false, + event => panic!("unexpected admission event: {event:?}"), + } + }; + let unchanged = length() == before && offset(db, &key) == original; + release.send(draining && blocked && unchanged).unwrap(); + let old_valid = old_reader.join().unwrap(); + let new_valid = new_reader.join().unwrap(); + snapshot.join().unwrap().unwrap(); + + assert!( + draining, + "snapshot bypassed the active reader instead of draining it" + ); + assert!(blocked, "a new reader must observe closed admission"); + assert!( + unchanged, + "relocation or truncation happened while the old view was live" + ); + assert_eq!(old_valid, Some(true)); + assert_eq!(new_valid, Some(true)); + assert!( + offset(db, &key) != original, + "the fixture must actually relocate the image" + ); + assert!( + length() < before, + "the fixture must actually truncate the file" + ); + // The retained pre-compaction offset would address a removed mmap page. + assert!( + old_bytes.unwrap() > length(), + "the stale image must lie beyond the new EOF" + ); + }); + assert_eq!(db.loader().read(&key, valid).unwrap(), Some(true)); +} diff --git a/engine/src/testkit.rs b/engine/src/testkit.rs index 3f88bd12..b94de6fb 100644 --- a/engine/src/testkit.rs +++ b/engine/src/testkit.rs @@ -119,7 +119,7 @@ impl TestEngine { /// Full committed account, or `None` when absent/closed. pub fn get_account(&self, key: Pubkey) -> Option { - self.engine.accounts().loader().load(&key).unwrap() + self.engine.accounts().loader().read(&key, Clone::clone).unwrap() } /// Executes instructions and returns the committed transaction result. diff --git a/keeper/README.md b/keeper/README.md index 469da2c5..6824688a 100644 --- a/keeper/README.md +++ b/keeper/README.md @@ -58,10 +58,11 @@ rejects a non-empty deployment whose configured authority account is absent. ## Superblock finalization -`Keeper::finalize_superblock` requires quiesced execution. It snapshots accountsdb -to refresh the checksum, then signs the reconstructed seal or compares it with -an authenticated upstream payload and retains its signature. The snapshot is -archived in the successor directory; completion acknowledges durable sealing +`Keeper::finalize_superblock` requires quiesced execution; accountsdb separately +drains scoped readers while relocating and truncating storage. It snapshots +accountsdb to refresh the checksum, then signs the reconstructed seal or compares +it with an authenticated upstream payload and retains its signature. The snapshot +is archived in the successor directory; completion acknowledges durable sealing and rotation, not archive completion. `SuperblockAccessor::sealed` returns the unsigned accountsdb state. @@ -88,7 +89,10 @@ delegated, ephemeral, and unresolved transient state remains outside it. Dedicated channels publish account and program updates, signature results, logs, processed transactions, blocks, cache evictions, completed snapshots, and -service messages. Signatures have terminal oneshot fanout; persistent multicast +service messages. Processed-transaction accounts are made owned before queueing, +so subscribers never retain mmap views across compaction. This copy is skipped +when there is no live processed-transaction subscriber. +Signatures have terminal oneshot fanout; persistent multicast streams give each receiver a bounded queue and disconnect a receiver that falls behind. Processed transactions, service messages, and cache evictions each have one process-lifetime receiver and apply producer backpressure when full. diff --git a/keeper/src/accessor.rs b/keeper/src/accessor.rs index 8b351703..61d81472 100644 --- a/keeper/src/accessor.rs +++ b/keeper/src/accessor.rs @@ -83,13 +83,13 @@ impl<'a> AccountsAccessor<'a> { pub fn update_sysvars(&self, block: Block) -> Result<()> { let loader = self.loader(); - let Some(mut hacc) = loader.load(&SlotHashes::id())? else { + let Some(mut hacc) = loader.read(&SlotHashes::id(), Clone::clone)? else { return Ok(()); }; let mut hashes: SlotHashes = hacc.deserialize_data().map_err(AccountsDBError::from)?; hashes.add(block.slot, block.hash); hacc.serialize_data(&hashes).map_err(AccountsDBError::from)?; - let Some(mut cacc) = loader.load(&Clock::id())? else { + let Some(mut cacc) = loader.read(&Clock::id(), Clone::clone)? else { return Ok(()); }; let mut clock: Clock = cacc.deserialize_data().map_err(AccountsDBError::from)?; @@ -219,10 +219,22 @@ impl<'a> TransactionsAccessor<'a> { } } while let Some(msg) = TlsManager::dequeue() { - subs.services.blocking_send(msg); + subs.services.blocking_send(|| msg); } } - subs.transactions.blocking_send(txn); + subs.transactions.blocking_send(|| { + // Subscriber queues outlive the execution barrier. Materialize + // borrowed accounts before handing off, while their mmap images + // are still protected by execution's account ownership. + if let Ok(execution) = &mut txn.execution.result { + for (_, account) in &mut execution.loaded_transaction.accounts { + if matches!(account.cow(), solana_account::CoWAccount::Borrowed(_)) { + *account = account.clone(); + } + } + } + txn + }); // Clear TLS unconditionally so unsent messages cannot leak into the next transaction. TlsManager::clear(); subs.signatures.send(&commit.signature, &commit.status); diff --git a/keeper/src/builder.rs b/keeper/src/builder.rs index f8775225..8f865348 100644 --- a/keeper/src/builder.rs +++ b/keeper/src/builder.rs @@ -101,7 +101,7 @@ impl KeeperBuilder { self.seed_programs(&mut accounts)?; let caches = self.seed_sysvars(accountsdb, ledger, &mut accounts).await?; let authority = self.authority.pubkey(); - if accountsdb.loader().load(&authority)?.is_none() { + if !accountsdb.loader().contains(&authority)? { let sponsor = AccountBuilder::default() .lamports(SPONSOR_INIT_BALANCE) .mode(AccountMode::Ephemeral); @@ -182,14 +182,14 @@ impl KeeperBuilder { accounts: &mut Vec, ) -> Result { let slot = accountsdb.slot(); - let loader = accountsdb.loader(); let id = SlotHashes::id(); // AccountsDB starts at slot 1 and the ledger starts at head 1. let genesis = slot == 1 && ledger.head() == 1; - let slothashes = match loader.load(&id)? { - Some(account) => { - account.deserialize_data::().map_err(AccountsDBError::from)? - } + let slothashes = accountsdb + .loader() + .read(&id, |account| account.deserialize_data::())?; + let slothashes = match slothashes { + Some(hashes) => hashes.map_err(AccountsDBError::from)?, None => { // Keep the sysvar account at its fixed serialized capacity so live // updates can replace entries without resizing the account. diff --git a/keeper/src/lib.rs b/keeper/src/lib.rs index 8a97644c..ed7ce17b 100644 --- a/keeper/src/lib.rs +++ b/keeper/src/lib.rs @@ -149,9 +149,9 @@ impl Keeper { /// Snapshot the current superblock, enqueue its seal, and archive the snapshot. /// - /// Must run only when no account store can race the snapshot export; the - /// in-body `SAFETY` note relies on this exclusivity. The returned signal - /// resolves after the appender durably seals and rotates the ledger. + /// Must run only when account writes and unscoped borrowed views are + /// quiesced. Accountsdb drains scoped readers during compaction. The returned + /// signal resolves after the appender durably seals and rotates the ledger. /// An upstream seal must already be authenticated; its payload must match /// the snapshot state, and its signature is retained unchanged. pub fn finalize_superblock( @@ -160,8 +160,8 @@ impl Keeper { ) -> Result> { let _timer = metrics::time(Operation::FinalizeSuperblock); let head = self.ledger.head(); - // SAFETY: the caller must ensure exclusive account-store - // access. Snapshot also publishes the checksum for this superblock id; + // SAFETY: the caller excludes writes and unscoped borrowed views; + // accountsdb drains scoped readers. Snapshot also publishes the checksum; // read/sign/compare it only after that refresh. let snapshot = unsafe { self.accountsdb.snapshot(head) }?; let payload = self.superblocks().sealed(); @@ -226,7 +226,7 @@ impl Keeper { pub fn apply_reset(&self, reset: Reset) -> Result<()> { self.accountsdb.reset(); let authority = self.authority(); - let account = self.accounts().loader().load(&authority)?; + let account = self.accounts().loader().read(&authority, Clone::clone)?; if let Some(account) = account { let acc = AccountBuilder::from(account).lamports(SPONSOR_INIT_BALANCE); self.accounts().store(&[(authority, acc.build())])?; diff --git a/keeper/src/subscriptions.rs b/keeper/src/subscriptions.rs index 1fe11373..95196865 100644 --- a/keeper/src/subscriptions.rs +++ b/keeper/src/subscriptions.rs @@ -154,12 +154,13 @@ impl Unicast { let _ = sender.send(value).await; } - /// Sends from a synchronous worker, waiting until the receiver has capacity. - pub(crate) fn blocking_send(&self, value: V) { - let Some(sender) = self.sender.get() else { - return; - }; - let _ = sender.blocking_send(value); + /// Prepares a value only for a live receiver, then waits for queue capacity. + pub(crate) fn blocking_send(&self, prepare: impl FnOnce() -> V) { + if let Some(sender) = self.sender.get() + && !sender.is_closed() + { + let _ = sender.blocking_send(prepare()); + } } } diff --git a/keeper/src/tests/recovery.rs b/keeper/src/tests/recovery.rs index c165a60f..c41df161 100644 --- a/keeper/src/tests/recovery.rs +++ b/keeper/src/tests/recovery.rs @@ -54,7 +54,7 @@ async fn seeds_features_programs_and_sysvars() { } for (&id, &slot) in keeper.features().active() { assert_eq!(slot, 0, "features activate at slot 0"); - let acc = loader.load(&id).unwrap().expect("feature account seeded"); + let acc = loader.read(&id, Clone::clone).unwrap().expect("feature account seeded"); assert_eq!(acc.owner(), &solana_feature_gate_interface::ID); assert!(acc.lamports() >= rent.minimum_balance(acc.data().len())); } @@ -63,7 +63,7 @@ async fn seeds_features_programs_and_sysvars() { // owned by loader_v4 (not the BPF upgradeable loader), and rent-exempt. // Builtins are seeded through the same path with an executable native-loader // account, so they share this shape. - let acc = loader.load(&program).unwrap().expect("program seeded"); + let acc = loader.read(&program, Clone::clone).unwrap().expect("program seeded"); assert!(acc.executable()); assert_eq!(acc.owner(), &loader_v4::ID); assert_eq!(acc.data(), elf.as_slice()); @@ -72,7 +72,7 @@ async fn seeds_features_programs_and_sysvars() { // The Clock is seeded one slot ahead of the last block; a fresh ledger's last // block defaults to slot 0, so the clock starts at slot 1. let clock: Clock = loader - .load(&Clock::id()) + .read(&Clock::id(), Clone::clone) .unwrap() .expect("clock seeded") .deserialize_data() @@ -81,7 +81,7 @@ async fn seeds_features_programs_and_sysvars() { // Rent and EpochSchedule sysvars are present and sysvar-owned. for id in [Rent::id(), EpochSchedule::id()] { - let acc = loader.load(&id).unwrap().expect("sysvar seeded"); + let acc = loader.read(&id, Clone::clone).unwrap().expect("sysvar seeded"); assert_eq!(acc.owner(), &sysvar::ID); } drop(loader); @@ -123,7 +123,12 @@ async fn recovers_the_newest_snapshot() { let keeper = TestKeeper::from_builder(dirs, builder).await; keeper.accounts().validate().expect("restored store validates"); - let restored = keeper.accounts().loader().load(&marker).unwrap().expect("marker restored"); + let restored = keeper + .accounts() + .loader() + .read(&marker, Clone::clone) + .unwrap() + .expect("marker restored"); assert_eq!(restored.lamports(), 2, "newest snapshot wins"); // The corrupt tree saved for inspection is removed on successful recovery. assert!(!keeper.dirs.accounts.path().join("CURRENT.bkp").exists()); diff --git a/keeper/src/tests/subscriptions.rs b/keeper/src/tests/subscriptions.rs index ee063695..c854d06c 100644 --- a/keeper/src/tests/subscriptions.rs +++ b/keeper/src/tests/subscriptions.rs @@ -30,7 +30,7 @@ async fn subscribers_send_semantics() { unicast.send(3).await; let sender = unicast.clone(); - let send = std::thread::spawn(move || sender.blocking_send(4)); + let send = std::thread::spawn(move || sender.blocking_send(|| 4)); assert_eq!(unicast_rx.recv().await, Some(3)); send.join().unwrap(); assert_eq!(unicast_rx.recv().await, Some(4)); diff --git a/processor/README.md b/processor/README.md index 254c0564..8069c869 100644 --- a/processor/README.md +++ b/processor/README.md @@ -34,6 +34,12 @@ coherent superblock snapshots, replay seal checks, replication handshakes, and shutdown. A superblock checkpoint finalizes its block and enters that pause as one sequencer message, so later transactions cannot enter the sealed snapshot. +Transaction execution uses an unguarded AccountsDB loader: the sequencer already +excludes compaction through execution, commit, and owned subscription fanout. +This avoids reader registration, slot updates, and admission fences on the +transaction-loading path. Simulation retains guarded reads because it runs +independently of the sequencer barrier. + ## Simulation Simulation has a separate worker and SVM context. It resolves a transaction, diff --git a/processor/src/callback.rs b/processor/src/callback.rs index fa5b0b9c..d654192a 100644 --- a/processor/src/callback.rs +++ b/processor/src/callback.rs @@ -21,16 +21,18 @@ pub(crate) struct LoadCallback<'a, const LOAD_OWNED: bool> { impl TransactionProcessingCallback for LoadCallback<'_, LOAD_OWNED> { fn get_account_shared_data(&self, pubkey: &Pubkey) -> Option<(AccountSharedData, Slot)> { - self.loader - .load(pubkey) + let account = if LOAD_OWNED { + self.loader.read(pubkey, Clone::clone) + } else { + // SAFETY: execution owns the account until commit/fanout and is + // drained before compaction. Startup cache seeding is exclusive. + // Borrowed results are made owned before leaving that boundary. + unsafe { self.loader.load(pubkey) } + }; + account .inspect_err(|error| error!(?error, "accountsdb load error")) .unwrap_or_default() - .map(|mut acc| { - if LOAD_OWNED { - acc = acc.owned().into(); - } - (acc, 0) - }) + .map(|acc| (acc, 0)) } } diff --git a/processor/src/executor.rs b/processor/src/executor.rs index a11dfa0b..9ef4d43f 100644 --- a/processor/src/executor.rs +++ b/processor/src/executor.rs @@ -8,6 +8,7 @@ use std::{ thread::{self, JoinHandle}, }; +use accountsdb::AccountLoader; use keeper::{ExecutionRecord, FullTransaction, Keeper, ResolvedTransaction}; use nucleus::{ ledger::Block, @@ -169,7 +170,11 @@ impl TransactionExecutor { /// its raw state transition (replay) or full execution. fn process(&mut self, txn: ResolvedTransaction) -> Result<()> { let accounts = self.state.accounts(); - let output = self.svm.execute::(accounts.loader(), &txn, self.state.features()); + // SAFETY: the sequencer cannot acknowledge compaction until this + // executor completes processing, including commit and owned fanout. + // Replay follows the same drain protocol; no borrowed result escapes. + let loader = unsafe { AccountLoader::unguarded(&accounts) }; + let output = self.svm.execute::(loader, &txn, self.state.features()); if !output.processing_result.was_processed_with_successful_result() { metrics::failed_transaction(FailureKind::Execution); } diff --git a/processor/src/tests.rs b/processor/src/tests.rs index 3f275170..31b79b9a 100644 --- a/processor/src/tests.rs +++ b/processor/src/tests.rs @@ -218,7 +218,12 @@ async fn block_hash_includes_appended_transaction_signature() { #[tokio::test(flavor = "current_thread")] async fn seeded_program_and_recursive_cpi_return_data_work() { let harness = Harness::new(false).await; - let account = harness.accounts().loader().load(&V42_ID).unwrap().expect("v42 program seeded"); + let account = harness + .accounts() + .loader() + .read(&V42_ID, Clone::clone) + .unwrap() + .expect("v42 program seeded"); assert!(account.executable()); assert_eq!(*account.owner(), loader_v4::ID); From 05b8010133efda117bf456d4c04390208a0e7fc0 Mon Sep 17 00:00:00 2001 From: Babur Makhmudov Date: Wed, 16 Sep 2026 14:37:22 +0700 Subject: [PATCH 2/2] docs: clarify lazy subscription preparation guarantee --- keeper/src/subscriptions.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/keeper/src/subscriptions.rs b/keeper/src/subscriptions.rs index 95196865..c98834a3 100644 --- a/keeper/src/subscriptions.rs +++ b/keeper/src/subscriptions.rs @@ -154,7 +154,10 @@ impl Unicast { let _ = sender.send(value).await; } - /// Prepares a value only for a live receiver, then waits for queue capacity. + /// Prepares and sends a value, waiting for queue capacity. + /// + /// Skips preparation if no sender exists or it is observed closed. + /// The receiver may close after this check. pub(crate) fn blocking_send(&self, prepare: impl FnOnce() -> V) { if let Some(sender) = self.sender.get() && !sender.is_closed()