Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 3 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ anyhow = "1.0.40"
async-stream = "0.3.5"
async-trait = "0.1.80"
base64 = "0.22.1"
bytes = { version = "1.6", default-features = false }
bytes = { version = "1.7.1", default-features = false }
cid = { version = "0.11.2", default-features = false }
dotenvy = "0.15"
futures = "0.3.30"
Expand All @@ -57,6 +57,7 @@ gloo-timers = { version = "0.3", features = ["futures"] }
hex = "0.4.3"
http = "1"
http-body = "1"
hyper-util = "0.1.19"
js-sys = "0.3.77"
jsonrpsee = "0.26"
jsonrpsee-core = "0.26"
Expand All @@ -75,6 +76,7 @@ rand = "0.8.5"
redb = "2.6"
rexie = "0.6"
rust-embed = "8"
rustls-pki-types = "1.11"
send_wrapper = { version = "0.6", features = ["futures"] }
serde = "1.0.194"
serde_json = "1.0.142"
Expand Down
9 changes: 7 additions & 2 deletions client/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,15 @@ pub mod api {

/// Fibre API related types.
pub mod fibre {
#[cfg(target_arch = "wasm32")]
pub use celestia_fibre::BrowserWebSocketConnector;
#[cfg(not(target_arch = "wasm32"))]
pub use celestia_fibre::NativeTcpConnector;
#[doc(inline)]
pub use celestia_fibre::{
Blob as FibreBlob, BlobID, DownloadOptions, FibreClient, FibreClientConfig, FibreError,
PaymentPromise, SignedPaymentPromise,
Blob as FibreBlob, BlobID, BoxedFibreIo, DownloadOptions, FibreClient,
FibreClientConfig, FibreError, FibreIo, FibreIoConnector, PaymentPromise,
SignedPaymentPromise,
};
}

Expand Down
19 changes: 18 additions & 1 deletion fibre/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,14 @@ futures.workspace = true
lumina-utils = { workspace = true, features = ["executor"] }

# gRPC
der = { version = "0.7.2", features = ["derive"] }
hyper = { version = "1.8", features = ["client", "http2"] }
hyper-util = { workspace = true, features = ["tokio"] }
http.workspace = true
tokio-rustls = { version = "0.26.1", default-features = false, features = ["logging", "ring"] }
tonic = { workspace = true, features = ["codegen"] }
tower = { version = "0.5.2", features = ["util"] }
x509-parser = "0.17"
prost.workspace = true
tendermint-proto.workspace = true

Expand All @@ -41,7 +48,17 @@ async-trait.workspace = true
bech32 = "0.11"

[target.'cfg(not(target_arch = "wasm32"))'.dependencies]
tonic = { workspace = true, features = ["transport"] }
tokio = { workspace = true, features = ["net"] }

[target.'cfg(target_arch = "wasm32")'.dependencies]
gloo-net = { version = "0.6", default-features = false, features = ["websocket"] }
ring = { version = "0.17.14", features = ["wasm32_unknown_unknown_js"] }
rustls-pki-types = { workspace = true, features = ["web"] }
send_wrapper.workspace = true

[dev-dependencies]
serde = { workspace = true, features = ["derive"] }
serde_json.workspace = true

[target.'cfg(not(target_arch = "wasm32"))'.dev-dependencies]
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "test-util"] }
52 changes: 50 additions & 2 deletions fibre/src/client/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,24 +69,72 @@ impl FibreClient {
pub fn from_endpoint(
endpoint: impl Into<celestia_grpc::Endpoint>,
config: FibreClientConfig,
) -> Result<Self, FibreError> {
#[cfg(target_arch = "wasm32")]
{
let _ = (endpoint, config);
Err(FibreError::IoConnectorRequired)
}

#[cfg(not(target_arch = "wasm32"))]
{
let grpc_client = celestia_grpc::GrpcClient::builder()
.endpoint(endpoint)
.build()
.map_err(|e| FibreError::Other(format!("failed to build GrpcClient: {e}")))?;
Self::from_grpc_client(grpc_client, config)
}
}

/// Build a [`FibreClient`] from a single gRPC endpoint and byte-stream transport.
pub fn from_endpoint_with_io_connector(
endpoint: impl Into<celestia_grpc::Endpoint>,
config: FibreClientConfig,
io_connector: Arc<dyn crate::transport::io_connector::FibreIoConnector>,
) -> Result<Self, FibreError> {
let grpc_client = celestia_grpc::GrpcClient::builder()
.endpoint(endpoint)
.build()
.map_err(|e| FibreError::Other(format!("failed to build GrpcClient: {e}")))?;

Self::from_grpc_client(grpc_client, config)
Self::from_grpc_client_with_io_connector(grpc_client, config, io_connector)
}

/// Build a [`FibreClient`] from an existing [`celestia_grpc::GrpcClient`].
pub fn from_grpc_client(
grpc_client: celestia_grpc::GrpcClient,
config: FibreClientConfig,
) -> Result<Self, FibreError> {
#[cfg(target_arch = "wasm32")]
{
let _ = (grpc_client, config);
Err(FibreError::IoConnectorRequired)
}

#[cfg(not(target_arch = "wasm32"))]
{
Self::from_grpc_client_with_io_connector(
grpc_client,
config,
Arc::new(crate::transport::io_connector::NativeTcpConnector),
)
}
}

/// Build a [`FibreClient`] from an existing gRPC client and byte-stream transport.
pub fn from_grpc_client_with_io_connector(
grpc_client: celestia_grpc::GrpcClient,
config: FibreClientConfig,
io_connector: Arc<dyn crate::transport::io_connector::FibreIoConnector>,
) -> Result<Self, FibreError> {
let host_registry = Arc::new(crate::host_registry::GrpcHostRegistry::new(
grpc_client.clone(),
));
let connector = crate::grpc_validator_client::GrpcValidatorConnector::new(host_registry);
let connector = crate::grpc_validator_client::GrpcValidatorConnector::new_with_io_connector(
host_registry,
config.chain_id.clone(),
io_connector,
);

Self::builder()
.config(config)
Expand Down
4 changes: 4 additions & 0 deletions fibre/src/domain/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,10 @@ pub enum FibreError {
#[error("operation cancelled")]
Cancelled,

/// WASM transports require an explicit byte-stream connector.
#[error("a Fibre I/O connector is required on WASM")]
IoConnectorRequired,

// -- Wrapped errors --
/// A gRPC status error from tonic (per-validator connections).
#[error("gRPC error: {0}")]
Expand Down
6 changes: 5 additions & 1 deletion fibre/src/domain/payment_promise.rs
Original file line number Diff line number Diff line change
Expand Up @@ -249,7 +249,11 @@ pub struct SignedPaymentPromise {
const COMET_RAW_BYTES_PREFIX: &[u8] = b"COMET::RAW_BYTES::SIGN";

/// Construct the CometBFT domain-separated sign bytes for raw byte messages.
fn raw_bytes_message_sign_bytes(chain_id: &str, unique_id: &[u8], raw_bytes: &[u8]) -> Vec<u8> {
pub(crate) fn raw_bytes_message_sign_bytes(
chain_id: &str,
unique_id: &[u8],
raw_bytes: &[u8],
) -> Vec<u8> {
use celestia_proto::tendermint_celestia_mods::privval::SignRawBytesRequest;
use prost::Message;

Expand Down
5 changes: 5 additions & 0 deletions fibre/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,5 +33,10 @@ pub use domain::payment_promise::{PaymentPromise, SignedPaymentPromise};
pub use error::{FibreError, Result};
pub use transport::grpc_validator_client::GrpcValidatorConnector;
pub use transport::host_registry::{GrpcHostRegistry, Host, HostRegistry};
#[cfg(target_arch = "wasm32")]
pub use transport::io_connector::BrowserWebSocketConnector;
#[cfg(not(target_arch = "wasm32"))]
pub use transport::io_connector::NativeTcpConnector;
pub use transport::io_connector::{BoxedFibreIo, FibreIo, FibreIoConnector};
pub use transport::validator_client::{DownloadResponse, ValidatorConnection, ValidatorConnector};
pub use validator::{GrpcSetGetter, SetGetter, ValidatorInfo, ValidatorSet};
68 changes: 42 additions & 26 deletions fibre/src/transport/grpc_validator_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ use crate::error::FibreError;
use crate::host_registry::HostRegistry;
use crate::payment_promise::PaymentPromise;
use crate::proto_conv;
use crate::transport::io_connector::FibreIoConnector;
use crate::validator::ValidatorInfo;
use crate::validator_client::{
DownloadResponse, UploadResponse, ValidatorConnection, ValidatorConnector,
Expand All @@ -23,14 +24,32 @@ use crate::validator_client::{
/// Factory that resolves validator hosts and caches gRPC connections.
pub struct GrpcValidatorConnector {
host_registry: Arc<dyn HostRegistry>,
chain_id: String,
io_connector: Arc<dyn FibreIoConnector>,
connections: tokio::sync::Mutex<HashMap<[u8; 20], Arc<GrpcValidatorConnection>>>,
}

impl GrpcValidatorConnector {
/// Create a new connector backed by the given host registry.
pub fn new(host_registry: Arc<dyn HostRegistry>) -> Self {
#[cfg(not(target_arch = "wasm32"))]
pub fn new(host_registry: Arc<dyn HostRegistry>, chain_id: impl Into<String>) -> Self {
Self::new_with_io_connector(
host_registry,
chain_id,
Arc::new(crate::transport::io_connector::NativeTcpConnector),
)
}

/// Create a new connector over a caller-provided byte-stream transport.
pub fn new_with_io_connector(
host_registry: Arc<dyn HostRegistry>,
chain_id: impl Into<String>,
io_connector: Arc<dyn FibreIoConnector>,
) -> Self {
Self {
host_registry,
chain_id: chain_id.into(),
io_connector,
connections: tokio::sync::Mutex::new(HashMap::new()),
}
}
Expand Down Expand Up @@ -59,10 +78,12 @@ impl ValidatorConnector for GrpcValidatorConnector {
// URI, so normalise the address first.
let url = normalize_host(&host.0);

let client = GrpcClient::builder()
.url(&url)
.build()
.map_err(|e| FibreError::Other(format!("invalid endpoint '{}': {e}", url)))?;
let client = crate::transport::tls::grpc_client(
&url,
validator.pubkey,
self.chain_id.clone(),
self.io_connector.clone(),
)?;

let conn = Arc::new(GrpcValidatorConnection { client });

Expand All @@ -75,21 +96,16 @@ impl ValidatorConnector for GrpcValidatorConnector {
}
}

/// Normalise a host string into a standard `http://` URI.
/// Normalise a host string into a standard `https://` URI.
fn normalize_host(raw: &str) -> String {
// Strip the "dns:///" (or "dns://") prefix if present.
if let Some(rest) = raw
let authority = raw
.strip_prefix("dns:///")
.or_else(|| raw.strip_prefix("dns://"))
{
return format!("http://{rest}");
}
// Already a normal URI (http:// or https://).
if raw.starts_with("http://") || raw.starts_with("https://") {
return raw.to_string();
}
// Bare host:port — assume http.
format!("http://{raw}")
.or_else(|| raw.strip_prefix("http://"))
.or_else(|| raw.strip_prefix("https://"))
.unwrap_or(raw);

format!("https://{authority}")
}

/// A connection to a single validator's Fibre gRPC service.
Expand Down Expand Up @@ -181,7 +197,7 @@ mod tests {
call_count: AtomicUsize::new(0),
});

let connector = GrpcValidatorConnector::new(registry.clone());
let connector = GrpcValidatorConnector::new(registry.clone(), "test-chain");

// First connect should call the registry.
let conn1 = connector.connect(&validator).await;
Expand Down Expand Up @@ -219,7 +235,7 @@ mod tests {
call_count: AtomicUsize::new(0),
});

let connector = GrpcValidatorConnector::new(registry.clone());
let connector = GrpcValidatorConnector::new(registry.clone(), "test-chain");

let conn_a = connector.connect(&validator_a).await;
assert!(conn_a.is_ok(), "connect to validator A should succeed");
Expand All @@ -245,7 +261,7 @@ mod tests {
call_count: AtomicUsize::new(0),
});

let connector = GrpcValidatorConnector::new(registry.clone());
let connector = GrpcValidatorConnector::new(registry.clone(), "test-chain");

let result = connector.connect(&validator).await;
match result {
Expand All @@ -265,19 +281,19 @@ mod tests {
fn normalize_host_strips_dns_prefix() {
assert_eq!(
normalize_host("dns:///138.68.236.99:9091"),
"http://138.68.236.99:9091"
"https://138.68.236.99:9091"
);
assert_eq!(
normalize_host("dns://138.68.236.99:9091"),
"http://138.68.236.99:9091"
"https://138.68.236.99:9091"
);
}

#[test]
fn normalize_host_preserves_http() {
fn normalize_host_uses_https() {
assert_eq!(
normalize_host("http://127.0.0.1:9090"),
"http://127.0.0.1:9090"
"https://127.0.0.1:9090"
);
assert_eq!(
normalize_host("https://validator.example.com:9090"),
Expand All @@ -286,10 +302,10 @@ mod tests {
}

#[test]
fn normalize_host_adds_http_to_bare() {
fn normalize_host_adds_https_to_bare() {
assert_eq!(
normalize_host("138.68.236.99:9091"),
"http://138.68.236.99:9091"
"https://138.68.236.99:9091"
);
}
}
Loading
Loading