1use alloy_primitives::Address;
2use alloy_sol_types::sol;
3use chio_egress_contract::HttpEgressContract;
4use reqwest::Url;
5
6use crate::chainlink::contract_backed_provider;
7use crate::config::ChainlinkNetworkConfig;
8use crate::PriceOracleError;
9
10sol! {
11 #[sol(rpc)]
12 contract AggregatorV3Interface {
13 function latestRoundData() external view returns (
14 uint80 roundId,
15 int256 answer,
16 uint256 startedAt,
17 uint256 updatedAt,
18 uint80 answeredInRound
19 );
20 }
21}
22
23#[derive(Debug, Clone, Copy, PartialEq, Eq)]
24pub enum SequencerAvailability {
25 Up,
26 Down,
27 Recovering,
28}
29
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub struct SequencerStatus {
32 pub chain_id: u64,
33 pub feed_address: String,
34 pub checked_at: u64,
35 pub status_started_at: u64,
36 pub availability: SequencerAvailability,
37}
38
39pub async fn read_sequencer_status(
45 chain: &ChainlinkNetworkConfig,
46 now: u64,
47 egress_contract: &HttpEgressContract,
48) -> Result<Option<SequencerStatus>, PriceOracleError> {
49 let Some(feed_address) = chain.sequencer_uptime_feed.as_ref() else {
50 return Ok(None);
51 };
52 egress_contract
56 .enforce_url(&chain.rpc_endpoint, 0)
57 .map_err(|err| {
58 PriceOracleError::InvalidConfiguration(format!(
59 "HttpEgressContract rejects sequencer RPC endpoint: {err}"
60 ))
61 })?;
62 let url = chain.rpc_endpoint.parse::<Url>().map_err(|err| {
63 PriceOracleError::InvalidConfiguration(format!(
64 "invalid RPC endpoint {} for chain {}: {err}",
65 chain.rpc_endpoint, chain.chain_id
66 ))
67 })?;
68 let address = feed_address.parse::<Address>().map_err(|err| {
69 PriceOracleError::InvalidConfiguration(format!(
70 "invalid sequencer uptime feed {} for chain {}: {err}",
71 feed_address, chain.chain_id
72 ))
73 })?;
74 let provider = contract_backed_provider(url, egress_contract)?;
75 let contract = AggregatorV3Interface::new(address, &provider);
76 let latest = contract.latestRoundData().call().await.map_err(|err| {
77 PriceOracleError::Unavailable(format!(
78 "sequencer uptime read failed for chain {} at {}: {err}",
79 chain.chain_id, feed_address
80 ))
81 })?;
82 let answer = u8::try_from(latest.answer).map_err(|_| {
83 PriceOracleError::InvalidFeed(format!(
84 "sequencer uptime answer was invalid for chain {} at {}",
85 chain.chain_id, feed_address
86 ))
87 })?;
88 if answer > 1 {
89 return Err(PriceOracleError::InvalidFeed(format!(
90 "sequencer uptime answer {} was unsupported for chain {} at {}",
91 answer, chain.chain_id, feed_address
92 )));
93 }
94 let status_started_at = u64::try_from(latest.startedAt).map_err(|_| {
95 PriceOracleError::InvalidFeed(format!(
96 "sequencer startedAt overflowed for chain {} at {}",
97 chain.chain_id, feed_address
98 ))
99 })?;
100 let availability = if answer == 1 {
101 SequencerAvailability::Down
102 } else if status_started_at > 0
103 && now.saturating_sub(status_started_at) < chain.sequencer_grace_period_seconds
104 {
105 SequencerAvailability::Recovering
106 } else {
107 SequencerAvailability::Up
108 };
109 Ok(Some(SequencerStatus {
110 chain_id: chain.chain_id,
111 feed_address: feed_address.clone(),
112 checked_at: now,
113 status_started_at,
114 availability,
115 }))
116}
117
118#[cfg(test)]
119mod tests {
120 use super::*;
121 use crate::config::{ChainlinkNetworkConfig, BASE_MAINNET_CAIP2, BASE_MAINNET_CHAIN_ID};
122 use crate::test_support::{TestUnwrap, TestUnwrapErr};
123 use std::io::{Read, Write};
124 use std::net::{SocketAddr, TcpListener, TcpStream};
125 use std::sync::mpsc::{self, Receiver};
126 use std::thread::{self, JoinHandle};
127
128 fn authority(addr: SocketAddr) -> String {
129 format!("{}:{}", addr.ip(), addr.port())
130 }
131
132 fn read_http_request(stream: &mut TcpStream) -> String {
133 let mut request = Vec::new();
134 let mut buffer = [0_u8; 1024];
135 loop {
136 let read = stream
137 .read(&mut buffer)
138 .unwrap_or_else(|err| panic!("read request: {err}"));
139 if read == 0 {
140 break;
141 }
142 request.extend_from_slice(&buffer[..read]);
143 if request.windows(4).any(|window| window == b"\r\n\r\n") {
144 break;
145 }
146 }
147 String::from_utf8_lossy(&request).into_owned()
148 }
149
150 fn spawn_single_response_server<F>(
151 response_for_request: F,
152 ) -> (SocketAddr, Receiver<String>, JoinHandle<()>)
153 where
154 F: FnOnce(String) -> String + Send + 'static,
155 {
156 let listener =
157 TcpListener::bind("127.0.0.1:0").unwrap_or_else(|err| panic!("bind server: {err}"));
158 let addr = listener
159 .local_addr()
160 .unwrap_or_else(|err| panic!("read server addr: {err}"));
161 let (tx, rx) = mpsc::channel();
162 let handle = thread::spawn(move || {
163 let (mut stream, _) = listener
164 .accept()
165 .unwrap_or_else(|err| panic!("accept request: {err}"));
166 let request = read_http_request(&mut stream);
167 tx.send(request.clone())
168 .unwrap_or_else(|err| panic!("send captured request: {err}"));
169 let response = response_for_request(request);
170 stream
171 .write_all(response.as_bytes())
172 .unwrap_or_else(|err| panic!("write response: {err}"));
173 });
174 (addr, rx, handle)
175 }
176
177 fn base_chain(rpc_endpoint: &str, feed: Option<&str>) -> ChainlinkNetworkConfig {
178 ChainlinkNetworkConfig {
179 chain_id: BASE_MAINNET_CHAIN_ID,
180 label: "base-mainnet".to_string(),
181 caip2: BASE_MAINNET_CAIP2.to_string(),
182 rpc_endpoint: rpc_endpoint.to_string(),
183 enabled: true,
184 sequencer_uptime_feed: feed.map(ToString::to_string),
185 sequencer_grace_period_seconds: 300,
186 }
187 }
188
189 fn permissive_contract() -> HttpEgressContract {
190 HttpEgressContract::permissive_for_tests("rpc.example")
191 }
192
193 #[tokio::test]
194 async fn returns_none_when_no_sequencer_feed_is_configured() {
195 let result = read_sequencer_status(
196 &base_chain("https://rpc.example", None),
197 1_743_292_780,
198 &permissive_contract(),
199 )
200 .await
201 .test_unwrap("no feed configured");
202
203 assert_eq!(result, None);
204 }
205
206 #[tokio::test]
207 async fn rejects_invalid_rpc_endpoints() {
208 let error = read_sequencer_status(
209 &base_chain(
210 "not a url",
211 Some("0xFdB631F5EE196F0ed6FAa767959853A9F217697D"),
212 ),
213 1_743_292_780,
214 &permissive_contract(),
215 )
216 .await
217 .test_unwrap_err("invalid rpc endpoint");
218
219 assert!(matches!(error, PriceOracleError::InvalidConfiguration(_)));
220 }
221
222 #[tokio::test]
223 async fn accepts_hostname_rpc_endpoints_for_contract_backed_dispatch() {
224 let error = read_sequencer_status(
225 &base_chain(
226 "https://rpc.example",
227 Some("0xFdB631F5EE196F0ed6FAa767959853A9F217697D"),
228 ),
229 1_743_292_780,
230 &permissive_contract(),
231 )
232 .await
233 .test_unwrap_err("unreachable hostname RPC endpoint should fail during dispatch");
234
235 let message = error.to_string();
236 assert!(
237 message.contains("rpc.example") && message.contains("oracle backend unavailable"),
238 "unexpected sequencer hostname dispatch error: {message}"
239 );
240 }
241
242 #[tokio::test]
243 async fn rpc_transport_rejects_redirects_outside_contract() {
244 let denied_target = "http://127.0.0.1:9/final";
245 let (addr, _rx, handle) = spawn_single_response_server({
246 let denied_target = denied_target.to_string();
247 move |_| {
248 format!(
249 "HTTP/1.1 302 Found\r\nLocation: {denied_target}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
250 )
251 }
252 });
253 let endpoint = format!("http://{}/rpc", authority(addr));
254 let contract = HttpEgressContract::permissive_for_tests(&authority(addr));
255
256 let error = read_sequencer_status(
257 &base_chain(
258 &endpoint,
259 Some("0xFdB631F5EE196F0ed6FAa767959853A9F217697D"),
260 ),
261 1_743_292_780,
262 &contract,
263 )
264 .await
265 .test_unwrap_err("redirect target should be denied");
266 let message = error.to_string();
267 assert!(
268 message.contains("sequencer uptime read failed")
269 && message.contains("authority")
270 && message.contains("not allowed"),
271 "unexpected sequencer redirect denial: {message}"
272 );
273 handle
274 .join()
275 .unwrap_or_else(|_| panic!("join redirect server"));
276 }
277
278 #[tokio::test]
279 async fn rpc_transport_rejects_oversized_responses() {
280 let (addr, _rx, handle) = spawn_single_response_server(|_| {
281 "HTTP/1.1 200 OK\r\nContent-Length: 6\r\nConnection: close\r\n\r\nabcdef".to_string()
282 });
283 let endpoint = format!("http://{}/rpc", authority(addr));
284 let mut contract = HttpEgressContract::permissive_for_tests(&authority(addr));
285 contract.max_response_bytes = 5;
286
287 let error = read_sequencer_status(
288 &base_chain(
289 &endpoint,
290 Some("0xFdB631F5EE196F0ed6FAa767959853A9F217697D"),
291 ),
292 1_743_292_780,
293 &contract,
294 )
295 .await
296 .test_unwrap_err("oversized response should be denied");
297 let message = error.to_string();
298 assert!(
299 message.contains("sequencer uptime read failed")
300 && message.contains("response size 6 exceeds maximum 5"),
301 "unexpected sequencer response-size denial: {message}"
302 );
303 handle.join().unwrap_or_else(|_| panic!("join body server"));
304 }
305
306 #[tokio::test]
307 async fn rejects_invalid_feed_addresses() {
308 let error = read_sequencer_status(
309 &base_chain("http://203.0.113.10", Some("not-an-address")),
310 1_743_292_780,
311 &HttpEgressContract::permissive_for_tests("203.0.113.10"),
312 )
313 .await
314 .test_unwrap_err("invalid sequencer feed");
315
316 assert!(matches!(error, PriceOracleError::InvalidConfiguration(_)));
317 }
318}