use std::sync::atomic::{AtomicU32, Ordering};
#[derive(Debug, Clone)]
pub struct ResourceLimits {
pub ring_entries: u32,
pub max_in_flight: u32,
pub memory_limit: Option<u64>,
}
impl Default for ResourceLimits {
fn default() -> Self {
Self::detect()
}
}
impl ResourceLimits {
pub fn detect() -> Self {
let (memory_limit, memory_available) = Self::read_cgroup_memory_limit();
let ring_entries = if let Some(available) = memory_available {
let ring_memory = available / 100;
let entries = (ring_memory / 80) as u32; entries.clamp(256, 4096).next_power_of_two()
} else {
1024 };
let max_in_flight = if let Some(available) = memory_available {
let tracking_memory = available / 20;
let count = (tracking_memory / 128) as u32; count.clamp(64, 16384)
} else {
1024 };
Self {
ring_entries,
max_in_flight,
memory_limit,
}
}
fn read_cgroup_memory_limit() -> (Option<u64>, Option<u64>) {
#[cfg(target_os = "linux")]
{
if let Ok(contents) = std::fs::read_to_string("/sys/fs/cgroup/memory.max") {
let contents = contents.trim();
if contents != "max" {
if let Ok(limit) = contents.parse::<u64>() {
let usable = limit / 2;
return (Some(limit), Some(usable));
}
}
}
if let Ok(contents) =
std::fs::read_to_string("/sys/fs/cgroup/memory/memory.limit_in_bytes")
{
if let Ok(limit) = contents.trim().parse::<u64>() {
if limit < u64::MAX / 2 {
let usable = limit / 2;
return (Some(limit), Some(usable));
}
}
}
if let Ok(contents) = std::fs::read_to_string("/proc/meminfo") {
for line in contents.lines() {
if line.starts_with("MemAvailable:") {
if let Some(kb) = line.split_whitespace().nth(1) {
if let Ok(kb) = kb.parse::<u64>() {
let bytes = kb * 1024;
let usable = bytes / 10;
return (None, Some(usable));
}
}
}
}
}
(None, None)
}
#[cfg(not(target_os = "linux"))]
{
(None, None)
}
}
}
pub struct ResourceLimiter {
limits: ResourceLimits,
in_flight: AtomicU32,
}
impl ResourceLimiter {
pub fn new() -> Self {
Self {
limits: ResourceLimits::detect(),
in_flight: AtomicU32::new(0),
}
}
pub fn with_limits(limits: ResourceLimits) -> Self {
Self {
limits,
in_flight: AtomicU32::new(0),
}
}
pub fn can_submit(&self) -> bool {
self.in_flight.load(Ordering::Relaxed) < self.limits.max_in_flight
}
pub fn try_reserve(&self) -> bool {
let current = self.in_flight.load(Ordering::Relaxed);
if current >= self.limits.max_in_flight {
return false;
}
self.in_flight.fetch_add(1, Ordering::AcqRel) < self.limits.max_in_flight
}
pub fn release(&self) {
self.in_flight.fetch_sub(1, Ordering::AcqRel);
}
pub fn in_flight(&self) -> u32 {
self.in_flight.load(Ordering::Relaxed)
}
pub fn limits(&self) -> &ResourceLimits {
&self.limits
}
pub fn recommended_ring_entries(&self) -> u32 {
self.limits.ring_entries
}
}
impl Default for ResourceLimiter {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_resource_limits_detect() {
let limits = ResourceLimits::detect();
assert!(limits.ring_entries.is_power_of_two());
assert!(limits.ring_entries >= 256);
assert!(limits.max_in_flight >= 64);
}
#[test]
fn test_resource_limiter() {
let limiter = ResourceLimiter::with_limits(ResourceLimits {
ring_entries: 512,
max_in_flight: 10,
memory_limit: None,
});
assert!(limiter.can_submit());
assert_eq!(limiter.in_flight(), 0);
for _ in 0..10 {
assert!(limiter.try_reserve());
}
assert_eq!(limiter.in_flight(), 10);
assert!(!limiter.can_submit());
assert!(!limiter.try_reserve());
limiter.release();
assert_eq!(limiter.in_flight(), 9);
assert!(limiter.can_submit());
}
#[test]
fn test_resource_limiter_default() {
let limiter = ResourceLimiter::new();
assert!(limiter.recommended_ring_entries().is_power_of_two());
assert!(limiter.limits().max_in_flight > 0);
}
}