Skip to main content

espresso_node/
persistence.rs

1//! Sequencer node persistence.
2//!
3//! This module implements the persistence required for a sequencer node to rejoin the network and
4//! resume participating in consensus, in the event that its process crashes or is killed and loses
5//! all in-memory state.
6//!
7//! This is distinct from the query service persistent storage found in the `api` module, which is
8//! an extension that node operators can opt into. This module defines the minimum level of
9//! persistence which is _required_ to run a node.
10
11use std::collections::HashMap;
12
13use alloy::primitives::{Address, U256};
14use anyhow::Context;
15use async_trait::async_trait;
16use espresso_types::{
17    PubKey,
18    v0_3::{ChainConfig, RegisteredValidator},
19};
20
21pub mod fs;
22pub mod no_storage;
23mod persistence_metrics;
24pub mod sql;
25
26/// RegisteredValidator without x25519_key/p2p_addr fields.
27/// Used for migrating data written before x25519 support was added.
28#[derive(serde::Serialize, serde::Deserialize)]
29pub(crate) struct RegisteredValidatorNoX25519 {
30    pub account: Address,
31    pub stake_table_key: PubKey,
32    pub state_ver_key: hotshot_types::light_client::StateVerKey,
33    pub stake: U256,
34    pub commission: u16,
35    pub delegators: HashMap<Address, U256>,
36    pub authenticated: bool,
37}
38
39impl RegisteredValidatorNoX25519 {
40    pub fn migrate(self) -> RegisteredValidator<PubKey> {
41        RegisteredValidator {
42            account: self.account,
43            stake_table_key: Some(self.stake_table_key),
44            state_ver_key: Some(self.state_ver_key),
45            stake: self.stake,
46            commission: self.commission,
47            delegators: self.delegators,
48            authenticated: self.authenticated,
49            x25519_key: None,
50            p2p_addr: None,
51        }
52    }
53}
54
55// TODO: replace with a proper data migration if it is ever deemed worthwhile.
56// Field order must match `RegisteredValidator<KEY>` in v0_3::stake_table exactly (bincode is order-sensitive).
57#[derive(serde::Serialize, serde::Deserialize)]
58pub(crate) struct RegisteredValidatorPreOption {
59    pub account: Address,
60    pub stake_table_key: PubKey,
61    pub state_ver_key: hotshot_types::light_client::StateVerKey,
62    pub stake: U256,
63    pub commission: u16,
64    pub delegators: HashMap<Address, U256>,
65    pub authenticated: bool,
66    pub x25519_key: Option<hotshot_types::x25519::PublicKey>,
67    pub p2p_addr: Option<hotshot_types::addr::NetAddr>,
68}
69
70impl RegisteredValidatorPreOption {
71    pub fn migrate(self) -> RegisteredValidator<PubKey> {
72        RegisteredValidator {
73            account: self.account,
74            stake_table_key: Some(self.stake_table_key),
75            state_ver_key: Some(self.state_ver_key),
76            stake: self.stake,
77            commission: self.commission,
78            delegators: self.delegators,
79            authenticated: self.authenticated,
80            x25519_key: self.x25519_key,
81            p2p_addr: self.p2p_addr,
82        }
83    }
84}
85
86// TODO: replace with a proper data migration if it is ever deemed worthwhile.
87// Layout written after BLS fix (381b0f9207) but before Schnorr Option fix.
88// stake_table_key is Option<KEY>, state_ver_key is raw StateVerKey.
89// Field order must match exactly (bincode is order-sensitive).
90#[derive(serde::Serialize, serde::Deserialize)]
91pub(crate) struct RegisteredValidatorPreSchnorrOption {
92    pub account: Address,
93    pub stake_table_key: Option<PubKey>,
94    pub state_ver_key: hotshot_types::light_client::StateVerKey,
95    pub stake: U256,
96    pub commission: u16,
97    pub delegators: HashMap<Address, U256>,
98    pub authenticated: bool,
99    pub x25519_key: Option<hotshot_types::x25519::PublicKey>,
100    pub p2p_addr: Option<hotshot_types::addr::NetAddr>,
101}
102
103impl RegisteredValidatorPreSchnorrOption {
104    pub fn migrate(self) -> RegisteredValidator<PubKey> {
105        RegisteredValidator {
106            account: self.account,
107            stake_table_key: self.stake_table_key,
108            state_ver_key: Some(self.state_ver_key),
109            stake: self.stake,
110            commission: self.commission,
111            delegators: self.delegators,
112            authenticated: self.authenticated,
113            x25519_key: self.x25519_key,
114            p2p_addr: self.p2p_addr,
115        }
116    }
117}
118
119/// Update a `NetworkConfig` that may have originally been persisted with an old version.
120fn migrate_network_config(
121    mut network_config: serde_json::Value,
122) -> anyhow::Result<serde_json::Value> {
123    let config = network_config
124        .get_mut("config")
125        .context("missing field `config`")?
126        .as_object_mut()
127        .context("`config` must be an object")?;
128
129    if !config.contains_key("builder_urls") {
130        // When multi-builder support was added, the configuration field `builder_url: Url` was
131        // replaced by an array `builder_urls: Vec<Url>`. If the saved config has no `builder_urls`
132        // field, it is older than this change. Populate `builder_urls` with a singleton array
133        // formed from the old value of `builder_url`, and delete the no longer used `builder_url`.
134        let url = config
135            .remove("builder_url")
136            .context("missing field `builder_url`")?;
137        config.insert("builder_urls".into(), vec![url].into());
138    }
139
140    // HotShotConfig was upgraded to include parameters for proposing and voting on upgrades.
141    // Configs which were persisted before this upgrade may be missing these parameters. This
142    // migration initializes them with a default. By default, we use JS MAX_SAFE_INTEGER for the
143    // start parameters so that nodes will never do an upgrade, unless explicitly configured
144    // otherwise.
145    if !config.contains_key("start_proposing_view") {
146        config.insert("start_proposing_view".into(), 9007199254740991u64.into());
147    }
148    if !config.contains_key("stop_proposing_view") {
149        config.insert("stop_proposing_view".into(), 0.into());
150    }
151    if !config.contains_key("start_voting_view") {
152        config.insert("start_voting_view".into(), 9007199254740991u64.into());
153    }
154    if !config.contains_key("stop_voting_view") {
155        config.insert("stop_voting_view".into(), 0.into());
156    }
157    if !config.contains_key("start_proposing_time") {
158        config.insert("start_proposing_time".into(), 9007199254740991u64.into());
159    }
160    if !config.contains_key("stop_proposing_time") {
161        config.insert("stop_proposing_time".into(), 0.into());
162    }
163    if !config.contains_key("start_voting_time") {
164        config.insert("start_voting_time".into(), 9007199254740991u64.into());
165    }
166    if !config.contains_key("stop_voting_time") {
167        config.insert("stop_voting_time".into(), 0.into());
168    }
169
170    // HotShotConfig was upgraded to include an `epoch_height` parameter. Initialize with a default
171    // if missing.
172    if !config.contains_key("epoch_height") {
173        config.insert("epoch_height".into(), 0.into());
174    }
175
176    // HotShotConfig was upgraded to include `drb_difficulty` and `drb_upgrade_difficulty` parameters. Initialize with a default
177    // if missing.
178    if !config.contains_key("drb_difficulty") {
179        config.insert("drb_difficulty".into(), 0.into());
180    }
181    if !config.contains_key("drb_upgrade_difficulty") {
182        config.insert("drb_upgrade_difficulty".into(), 0.into());
183    }
184
185    // HotShotConfig was upgraded to include `da_committeees`. Initialize with an empty `da_committees` if missing.
186    if !config.contains_key("da_committees") {
187        config.insert("da_committees".into(), serde_json::json!([]));
188    }
189
190    Ok(network_config)
191}
192
193#[async_trait]
194pub trait ChainConfigPersistence: Sized + Send + Sync {
195    async fn insert_chain_config(&mut self, chain_config: ChainConfig) -> anyhow::Result<()>;
196}
197
198#[cfg(test)]
199mod tests {
200    use std::{cmp::max, collections::BTreeMap, marker::PhantomData, sync::Arc, time::Duration};
201
202    use alloy::{
203        network::EthereumWallet,
204        primitives::{Address, U256},
205        providers::{Provider, ProviderBuilder, ext::AnvilApi},
206    };
207    use anyhow::bail;
208    use async_lock::{Mutex, RwLock};
209    use async_trait::async_trait;
210    use committable::{Commitment, Committable};
211    use espresso_contract_deployer::{
212        Contract, Contracts, DEFAULT_EXIT_ESCROW_PERIOD_SECONDS, builder::DeployerArgsBuilder,
213        network_config::light_client_genesis_from_stake_table,
214    };
215    use espresso_types::{
216        Event, L1Client, L1ClientOptions, Leaf, Leaf2, NodeState, PubKey, SeqTypes, ValidatedState,
217        traits::{
218            EventConsumer, EventsPersistenceRead, MembershipPersistence, NullEventConsumer,
219            PersistenceOptions, SequencerPersistence,
220        },
221        v0_3::{AuthenticatedValidator, EventKey, Fetcher, RegisteredValidator, StakeTableEvent},
222    };
223    use futures::{StreamExt, TryStreamExt, future::join_all};
224    use hotshot::{
225        InitializerEpochInfo,
226        types::{BLSPubKey, SignatureKey},
227    };
228    use hotshot_contract_adapter::{
229        sol_types::StakeTableV3::Delegated, stake_table::StakeTableContractVersion,
230    };
231    use hotshot_example_types::node_types::TEST_VERSIONS;
232    use hotshot_new_protocol::message::Certificate2;
233    use hotshot_query_service::{availability::BlockQueryData, testing::mocks::MOCK_UPGRADE};
234    use hotshot_types::{
235        data::{
236            DaProposal2, EpochNumber, QuorumProposal2, QuorumProposalWrapper, VidCommitment,
237            ViewNumber, ns_table::parse_ns_table, vid_commitment, vid_disperse::AvidMDisperseShare,
238        },
239        event::{EventType, HotShotAction, LeafInfo},
240        light_client::StateKeyPair,
241        message::{Proposal, UpgradeLock, convert_proposal},
242        new_protocol::CoordinatorEvent,
243        simple_certificate::{
244            CertificatePair, NextEpochQuorumCertificate2, QuorumCertificate, QuorumCertificate2,
245            UpgradeCertificate,
246        },
247        simple_vote::{
248            NextEpochQuorumData2, QuorumData2, UpgradeProposalData, VersionedVoteData, Vote2Data,
249        },
250        traits::{EncodeBytes, block_contents::BlockHeader},
251        utils::EpochTransitionIndicator,
252        vid::avidm::{AvidMScheme, init_avidm_param},
253        vote::HasViewNumber,
254    };
255    use http_client::{Client, error::ClientErr};
256    use indexmap::IndexMap;
257    use staking_cli::demo::{DelegationConfig, StakingTransactions};
258    use test_utils::reserve_tcp_port;
259    use tokio::{spawn, time::sleep};
260    use vbs::version::Version;
261    use versions::{Upgrade, version};
262
263    use crate::{
264        RECENT_STAKE_TABLES_LIMIT, SequencerApiVersion,
265        api::{
266            Options,
267            test_helpers::{STAKE_TABLE_CAPACITY_FOR_TEST, TestNetwork, TestNetworkConfigBuilder},
268        },
269        catchup::NullStateCatchup,
270        testing::{TestConfigBuilder, staking_priv_keys},
271    };
272
273    #[async_trait]
274    pub trait TestablePersistence: SequencerPersistence + MembershipPersistence {
275        type Storage: Sync;
276
277        async fn tmp_storage() -> Self::Storage;
278        fn options(storage: &Self::Storage) -> impl PersistenceOptions<Persistence = Self>;
279
280        async fn connect(storage: &Self::Storage) -> Self {
281            Self::options(storage).create().await.unwrap()
282        }
283    }
284
285    #[rstest_reuse::template]
286    #[rstest::rstest]
287    #[case(PhantomData::<crate::persistence::sql::Persistence>)]
288    #[case(PhantomData::<crate::persistence::fs::Persistence>)]
289    #[test_log::test(tokio::test(flavor = "multi_thread"))]
290    pub fn persistence_types<P: TestablePersistence>(#[case] _p: PhantomData<P>) {}
291
292    #[derive(Clone, Debug, Default)]
293    struct EventCollector {
294        events: Arc<RwLock<Vec<Event>>>,
295    }
296
297    impl EventCollector {
298        async fn leaf_chain(&self) -> Vec<LeafInfo<SeqTypes>> {
299            self.events
300                .read()
301                .await
302                .iter()
303                .flat_map(|event| {
304                    let EventType::Decide { leaf_chain, .. } = &event.event else {
305                        panic!("expected decide event, got {event:?}");
306                    };
307                    leaf_chain.iter().cloned().rev()
308                })
309                .collect::<Vec<_>>()
310        }
311    }
312
313    #[async_trait]
314    impl EventConsumer for EventCollector {
315        async fn handle_event(&self, event: &CoordinatorEvent<SeqTypes>) -> anyhow::Result<()> {
316            if let CoordinatorEvent::LegacyEvent(event) = event {
317                self.events.write().await.push(event.clone());
318            }
319            Ok(())
320        }
321    }
322
323    #[derive(Clone, Copy, Debug)]
324    struct FailConsumer;
325
326    #[async_trait]
327    impl EventConsumer for FailConsumer {
328        async fn handle_event(&self, _: &CoordinatorEvent<SeqTypes>) -> anyhow::Result<()> {
329            bail!("mock error injection");
330        }
331    }
332
333    #[rstest_reuse::apply(persistence_types)]
334    pub async fn test_voted_view<P: TestablePersistence>(_p: PhantomData<P>) {
335        let tmp = P::tmp_storage().await;
336        let storage = P::connect(&tmp).await;
337
338        // Initially, there is no saved view.
339        assert_eq!(storage.load_latest_acted_view().await.unwrap(), None);
340
341        // Store a view.
342        let view1 = ViewNumber::genesis();
343        storage
344            .record_action(view1, None, HotShotAction::Vote)
345            .await
346            .unwrap();
347        assert_eq!(
348            storage.load_latest_acted_view().await.unwrap().unwrap(),
349            view1
350        );
351
352        // Store a newer view, make sure storage gets updated.
353        let view2 = view1 + 1;
354        storage
355            .record_action(view2, None, HotShotAction::Vote)
356            .await
357            .unwrap();
358        assert_eq!(
359            storage.load_latest_acted_view().await.unwrap().unwrap(),
360            view2
361        );
362
363        // Store an old view, make sure storage is unchanged.
364        storage
365            .record_action(view1, None, HotShotAction::Vote)
366            .await
367            .unwrap();
368        assert_eq!(
369            storage.load_latest_acted_view().await.unwrap().unwrap(),
370            view2
371        );
372    }
373
374    #[rstest_reuse::apply(persistence_types)]
375    pub async fn test_high_qc2_monotonic<P: TestablePersistence>(_p: PhantomData<P>) {
376        let tmp = P::tmp_storage().await;
377        let storage = P::connect(&tmp).await;
378
379        // Initially there is no persisted lock.
380        assert_eq!(storage.load_high_qc2().await.unwrap(), None);
381
382        let mut qc = QuorumCertificate2::genesis(
383            &ValidatedState::default(),
384            &NodeState::mock(),
385            TEST_VERSIONS.test,
386        )
387        .await;
388
389        // Persist a lock at view 5.
390        qc.view_number = ViewNumber::new(5);
391        storage.append_high_qc2(qc.clone()).await.unwrap();
392        assert_eq!(
393            storage.load_high_qc2().await.unwrap().unwrap().view_number,
394            ViewNumber::new(5)
395        );
396
397        // A newer lock advances the stored view.
398        qc.view_number = ViewNumber::new(6);
399        storage.append_high_qc2(qc.clone()).await.unwrap();
400        assert_eq!(
401            storage.load_high_qc2().await.unwrap().unwrap().view_number,
402            ViewNumber::new(6)
403        );
404
405        // A stale (older) lock is a no-op: the compare-and-set never regresses
406        // the persisted view, even though the file/row write is unconditional.
407        qc.view_number = ViewNumber::new(4);
408        storage.append_high_qc2(qc.clone()).await.unwrap();
409        assert_eq!(
410            storage.load_high_qc2().await.unwrap().unwrap().view_number,
411            ViewNumber::new(6)
412        );
413    }
414
415    #[rstest_reuse::apply(persistence_types)]
416    pub async fn test_restart_view<P: TestablePersistence>(_p: PhantomData<P>) {
417        let tmp = P::tmp_storage().await;
418        let storage = P::connect(&tmp).await;
419
420        // Initially, there is no saved view.
421        assert_eq!(storage.load_restart_view().await.unwrap(), None);
422
423        // Store a view.
424        let view1 = ViewNumber::genesis();
425        storage
426            .record_action(view1, None, HotShotAction::Vote)
427            .await
428            .unwrap();
429        assert_eq!(
430            storage.load_restart_view().await.unwrap().unwrap(),
431            view1 + 1
432        );
433
434        // Store a newer view, make sure storage gets updated.
435        let view2 = view1 + 1;
436        storage
437            .record_action(view2, None, HotShotAction::Vote)
438            .await
439            .unwrap();
440        assert_eq!(
441            storage.load_restart_view().await.unwrap().unwrap(),
442            view2 + 1
443        );
444
445        // Store an old view, make sure storage is unchanged.
446        storage
447            .record_action(view1, None, HotShotAction::Vote)
448            .await
449            .unwrap();
450        assert_eq!(
451            storage.load_restart_view().await.unwrap().unwrap(),
452            view2 + 1
453        );
454
455        // store a higher proposed view, make sure storage is unchanged.
456        storage
457            .record_action(view2 + 1, None, HotShotAction::Propose)
458            .await
459            .unwrap();
460        assert_eq!(
461            storage.load_restart_view().await.unwrap().unwrap(),
462            view2 + 1
463        );
464
465        // store a higher timeout vote view, make sure storage is unchanged.
466        storage
467            .record_action(view2 + 1, None, HotShotAction::TimeoutVote)
468            .await
469            .unwrap();
470        assert_eq!(
471            storage.load_restart_view().await.unwrap().unwrap(),
472            view2 + 1
473        );
474    }
475
476    #[rstest_reuse::apply(persistence_types)]
477    pub async fn test_store_drb_input<P: TestablePersistence>(_p: PhantomData<P>) {
478        use hotshot_types::drb::DrbInput;
479
480        let tmp = P::tmp_storage().await;
481        let storage = P::connect(&tmp).await;
482        let difficulty_level = 10;
483
484        // Initially, there is no saved info.
485        if storage.load_drb_input(10).await.is_ok() {
486            panic!("unexpected nonempty drb_input");
487        }
488
489        let drb_input_1 = DrbInput {
490            epoch: 10,
491            iteration: 10,
492            value: [0u8; 32],
493            difficulty_level,
494        };
495
496        let drb_input_2 = DrbInput {
497            epoch: 10,
498            iteration: 20,
499            value: [0u8; 32],
500            difficulty_level,
501        };
502
503        let drb_input_3 = DrbInput {
504            epoch: 10,
505            iteration: 30,
506            value: [0u8; 32],
507            difficulty_level,
508        };
509
510        let _ = storage.store_drb_input(drb_input_1.clone()).await;
511
512        assert_eq!(storage.load_drb_input(10).await.unwrap(), drb_input_1);
513
514        let _ = storage.store_drb_input(drb_input_3.clone()).await;
515
516        // check that the drb input is overwritten
517        assert_eq!(storage.load_drb_input(10).await.unwrap(), drb_input_3);
518
519        let _ = storage.store_drb_input(drb_input_2.clone()).await;
520
521        // check that the drb input is not overwritten by the older value
522        assert_eq!(storage.load_drb_input(10).await.unwrap(), drb_input_3);
523    }
524
525    #[rstest_reuse::apply(persistence_types)]
526    pub async fn test_epoch_info<P: TestablePersistence>(_p: PhantomData<P>) {
527        let tmp = P::tmp_storage().await;
528        let storage = P::connect(&tmp).await;
529
530        // Initially, there is no saved info.
531        assert_eq!(storage.load_start_epoch_info().await.unwrap(), Vec::new());
532
533        // Store a drb result.
534        storage
535            .store_drb_result(EpochNumber::new(1), [1; 32])
536            .await
537            .unwrap();
538        assert_eq!(
539            storage.load_start_epoch_info().await.unwrap(),
540            vec![InitializerEpochInfo::<SeqTypes> {
541                epoch: EpochNumber::new(1),
542                drb_result: [1; 32],
543                block_header: None,
544            }]
545        );
546
547        // Store a second DRB result
548        storage
549            .store_drb_result(EpochNumber::new(2), [3; 32])
550            .await
551            .unwrap();
552        assert_eq!(
553            storage.load_start_epoch_info().await.unwrap(),
554            vec![
555                InitializerEpochInfo::<SeqTypes> {
556                    epoch: EpochNumber::new(1),
557                    drb_result: [1; 32],
558                    block_header: None,
559                },
560                InitializerEpochInfo::<SeqTypes> {
561                    epoch: EpochNumber::new(2),
562                    drb_result: [3; 32],
563                    block_header: None,
564                }
565            ]
566        );
567
568        // Make a header
569        let instance_state = NodeState::mock();
570        let validated_state = hotshot_types::traits::ValidatedState::genesis(&instance_state).0;
571        let leaf: Leaf2 = Leaf::genesis(&validated_state, &instance_state, MOCK_UPGRADE.base)
572            .await
573            .into();
574        let header = leaf.block_header().clone();
575
576        // Test storing the header
577        storage
578            .store_epoch_root(EpochNumber::new(1), header.clone())
579            .await
580            .unwrap();
581        assert_eq!(
582            storage.load_start_epoch_info().await.unwrap(),
583            vec![
584                InitializerEpochInfo::<SeqTypes> {
585                    epoch: EpochNumber::new(1),
586                    drb_result: [1; 32],
587                    block_header: Some(header.clone()),
588                },
589                InitializerEpochInfo::<SeqTypes> {
590                    epoch: EpochNumber::new(2),
591                    drb_result: [3; 32],
592                    block_header: None,
593                }
594            ]
595        );
596
597        // Store more than the limit
598        let total_epochs = RECENT_STAKE_TABLES_LIMIT + 10;
599        for i in 0..total_epochs {
600            let epoch = EpochNumber::new(i);
601            let drb = [i as u8; 32];
602            storage
603                .store_drb_result(epoch, drb)
604                .await
605                .unwrap_or_else(|_| panic!("Failed to store DRB result for epoch {i}"));
606        }
607
608        let results = storage.load_start_epoch_info().await.unwrap();
609
610        // Check that only the most recent RECENT_STAKE_TABLES_LIMIT epochs are returned
611        assert_eq!(
612            results.len(),
613            RECENT_STAKE_TABLES_LIMIT as usize,
614            "Should return only the most recent {RECENT_STAKE_TABLES_LIMIT} epochs",
615        );
616
617        for (i, info) in results.iter().enumerate() {
618            let expected_epoch =
619                EpochNumber::new(total_epochs - RECENT_STAKE_TABLES_LIMIT + i as u64);
620            let expected_drb = [(total_epochs - RECENT_STAKE_TABLES_LIMIT + i as u64) as u8; 32];
621            assert_eq!(info.epoch, expected_epoch, "invalid epoch at index {i}",);
622            assert_eq!(info.drb_result, expected_drb, "invalid DRB at index {i}",);
623            assert!(info.block_header.is_none(), "Expected no block header");
624        }
625    }
626
627    fn leaf_info(leaf: Leaf2) -> LeafInfo<SeqTypes> {
628        LeafInfo {
629            leaf,
630            vid_share: None,
631            state: Default::default(),
632            delta: None,
633            state_cert: None,
634        }
635    }
636
637    #[rstest_reuse::apply(persistence_types)]
638    pub async fn test_append_and_decide<P: TestablePersistence>(_p: PhantomData<P>) {
639        let tmp = P::tmp_storage().await;
640        let storage = P::connect(&tmp).await;
641
642        // Test append VID
643        assert_eq!(
644            storage.load_vid_share(ViewNumber::new(0)).await.unwrap(),
645            None
646        );
647
648        let leaf: Leaf2 = Leaf2::genesis(
649            &ValidatedState::default(),
650            &NodeState::mock(),
651            TEST_VERSIONS.test.base,
652        )
653        .await;
654        let leaf_payload = leaf.block_payload().unwrap();
655        let leaf_payload_bytes_arc = leaf_payload.encode();
656
657        let avidm_param = init_avidm_param(2).unwrap();
658        let weights = vec![1u32; 2];
659
660        let ns_table = parse_ns_table(
661            leaf_payload.byte_len().as_usize(),
662            &leaf_payload.ns_table().encode(),
663        );
664        let (payload_commitment, shares) =
665            AvidMScheme::ns_disperse(&avidm_param, &weights, &leaf_payload_bytes_arc, ns_table)
666                .unwrap();
667
668        let (pubkey, privkey) = BLSPubKey::generated_from_seed_indexed([0; 32], 1);
669        let signature = PubKey::sign(&privkey, &[]).unwrap();
670        let mut vid = AvidMDisperseShare::<SeqTypes> {
671            view_number: ViewNumber::new(0),
672            payload_commitment,
673            share: shares[0].clone(),
674            recipient_key: pubkey,
675            epoch: Some(EpochNumber::new(0)),
676            target_epoch: Some(EpochNumber::new(0)),
677            common: avidm_param,
678        };
679        let mut quorum_proposal = Proposal {
680            data: QuorumProposalWrapper::<SeqTypes> {
681                proposal: QuorumProposal2::<SeqTypes> {
682                    epoch: None,
683                    block_header: leaf.block_header().clone(),
684                    view_number: ViewNumber::genesis(),
685                    justify_qc: QuorumCertificate2::genesis(
686                        &ValidatedState::default(),
687                        &NodeState::mock(),
688                        TEST_VERSIONS.test,
689                    )
690                    .await,
691                    upgrade_certificate: None,
692                    view_change_evidence: None,
693                    next_drb_result: None,
694                    next_epoch_justify_qc: None,
695                    state_cert: None,
696                },
697            },
698            signature,
699            _pd: Default::default(),
700        };
701
702        let vid_share0 = convert_proposal(vid.clone().to_proposal(&privkey).unwrap().clone());
703
704        storage.append_vid(&vid_share0).await.unwrap();
705
706        assert_eq!(
707            storage.load_vid_share(ViewNumber::new(0)).await.unwrap(),
708            Some(vid_share0.clone())
709        );
710
711        vid.view_number = ViewNumber::new(1);
712
713        let vid_share1 = convert_proposal(vid.clone().to_proposal(&privkey).unwrap().clone());
714        storage.append_vid(&vid_share1).await.unwrap();
715
716        assert_eq!(
717            storage.load_vid_share(vid.view_number()).await.unwrap(),
718            Some(vid_share1.clone())
719        );
720
721        vid.view_number = ViewNumber::new(2);
722
723        let vid_share2 = convert_proposal(vid.clone().to_proposal(&privkey).unwrap().clone());
724        storage.append_vid(&vid_share2).await.unwrap();
725
726        assert_eq!(
727            storage.load_vid_share(vid.view_number()).await.unwrap(),
728            Some(vid_share2.clone())
729        );
730
731        vid.view_number = ViewNumber::new(3);
732
733        let vid_share3 = convert_proposal(vid.clone().to_proposal(&privkey).unwrap().clone());
734        storage.append_vid(&vid_share3).await.unwrap();
735
736        assert_eq!(
737            storage.load_vid_share(vid.view_number()).await.unwrap(),
738            Some(vid_share3.clone())
739        );
740
741        let block_payload_signature = BLSPubKey::sign(&privkey, &leaf_payload_bytes_arc)
742            .expect("Failed to sign block payload");
743
744        let da_proposal_inner = DaProposal2::<SeqTypes> {
745            encoded_transactions: leaf_payload_bytes_arc.clone(),
746            metadata: leaf_payload.ns_table().clone(),
747            view_number: ViewNumber::new(0),
748            epoch: None,
749            epoch_transition_indicator: EpochTransitionIndicator::NotInTransition,
750        };
751
752        let da_proposal = Proposal {
753            data: da_proposal_inner,
754            signature: block_payload_signature,
755            _pd: Default::default(),
756        };
757
758        let vid_commitment = vid_commitment(
759            &leaf_payload_bytes_arc,
760            &leaf.block_header().metadata().encode(),
761            2,
762            TEST_VERSIONS.test.base,
763        );
764
765        storage
766            .append_da2(&da_proposal, vid_commitment)
767            .await
768            .unwrap();
769
770        assert_eq!(
771            storage.load_da_proposal(ViewNumber::new(0)).await.unwrap(),
772            Some(da_proposal.clone())
773        );
774
775        let mut da_proposal1 = da_proposal.clone();
776        da_proposal1.data.view_number = ViewNumber::new(1);
777        storage
778            .append_da2(&da_proposal1.clone(), vid_commitment)
779            .await
780            .unwrap();
781
782        assert_eq!(
783            storage
784                .load_da_proposal(da_proposal1.data.view_number)
785                .await
786                .unwrap(),
787            Some(da_proposal1.clone())
788        );
789
790        let mut da_proposal2 = da_proposal1.clone();
791        da_proposal2.data.view_number = ViewNumber::new(2);
792        storage
793            .append_da2(&da_proposal2.clone(), vid_commitment)
794            .await
795            .unwrap();
796
797        assert_eq!(
798            storage
799                .load_da_proposal(da_proposal2.data.view_number)
800                .await
801                .unwrap(),
802            Some(da_proposal2.clone())
803        );
804
805        let mut da_proposal3 = da_proposal2.clone();
806        da_proposal3.data.view_number = ViewNumber::new(3);
807        storage
808            .append_da2(&da_proposal3.clone(), vid_commitment)
809            .await
810            .unwrap();
811
812        assert_eq!(
813            storage
814                .load_da_proposal(da_proposal3.data.view_number)
815                .await
816                .unwrap(),
817            Some(da_proposal3.clone())
818        );
819
820        let quorum_proposal1 = quorum_proposal.clone();
821
822        storage
823            .append_quorum_proposal2(&quorum_proposal1)
824            .await
825            .unwrap();
826
827        assert_eq!(
828            storage.load_quorum_proposals().await.unwrap(),
829            BTreeMap::from_iter([(ViewNumber::genesis(), quorum_proposal1.clone())])
830        );
831
832        quorum_proposal.data.proposal.view_number = ViewNumber::new(1);
833        let quorum_proposal2 = quorum_proposal.clone();
834        storage
835            .append_quorum_proposal2(&quorum_proposal2)
836            .await
837            .unwrap();
838
839        assert_eq!(
840            storage.load_quorum_proposals().await.unwrap(),
841            BTreeMap::from_iter([
842                (ViewNumber::genesis(), quorum_proposal1.clone()),
843                (ViewNumber::new(1), quorum_proposal2.clone())
844            ])
845        );
846
847        quorum_proposal.data.proposal.view_number = ViewNumber::new(2);
848        quorum_proposal.data.proposal.justify_qc.view_number = ViewNumber::new(1);
849        let quorum_proposal3 = quorum_proposal.clone();
850        storage
851            .append_quorum_proposal2(&quorum_proposal3)
852            .await
853            .unwrap();
854
855        assert_eq!(
856            storage.load_quorum_proposals().await.unwrap(),
857            BTreeMap::from_iter([
858                (ViewNumber::genesis(), quorum_proposal1.clone()),
859                (ViewNumber::new(1), quorum_proposal2.clone()),
860                (ViewNumber::new(2), quorum_proposal3.clone())
861            ])
862        );
863
864        quorum_proposal.data.proposal.view_number = ViewNumber::new(3);
865        quorum_proposal.data.proposal.justify_qc.view_number = ViewNumber::new(2);
866
867        // This one should stick around after GC runs.
868        let quorum_proposal4 = quorum_proposal.clone();
869        storage
870            .append_quorum_proposal2(&quorum_proposal4)
871            .await
872            .unwrap();
873
874        assert_eq!(
875            storage.load_quorum_proposals().await.unwrap(),
876            BTreeMap::from_iter([
877                (ViewNumber::genesis(), quorum_proposal1.clone()),
878                (ViewNumber::new(1), quorum_proposal2.clone()),
879                (ViewNumber::new(2), quorum_proposal3.clone()),
880                (ViewNumber::new(3), quorum_proposal4.clone())
881            ])
882        );
883
884        // Test decide and garbage collection. Pass in a leaf chain with no VID shares or payloads,
885        // so we have to fetch the missing data from storage.
886        let leaves = [
887            Leaf2::from_quorum_proposal(&quorum_proposal1.data),
888            Leaf2::from_quorum_proposal(&quorum_proposal2.data),
889            Leaf2::from_quorum_proposal(&quorum_proposal3.data),
890            Leaf2::from_quorum_proposal(&quorum_proposal4.data),
891        ];
892        let mut final_qc = leaves[3].justify_qc();
893        final_qc.view_number += 1;
894        final_qc.data.leaf_commit = Committable::commit(&leaf);
895        let qcs = [
896            CertificatePair::non_epoch_change(leaves[1].justify_qc()),
897            CertificatePair::non_epoch_change(leaves[2].justify_qc()),
898            CertificatePair::non_epoch_change(leaves[3].justify_qc()),
899            CertificatePair::non_epoch_change(final_qc),
900        ];
901
902        assert_eq!(
903            storage.load_anchor_view().await.unwrap(),
904            ViewNumber::genesis()
905        );
906
907        let consumer = EventCollector::default();
908        let leaf_chain = leaves
909            .iter()
910            .take(3)
911            .map(|leaf| leaf_info(leaf.clone()))
912            .zip(&qcs)
913            .collect::<Vec<_>>();
914        tracing::info!(?leaf_chain, "decide view 2");
915        storage
916            .append_decided_leaves(
917                ViewNumber::new(2),
918                leaf_chain.iter().map(|(leaf, qc)| (leaf, (*qc).clone())),
919                None,
920                &consumer,
921            )
922            .await
923            .unwrap();
924        assert_eq!(
925            storage.load_anchor_view().await.unwrap(),
926            ViewNumber::new(2)
927        );
928
929        for i in 0..=2 {
930            assert_eq!(
931                storage.load_da_proposal(ViewNumber::new(i)).await.unwrap(),
932                None
933            );
934
935            assert_eq!(
936                storage.load_vid_share(ViewNumber::new(i)).await.unwrap(),
937                None
938            );
939        }
940
941        assert_eq!(
942            storage.load_da_proposal(ViewNumber::new(3)).await.unwrap(),
943            Some(da_proposal3)
944        );
945
946        assert_eq!(
947            storage.load_vid_share(ViewNumber::new(3)).await.unwrap(),
948            Some(convert_proposal(vid_share3.clone()))
949        );
950
951        let proposals = storage.load_quorum_proposals().await.unwrap();
952        assert_eq!(
953            proposals,
954            BTreeMap::from_iter([(ViewNumber::new(3), quorum_proposal4)])
955        );
956
957        // A decide event should have been processed.
958        for (leaf, info) in leaves.iter().zip(consumer.leaf_chain().await.iter()) {
959            assert_eq!(info.leaf, *leaf);
960            let decided_vid_share = info.vid_share.as_ref().unwrap();
961            assert_eq!(decided_vid_share.view_number(), leaf.view_number());
962        }
963
964        // The decided leaf should not have been garbage collected.
965        assert_eq!(
966            storage.load_anchor_leaf().await.unwrap(),
967            Some((leaves[2].clone(), qcs[2].clone()))
968        );
969        assert_eq!(
970            storage.load_anchor_view().await.unwrap(),
971            leaves[2].view_number()
972        );
973
974        // Process a second decide event.
975        let consumer = EventCollector::default();
976        tracing::info!(leaf = ?leaves[3], qc = ?qcs[3], "decide view 3");
977        storage
978            .append_decided_leaves(
979                ViewNumber::new(3),
980                vec![(&leaf_info(leaves[3].clone()), qcs[3].clone())],
981                None,
982                &consumer,
983            )
984            .await
985            .unwrap();
986        assert_eq!(
987            storage.load_anchor_view().await.unwrap(),
988            ViewNumber::new(3)
989        );
990
991        // A decide event should have been processed.
992        let events = consumer.events.read().await;
993        assert_eq!(events.len(), 1);
994        assert_eq!(events[0].view_number, ViewNumber::new(3));
995        let EventType::Decide {
996            committing_qc,
997            leaf_chain,
998            ..
999        } = &events[0].event
1000        else {
1001            panic!("expected decide event, got {:?}", events[0]);
1002        };
1003        assert_eq!(**committing_qc, qcs[3]);
1004        assert_eq!(leaf_chain.len(), 1);
1005        let info = &leaf_chain[0];
1006        assert_eq!(info.leaf, leaves[3]);
1007
1008        // The remaining data should have been GCed.
1009        assert_eq!(
1010            storage.load_da_proposal(ViewNumber::new(3)).await.unwrap(),
1011            None
1012        );
1013
1014        assert_eq!(
1015            storage.load_vid_share(ViewNumber::new(3)).await.unwrap(),
1016            None
1017        );
1018        assert_eq!(
1019            storage.load_quorum_proposals().await.unwrap(),
1020            BTreeMap::new()
1021        );
1022    }
1023
1024    #[rstest_reuse::apply(persistence_types)]
1025    pub async fn test_upgrade_certificate<P: TestablePersistence>(_p: PhantomData<P>) {
1026        let tmp = P::tmp_storage().await;
1027        let storage = P::connect(&tmp).await;
1028
1029        // Test get upgrade certificate
1030        assert_eq!(storage.load_upgrade_certificate().await.unwrap(), None);
1031
1032        let upgrade_data = UpgradeProposalData {
1033            old_version: Version { major: 0, minor: 1 },
1034            new_version: Version { major: 1, minor: 0 },
1035            decide_by: ViewNumber::genesis(),
1036            new_version_hash: Default::default(),
1037            old_version_last_view: ViewNumber::genesis(),
1038            new_version_first_view: ViewNumber::genesis(),
1039        };
1040
1041        let decide_upgrade_certificate = UpgradeCertificate::<SeqTypes>::new(
1042            upgrade_data.clone(),
1043            upgrade_data.commit(),
1044            ViewNumber::genesis(),
1045            Default::default(),
1046            Default::default(),
1047        );
1048        let res = storage
1049            .store_upgrade_certificate(Some(decide_upgrade_certificate.clone()))
1050            .await;
1051        assert!(res.is_ok());
1052
1053        let res = storage.load_upgrade_certificate().await.unwrap();
1054        let view_number = res.unwrap().view_number;
1055        assert_eq!(view_number, ViewNumber::genesis());
1056
1057        let new_view_number_for_certificate = ViewNumber::new(50);
1058        let mut new_upgrade_certificate = decide_upgrade_certificate.clone();
1059        new_upgrade_certificate.view_number = new_view_number_for_certificate;
1060
1061        let res = storage
1062            .store_upgrade_certificate(Some(new_upgrade_certificate.clone()))
1063            .await;
1064        assert!(res.is_ok());
1065
1066        let res = storage.load_upgrade_certificate().await.unwrap();
1067        let view_number = res.unwrap().view_number;
1068        assert_eq!(view_number, new_view_number_for_certificate);
1069    }
1070
1071    #[rstest_reuse::apply(persistence_types)]
1072    pub async fn test_next_epoch_quorum_certificate<P: TestablePersistence>(_p: PhantomData<P>) {
1073        let tmp = P::tmp_storage().await;
1074        let storage = P::connect(&tmp).await;
1075
1076        //  test that next epoch qc2 does not exist
1077        assert_eq!(
1078            storage.load_next_epoch_quorum_certificate().await.unwrap(),
1079            None
1080        );
1081
1082        let upgrade_lock = UpgradeLock::<SeqTypes>::new(TEST_VERSIONS.test);
1083
1084        let genesis_view = ViewNumber::genesis();
1085
1086        let leaf = Leaf2::genesis(
1087            &ValidatedState::default(),
1088            &NodeState::default(),
1089            TEST_VERSIONS.test.base,
1090        )
1091        .await;
1092        let data: NextEpochQuorumData2<SeqTypes> = QuorumData2 {
1093            leaf_commit: leaf.commit(),
1094            epoch: Some(EpochNumber::new(1)),
1095            block_number: Some(leaf.height()),
1096        }
1097        .into();
1098
1099        let versioned_data =
1100            VersionedVoteData::new_infallible(data.clone(), genesis_view, &upgrade_lock);
1101
1102        let bytes: [u8; 32] = versioned_data.commit().into();
1103
1104        let next_epoch_qc = NextEpochQuorumCertificate2::new(
1105            data,
1106            Commitment::from_raw(bytes),
1107            genesis_view,
1108            None,
1109            PhantomData,
1110        );
1111
1112        let res = storage
1113            .append_next_epoch_high_qc2(next_epoch_qc.clone())
1114            .await;
1115        assert!(res.is_ok());
1116
1117        let res = storage.load_next_epoch_quorum_certificate().await.unwrap();
1118        let view_number = res.unwrap().view_number;
1119        assert_eq!(view_number, ViewNumber::genesis());
1120
1121        let new_view_number_for_qc = ViewNumber::new(50);
1122        let mut new_qc = next_epoch_qc.clone();
1123        new_qc.view_number = new_view_number_for_qc;
1124
1125        let res = storage.append_next_epoch_high_qc2(new_qc.clone()).await;
1126        assert!(res.is_ok());
1127
1128        let res = storage.load_next_epoch_quorum_certificate().await.unwrap();
1129        let view_number = res.unwrap().view_number;
1130        assert_eq!(view_number, new_view_number_for_qc);
1131    }
1132
1133    /// `append_next_epoch_high_qc2` must apply an atomic monotonic compare-and-set (like
1134    /// `append_high_qc2`): a stale (older-view) write is a no-op and never regresses the stored view.
1135    #[rstest_reuse::apply(persistence_types)]
1136    pub async fn test_append_next_epoch_high_qc2_monotonic<P: TestablePersistence>(
1137        _p: PhantomData<P>,
1138    ) {
1139        let tmp = P::tmp_storage().await;
1140        let storage = P::connect(&tmp).await;
1141
1142        assert_eq!(
1143            storage.load_next_epoch_quorum_certificate().await.unwrap(),
1144            None
1145        );
1146
1147        let upgrade_lock = UpgradeLock::<SeqTypes>::new(TEST_VERSIONS.test);
1148        let leaf = Leaf2::genesis(
1149            &ValidatedState::default(),
1150            &NodeState::default(),
1151            TEST_VERSIONS.test.base,
1152        )
1153        .await;
1154        let data: NextEpochQuorumData2<SeqTypes> = QuorumData2 {
1155            leaf_commit: leaf.commit(),
1156            epoch: Some(EpochNumber::new(1)),
1157            block_number: Some(leaf.height()),
1158        }
1159        .into();
1160        let versioned_data =
1161            VersionedVoteData::new_infallible(data.clone(), ViewNumber::genesis(), &upgrade_lock);
1162        let bytes: [u8; 32] = versioned_data.commit().into();
1163        let mut qc = NextEpochQuorumCertificate2::new(
1164            data,
1165            Commitment::from_raw(bytes),
1166            ViewNumber::genesis(),
1167            None,
1168            PhantomData,
1169        );
1170
1171        // Persist at view 5, then advance to 6.
1172        qc.view_number = ViewNumber::new(5);
1173        storage
1174            .append_next_epoch_high_qc2(qc.clone())
1175            .await
1176            .unwrap();
1177        assert_eq!(
1178            storage
1179                .load_next_epoch_quorum_certificate()
1180                .await
1181                .unwrap()
1182                .unwrap()
1183                .view_number,
1184            ViewNumber::new(5)
1185        );
1186
1187        qc.view_number = ViewNumber::new(6);
1188        storage
1189            .append_next_epoch_high_qc2(qc.clone())
1190            .await
1191            .unwrap();
1192        assert_eq!(
1193            storage
1194                .load_next_epoch_quorum_certificate()
1195                .await
1196                .unwrap()
1197                .unwrap()
1198                .view_number,
1199            ViewNumber::new(6)
1200        );
1201
1202        // A stale (older) write is a no-op: the compare-and-set never regresses the stored view.
1203        qc.view_number = ViewNumber::new(4);
1204        storage
1205            .append_next_epoch_high_qc2(qc.clone())
1206            .await
1207            .unwrap();
1208        assert_eq!(
1209            storage
1210                .load_next_epoch_quorum_certificate()
1211                .await
1212                .unwrap()
1213                .unwrap()
1214                .view_number,
1215            ViewNumber::new(6)
1216        );
1217    }
1218
1219    /// The `Storage<SeqTypes>` impl on `Arc<P>` (the object the consensus tasks actually hold) used
1220    /// to no-op `update_high_qc2` / `update_next_epoch_high_qc2` — the persist-before-vote hook.
1221    /// They must now durably route to the `high_qc2` / `next_epoch_quorum_certificate` tables, and
1222    /// `update_high_qc2` must stay monotonic.
1223    #[rstest_reuse::apply(persistence_types)]
1224    pub async fn test_storage_update_high_qc2_persists<P: TestablePersistence>(_p: PhantomData<P>) {
1225        use hotshot_types::traits::storage::Storage;
1226
1227        let tmp = P::tmp_storage().await;
1228        let storage = P::connect(&tmp).await;
1229        // Production wraps the persistence as `Arc<P>`; that is the `Storage<SeqTypes>` impl we fixed.
1230        let arc = Arc::new(storage.clone());
1231
1232        assert_eq!(storage.load_high_qc2().await.unwrap(), None);
1233
1234        let mut high_qc = QuorumCertificate2::genesis(
1235            &ValidatedState::default(),
1236            &NodeState::mock(),
1237            TEST_VERSIONS.test,
1238        )
1239        .await;
1240        high_qc.view_number = ViewNumber::new(7);
1241        Storage::update_high_qc2(&arc, high_qc.clone())
1242            .await
1243            .unwrap();
1244        assert_eq!(
1245            storage.load_high_qc2().await.unwrap().unwrap().view_number,
1246            ViewNumber::new(7)
1247        );
1248
1249        // A stale write through the Storage trait must not regress the persisted view.
1250        let mut stale = high_qc.clone();
1251        stale.view_number = ViewNumber::new(3);
1252        Storage::update_high_qc2(&arc, stale).await.unwrap();
1253        assert_eq!(
1254            storage.load_high_qc2().await.unwrap().unwrap().view_number,
1255            ViewNumber::new(7)
1256        );
1257
1258        // The next-epoch high QC round-trips through update_next_epoch_high_qc2 as well.
1259        let upgrade_lock = UpgradeLock::<SeqTypes>::new(TEST_VERSIONS.test);
1260        let leaf = Leaf2::genesis(
1261            &ValidatedState::default(),
1262            &NodeState::default(),
1263            TEST_VERSIONS.test.base,
1264        )
1265        .await;
1266        let data: NextEpochQuorumData2<SeqTypes> = QuorumData2 {
1267            leaf_commit: leaf.commit(),
1268            epoch: Some(EpochNumber::new(1)),
1269            block_number: Some(leaf.height()),
1270        }
1271        .into();
1272        let versioned_data =
1273            VersionedVoteData::new_infallible(data.clone(), ViewNumber::new(7), &upgrade_lock);
1274        let bytes: [u8; 32] = versioned_data.commit().into();
1275        let next_epoch_qc = NextEpochQuorumCertificate2::new(
1276            data,
1277            Commitment::from_raw(bytes),
1278            ViewNumber::new(7),
1279            None,
1280            PhantomData,
1281        );
1282        Storage::update_next_epoch_high_qc2(&arc, next_epoch_qc)
1283            .await
1284            .unwrap();
1285        assert_eq!(
1286            storage
1287                .load_next_epoch_quorum_certificate()
1288                .await
1289                .unwrap()
1290                .unwrap()
1291                .view_number,
1292            ViewNumber::new(7)
1293        );
1294    }
1295
1296    /// `load_consensus_state` must fold the running high QC persisted via `update_high_qc2` back into
1297    /// the recovered `high_qc` (which `Consensus::new` then uses to restore `locked_view` past the
1298    /// decided anchor). With no decided anchor, recovery starts from the genesis QC at view 0 and
1299    /// must pick up a persisted high QC at a higher view.
1300    #[rstest_reuse::apply(persistence_types)]
1301    pub async fn test_load_consensus_state_recovers_high_qc<P: TestablePersistence>(
1302        _p: PhantomData<P>,
1303    ) {
1304        let tmp = P::tmp_storage().await;
1305        let storage = P::connect(&tmp).await;
1306
1307        let mut high_qc = QuorumCertificate2::genesis(
1308            &ValidatedState::default(),
1309            &NodeState::mock(),
1310            TEST_VERSIONS.test,
1311        )
1312        .await;
1313        high_qc.view_number = ViewNumber::new(9);
1314        storage.append_high_qc2(high_qc).await.unwrap();
1315
1316        let (initializer, _anchor_view) = storage
1317            .load_consensus_state(
1318                NodeState::mock(),
1319                Upgrade::trivial(versions::EPOCH_REWARD_VERSION),
1320            )
1321            .await
1322            .unwrap();
1323
1324        // Without the recovery fold this would be the genesis view (0); with it, the persisted lock.
1325        assert_eq!(initializer.high_qc().view_number, ViewNumber::new(9));
1326    }
1327
1328    /// Recovery must not hand `Consensus::new` a mismatched (high QC, next-epoch QC) pair: when a
1329    /// newer running high QC is adopted, a persisted next-epoch QC that does not correspond to it
1330    /// (`verify_next_epoch_qc`) is dropped to `None` instead of being recovered alongside it.
1331    #[rstest_reuse::apply(persistence_types)]
1332    pub async fn test_load_consensus_state_drops_noncorresponding_next_epoch_qc<
1333        P: TestablePersistence,
1334    >(
1335        _p: PhantomData<P>,
1336    ) {
1337        let tmp = P::tmp_storage().await;
1338        let storage = P::connect(&tmp).await;
1339
1340        // Build a next-epoch QC we can re-stamp at different views.
1341        let upgrade_lock = UpgradeLock::<SeqTypes>::new(TEST_VERSIONS.test);
1342        let leaf = Leaf2::genesis(
1343            &ValidatedState::default(),
1344            &NodeState::default(),
1345            TEST_VERSIONS.test.base,
1346        )
1347        .await;
1348        let data: NextEpochQuorumData2<SeqTypes> = QuorumData2 {
1349            leaf_commit: leaf.commit(),
1350            epoch: Some(EpochNumber::new(1)),
1351            block_number: Some(leaf.height()),
1352        }
1353        .into();
1354        let versioned_data =
1355            VersionedVoteData::new_infallible(data.clone(), ViewNumber::new(5), &upgrade_lock);
1356        let bytes: [u8; 32] = versioned_data.commit().into();
1357        let next_epoch_qc = NextEpochQuorumCertificate2::new(
1358            data,
1359            Commitment::from_raw(bytes),
1360            ViewNumber::new(5),
1361            None,
1362            PhantomData,
1363        );
1364
1365        // eQC checkpoint: high QC + corresponding next-epoch QC, both at view 5.
1366        let mut eqc_high_qc = QuorumCertificate2::genesis(
1367            &ValidatedState::default(),
1368            &NodeState::mock(),
1369            TEST_VERSIONS.test,
1370        )
1371        .await;
1372        eqc_high_qc.view_number = ViewNumber::new(5);
1373        storage
1374            .store_eqc(eqc_high_qc.clone(), next_epoch_qc.clone())
1375            .await
1376            .unwrap();
1377
1378        // A newer running high QC (view 9) — recovery adopts it...
1379        let mut running_high_qc = eqc_high_qc.clone();
1380        running_high_qc.view_number = ViewNumber::new(9);
1381        storage.append_high_qc2(running_high_qc).await.unwrap();
1382
1383        // ...and a standalone next-epoch QC (view 3) in its table. Neither it nor the eQC's view-5
1384        // next-epoch QC corresponds to the recovered high QC at view 9.
1385        let mut other_next_epoch_qc = next_epoch_qc.clone();
1386        other_next_epoch_qc.view_number = ViewNumber::new(3);
1387        storage
1388            .append_next_epoch_high_qc2(other_next_epoch_qc)
1389            .await
1390            .unwrap();
1391
1392        let (initializer, _anchor_view) = storage
1393            .load_consensus_state(
1394                NodeState::mock(),
1395                Upgrade::trivial(versions::EPOCH_REWARD_VERSION),
1396            )
1397            .await
1398            .unwrap();
1399
1400        assert_eq!(initializer.high_qc().view_number, ViewNumber::new(9));
1401        // The recovered high QC corresponds to neither persisted next-epoch QC, so the pair is left
1402        // without one rather than recovering a mismatched (high QC, next-epoch QC) pair.
1403        assert!(
1404            initializer.next_epoch_high_qc().is_none(),
1405            "a next-epoch QC that does not correspond to the recovered high QC must be dropped"
1406        );
1407    }
1408
1409    #[rstest_reuse::apply(persistence_types)]
1410    pub async fn test_decide_with_failing_event_consumer<P: TestablePersistence>(
1411        _p: PhantomData<P>,
1412    ) {
1413        let tmp = P::tmp_storage().await;
1414        let storage = P::connect(&tmp).await;
1415
1416        // Create a short blockchain.
1417        let mut chain = vec![];
1418
1419        let leaf: Leaf2 = Leaf::genesis(
1420            &ValidatedState::default(),
1421            &NodeState::mock(),
1422            MOCK_UPGRADE.base,
1423        )
1424        .await
1425        .into();
1426        let leaf_payload = leaf.block_payload().unwrap();
1427        let leaf_payload_bytes_arc = leaf_payload.encode();
1428        let avidm_param = init_avidm_param(2).unwrap();
1429        let weights = vec![1u32; 2];
1430        let ns_table = parse_ns_table(
1431            leaf_payload.byte_len().as_usize(),
1432            &leaf_payload.ns_table().encode(),
1433        );
1434        let (payload_commitment, shares) =
1435            AvidMScheme::ns_disperse(&avidm_param, &weights, &leaf_payload_bytes_arc, ns_table)
1436                .unwrap();
1437
1438        let (pubkey, privkey) = BLSPubKey::generated_from_seed_indexed([0; 32], 1);
1439        let mut vid = AvidMDisperseShare::<SeqTypes> {
1440            view_number: ViewNumber::new(0),
1441            payload_commitment,
1442            share: shares[0].clone(),
1443            recipient_key: pubkey,
1444            epoch: Some(EpochNumber::new(0)),
1445            target_epoch: Some(EpochNumber::new(0)),
1446            common: avidm_param,
1447        }
1448        .to_proposal(&privkey)
1449        .unwrap()
1450        .clone();
1451        let mut quorum_proposal = QuorumProposalWrapper::<SeqTypes> {
1452            proposal: QuorumProposal2::<SeqTypes> {
1453                block_header: leaf.block_header().clone(),
1454                view_number: ViewNumber::genesis(),
1455                justify_qc: QuorumCertificate::genesis(
1456                    &ValidatedState::default(),
1457                    &NodeState::mock(),
1458                    TEST_VERSIONS.test,
1459                )
1460                .await
1461                .to_qc2(),
1462                upgrade_certificate: None,
1463                view_change_evidence: None,
1464                next_drb_result: None,
1465                next_epoch_justify_qc: None,
1466                epoch: None,
1467                state_cert: None,
1468            },
1469        };
1470        let mut qc = QuorumCertificate2::genesis(
1471            &ValidatedState::default(),
1472            &NodeState::mock(),
1473            TEST_VERSIONS.test,
1474        )
1475        .await;
1476
1477        let block_payload_signature = BLSPubKey::sign(&privkey, &leaf_payload_bytes_arc)
1478            .expect("Failed to sign block payload");
1479        let mut da_proposal = Proposal {
1480            data: DaProposal2::<SeqTypes> {
1481                encoded_transactions: leaf_payload_bytes_arc.clone(),
1482                metadata: leaf_payload.ns_table().clone(),
1483                view_number: ViewNumber::new(0),
1484                epoch: Some(EpochNumber::new(0)),
1485                epoch_transition_indicator: EpochTransitionIndicator::NotInTransition,
1486            },
1487            signature: block_payload_signature,
1488            _pd: Default::default(),
1489        };
1490
1491        let vid_commitment = vid_commitment(
1492            &leaf_payload_bytes_arc,
1493            &leaf.block_header().metadata().encode(),
1494            2,
1495            TEST_VERSIONS.test.base,
1496        );
1497
1498        for i in 0..4 {
1499            quorum_proposal.proposal.view_number = ViewNumber::new(i);
1500            let leaf = Leaf2::from_quorum_proposal(&quorum_proposal);
1501            qc.view_number = leaf.view_number();
1502            qc.data.leaf_commit = Committable::commit(&leaf);
1503            vid.data.view_number = leaf.view_number();
1504            da_proposal.data.view_number = leaf.view_number();
1505            chain.push((leaf.clone(), qc.clone(), vid.clone(), da_proposal.clone()));
1506        }
1507
1508        // Add proposals.
1509        for (_, _, vid, da) in &chain {
1510            tracing::info!(?da, ?vid, "insert proposal");
1511            storage.append_da2(da, vid_commitment).await.unwrap();
1512            storage
1513                .append_vid(&convert_proposal(vid.clone()))
1514                .await
1515                .unwrap();
1516        }
1517
1518        // Decide 2 leaves, but fail in event processing.
1519        let leaf_chain = chain
1520            .iter()
1521            .take(2)
1522            .map(|(leaf, qc, ..)| (leaf_info(leaf.clone()), qc.clone()))
1523            .collect::<Vec<_>>();
1524        tracing::info!("decide with event handling failure");
1525        storage
1526            .append_decided_leaves(
1527                ViewNumber::new(1),
1528                leaf_chain
1529                    .iter()
1530                    .map(|(leaf, qc)| (leaf, CertificatePair::non_epoch_change(qc.clone()))),
1531                None,
1532                &FailConsumer,
1533            )
1534            .await
1535            .unwrap();
1536        // No garbage collection should have run.
1537        for i in 0..4 {
1538            tracing::info!(i, "check proposal availability");
1539            assert!(
1540                storage
1541                    .load_vid_share(ViewNumber::new(i))
1542                    .await
1543                    .unwrap()
1544                    .is_some()
1545            );
1546            assert!(
1547                storage
1548                    .load_da_proposal(ViewNumber::new(i))
1549                    .await
1550                    .unwrap()
1551                    .is_some()
1552            );
1553        }
1554        tracing::info!("check anchor leaf updated");
1555        assert_eq!(
1556            storage
1557                .load_anchor_leaf()
1558                .await
1559                .unwrap()
1560                .unwrap()
1561                .0
1562                .view_number(),
1563            ViewNumber::new(1)
1564        );
1565        assert_eq!(
1566            storage.load_anchor_view().await.unwrap(),
1567            ViewNumber::new(1)
1568        );
1569
1570        // Now decide remaining leaves successfully. We should now garbage collect and process a
1571        // decide event for all the leaves.
1572        let consumer = EventCollector::default();
1573        let leaf_chain = chain
1574            .iter()
1575            .skip(2)
1576            .map(|(leaf, qc, ..)| (leaf_info(leaf.clone()), qc.clone()))
1577            .collect::<Vec<_>>();
1578        tracing::info!("decide successfully");
1579        storage
1580            .append_decided_leaves(
1581                ViewNumber::new(3),
1582                leaf_chain
1583                    .iter()
1584                    .map(|(leaf, qc)| (leaf, CertificatePair::non_epoch_change(qc.clone()))),
1585                None,
1586                &consumer,
1587            )
1588            .await
1589            .unwrap();
1590        // Garbage collection should have run.
1591        for i in 0..4 {
1592            tracing::info!(i, "check proposal garbage collected");
1593            assert!(
1594                storage
1595                    .load_vid_share(ViewNumber::new(i))
1596                    .await
1597                    .unwrap()
1598                    .is_none()
1599            );
1600            assert!(
1601                storage
1602                    .load_da_proposal(ViewNumber::new(i))
1603                    .await
1604                    .unwrap()
1605                    .is_none()
1606            );
1607        }
1608        tracing::info!("check anchor leaf updated");
1609        assert_eq!(
1610            storage
1611                .load_anchor_leaf()
1612                .await
1613                .unwrap()
1614                .unwrap()
1615                .0
1616                .view_number(),
1617            ViewNumber::new(3)
1618        );
1619        assert_eq!(
1620            storage.load_anchor_view().await.unwrap(),
1621            ViewNumber::new(3)
1622        );
1623
1624        // Check decide event.
1625        tracing::info!("check decide event");
1626        let leaf_chain = consumer.leaf_chain().await;
1627        assert_eq!(leaf_chain.len(), 4, "{leaf_chain:#?}");
1628        for ((leaf, ..), info) in chain.iter().zip(leaf_chain.iter()) {
1629            assert_eq!(info.leaf, *leaf);
1630            let decided_vid_share = info.vid_share.as_ref().unwrap();
1631            assert_eq!(decided_vid_share.view_number(), leaf.view_number());
1632            assert!(info.leaf.block_payload().is_some());
1633        }
1634    }
1635
1636    /// Splitting a decide into `persist_decided_leaves` (persist only, no events/GC) and a later
1637    /// `process_decided_events` loses no data: multiple persists can accumulate, and one process
1638    /// pass at the latest view drains the whole backlog and runs GC.
1639    #[rstest_reuse::apply(persistence_types)]
1640    pub async fn test_deferred_decide_processing<P: TestablePersistence>(_p: PhantomData<P>) {
1641        let tmp = P::tmp_storage().await;
1642        let storage = P::connect(&tmp).await;
1643
1644        // Create a short blockchain (same setup as `test_decide_with_failing_event_consumer`).
1645        let mut chain = vec![];
1646
1647        let leaf: Leaf2 = Leaf::genesis(
1648            &ValidatedState::default(),
1649            &NodeState::mock(),
1650            MOCK_UPGRADE.base,
1651        )
1652        .await
1653        .into();
1654        let leaf_payload = leaf.block_payload().unwrap();
1655        let leaf_payload_bytes_arc = leaf_payload.encode();
1656        let avidm_param = init_avidm_param(2).unwrap();
1657        let weights = vec![1u32; 2];
1658        let ns_table = parse_ns_table(
1659            leaf_payload.byte_len().as_usize(),
1660            &leaf_payload.ns_table().encode(),
1661        );
1662        let (payload_commitment, shares) =
1663            AvidMScheme::ns_disperse(&avidm_param, &weights, &leaf_payload_bytes_arc, ns_table)
1664                .unwrap();
1665
1666        let (pubkey, privkey) = BLSPubKey::generated_from_seed_indexed([0; 32], 1);
1667        let mut vid = AvidMDisperseShare::<SeqTypes> {
1668            view_number: ViewNumber::new(0),
1669            payload_commitment,
1670            share: shares[0].clone(),
1671            recipient_key: pubkey,
1672            epoch: Some(EpochNumber::new(0)),
1673            target_epoch: Some(EpochNumber::new(0)),
1674            common: avidm_param,
1675        }
1676        .to_proposal(&privkey)
1677        .unwrap()
1678        .clone();
1679        let mut quorum_proposal = QuorumProposalWrapper::<SeqTypes> {
1680            proposal: QuorumProposal2::<SeqTypes> {
1681                block_header: leaf.block_header().clone(),
1682                view_number: ViewNumber::genesis(),
1683                justify_qc: QuorumCertificate::genesis(
1684                    &ValidatedState::default(),
1685                    &NodeState::mock(),
1686                    TEST_VERSIONS.test,
1687                )
1688                .await
1689                .to_qc2(),
1690                upgrade_certificate: None,
1691                view_change_evidence: None,
1692                next_drb_result: None,
1693                next_epoch_justify_qc: None,
1694                epoch: None,
1695                state_cert: None,
1696            },
1697        };
1698        let mut qc = QuorumCertificate2::genesis(
1699            &ValidatedState::default(),
1700            &NodeState::mock(),
1701            TEST_VERSIONS.test,
1702        )
1703        .await;
1704
1705        let block_payload_signature = BLSPubKey::sign(&privkey, &leaf_payload_bytes_arc)
1706            .expect("Failed to sign block payload");
1707        let mut da_proposal = Proposal {
1708            data: DaProposal2::<SeqTypes> {
1709                encoded_transactions: leaf_payload_bytes_arc.clone(),
1710                metadata: leaf_payload.ns_table().clone(),
1711                view_number: ViewNumber::new(0),
1712                epoch: Some(EpochNumber::new(0)),
1713                epoch_transition_indicator: EpochTransitionIndicator::NotInTransition,
1714            },
1715            signature: block_payload_signature,
1716            _pd: Default::default(),
1717        };
1718
1719        let vid_commitment = vid_commitment(
1720            &leaf_payload_bytes_arc,
1721            &leaf.block_header().metadata().encode(),
1722            2,
1723            TEST_VERSIONS.test.base,
1724        );
1725
1726        for i in 0..4 {
1727            quorum_proposal.proposal.view_number = ViewNumber::new(i);
1728            let leaf = Leaf2::from_quorum_proposal(&quorum_proposal);
1729            qc.view_number = leaf.view_number();
1730            qc.data.leaf_commit = Committable::commit(&leaf);
1731            vid.data.view_number = leaf.view_number();
1732            da_proposal.data.view_number = leaf.view_number();
1733            chain.push((leaf.clone(), qc.clone(), vid.clone(), da_proposal.clone()));
1734        }
1735
1736        // Per-view artifacts, as persisted at receive time (before any decide).
1737        for (_, _, vid, da) in &chain {
1738            storage.append_da2(da, vid_commitment).await.unwrap();
1739            storage
1740                .append_vid(&convert_proposal(vid.clone()))
1741                .await
1742                .unwrap();
1743        }
1744
1745        let consumer = EventCollector::default();
1746
1747        // Two decides with no processing in between (a lagging/coalescing processor): {0,1}, {2,3}.
1748        for (from, to, view) in [(0usize, 2usize, 1u64), (2, 4, 3)] {
1749            let leaf_chain = chain
1750                .iter()
1751                .take(to)
1752                .skip(from)
1753                .map(|(leaf, qc, ..)| (leaf_info(leaf.clone()), qc.clone()))
1754                .collect::<Vec<_>>();
1755            storage
1756                .persist_decided_leaves(
1757                    ViewNumber::new(view),
1758                    leaf_chain
1759                        .iter()
1760                        .map(|(leaf, qc)| (leaf, CertificatePair::non_epoch_change(qc.clone()))),
1761                    None,
1762                    &consumer,
1763                )
1764                .await
1765                .unwrap();
1766        }
1767
1768        // Persist-only: leaves saved (anchor advanced), but no events emitted and no GC.
1769        assert_eq!(
1770            storage.load_anchor_view().await.unwrap(),
1771            ViewNumber::new(3),
1772            "anchor view should reflect the latest persisted leaf"
1773        );
1774        assert!(
1775            consumer.leaf_chain().await.is_empty(),
1776            "persist_decided_leaves must not emit decide events"
1777        );
1778        for i in 0..4 {
1779            assert!(
1780                storage
1781                    .load_vid_share(ViewNumber::new(i))
1782                    .await
1783                    .unwrap()
1784                    .is_some(),
1785                "persist_decided_leaves must not garbage collect VID shares"
1786            );
1787            assert!(
1788                storage
1789                    .load_da_proposal(ViewNumber::new(i))
1790                    .await
1791                    .unwrap()
1792                    .is_some(),
1793                "persist_decided_leaves must not garbage collect DA proposals"
1794            );
1795        }
1796
1797        // A failing consumer propagates the error and leaves the cursor un-advanced: nothing is
1798        // GC'd and the range is retried below.
1799        storage
1800            .process_decided_events(ViewNumber::new(3), None, &FailConsumer)
1801            .await
1802            .unwrap_err();
1803        for i in 0..4 {
1804            assert!(
1805                storage
1806                    .load_da_proposal(ViewNumber::new(i))
1807                    .await
1808                    .unwrap()
1809                    .is_some(),
1810                "a failed process pass must not garbage collect anything"
1811            );
1812        }
1813
1814        // One process pass at the latest view drains the whole backlog, runs GC, and reports the
1815        // cursor it advanced to.
1816        let processed = storage
1817            .process_decided_events(ViewNumber::new(3), None, &consumer)
1818            .await
1819            .unwrap();
1820        assert_eq!(
1821            processed,
1822            Some(ViewNumber::new(3)),
1823            "process_decided_events should report the highest processed view"
1824        );
1825
1826        // All four leaves delivered, with payloads and VID shares reconstructed from storage.
1827        let leaf_chain = consumer.leaf_chain().await;
1828        assert_eq!(leaf_chain.len(), 4, "{leaf_chain:#?}");
1829        for ((leaf, ..), info) in chain.iter().zip(leaf_chain.iter()) {
1830            assert_eq!(info.leaf, *leaf);
1831            let decided_vid_share = info.vid_share.as_ref().unwrap();
1832            assert_eq!(decided_vid_share.view_number(), leaf.view_number());
1833            assert!(info.leaf.block_payload().is_some());
1834        }
1835
1836        // GC ran for the processed range.
1837        for i in 0..4 {
1838            assert!(
1839                storage
1840                    .load_vid_share(ViewNumber::new(i))
1841                    .await
1842                    .unwrap()
1843                    .is_none(),
1844                "process_decided_events should have garbage collected VID shares"
1845            );
1846            assert!(
1847                storage
1848                    .load_da_proposal(ViewNumber::new(i))
1849                    .await
1850                    .unwrap()
1851                    .is_none(),
1852                "process_decided_events should have garbage collected DA proposals"
1853            );
1854        }
1855
1856        // Re-processing with nothing new is a no-op.
1857        let consumer2 = EventCollector::default();
1858        storage
1859            .process_decided_events(ViewNumber::new(3), None, &consumer2)
1860            .await
1861            .unwrap();
1862        assert!(
1863            consumer2.leaf_chain().await.is_empty(),
1864            "re-processing an already-drained backlog must emit no events"
1865        );
1866        assert_eq!(
1867            storage.load_anchor_view().await.unwrap(),
1868            ViewNumber::new(3)
1869        );
1870    }
1871
1872    /// A chain of `n` V6 mock leaves at views `0..n` with real, consecutive block heights (the
1873    /// other mock chains in this module reuse the genesis header, so every leaf has height 0).
1874    async fn consecutive_height_chain(n: u64) -> Vec<(Leaf2, QuorumCertificate2<SeqTypes>)> {
1875        let node_state = NodeState::mock().with_genesis_version(versions::NEW_PROTOCOL_VERSION);
1876        let genesis_leaf = Leaf2::genesis(
1877            &ValidatedState::default(),
1878            &node_state,
1879            TEST_VERSIONS.test.base,
1880        )
1881        .await;
1882        let mut quorum_proposal = QuorumProposalWrapper::<SeqTypes> {
1883            proposal: QuorumProposal2::<SeqTypes> {
1884                block_header: genesis_leaf.block_header().clone(),
1885                view_number: ViewNumber::genesis(),
1886                justify_qc: QuorumCertificate2::genesis(
1887                    &ValidatedState::default(),
1888                    &node_state,
1889                    TEST_VERSIONS.test,
1890                )
1891                .await,
1892                upgrade_certificate: None,
1893                view_change_evidence: None,
1894                next_drb_result: None,
1895                next_epoch_justify_qc: None,
1896                epoch: None,
1897                state_cert: None,
1898            },
1899        };
1900        let mut qc = QuorumCertificate2::genesis(
1901            &ValidatedState::default(),
1902            &node_state,
1903            TEST_VERSIONS.test,
1904        )
1905        .await;
1906
1907        let mut chain = vec![];
1908        for i in 0..n {
1909            quorum_proposal.proposal.view_number = ViewNumber::new(i);
1910            *quorum_proposal.proposal.block_header.height_mut() = i;
1911            let leaf = Leaf2::from_quorum_proposal(&quorum_proposal);
1912            qc.view_number = leaf.view_number();
1913            qc.data.leaf_commit = Committable::commit(&leaf);
1914            chain.push((leaf, qc.clone()));
1915        }
1916        chain
1917    }
1918
1919    /// Decide the leaves of `chain` at indices (== views) `range`.
1920    async fn decide_range<P: TestablePersistence>(
1921        storage: &P,
1922        chain: &[(Leaf2, QuorumCertificate2<SeqTypes>)],
1923        range: std::ops::Range<usize>,
1924        consumer: &(impl EventConsumer + 'static),
1925    ) {
1926        let decided_view = ViewNumber::new(range.end as u64 - 1);
1927        let leaf_chain = chain[range]
1928            .iter()
1929            .map(|(leaf, qc)| (leaf_info(leaf.clone()), qc.clone()))
1930            .collect::<Vec<_>>();
1931        storage
1932            .append_decided_leaves(
1933                decided_view,
1934                leaf_chain
1935                    .iter()
1936                    .map(|(leaf, qc)| (leaf, CertificatePair::non_epoch_change(qc.clone()))),
1937                None,
1938                consumer,
1939            )
1940            .await
1941            .unwrap();
1942    }
1943
1944    /// Records the views delivered in decide events, in delivery order.
1945    #[derive(Clone, Debug, Default)]
1946    struct DecideViewCollector {
1947        views: Arc<RwLock<Vec<u64>>>,
1948    }
1949
1950    #[async_trait]
1951    impl EventConsumer for DecideViewCollector {
1952        async fn handle_event(&self, event: &CoordinatorEvent<SeqTypes>) -> anyhow::Result<()> {
1953            let mut views = self.views.write().await;
1954            match event {
1955                // Decide events carry their leaves newest first.
1956                CoordinatorEvent::NewDecide { leaf_infos, .. } => views.extend(
1957                    leaf_infos
1958                        .iter()
1959                        .rev()
1960                        .map(|info| info.leaf.view_number().u64()),
1961                ),
1962                CoordinatorEvent::LegacyEvent(Event {
1963                    event: EventType::Decide { leaf_chain, .. },
1964                    ..
1965                }) => views.extend(
1966                    leaf_chain
1967                        .iter()
1968                        .rev()
1969                        .map(|info| info.leaf.view_number().u64()),
1970                ),
1971                _ => {},
1972            }
1973            Ok(())
1974        }
1975    }
1976
1977    async fn received_views(consumer: &DecideViewCollector) -> Vec<u64> {
1978        consumer.views.read().await.clone()
1979    }
1980
1981    /// Until a gap-fill decide arrives, event processing holds its cursor at the gap: nothing
1982    /// past it is delivered, and the pending leaves are neither skipped nor dropped.
1983    ///
1984    /// sql-only: the fs backend keeps the pop-oldest cursor and skips gaps (the consumer
1985    /// re-fetches missing blocks), trading gap-fill delivery for cheap decide passes.
1986    #[rstest::rstest]
1987    #[case(PhantomData::<crate::persistence::sql::Persistence>)]
1988    #[test_log::test(tokio::test(flavor = "multi_thread"))]
1989    pub async fn test_decide_gap_holds_events<P: TestablePersistence>(#[case] _p: PhantomData<P>) {
1990        let tmp = P::tmp_storage().await;
1991        let storage = P::connect(&tmp).await;
1992        let consumer = DecideViewCollector::default();
1993        let chain = consecutive_height_chain(5).await;
1994
1995        decide_range(&storage, &chain, 0..2, &consumer).await;
1996        assert_eq!(received_views(&consumer).await, vec![0, 1]);
1997
1998        // Views 3..=4 decide while view 2 is still pending.
1999        decide_range(&storage, &chain, 3..5, &consumer).await;
2000        assert_eq!(
2001            received_views(&consumer).await,
2002            vec![0, 1],
2003            "leaves past a height gap must not be delivered before the gap is filled"
2004        );
2005
2006        // Nothing was lost: the pending leaves are still in storage...
2007        assert_eq!(
2008            storage
2009                .load_anchor_leaf()
2010                .await
2011                .unwrap()
2012                .unwrap()
2013                .0
2014                .view_number(),
2015            ViewNumber::new(4)
2016        );
2017
2018        // ...and re-processing neither skips the gap nor emits duplicates.
2019        storage
2020            .process_decided_events(ViewNumber::new(4), None, &consumer)
2021            .await
2022            .unwrap();
2023        assert_eq!(received_views(&consumer).await, vec![0, 1]);
2024    }
2025
2026    /// Once the gap-fill decide arrives, the consumer receives the gap leaf and everything held
2027    /// up behind it, in order, exactly once each.
2028    ///
2029    /// sql-only: see [`test_decide_gap_holds_events`].
2030    #[rstest::rstest]
2031    #[case(PhantomData::<crate::persistence::sql::Persistence>)]
2032    #[test_log::test(tokio::test(flavor = "multi_thread"))]
2033    pub async fn test_decide_gap_fill_delivers_all<P: TestablePersistence>(
2034        #[case] _p: PhantomData<P>,
2035    ) {
2036        let tmp = P::tmp_storage().await;
2037        let storage = P::connect(&tmp).await;
2038        let consumer = DecideViewCollector::default();
2039        let chain = consecutive_height_chain(5).await;
2040
2041        decide_range(&storage, &chain, 0..2, &consumer).await;
2042        decide_range(&storage, &chain, 3..5, &consumer).await;
2043        assert_eq!(received_views(&consumer).await, vec![0, 1]);
2044
2045        // The gap-fill decide of view 2 arrives.
2046        decide_range(&storage, &chain, 2..3, &consumer).await;
2047        // A re-process pass at the watermark must not skip or duplicate anything.
2048        storage
2049            .process_decided_events(ViewNumber::new(4), None, &consumer)
2050            .await
2051            .unwrap();
2052
2053        assert_eq!(
2054            received_views(&consumer).await,
2055            vec![0, 1, 2, 3, 4],
2056            "after the gap fills, every leaf is delivered in order, exactly once"
2057        );
2058    }
2059
2060    #[rstest_reuse::apply(persistence_types)]
2061    pub async fn test_pruning<P: TestablePersistence>(_p: PhantomData<P>) {
2062        let tmp = P::tmp_storage().await;
2063
2064        let mut options = P::options(&tmp);
2065        options.set_view_retention(1);
2066        let storage = options.create().await.unwrap();
2067
2068        // Add some "old" data, from view 0.
2069        let leaf = Leaf::genesis(
2070            &ValidatedState::default(),
2071            &NodeState::mock(),
2072            MOCK_UPGRADE.base,
2073        )
2074        .await;
2075        let leaf_payload = leaf.block_payload().unwrap();
2076        let leaf_payload_bytes_arc = leaf_payload.encode();
2077        let avidm_param = init_avidm_param(2).unwrap();
2078        let weights = vec![1u32; 2];
2079
2080        let ns_table = parse_ns_table(
2081            leaf_payload.byte_len().as_usize(),
2082            &leaf_payload.ns_table().encode(),
2083        );
2084        let (payload_commitment, shares) =
2085            AvidMScheme::ns_disperse(&avidm_param, &weights, &leaf_payload_bytes_arc, ns_table)
2086                .unwrap();
2087
2088        let (pubkey, privkey) = BLSPubKey::generated_from_seed_indexed([0; 32], 1);
2089        let vid_share = convert_proposal(
2090            AvidMDisperseShare::<SeqTypes> {
2091                view_number: ViewNumber::new(0),
2092                payload_commitment,
2093                share: shares[0].clone(),
2094                recipient_key: pubkey,
2095                epoch: None,
2096                target_epoch: None,
2097                common: avidm_param,
2098            }
2099            .to_proposal(&privkey)
2100            .unwrap()
2101            .clone(),
2102        );
2103
2104        let quorum_proposal = QuorumProposalWrapper::<SeqTypes> {
2105            proposal: QuorumProposal2::<SeqTypes> {
2106                block_header: leaf.block_header().clone(),
2107                view_number: ViewNumber::genesis(),
2108                justify_qc: QuorumCertificate::genesis(
2109                    &ValidatedState::default(),
2110                    &NodeState::mock(),
2111                    TEST_VERSIONS.test,
2112                )
2113                .await
2114                .to_qc2(),
2115                upgrade_certificate: None,
2116                view_change_evidence: None,
2117                next_drb_result: None,
2118                next_epoch_justify_qc: None,
2119                epoch: None,
2120                state_cert: None,
2121            },
2122        };
2123        let quorum_proposal_signature =
2124            BLSPubKey::sign(&privkey, &bincode::serialize(&quorum_proposal).unwrap())
2125                .expect("Failed to sign quorum proposal");
2126        let quorum_proposal = Proposal {
2127            data: quorum_proposal,
2128            signature: quorum_proposal_signature,
2129            _pd: Default::default(),
2130        };
2131
2132        let block_payload_signature = BLSPubKey::sign(&privkey, &leaf_payload_bytes_arc)
2133            .expect("Failed to sign block payload");
2134        let da_proposal = Proposal {
2135            data: DaProposal2::<SeqTypes> {
2136                encoded_transactions: leaf_payload_bytes_arc,
2137                metadata: leaf_payload.ns_table().clone(),
2138                view_number: ViewNumber::new(0),
2139                epoch: None,
2140                epoch_transition_indicator: EpochTransitionIndicator::NotInTransition,
2141            },
2142            signature: block_payload_signature,
2143            _pd: Default::default(),
2144        };
2145
2146        let leaf2: Leaf2 = leaf.clone().into();
2147        let vote_data = Vote2Data {
2148            leaf_commit: leaf2.commit(),
2149            epoch: EpochNumber::new(0),
2150            block_number: leaf2.height(),
2151        };
2152        let cert2 = Certificate2::new(
2153            vote_data.clone(),
2154            vote_data.commit(),
2155            ViewNumber::new(0),
2156            None,
2157            PhantomData,
2158        );
2159
2160        storage
2161            .append_da2(&da_proposal, VidCommitment::V1(payload_commitment))
2162            .await
2163            .unwrap();
2164        storage.append_vid(&vid_share).await.unwrap();
2165        storage
2166            .append_quorum_proposal2(&quorum_proposal)
2167            .await
2168            .unwrap();
2169        storage
2170            .append_cert2(ViewNumber::new(0), cert2)
2171            .await
2172            .unwrap();
2173
2174        // Decide a newer view, view 1.
2175        storage
2176            .append_decided_leaves(ViewNumber::new(1), [], None, &NullEventConsumer)
2177            .await
2178            .unwrap();
2179
2180        // The old data is not more than the retention period (1 view) old, so it should not be
2181        // GCed.
2182        assert_eq!(
2183            storage
2184                .load_da_proposal(ViewNumber::new(0))
2185                .await
2186                .unwrap()
2187                .unwrap(),
2188            da_proposal
2189        );
2190        assert_eq!(
2191            storage
2192                .load_vid_share(ViewNumber::new(0))
2193                .await
2194                .unwrap()
2195                .unwrap(),
2196            vid_share
2197        );
2198        assert_eq!(
2199            storage
2200                .load_quorum_proposal(ViewNumber::new(0))
2201                .await
2202                .unwrap(),
2203            quorum_proposal
2204        );
2205        assert!(
2206            storage
2207                .load_cert2(ViewNumber::new(0))
2208                .await
2209                .unwrap()
2210                .is_some()
2211        );
2212
2213        // Decide an even newer view, triggering GC of the old data.
2214        storage
2215            .append_decided_leaves(ViewNumber::new(2), [], None, &NullEventConsumer)
2216            .await
2217            .unwrap();
2218        assert!(
2219            storage
2220                .load_da_proposal(ViewNumber::new(0))
2221                .await
2222                .unwrap()
2223                .is_none()
2224        );
2225        assert!(
2226            storage
2227                .load_vid_share(ViewNumber::new(0))
2228                .await
2229                .unwrap()
2230                .is_none()
2231        );
2232        assert!(
2233            storage
2234                .load_quorum_proposal(ViewNumber::new(0))
2235                .await
2236                .is_err()
2237        );
2238        assert!(
2239            storage
2240                .load_cert2(ViewNumber::new(0))
2241                .await
2242                .unwrap()
2243                .is_none()
2244        );
2245    }
2246
2247    async fn assert_events_eq<P: TestablePersistence>(
2248        persistence: &P,
2249        block: u64,
2250        stake_table_fetcher: &Fetcher,
2251        l1_client: &L1Client,
2252        stake_table_contract: Address,
2253    ) -> anyhow::Result<()> {
2254        // Load persisted events
2255        let (stored_l1, events) = persistence.load_events(0, block).await?;
2256        assert!(!events.is_empty());
2257        assert!(stored_l1.is_some());
2258        assert!(events.iter().all(|((l1_block, _), _)| *l1_block <= block));
2259        // Fetch events directly from the contract and compare with persisted data
2260        let contract_events = Fetcher::fetch_events_from_contract(
2261            l1_client.clone(),
2262            stake_table_contract,
2263            None,
2264            block,
2265        )
2266        .await?;
2267        assert_eq!(
2268            contract_events, events,
2269            "Events from contract and persistence do not match"
2270        );
2271
2272        // Fetch events from stake table fetcher and compare with persisted data
2273        let fetched_events = stake_table_fetcher
2274            .fetch_and_store_stake_table_events(stake_table_contract, block)
2275            .await?;
2276        assert_eq!(fetched_events, events);
2277
2278        Ok(())
2279    }
2280
2281    // test for validating stake table event fetching from persistence,
2282    // ensuring that persisted data matches the on-chain events and that event fetcher work correctly.
2283    #[rstest_reuse::apply(persistence_types)]
2284    pub async fn test_stake_table_fetching_from_persistence<P: TestablePersistence>(
2285        #[values(
2286            StakeTableContractVersion::V1,
2287            StakeTableContractVersion::V2,
2288            StakeTableContractVersion::V3
2289        )]
2290        stake_table_version: StakeTableContractVersion,
2291        _p: PhantomData<P>,
2292    ) -> anyhow::Result<()> {
2293        let epoch_height = 20;
2294
2295        let network_config = TestConfigBuilder::default()
2296            .epoch_height(epoch_height)
2297            .build();
2298
2299        let anvil_provider = network_config.anvil().unwrap();
2300
2301        let query_service_port =
2302            reserve_tcp_port().expect("OS should have ephemeral ports available");
2303        let query_api_options = Options::with_port(query_service_port);
2304
2305        const NODE_COUNT: usize = 2;
2306
2307        let storage = join_all((0..NODE_COUNT).map(|_| P::tmp_storage())).await;
2308        let persistence_options: [_; NODE_COUNT] = storage
2309            .iter()
2310            .map(P::options)
2311            .collect::<Vec<_>>()
2312            .try_into()
2313            .unwrap();
2314
2315        let persistence = persistence_options[0].clone().create().await.unwrap();
2316
2317        // Build the config with PoS hook
2318        let l1_url = network_config.l1_url();
2319
2320        let upgrade = Upgrade::trivial(version(0, 3));
2321
2322        let testnet_config = TestNetworkConfigBuilder::with_num_nodes()
2323            .api_config(query_api_options)
2324            .network_config(network_config.clone())
2325            .persistences(persistence_options.clone())
2326            .pos_hook(
2327                DelegationConfig::MultipleDelegators,
2328                stake_table_version,
2329                upgrade,
2330            )
2331            .await
2332            .expect("Pos deployment failed")
2333            .build();
2334
2335        //start the network
2336        let test_network = TestNetwork::new(testnet_config, upgrade).await;
2337
2338        let client: Client<ClientErr, SequencerApiVersion> = Client::new(
2339            format!("http://localhost:{query_service_port}")
2340                .parse()
2341                .unwrap(),
2342        );
2343        client.connect(None).await;
2344        tracing::info!(query_service_port, "server running");
2345
2346        // wait until we enter in epoch 3
2347        let _initial_blocks = client
2348            .socket("availability/stream/blocks/0")
2349            .subscribe::<BlockQueryData<SeqTypes>>()
2350            .await
2351            .unwrap()
2352            .take(40)
2353            .try_collect::<Vec<_>>()
2354            .await
2355            .unwrap();
2356        // Load initial persisted events and validate they exist.
2357        let membership_coordinator = test_network
2358            .server
2359            .consensus_handle()
2360            .membership_coordinator()
2361            .await;
2362
2363        let l1_client = L1Client::new(vec![l1_url]).unwrap();
2364        let node_state = test_network.server.node_state();
2365        let chain_config = node_state.chain_config;
2366        let stake_table_contract = chain_config.stake_table_contract.unwrap();
2367
2368        let current_membership = membership_coordinator.membership();
2369        {
2370            let membership_state = current_membership;
2371            let stake_table_fetcher = membership_state.fetcher();
2372
2373            let block1 = anvil_provider
2374                .get_block_number()
2375                .await
2376                .expect("latest l1 block");
2377
2378            assert_events_eq(
2379                &persistence,
2380                block1,
2381                stake_table_fetcher,
2382                &l1_client,
2383                stake_table_contract,
2384            )
2385            .await?;
2386        }
2387        let _epoch_4_blocks = client
2388            .socket("availability/stream/blocks/0")
2389            .subscribe::<BlockQueryData<SeqTypes>>()
2390            .await
2391            .unwrap()
2392            .take(65)
2393            .try_collect::<Vec<_>>()
2394            .await
2395            .unwrap();
2396        let block2 = anvil_provider
2397            .get_block_number()
2398            .await
2399            .expect("latest l1 block");
2400
2401        {
2402            let membership_state = current_membership;
2403            let stake_table_fetcher = membership_state.fetcher();
2404
2405            assert_events_eq(
2406                &persistence,
2407                block2,
2408                stake_table_fetcher,
2409                &l1_client,
2410                stake_table_contract,
2411            )
2412            .await?;
2413        }
2414        Ok(())
2415    }
2416
2417    #[rstest_reuse::apply(persistence_types)]
2418    pub async fn test_stake_table_background_fetching<P: TestablePersistence>(
2419        #[values(
2420            StakeTableContractVersion::V1,
2421            StakeTableContractVersion::V2,
2422            StakeTableContractVersion::V3
2423        )]
2424        stake_table_version: StakeTableContractVersion,
2425        _p: PhantomData<P>,
2426    ) -> anyhow::Result<()> {
2427        use espresso_types::v0_3::ChainConfig;
2428        use hotshot_contract_adapter::stake_table::StakeTableContractVersion;
2429
2430        let blocks_per_epoch = 10;
2431
2432        let network_config = TestConfigBuilder::<1>::default()
2433            .epoch_height(blocks_per_epoch)
2434            .build();
2435
2436        let anvil_provider = network_config.anvil().unwrap();
2437
2438        let (genesis_state, genesis_stake) = light_client_genesis_from_stake_table(
2439            &network_config.hotshot_config().hotshot_stake_table(),
2440            STAKE_TABLE_CAPACITY_FOR_TEST,
2441        )
2442        .unwrap();
2443
2444        let (_, priv_keys): (Vec<_>, Vec<_>) = (0..20)
2445            .map(|i| <PubKey as SignatureKey>::generated_from_seed_indexed([1; 32], i as u64))
2446            .unzip();
2447        let state_key_pairs = (0..20)
2448            .map(|i| StateKeyPair::generate_from_seed_indexed([2; 32], i as u64))
2449            .collect::<Vec<_>>();
2450
2451        let validators = staking_priv_keys(&priv_keys, &state_key_pairs, &[], 20);
2452
2453        let deployer = ProviderBuilder::new()
2454            .wallet(EthereumWallet::from(network_config.signer().clone()))
2455            .connect_http(network_config.l1_url().clone());
2456
2457        let mut contracts = Contracts::new();
2458        let args = DeployerArgsBuilder::default()
2459            .deployer(deployer.clone())
2460            .rpc_url(network_config.l1_url().clone())
2461            .mock_light_client(true)
2462            .genesis_lc_state(genesis_state)
2463            .genesis_st_state(genesis_stake)
2464            .blocks_per_epoch(blocks_per_epoch)
2465            .epoch_start_block(1)
2466            .exit_escrow_period(U256::from(max(
2467                blocks_per_epoch * 15 + 100,
2468                DEFAULT_EXIT_ESCROW_PERIOD_SECONDS,
2469            )))
2470            .multisig_pauser(network_config.signer().address())
2471            .token_name("Espresso".to_string())
2472            .token_symbol("ESP".to_string())
2473            .initial_token_supply(U256::from(3590000000u64))
2474            .ops_timelock_delay(U256::from(0))
2475            .ops_timelock_admin(network_config.signer().address())
2476            .ops_timelock_proposers(vec![network_config.signer().address()])
2477            .ops_timelock_executors(vec![network_config.signer().address()])
2478            .safe_exit_timelock_delay(U256::from(10))
2479            .safe_exit_timelock_admin(network_config.signer().address())
2480            .safe_exit_timelock_proposers(vec![network_config.signer().address()])
2481            .safe_exit_timelock_executors(vec![network_config.signer().address()])
2482            .build()
2483            .unwrap();
2484
2485        match stake_table_version {
2486            StakeTableContractVersion::V1 => args.deploy_to_stake_table_v1(&mut contracts).await,
2487            StakeTableContractVersion::V2 => args.deploy_to_stake_table_v2(&mut contracts).await,
2488            StakeTableContractVersion::V3 => args.deploy_to_stake_table_v3(&mut contracts).await,
2489        }
2490        .expect("contracts deployed");
2491
2492        let st_addr = contracts
2493            .address(Contract::StakeTableProxy)
2494            .expect("StakeTableProxy deployed");
2495        let l1_url = network_config.l1_url().clone();
2496
2497        let mut planned_txns = StakingTransactions::create(
2498            l1_url.clone(),
2499            &deployer,
2500            st_addr,
2501            validators,
2502            None,
2503            DelegationConfig::MultipleDelegators,
2504        )
2505        .await
2506        .expect("stake table setup failed");
2507
2508        planned_txns
2509            .apply_prerequisites()
2510            .await
2511            .expect("prerequisites failed");
2512
2513        // Ensure we have at least one stake table affecting transaction
2514        planned_txns.apply_one().await.expect("send tx failed");
2515
2516        // new block every 1s
2517        anvil_provider
2518            .anvil_set_interval_mining(1)
2519            .await
2520            .expect("interval mining");
2521
2522        // spawn a separate task
2523        // this is going to keep registering validators and multiple delegators
2524        // the interval mining is set to 1s so each transaction finalization would take atleast 1s
2525        spawn({
2526            async move {
2527                {
2528                    while let Some(receipt) =
2529                        planned_txns.apply_one().await.expect("send tx failed")
2530                    {
2531                        tracing::debug!(?receipt, "transaction finalized");
2532                    }
2533                }
2534            }
2535        });
2536
2537        let storage = P::tmp_storage().await;
2538        let persistence = P::options(&storage).create().await.unwrap();
2539
2540        let l1_client = L1ClientOptions {
2541            stake_table_update_interval: Duration::from_secs(7),
2542            l1_retry_delay: Duration::from_millis(10),
2543            l1_events_max_block_range: 10000,
2544            ..Default::default()
2545        }
2546        .connect(vec![l1_url])
2547        .unwrap();
2548        l1_client.spawn_tasks().await;
2549
2550        let fetcher = Fetcher::new(
2551            Arc::new(NullStateCatchup::default()),
2552            Arc::new(Mutex::new(persistence.clone())),
2553            l1_client.clone(),
2554            ChainConfig {
2555                stake_table_contract: Some(st_addr),
2556                base_fee: 0.into(),
2557                ..Default::default()
2558            },
2559        );
2560
2561        // sleep so that we have enough events
2562        sleep(Duration::from_secs(20)).await;
2563
2564        fetcher.spawn_update_loop().await;
2565        let mut prev_l1_block = 0;
2566        let mut prev_events_len = 0;
2567        for _i in 0..10 {
2568            // Wait for more than update interval to assert that persistence was updated
2569            // L1 update interval is 7s in this test
2570
2571            tokio::time::sleep(std::time::Duration::from_secs(8)).await;
2572
2573            let block = anvil_provider
2574                .get_block_number()
2575                .await
2576                .expect("latest l1 block");
2577
2578            let (read_offset, persisted_events) = persistence.load_events(0, block).await?;
2579            let read_offset = read_offset.unwrap();
2580            let l1_block = match read_offset {
2581                EventsPersistenceRead::Complete => block,
2582                EventsPersistenceRead::UntilL1Block(block) => block,
2583            };
2584
2585            tracing::info!("{l1_block:?}, persistence events = {persisted_events:?}.");
2586            assert!(persisted_events.len() > prev_events_len);
2587
2588            assert!(l1_block > prev_l1_block, "events not updated");
2589
2590            let contract_events =
2591                Fetcher::fetch_events_from_contract(l1_client.clone(), st_addr, None, l1_block)
2592                    .await?;
2593            assert_eq!(persisted_events, contract_events);
2594
2595            prev_l1_block = l1_block;
2596            prev_events_len = persisted_events.len();
2597        }
2598
2599        Ok(())
2600    }
2601
2602    #[rstest_reuse::apply(persistence_types)]
2603    pub async fn test_membership_persistence<P: TestablePersistence>(
2604        _p: PhantomData<P>,
2605    ) -> anyhow::Result<()> {
2606        let tmp = P::tmp_storage().await;
2607        let mut opt = P::options(&tmp);
2608
2609        let storage = opt.create().await.unwrap();
2610
2611        let validator = AuthenticatedValidator::mock();
2612        let mut st = IndexMap::new();
2613        st.insert(validator.account, validator);
2614
2615        storage
2616            .store_stake(EpochNumber::new(10), st.clone(), None, None)
2617            .await?;
2618
2619        let (table, ..) = storage.load_stake(EpochNumber::new(10)).await?.unwrap();
2620        assert_eq!(st, table);
2621
2622        let val2 = AuthenticatedValidator::mock();
2623        let mut st2 = IndexMap::new();
2624        st2.insert(val2.account, val2);
2625        storage
2626            .store_stake(EpochNumber::new(11), st2.clone(), None, None)
2627            .await?;
2628
2629        let tables = storage.load_latest_stake(4).await?.unwrap();
2630        let mut iter = tables.iter();
2631        assert_eq!(
2632            Some(&(EpochNumber::new(11), (st2.clone(), None), None)),
2633            iter.next()
2634        );
2635        assert_eq!(Some(&(EpochNumber::new(10), (st, None), None)), iter.next());
2636        assert_eq!(None, iter.next());
2637
2638        for i in 0..=20 {
2639            storage
2640                .store_stake(EpochNumber::new(i), st2.clone(), None, None)
2641                .await?;
2642        }
2643
2644        let tables = storage.load_latest_stake(5).await?.unwrap();
2645        let mut iter = tables.iter();
2646        assert_eq!(
2647            Some(&(EpochNumber::new(20), (st2.clone(), None), None)),
2648            iter.next()
2649        );
2650        assert_eq!(
2651            Some(&(EpochNumber::new(19), (st2.clone(), None), None)),
2652            iter.next()
2653        );
2654        assert_eq!(
2655            Some(&(EpochNumber::new(18), (st2.clone(), None), None)),
2656            iter.next()
2657        );
2658        assert_eq!(
2659            Some(&(EpochNumber::new(17), (st2.clone(), None), None)),
2660            iter.next()
2661        );
2662        assert_eq!(
2663            Some(&(EpochNumber::new(16), (st2, None), None)),
2664            iter.next()
2665        );
2666        assert_eq!(None, iter.next());
2667
2668        Ok(())
2669    }
2670
2671    // Epoch data is stored piecemeal: `store_drb_result`, `store_epoch_root`
2672    // and `store_stake` each write only their part (one row in SQL, separate
2673    // files in FS). A loader must report data another writer created the
2674    // epoch entry for as absent — regression test for `load_stake` panicking
2675    // on a DRB-only row's NULL stake column.
2676    #[rstest_reuse::apply(persistence_types)]
2677    pub async fn test_load_partially_stored_epoch_data<P: TestablePersistence>(
2678        _p: PhantomData<P>,
2679    ) -> anyhow::Result<()> {
2680        let tmp = P::tmp_storage().await;
2681        let storage = P::connect(&tmp).await;
2682
2683        let instance_state = NodeState::mock();
2684        let validated_state = hotshot_types::traits::ValidatedState::genesis(&instance_state).0;
2685        let leaf: Leaf2 = Leaf::genesis(&validated_state, &instance_state, MOCK_UPGRADE.base)
2686            .await
2687            .into();
2688        let header = leaf.block_header().clone();
2689
2690        let validator = AuthenticatedValidator::mock();
2691        let mut stake = IndexMap::new();
2692        stake.insert(validator.account, validator);
2693
2694        // Nothing stored yet.
2695        let epoch = EpochNumber::new(1);
2696        assert!(storage.load_stake(epoch).await?.is_none());
2697        assert!(storage.load_drb_result(epoch).await?.is_none());
2698        assert!(storage.load_epoch_root(epoch).await?.is_none());
2699
2700        // A DRB result alone; the other loaders report their data as absent.
2701        storage.store_drb_result(epoch, [1; 32]).await?;
2702        assert!(storage.load_stake(epoch).await?.is_none());
2703        assert_eq!(storage.load_drb_result(epoch).await?, Some([1; 32]));
2704        assert!(storage.load_epoch_root(epoch).await?.is_none());
2705
2706        // A stake table alone, likewise.
2707        let epoch2 = EpochNumber::new(2);
2708        storage
2709            .store_stake(epoch2, stake.clone(), None, None)
2710            .await?;
2711        assert!(storage.load_stake(epoch2).await?.is_some());
2712        assert!(storage.load_drb_result(epoch2).await?.is_none());
2713        assert!(storage.load_epoch_root(epoch2).await?.is_none());
2714
2715        // With everything stored for one epoch, each loader returns its
2716        // value and none has clobbered the others.
2717        storage.store_epoch_root(epoch, header.clone()).await?;
2718        storage
2719            .store_stake(epoch, stake.clone(), None, None)
2720            .await?;
2721        let (table, ..) = storage.load_stake(epoch).await?.unwrap();
2722        assert_eq!(stake, table);
2723        assert_eq!(storage.load_drb_result(epoch).await?, Some([1; 32]));
2724        assert_eq!(storage.load_epoch_root(epoch).await?, Some(header));
2725
2726        Ok(())
2727    }
2728
2729    #[rstest_reuse::apply(persistence_types)]
2730    pub async fn test_delete_stake_tables<P: TestablePersistence>(
2731        _p: PhantomData<P>,
2732    ) -> anyhow::Result<()> {
2733        let tmp = P::tmp_storage().await;
2734        let mut opt = P::options(&tmp);
2735        let storage = opt.create().await.unwrap();
2736
2737        let event1 = StakeTableEvent::Delegate(Delegated {
2738            delegator: Address::ZERO,
2739            validator: Address::ZERO,
2740            amount: U256::from(100),
2741        });
2742        let event2 = StakeTableEvent::Delegate(Delegated {
2743            delegator: Address::ZERO,
2744            validator: Address::ZERO,
2745            amount: U256::from(200),
2746        });
2747
2748        let l1_block = 42u64;
2749        let events: Vec<(EventKey, StakeTableEvent)> =
2750            vec![((l1_block, 0), event1), ((l1_block, 1), event2)];
2751
2752        storage.store_events(l1_block, events.clone()).await?;
2753
2754        let (read_offset, loaded_events) = storage.load_events(0_u64, l1_block).await?;
2755        assert!(read_offset.is_some());
2756        assert_eq!(loaded_events.len(), 2);
2757        assert_eq!(loaded_events, events);
2758
2759        // Store some validators
2760        let v = RegisteredValidator::mock();
2761        let mut vmap = IndexMap::new();
2762        vmap.insert(v.account, v);
2763        storage
2764            .store_all_validators(EpochNumber::new(1), vmap.clone())
2765            .await?;
2766
2767        let loaded = storage
2768            .load_all_validators(EpochNumber::new(1), 0, 10)
2769            .await?;
2770        assert_eq!(loaded.len(), 1);
2771
2772        storage.delete_stake_tables().await?;
2773
2774        // Events cleared
2775        let (read_offset, loaded_events) = storage.load_events(0_u64, l1_block).await?;
2776        assert!(read_offset.is_none());
2777        assert!(loaded_events.is_empty());
2778
2779        // Validators cleared
2780        let loaded = storage
2781            .load_all_validators(EpochNumber::new(1), 0, 10)
2782            .await
2783            .unwrap_or_default();
2784        assert!(loaded.is_empty());
2785
2786        let event3 = StakeTableEvent::Delegate(Delegated {
2787            delegator: Address::ZERO,
2788            validator: Address::ZERO,
2789            amount: U256::from(300),
2790        });
2791        let new_events: Vec<(EventKey, StakeTableEvent)> = vec![((l1_block, 0), event3)];
2792        storage.store_events(l1_block, new_events.clone()).await?;
2793
2794        let (read_offset, loaded_events) = storage.load_events(0_u64, l1_block).await?;
2795        assert!(read_offset.is_some());
2796        assert_eq!(loaded_events.len(), 1);
2797        assert_eq!(loaded_events, new_events);
2798
2799        Ok(())
2800    }
2801
2802    #[rstest_reuse::apply(persistence_types)]
2803    pub async fn test_store_and_load_all_validators<P: TestablePersistence>(
2804        _p: PhantomData<P>,
2805    ) -> anyhow::Result<()> {
2806        let tmp = P::tmp_storage().await;
2807        let mut opt = P::options(&tmp);
2808        let storage = opt.create().await.unwrap();
2809
2810        let mut vmap1 = IndexMap::new();
2811        for _i in 0..25 {
2812            let v = RegisteredValidator::mock();
2813            vmap1.insert(v.account, v);
2814        }
2815        storage
2816            .store_all_validators(EpochNumber::new(10), vmap1.clone())
2817            .await?;
2818
2819        let mut expected_all: Vec<_> = vmap1.clone().into_values().collect();
2820        expected_all.sort_by_key(|v| v.account);
2821
2822        // Load all
2823        let loaded_all = storage
2824            .load_all_validators(EpochNumber::new(10), 0, 100)
2825            .await?;
2826        // SQLite returns a different ordered list even though there is an `ORDER BY address ASC` clause
2827        assert_eq!(expected_all, loaded_all);
2828
2829        // Load first 10
2830        let loaded_first_10 = storage
2831            .load_all_validators(EpochNumber::new(10), 0, 10)
2832            .await?;
2833
2834        assert_eq!(expected_all[..10], loaded_first_10);
2835
2836        // Load next 10
2837        let loaded_next_10 = storage
2838            .load_all_validators(EpochNumber::new(10), 10, 10)
2839            .await?;
2840
2841        assert_eq!(expected_all[10..20], loaded_next_10);
2842
2843        // Load remaining 5
2844        let loaded_last_5 = storage
2845            .load_all_validators(EpochNumber::new(10), 20, 10)
2846            .await?;
2847
2848        assert_eq!(expected_all[20..], loaded_last_5);
2849
2850        // offset beyond size should return empty
2851        let loaded_empty = storage
2852            .load_all_validators(EpochNumber::new(10), 100, 10)
2853            .await?;
2854        assert!(loaded_empty.is_empty());
2855
2856        // epoch 11
2857        let validator2 = RegisteredValidator::mock();
2858        let mut vmap2 = IndexMap::new();
2859        vmap2.insert(validator2.account, validator2.clone());
2860
2861        storage
2862            .store_all_validators(EpochNumber::new(11), vmap2.clone())
2863            .await?;
2864
2865        let mut expected_epoch11: Vec<_> = vmap2.clone().into_values().collect();
2866        expected_epoch11.sort_by_key(|v| v.account);
2867
2868        let loaded2 = storage
2869            .load_all_validators(EpochNumber::new(11), 0, 100)
2870            .await?;
2871
2872        assert_eq!(expected_epoch11, loaded2);
2873
2874        // Epoch 10 still there
2875        let loaded1_again = storage
2876            .load_all_validators(EpochNumber::new(10), 0, 100)
2877            .await?;
2878
2879        assert_eq!(expected_all, loaded1_again);
2880
2881        Ok(())
2882    }
2883
2884    #[rstest_reuse::apply(persistence_types)]
2885    pub async fn test_non_consecutive_decide<P: TestablePersistence>(_p: PhantomData<P>) {
2886        let tmp = P::tmp_storage().await;
2887        let storage = P::connect(&tmp).await;
2888
2889        let genesis_leaf: Leaf2 = Leaf2::genesis(
2890            &ValidatedState::default(),
2891            &NodeState::mock(),
2892            TEST_VERSIONS.test.base,
2893        )
2894        .await;
2895        let mut quorum_proposal = QuorumProposalWrapper::<SeqTypes> {
2896            proposal: QuorumProposal2::<SeqTypes> {
2897                epoch: None,
2898                block_header: genesis_leaf.block_header().clone(),
2899                view_number: genesis_leaf.view_number(),
2900                justify_qc: QuorumCertificate2::genesis(
2901                    &ValidatedState::default(),
2902                    &NodeState::mock(),
2903                    TEST_VERSIONS.test,
2904                )
2905                .await,
2906                upgrade_certificate: None,
2907                view_change_evidence: None,
2908                next_drb_result: None,
2909                next_epoch_justify_qc: None,
2910                state_cert: None,
2911            },
2912        };
2913
2914        let leaf0 = Leaf2::from_quorum_proposal(&quorum_proposal);
2915
2916        quorum_proposal.proposal.view_number = ViewNumber::new(2);
2917        *quorum_proposal.proposal.block_header.height_mut() = 2;
2918        quorum_proposal.proposal.justify_qc.view_number = ViewNumber::new(1);
2919        let leaf2 = Leaf2::from_quorum_proposal(&quorum_proposal);
2920
2921        let mut qc0 = leaf0.justify_qc();
2922        qc0.data.leaf_commit = Committable::commit(&leaf0);
2923
2924        let mut qc2 = leaf2.justify_qc();
2925        qc2.view_number += 1;
2926        qc2.data.leaf_commit = Committable::commit(&leaf2);
2927
2928        let mut deciding_qc = qc2.clone();
2929        deciding_qc.view_number += 1;
2930
2931        // Decide the first leaf, but fail to generate a decide event.
2932        storage
2933            .append_decided_leaves(
2934                ViewNumber::new(0),
2935                [(
2936                    &leaf_info(leaf0.clone()),
2937                    CertificatePair::non_epoch_change(qc0),
2938                )],
2939                None,
2940                &FailConsumer,
2941            )
2942            .await
2943            .unwrap();
2944
2945        // Later, decide a new leaf, but skipping some leaf in the middle. This should generate
2946        // decide events for both the leaves, correctly separating into two events since the leaves
2947        // are non-consecutive, and correctly applying `deciding_qc` only to the last event.
2948        let consumer = EventCollector::default();
2949        storage
2950            .append_decided_leaves(
2951                ViewNumber::new(2),
2952                [(
2953                    &leaf_info(leaf2.clone()),
2954                    CertificatePair::non_epoch_change(qc2),
2955                )],
2956                Some(Arc::new(CertificatePair::non_epoch_change(
2957                    deciding_qc.clone(),
2958                ))),
2959                &consumer,
2960            )
2961            .await
2962            .unwrap();
2963
2964        let events = consumer.events.read().await;
2965        assert_eq!(events.len(), 2);
2966
2967        let EventType::Decide {
2968            leaf_chain: leaf_chain0,
2969            deciding_qc: deciding_qc0,
2970            ..
2971        } = &events[0].event
2972        else {
2973            panic!("expected decide event, got {:?}", events[0].event);
2974        };
2975        assert_eq!(leaf_chain0.len(), 1);
2976        assert_eq!(leaf_chain0[0].leaf, leaf0);
2977        assert_eq!(*deciding_qc0, None);
2978
2979        let EventType::Decide {
2980            leaf_chain: leaf_chain2,
2981            deciding_qc: deciding_qc2,
2982            ..
2983        } = &events[1].event
2984        else {
2985            panic!("expected decide event, got {:?}", events[1].event);
2986        };
2987        assert_eq!(leaf_chain2.len(), 1);
2988        assert_eq!(leaf_chain2[0].leaf, leaf2);
2989        assert_eq!(deciding_qc2.as_ref().unwrap().qc(), &deciding_qc);
2990    }
2991}