diff --git a/crates/rds-agent/src/authz.rs b/crates/rds-agent/src/authz.rs index 4b9133f..57d3736 100644 --- a/crates/rds-agent/src/authz.rs +++ b/crates/rds-agent/src/authz.rs @@ -103,7 +103,11 @@ impl ConnAuthz { Self { audience, service_slots: Arc::new(tokio::sync::Semaphore::new( - streams.saturating_sub(usize::from(required)), + // One JoinSet lane is reserved whenever a control stream + // could need it: grant mode must always admit renewal, and + // any multi-stream connection must keep observability + // (Ping/Info) reachable past a saturated data plane. + streams.saturating_sub(usize::from(required || streams > 1)), )), state: Mutex::new(if required { State::Pending diff --git a/crates/rds-agent/src/lib.rs b/crates/rds-agent/src/lib.rs index 6770bbb..e45100e 100644 --- a/crates/rds-agent/src/lib.rs +++ b/crates/rds-agent/src/lib.rs @@ -698,21 +698,32 @@ async fn serve_stream( write_frame(&mut send, &HelloAck::Error { message: why }).await?; anyhow::bail!("stream outside grant scope"); } - let Some(_service_slot) = authz.try_service_slot() else { - rds_observe::request_refused(Reason::BudgetExhausted); - tokio::time::timeout( - policy.timeouts.hello, - write_frame( - &mut send, - &HelloAck::Error { - message: "service capacity reached; a slot is reserved for authorization" - .into(), - }, - ), - ) - .await??; - send.finish()?; - anyhow::bail!("service capacity reached"); + // 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. + let _service_slot = if matches!( + service_kind(&hello), + Some(ServiceKind::Tcp | ServiceKind::Desktop | ServiceKind::Sync | ServiceKind::Audio) + ) { + let Some(slot) = authz.try_service_slot() else { + rds_observe::request_refused(Reason::BudgetExhausted); + tokio::time::timeout( + policy.timeouts.hello, + write_frame( + &mut send, + &HelloAck::Error { + message: "service capacity reached; a lane is reserved for control traffic" + .into(), + }, + ), + ) + .await??; + send.finish()?; + anyhow::bail!("service capacity reached"); + }; + Some(slot) + } else { + None }; let span = info_span!("rds.stream", service = ?service_kind(&hello)); // Per-session frame route for `DesktopV2`; the shared `Desktop` diff --git a/crates/rds-agent/tests/grant_binding.rs b/crates/rds-agent/tests/grant_binding.rs index 6c086d8..3912a4c 100644 --- a/crates/rds-agent/tests/grant_binding.rs +++ b/crates/rds-agent/tests/grant_binding.rs @@ -218,8 +218,12 @@ async fn renewal_preserves_stream(backend: Backend) { .is_err() ); tokio::time::sleep(Duration::from_secs(5)).await; // original lease is now expired - // Its only service slot is held by TCP; the reserved slot still admitted renewal. - assert!(rds_client::ping(&conn, 3).await.is_err()); + // TCP holds the only data slot; the reserved lane still admitted + // renewal and keeps Ping reachable — the data plane cannot starve + // control traffic. + rds_client::ping(&conn, rand::random::()) + .await + .unwrap(); send.write_all(b"after").await.unwrap(); let mut last = [0; 5]; recv.read_exact(&mut last).await.unwrap(); diff --git a/crates/rds-agent/tests/lifecycle.rs b/crates/rds-agent/tests/lifecycle.rs index caf7175..3bbbda6 100644 --- a/crates/rds-agent/tests/lifecycle.rs +++ b/crates/rds-agent/tests/lifecycle.rs @@ -380,6 +380,59 @@ async fn process_fd_budget_above_usage_still_serves() { } } +/// The service pool keeps one JoinSet lane for control traffic: with two +/// streams, one long-lived data service holds the pool, a second data +/// service is refused, and Ping still answers — the data plane cannot +/// starve observability. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn saturated_data_plane_still_admits_control_streams() { + for backend in backends() { + let tcp = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = tcp.local_addr().unwrap().port(); + let (agent, clients) = fixture_target(backend, 1, 2, port).await; + let runner = start(&agent); + let conn = rds_cli::connect(&clients[0], agent.endpoint.addr()) + .await + .unwrap(); + // One valid service holds the single data pool slot. + let (mut forwarded, _reply) = rds_cli::open_tcp(&conn, "127.0.0.1", port).await.unwrap(); + forwarded.write_all(b"x").await.unwrap(); + let (mut remote, _) = tcp.accept().await.unwrap(); + use tokio::io::AsyncReadExt; + assert_eq!( + tokio::time::timeout(Duration::from_secs(2), remote.read_u8()) + .await + .unwrap() + .unwrap(), + b'x' + ); + // The data pool is full: a second service is refused politely. + let (mut raw_send, mut raw_recv) = conn.open_bi().await.unwrap(); + rds_net::write_frame( + &mut raw_send, + &rds_core::StreamHello::TcpConnect { + host: "127.0.0.1".into(), + port, + }, + ) + .await + .unwrap(); + let ack: rds_core::HelloAck = + tokio::time::timeout(Duration::from_secs(2), rds_net::read_frame(&mut raw_recv)) + .await + .expect("no ack for refused service") + .unwrap(); + assert!( + matches!(ack, rds_core::HelloAck::Error { ref message } if message.contains("service capacity")), + "{backend:?} over-capacity service was not refused: {ack:?}" + ); + // Control still arrives: Ping bypasses the data pool. + rds_cli::ping(&conn, rand::random::()).await.unwrap(); + state(&agent, 1, 1).await; + stop(&agent, &clients, runner).await; + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn metrics_track_real_admission_streams_and_do_not_own_agent_io() { for backend in backends() { diff --git a/crates/rds-agent/tests/local_manager.rs b/crates/rds-agent/tests/local_manager.rs index d459744..dee4488 100644 --- a/crates/rds-agent/tests/local_manager.rs +++ b/crates/rds-agent/tests/local_manager.rs @@ -992,9 +992,10 @@ async fn managed_desktop_reports_remote_refusal_without_leaking() { "expected clean remote refusal, got {result:?}" ); // The refused open must not park a stream permit: open_tcp still - // has its full budget. + // has its full data budget — every slot but the lane reserved + // for control traffic. let mut held = Vec::new(); - for _ in 0..64 { + for _ in 0..63 { held.push( client .open_tcp(session, _tcp.clone()) @@ -1002,6 +1003,13 @@ async fn managed_desktop_reports_remote_refusal_without_leaking() { .expect("stream slots leaked"), ); } + // One lane stays free of data services: a 64th TCP greeting is + // refused while that reserved slot keeps Ping reachable. + let overflow = client.open_tcp(session, _tcp.clone()).await.err(); + assert!( + matches!(overflow, Some(Error::Rejected(ErrorCode::Remote))), + "expected capacity refusal past the data budget, got {overflow:?}" + ); drop(held); server.close().await.unwrap(); tasks.abort_all(); diff --git a/crates/rds-sync/src/engine.rs b/crates/rds-sync/src/engine.rs index 6fcf541..3f28c82 100644 --- a/crates/rds-sync/src/engine.rs +++ b/crates/rds-sync/src/engine.rs @@ -39,6 +39,28 @@ const READ_STALL: Duration = Duration::from_secs(300); /// Call the `*_with_timeout` entry points to select a shorter or longer budget. pub const TRANSFER_TIMEOUT: Duration = Duration::from_secs(3600); +/// Process-wide bound on blocking filesystem work (W2.5). Waiting for a +/// permit happens on the async side, so a scan/store storm queues inside +/// this crate instead of filling Tokio's blocking pool ahead of identity, +/// revocation and announcement work. Generous enough that real transfers +/// never serialize on it. +const MAX_DISK_JOBS: usize = 32; +static DISK_JOBS: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(MAX_DISK_JOBS); + +/// Run `f` on the blocking pool once a disk-job permit frees. Cancelling +/// the future before a permit never reaches the blocking pool. +async fn disk_job(f: F) -> Result +where + F: FnOnce() -> R + Send + 'static, + R: Send + 'static, +{ + let _permit = DISK_JOBS + .acquire() + .await + .expect("disk-job semaphore never closes"); + tokio::task::spawn_blocking(f).await +} + /// Agent-side permissions, checked before any path or filesystem operation. /// Read means download from the agent; write means upload to the agent. #[derive(Debug, Clone, Default)] @@ -632,7 +654,7 @@ async fn serve_inner( // every handle again before any state or destination I/O. let preflight = { let (dir, rel) = (dir.clone(), rel.clone()); - tokio::task::spawn_blocking(move || { + disk_job(move || { match Directory::open_root(&dir, false).and_then(|root| root.read_path(&rel)) { Ok(_) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), @@ -694,11 +716,9 @@ async fn serve_inner( // same inode, even if the path is replaced after the offer. let source = { let rel = rel.clone(); - tokio::task::spawn_blocking(move || { - Directory::open_root(&dir, false)?.read_path(&rel) - }) - .await - .context("open source task")? + disk_job(move || Directory::open_root(&dir, false)?.read_path(&rel)) + .await + .context("open source task")? }; let source = match source { Ok(file) => Arc::new(file), @@ -777,7 +797,7 @@ async fn send_file_inner( } let path = path.to_path_buf(); let source = Arc::new( - tokio::task::spawn_blocking(move || -> anyhow::Result { + disk_job(move || -> anyhow::Result { use rustix::fs::{Mode, OFlags}; let file = File::from(rustix::fs::open( &path, @@ -954,7 +974,7 @@ where /// Manifest of a pinned file on the blocking pool — chunking + hashing a /// large file must not park an async worker. async fn manifest_from_file(file: Arc) -> anyhow::Result { - tokio::task::spawn_blocking(move || manifest_of_reader(&*file)) + disk_job(move || manifest_of_reader(&*file)) .await .context("manifest task")? .context("build manifest") @@ -990,14 +1010,22 @@ impl Drop for StoreCancellation { } impl JournalSink { - fn start(mut journal: Journal) -> (Self, tokio::sync::oneshot::Receiver<()>) { + async fn start(mut journal: Journal) -> (Self, tokio::sync::oneshot::Receiver<()>) { let (jobs, mut job_rx) = mpsc::channel::<(u32, Vec)>(FETCH_STREAMS * 4); let (finished, stopped) = tokio::sync::oneshot::channel(); let error = Arc::new(std::sync::Mutex::new(None)); let error_w = error.clone(); let canceled = Arc::new(AtomicBool::new(false)); let canceled_w = canceled.clone(); + // The permit waits on the async side, then lives inside the worker + // for the sink's whole lifetime — a transfer's store pump is one + // of the bounded disk jobs, not an uncounted blocking thread. + let permit = DISK_JOBS + .acquire() + .await + .expect("disk-job semaphore never closes"); let task = tokio::task::spawn_blocking(move || { + let _permit = permit; // Drop signals every exit, including panic, without depending on // a reader already waiting. Normal exit requires closing jobs. let _finished = finished; @@ -1084,7 +1112,7 @@ async fn receive( // work belongs on the blocking pool, not an async worker. let journal = { let (dir, rel, manifest) = (dir.to_path_buf(), rel.to_string(), manifest.clone()); - tokio::task::spawn_blocking(move || Journal::open(&dir, &rel, &manifest)) + disk_job(move || Journal::open(&dir, &rel, &manifest)) .await .context("journal open task")? }; @@ -1141,7 +1169,7 @@ async fn receive_chunks( // Chunk streams arrive on the already claimed transfer/service route; // neither other services nor different transfer IDs can consume them. let mut requested: std::collections::HashSet = journal.need().into_iter().collect(); - let (sink, stopped) = JournalSink::start(journal); + let (sink, stopped) = JournalSink::start(journal).await; // Remove only requested unique indices from the bounded set. Disk stores // remain asynchronous; after draining the sink, require complete verified // journal state before assembly or a success response. @@ -1236,7 +1264,7 @@ async fn receive_chunks( let fetched = journal.fetched(); // Assembly concatenates and rehashes every part — blocking pool. let dest = { - tokio::task::spawn_blocking(move || journal.assemble()) + disk_job(move || journal.assemble()) .await .context("assemble task")? .map_err(|e| anyhow::anyhow!("{e}"))? @@ -1301,7 +1329,7 @@ async fn push_chunks( // Positioned reads do not share a seek cursor across // streams and never reopen the peer-controlled path. let source = file.clone(); - buf = tokio::task::spawn_blocking(move || { + buf = disk_job(move || { source.read_exact_at(&mut buf, c.offset)?; Ok::<_, std::io::Error>(buf) }) diff --git a/crates/rds-sync/src/engine/tests.rs b/crates/rds-sync/src/engine/tests.rs index 136084b..151c688 100644 --- a/crates/rds-sync/src/engine/tests.rs +++ b/crates/rds-sync/src/engine/tests.rs @@ -32,7 +32,7 @@ fn queued_writer(cancel_finish: bool) { }); ready.await.unwrap(); let journal = Journal::open(&root.0, "data.bin", &manifest).unwrap(); - let (sink, stopped) = JournalSink::start(journal); + let (sink, stopped) = JournalSink::start(journal).await; sink.put(0, data).await.unwrap(); if cancel_finish { assert!( @@ -74,7 +74,7 @@ async fn normal_finish_drains_and_verifies_every_queued_store() { let data = vec![19u8; 4096]; let manifest = crate::manifest_of(&data); let journal = Journal::open(&root.0, "data.bin", &manifest).unwrap(); - let (sink, _stopped) = JournalSink::start(journal); + let (sink, _stopped) = JournalSink::start(journal).await; sink.put(0, data.clone()).await.unwrap(); let journal = sink.finish().await.unwrap(); assert!(journal.complete()); @@ -263,3 +263,27 @@ async fn session_open_refuses_wrong_version_with_clear_error() { .unwrap_err(); assert!(err.to_string().contains("version 99")); } + +/// W2.5: every `spawn_blocking` in this crate funnels through `disk_job`, +/// so a store/scan storm can never hold more than `MAX_DISK_JOBS` +/// blocking threads at once. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn disk_jobs_share_one_bounded_pool() { + use std::sync::atomic::AtomicUsize; + let active = Arc::new(AtomicUsize::new(0)); + let peak = Arc::new(AtomicUsize::new(0)); + let mut set = tokio::task::JoinSet::new(); + for _ in 0..(MAX_DISK_JOBS * 2) { + let (a, p) = (active.clone(), peak.clone()); + set.spawn(disk_job(move || { + let n = a.fetch_add(1, Ordering::SeqCst) + 1; + p.fetch_max(n, Ordering::SeqCst); + std::thread::sleep(Duration::from_millis(15)); + a.fetch_sub(1, Ordering::SeqCst); + })); + } + while set.join_next().await.is_some() {} + assert_eq!(active.load(Ordering::SeqCst), 0); + assert!(peak.load(Ordering::SeqCst) <= MAX_DISK_JOBS); + assert!(peak.load(Ordering::SeqCst) > 1, "jobs never overlapped"); +} diff --git a/docs/agent-configuration.md b/docs/agent-configuration.md index 3f3341f..4a9015b 100644 --- a/docs/agent-configuration.md +++ b/docs/agent-configuration.md @@ -80,6 +80,20 @@ services). `service.tcp_targets` permits extra `TcpConnect` destinations beyond the single SSH socket; the flag surface has no equivalent list flag. +## Stream budgets + +`limits.max_streams` bounds one connection's stream tasks. Long-lived +data services (`Tcp`, `Desktop`, `Sync`, `Audio`) hold a service slot for +the whole body; short control greetings (`Ping`, `Info`) and the +grant-mode `Authz`/`RenewAuthz` exchanges bypass that pool. Whenever a +connection can carry more than one stream, one lane is kept free of data +services — grant mode always keeps it for renewal — so a saturated data +plane cannot starve observability or authorization turnover. A data +service beyond capacity is refused with `HelloAck::Error` inside +`timeouts.hello_secs`. Sync filesystem work runs under one process-wide +bound on blocking disk jobs (32), so a transfer storm cannot fill the +blocking pool ahead of identity, announcement or other async work. + ## Timeout policy `timeouts.handshake_secs` bounds the inbound connection handshake; @@ -106,8 +120,8 @@ observed `rds_agent_process_fds` and `rds_agent_process_rss_bytes` keys appear in the metrics snapshot wherever the kernel reports them. Platforms without an observable quantity keep serving rather than gating on a guess. Flags `--max-fds` and `--max-rss-mb` override file -values. Per-service fairness and disk/media job budgets are still open -under W2.5. +values. Deeper per-service fairness and media/disk cancellation breadth +are still open under W2.5. ## Example diff --git a/docs/remediation-progress.md b/docs/remediation-progress.md index 914178c..75025a5 100644 --- a/docs/remediation-progress.md +++ b/docs/remediation-progress.md @@ -49,7 +49,7 @@ Neither increment closes these product gaps or any wave. | 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 + 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. Process-wide fd/RSS ceilings now gate connection admission (`limits.max_fds`/`max_rss_mb`, `--max-fds`/`--max-rss-mb`) with kernel-reported observability on Linux/macOS and honest ungated behavior elsewhere; per-service fairness, broader disk/media cancellation and storm-grade RSS/FD proof remain open. | +| 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. Process-wide fd/RSS ceilings now gate connection admission (`limits.max_fds`/`max_rss_mb`, `--max-fds`/`--max-rss-mb`) with kernel-reported observability on Linux/macOS and honest ungated behavior elsewhere. The stream budget now reserves a control lane universally (Ping/Info/Authz bypass the data-service pool, which refuses over-capacity greetings with a bounded `HelloAck::Error`), and all `rds-sync` blocking filesystem work funnels through one 32-permit disk-job bound including the long-lived journal store worker; deeper per-service fairness, broader disk/media cancellation and storm-grade RSS/FD proof 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. | @@ -2107,3 +2107,37 @@ above real usage admits and pings normally. Remaining W2.5: per-service fairness inside the stream budget, sync disk-job and desktop/media cancellation breadth, and storm-grade RSS/FD proof under adversarial slow peers. + +## 2026-09-27 — control-lane reservation and a process-wide disk bound (W2.5 partial) + +The per-connection stream budget now reserves one lane for control +traffic universally, not only in grant mode: `service_slots` is +`streams - 1` whenever a connection can carry more than one stream +(grant mode still reserves even a single-stream connection for +renewal), and the greetings that bypass the data pool are exactly the +short-lived control exchanges — `Ping`, `Info`, `Authz`, +`RenewAuthz`. `Tcp`, `Desktop`, `Sync` and `Audio` bodies hold a +service slot for their lifetime; a greeting that finds the data pool +full is refused with a `timeouts.hello_secs`-bounded `HelloAck::Error` +("service capacity reached; a lane is reserved for control traffic"). +A saturated data plane therefore cannot starve observability or +authorization turnover: at least one JoinSet lane always drains back +free for the next hello. `rds-sync` filesystem work is additionally +funneled through a single static semaphore of 32: every +`spawn_blocking` site — destination preflight, source and file opens, +manifest and journal writes, chunk reads — acquires a permit on the +async side before entering the blocking pool, and the long-lived +journal store worker holds its permit for its whole lifetime, so a +transfer storm cannot fill the blocking pool ahead of identity, +announcement or other async work. + +Tests: a saturated two-stream connection holds one live TCP forward, +refuses a second `TcpConnect` greeting with the capacity error and +still answers `Ping`; the renewal regression now asserts Ping succeeds +through the reserved lane rather than failing on a full pool; and a +64-job storm against `disk_job` never observes more than 32 +concurrently inside the blocking closure. + +Remaining W2.5: deeper per-service fairness (per-kind weights beyond +the single reserved lane), disk/media cancellation breadth, and +storm-grade RSS/FD proof under adversarial slow peers.