diff --git a/Cargo.lock b/Cargo.lock index 3bb6c1f..40ef768 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -80,6 +80,51 @@ version = "1.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ace50bade8e6234aa140d9a2f552bbee1db4d353f69b8217bc503490fc1a9f26" +[[package]] +name = "axum" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "021e862c184ae977658b36c4500f7feac3221ca5da43e3f25bd04ab6c79a29b5" +dependencies = [ + "axum-core", + "bytes", + "futures-util", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "itoa", + "matchit", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "rustversion", + "serde", + "sync_wrapper", + "tower 0.5.3", + "tower-layer", + "tower-service", +] + +[[package]] +name = "axum-core" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68464cd0412f486726fb3373129ef5d2993f90c34bc2bc1c1e9943b2f4fc7ca6" +dependencies = [ + "bytes", + "futures-core", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "mime", + "pin-project-lite", + "rustversion", + "sync_wrapper", + "tower-layer", + "tower-service", +] + [[package]] name = "base-x" version = "0.2.11" @@ -229,8 +274,8 @@ dependencies = [ [[package]] name = "celestia-proto" -version = "0.12.1" -source = "git+https://github.com/celestiaorg/lumina.git?rev=e114c6f#e114c6ffd605ae170c6013e9ef68f3a946404722" +version = "1.0.0-rc.2" +source = "git+https://github.com/celestiaorg/lumina.git?rev=6bf264d#6bf264dfc027da4440e4d0377aab156a39b1b288" dependencies = [ "bytes", "prost", @@ -240,12 +285,14 @@ dependencies = [ "serde", "subtle-encoding", "tendermint-proto", + "tonic", + "tonic-build", ] [[package]] name = "celestia-rpc" -version = "0.16.2" -source = "git+https://github.com/celestiaorg/lumina.git?rev=e114c6f#e114c6ffd605ae170c6013e9ef68f3a946404722" +version = "1.0.0-rc.2" +source = "git+https://github.com/celestiaorg/lumina.git?rev=6bf264d#6bf264dfc027da4440e4d0377aab156a39b1b288" dependencies = [ "async-stream", "base64", @@ -272,8 +319,8 @@ dependencies = [ [[package]] name = "celestia-types" -version = "0.20.0" -source = "git+https://github.com/celestiaorg/lumina.git?rev=e114c6f#e114c6ffd605ae170c6013e9ef68f3a946404722" +version = "1.0.0-rc.2" +source = "git+https://github.com/celestiaorg/lumina.git?rev=6bf264d#6bf264dfc027da4440e4d0377aab156a39b1b288" dependencies = [ "base64", "bech32", @@ -1035,6 +1082,19 @@ dependencies = [ "tower-service", ] +[[package]] +name = "hyper-timeout" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2b90d566bffbce6a75bd8b09a05aa8c2cb1fabb6cb348f8840c9e4c90a0d83b0" +dependencies = [ + "hyper", + "hyper-util", + "pin-project-lite", + "tokio", + "tower-service", +] + [[package]] name = "hyper-util" version = "0.1.20" @@ -1049,7 +1109,7 @@ dependencies = [ "hyper", "libc", "pin-project-lite", - "socket2", + "socket2 0.6.2", "tokio", "tower-service", "tracing", @@ -1475,6 +1535,7 @@ version = "0.1.0" dependencies = [ "anyhow", "async-trait", + "celestia-proto", "celestia-rpc", "celestia-types", "ed25519-consensus", @@ -1483,14 +1544,17 @@ dependencies = [ "jsonrpsee-core", "jsonrpsee-types", "once_cell", + "prost", "rand 0.8.5", "redis", "serde", "serde_json", "sha2 0.10.8", "tendermint", + "tendermint-proto", "thiserror 1.0.69", "tokio", + "tonic", "tower 0.4.13", "tower-http", "tracing", @@ -1550,12 +1614,18 @@ dependencies = [ [[package]] name = "lumina-utils" -version = "0.5.2" -source = "git+https://github.com/celestiaorg/lumina.git?rev=e114c6f#e114c6ffd605ae170c6013e9ef68f3a946404722" +version = "1.0.0-rc.2" +source = "git+https://github.com/celestiaorg/lumina.git?rev=6bf264d#6bf264dfc027da4440e4d0377aab156a39b1b288" dependencies = [ "js-sys", ] +[[package]] +name = "matchit" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3" + [[package]] name = "memchr" version = "2.7.4" @@ -1585,6 +1655,12 @@ dependencies = [ "syn", ] +[[package]] +name = "mime" +version = "0.3.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" + [[package]] name = "mio" version = "1.0.3" @@ -2049,7 +2125,7 @@ dependencies = [ "pin-project-lite", "ryu", "sha1_smol", - "socket2", + "socket2 0.6.2", "tokio", "tokio-util", "url", @@ -2487,6 +2563,16 @@ version = "1.14.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7fcf8323ef1faaee30a44a340193b1ac6814fd9b7b4e88e9d4519a3e4abe1cfd" +[[package]] +name = "socket2" +version = "0.5.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e22376abed350d73dd1cd119b57ffccad95b4e585a7cda43e286245ce23c0678" +dependencies = [ + "libc", + "windows-sys 0.52.0", +] + [[package]] name = "socket2" version = "0.6.2" @@ -2765,7 +2851,7 @@ dependencies = [ "parking_lot", "pin-project-lite", "signal-hook-registry", - "socket2", + "socket2 0.6.2", "tokio-macros", "windows-sys 0.61.2", ] @@ -2834,6 +2920,49 @@ dependencies = [ "winnow", ] +[[package]] +name = "tonic" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7e581ba15a835f4d9ea06c55ab1bd4dce26fc53752c69a04aac00703bfb49ba9" +dependencies = [ + "async-trait", + "axum", + "base64", + "bytes", + "h2", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "hyper", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "prost", + "socket2 0.5.10", + "tokio", + "tokio-stream", + "tower 0.5.3", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "tonic-build" +version = "0.13.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eac6f67be712d12f0b41328db3137e0d0757645d8904b4cb7d51cd9c2279e847" +dependencies = [ + "prettyplease", + "proc-macro2", + "prost-build", + "prost-types", + "quote", + "syn", +] + [[package]] name = "tower" version = "0.4.13" @@ -2853,10 +2982,15 @@ checksum = "ebe5ef63511595f1344e2d5cfa636d973292adc0eec1f0ad45fae9f0851ab1d4" dependencies = [ "futures-core", "futures-util", + "indexmap", "pin-project-lite", + "slab", "sync_wrapper", + "tokio", + "tokio-util", "tower-layer", "tower-service", + "tracing", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 88a92f5..ffad985 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,8 +7,12 @@ authors = ["Your Name "] license = "MIT OR Apache-2.0" [dependencies] -celestia-rpc = { git = "https://github.com/celestiaorg/lumina.git", rev = "e114c6f" } -celestia-types = { git = "https://github.com/celestiaorg/lumina.git", rev = "e114c6f" } +celestia-rpc = { git = "https://github.com/celestiaorg/lumina.git", rev = "6bf264d" } +celestia-types = { git = "https://github.com/celestiaorg/lumina.git", rev = "6bf264d" } +celestia-proto = { git = "https://github.com/celestiaorg/lumina.git", rev = "6bf264d", features = ["tonic", "tonic-server"] } +tonic = { version = "0.13", default-features = false, features = ["codegen", "prost", "transport", "router"] } +prost = "0.13" +tendermint-proto = "0.40.3" tendermint = "0.40.3" ed25519-consensus = "2.1.0" rand = "0.8.5" diff --git a/README.md b/README.md index 7c85ba5..acc40d7 100644 --- a/README.md +++ b/README.md @@ -33,6 +33,7 @@ Defaults if unset: - `REDIS_URL=redis://127.0.0.1:6379` - `LISTEN_ADDR=127.0.0.1:26658` +- `GRPC_ADDR=0.0.0.0:9090` Localestia clears the Redis database on startup. @@ -52,6 +53,29 @@ LOCALESTIA_REDIS_MODE=docker cargo test You can also use `LOCALESTIA_REDIS_MODE=auto` to prefer a local `REDIS_URL` if set and fall back to Docker when available. +### gRPC integration tests + +The `grpc` test suite starts a real localestia process on ephemeral ports and exercises all five gRPC services via generated tonic client stubs: + +```bash +# With a local Redis already running: +LOCALESTIA_REDIS_MODE=local cargo test --test grpc -- --test-threads=1 + +# Or let the test harness spin up a Redis container: +LOCALESTIA_REDIS_MODE=docker cargo test --test grpc -- --test-threads=1 +``` + +`--test-threads=1` is required because each test acquires a global lock (one localestia process at a time). The tests cover: + +| Test | gRPC service | What is verified | +|---|---|---| +| `test_get_node_info` | `cosmos.base.tendermint.v1beta1.Service` | moniker=`localestia`, network=`private` | +| `test_account` | `cosmos.auth.v1beta1.Query` | account_number=1, sequence=0 | +| `test_broadcast_tx_empty` | `cosmos.tx.v1beta1.Service` | code=0, height=0 for an empty tx | +| `test_tx_status_unknown_returns_committed` | `celestia.core.v1.tx.Tx` | status=`COMMITTED`, height=1 for unknown tx | +| `test_estimate_gas_price` | `celestia.core.v1.gas_estimation.GasEstimator` | price=0.002 | +| `test_estimate_gas_price_and_usage` | `celestia.core.v1.gas_estimation.GasEstimator` | price=0.002, gas_used=500 000 | + ## Supported RPC Methods Localestia implements the following JSON-RPC methods: diff --git a/src/grpc/auth.rs b/src/grpc/auth.rs new file mode 100644 index 0000000..b6df141 --- /dev/null +++ b/src/grpc/auth.rs @@ -0,0 +1,122 @@ +use std::sync::Arc; + +use celestia_proto::cosmos::auth::v1beta1::query_server::{Query, QueryServer}; +use celestia_proto::cosmos::auth::v1beta1::{ + AddressBytesToStringRequest, AddressBytesToStringResponse, AddressStringToBytesRequest, + AddressStringToBytesResponse, BaseAccount, Bech32PrefixRequest, Bech32PrefixResponse, + QueryAccountAddressByIdRequest, QueryAccountAddressByIdResponse, QueryAccountInfoRequest, + QueryAccountInfoResponse, QueryAccountRequest, QueryAccountResponse, QueryAccountsRequest, + QueryAccountsResponse, QueryModuleAccountByNameRequest, QueryModuleAccountByNameResponse, + QueryModuleAccountsRequest, QueryModuleAccountsResponse, QueryParamsRequest, + QueryParamsResponse, +}; +use prost::Message; +use tendermint_proto::google::protobuf::Any; +use tonic::{Request, Response, Status}; + +use crate::storage::RedisStorage; + +/// gRPC service: cosmos.auth.v1beta1.Query +/// +/// Handles the Account query that txclient calls to get sequence/account_number. +#[derive(Clone)] +pub struct AuthQueryService { + _storage: Arc, +} + +impl AuthQueryService { + pub fn new(storage: Arc) -> Self { + Self { _storage: storage } + } +} + +#[tonic::async_trait] +impl Query for AuthQueryService { + async fn account( + self: Arc, + request: Request, + ) -> Result, Status> { + let address = request.into_inner().address; + + let account = BaseAccount { + address, + pub_key: None, + account_number: 1, + sequence: 0, + }; + + let any = Any { + type_url: "/cosmos.auth.v1beta1.BaseAccount".to_string(), + value: account.encode_to_vec(), + }; + + Ok(Response::new(QueryAccountResponse { account: Some(any) })) + } + + async fn accounts( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn account_address_by_id( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn params( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn module_accounts( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn module_account_by_name( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn bech32_prefix( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn address_bytes_to_string( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn address_string_to_bytes( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn account_info( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } +} + +pub fn service(storage: Arc) -> QueryServer { + QueryServer::new(AuthQueryService::new(storage)) +} diff --git a/src/grpc/gas_estimator.rs b/src/grpc/gas_estimator.rs new file mode 100644 index 0000000..3ac8125 --- /dev/null +++ b/src/grpc/gas_estimator.rs @@ -0,0 +1,64 @@ +use std::sync::Arc; + +use celestia_proto::celestia::core::v1::gas_estimation::gas_estimator_server::{ + GasEstimator, GasEstimatorServer, +}; +use celestia_proto::celestia::core::v1::gas_estimation::{ + EstimateGasPriceAndUsageRequest, EstimateGasPriceAndUsageResponse, EstimateGasPriceRequest, + EstimateGasPriceResponse, +}; +use tonic::{Request, Response, Status}; +use tracing::info; + +use crate::storage::RedisStorage; + +// Fixed gas price returned to txclient (utia per gas unit) +const FIXED_GAS_PRICE: f64 = 0.002; +// Fixed gas estimate for blob transactions +const FIXED_GAS_USED: u64 = 500_000; + +/// gRPC service: celestia.core.v1.gas_estimation.GasEstimator +/// +/// Returns fixed gas price and usage estimates so the txclient can +/// construct a signed transaction without hitting a real Celestia node. +#[derive(Clone)] +pub struct GasEstimatorService { + _storage: Arc, +} + +impl GasEstimatorService { + pub fn new(storage: Arc) -> Self { + Self { _storage: storage } + } +} + +#[tonic::async_trait] +impl GasEstimator for GasEstimatorService { + async fn estimate_gas_price( + self: Arc, + _request: Request, + ) -> Result, Status> { + info!("EstimateGasPrice: returning fixed price {}", FIXED_GAS_PRICE); + Ok(Response::new(EstimateGasPriceResponse { + estimated_gas_price: FIXED_GAS_PRICE, + })) + } + + async fn estimate_gas_price_and_usage( + self: Arc, + _request: Request, + ) -> Result, Status> { + info!( + "EstimateGasPriceAndUsage: returning price={} gas={}", + FIXED_GAS_PRICE, FIXED_GAS_USED + ); + Ok(Response::new(EstimateGasPriceAndUsageResponse { + estimated_gas_price: FIXED_GAS_PRICE, + estimated_gas_used: FIXED_GAS_USED, + })) + } +} + +pub fn service(storage: Arc) -> GasEstimatorServer { + GasEstimatorServer::new(GasEstimatorService::new(storage)) +} diff --git a/src/grpc/mod.rs b/src/grpc/mod.rs new file mode 100644 index 0000000..35aa1e9 --- /dev/null +++ b/src/grpc/mod.rs @@ -0,0 +1,29 @@ +use std::sync::Arc; + +use tonic::transport::Server; +use tracing::info; + +use crate::storage::RedisStorage; + +mod auth; +mod gas_estimator; +mod node_info; +mod tx_service; +mod tx_status; + +pub async fn serve( + storage: Arc, + addr: std::net::SocketAddr, + shutdown: impl std::future::Future, +) -> Result<(), tonic::transport::Error> { + info!("Starting gRPC server on {}", addr); + + Server::builder() + .add_service(node_info::service(storage.clone())) + .add_service(auth::service(storage.clone())) + .add_service(tx_service::service(storage.clone())) + .add_service(tx_status::service(storage.clone())) + .add_service(gas_estimator::service(storage.clone())) + .serve_with_shutdown(addr, shutdown) + .await +} diff --git a/src/grpc/node_info.rs b/src/grpc/node_info.rs new file mode 100644 index 0000000..4d02764 --- /dev/null +++ b/src/grpc/node_info.rs @@ -0,0 +1,102 @@ +use std::sync::Arc; + +use celestia_proto::cosmos::base::tendermint::v1beta1::service_server::{ + Service as TendermintService, ServiceServer as TendermintServiceServer, +}; +use celestia_proto::cosmos::base::tendermint::v1beta1::{ + AbciQueryRequest, AbciQueryResponse, GetBlockByHeightRequest, GetBlockByHeightResponse, + GetLatestBlockRequest, GetLatestBlockResponse, GetLatestValidatorSetRequest, + GetLatestValidatorSetResponse, GetNodeInfoRequest, GetNodeInfoResponse, GetSyncingRequest, + GetSyncingResponse, GetValidatorSetByHeightRequest, GetValidatorSetByHeightResponse, + VersionInfo, +}; +use tendermint_proto::v0_38::p2p::DefaultNodeInfo; +use tonic::{Request, Response, Status}; + +use crate::storage::RedisStorage; + +/// gRPC service: cosmos.base.tendermint.v1beta1.Service +/// +/// Only handles GetNodeInfo, which txclient calls during initialization. +#[derive(Clone)] +pub struct NodeInfoService { + _storage: Arc, +} + +impl NodeInfoService { + pub fn new(storage: Arc) -> Self { + Self { _storage: storage } + } +} + +#[tonic::async_trait] +impl TendermintService for NodeInfoService { + async fn get_node_info( + self: Arc, + _request: Request, + ) -> Result, Status> { + Ok(Response::new(GetNodeInfoResponse { + default_node_info: Some(DefaultNodeInfo { + network: "private".to_string(), + moniker: "localestia".to_string(), + ..Default::default() + }), + application_version: Some(VersionInfo { + name: "localestia".to_string(), + app_name: "localestia".to_string(), + version: "0.1.0".to_string(), + git_commit: String::new(), + build_tags: String::new(), + go_version: String::new(), + build_deps: vec![], + cosmos_sdk_version: String::new(), + }), + })) + } + + async fn get_syncing( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_latest_block( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_block_by_height( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_latest_validator_set( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_validator_set_by_height( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn abci_query( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } +} + +pub fn service(storage: Arc) -> TendermintServiceServer { + TendermintServiceServer::new(NodeInfoService::new(storage)) +} diff --git a/src/grpc/tx_service.rs b/src/grpc/tx_service.rs new file mode 100644 index 0000000..c64c0ba --- /dev/null +++ b/src/grpc/tx_service.rs @@ -0,0 +1,185 @@ +use std::sync::Arc; + +use celestia_proto::cosmos::base::abci::v1beta1::TxResponse; +use celestia_proto::cosmos::tx::v1beta1::service_server::{Service, ServiceServer}; +use celestia_proto::cosmos::tx::v1beta1::{ + BroadcastTxRequest, BroadcastTxResponse, GetBlockWithTxsRequest, GetBlockWithTxsResponse, + GetTxRequest, GetTxResponse, GetTxsEventRequest, GetTxsEventResponse, SimulateRequest, + SimulateResponse, TxDecodeAminoRequest, TxDecodeAminoResponse, TxDecodeRequest, + TxDecodeResponse, TxEncodeAminoRequest, TxEncodeAminoResponse, TxEncodeRequest, + TxEncodeResponse, +}; +use celestia_proto::proto::blob::v2::BlobTx; +use celestia_types::consts::appconsts::AppVersion; +use celestia_types::nmt::Namespace; +use celestia_types::Blob; +use prost::Message; +use sha2::{Digest, Sha256}; +use tonic::{Request, Response, Status}; +use tracing::{error, info}; + +use crate::storage::RedisStorage; + +/// gRPC service: cosmos.tx.v1beta1.Service +/// +/// Handles BroadcastTx — decodes BlobTx, stores blobs in Redis, returns tx hash. +#[derive(Clone)] +pub struct TxService { + storage: Arc, +} + +impl TxService { + pub fn new(storage: Arc) -> Self { + Self { storage } + } +} + +#[tonic::async_trait] +impl Service for TxService { + async fn broadcast_tx( + self: Arc, + request: Request, + ) -> Result, Status> { + let tx_bytes = request.into_inner().tx_bytes; + + // Compute tx hash (SHA-256 of raw tx bytes) + let tx_hash = { + let mut hasher = Sha256::new(); + hasher.update(&tx_bytes); + hex::encode(hasher.finalize()).to_uppercase() + }; + + info!("BroadcastTx: tx_hash={}", tx_hash); + + // Try to decode as BlobTx (go-square proto.blob.v2.BlobTx) + let blobs = match BlobTx::decode(tx_bytes.as_ref()) { + Ok(blob_tx) if !blob_tx.blobs.is_empty() => { + info!("Decoded BlobTx with {} blobs", blob_tx.blobs.len()); + let mut celestia_blobs = Vec::new(); + for blob_proto in blob_tx.blobs { + let ns_version = + u8::try_from(blob_proto.namespace_version).map_err(|_| { + Status::invalid_argument(format!( + "invalid namespace version: {}", + blob_proto.namespace_version + )) + })?; + let ns = Namespace::new(ns_version, &blob_proto.namespace_id) + .map_err(|e| Status::invalid_argument(format!("invalid namespace: {e}")))?; + + match Blob::new(ns, blob_proto.data, None, AppVersion::V2) { + Ok(blob) => celestia_blobs.push(blob), + Err(e) => { + error!("Failed to construct Blob: {}", e); + return Err(Status::invalid_argument(format!("invalid blob: {e}"))); + } + } + } + celestia_blobs + } + _ => { + // No blobs (or not a BlobTx) — treat as plain cosmos tx, nothing to store + info!("BroadcastTx: no blobs to store"); + vec![] + } + }; + + let tx_height = if !blobs.is_empty() { + let height = self.storage.store_blobs(&blobs).await.map_err(|e| { + error!("Failed to store blobs: {}", e); + Status::internal(format!("storage error: {e}")) + })?; + info!("Stored {} blobs at height {}", blobs.len(), height); + self.storage + .store_tx_height(&tx_hash, height) + .await + .map_err(|e| { + error!("Failed to store tx height mapping: {}", e); + Status::internal(format!("storage error: {e}")) + })?; + height + } else { + 0 + }; + + let tx_response = TxResponse { + height: tx_height as i64, + txhash: tx_hash, + codespace: String::new(), + code: 0, + data: String::new(), + raw_log: String::new(), + logs: vec![], + info: String::new(), + gas_wanted: 0, + gas_used: 0, + tx: None, + timestamp: String::new(), + events: vec![], + }; + + Ok(Response::new(BroadcastTxResponse { + tx_response: Some(tx_response), + })) + } + + async fn simulate( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_tx( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_txs_event( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn get_block_with_txs( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn tx_decode( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn tx_encode( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn tx_encode_amino( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } + + async fn tx_decode_amino( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } +} + +pub fn service(storage: Arc) -> ServiceServer { + ServiceServer::new(TxService::new(storage)) +} diff --git a/src/grpc/tx_status.rs b/src/grpc/tx_status.rs new file mode 100644 index 0000000..7e50866 --- /dev/null +++ b/src/grpc/tx_status.rs @@ -0,0 +1,68 @@ +use std::sync::Arc; + +use celestia_proto::celestia::core::v1::tx::tx_server::{Tx, TxServer}; +use celestia_proto::celestia::core::v1::tx::{ + TxStatusBatchRequest, TxStatusBatchResponse, TxStatusRequest, TxStatusResponse, +}; +use tonic::{Request, Response, Status}; +use tracing::{error, info}; + +use crate::storage::RedisStorage; + +/// gRPC service: celestia.core.v1.tx.Tx +/// +/// Handles TxStatus — looks up the stored height for the tx and returns COMMITTED. +#[derive(Clone)] +pub struct TxStatusService { + storage: Arc, +} + +impl TxStatusService { + pub fn new(storage: Arc) -> Self { + Self { storage } + } +} + +#[tonic::async_trait] +impl Tx for TxStatusService { + async fn tx_status( + self: Arc, + request: Request, + ) -> Result, Status> { + let tx_id = request.into_inner().tx_id; + info!("TxStatus: tx_id={}", tx_id); + + let height = match self.storage.get_tx_height(&tx_id).await { + Ok(Some(h)) => h, + Ok(None) => 1, + Err(e) => { + error!("TxStatus: storage error for tx_id={}: {}", tx_id, e); + return Err(Status::internal(format!("storage error: {e}"))); + } + }; + info!("TxStatus: resolved height={}", height); + + Ok(Response::new(TxStatusResponse { + height: height as i64, + index: 0, + execution_code: 0, + error: String::new(), + status: "COMMITTED".to_string(), + codespace: String::new(), + gas_wanted: 0, + gas_used: 0, + signers: vec![], + })) + } + + async fn tx_status_batch( + self: Arc, + _request: Request, + ) -> Result, Status> { + Err(Status::unimplemented("not supported by localestia")) + } +} + +pub fn service(storage: Arc) -> TxServer { + TxServer::new(TxStatusService::new(storage)) +} diff --git a/src/main.rs b/src/main.rs index e03d725..b6ecf14 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use tracing::{error, info}; mod error; +mod grpc; mod rpc; mod storage; mod types; @@ -26,8 +27,9 @@ async fn main() -> Result<(), Box> { std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); let listen_addr = std::env::var("LISTEN_ADDR").unwrap_or_else(|_| "127.0.0.1:26658".to_string()); + let grpc_addr = + std::env::var("GRPC_ADDR").unwrap_or_else(|_| "0.0.0.0:9090".to_string()); - // Create a new error variant for AddrParseError let addr: SocketAddr = listen_addr.parse().map_err(|e| { Box::new(LocalError::TransactionError(format!( "Failed to parse address: {}", @@ -35,9 +37,17 @@ async fn main() -> Result<(), Box> { ))) })?; + let grpc_socket: SocketAddr = grpc_addr.parse().map_err(|e| { + Box::new(LocalError::TransactionError(format!( + "Failed to parse gRPC address: {}", + e + ))) + })?; + info!("Starting local Celestia Blob RPC server..."); info!("Redis URL: {}", redis_url); - info!("Listening on: {}", listen_addr); + info!("JSON-RPC listening on: {}", listen_addr); + info!("gRPC listening on: {}", grpc_addr); // Initialize Redis storage let storage = match RedisStorage::new(&redis_url) { @@ -88,33 +98,42 @@ async fn main() -> Result<(), Box> { ))) })?; - // Register our RPC methods (both Blob and Header) let server_handle = server.start(module); - info!("Server started successfully"); - info!( - "The server is ready to accept RPC calls at ws://{}", - listen_addr - ); - info!( - "Connect using: celestia-rpc::Client::new(\"ws://{}\", None)", - listen_addr - ); - - // Keep the server running until Ctrl+C - match tokio::signal::ctrl_c().await { - Ok(_) => info!("Received shutdown signal"), - Err(e) => error!("Failed to listen for ctrl-c: {}", e), + info!("JSON-RPC server started at ws://{}", listen_addr); + + // Start gRPC server concurrently + let (grpc_shutdown_tx, grpc_shutdown_rx) = tokio::sync::oneshot::channel::<()>(); + let grpc_storage = storage.clone(); + let mut grpc_task = tokio::spawn(async move { + if let Err(e) = + grpc::serve(grpc_storage, grpc_socket, async { let _ = grpc_shutdown_rx.await; }).await + { + error!("gRPC server error: {}", e); + } + }); + + // Keep both servers running until Ctrl+C or gRPC task exits + tokio::select! { + _ = tokio::signal::ctrl_c() => { + info!("Received shutdown signal"); + } + _ = &mut grpc_task => { + error!("gRPC server exited unexpectedly"); + } } - info!("Shutting down server..."); + info!("Shutting down..."); + + // Signal gRPC server to stop and wait for it to drain + let _ = grpc_shutdown_tx.send(()); + grpc_task.await.ok(); - // Stop the server if let Err(e) = server_handle.stop() { - error!("Error stopping server: {}", e); + error!("Error stopping JSON-RPC server: {}", e); } - info!("Server shutdown complete"); + info!("Shutdown complete"); Ok(()) } diff --git a/src/storage/redis_storage.rs b/src/storage/redis_storage.rs index 5ccdf16..ad1832e 100644 --- a/src/storage/redis_storage.rs +++ b/src/storage/redis_storage.rs @@ -1213,6 +1213,28 @@ impl RedisStorage { Ok(response) } + pub async fn store_tx_height(&self, tx_hash: &str, height: u64) -> Result<(), LocalError> { + let mut conn = self + .client + .get_multiplexed_async_connection() + .await + .map_err(LocalError::RedisError)?; + let key = format!("tx:height:{}", tx_hash); + let _: () = conn.set(key, height).await.map_err(LocalError::RedisError)?; + Ok(()) + } + + pub async fn get_tx_height(&self, tx_hash: &str) -> Result, LocalError> { + let mut conn = self + .client + .get_multiplexed_async_connection() + .await + .map_err(LocalError::RedisError)?; + let key = format!("tx:height:{}", tx_hash); + let val: Option = conn.get(key).await.map_err(LocalError::RedisError)?; + Ok(val) + } + pub async fn clear_database(&self) -> Result<(), LocalError> { info!("Clearing Redis database"); diff --git a/tests/grpc.rs b/tests/grpc.rs new file mode 100644 index 0000000..c8291b0 --- /dev/null +++ b/tests/grpc.rs @@ -0,0 +1,151 @@ +mod utils; + +use celestia_proto::celestia::core::v1::gas_estimation::{ + gas_estimator_client::GasEstimatorClient, EstimateGasPriceAndUsageRequest, + EstimateGasPriceRequest, +}; +use celestia_proto::celestia::core::v1::tx::{tx_client::TxClient, TxStatusRequest}; +use celestia_proto::cosmos::auth::v1beta1::{ + query_client::QueryClient as AuthQueryClient, QueryAccountRequest, +}; +use celestia_proto::cosmos::base::tendermint::v1beta1::{ + service_client::ServiceClient as TendermintClient, GetNodeInfoRequest, +}; +use celestia_proto::cosmos::tx::v1beta1::{ + service_client::ServiceClient as TxServiceClient, BroadcastTxRequest, +}; +use prost::Message; +use tonic::transport::Channel; + +async fn grpc_channel(grpc_addr: &str) -> Channel { + tonic::transport::Endpoint::new(format!("http://{}", grpc_addr)) + .expect("invalid gRPC addr") + .connect() + .await + .expect("failed to connect to gRPC server") +} + +#[tokio::test] +async fn test_get_node_info() { + let ctx = utils::setup_process_grpc().await; + let mut client = TendermintClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .get_node_info(GetNodeInfoRequest {}) + .await + .expect("GetNodeInfo failed"); + let body = resp.into_inner(); + + let node_info = body.default_node_info.expect("missing default_node_info"); + assert_eq!(node_info.moniker, "localestia"); + assert_eq!(node_info.network, "private"); + + let version = body.application_version.expect("missing application_version"); + assert_eq!(version.app_name, "localestia"); +} + +#[tokio::test] +async fn test_account() { + let ctx = utils::setup_process_grpc().await; + let mut client = AuthQueryClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .account(QueryAccountRequest { + address: "cosmos1testaddress".to_string(), + }) + .await + .expect("Account failed"); + let body = resp.into_inner(); + + let any = body.account.expect("missing account Any"); + assert_eq!( + any.type_url, + "/cosmos.auth.v1beta1.BaseAccount", + "unexpected type_url" + ); + + use celestia_proto::cosmos::auth::v1beta1::BaseAccount; + let base = BaseAccount::decode(any.value.as_slice()).expect("failed to decode BaseAccount"); + assert_eq!(base.account_number, 1); + assert_eq!(base.sequence, 0); +} + +#[tokio::test] +async fn test_broadcast_tx_empty() { + let ctx = utils::setup_process_grpc().await; + let mut client = TxServiceClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .broadcast_tx(BroadcastTxRequest { + tx_bytes: vec![], + mode: 0, + }) + .await + .expect("BroadcastTx failed"); + let body = resp.into_inner(); + + let tx_resp = body.tx_response.expect("missing tx_response"); + assert_eq!(tx_resp.code, 0); + // Empty tx has no blobs → height 0 + assert_eq!(tx_resp.height, 0); +} + +#[tokio::test] +async fn test_tx_status_unknown_returns_committed() { + let ctx = utils::setup_process_grpc().await; + let mut client = TxClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .tx_status(TxStatusRequest { + tx_id: "deadbeef1234".to_string(), + }) + .await + .expect("TxStatus failed"); + let body = resp.into_inner(); + + // Unknown tx_id falls back to height 1, status always COMMITTED + assert_eq!(body.status, "COMMITTED"); + assert_eq!(body.height, 1); +} + +#[tokio::test] +async fn test_estimate_gas_price() { + let ctx = utils::setup_process_grpc().await; + let mut client = GasEstimatorClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .estimate_gas_price(EstimateGasPriceRequest { + tx_priority: 0, + }) + .await + .expect("EstimateGasPrice failed"); + let body = resp.into_inner(); + + assert!( + (body.estimated_gas_price - 0.002).abs() < 1e-9, + "expected gas price 0.002, got {}", + body.estimated_gas_price + ); +} + +#[tokio::test] +async fn test_estimate_gas_price_and_usage() { + let ctx = utils::setup_process_grpc().await; + let mut client = GasEstimatorClient::new(grpc_channel(&ctx.grpc_addr).await); + + let resp = client + .estimate_gas_price_and_usage(EstimateGasPriceAndUsageRequest { + tx_bytes: vec![], + tx_priority: 0, + }) + .await + .expect("EstimateGasPriceAndUsage failed"); + let body = resp.into_inner(); + + assert!( + (body.estimated_gas_price - 0.002).abs() < 1e-9, + "expected gas price 0.002, got {}", + body.estimated_gas_price + ); + assert_eq!(body.estimated_gas_used, 500_000); +} diff --git a/tests/utils/mod.rs b/tests/utils/mod.rs index 2a8269a..e35b1ef 100644 --- a/tests/utils/mod.rs +++ b/tests/utils/mod.rs @@ -505,6 +505,54 @@ async fn wait_for_port(listen_addr: &str) { } } +pub struct GrpcProcessTestContext { + pub client: Client, + pub http_url: String, + pub grpc_addr: String, + _child: Child, + _redis: RedisTestGuard, + _lock: MutexGuard<'static, ()>, +} + +impl Drop for GrpcProcessTestContext { + fn drop(&mut self) { + let _ = self._child.kill(); + let _ = self._child.wait(); + } +} + +pub async fn setup_process_grpc() -> GrpcProcessTestContext { + let lock = TEST_LOCK.lock().await; + let redis = setup_redis().await; + let redis_url = redis.url.clone(); + let listen_addr = reserve_listen_addr(); + let grpc_addr = reserve_listen_addr(); + let bin_path = resolve_localestia_bin(); + + let child = Command::new(&bin_path) + .env("REDIS_URL", &redis_url) + .env("LISTEN_ADDR", &listen_addr) + .env("GRPC_ADDR", &grpc_addr) + .stdout(Stdio::inherit()) + .stderr(Stdio::inherit()) + .spawn() + .expect("failed to start localestia process"); + + wait_for_port(&listen_addr).await; + wait_for_port(&grpc_addr).await; + let ws_url = format!("ws://{}", listen_addr); + let client = wait_for_client(&ws_url).await; + + GrpcProcessTestContext { + client, + http_url: format!("http://{}", listen_addr), + grpc_addr, + _child: child, + _redis: redis, + _lock: lock, + } +} + async fn wait_for_client(ws_url: &str) -> Client { let start = Instant::now(); let timeout = Duration::from_secs(10);