diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 735ffed..015cfc1 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -716,7 +716,9 @@ single `PrivateAccess` set is shared by the HTTP and WebSocket services and the announcement admission policy. Push authorization remains the GRASP-01 policy; a private credential proves service membership but never grants push rights. The effective set combines operator-configured members with NIP-11 -owner pubkeys learned for relays referenced by accepted announcements. +owner pubkeys learned for relays referenced by accepted announcements, +provided those relays' NIP-11 also advertises GRASP-08 (a public relay's +owner gains nothing legitimate from private membership). Purgatory-only announcements are excluded. Reconciliation reuses the accepted repository index and the NIP-11 fetch already performed once per connection session, so private mode adds neither outbound connections nor subscriptions. @@ -754,6 +756,18 @@ identity publication: the kind 0/10002 events are seeded and served locally but never sent to the configured user-index relays, so a private relay does not advertise its existence. +Outbound, every sync connection (public or private instance) answers NIP-42 +challenges with the relay owner key when available; a `restricted:` refusal +after authentication parks the subscription via the existing policy-refusal +machinery. Peers advertising GRASP-08 in NIP-11 split by our own mode: a +public instance detects them with a pre-dial NIP-11 fetch and parks them +without ever opening the WebSocket, while a private instance treats them as +peers — NIP-42 on the WebSocket plus the GRASP-08 repository-root NIP-98 +credential on purgatory Git fetches from that peer's host, both signed with +the relay owner key. Relays without a readable `supported_grasps` are +ordinary sync targets. See +[GRASP-08 design](grasp-08-private-service.md#outbound-authentication-and-sync-policy). + One process currently represents one private collaborator service. Operators can run several independently configured instances for different groups. Fleet provisioning and lifecycle automation are deliberately left to a later change; diff --git a/docs/explanation/grasp-08-private-service.md b/docs/explanation/grasp-08-private-service.md index 93d12f6..6a9dc36 100644 --- a/docs/explanation/grasp-08-private-service.md +++ b/docs/explanation/grasp-08-private-service.md @@ -100,10 +100,11 @@ hints** (who to connect to), while the service-wide member set is the **ACL** ### Why accepted-relay owners are admitted dynamically Two private services mirroring the same repository must be able to read from -each other. When an accepted announcement references another relay, that -relay's NIP-11 `pubkey` is added to the member set, so a peer service (or the -operator of an ordinary relay the team uses) can authenticate without manual -whitelisting on both sides. The consequence — accepting one announcement +each other. When an accepted announcement references another relay whose +NIP-11 also advertises `GRASP-08`, that relay's NIP-11 `pubkey` is added to +the member set, so a peer service can authenticate without manual +whitelisting on both sides. Owners of referenced *public* relays are not +admitted (see "Outbound authentication and sync policy" below). The consequence — accepting one announcement grants its referenced relay operators read access to the whole service — follows directly from the one-trust-domain model and is the operator's opt-in via announcement admission. @@ -160,7 +161,8 @@ lock a user out. - **Accepted-relay owners** are derived members, trusted transitively via announcement admission plus the referenced relay's NIP-11 self-assertion (the same HTTPS-from-domain trust anchor a `_@domain` NIP-05 lookup would - provide, without an extra fetch or format). + provide, without an extra fetch or format). Only relays whose NIP-11 also + advertises `GRASP-08` qualify. - **Membership grants read access only.** Push authorization remains GRASP-01's maintainer model; repository admission remains announcement policy, which in private mode additionally requires the announcement @@ -169,13 +171,63 @@ lock a user out. invalidates future credentials, but repositories admitted while they were a member remain hosted until the operator curates them. +## Outbound authentication and sync policy + +The outbound half of private-service support decides how this instance, as a +*client*, treats the relays it syncs from: + +| Situation | Behavior | +| --- | --- | +| Any relay issues a NIP-42 challenge | Answer with the relay owner key (public and private instances alike). The SDK retries the refused subscription once after authenticating. Without an owner key, authentication is skipped and auth-demanding subscriptions park immediately. | +| A relay answers `restricted:` after valid authentication | Terminal: the subscription parks through the policy-refusal machinery (24-hour probe), no retry storm. | +| Peer NIP-11 advertises `GRASP-08`, this instance is **public** | Not a sync target at all: detected by a pre-dial NIP-11 fetch and parked without ever opening the WebSocket, so no AUTH exchange happens and no credential could leak. | +| Peer NIP-11 advertises `GRASP-08`, this instance is **private** | A peer: NIP-42 on the WebSocket *plus* the GRASP-08 repository-root NIP-98 credential attached to purgatory Git fetches from that peer's host, both signed with the relay owner key. | +| NIP-11 missing, unreadable, or without `supported_grasps` | An ordinary relay. | + +The pre-dial NIP-11 fetch re-runs the outbound target policy for +event-directed URLs first, so the SSRF gate covers it like the dial itself. + +**Why NIP-42 everywhere?** NIP-42 is identification, not confidentiality. A +gated relay admitting our pubkey grants a *known* service read access — our +pubkey is already published via NIP-11 and the NIP-05 root identity. A +private instance authenticating outbound discloses its identity to the relays +it syncs from, which is consistent with GRASP-08's public-discovery stance: +private mode hides repository content, not the service. + +**Why derived membership requires GRASP-08 (see above)?** For the same +asymmetry: a public relay's owner gains nothing legitimate from private +membership, because their relay enforces no confidentiality for the +repositories it mirrors. Only relays advertising `GRASP-08` in their NIP-11 +`supported_grasps` mint derived members; configured `NGIT_PRIVATE_MEMBERS` +are unaffected. + +## Disclosure and outside contributions + +Teams running a private service still receive security reports (CVEs, +vulnerability disclosures) from people outside the member set. GRASP-08 keeps +that path open without weakening the access boundary: + +- **Finding the contact**: the service's existence and operator identity are + deliberately public. The security contact is discoverable through the + NIP-11 `contact` field and the NIP-05 root identity (`_@domain`), both + served unauthenticated. +- **Sending a report**: reports arrive as NIP-17 encrypted direct messages on + the maintainers' public mailbox relays. Nothing about a private repository + needs to be readable for a reporter to reach its maintainers privately. +- **Granting collaboration access**: when a report leads to joint work, the + operator adds the reporter to `NGIT_PRIVATE_MEMBERS`. Membership grants + read access to the whole trust domain (see "Why membership is + service-wide"), which is the intended granularity: triaging a + vulnerability together means trusting the reporter with the codebase. +- **No non-member submission lane**: a dedicated unauthenticated inbox for + outside patches or PR events is deliberately not implemented. It would + reopen exactly the unauthenticated write surface GRASP-08 exists to close — + the same reasoning that makes private mode refuse to combine with + GRASP-06. + ## Follow-up scope Deliberately excluded from the initial single-service implementation: -- **Outbound authentication**: presenting NIP-42 and GRASP-08 NIP-98 - credentials when syncing *from* other private services, using a service - identity key. This is the missing half of zero-configuration private - mirroring. - **Multi-service fleet orchestration** and **encrypted kind-10318 client discovery**, which belong to future GRASP proposals. diff --git a/src/private/mod.rs b/src/private/mod.rs index 1a5596f..a01a0df 100644 --- a/src/private/mod.rs +++ b/src/private/mod.rs @@ -5,6 +5,8 @@ pub mod access; pub mod nip98; +pub mod peers; pub mod ws_auth; pub use access::PrivateAccess; +pub use peers::Grasp08Peers; diff --git a/src/private/nip98.rs b/src/private/nip98.rs index b61d36d..f526e92 100644 --- a/src/private/nip98.rs +++ b/src/private/nip98.rs @@ -3,7 +3,7 @@ use std::time::{SystemTime, UNIX_EPOCH}; use base64::Engine; use hyper::header::{AUTHORIZATION, WWW_AUTHENTICATE}; use hyper::{Request, Response, StatusCode}; -use nostr_sdk::prelude::{Event, Kind}; +use nostr_sdk::prelude::{Event, EventBuilder, FinalizeEvent, Keys, Kind, Tag}; use crate::config::Config; use crate::git::{empty_body, GitResponseBody}; @@ -78,6 +78,38 @@ pub fn validate_request( Ok(()) } +/// Build the outbound Authorization header for a GRASP-08 peer fetch. +/// +/// Signs exactly the profile [`validate_request`] checks: kind 27235 with one +/// `u` tag naming the repository root, one `method` tag of `GET`, and a fresh +/// timestamp (the peer allows 60 seconds of skew, so callers generate a new +/// header per subprocess invocation rather than caching one). +pub fn repository_credential_header( + keys: &Keys, + repository_root_url: &str, +) -> anyhow::Result { + let event = EventBuilder::new(Kind::HttpAuth, "") + .tags([ + Tag::parse(["u", repository_root_url])?, + Tag::parse(["method", "GET"])?, + ]) + .finalize(keys)?; + Ok(format!( + "Nostr {}", + base64::engine::general_purpose::STANDARD.encode(event.as_json()) + )) +} + +/// Truncate a Smart HTTP fetch URL at the repository root the GRASP-08 +/// credential must sign. +/// +/// Identifiers may themselves contain `.git`, so like +/// [`canonical_repository_url`] the root ends at the LAST occurrence. +pub fn repository_root_from_fetch_url(url: &str) -> Option { + let end = url.rfind(".git")?.checked_add(4)?; + url.get(..end).map(str::to_string) +} + /// Build the canonical absolute root URL signed by Git credentials. pub fn canonical_repository_url(config: &Config, request_path: &str) -> Option { // Identifiers may themselves contain `.git`; the route suffix is the last @@ -204,6 +236,33 @@ mod tests { } } + #[test] + fn outbound_credential_round_trips_through_validation() { + let keys = Keys::generate(); + let access = PrivateAccess::new([keys.public_key()]); + let root = repository_root_from_fetch_url( + "http://127.0.0.1:7334/npub/repo.git/info/refs?service=git-upload-pack", + ) + .expect("repository root"); + assert_eq!(root, "http://127.0.0.1:7334/npub/repo.git"); + assert_eq!( + repository_root_from_fetch_url("https://h.example/npub/repo.git-tools.git/info/refs") + .as_deref(), + Some("https://h.example/npub/repo.git-tools.git") + ); + + let header = repository_credential_header(&keys, &root).expect("credential header"); + let request = Request::get("/npub/repo.git/info/refs") + .header(AUTHORIZATION, header) + .body(()) + .unwrap(); + assert_eq!(validate_request(&request, &root, &access), Ok(())); + // The same credential must not authorize a different repository root. + assert!( + validate_request(&request, "http://127.0.0.1:7334/npub/other.git", &access).is_err() + ); + } + #[test] fn canonical_url_uses_http_only_for_loopback() { let mut config = Config::for_testing(); diff --git a/src/private/peers.rs b/src/private/peers.rs new file mode 100644 index 0000000..767b2d5 --- /dev/null +++ b/src/private/peers.rs @@ -0,0 +1,93 @@ +//! Registry of GRASP-08 private-service peers learned from NIP-11. +//! +//! A private instance attaches its outbound NIP-98 credential only to Git +//! fetches from hosts it has confirmed to be GRASP-08 private services, so +//! the credential is never presented to an ordinary public server. + +use std::collections::HashSet; +use std::sync::{Arc, RwLock}; + +/// Shared set of confirmed GRASP-08 peers, keyed by canonical `host:port`. +/// +/// The sync manager writes entries when a relay's NIP-11 advertises (or stops +/// advertising) GRASP-08; the purgatory Git fetch path reads them. Keys always +/// carry an explicit port - tests run several services on one loopback host, +/// so the host alone would collide. +#[derive(Clone, Debug, Default)] +pub struct Grasp08Peers { + inner: Arc>>, +} + +impl Grasp08Peers { + /// Record `key` as a confirmed GRASP-08 peer. Returns true when new. + pub fn insert(&self, key: String) -> bool { + self.inner + .write() + .expect("GRASP-08 peer set poisoned") + .insert(key) + } + + /// Forget `key`. Returns true when it was present. + pub fn remove(&self, key: &str) -> bool { + self.inner + .write() + .expect("GRASP-08 peer set poisoned") + .remove(key) + } + + /// Whether `key` is a confirmed GRASP-08 peer. + pub fn contains(&self, key: &str) -> bool { + self.inner + .read() + .expect("GRASP-08 peer set poisoned") + .contains(key) + } + + /// Canonical `host:port` registry key for a ws/wss/http/https URL. + /// + /// The same service is reached over WebSocket (relay URL) and HTTP + /// (clone URL); scheme-default ports are made explicit so both spellings + /// map to one key. + pub fn peer_key(url: &str) -> Option { + let parsed = reqwest::Url::parse(url).ok()?; + let host = parsed.host_str()?.to_ascii_lowercase(); + let port = parsed.port_or_known_default()?; + Some(format!("{host}:{port}")) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn ws_and_http_urls_for_one_service_share_a_key() { + let key = Grasp08Peers::peer_key("ws://127.0.0.1:7334").expect("ws key"); + assert_eq!(key, "127.0.0.1:7334"); + assert_eq!( + Grasp08Peers::peer_key("http://127.0.0.1:7334/npub/repo.git").as_deref(), + Some("127.0.0.1:7334") + ); + // Scheme-default ports become explicit for both transports. + assert_eq!( + Grasp08Peers::peer_key("wss://Relay.Example").as_deref(), + Some("relay.example:443") + ); + assert_eq!( + Grasp08Peers::peer_key("https://relay.example/x.git").as_deref(), + Some("relay.example:443") + ); + assert_eq!(Grasp08Peers::peer_key("not a url"), None); + } + + #[test] + fn registry_tracks_insert_and_remove() { + let peers = Grasp08Peers::default(); + assert!(!peers.contains("relay.example:443")); + assert!(peers.insert("relay.example:443".into())); + assert!(!peers.insert("relay.example:443".into())); + assert!(peers.contains("relay.example:443")); + assert!(peers.remove("relay.example:443")); + assert!(!peers.contains("relay.example:443")); + } +} diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index f83e423..07f9835 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -273,6 +273,14 @@ pub struct RealSyncContext { /// Outbound target policy applied before every event-directed git fetch outbound_policy: OutboundTargetPolicy, + + /// Confirmed GRASP-08 peers (by canonical `host:port`), fed by the sync + /// manager's NIP-11 fetches. Only set for private instances. + grasp08_peers: Option, + + /// Keys signing outbound GRASP-08 repository credentials. Only set for + /// private instances; fetches stay unauthenticated without them. + credential_keys: Option, } impl RealSyncContext { @@ -287,6 +295,9 @@ impl RealSyncContext { /// * `write_policy` - Write policy for promotion-time recovery hooks /// * `git_naughty_list` - Naughty list tracker for git remote domains /// * `outbound_policy` - Policy vetting event-directed git fetch targets + /// * `grasp08_peers` - Confirmed GRASP-08 peer registry (private mode only) + /// * `credential_keys` - Keys signing outbound GRASP-08 credentials + /// (private mode only) #[allow(clippy::too_many_arguments)] pub fn new( purgatory: Arc, @@ -297,6 +308,8 @@ impl RealSyncContext { write_policy: Option, git_naughty_list: Arc, outbound_policy: OutboundTargetPolicy, + grasp08_peers: Option, + credential_keys: Option, ) -> Self { Self { purgatory, @@ -308,6 +321,8 @@ impl RealSyncContext { git_naughty_list, miss_memo: Arc::new(Mutex::new(HashMap::new())), outbound_policy, + grasp08_peers, + credential_keys, } } @@ -410,7 +425,19 @@ fn resolve_pin_entry(resolved: &ResolvedTarget) -> Option { /// of requests to event-directed servers; /// - `http.curloptResolve` pins the vetted DNS answers so the fetch cannot be /// re-bound to a different address between authorization and connection. -fn hardened_git_command(repo_path: &Path, resolve_pin: Option<&str>, args: &[String]) -> Command { +/// +/// `auth_header` attaches an explicit `Authorization` header (the GRASP-08 +/// repository credential for a confirmed private peer). This deliberately +/// does not conflict with the `credential.helper=` hardening: that control +/// keeps *ambient operator* credentials away from event-directed servers, +/// while this header is a peer-scoped credential minted for exactly this +/// fetch target. +fn hardened_git_command( + repo_path: &Path, + resolve_pin: Option<&str>, + auth_header: Option<&str>, + args: &[String], +) -> Command { let mut command = Command::new("git"); command .arg("-c") @@ -422,6 +449,11 @@ fn hardened_git_command(repo_path: &Path, resolve_pin: Option<&str>, args: &[Str if let Some(pin) = resolve_pin { command.arg("-c").arg(format!("http.curloptResolve={pin}")); } + if let Some(header) = auth_header { + command + .arg("-c") + .arg(format!("http.extraHeader=Authorization: {header}")); + } command .args(args) .env("GIT_ALLOW_PROTOCOL", "http:https") @@ -949,6 +981,41 @@ impl SyncContext for RealSyncContext { let naughty_list = self.git_naughty_list.clone(); let miss_memo = self.miss_memo.clone(); + // GRASP-08: fetches from a confirmed private peer carry the + // repository-root credential. The registry key is host:port derived + // from the URL itself (`extract_domain` drops the port, which would + // collide loopback services). Headers are minted fresh before each + // subprocess: the peer's 60-second validity window must not expire + // mid-pass on a long batch fetch. + let credential_signer = self.credential_keys.clone().filter(|_| { + self.grasp08_peers.as_ref().is_some_and(|peers| { + crate::private::Grasp08Peers::peer_key(&url).is_some_and(|key| peers.contains(&key)) + }) + }); + let repository_root = credential_signer + .as_ref() + .and_then(|_| crate::private::nip98::repository_root_from_fetch_url(&url)); + let mut credential_warning_logged = false; + let credential_url = url.clone(); + let mut fresh_auth_header = move || -> Option { + let keys = credential_signer.as_ref()?; + let root = repository_root.as_deref()?; + match crate::private::nip98::repository_credential_header(keys, root) { + Ok(header) => Some(header), + Err(error) => { + if !credential_warning_logged { + credential_warning_logged = true; + tracing::warn!( + url = %credential_url, + error = %error, + "Failed to sign GRASP-08 credential; fetching unauthenticated" + ); + } + None + } + } + }; + // Phase 1: compare the remote's advertised refs against our // needs. Most needed OIDs are ref tips declared by state events // (PR tips appear under `refs/nostr/`), so the @@ -957,7 +1024,12 @@ impl SyncContext for RealSyncContext { // failed "not our ref" upload-pack round trip. let ls_remote_args = vec!["ls-remote".to_string(), url.clone()]; let advertised = match run_observed_git_command( - hardened_git_command(&repo_path, resolve_pin.as_deref(), &ls_remote_args), + hardened_git_command( + &repo_path, + resolve_pin.as_deref(), + fresh_auth_header().as_deref(), + &ls_remote_args, + ), &domain, "ls_remote", role, @@ -1022,7 +1094,12 @@ impl SyncContext for RealSyncContext { args.extend(advertised_tips.iter().cloned()); match run_observed_git_command( - hardened_git_command(&repo_path, resolve_pin.as_deref(), &args), + hardened_git_command( + &repo_path, + resolve_pin.as_deref(), + fresh_auth_header().as_deref(), + &args, + ), &domain, "fetch_batch", role, @@ -1087,7 +1164,12 @@ impl SyncContext for RealSyncContext { oid.clone(), ]; match run_observed_git_command( - hardened_git_command(&repo_path, resolve_pin.as_deref(), &args), + hardened_git_command( + &repo_path, + resolve_pin.as_deref(), + fresh_auth_header().as_deref(), + &args, + ), &domain, "fetch_residual", role, diff --git a/src/server.rs b/src/server.rs index a8fd6cb..3ac3620 100644 --- a/src/server.rs +++ b/src/server.rs @@ -256,6 +256,26 @@ impl RelayServer { // Start SyncManager for proactive sync (Phase 2: multi-relay support, Phase 3: health tracking) // Even without bootstrap relay, SyncManager discovers relays from stored announcements // Pass the already-registered sync metrics from Metrics to avoid duplicate registration + // GRASP-08 peer registry and credential signer exist only on private + // instances: a public mirror never authenticates its Git fetches. + let grasp08_peers = config + .private_mode + .then(crate::private::Grasp08Peers::default); + let outbound_credential_keys = if config.private_mode { + match config.relay_owner_keys() { + Ok(keys) => Some(keys), + Err(error) => { + warn!( + %error, + "Relay owner key unavailable; outbound GRASP-08 credentials disabled" + ); + None + } + } + } else { + None + }; + let sync_manager = SyncManager::new( config.sync_bootstrap_relay_url.clone(), config.domain.clone(), @@ -266,6 +286,7 @@ impl RelayServer { PathBuf::from(config.effective_git_data_path()), metrics.as_ref().and_then(|m| m.sync_metrics().cloned()), private_access.clone(), + grasp08_peers.clone(), ); if config.sync_bootstrap_relay_url.is_some() { @@ -397,6 +418,8 @@ impl RelayServer { OutboundTargetPolicy { allow_non_global: config.sync_allow_non_global_targets, }, + grasp08_peers, + outbound_credential_keys, )); // Create throttle manager for rate limiting remote git servers diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 8e94f74..40d9c27 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1785,7 +1785,12 @@ enum ConnectAttemptOutcome { advertised_default_limit: Option, advertised_max_subscriptions: Option, 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), } @@ -2464,9 +2469,19 @@ pub struct SyncManager { proactive_participant_authors: crate::nostr::policy::SharedProactiveParticipantAuthorIndex, /// GRASP-08 access shared with the inbound HTTP/WebSocket boundary. private_access: Option, + /// Confirmed GRASP-08 peers shared with the purgatory Git fetch path. + /// + /// Only populated on private instances (the server passes `Some`): a + /// public mirror never presents credentials, and its GRASP-08 targets are + /// parked pre-dial instead. + grasp08_peers: Option, /// Operator-configured members form the permanent base of private access. configured_private_members: HashSet, - /// Latest NIP-11 owner learned for each connected repository relay. + /// Latest NIP-11 owner learned for each connected repository relay whose + /// NIP-11 also advertises GRASP-08. Only private-service owners can mint + /// derived membership: the owner of a public relay gains nothing + /// legitimate from private membership, since their relay enforces no + /// confidentiality for the repositories it mirrors. relay_owners: HashMap, /// What we've confirmed syncing + connection state relay_sync_index: RelaySyncIndex, @@ -2492,6 +2507,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 @@ -2570,6 +2592,7 @@ impl SyncManager { data_path: PathBuf, sync_metrics: Option, private_access: Option, + grasp08_peers: Option, ) -> Self { // Extract purgatory from write_policy for read-only access let purgatory = write_policy.purgatory().clone(); @@ -2622,6 +2645,7 @@ impl SyncManager { root_candidate_index: Arc::new(RwLock::new(HashMap::new())), proactive_participant_authors, private_access, + grasp08_peers, configured_private_members, relay_owners: HashMap::new(), relay_sync_index: Arc::new(RwLock::new(HashMap::new())), @@ -2631,6 +2655,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(), @@ -5271,6 +5296,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. @@ -5305,11 +5340,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(), @@ -5399,6 +5443,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; @@ -5445,6 +5496,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"); @@ -5459,17 +5511,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, + 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; @@ -6253,8 +6317,13 @@ impl SyncManager { advertised_default_limit, advertised_max_subscriptions, advertised_owner, + advertised_grasp08, } => { - match advertised_owner { + // Only GRASP-08-advertising relays mint derived private + // membership: a public relay's owner gains nothing legitimate + // from private membership because their relay enforces no + // confidentiality for the mirrored repositories. + match advertised_owner.filter(|_| advertised_grasp08) { Some(owner) => { self.relay_owners.insert(result.relay_url.clone(), owner); } @@ -6262,6 +6331,23 @@ impl SyncManager { self.relay_owners.remove(&result.relay_url); } } + // On a private instance, remember which hosts are GRASP-08 + // peers so purgatory Git fetches from them carry the outbound + // NIP-98 credential. + if let Some(peers) = &self.grasp08_peers { + if let Some(key) = crate::private::Grasp08Peers::peer_key(&result.relay_url) { + if advertised_grasp08 { + if peers.insert(key) { + tracing::info!( + relay = %result.relay_url, + "Confirmed GRASP-08 peer; Git fetches will carry credentials" + ); + } + } else { + peers.remove(&key); + } + } + } self.reconcile_private_membership().await; if let Some(connection) = self.connections.get(&result.relay_url) { connection.reset_subscription_budget(advertised_max_subscriptions); @@ -6288,6 +6374,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() { @@ -7775,11 +7886,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 98d4930..6e19181 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -442,12 +442,26 @@ pub struct RelayLimitHints { /// grant this identity access only when the relay is referenced by an /// accepted repository announcement. pub owner: Option, + /// Whether the document's `supported_grasps` array (a GRASP extension + /// field, parsed from the raw JSON because the SDK's NIP-11 type does not + /// carry it) advertises the "GRASP-08" private-service extension. + pub grasp08: bool, } fn parse_relay_limit_hints(body: &str) -> RelayLimitHints { let Some(document) = nostr::nips::nip11::RelayInformationDocument::from_json(body).ok() else { return RelayLimitHints::default(); }; + let grasp08 = serde_json::from_str::(body) + .ok() + .and_then(|value| { + value.get("supported_grasps")?.as_array().map(|grasps| { + grasps + .iter() + .any(|grasp| grasp.as_str() == Some("GRASP-08")) + }) + }) + .unwrap_or(false); let limitation = document.limitation.unwrap_or_default(); RelayLimitHints { default_limit: limitation @@ -459,6 +473,7 @@ fn parse_relay_limit_hints(body: &str) -> RelayLimitHints { .and_then(|limit| usize::try_from(limit).ok()) .filter(|limit| *limit > 0), owner: document.pubkey, + grasp08, } } @@ -573,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) @@ -697,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)), @@ -763,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)), @@ -966,6 +995,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 @@ -1665,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; @@ -1766,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(); } @@ -2176,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. @@ -2763,12 +2830,28 @@ mod tests { assert_eq!(hints.max_subscriptions, None); } + #[test] + fn grasp08_flag_requires_supported_grasps_entry() { + assert!(parse_relay_limit_hints(r#"{"supported_grasps":["GRASP-01","GRASP-08"]}"#).grasp08); + assert!(!parse_relay_limit_hints(r#"{"supported_grasps":["GRASP-01"]}"#).grasp08); + assert!(!parse_relay_limit_hints(r#"{"name":"relay without grasps"}"#).grasp08); + } + + #[test] + fn grasp08_flag_defaults_to_false_for_malformed_documents() { + // Non-array supported_grasps and unparseable bodies both mean "not a + // known private service", never an error. + assert!(!parse_relay_limit_hints(r#"{"supported_grasps":"GRASP-08"}"#).grasp08); + assert!(!parse_relay_limit_hints(r#"{"supported_grasps":8}"#).grasp08); + assert!(!parse_relay_limit_hints("not json at all").grasp08); + } + /// Event-directed connection with the permissive policy used by tests /// that dial loopback fixtures. fn permissive_connection(url: &str, keys: Keys) -> RelayConnection { RelayConnection::new( url.to_string(), - keys, + Some(keys), RelayTargetSource::EventDirected, OutboundTargetPolicy { allow_non_global: true, @@ -3016,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(), ); @@ -3074,7 +3157,7 @@ mod tests { let connection = RelayConnection::new( relay.url().await.to_string(), - Keys::generate(), + Some(Keys::generate()), RelayTargetSource::OperatorConfigured, OutboundTargetPolicy::default(), ); @@ -3121,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(), ); @@ -3168,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(), ); @@ -3203,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(), ); @@ -3366,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(), ); @@ -3708,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(), ); @@ -3771,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(), ); @@ -3863,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/mock_relay.rs b/tests/common/mock_relay.rs index 43befa1..ebbed89 100644 --- a/tests/common/mock_relay.rs +++ b/tests/common/mock_relay.rs @@ -122,6 +122,24 @@ impl MockRelay { rate_limit: RateLimit, pagination: Option, max_filters: Option, + ) -> Self { + Self::start_with_options_and_nip11(rate_limit, pagination, max_filters, None).await + } + + /// Start a mock relay that serves a caller-provided NIP-11 document for + /// `Accept: application/nostr+json` requests. + /// + /// This lets tests control fields the default document never emits, such + /// as the GRASP `supported_grasps` extension array or the `pubkey` owner. + pub async fn start_with_nip11_document(document: serde_json::Value) -> Self { + Self::start_with_options_and_nip11(RateLimit::default(), None, None, Some(document)).await + } + + async fn start_with_options_and_nip11( + rate_limit: RateLimit, + pagination: Option, + max_filters: Option, + custom_nip11: Option, ) -> Self { // Create and bind listener (eliminates port race condition) let std_listener = @@ -145,6 +163,7 @@ impl MockRelay { pagination, max_filters, Vec::new(), + custom_nip11, ) .await } @@ -164,10 +183,20 @@ impl MockRelay { let listener = TcpListener::bind(addr) .await .expect("Failed to bind to address"); - Self::start_with_listener(listener, port, RateLimit::default(), None, None, events).await + Self::start_with_listener( + listener, + port, + RateLimit::default(), + None, + None, + events, + None, + ) + .await } /// Internal method to start the relay with an existing listener. + #[allow(clippy::too_many_arguments)] async fn start_with_listener( listener: TcpListener, port: u16, @@ -175,6 +204,7 @@ impl MockRelay { pagination: Option, max_filters: Option, initial_events: Vec, + custom_nip11: Option, ) -> Self { // Create a simple relay with no write policy (accepts all events) let mut builder = LocalRelayBuilder::default().rate_limit(rate_limit); @@ -209,13 +239,22 @@ impl MockRelay { Ok((stream, remote_addr)) => { let relay = server_relay.clone(); let pagination = pagination; + let custom_nip11 = custom_nip11.clone(); let io = TokioIo::new(stream); tokio::spawn(async move { let service = service_fn(move |req| { let relay = relay.clone(); + let custom_nip11 = custom_nip11.clone(); async move { - handle_request(req, relay, remote_addr, pagination).await + handle_request( + req, + relay, + remote_addr, + pagination, + custom_nip11, + ) + .await } }); @@ -302,6 +341,7 @@ async fn handle_request( relay: LocalRelay, addr: SocketAddr, pagination: Option, + custom_nip11: Option, ) -> Result>, hyper::Error> { // Check for WebSocket upgrade request let is_websocket = req @@ -350,6 +390,13 @@ async fn handle_request( .and_then(|value| value.to_str().ok()) .is_some_and(|value| value.contains("application/nostr+json")) { + if let Some(document) = custom_nip11 { + return Ok(Response::builder() + .status(StatusCode::OK) + .header("Content-Type", "application/nostr+json") + .body(Full::new(Bytes::from(document.to_string()))) + .unwrap()); + } let limitation = pagination.map(|config| { serde_json::json!({ "default_limit": config.advertised_default_limit, 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/purgatory_helpers.rs b/tests/common/purgatory_helpers.rs index 07fa844..b9dd77f 100644 --- a/tests/common/purgatory_helpers.rs +++ b/tests/common/purgatory_helpers.rs @@ -587,6 +587,21 @@ pub fn push_to_relay( relay_domain: &str, npub: &str, repo_id: &str, +) -> Result<(), String> { + push_to_relay_with_auth_header(local_path, relay_domain, npub, repo_id, None) +} + +/// Push a local repository to a relay with an optional `Authorization` header. +/// +/// GRASP-08 private relays require every Smart HTTP request to carry the +/// repository-scoped NIP-98 credential, which git can attach via +/// `http.extraHeader`. Pass the full header value (e.g. `Nostr `). +pub fn push_to_relay_with_auth_header( + local_path: &Path, + relay_domain: &str, + npub: &str, + repo_id: &str, + auth_header: Option<&str>, ) -> Result<(), String> { let remote_url = format!("http://{}/{}/{}.git", relay_domain, npub, repo_id); @@ -606,7 +621,13 @@ pub fn push_to_relay( } // Push all refs - let output = git_command() + let mut command = git_command(); + if let Some(header) = auth_header { + command + .arg("-c") + .arg(format!("http.extraHeader=Authorization: {header}")); + } + let output = command .args(["push", "-u", "origin", "--all"]) .current_dir(local_path) .output() diff --git a/tests/common/relay.rs b/tests/common/relay.rs index 3ecfca6..13f8bfd 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -155,6 +155,36 @@ impl TestRelay { .await } + /// Start a GRASP-08 private service with persistent LMDB storage in + /// caller-owned directories, so [`Self::restart`] resumes from the same + /// events and git data. + /// + /// The private NIP-42 gate applies to the relay's own internal + /// self-subscription, so sync targets for locally published + /// announcements are derived from the database at startup; tests + /// exercising that pipeline publish, then restart. + pub async fn start_private_with_lmdb_paths( + member: &nostr_sdk::prelude::PublicKey, + git_data_path: PathBuf, + relay_data_path: PathBuf, + ) -> Self { + Self::start_internal( + port::reserve_port(), + RelayOptions { + private_members: Some( + member + .to_bech32() + .expect("Failed to encode private test member"), + ), + lmdb_backend: true, + git_data_path: Some(git_data_path), + relay_data_path: Some(relay_data_path), + ..RelayOptions::default() + }, + ) + .await + } + /// Start relay with sync from another relay (bootstrap relay) /// /// # Example @@ -182,6 +212,43 @@ 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 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 @@ -202,6 +269,38 @@ impl TestRelay { .await } + /// Start a GRASP-08 private service with an explicit member list, an + /// optional sync bootstrap, and a caller-chosen relay-owner identity. + /// + /// Lets private-to-private sync tests model two services whose member + /// sets deliberately differ (e.g. the source service admits the + /// mirroring service's owner identity). + pub async fn start_private_with_sync_owner_keys_and_members( + bootstrap_relay_url: Option, + owner_keys: Keys, + members: &[nostr_sdk::prelude::PublicKey], + ) -> Self { + let members = members + .iter() + .map(|member| { + member + .to_bech32() + .expect("Failed to encode private test member") + }) + .collect::>() + .join(","); + Self::start_internal( + port::reserve_port(), + RelayOptions { + bootstrap_relay_url, + owner_keys: Some(owner_keys), + private_members: Some(members), + ..RelayOptions::default() + }, + ) + .await + } + /// Start a GRASP-08 private relay whose sole configured member is also /// the relay-owner identity, with user-index relays configured. /// 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 4084a05..e7dd955 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/private_mode.rs b/tests/private_mode.rs index 0d7fadd..0c7b122 100644 --- a/tests/private_mode.rs +++ b/tests/private_mode.rs @@ -3,7 +3,10 @@ mod common; use std::time::Duration; use base64::Engine; -use common::TestRelay; +use common::{ + create_state_event, create_test_repo_with_commit, push_to_relay_with_auth_header, + wait_for_sync_connection, CommitVariant, MockRelay, TestClient, TestRelay, +}; use futures_util::{SinkExt, StreamExt}; use nostr_sdk::prelude::{EventBuilder, FinalizeEvent, Keys, Kind, Tag, Timestamp, ToBech32}; use reqwest::header::{ACCEPT, AUTHORIZATION, WWW_AUTHENTICATE}; @@ -428,6 +431,322 @@ async fn private_announcement_admission_requires_member_author() { relay.stop().await; } +/// Poll until `event_id` is served to an authenticated session using `keys`. +async fn wait_for_event_as( + relay_url: &str, + keys: &Keys, + event_id: nostr_sdk::prelude::EventId, + timeout: Duration, +) -> bool { + use nostr_sdk::prelude::{Client, Filter, SignerAuthenticator}; + let deadline = tokio::time::Instant::now() + timeout; + loop { + let client = Client::builder() + .authenticator(SignerAuthenticator::new(keys.clone())) + .build(); + if client.add_relay(relay_url).await.is_ok() { + client.connect().await; + let result = client + .fetch_events(Filter::new().id(event_id)) + .timeout(Duration::from_secs(2)) + .await; + client.disconnect().await; + if let Ok(events) = result { + if !events.is_empty() { + return true; + } + } + } + if tokio::time::Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(250)).await; + } +} + +/// One bounded probe of whether `keys` is currently an effective member: +/// authenticate over a fresh WebSocket session and, when admitted, prove the +/// session is bridged by seeing a REQ answered with EOSE. +async fn membership_admitted(relay: &TestRelay, keys: &Keys) -> bool { + let (mut stream, challenge) = connect_and_challenge(relay).await; + send_text(&mut stream, auth_message(keys, &relay.domain(), &challenge)).await; + let ok: serde_json::Value = + serde_json::from_str(&next_text(&mut stream).await).expect("OK JSON"); + assert_eq!(ok[0], "OK"); + if ok[2] != true { + return false; + } + send_text( + &mut stream, + r#"["REQ","bridge",{"kinds":[1],"limit":1}]"#.to_string(), + ) + .await; + loop { + let frame: serde_json::Value = + serde_json::from_str(&next_text(&mut stream).await).expect("relay JSON"); + if frame[0] == "EOSE" && frame[1] == "bridge" { + return true; + } + assert_ne!( + frame[0], "CLOSED", + "admitted member REQ was rejected: {frame}" + ); + } +} + +#[tokio::test] +async fn derived_membership_requires_grasp08_advertising_relay() { + let member = Keys::generate(); + let owner_plain = Keys::generate(); + let owner_private = Keys::generate(); + let git_data_dir = tempfile::tempdir().expect("git data dir"); + let relay_data_dir = tempfile::tempdir().expect("relay data dir"); + let relay = TestRelay::start_private_with_lmdb_paths( + &member.public_key(), + git_data_dir.path().to_path_buf(), + relay_data_dir.path().to_path_buf(), + ) + .await; + + // Two referenced relays that differ only in whether their NIP-11 + // advertises the GRASP-08 private-service extension. + let plain = MockRelay::start_with_nip11_document(serde_json::json!({ + "name": "public mirror", + "pubkey": owner_plain.public_key().to_hex(), + "supported_nips": [1, 11], + "supported_grasps": ["GRASP-01"], + })) + .await; + let private_peer = MockRelay::start_with_nip11_document(serde_json::json!({ + "name": "private peer", + "pubkey": owner_private.public_key().to_hex(), + "supported_nips": [1, 11], + "supported_grasps": ["GRASP-01", "GRASP-08"], + })) + .await; + + // Member-authored announcement listing our own service (so it is + // admitted) plus both mock relays; the state event and git push promote + // it out of purgatory so it can mint derived membership. + let identifier = "grasp08-derived-membership"; + let npub = member.public_key().to_bech32().expect("member npub"); + let clone_url = format!("http://{}/{npub}/{identifier}.git", relay.domain()); + let relay_urls = vec![ + format!("ws://{}", relay.domain()), + plain.url().to_string(), + private_peer.url().to_string(), + ]; + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom("clone", vec![clone_url.clone()]), + Tag::custom("relays", relay_urls.clone()), + ]) + .finalize(&member) + .expect("signed announcement"); + + let git_dir = tempfile::tempdir().expect("git temp dir"); + let commit = create_test_repo_with_commit(git_dir.path(), CommitVariant::StateTest) + .expect("test repository"); + let relay_url_refs: Vec<&str> = relay_urls.iter().map(String::as_str).collect(); + let state_event = create_state_event( + &member, + identifier, + &[("main", &commit)], + &[], + &[clone_url.as_str()], + &relay_url_refs, + ) + .expect("state event"); + + // The member passes the inbound NIP-42 gate via the client authenticator + // and the inbound NIP-98 gate via the repository-root credential. + let client = TestClient::new(relay.url(), member.clone()) + .await + .expect("authenticated member client"); + client + .send_event(&announcement) + .await + .expect("announcement admitted"); + client + .send_event(&state_event) + .await + .expect("state event admitted"); + push_to_relay_with_auth_header( + git_dir.path(), + &relay.domain(), + &npub, + identifier, + Some(&credential(&member, &clone_url)), + ) + .expect("authenticated git push"); + client.disconnect().await; + + assert!( + wait_for_event_as( + relay.url(), + &member, + announcement.id, + Duration::from_secs(30) + ) + .await, + "announcement must be promoted out of purgatory" + ); + + // The internal live self-subscription cannot pass the private NIP-42 + // gate, so sync targets for locally promoted announcements are derived + // from the database at startup; restart to exercise that path. + let relay = relay.restart().await; + + wait_for_sync_connection(relay.url(), 2, Duration::from_secs(60)) + .await + .expect("sync connections to both referenced relays"); + + // The GRASP-08-advertising relay's owner becomes an effective member. + let deadline = tokio::time::Instant::now() + Duration::from_secs(60); + loop { + if membership_admitted(&relay, &owner_private).await { + break; + } + assert!( + tokio::time::Instant::now() < deadline, + "GRASP-08 relay owner was not admitted before the deadline" + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + + // The plain relay's owner authenticates validly but is never minted. + let (mut stream, challenge) = connect_and_challenge(&relay).await; + send_text( + &mut stream, + auth_message(&owner_plain, &relay.domain(), &challenge), + ) + .await; + let ok: serde_json::Value = + serde_json::from_str(&next_text(&mut stream).await).expect("OK JSON"); + assert_eq!(ok[0], "OK"); + assert_eq!( + ok[2], false, + "non-GRASP-08 relay owner must not gain membership: {ok}" + ); + assert!( + ok[3] + .as_str() + .expect("OK message") + .starts_with("restricted:"), + "{ok}" + ); + expect_closed(&mut stream).await; + + relay.stop().await; + plain.stop().await; + private_peer.stop().await; +} + +/// End-to-end private-to-private mirroring: the mirroring instance must +/// authenticate its WebSocket session with NIP-42 AND attach the GRASP-08 +/// repository credential to its purgatory Git fetches. The state event is +/// only promoted (and served) once the git data arrived, which requires both +/// credentials to have worked against the private source. +#[tokio::test] +async fn private_instance_syncs_from_private_peer_with_outbound_credentials() { + let member = Keys::generate(); + let a_owner = Keys::generate(); + + // Source service B admits the member and A's owner identity. + let relay_b = TestRelay::start_private_with_sync_owner_keys_and_members( + None, + Keys::generate(), + &[member.public_key(), a_owner.public_key()], + ) + .await; + // Mirroring service A bootstraps from B with its own owner identity. + let relay_a = TestRelay::start_private_with_sync_owner_keys_and_members( + Some(relay_b.url().to_string()), + a_owner.clone(), + &[member.public_key()], + ) + .await; + + let identifier = "private-peer-sync"; + let npub = member.public_key().to_bech32().expect("member npub"); + let b_clone_url = format!("http://{}/{npub}/{identifier}.git", relay_b.domain()); + // Both services must appear in clone AND relays tags for admission; + // each instance excludes its own domain from fetch targets, so A only + // ever fetches git data from B. + let clone_urls = vec![ + b_clone_url.clone(), + format!("http://{}/{npub}/{identifier}.git", relay_a.domain()), + ]; + let relay_urls = vec![ + format!("ws://{}", relay_b.domain()), + format!("ws://{}", relay_a.domain()), + ]; + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom("clone", clone_urls.clone()), + Tag::custom("relays", relay_urls.clone()), + ]) + .finalize(&member) + .expect("signed announcement"); + + let git_dir = tempfile::tempdir().expect("git temp dir"); + let commit = create_test_repo_with_commit(git_dir.path(), CommitVariant::StateTest) + .expect("test repository"); + let clone_url_refs: Vec<&str> = clone_urls.iter().map(String::as_str).collect(); + let relay_url_refs: Vec<&str> = relay_urls.iter().map(String::as_str).collect(); + let state_event = create_state_event( + &member, + identifier, + &[("main", &commit)], + &[], + &clone_url_refs, + &relay_url_refs, + ) + .expect("state event"); + + // Publish to B through the member's authenticated session and push the + // git data with the member's inbound GRASP-08 credential. + let client = TestClient::new(relay_b.url(), member.clone()) + .await + .expect("authenticated member client"); + client + .send_event(&announcement) + .await + .expect("announcement admitted on B"); + client + .send_event(&state_event) + .await + .expect("state event admitted on B"); + push_to_relay_with_auth_header( + git_dir.path(), + &relay_b.domain(), + &npub, + identifier, + Some(&credential(&member, &b_clone_url)), + ) + .expect("authenticated git push to B"); + client.disconnect().await; + + // Promotion on A requires the git fetch from B to have succeeded, which + // in turn requires the outbound NIP-98 credential; a generous deadline + // covers connect, sync, and the purgatory fetch pass. + assert!( + wait_for_event_as( + relay_a.url(), + &member, + state_event.id, + Duration::from_secs(120) + ) + .await, + "state event must be promoted on the mirroring private instance" + ); + + relay_a.stop().await; + relay_b.stop().await; +} + #[tokio::test] async fn private_websocket_bounds_invalid_authentication_attempts() { let member = Keys::generate(); 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..3b77ffe --- /dev/null +++ b/tests/sync/outbound_auth.rs @@ -0,0 +1,183 @@ +//! 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. +//! - 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::{ + 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 +/// 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; +} + +/// 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; +}