1- use crate :: providers:: provider:: { ProviderType , RpcProvider , SharedRpcProviders } ;
1+ use crate :: providers:: provider:: { RpcProvider , SharedRpcProviders } ;
22use actix_web:: rt:: time:: timeout;
33use alloy:: hex;
4- use alloy:: network:: EthereumWallet ;
54use alloy:: providers:: ProviderBuilder ;
65use alloy:: rpc:: types:: Block ;
76use alloy:: transports:: ws:: WsConnect ;
@@ -26,8 +25,8 @@ use blocksense_utils::await_time;
2625use crate :: providers:: eth_send_utils:: { try_to_sync, BatchOfUpdatesToProcess } ;
2726
2827async 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
4039async 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-
6552struct 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