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
8 changes: 3 additions & 5 deletions crates/rds-agent/src/authz.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,6 @@ use tokio::task::JoinHandle;

use crate::{AgentPolicy, lock};

const AUTHZ_REPLY_TIMEOUT: Duration = Duration::from_secs(15);

enum State {
Open,
Pending,
Expand Down Expand Up @@ -317,7 +315,7 @@ impl ConnAuthz {
) -> Result<Option<Arc<VerifiedGrant>>, ScopeError> {
// Subscribe before reading the state so commit/close cannot be missed.
let mut changed = self.changed.subscribe();
tokio::time::timeout(AUTHZ_REPLY_TIMEOUT, async {
tokio::time::timeout(policy.timeouts.authz, async {
loop {
match self.scope(policy) {
Err(ScopeError::Authorizing) => {
Expand Down Expand Up @@ -444,7 +442,7 @@ pub(crate) async fn authorize(
if let Err(message) = begin {
rds_observe::request_refused(rds_observe::Reason::Denied);
let _ = tokio::time::timeout(
AUTHZ_REPLY_TIMEOUT,
policy.timeouts.authz,
write_frame(
&mut send,
&HelloAck::Error {
Expand Down Expand Up @@ -506,7 +504,7 @@ pub(crate) async fn authorize(
};
// Watcher and reservation are already owned while the reply is in
// flight. Service admission remains pending until this completes.
tokio::time::timeout(AUTHZ_REPLY_TIMEOUT, write_frame(&mut send, &HelloAck::Ok)).await??;
tokio::time::timeout(policy.timeouts.authz, write_frame(&mut send, &HelloAck::Ok)).await??;
if let Some(next) = next {
authz
.commit_renewal(next, policy, conn)
Expand Down
11 changes: 10 additions & 1 deletion crates/rds-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ use rds_observe::{Reason, conn_span, next_session_id};
/// bounded here so silent streams cost seconds, not the session.
const HELLO_TIMEOUT: Duration = Duration::from_secs(15);
const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(15);
const AUTHZ_REPLY_TIMEOUT: Duration = Duration::from_secs(15);
const SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);

/// Operational deadlines on the serving side. Deployments tune these at
Expand All @@ -67,13 +68,21 @@ pub struct TimeoutPolicy {
/// Per-stream `StreamHello` read deadline — and the reply deadline for
/// greeting refusals that must not outlive a parked peer task.
pub hello: Duration,
/// Authorization-path reply budget: refusal answers and the final
/// `HelloAck` write while watcher/reservation state is already owned.
pub authz: Duration,
/// Join budget for established connection tasks during shutdown; the
/// runner never waits unboundedly for drained peers.
pub shutdown: Duration,
}

impl Default for TimeoutPolicy {
fn default() -> Self {
Self {
handshake: HANDSHAKE_TIMEOUT,
hello: HELLO_TIMEOUT,
authz: AUTHZ_REPLY_TIMEOUT,
shutdown: SHUTDOWN_TIMEOUT,
}
}
}
Expand Down Expand Up @@ -459,7 +468,7 @@ impl Agent {
}
// Endpoint closure wakes established connections and handshakes. Let
// their normal paths join service workers before this runner returns.
if tokio::time::timeout(SHUTDOWN_TIMEOUT, async {
if tokio::time::timeout(self.policy.timeouts.shutdown, async {
while let Some(result) = connections.join_next().await {
if let Err(error) = result {
debug!(%error, "connection task ended");
Expand Down
12 changes: 11 additions & 1 deletion crates/rds-agent/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,8 @@ use rds_agent::{
Agent, AgentLimits, AgentOverrides, AgentPolicy, AgentSettings, Role, ServiceName,
};
use rds_net::{
EndpointOverrides, EndpointSettings, Ticket, acquire_key, bind_endpoint, default_key_path,
EndpointOverrides, EndpointSettings, RetryPolicy, Ticket, acquire_key, bind_endpoint,
default_key_path,
};

#[derive(Parser)]
Expand Down Expand Up @@ -83,6 +84,12 @@ struct Cli {
/// Per-stream greeting read deadline in seconds (1..=3600).
#[arg(long)]
hello_timeout: Option<u64>,
/// Authorization-path reply budget in seconds (1..=3600).
#[arg(long)]
authz_timeout: Option<u64>,
/// Join budget for established connections during shutdown (1..=3600).
#[arg(long)]
shutdown_timeout: Option<u64>,
/// Directory HTTP(S) origin or legacy IP:port; the agent publishes its
/// signed record and keeps it fresh.
#[arg(long)]
Expand Down Expand Up @@ -207,6 +214,8 @@ async fn run(cli: Cli) -> anyhow::Result<()> {
max_streams: cli.max_streams,
handshake_timeout: cli.handshake_timeout,
hello_timeout: cli.hello_timeout,
authz_timeout: cli.authz_timeout,
shutdown_timeout: cli.shutdown_timeout,
});
merged.validate()?;
// Budget and authority cross-checks run on the merged document before
Expand Down Expand Up @@ -375,6 +384,7 @@ async fn run(cli: Cli) -> anyhow::Result<()> {
directory: client,
services,
ttl: std::time::Duration::from_secs(resolved.record_ttl.unwrap_or(300)),
retry: RetryPolicy::default(),
},
);
match announced {
Expand Down
92 changes: 69 additions & 23 deletions crates/rds-agent/src/settings.rs
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,10 @@ pub struct LimitSettings {
pub struct TimeoutSettings {
pub handshake_secs: Option<u64>,
pub hello_secs: Option<u64>,
/// Authorization-path reply budget (refusal and `HelloAck` writes).
pub authz_secs: Option<u64>,
/// Join budget for established connection tasks during shutdown.
pub shutdown_secs: Option<u64>,
}

/// Agent-level file schema; holds policy identities, never secret key
Expand Down Expand Up @@ -261,6 +265,8 @@ pub struct AgentOverrides {
pub max_streams: Option<NonZeroU16>,
pub handshake_timeout: Option<u64>,
pub hello_timeout: Option<u64>,
pub authz_timeout: Option<u64>,
pub shutdown_timeout: Option<u64>,
}

/// The merged, typed configuration the binary consumes.
Expand Down Expand Up @@ -467,6 +473,12 @@ impl AgentSettings {
if flags.hello_timeout.is_some() {
self.timeouts.hello_secs = flags.hello_timeout;
}
if flags.authz_timeout.is_some() {
self.timeouts.authz_secs = flags.authz_timeout;
}
if flags.shutdown_timeout.is_some() {
self.timeouts.shutdown_secs = flags.shutdown_timeout;
}
self
}

Expand Down Expand Up @@ -562,9 +574,14 @@ impl AgentSettings {
));
}
}
for secs in [self.timeouts.handshake_secs, self.timeouts.hello_secs]
.into_iter()
.flatten()
for secs in [
self.timeouts.handshake_secs,
self.timeouts.hello_secs,
self.timeouts.authz_secs,
self.timeouts.shutdown_secs,
]
.into_iter()
.flatten()
{
if secs == 0 || secs > MAX_TIMEOUT.as_secs() {
return Err(AgentConfigError::Invalid(
Expand Down Expand Up @@ -646,22 +663,29 @@ impl AgentSettings {
let authority = &self.authority;
let registry = authority.registry.as_ref();
let revocations = authority.revocations.as_ref();
let timeouts =
if self.timeouts.handshake_secs.is_some() || self.timeouts.hello_secs.is_some() {
let default = TimeoutPolicy::default();
Some(TimeoutPolicy {
handshake: self
.timeouts
.handshake_secs
.map_or(default.handshake, Duration::from_secs),
hello: self
.timeouts
.hello_secs
.map_or(default.hello, Duration::from_secs),
})
} else {
None
};
let timeouts = if self.timeouts != TimeoutSettings::default() {
let default = TimeoutPolicy::default();
Some(TimeoutPolicy {
handshake: self
.timeouts
.handshake_secs
.map_or(default.handshake, Duration::from_secs),
hello: self
.timeouts
.hello_secs
.map_or(default.hello, Duration::from_secs),
authz: self
.timeouts
.authz_secs
.map_or(default.authz, Duration::from_secs),
shutdown: self
.timeouts
.shutdown_secs
.map_or(default.shutdown, Duration::from_secs),
})
} else {
None
};
Ok(ResolvedAgent {
services,
ssh_target,
Expand Down Expand Up @@ -831,6 +855,8 @@ mod tests {
let bad = [
r#"{"schema_version":1,"timeouts":{"handshake_secs":0}}"#,
r#"{"schema_version":1,"timeouts":{"hello_secs":3601}}"#,
r#"{"schema_version":1,"timeouts":{"authz_secs":0}}"#,
r#"{"schema_version":1,"timeouts":{"shutdown_secs":3601}}"#,
r#"{"schema_version":1,"authority":{"grant_ttl_secs":0}}"#,
r#"{"schema_version":1,"authority":{"grant_ttl_secs":86401}}"#,
r#"{"schema_version":1,"authority":{"directory":"http://localhost:9","issuers":["a"],"revocations":{"key":"k","interval_secs":0}}}"#,
Expand Down Expand Up @@ -943,13 +969,33 @@ mod tests {
assert_eq!(policy.hello, Duration::from_secs(3));
// Partial timeout config preserves the other default.
let resolved = resolve_ok(r#"{"schema_version":1,"timeouts":{"hello_secs":3}}"#);
assert_eq!(
resolved.timeouts.unwrap().handshake,
TimeoutPolicy::default().handshake
);
let policy = resolved.timeouts.unwrap();
assert_eq!(policy.handshake, TimeoutPolicy::default().handshake);
assert_eq!(policy.authz, TimeoutPolicy::default().authz);
assert_eq!(policy.shutdown, TimeoutPolicy::default().shutdown);
// The reply/shutdown classes resolve identically.
let resolved =
resolve_ok(r#"{"schema_version":1,"timeouts":{"authz_secs":7,"shutdown_secs":2}}"#);
let policy = resolved.timeouts.unwrap();
assert_eq!(policy.authz, Duration::from_secs(7));
assert_eq!(policy.shutdown, Duration::from_secs(2));
assert_eq!(policy.handshake, TimeoutPolicy::default().handshake);
assert_eq!(resolve_ok(r#"{"schema_version":1}"#).timeouts, None);
}

#[test]
fn timeout_flags_override_file() {
let settings =
parse(r#"{"schema_version":1,"timeouts":{"authz_secs":3}}"#).apply(AgentOverrides {
authz_timeout: Some(11),
shutdown_timeout: Some(4),
..Default::default()
});
let policy = settings.resolve().unwrap().timeouts.unwrap();
assert_eq!(policy.authz, Duration::from_secs(11));
assert_eq!(policy.shutdown, Duration::from_secs(4));
}

#[test]
fn tenant_and_revision_binding() {
// Binding claims without issuers would never evaluate — refused
Expand Down
1 change: 1 addition & 0 deletions crates/rds-bench/src/scenario.rs
Original file line number Diff line number Diff line change
Expand Up @@ -566,6 +566,7 @@ async fn resolve_connect(p: &Params) -> anyhow::Result<BenchReport> {
directory: directory.clone(),
services: vec![rds_discovery::Service::Ping],
ttl: Duration::from_secs(120),
retry: rds_net::RetryPolicy::default(),
},
)?;
let mut policy = AgentPolicy::ssh_only(("127.0.0.1".into(), 9));
Expand Down
5 changes: 3 additions & 2 deletions crates/rds-net/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,9 @@ noq = { workspace = true, optional = true }
# Non-optional: iroh runs on noq internally, so it is in the tree either
# way — needed for transport tuning (BBRv3, windows) on the iroh backend.
noq-proto.workspace = true
rand = { workspace = true, optional = true }
# Non-optional: the announce loop's publish-retry jitter uses it, which is
# on in every build — and iroh/noq already pull rand transitively anyway.
rand.workspace = true
postcard = { workspace = true, features = ["alloc", "use-std"] }
rds-core.workspace = true
rds-discovery.workspace = true
Expand Down Expand Up @@ -47,7 +49,6 @@ transport-noq = [
"dep:tokio-stream",
"dep:tokio-util",
"dep:blake3",
"dep:rand",
]

[lints]
Expand Down
Loading
Loading