Skip to main content

evm_oracle_state/pending/
mev_share.rs

1//! Flashbots MEV-Share SSE source for transaction and bundle hints.
2
3use 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/// Flashbots MEV-Share event-stream source.
18#[derive(Clone)]
19pub struct MevSharePendingTransactionSource {
20    stream_url: String,
21    source_id: PendingOracleSourceId,
22    chain_id: u64,
23}
24
25impl MevSharePendingTransactionSource {
26    /// Construct a source for an MEV-Share SSE endpoint.
27    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    /// Construct a source for the public Ethereum MEV-Share event stream.
36    pub fn ethereum_mainnet() -> Self {
37        Self::new("https://mev-share.flashbots.net")
38    }
39
40    /// Override the stable source identifier used in health and provenance.
41    pub fn source_id(mut self, source_id: PendingOracleSourceId) -> Self {
42        self.source_id = source_id;
43        self
44    }
45
46    /// Set the declared chain id. It must match the pending runtime configuration.
47    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}