Skip to main content

phos_data_network/rpc/
http_rpc.rs

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