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
6 changes: 5 additions & 1 deletion crates/rds-agent/src/authz.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
41 changes: 26 additions & 15 deletions crates/rds-agent/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down
8 changes: 6 additions & 2 deletions crates/rds-agent/tests/grant_binding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<u64>())
.await
.unwrap();
send.write_all(b"after").await.unwrap();
let mut last = [0; 5];
recv.read_exact(&mut last).await.unwrap();
Expand Down
53 changes: 53 additions & 0 deletions crates/rds-agent/tests/lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<u64>()).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() {
Expand Down
12 changes: 10 additions & 2 deletions crates/rds-agent/tests/local_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -992,16 +992,24 @@ 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())
.await
.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();
Expand Down
54 changes: 41 additions & 13 deletions crates/rds-sync/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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, R>(f: F) -> Result<R, tokio::task::JoinError>
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)]
Expand Down Expand Up @@ -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(()),
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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<File> {
disk_job(move || -> anyhow::Result<File> {
use rustix::fs::{Mode, OFlags};
let file = File::from(rustix::fs::open(
&path,
Expand Down Expand Up @@ -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<File>) -> anyhow::Result<Manifest> {
tokio::task::spawn_blocking(move || manifest_of_reader(&*file))
disk_job(move || manifest_of_reader(&*file))
.await
.context("manifest task")?
.context("build manifest")
Expand Down Expand Up @@ -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<u8>)>(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;
Expand Down Expand Up @@ -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")?
};
Expand Down Expand Up @@ -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<u32> = 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.
Expand Down Expand Up @@ -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}"))?
Expand Down Expand Up @@ -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)
})
Expand Down
28 changes: 26 additions & 2 deletions crates/rds-sync/src/engine/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -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");
}
18 changes: 16 additions & 2 deletions docs/agent-configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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

Expand Down
Loading
Loading