use std::sync::Arc;
use tokio::sync::Semaphore;
pub struct BackpressureController {
max_concurrent_requests: usize,
semaphore: Arc<Semaphore>,
}
pub struct BackpressurePermit {
_permit: tokio::sync::OwnedSemaphorePermit,
}
impl BackpressureController {
pub fn new(max_concurrent_requests: usize) -> Self {
Self {
max_concurrent_requests,
semaphore: Arc::new(Semaphore::new(max_concurrent_requests)),
}
}
pub async fn acquire_permit(&self) -> Result<BackpressurePermit, BackpressureError> {
let permit = self
.semaphore
.clone()
.acquire_owned()
.await
.map_err(|_| BackpressureError::SemaphoreClosed)?;
Ok(BackpressurePermit { _permit: permit })
}
pub fn max_concurrent_requests(&self) -> usize {
self.max_concurrent_requests
}
pub fn available_permits(&self) -> usize {
self.semaphore.available_permits()
}
}
#[derive(Debug, thiserror::Error)]
pub enum BackpressureError {
#[error("Semaphore is closed")]
SemaphoreClosed,
#[error("No permits available")]
NoPermitsAvailable,
}