From 1d4926bcea19674f4ad15271f54cf8c36680014b Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 18:57:07 +0500 Subject: [PATCH 01/10] fix(desktop,client): close lifecycle races and enforce deadlines MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - mailbox: drain an item queued between a failed pop and the last sender's drop instead of closing over it (close/pop race). - local down-pump: track encoded/event tap liveness independently so a dead tap cannot report Finished while its sibling is still live. - session: file SetBitrate into a requested slot the pacing task drains and steers into the controller — a viewer target survives past one 250ms tick instead of being overwritten by the next adapt step. - session: bound input injection by the frame deadline and drop a wedged worker instead of stalling the whole control plane on a hung X server. - client: bound spawn_blocking decode by the frame-stream deadline; a decoder that never returns re-baselines on a fresh chain + IDR instead of stalling uni demux for every service on the connection. - capture/x11: refuse root depths without a 32bpp pixmap format at construction and fail honestly on a GetImage byte-count mismatch. - input/x11: drop the dead static-sink inject() footgun; scroll fan-out reuses one pointer_on_screen check instead of per-click re-queries. - codec/openh264: IDR detection now requires a real Annex-B start code. - run_desktop_client: headless builds print that decoding is absent instead of looking inert. --- crates/rds-client/src/local/desktop.rs | 15 +++-- crates/rds-desktop/src/capture/x11.rs | 22 +++++++ crates/rds-desktop/src/client.rs | 33 ++++++++-- crates/rds-desktop/src/codec/openh264.rs | 6 +- crates/rds-desktop/src/input/x11.rs | 20 ++---- crates/rds-desktop/src/mailbox.rs | 21 ++++++- crates/rds-desktop/src/session.rs | 77 +++++++++++++++++++++--- 7 files changed, 156 insertions(+), 38 deletions(-) diff --git a/crates/rds-client/src/local/desktop.rs b/crates/rds-client/src/local/desktop.rs index 41007ca..6f6097f 100644 --- a/crates/rds-client/src/local/desktop.rs +++ b/crates/rds-client/src/local/desktop.rs @@ -81,10 +81,13 @@ pub(super) async fn serve( let down = async { let encoded = session.encoded.as_mut().unwrap(); let events = &mut session.events; - let mut open = 2u8; - while open > 0 { + // Each tap's liveness is tracked separately: `recv` answers `None` + // immediately and repeatedly once its senders are gone, so a shared + // counter would end the pump while the surviving tap is still live. + let (mut enc_open, mut ev_open) = (true, true); + while enc_open || ev_open { tokio::select! { - frame = encoded.recv() => match frame { + frame = encoded.recv(), if enc_open => match frame { Some(f) => { if f.payload.len() > MAX_DESKTOP_PAYLOAD { return Err(invalid()); @@ -92,11 +95,11 @@ pub(super) async fn serve( write_frame(&mut writer, &DesktopDown::Frame { header: f.header }).await?; write_payload(&mut writer, &f.payload).await?; } - None => open -= 1, + None => enc_open = false, }, - event = events.recv() => match event { + event = events.recv(), if ev_open => match event { Some(e) => write_frame(&mut writer, &DesktopDown::Event(e)).await?, - None => open -= 1, + None => ev_open = false, }, } } diff --git a/crates/rds-desktop/src/capture/x11.rs b/crates/rds-desktop/src/capture/x11.rs index 8f6bd0d..123ce6f 100644 --- a/crates/rds-desktop/src/capture/x11.rs +++ b/crates/rds-desktop/src/capture/x11.rs @@ -69,6 +69,19 @@ impl X11Capturer { .ok_or_else(|| DesktopError::Capture("X11 display does not exist".into()))?; let root = display.root; let (width, height) = (display.width_in_pixels, display.height_in_pixels); + // Frames are decoded as packed 32bpp pixels; a server whose root + // depth has no 32bpp pixmap format cannot produce them — refuse + // honestly rather than panic or emit corrupt frames. + let ok = setup + .pixmap_formats + .iter() + .any(|f| f.depth == display.root_depth && f.bits_per_pixel == 32); + if !ok { + return Err(DesktopError::Capture(format!( + "X11 screen {idx} root depth {} has no 32bpp format", + display.root_depth + ))); + } let _ = default_screen; let shm = Self::try_shm(&conn, width, height); let damage = Self::try_damage(&conn, root); @@ -180,6 +193,15 @@ impl Capturer for X11Capturer { .map_err(|e| DesktopError::Capture(e.to_string()))? .reply() .map_err(|e| DesktopError::Capture(e.to_string()))?; + let expected = usize::from(self.width) * usize::from(self.height) * 4; + if reply.data.len() != expected { + return Err(DesktopError::Capture(format!( + "X11 GetImage returned {} bytes for {}x{} (expected {expected})", + reply.data.len(), + self.width, + self.height + ))); + } Ok(RawFrame { width: u32::from(self.width), height: u32::from(self.height), diff --git a/crates/rds-desktop/src/client.rs b/crates/rds-desktop/src/client.rs index 5cf6aed..1611e7b 100644 --- a/crates/rds-desktop/src/client.rs +++ b/crates/rds-desktop/src/client.rs @@ -617,12 +617,29 @@ async fn receive_frames(mut uni: rds_net::UniStreams, ctx: ReceiveContext) { // No control/queue/connection handles escape into native work. // An already running call may finish after cancellation, but // it cannot publish and holds both permits until it returns. - let decoded = tokio::task::spawn_blocking(move || { - let (_slot, _budget) = (slot, budget); - let result = delivery.decode(&header, body); - (delivery, result) - }).await; - let Ok((state, outcome)) = decoded else { break; }; + let decoded = tokio::time::timeout( + FRAME_STREAM_TIMEOUT, + tokio::task::spawn_blocking(move || { + let (_slot, _budget) = (slot, budget); + let result = delivery.decode(&header, body); + (delivery, result) + }), + ) + .await; + let (state, outcome) = match decoded { + Ok(Ok(pair)) => pair, + // A decoder that never returns would stall the demux + // for every service on this connection; the orphan + // still holds both permits until it finishes. The + // moved-out chain is gone — re-baseline on a fresh + // one instead of blocking uni routing forever. + Err(_) => { + delivery = Delivery::new(); + delivery.request_idr(&ctx.ctrl, ctx.clock.now_ms()); + continue; + } + Ok(Err(_)) => break, + }; delivery = state; match outcome { #[cfg(feature = "x11")] @@ -678,6 +695,10 @@ pub async fn run_desktop_client( ) -> Result<(), DesktopError> { let mut session = DesktopSession::connect(&conn, display, max_fps, Codec::H264).await?; println!("desktop caps: {:?}", session.caps()); + #[cfg(not(feature = "x11"))] + println!( + "note: headless build has no decoder — frames stay encoded (use the encoded relay tap)" + ); let mut count = 0u64; let start = std::time::Instant::now(); while let Some(frame) = session.frames.recv().await { diff --git a/crates/rds-desktop/src/codec/openh264.rs b/crates/rds-desktop/src/codec/openh264.rs index d22b006..521abf9 100644 --- a/crates/rds-desktop/src/codec/openh264.rs +++ b/crates/rds-desktop/src/codec/openh264.rs @@ -148,11 +148,13 @@ fn bitstream_has_idr(stream: &EncodedBitStream<'_>) -> bool { let Some(nal) = layer.nal_unit(n) else { continue; }; - // First non-zero byte is the start code's trailing 0x01; - // the byte right after it is the NAL header. + // Annex-B start code: ≥2 leading zeros then 0x01; the byte + // right after it is the NAL header. A NAL without that + // prefix is not a start we can interpret. let Some(hdr) = nal .iter() .position(|&b| b != 0) + .filter(|&i| i >= 2 && nal[i] == 1) .and_then(|i| nal.get(i + 1)) else { continue; diff --git a/crates/rds-desktop/src/input/x11.rs b/crates/rds-desktop/src/input/x11.rs index bd9132a..2965a6f 100644 --- a/crates/rds-desktop/src/input/x11.rs +++ b/crates/rds-desktop/src/input/x11.rs @@ -187,11 +187,14 @@ impl XtestInput { self.scroll_x -= f64::from(x); self.scroll_y -= f64::from(y); // Protocol convention: positive y is up, positive x is left. + // The single `pointer_on_screen` above covers the whole fan-out — + // raw fake calls skip `button`'s per-click re-query (≈128 round + // trips on a full scroll burst). for (steps, negative, positive) in [(y, 5, 4), (x, 7, 6)] { let button = if steps > 0 { positive } else { negative }; for _ in 0..steps.unsigned_abs() { - self.button(button, true)?; - self.button(button, false)?; + self.fake(BUTTON_PRESS, button, 0, 0)?; + self.fake(BUTTON_RELEASE, button, 0, 0)?; } } Ok(()) @@ -286,16 +289,3 @@ impl Drop for XtestInput { } } } - -/// Convenience one-shot injection used by the session control loop. -pub fn inject(event: &InputEvent) -> Result<(), DesktopError> { - use std::sync::Mutex; - static SINK: Mutex> = Mutex::new(None); - let mut guard = SINK - .lock() - .map_err(|_| DesktopError::Input("input sink poisoned".into()))?; - if guard.is_none() { - *guard = Some(XtestInput::new()?); - } - guard.as_mut().unwrap().inject(event) -} diff --git a/crates/rds-desktop/src/mailbox.rs b/crates/rds-desktop/src/mailbox.rs index da8cefc..20cc275 100644 --- a/crates/rds-desktop/src/mailbox.rs +++ b/crates/rds-desktop/src/mailbox.rs @@ -109,9 +109,10 @@ impl Receiver { return Some(item); } } - // All senders gone → the queue reads as closed. + // All senders gone → the queue reads as closed; drain any + // item that raced in between the failed pop and the check. if self.0.senders.load(Ordering::Acquire) == 0 { - return None; + return self.0.queue.lock().unwrap().pop_front(); } self.0.notify.notified().await; } @@ -184,6 +185,22 @@ mod tests { } } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn close_drains_item_racing_sender_drop() { + // The session-end tail: the last producer pushes its final item + // and exits. A receiver that had just observed an empty queue + // must still deliver that item, not close over it. + for _ in 0..200 { + let (tx, mut rx) = channel::(4); + let producer = std::thread::spawn(move || { + tx.send(9); + }); + assert_eq!(rx.recv().await, Some(9)); + producer.join().unwrap(); + assert_eq!(rx.recv().await, None); + } + } + #[tokio::test] async fn recv_awaits_produce() { let (tx, mut rx) = channel(4); diff --git a/crates/rds-desktop/src/session.rs b/crates/rds-desktop/src/session.rs index 3b7cb99..bf62bfd 100644 --- a/crates/rds-desktop/src/session.rs +++ b/crates/rds-desktop/src/session.rs @@ -65,12 +65,14 @@ pub struct Produced { } /// Live knobs the session applies while a producer runs: the pacing -/// controller writes `bitrate`, the control stream flips `idr`, and the -/// producer reports `deadline_misses` it observed. +/// controller writes `bitrate`, the control stream flips `idr` and files +/// `requested` bitrate targets (0 = none pending) for the controller to +/// drain, and the producer reports `deadline_misses` it observed. pub struct ProducerControls { pub bitrate: Arc, pub idr: Arc, pub deadline_misses: Arc, + pub requested: Arc, } impl ProducerControls { @@ -79,6 +81,7 @@ impl ProducerControls { bitrate: Arc::new(AtomicU64::new(initial_bps)), idr: Arc::new(AtomicBool::new(true)), deadline_misses: Arc::new(AtomicU64::new(0)), + requested: Arc::new(AtomicU64::new(0)), } } } @@ -161,6 +164,12 @@ impl BitrateController { self.current } + /// Steer to a viewer-requested target, clamped to the controller's + /// own floor and ceiling; adaptation resumes from there. + pub fn steer(&mut self, bps: u64) { + self.current = bps.clamp(self.floor, self.ceiling); + } + /// One pacing step. `path` is the selected path's counters (delta'd /// against the previous call); `deadline_misses` is the count of /// produce calls that missed their cadence slot since the last step. @@ -249,12 +258,17 @@ pub async fn serve_desktop_with( let bitrate = Arc::clone(&controls.bitrate); let idr = Arc::clone(&controls.idr); let misses = Arc::clone(&controls.deadline_misses); + let requested = Arc::clone(&controls.requested); let mut producer = config.producer; capture.spawn_blocking(move || { let producer_controls = ProducerControls { bitrate, idr, deadline_misses: misses, + // Share the session's target slot — a producer observing + // `requested` must see what the control arm filed, not a + // private always-empty copy. + requested, }; let mut source = match producer.take() { Some(p) => p, @@ -292,6 +306,7 @@ pub async fn serve_desktop_with( { let conn = conn.clone(); let bitrate = Arc::clone(&controls.bitrate); + let requested = Arc::clone(&controls.requested); let misses = Arc::clone(&controls.deadline_misses); let mut controller = BitrateController::new(4_000_000, ceiling); workers.spawn(async move { @@ -299,6 +314,12 @@ pub async fn serve_desktop_with( tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); loop { tick.tick().await; + // A viewer-filed target survives past one tick: the + // controller is steered to it, then adaptation resumes. + let req = requested.swap(0, Ordering::Relaxed); + if req != 0 { + controller.steer(req); + } let missed = misses.swap(0, Ordering::Relaxed); let bps = controller.step(conn.current_path_stats(), missed); bitrate.store(bps.min(u64::from(u32::MAX)), Ordering::Relaxed); @@ -389,13 +410,27 @@ pub async fn serve_desktop_with( ); continue; } - let input = input.get_or_insert_with(|| { + let mut worker = input.take().unwrap_or_else(|| { super::input::worker::InputWorker::new(input_sink.take()) }); let seq = ev.seq; - if let Err(e) = input.inject(ev).await { - tracing::warn!("input injection failed: {e}"); - continue; + match tokio::time::timeout(FRAME_SEND_TIMEOUT, worker.inject(ev)).await { + Ok(Ok(())) => input = Some(worker), + Ok(Err(e)) => { + input = Some(worker); + tracing::warn!("input injection failed: {e}"); + continue; + } + Err(_) => { + // The platform input call never returned (a + // wedged X server). Drop the worker — its + // running syscall may still finish, per its + // contract — so the next event probes a fresh + // sink, and count this event unacked rather + // than stalling the whole control plane. + tracing::warn!("input injection timed out; dropping wedged worker"); + continue; + } } if acks { let ack = DesktopEvent::InputAck { @@ -419,7 +454,7 @@ pub async fn serve_desktop_with( } Ok(DesktopControl::SetBitrate(bps)) => { let bps = u64::from(bps.max(50_000)).min(ceiling); - controls.bitrate.store(bps, Ordering::Relaxed); + controls.requested.store(bps, Ordering::Relaxed); } Ok(DesktopControl::Heartbeat { seq, ts_ms }) => { if !matches!( @@ -1248,4 +1283,32 @@ mod tests { let bps = c.step(Some(path(1_000_000, 500_000, 20, 0)), 0); assert_eq!(bps, 4_000_000, "first sample establishes baseline only"); } + + #[test] + fn steer_sets_target_inside_bounds() { + let mut c = BitrateController::new(4_000_000, 8_000_000); + c.steer(2_500_000); + assert_eq!(c.current(), 2_500_000); + // Out-of-range requests clamp to the controller's own contract. + c.steer(50_000); + assert_eq!(c.current(), 100_000); + c.steer(50_000_000); + assert_eq!(c.current(), 8_000_000); + } + + #[test] + fn steer_survives_adaptation_ticks() { + // The SetBitrate semantics the viewer sees: a filed target is not + // overwritten by the next pacing step — adaptation resumes *from* + // the steered point. + let mut c = BitrateController::new(4_000_000, 8_000_000); + c.step(Some(path(1000, 0, 20, 0)), 0); // prime baseline + c.steer(1_000_000); + let bps = c.step(Some(path(2000, 0, 20, 0)), 0); + // Clean window: one recovery step from 1Mbps, not a snap back to 4. + assert!( + bps > 1_000_000 && bps < 2_000_000, + "adaptation must resume from the steered target, got {bps}" + ); + } } From 745a80e3d36d3d3625acfd481aa557cce90a4de2 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:08:30 +0500 Subject: [PATCH 02/10] fix(sync): close control-frame races and bound disk jobs per store MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - push_chunks: the receiver's Done can legitimately arrive while chunk streams still finish; consume it as an early ack and validate the root through one shared recv_done path instead of reporting a spurious "unexpected control frame" abort. - Transfer::serve: a failed control.finish must not mask the transfer error that triggered it — warn and return the original cause. - session_open: a HelloAck with a wrong transfer ID now reports the ID mismatch rather than a phantom version failure. - JournalSink: the disk-job permit rides with each queued chunk instead of the sink's lifetime — N concurrent receives no longer pin N of the process's 32 blocking slots while waiting on the network. - Offer/Need/Done waits now use a 900s phase bound: they gate on the peer's heavy local work (manifest hashing, journal walk, assemble) rather than wire speed; the session deadline still bounds a dead peer. --- crates/rds-sync/src/engine.rs | 180 +++++++++++++++++++++++++--------- 1 file changed, 131 insertions(+), 49 deletions(-) diff --git a/crates/rds-sync/src/engine.rs b/crates/rds-sync/src/engine.rs index 3f28c82..fd1d5bd 100644 --- a/crates/rds-sync/src/engine.rs +++ b/crates/rds-sync/src/engine.rs @@ -35,6 +35,12 @@ use crate::{confined::Directory, journal::Journal}; /// reads gate on the peer's disk work (manifest scans, journal /// rescans); a dead connection ends them regardless. const READ_STALL: Duration = Duration::from_secs(300); +/// Stall bound for frames that gate on the peer's own heavy local work: +/// manifest hashing before `Offer`, journal walk before `Need`, +/// assembly+verify before `Done`. Those legitimately outgrow the +/// per-frame stall on large transfers; the session deadline still +/// bounds a peer that never finishes. +const PHASE_STALL: Duration = Duration::from_secs(900); /// Default absolute transfer budget, including local scans and all protocol I/O. /// Call the `*_with_timeout` entry points to select a shorter or longer budget. pub const TRANSFER_TIMEOUT: Duration = Duration::from_secs(3600); @@ -182,7 +188,15 @@ impl Wire { where R: tokio::io::AsyncRead + Unpin, { - let msg: SyncMsg = read_timed(stream).await?; + self.recv_within(stream, READ_STALL).await + } + + /// `recv` with a caller-chosen stall bound; see [`PHASE_STALL`]. + async fn recv_within(&self, stream: &mut R, stall: Duration) -> anyhow::Result + where + R: tokio::io::AsyncRead + Unpin, + { + let msg: SyncMsg = read_timed_within(stream, stall).await?; let Some(transfer_id) = self.transfer_id else { return Ok(msg); }; @@ -297,14 +311,19 @@ where Ok(SessionLimits::negotiate(SessionLimits::LOCAL, limits)?) } SyncMsg::Session { + transfer_id: got, msg: SessionMsg::HelloAck { version, .. }, + } if got == transfer_id => { + bail!("peer sync session version {version} unsupported (want {SESSION_VERSION})") + } + SyncMsg::Session { + msg: SessionMsg::HelloAck { .. }, .. - } => bail!("peer sync session version {version} unsupported (want {SESSION_VERSION})"), + } => bail!("sync session transfer ID mismatch"), SyncMsg::Session { msg: SessionMsg::Refuse { reason }, .. } => bail!("sync session refused: {reason}"), - SyncMsg::Session { .. } => bail!("sync session transfer ID mismatch"), other => bail!("expected sync session HelloAck, got {other:?}"), } } @@ -419,9 +438,12 @@ impl Transfer { result = serve_inner(conn, &mut control.send, &mut control.recv, dir, access, self.0, wire) => result, }; if let Err(error) = result { - // Preserve an explicitly written Refuse frame. Cancellation - // and deadlines still drop/reset the unfinished control pair. - control.finish(false).await?; + // Preserve an explicitly written Refuse frame. A finish + // failure here must not mask the transfer error that + // brought us here — report it and keep the real cause. + if let Err(finish) = control.finish(false).await { + tracing::warn!("sync control finish after error failed: {finish:#}"); + } return Err(error); } if self.0 != rds_core::UniHello::Sync { @@ -626,7 +648,9 @@ async fn serve_inner( route: rds_core::UniHello, wire: Wire, ) -> anyhow::Result<()> { - let first = wire.recv(recv).await?; + // An `Offer` first frame gates on the caller's manifest build — + // hashing a large source legitimately outlives the per-frame stall. + let first = wire.recv_within(recv, PHASE_STALL).await?; match first { SyncMsg::Offer { rel_path, @@ -736,19 +760,16 @@ async fn serve_inner( "sync pull serving" ); send_manifest(wire, send, &rel.to_string_lossy(), &manifest).await?; - let SyncMsg::Need { bits } = wire.recv(recv).await? else { + // `Need` gates on the receiver's journal open — a walk over + // every stored part — not on wire speed. + let SyncMsg::Need { bits } = wire.recv_within(recv, PHASE_STALL).await? else { bail!("expected Need"); }; let indices = bits_to_indices(&bits, manifest.chunks.len())?; - push_chunks(&conn, source, &manifest, &indices, route, wire, recv).await?; - match wire.recv(recv).await? { - SyncMsg::Done { root } if root == manifest.root => { - tracing::info!(sent = indices.len(), "sync pull complete"); - Ok(()) - } - SyncMsg::Cancel { reason } => bail!("receiver canceled transfer: {reason}"), - other => bail!("expected Done, got {other:?}"), - } + let early = push_chunks(&conn, source, &manifest, &indices, route, wire, recv).await?; + recv_done(early, wire, recv, manifest.root).await?; + tracing::info!(sent = indices.len(), "sync pull complete"); + Ok(()) } SyncMsg::Cancel { reason } => bail!("caller canceled before transfer started: {reason}"), other => bail!("unexpected first sync message {other:?}"), @@ -830,19 +851,16 @@ async fn send_file_inner( "sync push start" ); send_manifest(wire, send, &rel, &manifest).await?; - let indices = match wire.recv(recv).await? { + // `Need` gates on the receiver's journal open — a walk over every + // stored part — not on wire speed. + let indices = match wire.recv_within(recv, PHASE_STALL).await? { SyncMsg::Need { bits } => bits_to_indices(&bits, manifest.chunks.len())?, SyncMsg::Refuse { reason } => bail!("offer refused: {reason}"), SyncMsg::Cancel { reason } => bail!("receiver canceled before chunks: {reason}"), other => bail!("expected Need, got {other:?}"), }; - push_chunks(conn, source, &manifest, &indices, route, wire, recv).await?; - match wire.recv(recv).await? { - SyncMsg::Done { root } if root == manifest.root => {} - SyncMsg::Refuse { reason } => bail!("receiver refused: {reason}"), - SyncMsg::Cancel { reason } => bail!("receiver canceled transfer: {reason}"), - other => bail!("expected Done, got {other:?}"), - } + let early = push_chunks(conn, source, &manifest, &indices, route, wire, recv).await?; + recv_done(early, wire, recv, manifest.root).await?; let stats = Stats { fetched: indices.len() as u64, total: manifest.chunks.len() as u64, @@ -898,7 +916,8 @@ async fn recv_file_inner( }, ) .await?; - let (rel, size, root, chunk_count) = match wire.recv(recv).await? { + // The answer gates on the server's manifest build, not wire speed. + let (rel, size, root, chunk_count) = match wire.recv_within(recv, PHASE_STALL).await? { SyncMsg::Offer { rel_path, size, @@ -940,7 +959,16 @@ where S: tokio::io::AsyncRead + Unpin, T: serde::de::DeserializeOwned, { - match tokio::time::timeout(READ_STALL, read_frame(stream)).await { + read_timed_within(stream, READ_STALL).await +} + +/// `read_timed` with a caller-chosen stall bound. +async fn read_timed_within(stream: &mut S, stall: Duration) -> anyhow::Result +where + S: tokio::io::AsyncRead + Unpin, + T: serde::de::DeserializeOwned, +{ + match tokio::time::timeout(stall, read_frame(stream)).await { Ok(r) => r.map_err(Into::into), Err(_) => bail!("peer stalled mid-transfer"), } @@ -988,7 +1016,7 @@ struct JournalSink { // Fields drop in declaration order: publish cancellation before closing // the sender wakes the blocking receiver with its remaining queued data. cancel: StoreCancellation, - jobs: mpsc::Sender<(u32, Vec)>, + jobs: mpsc::Sender<(u32, Vec, tokio::sync::SemaphorePermit<'static>)>, /// First store failure, for error reporting across the task split. error: Arc>>, task: tokio::task::JoinHandle>, @@ -1011,25 +1039,24 @@ impl Drop for StoreCancellation { impl JournalSink { async fn start(mut journal: Journal) -> (Self, tokio::sync::oneshot::Receiver<()>) { - let (jobs, mut job_rx) = mpsc::channel::<(u32, Vec)>(FETCH_STREAMS * 4); + let (jobs, mut job_rx) = + mpsc::channel::<(u32, Vec, tokio::sync::SemaphorePermit<'static>)>( + 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"); + // The disk-job permit rides with each queued chunk (acquired in + // `put`), so a sink parked waiting for network data holds no slot. + // Holding one per sink lifetime let N concurrent receives starve + // every one-shot disk job in the process. 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; - while let Some((index, data)) = job_rx.blocking_recv() { + while let Some((index, data, _permit)) = job_rx.blocking_recv() { if canceled_w.load(Ordering::Acquire) { break; } @@ -1061,7 +1088,13 @@ impl JournalSink { /// Queue one bounded chunk for verification and storage; backpressures when the /// disk side falls behind. async fn put(&self, index: u32, data: Vec) -> anyhow::Result<()> { - self.jobs.send((index, data)).await.map_err(|_| { + // Permit per queued store, taken on the async side like every + // other disk job — it releases when the worker finishes the item. + let permit = DISK_JOBS + .acquire() + .await + .expect("disk-job semaphore never closes"); + self.jobs.send((index, data, permit)).await.map_err(|_| { let why = self .error .lock() @@ -1280,11 +1313,36 @@ async fn receive_chunks( )) } +/// Final `Done` read after a chunk push. `early` is the root the push +/// watcher consumed while chunk tasks were still draining — validating +/// it here keeps the root check in one place for both arrival orders. +/// The receiver's `Done` gates on its store drain + assemble, so the +/// read uses the phase bound. +async fn recv_done( + early: Option, + wire: Wire, + recv: &mut RecvStream, + expected: crate::ChunkHash, +) -> anyhow::Result<()> { + let msg = match early { + Some(root) => SyncMsg::Done { root }, + None => wire.recv_within(recv, PHASE_STALL).await?, + }; + match msg { + SyncMsg::Done { root } if root == expected => Ok(()), + SyncMsg::Done { .. } => bail!("receiver acknowledged a different manifest root"), + SyncMsg::Refuse { reason } => bail!("receiver refused: {reason}"), + SyncMsg::Cancel { reason } => bail!("receiver canceled transfer: {reason}"), + other => bail!("expected Done, got {other:?}"), + } +} + /// Holder half: open the negotiated number of uni streams (v1 uses /// [`FETCH_STREAMS`]), each walking an interleaved share of `indices` in /// `CHUNKSET_BATCH` batches. v2 watches the control stream so a typed /// `Cancel` stops chunk production immediately instead of surfacing as -/// write errors on half-closed streams. +/// write errors on half-closed streams. Returns the receiver's `Done` +/// root when it was consumed while chunk tasks were still draining. async fn push_chunks( conn: &Connection, file: Arc, @@ -1293,7 +1351,7 @@ async fn push_chunks( route: rds_core::UniHello, wire: Wire, recv: &mut RecvStream, -) -> anyhow::Result<()> { +) -> anyhow::Result> { let width = wire.limits.fetch_streams as usize; let mut tasks = tokio::task::JoinSet::new(); for k in 0..width { @@ -1361,6 +1419,12 @@ async fn push_chunks( std::future::pending().await } }); + // The receiver's `Done` can legitimately arrive while our last chunk + // streams are still finishing: its reads complete on stream FIN, not + // on this task draining. An early Done is consumed here and returned + // for the caller to validate — it must not fall into the abort path + // as an "unexpected control frame". + let mut early_done = None; loop { tokio::select! { result = tasks.join_next() => match result { @@ -1368,17 +1432,35 @@ async fn push_chunks( Some(result) => result??, }, msg = &mut cancelled => { - tasks.abort_all(); - let why = match msg { - Ok(SyncMsg::Cancel { reason }) => reason, - Ok(other) => format!("unexpected control frame {other:?}"), - Err(e) => format!("control read failed: {e}"), - }; - return Err(anyhow::anyhow!("receiver aborted chunk push: {why}")); + match msg { + Ok(SyncMsg::Done { root }) => { + early_done = Some(root); + } + Ok(SyncMsg::Cancel { reason }) => { + tasks.abort_all(); + return Err(anyhow::anyhow!("receiver aborted chunk push: {reason}")); + } + Ok(other) => { + tasks.abort_all(); + return Err(anyhow::anyhow!( + "receiver aborted chunk push: unexpected control frame {other:?}" + )); + } + Err(e) => { + tasks.abort_all(); + return Err(anyhow::anyhow!("receiver aborted chunk push: control read failed: {e}")); + } + } + break; } } } - Ok(()) + // Whether the watcher ended on an early Done or tasks drained first, + // every spawned stream still finishes its own writes before we leave. + while let Some(result) = tasks.join_next().await { + result??; + } + Ok(early_done) } /// Write `Offer` + `ManifestPart` frames for a built manifest. From 4b3adb1c09e54d220268a17b73542991f8a5d128 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:08:30 +0500 Subject: [PATCH 03/10] fix(agent): bound refusal writes and close admission-policy gaps MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Refusal HelloAck writes on the greeting path are wrapped in the documented deadlines (hello for policy/scope refusals, authz for authorization-path answers) and explicitly finished — a stalled peer can no longer park the task mid-refusal. - AgentPolicy::validate now checks authz and shutdown timeouts inside the same 1..=3600s contract as handshake/hello. - Agent::serve runs the same policy preflight as run — programmatic callers previously bypassed it entirely. - TCP target parse and permit checks move to preflight, before a service slot is consumed — a refused connect no longer holds a data lane while its refusal is written. - ResourceGate: the first sample no longer depends on subtracting an interval from Instant::now (fresh-process underflow); max_rss_mb saturates at the byte ceiling instead of silently disabling the gate. - macOS fd accounting uses proc_pidinfo(PROC_PIDLISTFDS) — the true fd table without /dev/fd's fdesc-mount dependency or the read_dir fd skew. - sync_paths scope entries normalize by splitting on both separators, matching check_scope_path's admitted spellings — a `a\b` scope can now actually match the `a/b` requests it was minted for. - Docs clarify Ping/Info are never gated by policy but remain subject to a grant's service scope. --- crates/rds-agent/src/lib.rs | 92 ++++++++++++++++++++++++++-------- crates/rds-agent/src/limits.rs | 15 +++--- crates/rds-agent/src/sys.rs | 35 ++++++++++--- docs/agent-configuration.md | 8 +-- 4 files changed, 111 insertions(+), 39 deletions(-) diff --git a/crates/rds-agent/src/lib.rs b/crates/rds-agent/src/lib.rs index e45100e..4e652aa 100644 --- a/crates/rds-agent/src/lib.rs +++ b/crates/rds-agent/src/lib.rs @@ -132,8 +132,9 @@ pub struct AgentPolicy { /// 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. + /// always-on control plane — never gated by deployment policy, though + /// a grant's service scope may still refuse them. A disabled service + /// is refused before any grant or service work runs. pub services: Option>, /// Admission and greeting deadlines applied per connection/stream. pub timeouts: TimeoutPolicy, @@ -263,8 +264,12 @@ impl AgentPolicy { } if self.timeouts.handshake.is_zero() || self.timeouts.hello.is_zero() + || self.timeouts.authz.is_zero() + || self.timeouts.shutdown.is_zero() || self.timeouts.handshake > MAX_TIMEOUT || self.timeouts.hello > MAX_TIMEOUT + || self.timeouts.authz > MAX_TIMEOUT + || self.timeouts.shutdown > MAX_TIMEOUT { return Err("timeouts must be between 1 and 3600 seconds"); } @@ -506,6 +511,12 @@ impl Agent { /// Serve a single already-established connection. pub async fn serve(&self, conn: Connection) -> anyhow::Result<()> { + // Programmatic callers reach serve() without run()'s preflight — + // the same policy contract applies either way. + if let Err(why) = self.policy.validate() { + conn.close(5u32.into(), b"invalid agent policy"); + anyhow::bail!("invalid agent policy: {why}"); + } if self.policy.grants_required() && self.limits.streams() < 2 { conn.close(5u32.into(), b"invalid grant stream budget"); anyhow::bail!("grant mode requires at least two stream slots"); @@ -665,13 +676,19 @@ async fn serve_stream( && !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"), - }, + // Greeting refusals are bounded by the hello deadline — a stalled + // peer must not park this task on a refusal write. + tokio::time::timeout( + policy.timeouts.hello, + write_frame( + &mut send, + &HelloAck::Error { + message: format!("service {kind:?} not enabled on this agent"), + }, + ), ) - .await?; + .await??; + send.finish()?; anyhow::bail!("service {kind:?} not enabled"); } let grant = match authz.service_scope(&policy).await { @@ -682,22 +699,54 @@ async fn serve_stream( conn.close(2u32.into(), why.message().as_bytes()); } let why = why.message(); - write_frame( - &mut send, - &HelloAck::Error { - message: why.into(), - }, + // Authorization-path answers are bounded by the authz budget. + tokio::time::timeout( + policy.timeouts.authz, + write_frame( + &mut send, + &HelloAck::Error { + message: why.into(), + }, + ), ) - .await?; + .await??; + send.finish()?; anyhow::bail!("stream refused: {why}"); } }; let scope_err = grant.as_ref().and_then(|g| scope_check(g, &hello).err()); if let Some(why) = scope_err { rds_observe::request_refused(Reason::Denied); - write_frame(&mut send, &HelloAck::Error { message: why }).await?; + tokio::time::timeout( + policy.timeouts.authz, + write_frame(&mut send, &HelloAck::Error { message: why }), + ) + .await??; + send.finish()?; anyhow::bail!("stream outside grant scope"); } + // TCP preflight resolves the target and policy before the slot is + // consumed — a refused connect must not hold a service lane while + // its refusal is written. The service arm re-validates on its own. + if let StreamHello::TcpConnect { host, port } = &hello { + let refusal = match rds_core::TcpTarget::new(host, *port) { + Ok(target) if !policy.permits_tcp_target(&target) => { + Some(format!("tcp target {host}:{port} not permitted")) + } + Ok(_) => None, + Err(error) => Some(format!("invalid TCP destination: {error}")), + }; + if let Some(message) = refusal { + rds_observe::request_refused(Reason::Denied); + tokio::time::timeout( + policy.timeouts.hello, + write_frame(&mut send, &HelloAck::Error { message }), + ) + .await??; + send.finish()?; + anyhow::bail!("tcp target refused at preflight"); + } + } // 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. @@ -903,16 +952,15 @@ async fn serve_stream( // Grant v3 `sync_paths` entries passed the decoder's // lexical checks; normalize `.`/empty components the // same way `check_rel_path` normalizes requests so - // prefix matching compares like with like. + // prefix matching compares like with like. Split on + // both separators — `check_scope_path` admits `\` but + // `Path::components` on Unix does not. let paths = g.payload.constraints.sync_paths.as_ref().map(|list| { list.iter() .map(|scope| { - std::path::Path::new(scope) - .components() - .filter_map(|c| match c { - std::path::Component::Normal(p) => Some(p), - _ => None, - }) + scope + .split(['/', '\\']) + .filter(|p| !p.is_empty() && *p != ".") .collect::() }) .collect::>() diff --git a/crates/rds-agent/src/limits.rs b/crates/rds-agent/src/limits.rs index e388ae6..4507e26 100644 --- a/crates/rds-agent/src/limits.rs +++ b/crates/rds-agent/src/limits.rs @@ -37,10 +37,11 @@ impl AgentLimits { /// Process-level ceiling: refuse admissions once the process holds /// `max_fds` descriptors or `max_rss_mb` resident MiB. `None` leaves - /// that quantity ungated. + /// that quantity ungated. An oversized MiB value saturates at the + /// byte ceiling instead of silently disabling the gate. pub fn with_process_budget(mut self, max_fds: Option, max_rss_mb: Option) -> Self { self.max_fds = max_fds; - self.max_rss_bytes = max_rss_mb.and_then(|mb| mb.checked_mul(1024 * 1024)); + self.max_rss_bytes = max_rss_mb.map(|mb| mb.saturating_mul(1024 * 1024)); self } @@ -98,7 +99,9 @@ const SAMPLE_INTERVAL: std::time::Duration = std::time::Duration::from_millis(20 pub(super) struct ResourceGate { max_fds: Option, max_rss_bytes: Option, - cache: std::sync::Mutex<(std::time::Instant, bool)>, + /// `None` forces the first check to sample immediately — subtracting + /// an interval from `Instant::now` can underflow on fresh processes. + cache: std::sync::Mutex<(Option, bool)>, } impl ResourceGate { @@ -109,7 +112,7 @@ impl ResourceGate { Some(Self { max_fds: limits.max_fds, max_rss_bytes: limits.max_rss_bytes, - cache: std::sync::Mutex::new((std::time::Instant::now() - 2 * SAMPLE_INTERVAL, true)), + cache: std::sync::Mutex::new((None, true)), }) } @@ -119,7 +122,7 @@ impl ResourceGate { pub fn allows(&self) -> bool { let mut cache = crate::lock(&self.cache); let (at, verdict) = *cache; - if at.elapsed() < SAMPLE_INTERVAL { + if at.is_some_and(|at| at.elapsed() < SAMPLE_INTERVAL) { return verdict; } let mut ok = true; @@ -138,7 +141,7 @@ impl ResourceGate { } // Nothing observable: keep serving rather than gate on a guess. let verdict = ok || !observed; - *cache = (std::time::Instant::now(), verdict); + *cache = (Some(std::time::Instant::now()), verdict); verdict } } diff --git a/crates/rds-agent/src/sys.rs b/crates/rds-agent/src/sys.rs index c9a2aef..b934362 100644 --- a/crates/rds-agent/src/sys.rs +++ b/crates/rds-agent/src/sys.rs @@ -3,15 +3,34 @@ //! Platform coverage is honest: an unobservable quantity reports `None` //! and the corresponding gate stays open rather than pretending a bound. -/// Open file descriptors owned by this process, where the kernel exposes -/// them (`/proc/self/fd` on Linux, `/dev/fd` on macOS). -#[cfg(any(target_os = "linux", target_os = "macos"))] +/// Open file descriptors owned by this process (`/proc/self/fd`). +#[cfg(target_os = "linux")] pub(crate) fn open_fds() -> Option { - #[cfg(target_os = "linux")] - const FD_DIR: &str = "/proc/self/fd"; - #[cfg(target_os = "macos")] - const FD_DIR: &str = "/dev/fd"; - Some(std::fs::read_dir(FD_DIR).ok()?.count()) + Some(std::fs::read_dir("/proc/self/fd").ok()?.count()) +} + +/// Open file descriptors via `proc_pidinfo` (`PROC_PIDLISTFDS`): a +/// zero-length query returns the fd-table byte size — the true count +/// without `/dev/fd`'s dependency on an fdesc mount or the +1 of the +/// `read_dir` descriptor itself. +#[cfg(target_os = "macos")] +#[allow(unsafe_code)] +pub(crate) fn open_fds() -> Option { + // SAFETY: a null buffer with a zero size is a documented size probe — + // the call writes nothing and returns the table's byte length. + let size = unsafe { + libc::proc_pidinfo( + libc::getpid(), + libc::PROC_PIDLISTFDS, + 0, + std::ptr::null_mut(), + 0, + ) + }; + if size <= 0 { + return None; + } + Some(size as usize / std::mem::size_of::()) } #[cfg(not(any(target_os = "linux", target_os = "macos")))] diff --git a/docs/agent-configuration.md b/docs/agent-configuration.md index 4a9015b..d5130c2 100644 --- a/docs/agent-configuration.md +++ b/docs/agent-configuration.md @@ -54,9 +54,11 @@ nothing. Flag equivalents: `--tenant`, `--policy-min-revision`. ## Service enablement -`Ping` and `Info` are the always-on control plane and are never gated. -The gateable data-plane services are `tcp`, `desktop` and `sync`; `audio` -is wire-reserved but unimplemented and is rejected rather than silently +`Ping` and `Info` are the always-on control plane: deployment policy +cannot disable them, though a grant's `services` scope still applies — +a grant that does not name `ping`/`info` cannot call them. The gateable +data-plane services are `tcp`, `desktop` and `sync`; `audio` is +wire-reserved but unimplemented and is rejected rather than silently accepted. - **No `role`/`services`/`disabled_services`** — the implicit set: `tcp`, From cb500875f69a8322437f010102584ea1ee178df6 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:13:20 +0500 Subject: [PATCH 04/10] fix(discovery): gate record versions before decode, accept hex keys MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Record/DeleteRequest verify now peeks the leading schema version after signature check but before committing to a typed decode — a payload minted by another schema reports "unsupported record version N" instead of an opaque serde failure (the root cause of the mixed- deployment `400 record malformed` incident). - EndpointKey::from_str accepts 64-digit hex alongside base32, matching rds_core::EndpointId's dual spelling — `--directory-allow` no longer rejects the hex that `rds id` prints. - Remove Client::metrics: it called `/v1/metrics`, a route the service dropped when metrics moved to the admin listener — it could only ever return 404. - Connection and worker task groups now count panicked tasks into a new rds_directory_task_panics_total metric instead of silently reaping them (maintenance already accounted failures via gc_failures). - decide() retirement semantics reviewed: the retired revision floor is deliberate and correct — same-revision+same-digest on a retired entry reports Expired, conflicts report Stale. --- crates/rds-discovery/src/client.rs | 7 ---- crates/rds-discovery/src/lib.rs | 56 +++++++++++++++++++++++++ crates/rds-discovery/src/record_wire.rs | 17 ++++++++ crates/rds-discovery/src/service.rs | 19 ++++++++- 4 files changed, 90 insertions(+), 9 deletions(-) diff --git a/crates/rds-discovery/src/client.rs b/crates/rds-discovery/src/client.rs index 61d3cde..636965d 100644 --- a/crates/rds-discovery/src/client.rs +++ b/crates/rds-discovery/src/client.rs @@ -297,13 +297,6 @@ impl Client { self.expect(resp, &[200]).map(|_| ()) } - /// `GET /v1/metrics` — raw prometheus text. - pub async fn metrics(&self) -> Result { - let resp = self.request("GET", "/v1/metrics", &[]).await?; - let resp = self.expect(resp, &[200])?; - String::from_utf8(resp.body).map_err(|e| DiscoveryError::InvalidRecord(e.to_string())) - } - async fn request( &self, method: &str, diff --git a/crates/rds-discovery/src/lib.rs b/crates/rds-discovery/src/lib.rs index 6977b28..c58434a 100644 --- a/crates/rds-discovery/src/lib.rs +++ b/crates/rds-discovery/src/lib.rs @@ -62,7 +62,31 @@ impl std::fmt::Display for EndpointKey { impl std::str::FromStr for EndpointKey { type Err = DiscoveryError; + /// Base32 (canonical display) or 64-digit hex — `rds id` prints hex, + /// so operators legitimately supply either spelling. Same dual + /// acceptance `rds_core::EndpointId` has. fn from_str(s: &str) -> Result { + if s.len() == 64 { + let mut key = [0u8; 32]; + for (i, &[hi, lo]) in s.as_bytes().as_chunks::<2>().0.iter().enumerate() { + let nib = |c: u8| -> Option { + match c { + b'0'..=b'9' => Some(c - b'0'), + b'a'..=b'f' => Some(c - b'a' + 10), + b'A'..=b'F' => Some(c - b'A' + 10), + _ => None, + } + }; + let (hi, lo) = (nib(hi), nib(lo)); + let Some((hi, lo)) = hi.zip(lo) else { + return Err(DiscoveryError::InvalidRecord( + "endpoint key is not hex".into(), + )); + }; + key[i] = (hi << 4) | lo; + } + return Ok(Self(key)); + } let bytes = data_encoding::BASE32_NOPAD .decode(s.to_uppercase().as_bytes()) .map_err(|_| DiscoveryError::InvalidRecord("endpoint key is not base32".into()))?; @@ -200,6 +224,38 @@ mod tests { assert!(matches!(record.verify(), Err(DiscoveryError::BadSignature))); } + #[test] + fn endpoint_key_accepts_hex_and_base32() { + use std::str::FromStr; + let bytes = [ + 0x85, 0xce, 0xb4, 0xd4, 0xab, 0xd6, 0x66, 0xb7, 0x2f, 0xce, 0x2d, 0xc9, 0x7e, 0x23, + 0x2c, 0xe8, 0xf9, 0x03, 0x91, 0x37, 0xa3, 0x95, 0x7d, 0x11, 0x5a, 0x7b, 0x5f, 0xb7, + 0x4f, 0xbb, 0x7e, 0xd4, + ]; + let key = EndpointKey(bytes); + // Canonical base32 display round-trips… + assert_eq!(EndpointKey::from_str(&key.to_string()).unwrap(), key); + // …and the hex spelling `rds id` prints parses identically. + let hex: String = bytes.iter().map(|b| format!("{b:02x}")).collect(); + assert_eq!(EndpointKey::from_str(&hex).unwrap(), key); + assert!(EndpointKey::from_str(&"g".repeat(64)).is_err()); + assert!(EndpointKey::from_str("not-a-key").is_err()); + } + + #[test] + fn signed_record_peeks_its_leading_version() { + let key = SigningKey::from_bytes(&[7u8; 32]); + let record = rec(&key); + assert_eq!( + record_wire::payload_version(&record.payload).unwrap(), + record_wire::RECORD_VERSION + ); + // A payload from a different schema decodes its own version + // rather than producing a serde error mid-struct. + let foreign = postcard::to_stdvec(&2u16).unwrap(); + assert_eq!(record_wire::payload_version(&foreign).unwrap(), 2); + } + #[test] fn memory_store_roundtrip_and_expiry() { let key = SigningKey::from_bytes(&[9u8; 32]); diff --git a/crates/rds-discovery/src/record_wire.rs b/crates/rds-discovery/src/record_wire.rs index 9d23956..a104dd3 100644 --- a/crates/rds-discovery/src/record_wire.rs +++ b/crates/rds-discovery/src/record_wire.rs @@ -101,6 +101,15 @@ fn decode Deserialize<'de>>(bytes: &[u8]) -> Result Result { + let (version, _rest) = + postcard::take_from_bytes::(bytes).map_err(|e| invalid(&e.to_string()))?; + Ok(version) +} + fn metadata( version: u16, revision: u64, @@ -232,6 +241,10 @@ impl EndpointRecord { &key, MAX_RECORD_BYTES, )?; + let version = payload_version(&self.payload)?; + if version != RECORD_VERSION { + return Err(invalid(&format!("unsupported record version {version}"))); + } let payload: Payload = decode(&self.payload)?; if payload.key != self.key { return Err(DiscoveryError::BadSignature); @@ -288,6 +301,10 @@ impl DeleteRequest { let key = VerifyingKey::from_bytes(&self.key.0).map_err(|_| DiscoveryError::BadSignature)?; authority::verify(DELETE_DOMAIN, &self.payload, &self.signature, &key, 128)?; + let version = payload_version(&self.payload)?; + if version != RECORD_VERSION { + return Err(invalid(&format!("unsupported record version {version}"))); + } let payload: DeletePayload = decode(&self.payload)?; if payload.key != self.key { return Err(DiscoveryError::BadSignature); diff --git a/crates/rds-discovery/src/service.rs b/crates/rds-discovery/src/service.rs index 5637881..6cb2ef1 100644 --- a/crates/rds-discovery/src/service.rs +++ b/crates/rds-discovery/src/service.rs @@ -126,6 +126,7 @@ struct Metrics { gc_failures: AtomicU64, connections_rejected: AtomicU64, workers_rejected: AtomicU64, + task_panics: AtomicU64, requests: AtomicU64, } @@ -199,6 +200,10 @@ impl DirectoryMetrics { "rds_directory_workers_rejected_total", m.workers_rejected.load(Ordering::Relaxed), ), + ( + "rds_directory_task_panics_total", + m.task_panics.load(Ordering::Relaxed), + ), ]); for (group, known, tasks, limit) in [ ( @@ -446,9 +451,19 @@ pub async fn serve( let accepted = tokio::select! { biased; _ = stop.changed() => break, - _ = connections.join_next(), if !connections.is_empty() => continue, + result = connections.join_next(), if !connections.is_empty() => { + if let Some(Err(error)) = result + && error.is_panic() + { state.metrics.task_panics.fetch_add(1, Ordering::Relaxed); } + continue; + }, _ = connections.changed() => continue, - _ = state.workers.join_next(), if !state.workers.is_empty() => continue, + result = state.workers.join_next(), if !state.workers.is_empty() => { + if let Some(Err(error)) = result + && error.is_panic() + { state.metrics.task_panics.fetch_add(1, Ordering::Relaxed); } + continue; + }, _ = state.workers.changed() => continue, result = maintenance.join_next(), if !maintenance.is_empty() => { if let Some(Err(error)) = result From dd15847391c3df14963a425e66053abe869526e3 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:22:18 +0500 Subject: [PATCH 05/10] fix(relay,server): fail closed on empty allowlist, bound iroh client rate MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The iroh relay backend admitted every endpoint when --allow was empty, while the owned backend required entries or an explicit dev-open opt-in. Copied production configs silently ran open relays. Both backends now share one admission contract: --allow entries or --development-open-relay. Allowlist keys are validated like directory enrollment — weak ed25519 keys and duplicates are rejected at startup instead of silently holding dead entries. The iroh backend also gains the owned relay's per-client envelope via upstream's implemented client_rx knob (64 MiB/s + 4 MiB burst); upstream exposes no connection-count cap, which the docs now state. Direct serve() callers get an explicit open-access warning. An allowlist e2e test proves the full path: a registered identity reaches its home relay while an unlisted one reports auth_denied_reason and never connects. Shared runtime fixtures opt into development-open where they exercise lifecycle or TLS validation rather than admission, and the former dev-open rejection case becomes the dev-open/allowlist conflict. Docs: relay-runtime backend table and deployment admission section describe the unified fail-closed contract. --- crates/rds-relay/src/lib.rs | 17 ++++++- crates/rds-relay/src/runtime.rs | 63 +++++++++++++++++------ crates/rds-relay/tests/allowlist_e2e.rs | 68 +++++++++++++++++++++++++ crates/rds-server/tests/lifecycle.rs | 3 ++ crates/rds-server/tests/runtime.rs | 4 +- docs/deployment.md | 6 ++- docs/relay-runtime.md | 7 +-- tests/support/admin_cli.rs | 9 ++++ tests/support/relay_runtime.rs | 19 +++++-- 9 files changed, 171 insertions(+), 25 deletions(-) create mode 100644 crates/rds-relay/tests/allowlist_e2e.rs diff --git a/crates/rds-relay/src/lib.rs b/crates/rds-relay/src/lib.rs index b39555e..604c010 100644 --- a/crates/rds-relay/src/lib.rs +++ b/crates/rds-relay/src/lib.rs @@ -33,6 +33,7 @@ pub const OWNED_BACKEND_COMPILED: bool = cfg!(feature = "owned-relay"); use std::collections::HashSet; use std::net::SocketAddr; +use std::num::NonZeroU32; use std::path::PathBuf; use std::sync::Arc; @@ -110,7 +111,14 @@ async fn serve_prepared( tls: Option, ) -> anyhow::Result { let mut relay_config = RelayConfig::new(addr); - if !allow.is_empty() { + if allow.is_empty() { + // Callers through RelayArgs must opt in with --development-open-relay; + // a direct API caller receives the same loud signal here. + tracing::warn!( + "relay serving with open access: every endpoint id is admitted. \ + Set --allow entries for production deployments" + ); + } else { relay_config.access = Arc::new(AllowList( allow .into_iter() @@ -120,6 +128,13 @@ async fn serve_prepared( .collect(), )); } + // Per-client RX limits mirror the owned relay's token bucket. Upstream has + // no implemented connection-count cap; the owned backend adds one there. + let mut client_rx = iroh_relay::server::ClientRateLimit::new( + NonZeroU32::new(64 * 1024 * 1024).expect("nonzero"), + ); + client_rx.max_burst_bytes = NonZeroU32::new(4 * 1024 * 1024); + relay_config.limits.client_rx = Some(client_rx); relay_config.tls = tls; let mut config = ServerConfig::default(); config.relay = Some(relay_config); diff --git a/crates/rds-relay/src/runtime.rs b/crates/rds-relay/src/runtime.rs index b84b69a..2204086 100644 --- a/crates/rds-relay/src/runtime.rs +++ b/crates/rds-relay/src/runtime.rs @@ -22,10 +22,11 @@ pub struct RelayArgs { /// Owned relay connections, including incomplete handshakes and registrations. #[arg(long)] pub relay_max_connections: Option, - /// Explicitly allow unknown peers on an owned relay for isolated development. + /// Explicitly allow unknown peers for isolated development. Both relay + /// backends run fail-closed without this flag or --allow entries. #[arg(long, conflicts_with = "allow")] pub development_open_relay: bool, - /// Allowed endpoint id (repeatable). Owned production mode requires entries. + /// Allowed endpoint id (repeatable). Production mode requires entries. #[arg(long = "allow")] pub allow: Vec, /// PEM certificate chain for the iroh HTTPS relay. @@ -47,6 +48,26 @@ pub struct RelayArgs { pub tls_acme_staging: bool, } +/// Relay allowlist entries are endpoint identities: reject weak ed25519 +/// keys like directory enrollment does, and duplicates that signal a +/// configuration mistake. +fn validate_allowlist(allow: &[EndpointId]) -> Result<(), RelayConfigError> { + let mut seen = std::collections::HashSet::new(); + for id in allow { + if id.verifying_key().is_weak() { + return Err(RelayConfigError::Invalid( + "relay allowlist contains a weak key", + )); + } + if !seen.insert(*id) { + return Err(RelayConfigError::Invalid( + "relay allowlist contains a duplicate key", + )); + } + } + Ok(()) +} + #[derive(Debug, thiserror::Error)] pub enum RelayConfigError { #[error("{0}")] @@ -78,12 +99,20 @@ impl RelayArgs { use RelayConfigError::Invalid; match self.relay_backend { RelayBackend::Iroh => { - if self.relay_key_file.is_some() - || self.relay_max_connections.is_some() - || self.development_open_relay - { + if self.relay_key_file.is_some() || self.relay_max_connections.is_some() { return Err(Invalid("owned-relay flags require --relay-backend noq")); } + if self.development_open_relay && !self.allow.is_empty() { + return Err(Invalid( + "development-open relay mode conflicts with an allowlist", + )); + } + if !self.development_open_relay && self.allow.is_empty() { + return Err(Invalid( + "iroh relay requires --allow entries or explicit --development-open-relay", + )); + } + validate_allowlist(&self.allow)?; let manual = self.tls_cert.is_some() || self.tls_key.is_some(); let acme = !self.tls_acme_domain.is_empty(); if manual && acme { @@ -154,6 +183,7 @@ impl RelayArgs { "owned relay requires --allow entries or explicit --development-open-relay", )); } + validate_allowlist(&self.allow)?; let key_file = self .relay_key_file .ok_or(Invalid("owned relay requires --relay-key-file"))?; @@ -442,15 +472,18 @@ mod tests { #[tokio::test] async fn canceled_iroh_observer_preserves_listener_and_shutdown() { - let mut relay = RelayArgs::default() - .prepare() - .unwrap() - .initialize() - .await - .unwrap() - .bind("127.0.0.1:0".parse().unwrap()) - .await - .unwrap(); + let mut relay = RelayArgs { + development_open_relay: true, + ..Default::default() + } + .prepare() + .unwrap() + .initialize() + .await + .unwrap() + .bind("127.0.0.1:0".parse().unwrap()) + .await + .unwrap(); assert!( tokio::time::timeout(Duration::from_millis(25), relay.stopped()) .await diff --git a/crates/rds-relay/tests/allowlist_e2e.rs b/crates/rds-relay/tests/allowlist_e2e.rs new file mode 100644 index 0000000..7491d7c --- /dev/null +++ b/crates/rds-relay/tests/allowlist_e2e.rs @@ -0,0 +1,68 @@ +//! End-to-end: the iroh relay's endpoint allowlist admits registered +//! identities and denies unknown ones, observed client-side through the +//! endpoint's home-relay connection status. +use std::{net::SocketAddr, str::FromStr, time::Duration}; + +use iroh::{Endpoint, RelayMap, RelayMode, RelayUrl, SecretKey, Watcher, endpoint::presets}; + +fn key(seed: u8) -> SecretKey { + SecretKey::from_bytes(&[seed; 32]) +} + +fn relay_url(addr: SocketAddr) -> RelayUrl { + RelayUrl::from_str(&format!("http://{addr}")).unwrap() +} + +async fn endpoint(key: SecretKey, url: &RelayUrl) -> Endpoint { + Endpoint::builder(presets::Minimal) + .secret_key(key) + .relay_mode(RelayMode::Custom(RelayMap::from_iter([url.clone()]))) + .bind() + .await + .unwrap() +} + +#[tokio::test] +async fn allowlist_admits_registered_and_denies_unknown_endpoints() { + let allowed = key(17); + let allowed_id = rds_core::EndpointId::from_bytes(allowed.public().as_bytes()).unwrap(); + let relay = rds_relay::serve("127.0.0.1:0".parse().unwrap(), vec![allowed_id], None) + .await + .unwrap(); + let url = relay_url(relay.http_addr().expect("http relay listener")); + + let admitted = endpoint(allowed, &url).await; + let denied = endpoint(key(42), &url).await; + + // The admitted identity reaches its home relay. + let mut admitted_status = admitted.home_relay_status(); + tokio::time::timeout(Duration::from_secs(15), async { + loop { + if admitted_status.get().iter().any(|s| s.is_connected()) { + break; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + }) + .await + .expect("allowed endpoint never reached its home relay"); + + // An unlisted identity is rejected at relay authentication: it reports an + // auth denial and must never appear connected. + let mut denied_status = denied.home_relay_status(); + tokio::time::timeout(Duration::from_secs(15), async { + loop { + let status = denied_status.get(); + assert!( + !status.iter().any(|s| s.is_connected()), + "denied endpoint must never reach its home relay" + ); + if status.iter().any(|s| s.auth_denied_reason().is_some()) { + break; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + }) + .await + .expect("denied endpoint never reported a relay auth denial"); +} diff --git a/crates/rds-server/tests/lifecycle.rs b/crates/rds-server/tests/lifecycle.rs index 9ce4724..836ec12 100644 --- a/crates/rds-server/tests/lifecycle.rs +++ b/crates/rds-server/tests/lifecycle.rs @@ -34,6 +34,9 @@ impl Scratch { let mut command = Command::new(env!("CARGO_BIN_EXE_rds-server")); command .args(["--http-addr", "127.0.0.1:0", "--relay-addr", relay]) + // These tests exercise service lifecycle, not relay admission: + // production runs fail closed without --allow entries. + .arg("--development-open-relay") .arg("--directory") .arg(self.0.join("records")) .env("RUST_LOG", "rds_server=info") diff --git a/crates/rds-server/tests/runtime.rs b/crates/rds-server/tests/runtime.rs index 8ef47d4..328b431 100644 --- a/crates/rds-server/tests/runtime.rs +++ b/crates/rds-server/tests/runtime.rs @@ -56,13 +56,15 @@ async fn invalid_host_inputs_are_rejected_before_identity_or_catalog_creation() // Probe the compiled relay library — feature unification can put the // owned backend into the spawned binary while this package's flag is // off. + // Open-relay opt-in keeps the iroh backend admissible so the + // assertion below is about the case's own invalid input. + values.extend(args(&["--development-open-relay"])); if rds_relay::OWNED_BACKEND_COMPILED { values.extend(args(&[ "--relay-backend", "noq", "--relay-key-file", &scratch.key(), - "--development-open-relay", ])); } assert!(!rejected(&scratch, values).await.is_empty()); diff --git a/docs/deployment.md b/docs/deployment.md index df9f4cb..fafca40 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -43,7 +43,11 @@ APIs. Give it a separate authorized key. Updated agent, direct CLI and owned relay binaries refuse a seed inode already in use before binding. Upgrade all local binaries together; old binaries and copied/manual-replaced keys are outside this cooperative guarantee. See [migration and ownership](local-sessions.md). -The relay's `--allow` covers every endpoint that may use the relay. +The relay's `--allow` covers every endpoint that may use the relay. Both +relay backends fail closed: without `--allow` entries a relay refuses to +start unless `--development-open-relay` explicitly selects open admission +for isolated development. Allowlist entries are validated like directory +enrollment — weak or duplicate keys are rejected at startup. Agent role, service, peer, authority, limit and timeout policy can live in a versioned JSON file instead of flags — `--agent-config /etc/rds/agent.json`; diff --git a/docs/relay-runtime.md b/docs/relay-runtime.md index 75ec526..857741f 100644 --- a/docs/relay-runtime.md +++ b/docs/relay-runtime.md @@ -15,9 +15,10 @@ dependency claim. | Bind | HTTP TCP; optional separate HTTPS | UDP QUIC | | Address flag | `--addr` in relay; `--relay-addr` in server | Same flags, interpreted as UDP | | Identity | Existing iroh relay behavior | Required `--relay-key-file`, separate from device/authority keys | -| Admission | Existing `--allow`; empty retains legacy open mode | Nonempty `--allow` required; unknown keys denied | -| Development | Existing iroh behavior | Explicit `--development-open-relay`, incompatible with `--allow` | -| Connection cap | Existing iroh implementation | `--relay-max-connections`, positive u16, default 256 | +| Admission | Nonempty `--allow` required; unknown keys denied; weak or duplicate keys rejected at startup | Same contract | +| Development | Explicit `--development-open-relay`, incompatible with `--allow` | Same contract | +| Per-client rate | 64 MiB/s RX with a 4 MiB burst per client | Token bucket at the same envelope | +| Connection cap | Not exposed by upstream | `--relay-max-connections`, positive u16, default 256 | | TLS flags | Manual PEM or in-process ACME | HTTP TLS/ACME options rejected; QUIC already authenticates its pinned key | The owned cap includes in-progress handshakes and registrations. Releasing a diff --git a/tests/support/admin_cli.rs b/tests/support/admin_cli.rs index b4ac0e6..feab6e4 100644 --- a/tests/support/admin_cli.rs +++ b/tests/support/admin_cli.rs @@ -50,6 +50,12 @@ impl Fixture { } "relay" => { cmd.args(["--addr", "127.0.0.1:0"]); + if !OWNED { + // The iroh backend fails closed without admission + // configuration; these fixtures are open development + // relays by contract. + cmd.arg("--development-open-relay"); + } } "server" => { cmd.args([ @@ -60,6 +66,9 @@ impl Fixture { "--directory", ]) .arg(self.0.join("records")); + if !OWNED { + cmd.arg("--development-open-relay"); + } } _ => panic!("unknown fixture role"), } diff --git a/tests/support/relay_runtime.rs b/tests/support/relay_runtime.rs index c9796d6..576f07e 100644 --- a/tests/support/relay_runtime.rs +++ b/tests/support/relay_runtime.rs @@ -239,9 +239,13 @@ async fn invalid_or_misapplied_relay_flags_do_not_initialize_state() { 0 => args(&["--relay-backend", "unknown"]), 1 => args(&["--relay-key-file", &key]), 2 => args(&["--relay-max-connections", "1"]), - 3 => args(&["--development-open-relay"]), - 4 => args(&["--tls-https-addr", "127.0.0.1:0"]), - 5 => args(&["--tls-acme-staging"]), + 3 => args(&["--development-open-relay", "--allow", &allow]), + 4 => args(&[ + "--development-open-relay", + "--tls-https-addr", + "127.0.0.1:0", + ]), + 5 => args(&["--development-open-relay", "--tls-acme-staging"]), 6 => args(&["--relay-backend", "noq", "--relay-key-file", &key]), 7 => args(&["--relay-backend", "noq", "--development-open-relay"]), 8 => args(&[ @@ -341,6 +345,7 @@ async fn malformed_or_oversized_relay_pem_is_rejected_before_catalog_creation() let error = rejected( &scratch, args(&[ + "--development-open-relay", "--tls-cert", cert.to_str().unwrap(), "--tls-key", @@ -357,7 +362,11 @@ async fn malformed_or_oversized_relay_pem_is_rejected_before_catalog_creation() #[tokio::test] async fn explicit_iroh_mode_preserves_listener_and_shutdown_contract() { let scratch = Scratch::new(); - let mut process = scratch.spawn(&args(&["--relay-backend", "iroh"])); + let mut process = scratch.spawn(&args(&[ + "--relay-backend", + "iroh", + "--development-open-relay", + ])); let ready = process.ready(false).await; assert!(ready.id.is_none()); if let Some(addr) = ready.directory { @@ -397,6 +406,7 @@ async fn mismatched_or_oversized_private_key_is_rejected_before_catalog_creation let error = rejected( &scratch, args(&[ + "--development-open-relay", "--tls-cert", cert.to_str().unwrap(), "--tls-key", @@ -415,6 +425,7 @@ async fn valid_manual_iroh_tls_starts_and_shuts_down() { let scratch = Scratch::new(); let (cert, key) = tls_fixture(&scratch); let mut process = scratch.spawn(&args(&[ + "--development-open-relay", "--tls-cert", cert.to_str().unwrap(), "--tls-key", From ab0b48161020c7e09d933ff7a24c26a6a96a9884 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:36:12 +0500 Subject: [PATCH 06/10] fix(core,client,cli): sanitize wire errors, bound tickets, harden file opens MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire-error hygiene: refusal reasons crossing the stream now use fixed vocabulary — tcp connect failures report io::ErrorKind, sync refusals report coarse causes, and desktop unavailability no longer leaks capability-probe internals. Real errors still reach the local log via the propagating bail. Structural config checks: rds --registry-epoch is Option (requires --registry-key, so an epoch can no longer be silently ignored), the managed-session rejection now tests presence rather than the 1 sentinel, and the agent schema validator rejects zero authority epochs from either flags or config file — Authority::new caught them only at bind time. effective_services is now the single honest service source: an explicit services list containing desktop on a headless build was reported as enabled while Info and the directory record disagreed; the explicit arm now filters Desktop by the build feature like the implicit arm does. File opens: read_rotations, --directory-ca and relay PEM reads used plain opens that block forever on FIFOs and never checked file type. All three now use NONBLOCK opens plus the regular-file check, matching the posture agent config and grant files already had. Tickets: parse_target accepted unbounded base32 bodies; anything over 4 KiB is refused before decode allocates an unbounded address set. --- crates/rds-agent/src/lib.rs | 22 +++++++++++++------ crates/rds-agent/src/settings.rs | 15 +++++++++++++ crates/rds-cli/Cargo.toml | 2 +- crates/rds-cli/src/main.rs | 34 +++++++++++++++++++++++------ crates/rds-discovery/Cargo.toml | 2 +- crates/rds-discovery/src/lib.rs | 2 +- crates/rds-discovery/src/policy.rs | 21 ++++++++++++++++-- crates/rds-net/src/backends/iroh.rs | 6 +++++ crates/rds-relay/Cargo.toml | 3 ++- crates/rds-relay/src/lib.rs | 17 ++++++++++++--- crates/rds-sync/src/engine.rs | 8 ++++--- 11 files changed, 106 insertions(+), 26 deletions(-) diff --git a/crates/rds-agent/src/lib.rs b/crates/rds-agent/src/lib.rs index 4e652aa..42d4862 100644 --- a/crates/rds-agent/src/lib.rs +++ b/crates/rds-agent/src/lib.rs @@ -210,11 +210,14 @@ impl AgentPolicy { pub fn effective_services(&self) -> BTreeSet { 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 - ) + // One honest source: an explicit `desktop` entry only counts when + // this binary can actually serve it, matching the implicit arm + // and the directory announcement. `validate` still rejects the + // flag combination at startup. + Some(explicit) => set.extend(explicit.iter().copied().filter(|k| match k { + ServiceKind::Tcp | ServiceKind::Sync => true, + ServiceKind::Desktop => cfg!(feature = "desktop"), + _ => false, })), None => { set.insert(ServiceKind::Tcp); @@ -847,13 +850,17 @@ async fn serve_stream( tokio::io::copy_bidirectional(&mut tcp, &mut quic).await?; } Err(e) => { + // ErrorKind is a fixed vocabulary — the raw OS error + // string (errno text, platform internals) never + // crosses the wire. The bail below logs it locally. write_frame( &mut send, &HelloAck::Error { - message: format!("connect {host}:{port} failed: {e}"), + message: format!("connect {host}:{port} failed: {}", e.kind()), }, ) .await?; + anyhow::bail!("tcp connect {host}:{port} failed: {e}"); } } } @@ -884,10 +891,11 @@ async fn serve_stream( write_frame( &mut send, &HelloAck::Error { - message: format!("desktop unavailable: {e}"), + message: "desktop unavailable".into(), }, ) .await?; + anyhow::bail!("desktop capability probe failed: {e}"); } } #[cfg(not(feature = "desktop"))] diff --git a/crates/rds-agent/src/settings.rs b/crates/rds-agent/src/settings.rs index a0a1ab6..99a9388 100644 --- a/crates/rds-agent/src/settings.rs +++ b/crates/rds-agent/src/settings.rs @@ -538,6 +538,21 @@ impl AgentSettings { return Err(AgentConfigError::Invalid("at most 256 allowed peers")); } let authority = &self.authority; + if authority + .registry + .as_ref() + .and_then(|r| r.epoch) + .is_some_and(|e| e == 0) + || authority + .revocations + .as_ref() + .and_then(|r| r.epoch) + .is_some_and(|e| e == 0) + { + return Err(AgentConfigError::Invalid( + "authority epochs must be positive", + )); + } if authority .grant_ttl_secs .is_some_and(|secs| secs == 0 || secs > MAX_GRANT_TTL_SECS) diff --git a/crates/rds-cli/Cargo.toml b/crates/rds-cli/Cargo.toml index 12a38ad..3b6d4e3 100644 --- a/crates/rds-cli/Cargo.toml +++ b/crates/rds-cli/Cargo.toml @@ -23,7 +23,7 @@ rds-net.workspace = true rds-discovery.workspace = true rds-sync.workspace = true rds-ssh.workspace = true -rustix = { workspace = true, features = ["termios", "event"] } +rustix = { workspace = true, features = ["fs", "termios", "event"] } serde_json.workspace = true tokio.workspace = true tracing.workspace = true diff --git a/crates/rds-cli/src/main.rs b/crates/rds-cli/src/main.rs index 939993a..2a07cab 100644 --- a/crates/rds-cli/src/main.rs +++ b/crates/rds-cli/src/main.rs @@ -66,9 +66,9 @@ struct Cli { /// Required for device names; tickets and pinned keys are independent. #[arg(long, global = true, requires = "server")] registry_key: Option, - /// Bootstrap registry authority epoch. - #[arg(long, global = true, default_value = "1")] - registry_epoch: u64, + /// Bootstrap registry authority epoch (default 1). + #[arg(long, global = true, requires = "registry_key")] + registry_epoch: Option, /// Private name-trust state directory; default is beside the endpoint key. #[arg(long, global = true, requires = "registry_key")] registry_state: Option, @@ -280,11 +280,13 @@ async fn run(cli: Cli) -> anyhow::Result<()> { .map(|origin| -> anyhow::Result<_> { let mut client = rds_discovery::client::Client::from_endpoint(origin)?; if let Some(path) = &cli.directory_ca { - client = client.with_ca_pem(&std::fs::read(path)?)?; + client = client.with_ca_pem(&read_pem_file(path)?)?; } if let Some(key) = cli.registry_key.as_deref() { - let authority = - rds_discovery::authority::Authority::from_base32(key, cli.registry_epoch)?; + let authority = rds_discovery::authority::Authority::from_base32( + key, + cli.registry_epoch.map_or(1, std::num::NonZeroU64::get), + )?; let path = cli .registry_state .clone() @@ -488,7 +490,7 @@ fn validate_managed(cli: &Cli) -> anyhow::Result<()> { && cli.server.is_none() && cli.directory_ca.is_none() && cli.registry_key.is_none() - && cli.registry_epoch == 1 + && cli.registry_epoch.is_none() && cli.registry_state.is_none() && cli.authority_rotation.is_empty(), "managed commands use agent configuration; configure the agent or use --direct with a separate identity" @@ -496,6 +498,24 @@ fn validate_managed(cli: &Cli) -> anyhow::Result<()> { Ok(()) } +/// Bounded operator-supplied PEM read. NONBLOCK plus the regular-file +/// check refuses FIFOs/devices before the size bound is applied — a +/// plain open on a FIFO blocks forever. +fn read_pem_file(path: &std::path::Path) -> anyhow::Result> { + use rustix::fs::{Mode, OFlags}; + use std::io::Read; + let file = std::fs::File::from(rustix::fs::open( + path, + OFlags::RDONLY | OFlags::NONBLOCK | OFlags::CLOEXEC, + Mode::empty(), + )?); + anyhow::ensure!(file.metadata()?.is_file(), "{path:?} is not a regular file"); + let mut bytes = Vec::new(); + file.take(1024 * 1024 + 1).read_to_end(&mut bytes)?; + anyhow::ensure!(bytes.len() <= 1024 * 1024, "{path:?} exceeds 1 MiB"); + Ok(bytes) +} + fn control_directory(cli: &Cli) -> anyhow::Result { match &cli.control_dir { Some(path) => Ok(path.clone()), diff --git a/crates/rds-discovery/Cargo.toml b/crates/rds-discovery/Cargo.toml index 5877cf9..e239e61 100644 --- a/crates/rds-discovery/Cargo.toml +++ b/crates/rds-discovery/Cargo.toml @@ -13,7 +13,7 @@ ed25519-dalek.workspace = true postcard = { workspace = true, features = ["use-std"] } rand.workspace = true redb.workspace = true -rustix = { workspace = true, features = ["time"] } +rustix = { workspace = true, features = ["fs", "time"] } rustls.workspace = true rustls-pki-types.workspace = true serde.workspace = true diff --git a/crates/rds-discovery/src/lib.rs b/crates/rds-discovery/src/lib.rs index c58434a..6d4b5c4 100644 --- a/crates/rds-discovery/src/lib.rs +++ b/crates/rds-discovery/src/lib.rs @@ -248,7 +248,7 @@ mod tests { let record = rec(&key); assert_eq!( record_wire::payload_version(&record.payload).unwrap(), - record_wire::RECORD_VERSION + RECORD_VERSION ); // A payload from a different schema decodes its own version // rather than producing a serde error mid-struct. diff --git a/crates/rds-discovery/src/policy.rs b/crates/rds-discovery/src/policy.rs index 3186969..a044f39 100644 --- a/crates/rds-discovery/src/policy.rs +++ b/crates/rds-discovery/src/policy.rs @@ -31,8 +31,25 @@ pub fn read_rotations(paths: &[std::path::PathBuf]) -> Result anyhow::Result { fn read_bounded(path: &std::path::Path, limit: usize) -> std::io::Result> { use std::io::Read; + // NONBLOCK plus the regular-file check refuses FIFOs/devices before the + // size bound is applied — a plain open on a FIFO blocks forever. + let file = std::fs::File::from(rustix::fs::open( + path, + rustix::fs::OFlags::RDONLY | rustix::fs::OFlags::NONBLOCK | rustix::fs::OFlags::CLOEXEC, + rustix::fs::Mode::empty(), + )?); + if !file.metadata()?.is_file() { + return Err(std::io::Error::new( + std::io::ErrorKind::InvalidInput, + "relay PEM path is not a regular file", + )); + } let mut bytes = Vec::new(); - std::fs::File::open(path)? - .take(limit as u64 + 1) - .read_to_end(&mut bytes)?; + file.take(limit as u64 + 1).read_to_end(&mut bytes)?; if bytes.len() > limit { return Err(std::io::Error::new( std::io::ErrorKind::InvalidData, diff --git a/crates/rds-sync/src/engine.rs b/crates/rds-sync/src/engine.rs index fd1d5bd..f5b6e3f 100644 --- a/crates/rds-sync/src/engine.rs +++ b/crates/rds-sync/src/engine.rs @@ -689,13 +689,15 @@ async fn serve_inner( .context("destination preflight task")? }; if let Err(e) = preflight { - refuse(wire, send, &e.to_string()).await?; + // Wire reasons are coarse by contract: host errno details + // stay in the local log, not on the peer's error chain. + refuse(wire, send, "sync destination not writable").await?; bail!("offer refused: {e}"); } let manifest = match read_manifest(recv, size, root, chunk_count, wire).await { Ok(m) => m, Err(e) => { - refuse(wire, send, &e.to_string()).await?; + refuse(wire, send, "invalid sync manifest").await?; return Err(e); } }; @@ -1152,7 +1154,7 @@ async fn receive( let journal = match journal { Ok(journal) => journal, Err(e) => { - refuse(wire, send, &e.to_string()).await?; + refuse(wire, send, "cannot open transfer journal").await?; return Err(e.into()); } }; From 5883daeb757f08e19b29f6431a60cb75b3788644 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:42:34 +0500 Subject: [PATCH 07/10] fix(net): stop fabricating incoming peer addr, surface ticket errors - Incoming::remote_address returns Option: the transient state between accepting a raw attempt and owning its handshake previously answered 0.0.0.0:0, which masquerades as a peer address - parse_target surfaces the real ticket decode error for rds1-prefixed input instead of shadowing it with an endpoint-id complaint - regression tests cover malformed, oversized and non-ticket targets Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- crates/rds-net/src/backends/iroh.rs | 41 +++++++++++++++++++++++--- crates/rds-net/src/backends/noq/mod.rs | 12 ++++---- 2 files changed, 44 insertions(+), 9 deletions(-) diff --git a/crates/rds-net/src/backends/iroh.rs b/crates/rds-net/src/backends/iroh.rs index 8a3139d..9635446 100644 --- a/crates/rds-net/src/backends/iroh.rs +++ b/crates/rds-net/src/backends/iroh.rs @@ -254,10 +254,16 @@ impl FromStr for Ticket { pub fn parse_target(target: &str) -> anyhow::Result { match Ticket::from_str(target) { Ok(ticket) => Ok(ticket.0), - Err(_) => Ok(EndpointAddr { - id: EndpointId::from_str(target)?, - addrs: Default::default(), - }), + Err(ticket_error) => match EndpointId::from_str(target) { + Ok(id) => Ok(EndpointAddr { + id, + addrs: Default::default(), + }), + // An rds1-prefixed string was meant as a ticket: surface the + // ticket decode error, not a misleading endpoint-id complaint. + Err(_) if target.starts_with("rds1") => Err(ticket_error), + Err(id_error) => Err(id_error.into()), + }, } } @@ -302,6 +308,33 @@ mod tests { assert!(addr.addrs.is_empty()); } + /// An `rds1` string that fails ticket decode must surface the ticket + /// error, not an endpoint-id complaint — the prefix already committed + /// the input to the ticket grammar. + #[test] + fn malformed_ticket_surfaces_ticket_error() { + let err = parse_target("rds1!!!not-base32!!!").unwrap_err(); + assert!( + err.to_string().contains("ticket"), + "unexpected error: {err}" + ); + + let oversized = format!("rds1{}", "a".repeat(MAX_TICKET_BODY + 1)); + let err = parse_target(&oversized).unwrap_err(); + assert_eq!(err.to_string(), "ticket too long"); + } + + /// A non-prefixed string that is neither ticket nor endpoint id keeps + /// the endpoint-id error — it was never a ticket attempt. + #[test] + fn plain_garbage_surfaces_endpoint_error() { + let err = parse_target("not-a-ticket-or-key").unwrap_err(); + assert!( + !err.to_string().contains("ticket"), + "unexpected error: {err}" + ); + } + /// The owned types must encode byte-identically to iroh-base: /// tickets and discovery records written by either side decode /// on the other. Guards the postcard layout contract. diff --git a/crates/rds-net/src/backends/noq/mod.rs b/crates/rds-net/src/backends/noq/mod.rs index c301703..1f96226 100644 --- a/crates/rds-net/src/backends/noq/mod.rs +++ b/crates/rds-net/src/backends/noq/mod.rs @@ -728,12 +728,14 @@ impl Incoming { } } - /// Address the attempt arrived from. - pub fn remote_address(&self) -> SocketAddr { + /// Address the attempt arrived from. `None` during the brief window + /// between accepting the raw attempt and owning its handshake — + /// fabricating a wildcard address there would masquerade as a peer. + pub fn remote_address(&self) -> Option { match (&self.incoming, &self.connecting) { - (Some(i), _) => i.remote_address(), - (None, Some(c)) => c.remote_address(), - (None, None) => SocketAddr::from(([0, 0, 0, 0], 0)), + (Some(i), _) => Some(i.remote_address()), + (None, Some(c)) => Some(c.remote_address()), + (None, None) => None, } } } From 52ce54b6e98b26ac39048f45da23b92a2e1e0cde Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:51:10 +0500 Subject: [PATCH 08/10] chore(deps): drop tls12, dead deps and features; compact record store - rustls/tokio-rustls lose the tls12 feature: every TLS peer here is our own rustls binary and QUIC is 1.3-only; 1.2 has no legitimate client - remove the unused workspace async-trait declaration and the never enabled, never gated rds-core/desktop feature - compact the redb record store on open so expired/deleted records reclaim pages instead of growing the file monotonically; a failed compact is maintenance, not a fatal open error - document the exact russh pin rationale in the manifest itself - semver-compatible lockfile refresh Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- Cargo.lock | 142 +++++++++++++--------------- Cargo.toml | 10 +- crates/rds-core/Cargo.toml | 1 - crates/rds-discovery/src/records.rs | 7 +- 4 files changed, 81 insertions(+), 79 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 42cb477..10cb7fd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -83,12 +83,6 @@ dependencies = [ "memchr", ] -[[package]] -name = "allocator-api2" -version = "0.2.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923" - [[package]] name = "android_system_properties" version = "0.1.6" @@ -193,7 +187,7 @@ dependencies = [ "nom", "num-traits", "rusticata-macros", - "thiserror 2.0.20", + "thiserror 2.0.21", "time", ] @@ -446,9 +440,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.4.7" +version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "54413ede23c2daf518f35156dfde027feb2374004d63bd497f983c8db9c0e313" +checksum = "f360145194ee8e21db5ee7f3fcd4fe52210864c75c985dae33218202c8bbe040" dependencies = [ "find-msvc-tools", "jobserver", @@ -575,7 +569,7 @@ version = "0.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0fa961b519f0b462e3a3b4a34b64d119eeaca1d59af726fe450bbba07a9fc0a1" dependencies = [ - "thiserror 2.0.20", + "thiserror 2.0.21", ] [[package]] @@ -1095,9 +1089,9 @@ checksum = "64cd1e32ddd350061ae6edb1b082d7c54915b5c672c389143b9a63403a109f24" [[package]] name = "find-msvc-tools" -version = "0.1.13" +version = "0.1.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ef25905e51abafe4dcea6c15fec58c57b601cdbd0ee53d22ea1d3016c587d39b" +checksum = "aedcfb3409746eddb02b9e19ebda1c3394f759a152e48ee875a0844d1b955484" [[package]] name = "fnv" @@ -1402,8 +1396,6 @@ version = "0.17.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a" dependencies = [ - "allocator-api2", - "equivalent", "foldhash", ] @@ -1553,16 +1545,17 @@ dependencies = [ [[package]] name = "hyper-util" -version = "0.1.20" +version = "0.1.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" +checksum = "ddc03d96684f9226b8a787cdb71488417b53ab5ea8fdb1dac946cb9431cc8bff" dependencies = [ - "base64 0.22.1", + "base64 0.23.1", "bytes", "futures-channel", "futures-util", "http", "http-body", + "httparse", "hyper", "ipnet", "libc", @@ -1870,9 +1863,9 @@ dependencies = [ [[package]] name = "iroh-metrics" -version = "1.0.1" +version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "291065721ad7c477b972e581bbc528df031dc8eb5e39fe1ff3300ae5dfb157ef" +checksum = "8ede55536349842337f7cdc63012ec33dc637945f3259c98a0763aac364cdf09" dependencies = [ "http-body-util", "hyper", @@ -2000,7 +1993,7 @@ dependencies = [ "jni-sys 0.4.1", "log", "simd_cesu8", - "thiserror 2.0.20", + "thiserror 2.0.21", "walkdir", "windows-link", ] @@ -2058,9 +2051,9 @@ dependencies = [ [[package]] name = "js-sys" -version = "0.3.105" +version = "0.3.106" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ce57d20d1ea864ce2ac172ab472d409214f4fd359f0b2a2775abdf522e2af99e" +checksum = "7883d941dae510fb2d978fc3fe018c71c9e2892fd38854de3e8b92c2e5ad9cc5" dependencies = [ "cfg-if", "futures-util", @@ -2147,9 +2140,9 @@ dependencies = [ [[package]] name = "lru" -version = "0.18.4" +version = "0.18.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ff9840bcc50b71349309900da0ce7279aa336ae71d73250b07998932c7d97c25" +checksum = "ef9ac18847474e638e3702b76c65d4eb93428471a74778ef0f1be711717f89b5" dependencies = [ "hashbrown 0.17.1", ] @@ -2434,7 +2427,7 @@ dependencies = [ "log", "netlink-packet-core 0.8.2", "netlink-sys 0.8.8", - "thiserror 2.0.20", + "thiserror 2.0.21", ] [[package]] @@ -2536,7 +2529,7 @@ dependencies = [ "rustc-hash", "rustls", "socket2", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-stream", "tracing", @@ -2566,7 +2559,7 @@ dependencies = [ "rustls-platform-verifier", "slab", "sorted-index-buffer", - "thiserror 2.0.20", + "thiserror 2.0.21", "tinyvec", "tracing", "web-time", @@ -2862,7 +2855,7 @@ dependencies = [ "log", "rand 0.10.3", "sha2", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "windows", "windows-strings", @@ -3388,7 +3381,7 @@ dependencies = [ "rustix", "serde", "serde_json", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tracing", ] @@ -3461,7 +3454,7 @@ dependencies = [ "rds-observe", "rds-sync", "rustix", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-util", "tracing", @@ -3479,7 +3472,7 @@ dependencies = [ "proptest", "rand 0.10.3", "serde", - "thiserror 2.0.20", + "thiserror 2.0.21", "url", ] @@ -3496,7 +3489,7 @@ dependencies = [ "rds-core", "rds-net", "rds-sync", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tracing", "x11rb", @@ -3520,7 +3513,7 @@ dependencies = [ "rustls-pki-types", "serde", "serde_json", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-rustls", "url", @@ -3548,7 +3541,7 @@ dependencies = [ "rustix", "serde", "serde_json", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-stream", "tokio-util", @@ -3569,7 +3562,7 @@ dependencies = [ "serde", "serde_json", "subtle", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tracing", "tracing-subscriber", @@ -3595,7 +3588,7 @@ dependencies = [ "rustix", "rustls", "rustls-pki-types", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-util", "tracing", @@ -3627,7 +3620,7 @@ name = "rds-ssh" version = "0.1.0" dependencies = [ "russh", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", ] @@ -3646,7 +3639,7 @@ dependencies = [ "rds-net", "rustix", "serde", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "tokio-util", "tracing", @@ -3829,7 +3822,7 @@ dependencies = [ "ssh-encoding", "ssh-key", "subtle", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", "typenum", "universal-hash 0.6.1", @@ -3920,7 +3913,7 @@ checksum = "8bb47c2a50fdfdaf95b0ac8b12620fc327da1fd4adbb30d0c56d866b005873ff" dependencies = [ "rustls-cert-read", "rustls-pki-types", - "thiserror 2.0.20", + "thiserror 2.0.21", "tokio", ] @@ -3943,7 +3936,7 @@ dependencies = [ "reloadable-state", "rustls", "rustls-cert-read", - "thiserror 2.0.20", + "thiserror 2.0.21", ] [[package]] @@ -3970,9 +3963,9 @@ dependencies = [ [[package]] name = "rustls-platform-verifier" -version = "0.7.0" +version = "0.7.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "26d1e2536ce4f35f4846aa13bff16bd0ff40157cdb14cc056c7b14ba41233ba0" +checksum = "1167586491e2b18b8bfbb293e8180ec17c201c4f076d7cb3070ca964e7598f98" dependencies = [ "core-foundation", "core-foundation-sys", @@ -3991,9 +3984,9 @@ dependencies = [ [[package]] name = "rustls-platform-verifier-android" -version = "0.1.1" +version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" +checksum = "eec689c0bc40ff2458a5977b6619cb718087084a18e02a131c599b62d05e1a5f" [[package]] name = "rustls-webpki" @@ -4329,9 +4322,9 @@ dependencies = [ [[package]] name = "siphasher" -version = "1.0.3" +version = "1.0.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8ee5873ec9cce0195efcb7a4e9507a04cd49aec9c83d0389df45b1ef7ba2e649" +checksum = "33f4fe9184a62d842c9ef383018f3306d8ba224fd9d836f56d7288308847c256" [[package]] name = "slab" @@ -4341,9 +4334,9 @@ checksum = "0c790de23124f9ab44544d7ac05d60440adc586479ce501c1d6d7da3cd8c9cf5" [[package]] name = "smallvec" -version = "1.16.1" +version = "1.16.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ba467056f1b547ed52077911161fc86985becbc60e8e1857c8a144dab0def891" +checksum = "f9395f0f0eee849a9b707b2f06bb92a6a422090e2123bb2ef8e87a0e61892a8e" [[package]] name = "socket2" @@ -4597,11 +4590,11 @@ dependencies = [ [[package]] name = "thiserror" -version = "2.0.20" +version = "2.0.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ec86235f5fcc2a73650310756d2ac5b138a5780bbbdfae3eeccec992c435ba4f" +checksum = "09e52cb86a36cede5cb101bf8908837b3e4c6e5e59fe7fd85c23fb56200d189e" dependencies = [ - "thiserror-impl 2.0.20", + "thiserror-impl 2.0.21", ] [[package]] @@ -4617,9 +4610,9 @@ dependencies = [ [[package]] name = "thiserror-impl" -version = "2.0.20" +version = "2.0.21" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" +checksum = "fe5197923287db20a58125f0bc85c062f7f2c892de97b18c356f9efb14b28524" dependencies = [ "proc-macro2", "quote", @@ -4740,7 +4733,7 @@ dependencies = [ "rustls", "serde", "serde_json", - "thiserror 2.0.20", + "thiserror 2.0.21", "time", "tokio", "tokio-rustls", @@ -5115,9 +5108,9 @@ dependencies = [ [[package]] name = "wasm-bindgen" -version = "0.2.128" +version = "0.2.129" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "aecb87a33d3b0c5e3b7aa46336eaf486cffafbd281b195e4c8b80d50df2351bf" +checksum = "9bb54f33acc68fd454578d9820b0bde1a1a3d17aa17bb7b6595806d02886d409" dependencies = [ "cfg-if", "once_cell", @@ -5128,19 +5121,20 @@ dependencies = [ [[package]] name = "wasm-bindgen-futures" -version = "0.4.78" +version = "0.4.79" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6ef4c5d3d2cdf5c54f4231181768f5510842e350db025faf1f7163b1030ed928" +checksum = "3cbab34de2d982e9b48e18d216d04c4a6f641066ff19ffb699980f591ee3610e" dependencies = [ "js-sys", + "tokio", "wasm-bindgen", ] [[package]] name = "wasm-bindgen-macro" -version = "0.2.128" +version = "0.2.129" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a690d511e3c1a8b3a55e33511e3c2c00c78415cd23650f32b808627f5696b9ed" +checksum = "2e29d0c35b16e224a7eeb5cd2d25e3e1968fbd65604117b44d3b789d00ee8535" dependencies = [ "quote", "wasm-bindgen-macro-support", @@ -5148,9 +5142,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-macro-support" -version = "0.2.128" +version = "0.2.129" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "411e4887f0071ef2d2164a9d5fdf2d20efbef78fccd3a78b0c10a1dc5295e48a" +checksum = "6f501a8bc3719dba86ef8ae4728879c08001bea749eb1333ac5b91e040e2a6b7" dependencies = [ "bumpalo", "proc-macro2", @@ -5161,9 +5155,9 @@ dependencies = [ [[package]] name = "wasm-bindgen-shared" -version = "0.2.128" +version = "0.2.129" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "81941cd78d0c92026c33e5e01312845a4cb1e9af3407f9134b100dd03144103e" +checksum = "23f0c9c52aa7cd7d77769a4cfe2a9adb1b331f489a41d912ce14513d5ab995c6" dependencies = [ "unicode-ident", ] @@ -5183,9 +5177,9 @@ dependencies = [ [[package]] name = "web-sys" -version = "0.3.105" +version = "0.3.106" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9fbddc4a036f00ec4f18c83445bd3115cb306a91da554919a099d9222fe4a7f8" +checksum = "88261b9deccee56594c11a3460c462c41f58d148598fe70ad77070126a68aba4" dependencies = [ "js-sys", "wasm-bindgen", @@ -5560,7 +5554,7 @@ dependencies = [ "futures", "log", "serde", - "thiserror 2.0.20", + "thiserror 2.0.21", "windows", "windows-core", ] @@ -5596,7 +5590,7 @@ dependencies = [ "pharos", "rustc_version", "send_wrapper", - "thiserror 2.0.20", + "thiserror 2.0.21", "wasm-bindgen", "wasm-bindgen-futures", "web-sys", @@ -5633,7 +5627,7 @@ dependencies = [ "oid-registry", "ring", "rusticata-macros", - "thiserror 2.0.20", + "thiserror 2.0.21", "time", ] @@ -5687,18 +5681,18 @@ dependencies = [ [[package]] name = "zerocopy" -version = "0.8.57" +version = "0.8.59" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d35102a9f36d089ccae9e4c6802bc118be4487b80aaffc0ab4e0cf5ce92d2873" +checksum = "6df92bf3d9227be3d53173901ddbffac2babc27ae50f397776ffd6dc33f800cb" dependencies = [ "zerocopy-derive", ] [[package]] name = "zerocopy-derive" -version = "0.8.57" +version = "0.8.59" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "146c01f5ab44258da43cf276c74a2763db2ff3969c9c652c3f2de07041d0b2bc" +checksum = "ac4f328cf2f05d084e496c3e9c3f33ed0a183656a16e1fcec4d464d8373aec82" dependencies = [ "proc-macro2", "quote", diff --git a/Cargo.toml b/Cargo.toml index 7b9323c..1cc6fb5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -40,7 +40,6 @@ repository = "https://github.com/NDDev-OpenNetwork/remote-device-sync" [workspace.dependencies] anyhow = "1" -async-trait = "0.1" blake3 = "1" bytes = "1" clap = { version = "4", features = ["derive"] } @@ -61,8 +60,13 @@ rand = "0.10" redb = "4.3" rcgen = { version = "0.14", default-features = false, features = ["crypto", "pem", "ring"] } rustix = { version = "1", features = ["fs"] } +# Pinned exactly: the release carries PTY encoding, sensitive-debug +# redaction, KEX/parser fixes and stalled-write timeout enforcement that +# the SSH surface relies on (docs/ssh.md). russh = { version = "=0.63.3", default-features = false, features = ["ring"] } -rustls = { version = "0.23", default-features = false, features = ["std", "logging", "tls12", "ring"] } +# No tls12: every TLS peer in the estate is our own rustls binary and +# QUIC is 1.3-only anyway — 1.2 has no legitimate client here. +rustls = { version = "0.23", default-features = false, features = ["std", "logging", "ring"] } rustls-pki-types = { version = "1", features = ["alloc"] } rds-agent = { path = "crates/rds-agent" } rds-audio = { path = "crates/rds-audio" } @@ -83,7 +87,7 @@ thiserror = "2" tokio = { version = "1", features = ["full"] } tokio-stream = "0.1" tokio-util = { version = "0.7.19", features = ["rt"] } -tokio-rustls = { version = "0.26", default-features = false, features = ["ring", "tls12"] } +tokio-rustls = { version = "0.26", default-features = false, features = ["ring"] } tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } url = "2" diff --git a/crates/rds-core/Cargo.toml b/crates/rds-core/Cargo.toml index 427aaf9..416134c 100644 --- a/crates/rds-core/Cargo.toml +++ b/crates/rds-core/Cargo.toml @@ -8,7 +8,6 @@ repository.workspace = true [features] default = [] -desktop = [] [dependencies] blake3.workspace = true diff --git a/crates/rds-discovery/src/records.rs b/crates/rds-discovery/src/records.rs index 8fa6a45..cbc0282 100644 --- a/crates/rds-discovery/src/records.rs +++ b/crates/rds-discovery/src/records.rs @@ -232,7 +232,7 @@ impl FileStore { } #[cfg(test)] let io_fault = std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)); - let db = Database::builder() + let mut db = Database::builder() .set_cache_size(16 * 1024 * 1024) .create_with_backend(BoundedFile { file, @@ -327,6 +327,11 @@ impl FileStore { if !anchor.initialized()? { anchor.seal()?; } + // Reclaim pages freed by expired/deleted records — without periodic + // compaction the file only ever grows. Maintenance, not an + // invariant: a failed compact must not refuse a valid store, and + // this crate carries no logging facade to report it through. + let _ = db.compact(); let observed = Published::new(record_snapshot(&metadata, true)); Ok(Self { observed, From 991cbcc63e6eb18adf275859e5367bdc349e7197 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 19:57:46 +0500 Subject: [PATCH 09/10] fix(ops,docs): admit session lifecycle events; correct relay admission docs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - vector.toml safe transform admitted only 12 of the 16 event names the pipeline emits — session_opened, session_closed, path_migrated and request_refused were silently dropped by the assert; add them plus the typed reason field they carry, with allowlist tests - rds-server.service still described --allow as "empty = open relay"; since the relay now fails closed, document the required allowlist and the --development-open-relay opt-in - rds-agent.service gains MemoryDenyWriteExecute for parity with the server unit (no JIT anywhere in the binary) - observability.md field table gains reason; deployment.md records that --directory-allow also accepts lowercase hex keys Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- deploy/systemd/rds-agent.service | 1 + deploy/systemd/rds-server.service | 7 +++++-- docs/deployment.md | 5 +++-- docs/observability.md | 2 +- ops/observability/vector.toml | 24 +++++++++++++++++++++++- 5 files changed, 33 insertions(+), 6 deletions(-) diff --git a/deploy/systemd/rds-agent.service b/deploy/systemd/rds-agent.service index 6e4f639..cb3c879 100644 --- a/deploy/systemd/rds-agent.service +++ b/deploy/systemd/rds-agent.service @@ -53,6 +53,7 @@ ProtectHostname=yes RestrictSUIDSGID=yes RestrictRealtime=yes LockPersonality=yes +MemoryDenyWriteExecute=yes # AF_NETLINK: interface/netmon watches (path migration on link change). RestrictAddressFamilies=AF_INET AF_INET6 AF_UNIX AF_NETLINK SystemCallFilter=@system-service diff --git a/deploy/systemd/rds-server.service b/deploy/systemd/rds-server.service index 5a48b46..de3da9d 100644 --- a/deploy/systemd/rds-server.service +++ b/deploy/systemd/rds-server.service @@ -22,8 +22,11 @@ ExecStart=/usr/local/bin/rds-server \ # --directory-allow (repeat per enrolled publisher; # empty denies directory record publish/fetch/delete) # --registry-key --registry /var/lib/rds/registry.json -# --allow … (empty = open relay — do NOT expose open -# relays on a public IP beyond short-lived bring-up) +# --allow … (relay packet admission — REQUIRED in +# production: an empty allowlist fails closed at startup. Only +# --development-open-relay opts into open admission, and it conflicts +# with --allow; it exists for isolated development, never for a +# public-facing relay) # Native directory TLS on --http-addr (docs/deployment.md): # --directory-tls-cert /var/lib/rds/directory-chain.pem # --directory-tls-key /var/lib/rds/directory-key.pem diff --git a/docs/deployment.md b/docs/deployment.md index fafca40..e0e94c8 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -150,8 +150,9 @@ or HTTP fallback. The registry verifying key remains separately provisioned; a TLS certificate cannot authorize a name binding. Directory membership is separately provisioned with repeated -`--directory-allow ` arguments (maximum 4096 distinct, -nonweak Ed25519 keys). An empty list denies record PUT, GET and DELETE, while +`--directory-allow ` arguments — base32 is the canonical form, +lowercase hexadecimal is also accepted (maximum 4096 distinct, nonweak +Ed25519 keys). An empty list denies record PUT, GET and DELETE, while health and configured policy routes remain available. This is independent of relay `--allow` and agent peer permissions. Removing a key and restarting denies its record access while preserving its replay floor; it does not revoke an diff --git a/docs/observability.md b/docs/observability.md index c811a39..72953ba 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -103,7 +103,7 @@ Each JSON record contains: | `run_id`, `sequence` | random 128-bit process ID and local event sequence; concurrent arrival order may differ | | `session_id` | numeric session ID minted by the accepting agent or the dialing client on its `rds.conn` span, scoped to `run_id`; never a peer key or distributed trace ID | | `level`, `target`, `line` | static source metadata | -| `event`, `operation`, `outcome`, `elapsed_us` | allowlisted operational fields; absent fields are null | +| `event`, `operation`, `outcome`, `reason`, `elapsed_us` | allowlisted operational fields; absent fields are null | | `telemetry_*_total` | cumulative queue rejection, oversized-record and output I/O error counts | No key, credential, address, user path, command, file content or remote output diff --git a/ops/observability/vector.toml b/ops/observability/vector.toml index 5a53c58..c791564 100644 --- a/ops/observability/vector.toml +++ b/ops/observability/vector.toml @@ -27,12 +27,13 @@ r = object!(parse_json!(string!(.message))) assert!(r.schema_version == 1) assert!(includes(["rds-agent", "rds-cli", "rds-relay", "rds-server", "rds-bench"], r.service)) assert!(includes(["TRACE", "DEBUG", "INFO", "WARN", "ERROR"], r.level)) -assert!(includes(["diagnostic", "process_started", "process_completed", "process_failed", "heartbeat", "listener_ready", "peer_accepted", "peer_rejected", "handshake_failed", "handshake_timed_out", "connection_budget_exhausted", "operation_completed"], r.event)) +assert!(includes(["diagnostic", "process_started", "process_completed", "process_failed", "heartbeat", "listener_ready", "peer_accepted", "peer_rejected", "handshake_failed", "handshake_timed_out", "connection_budget_exhausted", "operation_completed", "session_opened", "session_closed", "path_migrated", "request_refused"], r.event)) assert!(match(string!(r.run_id), r'^[0-9a-f]{32}$')) assert!(match(string!(r.target), r'^[a-zA-Z0-9_:.-]{1,128}$')) assert!(match(string!(r.version), r'^[0-9]+\.[0-9]+\.[0-9]+([+-][a-zA-Z0-9.-]+)?$')) assert!(includes([null, "connect", "service_stream", "ssh_connect", "ssh_session", "grant_authorize", "grant_renew", "sync_send", "sync_recv"], r.operation)) assert!(includes([null, "ok", "error", "cancelled"], r.outcome)) +assert!(includes([null, "completed", "cancelled", "timeout", "denied", "expired", "revoked", "budget_exhausted", "peer_closed", "local_closed", "reset", "transport", "protocol", "unsupported", "policy_unavailable", "aborted"], r.reason)) assert!(is_null(r.session_id) || is_integer(r.session_id)) assert!(is_null(r.elapsed_us) || is_integer(r.elapsed_us)) assert!(is_null(r.line) || is_integer(r.line)) @@ -55,6 +56,7 @@ assert!(match(node, r'^[a-zA-Z0-9_-]{1,32}$')) "run_id": r.run_id, "sequence": sequence, "level": r.level, "target": r.target, "line": r.line, "session_id": r.session_id, "event": r.event, "operation": r.operation, "outcome": r.outcome, + "reason": r.reason, "elapsed_us": r.elapsed_us, "node": node, "telemetry_dropped_total": dropped, "telemetry_oversize_total": oversize, "telemetry_write_errors_total": write_errors @@ -283,3 +285,23 @@ extract_from = "safe" [[tests.outputs.conditions]] type = "vrl" source = "assert!(.operation == \"sync_recv\")\nassert!(.outcome == \"cancelled\")\nassert!(!exists(.path) && !exists(.directory) && !exists(.key))" + +[[tests]] +name = "session lifecycle events carry typed reasons" +[[tests.inputs]] +insert_at = "safe" +type = "log" +log_fields.message = "{\"schema_version\": 1, \"timestamp_unix_us\": 1790300000000000, \"uptime_us\": 100, \"service\": \"rds-agent\", \"version\": \"0.1.0\", \"run_id\": \"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\", \"sequence\": 1, \"level\": \"WARN\", \"target\": \"rds_telemetry\", \"line\": 1, \"session_id\": 7, \"event\": \"request_refused\", \"reason\": \"denied\", \"telemetry_dropped_total\": 0, \"telemetry_oversize_total\": 0, \"telemetry_write_errors_total\": 0, \"peer\": \"PRIVATE_SENTINEL\"}" +[[tests.outputs]] +extract_from = "safe" +[[tests.outputs.conditions]] +type = "vrl" +source = "assert!(.event == \"request_refused\")\nassert!(.reason == \"denied\")\nassert!(!exists(.peer))" + +[[tests]] +name = "unknown reason text is rejected" +no_outputs_from = ["safe"] +[[tests.inputs]] +insert_at = "safe" +type = "log" +log_fields.message = "{\"schema_version\": 1, \"timestamp_unix_us\": 1790300000000000, \"uptime_us\": 100, \"service\": \"rds-agent\", \"version\": \"0.1.0\", \"run_id\": \"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa\", \"sequence\": 1, \"level\": \"WARN\", \"target\": \"rds_telemetry\", \"line\": 1, \"session_id\": 7, \"event\": \"session_closed\", \"reason\": \"PRIVATE_SENTINEL\", \"telemetry_dropped_total\": 0, \"telemetry_oversize_total\": 0, \"telemetry_write_errors_total\": 0}" From 7c490aaee544b28e20f3f573b8b35260ae52f015 Mon Sep 17 00:00:00 2001 From: rldyourmnd Date: Sun, 27 Sep 2026 20:15:20 +0500 Subject: [PATCH 10/10] fix(agent): satisfy CI clippy lanes - effective_services filter used a literal-bool match that newer clippy folds into matches! - sys.rs open_fds qualified std::mem::size_of, which is prelude in edition 2024 (the arm is macOS-cfg so only the macOS lane saw it) Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- crates/rds-agent/src/lib.rs | 7 +++---- crates/rds-agent/src/sys.rs | 2 +- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/crates/rds-agent/src/lib.rs b/crates/rds-agent/src/lib.rs index 42d4862..8005c3f 100644 --- a/crates/rds-agent/src/lib.rs +++ b/crates/rds-agent/src/lib.rs @@ -214,10 +214,9 @@ impl AgentPolicy { // this binary can actually serve it, matching the implicit arm // and the directory announcement. `validate` still rejects the // flag combination at startup. - Some(explicit) => set.extend(explicit.iter().copied().filter(|k| match k { - ServiceKind::Tcp | ServiceKind::Sync => true, - ServiceKind::Desktop => cfg!(feature = "desktop"), - _ => false, + Some(explicit) => set.extend(explicit.iter().copied().filter(|k| { + matches!(k, ServiceKind::Tcp | ServiceKind::Sync) + || (matches!(k, ServiceKind::Desktop) && cfg!(feature = "desktop")) })), None => { set.insert(ServiceKind::Tcp); diff --git a/crates/rds-agent/src/sys.rs b/crates/rds-agent/src/sys.rs index b934362..a414d2c 100644 --- a/crates/rds-agent/src/sys.rs +++ b/crates/rds-agent/src/sys.rs @@ -30,7 +30,7 @@ pub(crate) fn open_fds() -> Option { if size <= 0 { return None; } - Some(size as usize / std::mem::size_of::()) + Some(size as usize / size_of::()) } #[cfg(not(any(target_os = "linux", target_os = "macos")))]