evm_oracle_state/pending/
mev_share.rs1use std::{fmt, str::FromStr, time::Duration};
4
5use alloy_primitives::{Address, B256, Bytes, keccak256};
6use futures_util::StreamExt;
7use serde_json::Value;
8use tokio::sync::watch;
9
10use super::{
11 ETHEREUM_MAINNET_CHAIN_ID, PendingOracleCandidateSource, PendingOracleOrderingHandle,
12 PendingOracleSource, PendingOracleSourceDescriptor, PendingOracleSourceError,
13 PendingOracleSourceFuture, PendingOracleSourceId, PendingOracleSourceSink,
14 PendingOracleTransmissionId, PendingTransportCandidate,
15};
16
17#[derive(Clone)]
19pub struct MevSharePendingTransactionSource {
20 stream_url: String,
21 source_id: PendingOracleSourceId,
22 chain_id: u64,
23}
24
25impl MevSharePendingTransactionSource {
26 pub fn new(stream_url: impl Into<String>) -> Self {
28 Self {
29 stream_url: stream_url.into(),
30 source_id: PendingOracleSourceId::new("flashbots-mev-share"),
31 chain_id: ETHEREUM_MAINNET_CHAIN_ID,
32 }
33 }
34
35 pub fn ethereum_mainnet() -> Self {
37 Self::new("https://mev-share.flashbots.net")
38 }
39
40 pub fn source_id(mut self, source_id: PendingOracleSourceId) -> Self {
42 self.source_id = source_id;
43 self
44 }
45
46 pub fn chain_id(mut self, chain_id: u64) -> Self {
48 self.chain_id = chain_id;
49 self
50 }
51
52 async fn run_forever(
53 self,
54 sink: PendingOracleSourceSink,
55 mut shutdown: watch::Receiver<bool>,
56 ) -> Result<(), PendingOracleSourceError> {
57 let client = reqwest::Client::builder()
58 .user_agent("evm-oracle-state/pending-oracle-mev-share")
59 .build()
60 .map_err(transport_error)?;
61 let mut retry = Duration::from_secs(1);
62 loop {
63 if *shutdown.borrow() {
64 return Ok(());
65 }
66 match self.run_connection(&client, &sink, &mut shutdown).await {
67 Ok(()) if *shutdown.borrow() => return Ok(()),
68 Ok(()) => sink.coverage_gap("MEV-Share event stream ended"),
69 Err(error) => sink.coverage_gap(error.to_string()),
70 }
71 sink.reconnecting();
72 tokio::select! {
73 changed = shutdown.changed() => {
74 if changed.is_err() || *shutdown.borrow() {
75 return Ok(());
76 }
77 }
78 () = tokio::time::sleep(retry) => {}
79 }
80 retry = (retry * 2).min(Duration::from_secs(30));
81 }
82 }
83
84 async fn run_connection(
85 &self,
86 client: &reqwest::Client,
87 sink: &PendingOracleSourceSink,
88 shutdown: &mut watch::Receiver<bool>,
89 ) -> Result<(), PendingOracleSourceError> {
90 let response = client
91 .get(&self.stream_url)
92 .header(reqwest::header::ACCEPT, "text/event-stream")
93 .send()
94 .await
95 .map_err(transport_error)?
96 .error_for_status()
97 .map_err(transport_error)?;
98 sink.ready();
99 let mut stream = response.bytes_stream();
100 let mut buffer = String::new();
101 loop {
102 tokio::select! {
103 changed = shutdown.changed() => {
104 if changed.is_err() || *shutdown.borrow() {
105 return Ok(());
106 }
107 }
108 chunk = stream.next() => {
109 let Some(chunk) = chunk else {
110 return Err(PendingOracleSourceError::Transport(
111 "MEV-Share response body closed".to_string(),
112 ));
113 };
114 let chunk = chunk.map_err(transport_error)?;
115 sink.transport_message();
116 buffer.push_str(&String::from_utf8_lossy(&chunk).replace("\r\n", "\n"));
117 while let Some(boundary) = buffer.find("\n\n") {
118 let remainder = buffer.split_off(boundary + 2);
119 let frame = std::mem::replace(&mut buffer, remainder);
120 let frame = &frame[..boundary];
121 let Some(event) = parse_sse_frame(frame) else {
122 continue;
123 };
124 for candidate in self.event_candidates(&event) {
125 if sink
126 .runtime()
127 .interests()
128 .iter()
129 .any(|interest| interest.matches(&candidate))
130 {
131 sink.candidate();
132 match sink.runtime().observe_candidate(candidate) {
133 Ok(report) => {
134 for failure in report.failures {
135 tracing::debug!(
136 adapter_id = %failure.adapter_id,
137 error = %failure.error,
138 "pending oracle adapter rejected MEV-Share candidate"
139 );
140 }
141 }
142 Err(error) => {
143 tracing::debug!(%error, "pending oracle runtime rejected MEV-Share candidate");
144 }
145 }
146 }
147 }
148 }
149 }
150 }
151 }
152 }
153
154 fn event_candidates(&self, event: &Value) -> Vec<PendingTransportCandidate> {
155 let Some(hash) = event
156 .get("hash")
157 .and_then(Value::as_str)
158 .and_then(|value| B256::from_str(value).ok())
159 else {
160 return Vec::new();
161 };
162 event
163 .get("txs")
164 .and_then(Value::as_array)
165 .into_iter()
166 .flatten()
167 .enumerate()
168 .filter_map(|(index, transaction)| {
169 let to = transaction
170 .get("to")
171 .and_then(Value::as_str)
172 .and_then(|value| Address::from_str(value).ok())?;
173 let calldata = transaction
174 .get("callData")
175 .and_then(Value::as_str)
176 .and_then(|value| Bytes::from_str(value).ok())?;
177 if calldata.len() < 4 {
178 return None;
179 }
180 let mut identity = Vec::with_capacity(32 + 8 + 20 + calldata.len());
181 identity.extend_from_slice(hash.as_slice());
182 identity.extend_from_slice(&(index as u64).to_be_bytes());
183 identity.extend_from_slice(to.as_slice());
184 identity.extend_from_slice(&calldata);
185 Some(PendingTransportCandidate::new(
186 self.chain_id,
187 PendingOracleTransmissionId::from_hash(keccak256(identity)),
188 PendingOracleSource::MevShare,
189 self.source_id.clone(),
190 to,
191 calldata,
192 PendingOracleOrderingHandle::MevShare { hash },
193 None,
194 ))
195 })
196 .collect()
197 }
198}
199
200impl fmt::Debug for MevSharePendingTransactionSource {
201 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
202 formatter
203 .debug_struct("MevSharePendingTransactionSource")
204 .field("stream_url", &"<redacted>")
205 .field("source_id", &self.source_id)
206 .field("chain_id", &self.chain_id)
207 .finish()
208 }
209}
210
211impl PendingOracleCandidateSource for MevSharePendingTransactionSource {
212 fn descriptor(&self) -> PendingOracleSourceDescriptor {
213 PendingOracleSourceDescriptor::new(PendingOracleSource::MevShare, self.source_id.clone())
214 }
215
216 fn run(
217 self: Box<Self>,
218 sink: PendingOracleSourceSink,
219 shutdown: watch::Receiver<bool>,
220 ) -> PendingOracleSourceFuture {
221 Box::pin(async move { self.run_forever(sink, shutdown).await })
222 }
223}
224
225fn parse_sse_frame(frame: &str) -> Option<Value> {
226 let payload = frame
227 .lines()
228 .filter_map(|line| line.strip_prefix("data:"))
229 .map(str::trim_start)
230 .collect::<Vec<_>>()
231 .join("\n");
232 (!payload.is_empty())
233 .then(|| serde_json::from_str(&payload).ok())
234 .flatten()
235}
236
237fn transport_error(error: impl fmt::Display) -> PendingOracleSourceError {
238 PendingOracleSourceError::Transport(error.to_string())
239}
240
241#[cfg(test)]
242mod tests {
243 use super::parse_sse_frame;
244
245 #[test]
246 fn parses_multiline_sse_data() {
247 let parsed = parse_sse_frame("event: transaction\ndata: {\"hash\":\ndata: \"0x01\"}")
248 .expect("valid SSE payload");
249 assert_eq!(parsed["hash"], "0x01");
250 }
251}