Skip to content

Commit 718693e

Browse files
HristoStaykovreo101
authored andcommitted
refactor(sequencer/reorg_tracking): Remove wallet dependency in reorg_tracking
Reorg RPC helpers to accept &dyn Provider so both HTTP FillProvider and WS RootProvider handles can satisfy the minimal get_block_by_number/get_storage_at interface.
1 parent e426011 commit 718693e

1 file changed

Lines changed: 72 additions & 99 deletions

File tree

apps/sequencer/src/providers/reorg_tracking.rs

Lines changed: 72 additions & 99 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
1-
use crate::providers::provider::{ProviderType, RpcProvider, SharedRpcProviders};
1+
use crate::providers::provider::{RpcProvider, SharedRpcProviders};
22
use actix_web::rt::time::timeout;
33
use alloy::hex;
4-
use alloy::network::EthereumWallet;
54
use alloy::providers::ProviderBuilder;
65
use alloy::rpc::types::Block;
76
use alloy::transports::ws::WsConnect;
@@ -26,8 +25,8 @@ use blocksense_utils::await_time;
2625
use crate::providers::eth_send_utils::{try_to_sync, BatchOfUpdatesToProcess};
2726

2827
async fn rpc_get_block_by_number(
29-
rpc_http: &ProviderType,
30-
rpc_ws: Option<&ProviderType>,
28+
rpc_http: &dyn Provider,
29+
rpc_ws: Option<&dyn Provider>,
3130
block_number: BlockNumberOrTag,
3231
) -> eyre::Result<Option<Block>> {
3332
let provider = rpc_ws.unwrap_or(rpc_http);
@@ -38,8 +37,8 @@ async fn rpc_get_block_by_number(
3837
}
3938

4039
async fn rpc_get_storage_at(
41-
rpc_http: &ProviderType,
42-
rpc_ws: Option<&ProviderType>,
40+
rpc_http: &dyn Provider,
41+
rpc_ws: Option<&dyn Provider>,
4342
address: Address,
4443
slot: alloy_primitives::U256,
4544
) -> eyre::Result<alloy_primitives::U256> {
@@ -50,18 +49,6 @@ async fn rpc_get_storage_at(
5049
.map_err(Report::from)
5150
}
5251

53-
async fn ws_wallet_for_network(
54-
providers_mutex: &SharedRpcProviders,
55-
net: &str,
56-
) -> Option<EthereumWallet> {
57-
let provider_arc = {
58-
let providers = providers_mutex.read().await;
59-
providers.get(net).cloned()
60-
}?;
61-
let provider = provider_arc.lock().await;
62-
Some(EthereumWallet::from(provider.signer.clone()))
63-
}
64-
6552
struct ReconnectBackoff {
6653
backoff_idx: usize,
6754
next_retry_at: Option<Instant>,
@@ -141,8 +128,8 @@ impl ReorgTracker {
141128
// discarded observations, mirroring the existing log messages and structure.
142129
async fn handle_reorg(
143130
&mut self,
144-
rpc_handle: &ProviderType,
145-
rpc_ws: Option<&ProviderType>,
131+
rpc_handle: &dyn Provider,
132+
rpc_ws: Option<&dyn Provider>,
146133
provider_mutex: &Arc<Mutex<RpcProvider>>,
147134
observed_block_hashes: &HashMap<u64, B256>,
148135
observed_latest_height: u64,
@@ -287,7 +274,7 @@ impl ReorgTracker {
287274

288275
let websocket_url = self.websocket_url.clone();
289276

290-
let mut _provider_ws_opt = None;
277+
let mut provider_ws_opt = None;
291278
let mut stream_opt = None;
292279
let mut average_block_generation_time: u64 = 0;
293280

@@ -296,43 +283,38 @@ impl ReorgTracker {
296283
let mut reconnect_backoff_tracker = ReconnectBackoff::new();
297284

298285
if let Some(ref ws_url) = websocket_url {
299-
if let Some(ws_wallet) = ws_wallet_for_network(&providers_mutex, net.as_str()).await {
300-
info!("Attempting WS connect for {net} to {ws_url}");
301-
match ProviderBuilder::new()
302-
.disable_recommended_fillers()
303-
.wallet(ws_wallet)
304-
.connect_ws(WsConnect::new(ws_url))
305-
.await
306-
{
307-
Ok(provider_ws) => {
308-
info!("WS connected for {net}; subscribing to newHeads");
309-
match provider_ws.subscribe_blocks().await {
310-
Ok(sub) => {
311-
stream_opt = Some(sub.into_stream());
312-
_provider_ws_opt = Some(provider_ws);
313-
reconnect_backoff_tracker.connection_established();
314-
}
315-
Err(e) => {
316-
warn!("WS subscribe_blocks failed for {net}: {e:?}");
317-
reconnect_backoff_tracker.inc_retries_count();
318-
warn!(
319-
"Will retry WS subscribe for {net} in {jittered_ms:?}ms after error: {e:?}",
320-
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
321-
);
322-
}
286+
info!("Attempting WS connect for {net} to {ws_url}");
287+
match ProviderBuilder::new()
288+
.disable_recommended_fillers()
289+
.connect_ws(WsConnect::new(ws_url))
290+
.await
291+
{
292+
Ok(provider_ws) => {
293+
info!("WS connected for {net}; subscribing to newHeads");
294+
match provider_ws.subscribe_blocks().await {
295+
Ok(sub) => {
296+
stream_opt = Some(sub.into_stream());
297+
provider_ws_opt = Some(provider_ws);
298+
reconnect_backoff_tracker.connection_established();
299+
}
300+
Err(e) => {
301+
warn!("WS subscribe_blocks failed for {net}: {e:?}");
302+
reconnect_backoff_tracker.inc_retries_count();
303+
warn!(
304+
"Will retry WS subscribe for {net} in {jittered_ms:?}ms after error: {e:?}",
305+
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
306+
);
323307
}
324-
}
325-
Err(e) => {
326-
warn!("WS connect failed for {net}: {e:?}");
327-
reconnect_backoff_tracker.inc_retries_count();
328-
warn!(
329-
"Will retry WS connect for {net} in {jittered_ms:?}ms after error: {e:?}",
330-
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
331-
);
332308
}
333309
}
334-
} else {
335-
warn!("No signer available for {net}; skipping WS connect");
310+
Err(e) => {
311+
warn!("WS connect failed for {net}: {e:?}");
312+
reconnect_backoff_tracker.inc_retries_count();
313+
warn!(
314+
"Will retry WS connect for {net} in {jittered_ms:?}ms after error: {e:?}",
315+
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
316+
);
317+
}
336318
}
337319
} else {
338320
// Loop until block generation time is determined
@@ -364,49 +346,38 @@ impl ReorgTracker {
364346
};
365347
if should_attempt {
366348
info!("Attempting WS reconnect for {net} to {ws_url}");
367-
if let Some(ws_wallet) =
368-
ws_wallet_for_network(&providers_mutex, net.as_str()).await
349+
350+
match ProviderBuilder::new()
351+
.disable_recommended_fillers()
352+
.connect_ws(WsConnect::new(ws_url))
353+
.await
369354
{
370-
match ProviderBuilder::new()
371-
.disable_recommended_fillers()
372-
.wallet(ws_wallet)
373-
.connect_ws(WsConnect::new(ws_url))
374-
.await
375-
{
376-
Ok(provider_ws) => {
377-
info!("WS reconnected for {net}; subscribing to newHeads");
378-
match provider_ws.subscribe_blocks().await {
379-
Ok(sub) => {
380-
stream_opt = Some(sub.into_stream());
381-
_provider_ws_opt = Some(provider_ws);
382-
reconnect_backoff_tracker.connection_established();
383-
}
384-
Err(e) => {
385-
warn!("WS subscribe_blocks failed for {net}: {e:?}");
386-
reconnect_backoff_tracker.inc_retries_count();
387-
warn!(
388-
"Will retry WS subscribe for {net} in {jittered_ms:?}ms after error: {e:?}",
389-
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
390-
);
391-
}
355+
Ok(provider_ws) => {
356+
info!("WS reconnected for {net}; subscribing to newHeads");
357+
match provider_ws.subscribe_blocks().await {
358+
Ok(sub) => {
359+
stream_opt = Some(sub.into_stream());
360+
provider_ws_opt = Some(provider_ws);
361+
reconnect_backoff_tracker.connection_established();
362+
}
363+
Err(e) => {
364+
warn!("WS subscribe_blocks failed for {net}: {e:?}");
365+
reconnect_backoff_tracker.inc_retries_count();
366+
warn!(
367+
"Will retry WS subscribe for {net} in {jittered_ms:?}ms after error: {e:?}",
368+
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
369+
);
392370
}
393-
}
394-
Err(e) => {
395-
warn!("WS reconnect failed for {net}: {e:?}");
396-
reconnect_backoff_tracker.inc_retries_count();
397-
warn!(
398-
"Will retry WS connect for {net} in {jittered_ms:?}ms after error: {e:?}",
399-
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
400-
);
401371
}
402372
}
403-
} else {
404-
warn!("No signer available for {net}; skipping WS reconnect attempt");
405-
reconnect_backoff_tracker.inc_retries_count();
406-
warn!(
407-
"Will retry WS connect for {net} in {jittered_ms:?}ms after error: signer unavailable",
408-
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
409-
);
373+
Err(e) => {
374+
warn!("WS reconnect failed for {net}: {e:?}");
375+
reconnect_backoff_tracker.inc_retries_count();
376+
warn!(
377+
"Will retry WS connect for {net} in {jittered_ms:?}ms after error: {e:?}",
378+
jittered_ms = reconnect_backoff_tracker.get_next_retry_ms()
379+
);
380+
}
410381
}
411382
}
412383
}
@@ -430,7 +401,7 @@ impl ReorgTracker {
430401
loop_count = self.loop_count,
431402
);
432403
stream_opt = None;
433-
_provider_ws_opt = None;
404+
provider_ws_opt = None;
434405
reconnect_backoff_tracker.inc_retries_count();
435406
warn!(
436407
"Falling back to polling for {jittered_ms:?}ms before retrying WS in network {net}",
@@ -509,7 +480,9 @@ impl ReorgTracker {
509480
)
510481
};
511482

512-
let ws_provider = _provider_ws_opt.as_ref();
483+
let ws_provider = provider_ws_opt
484+
.as_ref()
485+
.map(|provider| provider as &dyn Provider);
513486

514487
let latest_block_result = timeout(
515488
self.rpc_timeout,
@@ -717,8 +690,8 @@ impl ReorgTracker {
717690

718691
async fn process_new_block(
719692
&mut self,
720-
rpc_handle: &ProviderType,
721-
rpc_ws: Option<&ProviderType>,
693+
rpc_handle: &dyn Provider,
694+
rpc_ws: Option<&dyn Provider>,
722695
provider_mutex: &Arc<Mutex<RpcProvider>>,
723696
observed_block_hashes: &HashMap<u64, B256>,
724697
provider_metrics: Arc<tokio::sync::RwLock<ProviderMetrics>>,
@@ -1011,7 +984,7 @@ mod tests {
1011984
use tokio::sync::RwLock;
1012985

1013986
use crate::providers::eth_send_utils::BatchOfUpdatesToProcess;
1014-
use crate::providers::provider::init_shared_rpc_providers;
987+
use crate::providers::provider::{init_shared_rpc_providers, ProviderType};
1015988

1016989
async fn mine_self_txs(
1017990
rpc: &ProviderType,

0 commit comments

Comments
 (0)