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
142 changes: 68 additions & 74 deletions Cargo.lock

Large diffs are not rendered by default.

10 changes: 7 additions & 3 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,6 @@ repository = "https://github.com/NDDev-OpenNetwork/remote-device-sync"

[workspace.dependencies]
anyhow = "1"
async-trait = "0.1"
blake3 = "1"
bytes = "1"
clap = { version = "4", features = ["derive"] }
Expand All @@ -61,8 +60,13 @@ rand = "0.10"
redb = "4.3"
rcgen = { version = "0.14", default-features = false, features = ["crypto", "pem", "ring"] }
rustix = { version = "1", features = ["fs"] }
# Pinned exactly: the release carries PTY encoding, sensitive-debug
# redaction, KEX/parser fixes and stalled-write timeout enforcement that
# the SSH surface relies on (docs/ssh.md).
russh = { version = "=0.63.3", default-features = false, features = ["ring"] }
rustls = { version = "0.23", default-features = false, features = ["std", "logging", "tls12", "ring"] }
# No tls12: every TLS peer in the estate is our own rustls binary and
# QUIC is 1.3-only anyway — 1.2 has no legitimate client here.
rustls = { version = "0.23", default-features = false, features = ["std", "logging", "ring"] }
rustls-pki-types = { version = "1", features = ["alloc"] }
rds-agent = { path = "crates/rds-agent" }
rds-audio = { path = "crates/rds-audio" }
Expand All @@ -83,7 +87,7 @@ thiserror = "2"
tokio = { version = "1", features = ["full"] }
tokio-stream = "0.1"
tokio-util = { version = "0.7.19", features = ["rt"] }
tokio-rustls = { version = "0.26", default-features = false, features = ["ring", "tls12"] }
tokio-rustls = { version = "0.26", default-features = false, features = ["ring"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
url = "2"
Expand Down
111 changes: 83 additions & 28 deletions crates/rds-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -132,8 +132,9 @@ pub struct AgentPolicy {
/// Data-plane services this agent answers. `None` keeps the implicit
/// set — `Tcp` plus `Desktop` when compiled and `Sync` when `sync_dir`
/// is configured. `Some` is the explicit set; `Ping`/`Info` are the
/// always-on control plane and are never gated. A disabled service is
/// refused before any grant or service work runs.
/// always-on control plane — never gated by deployment policy, though
/// a grant's service scope may still refuse them. A disabled service
/// is refused before any grant or service work runs.
pub services: Option<BTreeSet<ServiceKind>>,
/// Admission and greeting deadlines applied per connection/stream.
pub timeouts: TimeoutPolicy,
Expand Down Expand Up @@ -209,11 +210,13 @@ impl AgentPolicy {
pub fn effective_services(&self) -> BTreeSet<ServiceKind> {
let mut set = BTreeSet::from([ServiceKind::Ping, ServiceKind::Info]);
match &self.services {
// One honest source: an explicit `desktop` entry only counts when
// this binary can actually serve it, matching the implicit arm
// and the directory announcement. `validate` still rejects the
// flag combination at startup.
Some(explicit) => set.extend(explicit.iter().copied().filter(|k| {
matches!(
k,
ServiceKind::Tcp | ServiceKind::Desktop | ServiceKind::Sync
)
matches!(k, ServiceKind::Tcp | ServiceKind::Sync)
|| (matches!(k, ServiceKind::Desktop) && cfg!(feature = "desktop"))
})),
None => {
set.insert(ServiceKind::Tcp);
Expand Down Expand Up @@ -263,8 +266,12 @@ impl AgentPolicy {
}
if self.timeouts.handshake.is_zero()
|| self.timeouts.hello.is_zero()
|| self.timeouts.authz.is_zero()
|| self.timeouts.shutdown.is_zero()
|| self.timeouts.handshake > MAX_TIMEOUT
|| self.timeouts.hello > MAX_TIMEOUT
|| self.timeouts.authz > MAX_TIMEOUT
|| self.timeouts.shutdown > MAX_TIMEOUT
{
return Err("timeouts must be between 1 and 3600 seconds");
}
Expand Down Expand Up @@ -506,6 +513,12 @@ impl Agent {

/// Serve a single already-established connection.
pub async fn serve(&self, conn: Connection) -> anyhow::Result<()> {
// Programmatic callers reach serve() without run()'s preflight —
// the same policy contract applies either way.
if let Err(why) = self.policy.validate() {
conn.close(5u32.into(), b"invalid agent policy");
anyhow::bail!("invalid agent policy: {why}");
}
if self.policy.grants_required() && self.limits.streams() < 2 {
conn.close(5u32.into(), b"invalid grant stream budget");
anyhow::bail!("grant mode requires at least two stream slots");
Expand Down Expand Up @@ -665,13 +678,19 @@ async fn serve_stream(
&& !policy.service_enabled(kind)
{
rds_observe::request_refused(Reason::Denied);
write_frame(
&mut send,
&HelloAck::Error {
message: format!("service {kind:?} not enabled on this agent"),
},
// Greeting refusals are bounded by the hello deadline — a stalled
// peer must not park this task on a refusal write.
tokio::time::timeout(
policy.timeouts.hello,
write_frame(
&mut send,
&HelloAck::Error {
message: format!("service {kind:?} not enabled on this agent"),
},
),
)
.await?;
.await??;
send.finish()?;
anyhow::bail!("service {kind:?} not enabled");
}
let grant = match authz.service_scope(&policy).await {
Expand All @@ -682,22 +701,54 @@ async fn serve_stream(
conn.close(2u32.into(), why.message().as_bytes());
}
let why = why.message();
write_frame(
&mut send,
&HelloAck::Error {
message: why.into(),
},
// Authorization-path answers are bounded by the authz budget.
tokio::time::timeout(
policy.timeouts.authz,
write_frame(
&mut send,
&HelloAck::Error {
message: why.into(),
},
),
)
.await?;
.await??;
send.finish()?;
anyhow::bail!("stream refused: {why}");
}
};
let scope_err = grant.as_ref().and_then(|g| scope_check(g, &hello).err());
if let Some(why) = scope_err {
rds_observe::request_refused(Reason::Denied);
write_frame(&mut send, &HelloAck::Error { message: why }).await?;
tokio::time::timeout(
policy.timeouts.authz,
write_frame(&mut send, &HelloAck::Error { message: why }),
)
.await??;
send.finish()?;
anyhow::bail!("stream outside grant scope");
}
// TCP preflight resolves the target and policy before the slot is
// consumed — a refused connect must not hold a service lane while
// its refusal is written. The service arm re-validates on its own.
if let StreamHello::TcpConnect { host, port } = &hello {
let refusal = match rds_core::TcpTarget::new(host, *port) {
Ok(target) if !policy.permits_tcp_target(&target) => {
Some(format!("tcp target {host}:{port} not permitted"))
}
Ok(_) => None,
Err(error) => Some(format!("invalid TCP destination: {error}")),
};
if let Some(message) = refusal {
rds_observe::request_refused(Reason::Denied);
tokio::time::timeout(
policy.timeouts.hello,
write_frame(&mut send, &HelloAck::Error { message }),
)
.await??;
send.finish()?;
anyhow::bail!("tcp target refused at preflight");
}
}
// Control greetings (Ping/Info) are short-lived and bypass the service
// pool; only long-lived data services consume it, so a full pool drains
// back to one free JoinSet lane for whatever hello arrives next.
Expand Down Expand Up @@ -798,13 +849,17 @@ async fn serve_stream(
tokio::io::copy_bidirectional(&mut tcp, &mut quic).await?;
}
Err(e) => {
// ErrorKind is a fixed vocabulary — the raw OS error
// string (errno text, platform internals) never
// crosses the wire. The bail below logs it locally.
write_frame(
&mut send,
&HelloAck::Error {
message: format!("connect {host}:{port} failed: {e}"),
message: format!("connect {host}:{port} failed: {}", e.kind()),
},
)
.await?;
anyhow::bail!("tcp connect {host}:{port} failed: {e}");
}
}
}
Expand Down Expand Up @@ -835,10 +890,11 @@ async fn serve_stream(
write_frame(
&mut send,
&HelloAck::Error {
message: format!("desktop unavailable: {e}"),
message: "desktop unavailable".into(),
},
)
.await?;
anyhow::bail!("desktop capability probe failed: {e}");
}
}
#[cfg(not(feature = "desktop"))]
Expand Down Expand Up @@ -903,16 +959,15 @@ async fn serve_stream(
// Grant v3 `sync_paths` entries passed the decoder's
// lexical checks; normalize `.`/empty components the
// same way `check_rel_path` normalizes requests so
// prefix matching compares like with like.
// prefix matching compares like with like. Split on
// both separators — `check_scope_path` admits `\` but
// `Path::components` on Unix does not.
let paths = g.payload.constraints.sync_paths.as_ref().map(|list| {
list.iter()
.map(|scope| {
std::path::Path::new(scope)
.components()
.filter_map(|c| match c {
std::path::Component::Normal(p) => Some(p),
_ => None,
})
scope
.split(['/', '\\'])
.filter(|p| !p.is_empty() && *p != ".")
.collect::<PathBuf>()
})
.collect::<Vec<_>>()
Expand Down
15 changes: 9 additions & 6 deletions crates/rds-agent/src/limits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,11 @@ impl AgentLimits {

/// Process-level ceiling: refuse admissions once the process holds
/// `max_fds` descriptors or `max_rss_mb` resident MiB. `None` leaves
/// that quantity ungated.
/// that quantity ungated. An oversized MiB value saturates at the
/// byte ceiling instead of silently disabling the gate.
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.max_rss_bytes = max_rss_mb.map(|mb| mb.saturating_mul(1024 * 1024));
self
}

Expand Down Expand Up @@ -98,7 +99,9 @@ const SAMPLE_INTERVAL: std::time::Duration = std::time::Duration::from_millis(20
pub(super) struct ResourceGate {
max_fds: Option<u64>,
max_rss_bytes: Option<u64>,
cache: std::sync::Mutex<(std::time::Instant, bool)>,
/// `None` forces the first check to sample immediately — subtracting
/// an interval from `Instant::now` can underflow on fresh processes.
cache: std::sync::Mutex<(Option<std::time::Instant>, bool)>,
}

impl ResourceGate {
Expand All @@ -109,7 +112,7 @@ impl ResourceGate {
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)),
cache: std::sync::Mutex::new((None, true)),
})
}

Expand All @@ -119,7 +122,7 @@ impl ResourceGate {
pub fn allows(&self) -> bool {
let mut cache = crate::lock(&self.cache);
let (at, verdict) = *cache;
if at.elapsed() < SAMPLE_INTERVAL {
if at.is_some_and(|at| at.elapsed() < SAMPLE_INTERVAL) {
return verdict;
}
let mut ok = true;
Expand All @@ -138,7 +141,7 @@ impl ResourceGate {
}
// Nothing observable: keep serving rather than gate on a guess.
let verdict = ok || !observed;
*cache = (std::time::Instant::now(), verdict);
*cache = (Some(std::time::Instant::now()), verdict);
verdict
}
}
Expand Down
15 changes: 15 additions & 0 deletions crates/rds-agent/src/settings.rs
Original file line number Diff line number Diff line change
Expand Up @@ -538,6 +538,21 @@ impl AgentSettings {
return Err(AgentConfigError::Invalid("at most 256 allowed peers"));
}
let authority = &self.authority;
if authority
.registry
.as_ref()
.and_then(|r| r.epoch)
.is_some_and(|e| e == 0)
|| authority
.revocations
.as_ref()
.and_then(|r| r.epoch)
.is_some_and(|e| e == 0)
{
return Err(AgentConfigError::Invalid(
"authority epochs must be positive",
));
}
if authority
.grant_ttl_secs
.is_some_and(|secs| secs == 0 || secs > MAX_GRANT_TTL_SECS)
Expand Down
35 changes: 27 additions & 8 deletions crates/rds-agent/src/sys.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,15 +3,34 @@
//! Platform coverage is honest: an unobservable quantity reports `None`
//! and the corresponding gate stays open rather than pretending a bound.

/// Open file descriptors owned by this process, where the kernel exposes
/// them (`/proc/self/fd` on Linux, `/dev/fd` on macOS).
#[cfg(any(target_os = "linux", target_os = "macos"))]
/// Open file descriptors owned by this process (`/proc/self/fd`).
#[cfg(target_os = "linux")]
pub(crate) fn open_fds() -> Option<usize> {
#[cfg(target_os = "linux")]
const FD_DIR: &str = "/proc/self/fd";
#[cfg(target_os = "macos")]
const FD_DIR: &str = "/dev/fd";
Some(std::fs::read_dir(FD_DIR).ok()?.count())
Some(std::fs::read_dir("/proc/self/fd").ok()?.count())
}

/// Open file descriptors via `proc_pidinfo` (`PROC_PIDLISTFDS`): a
/// zero-length query returns the fd-table byte size — the true count
/// without `/dev/fd`'s dependency on an fdesc mount or the +1 of the
/// `read_dir` descriptor itself.
#[cfg(target_os = "macos")]
#[allow(unsafe_code)]
pub(crate) fn open_fds() -> Option<usize> {
// SAFETY: a null buffer with a zero size is a documented size probe —
// the call writes nothing and returns the table's byte length.
let size = unsafe {
libc::proc_pidinfo(
libc::getpid(),
libc::PROC_PIDLISTFDS,
0,
std::ptr::null_mut(),
0,
)
};
if size <= 0 {
return None;
}
Some(size as usize / size_of::<libc::proc_fdinfo>())
}

#[cfg(not(any(target_os = "linux", target_os = "macos")))]
Expand Down
2 changes: 1 addition & 1 deletion crates/rds-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ rds-net.workspace = true
rds-discovery.workspace = true
rds-sync.workspace = true
rds-ssh.workspace = true
rustix = { workspace = true, features = ["termios", "event"] }
rustix = { workspace = true, features = ["fs", "termios", "event"] }
serde_json.workspace = true
tokio.workspace = true
tracing.workspace = true
Expand Down
Loading
Loading