1use alloy_consensus::{BlockHeader, Transaction};
4use alloy_primitives::Bytes;
5use alloy_rpc_types_engine::{ForkchoiceState, PayloadStatus};
6use futures::{stream::FuturesUnordered, Stream, StreamExt, TryFutureExt};
7use itertools::Either;
8use reth_chainspec::{ChainSpecProvider, EthChainSpec};
9use reth_engine_primitives::{
10 BeaconEngineMessage, BeaconOnNewPayloadError, ExecutionPayload as _, OnForkChoiceUpdated,
11};
12use reth_engine_tree::tree::EngineValidator;
13use reth_errors::{BlockExecutionError, BlockValidationError, RethError, RethResult};
14use reth_evm::{
15 execute::{BlockBuilder, BlockBuilderOutcome},
16 ConfigureEvm,
17};
18use reth_payload_primitives::{BuiltPayload, PayloadTypes};
19use reth_primitives_traits::{
20 block::Block as _, BlockBody as _, BlockTy, HeaderTy, SealedBlock, SignedTransaction,
21};
22use reth_revm::{database::StateProviderDatabase, db::State};
23use reth_storage_api::{errors::ProviderError, BlockReader, StateProvider, StateProviderFactory};
24use std::{
25 collections::VecDeque,
26 future::Future,
27 pin::Pin,
28 task::{ready, Context, Poll},
29};
30use tokio::sync::oneshot;
31use tracing::*;
32
33#[derive(Debug)]
34enum EngineReorgState<T: PayloadTypes> {
35 Forward,
36 Reorg { queue: VecDeque<BeaconEngineMessage<T>> },
37}
38
39type EngineReorgResponse = Result<
40 Either<Result<PayloadStatus, BeaconOnNewPayloadError>, RethResult<OnForkChoiceUpdated>>,
41 oneshot::error::RecvError,
42>;
43
44type ReorgResponseFut = Pin<Box<dyn Future<Output = EngineReorgResponse> + Send + Sync>>;
45
46#[derive(Debug)]
48#[pin_project::pin_project]
49pub struct EngineReorg<S, T: PayloadTypes, Provider, Evm, Validator> {
50 #[pin]
52 stream: S,
53 provider: Provider,
55 evm_config: Evm,
57 payload_validator: Validator,
59 frequency: usize,
61 depth: usize,
63 forkchoice_states_forwarded: usize,
66 state: EngineReorgState<T>,
68 last_forkchoice_state: Option<ForkchoiceState>,
70 reorg_responses: FuturesUnordered<ReorgResponseFut>,
72}
73
74impl<S, T: PayloadTypes, Provider, Evm, Validator> EngineReorg<S, T, Provider, Evm, Validator> {
75 pub fn new(
77 stream: S,
78 provider: Provider,
79 evm_config: Evm,
80 payload_validator: Validator,
81 frequency: usize,
82 depth: usize,
83 ) -> Self {
84 Self {
85 stream,
86 provider,
87 evm_config,
88 payload_validator,
89 frequency,
90 depth,
91 state: EngineReorgState::Forward,
92 forkchoice_states_forwarded: 0,
93 last_forkchoice_state: None,
94 reorg_responses: FuturesUnordered::new(),
95 }
96 }
97}
98
99impl<S, T, Provider, Evm, Validator> Stream for EngineReorg<S, T, Provider, Evm, Validator>
100where
101 S: Stream<Item = BeaconEngineMessage<T>>,
102 T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = Evm::Primitives>>,
103 Provider: BlockReader<Header = HeaderTy<Evm::Primitives>, Block = BlockTy<Evm::Primitives>>
104 + StateProviderFactory
105 + ChainSpecProvider,
106 Evm: ConfigureEvm,
107 Validator: EngineValidator<T, Evm::Primitives>,
108{
109 type Item = S::Item;
110
111 fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
112 let mut this = self.project();
113
114 loop {
115 if let Poll::Ready(Some(response)) = this.reorg_responses.poll_next_unpin(cx) {
116 match response {
117 Ok(Either::Left(Ok(payload_status))) => {
118 debug!(target: "engine::stream::reorg", ?payload_status, "Received response for reorg new payload");
119 }
120 Ok(Either::Left(Err(payload_error))) => {
121 error!(target: "engine::stream::reorg", %payload_error, "Error on reorg new payload");
122 }
123 Ok(Either::Right(Ok(fcu_status))) => {
124 debug!(target: "engine::stream::reorg", ?fcu_status, "Received response for reorg forkchoice update");
125 }
126 Ok(Either::Right(Err(fcu_error))) => {
127 error!(target: "engine::stream::reorg", %fcu_error, "Error on reorg forkchoice update");
128 }
129 Err(_) => {}
130 };
131 continue
132 }
133
134 if let EngineReorgState::Reorg { queue } = &mut this.state {
135 match queue.pop_front() {
136 Some(msg) => return Poll::Ready(Some(msg)),
137 None => {
138 *this.forkchoice_states_forwarded = 0;
139 *this.state = EngineReorgState::Forward;
140 }
141 }
142 }
143
144 let next = ready!(this.stream.poll_next_unpin(cx));
145 let item = match (next, &this.last_forkchoice_state) {
146 (
147 Some(BeaconEngineMessage::NewPayload { cause, payload, tx }),
148 Some(last_forkchoice_state),
149 ) if this.forkchoice_states_forwarded > this.frequency &&
150 last_forkchoice_state.head_block_hash == payload.parent_hash() =>
152 {
153 let (reorg_block, encoded_bal) = match create_reorg_head(
163 this.provider,
164 this.evm_config,
165 this.payload_validator,
166 *this.depth,
167 payload.clone(),
168 ) {
169 Ok(result) => result,
170 Err(error) => {
171 error!(target: "engine::stream::reorg", %error, "Error attempting to create reorg head");
172 return Poll::Ready(Some(BeaconEngineMessage::NewPayload {
175 cause,
176 payload,
177 tx,
178 }))
179 }
180 };
181 let reorg_forkchoice_state = ForkchoiceState {
182 finalized_block_hash: last_forkchoice_state.finalized_block_hash,
183 safe_block_hash: last_forkchoice_state.safe_block_hash,
184 head_block_hash: reorg_block.hash(),
185 };
186
187 let (reorg_payload_tx, reorg_payload_rx) = oneshot::channel();
188 let (reorg_fcu_tx, reorg_fcu_rx) = oneshot::channel();
189 this.reorg_responses.extend([
190 Box::pin(reorg_payload_rx.map_ok(Either::Left)) as ReorgResponseFut,
191 Box::pin(reorg_fcu_rx.map_ok(Either::Right)) as ReorgResponseFut,
192 ]);
193
194 let queue = VecDeque::from([
195 BeaconEngineMessage::NewPayload { cause: cause.clone(), payload, tx },
197 BeaconEngineMessage::NewPayload {
199 cause: cause.clone(),
200 payload: T::block_to_payload(reorg_block, encoded_bal),
201 tx: reorg_payload_tx,
202 },
203 BeaconEngineMessage::ForkchoiceUpdated {
205 cause,
206 state: reorg_forkchoice_state,
207 payload_attrs: None,
208 tx: reorg_fcu_tx,
209 },
210 ]);
211 *this.state = EngineReorgState::Reorg { queue };
212 continue
213 }
214 (
215 Some(BeaconEngineMessage::ForkchoiceUpdated {
216 cause,
217 state,
218 payload_attrs,
219 tx,
220 }),
221 _,
222 ) => {
223 *this.last_forkchoice_state = Some(state);
227 *this.forkchoice_states_forwarded += 1;
228 Some(BeaconEngineMessage::ForkchoiceUpdated { cause, state, payload_attrs, tx })
229 }
230 (item, _) => item,
231 };
232 return Poll::Ready(item)
233 }
234 }
235}
236
237#[allow(clippy::type_complexity)]
238fn create_reorg_head<Provider, Evm, T, Validator>(
239 provider: &Provider,
240 evm_config: &Evm,
241 payload_validator: &Validator,
242 mut depth: usize,
243 next_payload: T::ExecutionData,
244) -> RethResult<(SealedBlock<BlockTy<Evm::Primitives>>, Option<Bytes>)>
245where
246 Provider: BlockReader<Header = HeaderTy<Evm::Primitives>, Block = BlockTy<Evm::Primitives>>
247 + StateProviderFactory
248 + ChainSpecProvider<ChainSpec: EthChainSpec>,
249 Evm: ConfigureEvm,
250 T: PayloadTypes<BuiltPayload: BuiltPayload<Primitives = Evm::Primitives>>,
251 Validator: EngineValidator<T, Evm::Primitives>,
252{
253 let next_block =
255 payload_validator.convert_payload_to_block(next_payload).map_err(RethError::msg)?;
256
257 let mut previous_hash = next_block.parent_hash();
259 let mut candidate_transactions = next_block.into_body().transactions().to_vec();
260 let reorg_target = 'target: {
261 loop {
262 let reorg_target = provider
263 .block_by_hash(previous_hash)?
264 .ok_or_else(|| ProviderError::HeaderNotFound(previous_hash.into()))?;
265 if depth == 0 {
266 break 'target reorg_target.seal_slow()
267 }
268
269 depth -= 1;
270 previous_hash = reorg_target.header().parent_hash();
271 candidate_transactions = reorg_target.into_body().into_transactions();
272 }
273 };
274 let reorg_target_parent = provider
275 .sealed_header_by_hash(reorg_target.header().parent_hash())?
276 .ok_or_else(|| ProviderError::HeaderNotFound(reorg_target.header().parent_hash().into()))?;
277
278 debug!(target: "engine::stream::reorg", number = reorg_target.header().number(), hash = %previous_hash, "Selected reorg target");
279
280 let has_bal = reorg_target.header().block_access_list_hash().is_some();
282 let state_provider = provider.state_by_block_hash(reorg_target.header().parent_hash())?;
283 let mut state = State::builder()
284 .with_database_ref(StateProviderDatabase::new((&state_provider).into_evm_state_provider()))
285 .with_bundle_update()
286 .with_bal_builder_if(has_bal)
287 .build();
288
289 let ctx = evm_config.context_for_block(&reorg_target).map_err(RethError::other)?;
290 let evm = evm_config.evm_for_block(&mut state, &reorg_target).map_err(RethError::other)?;
291 let mut builder = evm_config.create_block_builder(evm, &reorg_target_parent, ctx);
292
293 builder.apply_pre_execution_changes()?;
294
295 let mut cumulative_gas_used = 0;
296 for tx in candidate_transactions {
297 if cumulative_gas_used + tx.gas_limit() > reorg_target.gas_limit() {
299 continue
300 }
301
302 let tx_recovered =
303 tx.try_into_recovered().map_err(|_| ProviderError::SenderRecoveryError)?;
304 let gas_used = match builder.execute_transaction(tx_recovered) {
305 Ok(gas_used) => gas_used.tx_gas_used(),
306 Err(BlockExecutionError::Validation(BlockValidationError::InvalidTx {
307 hash,
308 error,
309 })) => {
310 trace!(target: "engine::stream::reorg", hash = %hash, ?error, "Error executing transaction from next block");
311 continue
312 }
313 Err(error) => return Err(RethError::Execution(error)),
315 };
316
317 cumulative_gas_used += gas_used;
318 }
319
320 let BlockBuilderOutcome { block, block_access_list, .. } =
321 builder.finish(&state_provider, None)?;
322
323 Ok((block.into_sealed_block(), block_access_list.map(|bal| bal.split().1)))
324}