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
11 changes: 9 additions & 2 deletions crates/common/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -254,9 +254,16 @@ pub struct BlockMergingConfig {
/// Maximum age of a merged bid before it is considered stale and discarded.
#[serde(default = "default_u64::<250>")]
pub max_merged_bid_age_ms: u64,
/// Flag to allow dry run mode.
/// Testing-only: marks every decoded submission as mergeable, without builders having to
/// submit real merge data themselves.
#[serde(default = "default_bool::<false>")]
pub is_dry_run: bool,
pub mark_all_txs_mergeable: bool,
/// Whether `get_header` serves merged blocks to the proposer. Only sets the startup value
/// of `LocalCache`'s live toggle (see `enable_merged_headers`/`disable_merged_headers`) —
/// a production safety valve that can be flipped via the admin API if merging misbehaves,
/// without affecting builder submissions or the merging process itself.
#[serde(default = "default_bool::<true>")]
pub serve_merged_headers: bool,
/// Builder-side merging over TCP. Tile is only spawned if set.
#[serde(default)]
pub tcp: Option<BlockMergingTcpConfig>,
Expand Down
12 changes: 6 additions & 6 deletions crates/common/src/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ pub struct SubmissionDecoderParams {
pub is_dehydrated: bool,
pub with_mergeable_data: bool,
pub with_adjustments: bool,
pub block_merging_dry_run: bool,
pub mark_all_txs_mergeable: bool,
pub fork_name: ForkName,
}

Expand All @@ -178,7 +178,7 @@ pub struct SubmissionDecoder {
is_dehydrated: bool,
with_mergeable_data: bool,
with_adjustments: bool,
block_merging_dry_run: bool,
mark_all_txs_mergeable: bool,
fork_name: ForkName,

bytes_before_decompress: usize,
Expand All @@ -198,7 +198,7 @@ impl SubmissionDecoder {
is_dehydrated: params.is_dehydrated,
with_mergeable_data: params.with_mergeable_data,
with_adjustments: params.with_adjustments,
block_merging_dry_run: params.block_merging_dry_run,
mark_all_txs_mergeable: params.mark_all_txs_mergeable,
fork_name: params.fork_name,
bytes_before_decompress: 0,
bytes_after_decompress: 0,
Expand Down Expand Up @@ -305,7 +305,7 @@ impl SubmissionDecoder {
Some(BlockMergingData::append_only(submission.fee_recipient()))
}
MergeType::None => {
if self.block_merging_dry_run {
if self.mark_all_txs_mergeable {
Some(BlockMergingData::allow_all(
submission.fee_recipient(),
submission.num_txs(),
Expand Down Expand Up @@ -369,7 +369,7 @@ impl SubmissionDecoder {
Some(BlockMergingData::append_only(submission.fee_recipient()))
}
MergeType::None => {
if self.block_merging_dry_run {
if self.mark_all_txs_mergeable {
Some(BlockMergingData::allow_all(
submission.fee_recipient(),
submission.num_txs(),
Expand Down Expand Up @@ -554,7 +554,7 @@ mod tests {
is_dehydrated: true,
with_mergeable_data: false,
with_adjustments: false,
block_merging_dry_run: false,
mark_all_txs_mergeable: false,
fork_name: ForkName::Fulu,
};
let mut decoder = SubmissionDecoder::new(&params);
Expand Down
36 changes: 36 additions & 0 deletions crates/common/src/local_cache.rs
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,10 @@ pub struct LocalCache {
pub api_key_cache: Arc<DashMap<String, Vec<BlsPublicKeyBytes>>>,
primev_proposers: Arc<DashSet<BlsPublicKeyBytes>>,
kill_switch: Arc<AtomicBool>,
/// Production safety valve: whether `get_header` serves merged blocks to the proposer.
/// Seeded at startup from `BlockMergingConfig::serve_merged_headers`, but can also be
/// toggled live via the admin API, same as `kill_switch`.
serve_merged_headers: Arc<AtomicBool>,
proposer_duties: Arc<RwLock<Vec<BuilderGetValidatorsResponseEntry>>>,
merged_blocks: Arc<DashMap<B256, MergedBlock>>,
pub validator_registration_cache:
Expand All @@ -115,6 +119,7 @@ impl LocalCache {
let api_key_cache = Arc::new(DashMap::with_capacity(ESTIMATED_BUILDER_INFOS_UPPER_BOUND));
let primev_proposers = Arc::new(DashSet::with_capacity(MAX_PRIMEV_PROPOSERS));
let kill_switch = Arc::new(AtomicBool::new(false));
let serve_merged_headers = Arc::new(AtomicBool::new(true));
let proposer_duties = Arc::new(RwLock::new(Vec::with_capacity(1000)));
let merged_blocks = Arc::new(DashMap::with_capacity(1000));
let validator_registration_cache = Arc::new(DashMap::with_capacity(1_800_000));
Expand All @@ -133,6 +138,7 @@ impl LocalCache {
api_key_cache,
primev_proposers,
kill_switch,
serve_merged_headers,
proposer_duties,
merged_blocks,
validator_registration_cache,
Expand Down Expand Up @@ -253,6 +259,18 @@ impl LocalCache {
self.kill_switch.store(false, Ordering::Relaxed);
}

pub fn merged_headers_enabled(&self) -> bool {
self.serve_merged_headers.load(Ordering::Relaxed)
}

pub fn enable_merged_headers(&self) {
self.serve_merged_headers.store(true, Ordering::Relaxed);
}

pub fn disable_merged_headers(&self) {
self.serve_merged_headers.store(false, Ordering::Relaxed);
}

pub fn update_current_inclusion_list(
&self,
inclusion_list: InclusionListWithMetadata,
Expand Down Expand Up @@ -442,4 +460,22 @@ mod tests {
let result = cache.kill_switch_enabled();
assert!(!result, "Kill switch should be disabled");
}

#[tokio::test]
pub async fn test_serve_merged_headers() {
let cache = LocalCache::new();

let result = cache.merged_headers_enabled();
assert!(result, "Merged headers should be served by default");

cache.disable_merged_headers();

let result = cache.merged_headers_enabled();
assert!(!result, "Merged headers should not be served");

cache.enable_merged_headers();

let result = cache.merged_headers_enabled();
assert!(result, "Merged headers should be served");
}
}
50 changes: 50 additions & 0 deletions crates/relay/src/api/admin_service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@ use tracing::{error, info};
pub async fn run_admin_service(auctioneer: Arc<LocalCache>, admin_token: String) {
let router = Router::new()
.route("/admin/v1/killswitch", post(enable_kill_switch).delete(disable_kill_switch))
.route(
"/admin/v1/merged-headers",
post(enable_merged_headers).delete(disable_merged_headers),
)
.layer(Extension(auctioneer))
.layer(ValidateRequestHeaderLayer::bearer(&admin_token));

Expand All @@ -34,6 +38,22 @@ async fn disable_kill_switch(
Ok((StatusCode::NO_CONTENT, ()))
}

async fn enable_merged_headers(
Extension(auctioneer): Extension<Arc<LocalCache>>,
) -> Result<impl IntoResponse, StatusCode> {
auctioneer.enable_merged_headers();
info!("Merged header serving enabled");
Ok((StatusCode::NO_CONTENT, ()))
}

async fn disable_merged_headers(
Extension(auctioneer): Extension<Arc<LocalCache>>,
) -> Result<impl IntoResponse, StatusCode> {
auctioneer.disable_merged_headers();
info!("Merged header serving disabled");
Ok((StatusCode::NO_CONTENT, ()))
}

#[cfg(test)]
#[allow(clippy::field_reassign_with_default)]
mod test {
Expand Down Expand Up @@ -73,6 +93,36 @@ mod test {
assert!(!auctioneer.kill_switch_enabled());
}

#[tokio::test]
#[serial]
async fn test_admin_service_merged_headers() {
let auctioneer = Arc::new(LocalCache::new());
assert!(auctioneer.merged_headers_enabled(), "should be enabled by default");

let admin_token = "test_token".into();
tokio::spawn(run_admin_service(auctioneer.clone(), admin_token));
tokio::time::sleep(std::time::Duration::from_secs(1)).await; // wait for server to start
let client = reqwest::Client::new();

let response = client
.delete("http://localhost:4050/admin/v1/merged-headers")
.bearer_auth("test_token")
.send()
.await
.unwrap();
assert_eq!(response.status(), 204);
assert!(!auctioneer.merged_headers_enabled());

let response = client
.post("http://localhost:4050/admin/v1/merged-headers")
.bearer_auth("test_token")
.send()
.await
.unwrap();
assert_eq!(response.status(), 204);
assert!(auctioneer.merged_headers_enabled());
}

#[tokio::test]
#[serial]
async fn test_admin_service_unauthorized() {
Expand Down
4 changes: 2 additions & 2 deletions crates/relay/src/auctioneer/block_merger.rs
Original file line number Diff line number Diff line change
Expand Up @@ -144,8 +144,8 @@ impl BlockMerger {
}
}

if self.config.block_merging_config.is_dry_run {
info!("dry run mode enabled, not returning merged header");
if !self.local_cache.merged_headers_enabled() {
info!("merged header serving disabled, not returning merged header");
return None;
}
Some(entry.bid.clone())
Expand Down
2 changes: 1 addition & 1 deletion crates/relay/src/bid_decoder/tile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -344,7 +344,7 @@ impl DecoderTile {
merge_type: header.merge_type,
with_mergeable_data,
with_adjustments,
block_merging_dry_run: config.block_merging_config.is_dry_run,
mark_all_txs_mergeable: config.block_merging_config.mark_all_txs_mergeable,
fork_name: chain_info.current_fork_name(),
};

Expand Down
3 changes: 3 additions & 0 deletions crates/relay/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,9 @@ async fn run(

let known_validators_loaded = Arc::new(AtomicBool::default());
let local_cache = Arc::new(LocalCache::new());
if !config.block_merging_config.serve_merged_headers {
local_cache.disable_merged_headers();
}

let (db_request_sender, db_request_receiver) = crossbeam_channel::bounded(10_000);
let (db_batch_request_sender, db_batch_request_receiver) = crossbeam_channel::bounded(10_000);
Expand Down
Loading