Skip to main content

chio_link/
sequencer.rs

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
39/// Read the L2 sequencer-uptime status while enforcing the supplied
40/// [`HttpEgressContract`] on the chain RPC endpoint URL before any HTTP
41/// connect. The contract is required: production paths must thread the
42/// tenant-scoped contract through this entry point so SSRF gates apply
43/// uniformly.
44pub 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    // Policy check only; the pinned ContractDnsResolver (contract_backed_provider)
53    // enforces hostname address-class at connect, so this does not resolve DNS
54    // here (a config-time lookup would be redundant and offline-fragile).
55    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}