Skip to content

Commit 8aa9ad1

Browse files
committed
http fallback
1 parent 08b141a commit 8aa9ad1

7 files changed

Lines changed: 147 additions & 57 deletions

File tree

crates/common/src/pbs/error.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,9 @@ pub enum PbsError {
4242
#[error("websocket error: {0}")]
4343
WebSocket(String),
4444

45+
#[error("websocket connect failed: {0}")]
46+
WebSocketConnect(String),
47+
4548
#[error("websocket timed out")]
4649
WebSocketTimeout,
4750
}

crates/pbs/src/mev_boost/get_header.rs

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -251,7 +251,29 @@ async fn get_header_from_relay(
251251

252252
match relay.get_header_request(params.slot, &params.parent_hash, &params.pubkey)? {
253253
GetHeaderRequest::Stream(url) => {
254-
get_header_ws(&request_info, &relay, url, timeout_left_ms).await
254+
let started = Instant::now();
255+
let err = match get_header_ws(&request_info, &relay, url, timeout_left_ms).await {
256+
Err(PbsError::WebSocketConnect(err)) => err,
257+
res => return res,
258+
};
259+
260+
let elapsed_ms = started.elapsed().as_millis() as u64;
261+
let timeout_left_ms = timeout_left_ms.saturating_sub(elapsed_ms);
262+
if timeout_left_ms == 0 {
263+
return Err(PbsError::WebSocketConnect(err));
264+
}
265+
266+
warn!(relay_id = relay.id.as_ref(), %err, timeout_left_ms, "stream failed, falling back to http get_header");
267+
268+
let url = relay.get_header_url(params.slot, &params.parent_hash, &params.pubkey)?;
269+
send_timed_get_header(
270+
request_info,
271+
relay,
272+
ms_into_slot + elapsed_ms,
273+
url,
274+
timeout_left_ms,
275+
)
276+
.await
255277
}
256278
GetHeaderRequest::Http(url) => {
257279
send_timed_get_header(request_info, relay, ms_into_slot, url, timeout_left_ms).await
@@ -271,7 +293,7 @@ async fn send_timed_get_header(
271293
// sleep until target time in slot
272294

273295
let delay = target_ms.saturating_sub(ms_into_slot);
274-
if delay > 0 {
296+
if delay > 0 && delay < timeout_left_ms {
275297
debug!(
276298
relay_id = relay.id.as_ref(),
277299
target_ms, ms_into_slot, "TG: waiting to send first header request"
@@ -281,7 +303,10 @@ async fn send_timed_get_header(
281303
} else {
282304
debug!(
283305
relay_id = relay.id.as_ref(),
284-
target_ms, ms_into_slot, "TG: request already late enough in slot"
306+
target_ms,
307+
ms_into_slot,
308+
timeout_left_ms,
309+
"TG: sending first header request now"
285310
);
286311
}
287312
}

crates/pbs/src/mev_boost/get_header_ws.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -101,7 +101,7 @@ pub(super) async fn get_header_ws(
101101
let status = rejected.as_ref().map_or(TRANSPORT_ERROR_CODE_STR, |code| code.as_str());
102102

103103
record_status(status, relay);
104-
return Err(PbsError::WebSocket(format!("connect failed: {err}")));
104+
return Err(PbsError::WebSocketConnect(err.to_string()));
105105
}
106106
Err(_) => {
107107
record_status(TIMEOUT_ERROR_CODE_STR, relay);
@@ -215,7 +215,7 @@ fn build_handshake_request(
215215
let mut request = url
216216
.as_str()
217217
.into_client_request()
218-
.map_err(|err| PbsError::WebSocket(format!("invalid ws url: {err}")))?;
218+
.map_err(|err| PbsError::WebSocketConnect(format!("invalid ws url: {err}")))?;
219219

220220
let headers = request.headers_mut();
221221
if let Some(user_agent) = request_info.headers.get(USER_AGENT) {

tests/src/mock_relay.rs

Lines changed: 26 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
use std::{
2-
collections::HashSet,
2+
collections::{HashMap, HashSet},
33
net::SocketAddr,
44
sync::{
55
Arc, RwLock,
@@ -36,7 +36,10 @@ use cb_common::{
3636
get_consensus_version_header, get_content_type,
3737
},
3838
};
39-
use cb_pbs::MAX_SIZE_SUBMIT_BLOCK_RESPONSE;
39+
use cb_pbs::{
40+
GET_HEADER_ENDPOINT_TAG, MAX_SIZE_SUBMIT_BLOCK_RESPONSE, REGISTER_VALIDATOR_ENDPOINT_TAG,
41+
STATUS_ENDPOINT_TAG, SUBMIT_BLINDED_BLOCK_ENDPOINT_TAG,
42+
};
4043
use lh_types::KzgProof;
4144
use reqwest::header::{ACCEPT, CONTENT_TYPE};
4245
use ssz::Encode;
@@ -95,7 +98,9 @@ pub struct MockRelayState {
9598
/// The raw `Accept` header PBS sent on the most recent get_header request,
9699
/// so a test can assert what encoding PBS asked the relay for.
97100
received_get_header_accept: RwLock<Option<String>>,
98-
last_register_api_key: RwLock<Option<String>>,
101+
/// Api key header seen per endpoint tag, so a test can assert the relay's
102+
/// configured key rides on every request PBS sends it.
103+
api_keys_seen: RwLock<HashMap<&'static str, String>>,
99104
}
100105

101106
impl MockRelayState {
@@ -135,8 +140,14 @@ impl MockRelayState {
135140
pub fn set_response_override(&self, status: StatusCode) {
136141
*self.response_override.write().unwrap() = Some(status);
137142
}
138-
pub fn last_register_api_key(&self) -> Option<String> {
139-
self.last_register_api_key.read().unwrap().clone()
143+
/// Api key seen on `endpoint`, one of the `*_ENDPOINT_TAG` constants
144+
pub fn api_key_seen(&self, endpoint: &str) -> Option<String> {
145+
self.api_keys_seen.read().unwrap().get(endpoint).cloned()
146+
}
147+
fn record_api_key(&self, endpoint: &'static str, headers: &HeaderMap) {
148+
if let Some(api_key) = headers.get(HEADER_API_KEY).and_then(|key| key.to_str().ok()) {
149+
self.api_keys_seen.write().unwrap().insert(endpoint, api_key.to_string());
150+
}
140151
}
141152
}
142153

@@ -158,7 +169,7 @@ impl MockRelayState {
158169
response_override: RwLock::new(None),
159170
bid_value: RwLock::new(U256::from(10)),
160171
received_get_header_accept: RwLock::new(None),
161-
last_register_api_key: RwLock::new(None),
172+
api_keys_seen: RwLock::new(HashMap::new()),
162173
supported_content_types: Arc::new(
163174
[EncodingType::Json, EncodingType::Ssz].iter().cloned().collect(),
164175
),
@@ -260,6 +271,7 @@ async fn handle_get_header(
260271
headers: HeaderMap,
261272
) -> Response {
262273
state.received_get_header.fetch_add(1, Ordering::Relaxed);
274+
state.record_api_key(GET_HEADER_ENDPOINT_TAG, &headers);
263275
*state.received_get_header_accept.write().unwrap() =
264276
headers.get(ACCEPT).and_then(|v| v.to_str().ok()).map(String::from);
265277
let accept_types = get_accept_types(&headers)
@@ -322,8 +334,12 @@ async fn handle_get_header(
322334
response
323335
}
324336

325-
async fn handle_get_status(State(state): State<Arc<MockRelayState>>) -> impl IntoResponse {
337+
async fn handle_get_status(
338+
State(state): State<Arc<MockRelayState>>,
339+
headers: HeaderMap,
340+
) -> impl IntoResponse {
326341
state.received_get_status.fetch_add(1, Ordering::Relaxed);
342+
state.record_api_key(STATUS_ENDPOINT_TAG, &headers);
327343
// Production `get_status` dispatches relays concurrently via `select_ok`,
328344
// which cancels losing futures as soon as any relay returns OK. On a
329345
// loaded runner this can abort a sibling relay's reqwest send before
@@ -341,10 +357,7 @@ async fn handle_register_validator(
341357
Json(validators): Json<Vec<ValidatorRegistration>>,
342358
) -> impl IntoResponse {
343359
state.received_register_validator.fetch_add(1, Ordering::Relaxed);
344-
*state.last_register_api_key.write().unwrap() = headers
345-
.get(HEADER_API_KEY)
346-
.and_then(|value| value.to_str().ok())
347-
.map(|value| value.to_string());
360+
state.record_api_key(REGISTER_VALIDATOR_ENDPOINT_TAG, &headers);
348361
debug!("Received {} registrations", validators.len());
349362

350363
if let Some(status) = state.response_override.read().unwrap().as_ref() {
@@ -363,6 +376,7 @@ async fn handle_submit_block_v1(
363376
return StatusCode::NOT_FOUND.into_response();
364377
}
365378
state.received_submit_block.fetch_add(1, Ordering::Relaxed);
379+
state.record_api_key(SUBMIT_BLINDED_BLOCK_ENDPOINT_TAG, &headers);
366380
// Short-circuit SSZ requests with an overridden status so tests can
367381
// drive the PBS SSZ→JSON retry logic. JSON requests still take the
368382
// normal path so a single mock run can exercise both attempts.
@@ -475,6 +489,7 @@ async fn handle_submit_block_v2(
475489
return StatusCode::NOT_FOUND.into_response();
476490
}
477491
state.received_submit_block.fetch_add(1, Ordering::Relaxed);
492+
state.record_api_key(SUBMIT_BLINDED_BLOCK_ENDPOINT_TAG, &headers);
478493
// See comment in `handle_submit_block_v1`. Override SSZ with the
479494
// injected status so C3 tests can assert retry / no-retry behavior.
480495
if let Some(status) = state.submit_block_ssz_status_override() &&

tests/src/utils.rs

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,9 @@ use url::Url;
2525

2626
pub const HEADER_API_KEY: &str = "x-api-key";
2727
pub const API_KEY: &str = "123e4567-e89b-12d3-a456-426614174000";
28+
/// Distinct from [`API_KEY`], which `MockValidator` also sends to PBS: a relay
29+
/// that sees this one can only have got it from its own config.
30+
pub const RELAY_API_KEY: &str = "f81d4fae-7dec-11d0-a765-00a0c91e6bf6";
2831

2932
pub fn get_local_address(port: u16) -> String {
3033
format!("http://0.0.0.0:{port}")
@@ -74,12 +77,36 @@ pub fn generate_mock_relay_with_batch_size(
7477
RelayClient::new(config)
7578
}
7679

80+
pub fn generate_mock_relay_with_api_key(
81+
port: u16,
82+
pubkey: BlsPublicKey,
83+
api_key: &str,
84+
) -> Result<RelayClient> {
85+
let mut config = mock_relay_config(port, pubkey)?;
86+
config.headers = Some(HashMap::from([(HEADER_API_KEY.into(), api_key.into())]));
87+
RelayClient::new(config)
88+
}
89+
7790
pub fn generate_mock_stream_relay(port: u16, pubkey: BlsPublicKey) -> Result<RelayClient> {
7891
let mut config = mock_relay_config(port, pubkey)?;
7992
config.get_header = GetHeaderTransport::Stream;
8093
RelayClient::new(config)
8194
}
8295

96+
pub fn generate_mock_stream_relay_with_timing_games(
97+
port: u16,
98+
pubkey: BlsPublicKey,
99+
target_first_request_ms: u64,
100+
frequency_get_header_ms: u64,
101+
) -> Result<RelayClient> {
102+
let mut config = mock_relay_config(port, pubkey)?;
103+
config.get_header = GetHeaderTransport::Stream;
104+
config.enable_timing_games = true;
105+
config.target_first_request_ms = Some(target_first_request_ms);
106+
config.frequency_get_header_ms = Some(frequency_get_header_ms);
107+
RelayClient::new(config)
108+
}
109+
83110
pub fn get_pbs_config(port: u16) -> PbsConfig {
84111
PbsConfig {
85112
host: Ipv4Addr::UNSPECIFIED,

tests/tests/pbs_get_header_ws.rs

Lines changed: 59 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,9 @@ use cb_tests::{
1919
mock_validator::MockValidator,
2020
mock_ws_relay::{MockWsRelayState, start_mock_ws_relay_service},
2121
utils::{
22-
API_KEY, generate_mock_relay, generate_mock_stream_relay, get_free_listener,
23-
get_pbs_config, setup_test_env, to_pbs_config,
22+
API_KEY, generate_mock_relay, generate_mock_stream_relay,
23+
generate_mock_stream_relay_with_timing_games, get_free_listener, get_pbs_config,
24+
setup_test_env, to_pbs_config,
2425
},
2526
};
2627
use eyre::Result;
@@ -290,6 +291,62 @@ async fn test_get_header_ws_unreachable_relay_returns_204() -> Result<()> {
290291
Ok(())
291292
}
292293

294+
/// A relay configured to stream that rejects the handshake — here by not
295+
/// serving the stream route at all — is still asked over HTTP, so its bid stays
296+
/// in the auction instead of the relay silently dropping out.
297+
#[tokio::test]
298+
async fn test_get_header_ws_falls_back_to_http() -> Result<()> {
299+
setup_test_env();
300+
let signer = random_secret();
301+
let pubkey = signer.public_key();
302+
let chain = Chain::Hoodi;
303+
304+
let listener = get_free_listener().await;
305+
let port = listener.local_addr()?.port();
306+
let relay_state =
307+
Arc::new(MockRelayState::new(chain, signer.clone()).with_bid_value(U256::from(42)));
308+
tokio::spawn(start_mock_relay_service_with_listener(relay_state.clone(), listener));
309+
310+
let relay = generate_mock_stream_relay(port, pubkey)?;
311+
let validator = start_pbs(chain, vec![relay], 1_000).await?;
312+
313+
let (code, res) = get_header_json(&validator).await?;
314+
assert_eq!(code, StatusCode::OK);
315+
assert_bid(&res.unwrap(), chain, &signer, U256::from(42));
316+
317+
// Timing games are off for this relay, so the fallback is a single request
318+
assert_eq!(relay_state.received_get_header(), 1);
319+
Ok(())
320+
}
321+
322+
/// The handshake can fail early in the slot, so the fallback is a normal http
323+
/// get_header and still plays this relay's timing games.
324+
#[tokio::test]
325+
async fn test_get_header_ws_fallback_runs_timing_games() -> Result<()> {
326+
setup_test_env();
327+
let signer = random_secret();
328+
let pubkey = signer.public_key();
329+
let chain = Chain::Hoodi;
330+
331+
let listener = get_free_listener().await;
332+
let port = listener.local_addr()?.port();
333+
let relay_state = Arc::new(MockRelayState::new(chain, signer.clone()));
334+
tokio::spawn(start_mock_relay_service_with_listener(relay_state.clone(), listener));
335+
336+
// Target is already behind us, so no wait, then one request per 200ms of
337+
// what is left of the budget
338+
let relay = generate_mock_stream_relay_with_timing_games(port, pubkey, 0, 200)?;
339+
let validator = start_pbs(chain, vec![relay], 1_000).await?;
340+
341+
let (code, res) = get_header_json(&validator).await?;
342+
assert_eq!(code, StatusCode::OK);
343+
assert_bid(&res.unwrap(), chain, &signer, U256::from(10));
344+
345+
let n_requests = relay_state.received_get_header();
346+
assert!(n_requests > 1, "fallback skipped the timing games loop: {n_requests} requests");
347+
Ok(())
348+
}
349+
293350
/// A streamed bid competes in the same auction as an HTTP one
294351
#[tokio::test]
295352
async fn test_get_header_ws_wins_auction_against_http() -> Result<()> {

tests/tests/pbs_post_validators.rs

Lines changed: 2 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,9 @@ use cb_common::{
77
};
88
use cb_pbs::{DefaultBuilderApi, PbsService, PbsState};
99
use cb_tests::{
10-
mock_relay::{
11-
MockRelayState, start_mock_relay_service, start_mock_relay_service_with_listener,
12-
},
10+
mock_relay::{MockRelayState, start_mock_relay_service},
1311
mock_validator::MockValidator,
14-
utils::{
15-
API_KEY, generate_mock_relay, get_free_listener, get_pbs_config, setup_test_env,
16-
to_pbs_config,
17-
},
12+
utils::{generate_mock_relay, get_pbs_config, setup_test_env, to_pbs_config},
1813
};
1914
use eyre::Result;
2015
use reqwest::StatusCode;
@@ -66,38 +61,6 @@ async fn test_register_validators() -> Result<()> {
6661
Ok(())
6762
}
6863

69-
/// The api key from `relay.config.headers` is forwarded on registrations.
70-
#[tokio::test]
71-
async fn test_register_validators_sends_api_key() -> Result<()> {
72-
setup_test_env();
73-
let signer = random_secret();
74-
let pubkey: BlsPublicKey = signer.public_key();
75-
let chain = Chain::Holesky;
76-
77-
let relay_listener = get_free_listener().await;
78-
let relay_port = relay_listener.local_addr()?.port();
79-
let mock_state = Arc::new(MockRelayState::new(chain, signer));
80-
tokio::spawn(start_mock_relay_service_with_listener(mock_state.clone(), relay_listener));
81-
82-
let pbs_listener = get_free_listener().await;
83-
let pbs_port = pbs_listener.local_addr()?.port();
84-
drop(pbs_listener);
85-
86-
let relays = vec![generate_mock_relay(relay_port, pubkey)?];
87-
let config = to_pbs_config(chain, get_pbs_config(pbs_port), relays);
88-
tokio::spawn(PbsService::run::<(), DefaultBuilderApi>(PbsState::new(config, PathBuf::new())));
89-
tokio::time::sleep(Duration::from_millis(100)).await;
90-
91-
let mock_validator = MockValidator::new(pbs_port)?;
92-
let res = mock_validator.do_register_validator().await?;
93-
assert_eq!(res.status(), StatusCode::OK);
94-
95-
assert_eq!(mock_state.received_register_validator(), 1);
96-
assert_eq!(mock_state.last_register_api_key().as_deref(), Some(API_KEY));
97-
98-
Ok(())
99-
}
100-
10164
#[tokio::test]
10265
async fn test_register_validators_does_not_retry_on_429() -> Result<()> {
10366
setup_test_env();

0 commit comments

Comments
 (0)