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 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 #[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 #[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 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}