phos_data_network/rpc/
http_rpc.rs1#[cfg(not(target_arch = "wasm32"))]
2use std::time::Duration;
3
4use async_trait::async_trait;
5use cometbft_rpc::{
6 self as rpc, client::Client, request::RequestMessage, response::Response, SimpleRequest,
7};
8use reqwest::header;
9
10use phos_light_client::{
11 components::io::{AtHeight, Io, IoError},
12 verifier::types::{Height, LightBlock, PeerId, SignedHeader, ValidatorAddress, ValidatorSet},
13};
14
15#[derive(Clone, Debug)]
16pub struct HttpClient {
17 inner: reqwest::Client,
18 url: reqwest::Url,
19}
20
21impl HttpClient {
22 pub fn new(url: reqwest::Url) -> Self {
23 #[cfg(not(target_arch = "wasm32"))]
24 let inner = reqwest::Client::builder()
25 .timeout(Duration::from_secs(30))
26 .connect_timeout(Duration::from_secs(10))
27 .pool_idle_timeout(Duration::from_secs(60))
28 .build()
29 .unwrap_or_else(|_| reqwest::Client::new());
30
31 #[cfg(target_arch = "wasm32")]
32 let inner = reqwest::Client::new();
33
34 Self { inner, url }
35 }
36
37 fn build_request<R>(&self, request: R) -> Result<reqwest::Request, rpc::Error>
38 where
39 R: RequestMessage,
40 {
41 let request_body = request.into_json();
42
43 self.inner
44 .post(self.url.clone())
45 .header(header::CONTENT_TYPE, "application/json")
46 .body(request_body.into_bytes())
47 .build()
48 .map_err(rpc::Error::http)
49 }
50}
51
52#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
53#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
54impl Client for HttpClient {
55 async fn perform<R>(&self, request: R) -> Result<R::Output, rpc::Error>
56 where
57 R: SimpleRequest,
58 {
59 let request = self.build_request(request)?;
60 let response = self
61 .inner
62 .execute(request)
63 .await
64 .map_err(rpc::Error::http)?;
65 let response_status = response.status();
66 let response_body = response.bytes().await.map_err(rpc::Error::http)?;
67
68 if response_status != reqwest::StatusCode::OK {
69 return Err(rpc::Error::http_request_failed(response_status));
70 }
71
72 R::Response::from_string(&response_body).map(Into::into)
73 }
74}
75
76#[derive(Clone, Debug)]
77pub struct ProdIo {
78 peer_id: PeerId,
79 rpc_client: HttpClient,
80}
81
82impl ProdIo {
83 pub fn new(peer_id: PeerId, rpc_client: HttpClient) -> Self {
84 Self {
85 peer_id,
86 rpc_client,
87 }
88 }
89
90 pub fn peer_id(&self) -> PeerId {
91 self.peer_id
92 }
93
94 pub async fn fetch_signed_header(&self, height: AtHeight) -> Result<SignedHeader, IoError> {
95 let response = match height {
96 AtHeight::Highest => self.rpc_client.latest_commit().await,
97 AtHeight::At(height) => self.rpc_client.commit(height).await,
98 }
99 .map_err(IoError::from_rpc)?;
100
101 Ok(response.signed_header)
102 }
103
104 pub async fn fetch_validator_set(
105 &self,
106 height: AtHeight,
107 proposer_address: Option<ValidatorAddress>,
108 ) -> Result<ValidatorSet, IoError> {
109 let height = match height {
110 AtHeight::Highest => return Err(IoError::invalid_height()),
111 AtHeight::At(height) => height,
112 };
113
114 let response = self
115 .rpc_client
116 .validators(height, rpc::Paging::All)
117 .await
118 .map_err(IoError::rpc)?;
119
120 match proposer_address {
121 Some(proposer_address) => {
122 ValidatorSet::with_proposer(response.validators, proposer_address)
123 .map_err(IoError::invalid_validator_set)
124 }
125 None => Ok(ValidatorSet::without_proposer(response.validators)),
126 }
127 }
128
129 pub async fn fetch_block(&self, height: Height) -> Result<cometbft::block::Block, IoError> {
130 self.rpc_client
131 .block(height)
132 .await
133 .map(|response| response.block)
134 .map_err(IoError::from_rpc)
135 }
136}
137
138#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
139#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
140impl Io for ProdIo {
141 async fn fetch_light_block(&self, height: AtHeight) -> Result<LightBlock, IoError> {
142 let signed_header = self.fetch_signed_header(height).await?;
143 let height = signed_header.header.height;
144 let proposer_address = signed_header.header.proposer_address;
145
146 let validator_set = self
147 .fetch_validator_set(height.into(), Some(proposer_address))
148 .await?;
149 let next_validator_set = self
150 .fetch_validator_set(height.increment().into(), None)
151 .await?;
152
153 Ok(LightBlock::new(
154 signed_header,
155 validator_set,
156 next_validator_set,
157 self.peer_id,
158 ))
159 }
160}