From 3e926fa9d40ffbc5267efbbcbcf953f3a0388860 Mon Sep 17 00:00:00 2001 From: Benjamin Saunders Date: Mon, 20 Apr 2026 20:01:24 -0700 Subject: [PATCH] feat(noq): Introduce noq::Connection::authenticated --- noq/src/connection.rs | 31 ++++++++++++++++++++++++++++++- 1 file changed, 30 insertions(+), 1 deletion(-) diff --git a/noq/src/connection.rs b/noq/src/connection.rs index 9466ccb0d..43f650805 100644 --- a/noq/src/connection.rs +++ b/noq/src/connection.rs @@ -741,6 +741,33 @@ impl Connection { .map(|c| c.clone_box()) } + /// Succeeds when an incoming connection is proven not to be a replay attack. + /// + /// Only interesting for `Connection`s obtained from [`Connecting::into_0rtt`]. On 1-RTT + /// connections, always completes immediately. Contrast + /// [`handshake_confirmed`](Self::handshake_confirmed), which waits longer on clients e.g. to + /// confirm client authentication. + /// + /// For incoming connections, reads from [`RecvStream`]s are guaranteed not to arise from replay + /// attacks after this succeeds, even for streams accepted or read during 0-RTT. For outgoing + /// connections, streams opened after this succeeds will never be discarded by the server due to + /// 0-RTT rejection. + pub async fn authenticated(&self) -> Result<(), ConnectionError> { + let notified = { + let conn = self.0.state.lock("connected"); + if let Some(e) = &conn.error { + return Err(e.clone()); + } + if conn.connected { + return Ok(()); + } + self.0.shared.connected.notified() + }; + notified.await; + let conn = self.0.state.lock("connected"); + conn.error.clone().map_or(Ok(()), Err) + } + /// Parameters negotiated during the handshake /// /// Guaranteed to return `Some` on fully established connections or after @@ -1396,6 +1423,7 @@ pub(crate) struct Shared { datagram_received: Notify, datagrams_unblocked: Notify, closed: Notify, + connected: Notify, /// Number of live handles that can be used to initiate or handle I/O; excludes the driver ref_count: AtomicUsize, } @@ -1598,6 +1626,7 @@ impl State { } Connected => { self.connected = true; + shared.connected.notify_waiters(); if let Some(x) = self.on_connected.take() { // We don't care if the on-connected future was dropped let _ = x.send(self.inner.accepted_0rtt()); @@ -1774,7 +1803,6 @@ impl State { shared.handshake_confirmed.notify_waiters(); wake_all_notify(&mut self.stopped); shared.closed.notify_waiters(); - // Send to the registered on_closed futures. if !self.on_closed.is_empty() { let closed = Closed::new(self, reason); @@ -1782,6 +1810,7 @@ impl State { tx.send(closed.clone()).ok(); } } + shared.connected.notify_waiters(); } fn close(&mut self, error_code: VarInt, reason: Bytes, shared: &Shared) {