exex_subscription/
main.rs1use 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#[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
25struct SubscriptionRequest {
27 address: Address,
29 response: oneshot::Sender<mpsc::UnboundedReceiver<StorageDiff>>,
31}
32
33type SubscriptionSender = mpsc::UnboundedSender<SubscriptionRequest>;
35
36#[rpc(server, namespace = "watcher")]
38pub trait StorageWatcherApi {
39 #[subscription(name = "subscribeStorageChanges", item = StorageDiff)]
42 fn subscribe_storage_changes(&self, address: Address) -> SubscriptionResult;
43}
44
45#[derive(Clone)]
47struct StorageWatcherRpc {
48 subscriptions: SubscriptionSender,
50}
51
52impl StorageWatcherRpc {
53 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 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 ¬ification {
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 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 }
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}