mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
refactor: extract relay runtime from main.rs into RelayServer
Salvaged from the closed in-process embedding PR (44eea10b) — the embedding API itself was rejected as unneeded complexity, but two pieces stand on their own: - src/server.rs: RelayServer owns the full relay runtime — binds the listener (supporting :0 with the resolved address back-filling bind_address and, when empty, domain), wires the nostr relay, purgatory, sync system, and background maintenance tasks, serves until a caller-supplied shutdown future resolves, then persists state and stops background tasks - src/main.rs: now only owns binary concerns (CLI/env config, tracing subscriber, SIGINT/SIGTERM wiring) - fix: create_relay errors were silently swallowed at startup (`if let Ok(...)` dropped the error and exited 0); they now propagate as a startup failure - src/http: split run_server so the runtime can serve on a pre-bound listener Deliberately not included from the PR: embed.rs, Config:: embedded_defaults, the programmatic fast-timers override, and the embedding docs/tests.
This commit is contained in:
@@ -48,25 +48,38 @@
|
||||
|
||||
## Component Design
|
||||
|
||||
### 1. Main Server ([`src/main.rs`](src/main.rs))
|
||||
### 1. Main Server ([`src/main.rs`](src/main.rs) + [`src/server.rs`](src/server.rs))
|
||||
|
||||
**Responsibilities:**
|
||||
Server startup is split across two layers: `main.rs` owns binary-only
|
||||
concerns, while `RelayServer` in `server.rs` owns the reusable relay
|
||||
runtime:
|
||||
|
||||
- Initialize configuration from environment (clap + dotenvy)
|
||||
- Set up Hyper HTTP server with request routing
|
||||
**[`src/main.rs`](src/main.rs) — binary-only concerns:**
|
||||
|
||||
- Initialize configuration from CLI/environment (clap + dotenvy)
|
||||
- Install the global tracing subscriber
|
||||
- Load/generate the relay owner key file (`.relay-owner.nsec`)
|
||||
- Translate OS signals (SIGINT/SIGTERM) into a shutdown future
|
||||
|
||||
**[`src/server.rs`](src/server.rs) — `RelayServer`, the shared runtime:**
|
||||
|
||||
- Bind the TCP listener (supports `:0` for kernel-assigned ports; the
|
||||
resolved address back-fills `bind_address` and, if empty, `domain`)
|
||||
- Initialize Nostr relay builder with custom [`Nip34WritePolicy`](src/nostr/builder.rs:51)
|
||||
- Set up shared storage (LMDB or Memory)
|
||||
- Handle WebSocket upgrades for Nostr relay
|
||||
- Handle graceful shutdown
|
||||
- Set up shared storage (LMDB or Memory), purgatory, sync manager, and
|
||||
background maintenance tasks
|
||||
- Serve HTTP + WebSocket until a caller-supplied shutdown future
|
||||
resolves, then persist state (purgatory, rejected-events cache,
|
||||
placeholder ref cleanup) and stop background tasks
|
||||
|
||||
**Key Dependencies:**
|
||||
|
||||
```rust
|
||||
hyper = "1"
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
nostr-relay-builder = "0.43"
|
||||
nostr-sdk = "0.43"
|
||||
nostr-lmdb = "0.43"
|
||||
nostr-relay-builder = "0.45.0-alpha.3"
|
||||
nostr-sdk = "0.45.0-alpha.3"
|
||||
nostr-lmdb = "0.45.0-alpha.3"
|
||||
```
|
||||
|
||||
### 2. HTTP Module ([`src/http/mod.rs`](src/http/mod.rs))
|
||||
|
||||
+46
-12
@@ -843,15 +843,11 @@ fn derive_accept_key(request_key: &[u8]) -> String {
|
||||
|
||||
/// Start the HTTP server with integrated Nostr relay
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `config` - Server configuration
|
||||
/// * `relay` - The LocalRelay for WebSocket connections
|
||||
/// * `database` - The database for direct queries (e.g., push authorization)
|
||||
/// * `metrics` - Optional metrics for Prometheus endpoint
|
||||
/// * `purgatory` - Purgatory for event/git coordination
|
||||
/// * `write_policy` - Write policy for re-processing hot-cache events after git push promotion
|
||||
/// * `lifecycle` - Repository lifecycle locks for Git HTTP serving
|
||||
/// * `rejected_events_index` - Rejected events index for hot-cache re-processing
|
||||
/// Binds a fresh [`TcpListener`] on `config.bind_address` and delegates to
|
||||
/// [`run_server_on_listener`]. Callers that need the actual bound address
|
||||
/// before serving (e.g. binding `:0` for a kernel-assigned port, as
|
||||
/// `RelayServer` does) should bind the listener themselves and call
|
||||
/// [`run_server_on_listener`] directly.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn run_server(
|
||||
config: Config,
|
||||
@@ -865,13 +861,51 @@ pub async fn run_server(
|
||||
repo_init_locks: RepoInitLocks,
|
||||
) -> anyhow::Result<()> {
|
||||
let bind_addr: SocketAddr = config.bind_address.parse()?;
|
||||
let listener = TcpListener::bind(&bind_addr).await?;
|
||||
run_server_on_listener(
|
||||
listener,
|
||||
config,
|
||||
relay,
|
||||
database,
|
||||
metrics,
|
||||
purgatory,
|
||||
write_policy,
|
||||
lifecycle,
|
||||
rejected_events_index,
|
||||
repo_init_locks,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
tracing::info!("Starting HTTP server on {}", bind_addr);
|
||||
/// Serve HTTP + WebSocket relay traffic on an already-bound listener.
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `listener` - Pre-bound TCP listener to accept connections on
|
||||
/// * `config` - Server configuration
|
||||
/// * `relay` - The LocalRelay for WebSocket connections
|
||||
/// * `database` - The database for direct queries (e.g., push authorization)
|
||||
/// * `metrics` - Optional metrics for Prometheus endpoint
|
||||
/// * `purgatory` - Purgatory for event/git coordination
|
||||
/// * `write_policy` - Write policy for re-processing hot-cache events after git push promotion
|
||||
/// * `lifecycle` - Repository lifecycle locks for Git HTTP serving
|
||||
/// * `rejected_events_index` - Rejected events index for hot-cache re-processing
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub async fn run_server_on_listener(
|
||||
listener: TcpListener,
|
||||
config: Config,
|
||||
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: RepoInitLocks,
|
||||
) -> anyhow::Result<()> {
|
||||
tracing::info!("Starting HTTP server on {}", listener.local_addr()?);
|
||||
tracing::info!("Relay name: {}", config.relay_name());
|
||||
tracing::info!("Domain: {}", config.domain);
|
||||
|
||||
let listener = TcpListener::bind(&bind_addr).await?;
|
||||
|
||||
loop {
|
||||
let (socket, addr) = listener.accept().await?;
|
||||
let io = TokioIo::new(socket);
|
||||
|
||||
@@ -8,4 +8,5 @@ pub mod metrics;
|
||||
pub mod nostr;
|
||||
pub mod purgatory;
|
||||
pub mod repair_deletion_requests;
|
||||
pub mod server;
|
||||
pub mod sync;
|
||||
|
||||
+42
-331
@@ -1,21 +1,11 @@
|
||||
use std::time::Duration;
|
||||
use std::{path::PathBuf, sync::Arc};
|
||||
|
||||
use anyhow::Result;
|
||||
use clap::Parser;
|
||||
use tokio::signal;
|
||||
use tracing::{error, info, warn};
|
||||
use tracing::info;
|
||||
use tracing_subscriber::{EnvFilter, FmtSubscriber};
|
||||
|
||||
use ngit_grasp::{
|
||||
audit_cleanup, cleanup_empty_repos,
|
||||
config::{Config, DatabaseBackend},
|
||||
git, grasp06, http,
|
||||
metrics::Metrics,
|
||||
nostr,
|
||||
purgatory::{sync::RealSyncContext, sync::ThrottleManager, Purgatory},
|
||||
repair_deletion_requests,
|
||||
sync::{naughty_list::NaughtyListTracker, SyncManager},
|
||||
cleanup_empty_repos, config::Config, nostr, repair_deletion_requests, server::RelayServer,
|
||||
};
|
||||
|
||||
/// Top-level CLI dispatcher.
|
||||
@@ -90,6 +80,11 @@ async fn main() -> Result<()> {
|
||||
}
|
||||
}
|
||||
|
||||
/// Run the relay until an OS shutdown signal arrives.
|
||||
///
|
||||
/// All relay wiring lives in [`ngit_grasp::server::RelayServer`];
|
||||
/// this function only owns concerns specific to the standalone binary:
|
||||
/// the global tracing subscriber and signal handling.
|
||||
async fn run_relay(config: Config) -> Result<()> {
|
||||
// Initialize tracing with configured log level
|
||||
let subscriber = FmtSubscriber::builder()
|
||||
@@ -99,324 +94,40 @@ async fn run_relay(config: Config) -> Result<()> {
|
||||
|
||||
info!("Starting ngit-grasp with log level: {}", config.log_level);
|
||||
|
||||
// Validate configuration and fail fast on fatal errors
|
||||
// Recoverable issues (e.g., malformed whitelist entries) are logged as warnings
|
||||
config.validate()?;
|
||||
let server = RelayServer::start(config).await?;
|
||||
|
||||
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.
|
||||
if let Ok(relay_runtime) =
|
||||
nostr::builder::create_relay(&config, purgatory.clone(), repo_init_locks.clone()).await
|
||||
{
|
||||
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;
|
||||
|
||||
// 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(ngit_grasp::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 shutdown_rejected_index = sync_manager.rejected_events_index();
|
||||
let http_rejected_index = shutdown_rejected_index.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
sync_manager.run().await;
|
||||
});
|
||||
|
||||
// Spawn background cleanup task for purgatory entries (60s interval)
|
||||
let cleanup_purgatory = purgatory.clone();
|
||||
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();
|
||||
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());
|
||||
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(),
|
||||
));
|
||||
|
||||
// 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
|
||||
let _sync_loop_handle =
|
||||
purgatory
|
||||
.clone()
|
||||
.start_sync_loop(sync_ctx, throttle_manager, git_naughty_list.clone());
|
||||
info!("Purgatory sync loop started (1s interval)");
|
||||
|
||||
// Setup shutdown handler for purgatory cleanup
|
||||
let shutdown_purgatory = purgatory.clone();
|
||||
let git_data_path = config.effective_git_data_path();
|
||||
|
||||
// Start HTTP server with integrated relay and database
|
||||
info!("Starting HTTP server on {}", config.bind_address);
|
||||
|
||||
let deletion_cleanup = deletion_runtime.spawn_cleanup_task();
|
||||
|
||||
// Wrap write_policy in Arc for sharing between HTTP server connections
|
||||
let http_write_policy = Arc::new(relay_runtime.write_policy.clone());
|
||||
let http_lifecycle = Arc::new(relay_runtime.lifecycle.clone());
|
||||
|
||||
// Run server until shutdown signal, then cleanup
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use tokio::signal::unix::{signal, SignalKind};
|
||||
let mut sigterm = signal(SignalKind::terminate())?;
|
||||
|
||||
tokio::select! {
|
||||
result = http::run_server(
|
||||
config,
|
||||
relay_runtime.relay,
|
||||
relay_runtime.stores.database,
|
||||
metrics,
|
||||
purgatory,
|
||||
http_write_policy,
|
||||
http_lifecycle,
|
||||
http_rejected_index,
|
||||
repo_init_locks,
|
||||
) => {
|
||||
result?
|
||||
}
|
||||
_ = signal::ctrl_c() => {
|
||||
info!("Received SIGINT (Ctrl+C), cleaning up...");
|
||||
}
|
||||
_ = sigterm.recv() => {
|
||||
info!("Received SIGTERM, cleaning up...");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
{
|
||||
tokio::select! {
|
||||
result = http::run_server(
|
||||
config,
|
||||
relay_runtime.relay,
|
||||
relay_runtime.stores.database,
|
||||
metrics,
|
||||
purgatory,
|
||||
http_write_policy,
|
||||
http_lifecycle,
|
||||
http_rejected_index,
|
||||
repo_init_locks,
|
||||
) => {
|
||||
result?
|
||||
}
|
||||
_ = signal::ctrl_c() => {
|
||||
info!("Received SIGINT (Ctrl+C), cleaning up...");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
deletion_cleanup.shutdown().await;
|
||||
|
||||
// Save purgatory state to disk
|
||||
let purgatory_save_path = PathBuf::from(&git_data_path).join("purgatory-state.json");
|
||||
if let Err(e) = shutdown_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(&git_data_path).join("rejected-events-cache.json");
|
||||
if let Err(e) = shutdown_rejected_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 = shutdown_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(&git_data_path, &placeholder_ids);
|
||||
}
|
||||
}
|
||||
|
||||
Ok(())
|
||||
server.run_until(shutdown_signal()).await
|
||||
}
|
||||
|
||||
/// Resolves when the process receives SIGINT (Ctrl+C) or, on unix, SIGTERM.
|
||||
async fn shutdown_signal() {
|
||||
#[cfg(unix)]
|
||||
{
|
||||
use tokio::signal::unix::{signal as unix_signal, SignalKind};
|
||||
let mut sigterm = match unix_signal(SignalKind::terminate()) {
|
||||
Ok(sigterm) => sigterm,
|
||||
Err(e) => {
|
||||
tracing::error!("Failed to install SIGTERM handler: {}", e);
|
||||
// Fall back to Ctrl+C only.
|
||||
let _ = signal::ctrl_c().await;
|
||||
info!("Received SIGINT (Ctrl+C), cleaning up...");
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
tokio::select! {
|
||||
_ = signal::ctrl_c() => {
|
||||
info!("Received SIGINT (Ctrl+C), cleaning up...");
|
||||
}
|
||||
_ = sigterm.recv() => {
|
||||
info!("Received SIGTERM, cleaning up...");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(unix))]
|
||||
{
|
||||
let _ = signal::ctrl_c().await;
|
||||
info!("Received SIGINT (Ctrl+C), cleaning up...");
|
||||
}
|
||||
}
|
||||
|
||||
+427
@@ -0,0 +1,427 @@
|
||||
//! 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_relay_builder::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},
|
||||
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;
|
||||
|
||||
// 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();
|
||||
|
||||
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(),
|
||||
));
|
||||
|
||||
// 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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user