1use std::io;
2use std::num::NonZeroUsize;
3use std::path::{Path, PathBuf};
4
5use miden_node_db::DatabaseError;
6use miden_node_db::sqlite::{DbReader, DbWriter};
7use miden_node_tracing::{info, miden_instrument};
8use miden_protocol::block::{BlockHeader, BlockNumber};
9use miden_protocol::protocol_config::ProtocolConfig;
10use miden_protocol::transaction::TransactionId;
11
12use crate::db::migrations::{bootstrap_database, migrate_database, verify_latest_schema};
13use crate::metrics::InitialMetrics;
14use crate::{COMPONENT, LOG_TARGET, StorageKeyEpoch, StoredPrivateRecord};
15
16mod migrations;
17mod queries;
18
19pub(crate) use queries::{ListTransactionsParams, ListedTransaction};
20
21#[derive(Clone)]
30pub struct ValidatorDbReader {
31 reader: DbReader,
32}
33
34impl ValidatorDbReader {
35 pub(crate) async fn transaction_exists(
37 &self,
38 tx_id: TransactionId,
39 ) -> Result<bool, DatabaseError> {
40 self.reader
41 .read("transaction_exists", move |tx| queries::transaction_exists(tx, tx_id))
42 .await
43 }
44
45 pub(crate) async fn find_unvalidated_transactions(
49 &self,
50 tx_ids: Vec<TransactionId>,
51 ) -> Result<Vec<TransactionId>, DatabaseError> {
52 self.reader
53 .read("find_unvalidated_transactions", move |tx| {
54 queries::find_unvalidated_transactions(tx, &tx_ids)
55 })
56 .await
57 }
58
59 #[miden_instrument(
62 target = COMPONENT,
63 )]
64 pub(crate) async fn load_chain_tip(&self) -> Result<Option<BlockHeader>, DatabaseError> {
65 self.reader.read("load_chain_tip", queries::load_chain_tip).await
66 }
67
68 pub(crate) async fn load_block_header(
70 &self,
71 block_num: BlockNumber,
72 ) -> Result<Option<BlockHeader>, DatabaseError> {
73 self.reader
74 .read("load_block_header", move |tx| queries::load_block_header(tx, block_num))
75 .await
76 }
77
78 pub async fn load_protocol_config(
80 &self,
81 commitment: miden_protocol::Word,
82 ) -> Result<Option<ProtocolConfig>, DatabaseError> {
83 self.reader
84 .read("load_protocol_config", move |tx| queries::load_protocol_config(tx, commitment))
85 .await
86 }
87
88 pub(crate) async fn load_initial_metrics(&self) -> Result<InitialMetrics, DatabaseError> {
91 self.reader
92 .read("load_initial_metrics", |tx| {
93 Ok(InitialMetrics {
94 chain_tip: queries::load_chain_tip(tx)?
95 .map_or(0, |header| header.block_num().as_u32()),
96 validated_transactions: u64::try_from(queries::count_validated_transactions(
97 tx,
98 )?)
99 .unwrap_or(0),
100 signed_blocks: u64::try_from(queries::count_signed_blocks(tx)?).unwrap_or(0),
101 })
102 })
103 .await
104 }
105
106 #[cfg(test)]
112 pub(crate) async fn count_validated_transactions(&self) -> Result<i64, DatabaseError> {
113 self.reader
114 .read("count_validated_transactions", queries::count_validated_transactions)
115 .await
116 }
117
118 pub async fn load_private_record(
120 &self,
121 transaction_id: TransactionId,
122 ) -> Result<Option<StoredPrivateRecord>, DatabaseError> {
123 self.reader
124 .read("load_private_record", move |tx| {
125 queries::load_private_record(tx, transaction_id)
126 })
127 .await
128 }
129
130 pub async fn load_private_records_by_key_epoch(
132 &self,
133 key_epoch: StorageKeyEpoch,
134 ) -> Result<Vec<StoredPrivateRecord>, DatabaseError> {
135 self.reader
136 .read("load_private_records_by_key_epoch", move |tx| {
137 queries::load_private_records_by_key_epoch(tx, key_epoch)
138 })
139 .await
140 }
141
142 pub async fn load_private_records_by_setup_context(
144 &self,
145 setup_context_id: [u8; 32],
146 ) -> Result<Vec<StoredPrivateRecord>, DatabaseError> {
147 self.reader
148 .read("load_private_records_by_setup_context", move |tx| {
149 queries::load_private_records_by_setup_context(tx, setup_context_id)
150 })
151 .await
152 }
153
154 pub(crate) async fn list_validated_transactions(
157 &self,
158 params: queries::ListTransactionsParams,
159 ) -> Result<queries::ListedTransactionsPage, DatabaseError> {
160 self.reader
161 .read("list_validated_transactions", move |tx| {
162 queries::list_validated_transactions(tx, ¶ms)
163 })
164 .await
165 }
166}
167
168pub struct ValidatorDbWriter {
174 writer: DbWriter,
175 reader: ValidatorDbReader,
176}
177
178impl std::ops::Deref for ValidatorDbWriter {
179 type Target = ValidatorDbReader;
180
181 fn deref(&self) -> &Self::Target {
182 &self.reader
183 }
184}
185
186impl ValidatorDbWriter {
187 pub fn reader(&self) -> ValidatorDbReader {
190 self.reader.clone()
191 }
192
193 #[miden_instrument(
196 target = COMPONENT,
197 )]
198 pub async fn insert_validated_private_transaction(
199 &self,
200 record: StoredPrivateRecord,
201 ) -> Result<usize, DatabaseError> {
202 self.writer
203 .write("insert_validated_private_transaction", move |tx| {
204 queries::insert_validated_private_transaction(tx, &record)
205 })
206 .await
207 }
208
209 #[cfg(test)]
219 #[miden_instrument(
220 target = COMPONENT,
221 )]
222 pub(crate) async fn upsert_block_header_with_protocol_config(
223 &self,
224 header: BlockHeader,
225 protocol_config: Option<ProtocolConfig>,
226 ) -> Result<(), DatabaseError> {
227 self.writer
228 .write("upsert_block_header_with_protocol_config", move |tx| {
229 record_protocol_config_activation(tx, &header, protocol_config)?;
230 queries::upsert_block_header(tx, &header)
231 })
232 .await
233 }
234
235 #[miden_instrument(
242 target = COMPONENT,
243 )]
244 pub(crate) async fn insert_signed_block(
245 &self,
246 header: BlockHeader,
247 protocol_config: ProtocolConfig,
248 transactions: Vec<TransactionId>,
249 ) -> Result<(), DatabaseError> {
250 self.writer
251 .write("insert_signed_block", move |tx| {
252 persist_signed_block(tx, &header, protocol_config, &transactions)
253 })
254 .await
255 }
256
257 #[miden_instrument(
262 target = COMPONENT,
263 )]
264 pub(crate) async fn replace_signed_block(
265 &self,
266 header: BlockHeader,
267 protocol_config: ProtocolConfig,
268 transactions: Vec<TransactionId>,
269 ) -> Result<(), DatabaseError> {
270 self.writer
271 .write("replace_signed_block", move |tx| {
272 queries::delete_block(tx, header.block_num())?;
273 persist_signed_block(tx, &header, protocol_config, &transactions)
274 })
275 .await
276 }
277}
278
279fn persist_signed_block(
283 tx: &miden_node_db::sqlite::WriteTx<'_>,
284 header: &BlockHeader,
285 protocol_config: ProtocolConfig,
286 transactions: &[TransactionId],
287) -> Result<(), DatabaseError> {
288 record_protocol_config_activation(tx, header, Some(protocol_config))?;
289 queries::insert_block_header(tx, header)?;
290 queries::link_block_transactions(tx, header.block_num(), transactions)
291}
292
293fn record_protocol_config_activation(
301 tx: &miden_node_db::sqlite::WriteTx<'_>,
302 header: &BlockHeader,
303 protocol_config: Option<ProtocolConfig>,
304) -> Result<(), DatabaseError> {
305 let commitment = header.protocol_config_commitment();
306 let config = if let Some(config) = protocol_config {
307 let calculated = config.to_commitment();
308 if calculated != commitment {
309 return Err(invalid_protocol_config(format!(
310 "protocol config commitment mismatch: expected {commitment}, got {calculated}"
311 )));
312 }
313 config
314 } else {
315 queries::load_protocol_config(tx, commitment)?.ok_or_else(|| {
316 invalid_protocol_config(format!("protocol config {commitment} is not stored"))
317 })?
318 };
319
320 let block_number = header.block_num();
321 let previous = queries::load_protocol_config_commitment_before(tx, block_number)?;
322 if previous != Some(commitment) {
323 queries::insert_protocol_config(tx, &config, block_number)?;
324 }
325
326 Ok(())
327}
328
329fn invalid_protocol_config(message: String) -> DatabaseError {
330 DatabaseError::deserialization(
331 "ProtocolConfig",
332 io::Error::new(io::ErrorKind::InvalidData, message),
333 )
334}
335
336#[cfg(test)]
338pub(crate) async fn delete_protocol_config_for_test(
339 db: &ValidatorDbWriter,
340 commitment: miden_protocol::Word,
341) -> Result<(), DatabaseError> {
342 db.writer
343 .write("delete_protocol_config_for_test", move |tx| {
344 tx.execute("DELETE FROM protocol_configs WHERE commitment = ?1", &[&commitment])?;
345 Ok::<_, DatabaseError>(())
346 })
347 .await
348}
349
350#[miden_instrument(
355 target = COMPONENT,
356)]
357pub async fn load(database_filepath: PathBuf) -> Result<ValidatorDbWriter, DatabaseError> {
358 load_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size()).await
359}
360
361#[miden_instrument(
364 target = COMPONENT,
365)]
366pub async fn load_with_pool_size(
367 database_filepath: PathBuf,
368 connection_pool_size: NonZeroUsize,
369) -> Result<ValidatorDbWriter, DatabaseError> {
370 verify_latest_schema(&database_filepath)?;
371
372 open_with_pool_size(&database_filepath, connection_pool_size)
373}
374
375#[miden_instrument(
377 target = COMPONENT,
378)]
379pub async fn setup(database_filepath: PathBuf) -> Result<ValidatorDbWriter, DatabaseError> {
380 setup_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size()).await
381}
382
383#[miden_instrument(
385 target = COMPONENT,
386)]
387async fn setup_with_pool_size(
388 database_filepath: PathBuf,
389 connection_pool_size: NonZeroUsize,
390) -> Result<ValidatorDbWriter, DatabaseError> {
391 bootstrap_database(&database_filepath)?;
392
393 open_with_pool_size(&database_filepath, connection_pool_size)
394}
395
396#[miden_instrument(
401 target = COMPONENT,
402 fields(path = database_filepath),
403 err,
404)]
405pub async fn bootstrap(
406 database_filepath: PathBuf,
407 connection_pool_size: NonZeroUsize,
408 genesis_header: BlockHeader,
409 protocol_config: ProtocolConfig,
410) -> Result<(), DatabaseError> {
411 let db = setup_with_pool_size(database_filepath, connection_pool_size).await?;
412
413 db.insert_signed_block(genesis_header, protocol_config, Vec::new()).await
414}
415
416#[miden_instrument(
418 target = COMPONENT,
419)]
420pub fn migrate(database_filepath: impl AsRef<Path>) -> Result<(), DatabaseError> {
421 migrate_database(database_filepath.as_ref())?;
422 Ok(())
423}
424
425fn open_with_pool_size(
426 database_filepath: &Path,
427 connection_pool_size: NonZeroUsize,
428) -> Result<ValidatorDbWriter, DatabaseError> {
429 let (writer, reader) =
430 miden_node_db::sqlite::open_with_pool_size(database_filepath, connection_pool_size)?;
431 info!(
432 target: LOG_TARGET,
433 "Connected to the database",
434 path = database_filepath,
435 db.sqlite.connection_pool_size = connection_pool_size.get()
436 );
437 Ok(ValidatorDbWriter {
438 writer,
439 reader: ValidatorDbReader { reader },
440 })
441}
442
443#[cfg(test)]
444mod tests {
445 mod protocol_config_history;
446
447 use miden_node_utils::fee::{test_fee_params, test_protocol_config};
448 use miden_protocol::Word;
449 use miden_protocol::asset::AssetId;
450 use miden_protocol::block::{BlockHeader, ValidatorConfig};
451 use miden_protocol::crypto::dsa::ecdsa_k256_keccak::SigningKey;
452 use miden_protocol::protocol_config::ProtocolConfig;
453 use miden_protocol::testing::account_id::ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET_1;
454 use miden_protocol::utils::serde::Deserializable;
455 use rand_chacha_03::ChaCha20Rng;
456 use rand_chacha_03::rand_core::SeedableRng;
457
458 use super::*;
459 use crate::private_record::test_private_record_sealer;
460 use crate::storage_key::tests::operator_keys;
461 use crate::{
462 PrivateRecordChainId,
463 PrivateRecordCombiner,
464 PrivateRecordContext,
465 PrivateRecordError,
466 PrivateRecordFormatVersion,
467 PrivateRecordId,
468 PrivateRecordSealer,
469 PrivateRecordShareRequest,
470 };
471
472 const CHAIN_ID: PrivateRecordChainId = PrivateRecordChainId::new([1; 32]);
473 const KEY_EPOCH: StorageKeyEpoch = StorageKeyEpoch::new([2; 32]);
474 const SETUP_CONTEXT_ID: [u8; 32] = [4; 32];
475
476 fn record_id(transaction_id: TransactionId) -> PrivateRecordId {
477 let signer = SigningKey::read_from_bytes(&[7; 32]).unwrap();
478 PrivateRecordId::new(transaction_id, &signer.public_key())
479 }
480
481 fn private_record(transaction_id: TransactionId, seed: u8) -> StoredPrivateRecord {
482 let context = PrivateRecordContext::new(CHAIN_ID, KEY_EPOCH, transaction_id);
483 let mut rng = ChaCha20Rng::from_seed([seed; 32]);
484 test_private_record_sealer(KEY_EPOCH, SETUP_CONTEXT_ID)
485 .seal(&mut rng, record_id(transaction_id), context, b"private transaction inputs")
486 .unwrap()
487 }
488
489 fn genesis_header(config: &miden_protocol::protocol_config::ProtocolConfig) -> BlockHeader {
490 miden_node_store::GenesisState::new(
491 vec![],
492 test_fee_params(),
493 0,
494 ValidatorConfig::new(vec![SigningKey::new().public_key()], 1).unwrap(),
495 config.clone(),
496 )
497 .into_block()
498 .unwrap()
499 .inner()
500 .header()
501 .clone()
502 }
503
504 fn header_with_next_timestamp(header: &BlockHeader) -> BlockHeader {
505 BlockHeader::new(
506 header.prev_block_commitment(),
507 header.block_num(),
508 header.chain_commitment(),
509 header.account_root(),
510 header.nullifier_root(),
511 header.note_root(),
512 header.tx_commitment(),
513 header.validator_config().clone(),
514 header.fee_parameters().clone(),
515 header.protocol_config_commitment(),
516 header.next_protocol_config().cloned(),
517 header.timestamp() + 1,
518 )
519 }
520
521 #[test]
522 fn migrate_rejects_missing_database() {
523 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
524 let db_path = temp_dir.path().join("validator.sqlite3");
525
526 let err = migrate(db_path.clone()).expect_err("missing database should fail");
527
528 assert!(matches!(err, DatabaseError::Migration(_)), "unexpected error: {err:?}");
529 assert!(!db_path.exists());
530 }
531
532 #[tokio::test]
534 async fn migration_preserves_headers_and_private_records() {
535 let temp_dir = tempfile::tempdir().unwrap();
536 let db_path = temp_dir.path().join("validator.sqlite3");
537 miden_node_db::migration::Migrator::builder()
538 .unwrap()
539 .push_sql("001_initial", include_str!("migrations/001_initial.sql"))
540 .unwrap()
541 .build()
542 .unwrap()
543 .bootstrap(&db_path)
544 .unwrap();
545
546 let config = test_protocol_config();
547 let header = genesis_header(&config);
548 let transaction_id = TransactionId::from_raw(Word::from([1u32, 2, 3, 4]));
549 let record = private_record(transaction_id, 1);
550 let db = open_with_pool_size(&db_path, NonZeroUsize::new(2).unwrap()).unwrap();
551 let stored_header = header.clone();
552 db.writer
553 .write("seed legacy header", move |tx| queries::upsert_block_header(tx, &stored_header))
554 .await
555 .unwrap();
556 db.insert_validated_private_transaction(record.clone()).await.unwrap();
557 drop(db);
558
559 assert!(load(db_path.clone()).await.is_err(), "the old schema requires migration");
560 migrate(&db_path).unwrap();
561 migrate(&db_path).expect("migration should also accept the latest schema");
562
563 let db = load(db_path).await.unwrap();
564 assert_eq!(db.load_chain_tip().await.unwrap(), Some(header.clone()));
565 assert_eq!(db.load_private_record(transaction_id).await.unwrap(), Some(record));
566 let migrated = db.load_private_record(transaction_id).await.unwrap().unwrap();
567 assert_eq!(migrated.context().format_version(), PrivateRecordFormatVersion::V1);
568 migrated.verify_encrypted_record_key().unwrap();
569 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), None);
570
571 db.upsert_block_header_with_protocol_config(header, Some(config.clone()))
572 .await
573 .unwrap();
574 assert_eq!(
575 protocol_config_history::history(&db).await,
576 vec![(0, config.to_commitment(), config.clone())]
577 );
578 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
579 }
580
581 #[tokio::test]
582 async fn setup_creates_database_that_load_accepts() {
583 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
584 let db_path = temp_dir.path().join("validator.sqlite3");
585
586 setup(db_path.clone()).await.expect("setup should bootstrap the database");
587 load(db_path).await.expect("load should accept a bootstrapped database");
588 }
589
590 #[tokio::test]
591 async fn setup_creates_protocol_config_storage() {
592 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
593 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
594
595 let row_count = db
596 .reader
597 .reader
598 .read("protocol_config_storage", |tx| {
599 Ok::<_, DatabaseError>(
600 tx.query("SELECT COUNT(*) FROM protocol_configs", &[], |row| {
601 row.get::<i64>(0)
602 })?
603 .into_iter()
604 .next()
605 .expect("COUNT always returns one row"),
606 )
607 })
608 .await
609 .unwrap();
610
611 assert_eq!(row_count, 0);
612 }
613
614 #[tokio::test]
615 async fn block_header_and_protocol_config_are_persisted_together() {
616 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
617 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
618 let config = test_protocol_config();
619 let header = genesis_header(&config);
620
621 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
622 .await
623 .unwrap();
624
625 assert_eq!(db.load_chain_tip().await.unwrap(), Some(header));
626 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
627 }
628
629 #[tokio::test]
630 async fn duplicate_supplied_protocol_config_is_accepted() {
631 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
632 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
633 let config = test_protocol_config();
634 let header = genesis_header(&config);
635 let replacement = header_with_next_timestamp(&header);
636
637 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
638 .await
639 .unwrap();
640 db.upsert_block_header_with_protocol_config(replacement.clone(), Some(config.clone()))
641 .await
642 .expect("a duplicate supplied protocol config should be accepted");
643
644 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), Some(replacement));
645 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
646 }
647
648 #[tokio::test]
649 async fn known_protocol_config_can_be_omitted() {
650 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
651 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
652 let config = test_protocol_config();
653 let header = genesis_header(&config);
654 let replacement = header_with_next_timestamp(&header);
655
656 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
657 .await
658 .unwrap();
659 db.upsert_block_header_with_protocol_config(replacement.clone(), None)
660 .await
661 .expect("a stored protocol config should not need to be supplied again");
662
663 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), Some(replacement));
664 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
665 }
666
667 #[tokio::test]
668 async fn unknown_protocol_config_rolls_back_block_header() {
669 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
670 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
671 let config = test_protocol_config();
672 let header = genesis_header(&config);
673
674 db.upsert_block_header_with_protocol_config(header.clone(), None)
675 .await
676 .expect_err("an unknown protocol config must reject the header");
677
678 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
679 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), None);
680 }
681
682 #[tokio::test]
683 async fn mismatched_protocol_config_rolls_back_block_header() {
684 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
685 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
686 let expected = test_protocol_config();
687 let header = genesis_header(&expected);
688 let mismatched = ProtocolConfig::current(AssetId::new_fungible(
689 ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET_1.try_into().unwrap(),
690 ))
691 .unwrap();
692
693 db.upsert_block_header_with_protocol_config(header.clone(), Some(mismatched.clone()))
694 .await
695 .expect_err("a mismatched config must reject the header transaction");
696
697 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
698 assert_eq!(db.load_protocol_config(expected.to_commitment()).await.unwrap(), None);
699 assert_eq!(db.load_protocol_config(mismatched.to_commitment()).await.unwrap(), None);
700 }
701
702 #[tokio::test]
703 async fn block_header_insertion_failure_rolls_back_new_protocol_config() {
704 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
705 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
706 let config = test_protocol_config();
707 let commitment = config.to_commitment();
708 let header = genesis_header(&config);
709
710 db.writer
711 .write("reject_block_header_inserts", |tx| {
712 tx.execute(
713 "CREATE TRIGGER reject_block_header_insert
714 BEFORE INSERT ON block_headers
715 BEGIN
716 SELECT RAISE(ABORT, 'block header insertion rejected');
717 END;",
718 &[],
719 )?;
720 Ok::<_, DatabaseError>(())
721 })
722 .await
723 .unwrap();
724
725 db.upsert_block_header_with_protocol_config(header.clone(), Some(config))
726 .await
727 .expect_err("a block header insertion failure must reject the transaction");
728
729 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
730 assert_eq!(db.load_protocol_config(commitment).await.unwrap(), None);
731 }
732
733 #[tokio::test]
734 async fn transaction_exists_detects_validated_transactions() {
735 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
736 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
737
738 let validated_id = TransactionId::from_raw(Word::try_from([1u64, 2, 3, 4]).unwrap());
739 let unknown_id = TransactionId::from_raw(Word::try_from([5u64, 6, 7, 8]).unwrap());
740
741 db.insert_validated_private_transaction(private_record(validated_id, 1))
742 .await
743 .unwrap();
744
745 assert!(
746 db.transaction_exists(validated_id).await.unwrap(),
747 "an inserted transaction id should be reported as existing"
748 );
749 assert!(
750 !db.transaction_exists(unknown_id).await.unwrap(),
751 "an unknown transaction id should not be reported as existing"
752 );
753 }
754
755 #[tokio::test]
758 async fn find_unvalidated_transactions_returns_only_missing_ids() {
759 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
760 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
761
762 let ids = (1u64..=4)
763 .map(|i| TransactionId::from_raw(Word::try_from([i, i, i, i]).unwrap()))
764 .collect::<Vec<_>>();
765
766 db.insert_validated_private_transaction(private_record(ids[1], 1))
768 .await
769 .unwrap();
770 db.insert_validated_private_transaction(private_record(ids[3], 2))
771 .await
772 .unwrap();
773
774 let unvalidated = db.find_unvalidated_transactions(ids.clone()).await.unwrap();
775 assert_eq!(unvalidated, vec![ids[0], ids[2]]);
776
777 assert!(db.find_unvalidated_transactions(vec![]).await.unwrap().is_empty());
779 }
780
781 #[tokio::test]
782 async fn load_initial_metrics_reports_persisted_state() {
783 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
784 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
785
786 let metrics = db.load_initial_metrics().await.unwrap();
788 assert_eq!(metrics.chain_tip, 0);
789 assert_eq!(metrics.validated_transactions, 0);
790 assert_eq!(metrics.signed_blocks, 0);
791 }
792
793 #[tokio::test]
794 async fn private_record_indexes_work() {
795 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
796 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
797 let transaction_id = TransactionId::from_raw(Word::from([5u32, 6, 7, 8]));
798 let record = private_record(transaction_id, 9);
799
800 let expected = record.clone();
801 db.insert_validated_private_transaction(record).await.unwrap();
802
803 let by_record = db.load_private_record(transaction_id).await.unwrap();
804 assert_eq!(by_record, Some(expected.clone()));
805
806 let by_epoch = db.load_private_records_by_key_epoch(KEY_EPOCH).await.unwrap();
807 assert_eq!(by_epoch, vec![expected.clone()]);
808
809 let by_setup = db.load_private_records_by_setup_context(SETUP_CONTEXT_ID).await.unwrap();
810 assert_eq!(by_setup, vec![expected.clone()]);
811 }
812
813 #[tokio::test]
814 async fn private_record_rejects_unsupported_formats() {
815 let temp_dir = tempfile::tempdir().unwrap();
816 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
817 let transaction_id = TransactionId::from_raw(Word::from([5u32, 6, 7, 8]));
818 db.insert_validated_private_transaction(private_record(transaction_id, 9))
819 .await
820 .unwrap();
821
822 for format_version in [2_u32, 3, u32::MAX] {
823 db.writer
824 .write("set unsupported record format", move |tx| {
825 tx.execute(
826 "UPDATE validated_transactions SET format_version = ? WHERE id = ?",
827 &[&i64::from(format_version), &transaction_id],
828 )
829 })
830 .await
831 .unwrap();
832
833 let error = db.load_private_record(transaction_id).await.unwrap_err();
834 assert!(matches!(
835 error,
836 DatabaseError::ConversionSqlToRust { inner: Some(source), .. }
837 if matches!(
838 source.downcast_ref::<PrivateRecordError>(),
839 Some(PrivateRecordError::UnsupportedFormat(version))
840 if *version == format_version,
841 ),
842 ));
843 }
844 }
845
846 #[tokio::test]
849 async fn uncommitted_transactions_are_not_listed() {
850 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
851 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
852 let transaction_ids = [
853 TransactionId::from_raw(Word::from([9u32, 0, 0, 0])),
854 TransactionId::from_raw(Word::from([1u32, 0, 0, 0])),
855 ];
856 for (transaction_id, seed) in transaction_ids.into_iter().zip([1u8, 2]) {
857 db.insert_validated_private_transaction(private_record(transaction_id, seed))
858 .await
859 .unwrap();
860 }
861
862 let params = ListTransactionsParams {
863 start: None,
864 block_to: None,
865 limit: 10,
866 include_records: false,
867 };
868 assert!(db.list_validated_transactions(params).await.unwrap().transactions.is_empty());
869 assert!(db.load_private_record(transaction_ids[0]).await.unwrap().is_some());
871 }
872
873 #[tokio::test]
876 async fn insert_signed_block_links_and_relinks_transactions() {
877 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
878 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
879 let transaction_ids = (1u64..=3)
880 .map(|i| TransactionId::from_raw(Word::try_from([i, i, i, i]).unwrap()))
881 .collect::<Vec<_>>();
882 for (seed, transaction_id) in [1u8, 2, 3].into_iter().zip(&transaction_ids) {
883 db.insert_validated_private_transaction(private_record(*transaction_id, seed))
884 .await
885 .unwrap();
886 }
887
888 let header = BlockHeader::mock(7, None, None, &[]);
889 db.insert_signed_block(
890 header.clone(),
891 ProtocolConfig::mock(),
892 vec![transaction_ids[0], transaction_ids[1]],
893 )
894 .await
895 .unwrap();
896
897 let params = ListTransactionsParams {
898 start: None,
899 block_to: None,
900 limit: 10,
901 include_records: false,
902 };
903 let listed = db.list_validated_transactions(params).await.unwrap().transactions;
904 assert_eq!(
905 listed
906 .iter()
907 .map(|item| (item.transaction_id, item.block_num, item.block_tx_index))
908 .collect::<Vec<_>>(),
909 vec![
910 (transaction_ids[0], BlockNumber::from(7u32), 0),
911 (transaction_ids[1], BlockNumber::from(7u32), 1),
912 ],
913 );
914
915 db.replace_signed_block(header, ProtocolConfig::mock(), vec![transaction_ids[2]])
919 .await
920 .unwrap();
921
922 let listed = db.list_validated_transactions(params).await.unwrap().transactions;
923 assert_eq!(
924 listed
925 .iter()
926 .map(|item| (item.transaction_id, item.block_num, item.block_tx_index))
927 .collect::<Vec<_>>(),
928 vec![(transaction_ids[2], BlockNumber::from(7u32), 0)],
929 );
930 }
931
932 #[tokio::test]
935 async fn listing_snapshot_survives_tip_replacement_and_advance() {
936 let directory = tempfile::tempdir().unwrap();
937 let db = setup(directory.path().join("validator.sqlite3")).await.unwrap();
938 let original_id = TransactionId::from_raw(Word::from([1u32, 0, 0, 0]));
939 let replacement_id = TransactionId::from_raw(Word::from([2u32, 0, 0, 0]));
940 let original_record = private_record(original_id, 1);
941 let replacement_record = private_record(replacement_id, 2);
942 for record in [&original_record, &replacement_record] {
943 db.insert_validated_private_transaction(record.clone()).await.unwrap();
944 }
945 let header = BlockHeader::mock(1, None, None, &[]);
946 db.insert_signed_block(header.clone(), ProtocolConfig::mock(), vec![original_id])
947 .await
948 .unwrap();
949
950 let snapshot = db.reader.reader.begin_read().await.unwrap();
951 let tip = snapshot.run("pin_listing_snapshot", queries::load_chain_tip).await.unwrap();
952 assert_eq!(tip.unwrap().block_num(), BlockNumber::from(1u32));
953
954 db.replace_signed_block(header, ProtocolConfig::mock(), vec![replacement_id])
955 .await
956 .unwrap();
957 db.insert_signed_block(
958 BlockHeader::mock(2, None, None, &[]),
959 ProtocolConfig::mock(),
960 Vec::new(),
961 )
962 .await
963 .unwrap();
964
965 let params = ListTransactionsParams {
966 start: None,
967 block_to: None,
968 limit: 10,
969 include_records: true,
970 };
971 let page = snapshot
972 .run("list_snapshot", move |tx| queries::list_validated_transactions(tx, ¶ms))
973 .await
974 .unwrap();
975 assert_eq!(page.chain_tip, BlockNumber::from(1u32));
976 assert_eq!(page.transactions.len(), 1);
977 assert_eq!(page.transactions[0].transaction_id, original_id);
978 assert_eq!(page.transactions[0].block_num, page.chain_tip);
979 assert_eq!(page.transactions[0].record.as_ref(), Some(&original_record));
980 snapshot.close().await.unwrap();
981
982 let current = db.list_validated_transactions(params).await.unwrap();
983 assert_eq!(current.chain_tip, BlockNumber::from(2u32));
984 assert_eq!(current.transactions.len(), 1);
985 assert_eq!(current.transactions[0].transaction_id, replacement_id);
986 assert_eq!(current.transactions[0].block_num, BlockNumber::from(1u32));
987 assert_eq!(current.transactions[0].record.as_ref(), Some(&replacement_record));
988 }
989
990 #[tokio::test]
993 async fn list_validated_transactions_pages_in_committed_order() {
994 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
995 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
996 let transaction_ids = (1u64..=5)
997 .map(|i| TransactionId::from_raw(Word::try_from([i, i, i, i]).unwrap()))
998 .collect::<Vec<_>>();
999 for (seed, transaction_id) in (1u8..=5).zip(&transaction_ids) {
1000 db.insert_validated_private_transaction(private_record(*transaction_id, seed))
1001 .await
1002 .unwrap();
1003 }
1004 let block_1 = BlockHeader::mock(1, None, None, &[]);
1008 let block_2 = BlockHeader::mock(2, None, None, &[]);
1009 db.insert_signed_block(
1010 block_1,
1011 ProtocolConfig::mock(),
1012 vec![transaction_ids[3], transaction_ids[0]],
1013 )
1014 .await
1015 .unwrap();
1016 db.insert_signed_block(
1017 block_2,
1018 ProtocolConfig::mock(),
1019 vec![transaction_ids[1], transaction_ids[2]],
1020 )
1021 .await
1022 .unwrap();
1023 let expected_order =
1024 [transaction_ids[3], transaction_ids[0], transaction_ids[1], transaction_ids[2]];
1025
1026 let mut swept = Vec::new();
1029 let mut start = None;
1030 loop {
1031 let params = ListTransactionsParams {
1032 start,
1033 block_to: None,
1034 limit: 3,
1035 include_records: false,
1036 };
1037 let page = db.list_validated_transactions(params).await.unwrap().transactions;
1038 let Some(last) = page.last() else { break };
1039 assert!(page.len() <= 3, "a page must honor the row limit");
1040 start = Some((last.block_num, last.block_tx_index + 1));
1041 swept.extend(page.into_iter().map(|item| item.transaction_id));
1042 }
1043 assert_eq!(swept, expected_order);
1044
1045 let from_block_2 = db
1047 .list_validated_transactions(ListTransactionsParams {
1048 start: Some((BlockNumber::from(2u32), 0)),
1049 block_to: None,
1050 limit: 10,
1051 include_records: false,
1052 })
1053 .await
1054 .unwrap();
1055 assert_eq!(
1056 from_block_2
1057 .transactions
1058 .iter()
1059 .map(|item| item.transaction_id)
1060 .collect::<Vec<_>>(),
1061 vec![transaction_ids[1], transaction_ids[2]],
1062 );
1063 let up_to_block_1 = db
1064 .list_validated_transactions(ListTransactionsParams {
1065 start: None,
1066 block_to: Some(BlockNumber::from(1u32)),
1067 limit: 10,
1068 include_records: false,
1069 })
1070 .await
1071 .unwrap();
1072 assert_eq!(
1073 up_to_block_1
1074 .transactions
1075 .iter()
1076 .map(|item| item.transaction_id)
1077 .collect::<Vec<_>>(),
1078 vec![transaction_ids[3], transaction_ids[0]],
1079 );
1080 }
1081
1082 #[tokio::test]
1085 async fn insert_signed_block_rejects_unvalidated_transactions() {
1086 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
1087 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
1088 let header = BlockHeader::mock(1, None, None, &[]);
1089 let unknown = TransactionId::from_raw(Word::from([1u32, 0, 0, 0]));
1090
1091 let result = db.insert_signed_block(header, ProtocolConfig::mock(), vec![unknown]).await;
1092 assert!(result.is_err(), "linking an unvalidated transaction must fail");
1093 }
1094
1095 #[tokio::test]
1096 async fn stored_private_record_opens_with_threshold_shares() {
1097 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
1098 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
1099 let operators = operator_keys();
1100 let transaction_id = TransactionId::from_raw(Word::from([9u32, 10, 11, 12]));
1101 let context = PrivateRecordContext::new(CHAIN_ID, operators[0].key_epoch(), transaction_id);
1102 let plaintext = b"private transaction inputs";
1103 let mut seal_rng = ChaCha20Rng::from_seed([40; 32]);
1104 let record = PrivateRecordSealer::from_operator_key(&operators[0])
1105 .seal(&mut seal_rng, record_id(transaction_id), context, plaintext)
1106 .unwrap();
1107 db.insert_validated_private_transaction(record).await.unwrap();
1108
1109 let stored = db.load_private_record(transaction_id).await.unwrap().unwrap();
1110 let request = PrivateRecordShareRequest::for_record(&stored);
1111 let mut first_rng = ChaCha20Rng::from_seed([41; 32]);
1112 let mut second_rng = ChaCha20Rng::from_seed([42; 32]);
1113 let shares = [
1114 operators[0]
1115 .issue_private_record_share(&mut first_rng, &request, &stored)
1116 .unwrap(),
1117 operators[1]
1118 .issue_private_record_share(&mut second_rng, &request, &stored)
1119 .unwrap(),
1120 ];
1121
1122 let opened = PrivateRecordCombiner::from_operator_key(&operators[2])
1123 .unwrap()
1124 .open(&request, &stored, &shares)
1125 .unwrap();
1126 assert_eq!(opened.as_slice(), plaintext);
1127 }
1128
1129 #[tokio::test]
1130 async fn private_record_schema_has_required_indexes() {
1131 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
1132 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
1133
1134 let schema = db
1135 .reader
1136 .reader
1137 .read("private_record_schema", |tx| {
1138 tx.query(
1139 "SELECT sql FROM sqlite_schema \
1140 WHERE tbl_name IN ('validated_transactions', 'block_transactions') \
1141 AND sql IS NOT NULL \
1142 ORDER BY name",
1143 &[],
1144 |row| row.get::<String>(0),
1145 )
1146 })
1147 .await
1148 .unwrap()
1149 .join("\n");
1150
1151 assert!(schema.contains("insertion_sequence INTEGER PRIMARY KEY AUTOINCREMENT"));
1152 assert!(schema.contains("id BLOB NOT NULL UNIQUE"));
1153 assert!(schema.contains("idx_validated_transactions_key_epoch"));
1154 assert!(schema.contains("idx_validated_transactions_setup_context_id"));
1155 assert!(schema.contains("PRIMARY KEY (block_num, block_tx_index)"));
1156 }
1157}