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