Skip to content

Commit 845b32b

Browse files
committed
V3 wait handle
1 parent 86329a1 commit 845b32b

4 files changed

Lines changed: 91 additions & 51 deletions

File tree

crates/rbuilder-operator/src/flashbots_config.rs

Lines changed: 19 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -177,18 +177,22 @@ impl LiveBuilderConfig for FlashbotsConfig {
177177
.create_bid_observer_and_submission_policy(&cancellation_token, &abort_token)
178178
.await?;
179179

180-
let (sink_factory, slot_info_provider, adjustment_fee_payers) =
181-
create_sink_factory_and_relays(
182-
&self.base_config,
183-
&self.l1_config,
184-
bidding_service.relay_sets().to_vec(),
185-
wallet_balance_watcher,
186-
bid_observer,
187-
submission_policy,
188-
bidding_service.clone(),
189-
cancellation_token.clone(),
190-
)
191-
.await?;
180+
let (
181+
sink_factory,
182+
slot_info_provider,
183+
adjustment_fee_payers,
184+
optimistic_v3_server_join_handle,
185+
) = create_sink_factory_and_relays(
186+
&self.base_config,
187+
&self.l1_config,
188+
bidding_service.relay_sets().to_vec(),
189+
wallet_balance_watcher,
190+
bid_observer,
191+
submission_policy,
192+
bidding_service.clone(),
193+
cancellation_token.clone(),
194+
)
195+
.await?;
192196

193197
let mut live_builder = create_builder_from_sink(
194198
&self.base_config,
@@ -205,7 +209,9 @@ impl LiveBuilderConfig for FlashbotsConfig {
205209
if let Some(handle) = clickhouse_shutdown_handle {
206210
live_builder.add_critical_task(handle);
207211
}
208-
212+
if let Some(optimistic_v3_server_join_handle) = optimistic_v3_server_join_handle {
213+
live_builder.add_critical_task(optimistic_v3_server_join_handle);
214+
}
209215
let mut module = RpcModule::new(());
210216
module.register_async_method("bid_subsidiseBlock", move |params, _| {
211217
handle_subsidise_block(bidding_service.clone(), params)

crates/rbuilder/src/live_builder/cli.rs

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,10 +17,7 @@ use crate::{
1717
builders::{BacktestSimulateBlockInput, Block},
1818
PartialBlockExecutionTracer,
1919
},
20-
live_builder::{
21-
process_killer::{ProcessKiller, MAX_WAIT_TIME},
22-
watchdog::spawn_watchdog_thread,
23-
},
20+
live_builder::{process_killer::ProcessKiller, watchdog::spawn_watchdog_thread},
2421
provider::StateProviderFactory,
2522
telemetry,
2623
utils::bls::generate_random_bls_address,

crates/rbuilder/src/live_builder/config.rs

Lines changed: 65 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ use std::{
7878
sync::Arc,
7979
time::Duration,
8080
};
81-
use tokio::sync::Mutex as TokioMutex;
81+
use tokio::{sync::Mutex as TokioMutex, task::JoinHandle};
8282
use tokio_util::sync::CancellationToken;
8383
use tracing::{info, warn};
8484
use url::Url;
@@ -395,7 +395,11 @@ impl L1Config {
395395
})
396396
}
397397

398-
/// Creates the RelaySubmitSinkFactory and also returns the associated relays (MevBoostRelaySlotInfoProvider).
398+
/// Creates:
399+
/// - RelaySubmitSinkFactory
400+
/// - The associated relays (MevBoostRelaySlotInfoProvider).
401+
/// - A hashmap of adjustment fee payers.
402+
/// - JoinHandle for the optimistic V3 server.
399403
#[allow(clippy::type_complexity)]
400404
pub fn create_relays_sealed_sink_factory(
401405
&self,
@@ -408,6 +412,7 @@ impl L1Config {
408412
RelaySubmitSinkFactory,
409413
Vec<MevBoostRelaySlotInfoProvider>,
410414
ahash::HashMap<MevBoostRelayID, Address>,
415+
Option<JoinHandle<()>>,
411416
)> {
412417
let signing_domain = get_signing_domain(
413418
chain_spec.chain,
@@ -416,7 +421,7 @@ impl L1Config {
416421
)?;
417422

418423
let mut optimistic_v3_config = None;
419-
if self
424+
let optimistic_v3_server_join_handle = if self
420425
.relays
421426
.iter()
422427
.any(|r| r.submit_config.as_ref().is_some_and(|c| c.optimistic_v3))
@@ -433,7 +438,7 @@ impl L1Config {
433438
}
434439

435440
let optimistic_v3_cache = OptimisticV3BlockCache::default();
436-
optimistic_v3::spawn_server(
441+
let optimistic_v3_server_join_handle = optimistic_v3::spawn_server(
437442
address,
438443
signing_domain,
439444
self.optimistic_v3_relay_pubkeys.clone(),
@@ -444,8 +449,11 @@ impl L1Config {
444449
optimistic_v3_config = Some(OptimisticV3Config {
445450
builder_url: builder_url.into_bytes(),
446451
cache: optimistic_v3_cache,
447-
})
448-
}
452+
});
453+
Some(optimistic_v3_server_join_handle)
454+
} else {
455+
None
456+
};
449457

450458
let submission_config = self.submission_config(
451459
chain_spec,
@@ -479,7 +487,12 @@ impl L1Config {
479487
})
480488
.collect();
481489

482-
Ok((sink_factory, slot_info_providers, adjustment_fee_payers))
490+
Ok((
491+
sink_factory,
492+
slot_info_providers,
493+
adjustment_fee_payers,
494+
optimistic_v3_server_join_handle,
495+
))
483496
}
484497

485498
pub fn registration_update_interval(&self) -> Duration {
@@ -516,20 +529,24 @@ impl LiveBuilderConfig for Config {
516529
let (wallet_balance_watcher, _) =
517530
create_wallet_balance_watcher(provider.clone(), &self.base_config).await?;
518531

519-
let (sink_factory, slot_info_provider, adjustment_fee_payers) =
520-
create_sink_factory_and_relays(
521-
&self.base_config,
522-
&self.l1_config,
523-
bidding_service.relay_sets(),
524-
wallet_balance_watcher,
525-
Box::new(NullBidObserver {}),
526-
Box::new(AlwaysSubmitPolicy {}),
527-
bidding_service,
528-
cancellation_token.clone(),
529-
)
530-
.await?;
531-
532-
let live_builder = create_builder_from_sink(
532+
let (
533+
sink_factory,
534+
slot_info_provider,
535+
adjustment_fee_payers,
536+
optimistic_v3_server_join_handle,
537+
) = create_sink_factory_and_relays(
538+
&self.base_config,
539+
&self.l1_config,
540+
bidding_service.relay_sets(),
541+
wallet_balance_watcher,
542+
Box::new(NullBidObserver {}),
543+
Box::new(AlwaysSubmitPolicy {}),
544+
bidding_service,
545+
cancellation_token.clone(),
546+
)
547+
.await?;
548+
549+
let mut live_builder = create_builder_from_sink(
533550
&self.base_config,
534551
&self.l1_config,
535552
provider,
@@ -544,6 +561,9 @@ impl LiveBuilderConfig for Config {
544561
self.live_builders()?,
545562
self.base_config.max_order_execution_duration_warning(),
546563
);
564+
if let Some(optimistic_v3_server_join_handle) = optimistic_v3_server_join_handle {
565+
live_builder.add_critical_task(optimistic_v3_server_join_handle);
566+
}
547567
Ok(live_builder.with_builders(builders))
548568
}
549569

@@ -1082,6 +1102,11 @@ where
10821102
.await??)
10831103
}
10841104

1105+
/// Returns:
1106+
/// - UnfinishedBuiltBlocksInputFactory
1107+
/// - The associated relays (MevBoostRelaySlotInfoProvider).
1108+
/// - A hashmap of adjustment fee payers.
1109+
/// - JoinHandle for the optimistic V3 server.
10851110
#[allow(clippy::too_many_arguments)]
10861111
pub async fn create_sink_factory_and_relays<P>(
10871112
base_config: &BaseConfig,
@@ -1096,18 +1121,23 @@ pub async fn create_sink_factory_and_relays<P>(
10961121
UnfinishedBuiltBlocksInputFactory<P>,
10971122
Vec<MevBoostRelaySlotInfoProvider>,
10981123
ahash::HashMap<MevBoostRelayID, Address>,
1124+
Option<JoinHandle<()>>,
10991125
)>
11001126
where
11011127
P: StateProviderFactory + Clone + 'static,
11021128
{
1103-
let (sink_sealed_factory, slot_info_provider, adjustment_fee_payers) = l1_config
1104-
.create_relays_sealed_sink_factory(
1105-
base_config.chain_spec()?,
1106-
relay_sets.clone(),
1107-
bid_observer,
1108-
submission_policy,
1109-
cancellation_token.clone(),
1110-
)?;
1129+
let (
1130+
sink_sealed_factory,
1131+
slot_info_provider,
1132+
adjustment_fee_payers,
1133+
optimistic_v3_server_join_handle,
1134+
) = l1_config.create_relays_sealed_sink_factory(
1135+
base_config.chain_spec()?,
1136+
relay_sets.clone(),
1137+
bid_observer,
1138+
submission_policy,
1139+
cancellation_token.clone(),
1140+
)?;
11111141

11121142
if !l1_config.relay_bid_scrapers.is_empty() {
11131143
let sender = Arc::new(BiddingService2BidSender::new(bidding_service.clone()));
@@ -1126,7 +1156,12 @@ where
11261156
relay_sets,
11271157
);
11281158

1129-
Ok((sink_factory, slot_info_provider, adjustment_fee_payers))
1159+
Ok((
1160+
sink_factory,
1161+
slot_info_provider,
1162+
adjustment_fee_payers,
1163+
optimistic_v3_server_join_handle,
1164+
))
11301165
}
11311166

11321167
/// Take the end of the pipeline (sink_factory) + pre-created slot_info_provider and creates an empty builder (it still needs the with_builders to be called)

crates/rbuilder/src/mev_boost/optimistic_v3.rs

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
use crate::{
2+
live_builder::process_killer::OPTIMISTIC_V3_CLOSE_TIME_SECONDS,
23
telemetry::{exponential_buckets_range, REGISTRY},
34
utils,
45
};
@@ -19,6 +20,7 @@ use std::{
1920
sync::Arc,
2021
time::{Duration, Instant, SystemTime},
2122
};
23+
use tokio::task::JoinHandle;
2224
use tokio_util::sync::CancellationToken;
2325
use tracing::*;
2426
use warp::{
@@ -88,7 +90,7 @@ pub fn spawn_server(
8890
relay_pubkeys: HashSet<BlsPublicKey>,
8991
blocks: OptimisticV3BlockCache,
9092
cancellation: CancellationToken,
91-
) -> eyre::Result<()> {
93+
) -> eyre::Result<JoinHandle<()>> {
9294
// Spawn relay server.
9395
let handler = Handler {
9496
domain,
@@ -111,13 +113,13 @@ pub fn spawn_server(
111113
let now = std::time::Instant::now();
112114
info!(target: "relay_server", "Received cancellation, initiating graceful shutdown");
113115
// Sleep for 12 seconds to avoid being demoted for the current slot.
114-
tokio::time::sleep(Duration::from_secs(12)).await;
116+
tokio::time::sleep(Duration::from_secs(OPTIMISTIC_V3_CLOSE_TIME_SECONDS)).await;
115117
info!(target: "relay_server", elapsed = ?now.elapsed(), "Graceful shutdown complete");
116118
});
117-
tokio::spawn(server);
119+
let res = tokio::spawn(server);
118120
info!(target: "relay_server", %address, "Relay server listening");
119121

120-
Ok(())
122+
Ok(res)
121123
}
122124

123125
#[derive(Clone, Debug)]

0 commit comments

Comments
 (0)