Removed local concurrency limiters, switch to global rate limiting

This commit is contained in:
mdecimus
2025-01-18 19:09:02 +01:00
parent e7c6be45d8
commit 8438435fbe
54 changed files with 439 additions and 11477 deletions

View File

@@ -24,7 +24,6 @@ tokio-rustls = { version = "0.26", default-features = false, features = ["ring",
parking_lot = "0.12"
ahash = { version = "0.8" }
md5 = "0.7.0"
dashmap = "6.0"
rand = "0.8.5"

View File

@@ -7,8 +7,8 @@
use std::{iter::Peekable, sync::Arc, vec::IntoIter};
use common::{
listener::{limiter::ConcurrencyLimiter, SessionResult, SessionStream},
ConcurrencyLimiters, KV_RATE_LIMIT_IMAP,
listener::{SessionResult, SessionStream},
KV_RATE_LIMIT_IMAP,
};
use imap_proto::{
receiver::{self, Request},
@@ -419,29 +419,6 @@ impl<T: SessionStream> Session<T> {
},
}
}
pub fn get_concurrency_limiter(&self, account_id: u32) -> Option<Arc<ConcurrencyLimiters>> {
let rate = self.server.core.imap.rate_concurrent?;
self.server
.inner
.data
.imap_limiter
.get(&account_id)
.map(|limiter| limiter.clone())
.unwrap_or_else(|| {
let limiter = Arc::new(ConcurrencyLimiters {
concurrent_requests: ConcurrencyLimiter::new(rate),
concurrent_uploads: ConcurrencyLimiter::new(rate),
});
self.server
.inner
.data
.imap_limiter
.insert(account_id, limiter.clone());
limiter
})
.into()
}
}
impl<T: SessionStream> State<T> {

View File

@@ -9,7 +9,7 @@ use common::{
sasl::{sasl_decode_challenge_oauth, sasl_decode_challenge_plain},
AuthRequest,
},
listener::SessionStream,
listener::{limiter::LimiterResult, SessionStream},
};
use directory::Permission;
use imap_proto::{
@@ -105,17 +105,14 @@ impl<T: SessionStream> Session<T> {
})?;
// Enforce concurrency limits
let in_flight = match self
.get_concurrency_limiter(access_token.primary_id())
.map(|limiter| limiter.concurrent_requests.is_allowed())
{
Some(Some(limiter)) => Some(limiter),
None => None,
Some(None) => {
let in_flight = match access_token.is_imap_request_allowed() {
LimiterResult::Allowed(in_flight) => Some(in_flight),
LimiterResult::Forbidden => {
return Err(trc::LimitEvent::ConcurrentRequest
.into_err()
.id(tag.clone()));
.id(tag.clone()))
}
LimiterResult::Disabled => None,
};
// Create session