use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
use crate::config::Config;
pub const WRITE_IDLE: Duration = Duration::from_secs(30);
pub const RESPONSE_LIFETIME: Duration = Duration::from_secs(15 * 60);
pub struct Limits {
downloads: Arc<Semaphore>,
active: Arc<Semaphore>,
active_permits: usize,
write_idle_millis: AtomicU64,
lifetime_millis: AtomicU64,
}
impl Limits {
pub fn new(config: &Config) -> Limits {
let active_permits = config.max_active_requests.get() as usize;
Limits {
downloads: Arc::new(Semaphore::new(config.max_artifact_downloads.get() as usize)),
active: Arc::new(Semaphore::new(active_permits)),
active_permits,
write_idle_millis: AtomicU64::new(WRITE_IDLE.as_millis() as u64),
lifetime_millis: AtomicU64::new(RESPONSE_LIFETIME.as_millis() as u64),
}
}
pub fn active_permit(&self) -> Option<OwnedSemaphorePermit> {
Arc::clone(&self.active).try_acquire_owned().ok()
}
pub fn active_available(&self) -> usize {
self.active.available_permits()
}
pub fn max_active(&self) -> usize {
self.active_permits
}
pub fn download_permit(&self) -> Option<OwnedSemaphorePermit> {
Arc::clone(&self.downloads).try_acquire_owned().ok()
}
pub fn download_available(&self) -> usize {
self.downloads.available_permits()
}
pub fn write_idle(&self) -> Duration {
Duration::from_millis(self.write_idle_millis.load(Ordering::Relaxed))
}
pub fn response_lifetime(&self) -> Duration {
Duration::from_millis(self.lifetime_millis.load(Ordering::Relaxed))
}
#[cfg(feature = "test-support")]
pub fn set_response_timeouts(&self, write_idle: Duration, lifetime: Duration) {
self.write_idle_millis
.store(write_idle.as_millis() as u64, Ordering::Relaxed);
self.lifetime_millis
.store(lifetime.as_millis() as u64, Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn config() -> Config {
Config::load(std::path::Path::new("config.sample.toml")).expect("the shipped sample")
}
#[test]
fn the_shipped_deadlines_are_thirty_seconds_and_fifteen_minutes() {
let limits = Limits::new(&config());
assert_eq!(limits.write_idle(), Duration::from_secs(30));
assert_eq!(limits.response_lifetime(), Duration::from_secs(900));
}
#[test]
fn an_exhausted_permit_is_refused_rather_than_queued() {
let mut config = config();
config.max_active_requests = std::num::NonZeroU32::new(1).expect("one");
let limits = Limits::new(&config);
let held = limits.active_permit().expect("the only permit");
assert!(limits.active_permit().is_none(), "the second is refused");
assert_eq!(limits.active_available(), 0);
drop(held);
assert_eq!(limits.active_available(), 1);
assert!(limits.active_permit().is_some());
}
}