Skip to main content

exex_subscription/
main.rs

1//! An ExEx example that installs a new RPC subscription endpoint that emits storage changes for a
2//! requested address.
3use alloy_primitives::{map::AddressMap, Address, U256};
4use futures::TryStreamExt;
5use jsonrpsee::{
6    core::SubscriptionResult, proc_macros::rpc, PendingSubscriptionSink, SubscriptionMessage,
7};
8use reth_ethereum::{
9    exex::{ExExContext, ExExEvent, ExExNotification},
10    node::{api::FullNodeComponents, builder::NodeHandleFor, EthereumNode},
11};
12use tokio::sync::{mpsc, oneshot};
13use tracing::{error, info};
14
15/// Subscription update format for storage changes.
16/// This is the format that will be sent to the client when a storage change occurs.
17#[derive(Debug, Clone, Copy, Default, serde::Serialize)]
18struct StorageDiff {
19    address: Address,
20    key: U256,
21    old_value: U256,
22    new_value: U256,
23}
24
25/// Subscription request format for storage changes.
26struct SubscriptionRequest {
27    /// The address to subscribe to.
28    address: Address,
29    /// The response channel to send the subscription updates to.
30    response: oneshot::Sender<mpsc::UnboundedReceiver<StorageDiff>>,
31}
32
33/// Subscription request format for storage changes.
34type SubscriptionSender = mpsc::UnboundedSender<SubscriptionRequest>;
35
36/// API to subscribe to storage changes for a specific Ethereum address.
37#[rpc(server, namespace = "watcher")]
38pub trait StorageWatcherApi {
39    /// Subscribes to storage changes for a given Ethereum address and streams `StorageDiff`
40    /// updates.
41    #[subscription(name = "subscribeStorageChanges", item = StorageDiff)]
42    fn subscribe_storage_changes(&self, address: Address) -> SubscriptionResult;
43}
44
45/// API implementation for the storage watcher.
46#[derive(Clone)]
47struct StorageWatcherRpc {
48    /// The subscription sender to send subscription requests to.
49    subscriptions: SubscriptionSender,
50}
51
52impl StorageWatcherRpc {
53    /// Creates a new [`StorageWatcherRpc`] instance with the given subscription sender.
54    fn new(subscriptions: SubscriptionSender) -> Self {
55        Self { subscriptions }
56    }
57}
58
59impl StorageWatcherApiServer for StorageWatcherRpc {
60    fn subscribe_storage_changes(
61        &self,
62        pending: PendingSubscriptionSink,
63        address: Address,
64    ) -> SubscriptionResult {
65        let subscription = self.subscriptions.clone();
66
67        tokio::spawn(async move {
68            let sink = match pending.accept().await {
69                Ok(sink) => sink,
70                Err(e) => {
71                    error!("failed to accept subscription: {e}");
72                    return;
73                }
74            };
75
76            let (resp_tx, resp_rx) = oneshot::channel();
77            subscription.send(SubscriptionRequest { address, response: resp_tx }).unwrap();
78
79            let Ok(mut rx) = resp_rx.await else { return };
80
81            loop {
82                // diffs only arrive when the address changes, so a failed send alone would not
83                // notice that the client went away
84                let diff = tokio::select! {
85                    _ = sink.closed() => break,
86                    diff = rx.recv() => diff,
87                };
88                let Some(diff) = diff else { break };
89                let msg = SubscriptionMessage::from(
90                    serde_json::value::to_raw_value(&diff).expect("serialize"),
91                );
92                if sink.send(msg).await.is_err() {
93                    break;
94                }
95            }
96        });
97
98        Ok(())
99    }
100}
101
102async fn my_exex<Node: FullNodeComponents>(
103    mut ctx: ExExContext<Node>,
104    mut subscription_requests: mpsc::UnboundedReceiver<SubscriptionRequest>,
105) -> eyre::Result<()> {
106    let mut subscriptions: AddressMap<Vec<mpsc::UnboundedSender<StorageDiff>>> =
107        AddressMap::default();
108
109    loop {
110        tokio::select! {
111            maybe_notification = ctx.notifications.try_next() => {
112                let notification = match maybe_notification? {
113                    Some(notification) => notification,
114                    None => break,
115                };
116
117                match &notification {
118                    ExExNotification::ChainCommitted { new } => {
119                        info!(committed_chain = ?new.range(), "Received commit");
120                        let execution_outcome = new.execution_outcome();
121
122                        for (address, senders) in subscriptions.iter_mut() {
123                            for change in &execution_outcome.bundle.state {
124                                if change.0 == address {
125                                    for (key, slot) in &change.1.storage {
126                                        let diff = StorageDiff {
127                                            address: *change.0,
128                                            key: *key,
129                                            old_value: slot.original_value(),
130                                            new_value: slot.present_value(),
131                                        };
132                                        // Send diff to all the active subscribers
133                                        senders.retain(|sender| sender.send(diff).is_ok());
134                                    }
135                                }
136                            }
137                        }
138                    }
139                    ExExNotification::ChainReorged { old, new } => {
140                        info!(from_chain = ?old.range(), to_chain = ?new.range(), "Received reorg");
141                    }
142                    ExExNotification::ChainReverted { old } => {
143                        info!(reverted_chain = ?old.range(), "Received revert");
144                    }
145                }
146
147                if let Some(committed_chain) = notification.committed_chain() {
148                    ctx.events.send(ExExEvent::FinishedHeight(committed_chain.tip().num_hash()))?;
149                }
150            }
151
152            maybe_subscription = subscription_requests.recv() => {
153                match maybe_subscription {
154                    Some(SubscriptionRequest { address, response }) => {
155                        let (tx, rx) = mpsc::unbounded_channel();
156                        subscriptions.entry(address).or_default().push(tx);
157                        let _ = response.send(rx);
158                    }
159                    None => {
160                        // channel closed
161                         }
162                }
163            }
164        }
165    }
166
167    Ok(())
168}
169
170fn main() -> eyre::Result<()> {
171    reth_ethereum::cli::Cli::parse_args().run(async move |builder, _| {
172        let (subscriptions_tx, subscriptions_rx) = mpsc::unbounded_channel::<SubscriptionRequest>();
173        let rpc = StorageWatcherRpc::new(subscriptions_tx);
174
175        let handle: NodeHandleFor<EthereumNode> = builder
176            .node(EthereumNode::default())
177            .extend_rpc_modules(move |ctx| {
178                ctx.modules.merge_configured(StorageWatcherApiServer::into_rpc(rpc))?;
179                Ok(())
180            })
181            .install_exex("my-exex", async move |ctx| Ok(my_exex(ctx, subscriptions_rx)))
182            .launch()
183            .await?;
184
185        handle.wait_for_node_exit().await
186    })
187}