From d9253aa5ca297c2a932c4c6bdb11190c825e14cf Mon Sep 17 00:00:00 2001 From: Vladimir Petrzhikovskii Date: Fri, 21 Feb 2025 11:58:46 +0100 Subject: [PATCH 1/2] refactor(core): centralize `blockchain_rpc` methods into a single enum Reduce boilerplate for metrics. --- .../{service.rs => service/mod.rs} | 135 +++++++----------- core/src/blockchain_rpc/service/util.rs | 43 ++++++ scripts/gen-dashboard.py | 50 +++++-- util/src/metrics/histogram_guard.rs | 54 +++++++ 4 files changed, 184 insertions(+), 98 deletions(-) rename core/src/blockchain_rpc/{service.rs => service/mod.rs} (84%) create mode 100644 core/src/blockchain_rpc/service/util.rs diff --git a/core/src/blockchain_rpc/service.rs b/core/src/blockchain_rpc/service/mod.rs similarity index 84% rename from core/src/blockchain_rpc/service.rs rename to core/src/blockchain_rpc/service/mod.rs index ba129dbfa2..8dfd3e7870 100644 --- a/core/src/blockchain_rpc/service.rs +++ b/core/src/blockchain_rpc/service/mod.rs @@ -1,3 +1,5 @@ +mod util; + use std::num::{NonZeroU32, NonZeroU64}; use std::sync::Arc; @@ -5,6 +7,7 @@ use anyhow::Context; use bytes::{Buf, Bytes}; use everscale_types::models::BlockId; use futures_util::Future; +use metrics::Label; use serde::{Deserialize, Serialize}; use tycho_block_util::message::validate_external_message; use tycho_network::{try_handle_prefix, InboundRequestMeta, Response, Service, ServiceRequest}; @@ -161,7 +164,19 @@ impl Service for BlockchainRpcService { } }; - tycho_network::match_tl_request!(body, tag = constructor, { + let method = util::Constructor::from_tl_id(constructor); + let label = vec![Label::new( + "method", + method.map_or("unknown", |m| m.as_str()), + )]; + let timer = + move || HistogramGuard::begin_with_labels_owned(RPC_METHOD_TIMINGS_METRIC, label); + + let inner = self.inner.clone(); + + // NOTE: update `constructor_to_string` after adding new methods + tycho_network::match_tl_request!(body, tag = constructor, + { overlay::Ping as _ => BoxFutureOrNoop::future(async { Some(Response::from_tl(overlay::Pong)) }), @@ -171,65 +186,44 @@ impl Service for BlockchainRpcService { max_size = req.max_size, "getNextKeyBlockIds", ); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_next_key_block_ids(&req); - Some(Response::from_tl(res)) + timed_future(timer, async move { + Some(Response::from_tl(inner.handle_get_next_key_block_ids(&req))) }) }, rpc::GetBlockFull as req => { tracing::debug!(block_id = %req.block_id, "getBlockFull"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_block_full(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_block_full(&req).await)) }) }, rpc::GetNextBlockFull as req => { tracing::debug!(prev_block_id = %req.prev_block_id, "getNextBlockFull"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_next_block_full(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_next_block_full(&req).await)) }) }, rpc::GetBlockDataChunk as req => { tracing::debug!(block_id = %req.block_id, offset = %req.offset, "getBlockDataChunk"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_block_data_chunk(&req); - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_block_data_chunk(&req))) }) }, rpc::GetKeyBlockProof as req => { tracing::debug!(block_id = %req.block_id, "getKeyBlockProof"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_key_block_proof(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_key_block_proof(&req).await)) }) }, rpc::GetPersistentShardStateInfo as req => { tracing::debug!(block_id = %req.block_id, "getPersistentShardStateInfo"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_persistent_state_info(&req); - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_persistent_state_info(&req))) }) }, rpc::GetPersistentQueueStateInfo as req => { tracing::debug!(block_id = %req.block_id, "getPersistentQueueStateInfo"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_queue_persistent_state_info(&req); - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_queue_persistent_state_info(&req))) }) }, rpc::GetPersistentShardStateChunk as req => { @@ -238,11 +232,8 @@ impl Service for BlockchainRpcService { offset = %req.offset, "getPersistentShardStateChunk" ); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_persistent_shard_state_chunk(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_persistent_shard_state_chunk(&req).await)) }) }, rpc::GetPersistentQueueStateChunk as req => { @@ -251,20 +242,14 @@ impl Service for BlockchainRpcService { offset = %req.offset, "getPersistentQueueStateChunk" ); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_persistent_queue_state_chunk(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_persistent_queue_state_chunk(&req).await)) }) }, rpc::GetArchiveInfo as req => { tracing::debug!(mc_seqno = %req.mc_seqno, "getArchiveInfo"); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_archive_info(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_archive_info(&req).await)) }) }, rpc::GetArchiveChunk as req => { @@ -273,11 +258,8 @@ impl Service for BlockchainRpcService { offset = %req.offset, "getArchiveChunk" ); - - let inner = self.inner.clone(); - BoxFutureOrNoop::future(async move { - let res = inner.handle_get_archive_chunk(&req).await; - Some(Response::from_tl(res)) + timed_future(timer,async move { + Some(Response::from_tl(inner.handle_get_archive_chunk(&req).await)) }) }, }, e => { @@ -344,6 +326,20 @@ impl Service for BlockchainRpcService { } } +pub fn timed_future(t: Timer, f: F) -> BoxFutureOrNoop +where + F: Future + Send + 'static, + Timer: FnOnce() -> TimerRet + Send + 'static, + TimerRet: Send + 'static, + T: 'static, +{ + let future = async move { + let _timer = t(); + f.await + }; + BoxFutureOrNoop::future(future) +} + struct Inner { storage: Storage, config: BlockchainRpcServiceConfig, @@ -359,9 +355,6 @@ impl Inner { &self, req: &rpc::GetNextKeyBlockIds, ) -> overlay::Response { - let label = [("method", "getNextKeyBlockIds")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let block_handle_storage = self.storage().block_handle_storage(); let limit = std::cmp::min(req.max_size as usize, self.config.max_key_blocks_list_len); @@ -405,9 +398,6 @@ impl Inner { } async fn handle_get_block_full(&self, req: &rpc::GetBlockFull) -> overlay::Response { - let label = [("method", "getBlockFull")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - match self.get_block_full(&req.block_id).await { Ok(block_full) => overlay::Response::Ok(block_full), Err(e) => { @@ -421,9 +411,6 @@ impl Inner { &self, req: &rpc::GetNextBlockFull, ) -> overlay::Response { - let label = [("method", "getNextBlockFull")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let block_handle_storage = self.storage().block_handle_storage(); let block_connection_storage = self.storage().block_connection_storage(); @@ -448,9 +435,6 @@ impl Inner { } fn handle_get_block_data_chunk(&self, req: &rpc::GetBlockDataChunk) -> overlay::Response { - let label = [("method", "getBlockDataChunk")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let block_storage = self.storage.block_storage(); match block_storage.get_block_data_chunk(&req.block_id, req.offset) { Ok(Some(data)) => overlay::Response::Ok(Data { @@ -468,9 +452,6 @@ impl Inner { &self, req: &rpc::GetKeyBlockProof, ) -> overlay::Response { - let label = [("method", "getKeyBlockProof")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let block_handle_storage = self.storage().block_handle_storage(); let block_storage = self.storage().block_storage(); @@ -502,9 +483,6 @@ impl Inner { let mc_seqno = req.mc_seqno; let node_state = self.storage.node_state(); - let label = [("method", "getArchiveInfo")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - match node_state.load_last_mc_block_id() { Some(last_applied_mc_block) => { if mc_seqno > last_applied_mc_block.seqno { @@ -540,9 +518,6 @@ impl Inner { &self, req: &rpc::GetArchiveChunk, ) -> overlay::Response { - let label = [("method", "getArchiveChunk")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let block_storage = self.storage.block_storage(); let get_archive_chunk = || async { @@ -578,8 +553,6 @@ impl Inner { &self, req: &rpc::GetPersistentQueueStateInfo, ) -> overlay::Response { - let label = [("method", "getQueuePersistentStateInfo")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); let res = self.read_persistent_state_info(&req.block_id, PersistentStateKind::Queue); overlay::Response::Ok(res) } @@ -588,8 +561,6 @@ impl Inner { &self, req: &rpc::GetPersistentShardStateChunk, ) -> overlay::Response { - let label = [("method", "getPersistentShardStateChunk")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Shard) .await } @@ -598,8 +569,6 @@ impl Inner { &self, req: &rpc::GetPersistentQueueStateChunk, ) -> overlay::Response { - let label = [("method", "getPersistentQueueStateChunk")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Queue) .await } diff --git a/core/src/blockchain_rpc/service/util.rs b/core/src/blockchain_rpc/service/util.rs new file mode 100644 index 0000000000..738624546f --- /dev/null +++ b/core/src/blockchain_rpc/service/util.rs @@ -0,0 +1,43 @@ +use crate::proto::blockchain::*; +use crate::proto::overlay; + +macro_rules! constructor_to_string { + ($($ty:path as $name:ident),* $(,)?) => { + #[derive(Debug, Clone, Copy, PartialEq, Eq)] + pub enum Constructor { + $($name),* + } + + impl Constructor { + pub fn from_tl_id(id: u32) -> Option { + match id { + $(<$ty>::TL_ID => Some(Self::$name)),*, + _ => None + } + } + + pub fn as_str(&self) -> &'static str { + match self { + $(Self::$name => stringify!($name)),* + } + } + + } + }; +} + +// update list in `def core_blockchain_rpc_per_method_stats() -> RowPanel:` after changing this +constructor_to_string! { + overlay::Ping as Ping, + rpc::GetNextKeyBlockIds as GetNextKeyBlockIds, + rpc::GetBlockFull as GetBlockFull, + rpc::GetNextBlockFull as GetNextBlockFull, + rpc::GetBlockDataChunk as GetBlockDataChunk, + rpc::GetKeyBlockProof as GetKeyBlockProof, + rpc::GetPersistentShardStateInfo as GetPersistentShardStateInfo, + rpc::GetPersistentQueueStateInfo as GetPersistentQueueStateInfo, + rpc::GetPersistentShardStateChunk as GetPersistentShardStateChunk, + rpc::GetPersistentQueueStateChunk as GetPersistentQueueStateChunk, + rpc::GetArchiveInfo as GetArchiveInfo, + rpc::GetArchiveChunk as GetArchiveChunk +} diff --git a/scripts/gen-dashboard.py b/scripts/gen-dashboard.py index 58af3043ef..1b36899a5a 100644 --- a/scripts/gen-dashboard.py +++ b/scripts/gen-dashboard.py @@ -521,18 +521,7 @@ def net_traffic() -> RowPanel: return create_row("network: Traffic", metrics) -def core_blockchain_rpc() -> RowPanel: - methods = [ - "getNextKeyBlockIds", - "getBlockFull", - "getBlockDataChunk", - "getNextBlockFull", - "getKeyBlockProof", - "getArchiveInfo", - "getArchiveChunk", - "getPersistentStateInfo", - "getPersistentStatePart", - ] +def core_blockchain_rpc_general() -> RowPanel: metrics = [ create_gauge_panel( "tycho_core_overlay_client_validators_to_resolve", @@ -558,7 +547,36 @@ def core_blockchain_rpc() -> RowPanel: legend_format="{{instance}} - {{kind}}", ), ] - metrics += [ + return create_row("blockchain: RPC - General Stats", metrics) + + +def core_blockchain_rpc_per_method_stats() -> RowPanel: + methods = [ + "Ping", + "GetNextKeyBlockIds", + "GetBlockFull", + "GetNextBlockFull", + "GetBlockDataChunk", + "GetKeyBlockProof", + "GetPersistentShardStateInfo", + "GetPersistentQueueStateInfo", + "GetPersistentShardStateChunk", + "GetPersistentQueueStateChunk", + "GetArchiveInfo", + "GetArchiveChunk", + ] + + counter_panels = [ + create_counter_panel( + expr="tycho_blockchain_rpc_method_time_count", + title=f"Blockchain RPC {method} calls/s", + labels_selectors=[f'method="{method}"'], + legend_format="{{instance}}", + ) + for method in methods + ] + + heatmap_panels = [ create_heatmap_panel( "tycho_blockchain_rpc_method_time", f"Blockchain RPC {method} time", @@ -566,7 +584,8 @@ def core_blockchain_rpc() -> RowPanel: ) for method in methods ] - return create_row("blockchain: RPC", metrics) + + return create_row("blockchain: RPC - Method Stats", counter_panels + heatmap_panels) def net_conn_manager() -> RowPanel: @@ -2588,7 +2607,8 @@ def templates() -> Templating: blockchain_stats(), core_bc(), core_block_strider(), - core_blockchain_rpc(), + core_blockchain_rpc_general(), + core_blockchain_rpc_per_method_stats(), storage(), collator_params_metrics(), collation_metrics(), diff --git a/util/src/metrics/histogram_guard.rs b/util/src/metrics/histogram_guard.rs index 135dde75e6..b4e8c8990b 100644 --- a/util/src/metrics/histogram_guard.rs +++ b/util/src/metrics/histogram_guard.rs @@ -24,6 +24,16 @@ impl HistogramGuard { HistogramGuardWithLabels::begin(name, labels) } + pub fn begin_with_labels_owned( + name: &'static str, + labels: T, + ) -> HistogramGuardWithLabelsOwned + where + T: metrics::IntoLabels, + { + HistogramGuardWithLabelsOwned::begin(name, labels) + } + pub fn finish(mut self) -> Duration { let duration = self.started_at.elapsed(); if let Some(name) = self.name.take() { @@ -82,3 +92,47 @@ where } } } + +#[must_use = "The guard is used to update the histogram when it is dropped"] +pub struct HistogramGuardWithLabelsOwned +where + T: metrics::IntoLabels, +{ + name: Option<&'static str>, + started_at: Instant, + labels: Option, +} + +impl HistogramGuardWithLabelsOwned +where + T: metrics::IntoLabels, +{ + pub fn begin(name: &'static str, labels: T) -> Self { + Self { + name: Some(name), + started_at: Instant::now(), + labels: Some(labels), + } + } + + pub fn finish(mut self) -> Duration { + let duration = self.started_at.elapsed(); + if let Some(name) = self.name.take() { + let labels = self.labels.take().unwrap(); + metrics::histogram!(name, labels).record(duration); + } + duration + } +} + +impl Drop for HistogramGuardWithLabelsOwned +where + T: metrics::IntoLabels, +{ + fn drop(&mut self) { + if let Some(name) = self.name.take() { + let labels = self.labels.take().unwrap(); + metrics::histogram!(name, labels).record(self.started_at.elapsed()); + } + } +} From 02b96f07efdc0896aca004260eca2944a31ba860 Mon Sep 17 00:00:00 2001 From: Vladimir Petrzhikovskii Date: Mon, 24 Feb 2025 11:58:36 +0100 Subject: [PATCH 2/2] feat(core): implement request-based rate-limiting for blockchain rpc --- Cargo.lock | 284 ++++++++++--- Cargo.toml | 1 + core/Cargo.toml | 3 + core/src/blockchain_rpc/service/handlers.rs | 338 ++++++++++++++++ core/src/blockchain_rpc/service/mod.rs | 426 +++++--------------- core/src/blockchain_rpc/service/util.rs | 9 + scripts/gen-dashboard.py | 5 + 7 files changed, 692 insertions(+), 374 deletions(-) create mode 100644 core/src/blockchain_rpc/service/handlers.rs diff --git a/Cargo.lock b/Cargo.lock index 86e277d0cc..674944ea88 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -34,10 +34,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e89da841a80418a9b391ebaea17f5c112ffaaa96f621d2c285b5174da76b9011" dependencies = [ "cfg-if", - "getrandom", + "getrandom 0.2.15", "once_cell", "version_check", - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -49,6 +49,12 @@ dependencies = [ "memchr", ] +[[package]] +name = "allocator-api2" +version = "0.2.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" + [[package]] name = "anes" version = "0.1.6" @@ -780,6 +786,20 @@ dependencies = [ "parking_lot_core", ] +[[package]] +name = "dashmap" +version = "6.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5041cc499144891f3790297212f32a74fb938e5136a14943f338ef9e0ae276cf" +dependencies = [ + "cfg-if", + "crossbeam-utils", + "hashbrown 0.14.5", + "lock_api", + "once_cell", + "parking_lot_core", +] + [[package]] name = "data-encoding" version = "2.6.0" @@ -947,7 +967,7 @@ dependencies = [ "curve25519-dalek", "generic-array", "hex", - "rand", + "rand 0.8.5", "serde", "sha2", "tl-proto", @@ -971,7 +991,7 @@ dependencies = [ "hex", "num-bigint", "num-traits", - "rand", + "rand 0.8.5", "rayon", "scc", "serde", @@ -1045,6 +1065,12 @@ version = "1.0.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3f9eec918d3f24069decb9af1554cad7c880e2da24a9afd88aca000531ab82c1" +[[package]] +name = "foldhash" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2" + [[package]] name = "form_urlencoded" version = "1.2.1" @@ -1131,6 +1157,12 @@ version = "0.3.31" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f90f7dce0722e95104fcb095585910c0977252f286e354b5e3bd38902cd99988" +[[package]] +name = "futures-timer" +version = "3.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f288b0a4f20f9a56b5d1da57e2227c661b7b16168e2f72365f57b63326e29b24" + [[package]] name = "futures-util" version = "0.3.31" @@ -1179,7 +1211,21 @@ checksum = "c4567c8db10ae91089c99af84c68c38da3ec2f087c3f82960bcdbf3656b6f4d7" dependencies = [ "cfg-if", "libc", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", +] + +[[package]] +name = "getrandom" +version = "0.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73fea8450eea4bac3940448fb7ae50d91f034f941199fcd9d909a5a07aa455f0" +dependencies = [ + "cfg-if", + "js-sys", + "libc", + "r-efi", + "wasi 0.14.2+wasi-0.2.4", + "wasm-bindgen", ] [[package]] @@ -1194,6 +1240,29 @@ version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a8d1add55171497b4705a648c6b583acafb01d58050a51727785f0b2c8e0a2b2" +[[package]] +name = "governor" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ae4a26ace7b5399e349df31c90afec7555b5af2e64ee8fe9f2a9b5a3204dca06" +dependencies = [ + "cfg-if", + "dashmap 6.1.0", + "futures-sink", + "futures-timer", + "futures-util", + "getrandom 0.3.2", + "hashbrown 0.15.2", + "nonzero_ext", + "parking_lot", + "portable-atomic", + "quanta", + "rand 0.9.0", + "smallvec", + "spinning_top", + "web-time", +] + [[package]] name = "h2" version = "0.4.6" @@ -1240,9 +1309,14 @@ dependencies = [ [[package]] name = "hashbrown" -version = "0.15.0" +version = "0.15.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e087f84d4f86bf4b218b927129862374b72199ae7d8657835f1e89000eea4fb" +checksum = "bf151400ff0baff5465007dd2f3e717f3fe502074ca563069ce3a6629d07b289" +dependencies = [ + "allocator-api2", + "equivalent", + "foldhash", +] [[package]] name = "heck" @@ -1281,7 +1355,7 @@ dependencies = [ "hickory-proto", "once_cell", "radix_trie", - "rand", + "rand 0.8.5", "thiserror 1.0.66", "tokio", "tracing", @@ -1303,7 +1377,7 @@ dependencies = [ "idna 0.4.0", "ipnet", "once_cell", - "rand", + "rand 0.8.5", "thiserror 1.0.66", "tinyvec", "tokio", @@ -1475,7 +1549,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "707907fe3c25f5424cce2cb7e1cbcafee6bdbe735ca90ef77c29e84591e5b9da" dependencies = [ "equivalent", - "hashbrown 0.15.0", + "hashbrown 0.15.2", ] [[package]] @@ -1545,10 +1619,11 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.70" +version = "0.3.77" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1868808506b929d7b0cfa8f75951347aa71bb21144b7791bae35d9bccfcfe37a" +checksum = "1cfaf33c695fc6e08064efbc1f72ec937429614f25eef83af942d0e227c3a28f" dependencies = [ + "once_cell", "wasm-bindgen", ] @@ -1745,7 +1820,7 @@ checksum = "80e04d1dcff3aae0704555fe5fee3bcfaf3d1fdf8a7e521d5b9d2b42acb52cec" dependencies = [ "hermit-abi 0.3.9", "libc", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", "windows-sys 0.52.0", ] @@ -1800,6 +1875,12 @@ dependencies = [ "minimal-lexical", ] +[[package]] +name = "nonzero_ext" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "38bf9645c8b145698bb0b18a4637dcacbc421ea49bef2317e4fd8065a387cf21" + [[package]] name = "ntapi" version = "0.4.1" @@ -1929,7 +2010,7 @@ dependencies = [ "once_cell", "opentelemetry_api", "percent-encoding", - "rand", + "rand 0.8.5", "thiserror 1.0.66", ] @@ -2143,7 +2224,7 @@ version = "0.2.20" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77957b295656769bb8ad2b6a6b09d897d94f05c41b069aede1fcdaa675eaea04" dependencies = [ - "zerocopy", + "zerocopy 0.7.35", ] [[package]] @@ -2228,7 +2309,7 @@ dependencies = [ "libc", "once_cell", "raw-cpuid", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", "web-sys", "winapi", ] @@ -2270,7 +2351,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fadfaed2cd7f389d0161bb73eeb07b7b78f8691047a6f3e73caaeae55310a4a6" dependencies = [ "bytes", - "rand", + "rand 0.8.5", "ring", "rustc-hash 2.1.0", "rustls", @@ -2302,6 +2383,12 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" + [[package]] name = "radix_trie" version = "0.2.1" @@ -2319,8 +2406,19 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" dependencies = [ "libc", - "rand_chacha", - "rand_core", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3779b94aeb87e8bd4e834cee3650289ee9e0d5677f976ecdb6d219e5f4f6cd94" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.3", + "zerocopy 0.8.24", ] [[package]] @@ -2330,7 +2428,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.3", ] [[package]] @@ -2339,7 +2447,16 @@ version = "0.6.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" dependencies = [ - "getrandom", + "getrandom 0.2.15", +] + +[[package]] +name = "rand_core" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" +dependencies = [ + "getrandom 0.3.2", ] [[package]] @@ -2348,7 +2465,7 @@ version = "0.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "59cad018caf63deb318e5a4586d99a24424a364f40f1e5778c29aca23f4fc73e" dependencies = [ - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -2395,7 +2512,7 @@ version = "0.4.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba009ff324d1fc1b900bd1fdb31564febe58a8ccc8a6fdbb93b543d33b13ca43" dependencies = [ - "getrandom", + "getrandom 0.2.15", "libredox", "thiserror 1.0.66", ] @@ -2495,7 +2612,7 @@ checksum = "c17fa4cb658e3583423e915b9f3acc01cceaee1860e33d59ebae66adc3a2dc0d" dependencies = [ "cc", "cfg-if", - "getrandom", + "getrandom 0.2.15", "libc", "spin", "untrusted", @@ -2793,7 +2910,7 @@ version = "2.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" dependencies = [ - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -2833,6 +2950,15 @@ version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" +[[package]] +name = "spinning_top" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d96d2d1d716fb500937168cc09353ffdc7a012be8475ac7308e1bdf0e3923300" +dependencies = [ + "lock_api", +] + [[package]] name = "spki" version = "0.7.3" @@ -2946,7 +3072,7 @@ dependencies = [ "humantime", "opentelemetry", "pin-project", - "rand", + "rand 0.8.5", "serde", "static_assertions", "tarpc-plugins", @@ -3499,7 +3625,7 @@ dependencies = [ "humantime", "metrics", "parking_lot", - "rand", + "rand 0.8.5", "rayon", "reqwest", "rustc_version", @@ -3544,7 +3670,7 @@ dependencies = [ "indexmap 2.6.0", "metrics", "parking_lot", - "rand", + "rand 0.8.5", "rayon", "scc", "scopeguard", @@ -3579,7 +3705,7 @@ dependencies = [ "blake3", "bytes", "clap", - "dashmap", + "dashmap 5.5.3", "everscale-crypto", "everscale-types", "futures-util", @@ -3588,7 +3714,7 @@ dependencies = [ "itertools 0.12.1", "metrics", "parking_lot", - "rand", + "rand 0.8.5", "rand_pcg", "rayon", "scopeguard", @@ -3637,19 +3763,22 @@ dependencies = [ name = "tycho-core" version = "0.2.7" dependencies = [ + "ahash", "anyhow", "arc-swap", "async-trait", "bytes", "bytesize", + "castaway", "everscale-crypto", "everscale-types", "futures-util", + "governor", "humantime", "metrics", "parking_lot", "pin-project-lite", - "rand", + "rand 0.8.5", "scopeguard", "serde", "tempfile", @@ -3692,7 +3821,7 @@ dependencies = [ "clap", "everscale-crypto", "everscale-types", - "rand", + "rand 0.8.5", "serde", "serde_json", "tokio", @@ -3717,7 +3846,7 @@ dependencies = [ "bytesize", "castaway", "clap", - "dashmap", + "dashmap 5.5.3", "ed25519", "everscale-crypto", "exponential-backoff", @@ -3730,7 +3859,7 @@ dependencies = [ "pin-project-lite", "pkcs8", "quinn", - "rand", + "rand 0.8.5", "ring", "rustls", "rustls-webpki", @@ -3786,7 +3915,7 @@ dependencies = [ "clap", "fs_extra", "hex", - "rand", + "rand 0.8.5", "serde", "serde_json", "tycho-util", @@ -3804,7 +3933,7 @@ dependencies = [ "bytes", "bytesize", "crc32c", - "dashmap", + "dashmap 5.5.3", "everscale-types", "fdlimit", "futures-util", @@ -3815,7 +3944,7 @@ dependencies = [ "parking_lot", "parking_lot_core", "quick_cache", - "rand", + "rand 0.8.5", "rlimit", "scopeguard", "serde", @@ -3843,14 +3972,14 @@ dependencies = [ "bytes", "castaway", "criterion", - "dashmap", + "dashmap 5.5.3", "futures-util", "getip", "humantime", "libc", "metrics", "metrics-exporter-prometheus", - "rand", + "rand 0.8.5", "rayon", "serde", "serde_json", @@ -3974,7 +4103,7 @@ version = "1.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "81dfa00651efa65069b0b6b651f4aaa31ba9e3c3ce0137aaad053604ee7e0314" dependencies = [ - "getrandom", + "getrandom 0.2.15", ] [[package]] @@ -4020,26 +4149,35 @@ version = "0.11.0+wasi-snapshot-preview1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" +[[package]] +name = "wasi" +version = "0.14.2+wasi-0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9683f9a5a998d873c0d21fcbe3c083009670149a8fab228644b8bd36b2c48cb3" +dependencies = [ + "wit-bindgen-rt", +] + [[package]] name = "wasm-bindgen" -version = "0.2.93" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a82edfc16a6c469f5f44dc7b571814045d60404b55a0ee849f9bcfa2e63dd9b5" +checksum = "1edc8929d7499fc4e8f0be2262a241556cfc54a0bea223790e71446f2aab1ef5" dependencies = [ "cfg-if", "once_cell", + "rustversion", "wasm-bindgen-macro", ] [[package]] name = "wasm-bindgen-backend" -version = "0.2.93" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9de396da306523044d3302746f1208fa71d7532227f15e347e2d93e4145dd77b" +checksum = "2f0a0651a5c2bc21487bde11ee802ccaf4c51935d0d3d42a6101f98161700bc6" dependencies = [ "bumpalo", "log", - "once_cell", "proc-macro2", "quote", "syn 2.0.96", @@ -4060,9 +4198,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro" -version = "0.2.93" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "585c4c91a46b072c92e908d99cb1dcdf95c5218eeb6f3bf1efa991ee7a68cccf" +checksum = "7fe63fc6d09ed3792bd0897b314f53de8e16568c2b3f7982f468c0bf9bd0b407" dependencies = [ "quote", "wasm-bindgen-macro-support", @@ -4070,9 +4208,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.93" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "afc340c74d9005395cf9dd098506f7f44e38f2b4a21c6aaacf9a105ea5e1e836" +checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de" dependencies = [ "proc-macro2", "quote", @@ -4083,9 +4221,12 @@ dependencies = [ [[package]] name = "wasm-bindgen-shared" -version = "0.2.93" +version = "0.2.100" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c62a0a307cb4a311d3a07867860911ca130c3494e8c2719593806c08bc5d0484" +checksum = "1a05d73b933a847d6cccdda8f838a22ff101ad9bf93e33684f39c1f5f0eece3d" +dependencies = [ + "unicode-ident", +] [[package]] name = "web-sys" @@ -4097,6 +4238,16 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "webpki-roots" version = "0.26.6" @@ -4359,6 +4510,15 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" +[[package]] +name = "wit-bindgen-rt" +version = "0.39.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6f42320e61fe2cfd34354ecb597f86f413484a798ba44a8ca1165c58d42da6c1" +dependencies = [ + "bitflags", +] + [[package]] name = "zerocopy" version = "0.7.35" @@ -4366,7 +4526,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1b9b4fd18abc82b8136838da5d50bae7bdea537c574d8dc1a34ed098d6c166f0" dependencies = [ "byteorder", - "zerocopy-derive", + "zerocopy-derive 0.7.35", +] + +[[package]] +name = "zerocopy" +version = "0.8.24" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2586fea28e186957ef732a5f8b3be2da217d65c5969d4b1e17f973ebbe876879" +dependencies = [ + "zerocopy-derive 0.8.24", ] [[package]] @@ -4380,6 +4549,17 @@ dependencies = [ "syn 2.0.96", ] +[[package]] +name = "zerocopy-derive" +version = "0.8.24" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a996a8f63c5c4448cd959ac1bab0aaa3306ccfd060472f85943ee0750f0169be" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.96", +] + [[package]] name = "zeroize" version = "1.8.1" diff --git a/Cargo.toml b/Cargo.toml index 2a3aa81d0f..c060b79651 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -53,6 +53,7 @@ exponential-backoff = "1" fdlimit = "0.3.0" futures-util = "0.3" getip = "0.1" +governor = { version = "0.9" } hex = "0.4" humantime = "2" indexmap = "2.2" diff --git a/core/Cargo.toml b/core/Cargo.toml index 1d251fd53d..ae60247df9 100644 --- a/core/Cargo.toml +++ b/core/Cargo.toml @@ -10,13 +10,16 @@ repository.workspace = true license.workspace = true [dependencies] +ahash = { workspace = true } anyhow = { workspace = true } arc-swap = { workspace = true } async-trait = { workspace = true } bytes = { workspace = true, features = ["serde"] } bytesize = { workspace = true } +castaway = { workspace = true } everscale-types = { workspace = true, features = ["blake3", "rayon"] } futures-util = { workspace = true } +governor = { workspace = true } humantime = { workspace = true } metrics = { workspace = true } parking_lot = { workspace = true } diff --git a/core/src/blockchain_rpc/service/handlers.rs b/core/src/blockchain_rpc/service/handlers.rs new file mode 100644 index 0000000000..92ff29f681 --- /dev/null +++ b/core/src/blockchain_rpc/service/handlers.rs @@ -0,0 +1,338 @@ +use std::num::{NonZeroU32, NonZeroU64}; + +use anyhow::Context; +use bytes::Bytes; +use everscale_types::models::BlockId; +use tycho_storage::{ArchiveId, BlockConnection, KeyBlocksDirection, PersistentStateKind}; +use tycho_util::metrics::HistogramGuard; + +use super::{Inner, RPC_METHOD_TIMINGS_METRIC}; +use crate::blockchain_rpc::{BAD_REQUEST_ERROR_CODE, INTERNAL_ERROR_CODE, NOT_FOUND_ERROR_CODE}; +use crate::proto::blockchain::{ + rpc, ArchiveInfo, BlockData, BlockFull, Data, KeyBlockIds, KeyBlockProof, PersistentStateInfo, +}; +use crate::proto::overlay; + +impl Inner { + pub(super) fn handle_get_next_key_block_ids( + &self, + req: &rpc::GetNextKeyBlockIds, + ) -> overlay::Response { + let block_handle_storage = self.storage().block_handle_storage(); + + let limit = std::cmp::min(req.max_size as usize, self.config.max_key_blocks_list_len); + + let get_next_key_block_ids = || { + if !req.block_id.shard.is_masterchain() { + anyhow::bail!("first block id is not from masterchain"); + } + + let mut iterator = block_handle_storage + .key_blocks_iterator(KeyBlocksDirection::ForwardFrom(req.block_id.seqno)) + .take(limit + 1); + + if let Some(id) = iterator.next() { + anyhow::ensure!( + id.root_hash == req.block_id.root_hash, + "first block root hash mismatch" + ); + anyhow::ensure!( + id.file_hash == req.block_id.file_hash, + "first block file hash mismatch" + ); + } + + Ok::<_, anyhow::Error>(iterator.take(limit).collect::>()) + }; + + match get_next_key_block_ids() { + Ok(ids) => { + let incomplete = ids.len() < limit; + overlay::Response::Ok(KeyBlockIds { + block_ids: ids, + incomplete, + }) + } + Err(e) => { + tracing::warn!("get_next_key_block_ids failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) async fn handle_get_block_full( + &self, + req: &rpc::GetBlockFull, + ) -> overlay::Response { + match self.get_block_full(&req.block_id).await { + Ok(block_full) => overlay::Response::Ok(block_full), + Err(e) => { + tracing::warn!("get_block_full failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) async fn handle_get_next_block_full( + &self, + req: &rpc::GetNextBlockFull, + ) -> overlay::Response { + let block_handle_storage = self.storage().block_handle_storage(); + let block_connection_storage = self.storage().block_connection_storage(); + + let get_next_block_full = async { + let next_block_id = match block_handle_storage.load_handle(&req.prev_block_id) { + Some(handle) if handle.has_next1() => block_connection_storage + .load_connection(&req.prev_block_id, BlockConnection::Next1) + .context("connection not found")?, + _ => return Ok(BlockFull::NotFound), + }; + + self.get_block_full(&next_block_id).await + }; + + match get_next_block_full.await { + Ok(block_full) => overlay::Response::Ok(block_full), + Err(e) => { + tracing::warn!("get_next_block_full failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) fn handle_get_block_data_chunk( + &self, + req: &rpc::GetBlockDataChunk, + ) -> overlay::Response { + let block_storage = self.storage.block_storage(); + match block_storage.get_block_data_chunk(&req.block_id, req.offset) { + Ok(Some(data)) => overlay::Response::Ok(Data { + data: Bytes::from_owner(data), + }), + Ok(None) => overlay::Response::Err(NOT_FOUND_ERROR_CODE), + Err(e) => { + tracing::warn!("get_block_data_chunk failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) async fn handle_get_key_block_proof( + &self, + req: &rpc::GetKeyBlockProof, + ) -> overlay::Response { + let block_handle_storage = self.storage().block_handle_storage(); + let block_storage = self.storage().block_storage(); + + let get_key_block_proof = async { + match block_handle_storage.load_handle(&req.block_id) { + Some(handle) if handle.has_proof() => { + let data = block_storage.load_block_proof_raw(&handle).await?; + Ok::<_, anyhow::Error>(KeyBlockProof::Found { + proof: Bytes::from_owner(data), + }) + } + _ => Ok(KeyBlockProof::NotFound), + } + }; + + match get_key_block_proof.await { + Ok(key_block_proof) => overlay::Response::Ok(key_block_proof), + Err(e) => { + tracing::warn!("get_key_block_proof failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) async fn handle_get_archive_info( + &self, + req: &rpc::GetArchiveInfo, + ) -> overlay::Response { + let mc_seqno = req.mc_seqno; + let node_state = self.storage.node_state(); + + match node_state.load_last_mc_block_id() { + Some(last_applied_mc_block) => { + if mc_seqno > last_applied_mc_block.seqno { + return overlay::Response::Ok(ArchiveInfo::TooNew); + } + + let block_storage = self.storage().block_storage(); + + let id = block_storage.get_archive_id(mc_seqno); + let size_res = match id { + ArchiveId::Found(id) => block_storage.get_archive_size(id), + ArchiveId::TooNew | ArchiveId::NotFound => Ok(None), + }; + + overlay::Response::Ok(match (id, size_res) { + (ArchiveId::Found(id), Ok(Some(size))) if size > 0 => ArchiveInfo::Found { + id: id as u64, + size: NonZeroU64::new(size as _).unwrap(), + chunk_size: block_storage.archive_chunk_size(), + }, + (ArchiveId::TooNew, Ok(None)) => ArchiveInfo::TooNew, + _ => ArchiveInfo::NotFound, + }) + } + None => { + tracing::warn!("get_archive_id failed: no blocks applied"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) async fn handle_get_archive_chunk( + &self, + req: &rpc::GetArchiveChunk, + ) -> overlay::Response { + let block_storage = self.storage.block_storage(); + + let get_archive_chunk = || async { + let archive_slice = block_storage + .get_archive_chunk(req.archive_id as u32, req.offset) + .await?; + + Ok::<_, anyhow::Error>(archive_slice) + }; + + match get_archive_chunk().await { + Ok(data) => overlay::Response::Ok(Data { + data: Bytes::from_owner(data), + }), + Err(e) => { + tracing::warn!("get_archive_chunk failed: {e:?}"); + overlay::Response::Err(INTERNAL_ERROR_CODE) + } + } + } + + pub(super) fn handle_get_persistent_state_info( + &self, + req: &rpc::GetPersistentShardStateInfo, + ) -> overlay::Response { + let label = [("method", "getPersistentStateInfo")]; + let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); + let res = self.read_persistent_state_info(&req.block_id, PersistentStateKind::Shard); + overlay::Response::Ok(res) + } + + pub(super) fn handle_get_queue_persistent_state_info( + &self, + req: &rpc::GetPersistentQueueStateInfo, + ) -> overlay::Response { + let res = self.read_persistent_state_info(&req.block_id, PersistentStateKind::Queue); + overlay::Response::Ok(res) + } + + pub(super) async fn handle_get_persistent_shard_state_chunk( + &self, + req: &rpc::GetPersistentShardStateChunk, + ) -> overlay::Response { + self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Shard) + .await + } + + pub(super) async fn handle_get_persistent_queue_state_chunk( + &self, + req: &rpc::GetPersistentQueueStateChunk, + ) -> overlay::Response { + self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Queue) + .await + } + + pub(super) async fn get_block_full(&self, block_id: &BlockId) -> anyhow::Result { + let block_handle_storage = self.storage().block_handle_storage(); + let block_storage = self.storage().block_storage(); + + let handle = match block_handle_storage.load_handle(block_id) { + Some(handle) if handle.has_all_block_parts() => handle, + _ => return Ok(BlockFull::NotFound), + }; + + let Some(data) = block_storage.get_block_data_chunk(block_id, 0)? else { + return Ok(BlockFull::NotFound); + }; + + let data_chunk_size = block_storage.block_data_chunk_size(); + let data_size = if data.len() < data_chunk_size.get() as usize { + // NOTE: Skip one RocksDB read for relatively small blocks + // Average block size is 4KB, while the chunk size is 1MB. + data.len() as u32 + } else { + match block_storage.get_block_data_size(block_id)? { + Some(size) => size, + None => return Ok(BlockFull::NotFound), + } + }; + + let block = BlockData { + data: Bytes::from_owner(data), + size: NonZeroU32::new(data_size).expect("shouldn't happen"), + chunk_size: data_chunk_size, + }; + + let (proof, queue_diff) = tokio::join!( + block_storage.load_block_proof_raw(&handle), + block_storage.load_queue_diff_raw(&handle) + ); + + Ok(BlockFull::Found { + block_id: *block_id, + block, + proof: Bytes::from_owner(proof?), + queue_diff: Bytes::from_owner(queue_diff?), + }) + } + + pub(super) fn read_persistent_state_info( + &self, + block_id: &BlockId, + state_kind: PersistentStateKind, + ) -> PersistentStateInfo { + let persistent_state_storage = self.storage().persistent_state_storage(); + if self.config.serve_persistent_states { + if let Some(info) = persistent_state_storage.get_state_info(block_id, state_kind) { + return PersistentStateInfo::Found { + size: info.size, + chunk_size: info.chunk_size, + }; + } + } + PersistentStateInfo::NotFound + } + + pub(super) async fn read_persistent_state_chunk( + &self, + block_id: &BlockId, + offset: u64, + state_kind: PersistentStateKind, + ) -> overlay::Response { + let persistent_state_storage = self.storage().persistent_state_storage(); + + let persistent_state_request_validation = || { + anyhow::ensure!( + self.config.serve_persistent_states, + "persistent states are disabled" + ); + Ok::<_, anyhow::Error>(()) + }; + + if let Err(e) = persistent_state_request_validation() { + tracing::debug!("persistent state request validation failed: {e:?}"); + return overlay::Response::Err(BAD_REQUEST_ERROR_CODE); + } + + match persistent_state_storage + .read_state_part(block_id, offset, state_kind) + .await + { + Some(data) => overlay::Response::Ok(Data { data: data.into() }), + None => { + tracing::debug!("failed to read persistent state part"); + overlay::Response::Err(NOT_FOUND_ERROR_CODE) + } + } + } +} diff --git a/core/src/blockchain_rpc/service/mod.rs b/core/src/blockchain_rpc/service/mod.rs index 8dfd3e7870..36b3e26644 100644 --- a/core/src/blockchain_rpc/service/mod.rs +++ b/core/src/blockchain_rpc/service/mod.rs @@ -1,21 +1,20 @@ +mod handlers; mod util; -use std::num::{NonZeroU32, NonZeroU64}; +use std::num::NonZeroU32; use std::sync::Arc; -use anyhow::Context; use bytes::{Buf, Bytes}; -use everscale_types::models::BlockId; use futures_util::Future; use metrics::Label; use serde::{Deserialize, Serialize}; use tycho_block_util::message::validate_external_message; use tycho_network::{try_handle_prefix, InboundRequestMeta, Response, Service, ServiceRequest}; -use tycho_storage::{ArchiveId, BlockConnection, KeyBlocksDirection, PersistentStateKind, Storage}; +use tycho_storage::Storage; use tycho_util::futures::BoxFutureOrNoop; use tycho_util::metrics::HistogramGuard; -use crate::blockchain_rpc::{BAD_REQUEST_ERROR_CODE, INTERNAL_ERROR_CODE, NOT_FOUND_ERROR_CODE}; +use crate::blockchain_rpc::service::util::{Constructor, RateLimiter}; use crate::proto::blockchain::*; use crate::proto::overlay; @@ -29,6 +28,10 @@ pub trait BroadcastListener: Send + Sync + 'static { meta: Arc, message: Bytes, ) -> Self::HandleMessageFut<'_>; + + fn is_noop(&self) -> bool { + false + } } #[derive(Debug, Default, Clone, Copy, Eq, PartialEq)] @@ -41,6 +44,11 @@ impl BroadcastListener for NoopBroadcastListener { fn handle_message(&self, _: Arc, _: Bytes) -> Self::HandleMessageFut<'_> { futures_util::future::ready(()) } + + #[inline] + fn is_noop(&self) -> bool { + true + } } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -56,6 +64,8 @@ pub struct BlockchainRpcServiceConfig { /// /// Default: yes. pub serve_persistent_states: bool, + + pub rate_limits: RateLimits, } impl Default for BlockchainRpcServiceConfig { @@ -63,6 +73,30 @@ impl Default for BlockchainRpcServiceConfig { Self { max_key_blocks_list_len: 8, serve_persistent_states: true, + rate_limits: Default::default(), + } + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(default)] +#[non_exhaustive] +pub struct RateLimits { + /// rate limits for methods like `GetPersistentQueueStateInfo`, `GetArchiveInfo`, etc. + pub info_method_rps: NonZeroU32, + /// rate limits for methods like `GetPersistentQueueStateChunk`, `GetArchiveChunk`, etc. + pub chunk_method_rps: NonZeroU32, + + /// Message broadcast rate limits + pub send_message: NonZeroU32, +} + +impl Default for RateLimits { + fn default() -> Self { + Self { + info_method_rps: NonZeroU32::new(100).unwrap(), + chunk_method_rps: NonZeroU32::new(100).unwrap(), + send_message: NonZeroU32::new(10_000).unwrap(), } } } @@ -82,6 +116,18 @@ where BlockchainRpcService { inner: Arc::new(Inner { storage, + info_rate_limiter: RateLimiter::dashmap_with_hasher( + governor::Quota::per_second(self.config.rate_limits.info_method_rps), + Default::default(), + ), + chunk_rate_limiter: RateLimiter::dashmap_with_hasher( + governor::Quota::per_second(self.config.rate_limits.chunk_method_rps), + Default::default(), + ), + send_message_rate_limiter: RateLimiter::dashmap_with_hasher( + governor::Quota::per_second(self.config.rate_limits.send_message), + Default::default(), + ), config: self.config, broadcast_listener, }), @@ -164,13 +210,20 @@ impl Service for BlockchainRpcService { } }; - let method = util::Constructor::from_tl_id(constructor); + let method = Constructor::from_tl_id(constructor); let label = vec![Label::new( "method", method.map_or("unknown", |m| m.as_str()), )]; - let timer = - move || HistogramGuard::begin_with_labels_owned(RPC_METHOD_TIMINGS_METRIC, label); + let timer = { + let label = label.clone(); + move || HistogramGuard::begin_with_labels_owned(RPC_METHOD_TIMINGS_METRIC, label) + }; + + if let Some(value) = self.inner.check_rate_limit(&req, method) { + metrics::counter!("tycho_rpc_rate_limit_exceeded_total", label).increment(1); + return value; + } let inner = self.inner.clone(); @@ -272,7 +325,20 @@ impl Service for BlockchainRpcService { fn on_message(&self, mut req: ServiceRequest) -> Self::OnMessageFuture { use tl_proto::{BytesMeta, TlRead}; - // TODO: Do nothing if `B` is `NoopBroadcastListener` via `castaway` ? + if self.inner.broadcast_listener.is_noop() { + return BoxFutureOrNoop::Noop; + } + + if self + .inner + .send_message_rate_limiter + .check_key(&req.metadata.peer_id) + .is_err() + { + metrics::counter!("tycho_rpc_rate_limit_exceeded_total", "method" => "sendMessage") + .increment(1); + return BoxFutureOrNoop::Noop; + } // Require message body to contain at least two constructors. if req.body.len() < 8 { @@ -344,6 +410,9 @@ struct Inner { storage: Storage, config: BlockchainRpcServiceConfig, broadcast_listener: B, + info_rate_limiter: RateLimiter, + chunk_rate_limiter: RateLimiter, + send_message_rate_limiter: RateLimiter, } impl Inner { @@ -351,321 +420,34 @@ impl Inner { &self.storage } - fn handle_get_next_key_block_ids( - &self, - req: &rpc::GetNextKeyBlockIds, - ) -> overlay::Response { - let block_handle_storage = self.storage().block_handle_storage(); - - let limit = std::cmp::min(req.max_size as usize, self.config.max_key_blocks_list_len); - - let get_next_key_block_ids = || { - if !req.block_id.shard.is_masterchain() { - anyhow::bail!("first block id is not from masterchain"); - } - - let mut iterator = block_handle_storage - .key_blocks_iterator(KeyBlocksDirection::ForwardFrom(req.block_id.seqno)) - .take(limit + 1); - - if let Some(id) = iterator.next() { - anyhow::ensure!( - id.root_hash == req.block_id.root_hash, - "first block root hash mismatch" - ); - anyhow::ensure!( - id.file_hash == req.block_id.file_hash, - "first block file hash mismatch" - ); - } - - Ok::<_, anyhow::Error>(iterator.take(limit).collect::>()) - }; - - match get_next_key_block_ids() { - Ok(ids) => { - let incomplete = ids.len() < limit; - overlay::Response::Ok(KeyBlockIds { - block_ids: ids, - incomplete, - }) - } - Err(e) => { - tracing::warn!("get_next_key_block_ids failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - async fn handle_get_block_full(&self, req: &rpc::GetBlockFull) -> overlay::Response { - match self.get_block_full(&req.block_id).await { - Ok(block_full) => overlay::Response::Ok(block_full), - Err(e) => { - tracing::warn!("get_block_full failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - async fn handle_get_next_block_full( - &self, - req: &rpc::GetNextBlockFull, - ) -> overlay::Response { - let block_handle_storage = self.storage().block_handle_storage(); - let block_connection_storage = self.storage().block_connection_storage(); - - let get_next_block_full = async { - let next_block_id = match block_handle_storage.load_handle(&req.prev_block_id) { - Some(handle) if handle.has_next1() => block_connection_storage - .load_connection(&req.prev_block_id, BlockConnection::Next1) - .context("connection not found")?, - _ => return Ok(BlockFull::NotFound), - }; - - self.get_block_full(&next_block_id).await - }; - - match get_next_block_full.await { - Ok(block_full) => overlay::Response::Ok(block_full), - Err(e) => { - tracing::warn!("get_next_block_full failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - fn handle_get_block_data_chunk(&self, req: &rpc::GetBlockDataChunk) -> overlay::Response { - let block_storage = self.storage.block_storage(); - match block_storage.get_block_data_chunk(&req.block_id, req.offset) { - Ok(Some(data)) => overlay::Response::Ok(Data { - data: Bytes::from_owner(data), - }), - Ok(None) => overlay::Response::Err(NOT_FOUND_ERROR_CODE), - Err(e) => { - tracing::warn!("get_block_data_chunk failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - async fn handle_get_key_block_proof( - &self, - req: &rpc::GetKeyBlockProof, - ) -> overlay::Response { - let block_handle_storage = self.storage().block_handle_storage(); - let block_storage = self.storage().block_storage(); - - let get_key_block_proof = async { - match block_handle_storage.load_handle(&req.block_id) { - Some(handle) if handle.has_proof() => { - let data = block_storage.load_block_proof_raw(&handle).await?; - Ok::<_, anyhow::Error>(KeyBlockProof::Found { - proof: Bytes::from_owner(data), - }) - } - _ => Ok(KeyBlockProof::NotFound), - } - }; - - match get_key_block_proof.await { - Ok(key_block_proof) => overlay::Response::Ok(key_block_proof), - Err(e) => { - tracing::warn!("get_key_block_proof failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - async fn handle_get_archive_info( - &self, - req: &rpc::GetArchiveInfo, - ) -> overlay::Response { - let mc_seqno = req.mc_seqno; - let node_state = self.storage.node_state(); - - match node_state.load_last_mc_block_id() { - Some(last_applied_mc_block) => { - if mc_seqno > last_applied_mc_block.seqno { - return overlay::Response::Ok(ArchiveInfo::TooNew); - } - - let block_storage = self.storage().block_storage(); - - let id = block_storage.get_archive_id(mc_seqno); - let size_res = match id { - ArchiveId::Found(id) => block_storage.get_archive_size(id), - ArchiveId::TooNew | ArchiveId::NotFound => Ok(None), - }; - - overlay::Response::Ok(match (id, size_res) { - (ArchiveId::Found(id), Ok(Some(size))) if size > 0 => ArchiveInfo::Found { - id: id as u64, - size: NonZeroU64::new(size as _).unwrap(), - chunk_size: block_storage.archive_chunk_size(), - }, - (ArchiveId::TooNew, Ok(None)) => ArchiveInfo::TooNew, - _ => ArchiveInfo::NotFound, - }) - } - None => { - tracing::warn!("get_archive_id failed: no blocks applied"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - async fn handle_get_archive_chunk( - &self, - req: &rpc::GetArchiveChunk, - ) -> overlay::Response { - let block_storage = self.storage.block_storage(); - - let get_archive_chunk = || async { - let archive_slice = block_storage - .get_archive_chunk(req.archive_id as u32, req.offset) - .await?; - - Ok::<_, anyhow::Error>(archive_slice) - }; - - match get_archive_chunk().await { - Ok(data) => overlay::Response::Ok(Data { - data: Bytes::from_owner(data), - }), - Err(e) => { - tracing::warn!("get_archive_chunk failed: {e:?}"); - overlay::Response::Err(INTERNAL_ERROR_CODE) - } - } - } - - fn handle_get_persistent_state_info( - &self, - req: &rpc::GetPersistentShardStateInfo, - ) -> overlay::Response { - let label = [("method", "getPersistentStateInfo")]; - let _hist = HistogramGuard::begin_with_labels(RPC_METHOD_TIMINGS_METRIC, &label); - let res = self.read_persistent_state_info(&req.block_id, PersistentStateKind::Shard); - overlay::Response::Ok(res) - } - - fn handle_get_queue_persistent_state_info( - &self, - req: &rpc::GetPersistentQueueStateInfo, - ) -> overlay::Response { - let res = self.read_persistent_state_info(&req.block_id, PersistentStateKind::Queue); - overlay::Response::Ok(res) - } - - async fn handle_get_persistent_shard_state_chunk( + fn check_rate_limit( &self, - req: &rpc::GetPersistentShardStateChunk, - ) -> overlay::Response { - self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Shard) - .await - } - - async fn handle_get_persistent_queue_state_chunk( - &self, - req: &rpc::GetPersistentQueueStateChunk, - ) -> overlay::Response { - self.read_persistent_state_chunk(&req.block_id, req.offset, PersistentStateKind::Queue) - .await - } -} - -impl Inner { - async fn get_block_full(&self, block_id: &BlockId) -> anyhow::Result { - let block_handle_storage = self.storage().block_handle_storage(); - let block_storage = self.storage().block_storage(); - - let handle = match block_handle_storage.load_handle(block_id) { - Some(handle) if handle.has_all_block_parts() => handle, - _ => return Ok(BlockFull::NotFound), + req: &ServiceRequest, + method: Option, + ) -> Option>> { + let rate_limiter = match method { + Some( + Constructor::GetPersistentQueueStateChunk + | Constructor::GetArchiveChunk + | Constructor::GetPersistentShardStateChunk + | Constructor::GetBlockDataChunk + | Constructor::GetNextBlockFull + | Constructor::GetBlockFull + | Constructor::GetNextKeyBlockIds + | Constructor::GetKeyBlockProof, + ) => &self.chunk_rate_limiter, + Some( + Constructor::GetPersistentShardStateInfo + | Constructor::GetPersistentQueueStateInfo + | Constructor::GetArchiveInfo + | Constructor::Ping, + ) => &self.info_rate_limiter, + None => return Some(BoxFutureOrNoop::Noop), }; - let Some(data) = block_storage.get_block_data_chunk(block_id, 0)? else { - return Ok(BlockFull::NotFound); - }; - - let data_chunk_size = block_storage.block_data_chunk_size(); - let data_size = if data.len() < data_chunk_size.get() as usize { - // NOTE: Skip one RocksDB read for relatively small blocks - // Average block size is 4KB, while the chunk size is 1MB. - data.len() as u32 - } else { - match block_storage.get_block_data_size(block_id)? { - Some(size) => size, - None => return Ok(BlockFull::NotFound), - } - }; - - let block = BlockData { - data: Bytes::from_owner(data), - size: NonZeroU32::new(data_size).expect("shouldn't happen"), - chunk_size: data_chunk_size, - }; - - let (proof, queue_diff) = tokio::join!( - block_storage.load_block_proof_raw(&handle), - block_storage.load_queue_diff_raw(&handle) - ); - - Ok(BlockFull::Found { - block_id: *block_id, - block, - proof: Bytes::from_owner(proof?), - queue_diff: Bytes::from_owner(queue_diff?), - }) - } - - fn read_persistent_state_info( - &self, - block_id: &BlockId, - state_kind: PersistentStateKind, - ) -> PersistentStateInfo { - let persistent_state_storage = self.storage().persistent_state_storage(); - if self.config.serve_persistent_states { - if let Some(info) = persistent_state_storage.get_state_info(block_id, state_kind) { - return PersistentStateInfo::Found { - size: info.size, - chunk_size: info.chunk_size, - }; - } - } - PersistentStateInfo::NotFound - } - - async fn read_persistent_state_chunk( - &self, - block_id: &BlockId, - offset: u64, - state_kind: PersistentStateKind, - ) -> overlay::Response { - let persistent_state_storage = self.storage().persistent_state_storage(); - - let persistent_state_request_validation = || { - anyhow::ensure!( - self.config.serve_persistent_states, - "persistent states are disabled" - ); - Ok::<_, anyhow::Error>(()) - }; - - if let Err(e) = persistent_state_request_validation() { - tracing::debug!("persistent state request validation failed: {e:?}"); - return overlay::Response::Err(BAD_REQUEST_ERROR_CODE); - } - - match persistent_state_storage - .read_state_part(block_id, offset, state_kind) - .await - { - Some(data) => overlay::Response::Ok(Data { data: data.into() }), - None => { - tracing::debug!("failed to read persistent state part"); - overlay::Response::Err(NOT_FOUND_ERROR_CODE) - } - } + rate_limiter + .check_key(&req.metadata.peer_id) + .err() + .map(|_| BoxFutureOrNoop::Noop) } } diff --git a/core/src/blockchain_rpc/service/util.rs b/core/src/blockchain_rpc/service/util.rs index 738624546f..7fc4583aea 100644 --- a/core/src/blockchain_rpc/service/util.rs +++ b/core/src/blockchain_rpc/service/util.rs @@ -1,3 +1,6 @@ +use governor::state::keyed::DefaultKeyedStateStore; +use tycho_network::PeerId; + use crate::proto::blockchain::*; use crate::proto::overlay; @@ -41,3 +44,9 @@ constructor_to_string! { rpc::GetArchiveInfo as GetArchiveInfo, rpc::GetArchiveChunk as GetArchiveChunk } + +pub type RateLimiter = governor::RateLimiter< + PeerId, + DefaultKeyedStateStore, + governor::clock::DefaultClock, +>; diff --git a/scripts/gen-dashboard.py b/scripts/gen-dashboard.py index 1b36899a5a..2ab93ae76e 100644 --- a/scripts/gen-dashboard.py +++ b/scripts/gen-dashboard.py @@ -546,6 +546,11 @@ def core_blockchain_rpc_general() -> RowPanel: unit_format=UNITS.SECONDS, legend_format="{{instance}} - {{kind}}", ), + create_counter_panel( + "tycho_rpc_rate_limit_exceeded_total", + "RPC Rate Limit Exceeded", + legend_format="{{instance}} - {{method}}", + ), ] return create_row("blockchain: RPC - General Stats", metrics)