rpc server: add prometheus label is_rate_limited (#3504)

After some discussion with @kogeler after the we added the rate-limit
middleware it may slow down
the rpc call timings metrics significantly because it works as follows:

1. The rate limit guard is checked when the call comes and if a slot is
available -> process the call
2. If no free spot is available then the call will be sleeping
`jitter_delay + min_time_rate_guard` then woken up and checked at most
ten times
3. If no spot is available after 10 iterations -> the call is rejected
(this may take tens of seconds)

Thus, this PR adds a label "is_rate_limited" to filter those out on the
metrics "substrate_rpc_calls_time" and "substrate_rpc_calls_finished".

I had to merge two middleware layers Metrics and RateLimit to avoid
shared state in a hacky way.

---------

Co-authored-by: James Wilson <james@jsdw.me>
This commit is contained in:
Niklas Adolfsson
2024-03-05 10:50:57 +01:00
committed by GitHub
parent 0eda0b3f17
commit efcea0edab
7 changed files with 240 additions and 213 deletions
@@ -16,92 +16,32 @@
// You should have received a copy of the GNU General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/>.
//! RPC rate limiting middleware.
//! RPC rate limit.
use std::{num::NonZeroU32, sync::Arc, time::Duration};
use futures::future::{BoxFuture, FutureExt};
use governor::{
clock::{Clock, DefaultClock, QuantaClock},
clock::{DefaultClock, QuantaClock},
middleware::NoOpMiddleware,
state::{InMemoryState, NotKeyed},
Jitter,
};
use jsonrpsee::{
server::middleware::rpc::RpcServiceT,
types::{ErrorObject, Id, Request},
MethodResponse,
Quota,
};
use std::{num::NonZeroU32, sync::Arc};
type RateLimitInner = governor::RateLimiter<NotKeyed, InMemoryState, DefaultClock, NoOpMiddleware>;
const MAX_JITTER: Duration = Duration::from_millis(50);
const MAX_RETRIES: usize = 10;
/// JSON-RPC rate limit middleware layer.
/// Rate limit.
#[derive(Debug, Clone)]
pub struct RateLimitLayer(governor::Quota);
pub struct RateLimit {
pub(crate) inner: Arc<RateLimitInner>,
pub(crate) clock: QuantaClock,
}
impl RateLimitLayer {
/// Create new rate limit enforced per minute.
impl RateLimit {
/// Create a new `RateLimit` per minute.
pub fn per_minute(n: NonZeroU32) -> Self {
Self(governor::Quota::per_minute(n))
}
}
/// JSON-RPC rate limit middleware
pub struct RateLimit<S> {
service: S,
rate_limit: Arc<RateLimitInner>,
clock: QuantaClock,
}
impl<S> tower::Layer<S> for RateLimitLayer {
type Service = RateLimit<S>;
fn layer(&self, service: S) -> Self::Service {
let clock = QuantaClock::default();
RateLimit {
service,
rate_limit: Arc::new(RateLimitInner::direct_with_clock(self.0, &clock)),
Self {
inner: Arc::new(RateLimitInner::direct_with_clock(Quota::per_minute(n), &clock)),
clock,
}
}
}
impl<'a, S> RpcServiceT<'a> for RateLimit<S>
where
S: Send + Sync + RpcServiceT<'a> + Clone + 'static,
{
type Future = BoxFuture<'a, MethodResponse>;
fn call(&self, req: Request<'a>) -> Self::Future {
let service = self.service.clone();
let rate_limit = self.rate_limit.clone();
let clock = self.clock.clone();
async move {
let mut attempts = 0;
let jitter = Jitter::up_to(MAX_JITTER);
loop {
if attempts >= MAX_RETRIES {
break reject_too_many_calls(req.id);
}
if let Err(rejected) = rate_limit.check() {
tokio::time::sleep(jitter + rejected.wait_time_from(clock.now())).await;
} else {
break service.call(req).await;
}
attempts += 1;
}
}
.boxed()
}
}
fn reject_too_many_calls(id: Id) -> MethodResponse {
MethodResponse::error(id, ErrorObject::owned(-32999, "RPC rate limit exceeded", None::<()>))
}