Merge #0a6a0d75: feat(private-repos): authenticate outbound sync betwee…

feat(private-repos): authenticate outbound sync between GRASP-08 services

nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsq56sdwh0cccm69uqln5ar58hf38xssqwpzs9fz87kfdtlv3x45qqr462cs

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

PR description:

Implements the outbound half of GRASP-08 private-service support: how an instance behaves as a sync client toward relays that demand authentication or advertise GRASP-08.

Policy: (1) every sync connection answers NIP-42 AUTH challenges with the relay owner key when available, on public and private instances alike; after successful authentication the refused subscription is retried once, and without an owner key authentication is skipped and auth-demanding subscriptions park immediately. (2) A restricted: rejection after valid authentication is terminal and parks subscription work via the existing policy-refusal machinery, so there are no retry storms. (3) A relay whose NIP-11 supported_grasps advertises GRASP-08 is recognised with a pre-dial NIP-11 fetch that re-runs the outbound target policy, keeping the SSRF gate intact: a public instance parks such a relay without ever opening the WebSocket (no AUTH exchange, nothing to launder), while a private instance treats it as a peer and attaches the GRASP-08 repository-root NIP-98 credential (60-second validity, minted fresh per git subprocess) to purgatory Git fetches from that host, both signed with the relay owner key. (4) Relays with missing or unreadable NIP-11 remain ordinary sync targets. (5) Membership tightening: the NIP-11 owner of a relay referenced by an accepted announcement is only admitted as a derived member when that relay also advertises GRASP-08; configured NGIT_PRIVATE_MEMBERS are unaffected.

NIP-42 here is identification, not confidentiality: our pubkey is already published via NIP-11 and NIP-05, and a private instance authenticating outbound discloses its identity consistently with the public-discovery stance of GRASP-08. Also adds disclosure guidance for private services: security contact via NIP-11 contact and the NIP-05 root identity, reports as NIP-17 encrypted DMs on the public mailbox relays of maintainers, collaboration access by adding the reporter to NGIT_PRIVATE_MEMBERS, and deliberately no non-member submission lane.

Validation: cargo fmt --check, cargo clippy --all-targets -D warnings, cargo test --lib (792 passed), and the private_mode, sync, outbound_policy, purgatory_sync, nip11_document and relay_identity integration suites all pass.
This commit is contained in:
DanConwayDev
2026-08-15 19:50:34 +01:00
21 changed files with 1749 additions and 102 deletions
+15 -1
View File
@@ -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;
+61 -9
View File
@@ -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.
+2
View File
@@ -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;
+60 -1
View File
@@ -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<B>(
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<String> {
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<String> {
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<String> {
// 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();
+93
View File
@@ -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<RwLock<HashSet<String>>>,
}
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<String> {
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"));
}
}
+86 -4
View File
@@ -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<crate::private::Grasp08Peers>,
/// Keys signing outbound GRASP-08 repository credentials. Only set for
/// private instances; fetches stay unauthenticated without them.
credential_keys: Option<nostr_sdk::prelude::Keys>,
}
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<Purgatory>,
@@ -297,6 +308,8 @@ impl RealSyncContext {
write_policy: Option<Nip34WritePolicy>,
git_naughty_list: Arc<NaughtyListTracker>,
outbound_policy: OutboundTargetPolicy,
grasp08_peers: Option<crate::private::Grasp08Peers>,
credential_keys: Option<nostr_sdk::prelude::Keys>,
) -> 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<String> {
/// 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<String> {
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/<event-id>`), 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,
+23
View File
@@ -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
+139 -22
View File
@@ -1785,7 +1785,12 @@ enum ConnectAttemptOutcome {
advertised_default_limit: Option<usize>,
advertised_max_subscriptions: Option<usize>,
advertised_owner: Option<PublicKey>,
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<PrivateAccess>,
/// 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<crate::private::Grasp08Peers>,
/// Operator-configured members form the permanent base of private access.
configured_private_members: HashSet<PublicKey>,
/// 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<String, PublicKey>,
/// 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<String>,
/// 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<String>,
/// Last exact-ID dependency recovery attempt, used to bound retries.
dependency_refetch_attempts: Arc<std::sync::Mutex<HashMap<EventId, Instant>>>,
/// Events relays reported during negentropy reconciliation but failed to
@@ -2570,6 +2592,7 @@ impl SyncManager {
data_path: PathBuf,
sync_metrics: Option<SyncMetrics>,
private_access: Option<PrivateAccess>,
grasp08_peers: Option<crate::private::Grasp08Peers>,
) -> 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,
+117 -34
View File
@@ -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<PublicKey>,
/// 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::<serde_json::Value>(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<SharedDatabase>,
/// 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<Keys>,
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<Keys>,
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(),
);
+429
View File
@@ -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", <challenge>]`. 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<Mutex<HashSet<PublicKey>>>,
req_count: Arc<AtomicUsize>,
}
/// NIP-42 gate in front of a backend relay. See the module docs.
pub struct AuthGatingRelay {
url: String,
state: GateState,
shutdown_tx: Option<oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
}
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<PublicKey> {
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<hyper::body::Incoming>,
backend_url: String,
mode: GateMode,
state: GateState,
) -> Result<Response<Full<Bytes>>, 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<TokioIo<Upgraded>>;
type SharedClientSink = Arc<tokio::sync::Mutex<SplitSink<ClientWs, Message>>>;
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<SplitSink<_, Message>> = None;
let mut backend_task: Option<tokio::task::JoinHandle<()>> = 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::<serde_json::Value>(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<tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>>,
>,
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(),
)
}
+49 -2
View File
@@ -122,6 +122,24 @@ impl MockRelay {
rate_limit: RateLimit,
pagination: Option<PaginationConfig>,
max_filters: Option<usize>,
) -> 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<PaginationConfig>,
max_filters: Option<usize>,
custom_nip11: Option<serde_json::Value>,
) -> 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<PaginationConfig>,
max_filters: Option<usize>,
initial_events: Vec<Event>,
custom_nip11: Option<serde_json::Value>,
) -> 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<PaginationConfig>,
custom_nip11: Option<serde_json::Value>,
) -> Result<Response<Full<Bytes>>, 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,
+2
View File
@@ -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::*;
+22 -1
View File
@@ -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 <base64>`).
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()
+99
View File
@@ -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<String>) -> 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<String>,
) -> 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<String>,
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::<Vec<_>>()
.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.
///
+19 -3
View File
@@ -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<bool>,
dropped_tx: Arc<watch::Sender<bool>>,
mut dropped_rx: watch::Receiver<bool>,
active_websockets: Arc<AtomicUsize>,
) {
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{}",
+23
View File
@@ -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<F>(
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)
+1 -21
View File
@@ -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<AtomicUsize>) {
(port, count)
}
/// Wait until some line of the relay subprocess log satisfies the predicate.
async fn wait_for_log_line<F>(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,
+320 -1
View File
@@ -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();
+1
View File
@@ -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;
+5 -3
View File
@@ -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())
+183
View File
@@ -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;
}