diff --git a/crates/rds-agent/src/authz.rs b/crates/rds-agent/src/authz.rs index 6d2c23b..4b9133f 100644 --- a/crates/rds-agent/src/authz.rs +++ b/crates/rds-agent/src/authz.rs @@ -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, @@ -317,7 +315,7 @@ impl ConnAuthz { ) -> Result>, 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) => { @@ -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 { @@ -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) diff --git a/crates/rds-agent/src/lib.rs b/crates/rds-agent/src/lib.rs index 3c869ef..3493674 100644 --- a/crates/rds-agent/src/lib.rs +++ b/crates/rds-agent/src/lib.rs @@ -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 @@ -67,6 +68,12 @@ 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 { @@ -74,6 +81,8 @@ impl Default for TimeoutPolicy { Self { handshake: HANDSHAKE_TIMEOUT, hello: HELLO_TIMEOUT, + authz: AUTHZ_REPLY_TIMEOUT, + shutdown: SHUTDOWN_TIMEOUT, } } } @@ -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"); diff --git a/crates/rds-agent/src/main.rs b/crates/rds-agent/src/main.rs index 4bca130..b6e9d9e 100644 --- a/crates/rds-agent/src/main.rs +++ b/crates/rds-agent/src/main.rs @@ -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)] @@ -83,6 +84,12 @@ struct Cli { /// Per-stream greeting read deadline in seconds (1..=3600). #[arg(long)] hello_timeout: Option, + /// Authorization-path reply budget in seconds (1..=3600). + #[arg(long)] + authz_timeout: Option, + /// Join budget for established connections during shutdown (1..=3600). + #[arg(long)] + shutdown_timeout: Option, /// Directory HTTP(S) origin or legacy IP:port; the agent publishes its /// signed record and keeps it fresh. #[arg(long)] @@ -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 @@ -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 { diff --git a/crates/rds-agent/src/settings.rs b/crates/rds-agent/src/settings.rs index b45b692..1745831 100644 --- a/crates/rds-agent/src/settings.rs +++ b/crates/rds-agent/src/settings.rs @@ -199,6 +199,10 @@ pub struct LimitSettings { pub struct TimeoutSettings { pub handshake_secs: Option, pub hello_secs: Option, + /// Authorization-path reply budget (refusal and `HelloAck` writes). + pub authz_secs: Option, + /// Join budget for established connection tasks during shutdown. + pub shutdown_secs: Option, } /// Agent-level file schema; holds policy identities, never secret key @@ -261,6 +265,8 @@ pub struct AgentOverrides { pub max_streams: Option, pub handshake_timeout: Option, pub hello_timeout: Option, + pub authz_timeout: Option, + pub shutdown_timeout: Option, } /// The merged, typed configuration the binary consumes. @@ -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 } @@ -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( @@ -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, @@ -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}}}"#, @@ -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 diff --git a/crates/rds-bench/src/scenario.rs b/crates/rds-bench/src/scenario.rs index cbcb7e8..c7e3dd3 100644 --- a/crates/rds-bench/src/scenario.rs +++ b/crates/rds-bench/src/scenario.rs @@ -566,6 +566,7 @@ async fn resolve_connect(p: &Params) -> anyhow::Result { 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)); diff --git a/crates/rds-net/Cargo.toml b/crates/rds-net/Cargo.toml index 65c859e..2c31ad0 100644 --- a/crates/rds-net/Cargo.toml +++ b/crates/rds-net/Cargo.toml @@ -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 @@ -47,7 +49,6 @@ transport-noq = [ "dep:tokio-stream", "dep:tokio-util", "dep:blake3", - "dep:rand", ] [lints] diff --git a/crates/rds-net/src/announce.rs b/crates/rds-net/src/announce.rs index c03c94b..9850c26 100644 --- a/crates/rds-net/src/announce.rs +++ b/crates/rds-net/src/announce.rs @@ -9,6 +9,8 @@ use std::time::Duration; +use rand::RngExt; + use crate::{EndpointAddr, TransportAddr}; use rds_discovery::client::Client as DirectoryClient; use rds_discovery::publisher::{RecordDraft, RecordIssuer}; @@ -26,6 +28,45 @@ pub struct AnnounceConfig { pub services: Vec, /// Record TTL; the loop refreshes at `ttl/3`. pub ttl: Duration, + /// Publish-failure retry bounds. + pub retry: RetryPolicy, +} + +/// Bounded exponential backoff for retried directory publishes. +/// +/// Delays grow `base × 2^failures` up to `cap`; the applied sleep uses +/// equal jitter — uniform inside `[delay/2, delay]` — so the minimum +/// cadence remains provable (`base/2`) while fleet retries decorrelate. +/// The healthy poll cadence is untouched: backoff replaces one poll only +/// after a retryable failure. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct RetryPolicy { + /// Delay used by the first retry after a failure. + pub base: Duration, + /// Longest delay between retries. + pub cap: Duration, +} + +impl Default for RetryPolicy { + fn default() -> Self { + Self { + base: Duration::from_secs(1), + cap: Duration::from_secs(30), + } + } +} + +impl RetryPolicy { + /// Sleep duration for the `failures`-th consecutive retryable + /// failure (1-indexed): `min(cap, base × 2^(failures-1))` with equal + /// jitter inside `[delay/2, delay]`. + fn delay(&self, failures: u32) -> Duration { + let shift = failures.saturating_sub(1).min(16); + let delay = self.base.saturating_mul(1 << shift).min(self.cap); + let half_ms = (delay.as_millis() / 2).min(u64::MAX as u128) as u64; + Duration::from_millis(half_ms) + + Duration::from_millis(rand::rng().random_range(0..=half_ms)) + } } /// Advertised reachability: direct socket addrs + relay urls. @@ -68,16 +109,23 @@ pub fn announce(endpoint: Endpoint, config: AnnounceConfig) -> Result config.retry.cap { + return Err(DiscoveryError::Configuration( + "publish retry needs 0 < base <= cap".into(), + )); + } let task = tokio::spawn(async move { let AnnounceConfig { mut issuer, directory, services, ttl, + retry, } = config; let poll = (ttl / 6).clamp(Duration::from_millis(250), Duration::from_secs(1)); let mut last: Option = None; let mut renew = false; + let mut failures = 0u32; loop { let current = split_addrs(&endpoint.addr()); let draft = RecordDraft { @@ -107,6 +155,7 @@ pub fn announce(endpoint: Endpoint, config: AnnounceConfig) -> Result { last = Some(record); + failures = 0; tracing::debug!("endpoint record published"); } Err(DiscoveryError::Http { status: 410, .. }) => { @@ -119,7 +168,14 @@ pub fn announce(endpoint: Endpoint, config: AnnounceConfig) -> Result return Err(e), - Err(e) => tracing::warn!("directory publish failed: {e}"), + Err(e) => { + tracing::warn!("directory publish failed: {e}"); + failures = failures.saturating_add(1); + // Bounded backoff, not the healthy poll: a sustained + // outage must not publish-storm the directory. + tokio::time::sleep(retry.delay(failures)).await; + continue; + } } } tokio::time::sleep(poll).await; @@ -141,3 +197,46 @@ fn split_addrs(addr: &EndpointAddr) -> Advertised { } (addrs, relays) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn retry_delay_stays_inside_equal_jitter_bounds() { + let policy = RetryPolicy { + base: Duration::from_millis(100), + cap: Duration::from_secs(2), + }; + // Pre-jitter delays: 100ms, 200ms, 400ms, 800ms, 1600ms, then capped 2s. + for (failures, expected) in [ + (1, 100u64), + (2, 200), + (3, 400), + (4, 800), + (5, 1600), + (6, 2000), + (60, 2000), + ] { + let nominal = Duration::from_millis(expected); + for _ in 0..64 { + let delay = policy.delay(failures); + assert!( + delay >= nominal / 2 && delay <= nominal, + "failures={failures}: {delay:?} outside [nominal/2, nominal]" + ); + } + } + } + + #[test] + fn retry_delay_never_exceeds_cap_even_after_many_failures() { + let policy = RetryPolicy { + base: Duration::from_secs(1), + cap: Duration::from_secs(30), + }; + for failures in (0..=u32::MAX).step_by(1 << 20) { + assert!(policy.delay(failures) <= policy.cap); + } + } +} diff --git a/crates/rds-net/src/lib.rs b/crates/rds-net/src/lib.rs index c994350..293722c 100644 --- a/crates/rds-net/src/lib.rs +++ b/crates/rds-net/src/lib.rs @@ -64,7 +64,7 @@ pub use iroh::endpoint::{ ReadExactError, RecvStream, SendDatagramError, SendStream, VarInt, WriteError, }; -pub use announce::{Announce, AnnounceConfig, announce}; +pub use announce::{Announce, AnnounceConfig, RetryPolicy, announce}; pub use backends::iroh::{Ticket, parse_target, relay_url_of}; pub use config::{ConfigError, EndpointOverrides, EndpointSettings, RelayLimits, RelaySettings}; pub use identity::{KeyOwner, KeyStoreError, acquire_key, default_key_path, load_or_create_key}; diff --git a/crates/rds-net/tests/announce_e2e.rs b/crates/rds-net/tests/announce_e2e.rs index 5924132..ebd420d 100644 --- a/crates/rds-net/tests/announce_e2e.rs +++ b/crates/rds-net/tests/announce_e2e.rs @@ -7,7 +7,7 @@ use std::time::Duration; use rds_discovery::client::Client; use rds_discovery::service::{self, ServiceConfig}; use rds_discovery::{EndpointKey, MemoryStore, Service}; -use rds_net::{AnnounceConfig, EndpointConfig, announce, bind_endpoint}; +use rds_net::{AnnounceConfig, EndpointConfig, RetryPolicy, announce, bind_endpoint}; #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn announce_publishes_and_keeps_record_live() { @@ -38,6 +38,7 @@ async fn announce_publishes_and_keeps_record_live() { directory: client.clone(), services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -109,6 +110,7 @@ async fn announce_republishes_when_addrs_change() { directory: client.clone(), services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -332,6 +334,7 @@ async fn resolve_then_connect_by_bare_key() { directory: client.clone(), services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -369,3 +372,100 @@ async fn resolve_then_connect_by_bare_key() { let conn_b = accept.await.unwrap().expect("agent handshake"); assert_eq!(conn_b.remote_id(), cli_ep.id()); } + +/// Publish retries must back off, not storm: a directory that answers +/// every publish with 500 sees a bounded attempt count, and the loop +/// keeps retrying rather than exiting. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn announce_retry_is_bounded_backoff_under_outage() { + use std::sync::atomic::{AtomicU64, Ordering}; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let attempts = Arc::new(AtomicU64::new(0)); + tokio::spawn({ + let attempts = attempts.clone(); + async move { + while let Ok((mut sock, _)) = listener.accept().await { + attempts.fetch_add(1, Ordering::SeqCst); + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + let mut buf = [0u8; 8192]; + let _ = sock.read(&mut buf).await; + let _ = sock + .write_all(b"HTTP/1.1 500 Internal Server Error\r\ncontent-length: 0\r\n\r\n") + .await; + } + } + }); + + let key = rds_net::SecretKey::from_bytes(&[11u8; 32]); + let endpoint = bind_endpoint(EndpointConfig { + secret_key: Some(key.clone()), + ..Default::default() + }) + .await + .unwrap(); + let mut announce = announce( + endpoint, + AnnounceConfig { + issuer: rds_discovery::publisher::RecordIssuer::memory( + ed25519_dalek::SigningKey::from_bytes(&key.to_bytes()), + ), + directory: Client::new(addr.to_string().parse().unwrap()), + services: vec![Service::Ping], + ttl: Duration::from_secs(120), + retry: RetryPolicy { + base: Duration::from_millis(500), + cap: Duration::from_secs(2), + }, + }, + ) + .unwrap(); + + tokio::time::sleep(Duration::from_millis(3500)).await; + let count = attempts.load(Ordering::SeqCst); + // Minimum sleeps are 250ms, 500ms, 1s, 2s, 2s… — the loop cannot have + // attempted more often than that even at the jitter floor, while the + // old fixed ~1s cadence would land ~14 attempts in this window. + assert!( + (2..=7).contains(&count), + "bounded retries: {count} publish attempts in 3.5s" + ); + // A retrying failure must not look fatal: `wait` stays pending. + assert!( + tokio::time::timeout(Duration::from_millis(50), announce.wait()) + .await + .is_err() + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn announce_rejects_inverted_retry_policy() { + let client = Client::new("127.0.0.1:1".parse().unwrap()); + let key = rds_net::SecretKey::from_bytes(&[13u8; 32]); + let endpoint = bind_endpoint(EndpointConfig { + secret_key: Some(key.clone()), + ..Default::default() + }) + .await + .unwrap(); + let result = announce( + endpoint, + AnnounceConfig { + issuer: rds_discovery::publisher::RecordIssuer::memory( + ed25519_dalek::SigningKey::from_bytes(&key.to_bytes()), + ), + directory: client, + services: vec![Service::Ping], + ttl: Duration::from_secs(120), + retry: RetryPolicy { + base: Duration::from_secs(30), + cap: Duration::from_secs(1), + }, + }, + ); + assert!(matches!( + result, + Err(rds_discovery::DiscoveryError::Configuration(_)) + )); +} diff --git a/crates/rds-net/tests/publisher_lifecycle.rs b/crates/rds-net/tests/publisher_lifecycle.rs index 96ee6dc..2522a87 100644 --- a/crates/rds-net/tests/publisher_lifecycle.rs +++ b/crates/rds-net/tests/publisher_lifecycle.rs @@ -6,7 +6,7 @@ use rds_discovery::{ publisher::RecordIssuer, service::{self, ServiceConfig}, }; -use rds_net::{AnnounceConfig, EndpointConfig, SecretKey, announce, bind_endpoint}; +use rds_net::{AnnounceConfig, EndpointConfig, RetryPolicy, SecretKey, announce, bind_endpoint}; use std::{path::PathBuf, sync::Arc, time::Duration}; use tokio::net::{TcpListener, TcpStream}; @@ -85,6 +85,7 @@ async fn lost_success_reply_retries_identical_signed_bytes_through_real_director directory: client, services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -112,6 +113,7 @@ async fn missing_publisher_history_reaches_supervisor_without_network_publicatio directory: Client::new(listener.local_addr().unwrap()), services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -156,6 +158,7 @@ async fn stale_publisher_history_is_fatal_instead_of_guessing_server_revision() directory: client, services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); @@ -221,6 +224,7 @@ async fn expired_server_lease_allocates_one_durable_successor_then_retries_exact directory: client, services: vec![Service::Ping], ttl: Duration::from_secs(120), + retry: RetryPolicy::default(), }, ) .unwrap(); diff --git a/crates/rds-relay/tests/owned_e2e.rs b/crates/rds-relay/tests/owned_e2e.rs index 942cc36..834802b 100644 --- a/crates/rds-relay/tests/owned_e2e.rs +++ b/crates/rds-relay/tests/owned_e2e.rs @@ -203,6 +203,7 @@ async fn relay_forwards_handshake_and_datagrams() { directory: client.clone(), services: vec![rds_discovery::Service::Ping], ttl: std::time::Duration::from_secs(120), + retry: rds_net::RetryPolicy::default(), }, ) .unwrap(); diff --git a/docs/agent-configuration.md b/docs/agent-configuration.md index da75b63..6babf1d 100644 --- a/docs/agent-configuration.md +++ b/docs/agent-configuration.md @@ -35,7 +35,7 @@ identity creation or socket binding, alongside the endpoint preflight. | `peers` | `allow`: endpoint-id strings, at most 256 | | `authority` | `issuers`, `grant_ttl_secs` (1–86400), `tenant`, `policy_min_revision`, `directory`, `directory_ca`, `record_ttl_secs`, `record_state`, `registry`, `revocations` | | `limits` | `max_connections`, `max_streams`; positive 16-bit | -| `timeouts` | `handshake_secs`, `hello_secs`; each 1–3600 | +| `timeouts` | `handshake_secs`, `hello_secs`, `authz_secs`, `shutdown_secs`; each 1–3600 | `registry` holds `key` (required when present), `epoch`, `state` and `rotations`. `revocations` holds `key` (required when present), `epoch`, @@ -84,10 +84,14 @@ single SSH socket; the flag surface has no equivalent list flag. `timeouts.handshake_secs` bounds the inbound connection handshake; `timeouts.hello_secs` bounds the `StreamHello` read on every new stream. -Both default to 15 seconds and accept 1–3600. Flags `--handshake-timeout` -and `--hello-timeout` override file values. Per-service budgets (frame -streams, transfer deadlines, renewal windows) remain service-internal and -are not set from this file; global timeout classes are W2.6 work. +`timeouts.authz_secs` bounds the authorization-path replies (refusal and +final `HelloAck` writes); `timeouts.shutdown_secs` bounds the join wait +for established connection tasks when the agent stops. Defaults: 15s, +15s, 15s, 5s; each accepts 1–3600. Flags `--handshake-timeout`, +`--hello-timeout`, `--authz-timeout`, `--shutdown-timeout` override file +values. Per-service budgets (frame streams, transfer deadlines, renewal +windows) remain service-internal and are not set from this file; client +dial, idle and media/progress classes are still W2.6 open items. ## Example diff --git a/docs/remediation-progress.md b/docs/remediation-progress.md index b50339d..11db746 100644 --- a/docs/remediation-progress.md +++ b/docs/remediation-progress.md @@ -48,9 +48,9 @@ Neither increment closes these product gaps or any wave. | W2.1 | Implemented; endpoint + agent settings checked on Linux | Shared version-1 endpoint JSON, explicit file/flag precedence, typed backend/relay validation and preflight before identity creation are implemented. The version-1 agent JSON now carries role/service/peers/authority/limits/timeouts with the same precedence and preflight; `--role`/`--service`/`--no-service` select the gateable service set, disabled services are refused by name ahead of grant machinery, `Info` and directory announcements advertise exactly the served set, and handshake/hello deadlines come from `TimeoutPolicy`. See [agent configuration](agent-configuration.md). | | W2.2 | Partial; exact ALPN selection + sync/desktop session routing | Immutable per-protocol TLS offers prevent silent fallback and concurrent request interference. Managed single-file transfers use fresh control/uni routing IDs and negotiate version/limits in the `SyncTransferV2` session envelope before any filesystem operation. Desktop sessions mint random per-session IDs in `StreamHello::DesktopV2` and route frames through `UniHello::DesktopFrames { id }`, isolating stale streams and allowing concurrent sessions; the same display grant scope check covers both greetings. Service-wide capability negotiation remains open. | | W2.3 | Partial; destination-bound renewable grants, directional scopes, tenant/policy binding and per-path sync scopes | Grant v2 adds a strict signature domain, audience and stable session ID across positive lease revisions. Same-scope renewal preserves streams/revocation, retains one replay slot/watchdog and enforces wall/continuous expiry. Explicit managed renewal uses IPC v3 and a control-completion barrier. `SyncRead`/`SyncWrite` and `DesktopView`/`DesktopControl` are enforced before filesystem/input operations. Grant v3 adds the `tenant`/`policy_revision` claims and the `constraints.sync_paths` subtree scope, checked after `rel_path` normalization before any filesystem work; v2 payloads still verify with all claims absent, and agents pin the binding via `authority.tenant`/`policy_min_revision` (`--tenant`/`--policy-min-revision`), refusing unscoped or stale grants. Account-level scopes and automatic GDS issuer integration remain open. See [contract](grant-leases.md). | -| W2.4 | Partial; default connectivity manager | Agent local control is enabled by default; ordinary ticket/ping/info/SSH/forward/send/recv commands and keyless `rds session` reuse its endpoint. Same-UID IPC, pinned streams, cancellation and aggregate metrics are implemented. Agent/direct CLI/owned relay acquire exclusive ownership of a validated seed inode. Viewer manager APIs, coordinated installed-binary migration, native macOS and real multi-user/relay qualification remain open. See [contract](local-sessions.md) and [migration receipt](reports/rds-identity-migration-20260925.md). | +| W2.4 | Partial; default connectivity manager + managed desktop channel | Agent local control is enabled by default; ordinary ticket/ping/info/SSH/forward/send/recv/desktop commands and keyless `rds session` reuse its endpoint (local wire v5; `desktop --direct` bypasses). Same-UID IPC, pinned streams, cancellation and aggregate metrics are implemented. Agent/direct CLI/owned relay acquire exclusive ownership of a validated seed inode. Coordinated installed-binary migration, native macOS and real multi-user/relay qualification remain open. See [contract](local-sessions.md) and [migration receipt](reports/rds-identity-migration-20260925.md). | | W2.5 | Partial; transport and agent task ownership | Owned policy tasks terminate, including explicit shutdown after stopped protocol I/O; uni routing is bounded and acyclic. Agent and client forwarding groups own cancellation, normal joins and positive admission budgets. Client relay queues/peer leases and server admission/owned shutdown are bounded. Metric samplers use weak backend observations, release their gauges on drop and wake on closure independently of the sampling interval. Global RSS/FD bounds, per-service fairness and broader disk/media cancellation remain open. | -| W2.6 | Partial; client preludes bounded | One request deadline covers stream credit, writes, replies and Ping echo; canceled Authz closes its connection. Agent and owned relay handshake/shutdown budgets exist; canceling relay drain does not cancel cleanup. Agent local startup no longer waits indefinitely for an iroh relay, including disabled/unavailable relay mode. Global timeout classes, retry jitter, broader startup recovery and desktop/media deadlines remain open. | +| W2.6 | Partial; agent timeout classes complete, publish retry bounded | One request deadline covers stream credit, writes, replies and Ping echo; canceled Authz closes its connection. Agent policy now owns all four server-side classes — handshake, hello, authz reply and shutdown join — as `TimeoutPolicy` tunables (`timeouts.*_secs`, `--*-timeout`, 1..=3600). Directory publish retries use bounded exponential backoff with equal jitter (`RetryPolicy`, 1s→30s default) instead of the fixed ~1s poll cadence; fatal 4xx still fails closed. Agent and owned relay handshake/shutdown budgets exist; agent local startup no longer waits indefinitely for an iroh relay. Client dial/idle classes, transport-level retry policy reuse, desktop/media deadlines and broader startup recovery remain open. | | W3.1 | Partial; fair bounded candidate race | Canonical direct candidates alternate supported families under one eight-address cap, plus attached relay; attempts share a deadline and one authenticated winner. Independent relay bootstrap, progressive probing, remote scope/interface discovery and real topology qualification remain open. | | W3.2 | Partial; owned binary runtime checked on Linux | Both server binaries share strict backend/allow/key/limit/TLS config, persistent relay identity, local readiness and checked joined shutdown. Real processes forward inner authenticated traffic and retain identity/catalog across restart. Malformed datagrams are charged before parsing, and routing uses authenticated key-table lookup. Unexpected service-runner completion now initiates joined host shutdown with retained failure. Hung-task/recovery policy, global/reconnect/control budgets and platform/network qualification remain open. | | W3.3 | Done on Linux loopback; real-network open | Multiple owned-relay attachments (≤8 slots), slot-scoped synthetic routing, drain-before-death mask, immediate `PeerGone` invalidation, endpoint-scoped watchers; measured drain ~10ms / kill ~120-190ms recovery with 0 lost probes (`docs/reports/noq-relay-failover-20260926.md`). WAN/lossy migration timing remains unqualified. | @@ -2048,3 +2048,33 @@ stream-permit leak on peers that cannot serve desktop. Remaining W2.4: coordinated installed-binary migration and native/installed qualification — cross-platform, deferred. + +## 2026-09-27 — agent timeout classes + bounded publish retry (W2.6) + +`TimeoutPolicy` now owns all four server-side deadline classes: the +existing `handshake`/`hello` plus `authz` (authorization refusal and +final `HelloAck` writes, previously the `AUTHZ_REPLY_TIMEOUT` constant) +and `shutdown` (the connection-task join budget, previously +`SHUTDOWN_TIMEOUT`). Both are configurable via `timeouts.authz_secs` / +`timeouts.shutdown_secs` and `--authz-timeout` / `--shutdown-timeout`, +validated 1..=3600 like the existing classes, and defaulted +`TimeoutPolicy` keeps the prior constants so flag/file absence changes +nothing. + +The announce loop's publish-failure path no longer retries on the healthy +`min(1s, ttl/6)` poll: consecutive retryable failures sleep +`RetryPolicy::delay` — `base × 2^(n-1)` capped at `cap` (defaults 1s→30s) +with equal jitter inside `[delay/2, delay]`, so minimum cadence stays +provable (`base/2`) while fleet retries decorrelate. Fatal 4xx classes, +the 410 lease-renew path and issuer/disk failure surfacing are unchanged; +success or lease renewal resets the backoff. `AnnounceConfig` carries the +policy and refuses `base > cap`/`base == 0` at construction. + +Tests: `RetryPolicy::delay` bounds (per-failure jitter range, cap +saturation past overflow-scale failure counts), `announce` rejecting an +inverted policy, and an e2e where a directory answering every publish +with HTTP 500 sees a bounded attempt count over 3.5s while the task stays +alive — the fixed-cadence storm the criterion rules out. + +Remaining W2.6: client dial/idle classes, transport-level retry reuse +beyond announce, desktop/media deadlines, broader startup recovery.