mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-06 07:28:23 +00:00
Repository counting traverses the filesystem during each metrics scrape. Run rendering on a blocking worker so slow metadata I/O cannot occupy an async worker serving WebSocket and Git traffic. Acquire a shared permit before spawning and retain it inside the worker, including after HTTP cancellation, to prevent overlapping filesystem scans. The synchronous rendering API and metric contents remain unchanged. This does not add caching or alter monitoring configuration. Validation: workspace library suite passes all 907 tests, including a current-thread runtime check that queued scrapes yield and release capacity and that asynchronous rendering preserves repository counts.
1173 lines
50 KiB
Rust
1173 lines
50 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 nip05;
|
|
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::error::ErrorKind as RelayErrorKind;
|
|
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::private::{nip98, PrivateAccess};
|
|
use crate::purgatory::promotion_hooks::NostrPurgatoryPromotionHooks;
|
|
use crate::purgatory::Purgatory;
|
|
use crate::sync::rejected_index::RejectedEventsIndex;
|
|
|
|
type HttpBody = GitResponseBody;
|
|
|
|
fn is_peer_scoped_relay_error(kind: RelayErrorKind) -> bool {
|
|
matches!(
|
|
kind,
|
|
RelayErrorKind::Protocol
|
|
| RelayErrorKind::IO
|
|
| RelayErrorKind::Transport
|
|
| RelayErrorKind::Timeout
|
|
| RelayErrorKind::Rejected
|
|
| RelayErrorKind::LimitExceeded
|
|
| RelayErrorKind::Other
|
|
)
|
|
}
|
|
|
|
fn is_peer_scoped_http_error(error: &hyper::Error) -> bool {
|
|
error.is_parse()
|
|
|| error.is_canceled()
|
|
|| error.is_closed()
|
|
|| error.is_incomplete_message()
|
|
|| error.is_body_write_aborted()
|
|
}
|
|
|
|
/// 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, Authorization, Git-Protocol";
|
|
const CORS_EXPOSE_HEADERS: &str = "WWW-Authenticate";
|
|
|
|
/// 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;
|
|
}
|
|
git::parse_repository_root_url(path)
|
|
}
|
|
|
|
/// 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)
|
|
.header("Access-Control-Expose-Headers", CORS_EXPOSE_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,
|
|
/// GRASP-08 access list. Present only when private mode is enabled.
|
|
private_access: Option<PrivateAccess>,
|
|
}
|
|
|
|
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,
|
|
private_access: Option<PrivateAccess>,
|
|
) -> Self {
|
|
Self {
|
|
relay,
|
|
config,
|
|
remote,
|
|
database,
|
|
metrics,
|
|
purgatory,
|
|
write_policy,
|
|
lifecycle,
|
|
rejected_events_index,
|
|
repo_init_locks,
|
|
private_access,
|
|
}
|
|
}
|
|
}
|
|
|
|
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 request_path = req.uri().path().to_string();
|
|
let Some(path) = self
|
|
.config
|
|
.strip_base_path(&request_path)
|
|
.map(str::to_string)
|
|
else {
|
|
let config = self.config.clone();
|
|
return Box::pin(async move {
|
|
let html = landing::get_generic_404_html(&config, &request_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(),
|
|
)
|
|
});
|
|
};
|
|
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(),
|
|
)
|
|
});
|
|
}
|
|
|
|
// NIP-05 root identity: `_@domain` resolves to the relay operator's
|
|
// public key. This exact-path route must remain ahead of generic NIP-11
|
|
// content negotiation so clients receive the well-known document even
|
|
// if they send a broad or unusual Accept header.
|
|
if self.config.is_domain_root()
|
|
&& path == "/.well-known/nostr.json"
|
|
&& (method == Method::GET || method == Method::HEAD)
|
|
{
|
|
let method = method.clone();
|
|
let document = nip05::Nip05Document::from_config(&self.config);
|
|
|
|
return Box::pin(async move {
|
|
match document.and_then(|document| document.to_json()) {
|
|
Ok(json) => {
|
|
let body = if method == Method::HEAD {
|
|
empty_body()
|
|
} else {
|
|
full_body(json)
|
|
};
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(200)
|
|
.header("content-type", "application/json; charset=utf-8")
|
|
.body(body)
|
|
.unwrap(),
|
|
)
|
|
}
|
|
Err(error) => {
|
|
tracing::error!(%error, "Failed to build NIP-05 root identity document");
|
|
Ok(
|
|
add_cors_headers(Response::builder().header("server", "ngit-grasp"))
|
|
.status(500)
|
|
.body(full_body("Failed to build NIP-05 document"))
|
|
.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.service_address(),
|
|
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) {
|
|
if let Some(access) = &self.private_access {
|
|
let authorized = nip98::canonical_repository_url(&self.config, &request_path)
|
|
.is_some_and(|url| nip98::validate_request(&req, &url, access).is_ok());
|
|
if !authorized {
|
|
let response = nip98::unauthorized_response(&self.config);
|
|
return Box::pin(async move {
|
|
let (parts, body) = response.into_parts();
|
|
Ok(add_cors_headers(Response::builder().status(parts.status))
|
|
.header(
|
|
"www-authenticate",
|
|
parts.headers["www-authenticate"].clone(),
|
|
)
|
|
.body(body)
|
|
.unwrap())
|
|
});
|
|
}
|
|
}
|
|
|
|
// 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)
|
|
.expect("parse_git_url validated repository coordinate");
|
|
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)) = (path.as_str(), 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) {
|
|
if let Some(access) = &self.private_access {
|
|
let authorized = nip98::canonical_repository_url(&self.config, &request_path)
|
|
.is_some_and(|url| nip98::validate_request(&req, &url, access).is_ok());
|
|
if !authorized {
|
|
let response = nip98::unauthorized_response(&self.config);
|
|
return Box::pin(async move {
|
|
let (parts, body) = response.into_parts();
|
|
Ok(add_cors_headers(Response::builder().status(parts.status))
|
|
.header(
|
|
"www-authenticate",
|
|
parts.headers["www-authenticate"].clone(),
|
|
)
|
|
.body(body)
|
|
.unwrap())
|
|
});
|
|
}
|
|
}
|
|
|
|
let config = self.config.clone();
|
|
let repo_path = git::resolve_repo_path(&git_data_path, &npub, &identifier)
|
|
.expect("parse_repo_url validated repository coordinate");
|
|
|
|
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)) = (
|
|
path.as_str(),
|
|
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();
|
|
let private_access = self.private_access.clone();
|
|
let relay_domain = self.config.service_address();
|
|
|
|
tokio::spawn(async move {
|
|
match hyper::upgrade::on(req).await {
|
|
Ok(upgraded) => {
|
|
tracing::debug!(
|
|
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 Some(access) = private_access {
|
|
if let Err(e) = crate::private::ws_auth::authenticate_and_bridge(
|
|
TokioIo::new(upgraded),
|
|
addr,
|
|
relay,
|
|
&relay_domain,
|
|
access,
|
|
)
|
|
.await
|
|
{
|
|
tracing::debug!(
|
|
"Private relay connection ended for {}: {}",
|
|
addr,
|
|
e
|
|
);
|
|
}
|
|
} else if let Err(e) =
|
|
relay.take_connection(TokioIo::new(upgraded), addr).await
|
|
{
|
|
if is_peer_scoped_relay_error(e.kind()) {
|
|
tracing::debug!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
error_kind = ?e.kind(),
|
|
error = %e,
|
|
"Relay client session ended with an error"
|
|
);
|
|
} else {
|
|
tracing::error!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
error_kind = ?e.kind(),
|
|
error = %e,
|
|
"Relay connection failed"
|
|
);
|
|
}
|
|
}
|
|
tracing::debug!(
|
|
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) => {
|
|
if is_peer_scoped_http_error(&e) {
|
|
tracing::debug!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
error = %e,
|
|
"WebSocket upgrade ended before completion"
|
|
);
|
|
} else {
|
|
tracing::error!(
|
|
client_ip = %addr.ip(),
|
|
peer = %peer,
|
|
error = %e,
|
|
"WebSocket upgrade failed"
|
|
);
|
|
}
|
|
}
|
|
}
|
|
});
|
|
|
|
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 = match metrics.render_async().await {
|
|
Ok(output) => output,
|
|
Err(error) => {
|
|
tracing::error!(%error, "Metrics render worker failed");
|
|
return Ok(add_cors_headers(
|
|
Response::builder().header("server", "ngit-grasp"),
|
|
)
|
|
.status(500)
|
|
.body(full_body("Metrics temporarily unavailable"))
|
|
.unwrap());
|
|
}
|
|
};
|
|
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, &request_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,
|
|
private_access: Option<PrivateAccess>,
|
|
) -> 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,
|
|
private_access,
|
|
)
|
|
.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,
|
|
private_access: Option<PrivateAccess>,
|
|
) -> anyhow::Result<()> {
|
|
tracing::info!("Starting HTTP server on {}", listener.local_addr()?);
|
|
tracing::info!("Relay name: {}", config.relay_name());
|
|
tracing::info!("Public service: {}", config.service_address());
|
|
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?;
|
|
// WebSocket EVENT/EOSE and streaming Git responses flush small writes.
|
|
// Nagle can hold the tail until the peer's delayed ACK (about 40 ms
|
|
// even through a loopback proxy), multiplying multi-query latency.
|
|
if let Err(error) = socket.set_nodelay(true) {
|
|
tracing::warn!(peer = %addr, %error, "Could not disable Nagle for client socket");
|
|
continue;
|
|
}
|
|
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(),
|
|
private_access.clone(),
|
|
);
|
|
|
|
tokio::spawn(async move {
|
|
if let Err(e) = http1::Builder::new()
|
|
.serve_connection(io, service)
|
|
.with_upgrades()
|
|
.await
|
|
{
|
|
if is_peer_scoped_http_error(&e) {
|
|
tracing::debug!(
|
|
peer = %addr,
|
|
error = %e,
|
|
"HTTP client connection ended before completion"
|
|
);
|
|
} else {
|
|
tracing::error!(
|
|
peer = %addr,
|
|
error = %e,
|
|
"Failed to handle HTTP request"
|
|
);
|
|
}
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn relay_client_failures_are_peer_scoped() {
|
|
for kind in [
|
|
RelayErrorKind::Protocol,
|
|
RelayErrorKind::IO,
|
|
RelayErrorKind::Transport,
|
|
RelayErrorKind::Timeout,
|
|
RelayErrorKind::Rejected,
|
|
RelayErrorKind::LimitExceeded,
|
|
RelayErrorKind::Other,
|
|
] {
|
|
assert!(
|
|
is_peer_scoped_relay_error(kind),
|
|
"unexpected kind: {kind:?}"
|
|
);
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn relay_internal_failures_remain_operator_actionable() {
|
|
for kind in [
|
|
RelayErrorKind::Database,
|
|
RelayErrorKind::Gossip,
|
|
RelayErrorKind::Policy,
|
|
RelayErrorKind::NotFound,
|
|
RelayErrorKind::Invalid,
|
|
RelayErrorKind::State,
|
|
RelayErrorKind::Unsupported,
|
|
] {
|
|
assert!(
|
|
!is_peer_scoped_relay_error(kind),
|
|
"unexpected kind: {kind:?}"
|
|
);
|
|
}
|
|
}
|
|
}
|