Skip to main content

miden_validator/db/
mod.rs

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// VALIDATOR DATABASE
20// ================================================================================================
21
22/// Read-only handle to the validator database.
23///
24/// Wraps the framework [`DbReader`] and exposes every read query as a method. Cloneable, and handed
25/// to read-only components (the administration API); it has no write methods, so those components
26/// cannot mutate the database.
27#[derive(Clone)]
28pub struct ValidatorDbReader {
29    reader: DbReader,
30}
31
32impl ValidatorDbReader {
33    /// Returns whether a transaction with the given id has already been validated.
34    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    /// Returns the subset of `tx_ids` that this validator has not validated yet.
44    ///
45    /// An empty result means all supplied transaction ids have been validated in the past.
46    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    /// Loads the chain tip, or `None` if no block header has been persisted yet (i.e. bootstrap has
58    /// not been run).
59    #[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    /// Loads the block header at the given height, or `None` if no block header is stored there.
67    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    /// Loads the protocol configuration with the given commitment.
77    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    /// Reads the values the server's in-memory counters start from, all within a single read
87    /// transaction so that they describe one consistent database state.
88    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    /// Returns the total number of validated transactions.
105    ///
106    /// Production code seeds its counter from [`Self::load_initial_metrics`] and tracks it in
107    /// memory from there, so this standalone count only backs test assertions about what was
108    /// actually persisted.
109    #[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    /// Loads one encrypted private record by transaction id.
117    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    /// Loads the encrypted private records sealed under one storage key epoch.
129    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    /// Loads the encrypted private records belonging to one Golden setup context.
141    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    /// Loads all validated private transactions in insertion order.
153    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
160/// Write handle to the validator database.
161///
162/// Wraps the framework [`DbWriter`] and additionally holds a [`ValidatorDbReader`], so it exposes
163/// the write queries directly and every read query through `Deref`. **Not `Clone`**: writes have a
164/// single owner, matching SQLite's single-writer model.
165pub 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    /// Returns a read-only handle onto the same connection pool, for handing to components that
180    /// must not be able to write.
181    pub fn reader(&self) -> ValidatorDbReader {
182        self.reader.clone()
183    }
184
185    /// Inserts a validated transaction and its encrypted private record, returning the number of
186    /// inserted rows. The count is zero if the transaction was already recorded.
187    #[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    /// Persists a block header and its configuration activation in one transaction.
202    ///
203    /// Records an activation if the configuration differs from the preceding activation.
204    /// A replacement at the current tip must retain its active configuration.
205    /// Callers must validate block order before this method runs.
206    ///
207    /// If `protocol_config` is absent, the configuration must already be stored
208    /// otherwise an error is returned.
209    #[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/// Deletes a stored protocol configuration for a test.
257#[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// LIFECYCLE
271// ================================================================================================
272
273/// Opens a connection pool after verifying that the database is at the latest schema version.
274#[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/// Opens a connection pool with a specific pool size after verifying that the database is at the
282/// latest schema version.
283#[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/// Creates a new database, applies all migrations, and opens a connection pool.
296#[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/// Creates a new database with a specific pool size and applies all migrations.
304#[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/// Creates and initializes the database, then seeds it with the genesis block header as the chain
317/// tip.
318///
319/// Returns an error if the database has already been bootstrapped.
320#[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/// Applies all pending migrations to an existing DB.
338#[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    /// The protocol configuration migration must preserve headers and private records.
454    #[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    /// The `rarray`-based lookup must return exactly the ids that are absent, preserving the order
677    /// they were supplied in.
678    #[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        // Validate the second and fourth ids only.
688        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        // An empty request must not error and must return nothing.
699        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        // A freshly bootstrapped database is empty.
708        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}