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
19#[derive(Clone)]
28pub struct ValidatorDbReader {
29 reader: DbReader,
30}
31
32impl ValidatorDbReader {
33 pub(crate) async fn transaction_exists(
35 &self,
36 tx_id: TransactionId,
37 ) -> Result<bool, DatabaseError> {
38 self.reader
39 .read("transaction_exists", move |tx| queries::transaction_exists(tx, tx_id))
40 .await
41 }
42
43 pub(crate) async fn find_unvalidated_transactions(
47 &self,
48 tx_ids: Vec<TransactionId>,
49 ) -> Result<Vec<TransactionId>, DatabaseError> {
50 self.reader
51 .read("find_unvalidated_transactions", move |tx| {
52 queries::find_unvalidated_transactions(tx, &tx_ids)
53 })
54 .await
55 }
56
57 #[miden_instrument(
60 target = COMPONENT,
61 )]
62 pub(crate) async fn load_chain_tip(&self) -> Result<Option<BlockHeader>, DatabaseError> {
63 self.reader.read("load_chain_tip", queries::load_chain_tip).await
64 }
65
66 pub(crate) async fn load_block_header(
68 &self,
69 block_num: BlockNumber,
70 ) -> Result<Option<BlockHeader>, DatabaseError> {
71 self.reader
72 .read("load_block_header", move |tx| queries::load_block_header(tx, block_num))
73 .await
74 }
75
76 pub async fn load_protocol_config(
78 &self,
79 commitment: miden_protocol::Word,
80 ) -> Result<Option<ProtocolConfig>, DatabaseError> {
81 self.reader
82 .read("load_protocol_config", move |tx| queries::load_protocol_config(tx, commitment))
83 .await
84 }
85
86 pub(crate) async fn load_initial_metrics(&self) -> Result<InitialMetrics, DatabaseError> {
89 self.reader
90 .read("load_initial_metrics", |tx| {
91 Ok(InitialMetrics {
92 chain_tip: queries::load_chain_tip(tx)?
93 .map_or(0, |header| header.block_num().as_u32()),
94 validated_transactions: u64::try_from(queries::count_validated_transactions(
95 tx,
96 )?)
97 .unwrap_or(0),
98 signed_blocks: u64::try_from(queries::count_signed_blocks(tx)?).unwrap_or(0),
99 })
100 })
101 .await
102 }
103
104 #[cfg(test)]
110 pub(crate) async fn count_validated_transactions(&self) -> Result<i64, DatabaseError> {
111 self.reader
112 .read("count_validated_transactions", queries::count_validated_transactions)
113 .await
114 }
115
116 pub async fn load_private_record(
118 &self,
119 transaction_id: TransactionId,
120 ) -> Result<Option<StoredPrivateRecord>, DatabaseError> {
121 self.reader
122 .read("load_private_record", move |tx| {
123 queries::load_private_record(tx, transaction_id)
124 })
125 .await
126 }
127
128 pub async fn load_private_records_by_key_epoch(
130 &self,
131 key_epoch: StorageKeyEpoch,
132 ) -> Result<Vec<StoredPrivateRecord>, DatabaseError> {
133 self.reader
134 .read("load_private_records_by_key_epoch", move |tx| {
135 queries::load_private_records_by_key_epoch(tx, key_epoch)
136 })
137 .await
138 }
139
140 pub async fn load_private_records_by_setup_context(
142 &self,
143 setup_context_id: [u8; 32],
144 ) -> Result<Vec<StoredPrivateRecord>, DatabaseError> {
145 self.reader
146 .read("load_private_records_by_setup_context", move |tx| {
147 queries::load_private_records_by_setup_context(tx, setup_context_id)
148 })
149 .await
150 }
151
152 pub(crate) async fn load_all_transactions(
154 &self,
155 ) -> Result<Vec<StoredPrivateRecord>, DatabaseError> {
156 self.reader.read("load_all_transactions", queries::load_all_transactions).await
157 }
158}
159
160pub struct ValidatorDbWriter {
166 writer: DbWriter,
167 reader: ValidatorDbReader,
168}
169
170impl std::ops::Deref for ValidatorDbWriter {
171 type Target = ValidatorDbReader;
172
173 fn deref(&self) -> &Self::Target {
174 &self.reader
175 }
176}
177
178impl ValidatorDbWriter {
179 pub fn reader(&self) -> ValidatorDbReader {
182 self.reader.clone()
183 }
184
185 #[miden_instrument(
188 target = COMPONENT,
189 )]
190 pub async fn insert_validated_private_transaction(
191 &self,
192 record: StoredPrivateRecord,
193 ) -> Result<usize, DatabaseError> {
194 self.writer
195 .write("insert_validated_private_transaction", move |tx| {
196 queries::insert_validated_private_transaction(tx, &record)
197 })
198 .await
199 }
200
201 #[miden_instrument(
210 target = COMPONENT,
211 )]
212 pub(crate) async fn upsert_block_header_with_protocol_config(
213 &self,
214 header: BlockHeader,
215 protocol_config: Option<ProtocolConfig>,
216 ) -> Result<(), DatabaseError> {
217 self.writer
218 .write("upsert_block_header_with_protocol_config", move |tx| {
219 let commitment = header.protocol_config_commitment();
220 let config = if let Some(config) = protocol_config {
221 let calculated = config.to_commitment();
222 if calculated != commitment {
223 return Err(invalid_protocol_config(format!(
224 "protocol config commitment mismatch: expected {commitment}, got \
225 {calculated}"
226 )));
227 }
228 config
229 } else {
230 queries::load_protocol_config(tx, commitment)?.ok_or_else(|| {
231 invalid_protocol_config(format!(
232 "protocol config {commitment} is not stored"
233 ))
234 })?
235 };
236
237 let block_number = header.block_num();
238 let previous = queries::load_protocol_config_commitment_before(tx, block_number)?;
239 if previous != Some(commitment) {
240 queries::insert_protocol_config(tx, &config, block_number)?;
241 }
242
243 queries::upsert_block_header(tx, &header)
244 })
245 .await
246 }
247}
248
249fn invalid_protocol_config(message: String) -> DatabaseError {
250 DatabaseError::deserialization(
251 "ProtocolConfig",
252 io::Error::new(io::ErrorKind::InvalidData, message),
253 )
254}
255
256#[cfg(test)]
258pub(crate) async fn delete_protocol_config_for_test(
259 db: &ValidatorDbWriter,
260 commitment: miden_protocol::Word,
261) -> Result<(), DatabaseError> {
262 db.writer
263 .write("delete_protocol_config_for_test", move |tx| {
264 tx.execute("DELETE FROM protocol_configs WHERE commitment = ?1", &[&commitment])?;
265 Ok::<_, DatabaseError>(())
266 })
267 .await
268}
269
270#[miden_instrument(
275 target = COMPONENT,
276)]
277pub async fn load(database_filepath: PathBuf) -> Result<ValidatorDbWriter, DatabaseError> {
278 load_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size()).await
279}
280
281#[miden_instrument(
284 target = COMPONENT,
285)]
286pub async fn load_with_pool_size(
287 database_filepath: PathBuf,
288 connection_pool_size: NonZeroUsize,
289) -> Result<ValidatorDbWriter, DatabaseError> {
290 verify_latest_schema(&database_filepath)?;
291
292 open_with_pool_size(&database_filepath, connection_pool_size)
293}
294
295#[miden_instrument(
297 target = COMPONENT,
298)]
299pub async fn setup(database_filepath: PathBuf) -> Result<ValidatorDbWriter, DatabaseError> {
300 setup_with_pool_size(database_filepath, miden_node_db::default_connection_pool_size()).await
301}
302
303#[miden_instrument(
305 target = COMPONENT,
306)]
307async fn setup_with_pool_size(
308 database_filepath: PathBuf,
309 connection_pool_size: NonZeroUsize,
310) -> Result<ValidatorDbWriter, DatabaseError> {
311 bootstrap_database(&database_filepath)?;
312
313 open_with_pool_size(&database_filepath, connection_pool_size)
314}
315
316#[miden_instrument(
321 target = COMPONENT,
322 fields(path = database_filepath),
323 err,
324)]
325pub async fn bootstrap(
326 database_filepath: PathBuf,
327 connection_pool_size: NonZeroUsize,
328 genesis_header: BlockHeader,
329 protocol_config: ProtocolConfig,
330) -> Result<(), DatabaseError> {
331 let db = setup_with_pool_size(database_filepath, connection_pool_size).await?;
332
333 db.upsert_block_header_with_protocol_config(genesis_header, Some(protocol_config))
334 .await
335}
336
337#[miden_instrument(
339 target = COMPONENT,
340)]
341pub fn migrate(database_filepath: impl AsRef<Path>) -> Result<(), DatabaseError> {
342 migrate_database(database_filepath.as_ref())?;
343 Ok(())
344}
345
346fn open_with_pool_size(
347 database_filepath: &Path,
348 connection_pool_size: NonZeroUsize,
349) -> Result<ValidatorDbWriter, DatabaseError> {
350 let (writer, reader) =
351 miden_node_db::sqlite::open_with_pool_size(database_filepath, connection_pool_size)?;
352 info!(
353 target: LOG_TARGET,
354 "Connected to the database",
355 path = database_filepath,
356 db.sqlite.connection_pool_size = connection_pool_size.get()
357 );
358 Ok(ValidatorDbWriter {
359 writer,
360 reader: ValidatorDbReader { reader },
361 })
362}
363
364#[cfg(test)]
365mod tests {
366 mod protocol_config_history;
367
368 use miden_node_utils::fee::{test_fee_params, test_protocol_config};
369 use miden_protocol::Word;
370 use miden_protocol::asset::AssetId;
371 use miden_protocol::block::{BlockHeader, ValidatorConfig};
372 use miden_protocol::crypto::dsa::ecdsa_k256_keccak::SigningKey;
373 use miden_protocol::protocol_config::ProtocolConfig;
374 use miden_protocol::testing::account_id::ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET_1;
375 use miden_protocol::utils::serde::Deserializable;
376 use rand_chacha_03::ChaCha20Rng;
377 use rand_chacha_03::rand_core::SeedableRng;
378
379 use super::*;
380 use crate::private_record::test_private_record_sealer;
381 use crate::storage_key::tests::operator_keys;
382 use crate::{
383 PrivateRecordChainId,
384 PrivateRecordCombiner,
385 PrivateRecordContext,
386 PrivateRecordError,
387 PrivateRecordFormatVersion,
388 PrivateRecordId,
389 PrivateRecordSealer,
390 PrivateRecordShareRequest,
391 };
392
393 const CHAIN_ID: PrivateRecordChainId = PrivateRecordChainId::new([1; 32]);
394 const KEY_EPOCH: StorageKeyEpoch = StorageKeyEpoch::new([2; 32]);
395 const SETUP_CONTEXT_ID: [u8; 32] = [4; 32];
396
397 fn record_id(transaction_id: TransactionId) -> PrivateRecordId {
398 let signer = SigningKey::read_from_bytes(&[7; 32]).unwrap();
399 PrivateRecordId::new(transaction_id, &signer.public_key())
400 }
401
402 fn private_record(transaction_id: TransactionId, seed: u8) -> StoredPrivateRecord {
403 let context = PrivateRecordContext::new(CHAIN_ID, KEY_EPOCH, transaction_id);
404 let mut rng = ChaCha20Rng::from_seed([seed; 32]);
405 test_private_record_sealer(KEY_EPOCH, SETUP_CONTEXT_ID)
406 .seal(&mut rng, record_id(transaction_id), context, b"private transaction inputs")
407 .unwrap()
408 }
409
410 fn genesis_header(config: &miden_protocol::protocol_config::ProtocolConfig) -> BlockHeader {
411 miden_node_store::GenesisState::new(
412 vec![],
413 test_fee_params(),
414 0,
415 ValidatorConfig::new(vec![SigningKey::new().public_key()], 1).unwrap(),
416 config.clone(),
417 )
418 .into_block()
419 .unwrap()
420 .inner()
421 .header()
422 .clone()
423 }
424
425 fn header_with_next_timestamp(header: &BlockHeader) -> BlockHeader {
426 BlockHeader::new(
427 header.prev_block_commitment(),
428 header.block_num(),
429 header.chain_commitment(),
430 header.account_root(),
431 header.nullifier_root(),
432 header.note_root(),
433 header.tx_commitment(),
434 header.validator_config().clone(),
435 header.fee_parameters().clone(),
436 header.protocol_config_commitment(),
437 header.next_protocol_config().cloned(),
438 header.timestamp() + 1,
439 )
440 }
441
442 #[test]
443 fn migrate_rejects_missing_database() {
444 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
445 let db_path = temp_dir.path().join("validator.sqlite3");
446
447 let err = migrate(db_path.clone()).expect_err("missing database should fail");
448
449 assert!(matches!(err, DatabaseError::Migration(_)), "unexpected error: {err:?}");
450 assert!(!db_path.exists());
451 }
452
453 #[tokio::test]
455 async fn migration_preserves_headers_and_private_records() {
456 let temp_dir = tempfile::tempdir().unwrap();
457 let db_path = temp_dir.path().join("validator.sqlite3");
458 miden_node_db::migration::Migrator::builder()
459 .unwrap()
460 .push_sql("001_initial", include_str!("migrations/001_initial.sql"))
461 .unwrap()
462 .build()
463 .unwrap()
464 .bootstrap(&db_path)
465 .unwrap();
466
467 let config = test_protocol_config();
468 let header = genesis_header(&config);
469 let transaction_id = TransactionId::from_raw(Word::from([1u32, 2, 3, 4]));
470 let record = private_record(transaction_id, 1);
471 let db = open_with_pool_size(&db_path, NonZeroUsize::new(2).unwrap()).unwrap();
472 let stored_header = header.clone();
473 db.writer
474 .write("seed legacy header", move |tx| queries::upsert_block_header(tx, &stored_header))
475 .await
476 .unwrap();
477 db.insert_validated_private_transaction(record.clone()).await.unwrap();
478 drop(db);
479
480 assert!(load(db_path.clone()).await.is_err(), "the old schema requires migration");
481 migrate(&db_path).unwrap();
482 migrate(&db_path).expect("migration should also accept the latest schema");
483
484 let db = load(db_path).await.unwrap();
485 assert_eq!(db.load_chain_tip().await.unwrap(), Some(header.clone()));
486 assert_eq!(db.load_all_transactions().await.unwrap(), vec![record]);
487 let migrated = db.load_private_record(transaction_id).await.unwrap().unwrap();
488 assert_eq!(migrated.context().format_version(), PrivateRecordFormatVersion::V1);
489 migrated.verify_encrypted_record_key().unwrap();
490 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), None);
491
492 db.upsert_block_header_with_protocol_config(header, Some(config.clone()))
493 .await
494 .unwrap();
495 assert_eq!(
496 protocol_config_history::history(&db).await,
497 vec![(0, config.to_commitment(), config.clone())]
498 );
499 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
500 }
501
502 #[tokio::test]
503 async fn setup_creates_database_that_load_accepts() {
504 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
505 let db_path = temp_dir.path().join("validator.sqlite3");
506
507 setup(db_path.clone()).await.expect("setup should bootstrap the database");
508 load(db_path).await.expect("load should accept a bootstrapped database");
509 }
510
511 #[tokio::test]
512 async fn setup_creates_protocol_config_storage() {
513 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
514 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
515
516 let row_count = db
517 .reader
518 .reader
519 .read("protocol_config_storage", |tx| {
520 Ok::<_, DatabaseError>(
521 tx.query("SELECT COUNT(*) FROM protocol_configs", &[], |row| {
522 row.get::<i64>(0)
523 })?
524 .into_iter()
525 .next()
526 .expect("COUNT always returns one row"),
527 )
528 })
529 .await
530 .unwrap();
531
532 assert_eq!(row_count, 0);
533 }
534
535 #[tokio::test]
536 async fn block_header_and_protocol_config_are_persisted_together() {
537 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
538 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
539 let config = test_protocol_config();
540 let header = genesis_header(&config);
541
542 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
543 .await
544 .unwrap();
545
546 assert_eq!(db.load_chain_tip().await.unwrap(), Some(header));
547 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
548 }
549
550 #[tokio::test]
551 async fn duplicate_supplied_protocol_config_is_accepted() {
552 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
553 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
554 let config = test_protocol_config();
555 let header = genesis_header(&config);
556 let replacement = header_with_next_timestamp(&header);
557
558 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
559 .await
560 .unwrap();
561 db.upsert_block_header_with_protocol_config(replacement.clone(), Some(config.clone()))
562 .await
563 .expect("a duplicate supplied protocol config should be accepted");
564
565 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), Some(replacement));
566 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
567 }
568
569 #[tokio::test]
570 async fn known_protocol_config_can_be_omitted() {
571 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
572 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
573 let config = test_protocol_config();
574 let header = genesis_header(&config);
575 let replacement = header_with_next_timestamp(&header);
576
577 db.upsert_block_header_with_protocol_config(header.clone(), Some(config.clone()))
578 .await
579 .unwrap();
580 db.upsert_block_header_with_protocol_config(replacement.clone(), None)
581 .await
582 .expect("a stored protocol config should not need to be supplied again");
583
584 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), Some(replacement));
585 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), Some(config));
586 }
587
588 #[tokio::test]
589 async fn unknown_protocol_config_rolls_back_block_header() {
590 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
591 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
592 let config = test_protocol_config();
593 let header = genesis_header(&config);
594
595 db.upsert_block_header_with_protocol_config(header.clone(), None)
596 .await
597 .expect_err("an unknown protocol config must reject the header");
598
599 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
600 assert_eq!(db.load_protocol_config(config.to_commitment()).await.unwrap(), None);
601 }
602
603 #[tokio::test]
604 async fn mismatched_protocol_config_rolls_back_block_header() {
605 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
606 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
607 let expected = test_protocol_config();
608 let header = genesis_header(&expected);
609 let mismatched = ProtocolConfig::current(AssetId::new_fungible(
610 ACCOUNT_ID_PUBLIC_FUNGIBLE_FAUCET_1.try_into().unwrap(),
611 ))
612 .unwrap();
613
614 db.upsert_block_header_with_protocol_config(header.clone(), Some(mismatched.clone()))
615 .await
616 .expect_err("a mismatched config must reject the header transaction");
617
618 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
619 assert_eq!(db.load_protocol_config(expected.to_commitment()).await.unwrap(), None);
620 assert_eq!(db.load_protocol_config(mismatched.to_commitment()).await.unwrap(), None);
621 }
622
623 #[tokio::test]
624 async fn block_header_insertion_failure_rolls_back_new_protocol_config() {
625 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
626 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
627 let config = test_protocol_config();
628 let commitment = config.to_commitment();
629 let header = genesis_header(&config);
630
631 db.writer
632 .write("reject_block_header_inserts", |tx| {
633 tx.execute(
634 "CREATE TRIGGER reject_block_header_insert
635 BEFORE INSERT ON block_headers
636 BEGIN
637 SELECT RAISE(ABORT, 'block header insertion rejected');
638 END;",
639 &[],
640 )?;
641 Ok::<_, DatabaseError>(())
642 })
643 .await
644 .unwrap();
645
646 db.upsert_block_header_with_protocol_config(header.clone(), Some(config))
647 .await
648 .expect_err("a block header insertion failure must reject the transaction");
649
650 assert_eq!(db.load_block_header(header.block_num()).await.unwrap(), None);
651 assert_eq!(db.load_protocol_config(commitment).await.unwrap(), None);
652 }
653
654 #[tokio::test]
655 async fn transaction_exists_detects_validated_transactions() {
656 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
657 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
658
659 let validated_id = TransactionId::from_raw(Word::try_from([1u64, 2, 3, 4]).unwrap());
660 let unknown_id = TransactionId::from_raw(Word::try_from([5u64, 6, 7, 8]).unwrap());
661
662 db.insert_validated_private_transaction(private_record(validated_id, 1))
663 .await
664 .unwrap();
665
666 assert!(
667 db.transaction_exists(validated_id).await.unwrap(),
668 "an inserted transaction id should be reported as existing"
669 );
670 assert!(
671 !db.transaction_exists(unknown_id).await.unwrap(),
672 "an unknown transaction id should not be reported as existing"
673 );
674 }
675
676 #[tokio::test]
679 async fn find_unvalidated_transactions_returns_only_missing_ids() {
680 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
681 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
682
683 let ids = (1u64..=4)
684 .map(|i| TransactionId::from_raw(Word::try_from([i, i, i, i]).unwrap()))
685 .collect::<Vec<_>>();
686
687 db.insert_validated_private_transaction(private_record(ids[1], 1))
689 .await
690 .unwrap();
691 db.insert_validated_private_transaction(private_record(ids[3], 2))
692 .await
693 .unwrap();
694
695 let unvalidated = db.find_unvalidated_transactions(ids.clone()).await.unwrap();
696 assert_eq!(unvalidated, vec![ids[0], ids[2]]);
697
698 assert!(db.find_unvalidated_transactions(vec![]).await.unwrap().is_empty());
700 }
701
702 #[tokio::test]
703 async fn load_initial_metrics_reports_persisted_state() {
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
707 let metrics = db.load_initial_metrics().await.unwrap();
709 assert_eq!(metrics.chain_tip, 0);
710 assert_eq!(metrics.validated_transactions, 0);
711 assert_eq!(metrics.signed_blocks, 0);
712 }
713
714 #[tokio::test]
715 async fn private_record_indexes_work() {
716 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
717 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
718 let transaction_id = TransactionId::from_raw(Word::from([5u32, 6, 7, 8]));
719 let record = private_record(transaction_id, 9);
720
721 let expected = record.clone();
722 db.insert_validated_private_transaction(record).await.unwrap();
723
724 let by_record = db.load_private_record(transaction_id).await.unwrap();
725 assert_eq!(by_record, Some(expected.clone()));
726
727 let by_epoch = db.load_private_records_by_key_epoch(KEY_EPOCH).await.unwrap();
728 assert_eq!(by_epoch, vec![expected.clone()]);
729
730 let by_setup = db.load_private_records_by_setup_context(SETUP_CONTEXT_ID).await.unwrap();
731 assert_eq!(by_setup, vec![expected.clone()]);
732 }
733
734 #[tokio::test]
735 async fn private_record_rejects_unsupported_formats() {
736 let temp_dir = tempfile::tempdir().unwrap();
737 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
738 let transaction_id = TransactionId::from_raw(Word::from([5u32, 6, 7, 8]));
739 db.insert_validated_private_transaction(private_record(transaction_id, 9))
740 .await
741 .unwrap();
742
743 for format_version in [2_u32, 3, u32::MAX] {
744 db.writer
745 .write("set unsupported record format", move |tx| {
746 tx.execute(
747 "UPDATE validated_transactions SET format_version = ? WHERE id = ?",
748 &[&i64::from(format_version), &transaction_id],
749 )
750 })
751 .await
752 .unwrap();
753
754 let error = db.load_private_record(transaction_id).await.unwrap_err();
755 assert!(matches!(
756 error,
757 DatabaseError::ConversionSqlToRust { inner: Some(source), .. }
758 if matches!(
759 source.downcast_ref::<PrivateRecordError>(),
760 Some(PrivateRecordError::UnsupportedFormat(version))
761 if *version == format_version,
762 ),
763 ));
764 }
765 }
766
767 #[tokio::test]
768 async fn validated_private_transactions_are_loaded_in_insertion_order() {
769 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
770 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
771 let transaction_ids = [
772 TransactionId::from_raw(Word::from([9u32, 0, 0, 0])),
773 TransactionId::from_raw(Word::from([1u32, 0, 0, 0])),
774 TransactionId::from_raw(Word::from([5u32, 0, 0, 0])),
775 ];
776 let records = transaction_ids
777 .into_iter()
778 .zip([1u8, 2, 3])
779 .map(|(transaction_id, seed)| private_record(transaction_id, seed))
780 .collect::<Vec<_>>();
781
782 for record in records.clone() {
783 db.insert_validated_private_transaction(record).await.unwrap();
784 }
785
786 let loaded = db.load_all_transactions().await.unwrap();
787
788 assert_eq!(loaded, records);
789 }
790
791 #[tokio::test]
792 async fn stored_private_record_opens_with_threshold_shares() {
793 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
794 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
795 let operators = operator_keys();
796 let transaction_id = TransactionId::from_raw(Word::from([9u32, 10, 11, 12]));
797 let context = PrivateRecordContext::new(CHAIN_ID, operators[0].key_epoch(), transaction_id);
798 let plaintext = b"private transaction inputs";
799 let mut seal_rng = ChaCha20Rng::from_seed([40; 32]);
800 let record = PrivateRecordSealer::from_operator_key(&operators[0])
801 .seal(&mut seal_rng, record_id(transaction_id), context, plaintext)
802 .unwrap();
803 db.insert_validated_private_transaction(record).await.unwrap();
804
805 let stored = db.load_private_record(transaction_id).await.unwrap().unwrap();
806 let request = PrivateRecordShareRequest::for_record(&stored);
807 let mut first_rng = ChaCha20Rng::from_seed([41; 32]);
808 let mut second_rng = ChaCha20Rng::from_seed([42; 32]);
809 let shares = [
810 operators[0]
811 .issue_private_record_share(&mut first_rng, &request, &stored)
812 .unwrap(),
813 operators[1]
814 .issue_private_record_share(&mut second_rng, &request, &stored)
815 .unwrap(),
816 ];
817
818 let opened = PrivateRecordCombiner::from_operator_key(&operators[2])
819 .unwrap()
820 .open(&request, &stored, &shares)
821 .unwrap();
822 assert_eq!(opened.as_slice(), plaintext);
823 }
824
825 #[tokio::test]
826 async fn private_record_schema_has_required_indexes() {
827 let temp_dir = tempfile::tempdir().expect("failed to create temp directory");
828 let db = setup(temp_dir.path().join("validator.sqlite3")).await.unwrap();
829
830 let schema = db
831 .reader
832 .reader
833 .read("private_record_schema", |tx| {
834 tx.query(
835 "SELECT sql FROM sqlite_schema \
836 WHERE tbl_name = 'validated_transactions' AND sql IS NOT NULL \
837 ORDER BY name",
838 &[],
839 |row| row.get::<String>(0),
840 )
841 })
842 .await
843 .unwrap()
844 .join("\n");
845
846 assert!(schema.contains("insertion_sequence INTEGER PRIMARY KEY AUTOINCREMENT"));
847 assert!(schema.contains("id BLOB NOT NULL UNIQUE"));
848 assert!(schema.contains("idx_validated_transactions_key_epoch"));
849 assert!(schema.contains("idx_validated_transactions_setup_context_id"));
850 }
851}