Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 3 additions & 0 deletions crates/rds-agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,9 @@ thiserror.workspace = true
tokio.workspace = true
tracing.workspace = true

[target.'cfg(target_os = "macos")'.dependencies]
libc = "0.2"

[dev-dependencies]
rustix = { workspace = true, features = ["process"] }
iroh-relay = { workspace = true, features = ["server"] }
Expand Down
25 changes: 25 additions & 0 deletions crates/rds-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ mod authz;
mod limits;
mod revocations;
pub mod settings;
mod sys;
use authz::{ConnAuthz, ConnectionLifetime, authorize};
pub use limits::AgentLimits;
pub use revocations::{RevocationFeed, RevocationPolicy, watch_revocations};
Expand Down Expand Up @@ -317,6 +318,7 @@ pub struct Agent {
limits: AgentLimits,
admission: Arc<Semaphore>,
stream_counter: limits::StreamCounter,
gate: Option<limits::ResourceGate>,
}

/// Weak, in-memory observation: keeping an exporter alive never owns agent I/O.
Expand Down Expand Up @@ -349,6 +351,14 @@ impl AgentMetrics {
u64::from(agent.policy.grants_required()),
),
]);
// Process footprint where the kernel reports it; an unobservable
// platform simply omits the key rather than inventing a number.
if let Some(fds) = sys::open_fds() {
values.insert("rds_agent_process_fds", fds as u64);
}
if let Some(rss) = sys::rss_bytes() {
values.insert("rds_agent_process_rss_bytes", rss);
}
let grants = agent.policy.active_grants.try_lock().ok();
values.insert("rds_agent_active_grants_known", u64::from(grants.is_some()));
if let Some(grants) = grants {
Expand Down Expand Up @@ -379,6 +389,7 @@ impl Agent {
limits,
admission: Arc::new(Semaphore::new(limits.connections())),
stream_counter: Default::default(),
gate: limits::ResourceGate::new(limits),
}
}

Expand All @@ -387,6 +398,7 @@ impl Agent {
pub fn with_limits(mut self, limits: AgentLimits) -> Self {
self.limits = limits;
self.admission = Arc::new(Semaphore::new(limits.connections()));
self.gate = limits::ResourceGate::new(limits);
self
}

Expand Down Expand Up @@ -428,6 +440,14 @@ impl Agent {
}
incoming = self.endpoint.accept() => {
let Some(incoming) = incoming else { break; };
if self.gate.as_ref().is_some_and(|gate| !gate.allows()) {
rds_observe::emit(rds_observe::Event::ConnectionBudgetExhausted);
// Dropping Incoming refuses the handshake without a
// parked application task or a new connection slot.
drop(incoming);
debug!("connection refused: process resource budget exceeded");
continue;
}
let Ok(permit) = self.admission.clone().try_acquire_owned() else {
rds_observe::emit(rds_observe::Event::ConnectionBudgetExhausted);
// Dropping Incoming refuses the handshake without a
Expand Down Expand Up @@ -490,6 +510,11 @@ impl Agent {
conn.close(5u32.into(), b"invalid grant stream budget");
anyhow::bail!("grant mode requires at least two stream slots");
}
if self.gate.as_ref().is_some_and(|gate| !gate.allows()) {
rds_observe::emit(rds_observe::Event::ConnectionBudgetExhausted);
conn.close(5u32.into(), b"agent process resource budget exceeded");
anyhow::bail!("agent process resource budget exceeded");
}
let Ok(_permit) = self.admission.clone().try_acquire_owned() else {
rds_observe::emit(rds_observe::Event::ConnectionBudgetExhausted);
conn.close(5u32.into(), b"agent connection budget exhausted");
Expand Down
113 changes: 111 additions & 2 deletions crates/rds-agent/src/limits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,19 +4,24 @@ use std::num::NonZeroU16;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};

/// Per-agent connection slots and per-connection service tasks. The constructor
/// requires positive values; callers may choose smaller budgets for their host.
/// Per-agent connection slots and per-connection service tasks, plus an
/// optional process resource ceiling. The constructor requires positive
/// values; callers may choose smaller budgets for their host.
#[derive(Clone, Copy, Debug)]
pub struct AgentLimits {
connections: usize,
streams: usize,
max_fds: Option<u64>,
max_rss_bytes: Option<u64>,
}

impl Default for AgentLimits {
fn default() -> Self {
Self {
connections: 32,
streams: 64,
max_fds: None,
max_rss_bytes: None,
}
}
}
Expand All @@ -26,9 +31,29 @@ impl AgentLimits {
Self {
connections: usize::from(connections.get()),
streams: usize::from(streams.get()),
..Default::default()
}
}

/// Process-level ceiling: refuse admissions once the process holds
/// `max_fds` descriptors or `max_rss_mb` resident MiB. `None` leaves
/// that quantity ungated.
pub fn with_process_budget(mut self, max_fds: Option<u64>, max_rss_mb: Option<u64>) -> Self {
self.max_fds = max_fds;
self.max_rss_bytes = max_rss_mb.and_then(|mb| mb.checked_mul(1024 * 1024));
self
}

/// Configured fd ceiling, if any.
pub fn max_fds(self) -> Option<u64> {
self.max_fds
}

/// Configured resident-set ceiling in bytes, if any.
pub fn max_rss_bytes(self) -> Option<u64> {
self.max_rss_bytes
}

/// Pending handshakes plus admitted connections, shared by run and serve.
pub fn connections(self) -> usize {
self.connections
Expand Down Expand Up @@ -62,3 +87,87 @@ impl Drop for StreamTask {
self.0.fetch_sub(1, Ordering::Relaxed);
}
}

/// How often the process budget re-observes the kernel's view. Re-statting
/// per accepted connection would put procfs/`dev` scans on the hot path.
const SAMPLE_INTERVAL: std::time::Duration = std::time::Duration::from_millis(200);

/// Process-level admission ceiling. Ungated (`None` limits) or
/// unobservable (platform reports neither fds nor RSS) configurations
/// never refuse — a bound only exists where the kernel actually reports it.
pub(super) struct ResourceGate {
max_fds: Option<u64>,
max_rss_bytes: Option<u64>,
cache: std::sync::Mutex<(std::time::Instant, bool)>,
}

impl ResourceGate {
pub fn new(limits: AgentLimits) -> Option<Self> {
if limits.max_fds.is_none() && limits.max_rss_bytes.is_none() {
return None;
}
Some(Self {
max_fds: limits.max_fds,
max_rss_bytes: limits.max_rss_bytes,
cache: std::sync::Mutex::new((std::time::Instant::now() - 2 * SAMPLE_INTERVAL, true)),
})
}

/// True while the observed process usage fits the configured budget.
/// An unobservable quantity contributes nothing — the gate reports
/// only what the kernel actually showed it.
pub fn allows(&self) -> bool {
let mut cache = crate::lock(&self.cache);
let (at, verdict) = *cache;
if at.elapsed() < SAMPLE_INTERVAL {
return verdict;
}
let mut ok = true;
let mut observed = false;
if let Some(max) = self.max_fds
&& let Some(fds) = crate::sys::open_fds()
{
observed = true;
ok &= (fds as u64) < max;
}
if let Some(max) = self.max_rss_bytes
&& let Some(rss) = crate::sys::rss_bytes()
{
observed = true;
ok &= rss < max;
}
// Nothing observable: keep serving rather than gate on a guess.
let verdict = ok || !observed;
*cache = (std::time::Instant::now(), verdict);
verdict
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn no_budget_means_no_gate() {
assert!(ResourceGate::new(AgentLimits::default()).is_none());
assert!(
ResourceGate::new(AgentLimits::default().with_process_budget(None, None)).is_none()
);
}

#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn impossible_fd_ceiling_refuses() {
let limits = AgentLimits::default().with_process_budget(Some(1), None);
let gate = ResourceGate::new(limits).expect("budget configured");
assert!(!gate.allows());
}

#[cfg(any(target_os = "linux", target_os = "macos"))]
#[test]
fn generous_ceiling_allows() {
let limits = AgentLimits::default().with_process_budget(Some(u64::MAX), Some(u64::MAX));
let gate = ResourceGate::new(limits).expect("budget configured");
assert!(gate.allows());
}
}
28 changes: 22 additions & 6 deletions crates/rds-agent/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,14 @@ struct Cli {
/// one slot is reserved from service bodies for authorization/renewal.
#[arg(long)]
max_streams: Option<std::num::NonZeroU16>,
/// Process-wide open-descriptor ceiling; new connections are refused
/// while the process holds this many or more (Linux/macOS).
#[arg(long)]
max_fds: Option<std::num::NonZeroU64>,
/// Process-wide resident-set ceiling in MiB; new connections are
/// refused while resident memory meets or exceeds it.
#[arg(long)]
max_rss_mb: Option<std::num::NonZeroU64>,
/// Inbound connection handshake deadline in seconds (1..=3600).
#[arg(long)]
handshake_timeout: Option<u64>,
Expand Down Expand Up @@ -212,6 +220,8 @@ async fn run(cli: Cli) -> anyhow::Result<()> {
revocations_interval: cli.revocations_interval,
max_connections: cli.max_connections,
max_streams: cli.max_streams,
max_fds: cli.max_fds,
max_rss_mb: cli.max_rss_mb,
handshake_timeout: cli.handshake_timeout,
hello_timeout: cli.hello_timeout,
authz_timeout: cli.authz_timeout,
Expand Down Expand Up @@ -399,12 +409,18 @@ async fn run(cli: Cli) -> anyhow::Result<()> {
};

let agent = std::sync::Arc::new(
Agent::new(endpoint, policy).with_limits(AgentLimits::new(
resolved
.max_connections
.unwrap_or(std::num::NonZeroU16::new(32).expect("positive limit")),
max_streams,
)),
Agent::new(endpoint, policy).with_limits(
AgentLimits::new(
resolved
.max_connections
.unwrap_or(std::num::NonZeroU16::new(32).expect("positive limit")),
max_streams,
)
.with_process_budget(
resolved.max_fds.map(std::num::NonZeroU64::get),
resolved.max_rss_mb.map(std::num::NonZeroU64::get),
),
),
);
let metrics = agent.metrics();
let mut control =
Expand Down
53 changes: 52 additions & 1 deletion crates/rds-agent/src/settings.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@

use std::collections::BTreeSet;
use std::io::Read;
use std::num::NonZeroU16;
use std::num::{NonZeroU16, NonZeroU64};
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::time::Duration;
Expand Down Expand Up @@ -190,6 +190,13 @@ pub struct AuthoritySettings {
pub struct LimitSettings {
pub max_connections: Option<NonZeroU16>,
pub max_streams: Option<NonZeroU16>,
/// Process-wide open-descriptor ceiling; once the kernel reports the
/// process at or above it, new connections are refused until usage
/// falls. Absent = ungated.
pub max_fds: Option<NonZeroU64>,
/// Process-wide resident-set ceiling in MiB, same admission rule as
/// `max_fds`. Absent = ungated.
pub max_rss_mb: Option<NonZeroU64>,
}

/// Connection-admission and stream-greeting deadlines in seconds, each
Expand Down Expand Up @@ -263,6 +270,8 @@ pub struct AgentOverrides {
pub revocations_interval: Option<u64>,
pub max_connections: Option<NonZeroU16>,
pub max_streams: Option<NonZeroU16>,
pub max_fds: Option<NonZeroU64>,
pub max_rss_mb: Option<NonZeroU64>,
pub handshake_timeout: Option<u64>,
pub hello_timeout: Option<u64>,
pub authz_timeout: Option<u64>,
Expand Down Expand Up @@ -299,6 +308,10 @@ pub struct ResolvedAgent {
pub revocations_interval: Option<u64>,
pub max_connections: Option<NonZeroU16>,
pub max_streams: Option<NonZeroU16>,
/// Process fd ceiling (see `LimitSettings::max_fds`).
pub max_fds: Option<NonZeroU64>,
/// Process resident-set ceiling in MiB.
pub max_rss_mb: Option<NonZeroU64>,
pub timeouts: Option<TimeoutPolicy>,
}

Expand Down Expand Up @@ -467,6 +480,12 @@ impl AgentSettings {
if flags.max_streams.is_some() {
self.limits.max_streams = flags.max_streams;
}
if flags.max_fds.is_some() {
self.limits.max_fds = flags.max_fds;
}
if flags.max_rss_mb.is_some() {
self.limits.max_rss_mb = flags.max_rss_mb;
}
if flags.handshake_timeout.is_some() {
self.timeouts.handshake_secs = flags.handshake_timeout;
}
Expand Down Expand Up @@ -712,6 +731,8 @@ impl AgentSettings {
revocations_interval: revocations.and_then(|r| r.interval_secs),
max_connections: self.limits.max_connections,
max_streams: self.limits.max_streams,
max_fds: self.limits.max_fds,
max_rss_mb: self.limits.max_rss_mb,
timeouts,
})
}
Expand Down Expand Up @@ -870,6 +891,36 @@ mod tests {
AgentSettings::from_json(br#"{"schema_version":1,"limits":{"max_streams":0}}"#)
.is_err()
);
assert!(
AgentSettings::from_json(br#"{"schema_version":1,"limits":{"max_fds":0}}"#).is_err()
);
assert!(
AgentSettings::from_json(br#"{"schema_version":1,"limits":{"max_rss_mb":0}}"#).is_err()
);
}

#[test]
fn resource_budgets_merge_and_resolve() {
let resolved =
resolve_ok(r#"{"schema_version":1,"limits":{"max_fds":512,"max_rss_mb":256}}"#);
assert_eq!(resolved.max_fds.unwrap().get(), 512);
assert_eq!(resolved.max_rss_mb.unwrap().get(), 256);

// Flag values replace file values per-field.
let settings =
parse(r#"{"schema_version":1,"limits":{"max_fds":128}}"#).apply(AgentOverrides {
max_fds: NonZeroU64::new(256),
max_rss_mb: NonZeroU64::new(64),
..Default::default()
});
let resolved = settings.resolve().unwrap();
assert_eq!(resolved.max_fds.unwrap().get(), 256);
assert_eq!(resolved.max_rss_mb.unwrap().get(), 64);

// Absent everywhere resolves to no process gate.
let resolved = resolve_ok(r#"{"schema_version":1}"#);
assert!(resolved.max_fds.is_none());
assert!(resolved.max_rss_mb.is_none());
}

#[test]
Expand Down
Loading
Loading