diff --git a/Cargo.lock b/Cargo.lock index 7fb871df..f1359fa9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1928,6 +1928,30 @@ version = "0.8.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e764a1d40d510daf35e07be9eb06e75770908c27d411ee6c92109c9840eaaf7" +[[package]] +name = "bitcode" +version = "0.6.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0a6ed1b54d8dc333e7be604d00fa9262f4635485ffea923647b6521a5fff045d" +dependencies = [ + "arrayvec", + "bitcode_derive", + "bytemuck", + "glam", + "serde", +] + +[[package]] +name = "bitcode_derive" +version = "0.6.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "238b90427dfad9da4a9abd60f3ec1cdee6b80454bde49ed37f1781dd8e9dc7f9" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.111", +] + [[package]] name = "bitcoin-io" version = "0.1.4" @@ -3370,6 +3394,15 @@ dependencies = [ "subtle", ] +[[package]] +name = "directories" +version = "5.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a49173b84e034382284f27f1af4dcbbd231ffa358c0fe316541a7337f376a35" +dependencies = [ + "dirs-sys 0.4.1", +] + [[package]] name = "dirs" version = "5.0.1" @@ -4442,6 +4475,86 @@ dependencies = [ "miniz_oxide", ] +[[package]] +name = "flux" +version = "0.0.36" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "bitcode", + "core_affinity", + "flux-communication", + "flux-network", + "flux-timing", + "flux-utils", + "libc", + "once_cell", + "serde", + "signal-hook", + "spine-derive", + "tinyvec", + "tracing", + "type-hash", + "type-hash-derive", + "zstd 0.13.3", +] + +[[package]] +name = "flux-communication" +version = "0.0.26" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "directories", + "flux-timing", + "flux-utils", + "once_cell", + "shared_memory", + "thiserror 1.0.69", + "tracing", +] + +[[package]] +name = "flux-network" +version = "0.0.7" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "flux-communication", + "flux-timing", + "flux-utils", + "mio", + "tracing", +] + +[[package]] +name = "flux-timing" +version = "0.0.14" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "bitcode", + "chrono", + "governor 0.6.3", + "humantime", + "once_cell", + "quanta", + "serde", + "serde_json", + "tracing", + "type-hash", + "type-hash-derive", +] + +[[package]] +name = "flux-utils" +version = "0.0.22" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "core_affinity", + "directories", + "flux-timing", + "libc", + "serde", + "tracing", +] + [[package]] name = "fnv" version = "1.0.7" @@ -4723,6 +4836,12 @@ dependencies = [ "url", ] +[[package]] +name = "glam" +version = "0.31.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74a4d85559e2637d3d839438b5b3d75c31e655276f9544d72475c36b92fabbed" + [[package]] name = "glob" version = "0.3.3" @@ -4797,6 +4916,26 @@ dependencies = [ "windows-sys 0.60.2", ] +[[package]] +name = "governor" +version = "0.6.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68a7f542ee6b35af73b06abc0dad1c1bae89964e4e253bc4b587b91c9637867b" +dependencies = [ + "cfg-if", + "dashmap 5.5.3", + "futures", + "futures-timer", + "no-std-compat", + "nonzero_ext", + "parking_lot", + "portable-atomic", + "quanta", + "rand 0.8.5", + "smallvec", + "spinning_top", +] + [[package]] name = "governor" version = "0.8.1" @@ -5044,6 +5183,7 @@ dependencies = [ "ethers", "eyre", "flate2", + "flux", "futures", "helix-common", "helix-types", @@ -6602,6 +6742,15 @@ dependencies = [ "libc", ] +[[package]] +name = "memoffset" +version = "0.6.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5aa361d4faea93603064a027415f07bd8e1d5c88c9fbf68bf56a285428fd79ce" +dependencies = [ + "autocfg", +] + [[package]] name = "merkle_proof" version = "0.2.0" @@ -6933,6 +7082,19 @@ version = "1.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" +[[package]] +name = "nix" +version = "0.23.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f3790c00a0150112de0f4cd161e3d7fc4b2d8a5542ffc35f099a2562aecb35c" +dependencies = [ + "bitflags 1.3.2", + "cc", + "cfg-if", + "libc", + "memoffset", +] + [[package]] name = "no-std-compat" version = "0.4.1" @@ -12454,6 +12616,19 @@ dependencies = [ "lazy_static", ] +[[package]] +name = "shared_memory" +version = "0.12.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba8593196da75d9dc4f69349682bd4c2099f8cde114257d1ef7ef1b33d1aba54" +dependencies = [ + "cfg-if", + "libc", + "nix", + "rand 0.8.5", + "win-sys", +] + [[package]] name = "shellexpand" version = "3.1.1" @@ -12654,6 +12829,16 @@ version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +[[package]] +name = "spine-derive" +version = "0.0.8" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.111", +] + [[package]] name = "spinning_top" version = "0.3.0" @@ -13689,7 +13874,7 @@ checksum = "57a2ccff6830fa835371af7541e561a90e4c07b84f72991ebac4b3cb6790dc0d" dependencies = [ "axum 0.8.7", "forwarded-header-value", - "governor", + "governor 0.8.1", "http 1.4.0", "pin-project", "thiserror 2.0.17", @@ -13977,6 +14162,22 @@ dependencies = [ "utf-8", ] +[[package]] +name = "type-hash" +version = "0.0.1" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" + +[[package]] +name = "type-hash-derive" +version = "0.0.4" +source = "git+https://github.com/gattaca-com/flux#16973424f7eb3eac68f80cc305a03b948486805d" +dependencies = [ + "proc-macro-crate", + "proc-macro2", + "quote", + "syn 2.0.111", +] + [[package]] name = "typenum" version = "1.19.0" @@ -14505,6 +14706,15 @@ version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72069c3113ab32ab29e5584db3c6ec55d416895e60715417b5b883a357c3e471" +[[package]] +name = "win-sys" +version = "0.3.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b7b128a98c1cfa201b09eb49ba285887deb3cbe7466a98850eb1adabb452be5" +dependencies = [ + "windows 0.34.0", +] + [[package]] name = "winapi" version = "0.3.9" @@ -14536,6 +14746,19 @@ version = "0.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "712e227841d057c1ee1cd2fb22fa7e5a5461ae8e48fa2ca79ec42cfc1931183f" +[[package]] +name = "windows" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "45296b64204227616fdbf2614cefa4c236b98ee64dfaaaa435207ed99fe7829f" +dependencies = [ + "windows_aarch64_msvc 0.34.0", + "windows_i686_gnu 0.34.0", + "windows_i686_msvc 0.34.0", + "windows_x86_64_gnu 0.34.0", + "windows_x86_64_msvc 0.34.0", +] + [[package]] name = "windows" version = "0.57.0" @@ -14840,6 +15063,12 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53" +[[package]] +name = "windows_aarch64_msvc" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "17cffbe740121affb56fad0fc0e421804adf0ae00891205213b5cecd30db881d" + [[package]] name = "windows_aarch64_msvc" version = "0.42.2" @@ -14864,6 +15093,12 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006" +[[package]] +name = "windows_i686_gnu" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2564fde759adb79129d9b4f54be42b32c89970c18ebf93124ca8870a498688ed" + [[package]] name = "windows_i686_gnu" version = "0.42.2" @@ -14900,6 +15135,12 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c" +[[package]] +name = "windows_i686_msvc" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9cd9d32ba70453522332c14d38814bceeb747d80b3958676007acadd7e166956" + [[package]] name = "windows_i686_msvc" version = "0.42.2" @@ -14924,6 +15165,12 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2" +[[package]] +name = "windows_x86_64_gnu" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cfce6deae227ee8d356d19effc141a509cc503dfd1f850622ec4b0f84428e1f4" + [[package]] name = "windows_x86_64_gnu" version = "0.42.2" @@ -14972,6 +15219,12 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1" +[[package]] +name = "windows_x86_64_msvc" +version = "0.34.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d19538ccc21819d01deaf88d6a17eae6596a12e9aafdbb97916fb49896d89de9" + [[package]] name = "windows_x86_64_msvc" version = "0.42.2" diff --git a/Cargo.toml b/Cargo.toml index ce15321a..df7a8378 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,7 +7,7 @@ resolver = "2" edition = "2024" license = "MIT OR Apache-2.0" repository = "https://github.com/gattaca-com/helix" -rust-version = "1.90.0" +rust-version = "1.91.0" version = "0.0.1" [workspace.dependencies] @@ -40,6 +40,7 @@ ethereum_ssz_derive = "0.9" ethers = "2.0.14" eyre = "0.6.12" flate2 = "1.0" +flux = { git = "https://github.com/gattaca-com/flux" } futures = "0.3" futures-util = { version = "0.3", features = ["compat"] } helix-common = { path = "crates/common" } diff --git a/crates/relay/Cargo.toml b/crates/relay/Cargo.toml index 945e6181..dbca1093 100644 --- a/crates/relay/Cargo.toml +++ b/crates/relay/Cargo.toml @@ -25,6 +25,7 @@ ethereum_ssz_derive.workspace = true ethers.workspace = true eyre.workspace = true flate2.workspace = true +flux.workspace = true futures.workspace = true helix-common.workspace = true helix-types.workspace = true diff --git a/crates/relay/src/api/service.rs b/crates/relay/src/api/service.rs index 1e3f419b..653bce5f 100644 --- a/crates/relay/src/api/service.rs +++ b/crates/relay/src/api/service.rs @@ -12,7 +12,7 @@ use moka::sync::Cache; use tracing::{error, info}; use crate::{ - BidAdjustor, + AuctioneerHandle, RegWorkerHandle, api::{ Api, builder::api::BuilderApi, @@ -20,7 +20,6 @@ use crate::{ relay_data::{BidsCache, DataApi, DeliveredPayloadsCache, SelectiveExpiry}, router::build_router, }, - auctioneer::{Event, spawn_workers}, beacon::multi_beacon_client::MultiBeaconClient, database::postgres::postgres_db_service::PostgresDatabaseService, gossip::{GrpcGossiperClientManager, process_gossip_messages}, @@ -31,7 +30,41 @@ use crate::{ pub(crate) const API_REQUEST_TIMEOUT: Duration = Duration::from_secs(5); pub(crate) const SIMULATOR_REQUEST_TIMEOUT: Duration = Duration::from_secs(20); -pub async fn start_api_service( +pub fn start_api_service( + config: RelayConfig, + db: Arc, + local_cache: Arc, + current_slot_info: CurrentSlotInfo, + chain_info: Arc, + relay_signing_context: Arc, + multi_beacon_client: Arc, + api_provider: Arc, + known_validators_loaded: Arc, + terminating: Arc, + top_bid_tx: tokio::sync::broadcast::Sender, + relay_network_api: RelayNetworkApi, + auctioneer_handle: AuctioneerHandle, + registrations_handle: RegWorkerHandle, +) { + tokio::spawn(run_api_service::( + config.clone(), + db.clone(), + local_cache.clone(), + current_slot_info, + chain_info.clone(), + relay_signing_context, + multi_beacon_client, + api_provider, + known_validators_loaded, + terminating.clone(), + top_bid_tx.clone(), + relay_network_api, + auctioneer_handle, + registrations_handle, + )); +} + +pub async fn run_api_service( mut config: RelayConfig, db: Arc, local_cache: Arc, @@ -40,12 +73,12 @@ pub async fn start_api_service( relay_signing_context: Arc, multi_beacon_client: Arc, api_provider: Arc, - bid_adjustor: impl BidAdjustor, known_validators_loaded: Arc, terminating: Arc, top_bid_tx: tokio::sync::broadcast::Sender, - event_channel: (crossbeam_channel::Sender, crossbeam_channel::Receiver), relay_network_api: RelayNetworkApi, + auctioneer_handle: AuctioneerHandle, + registrations_handle: RegWorkerHandle, ) { let gossiper = Arc::new( GrpcGossiperClientManager::new(config.relays.iter().map(|cfg| cfg.url.clone()).collect()) @@ -57,17 +90,6 @@ pub async fn start_api_service( let (gossip_sender, gossip_receiver) = tokio::sync::mpsc::channel(10_000); - // spawn auctioneer - let (auctioneer_handle, registrations_handle) = spawn_workers( - Arc::unwrap_or_clone(chain_info.clone()), - config.clone(), - db.clone(), - Arc::unwrap_or_clone(local_cache.clone()), - bid_adjustor, - top_bid_tx.clone(), - event_channel, - ); - let builder_api = BuilderApi::::new( local_cache.clone(), db.clone(), diff --git a/crates/relay/src/auctioneer/mod.rs b/crates/relay/src/auctioneer/mod.rs index 997421fb..4c7edff5 100644 --- a/crates/relay/src/auctioneer/mod.rs +++ b/crates/relay/src/auctioneer/mod.rs @@ -20,128 +20,89 @@ use std::{ use alloy_primitives::B256; pub use block_merger::OrderValidationError; +use flux::tile::Tile; pub use handle::{AuctioneerHandle, RegWorkerHandle}; use helix_common::{ RelayConfig, - api::builder_api::{ - BuilderGetValidatorsResponseEntry, InclusionListWithMetadata, TopBidUpdate, - }, + api::builder_api::{BuilderGetValidatorsResponseEntry, InclusionListWithMetadata}, chain_info::ChainInfo, local_cache::LocalCache, metrics::{STATE_TRANSITION_COUNT, STATE_TRANSITION_LATENCY, WORKER_QUEUE_LEN, WORKER_UTIL}, record_submission_step, - utils::pin_thread_to_core, }; use helix_types::Slot; use rustc_hash::FxHashMap; pub use simulator::*; -use tracing::{debug, error, info, info_span, trace, warn}; +use tracing::{debug, error, info, trace, warn}; pub use types::{ Event, GetPayloadResultData, PayloadBidData, PayloadEntry, SlotData, SubmissionPayload, }; -use worker::{RegWorker, SubWorker}; +pub use worker::{RegWorker, SubWorker}; pub use crate::auctioneer::{ bid_adjustor::{BidAdjustor, DefaultBidAdjustor}, - simulator::{SimulatorRequest, client::SimulatorClient}, + bid_sorter::BidSorter, + context::Context, + simulator::{SimulatorRequest, client::SimulatorClient, manager::SimulatorManager}, }; use crate::{ - PostgresDatabaseService, + HelixSpine, PostgresDatabaseService, api::{builder::error::BuilderApiError, proposer::ProposerApiError}, - auctioneer::{ - bid_sorter::BidSorter, context::Context, manager::SimulatorManager, types::PendingPayload, - }, + auctioneer::types::PendingPayload, housekeeper::PayloadAttributesUpdate, }; -// TODO: tidy up builder and proposer api state, and spawn in a separate function -pub fn spawn_workers( - chain_info: ChainInfo, - config: RelayConfig, - db: Arc, - cache: LocalCache, - bid_adjustor: B, - top_bid_tx: tokio::sync::broadcast::Sender, - event_channel: (crossbeam_channel::Sender, crossbeam_channel::Receiver), -) -> (AuctioneerHandle, RegWorkerHandle) { - let (sub_worker_tx, sub_worker_rx) = crossbeam_channel::bounded(10_000); - let (reg_worker_tx, reg_worker_rx) = crossbeam_channel::bounded(100_000); - let (event_tx, event_rx) = event_channel; - - if config.is_registration_instance { - for core in config.cores.reg_workers.clone() { - let worker = RegWorker::new(core, chain_info.clone()); - let rx = reg_worker_rx.clone(); - - std::thread::Builder::new() - .name(format!("worker-{core}")) - .spawn(move || { - pin_thread_to_core(core); - worker.run(rx) - }) - .unwrap(); - } - } - - if config.is_submission_instance { - for core in config.cores.sub_workers.clone() { - let worker = SubWorker::new( - core, - event_tx.clone(), - cache.clone(), - chain_info.clone(), - config.clone(), - ); - let rx = sub_worker_rx.clone(); - - std::thread::Builder::new() - .name(format!("worker-{core}")) - .spawn(move || { - pin_thread_to_core(core); - worker.run(rx) - }) - .unwrap(); - } - - let auctioneer_core = config.cores.auctioneer; - - let bid_sorter = BidSorter::new(top_bid_tx); - let sim_manager = SimulatorManager::new(config.simulators.clone(), event_tx.clone()); - let ctx = - Context::new(chain_info, config, sim_manager, db, bid_sorter, cache, bid_adjustor); - let id = format!("auctioneer_{auctioneer_core}"); - let auctioneer = Auctioneer { ctx, state: State::default(), tel: Telemetry::new(id) }; - - std::thread::Builder::new() - .name("auctioneer".to_string()) - .spawn(move || { - pin_thread_to_core(auctioneer_core); - auctioneer.run(event_rx) - }) - .unwrap(); - } - - (AuctioneerHandle::new(sub_worker_tx, event_tx), RegWorkerHandle::new(reg_worker_tx)) -} - -struct Auctioneer { +pub struct Auctioneer { ctx: Context, state: State, tel: Telemetry, + event_rx: crossbeam_channel::Receiver, } impl Auctioneer { - fn run(mut self, rx: crossbeam_channel::Receiver) { - let _span = info_span!("auctioneer").entered(); - info!("starting"); + pub fn new( + chain_info: ChainInfo, + config: RelayConfig, + db: Arc, + bid_sorter: BidSorter, + local_cache: LocalCache, + bid_adjustor: B, + event_tx: crossbeam_channel::Sender, + event_rx: crossbeam_channel::Receiver, + id: usize, + ) -> Self { + let sim_manager = SimulatorManager::new(config.simulators.clone(), event_tx.clone()); - loop { - for evt in rx.try_iter() { - self.state.step(evt, &mut self.ctx, &mut self.tel); - } + let ctx = Context::new( + chain_info, + config, + sim_manager, + db, + bid_sorter, + local_cache, + bid_adjustor, + ); + Self { + ctx, + state: State::default(), + tel: Telemetry::new(format!("auctioneer_{id}")), + event_rx, + } + } +} - self.tel.telemetry(&rx); +impl Tile for Auctioneer { + fn loop_body(&mut self, _adapter: &mut flux::spine::SpineAdapter) { + for event in self.event_rx.try_iter() { + self.state.step(event, &mut self.ctx, &mut self.tel); } + + self.tel.telemetry(&self.event_rx); + } + + fn try_init(&mut self, _adapter: &mut flux::spine::SpineAdapter) -> bool { + info!("starting"); + true } } diff --git a/crates/relay/src/auctioneer/worker.rs b/crates/relay/src/auctioneer/worker.rs index 7f8b69d8..575d88e8 100644 --- a/crates/relay/src/auctioneer/worker.rs +++ b/crates/relay/src/auctioneer/worker.rs @@ -1,6 +1,10 @@ use std::time::{Duration, Instant}; use alloy_primitives::B256; +use flux::{ + tile::{Tile, TileName}, + utils::{ShortTypename, short_typename}, +}; use helix_common::{ GetPayloadTrace, RelayConfig, SubmissionTrace, chain_info::ChainInfo, @@ -17,9 +21,10 @@ use helix_types::{ SubmissionVersion, }; use http::HeaderValue; -use tracing::{error, info, info_span, trace}; +use tracing::{error, trace}; use crate::{ + HelixSpine, api::{ HEADER_API_KEY, HEADER_HYDRATE, HEADER_IS_MERGEABLE, HEADER_SEQUENCE, HEADER_WITH_ADJUSTMENTS, builder::error::BuilderApiError, proposer::ProposerApiError, @@ -87,9 +92,11 @@ impl Default for Telemetry { } // TODO: spans -pub(super) struct SubWorker { - id: String, +pub struct SubWorker { + core_id: usize, + id: ShortTypename, tx: crossbeam_channel::Sender, + rx: crossbeam_channel::Receiver, cache: LocalCache, chain_info: ChainInfo, tel: Telemetry, @@ -98,48 +105,15 @@ pub(super) struct SubWorker { impl SubWorker { pub fn new( - id: usize, + core_id: usize, tx: crossbeam_channel::Sender, + rx: crossbeam_channel::Receiver, cache: LocalCache, chain_info: ChainInfo, config: RelayConfig, ) -> Self { - Self { - id: format!("submission_{id}"), - tx, - cache, - chain_info, - tel: Telemetry::default(), - config, - } - } - - pub(super) fn run(mut self, rx: crossbeam_channel::Receiver) { - let _span = info_span!("worker", id = self.id).entered(); - info!("starting"); - - loop { - for task in rx.try_iter() { - self.handle_task(task); - } - - self.tel.telemetry(&self.id, "submission", &rx); - } - } - - fn handle_task(&mut self, task: SubWorkerJob) { - let start_task = Instant::now(); - let task_name = task.as_str(); - - self._handle_task(task); - - let task_dur = start_task.elapsed(); - self.tel.loop_worked += task_dur; - - WORKER_TASK_COUNT.with_label_values(&[task_name, &self.id]).inc(); - WORKER_TASK_LATENCY_US - .with_label_values(&[task_name, &self.id]) - .observe(task_dur.as_micros() as f64); + let id = ShortTypename::from_str_truncate(&format!("submission_{core_id}")); + Self { core_id, id, tx, rx, cache, chain_info, tel: Telemetry::default(), config } } fn _handle_task(&self, task: SubWorkerJob) { @@ -243,7 +217,7 @@ impl SubWorker { }) .is_err() { - error!("failed sending get_payload to auctioneer"); + error!("failed to send get_payload to auctioneer"); }; } Err(err) => { @@ -348,29 +322,48 @@ impl SubWorker { } } +impl Tile for SubWorker { + fn loop_body(&mut self, _adapter: &mut flux::spine::SpineAdapter) { + for task in self.rx.try_iter() { + let start_task = Instant::now(); + let task_name = task.as_str(); + + self._handle_task(task); + + let task_dur = start_task.elapsed(); + self.tel.loop_worked += task_dur; + + WORKER_TASK_COUNT.with_label_values(&[task_name, &self.id]).inc(); + WORKER_TASK_LATENCY_US + .with_label_values(&[task_name, &self.id]) + .observe(task_dur.as_micros() as f64); + } + + self.tel.telemetry(&self.id, "submission", &self.rx); + } + + fn name(&self) -> TileName { + TileName::from_str_truncate(&format!("{}_{}", short_typename::(), self.core_id)) + } +} + /// Worker to process registrations verifications -pub(super) struct RegWorker { - id: String, +pub struct RegWorker { + core_id: usize, + id: ShortTypename, chain_info: ChainInfo, tel: Telemetry, + rx: crossbeam_channel::Receiver, } impl RegWorker { - pub(super) fn new(id: usize, chain_info: ChainInfo) -> Self { - Self { id: format!("registration_{id}"), chain_info, tel: Default::default() } - } - - pub(super) fn run(mut self, rx: crossbeam_channel::Receiver) { - let _span = info_span!("worker", id = self.id).entered(); - info!("starting"); - - loop { - if let Ok(task) = rx.recv_timeout(Duration::from_millis(50)) { - self.handle_task(task); - } - - self.tel.telemetry(&self.id, "registration", &rx); - } + pub fn new( + core_id: usize, + chain_info: ChainInfo, + rx: crossbeam_channel::Receiver, + ) -> Self { + let id = ShortTypename::from_str_truncate(&format!("registration_{core_id}")); + Self { core_id, id, chain_info, tel: Default::default(), rx } } fn handle_task(&mut self, task: RegWorkerJob) { @@ -414,6 +407,20 @@ impl RegWorker { } } +impl Tile for RegWorker { + fn loop_body(&mut self, _adapter: &mut flux::spine::SpineAdapter) { + if let Ok(task) = self.rx.recv_timeout(Duration::from_millis(50)) { + self.handle_task(task); + } + + self.tel.telemetry(&self.id, "registration", &self.rx); + } + + fn name(&self) -> TileName { + TileName::from_str_truncate(&format!("{}_{}", short_typename::(), self.core_id)) + } +} + struct DecodeFlags { skip_sigverify: bool, merge_type: Option, diff --git a/crates/relay/src/beacon/types/chain.rs b/crates/relay/src/beacon/types/chain.rs index fd2c7c8f..738a3368 100644 --- a/crates/relay/src/beacon/types/chain.rs +++ b/crates/relay/src/beacon/types/chain.rs @@ -62,7 +62,6 @@ pub struct SyncStatus { pub is_syncing: bool, } -// HeadEventData represents the data of a head event // {"slot":"827256","block":"0x56b683afa68170c775f3c9debc18a6a72caea9055584d037333a6fe43c8ceb83","state":"0x419e2965320d69c4213782dae73941de802a4f436408fddd6f68b671b3ff4e55","epoch_transition":false,"execution_optimistic":false,"previous_duty_dependent_root":"0x5b81a526839b7fb67c3896f1125451755088fb578ad27c2690b3209f3d7c6b54","current_duty_dependent_root":"0x5f3232c0d5741e27e13754e1d88285c603b07dd6164b35ca57e94344a9e42942"} #[derive(Serialize, Deserialize, Debug, Clone, Default)] pub struct HeadEventData { diff --git a/crates/relay/src/housekeeper/chain_event_updater.rs b/crates/relay/src/housekeeper/chain_event_updater.rs index 6492ff70..50fc08a9 100644 --- a/crates/relay/src/housekeeper/chain_event_updater.rs +++ b/crates/relay/src/housekeeper/chain_event_updater.rs @@ -21,7 +21,7 @@ use crate::{ }; // Do not accept slots more than 60 seconds in the future -const MAX_DISTANCE_FOR_FUTURE_SLOT: u64 = 60; +const MAX_DISTANCE_FOR_SLOT_SEC: u64 = 30; const CUTOFF_TIME: u64 = 4; /// Payload for a new payload attribute event sent to subscribers. @@ -174,15 +174,18 @@ impl ChainEventUpdater { info!(head_slot =% slot, "processing slot"); - // Validate this isn't a faulty head slot let slot_timestamp = self.chain_info.genesis_time_in_secs + (slot * self.chain_info.seconds_per_slot()); - if slot_timestamp > utcnow_sec() + MAX_DISTANCE_FOR_FUTURE_SLOT { + + let now = utcnow_sec(); + if slot_timestamp < now - MAX_DISTANCE_FOR_SLOT_SEC { + warn!(head_slot = slot, "slot is too old"); + return; + } else if slot_timestamp > now + MAX_DISTANCE_FOR_SLOT_SEC { warn!(head_slot = slot, "slot is too far in the future"); return; } - // Log missed slots if self.head_slot != 0 { for s in (self.head_slot + 1)..slot { warn!(missed_slot = s, "missed slot"); diff --git a/crates/relay/src/lib.rs b/crates/relay/src/lib.rs index 067c5906..9e563fcf 100644 --- a/crates/relay/src/lib.rs +++ b/crates/relay/src/lib.rs @@ -5,14 +5,20 @@ mod database; mod gossip; mod housekeeper; mod network; +mod spine; mod website; pub use crate::{ api::{Api, BidAdjustor, DefaultBidAdjustor, start_admin_service, start_api_service}, - auctioneer::{PayloadEntry, SimulatorClient, SimulatorRequest, SlotData, SubmissionPayload}, + auctioneer::{ + Auctioneer, AuctioneerHandle, BidSorter, Context, PayloadEntry, RegWorker, RegWorkerHandle, + SimulatorClient, SimulatorManager, SimulatorRequest, SlotData, SubWorker, + SubmissionPayload, + }, beacon::start_beacon_client, database::{postgres::postgres_db_service::PostgresDatabaseService, start_db_service}, housekeeper::start_housekeeper, network::RelayNetworkManager, + spine::HelixSpine, website::WebsiteService, }; diff --git a/crates/relay/src/main.rs b/crates/relay/src/main.rs index f38fc7ea..39c3753f 100644 --- a/crates/relay/src/main.rs +++ b/crates/relay/src/main.rs @@ -7,6 +7,10 @@ use std::{ }; use eyre::eyre; +use flux::{ + tile::{TileConfig, attach_tile}, + utils::ThreadPriority, +}; use helix_common::{ RelayConfig, api_provider::DefaultApiProvider, @@ -18,13 +22,14 @@ use helix_common::{ utils::{init_panic_hook, init_tracing_log}, }; use helix_relay::{ - Api, DefaultBidAdjustor, PostgresDatabaseService, RelayNetworkManager, WebsiteService, - start_admin_service, start_api_service, start_beacon_client, start_db_service, + Api, Auctioneer, AuctioneerHandle, BidSorter, DefaultBidAdjustor, HelixSpine, + PostgresDatabaseService, RegWorker, RegWorkerHandle, RelayNetworkManager, SubWorker, + WebsiteService, start_admin_service, start_api_service, start_beacon_client, start_db_service, start_housekeeper, }; use helix_types::BlsKeypair; use tikv_jemallocator::Jemalloc; -use tokio::signal::unix::SignalKind; +use tokio::signal::unix::{SignalKind, signal}; use tracing::{error, info}; #[global_allocator] @@ -74,8 +79,7 @@ fn main() { async fn run(instance_id: String, config: RelayConfig, keypair: BlsKeypair) -> eyre::Result<()> { let beacon_client = start_beacon_client(&config); - let chain_info = beacon_client.load_chain_info().await; - let chain_info = Arc::new(chain_info); + let chain_info = Arc::new(beacon_client.load_chain_info().await); info!( instance_id, @@ -85,88 +89,137 @@ async fn run(instance_id: String, config: RelayConfig, keypair: BlsKeypair) -> e "starting relay" ); - let relay_signing_context = Arc::new(RelaySigningContext::new(keypair, chain_info.clone())); - let known_validators_loaded = Arc::new(AtomicBool::default()); - let db = start_db_service(&config, known_validators_loaded.clone()).await?; - let local_cache = start_auctioneer(db.clone()).await?; + let local_cache = start_local_cache_reloader(db.clone()).await?; + + let relay_signing_context = Arc::new(RelaySigningContext::new(keypair, chain_info.clone())); - let event_channel = crossbeam_channel::bounded(10_000); let relay_network_api = RelayNetworkManager::new(config.relay_network.clone(), relay_signing_context.clone()); - let (top_bid_tx, _) = tokio::sync::broadcast::channel(100); - config.router_config.validate_bid_sorter()?; + let (sub_worker_tx, sub_worker_rx) = crossbeam_channel::bounded(10_000); + let (reg_worker_tx, reg_worker_rx) = crossbeam_channel::bounded(100_000); + + let (event_tx, event_rx) = crossbeam_channel::bounded(10_000); + let current_slot_info = start_housekeeper( db.clone(), local_cache.clone(), &config, beacon_client.clone(), chain_info.clone(), - event_channel.0.clone(), + event_tx.clone(), relay_network_api.clone(), ) .await .map_err(|e| eyre!("housekeeper init: {e}"))?; let terminating = Arc::new(AtomicBool::default()); + let termination_grace_period = config.router_config.shutdown_delay_ms; - start_admin_service(local_cache.clone(), expect_env_var(ADMIN_TOKEN_ENV_VAR)); - - tokio::spawn(start_api_service::( - config.clone(), - db.clone(), - local_cache, - current_slot_info, - chain_info, - relay_signing_context, - beacon_client, - Arc::new(DefaultApiProvider {}), - DefaultBidAdjustor {}, - known_validators_loaded, - terminating.clone(), - top_bid_tx, - event_channel, - relay_network_api.api(), - )); + let spine = HelixSpine::new(None); + spine.start(None, |spine| { + start_admin_service(local_cache.clone(), expect_env_var(ADMIN_TOKEN_ENV_VAR)); + + let auctioneer_handle = AuctioneerHandle::new(sub_worker_tx, event_tx.clone()); + let registrations_handle = RegWorkerHandle::new(reg_worker_tx); + + let (top_bid_tx, _) = tokio::sync::broadcast::channel(100); + + start_api_service::( + config.clone(), + db.clone(), + local_cache.clone(), + current_slot_info, + chain_info.clone(), + relay_signing_context, + beacon_client, + Arc::new(DefaultApiProvider {}), + known_validators_loaded, + terminating.clone(), + top_bid_tx.clone(), + relay_network_api.api(), + auctioneer_handle, + registrations_handle, + ); + + if config.website.enabled { + tokio::spawn(WebsiteService::run_loop(config.clone(), db.clone())); + } - let termination_grace_period = config.router_config.shutdown_delay_ms; + if config.is_registration_instance { + for core in config.cores.reg_workers.clone() { + let worker = + RegWorker::new(core, chain_info.as_ref().clone(), reg_worker_rx.clone()); - if config.website.enabled { - tokio::spawn(WebsiteService::run_loop(config, db)); - } + attach_tile(worker, spine, TileConfig::new(core, ThreadPriority::OSDefault)); + } + } - // wait for SIGTERM or SIGINT - let mut sigint = tokio::signal::unix::signal(SignalKind::interrupt())?; - let mut sigterm = tokio::signal::unix::signal(SignalKind::terminate())?; + if config.is_submission_instance { + for core in config.cores.sub_workers.clone() { + let worker = SubWorker::new( + core, + event_tx.clone(), + sub_worker_rx.clone(), + local_cache.as_ref().clone(), + chain_info.as_ref().clone(), + config.clone(), + ); + + attach_tile(worker, spine, TileConfig::new(core, ThreadPriority::OSDefault)); + } + + let auctioneer_core = config.cores.auctioneer; + let auctioneer = Auctioneer::new( + chain_info.as_ref().clone(), + config, + db, + BidSorter::new(top_bid_tx), + local_cache.as_ref().clone(), + DefaultBidAdjustor {}, + event_tx, + event_rx, + auctioneer_core, + ); + attach_tile( + auctioneer, + spine, + TileConfig::new(auctioneer_core, ThreadPriority::OSDefault), + ); + } + }); + let mut sigint = signal(SignalKind::interrupt())?; + let mut sigterm = signal(SignalKind::terminate())?; tokio::select! { _ = sigint.recv() => {} _ = sigterm.recv() => {} } - // Set terminating flag. terminating.store(true, Ordering::Relaxed); if termination_grace_period != 0 { - // Wait for the grace period to expire before exiting. tracing::info!("Pausing for {termination_grace_period}ms before exit"); + tokio::time::sleep(Duration::from_millis(termination_grace_period)).await; } Ok(()) } -pub async fn start_auctioneer(db: Arc) -> eyre::Result> { - let auctioneer = Arc::new(LocalCache::new()); - let auctioneer_clone = auctioneer.clone(); +pub async fn start_local_cache_reloader( + db: Arc, +) -> eyre::Result> { + let local_cache = Arc::new(LocalCache::new()); + let local_cache_clone = local_cache.clone(); tokio::spawn(async move { let builder_infos = db.get_all_builder_infos().await.expect("failed to load builder infos"); - auctioneer_clone.update_builder_infos(&builder_infos, true); + local_cache_clone.update_builder_infos(&builder_infos, true); }); - Ok(auctioneer) + Ok(local_cache) } diff --git a/crates/relay/src/spine/messages.rs b/crates/relay/src/spine/messages.rs new file mode 100644 index 00000000..c273ff10 --- /dev/null +++ b/crates/relay/src/spine/messages.rs @@ -0,0 +1,3 @@ +#[derive(Copy, Clone, Debug)] +#[repr(C)] +pub struct Tmp {} diff --git a/crates/relay/src/spine/mod.rs b/crates/relay/src/spine/mod.rs new file mode 100644 index 00000000..de765cc0 --- /dev/null +++ b/crates/relay/src/spine/mod.rs @@ -0,0 +1,12 @@ +pub mod messages; + +use flux::{communication::ShmemData, spine::SpineQueue, spine_derive::from_spine, tile::TileInfo}; + +#[from_spine("helix")] +#[derive(Debug)] +pub struct HelixSpine { + pub tile_info: ShmemData, + + #[queue(size(2usize.pow(16)))] + pub filler: SpineQueue, +} diff --git a/relay.Dockerfile b/relay.Dockerfile index c9d8ff7e..2f9a6065 100644 --- a/relay.Dockerfile +++ b/relay.Dockerfile @@ -1,4 +1,4 @@ -FROM lukemathwalker/cargo-chef:latest-rust-1.90 AS chef +FROM lukemathwalker/cargo-chef:latest-rust-1.91 AS chef WORKDIR /app FROM chef AS planner diff --git a/rust-toolchain.toml b/rust-toolchain.toml index 2ce412d5..1612370a 100644 --- a/rust-toolchain.toml +++ b/rust-toolchain.toml @@ -1,3 +1,3 @@ [toolchain] -channel = "1.90.0" +channel = "1.91.0" profile = "default" diff --git a/simulator.Dockerfile b/simulator.Dockerfile index 78022ba5..e47ebcf7 100644 --- a/simulator.Dockerfile +++ b/simulator.Dockerfile @@ -1,4 +1,4 @@ -FROM lukemathwalker/cargo-chef:latest-rust-1.90 AS chef +FROM lukemathwalker/cargo-chef:latest-rust-1.91 AS chef WORKDIR /app # Install libclang and dependencies required by bindgen