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
4 changes: 4 additions & 0 deletions Cargo.lock

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

6 changes: 6 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,12 @@ bind access to the controlled endpoint. `rds session renew --session <id>
--grant-file <path>` extends a live same-scope connection; automatic GDS issuance
is still pending. Old grants must be reissued, and CLI/agent IPC v4 upgraded together.

The agent's role, data-plane services, authority posture and deadlines can
also live in one versioned JSON document (`--agent-config`), with
`--role`/`--service`/`--no-service` as the flag surface; disabled services
are refused before grant machinery runs and `Info` advertises exactly what
is served. See [agent configuration](docs/agent-configuration.md).

Endpoint publication also persists its revision counter and exact retry bytes
(`rds-agent --record-state`, default beside the identity key). Records and
deletes use a versioned signature contract; old directories require the explicit
Expand Down
5 changes: 5 additions & 0 deletions crates/rds-agent/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,12 +25,17 @@ rds-desktop = { workspace = true, optional = true }
rds-net.workspace = true
rds-discovery.workspace = true
rds-sync.workspace = true
rustix.workspace = true
serde.workspace = true
serde_json.workspace = true
thiserror.workspace = true
tokio.workspace = true
tracing.workspace = true

[dev-dependencies]
rustix = { workspace = true, features = ["process"] }
iroh-relay = { workspace = true, features = ["server"] }
rand = { workspace = true }
rds-cli = { path = "../rds-cli" }

[lints]
Expand Down
171 changes: 152 additions & 19 deletions crates/rds-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@
//! TCP forwarding is restricted to an explicit set of `(host, port)`
//! targets; the default set is exactly the configured SSH socket.

use std::collections::HashSet;
use std::collections::{BTreeSet, HashSet};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
use std::time::Duration;
Expand All @@ -38,9 +38,13 @@ use tracing::{Instrument, debug, info, info_span, warn};
mod authz;
mod limits;
mod revocations;
pub mod settings;
use authz::{ConnAuthz, ConnectionLifetime, authorize};
pub use limits::AgentLimits;
pub use revocations::{RevocationFeed, RevocationPolicy, watch_revocations};
pub use settings::{
AgentConfigError, AgentOverrides, AgentSettings, ResolvedAgent, Role, ServiceName,
};

/// Session ids and the `rds.conn` span shape are minted by rds-observe so
/// both sides of a connection share the correlation convention.
Expand All @@ -53,6 +57,32 @@ const HELLO_TIMEOUT: Duration = Duration::from_secs(15);
const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(15);
const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);

/// Operational deadlines on the serving side. Deployments tune these at
/// the policy level; session and transfer internals keep their own
/// service-scoped budgets.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TimeoutPolicy {
/// Inbound connection-handshake budget.
pub handshake: Duration,
/// Per-stream `StreamHello` read deadline — and the reply deadline for
/// greeting refusals that must not outlive a parked peer task.
pub hello: Duration,
}

impl Default for TimeoutPolicy {
fn default() -> Self {
Self {
handshake: HANDSHAKE_TIMEOUT,
hello: HELLO_TIMEOUT,
}
}
}

/// The largest accepted timeout; guards against effectively-unbounded
/// deadline configuration (roughly one workday of idle handshake budget
/// is never the intent).
pub const MAX_TIMEOUT: Duration = Duration::from_secs(3600);

/// Mutex acquisition that survives a poisoned lock: every mutex here
/// guards plain data (a state word, an `Option`, a `HashSet`) whose
/// invariants a panic cannot corrupt, so refusing service forever
Expand Down Expand Up @@ -89,6 +119,14 @@ pub struct AgentPolicy {
/// Directory the `Sync` service may read/write under (WS6). `None`
/// disables sync entirely.
pub sync_dir: Option<PathBuf>,
/// 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.
pub services: Option<BTreeSet<ServiceKind>>,
/// Admission and greeting deadlines applied per connection/stream.
pub timeouts: TimeoutPolicy,
}

impl std::fmt::Debug for AgentPolicy {
Expand All @@ -100,6 +138,8 @@ impl std::fmt::Debug for AgentPolicy {
.field("issuers", &self.issuers.len())
.field("grant_max_ttl", &self.grant_max_ttl)
.field("sync_dir", &self.sync_dir)
.field("services", &self.services)
.field("timeouts", &self.timeouts)
.finish_non_exhaustive()
}
}
Expand All @@ -115,6 +155,8 @@ impl AgentPolicy {
denylist: watch::channel(Arc::new(RevocationPolicy::default())).0,
active_grants: Arc::new(Mutex::new(HashSet::new())),
sync_dir: None,
services: None,
timeouts: TimeoutPolicy::default(),
}
}

Expand All @@ -137,6 +179,76 @@ impl AgentPolicy {
!self.issuers.is_empty()
}

/// Every service this agent answers: the always-on `Ping`/`Info`
/// control plane plus either the explicit `services` set or the
/// implicit set derived from build features and configuration.
/// `Audio` is wire-reserved but unimplemented, so it is never in the
/// effective set even if it slips into an explicit one.
pub fn effective_services(&self) -> BTreeSet<ServiceKind> {
let mut set = BTreeSet::from([ServiceKind::Ping, ServiceKind::Info]);
match &self.services {
Some(explicit) => set.extend(explicit.iter().copied().filter(|k| {
matches!(
k,
ServiceKind::Tcp | ServiceKind::Desktop | ServiceKind::Sync
)
})),
None => {
set.insert(ServiceKind::Tcp);
if cfg!(feature = "desktop") {
set.insert(ServiceKind::Desktop);
}
if self.sync_dir.is_some() {
set.insert(ServiceKind::Sync);
}
}
}
set
}

/// Whether a service stream is admitted by deployment policy, before
/// any grant scope or per-service check.
fn service_enabled(&self, kind: ServiceKind) -> bool {
self.effective_services().contains(&kind)
}

/// Deployment preflight: an explicit `services` set may only name
/// implemented services whose prerequisites are configured, and the
/// timeout policy must stay inside sane bounds. Programmatic callers
/// should run this before binding; `Agent::run` runs it too.
pub fn validate(&self) -> Result<(), &'static str> {
if let Some(set) = &self.services {
for kind in set {
match kind {
ServiceKind::Tcp => {}
ServiceKind::Desktop if cfg!(feature = "desktop") => {}
ServiceKind::Desktop => {
return Err("desktop service requires the `desktop` build feature");
}
ServiceKind::Sync if self.sync_dir.is_some() => {}
ServiceKind::Sync => {
return Err("sync service requires a configured sync directory");
}
ServiceKind::Audio => {
return Err("audio service is reserved but not implemented");
}
_ => {
return Err("service set may only contain tcp, desktop or sync; \
ping and info are always served");
}
}
}
}
if self.timeouts.handshake.is_zero()
|| self.timeouts.hello.is_zero()
|| self.timeouts.handshake > MAX_TIMEOUT
|| self.timeouts.hello > MAX_TIMEOUT
{
return Err("timeouts must be between 1 and 3600 seconds");
}
Ok(())
}

/// Revoke a grant id — pushes onto the denylist and notifies every
/// live connection watcher. The value is retained without subscribers.
pub fn revoke(&self, id: GrantId) {
Expand Down Expand Up @@ -263,6 +375,9 @@ impl Agent {
!self.policy.grants_required() || self.limits.streams() >= 2,
"grant mode requires at least two stream slots"
);
self.policy
.validate()
.map_err(|why| anyhow::anyhow!("invalid agent policy: {why}"))?;
info!(id = %self.endpoint.id(), "agent listening");
rds_observe::emit(rds_observe::Event::ListenerReady);
let mut connections = JoinSet::new();
Expand Down Expand Up @@ -292,7 +407,7 @@ impl Agent {
let stream_counter = self.stream_counter.clone();
connections.spawn(async move {
let _permit = permit;
match tokio::time::timeout(HANDSHAKE_TIMEOUT, incoming).await {
match tokio::time::timeout(policy.timeouts.handshake, incoming).await {
Ok(Ok(conn)) => {
let span = conn_span(next_session_id());
span.record("peer", tracing::field::display(conn.remote_id()));
Expand Down Expand Up @@ -460,14 +575,14 @@ async fn serve_stream(
authz: Arc<ConnAuthz>,
desktop: bool,
) -> anyhow::Result<()> {
let hello: StreamHello = match tokio::time::timeout(HELLO_TIMEOUT, read_frame(&mut recv)).await
{
Ok(h) => h?,
Err(_) => {
rds_observe::request_refused(Reason::Timeout);
anyhow::bail!("stream hello timed out");
}
};
let hello: StreamHello =
match tokio::time::timeout(policy.timeouts.hello, read_frame(&mut recv)).await {
Ok(h) => h?,
Err(_) => {
rds_observe::request_refused(Reason::Timeout);
anyhow::bail!("stream hello timed out");
}
};
if let StreamHello::Authz(grant) = hello {
return rds_observe::observe(
rds_observe::Operation::GrantAuthorize,
Expand All @@ -482,6 +597,21 @@ async fn serve_stream(
)
.await;
}
// Deployment policy answers first: a disabled service is refused before
// grant machinery or per-service work runs.
if let Some(kind) = service_kind(&hello)
&& !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"),
},
)
.await?;
anyhow::bail!("service {kind:?} not enabled");
}
let grant = match authz.service_scope(&policy).await {
Ok(g) => g,
Err(why) => {
Expand Down Expand Up @@ -509,7 +639,7 @@ async fn serve_stream(
let Some(_service_slot) = authz.try_service_slot() else {
rds_observe::request_refused(Reason::BudgetExhausted);
tokio::time::timeout(
HELLO_TIMEOUT,
policy.timeouts.hello,
write_frame(
&mut send,
&HelloAck::Error {
Expand Down Expand Up @@ -546,14 +676,17 @@ async fn serve_stream(
version: env!("CARGO_PKG_VERSION").to_string(),
hostname: hostname(),
services: {
let mut s = vec![ServiceKind::Ping, ServiceKind::Info, ServiceKind::Tcp];
if desktop {
s.push(ServiceKind::Desktop);
}
if policy.sync_dir.is_some() {
s.push(ServiceKind::Sync);
}
s
let enabled = policy.effective_services();
[
ServiceKind::Ping,
ServiceKind::Info,
ServiceKind::Tcp,
ServiceKind::Desktop,
ServiceKind::Sync,
]
.into_iter()
.filter(|k| enabled.contains(k) && (*k != ServiceKind::Desktop || desktop))
.collect()
},
desktop: desktop_caps(desktop),
};
Expand Down
Loading
Loading