Skip to main content

phos_data_network/
consensus.rs

1use std::{marker::PhantomData, sync::Arc, time::Duration};
2
3use alloy::{
4    consensus::{
5        transaction::{SignerRecoverable, TransactionInfo},
6        TxEnvelope,
7    },
8    eips::{eip4895::Withdrawal, eip7685::Requests},
9    primitives::{Address, Bloom, Bytes, B256, U256},
10    rpc::types::{
11        engine::{
12            CancunPayloadFields, ExecutionPayload, ExecutionPayloadSidecar, ExecutionPayloadV1,
13            ExecutionPayloadV2, ExecutionPayloadV3, PraguePayloadFields,
14        },
15        Block, Transaction,
16    },
17};
18use cometbft::{
19    block::Block as CometBlock, crypto::default::Sha256, merkle::simple_hash_from_byte_vectors,
20};
21use cometbft_rpc::client::Client;
22use eyre::{bail, eyre, Context, Result};
23use prost::Message;
24use tokio::sync::{
25    mpsc::{channel, Receiver, Sender},
26    watch, Mutex,
27};
28use tracing::{error, info, warn};
29
30use helios_common::fork_schedule::ForkSchedule;
31use helios_core::{consensus::Consensus, time::interval};
32use phos_light_client::{
33    builder::LightClientBuilder,
34    components::{clock::SystemClock, io::IoError, scheduler::basic_bisecting_schedule},
35    errors::Error as LightClientError,
36    instance::Instance,
37    store::memory::MemoryStore,
38    verifier::{
39        predicates::ProdPredicates,
40        types::{Height, LightBlock},
41        ProdVerifier,
42    },
43};
44use phos_proto::{
45    cosmos::tx::v1beta1::{TxBody, TxRaw},
46    story::evmengine::v1::types::{ExecutionPayloadDeneb, MsgExecutionPayload},
47};
48
49use crate::{
50    config::{Config, TrustOptions},
51    database::Database,
52    rpc::http_rpc::{HttpClient, ProdIo},
53};
54
55// DATA blocks retain the legacy protobuf package name on the wire.
56const EXECUTION_PAYLOAD_TYPE_URL: &str = "/story.evmengine.v1.types.MsgExecutionPayload";
57
58#[derive(Debug, Clone)]
59enum ConsensusSyncStatus {
60    Syncing,
61    Synced,
62    Error(String),
63}
64
65pub struct ConsensusClient<DB: Database> {
66    pub block_recv: Option<Receiver<Block<Transaction>>>,
67    pub finalized_block_recv: Option<watch::Receiver<Option<Block<Transaction>>>>,
68    sync_status_recv: Mutex<watch::Receiver<ConsensusSyncStatus>>,
69    shutdown_send: watch::Sender<bool>,
70    config: Arc<Config>,
71    phantom: PhantomData<DB>,
72}
73
74#[derive(Debug)]
75pub(crate) struct Provider {
76    instance: Instance,
77    rpc: ProdIo,
78}
79
80impl Provider {
81    pub fn new(instance: Instance, rpc: ProdIo) -> Self {
82        Self { instance, rpc }
83    }
84
85    pub async fn verify_to_highest(&mut self) -> Result<LightBlock, LightClientError> {
86        self.instance
87            .light_client
88            .verify_to_highest(&mut self.instance.state)
89            .await
90    }
91
92    pub async fn fetch_block(&self, height: Height) -> Result<CometBlock, IoError> {
93        self.rpc.fetch_block(height).await
94    }
95}
96
97struct Inner<DB: Database> {
98    provider: Provider,
99    last_comet_height: Option<Height>,
100    execution_forks: ForkSchedule,
101    block_send: Sender<Block<Transaction>>,
102    finalized_block_send: watch::Sender<Option<Block<Transaction>>>,
103    db: Arc<DB>,
104}
105
106impl<DB: Database> ConsensusClient<DB> {
107    pub fn new(config: Arc<Config>) -> Result<Self> {
108        let (block_send, block_recv) = channel(256);
109        let (finalized_block_send, finalized_block_recv) = watch::channel(None);
110        let (sync_status_send, sync_status_recv) = watch::channel(ConsensusSyncStatus::Syncing);
111        let (shutdown_send, mut shutdown_recv) = watch::channel(false);
112        let config_clone = config.clone();
113        let db = Arc::new(DB::new(&config)?);
114        let trust_options = db.load_trust_options()?;
115        let consensus_rpc = config.consensus_rpc.clone();
116        let verifier_options = config.verifier_options;
117        let execution_forks = config.execution_forks;
118
119        #[cfg(not(target_arch = "wasm32"))]
120        let run = tokio::spawn;
121
122        #[cfg(target_arch = "wasm32")]
123        let run = wasm_bindgen_futures::spawn_local;
124
125        run(async move {
126            let mut inner = match Inner::new(
127                consensus_rpc,
128                trust_options,
129                verifier_options,
130                execution_forks,
131                block_send,
132                finalized_block_send,
133                db,
134            )
135            .await
136            {
137                Ok(inner) => inner,
138                Err(err) => {
139                    error!(target: "helios::data_network", err = %err, "sync failed");
140                    _ = sync_status_send.send(ConsensusSyncStatus::Error(err.to_string()));
141                    return;
142                }
143            };
144
145            if let Err(err) = inner.advance().await {
146                error!(target: "helios::data_network", err = %err, "sync failed");
147                _ = sync_status_send.send(ConsensusSyncStatus::Error(err.to_string()));
148                return;
149            }
150            _ = sync_status_send.send(ConsensusSyncStatus::Synced);
151
152            let mut interval = interval(Duration::from_secs(1));
153            loop {
154                tokio::select! {
155                    result = shutdown_recv.changed() => {
156                        if result.is_err() || *shutdown_recv.borrow_and_update() {
157                            info!(target: "helios::data_network", "shutting down consensus client");
158                            break;
159                        }
160                    }
161                    _ = interval.tick() => {
162                        if let Err(err) = inner.advance().await {
163                            warn!(target: "helios::data_network", err = %err, "advance failed");
164                        }
165                    }
166                }
167            }
168        });
169
170        Ok(Self {
171            block_recv: Some(block_recv),
172            finalized_block_recv: Some(finalized_block_recv),
173            sync_status_recv: Mutex::new(sync_status_recv),
174            shutdown_send,
175            config: config_clone,
176            phantom: PhantomData,
177        })
178    }
179}
180
181#[async_trait::async_trait]
182impl<DB: Database> Consensus<Block<Transaction>> for ConsensusClient<DB> {
183    fn block_recv(&mut self) -> Option<Receiver<Block<Transaction>>> {
184        self.block_recv.take()
185    }
186
187    fn finalized_block_recv(&mut self) -> Option<watch::Receiver<Option<Block<Transaction>>>> {
188        self.finalized_block_recv.take()
189    }
190
191    fn checkpoint_recv(&self) -> Option<watch::Receiver<Option<B256>>> {
192        None
193    }
194
195    fn expected_highest_block(&self) -> u64 {
196        u64::MAX
197    }
198
199    fn chain_id(&self) -> u64 {
200        self.config.chain.chain_id
201    }
202
203    fn shutdown(&self) -> Result<()> {
204        self.shutdown_send.send(true)?;
205        Ok(())
206    }
207
208    async fn wait_synced(&self) -> Result<()> {
209        let mut sync_status_recv = self.sync_status_recv.lock().await;
210
211        loop {
212            let status = sync_status_recv.borrow().clone();
213            match status {
214                ConsensusSyncStatus::Synced => return Ok(()),
215                ConsensusSyncStatus::Error(err) => return Err(eyre!("sync failed: {err}")),
216                ConsensusSyncStatus::Syncing => sync_status_recv.changed().await?,
217            }
218        }
219    }
220}
221
222impl<DB: Database> Inner<DB> {
223    async fn new(
224        consensus_rpc: reqwest::Url,
225        trust_options: TrustOptions,
226        verifier_options: phos_light_client::verifier::options::Options,
227        execution_forks: ForkSchedule,
228        block_send: Sender<Block<Transaction>>,
229        finalized_block_send: watch::Sender<Option<Block<Transaction>>>,
230        db: Arc<DB>,
231    ) -> Result<Self> {
232        let rpc_client = HttpClient::new(consensus_rpc);
233        let peer_id = rpc_client.status().await?.node_info.id;
234        let rpc = ProdIo::new(peer_id, rpc_client);
235        let instance = LightClientBuilder::custom(
236            peer_id,
237            verifier_options,
238            Box::new(MemoryStore::new()),
239            Box::new(rpc.clone()),
240            Box::new(SystemClock),
241            Box::new(ProdVerifier::default()),
242            Box::new(basic_bisecting_schedule),
243            Box::new(ProdPredicates),
244        )
245        .trust_primary_at(trust_options.height, trust_options.hash)
246        .await?
247        .build();
248
249        let provider = Provider::new(instance, rpc);
250
251        Ok(Self {
252            provider,
253            last_comet_height: None,
254            execution_forks,
255            block_send,
256            finalized_block_send,
257            db,
258        })
259    }
260
261    async fn advance(&mut self) -> Result<()> {
262        let light_block = self.provider.verify_to_highest().await?;
263
264        if self.last_comet_height == Some(light_block.height()) {
265            return Ok(());
266        }
267
268        let block =
269            fetch_execution_block(&self.provider, &light_block, &self.execution_forks).await?;
270        self.last_comet_height = Some(light_block.height());
271        self.db.save_trust_options(&TrustOptions {
272            height: light_block.height(),
273            hash: light_block.signed_header.header.hash_with::<Sha256>(),
274        })?;
275
276        let Some(block) = block else {
277            return Ok(());
278        };
279
280        self.block_send
281            .send(block.clone())
282            .await
283            .map_err(|_| eyre!("block receiver closed"))?;
284        self.finalized_block_send
285            .send(Some(block))
286            .map_err(|_| eyre!("finalized block receiver closed"))?;
287
288        Ok(())
289    }
290}
291
292pub(crate) async fn fetch_execution_block(
293    provider: &Provider,
294    light_block: &LightBlock,
295    execution_forks: &ForkSchedule,
296) -> Result<Option<Block<Transaction>>> {
297    let trusted_header = &light_block.signed_header.header;
298    let block = provider.fetch_block(trusted_header.height).await?;
299
300    verify_block_data(&block, light_block)?;
301
302    let Some(tx) = block.data.first() else {
303        return Ok(None);
304    };
305
306    let payload = decode_execution_payload(tx)?;
307    let app_hash = B256::try_from(trusted_header.app_hash.as_bytes())
308        .map_err(|_| eyre!("CometBFT app hash is not 32 bytes"))?;
309
310    payload_to_block(payload, app_hash, execution_forks).map(Some)
311}
312
313fn verify_block_data(block: &CometBlock, light_block: &LightBlock) -> Result<()> {
314    let trusted_header = &light_block.signed_header.header;
315
316    if block.header.height != trusted_header.height {
317        bail!(
318            "fetched block height {} does not match verified height {}",
319            block.header.height,
320            trusted_header.height
321        );
322    }
323
324    if block.header.hash_with::<Sha256>() != trusted_header.hash_with::<Sha256>() {
325        bail!("fetched block header does not match verified header");
326    }
327
328    let data_hash = block_data_hash(&block.data);
329    if data_hash.as_slice() != block.header.data_hash.unwrap_or_default().as_ref() {
330        bail!("fetched block data does not match verified data hash");
331    }
332
333    let expected_transactions = if block.header.height.value() == 1 {
334        0
335    } else {
336        1
337    };
338
339    if block.data.len() != expected_transactions {
340        bail!(
341            "DATA block at height {} must contain exactly {} execution payload transactions, found {}",
342            block.header.height,
343            expected_transactions,
344            block.data.len()
345        );
346    }
347
348    Ok(())
349}
350
351fn block_data_hash(data: &[impl AsRef<[u8]>]) -> [u8; 32] {
352    let transaction_hashes = data
353        .iter()
354        .map(|tx| <Sha256 as cometbft::crypto::Sha256>::digest(tx))
355        .collect::<Vec<_>>();
356
357    simple_hash_from_byte_vectors::<Sha256>(&transaction_hashes)
358}
359
360fn decode_execution_payload(tx: &[u8]) -> Result<ExecutionPayloadV3> {
361    let tx = TxRaw::decode(tx).wrap_err("failed to decode Cosmos SDK TxRaw")?;
362    let body =
363        TxBody::decode(tx.body_bytes.as_slice()).wrap_err("failed to decode Cosmos SDK TxBody")?;
364
365    if body.messages.len() != 1 {
366        bail!(
367            "DATA transaction must contain exactly one message, found {}",
368            body.messages.len()
369        );
370    }
371
372    let message = &body.messages[0];
373    if message.type_url != EXECUTION_PAYLOAD_TYPE_URL {
374        bail!("unexpected DATA message type: {}", message.type_url);
375    }
376
377    let payload = MsgExecutionPayload::decode(message.value.as_slice())
378        .wrap_err("failed to decode MsgExecutionPayload")?;
379
380    match (
381        payload.execution_payload.is_empty(),
382        payload.execution_payload_deneb,
383    ) {
384        (false, None) => serde_json::from_slice(&payload.execution_payload)
385            .wrap_err("failed to decode legacy JSON execution payload"),
386        (true, Some(payload)) => payload_from_proto(payload),
387        (false, Some(_)) => bail!("execution payload contains both legacy and protobuf forms"),
388        (true, None) => bail!("execution payload is missing"),
389    }
390}
391
392fn payload_to_block(
393    execution_payload: ExecutionPayloadV3,
394    app_hash: B256,
395    execution_forks: &ForkSchedule,
396) -> Result<Block<Transaction>> {
397    let payload = ExecutionPayload::V3(execution_payload);
398    let expected_block_hash = payload.block_hash();
399
400    let versioned_hashes = Vec::new();
401    let parent_beacon_block_root = app_hash;
402    let cancun = CancunPayloadFields::new(parent_beacon_block_root, versioned_hashes);
403
404    let sidecar = if payload.timestamp() >= execution_forks.prague_timestamp {
405        let requests = Requests::default();
406        let prague = PraguePayloadFields::new(requests);
407        ExecutionPayloadSidecar::v4(cancun, prague)
408    } else {
409        ExecutionPayloadSidecar::v3(cancun)
410    };
411
412    let consensus_block = payload
413        .try_into_block_with_sidecar::<TxEnvelope>(&sidecar)
414        .wrap_err("failed to construct execution block from payload")?;
415
416    let block_hash = consensus_block.header.hash_slow();
417    if block_hash != expected_block_hash {
418        bail!(
419            "execution payload block hash mismatch: expected {expected_block_hash}, got {block_hash}"
420        );
421    }
422
423    let block_number = consensus_block.header.number;
424    let base_fee = consensus_block.header.base_fee_per_gas;
425    let mut transaction_index = 0;
426
427    let rpc_block = Block::from_consensus(consensus_block, Some(U256::ZERO));
428    let rpc_block = rpc_block.try_map_transactions(|transaction| {
429        let transaction_info = TransactionInfo {
430            hash: Some(*transaction.tx_hash()),
431            index: Some(transaction_index),
432            block_hash: Some(block_hash),
433            block_number: Some(block_number),
434            base_fee,
435        };
436        transaction_index += 1;
437
438        transaction
439            .try_into_recovered()
440            .map(|transaction| Transaction::from_transaction(transaction, transaction_info))
441    })?;
442
443    Ok(rpc_block)
444}
445
446fn payload_from_proto(payload: ExecutionPayloadDeneb) -> Result<ExecutionPayloadV3> {
447    let withdrawals = payload
448        .withdrawals
449        .into_iter()
450        .map(|withdrawal| {
451            Ok(Withdrawal {
452                index: withdrawal.index,
453                validator_index: withdrawal.validator_index,
454                address: Address::from(fixed_bytes("withdrawal address", withdrawal.address)?),
455                amount: withdrawal.amount,
456            })
457        })
458        .collect::<Result<Vec<_>>>()?;
459
460    Ok(ExecutionPayloadV3 {
461        payload_inner: ExecutionPayloadV2 {
462            payload_inner: ExecutionPayloadV1 {
463                parent_hash: B256::from(fixed_bytes("parent hash", payload.parent_hash)?),
464                fee_recipient: Address::from(fixed_bytes("fee recipient", payload.fee_recipient)?),
465                state_root: B256::from(fixed_bytes("state root", payload.state_root)?),
466                receipts_root: B256::from(fixed_bytes("receipts root", payload.receipts_root)?),
467                logs_bloom: Bloom::from(fixed_bytes("logs bloom", payload.logs_bloom)?),
468                prev_randao: B256::from(fixed_bytes("prev randao", payload.prev_randao)?),
469                block_number: payload.block_number,
470                gas_limit: payload.gas_limit,
471                gas_used: payload.gas_used,
472                timestamp: payload.timestamp,
473                extra_data: Bytes::from(payload.extra_data),
474                base_fee_per_gas: U256::from_be_bytes(fixed_bytes::<32>(
475                    "base fee per gas",
476                    payload.base_fee_per_gas,
477                )?),
478                block_hash: B256::from(fixed_bytes("block hash", payload.block_hash)?),
479                transactions: payload.transactions.into_iter().map(Bytes::from).collect(),
480            },
481            withdrawals,
482        },
483        blob_gas_used: payload.blob_gas_used,
484        excess_blob_gas: payload.excess_blob_gas,
485    })
486}
487
488fn fixed_bytes<const N: usize>(field: &str, bytes: Vec<u8>) -> Result<[u8; N]> {
489    bytes
490        .try_into()
491        .map_err(|bytes: Vec<u8>| eyre!("{field} must be {N} bytes, found {}", bytes.len()))
492}