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