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
19pub(crate) use queries::{ListTransactionsParams, ListedTransaction};
20
21// VALIDATOR DATABASE
22// ================================================================================================
23
24/// Read-only handle to the validator database.
25///
26/// Wraps the framework [`DbReader`] and exposes every read query as a method. Cloneable, and handed
27/// to read-only components (the administration API); it has no write methods, so those components
28/// cannot mutate the database.
29#[derive(Clone)]
30pub struct ValidatorDbReader {
31    reader: DbReader,
32}
33
34impl ValidatorDbReader {
35    /// Returns whether a transaction with the given id has already been validated.
36    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    /// Returns the subset of `tx_ids` that this validator has not validated yet.
46    ///
47    /// An empty result means all supplied transaction ids have been validated in the past.
48    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    /// Loads the chain tip, or `None` if no block header has been persisted yet (i.e. bootstrap has
60    /// not been run).
61    #[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    /// Loads the block header at the given height, or `None` if no block header is stored there.
69    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    /// Loads the protocol configuration with the given commitment.
79    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    /// Reads the values the server's in-memory counters start from, all within a single read
89    /// transaction so that they describe one consistent database state.
90    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    /// Returns the total number of validated transactions.
107    ///
108    /// Production code seeds its counter from [`Self::load_initial_metrics`] and tracks it in
109    /// memory from there, so this standalone count only backs test assertions about what was
110    /// actually persisted.
111    #[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    /// Loads one encrypted private record by transaction id.
119    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    /// Loads the encrypted private records sealed under one storage key epoch.
131    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    /// Loads the encrypted private records belonging to one Golden setup context.
143    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    /// Loads one page of committed transactions in chronological order i.e. `(block_num,
155    /// block_tx_index)`, with the chain tip and optional records from the same snapshot.
156    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, &params)
163            })
164            .await
165    }
166}
167
168/// Write handle to the validator database.
169///
170/// Wraps the framework [`DbWriter`] and additionally holds a [`ValidatorDbReader`], so it exposes
171/// the write queries directly and every read query through `Deref`. **Not `Clone`**: writes have a
172/// single owner, matching SQLite's single-writer model.
173pub 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    /// Returns a read-only handle onto the same connection pool, for handing to components that
188    /// must not be able to write.
189    pub fn reader(&self) -> ValidatorDbReader {
190        self.reader.clone()
191    }
192
193    /// Inserts a validated transaction and its encrypted private record, returning the number of
194    /// inserted rows. The count is zero if the transaction was already recorded.
195    #[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    /// Persists a block header and its configuration activation in one transaction.
210    ///
211    /// See [`record_protocol_config_activation`] for how the activation is recorded.
212    /// Callers must validate block order before this method runs.
213    ///
214    /// The write replaces the header row. If the height holds a block, the `ON DELETE CASCADE` on
215    /// `block_transactions` deletes the links of that block, and this method does not link the
216    /// transactions of the new block. Only tests use this method. Server code uses
217    /// [`Self::insert_signed_block`] and [`Self::replace_signed_block`].
218    #[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    /// Persists a signed block's header, its configuration activation, and the links from the
236    /// block's transactions to their in-block positions, all in one database transaction so the
237    /// three stay consistent.
238    ///
239    /// The height must not already hold a block; use [`Self::replace_signed_block`] to replace
240    /// one.
241    #[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    /// Replaces the block signed at `header`'s height, in one database transaction: deletes the
258    /// replaced block — unlinking its transactions via the `ON DELETE CASCADE` on
259    /// `block_transactions`, so transactions dropped by the replacement do not keep a stale link —
260    /// then persists the new header and links exactly as [`Self::insert_signed_block`] does.
261    #[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
279/// Persists a signed block's header, records its configuration activation, and links its
280/// transactions, within the caller's database transaction: the shared tail of
281/// [`ValidatorDbWriter::insert_signed_block`] and [`ValidatorDbWriter::replace_signed_block`].
282fn 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
293/// Records the configuration activation for a block, within the caller's database transaction.
294///
295/// An activation is recorded only if the configuration differs from the preceding activation.
296/// A replacement at the current tip therefore retains its active configuration.
297///
298/// If `protocol_config` is absent, the configuration must already be stored, otherwise an error is
299/// returned.
300fn 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/// Deletes a stored protocol configuration for a test.
337#[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// LIFECYCLE
351// ================================================================================================
352
353/// Opens a connection pool after verifying that the database is at the latest schema version.
354#[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/// Opens a connection pool with a specific pool size after verifying that the database is at the
362/// latest schema version.
363#[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/// Creates a new database, applies all migrations, and opens a connection pool.
376#[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/// Creates a new database with a specific pool size and applies all migrations.
384#[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/// Creates and initializes the database, then seeds it with the genesis block header as the chain
397/// tip.
398///
399/// Returns an error if the database has already been bootstrapped.
400#[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/// Applies all pending migrations to an existing DB.
417#[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    /// The protocol configuration migration must preserve headers and private records.
533    #[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    /// The `rarray`-based lookup must return exactly the ids that are absent, preserving the order
756    /// they were supplied in.
757    #[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        // Validate the second and fourth ids only.
767        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        // An empty request must not error and must return nothing.
778        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        // A freshly bootstrapped database is empty.
787        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    /// Validated transactions that are not part of a signed block have no position in the committed
847    /// order, so the listing does not surface them at all.
848    #[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        // They remain reachable by transaction id.
870        assert!(db.load_private_record(transaction_ids[0]).await.unwrap().is_some());
871    }
872
873    /// A signed block links its transactions in block order; replacing the block at the same height
874    /// deletes the replaced block's links along with its header.
875    #[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        // Replace the block at the same height with one that only includes the third transaction.
916        // The two transactions the replacement drops go back to being uncommitted, and stop being
917        // listed.
918        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    /// A listing keeps its records and chain tip consistent even when the tip is replaced and
933    /// advanced while its read snapshot is open.
934    #[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, &params))
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    /// A full sweep pages through committed transactions in committed order, honoring the row limit
991    /// exactly, with each page resuming one position past the last row of the previous one.
992    #[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        // Blocks 1 and 2 include two transactions each; the fifth is never committed. The later
1005        // insertion is committed in the earlier block, to prove the listing follows committed order
1006        // rather than insertion order.
1007        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        // Sweep with a limit of three: the first page ends mid-block, and the next page resumes one
1027        // position past the last row returned.
1028        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        // Bounding by block number: `start` excludes block 1, `block_to` excludes block 2.
1046        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    /// Linking a transaction this validator never validated fails the foreign key on
1083    /// `block_transactions`, so a signed block cannot silently reference unknown transactions.
1084    #[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}