mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Reverse-proxied production deployments currently attribute every WebSocket connection to the proxy peer, collapsing per-IP limits, abuse metrics, and logs onto localhost. Blindly trusting forwarding headers would let direct clients spoof the same controls. Add an opt-in CIDR trust boundary and resolve X-Forwarded-For, Forwarded, or X-Real-IP only when the TCP peer is trusted. Walk proxy chains from the peer inward, stop at the first untrusted hop, and fall back to the peer on malformed input. Feed the resolved address consistently into rust-nostr connection policy, connection metrics, and lifecycle logs. Expose the setting across CLI/environment, NixOS, examples, and reference documentation. Multi-hop correctness assumes every trusted proxy appends or overwrites the forwarding chain and the backend is unreachable from untrusted networks. HTTP Git request accounting and PROXY protocol support remain out of scope. Unit coverage exercises direct clients, trusted single and multi-proxy paths, spoofed headers, malformed chains, IPv6, and configuration parsing.
963 lines
41 KiB
Rust
963 lines
41 KiB
Rust
//! HTTP Server Module
|
|
//!
|
|
//! Provides hyper HTTP server with WebSocket upgrade support for the Nostr relay.
|
|
mod client_ip;
|
|
pub mod landing;
|
|
pub mod nip11;
|
|
|
|
use std::future::Future;
|
|
use std::net::SocketAddr;
|
|
use std::pin::Pin;
|
|
use std::sync::Arc;
|
|
|
|
use base64::Engine;
|
|
use bitcoin_hashes::sha1::Hash as Sha1Hash;
|
|
use bitcoin_hashes::{Hash, HashEngine};
|
|
use http_body_util::BodyExt;
|
|
use hyper::body::{Bytes, Incoming};
|
|
use hyper::header::{CONNECTION, SEC_WEBSOCKET_ACCEPT, UPGRADE};
|
|
use hyper::server::conn::http1;
|
|
use hyper::service::Service;
|
|
use hyper::{Method, Request, Response};
|
|
use hyper_util::rt::TokioIo;
|
|
use nostr_sdk::local_relay::LocalRelay;
|
|
use nostr_sdk::prelude::PublicKey;
|
|
use tokio::net::TcpListener;
|
|
|
|
use crate::config::Config;
|
|
use crate::git;
|
|
use crate::git::{empty_body, full_body, GitResponseBody};
|
|
use crate::grasp06::receive::RepoInitLocks;
|
|
use crate::metrics::Metrics;
|
|
use crate::nostr::builder::Nip34WritePolicy;
|
|
use crate::nostr::lifecycle::RepositoryLifecycle;
|
|
use crate::nostr::SharedDatabase;
|
|
use crate::purgatory::promotion_hooks::NostrPurgatoryPromotionHooks;
|
|
use crate::purgatory::Purgatory;
|
|
use crate::sync::rejected_index::RejectedEventsIndex;
|
|
|
|
type HttpBody = GitResponseBody;
|
|
|
|
/// CORS headers required by GRASP-01 specification (lines 48-51)
|
|
const CORS_ALLOW_ORIGIN: &str = "*";
|
|
const CORS_ALLOW_METHODS: &str = "GET, POST";
|
|
const CORS_ALLOW_HEADERS: &str = "Content-Type";
|
|
|
|
/// Embedded icon image (Grasp logo)
|
|
const ICON_PNG: &[u8] = include_bytes!("../../static/icon.png");
|
|
|
|
/// Extract npub and identifier from a repository URL path (no git subpath required)
|
|
///
|
|
/// Parses paths like `/<npub>/<identifier>.git` (for repository webpage/404)
|
|
///
|
|
/// The identifier is percent-decoded so that URLs like `/npub1.../my%20repo.git`
|
|
/// resolve to the correct filesystem path.
|
|
///
|
|
/// Returns (npub, identifier) if the path matches a repository URL pattern
|
|
fn parse_repo_url(path: &str) -> Option<(String, String)> {
|
|
// Defensive guard: never match the GRASP-06 `/prs/` namespace as a
|
|
// standard repo URL. `HttpService::call` already intercepts `/prs/*`
|
|
// before reaching the landing branch when grasp06_enable is on, but
|
|
// this early-return ensures a future routing change cannot regress us
|
|
// into serving repo landing HTML at `/prs/<npub>/<id>.git` regardless
|
|
// of the feature flag.
|
|
if path.starts_with("/prs/") || path.starts_with("prs/") {
|
|
return None;
|
|
}
|
|
|
|
// Remove leading slash
|
|
let path = path.strip_prefix('/').unwrap_or(path);
|
|
|
|
// Split into components
|
|
let parts: Vec<&str> = path.split('/').collect();
|
|
|
|
// Must be exactly 2 parts: npub and repo.git (no subpath)
|
|
if parts.len() != 2 {
|
|
return None;
|
|
}
|
|
|
|
let npub = parts[0];
|
|
let repo_part = git::percent_decode(parts[1]);
|
|
|
|
// The repo part must end with .git
|
|
if !repo_part.ends_with(".git") {
|
|
return None;
|
|
}
|
|
|
|
// Must have an npub that looks valid (starts with npub1)
|
|
if !npub.starts_with("npub1") {
|
|
return None;
|
|
}
|
|
|
|
// Extract identifier (remove .git suffix)
|
|
let identifier = repo_part
|
|
.strip_suffix(".git")
|
|
.unwrap_or(&repo_part)
|
|
.to_string();
|
|
|
|
// Identifier must not be empty
|
|
if identifier.is_empty() {
|
|
return None;
|
|
}
|
|
|
|
Some((npub.to_string(), identifier))
|
|
}
|
|
|
|
/// Add CORS headers to a response builder
|
|
fn add_cors_headers(builder: hyper::http::response::Builder) -> hyper::http::response::Builder {
|
|
builder
|
|
.header("Access-Control-Allow-Origin", CORS_ALLOW_ORIGIN)
|
|
.header("Access-Control-Allow-Methods", CORS_ALLOW_METHODS)
|
|
.header("Access-Control-Allow-Headers", CORS_ALLOW_HEADERS)
|
|
}
|
|
|
|
/// HTTP Service that serves both WebSocket (relay) and HTML landing page
|
|
struct HttpService {
|
|
relay: LocalRelay,
|
|
config: Config,
|
|
remote: SocketAddr,
|
|
/// Database reference for direct queries (e.g., push authorization)
|
|
database: SharedDatabase,
|
|
/// Optional metrics for Prometheus endpoint
|
|
metrics: Option<Arc<Metrics>>,
|
|
/// Purgatory for event/git coordination
|
|
purgatory: Arc<Purgatory>,
|
|
/// Write policy for re-processing hot-cache events after git push promotion
|
|
write_policy: Arc<Nip34WritePolicy>,
|
|
/// Repository lifecycle locks shared with deletion/recovery.
|
|
lifecycle: Arc<RepositoryLifecycle>,
|
|
/// Rejected events index for hot-cache re-processing after git push promotion
|
|
rejected_events_index: Arc<RejectedEventsIndex>,
|
|
/// Per-path init mutexes for GRASP-06 `/prs/` on-demand bare-repo
|
|
/// creation. See [`crate::grasp06::receive::RepoInitLocks`].
|
|
repo_init_locks: RepoInitLocks,
|
|
}
|
|
|
|
impl HttpService {
|
|
#[allow(clippy::too_many_arguments)]
|
|
fn new(
|
|
relay: LocalRelay,
|
|
config: Config,
|
|
remote: SocketAddr,
|
|
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,
|
|
) -> Self {
|
|
Self {
|
|
relay,
|
|
config,
|
|
remote,
|
|
database,
|
|
metrics,
|
|
purgatory,
|
|
write_policy,
|
|
lifecycle,
|
|
rejected_events_index,
|
|
repo_init_locks,
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Service<Request<Incoming>> for HttpService {
|
|
type Response = Response<HttpBody>;
|
|
type Error = String;
|
|
type Future = Pin<Box<dyn Future<Output = Result<Self::Response, Self::Error>> + Send>>;
|
|
|
|
fn call(&self, req: Request<Incoming>) -> Self::Future {
|
|
let base = add_cors_headers(Response::builder().header("server", "ngit-grasp"));
|
|
let path = req.uri().path().to_string();
|
|
let query = req.uri().query().map(|s| s.to_string());
|
|
let method = req.method().clone();
|
|
let git_data_path = self.config.effective_git_data_path();
|
|
let database = self.database.clone();
|
|
let purgatory = self.purgatory.clone();
|
|
let write_policy = self.write_policy.clone();
|
|
let lifecycle = self.lifecycle.clone();
|
|
let rejected_events_index = self.rejected_events_index.clone();
|
|
let repo_init_locks = self.repo_init_locks.clone();
|
|
|
|
// Handle OPTIONS preflight requests (CORS)
|
|
// GRASP-01 spec line 51: Respond to OPTIONS with 204 No Content
|
|
if method == Method::OPTIONS {
|
|
return Box::pin(async move {
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(204)
|
|
.body(empty_body())
|
|
.unwrap(),
|
|
)
|
|
});
|
|
}
|
|
|
|
// GRASP-06: route /prs/<npub>/<id>.git/* before the standard git URL
|
|
// parser. When disabled, the path falls through to existing 404
|
|
// handling (preserving the discovery-gate contract).
|
|
if self.config.grasp06_enable {
|
|
if let Some(prs) = crate::grasp06::endpoint::parse_prs_url(&path) {
|
|
let git_protocol = req
|
|
.headers()
|
|
.get("git-protocol")
|
|
.and_then(|v| v.to_str().ok())
|
|
.map(|s| s.to_string());
|
|
let content_encoding = req
|
|
.headers()
|
|
.get("content-encoding")
|
|
.and_then(|v| v.to_str().ok())
|
|
.map(|s| s.to_lowercase());
|
|
|
|
tracing::debug!(
|
|
"/prs/ request: {} {} (submitter={}, id={}, subpath={}, protocol={:?})",
|
|
method,
|
|
path,
|
|
prs.submitter.to_hex(),
|
|
prs.identifier,
|
|
prs.subpath,
|
|
git_protocol
|
|
);
|
|
|
|
let subpath = prs.subpath.clone();
|
|
let method_clone = method.clone();
|
|
let metrics_clone = self.metrics.clone();
|
|
let relay_clone = self.relay.clone();
|
|
let config_clone = self.config.clone();
|
|
|
|
return Box::pin(async move {
|
|
// Collect (and gunzip if needed) the request body just like
|
|
// the standard git branch does.
|
|
let raw_body = req
|
|
.collect()
|
|
.await
|
|
.map(|collected| collected.to_bytes())
|
|
.unwrap_or_else(|_| Bytes::new());
|
|
let body_bytes = if content_encoding.as_deref() == Some("gzip") {
|
|
use flate2::read::GzDecoder;
|
|
use std::io::Read;
|
|
let mut decoder = GzDecoder::new(&raw_body[..]);
|
|
let mut decompressed = Vec::new();
|
|
match decoder.read_to_end(&mut decompressed) {
|
|
Ok(_) => Bytes::from(decompressed),
|
|
Err(e) => {
|
|
tracing::warn!("/prs/ gzip decompress failed: {}", e);
|
|
raw_body
|
|
}
|
|
}
|
|
} else {
|
|
raw_body
|
|
};
|
|
|
|
let result: Result<Response<HttpBody>, git::handlers::GitError> =
|
|
match (method_clone.as_ref(), subpath.as_str()) {
|
|
// GET|HEAD /info/refs?service=git-{upload,receive}-pack
|
|
// HEAD must mirror GET headers with no body (RFC 9110 §9.3.2)
|
|
(m, sp)
|
|
if (m == Method::GET || m == Method::HEAD)
|
|
&& sp.starts_with("info/refs") =>
|
|
{
|
|
let service = query
|
|
.as_deref()
|
|
.unwrap_or("")
|
|
.strip_prefix("service=")
|
|
.and_then(git::protocol::GitService::from_query_param);
|
|
match service {
|
|
Some(svc) => {
|
|
let r = crate::grasp06::fetch::handle_prs_info_refs(
|
|
&prs,
|
|
&git_data_path,
|
|
svc,
|
|
git_protocol.as_deref(),
|
|
)
|
|
.await;
|
|
if let Some(ref m) = metrics_clone {
|
|
let status =
|
|
if r.is_ok() { "success" } else { "error" };
|
|
let op = match svc {
|
|
git::protocol::GitService::UploadPack => "fetch",
|
|
git::protocol::GitService::ReceivePack => "push",
|
|
};
|
|
m.record_git_operation(op, status);
|
|
}
|
|
r
|
|
}
|
|
None => Err(git::handlers::GitError::RepositoryNotFound),
|
|
}
|
|
}
|
|
|
|
// POST /git-upload-pack — clone/fetch.
|
|
(m, "git-upload-pack") if m == Method::POST => {
|
|
let streams_real_repo = crate::grasp06::paths::prs_repo_path(
|
|
std::path::Path::new(&git_data_path),
|
|
&prs.submitter.to_hex(),
|
|
&prs.identifier,
|
|
)
|
|
.exists();
|
|
let r = crate::grasp06::fetch::handle_prs_upload_pack(
|
|
&prs,
|
|
&git_data_path,
|
|
body_bytes,
|
|
git_protocol.as_deref(),
|
|
metrics_clone.clone(),
|
|
)
|
|
.await;
|
|
if let Some(ref m) = metrics_clone {
|
|
let status = if r.is_ok() { "success" } else { "error" };
|
|
if streams_real_repo && r.is_ok() {
|
|
// Real /prs/ repos delegate to the standard streaming
|
|
// upload-pack handler. The streaming task records the final
|
|
// git-upload-pack success/failure when Git exits; this
|
|
// `Ok(Response)` only means Hyper received a body stream.
|
|
} else {
|
|
// Synthesized empty repos are still buffered, so `Ok` here
|
|
// means git-upload-pack has already completed. For streamed
|
|
// setup errors, there is no task to record the failure.
|
|
m.record_git_operation("clone", status);
|
|
}
|
|
}
|
|
r
|
|
}
|
|
|
|
// POST /git-receive-pack — accept pushes to
|
|
// refs/nostr/<event-id>, reject anything else,
|
|
// per GRASP-06 06.md line 13.
|
|
(m, "git-receive-pack") if m == Method::POST => {
|
|
// Metrics are recorded inside the `/prs/`
|
|
// handler and its detached streaming task. For
|
|
// successful setup this `await` only means the
|
|
// response stream was created, not that Git has
|
|
// finished receiving the push.
|
|
let r = crate::grasp06::receive::handle_prs_receive_pack(
|
|
&prs,
|
|
body_bytes,
|
|
database.clone(),
|
|
relay_clone.clone(),
|
|
purgatory.clone(),
|
|
write_policy.clone(),
|
|
rejected_events_index.clone(),
|
|
&git_data_path,
|
|
git_protocol.as_deref(),
|
|
repo_init_locks.clone(),
|
|
&config_clone.domain,
|
|
metrics_clone.clone(),
|
|
)
|
|
.await;
|
|
r
|
|
}
|
|
|
|
_ => Err(git::handlers::GitError::RepositoryNotFound),
|
|
};
|
|
|
|
match result {
|
|
Ok(response) => {
|
|
let (parts, body) = response.into_parts();
|
|
// RFC 9110 §9.3.2: HEAD response must have same headers as GET
|
|
// but no body.
|
|
let body = if method_clone == Method::HEAD {
|
|
empty_body()
|
|
} else {
|
|
body
|
|
};
|
|
Ok(add_cors_headers(Response::builder().status(parts.status))
|
|
.header(
|
|
"content-type",
|
|
parts
|
|
.headers
|
|
.get("content-type")
|
|
.and_then(|v| v.to_str().ok())
|
|
.unwrap_or("application/octet-stream"),
|
|
)
|
|
.header(
|
|
"cache-control",
|
|
parts
|
|
.headers
|
|
.get("cache-control")
|
|
.and_then(|v| v.to_str().ok())
|
|
.unwrap_or("no-cache"),
|
|
)
|
|
.body(body)
|
|
.unwrap())
|
|
}
|
|
Err(e) => {
|
|
let error_msg = format!("Git error: {}", e);
|
|
Ok(add_cors_headers(Response::builder())
|
|
.status(e.status_code())
|
|
.body(full_body(error_msg))
|
|
.unwrap())
|
|
}
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
// Check for Git HTTP requests first
|
|
if let Some((npub, identifier, subpath)) = git::parse_git_url(&path) {
|
|
// Extract Git-Protocol header for protocol v2 support
|
|
let git_protocol = req
|
|
.headers()
|
|
.get("git-protocol")
|
|
.and_then(|v| v.to_str().ok())
|
|
.map(|s| s.to_string());
|
|
|
|
// Extract Content-Encoding header to handle gzip-compressed request bodies
|
|
// Modern git clients send gzip-compressed POST bodies for efficiency
|
|
let content_encoding = req
|
|
.headers()
|
|
.get("content-encoding")
|
|
.and_then(|v| v.to_str().ok())
|
|
.map(|s| s.to_lowercase());
|
|
|
|
tracing::debug!(
|
|
"Git request: {} {} (npub={}, id={}, subpath={}, protocol={:?}, encoding={:?})",
|
|
method,
|
|
path,
|
|
npub,
|
|
identifier,
|
|
subpath,
|
|
git_protocol,
|
|
content_encoding
|
|
);
|
|
|
|
let repo_path = git::resolve_repo_path(&git_data_path, &npub, &identifier);
|
|
let metrics_clone = self.metrics.clone();
|
|
let relay = self.relay.clone();
|
|
|
|
return Box::pin(async move {
|
|
// Collect request body once before the match statement
|
|
let raw_body = req
|
|
.collect()
|
|
.await
|
|
.map(|collected| collected.to_bytes())
|
|
.unwrap_or_else(|_| Bytes::new());
|
|
|
|
// Decompress gzip-encoded request bodies
|
|
// Git clients send Content-Encoding: gzip for POST requests
|
|
let body_bytes = if content_encoding.as_deref() == Some("gzip") {
|
|
use flate2::read::GzDecoder;
|
|
use std::io::Read;
|
|
|
|
let mut decoder = GzDecoder::new(&raw_body[..]);
|
|
let mut decompressed = Vec::new();
|
|
match decoder.read_to_end(&mut decompressed) {
|
|
Ok(_) => {
|
|
tracing::debug!(
|
|
"Decompressed gzip body: {} -> {} bytes",
|
|
raw_body.len(),
|
|
decompressed.len()
|
|
);
|
|
Bytes::from(decompressed)
|
|
}
|
|
Err(e) => {
|
|
tracing::warn!("Failed to decompress gzip body: {}", e);
|
|
// Fall back to raw body (might work if not actually gzip)
|
|
raw_body
|
|
}
|
|
}
|
|
} else {
|
|
raw_body
|
|
};
|
|
|
|
let result = match (method.as_ref(), subpath.as_str()) {
|
|
// GET|HEAD /info/refs?service=git-upload-pack or git-receive-pack
|
|
// HEAD must mirror GET headers with no body (RFC 9110 §9.3.2)
|
|
(m, sp)
|
|
if (m == Method::GET || m == Method::HEAD)
|
|
&& sp.starts_with("info/refs") =>
|
|
{
|
|
// Parse query string for service parameter
|
|
let service = query
|
|
.as_deref()
|
|
.unwrap_or("")
|
|
.strip_prefix("service=")
|
|
.and_then(git::protocol::GitService::from_query_param);
|
|
|
|
match service {
|
|
Some(svc) => {
|
|
let _repo_lifecycle_guard = match PublicKey::parse(&npub) {
|
|
Ok(owner_pk) => Some(
|
|
lifecycle
|
|
.read_repository(&owner_pk.to_hex(), &identifier)
|
|
.await,
|
|
),
|
|
Err(_) => None,
|
|
};
|
|
let result = git::handlers::handle_info_refs(
|
|
repo_path,
|
|
svc,
|
|
git_protocol.as_deref(),
|
|
)
|
|
.await;
|
|
// Track operation
|
|
if let Some(ref m) = metrics_clone {
|
|
let status = if result.is_ok() { "success" } else { "error" };
|
|
let operation = match svc {
|
|
git::protocol::GitService::UploadPack => "fetch",
|
|
git::protocol::GitService::ReceivePack => "push",
|
|
};
|
|
m.record_git_operation(operation, status);
|
|
}
|
|
result
|
|
}
|
|
None => Err(git::handlers::GitError::RepositoryNotFound),
|
|
}
|
|
}
|
|
|
|
// POST /git-upload-pack (clone/fetch)
|
|
(m, "git-upload-pack") if m == Method::POST => {
|
|
let repo_lifecycle_guard = match PublicKey::parse(&npub) {
|
|
Ok(owner_pk) => Some(
|
|
lifecycle
|
|
.read_repository(&owner_pk.to_hex(), &identifier)
|
|
.await,
|
|
),
|
|
Err(_) => None,
|
|
};
|
|
let result = git::handlers::handle_upload_pack(
|
|
repo_path,
|
|
body_bytes,
|
|
git_protocol.as_deref(),
|
|
repo_lifecycle_guard,
|
|
metrics_clone.clone(),
|
|
)
|
|
.await;
|
|
if let Some(ref m) = metrics_clone {
|
|
if result.is_err() {
|
|
// On `Ok(Response)`, the streaming task owns the actual
|
|
// git-upload-pack outcome metric. The synchronous handler result
|
|
// only says whether stream setup succeeded.
|
|
m.record_git_operation("clone", "error");
|
|
}
|
|
}
|
|
result
|
|
}
|
|
|
|
// POST /git-receive-pack (push) - with GRASP authorization via database
|
|
(m, "git-receive-pack") if m == Method::POST => {
|
|
// Convert npub (bech32) to hex pubkey for authorization
|
|
let owner_pubkey_hex = match PublicKey::parse(&npub) {
|
|
Ok(pk) => pk.to_hex(),
|
|
Err(e) => {
|
|
tracing::warn!("Invalid npub in URL {}: {}", npub, e);
|
|
// Track failed push due to invalid npub
|
|
if let Some(ref m) = metrics_clone {
|
|
m.record_git_operation("push", "error");
|
|
}
|
|
return Ok(add_cors_headers(Response::builder())
|
|
.status(hyper::StatusCode::BAD_REQUEST)
|
|
.body(full_body(format!("Invalid npub: {}", e)))
|
|
.unwrap());
|
|
}
|
|
};
|
|
|
|
let repo_lifecycle_guard = Some(
|
|
lifecycle
|
|
.read_repository(&owner_pubkey_hex, &identifier)
|
|
.await,
|
|
);
|
|
|
|
let promotion_hooks = NostrPurgatoryPromotionHooks::git_push(
|
|
write_policy.clone(),
|
|
rejected_events_index.clone(),
|
|
Some(relay.clone()),
|
|
);
|
|
|
|
let result = git::handlers::handle_receive_pack(
|
|
repo_path,
|
|
body_bytes.clone(),
|
|
database.clone(),
|
|
relay.clone(),
|
|
&identifier,
|
|
&owner_pubkey_hex,
|
|
purgatory.clone(),
|
|
&git_data_path,
|
|
git_protocol.as_deref(),
|
|
repo_lifecycle_guard,
|
|
Some(Arc::new(promotion_hooks)),
|
|
metrics_clone.clone(),
|
|
)
|
|
.await;
|
|
|
|
if let Some(ref m) = metrics_clone {
|
|
if result.is_err() {
|
|
// On `Ok(Response)`, the operation may be a buffered protocol
|
|
// rejection already recorded by `handle_receive_pack`, or a stream
|
|
// whose final git-receive-pack outcome will be recorded by the
|
|
// streaming task. Do not treat this setup result as success.
|
|
m.record_git_operation("push", "error");
|
|
}
|
|
}
|
|
|
|
result
|
|
}
|
|
|
|
_ => Err(git::handlers::GitError::RepositoryNotFound),
|
|
};
|
|
|
|
match result {
|
|
Ok(response) => {
|
|
// Add CORS headers to successful Git responses
|
|
let (parts, body) = response.into_parts();
|
|
// RFC 9110 §9.3.2: HEAD response must have same headers as GET
|
|
// but no body.
|
|
let body = if method == Method::HEAD {
|
|
empty_body()
|
|
} else {
|
|
body
|
|
};
|
|
Ok(add_cors_headers(Response::builder().status(parts.status))
|
|
.header(
|
|
"content-type",
|
|
parts
|
|
.headers
|
|
.get("content-type")
|
|
.and_then(|v| v.to_str().ok())
|
|
.unwrap_or("application/octet-stream"),
|
|
)
|
|
.header(
|
|
"cache-control",
|
|
parts
|
|
.headers
|
|
.get("cache-control")
|
|
.and_then(|v| v.to_str().ok())
|
|
.unwrap_or("no-cache"),
|
|
)
|
|
.body(body)
|
|
.unwrap())
|
|
}
|
|
Err(e) => {
|
|
// Errors are already logged at their source with full context
|
|
let error_msg = format!("Git error: {}", e);
|
|
Ok(add_cors_headers(Response::builder())
|
|
.status(e.status_code())
|
|
.body(full_body(error_msg))
|
|
.unwrap())
|
|
}
|
|
}
|
|
});
|
|
}
|
|
|
|
// Check for NIP-11 relay information request (Accept: application/nostr+json)
|
|
if let Some(accept) = req.headers().get("accept") {
|
|
if accept
|
|
.to_str()
|
|
.map(|s| s.contains("application/nostr+json"))
|
|
.unwrap_or(false)
|
|
{
|
|
let doc = nip11::RelayInformationDocument::from_config(&self.config);
|
|
let json = doc.to_json().unwrap_or_else(|e| {
|
|
tracing::error!("Failed to serialize NIP-11 document: {}", e);
|
|
"{}".to_string()
|
|
});
|
|
|
|
tracing::debug!(
|
|
"Serving NIP-11 relay information document to {}",
|
|
self.remote
|
|
);
|
|
|
|
return Box::pin(async move {
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "application/nostr+json")
|
|
.body(full_body(json))
|
|
.unwrap(),
|
|
)
|
|
});
|
|
}
|
|
}
|
|
|
|
// Check for repository URL pattern (e.g., /npub/repo.git without subpath)
|
|
// GRASP-01: "SHOULD serve a webpage at the same endpoint linking to git nostr client(s)
|
|
// to browse the repository and a 404 page for repositories it doesn't host"
|
|
if let Some((npub, identifier)) = parse_repo_url(&path) {
|
|
let config = self.config.clone();
|
|
let repo_path = git::resolve_repo_path(&git_data_path, &npub, &identifier);
|
|
|
|
tracing::debug!(
|
|
"Repository URL request: {} (npub={}, id={}, path={:?})",
|
|
path,
|
|
npub,
|
|
identifier,
|
|
repo_path
|
|
);
|
|
|
|
return Box::pin(async move {
|
|
// Check if repository exists
|
|
if repo_path.exists() {
|
|
// Serve repository webpage
|
|
let html = landing::get_repo_html(&config, &npub, &identifier);
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "text/html; charset=utf-8")
|
|
.body(full_body(html))
|
|
.unwrap(),
|
|
)
|
|
} else {
|
|
// Serve 404 page for non-existent repository
|
|
let html = landing::get_404_html(&config, &npub, &identifier);
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(404)
|
|
.header("content-type", "text/html; charset=utf-8")
|
|
.body(full_body(html))
|
|
.unwrap(),
|
|
)
|
|
}
|
|
});
|
|
}
|
|
|
|
// Check if this is a WebSocket upgrade request
|
|
if let (Some(c), Some(w)) = (
|
|
req.headers().get("connection"),
|
|
req.headers().get("upgrade"),
|
|
) {
|
|
if c.to_str()
|
|
.map(|s| s.to_lowercase() == "upgrade")
|
|
.unwrap_or(false)
|
|
&& w.to_str()
|
|
.map(|s| s.to_lowercase() == "websocket")
|
|
.unwrap_or(false)
|
|
{
|
|
let key = req.headers().get("sec-websocket-key");
|
|
let derived = key.map(|k| derive_accept_key(k.as_bytes()));
|
|
|
|
let peer = self.remote;
|
|
let addr = client_ip::resolve_client_addr(
|
|
peer,
|
|
req.headers(),
|
|
&self.config.trusted_proxy_cidrs,
|
|
);
|
|
let relay = self.relay.clone();
|
|
let metrics_clone = self.metrics.clone();
|
|
|
|
tokio::spawn(async move {
|
|
match hyper::upgrade::on(req).await {
|
|
Ok(upgraded) => {
|
|
tracing::info!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
"WebSocket connection established"
|
|
);
|
|
// Track connection
|
|
let _connection_timer =
|
|
metrics_clone.as_ref().map(|m| m.start_connection_timer());
|
|
if let Some(ref m) = metrics_clone {
|
|
m.connection_tracker().on_connect(addr.ip());
|
|
m.record_websocket_connection();
|
|
}
|
|
if let Err(e) =
|
|
relay.take_connection(TokioIo::new(upgraded), addr).await
|
|
{
|
|
tracing::error!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
error = %e,
|
|
"Relay connection failed"
|
|
);
|
|
}
|
|
tracing::info!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
"WebSocket connection closed"
|
|
);
|
|
// Untrack connection
|
|
if let Some(ref m) = metrics_clone {
|
|
m.connection_tracker().on_disconnect(addr.ip());
|
|
}
|
|
}
|
|
Err(e) => tracing::error!("Upgrade error: {}", e),
|
|
}
|
|
});
|
|
|
|
return Box::pin(async move {
|
|
Ok(base
|
|
.status(101)
|
|
.header(CONNECTION, "upgrade")
|
|
.header(UPGRADE, "websocket")
|
|
.header(SEC_WEBSOCKET_ACCEPT, derived.unwrap())
|
|
.body(empty_body())
|
|
.unwrap())
|
|
});
|
|
}
|
|
}
|
|
|
|
// Serve Prometheus metrics if enabled
|
|
if path == "/metrics" {
|
|
if let Some(ref metrics) = self.metrics {
|
|
let metrics = metrics.clone();
|
|
return Box::pin(async move {
|
|
let output = metrics.render();
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "text/plain; version=0.0.4; charset=utf-8")
|
|
.body(full_body(output))
|
|
.unwrap(),
|
|
)
|
|
});
|
|
} else {
|
|
// Metrics disabled
|
|
return Box::pin(async move {
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(404)
|
|
.body(full_body("Metrics disabled"))
|
|
.unwrap(),
|
|
)
|
|
});
|
|
}
|
|
}
|
|
|
|
// Serve static icon at /icon.png
|
|
if path == "/icon.png" {
|
|
return Box::pin(async move {
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "image/png")
|
|
.header("cache-control", "public, max-age=86400")
|
|
.body(full_body(Bytes::from_static(ICON_PNG)))
|
|
.unwrap(),
|
|
)
|
|
});
|
|
}
|
|
|
|
// Only serve landing page for root path "/", 404 for everything else
|
|
let config = self.config.clone();
|
|
Box::pin(async move {
|
|
if path == "/" {
|
|
// Serve landing page for root
|
|
let html = landing::get_html(&config);
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "text/html; charset=utf-8")
|
|
.body(full_body(html))
|
|
.unwrap(),
|
|
)
|
|
} else {
|
|
// Serve generic 404 for unknown paths
|
|
let html = landing::get_generic_404_html(&config, &path);
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(404)
|
|
.header("content-type", "text/html; charset=utf-8")
|
|
.body(full_body(html))
|
|
.unwrap(),
|
|
)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
/// Derive the `Sec-WebSocket-Accept` response header from a `Sec-WebSocket-Key` request header
|
|
fn derive_accept_key(request_key: &[u8]) -> String {
|
|
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 = Sha1Hash::from_engine(engine);
|
|
base64::prelude::BASE64_STANDARD.encode(hash)
|
|
}
|
|
|
|
/// Start the HTTP server with integrated Nostr relay
|
|
///
|
|
/// 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,
|
|
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<()> {
|
|
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
|
|
}
|
|
|
|
/// 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);
|
|
if !config.trusted_proxy_cidrs.is_empty() {
|
|
tracing::info!(
|
|
trusted_proxy_cidrs = ?config.trusted_proxy_cidrs,
|
|
"Trusted proxy client IP resolution enabled"
|
|
);
|
|
}
|
|
|
|
loop {
|
|
let (socket, addr) = listener.accept().await?;
|
|
let io = TokioIo::new(socket);
|
|
let service = HttpService::new(
|
|
relay.clone(),
|
|
config.clone(),
|
|
addr,
|
|
database.clone(),
|
|
metrics.clone(),
|
|
purgatory.clone(),
|
|
write_policy.clone(),
|
|
lifecycle.clone(),
|
|
rejected_events_index.clone(),
|
|
repo_init_locks.clone(),
|
|
);
|
|
|
|
tokio::spawn(async move {
|
|
if let Err(e) = http1::Builder::new()
|
|
.serve_connection(io, service)
|
|
.with_upgrades()
|
|
.await
|
|
{
|
|
tracing::error!("Failed to handle request from {}: {}", addr, e);
|
|
}
|
|
});
|
|
}
|
|
}
|