Skip to main content

miden_node_store/db/
mod.rs

1use std::collections::{BTreeMap, BTreeSet, HashSet};
2use std::mem::size_of;
3use std::num::NonZeroUsize;
4use std::ops::{Deref, DerefMut};
5use std::path::{Path, PathBuf};
6use std::sync::Arc;
7
8use anyhow::Context;
9use miden_node_db::sqlite::{DbReader, DbWriter, WriteTx};
10use miden_node_proto::domain::account::AccountInfo;
11use miden_node_tracing::{info, miden_instrument, warn};
12use miden_node_utils::limiter::{
13    MAX_RESPONSE_PAYLOAD_BYTES,
14    QueryParamLimiter,
15    QueryParamNoteCommitmentLimit,
16};
17use miden_protocol::Word;
18use miden_protocol::account::{AccountHeader, AccountId, AccountStorageHeader, StorageMapKey};
19use miden_protocol::asset::{Asset, AssetId};
20use miden_protocol::block::{
21    BlockAccountUpdate,
22    BlockHeader,
23    BlockNoteIndex,
24    BlockNumber,
25    BlockSignatures,
26    SignedBlock,
27};
28use miden_protocol::crypto::merkle::SparseMerklePath;
29use miden_protocol::note::{
30    NoteAttachments,
31    NoteDetails,
32    NoteId,
33    NoteInclusionProof,
34    NoteMetadata,
35    NoteScript,
36    Nullifier,
37};
38use miden_protocol::protocol_config::ProtocolConfig;
39use miden_protocol::transaction::TransactionHeader;
40use miden_protocol::utils::serde::Deserializable;
41
42use crate::db::migrations::{migrate_database, verify_latest_schema};
43use crate::db::models::conv::SqlTypeConvert;
44use crate::db::models::queries as diesel_queries;
45use crate::db::models::queries::StorageMapValuesPage;
46pub use crate::db::models::queries::{
47    AccountCommitmentsPage,
48    NullifiersPage,
49    PublicAccountIdsPage,
50    PublicAccountStateRootsPage,
51};
52pub use crate::db::queries::{
53    HISTORICAL_BLOCK_RETENTION,
54    PrecomputedPublicAccountState,
55    PrecomputedPublicAccountStates,
56};
57use crate::errors::{DatabaseError, NoteSyncError};
58use crate::genesis::GenesisBlock;
59use crate::state::{ScopedBlockNum, ScopedBlockRange};
60use crate::{COMPONENT, LOG_TARGET};
61
62const STORAGE_MAP_VALUE_PER_ROW_BYTES: usize =
63    2 * size_of::<Word>() + size_of::<u32>() + size_of::<u8>();
64
65fn default_storage_map_entries_limit() -> usize {
66    MAX_RESPONSE_PAYLOAD_BYTES / STORAGE_MAP_VALUE_PER_ROW_BYTES
67}
68
69mod migrations;
70#[cfg(test)]
71pub(crate) use migrations::bootstrap_database;
72
73#[cfg(test)]
74mod tests;
75
76#[cfg(test)]
77mod test_db;
78#[cfg(test)]
79pub(crate) use test_db::TestDb;
80
81/// Query functions on the `miden-node-db` SQLite framework.
82///
83/// All writes run here; reads are migrated from [`models`] incrementally.
84pub(crate) mod queries;
85
86mod utils;
87
88pub(crate) mod models;
89
90/// [diesel](https://diesel.rs) generated schema
91///
92/// The ignored `diesel_schema_is_in_sync_with_migrations` test verifies that this file matches the
93/// schema produced by the current migrations.
94pub(crate) mod schema;
95
96pub type Result<T, E = DatabaseError> = std::result::Result<T, E>;
97
98/// Database options used by the store state.
99#[derive(Copy, Clone, Debug, PartialEq, Eq)]
100pub struct DatabaseOptions {
101    /// Maximum number of SQLite connections in the connection pool.
102    pub connection_pool_size: NonZeroUsize,
103}
104
105impl Default for DatabaseOptions {
106    fn default() -> Self {
107        Self {
108            connection_pool_size: miden_node_db::default_connection_pool_size(),
109        }
110    }
111}
112
113/// The Store's database.
114///
115/// Extends the underlying [`miden_node_db::Db`] type with functionality specific to the Store.
116///
117/// The store is mid-migration to the `miden-node-db` SQLite framework: every write serializes on
118/// the single framework writer connection, while most reads still run on the diesel pool. Reads
119/// move to the framework reader pool one batch at a time until the diesel pool is removed.
120pub struct Db {
121    diesel: miden_node_db::Db,
122    writer: DbWriter,
123    reader: DbReader,
124}
125
126/// Inserts the genesis block and the protocol configuration that it activates.
127fn insert_genesis(tx: &WriteTx<'_>, genesis: GenesisBlock) -> Result<()> {
128    let (genesis_block, protocol_config) = genesis.into_parts();
129    // The genesis block has no transactions, but it creates every account it contains.
130    let new_account_ids = genesis_block
131        .body()
132        .updated_accounts()
133        .iter()
134        .map(BlockAccountUpdate::account_id)
135        .collect();
136    queries::insert_protocol_config(tx, &protocol_config, BlockNumber::GENESIS)?;
137    queries::apply_block(
138        tx,
139        &genesis_block,
140        &[],
141        &PrecomputedPublicAccountStates::new(),
142        &new_account_ids,
143    )?;
144    Ok(())
145}
146
147impl Deref for Db {
148    type Target = miden_node_db::Db;
149
150    fn deref(&self) -> &Self::Target {
151        &self.diesel
152    }
153}
154
155impl DerefMut for Db {
156    fn deref_mut(&mut self) -> &mut Self::Target {
157        &mut self.diesel
158    }
159}
160
161/// The commitment of a [`BlockHeader`], stored alongside the header it belongs to.
162///
163/// Keeping it in its own column lets the chain MMR be rebuilt at startup without deserializing
164/// every header.
165#[derive(Debug, Clone, Copy, PartialEq, Eq)]
166#[repr(transparent)]
167pub struct BlockHeaderCommitment(pub(crate) Word);
168
169impl BlockHeaderCommitment {
170    pub fn new(header: &BlockHeader) -> Self {
171        Self(header.commitment())
172    }
173
174    pub fn word(self) -> Word {
175        self.0
176    }
177}
178
179/// Describes the value of an asset for an account ID at `block_num` specifically.
180///
181/// If `asset` is `None`, the asset was removed.
182#[derive(Debug, Clone)]
183pub struct AccountVaultValue {
184    pub block_num: BlockNumber,
185    pub vault_key: AssetId,
186    /// None if the asset was removed
187    pub asset: Option<Asset>,
188}
189
190impl AccountVaultValue {
191    pub fn from_raw_row(row: (i64, Vec<u8>, Option<Vec<u8>>)) -> Result<Self, DatabaseError> {
192        let (block_num, vault_key, asset) = row;
193        let vault_key = Word::read_from_bytes(&vault_key)?;
194        Ok(Self {
195            block_num: BlockNumber::from_raw_sql(block_num)?,
196            vault_key: AssetId::try_from(vault_key)?,
197            asset: asset.map(|b| miden_node_persistence::decode::<Asset>(&b)).transpose()?,
198        })
199    }
200}
201
202#[derive(Debug, PartialEq)]
203pub struct NullifierInfo {
204    pub nullifier: Nullifier,
205    pub block_num: BlockNumber,
206}
207
208impl PartialEq<(Nullifier, BlockNumber)> for NullifierInfo {
209    fn eq(&self, (nullifier, block_num): &(Nullifier, BlockNumber)) -> bool {
210        &self.nullifier == nullifier && &self.block_num == block_num
211    }
212}
213
214#[derive(Debug, PartialEq)]
215pub struct TransactionRecord {
216    pub block_num: BlockNumber,
217    pub header: TransactionHeader,
218    /// Inclusion proofs for committed output notes. Notes in `header.output_notes()` without a
219    /// corresponding proof here were erased (created and consumed within the same batch).
220    pub output_note_proofs: Vec<NoteSyncRecord>,
221    /// Maps each consumed input note's nullifier to its note ID, for public notes the node could
222    /// resolve. This is to enable the recover of notes by their id.
223    pub consumed_note_refs: Vec<(Nullifier, NoteId)>,
224}
225
226#[derive(Debug, Clone, PartialEq)]
227pub struct NoteRecord {
228    pub block_num: BlockNumber,
229    pub note_index: BlockNoteIndex,
230    pub note_id: Word,
231    pub metadata: NoteMetadata,
232    pub details: Option<NoteDetails>,
233    pub attachments: NoteAttachments,
234    pub inclusion_path: SparseMerklePath,
235}
236
237#[derive(Debug, PartialEq)]
238pub struct NoteSyncUpdate {
239    pub notes: Vec<NoteSyncRecord>,
240    pub block_header: BlockHeader,
241}
242
243#[derive(Debug, Clone, PartialEq)]
244pub struct NoteSyncRecord {
245    pub block_num: BlockNumber,
246    pub note_index: BlockNoteIndex,
247    pub note_id: NoteId,
248    pub metadata: NoteMetadata,
249    pub attachments: NoteAttachments,
250    pub inclusion_path: SparseMerklePath,
251}
252
253impl From<NoteRecord> for NoteSyncRecord {
254    fn from(note: NoteRecord) -> Self {
255        Self {
256            block_num: note.block_num,
257            note_index: note.note_index,
258            note_id: NoteId::from_raw(note.note_id),
259            metadata: note.metadata,
260            attachments: note.attachments,
261            inclusion_path: note.inclusion_path,
262        }
263    }
264}
265
266impl Db {
267    /// Creates a new database and inserts the genesis block.
268    #[miden_instrument(
269        target = COMPONENT,
270        name = "store.database.bootstrap",
271        fields(path = database_filepath),
272        err,
273    )]
274    pub async fn bootstrap(
275        database_filepath: PathBuf,
276        genesis: GenesisBlock,
277    ) -> anyhow::Result<()> {
278        migrations::bootstrap_database(&database_filepath)
279            .context("failed to bootstrap database schema")?;
280
281        let (writer, _reader) = miden_node_db::sqlite::open(&database_filepath)
282            .context("failed to open a database connection")?;
283
284        // Insert genesis block data.
285        writer
286            .write("insert genesis block", move |tx| insert_genesis(tx, genesis))
287            .await
288            .context("failed to insert genesis block")?;
289        Ok(())
290    }
291
292    /// Open a connection to the DB after verifying that it is at the latest schema version.
293    #[miden_instrument(
294        target = COMPONENT,
295    )]
296    pub async fn load(database_filepath: PathBuf) -> Result<Self, DatabaseError> {
297        Self::load_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size())
298            .await
299    }
300
301    /// Open a connection to the DB with a specific pool size after verifying that it is at the
302    /// latest schema version.
303    #[miden_instrument(
304        target = COMPONENT,
305    )]
306    pub async fn load_with_pool_size(
307        database_filepath: PathBuf,
308        connection_pool_size: NonZeroUsize,
309    ) -> Result<Self, DatabaseError> {
310        verify_latest_schema(&database_filepath)?;
311
312        let db = miden_node_db::Db::new_with_pool_size(&database_filepath, connection_pool_size)?;
313        let (writer, reader) =
314            miden_node_db::sqlite::open_with_pool_size(&database_filepath, connection_pool_size)?;
315        info!(
316            target: LOG_TARGET,
317            "Connected to the database",
318            path = database_filepath,
319            db.sqlite.connection_pool_size = connection_pool_size.get()
320        );
321
322        Ok(Self { diesel: db, writer, reader })
323    }
324
325    /// The write handle, for tests that need to seed or corrupt rows no production method writes.
326    #[cfg(test)]
327    pub(crate) fn writer(&self) -> &DbWriter {
328        &self.writer
329    }
330
331    /// Selects a protocol configuration by its commitment.
332    #[miden_instrument(
333        level = "debug",
334        target = COMPONENT,
335        err,
336    )]
337    pub async fn select_protocol_config_by_commitment(
338        &self,
339        commitment: Word,
340    ) -> Result<Option<ProtocolConfig>> {
341        self.transact("protocol config by commitment", move |conn| {
342            diesel_queries::select_protocol_config_by_commitment(conn, commitment)
343        })
344        .await
345    }
346
347    /// Selects the configuration commitment active at the specified block.
348    pub async fn select_protocol_config_commitment_at(
349        &self,
350        block_number: ScopedBlockNum,
351    ) -> Result<Option<Word>> {
352        self.transact("protocol config commitment at block", move |conn| {
353            diesel_queries::select_protocol_config_commitment_at(conn, *block_number)
354        })
355        .await
356    }
357
358    /// Applies all pending migrations to an existing DB.
359    #[miden_instrument(
360        target = COMPONENT,
361    )]
362    pub fn migrate(database_filepath: impl AsRef<Path>) -> Result<(), DatabaseError> {
363        migrate_database(database_filepath.as_ref())?;
364        Ok(())
365    }
366
367    /// Returns a page of nullifiers for tree rebuilding.
368    #[miden_instrument(
369        level = "debug",
370        target = COMPONENT,
371        err,
372    )]
373    pub async fn select_nullifiers_paged(
374        &self,
375        page_size: std::num::NonZeroUsize,
376        after_nullifier: Option<Nullifier>,
377    ) -> Result<NullifiersPage> {
378        self.transact("read nullifiers paged", move |conn| {
379            diesel_queries::select_nullifiers_paged(conn, page_size, after_nullifier)
380        })
381        .await
382    }
383
384    /// Loads the nullifiers that match the prefixes from the DB.
385    #[miden_instrument(
386        level = "debug",
387        target = COMPONENT,
388        fields(
389            prefix_len,
390            prefix.count = nullifier_prefixes.len(),
391        ),
392        err,
393    )]
394    pub async fn select_nullifiers_by_prefix(
395        &self,
396        prefix_len: u32,
397        nullifier_prefixes: Vec<u32>,
398        block_range: ScopedBlockRange,
399    ) -> Result<(Vec<NullifierInfo>, BlockNumber)> {
400        let block_range = block_range.into_inner();
401        assert_eq!(prefix_len, 16, "Only 16-bit prefixes are supported");
402
403        self.transact("nullifieres by prefix", move |conn| {
404            let nullifier_prefixes =
405                nullifier_prefixes.into_iter().map(|prefix| prefix as u16).collect::<Vec<_>>();
406            diesel_queries::select_nullifiers_by_prefix(
407                conn,
408                prefix_len as u8,
409                &nullifier_prefixes[..],
410                block_range,
411            )
412        })
413        .await
414    }
415
416    /// Search for a [`BlockHeader`] from the database by its `block_num`.
417    ///
418    /// When `block_number` is [None], the latest block header is returned.
419    #[miden_instrument(
420        level = "debug",
421        target = COMPONENT,
422        err,
423    )]
424    pub async fn select_block_header_by_block_num(
425        &self,
426        maybe_block_number: Option<ScopedBlockNum>,
427    ) -> Result<Option<BlockHeader>> {
428        self.transact("block headers by block number", move |conn| {
429            let val = diesel_queries::select_block_header_by_block_num(
430                conn,
431                maybe_block_number.map(|block_number| *block_number),
432            )?;
433            Ok(val)
434        })
435        .await
436    }
437
438    /// Selects the genesis block header for state initialization.
439    pub(crate) async fn select_genesis_block_header(&self) -> Result<Option<BlockHeader>> {
440        self.transact("genesis block header", |conn| {
441            diesel_queries::select_block_header_by_block_num(conn, Some(BlockNumber::GENESIS))
442        })
443        .await
444    }
445
446    /// Search for a [`BlockHeader`] and its [`BlockSignatures`] from the database by its
447    /// `block_num`.
448    #[miden_instrument(
449        level = "debug",
450        target = COMPONENT,
451        err,
452    )]
453    pub async fn select_block_header_and_signatures_by_block_num(
454        &self,
455        block_number: ScopedBlockNum,
456    ) -> Result<Option<(BlockHeader, BlockSignatures)>> {
457        self.transact("block headers and signatures by block number", move |conn| {
458            let val = diesel_queries::select_block_header_and_signatures_by_block_num(
459                conn,
460                *block_number,
461            )?;
462            Ok(val)
463        })
464        .await
465    }
466
467    /// Loads multiple block headers from the DB.
468    #[miden_instrument(
469        level = "debug",
470        target = COMPONENT,
471        err,
472    )]
473    pub async fn select_block_headers(
474        &self,
475        blocks: impl Iterator<Item = ScopedBlockNum> + Send + 'static,
476    ) -> Result<Vec<BlockHeader>> {
477        self.transact("block headers from given block numbers", move |conn| {
478            let raw = diesel_queries::select_block_headers(conn, blocks.map(|block| *block))?;
479            Ok(raw)
480        })
481        .await
482    }
483
484    /// Loads all the block headers from the DB.
485    #[miden_instrument(
486        level = "debug",
487        target = COMPONENT,
488        err,
489    )]
490    pub async fn select_all_block_header_commitments(&self) -> Result<Vec<BlockHeaderCommitment>> {
491        self.transact("all block headers", |conn| {
492            let raw = diesel_queries::select_all_block_header_commitments(conn)?;
493            Ok(raw)
494        })
495        .await
496    }
497
498    /// Returns a page of account commitments for tree rebuilding.
499    #[miden_instrument(
500        level = "debug",
501        target = COMPONENT,
502        err,
503    )]
504    pub async fn select_account_commitments_paged(
505        &self,
506        page_size: std::num::NonZeroUsize,
507        after_account_id: Option<AccountId>,
508    ) -> Result<AccountCommitmentsPage> {
509        self.transact("read account commitments paged", move |conn| {
510            diesel_queries::select_account_commitments_paged(conn, page_size, after_account_id)
511        })
512        .await
513    }
514
515    /// Returns a page of public account IDs for forest rebuilding.
516    #[miden_instrument(
517        level = "debug",
518        target = COMPONENT,
519        err,
520    )]
521    pub async fn select_public_account_ids_paged(
522        &self,
523        page_size: std::num::NonZeroUsize,
524        after_account_id: Option<AccountId>,
525    ) -> Result<PublicAccountIdsPage> {
526        self.transact("read public account IDs paged", move |conn| {
527            diesel_queries::select_public_account_ids_paged(conn, page_size, after_account_id)
528        })
529        .await
530    }
531
532    /// Returns a page of public account state roots for forest consistency verification.
533    #[miden_instrument(
534        level = "debug",
535        target = COMPONENT,
536        err,
537    )]
538    pub async fn select_public_account_state_roots_paged(
539        &self,
540        page_size: std::num::NonZeroUsize,
541        after_account_id: Option<AccountId>,
542    ) -> Result<PublicAccountStateRootsPage> {
543        self.transact("read public account state roots paged", move |conn| {
544            diesel_queries::select_public_account_state_roots_paged(
545                conn,
546                page_size,
547                after_account_id,
548            )
549        })
550        .await
551    }
552
553    /// Loads public account details from the DB.
554    #[miden_instrument(
555        level = "debug",
556        target = COMPONENT,
557        err,
558    )]
559    pub async fn select_account(&self, id: AccountId) -> Result<AccountInfo> {
560        self.transact("Get account details", move |conn| diesel_queries::select_account(conn, id))
561            .await
562    }
563
564    /// Returns the subset of the provided account IDs that classify as network accounts.
565    #[miden_instrument(
566        level = "debug",
567        target = COMPONENT,
568        err,
569    )]
570    pub async fn filter_network_accounts(
571        &self,
572        account_ids: Vec<AccountId>,
573    ) -> Result<HashSet<AccountId>> {
574        self.reader
575            .read("Filter network accounts", move |tx| {
576                queries::filter_network_accounts(tx, &account_ids)
577            })
578            .await
579    }
580
581    /// Queries the account code by its commitment hash.
582    ///
583    /// Returns `None` if no code exists with that commitment.
584    #[miden_instrument(
585        target = COMPONENT,
586    )]
587    pub async fn select_account_code_by_commitment(
588        &self,
589        code_commitment: Word,
590    ) -> Result<Option<miden_protocol::account::AccountCode>> {
591        self.transact("Get account code by commitment", move |conn| {
592            diesel_queries::select_account_code_by_commitment(conn, code_commitment)?
593                .map(|bytes| {
594                    miden_node_persistence::decode::<miden_protocol::account::AccountCode>(&bytes)
595                })
596                .transpose()
597                .map_err(DatabaseError::from)
598        })
599        .await
600    }
601
602    /// Queries the account header and storage header for a specific account at a block.
603    ///
604    /// Returns both in a single query to avoid querying the database twice.
605    /// Returns `None` if the account doesn't exist at that block.
606    #[miden_instrument(
607        target = COMPONENT,
608    )]
609    pub async fn select_account_header_with_storage_header_at_block(
610        &self,
611        account_id: AccountId,
612        block_num: ScopedBlockNum,
613    ) -> Result<Option<(AccountHeader, AccountStorageHeader)>> {
614        self.reader
615            .read("Get account header with storage header at block", move |tx| {
616                queries::select_account_header_with_storage_header_at_block(
617                    tx, account_id, *block_num,
618                )
619            })
620            .await
621    }
622
623    #[miden_instrument(
624        level = "debug",
625        target = COMPONENT,
626        err,
627    )]
628    pub async fn get_note_sync_multi(
629        &self,
630        block_range: ScopedBlockRange,
631        note_tags: Arc<[u32]>,
632    ) -> Result<Vec<NoteSyncUpdate>, NoteSyncError> {
633        let block_range = block_range.into_inner();
634        self.transact("notes sync task", move |conn| {
635            diesel_queries::get_note_sync_multi(
636                conn,
637                &note_tags,
638                block_range,
639                MAX_RESPONSE_PAYLOAD_BYTES,
640            )
641        })
642        .await
643    }
644
645    /// Loads all the [`miden_protocol::note::Note`]s matching a certain [`NoteId`] from the
646    /// database.
647    #[miden_instrument(
648        level = "debug",
649        target = COMPONENT,
650        err,
651    )]
652    pub async fn select_notes_by_id(&self, note_ids: Vec<NoteId>) -> Result<Vec<NoteRecord>> {
653        self.transact("note by id", move |conn| {
654            diesel_queries::select_notes_by_id(conn, note_ids.as_slice())
655        })
656        .await
657    }
658
659    /// Returns the requested note IDs that the database contains at or before `up_to_block`.
660    #[miden_instrument(
661        level = "debug",
662        target = COMPONENT,
663        err,
664    )]
665    pub async fn select_existing_note_ids(
666        &self,
667        note_ids: Vec<NoteId>,
668        up_to_block: ScopedBlockNum,
669    ) -> Result<HashSet<NoteId>> {
670        self.transact("existing note IDs", move |conn| {
671            diesel_queries::select_existing_note_ids(conn, note_ids.as_slice(), *up_to_block)
672        })
673        .await
674    }
675
676    /// Loads inclusion proofs for notes matching the given note commitments that were committed at
677    /// or before `up_to_block`.
678    #[miden_instrument(
679        level = "debug",
680        target = COMPONENT,
681        err,
682    )]
683    pub async fn select_note_inclusion_proofs(
684        &self,
685        note_commitments: BTreeSet<Word>,
686        up_to_block: ScopedBlockNum,
687    ) -> Result<BTreeMap<NoteId, NoteInclusionProof>> {
688        self.transact("block note inclusion proofs by commitment", move |conn| {
689            diesel_queries::select_note_inclusion_proofs(conn, &note_commitments, *up_to_block)
690        })
691        .await
692    }
693
694    /// Inserts the data of a new block into the DB.
695    ///
696    /// The transaction is committed when this method returns. Synchronization with the in-memory
697    /// trees is handled by the block writer task; see [`super::state::State::apply_block`].
698    ///
699    /// Account history is pruned in the same transaction against `prune_tip`: the effective tip
700    /// for retention, which lags the actual tip while old snapshot generations are still pinned
701    /// by readers (SQLite reads have no point-in-time protection, unlike the `RocksDB`-backed
702    /// trees).
703    ///
704    /// Consumed note IDs omitted from transaction headers are resolved from
705    /// `unresolved_note_nullifiers` on a best-effort basis. The returned mapping is used only for
706    /// lifecycle events and never affects block application. `unresolved_note_nullifiers` is empty
707    /// when neither INFO nor DEBUG lifecycle events are enabled.
708    // TODO: This span is logged in a root span, we should connect it to the parent one.
709    #[expect(
710        clippy::too_many_arguments,
711        reason = "the arguments are the block and the state that the writer precomputed for it"
712    )]
713    #[miden_instrument(
714        target = COMPONENT,
715        err,
716    )]
717    pub(crate) async fn apply_block(
718        &self,
719        signed_block: SignedBlock,
720        activated_protocol_config: Option<ProtocolConfig>,
721        notes: Vec<(NoteRecord, Option<Nullifier>)>,
722        precomputed_public_states: PrecomputedPublicAccountStates,
723        new_account_ids: BTreeSet<AccountId>,
724        unresolved_note_nullifiers: Vec<Nullifier>,
725        prune_tip: BlockNumber,
726    ) -> Result<BTreeMap<Nullifier, NoteId>> {
727        self.writer
728            .write::<_, DatabaseError, _>("apply block", move |tx| {
729                if let Some(protocol_config) = activated_protocol_config.as_ref() {
730                    queries::insert_protocol_config(
731                        tx,
732                        protocol_config,
733                        signed_block.header().block_num(),
734                    )?;
735                }
736                queries::apply_block(
737                    tx,
738                    &signed_block,
739                    &notes,
740                    &precomputed_public_states,
741                    &new_account_ids,
742                )?;
743                queries::prune_history(tx, prune_tip)?;
744                Ok(())
745            })
746            .await?;
747
748        Ok(self.resolve_consumed_note_ids(unresolved_note_nullifiers).await)
749    }
750
751    /// Maps consumed nullifiers back to their note IDs for lifecycle events, on a best-effort
752    /// basis.
753    ///
754    /// A failed lookup is logged and abandoned: the caller uses this only for reporting.
755    async fn resolve_consumed_note_ids(
756        &self,
757        nullifiers: Vec<Nullifier>,
758    ) -> BTreeMap<Nullifier, NoteId> {
759        let mut resolved_note_ids = BTreeMap::new();
760        for chunk in nullifiers.chunks(QueryParamNoteCommitmentLimit::LIMIT) {
761            let chunk = chunk.to_vec();
762            let count = chunk.len();
763            let result = self
764                .transact("resolve consumed note ids", move |conn| {
765                    diesel_queries::select_note_ids_by_nullifier(conn, &chunk)
766                })
767                .await;
768
769            match result {
770                Ok(note_ids) => resolved_note_ids.extend(note_ids),
771                Err(err) => {
772                    warn!(
773                        &err,
774                        target: COMPONENT,
775                        "Failed to resolve consumed note IDs for lifecycle events",
776                        note.nullifier.count = count
777                    );
778                    break;
779                },
780            }
781        }
782
783        resolved_note_ids
784    }
785
786    /// Selects storage map values for syncing storage maps for a specific account ID.
787    ///
788    /// The returned values are the latest known values up to `block_range.end()`, and no values
789    /// earlier than `block_range.start()` are returned.
790    pub(crate) async fn select_storage_map_sync_values(
791        &self,
792        account_id: AccountId,
793        block_range: ScopedBlockRange,
794        entries_limit: Option<usize>,
795    ) -> Result<StorageMapValuesPage> {
796        let block_range = block_range.into_inner();
797        let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
798
799        self.transact("select storage map sync values", move |conn| {
800            diesel_queries::select_account_storage_map_values_paged(
801                conn,
802                account_id,
803                block_range,
804                entries_limit,
805            )
806        })
807        .await
808    }
809
810    /// Reconstructs storage map details from the database for a specific slot at a block.
811    ///
812    /// Used as fallback when `AccountStateForest` cache misses (historical or evicted queries).
813    /// Rebuilds all entries by querying the DB and filtering to the specific slot.
814    ///
815    /// Returns:
816    ///     - `::LimitExceeded` when too many entries are present
817    ///     - `::AllEntries` if the size is less than or equal given `entries_limit`, if any
818    #[miden_instrument(
819        target = COMPONENT,
820    )]
821    pub(crate) async fn reconstruct_storage_map_from_db(
822        &self,
823        account_id: AccountId,
824        slot_name: miden_protocol::account::StorageSlotName,
825        block_num: ScopedBlockNum,
826        entries_limit: Option<usize>,
827    ) -> Result<miden_node_proto::domain::account::AccountStorageMapDetails> {
828        use miden_node_proto::domain::account::{AccountStorageMapDetails, StorageMapEntries};
829        use miden_protocol::EMPTY_WORD;
830
831        // TODO this remains expensive with a large history until we implement pruning for DB
832        // columns
833        let mut values = Vec::new();
834        let mut block_range_start = BlockNumber::GENESIS;
835        let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
836
837        let mut page = self
838            .select_storage_map_sync_values(
839                account_id,
840                block_num.range_from(block_range_start),
841                Some(entries_limit),
842            )
843            .await?;
844
845        values.extend(page.values);
846        let mut last_block_included = page.last_block_included;
847
848        // If the first page returned no values, the block at block_range_start has more entries
849        // than the limit allows (e.g. genesis accounts with large storage maps).
850        if values.is_empty() && last_block_included == block_range_start {
851            return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
852        }
853
854        loop {
855            if page.last_block_included == *block_num
856                || page.last_block_included < block_range_start
857            {
858                break;
859            }
860
861            block_range_start = page.last_block_included.child();
862            page = self
863                .select_storage_map_sync_values(
864                    account_id,
865                    block_num.range_from(block_range_start),
866                    Some(entries_limit),
867                )
868                .await?;
869
870            if page.last_block_included <= last_block_included {
871                return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
872            }
873
874            last_block_included = page.last_block_included;
875            values.extend(page.values);
876        }
877
878        if page.last_block_included != *block_num {
879            return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
880        }
881
882        // Filter to the specific slot and collect latest values per key
883        let mut latest_values = BTreeMap::<StorageMapKey, Word>::new();
884        for value in values {
885            if value.slot_name == slot_name {
886                let raw_key = value.key;
887                latest_values.insert(raw_key, value.value);
888            }
889        }
890
891        // Remove EMPTY_WORD entries (deletions)
892        latest_values.retain(|_, v| *v != EMPTY_WORD);
893
894        if latest_values.len() > AccountStorageMapDetails::MAX_RETURN_ENTRIES {
895            return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
896        }
897
898        let entries = latest_values.into_iter().collect::<Vec<_>>();
899        Ok(AccountStorageMapDetails {
900            slot_name,
901            entries: StorageMapEntries::AllEntries(entries),
902        })
903    }
904
905    /// Reconstructs the account vault from the database for a specific account at a block.
906    ///
907    /// Used as fallback when the `AccountStateForest` vault-key cache misses (historical or evicted
908    /// queries). Returns the latest asset for each vault key at or before `block_num`.
909    #[miden_instrument(
910        target = COMPONENT,
911    )]
912    pub async fn select_vault_at_block(
913        &self,
914        account_id: AccountId,
915        block_num: ScopedBlockNum,
916    ) -> Result<Vec<Asset>, DatabaseError> {
917        self.reader
918            .read("select vault at block", move |tx| {
919                queries::select_vault_at_block(tx, account_id, *block_num)
920            })
921            .await
922    }
923
924    pub async fn get_account_vault_sync(
925        &self,
926        account_id: AccountId,
927        block_range: ScopedBlockRange,
928    ) -> Result<(BlockNumber, Vec<AccountVaultValue>)> {
929        let block_range = block_range.into_inner();
930        self.transact("account vault sync", move |conn| {
931            diesel_queries::select_account_vault_assets(conn, account_id, block_range)
932        })
933        .await
934    }
935
936    /// Returns the script for a note by its root.
937    pub async fn select_note_script_by_root(&self, root: Word) -> Result<Option<NoteScript>> {
938        self.transact("note script by root", move |conn| {
939            diesel_queries::select_note_script_by_root(conn, root)
940        })
941        .await
942    }
943
944    /// Returns the complete transaction records for the specified accounts within the specified
945    /// block range, including state commitments and note IDs.
946    ///
947    /// Note: This method is size-limited (~5MB) and may not return all matching transactions
948    /// if the limit is exceeded. Transactions from partial blocks are excluded to maintain
949    /// consistency.
950    pub async fn select_transactions_records(
951        &self,
952        account_ids: Vec<AccountId>,
953        block_range: ScopedBlockRange,
954    ) -> Result<(BlockNumber, Vec<TransactionRecord>)> {
955        let block_range = block_range.into_inner();
956        self.transact("full transactions records", move |conn| {
957            diesel_queries::select_transactions_records(conn, &account_ids, block_range)
958        })
959        .await
960    }
961}