Applies the extension to the supplied hooks.
(self: Box<Self>, hooks: NodeHooks)
| 29 | impl BaseNodeExtension for FlashblocksExtension { |
| 30 | /// Applies the extension to the supplied hooks. |
| 31 | fn apply(self: Box<Self>, hooks: NodeHooks) -> NodeHooks { |
| 32 | let Some(cfg) = self.config else { |
| 33 | info!(message = "flashblocks integration is disabled"); |
| 34 | return hooks; |
| 35 | }; |
| 36 | |
| 37 | let state = cfg.state; |
| 38 | let mut subscriber = FlashblocksSubscriber::new( |
| 39 | Arc::clone(&state), |
| 40 | cfg.websocket_url, |
| 41 | cfg.subscriber_ping_interval, |
| 42 | ); |
| 43 | |
| 44 | let state_for_canonical = Arc::clone(&state); |
| 45 | let state_for_rpc = Arc::clone(&state); |
| 46 | let state_for_start = state; |
| 47 | |
| 48 | // Start state processor, subscriber, and canonical subscription after node is started |
| 49 | let hooks = hooks.add_node_started_hook(move |ctx| { |
| 50 | info!(message = "Starting Flashblocks state processor"); |
| 51 | state_for_start.start(ctx.provider().clone()); |
| 52 | subscriber.start(); |
| 53 | |
| 54 | let mut canonical_stream = |
| 55 | BroadcastStream::new(ctx.provider().subscribe_to_canonical_state()); |
| 56 | tokio::spawn(async move { |
| 57 | while let Some(Ok(notification)) = canonical_stream.next().await { |
| 58 | let committed = notification.committed(); |
| 59 | for block in committed.blocks_iter() { |
| 60 | state_for_canonical.on_canonical_block_received(block.as_ref().clone()); |
| 61 | } |
| 62 | } |
| 63 | }); |
| 64 | |
| 65 | Ok(()) |
| 66 | }); |
| 67 | |
| 68 | // Extend with RPC modules |
| 69 | hooks.add_rpc_module(move |ctx| { |
| 70 | info!(message = "Starting Flashblocks RPC"); |
| 71 | |
| 72 | let api_ext = EthApiExt::new( |
| 73 | ctx.registry.eth_api().clone(), |
| 74 | ctx.registry.eth_handlers().filter.clone(), |
| 75 | Arc::clone(&state_for_rpc), |
| 76 | ); |
| 77 | ctx.modules.replace_configured(api_ext.into_rpc())?; |
| 78 | |
| 79 | // Register the eth_subscribe subscription endpoint for flashblocks |
| 80 | // Uses replace_configured since eth_subscribe already exists from reth's standard module |
| 81 | // Pass eth_api to enable proxying standard subscription types to reth's implementation |
| 82 | let eth_pubsub = EthPubSub::new( |
| 83 | ctx.registry.eth_api().clone(), |
| 84 | ctx.node().task_executor.clone(), |
| 85 | state_for_rpc, |
| 86 | ); |
| 87 | ctx.modules.replace_configured(eth_pubsub.into_rpc())?; |
| 88 |
nothing calls this directly
no test coverage detected