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
55const 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}