Files

695 lines
24 KiB
Rust

use std::sync::Arc;
use futures::{Stream, StreamExt as _, TryStreamExt as _};
pub use lb_chain_broadcast_service::BlockInfo;
pub use lb_chain_service::{ChainServiceInfo, CryptarchiaInfo, PhaseTag, Slot, State};
pub use lb_core::events::{Event, Events, TxEventPayload};
use lb_core::{
block::MAX_BLOCK_TRANSACTIONS_SIZE,
header::{ContentId, HeaderId},
mantle::{
SignedMantleTx, channel::ChannelState, ops::channel::ChannelId,
transactions::states::Unverified,
},
proofs::leader_proof::Groth16LeaderProof,
sdp::{DeclarationId, DeclarationMessage},
};
use lb_groth16::fr_to_bytes;
pub use lb_http_api_common::TimeInfo;
use lb_http_api_common::{
MAX_BLOCKS_STREAM_BLOCKS, MAX_BLOCKS_STREAM_CHUNK_SIZE,
bodies::{
blend::JoinBlendRequestBody,
mantle::GasPricesResponseBody,
wallet::{
balance::WalletBalanceResponseBody,
claimable_vouchers::WalletClaimableVouchersResponseBody,
fund::{WalletFundRequestBody, WalletFundResponseBody},
transfer_funds::{WalletTransferFundsRequestBody, WalletTransferFundsResponseBody},
},
},
paths::{
BLEND_DISPERSE_TRANSACTION, BLEND_JOIN_NETWORK, BLEND_PENDING_TRANSACTIONS, BLOCK_EVENTS,
BLOCKS, BLOCKS_DETAIL, BLOCKS_RANGE_STREAM, BLOCKS_STREAM, CHANNEL, CRYPTARCHIA_INFO,
CRYPTARCHIA_LIB_STREAM, LEADER_CLAIM_VOUCHERS, MANTLE_GAS_PRICES, MEMPOOL_ADD_TX,
NODE_VERSION, SDP_POST_DECLARATION, TIME_INFO,
wallet::{BALANCE, FUND, TRANSACTIONS_TRANSFER_FUNDS},
},
queries::BlocksStreamQuery,
settings::default_max_body_size,
};
use lb_key_management_system_keys::keys::{Ed25519Signature, ZkPublicKey};
use lb_log_targets::http_client;
use log::warn;
use reqwest::{Client, ClientBuilder, RequestBuilder, StatusCode, Url};
use serde::{Deserialize, Serialize, de::DeserializeOwned};
use tokio_util::{
codec::{FramedRead, LinesCodec},
io::StreamReader,
};
const LOG_TARGET: &str = http_client::ROOT;
/// Client-side header representation matching the server's
/// `ApiHeaderSerializer`.
#[derive(Clone, Debug, Deserialize)]
pub struct ApiHeader {
pub id: HeaderId,
pub parent_block: HeaderId,
pub slot: Slot,
pub body_root: ContentId,
pub proof_of_leadership: Groth16LeaderProof,
}
/// Client-side block representation matching the server's `ApiBlockSerializer`.
/// Note: The server omits the signature field.
#[derive(Clone, Debug, Deserialize)]
pub struct ApiBlock {
pub header: ApiHeader,
pub uncle_headers: Vec<ApiSignedHeader>,
pub transactions: Vec<SignedMantleTx<Unverified>>,
}
/// Client-side signed header representation matching the server's
/// `ApiSignedHeader`.
#[derive(Clone, Debug, Deserialize)]
pub struct ApiSignedHeader {
pub header: ApiHeader,
pub signature: Ed25519Signature,
}
/// Processed block event from the block stream.
/// Matches the server's `ApiProcessedBlockEvent`.
#[derive(Clone, Debug, Deserialize)]
pub struct ProcessedBlockEvent {
pub block: ApiBlock,
pub tip: HeaderId,
pub tip_slot: Slot,
pub lib: HeaderId,
pub lib_slot: Slot,
}
#[derive(thiserror::Error, Debug)]
pub enum Error {
#[error("Internal server error: {0}")]
Server(String),
#[error("Failed to execute request: {0}")]
Client(String),
#[error(transparent)]
Request(#[from] reqwest::Error),
#[error(transparent)]
Url(#[from] url::ParseError),
}
#[derive(Clone)]
pub struct BasicAuthCredentials {
username: String,
password: Option<String>,
}
impl BasicAuthCredentials {
#[must_use]
pub const fn new(username: String, password: Option<String>) -> Self {
Self { username, password }
}
}
#[derive(Clone)]
pub struct CommonHttpClient {
client: Arc<Client>,
basic_auth: Option<BasicAuthCredentials>,
}
impl CommonHttpClient {
fn is_channel_not_found(body: &str) -> bool {
body.to_ascii_lowercase().contains("channel not found")
}
#[must_use]
pub fn new(basic_auth: Option<BasicAuthCredentials>) -> Self {
let initial_stream_window_size: u32 =
u32::try_from(6 * default_max_body_size() / 10).unwrap_or(4 * 1025);
let client = ClientBuilder::new()
.http2_initial_stream_window_size(initial_stream_window_size)
.build()
.expect("Client from default settings should be able to build");
Self {
client: Arc::new(client),
basic_auth,
}
}
pub async fn post<Req, Res>(&self, request_url: Url, request_body: &Req) -> Result<Res, Error>
where
Req: Serialize + ?Sized + Send + Sync,
Res: DeserializeOwned + Send + Sync,
{
let request = self.client.post(request_url).json(request_body);
self.execute_request::<Res>(request).await
}
pub async fn get<Req, Res>(
&self,
request_url: Url,
request_body: Option<&Req>,
) -> Result<Res, Error>
where
Req: Serialize + ?Sized + Send + Sync,
Res: DeserializeOwned + Send + Sync,
{
let mut request = self.client.get(request_url);
if let Some(request_body) = request_body {
request = request.json(request_body);
}
self.execute_request::<Res>(request).await
}
pub async fn put<Req, Res>(
&self,
request_url: Url,
request_body: Option<&Req>,
) -> Result<Res, Error>
where
Req: Serialize + ?Sized + Send + Sync,
Res: DeserializeOwned + Send + Sync,
{
let mut request = self.client.put(request_url);
if let Some(request_body) = request_body {
request = request.json(request_body);
}
self.execute_request::<Res>(request).await
}
async fn execute_request<Res: DeserializeOwned>(
&self,
mut request: RequestBuilder,
) -> Result<Res, Error> {
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let body = response.text().await.map_err(Error::Request)?;
match status {
StatusCode::OK | StatusCode::CREATED => serde_json::from_str(&body)
.map_err(|e| Error::Server(format!("Failed to parse response: {e}"))),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server(body)),
_ => Err(Error::Server(format!(
"Unexpected response [{status}]: {body}",
))),
}
}
pub async fn get_lib_stream(
&self,
base_url: Url,
) -> Result<impl Stream<Item = BlockInfo> + use<>, Error> {
let request_url = base_url
.join(CRYPTARCHIA_LIB_STREAM.trim_start_matches('/'))
.map_err(Error::Url)?;
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let lib_stream = response.bytes_stream().filter_map(async |item| {
let bytes = item.ok()?;
serde_json::from_slice::<BlockInfo>(&bytes).ok()
});
match status {
StatusCode::OK => Ok(lib_stream),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server("Error".to_owned())),
_ => Err(Error::Server(format!("Unexpected response [{status}]"))),
}
}
pub async fn get_block_by_id(
&self,
base_url: Url,
id: HeaderId,
) -> Result<Option<ApiBlock>, Error> {
let path = BLOCKS_DETAIL
.trim_start_matches('/')
.replace(":id", &id.to_string());
let request_url = base_url.join(path.as_str()).map_err(Error::Url)?;
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let body = response.text().await.map_err(Error::Request)?;
match status {
StatusCode::OK => serde_json::from_str::<ApiBlock>(&body)
.map(Some)
.map_err(|e| Error::Server(format!("Failed to parse response: {e}"))),
StatusCode::NOT_FOUND => Ok(None),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server(body)),
_ => Err(Error::Server(format!(
"Unexpected response [{status}]: {body}",
))),
}
}
/// Get the events emitted by execution of a block by its id.
pub async fn get_block_events(
&self,
base_url: Url,
id: HeaderId,
) -> Result<Option<Events>, Error> {
let path = BLOCK_EVENTS
.trim_start_matches('/')
.replace(":id", &id.to_string());
let request_url = base_url.join(path.as_str()).map_err(Error::Url)?;
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let body = response.text().await.map_err(Error::Request)?;
match status {
StatusCode::OK => serde_json::from_str::<Events>(&body)
.map(Some)
.map_err(|e| Error::Server(format!("Failed to parse response: {e}"))),
StatusCode::NOT_FOUND => Ok(None),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server(body)),
_ => Err(Error::Server(format!(
"Unexpected response [{status}]: {body}",
))),
}
}
pub async fn post_transaction<Tx>(&self, base_url: Url, transaction: Tx) -> Result<(), Error>
where
Tx: Serialize + Send + Sync + 'static,
{
let request_url = base_url
.join(MEMPOOL_ADD_TX.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, &transaction).await
}
/// Send a transaction through the Blend network, without adding it to the
/// node's own mempool.
///
/// The node it is submitted to never gossips the transaction itself — that
/// would say where it came from — so it reaches the network blended, and
/// comes back to this node's mempool like any other transaction once
/// whichever node exits it gossips it on.
///
/// Returns the transaction's id.
pub async fn blend_transaction<Tx, Id>(
&self,
base_url: Url,
transaction: Tx,
) -> Result<Id, Error>
where
Tx: Serialize + Send + Sync + 'static,
Id: for<'de> Deserialize<'de> + Send + Sync,
{
let request_url = base_url
.join(BLEND_DISPERSE_TRANSACTION.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, &transaction).await
}
/// The ids of the transactions the node is still waiting on a `PoW`
/// solution for before it can blend them.
pub async fn blend_pending_transactions<Id>(&self, base_url: Url) -> Result<Vec<Id>, Error>
where
Id: for<'de> Deserialize<'de> + Send + Sync,
{
let request_url = base_url
.join(BLEND_PENDING_TRANSACTIONS.trim_start_matches('/'))
.map_err(Error::Url)?;
self.get::<(), _>(request_url, None).await
}
/// Post a service declaration to the SDP endpoint.
pub async fn post_declaration(
&self,
base_url: Url,
declaration: &DeclarationMessage,
) -> Result<DeclarationId, Error> {
let request_url = base_url
.join(SDP_POST_DECLARATION.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, declaration).await
}
/// Get the version of the node, e.g. `0.1.2 (abcdefaa)`.
pub async fn get_node_version(&self, base_url: Url) -> Result<String, Error> {
let request_url = base_url
.join(NODE_VERSION.trim_start_matches('/'))
.map_err(Error::Url)?;
self.get::<(), String>(request_url, None).await
}
/// Get consensus info (tip, height, etc.)
pub async fn consensus_info(&self, base_url: Url) -> Result<ChainServiceInfo, Error> {
let request_url = base_url
.join(CRYPTARCHIA_INFO.trim_start_matches('/'))
.map_err(Error::Url)?;
self.get::<(), ChainServiceInfo>(request_url, None).await
}
/// Get the gas prices from the ledger state at `tip`, or at the current
/// tip when `tip` is `None`.
pub async fn gas_prices(
&self,
base_url: Url,
tip: Option<HeaderId>,
) -> Result<GasPricesResponseBody, Error> {
let mut request_url = base_url
.join(MANTLE_GAS_PRICES.trim_start_matches('/'))
.map_err(Error::Url)?;
if let Some(t) = tip {
request_url
.query_pairs_mut()
.append_pair("tip", &t.to_string());
}
self.get::<(), GasPricesResponseBody>(request_url, None)
.await
}
/// Get time service info derived from deployment settings.
pub async fn time_info(&self, base_url: Url) -> Result<TimeInfo, Error> {
let request_url = base_url
.join(TIME_INFO.trim_start_matches('/'))
.map_err(Error::Url)?;
self.get::<(), TimeInfo>(request_url, None).await
}
/// Get channel state for a specific channel id.
pub async fn channel_state(
&self,
base_url: Url,
channel_id: ChannelId,
) -> Result<Option<ChannelState>, Error> {
let path = CHANNEL
.trim_start_matches('/')
.replace(":id", &hex::encode(channel_id.as_ref()));
let request_url = base_url.join(path.as_str()).map_err(Error::Url)?;
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let body = response.text().await.map_err(Error::Request)?;
match status {
StatusCode::OK => serde_json::from_str::<ChannelState>(&body)
.map(Some)
.map_err(|e| Error::Server(format!("Failed to parse response: {e}"))),
StatusCode::NOT_FOUND => Ok(None),
StatusCode::INTERNAL_SERVER_ERROR if Self::is_channel_not_found(&body) => Ok(None),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server(body)),
_ if Self::is_channel_not_found(&body) => Ok(None),
_ => Err(Error::Server(format!(
"Unexpected response [{status}]: {body}",
))),
}
}
/// Get immutable blocks in a slot range.
pub async fn get_immutable_blocks(
&self,
base_url: Url,
slot_from: u64,
slot_to: u64,
) -> Result<Vec<ApiBlock>, Error> {
let mut request_url = base_url
.join(BLOCKS.trim_start_matches('/'))
.map_err(Error::Url)?;
request_url
.query_pairs_mut()
.append_pair("slot_from", &slot_from.to_string())
.append_pair("slot_to", &slot_to.to_string());
self.get::<(), Vec<ApiBlock>>(request_url, None).await
}
/// Subscribe to the processed blocks stream.
/// Each event contains the block, current tip, and current LIB.
pub async fn get_blocks_stream(
&self,
base_url: Url,
) -> Result<impl Stream<Item = ProcessedBlockEvent> + use<>, Error> {
let request_url = base_url
.join(BLOCKS_STREAM.trim_start_matches('/'))
.map_err(Error::Url)?;
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let blocks_stream = response.bytes_stream().filter_map(async |item| {
let bytes = item.ok()?;
serde_json::from_slice::<ProcessedBlockEvent>(&bytes).ok()
});
match status {
StatusCode::OK => Ok(blocks_stream),
StatusCode::INTERNAL_SERVER_ERROR => Err(Error::Server("Error".to_owned())),
_ => Err(Error::Server(format!("Unexpected response [{status}]"))),
}
}
/// Build the request URL for the `get_blocks_range_stream` method with the
/// given parameters.
pub fn build_blocks_range_stream_request_url(
base_url: &Url,
params: &BlocksStreamQuery,
) -> Result<Url, Error> {
let mut request_url = base_url
.join(BLOCKS_RANGE_STREAM.trim_start_matches('/'))
.map_err(Error::Url)?;
params.append_to_url(&mut request_url);
Ok(request_url)
}
async fn send_blocks_range_stream_request(
&self,
request_url: Url,
) -> Result<reqwest::Response, Error> {
let mut request = self.client.get(request_url);
if let Some(basic_auth) = &self.basic_auth {
request = request.basic_auth(&basic_auth.username, basic_auth.password.as_deref());
}
let response = request.send().await.map_err(Error::Request)?;
let status = response.status();
let response_url = response.url().clone();
if status != StatusCode::OK {
let body = response.text().await.map_err(Error::Request)?;
return match status {
StatusCode::INTERNAL_SERVER_ERROR => {
Err(Error::Server(format!("{body} [{response_url}]")))
}
_ => Err(Error::Server(format!(
"Unexpected response [{status}] at [{response_url}]: {body}"
))),
};
}
Ok(response)
}
// Helper function to validate inputs for block streaming methods.
fn verify_inputs(
blocks_limit: Option<std::num::NonZero<usize>>,
slot_from: Option<u64>,
slot_to: Option<u64>,
server_batch_size: Option<std::num::NonZero<usize>>,
) -> Result<(), Error> {
if let Some(blocks) = blocks_limit
&& blocks.get() > MAX_BLOCKS_STREAM_BLOCKS
{
return Err(Error::Client(format!(
"'blocks_limit' must be <= {MAX_BLOCKS_STREAM_BLOCKS}, got {blocks}"
)));
}
if let Some(size) = server_batch_size
&& size.get() > MAX_BLOCKS_STREAM_CHUNK_SIZE
{
return Err(Error::Client(format!(
"'server_batch_size' must be <= {MAX_BLOCKS_STREAM_CHUNK_SIZE}, got {size}"
)));
}
if let (Some(slot_from), Some(slot_to)) = (slot_from, slot_to)
&& slot_from > slot_to
{
return Err(Error::Client(format!(
"'slot_from' must be <= 'slot_to', got slot_from={slot_from}, slot_to={slot_to}"
)));
}
Ok(())
}
/// Stream processed blocks in a slot-bounded window.
///
/// `server_batch_size` lets callers request smaller chunks; the server
/// still enforces its own upper bound.
pub async fn get_blocks_range_stream(
&self,
base_url: Url,
params: BlocksStreamQuery,
) -> Result<impl Stream<Item = ProcessedBlockEvent> + use<>, Error> {
let request_url = Self::build_blocks_range_stream_request_url(&base_url, &params)?;
Self::verify_inputs(
params.blocks_limit,
params.slot_from,
params.slot_to,
params.server_batch_size,
)?;
let response = self.send_blocks_range_stream_request(request_url).await?;
Ok(Self::parse_processed_blocks_range_event_stream(response))
}
fn parse_processed_blocks_range_event_stream(
response: reqwest::Response,
) -> impl Stream<Item = ProcessedBlockEvent> {
// NDJSON event upper bound; margin above max serialized single event line
const MAX_NDJSON_LINE_BYTES: usize = MAX_BLOCK_TRANSACTIONS_SIZE * 3 / 2;
const LOG_LINE_PREVIEW_CHARS: usize = 256;
let byte_stream = response.bytes_stream().map_err(std::io::Error::other);
let reader = StreamReader::new(byte_stream);
let codec = LinesCodec::new_with_max_length(MAX_NDJSON_LINE_BYTES);
let lines = FramedRead::new(reader, codec);
lines.filter_map(async |line_result| match line_result {
Ok(line) => {
if line.is_empty() {
return None;
}
match serde_json::from_str::<ProcessedBlockEvent>(&line) {
Ok(event) => Some(event),
Err(err) => {
let preview: String = line.chars().take(LOG_LINE_PREVIEW_CHARS).collect();
warn!(
target: LOG_TARGET,
"blocks stream JSON decode failed: {err}; line_preview={preview:?}"
);
None
}
}
}
Err(err) => {
warn!(target: LOG_TARGET, "blocks stream line decode failed: {err}");
None
}
})
}
/// Get the balance for a specific `ZkPublicKey`.
pub async fn get_wallet_balance(
&self,
base_url: Url,
zk_pk: ZkPublicKey,
tip: Option<HeaderId>,
) -> Result<WalletBalanceResponseBody, Error> {
let key_id = hex::encode(fr_to_bytes(zk_pk.as_fr()));
let mut request_url = base_url
.join(&BALANCE.replace(":public_key", &key_id))
.map_err(Error::Url)?;
if let Some(t) = tip {
request_url
.query_pairs_mut()
.append_pair("tip", &t.to_string());
}
self.get::<(), WalletBalanceResponseBody>(request_url, None)
.await
}
/// Get claimable reward vouchers tracked by the wallet.
pub async fn get_claimable_vouchers(
&self,
base_url: Url,
tip: Option<HeaderId>,
) -> Result<WalletClaimableVouchersResponseBody, Error> {
let mut request_url = base_url
.join(LEADER_CLAIM_VOUCHERS.trim_start_matches('/'))
.map_err(Error::Url)?;
if let Some(t) = tip {
request_url
.query_pairs_mut()
.append_pair("tip", &t.to_string());
}
self.get::<(), WalletClaimableVouchersResponseBody>(request_url, None)
.await
}
/// Post a request to transfer funds.
pub async fn transfer_funds(
&self,
base_url: Url,
body: WalletTransferFundsRequestBody,
) -> Result<WalletTransferFundsResponseBody, Error> {
let request_url = base_url
.join(TRANSACTIONS_TRANSFER_FUNDS.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, &body).await
}
/// Post a request to fund a transaction from the node's wallet.
///
/// The node adds fee inputs and change from its own wallet, reserving the
/// request's percentage-based priority-fee reserve, then signs only the
/// appended fee transfer and returns the funded (still unsigned)
/// transaction together with the transfer proof.
pub async fn fund_tx(
&self,
base_url: Url,
body: WalletFundRequestBody,
) -> Result<WalletFundResponseBody, Error> {
let request_url = base_url
.join(FUND.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, &body).await
}
/// Post a request via an SDP declaration to join the blend network and
/// returns its declaration ID if successful.
pub async fn join_blend_network(
&self,
base_url: &Url,
body: JoinBlendRequestBody,
) -> Result<DeclarationId, Error> {
let request_url = base_url
.join(BLEND_JOIN_NETWORK.trim_start_matches('/'))
.map_err(Error::Url)?;
self.post(request_url, &body).await
}
}