diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 0f27c8b..4f50ddd 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1787,6 +1787,10 @@ enum ConnectAttemptOutcome { advertised_owner: Option, advertised_grasp08: bool, }, + /// The relay's NIP-11 advertises the GRASP-08 private-service extension + /// while this instance is public. Detected before the WebSocket dial, so + /// no connection or AUTH exchange ever happened. + PrivateService, Failed(String), } @@ -2497,6 +2501,13 @@ pub struct SyncManager { /// re-log the same forbidden target. Bounded by the set of distinct relay /// URLs in stored/purgatory events, which the indexes already carry. rejected_relay_targets: HashSet, + /// Relays whose NIP-11 advertises GRASP-08 while this instance is public. + /// + /// Laundering guard: a public mirror must not present private credentials + /// to such a relay nor hammer a service that will never admit it, so the + /// target is parked before the dial (no WebSocket, no AUTH exchange). + /// Held in memory only, re-probed at most once per process lifetime. + private_service_relays: HashSet, /// Last exact-ID dependency recovery attempt, used to bound retries. dependency_refetch_attempts: Arc>>, /// Events relays reported during negentropy reconciliation but failed to @@ -2636,6 +2647,7 @@ impl SyncManager { nip65_discovery_only_relays: HashSet::new(), pagination_sessions: HashMap::new(), rejected_relay_targets: HashSet::new(), + private_service_relays: HashSet::new(), dependency_refetch_attempts: Arc::new(std::sync::Mutex::new(HashMap::new())), missing_event_recovery: Arc::new(std::sync::Mutex::new( missing_events::MissingEventRecoveryIndex::default(), @@ -5276,6 +5288,16 @@ impl SyncManager { } }; + // A GRASP-08 private service is not a sync target for a public + // instance; never re-enter the connection lifecycle for it. + if !self.config.private_mode && self.private_service_relays.contains(&relay_url) { + tracing::trace!( + relay = %relay_url, + "Skipping registration of GRASP-08 private-service relay" + ); + return false; + } + // An ordinary sync registration upgrades a connection that was first // opened only for NIP-65 discovery. Discovery must never downgrade an // existing repository source. @@ -5404,6 +5426,13 @@ impl SyncManager { ); return; } + if !self.config.private_mode && self.private_service_relays.contains(&relay_url) { + tracing::debug!( + relay = %relay_url, + "Suppressing connection attempt for GRASP-08 private-service relay" + ); + return; + } let Some(result_tx) = self.connect_attempt_result_tx.clone() else { tracing::error!(relay = %relay_url, "Connection scheduler is not running"); return; @@ -5450,6 +5479,7 @@ impl SyncManager { let health_tracker = Arc::clone(&self.health_tracker); let semaphore = Arc::clone(&self.connect_attempt_semaphore); let timeout = self.health_tracker.base_backoff_secs(); + let private_mode = self.config.private_mode; let Some(mut shutdown_rx) = self.shutdown_tx.as_ref().map(|sender| sender.subscribe()) else { tracing::error!(relay = %relay_url, "Connection scheduler has no shutdown signal"); @@ -5464,18 +5494,29 @@ impl SyncManager { return; }; let outcome = tokio::select! { - result = connection.connect(timeout) => match result { - Ok(()) => { - let hints = connection.fetch_limit_hints().await; - ConnectAttemptOutcome::Connected { - advertised_default_limit: hints.default_limit, - advertised_max_subscriptions: hints.max_subscriptions, - advertised_owner: hints.owner, - advertised_grasp08: hints.grasp08, + outcome = async { + // A GRASP-08 private service must be recognized before + // the dial so no WebSocket or AUTH exchange ever reaches + // it. Only public instances park, so only they pay the + // extra pre-dial probe; session hints still come from + // the post-connect fetch below, which runs on every + // attempt so they stay per-session. + if !private_mode && connection.preflight_limit_hints().await.grasp08 { + return ConnectAttemptOutcome::PrivateService; + } + match connection.connect(timeout).await { + Ok(()) => { + let hints = connection.fetch_limit_hints().await; + ConnectAttemptOutcome::Connected { + advertised_default_limit: hints.default_limit, + advertised_max_subscriptions: hints.max_subscriptions, + advertised_owner: hints.owner, + advertised_grasp08: hints.grasp08, + } } - }, - Err(error) => ConnectAttemptOutcome::Failed(error), - }, + Err(error) => ConnectAttemptOutcome::Failed(error), + } + } => outcome, _ = shutdown_rx.recv() => { connection.disconnect().await; return; @@ -6299,6 +6340,31 @@ impl SyncManager { } self.handle_connect_or_reconnect(&result.relay_url).await; } + ConnectAttemptOutcome::PrivateService => { + if self.private_service_relays.insert(result.relay_url.clone()) { + tracing::warn!( + relay = %result.relay_url, + "Relay advertises GRASP-08 private service; excluding it from public sync" + ); + } + // Retire the target like `complete_ended_session`, but without + // the re-registration path: the park is permanent for this + // process. No connected gauge to decrement - we never dialed. + self.relay_sync_index + .write() + .await + .remove(&result.relay_url); + self.pending_sync_index + .write() + .await + .remove(&result.relay_url); + self.connections.remove(&result.relay_url); + self.nip65_discovery_only_relays.remove(&result.relay_url); + self.health_tracker.forget_relay(&result.relay_url); + if let Some(ref metrics) = self.metrics { + metrics.forget_relay(&result.relay_url); + } + } ConnectAttemptOutcome::Failed(error) => { if let Some(category) = naughty_list::NaughtyListTracker::classify_error(&error) { if let Some(ref naughty_list) = self.health_tracker.naughty_list() { diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 4c81879..b961e08 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -981,6 +981,26 @@ impl RelayConnection { parse_relay_limit_hints(&body) } + /// Fetch NIP-11 hints before the WebSocket dial. + /// + /// The pre-dial NIP-11 fetch is itself an outbound TCP connection, so it + /// must not bypass the SSRF gate: event-directed targets are authorized + /// first, and on rejection default hints are returned without any HTTP + /// request — the subsequent `connect()` then fails with the same policy + /// rejection through its own pre-dial check. + pub async fn preflight_limit_hints(&self) -> RelayLimitHints { + if self.source == RelayTargetSource::EventDirected + && self + .policy + .authorize_resolved(OutboundTargetKind::EventRelay, &self.url) + .await + .is_err() + { + return RelayLimitHints::default(); + } + self.fetch_limit_hints().await + } + /// Whether the SDK still considers this relay's WebSocket established. /// /// Connection setup performs bounded HTTP work after the handshake. The diff --git a/tests/common/relay.rs b/tests/common/relay.rs index 2200876..6f70dcb 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -212,6 +212,23 @@ impl TestRelay { .await } + /// Start a syncing relay with user-index identity publication disabled. + /// + /// `start_with_sync` points identity publication at the bootstrap relay. + /// Tests making log-based assertions about the bootstrap connection use + /// this variant so identity-publication traffic cannot confound them. + pub async fn start_with_sync_without_user_index(bootstrap_relay_url: Option) -> Self { + Self::start_internal( + port::reserve_port(), + 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/common/setup_drop_relay.rs b/tests/common/setup_drop_relay.rs index 1c23d4f..ab27b5e 100644 --- a/tests/common/setup_drop_relay.rs +++ b/tests/common/setup_drop_relay.rs @@ -1,5 +1,11 @@ //! Relay fixture that disconnects the WebSocket during its NIP-11 setup fetch. +//! +//! Public syncing instances also probe NIP-11 *before* dialing (to detect +//! GRASP-08 private services); that pre-dial probe is answered immediately. +//! Only a NIP-11 request arriving while a WebSocket session is live triggers +//! the drop-during-setup choreography under test. +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::time::Duration; @@ -24,6 +30,7 @@ impl SetupDropRelay { let (dropped_tx, dropped_rx) = watch::channel(false); let nip11_tx = Arc::new(nip11_tx); let dropped_tx = Arc::new(dropped_tx); + let active_websockets = Arc::new(AtomicUsize::new(0)); let (shutdown_tx, mut shutdown_rx) = oneshot::channel(); let handle = tokio::spawn(async move { @@ -35,6 +42,7 @@ impl SetupDropRelay { let nip11_rx = nip11_rx.clone(); let dropped_tx = dropped_tx.clone(); let dropped_rx = dropped_rx.clone(); + let active_websockets = active_websockets.clone(); tokio::spawn(async move { handle_connection( stream, @@ -42,6 +50,7 @@ impl SetupDropRelay { nip11_rx, dropped_tx, dropped_rx, + active_websockets, ) .await; }); @@ -78,6 +87,7 @@ async fn handle_connection( mut nip11_rx: watch::Receiver, dropped_tx: Arc>, mut dropped_rx: watch::Receiver, + active_websockets: Arc, ) { let mut header = [0_u8; 4096]; let header_len = loop { @@ -101,19 +111,25 @@ async fn handle_connection( let Ok(_websocket) = tokio_tungstenite::accept_async(stream).await else { return; }; + active_websockets.fetch_add(1, Ordering::SeqCst); while !*nip11_rx.borrow() && nip11_rx.changed().await.is_ok() {} // Dropping the live socket is the behavior under test. Give the SDK a // bounded propagation window before allowing the HTTP setup to finish. drop(_websocket); + active_websockets.fetch_sub(1, Ordering::SeqCst); let _ = dropped_tx.send(true); } else { let mut request_bytes = vec![0_u8; header_len]; if stream.read_exact(&mut request_bytes).await.is_err() { return; } - let _ = nip11_tx.send(true); - while !*dropped_rx.borrow() && dropped_rx.changed().await.is_ok() {} - tokio::time::sleep(Duration::from_millis(100)).await; + // Pre-dial NIP-11 probes (no live WebSocket) are answered without + // choreography; only a fetch during a live session drops it. + if active_websockets.load(Ordering::SeqCst) > 0 { + let _ = nip11_tx.send(true); + while !*dropped_rx.borrow() && dropped_rx.changed().await.is_ok() {} + tokio::time::sleep(Duration::from_millis(100)).await; + } let body = r#"{"limitation":{"max_subscriptions":20}}"#; let response = format!( "HTTP/1.1 200 OK\r\nContent-Type: application/nostr+json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", diff --git a/tests/common/sync_helpers.rs b/tests/common/sync_helpers.rs index 668b5fc..6b1fb23 100644 --- a/tests/common/sync_helpers.rs +++ b/tests/common/sync_helpers.rs @@ -63,6 +63,29 @@ pub async fn wait_for_full_repo_live_coverage(relay: &TestRelay, timeout: Durati } } +/// Wait until some line of a relay subprocess log satisfies the predicate. +pub async fn wait_for_log_line( + log_path: &std::path::Path, + timeout: Duration, + predicate: F, +) -> bool +where + F: Fn(&str) -> bool, +{ + let deadline = tokio::time::Instant::now() + timeout; + loop { + if let Ok(content) = tokio::fs::read_to_string(log_path).await { + if content.lines().any(&predicate) { + return true; + } + } + if tokio::time::Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + // NOTE: Using rust-nostr Kind variants: // - Kind::GitIssue.as_u16() -> Kind::GitIssue (1621) // - Kind::Comment.as_u16() -> Kind::Comment (1111) diff --git a/tests/outbound_policy.rs b/tests/outbound_policy.rs index dde79ea..cf2a139 100644 --- a/tests/outbound_policy.rs +++ b/tests/outbound_policy.rs @@ -20,12 +20,11 @@ mod common; -use std::path::Path; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::time::Duration; -use common::{create_state_event, MockRelay, TestClient, TestRelay}; +use common::{create_state_event, wait_for_log_line, MockRelay, TestClient, TestRelay}; use nostr_sdk::prelude::*; /// A commit hash that exists nowhere, keeping state events in purgatory so @@ -56,25 +55,6 @@ async fn start_counting_listener() -> (u16, Arc) { (port, count) } -/// Wait until some line of the relay subprocess log satisfies the predicate. -async fn wait_for_log_line(log_path: &Path, timeout: Duration, predicate: F) -> bool -where - F: Fn(&str) -> bool, -{ - let deadline = tokio::time::Instant::now() + timeout; - loop { - if let Ok(content) = tokio::fs::read_to_string(log_path).await { - if content.lines().any(&predicate) { - return true; - } - } - if tokio::time::Instant::now() >= deadline { - return false; - } - tokio::time::sleep(Duration::from_millis(100)).await; - } -} - /// Build a repository announcement with explicit clone and relay URL lists. fn announcement_with_urls( keys: &Keys, diff --git a/tests/sync.rs b/tests/sync.rs index 0c1f56d..9ce8a73 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -42,6 +42,7 @@ mod sync { pub mod metrics; pub mod naughty_list_scheduling; pub mod neg_concurrency; + pub mod outbound_auth; pub mod proactive_sync_plus; pub mod purgatory_fetch; pub mod reconnect_backoff; diff --git a/tests/sync/naughty_list_scheduling.rs b/tests/sync/naughty_list_scheduling.rs index add65f8..bc21231 100644 --- a/tests/sync/naughty_list_scheduling.rs +++ b/tests/sync/naughty_list_scheduling.rs @@ -85,15 +85,17 @@ async fn naughty_relay_is_not_scheduled_for_reconnection() { .await, "broken endpoint must be classified as naughty" ); - assert_eq!(accepted.load(Ordering::SeqCst), 1, "first dial is required"); + // The first attempt opens two connections: the pre-dial NIP-11 probe + // (GRASP-08 private-service detection) and the WebSocket dial itself. + assert_eq!(accepted.load(Ordering::SeqCst), 2, "first dial is required"); let reconnect_deadline = tokio::time::Instant::now() + RECONNECT_OBSERVATION; - while tokio::time::Instant::now() < reconnect_deadline && accepted.load(Ordering::SeqCst) == 1 { + while tokio::time::Instant::now() < reconnect_deadline && accepted.load(Ordering::SeqCst) == 2 { tokio::time::sleep(Duration::from_millis(100)).await; } assert_eq!( accepted.load(Ordering::SeqCst), - 1, + 2, "a naughty relay must not receive another dial across reconnect ticks" ); let log = tokio::fs::read_to_string(relay.log_path()) diff --git a/tests/sync/outbound_auth.rs b/tests/sync/outbound_auth.rs new file mode 100644 index 0000000..40f93c8 --- /dev/null +++ b/tests/sync/outbound_auth.rs @@ -0,0 +1,70 @@ +//! Outbound Authentication and GRASP-08 Private-Service Sync Policy +//! +//! These tests cover how a syncing instance treats relays that demand +//! authentication or advertise the GRASP-08 private-service extension: +//! +//! - 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. + +use std::time::Duration; + +use crate::common::{wait_for_log_line, TestRelay}; +use nostr_sdk::prelude::*; + +/// The stable park warning emitted when a public instance excludes a +/// GRASP-08 private service from sync. +const PARK_LOG: &str = "Relay advertises GRASP-08 private service; excluding it from public sync"; + +/// A public instance must never dial a relay whose NIP-11 advertises +/// GRASP-08: the private service would only refuse it, and dialing would +/// leak an AUTH exchange to a service that never admits this mirror. +#[tokio::test] +async fn public_instance_parks_grasp08_relay_without_dialing() { + let member = Keys::generate(); + let private_service = TestRelay::start_private(&member.public_key()).await; + let syncing = + TestRelay::start_with_sync_without_user_index(Some(private_service.url().to_string())) + .await; + + let parked = wait_for_log_line(&syncing.log_path(), Duration::from_secs(30), |line| { + line.contains(PARK_LOG) && line.contains(private_service.url()) + }) + .await; + assert!( + parked, + "public instance must park its GRASP-08 bootstrap relay" + ); + + // Absence over time (see relay_identity.rs for the sanctioned pattern): + // across a 2s observation window the parked relay is never connected to, + // no NIP-42 authentication happens, and the park warning is not repeated. + let window_end = tokio::time::Instant::now() + Duration::from_secs(2); + loop { + let log = tokio::fs::read_to_string(syncing.log_path()) + .await + .unwrap_or_default(); + assert!( + !log.lines() + .any(|line| line.contains("Connected") && line.contains(private_service.url())), + "parked GRASP-08 relay must never be dialed" + ); + assert!( + !log.lines() + .any(|line| line.contains("Authenticated to relay")), + "no NIP-42 exchange may happen with a parked relay" + ); + assert_eq!( + log.lines().filter(|line| line.contains(PARK_LOG)).count(), + 1, + "park warning must be logged exactly once" + ); + if tokio::time::Instant::now() >= window_end { + break; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + + syncing.stop().await; + private_service.stop().await; +}