Files
ngit-grasp/src/server.rs
T
DanConwayDev 6a3b723eaf fix(sync): restrict event-directed sync targets to globally reachable endpoints
After deploying a8964bb to gitnostr.com, production logs showed
event-directed proactive sync dialling ws://localhost:3334,
ws://127.0.0.1:7334, and ws://100.125.184.46:7334 (CGNAT). Repository
announcements, state events, and PR events are untrusted - anyone can
publish them - yet their relays/clone tags reached the outbound
WebSocket and git-fetch sinks after syntax-only checks, letting a
crafted event point a public relay at loopback, private, link-local,
or local-name infrastructure (SSRF).

Add one fail-closed outbound target policy (src/outbound.rs) applied
immediately before every event-directed sink so no call path can
bypass it:

- RelayConnection::connect re-authorizes (with DNS vetting) before
  every dial and reconnect. SyncManager::register_relay additionally
  refuses to register forbidden targets so they never enter the
  reconnect lifecycle, and memoizes rejections so stored events cannot
  spam logs or starve the bounded purgatory sync tick.
- RealSyncContext::fetch_oids authorizes clone URLs from announcements
  and purgatory PR events immediately before spawning git fetch, then
  pins the vetted DNS answers via http.curloptResolve and confines the
  subprocess with GIT_ALLOW_PROTOCOL=http:https,
  http.followRedirects=false, cleared proxy config/environment, and an
  empty credential helper, so redirects, proxies, or alternate
  protocols cannot escape the authorized target.

The policy enforces per-sink scheme allowlists (ws/wss for relays,
http/https for git), rejects embedded credentials and local hostnames
(localhost, single-label names, IANA special-use suffixes), and
requires IP literals and every DNS answer to be globally reachable.
Service admission (lists_service) and the don't-fetch-from-ourselves
filter now compare parsed host and port instead of substrings, so
gitnostr.com.attacker.example or a path containing the domain no
longer satisfies a check for gitnostr.com.

The operator-configured bootstrap relay stays usable even when local:
trust is carried by RelayTargetSource::OperatorConfigured at
construction, never by comparing event URLs against the configured
value, so event URLs that merely resemble the bootstrap relay are
still rejected. The new NGIT_SYNC_ALLOW_NON_GLOBAL_TARGETS option
(default false; documented in configuration.md, module.nix, and
.env.example) relaxes only the reachability checks for integration
tests and closed development networks; the TestRelay fixture sets it
because the test infrastructure lives on loopback, while the new
regression tests opt back into production behaviour.

Known limitation: nostr-sdk's connect API takes a URL, not a
pre-resolved address, so relay DNS is re-validated before every dial
but re-resolved by the SDK during connection, leaving a narrow
DNS-rebinding window (documented in defensive-measures.md). Git
fetches do not share this window because their DNS answers are pinned.
Closing it requires upstream connector support rather than a custom
connector here.

Validation: tests/outbound_policy.rs adds integration scenarios
against the real relay binary proving that loopback relay URLs and
loopback git clone URLs produce no outbound connection (counting TCP
listeners stand in for attacker infrastructure), that private,
link-local, CGNAT, unspecified, and multicast literals plus localhost
and credential URLs are rejected, that the local bootstrap relay still
connects while a resembling event URL is rejected, and that
substring-embedded domains are no longer admitted. src/outbound.rs
unit tests cover the reachability matrix (including 100.125.184.46)
and exact service matching. cargo fmt, cargo clippy (workspace, zero
warnings), the full cargo test suite, and cargo test -p grasp-audit
--lib all pass in the nix dev shell.
2026-08-01 19:39:36 +00:00

438 lines
18 KiB
Rust

//! In-process relay server lifecycle.
//!
//! [`RelayServer`] owns the full ngit-grasp runtime: it binds the TCP
//! listener (natively supporting `:0` for kernel-assigned ports), wires up
//! the nostr relay, purgatory, sync system, and background maintenance
//! tasks, then serves until a caller-supplied shutdown future resolves and
//! finally persists state to disk.
//!
//! The `ngit-grasp serve` binary (`src/main.rs`) drives this module,
//! supplying OS signal futures (SIGINT/SIGTERM) as the shutdown trigger;
//! `main.rs` itself only owns binary concerns (CLI/env config, tracing
//! subscriber, owner-key file, signal wiring).
use std::future::Future;
use std::net::SocketAddr;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use nostr_sdk::local_relay::LocalRelay;
use tokio::net::TcpListener;
use tokio::task::JoinHandle;
use tracing::{error, info, warn};
use crate::{
audit_cleanup,
config::{Config, DatabaseBackend},
git, grasp06, http,
metrics::Metrics,
nostr::{self, builder::Nip34WritePolicy, lifecycle::RepositoryLifecycle, SharedDatabase},
outbound::OutboundTargetPolicy,
purgatory::{sync::RealSyncContext, sync::ThrottleManager, Purgatory},
sync::{naughty_list::NaughtyListTracker, rejected_index::RejectedEventsIndex, SyncManager},
};
/// A fully-wired relay that has bound its listener but not yet started
/// accepting connections.
///
/// Construct with [`RelayServer::start`], inspect the actual bound address
/// with [`RelayServer::local_addr`], then drive it with
/// [`RelayServer::run_until`].
pub struct RelayServer {
config: Config,
listener: TcpListener,
local_addr: SocketAddr,
relay: LocalRelay,
database: SharedDatabase,
metrics: Option<Arc<Metrics>>,
purgatory: Arc<Purgatory>,
write_policy: Arc<Nip34WritePolicy>,
lifecycle: Arc<RepositoryLifecycle>,
rejected_events_index: Arc<RejectedEventsIndex>,
repo_init_locks: crate::grasp06::receive::RepoInitLocks,
deletion_cleanup: nostr::lifecycle::DeletionCleanupTask,
/// Detached background loops (sync manager, purgatory cleanup, audit
/// cleanup, purgatory sync loop). Aborted on shutdown so the host
/// process does not leak tasks per relay instance.
background_tasks: Vec<JoinHandle<()>>,
git_data_path: String,
}
impl RelayServer {
/// Validate the config, bind the listener, and wire up the full relay
/// runtime. Does not accept connections yet — call
/// [`RelayServer::run_until`] for that.
///
/// `config.bind_address` may use port `0` for a kernel-assigned port;
/// the resolved address is available via [`RelayServer::local_addr`]
/// and is written back into `config.bind_address`. If `config.domain`
/// is empty it is set to the resolved `host:port`.
///
/// Expects `config.relay_owner_nsec` to be set (see [`Config::load`])
/// and does **not** install a tracing subscriber — that is the
/// caller's concern.
pub async fn start(mut config: Config) -> Result<Self> {
// Bind first so kernel-assigned ports are resolved before any
// component captures the domain / bind address.
let requested: SocketAddr = config
.bind_address
.parse()
.with_context(|| format!("invalid bind address '{}'", config.bind_address))?;
let listener = TcpListener::bind(&requested)
.await
.with_context(|| format!("failed to bind {}", requested))?;
let local_addr = listener.local_addr()?;
config.bind_address = local_addr.to_string();
if config.domain.is_empty() {
config.domain = local_addr.to_string();
}
// Validate configuration and fail fast on fatal errors.
// Recoverable issues (e.g., malformed whitelist entries) are logged
// as warnings.
config.validate()?;
info!(
"Configuration loaded and validated: {}",
config.bind_address
);
info!("Domain: {}", config.domain);
info!("Relay name: {}", config.relay_name());
info!("Git data directory: {}", config.effective_git_data_path());
if config.database_backend != DatabaseBackend::Memory {
info!("Relay data directory: {}", config.relay_data_path);
}
info!("Database backend: {}", config.database_backend);
// Initialize metrics if enabled
let metrics = if config.metrics_enabled {
info!("Metrics enabled on /metrics endpoint");
let m = Arc::new(Metrics::new(
config.metrics_connection_per_ip_abuse_threshold,
Some(config.effective_git_data_path()),
));
info!("Repository count will be updated on each metrics request");
Some(m)
} else {
info!("Metrics disabled");
None
};
// Create purgatory for event/git coordination
let purgatory = Arc::new(Purgatory::new(PathBuf::from(
config.effective_git_data_path(),
)));
info!("Purgatory initialized for event coordination");
// Restore purgatory state from disk if available
let purgatory_path =
PathBuf::from(config.effective_git_data_path()).join("purgatory-state.json");
if purgatory_path.exists() {
match purgatory.restore_from_disk(&purgatory_path) {
Ok(()) => {
info!("Restored purgatory state from disk");
// Re-queueing will happen later after sync system is created
}
Err(e) => {
warn!("Failed to restore purgatory state: {}, starting empty", e);
}
}
}
// Shared per-path state for the GRASP-06 `/prs/` endpoint (mutex +
// in-flight counter, see [`grasp06::receive::PrsPathState`]). Lives
// for the lifetime of the process so concurrent pushes to the same
// `/prs/<submitter>/<id>.git` path — and policy / purgatory code
// paths that may delete a `/prs/` bare repo — coordinate through
// the same DashMap.
let repo_init_locks = grasp06::receive::new_repo_init_locks();
// Startup recovery for the GRASP-06 `/prs/` subtree: remove any
// zero-ref bare repos and empty submitter dirs left behind by a
// previous run (crashes mid-push, mid-cleanup, or shutdowns with
// unresolved scoped placeholders). Inline cleanup paths only fire
// while the process is running, so without this scan abandoned
// dirs persist indefinitely — the `cleanup-empty-repos` CLI tool
// explicitly skips `/prs/`. Gated on `grasp06_enable` so an
// operator who has turned the feature off does not start removing
// their existing `/prs/` data on the next restart. Runs before
// anything that could write to `/prs/` so no locking is needed.
if config.grasp06_enable {
let git_data_path = PathBuf::from(config.effective_git_data_path());
let (repos, dirs) = grasp06::cleanup::scan_on_startup(&git_data_path);
if repos > 0 || dirs > 0 {
info!(
"GRASP-06 /prs/ startup recovery: removed {} zero-ref repo(s), {} empty submitter dir(s)",
repos, dirs
);
}
}
// Create Nostr relay runtime with NIP-34 validation and shared stores.
let relay_runtime =
nostr::builder::create_relay(&config, purgatory.clone(), repo_init_locks.clone())
.await
.context("failed to create relay runtime")?;
info!(
"Relay created with NIP-34 validation for domain: {}",
config.domain
);
// Set the local relay on the write policy for purgatory notifications
// This must be done after relay creation since the relay depends on the policy
relay_runtime
.write_policy
.set_local_relay(relay_runtime.relay.clone());
let deletion_runtime = relay_runtime.deletion_runtime.clone();
deletion_runtime
.run_startup_tasks()
.await
.context("failed deletion lifecycle startup reconciliation")?;
// Wire the GRASP-06 `/prs/` filesystem cleanup context into
// purgatory so the standard expiry sweep can delete dangling
// refs/nostr/<event-id> refs (and zero-ref bare repos) when a
// scoped placeholder expires without a matching PR event. Held
// under the same per-path lock the receive handler uses so the
// cleanup can never race with an in-flight push.
purgatory.set_prs_cleanup_ctx(crate::purgatory::PrsCleanupCtx {
git_data_path: PathBuf::from(config.effective_git_data_path()),
repo_init_locks: repo_init_locks.clone(),
});
// 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
let sync_manager = SyncManager::new(
config.sync_bootstrap_relay_url.clone(),
config.domain.clone(),
relay_runtime.stores.database.clone(),
relay_runtime.write_policy.clone(),
relay_runtime.relay.clone(),
&config,
PathBuf::from(config.effective_git_data_path()),
metrics.as_ref().and_then(|m| m.sync_metrics().cloned()),
);
if config.sync_bootstrap_relay_url.is_some() {
info!(
"Starting proactive sync with bootstrap relay: {:?}",
config.sync_bootstrap_relay_url
);
} else {
info!("Proactive sync enabled (will discover relays from stored announcements)");
}
// Re-queue all restored purgatory repos for sync
let restored_identifiers = purgatory.get_all_identifiers();
if !restored_identifiers.is_empty() {
info!(
"Re-queueing {} restored repositories for sync",
restored_identifiers.len()
);
for identifier in restored_identifiers {
purgatory.enqueue_sync_immediate(&identifier);
}
}
// Get a reference to the rejected events index for shutdown persistence
// and for the HTTP server's git push path (hot-cache re-processing)
let rejected_events_index = sync_manager.rejected_events_index();
relay_runtime
.write_policy
.set_rejected_events_index(rejected_events_index.clone());
let mut background_tasks = Vec::new();
background_tasks.push(tokio::spawn(async move {
sync_manager.run().await;
}));
// Spawn background cleanup task for purgatory entries (60s interval)
let cleanup_purgatory = purgatory.clone();
background_tasks.push(tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(60));
loop {
interval.tick().await;
let (announcement_removed, state_removed, pr_removed) = cleanup_purgatory.cleanup();
if announcement_removed > 0 || state_removed > 0 || pr_removed > 0 {
info!(
"Purgatory cleanup: removed {} announcements, {} state events, {} PR events",
announcement_removed, state_removed, pr_removed
);
}
}
}));
info!("Purgatory cleanup task started (60s interval)");
// Spawn daily cleanup task for old expired event records (prevent unbounded growth)
let expired_cleanup_purgatory = purgatory.clone();
background_tasks.push(tokio::spawn(async move {
// Run immediately on startup, then every 24 hours
let mut interval = tokio::time::interval(Duration::from_secs(24 * 3600));
loop {
interval.tick().await;
// Remove expired event records older than 7 days
let removed = expired_cleanup_purgatory
.cleanup_expired_events(Duration::from_secs(7 * 24 * 3600));
if removed > 0 {
info!(
"Expired event cleanup: removed {} old expired event records (>7 days)",
removed
);
}
}
}));
info!("Expired event cleanup task started (24h interval, keeps 7 days)");
// Spawn audit event cleanup task (30m interval, removes events >2h old)
let audit_db = relay_runtime.stores.database.clone();
let audit_git_path = PathBuf::from(config.effective_git_data_path());
background_tasks.push(tokio::spawn(async move {
audit_cleanup::run_audit_cleanup_loop(audit_db, audit_git_path).await;
}));
info!("Audit event cleanup task started (30m interval, removes events >2h old)");
// Start purgatory sync loop for background git data fetching
// Create naughty list tracker for git remote domains with persistent errors (12h expiration)
let git_naughty_list = Arc::new(NaughtyListTracker::with_defaults());
let sync_ctx = Arc::new(RealSyncContext::new(
purgatory.clone(),
relay_runtime.stores.database.clone(),
PathBuf::from(config.effective_git_data_path()),
Some(config.domain.clone()),
Some(relay_runtime.relay.clone()),
Some(relay_runtime.write_policy.clone()),
git_naughty_list.clone(),
OutboundTargetPolicy {
allow_non_global: config.sync_allow_non_global_targets,
},
));
// Create throttle manager for rate limiting remote git servers
// Default: 5 concurrent requests per domain, 60 requests per minute per domain
let throttle_manager = Arc::new(ThrottleManager::new(5, 60));
throttle_manager.set_context(sync_ctx.clone());
throttle_manager.set_git_naughty_list(git_naughty_list.clone());
// Start the sync loop
background_tasks.push(purgatory.clone().start_sync_loop(
sync_ctx,
throttle_manager,
git_naughty_list.clone(),
));
info!("Purgatory sync loop started (1s interval)");
let deletion_cleanup = deletion_runtime.spawn_cleanup_task();
let git_data_path = config.effective_git_data_path();
Ok(Self {
relay: relay_runtime.relay,
database: relay_runtime.stores.database,
write_policy: Arc::new(relay_runtime.write_policy),
lifecycle: Arc::new(relay_runtime.lifecycle),
config,
listener,
local_addr,
metrics,
purgatory,
rejected_events_index,
repo_init_locks,
deletion_cleanup,
background_tasks,
git_data_path,
})
}
/// The actual bound socket address (with the kernel-assigned port when
/// the configured bind address used port `0`).
pub fn local_addr(&self) -> SocketAddr {
self.local_addr
}
/// The effective configuration (bind address and domain reflect the
/// resolved listener address).
pub fn config(&self) -> &Config {
&self.config
}
/// Serve HTTP + relay traffic until `shutdown` resolves (or the server
/// itself fails), then persist state and tear down background tasks.
///
/// The shutdown sequence mirrors the binary's signal handling: stop the
/// holding-cleanup task gracefully, save purgatory state and the
/// rejected-events cache to disk, remove placeholder `refs/nostr/`
/// refs, and abort the remaining background loops.
pub async fn run_until(self, shutdown: impl Future<Output = ()>) -> Result<()> {
info!("Starting HTTP server on {}", self.config.bind_address);
let serve = http::run_server_on_listener(
self.listener,
self.config,
self.relay,
self.database,
self.metrics,
self.purgatory.clone(),
self.write_policy,
self.lifecycle,
self.rejected_events_index.clone(),
self.repo_init_locks,
);
let result = tokio::select! {
result = serve => result,
_ = shutdown => {
info!("Shutdown signal received, cleaning up...");
Ok(())
}
};
self.deletion_cleanup.shutdown().await;
// Save purgatory state to disk
let purgatory_save_path = PathBuf::from(&self.git_data_path).join("purgatory-state.json");
if let Err(e) = self.purgatory.save_to_disk(&purgatory_save_path) {
error!("Failed to save purgatory state: {}", e);
} else {
info!("Purgatory state saved to disk");
}
// Save rejected events cache to disk
let rejected_cache_path =
PathBuf::from(&self.git_data_path).join("rejected-events-cache.json");
if let Err(e) = self
.rejected_events_index
.save_to_disk(&rejected_cache_path)
{
error!("Failed to save rejected events cache: {}", e);
} else {
info!("Rejected events cache saved to disk");
}
// Cleanup placeholder refs on shutdown
let placeholder_ids = self.purgatory.get_placeholder_event_ids();
if !placeholder_ids.is_empty() {
info!(
"Cleaning up {} placeholder refs/nostr/ refs on shutdown",
placeholder_ids.len()
);
git::cleanup_placeholder_refs(&self.git_data_path, &placeholder_ids);
}
// Abort detached background loops so the host process does not
// accumulate live tasks (and the state they capture) per instance.
for task in self.background_tasks {
task.abort();
}
result
}
}