Skip to main content

chio_link/
chainlink.rs

1use alloy_json_rpc::{RequestPacket, ResponsePacket};
2use alloy_primitives::Address;
3use alloy_provider::{Provider, ProviderBuilder};
4use alloy_rpc_client::ClientBuilder as AlloyClientBuilder;
5use alloy_sol_types::sol;
6use alloy_transport::{TransportError, TransportErrorKind, TransportFut};
7use chio_egress_contract::{client_builder_with_contract, send_with_contract, HttpEgressContract};
8use reqwest::{Client, Url};
9use std::task;
10use tower::Service;
11
12use crate::config::{ChainlinkFeedConfig, ChainlinkNetworkConfig, PairConfig};
13use crate::{ExchangeRate, OracleBackend, OracleBackendKind, OracleFuture, PriceOracleError};
14
15sol! {
16    #[sol(rpc)]
17    contract AggregatorV3Interface {
18        function latestRoundData() external view returns (
19            uint80 roundId,
20            int256 answer,
21            uint256 startedAt,
22            uint256 updatedAt,
23            uint80 answeredInRound
24        );
25        function decimals() external view returns (uint8 decimalsValue);
26    }
27}
28
29#[derive(Clone, Debug)]
30pub(crate) struct ContractJsonRpcTransport {
31    url: Url,
32    http_client: Client,
33    egress_contract: HttpEgressContract,
34}
35
36impl ContractJsonRpcTransport {
37    fn new(url: Url, egress_contract: &HttpEgressContract) -> Result<Self, PriceOracleError> {
38        // Policy check only; the pinned ContractDnsResolver on the client built
39        // below enforces hostname address-class at connect, so this does not
40        // resolve DNS here (a config-time lookup would be redundant + offline-fragile).
41        egress_contract
42            .enforce_url(url.as_str(), 0)
43            .map_err(|err| {
44                PriceOracleError::InvalidConfiguration(format!(
45                    "HttpEgressContract rejects JSON-RPC endpoint: {err}"
46                ))
47            })?;
48        let http_client = client_builder_with_contract(egress_contract)
49            .build()
50            .map_err(|err| {
51                PriceOracleError::Unavailable(format!(
52                    "building contract-backed JSON-RPC client failed: {err}"
53                ))
54            })?;
55        Ok(Self {
56            url,
57            http_client,
58            egress_contract: egress_contract.clone(),
59        })
60    }
61
62    async fn send_json_rpc(self, req: RequestPacket) -> Result<ResponsePacket, TransportError> {
63        let request = self
64            .http_client
65            .post(self.url.clone())
66            .json(&req)
67            .headers(req.headers())
68            .build()
69            .map_err(TransportErrorKind::custom)?;
70        let response = send_with_contract(&self.egress_contract, &self.http_client, request)
71            .await
72            .map_err(TransportErrorKind::custom)?;
73        let status = response.status();
74        let body = response.body();
75        if !status.is_success() {
76            return Err(TransportErrorKind::http_error(
77                status.as_u16(),
78                String::from_utf8_lossy(body).into_owned(),
79            ));
80        }
81        serde_json::from_slice(body)
82            .map_err(|err| TransportError::deser_err(err, String::from_utf8_lossy(body)))
83    }
84}
85
86impl Service<RequestPacket> for ContractJsonRpcTransport {
87    type Response = ResponsePacket;
88    type Error = TransportError;
89    type Future = TransportFut<'static>;
90
91    fn poll_ready(&mut self, _cx: &mut task::Context<'_>) -> task::Poll<Result<(), Self::Error>> {
92        task::Poll::Ready(Ok(()))
93    }
94
95    fn call(&mut self, req: RequestPacket) -> Self::Future {
96        let transport = self.clone();
97        Box::pin(async move { transport.send_json_rpc(req).await })
98    }
99}
100
101pub(crate) fn contract_backed_provider(
102    url: Url,
103    egress_contract: &HttpEgressContract,
104) -> Result<impl Provider, PriceOracleError> {
105    let is_local = rpc_url_is_local(&url);
106    let transport = ContractJsonRpcTransport::new(url, egress_contract)?;
107    let client = AlloyClientBuilder::default().transport(transport, is_local);
108    Ok(ProviderBuilder::new().connect_client(client))
109}
110
111fn rpc_url_is_local(url: &Url) -> bool {
112    matches!(
113        url.host_str(),
114        Some("localhost" | "localhost.localdomain" | "127.0.0.1" | "::1")
115    )
116}
117
118#[derive(Debug)]
119pub struct ChainlinkFeedReader {
120    networks: Vec<ChainlinkNetworkConfig>,
121    egress_contract: HttpEgressContract,
122}
123
124impl ChainlinkFeedReader {
125    /// Construct a [`ChainlinkFeedReader`] that enforces the supplied
126    /// [`HttpEgressContract`] on every RPC endpoint URL before any HTTP
127    /// connect through `alloy`. The contract is required: the typed
128    /// egress gate is the only thing standing between substrate
129    /// dispatches and unconstrained RPC URLs.
130    #[must_use]
131    pub fn new(networks: Vec<ChainlinkNetworkConfig>, egress_contract: HttpEgressContract) -> Self {
132        Self {
133            networks,
134            egress_contract,
135        }
136    }
137
138    pub fn try_new(
139        networks: Vec<ChainlinkNetworkConfig>,
140        egress_contract: HttpEgressContract,
141    ) -> Result<Self, PriceOracleError> {
142        egress_contract
143            .validate_dispatchable_with_pinned_dns()
144            .map_err(|err| {
145                PriceOracleError::InvalidConfiguration(format!(
146                    "Chainlink HttpEgressContract is not dispatchable with pinned DNS: {err}"
147                ))
148            })?;
149        Ok(Self::new(networks, egress_contract))
150    }
151
152    /// Alias for [`ChainlinkFeedReader::new`].
153    #[must_use]
154    pub fn with_contract(
155        networks: Vec<ChainlinkNetworkConfig>,
156        egress_contract: HttpEgressContract,
157    ) -> Self {
158        Self::new(networks, egress_contract)
159    }
160
161    fn network_for_pair(
162        &self,
163        pair: &PairConfig,
164    ) -> Result<&ChainlinkNetworkConfig, PriceOracleError> {
165        self.networks
166            .iter()
167            .find(|network| network.chain_id == pair.chain_id)
168            .ok_or_else(|| {
169                PriceOracleError::InvalidConfiguration(format!(
170                    "no Chainlink network is configured for {} chain_id {}",
171                    pair.pair(),
172                    pair.chain_id
173                ))
174            })
175    }
176}
177
178impl OracleBackend for ChainlinkFeedReader {
179    fn kind(&self) -> OracleBackendKind {
180        OracleBackendKind::Chainlink
181    }
182
183    fn read_rate<'a>(&'a self, pair: &'a PairConfig, now: u64) -> OracleFuture<'a> {
184        Box::pin(async move {
185            let feed =
186                pair.chainlink
187                    .as_ref()
188                    .ok_or_else(|| PriceOracleError::NoPairAvailable {
189                        base: pair.base.clone(),
190                        quote: pair.quote.clone(),
191                    })?;
192            let network = self.network_for_pair(pair)?;
193            // Policy check only (scheme/authority + IP-literal class). The pinned
194            // ContractDnsResolver on the contract-backed transport enforces
195            // hostname address-class at connect, so this does not resolve DNS
196            // here (a config-time lookup would be redundant + offline-fragile).
197            // The contract is non-optional in production paths; bypass is not
198            // possible without recompiling.
199            self.egress_contract
200                .enforce_url(&network.rpc_endpoint, 0)
201                .map_err(|err| {
202                    PriceOracleError::InvalidConfiguration(format!(
203                        "HttpEgressContract rejects Chainlink RPC endpoint: {err}"
204                    ))
205                })?;
206            read_chainlink_rate(
207                &network.rpc_endpoint,
208                pair,
209                feed,
210                now,
211                &self.egress_contract,
212            )
213            .await
214        })
215    }
216}
217
218async fn read_chainlink_rate(
219    rpc_endpoint: &str,
220    pair: &PairConfig,
221    feed: &ChainlinkFeedConfig,
222    now: u64,
223    egress_contract: &HttpEgressContract,
224) -> Result<ExchangeRate, PriceOracleError> {
225    let url = rpc_endpoint.parse::<Url>().map_err(|err| {
226        PriceOracleError::InvalidConfiguration(format!(
227            "invalid Chainlink RPC endpoint {rpc_endpoint}: {err}"
228        ))
229    })?;
230    let address = feed.address.parse::<Address>().map_err(|err| {
231        PriceOracleError::InvalidConfiguration(format!(
232            "invalid Chainlink feed address {} for {}: {err}",
233            feed.address,
234            pair.pair()
235        ))
236    })?;
237    let provider = contract_backed_provider(url, egress_contract)?;
238    let contract = AggregatorV3Interface::new(address, &provider);
239    let latest = contract.latestRoundData().call().await.map_err(|err| {
240        PriceOracleError::Unavailable(format!(
241            "Chainlink latestRoundData failed for {} at {}: {err}",
242            pair.pair(),
243            feed.address
244        ))
245    })?;
246    let decimals = contract.decimals().call().await.map_err(|err| {
247        PriceOracleError::Unavailable(format!(
248            "Chainlink decimals failed for {} at {}: {err}",
249            pair.pair(),
250            feed.address
251        ))
252    })?;
253    if decimals != feed.decimals {
254        return Err(PriceOracleError::InvalidFeed(format!(
255            "Chainlink decimals mismatch for {} at {}: configured {}, contract returned {}",
256            pair.pair(),
257            feed.address,
258            feed.decimals,
259            decimals
260        )));
261    }
262    let answer = u128::try_from(latest.answer).map_err(|_| {
263        PriceOracleError::InvalidFeed(format!(
264            "Chainlink returned a negative or oversized answer for {} at {}",
265            pair.pair(),
266            feed.address
267        ))
268    })?;
269    if answer == 0 {
270        return Err(PriceOracleError::InvalidFeed(format!(
271            "Chainlink returned zero for {} at {}",
272            pair.pair(),
273            feed.address
274        )));
275    }
276    let updated_at = u64::try_from(latest.updatedAt).map_err(|_| {
277        PriceOracleError::InvalidFeed(format!(
278            "Chainlink updatedAt overflowed u64 for {} at {}",
279            pair.pair(),
280            feed.address
281        ))
282    })?;
283    if updated_at == 0 {
284        return Err(PriceOracleError::InvalidFeed(format!(
285            "Chainlink updatedAt was zero for {} at {}",
286            pair.pair(),
287            feed.address
288        )));
289    }
290    let denominator = 10_u128
291        .checked_pow(u32::from(feed.decimals))
292        .ok_or_else(|| {
293            PriceOracleError::ArithmeticOverflow(format!(
294                "decimal normalization overflowed for {} at {}",
295                pair.pair(),
296                feed.address
297            ))
298        })?;
299    let max_age_seconds = pair.policy.max_age_seconds.min(feed.heartbeat_seconds);
300    let rate = ExchangeRate {
301        base: pair.base.clone(),
302        quote: pair.quote.clone(),
303        rate_numerator: answer,
304        rate_denominator: denominator,
305        updated_at,
306        fetched_at: now,
307        source: "chainlink".to_string(),
308        feed_reference: feed.address.clone(),
309        max_age_seconds,
310        conversion_margin_bps: pair.policy.exchange_rate_margin_bps,
311        confidence_numerator: None,
312        confidence_denominator: None,
313    };
314    rate.ensure_fresh(now)?;
315    Ok(rate)
316}
317
318#[cfg(test)]
319mod tests {
320    use std::io::{Read, Write};
321    use std::net::{SocketAddr, TcpListener, TcpStream};
322    use std::sync::mpsc::{self, Receiver};
323    use std::thread::{self, JoinHandle};
324
325    use crate::config::{
326        ChainlinkFeedConfig, ChainlinkNetworkConfig, PairConfig, PairPolicy, BASE_MAINNET_CAIP2,
327        BASE_MAINNET_CHAIN_ID,
328    };
329    use crate::test_support::{TestUnwrap, TestUnwrapErr};
330    use crate::OracleBackend;
331
332    use super::{read_chainlink_rate, ChainlinkFeedReader};
333
334    fn authority(addr: SocketAddr) -> String {
335        format!("{}:{}", addr.ip(), addr.port())
336    }
337
338    fn read_http_request(stream: &mut TcpStream) -> String {
339        let mut request = Vec::new();
340        let mut buffer = [0_u8; 1024];
341        loop {
342            let read = stream
343                .read(&mut buffer)
344                .unwrap_or_else(|err| panic!("read request: {err}"));
345            if read == 0 {
346                break;
347            }
348            request.extend_from_slice(&buffer[..read]);
349            if request.windows(4).any(|window| window == b"\r\n\r\n") {
350                break;
351            }
352        }
353        String::from_utf8_lossy(&request).into_owned()
354    }
355
356    fn spawn_single_response_server<F>(
357        response_for_request: F,
358    ) -> (SocketAddr, Receiver<String>, JoinHandle<()>)
359    where
360        F: FnOnce(String) -> String + Send + 'static,
361    {
362        let listener =
363            TcpListener::bind("127.0.0.1:0").unwrap_or_else(|err| panic!("bind server: {err}"));
364        let addr = listener
365            .local_addr()
366            .unwrap_or_else(|err| panic!("read server addr: {err}"));
367        let (tx, rx) = mpsc::channel();
368        let handle = thread::spawn(move || {
369            let (mut stream, _) = listener
370                .accept()
371                .unwrap_or_else(|err| panic!("accept request: {err}"));
372            let request = read_http_request(&mut stream);
373            tx.send(request.clone())
374                .unwrap_or_else(|err| panic!("send captured request: {err}"));
375            let response = response_for_request(request);
376            stream
377                .write_all(response.as_bytes())
378                .unwrap_or_else(|err| panic!("write response: {err}"));
379        });
380        (addr, rx, handle)
381    }
382
383    fn pair_with_chainlink(address: &str) -> PairConfig {
384        PairConfig {
385            base: "ETH".to_string(),
386            quote: "USD".to_string(),
387            chain_id: BASE_MAINNET_CHAIN_ID,
388            chainlink: Some(ChainlinkFeedConfig {
389                address: address.to_string(),
390                decimals: 8,
391                heartbeat_seconds: 300,
392            }),
393            pyth: None,
394            policy: PairPolicy::volatile_default(),
395        }
396    }
397
398    fn base_network(rpc_endpoint: &str) -> ChainlinkNetworkConfig {
399        ChainlinkNetworkConfig {
400            chain_id: BASE_MAINNET_CHAIN_ID,
401            label: "base-mainnet".to_string(),
402            caip2: BASE_MAINNET_CAIP2.to_string(),
403            rpc_endpoint: rpc_endpoint.to_string(),
404            enabled: true,
405            sequencer_uptime_feed: None,
406            sequencer_grace_period_seconds: 300,
407        }
408    }
409
410    #[tokio::test]
411    async fn rejects_invalid_rpc_endpoints() {
412        let pair = pair_with_chainlink("0x71041dddad3595F9CEd3DcCFBe3D1F4b0a16Bb70");
413        let error = read_chainlink_rate(
414            "not a url",
415            &pair,
416            pair.chainlink.as_ref().test_unwrap("feed"),
417            1_743_292_780,
418            &chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
419        )
420        .await
421        .test_unwrap_err("invalid endpoint");
422        assert!(matches!(
423            error,
424            crate::PriceOracleError::InvalidConfiguration(_)
425        ));
426    }
427
428    #[test]
429    fn rejects_pairs_without_a_configured_network() {
430        let reader = ChainlinkFeedReader::new(
431            vec![base_network("https://rpc.example")],
432            chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
433        );
434        let mut pair = pair_with_chainlink("0x71041dddad3595F9CEd3DcCFBe3D1F4b0a16Bb70");
435        pair.chain_id = 1;
436
437        let error = reader
438            .network_for_pair(&pair)
439            .test_unwrap_err("missing network");
440
441        assert!(matches!(
442            error,
443            crate::PriceOracleError::InvalidConfiguration(_)
444        ));
445    }
446
447    #[test]
448    fn try_new_accepts_hostname_contract_with_pinned_dns() {
449        let reader = ChainlinkFeedReader::try_new(
450            vec![base_network("https://rpc.example")],
451            chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
452        )
453        .test_unwrap("hostname contract is resolver-enforced at dispatch");
454        assert_eq!(reader.networks.len(), 1);
455    }
456
457    #[tokio::test]
458    async fn backend_accepts_hostname_rpc_for_contract_backed_dispatch() {
459        let reader = ChainlinkFeedReader::new(
460            vec![base_network("https://rpc.example")],
461            chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
462        );
463        let pair = pair_with_chainlink("0x71041dddad3595F9CEd3DcCFBe3D1F4b0a16Bb70");
464
465        let error = reader
466            .read_rate(&pair, 1_743_292_780)
467            .await
468            .test_unwrap_err("unreachable hostname RPC should fail during dispatch");
469        let message = error.to_string();
470        assert!(
471            message.contains("rpc.example") && message.contains("oracle backend unavailable"),
472            "unexpected Chainlink hostname dispatch error: {message}"
473        );
474    }
475
476    #[tokio::test]
477    async fn rpc_transport_rejects_redirects_outside_contract() {
478        let denied_target = "http://127.0.0.1:9/final";
479        let (addr, _rx, handle) = spawn_single_response_server({
480            let denied_target = denied_target.to_string();
481            move |_| {
482                format!(
483                    "HTTP/1.1 302 Found\r\nLocation: {denied_target}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
484                )
485            }
486        });
487        let endpoint = format!("http://{}/rpc", authority(addr));
488        let contract =
489            chio_egress_contract::HttpEgressContract::permissive_for_tests(&authority(addr));
490        let pair = pair_with_chainlink("0x71041dddad3595F9CEd3DcCFBe3D1F4b0a16Bb70");
491
492        let error = read_chainlink_rate(
493            &endpoint,
494            &pair,
495            pair.chainlink.as_ref().test_unwrap("feed"),
496            1_743_292_780,
497            &contract,
498        )
499        .await
500        .test_unwrap_err("redirect target should be denied");
501        let message = error.to_string();
502        assert!(
503            message.contains("Chainlink latestRoundData failed")
504                && message.contains("authority")
505                && message.contains("not allowed"),
506            "unexpected Chainlink redirect denial: {message}"
507        );
508        handle
509            .join()
510            .unwrap_or_else(|_| panic!("join redirect server"));
511    }
512
513    #[tokio::test]
514    async fn rpc_transport_rejects_oversized_responses() {
515        let (addr, _rx, handle) = spawn_single_response_server(|_| {
516            "HTTP/1.1 200 OK\r\nContent-Length: 6\r\nConnection: close\r\n\r\nabcdef".to_string()
517        });
518        let endpoint = format!("http://{}/rpc", authority(addr));
519        let mut contract =
520            chio_egress_contract::HttpEgressContract::permissive_for_tests(&authority(addr));
521        contract.max_response_bytes = 5;
522        let pair = pair_with_chainlink("0x71041dddad3595F9CEd3DcCFBe3D1F4b0a16Bb70");
523
524        let error = read_chainlink_rate(
525            &endpoint,
526            &pair,
527            pair.chainlink.as_ref().test_unwrap("feed"),
528            1_743_292_780,
529            &contract,
530        )
531        .await
532        .test_unwrap_err("oversized response should be denied");
533        let message = error.to_string();
534        assert!(
535            message.contains("Chainlink latestRoundData failed")
536                && message.contains("response size 6 exceeds maximum 5"),
537            "unexpected Chainlink response-size denial: {message}"
538        );
539        handle.join().unwrap_or_else(|_| panic!("join body server"));
540    }
541
542    #[tokio::test]
543    async fn rejects_invalid_feed_addresses() {
544        let pair = pair_with_chainlink("not-an-address");
545
546        let error = read_chainlink_rate(
547            "https://rpc.example",
548            &pair,
549            pair.chainlink.as_ref().test_unwrap("feed"),
550            1_743_292_780,
551            &chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
552        )
553        .await
554        .test_unwrap_err("invalid feed address");
555
556        assert!(matches!(
557            error,
558            crate::PriceOracleError::InvalidConfiguration(_)
559        ));
560    }
561
562    #[tokio::test]
563    async fn backend_rejects_pairs_without_chainlink_feeds() {
564        let reader = ChainlinkFeedReader::new(
565            vec![base_network("https://rpc.example")],
566            chio_egress_contract::HttpEgressContract::permissive_for_tests("rpc.example"),
567        );
568        let pair = PairConfig {
569            base: "ETH".to_string(),
570            quote: "USD".to_string(),
571            chain_id: BASE_MAINNET_CHAIN_ID,
572            chainlink: None,
573            pyth: None,
574            policy: PairPolicy::volatile_default(),
575        };
576
577        let error = reader
578            .read_rate(&pair, 1_743_292_780)
579            .await
580            .test_unwrap_err("missing feed");
581
582        assert!(matches!(
583            error,
584            crate::PriceOracleError::NoPairAvailable { .. }
585        ));
586    }
587}