use std::collections::{BTreeMap, BTreeSet, HashSet};
use std::mem::size_of;
use std::num::NonZeroUsize;
use std::ops::{Deref, DerefMut};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use anyhow::Context;
use miden_node_db::sqlite::{DbReader, DbWriter, WriteTx};
use miden_node_proto::domain::account::AccountInfo;
use miden_node_tracing::{info, miden_instrument, warn};
use miden_node_utils::limiter::{
MAX_RESPONSE_PAYLOAD_BYTES,
QueryParamLimiter,
QueryParamNoteCommitmentLimit,
};
use miden_protocol::Word;
use miden_protocol::account::{AccountHeader, AccountId, AccountStorageHeader, StorageMapKey};
use miden_protocol::asset::{Asset, AssetId};
use miden_protocol::block::{
BlockAccountUpdate,
BlockHeader,
BlockNoteIndex,
BlockNumber,
BlockSignatures,
SignedBlock,
};
use miden_protocol::crypto::merkle::SparseMerklePath;
use miden_protocol::note::{
NoteAttachments,
NoteDetails,
NoteId,
NoteInclusionProof,
NoteMetadata,
NoteScript,
Nullifier,
};
use miden_protocol::protocol_config::ProtocolConfig;
use miden_protocol::transaction::TransactionHeader;
use miden_protocol::utils::serde::Deserializable;
use crate::db::migrations::{migrate_database, verify_latest_schema};
use crate::db::models::conv::SqlTypeConvert;
use crate::db::models::queries as diesel_queries;
use crate::db::models::queries::StorageMapValuesPage;
pub use crate::db::models::queries::{
AccountCommitmentsPage,
NullifiersPage,
PublicAccountIdsPage,
PublicAccountStateRootsPage,
};
pub use crate::db::queries::{
HISTORICAL_BLOCK_RETENTION,
PrecomputedPublicAccountState,
PrecomputedPublicAccountStates,
};
use crate::errors::{DatabaseError, NoteSyncError};
use crate::genesis::GenesisBlock;
use crate::state::{ScopedBlockNum, ScopedBlockRange};
use crate::{COMPONENT, LOG_TARGET};
const STORAGE_MAP_VALUE_PER_ROW_BYTES: usize =
2 * size_of::<Word>() + size_of::<u32>() + size_of::<u8>();
fn default_storage_map_entries_limit() -> usize {
MAX_RESPONSE_PAYLOAD_BYTES / STORAGE_MAP_VALUE_PER_ROW_BYTES
}
mod migrations;
#[cfg(test)]
pub(crate) use migrations::bootstrap_database;
#[cfg(test)]
mod tests;
#[cfg(test)]
mod test_db;
#[cfg(test)]
pub(crate) use test_db::TestDb;
pub(crate) mod queries;
mod utils;
pub(crate) mod models;
pub(crate) mod schema;
pub type Result<T, E = DatabaseError> = std::result::Result<T, E>;
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub struct DatabaseOptions {
pub connection_pool_size: NonZeroUsize,
}
impl Default for DatabaseOptions {
fn default() -> Self {
Self {
connection_pool_size: miden_node_db::default_connection_pool_size(),
}
}
}
pub struct Db {
diesel: miden_node_db::Db,
writer: DbWriter,
reader: DbReader,
}
fn insert_genesis(tx: &WriteTx<'_>, genesis: GenesisBlock) -> Result<()> {
let (genesis_block, protocol_config) = genesis.into_parts();
let new_account_ids = genesis_block
.body()
.updated_accounts()
.iter()
.map(BlockAccountUpdate::account_id)
.collect();
queries::insert_protocol_config(tx, &protocol_config, BlockNumber::GENESIS)?;
queries::apply_block(
tx,
&genesis_block,
&[],
&PrecomputedPublicAccountStates::new(),
&new_account_ids,
)?;
Ok(())
}
impl Deref for Db {
type Target = miden_node_db::Db;
fn deref(&self) -> &Self::Target {
&self.diesel
}
}
impl DerefMut for Db {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.diesel
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(transparent)]
pub struct BlockHeaderCommitment(pub(crate) Word);
impl BlockHeaderCommitment {
pub fn new(header: &BlockHeader) -> Self {
Self(header.commitment())
}
pub fn word(self) -> Word {
self.0
}
}
#[derive(Debug, Clone)]
pub struct AccountVaultValue {
pub block_num: BlockNumber,
pub vault_key: AssetId,
pub asset: Option<Asset>,
}
impl AccountVaultValue {
pub fn from_raw_row(row: (i64, Vec<u8>, Option<Vec<u8>>)) -> Result<Self, DatabaseError> {
let (block_num, vault_key, asset) = row;
let vault_key = Word::read_from_bytes(&vault_key)?;
Ok(Self {
block_num: BlockNumber::from_raw_sql(block_num)?,
vault_key: AssetId::try_from(vault_key)?,
asset: asset.map(|b| miden_node_persistence::decode::<Asset>(&b)).transpose()?,
})
}
}
#[derive(Debug, PartialEq)]
pub struct NullifierInfo {
pub nullifier: Nullifier,
pub block_num: BlockNumber,
}
impl PartialEq<(Nullifier, BlockNumber)> for NullifierInfo {
fn eq(&self, (nullifier, block_num): &(Nullifier, BlockNumber)) -> bool {
&self.nullifier == nullifier && &self.block_num == block_num
}
}
#[derive(Debug, PartialEq)]
pub struct TransactionRecord {
pub block_num: BlockNumber,
pub header: TransactionHeader,
pub output_note_proofs: Vec<NoteSyncRecord>,
pub consumed_note_refs: Vec<(Nullifier, NoteId)>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct NoteRecord {
pub block_num: BlockNumber,
pub note_index: BlockNoteIndex,
pub note_id: Word,
pub metadata: NoteMetadata,
pub details: Option<NoteDetails>,
pub attachments: NoteAttachments,
pub inclusion_path: SparseMerklePath,
}
#[derive(Debug, PartialEq)]
pub struct NoteSyncUpdate {
pub notes: Vec<NoteSyncRecord>,
pub block_header: BlockHeader,
}
#[derive(Debug, Clone, PartialEq)]
pub struct NoteSyncRecord {
pub block_num: BlockNumber,
pub note_index: BlockNoteIndex,
pub note_id: NoteId,
pub metadata: NoteMetadata,
pub attachments: NoteAttachments,
pub inclusion_path: SparseMerklePath,
}
impl From<NoteRecord> for NoteSyncRecord {
fn from(note: NoteRecord) -> Self {
Self {
block_num: note.block_num,
note_index: note.note_index,
note_id: NoteId::from_raw(note.note_id),
metadata: note.metadata,
attachments: note.attachments,
inclusion_path: note.inclusion_path,
}
}
}
impl Db {
#[miden_instrument(
target = COMPONENT,
name = "store.database.bootstrap",
fields(path = database_filepath),
err,
)]
pub async fn bootstrap(
database_filepath: PathBuf,
genesis: GenesisBlock,
) -> anyhow::Result<()> {
migrations::bootstrap_database(&database_filepath)
.context("failed to bootstrap database schema")?;
let (writer, _reader) = miden_node_db::sqlite::open(&database_filepath)
.context("failed to open a database connection")?;
writer
.write("insert genesis block", move |tx| insert_genesis(tx, genesis))
.await
.context("failed to insert genesis block")?;
Ok(())
}
#[miden_instrument(
target = COMPONENT,
)]
pub async fn load(database_filepath: PathBuf) -> Result<Self, DatabaseError> {
Self::load_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size())
.await
}
#[miden_instrument(
target = COMPONENT,
)]
pub async fn load_with_pool_size(
database_filepath: PathBuf,
connection_pool_size: NonZeroUsize,
) -> Result<Self, DatabaseError> {
verify_latest_schema(&database_filepath)?;
let db = miden_node_db::Db::new_with_pool_size(&database_filepath, connection_pool_size)?;
let (writer, reader) =
miden_node_db::sqlite::open_with_pool_size(&database_filepath, connection_pool_size)?;
info!(
target: LOG_TARGET,
"Connected to the database",
path = database_filepath,
db.sqlite.connection_pool_size = connection_pool_size.get()
);
Ok(Self { diesel: db, writer, reader })
}
#[cfg(test)]
pub(crate) fn writer(&self) -> &DbWriter {
&self.writer
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_protocol_config_by_commitment(
&self,
commitment: Word,
) -> Result<Option<ProtocolConfig>> {
self.transact("protocol config by commitment", move |conn| {
diesel_queries::select_protocol_config_by_commitment(conn, commitment)
})
.await
}
pub async fn select_protocol_config_commitment_at(
&self,
block_number: ScopedBlockNum,
) -> Result<Option<Word>> {
self.transact("protocol config commitment at block", move |conn| {
diesel_queries::select_protocol_config_commitment_at(conn, *block_number)
})
.await
}
#[miden_instrument(
target = COMPONENT,
)]
pub fn migrate(database_filepath: impl AsRef<Path>) -> Result<(), DatabaseError> {
migrate_database(database_filepath.as_ref())?;
Ok(())
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_nullifiers_paged(
&self,
page_size: std::num::NonZeroUsize,
after_nullifier: Option<Nullifier>,
) -> Result<NullifiersPage> {
self.transact("read nullifiers paged", move |conn| {
diesel_queries::select_nullifiers_paged(conn, page_size, after_nullifier)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
fields(
prefix_len,
prefix.count = nullifier_prefixes.len(),
),
err,
)]
pub async fn select_nullifiers_by_prefix(
&self,
prefix_len: u32,
nullifier_prefixes: Vec<u32>,
block_range: ScopedBlockRange,
) -> Result<(Vec<NullifierInfo>, BlockNumber)> {
let block_range = block_range.into_inner();
assert_eq!(prefix_len, 16, "Only 16-bit prefixes are supported");
self.transact("nullifieres by prefix", move |conn| {
let nullifier_prefixes =
nullifier_prefixes.into_iter().map(|prefix| prefix as u16).collect::<Vec<_>>();
diesel_queries::select_nullifiers_by_prefix(
conn,
prefix_len as u8,
&nullifier_prefixes[..],
block_range,
)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_block_header_by_block_num(
&self,
maybe_block_number: Option<ScopedBlockNum>,
) -> Result<Option<BlockHeader>> {
self.transact("block headers by block number", move |conn| {
let val = diesel_queries::select_block_header_by_block_num(
conn,
maybe_block_number.map(|block_number| *block_number),
)?;
Ok(val)
})
.await
}
pub(crate) async fn select_genesis_block_header(&self) -> Result<Option<BlockHeader>> {
self.transact("genesis block header", |conn| {
diesel_queries::select_block_header_by_block_num(conn, Some(BlockNumber::GENESIS))
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_block_header_and_signatures_by_block_num(
&self,
block_number: ScopedBlockNum,
) -> Result<Option<(BlockHeader, BlockSignatures)>> {
self.transact("block headers and signatures by block number", move |conn| {
let val = diesel_queries::select_block_header_and_signatures_by_block_num(
conn,
*block_number,
)?;
Ok(val)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_block_headers(
&self,
blocks: impl Iterator<Item = ScopedBlockNum> + Send + 'static,
) -> Result<Vec<BlockHeader>> {
self.transact("block headers from given block numbers", move |conn| {
let raw = diesel_queries::select_block_headers(conn, blocks.map(|block| *block))?;
Ok(raw)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_all_block_header_commitments(&self) -> Result<Vec<BlockHeaderCommitment>> {
self.transact("all block headers", |conn| {
let raw = diesel_queries::select_all_block_header_commitments(conn)?;
Ok(raw)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_account_commitments_paged(
&self,
page_size: std::num::NonZeroUsize,
after_account_id: Option<AccountId>,
) -> Result<AccountCommitmentsPage> {
self.transact("read account commitments paged", move |conn| {
diesel_queries::select_account_commitments_paged(conn, page_size, after_account_id)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_public_account_ids_paged(
&self,
page_size: std::num::NonZeroUsize,
after_account_id: Option<AccountId>,
) -> Result<PublicAccountIdsPage> {
self.transact("read public account IDs paged", move |conn| {
diesel_queries::select_public_account_ids_paged(conn, page_size, after_account_id)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_public_account_state_roots_paged(
&self,
page_size: std::num::NonZeroUsize,
after_account_id: Option<AccountId>,
) -> Result<PublicAccountStateRootsPage> {
self.transact("read public account state roots paged", move |conn| {
diesel_queries::select_public_account_state_roots_paged(
conn,
page_size,
after_account_id,
)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_account(&self, id: AccountId) -> Result<AccountInfo> {
self.transact("Get account details", move |conn| diesel_queries::select_account(conn, id))
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn filter_network_accounts(
&self,
account_ids: Vec<AccountId>,
) -> Result<HashSet<AccountId>> {
self.reader
.read("Filter network accounts", move |tx| {
queries::filter_network_accounts(tx, &account_ids)
})
.await
}
#[miden_instrument(
target = COMPONENT,
)]
pub async fn select_account_code_by_commitment(
&self,
code_commitment: Word,
) -> Result<Option<miden_protocol::account::AccountCode>> {
self.transact("Get account code by commitment", move |conn| {
diesel_queries::select_account_code_by_commitment(conn, code_commitment)?
.map(|bytes| {
miden_node_persistence::decode::<miden_protocol::account::AccountCode>(&bytes)
})
.transpose()
.map_err(DatabaseError::from)
})
.await
}
#[miden_instrument(
target = COMPONENT,
)]
pub async fn select_account_header_with_storage_header_at_block(
&self,
account_id: AccountId,
block_num: ScopedBlockNum,
) -> Result<Option<(AccountHeader, AccountStorageHeader)>> {
self.reader
.read("Get account header with storage header at block", move |tx| {
queries::select_account_header_with_storage_header_at_block(
tx, account_id, *block_num,
)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn get_note_sync_multi(
&self,
block_range: ScopedBlockRange,
note_tags: Arc<[u32]>,
) -> Result<Vec<NoteSyncUpdate>, NoteSyncError> {
let block_range = block_range.into_inner();
self.transact("notes sync task", move |conn| {
diesel_queries::get_note_sync_multi(
conn,
¬e_tags,
block_range,
MAX_RESPONSE_PAYLOAD_BYTES,
)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_notes_by_id(&self, note_ids: Vec<NoteId>) -> Result<Vec<NoteRecord>> {
self.transact("note by id", move |conn| {
diesel_queries::select_notes_by_id(conn, note_ids.as_slice())
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_existing_note_ids(
&self,
note_ids: Vec<NoteId>,
up_to_block: ScopedBlockNum,
) -> Result<HashSet<NoteId>> {
self.transact("existing note IDs", move |conn| {
diesel_queries::select_existing_note_ids(conn, note_ids.as_slice(), *up_to_block)
})
.await
}
#[miden_instrument(
level = "debug",
target = COMPONENT,
err,
)]
pub async fn select_note_inclusion_proofs(
&self,
note_commitments: BTreeSet<Word>,
up_to_block: ScopedBlockNum,
) -> Result<BTreeMap<NoteId, NoteInclusionProof>> {
self.transact("block note inclusion proofs by commitment", move |conn| {
diesel_queries::select_note_inclusion_proofs(conn, ¬e_commitments, *up_to_block)
})
.await
}
#[expect(
clippy::too_many_arguments,
reason = "the arguments are the block and the state that the writer precomputed for it"
)]
#[miden_instrument(
target = COMPONENT,
err,
)]
pub(crate) async fn apply_block(
&self,
signed_block: SignedBlock,
activated_protocol_config: Option<ProtocolConfig>,
notes: Vec<(NoteRecord, Option<Nullifier>)>,
precomputed_public_states: PrecomputedPublicAccountStates,
new_account_ids: BTreeSet<AccountId>,
unresolved_note_nullifiers: Vec<Nullifier>,
prune_tip: BlockNumber,
) -> Result<BTreeMap<Nullifier, NoteId>> {
self.writer
.write::<_, DatabaseError, _>("apply block", move |tx| {
if let Some(protocol_config) = activated_protocol_config.as_ref() {
queries::insert_protocol_config(
tx,
protocol_config,
signed_block.header().block_num(),
)?;
}
queries::apply_block(
tx,
&signed_block,
¬es,
&precomputed_public_states,
&new_account_ids,
)?;
queries::prune_history(tx, prune_tip)?;
Ok(())
})
.await?;
Ok(self.resolve_consumed_note_ids(unresolved_note_nullifiers).await)
}
async fn resolve_consumed_note_ids(
&self,
nullifiers: Vec<Nullifier>,
) -> BTreeMap<Nullifier, NoteId> {
let mut resolved_note_ids = BTreeMap::new();
for chunk in nullifiers.chunks(QueryParamNoteCommitmentLimit::LIMIT) {
let chunk = chunk.to_vec();
let count = chunk.len();
let result = self
.transact("resolve consumed note ids", move |conn| {
diesel_queries::select_note_ids_by_nullifier(conn, &chunk)
})
.await;
match result {
Ok(note_ids) => resolved_note_ids.extend(note_ids),
Err(err) => {
warn!(
&err,
target: COMPONENT,
"Failed to resolve consumed note IDs for lifecycle events",
note.nullifier.count = count
);
break;
},
}
}
resolved_note_ids
}
pub(crate) async fn select_storage_map_sync_values(
&self,
account_id: AccountId,
block_range: ScopedBlockRange,
entries_limit: Option<usize>,
) -> Result<StorageMapValuesPage> {
let block_range = block_range.into_inner();
let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
self.transact("select storage map sync values", move |conn| {
diesel_queries::select_account_storage_map_values_paged(
conn,
account_id,
block_range,
entries_limit,
)
})
.await
}
#[miden_instrument(
target = COMPONENT,
)]
pub(crate) async fn reconstruct_storage_map_from_db(
&self,
account_id: AccountId,
slot_name: miden_protocol::account::StorageSlotName,
block_num: ScopedBlockNum,
entries_limit: Option<usize>,
) -> Result<miden_node_proto::domain::account::AccountStorageMapDetails> {
use miden_node_proto::domain::account::{AccountStorageMapDetails, StorageMapEntries};
use miden_protocol::EMPTY_WORD;
let mut values = Vec::new();
let mut block_range_start = BlockNumber::GENESIS;
let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
let mut page = self
.select_storage_map_sync_values(
account_id,
block_num.range_from(block_range_start),
Some(entries_limit),
)
.await?;
values.extend(page.values);
let mut last_block_included = page.last_block_included;
if values.is_empty() && last_block_included == block_range_start {
return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
}
loop {
if page.last_block_included == *block_num
|| page.last_block_included < block_range_start
{
break;
}
block_range_start = page.last_block_included.child();
page = self
.select_storage_map_sync_values(
account_id,
block_num.range_from(block_range_start),
Some(entries_limit),
)
.await?;
if page.last_block_included <= last_block_included {
return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
}
last_block_included = page.last_block_included;
values.extend(page.values);
}
if page.last_block_included != *block_num {
return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
}
let mut latest_values = BTreeMap::<StorageMapKey, Word>::new();
for value in values {
if value.slot_name == slot_name {
let raw_key = value.key;
latest_values.insert(raw_key, value.value);
}
}
latest_values.retain(|_, v| *v != EMPTY_WORD);
if latest_values.len() > AccountStorageMapDetails::MAX_RETURN_ENTRIES {
return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
}
let entries = latest_values.into_iter().collect::<Vec<_>>();
Ok(AccountStorageMapDetails {
slot_name,
entries: StorageMapEntries::AllEntries(entries),
})
}
#[miden_instrument(
target = COMPONENT,
)]
pub async fn select_vault_at_block(
&self,
account_id: AccountId,
block_num: ScopedBlockNum,
) -> Result<Vec<Asset>, DatabaseError> {
self.reader
.read("select vault at block", move |tx| {
queries::select_vault_at_block(tx, account_id, *block_num)
})
.await
}
pub async fn get_account_vault_sync(
&self,
account_id: AccountId,
block_range: ScopedBlockRange,
) -> Result<(BlockNumber, Vec<AccountVaultValue>)> {
let block_range = block_range.into_inner();
self.transact("account vault sync", move |conn| {
diesel_queries::select_account_vault_assets(conn, account_id, block_range)
})
.await
}
pub async fn select_note_script_by_root(&self, root: Word) -> Result<Option<NoteScript>> {
self.transact("note script by root", move |conn| {
diesel_queries::select_note_script_by_root(conn, root)
})
.await
}
pub async fn select_transactions_records(
&self,
account_ids: Vec<AccountId>,
block_range: ScopedBlockRange,
) -> Result<(BlockNumber, Vec<TransactionRecord>)> {
let block_range = block_range.into_inner();
self.transact("full transactions records", move |conn| {
diesel_queries::select_transactions_records(conn, &account_ids, block_range)
})
.await
}
}