From e3c3a73e6df4c37c7d2c932e1a2b4d2e1fe8afb6 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Sat, 15 Aug 2026 14:02:01 +0000 Subject: [PATCH] feat(sync): make outbound NIP-42 authentication optional and terminal on refusal Outbound sync answers NIP-42 challenges with the relay owner key on every connection. When no owner key is available, the previous code panicked at registration (`.expect`), and the retry machinery would still have reserved a one-shot authentication retry that nothing could ever fulfil. Approach: `RelayConnection::new{,_with_database}` now take `Option` and only attach the SDK authenticator when present, exposing `answers_auth_challenges()`. Without an authenticator an auth-required CLOSED is terminal like any other CLOSED: the terminal listener retires the subscription, the data lane releases its live permit, and `handle_subscription_closed` skips the one-retry reservation and goes straight to retirement plus the AuthenticationRequired policy refusal (24h probe). `register_relay` degrades gracefully to an unauthenticated connection with a warning instead of panicking. Correctness assumptions: rust-nostr only retains auth-refused subscriptions for post-authentication resubscription when an authenticator is configured, so every has_authenticator branch mirrors an SDK behavior split; a reserved retry without an authenticator would dangle until disconnect cleanup. Test infrastructure: new AuthGatingRelay helper - a NIP-42 gate in front of a backend relay that serves a plain NIP-11 document, challenges every session, refuses queries pre-auth, marks negentropy unsupported, and either bridges (Admit) or answers `restricted:` (Restricted) after a valid AUTH, recording authenticated pubkeys and REQ counts. Validation: integration tests prove (1) a public instance authenticates to a gated ordinary relay with its owner key and the retained subscription is answered after AUTH (announcement reaches purgatory through the gate, the instance's only event source), and (2) a restricted refusal after successful authentication parks the work - the gate's REQ count holds still for a full 2s observation window. cargo test --lib (789 passed) and --test sync sync::outbound_auth pass. --- src/sync/mod.rs | 35 ++- src/sync/relay_connection.rs | 100 ++++--- tests/common/auth_gating_relay.rs | 429 ++++++++++++++++++++++++++++++ tests/common/mod.rs | 2 + tests/common/relay.rs | 20 ++ tests/sync/outbound_auth.rs | 115 +++++++- 6 files changed, 656 insertions(+), 45 deletions(-) create mode 100644 tests/common/auth_gating_relay.rs diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 4f50ddd..c18476b 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -5332,11 +5332,20 @@ impl SyncManager { } } - // Get relay owner keys for NIP-42 authentication - let keys = self - .config - .relay_owner_keys() - .expect("relay_owner_keys should be available"); + // The relay owner key answers outbound NIP-42 challenges. Sync + // must keep working without it, so a missing key degrades to an + // unauthenticated connection instead of aborting registration. + let keys = match self.config.relay_owner_keys() { + Ok(keys) => Some(keys), + Err(error) => { + tracing::warn!( + relay = %relay_url, + error = %error, + "Relay owner key unavailable; outbound NIP-42 authentication disabled for this connection" + ); + None + } + }; let connection = RelayConnection::new_with_database( relay_url.clone(), @@ -7852,11 +7861,17 @@ impl SyncManager { ); return; } - if reserve_authentication_retry( - &mut self.auth_required_attempts, - relay_url, - &subscription_id, - ) { + // A retry is only worth reserving when the SDK will actually + // answer the challenge and resubscribe; without an authenticator + // the subscription is already gone and a reserved retry would + // dangle until disconnect cleanup. + if connection.answers_auth_challenges() + && reserve_authentication_retry( + &mut self.auth_required_attempts, + relay_url, + &subscription_id, + ) + { tracing::info!( relay = %relay_url, sub_id = %subscription_id, diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index b961e08..6e19181 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -588,6 +588,10 @@ pub struct RelayConnection { policy: OutboundTargetPolicy, /// The underlying nostr-sdk client client: Client, + /// Whether a NIP-42 authenticator was attached to the client. Without + /// one the SDK never answers AUTH challenges and never retains + /// auth-refused subscriptions for a post-authentication retry. + has_authenticator: bool, /// Local database for negentropy comparison (used for NIP-77 sync) database: Option, /// Whether we've logged NIP-77 not supported for this relay (log once) @@ -712,24 +716,29 @@ impl RelayConnection { /// /// # Arguments /// * `url` - The relay URL to connect to (with or without scheme, e.g., "relay.example.com" or "wss://relay.example.com") - /// * `keys` - Cryptographic keys for NIP-42 authentication (typically the relay operator's keys) + /// * `keys` - Keys for answering NIP-42 challenges (typically the relay + /// operator's keys); `None` disables outbound authentication entirely /// * `source` - Whether the URL is operator-configured or event-directed /// * `policy` - Outbound target policy enforced before event-directed dials pub fn new( url: String, - keys: Keys, + keys: Option, source: RelayTargetSource, policy: OutboundTargetPolicy, ) -> Self { let normalized_url = Self::normalize_url(&url); - let client = Client::builder() - .authenticator(SignerAuthenticator::new(keys)) - .build(); + let has_authenticator = keys.is_some(); + let mut builder = Client::builder(); + if let Some(keys) = keys { + builder = builder.authenticator(SignerAuthenticator::new(keys)); + } + let client = builder.build(); Self { url: normalized_url, source, policy, client, + has_authenticator, database: None, nip77_warning_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)), nip77_supported: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)), @@ -778,25 +787,30 @@ impl RelayConnection { /// # Arguments /// * `url` - The relay URL to connect to (with or without scheme, e.g., "relay.example.com" or "wss://relay.example.com") /// * `database` - Shared database for local event comparison during negentropy sync - /// * `keys` - Cryptographic keys for NIP-42 authentication (typically the relay operator's keys) + /// * `keys` - Keys for answering NIP-42 challenges (typically the relay + /// operator's keys); `None` disables outbound authentication entirely /// * `source` - Whether the URL is operator-configured or event-directed /// * `policy` - Outbound target policy enforced before event-directed dials pub fn new_with_database( url: String, database: SharedDatabase, - keys: Keys, + keys: Option, source: RelayTargetSource, policy: OutboundTargetPolicy, ) -> Self { let normalized_url = Self::normalize_url(&url); - let client = Client::builder() - .authenticator(SignerAuthenticator::new(keys)) - .build(); + let has_authenticator = keys.is_some(); + let mut builder = Client::builder(); + if let Some(keys) = keys { + builder = builder.authenticator(SignerAuthenticator::new(keys)); + } + let client = builder.build(); Self { url: normalized_url, source, policy, client, + has_authenticator, database: Some(database), nip77_warning_logged: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)), nip77_supported: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(0)), @@ -1700,7 +1714,13 @@ impl RelayConnection { message, } => { let subscription_id = subscription_id.into_owned(); - if !is_auth_required_message(&message) { + // An auth-required CLOSED is only a retry signal + // when an authenticator can actually answer the + // challenge; otherwise it is terminal like any + // other CLOSED. + if !terminal_connection.has_authenticator + || !is_auth_required_message(&message) + { terminal_connection .retire_peer_closed_subscription(&subscription_id) .await; @@ -1801,19 +1821,22 @@ impl RelayConnection { // this processor-facing path retains live restoration. let subscription_id = subscription_id.into_owned(); // rust-nostr needs the same subscription and ledger - // slot alive while it answers a first NIP-42 challenge. - let released_live = if is_auth_required_message(&msg) { - self.live_req_permits_held - .lock() - .expect("live permit map poisoned") - .get(&subscription_id) - .map(|held| ReleasedLiveSubscription { - generation: held.generation, - filter_count: held.filters.len(), - }) - } else { - self.release_live_req_permit(&subscription_id) - }; + // slot alive while it answers a first NIP-42 + // challenge. Without an authenticator there is no + // challenge to answer, so release like any CLOSED. + let released_live = + if self.has_authenticator && is_auth_required_message(&msg) { + self.live_req_permits_held + .lock() + .expect("live permit map poisoned") + .get(&subscription_id) + .map(|held| ReleasedLiveSubscription { + generation: held.generation, + filter_count: held.filters.len(), + }) + } else { + self.release_live_req_permit(&subscription_id) + }; if is_query_rate_limit_message(&msg) { self.record_query_rate_limit(); } @@ -2211,6 +2234,15 @@ impl RelayConnection { &self.url } + /// Whether this connection can answer NIP-42 AUTH challenges. + /// + /// Without an authenticator the SDK removes auth-refused subscriptions + /// instead of retaining them, so no post-authentication retry can ever + /// happen and callers must treat auth-required CLOSED as terminal. + pub fn answers_auth_challenges(&self) -> bool { + self.has_authenticator + } + /// Get the number of active subscriptions on this connection /// /// Returns the count of subscriptions tracked by the underlying nostr-sdk client. @@ -2819,7 +2851,7 @@ mod tests { fn permissive_connection(url: &str, keys: Keys) -> RelayConnection { RelayConnection::new( url.to_string(), - keys, + Some(keys), RelayTargetSource::EventDirected, OutboundTargetPolicy { allow_non_global: true, @@ -3067,7 +3099,7 @@ mod tests { // the strict default policy, mirroring the bootstrap-relay exception. let connection = RelayConnection::new( configured.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3125,7 +3157,7 @@ mod tests { let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3172,7 +3204,7 @@ mod tests { relay.run().await.expect("start empty relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3219,7 +3251,7 @@ mod tests { relay.run().await.expect("start registry relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3254,7 +3286,7 @@ mod tests { relay.run().await.expect("start query-limited relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3417,7 +3449,7 @@ mod tests { async fn event_directed_connect_rejects_loopback_before_dialling() { let connection = RelayConnection::new( "ws://127.0.0.1:1".to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::EventDirected, OutboundTargetPolicy::default(), ); @@ -3759,7 +3791,7 @@ mod tests { relay.run().await.expect("start local relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3822,7 +3854,7 @@ mod tests { relay.run().await.expect("start local relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3914,7 +3946,7 @@ mod tests { relay.run().await.expect("start local relay"); let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); diff --git a/tests/common/auth_gating_relay.rs b/tests/common/auth_gating_relay.rs new file mode 100644 index 0000000..647ba23 --- /dev/null +++ b/tests/common/auth_gating_relay.rs @@ -0,0 +1,429 @@ +//! NIP-42 Gating Relay for Outbound Authentication Tests +//! +//! A WebSocket front-end that demands NIP-42 authentication before doing +//! anything, standing in for an authenticated (but NOT GRASP-08) relay: +//! +//! - HTTP requests with `Accept: application/nostr+json` receive a minimal +//! NIP-11 document without a `supported_grasps` field, so a syncing +//! instance treats the gate as an ordinary relay. +//! - Every WebSocket session is greeted with `["AUTH", ]`. Until +//! a valid AUTH event (kind 22242, verified signature, matching challenge +//! tag) arrives, `REQ`/`COUNT` receive an `auth-required:` CLOSED, `EVENT` +//! an `auth-required:` OK-false, and `NEG-OPEN` a `NEG-ERR` marking +//! negentropy unsupported (so sync falls back to plain REQs). +//! - After a valid AUTH the session either bridges transparently to the +//! backend relay ([`GateMode::Admit`]) or keeps answering every `REQ` with +//! a `restricted:` CLOSED ([`GateMode::Restricted`]). +//! +//! Authenticated pubkeys and the number of `REQ` frames received are shared +//! observable state for test assertions. + +use std::collections::HashSet; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; + +use futures_util::stream::{SplitSink, SplitStream}; +use futures_util::{SinkExt, StreamExt}; +use http_body_util::Full; +use hyper::body::Bytes; +use hyper::header::{ACCEPT, CONNECTION, SEC_WEBSOCKET_ACCEPT, SEC_WEBSOCKET_KEY, UPGRADE}; +use hyper::server::conn::http1; +use hyper::service::service_fn; +use hyper::upgrade::Upgraded; +use hyper::{Request, Response, StatusCode}; +use hyper_util::rt::TokioIo; +use nostr_sdk::prelude::{Event, Keys, Kind, PublicKey}; +use tokio::net::TcpListener; +use tokio::sync::oneshot; +use tokio_tungstenite::tungstenite::protocol::Role; +use tokio_tungstenite::tungstenite::Message; +use tokio_tungstenite::WebSocketStream; + +/// What an authenticated session is allowed to do. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum GateMode { + /// Bridge authenticated sessions transparently to the backend relay. + Admit, + /// Accept valid authentication but refuse every query with `restricted:`. + Restricted, +} + +#[derive(Clone)] +struct GateState { + authenticated: Arc>>, + req_count: Arc, +} + +/// NIP-42 gate in front of a backend relay. See the module docs. +pub struct AuthGatingRelay { + url: String, + state: GateState, + shutdown_tx: Option>, + handle: Option>, +} + +impl AuthGatingRelay { + /// Start the gate on a random loopback port in front of `backend_url`. + pub async fn start(backend_url: &str, mode: GateMode) -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("AuthGatingRelay failed to bind"); + let port = listener + .local_addr() + .expect("AuthGatingRelay local_addr") + .port(); + + let state = GateState { + authenticated: Arc::new(Mutex::new(HashSet::new())), + req_count: Arc::new(AtomicUsize::new(0)), + }; + let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>(); + + let backend_url = backend_url.to_string(); + let accept_state = state.clone(); + let handle = tokio::spawn(async move { + loop { + tokio::select! { + accepted = listener.accept() => { + let Ok((stream, _)) = accepted else { break }; + let backend_url = backend_url.clone(); + let state = accept_state.clone(); + let io = TokioIo::new(stream); + tokio::spawn(async move { + let service = service_fn(move |req| { + let backend_url = backend_url.clone(); + let state = state.clone(); + async move { handle_request(req, backend_url, mode, state).await } + }); + // Errors are expected when clients disconnect. + let _ = http1::Builder::new() + .serve_connection(io, service) + .with_upgrades() + .await; + }); + } + _ = &mut shutdown_rx => break, + } + } + }); + + Self { + url: format!("ws://127.0.0.1:{port}"), + state, + shutdown_tx: Some(shutdown_tx), + handle: Some(handle), + } + } + + /// The ws:// URL a syncing relay should use to reach the gate. + pub fn url(&self) -> &str { + &self.url + } + + /// Pubkeys that completed a valid NIP-42 authentication. + pub fn authenticated_pubkeys(&self) -> HashSet { + self.state + .authenticated + .lock() + .expect("authenticated set poisoned") + .clone() + } + + /// Total `REQ` frames received across all sessions. + pub fn req_count(&self) -> usize { + self.state.req_count.load(Ordering::SeqCst) + } + + /// Stop the gate. + pub async fn stop(mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + if let Some(handle) = self.handle.take() { + let _ = handle.await; + } + } +} + +impl Drop for AuthGatingRelay { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} + +async fn handle_request( + req: Request, + backend_url: String, + mode: GateMode, + state: GateState, +) -> Result>, hyper::Error> { + let is_websocket = req + .headers() + .get(UPGRADE) + .map(|v| v.to_str().unwrap_or("").eq_ignore_ascii_case("websocket")) + .unwrap_or(false); + + if is_websocket { + if let Some(key) = req + .headers() + .get(SEC_WEBSOCKET_KEY) + .and_then(|k| k.to_str().ok()) + .map(str::to_string) + { + let accept_key = derive_accept_key(key.as_bytes()); + tokio::spawn(async move { + match hyper::upgrade::on(req).await { + Ok(upgraded) => { + let ws = WebSocketStream::from_raw_socket( + TokioIo::new(upgraded), + Role::Server, + None, + ) + .await; + if let Err(error) = run_session(ws, &backend_url, mode, state).await { + eprintln!("AuthGatingRelay session ended: {error}"); + } + } + Err(error) => eprintln!("AuthGatingRelay upgrade error: {error}"), + } + }); + return Ok(Response::builder() + .status(StatusCode::SWITCHING_PROTOCOLS) + .header(CONNECTION, "upgrade") + .header(UPGRADE, "websocket") + .header(SEC_WEBSOCKET_ACCEPT, accept_key) + .body(Full::new(Bytes::new())) + .unwrap()); + } + } + + if req + .headers() + .get(ACCEPT) + .and_then(|value| value.to_str().ok()) + .is_some_and(|value| value.contains("application/nostr+json")) + { + // Deliberately no supported_grasps: the gate models an authenticated + // ordinary relay, not a GRASP-08 private service. + let document = serde_json::json!({ + "name": "auth gating relay", + "supported_nips": [1, 11, 42], + }); + return Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", "application/nostr+json") + .body(Full::new(Bytes::from(document.to_string()))) + .unwrap()); + } + + Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", "text/plain") + .body(Full::new(Bytes::from("AuthGatingRelay"))) + .unwrap()) +} + +type ClientWs = WebSocketStream>; +type SharedClientSink = Arc>>; + +async fn send_json(sink: &SharedClientSink, value: serde_json::Value) -> Result<(), String> { + sink.lock() + .await + .send(Message::Text(value.to_string().into())) + .await + .map_err(|e| format!("client write: {e}")) +} + +async fn run_session( + ws: ClientWs, + backend_url: &str, + mode: GateMode, + state: GateState, +) -> Result<(), String> { + let (client_tx, mut client_rx) = ws.split(); + let client_tx: SharedClientSink = Arc::new(tokio::sync::Mutex::new(client_tx)); + + // Challenge only needs to be unpredictable within the test process. + let challenge = Keys::generate().public_key().to_hex(); + send_json(&client_tx, serde_json::json!(["AUTH", challenge])).await?; + + let mut authenticated = false; + let mut backend: Option> = None; + let mut backend_task: Option> = None; + + while let Some(message) = client_rx.next().await { + let message = message.map_err(|e| format!("client read: {e}"))?; + match message { + Message::Text(text) => { + let Ok(frame) = serde_json::from_str::(text.as_str()) else { + continue; + }; + let Some(kind) = frame.get(0).and_then(|v| v.as_str()) else { + continue; + }; + + if kind == "AUTH" { + handle_auth(&frame, &challenge, &state, &mut authenticated, &client_tx).await?; + if authenticated && mode == GateMode::Admit && backend.is_none() { + let (backend_ws, _) = + tokio_tungstenite::connect_async(backend_url) + .await + .map_err(|e| format!("backend connect failed: {e}"))?; + let (backend_tx, backend_rx) = backend_ws.split(); + backend = Some(backend_tx); + backend_task = Some(spawn_backend_forwarder(backend_rx, client_tx.clone())); + } + continue; + } + + if kind == "REQ" { + state.req_count.fetch_add(1, Ordering::SeqCst); + } + + if authenticated && mode == GateMode::Admit { + if let Some(backend_tx) = backend.as_mut() { + backend_tx + .send(Message::Text(text)) + .await + .map_err(|e| format!("backend write: {e}"))?; + } + continue; + } + + // Pre-authentication, or authenticated in Restricted mode. + let refusal = if authenticated { + "restricted: not a member" + } else { + "auth-required: authentication required" + }; + let sub_id = frame.get(1).and_then(|v| v.as_str()).unwrap_or_default(); + match kind { + "REQ" | "COUNT" => { + send_json(&client_tx, serde_json::json!(["CLOSED", sub_id, refusal])) + .await?; + } + // The "not supported" wording makes sync classify + // negentropy as unsupported and fall back to REQs. + "NEG-OPEN" => { + send_json( + &client_tx, + serde_json::json!([ + "NEG-ERR", + sub_id, + "blocked: negentropy not supported" + ]), + ) + .await?; + } + "EVENT" => { + let id = frame + .get(1) + .and_then(|event| event.get("id")) + .and_then(|id| id.as_str()) + .unwrap_or_default(); + send_json(&client_tx, serde_json::json!(["OK", id, false, refusal])) + .await?; + } + _ => {} + } + } + Message::Ping(payload) => { + client_tx + .lock() + .await + .send(Message::Pong(payload)) + .await + .map_err(|e| format!("client write: {e}"))?; + } + Message::Close(_) => break, + _ => {} + } + } + + if let Some(task) = backend_task { + task.abort(); + } + Ok(()) +} + +async fn handle_auth( + frame: &serde_json::Value, + challenge: &str, + state: &GateState, + authenticated: &mut bool, + client_tx: &SharedClientSink, +) -> Result<(), String> { + let event = frame + .get(1) + .and_then(|value| Event::from_json(value.to_string()).ok()); + let Some(event) = event else { + return send_json( + client_tx, + serde_json::json!(["OK", "", false, "auth-required: malformed AUTH"]), + ) + .await; + }; + let challenge_matches = event.tags.iter().any(|tag| { + let values = tag.as_slice(); + values.first().is_some_and(|name| name == "challenge") + && values.get(1).is_some_and(|value| value == challenge) + }); + let valid = event.kind == Kind::Authentication && event.verify().is_ok() && challenge_matches; + if valid { + *authenticated = true; + state + .authenticated + .lock() + .expect("authenticated set poisoned") + .insert(event.pubkey); + } + send_json( + client_tx, + serde_json::json!([ + "OK", + event.id.to_hex(), + valid, + if valid { + "" + } else { + "auth-required: invalid AUTH" + } + ]), + ) + .await +} + +fn spawn_backend_forwarder( + mut backend_rx: SplitStream< + WebSocketStream>, + >, + client_tx: SharedClientSink, +) -> tokio::task::JoinHandle<()> { + tokio::spawn(async move { + while let Some(message) = backend_rx.next().await { + let Ok(message) = message else { break }; + if client_tx.lock().await.send(message).await.is_err() { + break; + } + } + }) +} + +/// Derive the Sec-WebSocket-Accept key from the request key. +fn derive_accept_key(request_key: &[u8]) -> String { + use bitcoin_hashes::sha1::Hash as Sha1Hash; + use bitcoin_hashes::{Hash, HashEngine}; + + const WS_GUID: &[u8] = b"258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; + + let mut engine = Sha1Hash::engine(); + engine.input(request_key); + engine.input(WS_GUID); + let hash = Sha1Hash::from_engine(engine); + base64::Engine::encode( + &base64::engine::general_purpose::STANDARD, + hash.as_byte_array(), + ) +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index e336a53..83baaef 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -2,6 +2,7 @@ #![allow(dead_code)] // Test helpers may not be used in all test configurations #![allow(unused_imports)] // Re-exports may not be used in all test configurations +pub mod auth_gating_relay; pub mod censoring_proxy; pub mod flapping_relay; pub mod git_server; @@ -16,6 +17,7 @@ pub mod setup_drop_relay; pub mod sync_helpers; pub mod upload_pack_counting_proxy; +pub use auth_gating_relay::{AuthGatingRelay, GateMode}; pub use git_server::{SimpleGitServer, SmartGitServer}; pub use mock_relay::MockRelay; pub use nip09_helpers::*; diff --git a/tests/common/relay.rs b/tests/common/relay.rs index 6f70dcb..f8d1c38 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -229,6 +229,26 @@ impl TestRelay { .await } + /// Start a syncing relay on a caller-reserved port with user-index + /// identity publication disabled. See + /// [`Self::start_with_sync_without_user_index`] and + /// [`Self::start_on_reservation_with_options`] for the two concerns + /// this combines. + pub async fn start_on_reservation_with_sync_without_user_index( + reservation: PortReservation, + bootstrap_relay_url: Option, + ) -> Self { + Self::start_internal( + reservation, + RelayOptions { + bootstrap_relay_url, + user_index_relays: Some(String::new()), + ..RelayOptions::default() + }, + ) + .await + } + /// Start a syncing relay with a caller-chosen relay-owner identity. /// /// Lets tests stage owner-signed events on other relays before this diff --git a/tests/sync/outbound_auth.rs b/tests/sync/outbound_auth.rs index 40f93c8..3b77ffe 100644 --- a/tests/sync/outbound_auth.rs +++ b/tests/sync/outbound_auth.rs @@ -6,10 +6,18 @@ //! - A PUBLIC instance recognizes a GRASP-08 private service from its NIP-11 //! document before dialing, and parks it without any WebSocket connection //! or AUTH exchange. +//! - A public instance answers NIP-42 challenges from an ordinary +//! authenticated relay with the relay owner key and syncs through it. +//! - A `restricted:` CLOSED after successful authentication is terminal: +//! subscription work parks via the policy-refusal machinery instead of +//! retrying. use std::time::Duration; -use crate::common::{wait_for_log_line, TestRelay}; +use crate::common::{ + reserve_port, send_to_relay_url, wait_for_log_line, AuthGatingRelay, GateMode, MockRelay, + TestRelay, +}; use nostr_sdk::prelude::*; /// The stable park warning emitted when a public instance excludes a @@ -68,3 +76,108 @@ async fn public_instance_parks_grasp08_relay_without_dialing() { syncing.stop().await; private_service.stop().await; } + +/// A public instance answers an ordinary relay's NIP-42 challenge with the +/// relay owner key and syncs through the authenticated session. +#[tokio::test] +async fn public_instance_authenticates_to_gated_relay_and_syncs() { + let maintainer = Keys::generate(); + let identifier = "outbound-auth-admit"; + + // The gate is the syncing relay's ONLY event source, so anything that + // reaches it must have crossed the authenticated bridge. + let backend = MockRelay::start().await; + let gate = AuthGatingRelay::start(backend.url(), GateMode::Admit).await; + + // The syncing relay's address must appear in the announcement before it + // boots, so its port is reserved up front. + let reservation = reserve_port(); + let syncing_domain = format!("127.0.0.1:{}", reservation.port()); + + // Announcement listing the syncing relay (in both clone and relays tags, + // for admission) plus the gate as the peer relay. + let npub = maintainer.public_key().to_bech32().expect("npub"); + let clone_url = format!("http://{syncing_domain}/{npub}/{identifier}.git"); + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom("clone", vec![clone_url]), + Tag::custom( + "relays", + vec![format!("ws://{syncing_domain}"), gate.url().to_string()], + ), + ]) + .finalize(&maintainer) + .expect("signed announcement"); + send_to_relay_url(backend.url(), &announcement) + .await + .expect("seed announcement on backend"); + + let syncing = TestRelay::start_on_reservation_with_sync_without_user_index( + reservation, + Some(gate.url().to_string()), + ) + .await; + + // The announcement is admitted to purgatory on the syncing relay - proof + // that a subscription refused pre-auth was answered after NIP-42. + let synced = wait_for_log_line(&syncing.log_path(), Duration::from_secs(60), |line| { + line.contains("Added announcement to purgatory") && line.contains(identifier) + }) + .await; + assert!( + synced, + "announcement must sync through the authenticated gate" + ); + assert!( + gate.authenticated_pubkeys() + .contains(&syncing.owner_keys().public_key()), + "syncing relay must authenticate with its owner key" + ); + + syncing.stop().await; + gate.stop().await; + backend.stop().await; +} + +/// A `restricted:` CLOSED after valid authentication parks subscription work +/// via the policy-refusal machinery: no retry storm follows. +#[tokio::test] +async fn restricted_after_authentication_is_terminal() { + let backend = MockRelay::start().await; + let gate = AuthGatingRelay::start(backend.url(), GateMode::Restricted).await; + let syncing = TestRelay::start_with_sync_without_user_index(Some(gate.url().to_string())).await; + + let refused = wait_for_log_line(&syncing.log_path(), Duration::from_secs(60), |line| { + line.contains("Relay policy refused subscription") && line.contains("restricted") + }) + .await; + assert!( + refused, + "restricted CLOSED must be recorded as a policy refusal" + ); + + // No retry storm: within a bounded settling deadline the gate's REQ count + // must hold still for one full 2s observation window. + let settle_deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let mut window_start = tokio::time::Instant::now(); + let mut last_count = gate.req_count(); + loop { + tokio::time::sleep(Duration::from_millis(100)).await; + let count = gate.req_count(); + if count != last_count { + last_count = count; + window_start = tokio::time::Instant::now(); + } else if window_start.elapsed() >= Duration::from_secs(2) { + break; + } + assert!( + tokio::time::Instant::now() < settle_deadline, + "REQ count kept growing after the restricted refusal (retry storm)" + ); + } + + syncing.stop().await; + gate.stop().await; + backend.stop().await; +}