Skip to content

Commit 45c836b

Browse files
committed
feat(beacon): add beacon RPC handlers
1 parent 275ed66 commit 45c836b

5 files changed

Lines changed: 434 additions & 11 deletions

File tree

‎src/config.rs‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,11 +57,23 @@ pub fn default_chains() -> Vec<ChainSpec> {
5757
]
5858
}
5959

60+
/// Default beacon light-client upstream. A SINGLE consistent node — Helios's
61+
/// `light_client/updates` walk breaks if fanned out across backends at different
62+
/// sync states (the whole reason the `/beacon` split exists). See `SERVER_CONSENSUS.md`.
63+
pub const DEFAULT_BEACON_LC_UPSTREAM: &str = "https://lodestar-mainnet.chainsafe.io";
64+
/// Default beacon blocks upstream. Serves the ~750 KB `/eth/v2/beacon/blocks/{slot}`
65+
/// payload Helios re-fetches every slot without rate-limiting it.
66+
pub const DEFAULT_BEACON_BLOCKS_UPSTREAM: &str = "https://ethereum-beacon-api.publicnode.com";
67+
6068
pub struct Config {
6169
pub bind: SocketAddr,
6270
pub drpc_base: String,
6371
pub key: Key,
6472
pub chains: Vec<ChainSpec>,
73+
/// Beacon light-client upstream (clean, in-order `light_client/updates`).
74+
pub beacon_lc_upstream: String,
75+
/// Beacon blocks upstream (un-rate-limited `/eth/v2/beacon/blocks/{slot}`).
76+
pub beacon_blocks_upstream: String,
6577
pub body_limit: usize,
6678
/// When true, refuse expensive/nonsensical-over-HTTP JSON-RPC methods so a
6779
/// leaked endpoint can't burn the shared key on archival traces. Off via
@@ -88,6 +100,14 @@ impl Config {
88100
drpc_base: drpc_base.trim_end_matches('/').to_string(),
89101
key: load_key()?,
90102
chains: default_chains(),
103+
beacon_lc_upstream: std::env::var("BEACON_LC_UPSTREAM")
104+
.unwrap_or_else(|_| DEFAULT_BEACON_LC_UPSTREAM.into())
105+
.trim_end_matches('/')
106+
.to_string(),
107+
beacon_blocks_upstream: std::env::var("BEACON_BLOCKS_UPSTREAM")
108+
.unwrap_or_else(|_| DEFAULT_BEACON_BLOCKS_UPSTREAM.into())
109+
.trim_end_matches('/')
110+
.to_string(),
91111
body_limit: env_usize("RPC_BODY_LIMIT_BYTES", 1 << 20), // 1 MiB
92112
method_denylist: std::env::var("RPC_METHOD_FILTER")
93113
.map(|v| v != "off")

‎src/main.rs‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,9 +21,17 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
2121
.connect_timeout(cfg.connect_timeout)
2222
.build()?;
2323

24+
// The beacon proxy reuses the same HTTPS-only client (cloning shares the
25+
// connection pool). Its ~20s timeout covers the large ~750 KB block fetch.
26+
let beacon = upstream::Beacon::new(
27+
client.clone(),
28+
cfg.beacon_lc_upstream,
29+
cfg.beacon_blocks_upstream,
30+
);
2431
let drpc = upstream::Drpc::new(client, cfg.drpc_base.clone(), cfg.key);
2532
let state = Arc::new(routes::AppState {
2633
drpc,
34+
beacon,
2735
chains: cfg.chains,
2836
method_denylist: cfg.method_denylist,
2937
});
@@ -39,6 +47,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
3947
"/v1/{chain}/history/{address}",
4048
get(routes::history_handler),
4149
)
50+
.route("/beacon/{*path}", get(routes::beacon_handler))
4251
.layer(DefaultBodyLimit::max(cfg.body_limit))
4352
.with_state(state);
4453

‎src/routes.rs‎

Lines changed: 52 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
//! POST /rpc/{chain} JSON-RPC passthrough (alias -> dRPC slug)
66
//! GET /v1/{chain}/balances/{address} Wallet-API balances for one chain
77
//! GET /v1/{chain}/history/{address} Wallet-API tx history for one chain
8+
//! GET /beacon/{*path} Beacon light-client proxy for Helios
89
//!
910
//! The wallet-API routes are a thin per-chain mirror of dRPC's Wallet API: the
1011
//! proxy resolves the chain alias to a dRPC slug, injects the shared key,
@@ -15,7 +16,7 @@
1516
//! source.
1617
1718
use crate::config::ChainSpec;
18-
use crate::upstream::{Drpc, UpstreamError};
19+
use crate::upstream::{Beacon, Drpc, UpstreamError};
1920
use axum::Json;
2021
use axum::body::Bytes;
2122
use axum::extract::{Path, RawQuery, State};
@@ -26,6 +27,7 @@ use std::sync::Arc;
2627

2728
pub struct AppState {
2829
pub drpc: Drpc,
30+
pub beacon: Beacon,
2931
pub chains: Vec<ChainSpec>,
3032
pub method_denylist: bool,
3133
}
@@ -50,6 +52,7 @@ pub async fn rpc_handler(
5052
Path(chain): Path<String>,
5153
body: Bytes,
5254
) -> Result<Response, AppError> {
55+
tracing::info!(chain, body_len = body.len(), "TEMP: rpc request received");
5356
let slug = state.resolve_chain(&chain).ok_or(AppError::UnknownChain)?;
5457
if state.method_denylist {
5558
reject_denied_methods(&body)?;
@@ -63,6 +66,7 @@ pub async fn balances_handler(
6366
Path((chain, address)): Path<(String, String)>,
6467
RawQuery(query): RawQuery,
6568
) -> Result<Response, AppError> {
69+
tracing::info!(chain, "TEMP: balances request received");
6670
let slug = state.resolve_chain(&chain).ok_or(AppError::UnknownChain)?;
6771
if !is_valid_eth_address(&address) {
6872
return Err(AppError::BadAddress);
@@ -79,6 +83,7 @@ pub async fn history_handler(
7983
Path((chain, address)): Path<(String, String)>,
8084
RawQuery(query): RawQuery,
8185
) -> Result<Response, AppError> {
86+
tracing::info!(chain, "TEMP: history request received");
8287
let slug = state.resolve_chain(&chain).ok_or(AppError::UnknownChain)?;
8388
if !is_valid_eth_address(&address) {
8489
return Err(AppError::BadAddress);
@@ -90,10 +95,43 @@ pub async fn history_handler(
9095
Ok(json_passthrough(code, body))
9196
}
9297

98+
/// Beacon light-client proxy. The wallet's Helios points its consensus RPC at
99+
/// `{KAO}/beacon` and makes plain HTTP GET / JSON requests. We strip the
100+
/// `/beacon` prefix, forward the path + query verbatim to the routed upstream,
101+
/// and return its status, `Content-Type`, and body unchanged.
102+
///
103+
/// On a pre-response upstream failure we return a clean 5xx (via `AppError`),
104+
/// which Helios retries — never a malformed 200 body, which it treats as fatal.
105+
pub async fn beacon_handler(
106+
State(state): State<Arc<AppState>>,
107+
Path(path): Path<String>,
108+
RawQuery(query): RawQuery,
109+
) -> Result<Response, AppError> {
110+
// Axum's `{*path}` capture drops the leading slash; the upstream wants an
111+
// absolute beacon-API path (e.g. `/eth/v2/beacon/blocks/14702657`).
112+
let path = format!("/{path}");
113+
let (code, content_type, body) = state.beacon.get(&path, query.as_deref()).await?;
114+
Ok(beacon_passthrough(code, content_type, body))
115+
}
116+
93117
/// Build a passthrough response: upstream status (clamped to a valid code),
94118
/// `application/json`, and the upstream body verbatim.
95119
fn json_passthrough(code: u16, body: Vec<u8>) -> Response {
96-
let status = StatusCode::from_u16(code).unwrap_or(StatusCode::BAD_GATEWAY);
120+
// TEMP: a clamp here is the only way this app emits a 502 with a non-JSON
121+
// (possibly empty) body — exactly the symptom being chased.
122+
let status = StatusCode::from_u16(code).unwrap_or_else(|_| {
123+
tracing::warn!(
124+
code,
125+
body_len = body.len(),
126+
"TEMP: invalid upstream status, clamping to 502"
127+
);
128+
StatusCode::BAD_GATEWAY
129+
});
130+
tracing::info!(
131+
status = code,
132+
body_len = body.len(),
133+
"TEMP: passthrough response"
134+
);
97135
(
98136
status,
99137
[(
@@ -105,6 +143,16 @@ fn json_passthrough(code: u16, body: Vec<u8>) -> Response {
105143
.into_response()
106144
}
107145

146+
/// Beacon passthrough: upstream status (clamped) and body verbatim, preserving
147+
/// the upstream `Content-Type` (Helios parses JSON; we don't rewrite the body).
148+
fn beacon_passthrough(code: u16, content_type: Option<String>, body: Vec<u8>) -> Response {
149+
let status = StatusCode::from_u16(code).unwrap_or(StatusCode::BAD_GATEWAY);
150+
let ct = content_type
151+
.and_then(|s| HeaderValue::from_str(&s).ok())
152+
.unwrap_or_else(|| HeaderValue::from_static("application/json"));
153+
(status, [(header::CONTENT_TYPE, ct)], body).into_response()
154+
}
155+
108156
// ---- validation & method filtering ------------------------------------------
109157

110158
fn is_valid_eth_address(s: &str) -> bool {
@@ -169,6 +217,8 @@ impl IntoResponse for AppError {
169217
AppError::BadRequest(m) => (StatusCode::BAD_REQUEST, m),
170218
AppError::Upstream(m) => (StatusCode::BAD_GATEWAY, m),
171219
};
220+
// TEMP:
221+
tracing::warn!(status = %status, msg, "TEMP: app error response");
172222
(
173223
status,
174224
Json(serde_json::json!({ "error": { "message": msg } })),

‎src/upstream.rs‎

Lines changed: 105 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -56,9 +56,24 @@ impl Drpc {
5656
.body(body.to_vec())
5757
.send()
5858
.await
59-
.map_err(|_| UpstreamError::Transport)?;
59+
.map_err(|e| {
60+
// TEMP: `.without_url()` strips the key-bearing URL before logging.
61+
tracing::error!(error = %e.without_url(), chain, "TEMP: rpc send failed");
62+
UpstreamError::Transport
63+
})?;
6064
let status = resp.status().as_u16();
61-
let bytes = resp.bytes().await.map_err(|_| UpstreamError::Transport)?;
65+
let bytes = resp.bytes().await.map_err(|e| {
66+
// TEMP: `.without_url()` strips the key-bearing URL before logging.
67+
tracing::error!(error = %e.without_url(), chain, status, "TEMP: rpc body read failed");
68+
UpstreamError::Transport
69+
})?;
70+
// TEMP:
71+
tracing::info!(
72+
chain,
73+
status,
74+
body_len = bytes.len(),
75+
"TEMP: rpc upstream response"
76+
);
6277
Ok((status, bytes.to_vec()))
6378
}
6479

@@ -108,13 +123,99 @@ impl Drpc {
108123
.header(ACCEPT, "application/json")
109124
.send()
110125
.await
111-
.map_err(|_| UpstreamError::Transport)?;
126+
.map_err(|e| {
127+
// TEMP: `.without_url()` strips the key-bearing URL before logging.
128+
tracing::error!(error = %e.without_url(), "TEMP: get send failed");
129+
UpstreamError::Transport
130+
})?;
112131
let status = resp.status().as_u16();
113-
let bytes = resp.bytes().await.map_err(|_| UpstreamError::Transport)?;
132+
let bytes = resp.bytes().await.map_err(|e| {
133+
// TEMP: `.without_url()` strips the key-bearing URL before logging.
134+
tracing::error!(error = %e.without_url(), status, "TEMP: get body read failed");
135+
UpstreamError::Transport
136+
})?;
137+
// TEMP:
138+
tracing::info!(
139+
status,
140+
body_len = bytes.len(),
141+
"TEMP: get upstream response"
142+
);
114143
Ok((status, bytes.to_vec()))
115144
}
116145
}
117146

147+
/// Beacon light-client proxy for the wallet's Helios light client.
148+
///
149+
/// Serves the six beacon-API paths under `/beacon` by routing each to whichever
150+
/// upstream handles *that* request well (see `SERVER_CONSENSUS.md`):
151+
/// - `/eth/v{1,2}/beacon/blocks/…` → `blocks_upstream` (large payload, must not 429)
152+
/// - everything else → `lc_upstream` (clean, in-order updates)
153+
///
154+
/// Like `Drpc`, requests are originated here — no client IP/UA/cookie reaches the
155+
/// upstream. Once the upstream answers, its status, `Content-Type`, and body are
156+
/// returned verbatim; only pre-response failures collapse into `Transport` (which
157+
/// the handler renders as a clean 5xx that Helios retries).
158+
pub struct Beacon {
159+
client: Client,
160+
lc_upstream: String,
161+
blocks_upstream: String,
162+
}
163+
164+
impl Beacon {
165+
pub fn new(client: Client, lc_upstream: String, blocks_upstream: String) -> Self {
166+
Self {
167+
client,
168+
lc_upstream,
169+
blocks_upstream,
170+
}
171+
}
172+
173+
/// Pick the upstream for a beacon path. Blocks go to the working-blocks
174+
/// upstream; everything else (light_client/*, config/spec, states/*) to the
175+
/// clean-updates upstream. `path` is absolute (leading `/`).
176+
fn upstream_for(&self, path: &str) -> &str {
177+
if path.starts_with("/eth/v1/beacon/blocks/") || path.starts_with("/eth/v2/beacon/blocks/")
178+
{
179+
&self.blocks_upstream
180+
} else {
181+
&self.lc_upstream
182+
}
183+
}
184+
185+
/// GET passthrough. `path` is the beacon-API path (leading `/`, no `/beacon`
186+
/// prefix); `query` is the caller's raw query string, forwarded verbatim.
187+
/// Returns the upstream status, `Content-Type`, and body unchanged.
188+
pub async fn get(
189+
&self,
190+
path: &str,
191+
query: Option<&str>,
192+
) -> Result<(u16, Option<String>, Vec<u8>), UpstreamError> {
193+
let mut url = format!("{}{}", self.upstream_for(path), path);
194+
append_query(&mut url, query);
195+
let resp = self
196+
.client
197+
.get(&url)
198+
.header(ACCEPT, "application/json")
199+
.send()
200+
.await
201+
.map_err(|e| {
202+
tracing::error!(error = %e.without_url(), path, "beacon send failed");
203+
UpstreamError::Transport
204+
})?;
205+
let status = resp.status().as_u16();
206+
let content_type = resp
207+
.headers()
208+
.get(CONTENT_TYPE)
209+
.and_then(|v| v.to_str().ok())
210+
.map(str::to_owned);
211+
let bytes = resp.bytes().await.map_err(|e| {
212+
tracing::error!(error = %e.without_url(), path, status, "beacon body read failed");
213+
UpstreamError::Transport
214+
})?;
215+
Ok((status, content_type, bytes.to_vec()))
216+
}
217+
}
218+
118219
/// Append `?query` to `url` when the caller supplied a non-empty query string.
119220
fn append_query(url: &mut String, query: Option<&str>) {
120221
if let Some(q) = query

0 commit comments

Comments
 (0)