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 diesel::{Connection, SqliteConnection};
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 BlockHeader,
22 BlockNoteIndex,
23 BlockNumber,
24 BlockSignatures,
25 SignedBlock,
26};
27use miden_protocol::crypto::merkle::SparseMerklePath;
28use miden_protocol::note::{
29 NoteAttachments,
30 NoteDetails,
31 NoteId,
32 NoteInclusionProof,
33 NoteMetadata,
34 NoteScript,
35 Nullifier,
36};
37use miden_protocol::transaction::TransactionHeader;
38use miden_protocol::utils::serde::Deserializable;
39
40use crate::db::migrations::{migrate_database, verify_latest_schema};
41use crate::db::models::conv::SqlTypeConvert;
42use crate::db::models::queries;
43pub use crate::db::models::queries::{
44 AccountCommitmentsPage,
45 NullifiersPage,
46 PublicAccountIdsPage,
47 PublicAccountStateRootsPage,
48};
49use crate::db::models::queries::{
50 BlockHeaderCommitment,
51 PrecomputedPublicAccountStates,
52 StorageMapValuesPage,
53};
54use crate::errors::{DatabaseError, NoteSyncError};
55use crate::genesis::GenesisBlock;
56use crate::state::{ScopedBlockNum, ScopedBlockRange};
57use crate::{COMPONENT, LOG_TARGET};
58
59const STORAGE_MAP_VALUE_PER_ROW_BYTES: usize =
60 2 * size_of::<Word>() + size_of::<u32>() + size_of::<u8>();
61
62fn default_storage_map_entries_limit() -> usize {
63 MAX_RESPONSE_PAYLOAD_BYTES / STORAGE_MAP_VALUE_PER_ROW_BYTES
64}
65
66mod migrations;
67#[cfg(test)]
68pub(crate) use migrations::bootstrap_database;
69
70#[cfg(test)]
71mod tests;
72
73pub(crate) mod models;
74
75pub(crate) mod schema;
80
81pub type Result<T, E = DatabaseError> = std::result::Result<T, E>;
82
83#[derive(Copy, Clone, Debug, PartialEq, Eq)]
85pub struct DatabaseOptions {
86 pub connection_pool_size: NonZeroUsize,
88}
89
90impl Default for DatabaseOptions {
91 fn default() -> Self {
92 Self {
93 connection_pool_size: miden_node_db::default_connection_pool_size(),
94 }
95 }
96}
97
98pub struct Db {
102 db: miden_node_db::Db,
103}
104
105impl Deref for Db {
106 type Target = miden_node_db::Db;
107
108 fn deref(&self) -> &Self::Target {
109 &self.db
110 }
111}
112
113impl DerefMut for Db {
114 fn deref_mut(&mut self) -> &mut Self::Target {
115 &mut self.db
116 }
117}
118
119#[derive(Debug, Clone)]
123pub struct AccountVaultValue {
124 pub block_num: BlockNumber,
125 pub vault_key: AssetId,
126 pub asset: Option<Asset>,
128}
129
130impl AccountVaultValue {
131 pub fn from_raw_row(row: (i64, Vec<u8>, Option<Vec<u8>>)) -> Result<Self, DatabaseError> {
132 let (block_num, vault_key, asset) = row;
133 let vault_key = Word::read_from_bytes(&vault_key)?;
134 Ok(Self {
135 block_num: BlockNumber::from_raw_sql(block_num)?,
136 vault_key: AssetId::try_from(vault_key)?,
137 asset: asset.map(|b| Asset::read_from_bytes(&b)).transpose()?,
138 })
139 }
140}
141
142#[derive(Debug, PartialEq)]
143pub struct NullifierInfo {
144 pub nullifier: Nullifier,
145 pub block_num: BlockNumber,
146}
147
148impl PartialEq<(Nullifier, BlockNumber)> for NullifierInfo {
149 fn eq(&self, (nullifier, block_num): &(Nullifier, BlockNumber)) -> bool {
150 &self.nullifier == nullifier && &self.block_num == block_num
151 }
152}
153
154#[derive(Debug, PartialEq)]
155pub struct TransactionRecord {
156 pub block_num: BlockNumber,
157 pub header: TransactionHeader,
158 pub output_note_proofs: Vec<NoteSyncRecord>,
161 pub consumed_note_refs: Vec<(Nullifier, NoteId)>,
164}
165
166#[derive(Debug, Clone, PartialEq)]
167pub struct NoteRecord {
168 pub block_num: BlockNumber,
169 pub note_index: BlockNoteIndex,
170 pub note_id: Word,
171 pub metadata: NoteMetadata,
172 pub details: Option<NoteDetails>,
173 pub attachments: NoteAttachments,
174 pub inclusion_path: SparseMerklePath,
175}
176
177#[derive(Debug, PartialEq)]
178pub struct NoteSyncUpdate {
179 pub notes: Vec<NoteSyncRecord>,
180 pub block_header: BlockHeader,
181}
182
183#[derive(Debug, Clone, PartialEq)]
184pub struct NoteSyncRecord {
185 pub block_num: BlockNumber,
186 pub note_index: BlockNoteIndex,
187 pub note_id: NoteId,
188 pub metadata: NoteMetadata,
189 pub attachments: NoteAttachments,
190 pub inclusion_path: SparseMerklePath,
191}
192
193impl From<NoteRecord> for NoteSyncRecord {
194 fn from(note: NoteRecord) -> Self {
195 Self {
196 block_num: note.block_num,
197 note_index: note.note_index,
198 note_id: NoteId::from_raw(note.note_id),
199 metadata: note.metadata,
200 attachments: note.attachments,
201 inclusion_path: note.inclusion_path,
202 }
203 }
204}
205
206impl Db {
207 #[miden_instrument(
209 target = COMPONENT,
210 name = "store.database.bootstrap",
211 fields(path = database_filepath),
212 err,
213 )]
214 pub fn bootstrap(database_filepath: PathBuf, genesis: GenesisBlock) -> anyhow::Result<()> {
215 migrations::bootstrap_database(&database_filepath)
216 .context("failed to bootstrap database schema")?;
217
218 let mut conn: SqliteConnection = diesel::sqlite::SqliteConnection::establish(
219 database_filepath.to_str().context("database filepath is invalid")?,
220 )
221 .context("failed to open a database connection")?;
222
223 miden_node_db::configure_connection_on_creation(&mut conn)?;
224
225 let genesis_block = genesis.into_inner();
227 conn.transaction(move |conn| {
228 models::queries::apply_block(
229 conn,
230 &genesis_block,
231 &[],
232 &PrecomputedPublicAccountStates::new(),
233 )
234 })
235 .context("failed to insert genesis block")?;
236 Ok(())
237 }
238
239 #[miden_instrument(
241 target = COMPONENT,
242 )]
243 pub async fn load(database_filepath: PathBuf) -> Result<Self, DatabaseError> {
244 Self::load_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size())
245 .await
246 }
247
248 #[miden_instrument(
251 target = COMPONENT,
252 )]
253 pub async fn load_with_pool_size(
254 database_filepath: PathBuf,
255 connection_pool_size: NonZeroUsize,
256 ) -> Result<Self, DatabaseError> {
257 verify_latest_schema(&database_filepath)?;
258
259 let db = miden_node_db::Db::new_with_pool_size(&database_filepath, connection_pool_size)?;
260 info!(
261 target: LOG_TARGET,
262 "Connected to the database",
263 path = database_filepath,
264 db.sqlite.connection_pool_size = connection_pool_size.get()
265 );
266
267 Ok(Self { db })
268 }
269
270 #[miden_instrument(
272 target = COMPONENT,
273 )]
274 pub fn migrate(database_filepath: impl AsRef<Path>) -> Result<(), DatabaseError> {
275 migrate_database(database_filepath.as_ref())?;
276 Ok(())
277 }
278
279 #[miden_instrument(
281 level = "debug",
282 target = COMPONENT,
283 err,
284 )]
285 pub async fn select_nullifiers_paged(
286 &self,
287 page_size: std::num::NonZeroUsize,
288 after_nullifier: Option<Nullifier>,
289 ) -> Result<NullifiersPage> {
290 self.transact("read nullifiers paged", move |conn| {
291 queries::select_nullifiers_paged(conn, page_size, after_nullifier)
292 })
293 .await
294 }
295
296 #[miden_instrument(
298 level = "debug",
299 target = COMPONENT,
300 fields(
301 prefix_len,
302 prefix.count = nullifier_prefixes.len(),
303 ),
304 err,
305 )]
306 pub async fn select_nullifiers_by_prefix(
307 &self,
308 prefix_len: u32,
309 nullifier_prefixes: Vec<u32>,
310 block_range: ScopedBlockRange,
311 ) -> Result<(Vec<NullifierInfo>, BlockNumber)> {
312 let block_range = block_range.into_inner();
313 assert_eq!(prefix_len, 16, "Only 16-bit prefixes are supported");
314
315 self.transact("nullifieres by prefix", move |conn| {
316 let nullifier_prefixes =
317 nullifier_prefixes.into_iter().map(|prefix| prefix as u16).collect::<Vec<_>>();
318 queries::select_nullifiers_by_prefix(
319 conn,
320 prefix_len as u8,
321 &nullifier_prefixes[..],
322 block_range,
323 )
324 })
325 .await
326 }
327
328 #[miden_instrument(
332 level = "debug",
333 target = COMPONENT,
334 err,
335 )]
336 pub async fn select_block_header_by_block_num(
337 &self,
338 maybe_block_number: Option<ScopedBlockNum>,
339 ) -> Result<Option<BlockHeader>> {
340 self.transact("block headers by block number", move |conn| {
341 let val = queries::select_block_header_by_block_num(
342 conn,
343 maybe_block_number.map(|block_number| *block_number),
344 )?;
345 Ok(val)
346 })
347 .await
348 }
349
350 #[miden_instrument(
353 level = "debug",
354 target = COMPONENT,
355 err,
356 )]
357 pub async fn select_block_header_and_signatures_by_block_num(
358 &self,
359 block_number: ScopedBlockNum,
360 ) -> Result<Option<(BlockHeader, BlockSignatures)>> {
361 self.transact("block headers and signatures by block number", move |conn| {
362 let val =
363 queries::select_block_header_and_signatures_by_block_num(conn, *block_number)?;
364 Ok(val)
365 })
366 .await
367 }
368
369 #[miden_instrument(
371 level = "debug",
372 target = COMPONENT,
373 err,
374 )]
375 pub async fn select_block_headers(
376 &self,
377 blocks: impl Iterator<Item = ScopedBlockNum> + Send + 'static,
378 ) -> Result<Vec<BlockHeader>> {
379 self.transact("block headers from given block numbers", move |conn| {
380 let raw = queries::select_block_headers(conn, blocks.map(|block| *block))?;
381 Ok(raw)
382 })
383 .await
384 }
385
386 #[miden_instrument(
388 level = "debug",
389 target = COMPONENT,
390 err,
391 )]
392 pub async fn select_all_block_header_commitments(&self) -> Result<Vec<BlockHeaderCommitment>> {
393 self.transact("all block headers", |conn| {
394 let raw = queries::select_all_block_header_commitments(conn)?;
395 Ok(raw)
396 })
397 .await
398 }
399
400 #[miden_instrument(
402 level = "debug",
403 target = COMPONENT,
404 err,
405 )]
406 pub async fn select_account_commitments_paged(
407 &self,
408 page_size: std::num::NonZeroUsize,
409 after_account_id: Option<AccountId>,
410 ) -> Result<AccountCommitmentsPage> {
411 self.transact("read account commitments paged", move |conn| {
412 queries::select_account_commitments_paged(conn, page_size, after_account_id)
413 })
414 .await
415 }
416
417 #[miden_instrument(
419 level = "debug",
420 target = COMPONENT,
421 err,
422 )]
423 pub async fn select_public_account_ids_paged(
424 &self,
425 page_size: std::num::NonZeroUsize,
426 after_account_id: Option<AccountId>,
427 ) -> Result<PublicAccountIdsPage> {
428 self.transact("read public account IDs paged", move |conn| {
429 queries::select_public_account_ids_paged(conn, page_size, after_account_id)
430 })
431 .await
432 }
433
434 #[miden_instrument(
436 level = "debug",
437 target = COMPONENT,
438 err,
439 )]
440 pub async fn select_public_account_state_roots_paged(
441 &self,
442 page_size: std::num::NonZeroUsize,
443 after_account_id: Option<AccountId>,
444 ) -> Result<PublicAccountStateRootsPage> {
445 self.transact("read public account state roots paged", move |conn| {
446 queries::select_public_account_state_roots_paged(conn, page_size, after_account_id)
447 })
448 .await
449 }
450
451 #[miden_instrument(
453 level = "debug",
454 target = COMPONENT,
455 err,
456 )]
457 pub async fn select_account(&self, id: AccountId) -> Result<AccountInfo> {
458 self.transact("Get account details", move |conn| queries::select_account(conn, id))
459 .await
460 }
461
462 #[miden_instrument(
464 level = "debug",
465 target = COMPONENT,
466 err,
467 )]
468 pub async fn select_network_accounts_subset(
469 &self,
470 account_ids: Vec<AccountId>,
471 ) -> Result<HashSet<AccountId>> {
472 self.transact("Filter network accounts subset", move |conn| {
473 queries::select_network_accounts_subset(conn, &account_ids)
474 })
475 .await
476 }
477
478 #[miden_instrument(
482 target = COMPONENT,
483 )]
484 pub async fn select_account_code_by_commitment(
485 &self,
486 code_commitment: Word,
487 ) -> Result<Option<Vec<u8>>> {
488 self.transact("Get account code by commitment", move |conn| {
489 queries::select_account_code_by_commitment(conn, code_commitment)
490 })
491 .await
492 }
493
494 #[miden_instrument(
499 target = COMPONENT,
500 )]
501 pub async fn select_account_header_with_storage_header_at_block(
502 &self,
503 account_id: AccountId,
504 block_num: ScopedBlockNum,
505 ) -> Result<Option<(AccountHeader, AccountStorageHeader)>> {
506 self.transact("Get account header with storage header at block", move |conn| {
507 queries::select_account_header_with_storage_header_at_block(
508 conn, account_id, *block_num,
509 )
510 })
511 .await
512 }
513
514 #[miden_instrument(
515 level = "debug",
516 target = COMPONENT,
517 err,
518 )]
519 pub async fn get_note_sync_multi(
520 &self,
521 block_range: ScopedBlockRange,
522 note_tags: Arc<[u32]>,
523 ) -> Result<Vec<NoteSyncUpdate>, NoteSyncError> {
524 let block_range = block_range.into_inner();
525 self.transact("notes sync task", move |conn| {
526 queries::get_note_sync_multi(conn, ¬e_tags, block_range, MAX_RESPONSE_PAYLOAD_BYTES)
527 })
528 .await
529 }
530
531 #[miden_instrument(
534 level = "debug",
535 target = COMPONENT,
536 err,
537 )]
538 pub async fn select_notes_by_id(&self, note_ids: Vec<NoteId>) -> Result<Vec<NoteRecord>> {
539 self.transact("note by id", move |conn| {
540 queries::select_notes_by_id(conn, note_ids.as_slice())
541 })
542 .await
543 }
544
545 #[miden_instrument(
548 level = "debug",
549 target = COMPONENT,
550 err,
551 )]
552 pub async fn select_existing_note_commitments(
553 &self,
554 note_commitments: Vec<Word>,
555 up_to_block: ScopedBlockNum,
556 ) -> Result<HashSet<Word>> {
557 self.transact("note by commitment", move |conn| {
558 queries::select_existing_note_commitments(
559 conn,
560 note_commitments.as_slice(),
561 *up_to_block,
562 )
563 })
564 .await
565 }
566
567 #[miden_instrument(
570 level = "debug",
571 target = COMPONENT,
572 err,
573 )]
574 pub async fn select_note_inclusion_proofs(
575 &self,
576 note_commitments: BTreeSet<Word>,
577 up_to_block: ScopedBlockNum,
578 ) -> Result<BTreeMap<NoteId, NoteInclusionProof>> {
579 self.transact("block note inclusion proofs by commitment", move |conn| {
580 models::queries::select_note_inclusion_proofs(conn, ¬e_commitments, *up_to_block)
581 })
582 .await
583 }
584
585 #[miden_instrument(
601 target = COMPONENT,
602 err,
603 )]
604 pub(crate) async fn apply_block(
605 &self,
606 signed_block: SignedBlock,
607 notes: Vec<(NoteRecord, Option<Nullifier>)>,
608 precomputed_public_states: PrecomputedPublicAccountStates,
609 unresolved_note_nullifiers: Vec<Nullifier>,
610 prune_tip: BlockNumber,
611 ) -> Result<BTreeMap<Nullifier, NoteId>> {
612 self.transact("apply block", move |conn| {
613 models::queries::apply_block(conn, &signed_block, ¬es, &precomputed_public_states)?;
614 models::queries::prune_history(conn, prune_tip)?;
615
616 let mut resolved_note_ids = BTreeMap::new();
617 for chunk in unresolved_note_nullifiers.chunks(QueryParamNoteCommitmentLimit::LIMIT) {
618 match queries::select_note_ids_by_nullifier(conn, chunk) {
619 Ok(note_ids) => resolved_note_ids.extend(note_ids),
620 Err(err) => {
621 warn!(
622 &err,
623 target: COMPONENT,
624 "Failed to resolve consumed note IDs for lifecycle events",
625 note.nullifier.count = chunk.len()
626 );
627 break;
628 },
629 }
630 }
631
632 Ok(resolved_note_ids)
633 })
634 .await
635 }
636
637 pub(crate) async fn select_storage_map_sync_values(
642 &self,
643 account_id: AccountId,
644 block_range: ScopedBlockRange,
645 entries_limit: Option<usize>,
646 ) -> Result<StorageMapValuesPage> {
647 let block_range = block_range.into_inner();
648 let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
649
650 self.transact("select storage map sync values", move |conn| {
651 models::queries::select_account_storage_map_values_paged(
652 conn,
653 account_id,
654 block_range,
655 entries_limit,
656 )
657 })
658 .await
659 }
660
661 #[miden_instrument(
670 target = COMPONENT,
671 )]
672 pub(crate) async fn reconstruct_storage_map_from_db(
673 &self,
674 account_id: AccountId,
675 slot_name: miden_protocol::account::StorageSlotName,
676 block_num: ScopedBlockNum,
677 entries_limit: Option<usize>,
678 ) -> Result<miden_node_proto::domain::account::AccountStorageMapDetails> {
679 use miden_node_proto::domain::account::{AccountStorageMapDetails, StorageMapEntries};
680 use miden_protocol::EMPTY_WORD;
681
682 let mut values = Vec::new();
685 let mut block_range_start = BlockNumber::GENESIS;
686 let entries_limit = entries_limit.unwrap_or_else(default_storage_map_entries_limit);
687
688 let mut page = self
689 .select_storage_map_sync_values(
690 account_id,
691 block_num.range_from(block_range_start),
692 Some(entries_limit),
693 )
694 .await?;
695
696 values.extend(page.values);
697 let mut last_block_included = page.last_block_included;
698
699 if values.is_empty() && last_block_included == block_range_start {
702 return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
703 }
704
705 loop {
706 if page.last_block_included == *block_num
707 || page.last_block_included < block_range_start
708 {
709 break;
710 }
711
712 block_range_start = page.last_block_included.child();
713 page = self
714 .select_storage_map_sync_values(
715 account_id,
716 block_num.range_from(block_range_start),
717 Some(entries_limit),
718 )
719 .await?;
720
721 if page.last_block_included <= last_block_included {
722 return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
723 }
724
725 last_block_included = page.last_block_included;
726 values.extend(page.values);
727 }
728
729 if page.last_block_included != *block_num {
730 return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
731 }
732
733 let mut latest_values = BTreeMap::<StorageMapKey, Word>::new();
735 for value in values {
736 if value.slot_name == slot_name {
737 let raw_key = value.key;
738 latest_values.insert(raw_key, value.value);
739 }
740 }
741
742 latest_values.retain(|_, v| *v != EMPTY_WORD);
744
745 if latest_values.len() > AccountStorageMapDetails::MAX_RETURN_ENTRIES {
746 return Ok(AccountStorageMapDetails::limit_exceeded(slot_name));
747 }
748
749 let entries = latest_values.into_iter().collect::<Vec<_>>();
750 Ok(AccountStorageMapDetails {
751 slot_name,
752 entries: StorageMapEntries::AllEntries(entries),
753 })
754 }
755
756 #[miden_instrument(
761 target = COMPONENT,
762 )]
763 pub async fn select_account_vault_at_block(
764 &self,
765 account_id: AccountId,
766 block_num: ScopedBlockNum,
767 ) -> Result<Vec<Asset>, DatabaseError> {
768 self.transact("select account vault at block", move |conn| {
769 queries::select_account_vault_at_block(conn, account_id, *block_num)
770 })
771 .await
772 }
773
774 pub async fn get_account_vault_sync(
775 &self,
776 account_id: AccountId,
777 block_range: ScopedBlockRange,
778 ) -> Result<(BlockNumber, Vec<AccountVaultValue>)> {
779 let block_range = block_range.into_inner();
780 self.transact("account vault sync", move |conn| {
781 queries::select_account_vault_assets(conn, account_id, block_range)
782 })
783 .await
784 }
785
786 pub async fn select_note_script_by_root(&self, root: Word) -> Result<Option<NoteScript>> {
788 self.transact("note script by root", move |conn| {
789 queries::select_note_script_by_root(conn, root)
790 })
791 .await
792 }
793
794 pub async fn select_transactions_records(
801 &self,
802 account_ids: Vec<AccountId>,
803 block_range: ScopedBlockRange,
804 ) -> Result<(BlockNumber, Vec<TransactionRecord>)> {
805 let block_range = block_range.into_inner();
806 self.transact("full transactions records", move |conn| {
807 queries::select_transactions_records(conn, &account_ids, block_range)
808 })
809 .await
810 }
811}