Files
ngit-grasp/src/git/handlers.rs
T
DanConwayDev c8fee2ce6a fix(git): recover unsigned upload expiry after crashes
A crash between receive-pack and the purgatory checkpoint left a live staged ref with no expiry owner. Maintenance preserved that ref indefinitely, defeating reclamation of abandoned uploads.

Persist unsigned ref intentions and absolute deadlines in the staging registry before Git runs. Startup reconstructs placeholders from matching live refs before cleanup and request handling, preserving signed events and newer checkpoint deadlines. Legacy records receive one persisted grace period. Compaction drops stale journal entries while retaining owed history.

Recovery assumes exclusive startup access and an available event database; unreadable recovery metadata fails startup rather than authorizing deletion. This change does not add quotas or alter normal placeholder expiry policy.

Validation: 1,179 tests passed across the library, pending_upload_staging, purgatory and grasp06_pr_hosting suites. Regressions cover missing checkpoints, repeated restart deadlines, both endpoint scopes, signed-event preservation, failed pushes and legacy recovery. Workspace all-target clippy with warnings denied, cargo fmt --check and git diff --check passed.

Assisted-by: GPT-6
2026-09-29 11:25:37 +00:00

1481 lines
54 KiB
Rust

//! Git HTTP Protocol Handlers
//!
//! This module implements the HTTP handlers for Git Smart HTTP protocol.
use futures_util::stream;
use http_body_util::{BodyExt, StreamBody};
use hyper::{body::Bytes, body::Frame, Response, StatusCode};
use nostr_sdk::local_relay::LocalRelay;
use std::collections::HashSet;
use std::io;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::mpsc;
use tokio::time::MissedTickBehavior;
use tracing::{debug, error, info, warn, Instrument};
use super::protocol::{GitService, PktLine};
use super::receive_pack_plan::ReceivePackPlan;
use super::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
use super::subprocess::GitSubprocess;
use super::{full_body, GitResponseBody};
use crate::git::authorization::{authorize_push, normalize_applied_ref_updates, parse_pushed_refs};
use crate::git::sync::{process_newly_available_git_data, PurgatoryPromotionHooks};
use crate::metrics::Metrics;
use crate::nostr::lifecycle::LifecycleReadGuard;
use crate::nostr::SharedDatabase;
use crate::purgatory::Purgatory;
pub(crate) const STREAM_CHANNEL_DEPTH: usize = 8;
const STREAM_CHUNK_SIZE: usize = 8 * 1024;
const POST_PUSH_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(5);
const POST_PUSH_KEEPALIVE_MESSAGE: &[u8] = b"GRASP is finalizing the push\n";
/// Handle GET /info/refs?service=git-{upload,receive}-pack
///
/// This advertises the repository's refs to the client.
pub async fn handle_info_refs(
repo_path: PathBuf,
service: GitService,
git_protocol: Option<&str>,
) -> Result<Response<GitResponseBody>, GitError> {
debug!(
"Handling info/refs for {:?} with service {:?}",
repo_path, service
);
// Check if repository exists
if !repo_path.exists() {
debug!("Repository not found: {:?}", repo_path);
return Err(GitError::RepositoryNotFound);
}
// Spawn git with --advertise-refs
let mut git = GitSubprocess::spawn(service, &repo_path, true, git_protocol).map_err(|e| {
error!("Failed to spawn git process: {}", e);
GitError::ProcessSpawnFailed(e)
})?;
// Read the output from git
let mut output = Vec::new();
let mut stderr_output = Vec::new();
if let Some(stdout) = git.take_stdout() {
let mut stdout = stdout;
stdout.read_to_end(&mut output).await.map_err(|e| {
error!("Failed to read git output: {}", e);
GitError::IoError(e)
})?;
}
if let Some(stderr) = git.take_stderr() {
let mut stderr = stderr;
stderr.read_to_end(&mut stderr_output).await.map_err(|e| {
error!("Failed to read git stderr: {}", e);
GitError::IoError(e)
})?;
}
// Wait for process to complete
let status = git.wait().await.map_err(|e| {
error!("Failed to wait for git process: {}", e);
GitError::IoError(e)
})?;
if !status.success() {
let stderr_str = String::from_utf8_lossy(&stderr_output);
error!(
"Git process failed with status: {:?}, stderr: {}",
status, stderr_str
);
return Err(GitError::GitFailed(status.code()));
}
// Build response with pkt-line header
let mut response_body = Vec::new();
// First line: service advertisement
let service_line = format!("# service={}\n", service.as_str());
response_body.extend_from_slice(&PktLine::data(service_line.as_bytes()).encode());
response_body.extend_from_slice(&PktLine::flush().encode());
// Then the git output
response_body.extend_from_slice(&output);
Ok(Response::builder()
.status(StatusCode::OK)
.header("content-type", service.advertisement_content_type())
.header("cache-control", "no-cache")
.body(full_body(response_body))
.unwrap())
}
/// Detect whether a git-receive-pack client negotiated `side-band-64k`.
///
/// On a receive-pack POST, the client advertises its capabilities in the
/// first command pkt-line of the request body. Per the smart-HTTP protocol,
/// the format of that pkt-line is:
///
/// <old-oid> SP <new-oid> SP <ref-name> NUL <capabilities>
///
/// Where `<capabilities>` is a space-separated list that may include
/// `side-band-64k` (and/or the legacy `side-band`). We look at the bytes
/// after the first NUL inside the first non-flush pkt-line and check for
/// the capability token.
///
/// Returns `false` for any malformed or truncated body, which is the
/// conservative choice because a `false` answer means we do NOT band-wrap
/// the ERR pkt-line and clients without sideband-64k will still see it.
///
/// `pub(crate)` for testability.
pub(crate) fn client_negotiated_sideband_64k(request_body: &[u8]) -> bool {
// First 4 bytes are the hex length of the first pkt-line. "0000" is a
// flush packet (shouldn't appear here) — bail.
if request_body.len() < 4 {
return false;
}
let len_str = match std::str::from_utf8(&request_body[0..4]) {
Ok(s) => s,
Err(_) => return false,
};
let len = match u16::from_str_radix(len_str, 16) {
Ok(n) => n as usize,
Err(_) => return false,
};
if len < 4 || request_body.len() < len {
return false;
}
let payload = &request_body[4..len];
// Look for capability list after first NUL byte
let caps = match payload.iter().position(|&b| b == 0) {
Some(pos) => &payload[pos + 1..],
None => return false,
};
// Capabilities are space-separated tokens, possibly terminated by LF.
for tok in caps.split(|&b| b == b' ' || b == b'\n' || b == b'\0') {
if tok == b"side-band-64k" {
return true;
}
}
false
}
/// Build an HTTP 200 OK response with an ERR pkt-line for git protocol errors.
///
/// Per the git smart HTTP protocol spec, protocol-level errors (like "not our ref")
/// should be returned as HTTP 200 OK with the error message in pkt-line format:
/// `PKT-LINE("ERR" SP explanation-text)`
///
/// This allows git clients to properly parse and display the error message.
///
/// **Sideband wrapping (git-receive-pack):** modern `git push` clients
/// always advertise `side-band-64k` when posting to `git-receive-pack`.
/// When sideband is negotiated, the client demultiplexes the response via
/// `recv_sideband`, which reads the first byte of each pkt-line payload as
/// the band id (1=pack data, 2=progress, 3=error). A bare `ERR ...`
/// payload is therefore interpreted as band id `0x45` (`'E'`, decimal 69)
/// and the client aborts with `send-pack: protocol error: bad band #69`.
///
/// To stay compatible with sideband-aware clients, when this is a
/// receive-pack response and `request_body` shows the client advertised
/// `side-band-64k`, the ERR pkt-line payload is prefixed with band id 3
/// ("error") inside the pkt-line frame.
///
/// `pub(crate)` so the `/prs/` receive-pack handler in `crate::grasp06::receive`
/// can return identically-shaped rejections without duplicating the pkt-line
/// framing.
pub(crate) fn build_git_protocol_error_response(
service: GitService,
error_message: &str,
request_body: Option<&[u8]>,
) -> Response<GitResponseBody> {
let err_pktline = encode_err_pktline(service, error_message, request_body);
Response::builder()
.status(StatusCode::OK)
.header("content-type", service.result_content_type())
.header("cache-control", "no-cache")
.body(full_body(err_pktline))
.unwrap()
}
fn encode_err_pktline(
service: GitService,
error_message: &str,
request_body: Option<&[u8]>,
) -> Vec<u8> {
// Format: "ERR <message>\n"
let err_content = format!("ERR {}\n", error_message.trim());
// Wrap in sideband band 3 when the receive-pack client negotiated
// side-band-64k. See doc comment above for the protocol details.
let use_sideband = matches!(service, GitService::ReceivePack)
&& request_body
.map(client_negotiated_sideband_64k)
.unwrap_or(false);
if use_sideband {
let mut framed = Vec::with_capacity(1 + err_content.len());
framed.push(0x03); // band 3 = error
framed.extend_from_slice(err_content.as_bytes());
PktLine::data(framed).encode()
} else {
PktLine::data(err_content.as_bytes()).encode()
}
}
pub(crate) fn err_pktline_frame(
service: GitService,
error_message: &str,
request_body: Option<&[u8]>,
) -> Frame<Bytes> {
Frame::data(Bytes::from(encode_err_pktline(
service,
error_message,
request_body,
)))
}
/// Check if a git process failure is a protocol error (vs transport error).
///
/// Protocol errors are communicated via stderr when git exits with code 128.
/// These should be returned to the client as HTTP 200 with ERR pkt-line.
///
/// Transport errors (process spawn failures, I/O errors, signals) should
/// remain as HTTP 500 errors.
///
/// `pub(crate)` so the `/prs/` receive-pack handler in `crate::grasp06::receive`
/// can classify subprocess failures the same way without duplicating the rule.
pub(crate) fn is_git_protocol_error(exit_code: Option<i32>, stderr: &[u8]) -> bool {
// Git uses exit code 128 for protocol/usage errors
// If there's stderr content, it's a protocol error message
exit_code == Some(128) && !stderr.is_empty()
}
fn is_routine_missing_object_error(stderr: &[u8]) -> bool {
let stderr = String::from_utf8_lossy(stderr);
stderr.contains("not our ref")
|| stderr.contains("Server does not allow request for unadvertised object")
}
/// Build the common channel-backed Git response used by streaming handlers.
///
/// The handler returns this response before the Git child has necessarily
/// exited. A background task sends `Frame<Bytes>` values into `rx`; dropping the
/// sender closes the HTTP body. Kept `pub(crate)` so `/prs/` receive-pack can
/// use the same wire shape as the standard endpoints.
pub(crate) fn streaming_response(
service: GitService,
rx: mpsc::Receiver<Result<Frame<Bytes>, io::Error>>,
) -> Response<GitResponseBody> {
let body_stream = stream::unfold(rx, |mut rx| async {
rx.recv().await.map(|item| (item, rx))
});
let body = BodyExt::boxed(StreamBody::new(body_stream));
Response::builder()
.status(StatusCode::OK)
.header("content-type", service.result_content_type())
.header("cache-control", "no-cache")
.body(body)
.unwrap()
}
/// Record Git metrics from either the request handler or the detached streaming
/// task without repeating the `Option<Metrics>` plumbing at every callsite.
pub(crate) fn record_git_operation(metrics: &Option<Arc<Metrics>>, operation: &str, status: &str) {
if let Some(metrics) = metrics {
metrics.record_git_operation(operation, status);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PumpResult {
/// Stdout reached EOF. `sent_stdout` tells callers whether it is still safe
/// to synthesize a Git `ERR` pkt-line for a late process failure.
Eof { sent_stdout: bool },
/// The response receiver was dropped, usually because the HTTP client
/// disconnected. Callers should stop the child process and clean up state.
ClientDisconnected,
/// Reading stdout failed. The read error has already been sent through the
/// body channel so Hyper can terminate the response stream.
ReadError,
}
/// Forward Git stdout into the streaming response channel.
///
/// This is intentionally only the stdout pump. The caller still owns process
/// termination, stderr collection, protocol-error classification, and any
/// endpoint-specific cleanup after the stream ends.
pub(crate) async fn pump_stdout_to_channel<R>(
mut stdout: R,
tx: &mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
) -> PumpResult
where
R: tokio::io::AsyncRead + Unpin,
{
let mut read_buf = [0_u8; STREAM_CHUNK_SIZE];
let mut sent_stdout = false;
loop {
match stdout.read(&mut read_buf).await {
Ok(0) => return PumpResult::Eof { sent_stdout },
Ok(n) => {
sent_stdout = true;
if send_body_bytes(tx, read_buf[..n].to_vec()).await.is_err() {
return PumpResult::ClientDisconnected;
}
}
Err(e) => {
let _ = tx.send(Err(e)).await;
return PumpResult::ReadError;
}
}
}
}
/// Forward receive-pack progress while retaining its protocol terminator.
///
/// A successful receive-pack response ends with a `0000` flush pkt-line.
/// Git clients use that flush as the semantic end of the push and need not wait
/// for the HTTP body to reach EOF. Retaining the final four bytes lets the
/// caller finish GRASP post-push processing before making success visible,
/// without buffering the progress stream that keeps clients alive during
/// expensive pack processing.
pub(crate) async fn pump_receive_pack_stdout_to_channel<R>(
mut stdout: R,
tx: &mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
) -> (PumpResult, Option<Vec<u8>>)
where
R: tokio::io::AsyncRead + Unpin,
{
const FLUSH_PKT: &[u8; 4] = b"0000";
let mut read_buf = [0_u8; STREAM_CHUNK_SIZE];
let mut pending = Vec::with_capacity(4);
let mut sent_stdout = false;
loop {
match stdout.read(&mut read_buf).await {
Ok(0) => {
if pending.as_slice() == FLUSH_PKT {
return (
PumpResult::Eof { sent_stdout },
Some(std::mem::take(&mut pending)),
);
}
if !pending.is_empty() {
sent_stdout = true;
if send_body_bytes(tx, std::mem::take(&mut pending))
.await
.is_err()
{
return (PumpResult::ClientDisconnected, None);
}
}
return (PumpResult::Eof { sent_stdout }, None);
}
Ok(n) => {
pending.extend_from_slice(&read_buf[..n]);
if pending.len() > FLUSH_PKT.len() {
let retained = pending.split_off(pending.len() - FLUSH_PKT.len());
sent_stdout = true;
if send_body_bytes(tx, std::mem::replace(&mut pending, retained))
.await
.is_err()
{
return (PumpResult::ClientDisconnected, None);
}
}
}
Err(e) => {
let _ = tx.send(Err(e)).await;
return (PumpResult::ReadError, None);
}
}
}
}
/// Encode progress that keeps a sideband-aware receive-pack client alive.
///
/// Post-push purgatory promotion can copy objects and align several owner
/// repositories. The terminal flush remains withheld until that work finishes,
/// so send a valid band-2 pkt-line during a long finalization window instead of
/// leaving libgit2 with no response bytes until its per-recv timeout expires.
fn receive_pack_keepalive_pktline() -> Vec<u8> {
let mut payload = Vec::with_capacity(1 + POST_PUSH_KEEPALIVE_MESSAGE.len());
payload.push(0x02); // band 2 = progress
payload.extend_from_slice(POST_PUSH_KEEPALIVE_MESSAGE);
PktLine::data(payload).encode()
}
struct ReceivePackKeepalive {
task: Option<tokio::task::JoinHandle<()>>,
}
impl ReceivePackKeepalive {
fn spawn(tx: mpsc::Sender<Result<Frame<Bytes>, io::Error>>, period: Duration) -> Self {
let task = tokio::spawn(async move {
let mut ticker = tokio::time::interval(period);
ticker.set_missed_tick_behavior(MissedTickBehavior::Skip);
ticker.tick().await;
loop {
ticker.tick().await;
if send_body_bytes(&tx, receive_pack_keepalive_pktline())
.await
.is_err()
{
return;
}
}
});
Self { task: Some(task) }
}
async fn stop(mut self) {
if let Some(task) = self.task.take() {
task.abort();
let _ = task.await;
}
}
}
impl Drop for ReceivePackKeepalive {
fn drop(&mut self) {
if let Some(task) = self.task.take() {
task.abort();
}
}
}
/// Drain Git stderr for later logging or Git protocol error synthesis.
///
/// Streaming handlers run this concurrently with stdout pumping so the child
/// cannot block on a full stderr pipe while the HTTP body is still being read.
pub(crate) async fn read_stderr_to_end<R>(mut stderr: R) -> Vec<u8>
where
R: tokio::io::AsyncRead + Unpin,
{
let mut stderr_output = Vec::new();
if let Err(e) = stderr.read_to_end(&mut stderr_output).await {
warn!("Failed to read git subprocess stderr: {}", e);
}
stderr_output
}
/// Handle POST /git-upload-pack (clone/fetch)
pub async fn handle_upload_pack(
repo_path: PathBuf,
request_body: Bytes,
git_protocol: Option<&str>,
repo_lifecycle_guard: Option<LifecycleReadGuard>,
metrics: Option<Arc<Metrics>>,
) -> Result<Response<GitResponseBody>, GitError> {
debug!("Handling upload-pack for {:?}", repo_path);
if !repo_path.exists() {
return Err(GitError::RepositoryNotFound);
}
// Spawn git upload-pack
let mut git = GitSubprocess::spawn(GitService::UploadPack, &repo_path, false, git_protocol)
.map_err(GitError::ProcessSpawnFailed)?;
// Write request to git's stdin
if let Some(mut stdin) = git.take_stdin() {
stdin.write_all(&request_body).await.map_err(|e| {
error!("Failed to write to git upload-pack stdin: {}", e);
GitError::IoError(e)
})?;
// Close stdin to signal end of input
drop(stdin);
}
let stdout = git.take_stdout().ok_or_else(|| {
GitError::IoError(io::Error::new(
io::ErrorKind::BrokenPipe,
"git upload-pack stdout unavailable",
))
})?;
let stderr = git.take_stderr();
// Stream upload-pack stdout as Git produces it instead of buffering the
// whole fetch response in memory. The detached task below owns the child
// until EOF so the request future can return the HTTP body immediately.
let (tx, rx) = mpsc::channel::<Result<Frame<Bytes>, io::Error>>(STREAM_CHANNEL_DEPTH);
let span = tracing::info_span!("git_upload_pack", repository = %repo_path.display());
tokio::spawn(
async move {
stream_upload_pack_output(git, stdout, stderr, tx, repo_lifecycle_guard, metrics).await;
}
.instrument(span),
);
Ok(streaming_response(GitService::UploadPack, rx))
}
async fn stream_upload_pack_output<S, E>(
mut git: GitSubprocess,
stdout: S,
stderr: Option<E>,
tx: mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
repo_lifecycle_guard: Option<LifecycleReadGuard>,
metrics: Option<Arc<Metrics>>,
) where
S: tokio::io::AsyncRead + Unpin + Send + 'static,
E: tokio::io::AsyncRead + Unpin + Send + 'static,
{
let stderr_task = stderr.map(|stderr| tokio::spawn(read_stderr_to_end(stderr)));
let pump_result = pump_stdout_to_channel(stdout, &tx).await;
if !matches!(pump_result, PumpResult::Eof { .. }) {
let _ = git.kill().await;
}
let status = match git.wait().await {
Ok(status) => status,
Err(e) => {
let _ = tx.send(Err(e)).await;
drop(repo_lifecycle_guard);
record_git_operation(&metrics, "clone", "error");
return;
}
};
// Keep the repository lifecycle read lock until git-upload-pack exits so
// deletion/archive/restore cannot remove or replace the bare repository
// while Git is still reading it in the detached streaming task.
drop(repo_lifecycle_guard);
let stderr_output = match stderr_task {
Some(task) => task.await.unwrap_or_default(),
None => Vec::new(),
};
if matches!(pump_result, PumpResult::ClientDisconnected) {
record_git_operation(&metrics, "clone", "error");
debug!(exit_status = %status, "Git upload-pack cancelled after client disconnected");
return;
}
if !status.success() && matches!(pump_result, PumpResult::Eof { sent_stdout: false }) {
record_git_operation(&metrics, "clone", "error");
let stderr_str = String::from_utf8_lossy(&stderr_output);
let msg = if stderr_str.trim().is_empty() {
format!("git upload-pack failed with code {:?}", status.code())
} else {
stderr_str.to_string()
};
if is_routine_missing_object_error(&stderr_output) {
debug!(
exit_code = ?status.code(),
stderr = %stderr_str.trim(),
"Git upload-pack request referenced an unavailable object"
);
} else if is_git_protocol_error(status.code(), &stderr_output) {
warn!(
"Git upload-pack protocol error (returning ERR pkt-line): {}",
stderr_str
);
} else {
error!(exit_status = %status, stream_outcome = ?pump_result,
stderr = %stderr_str, "Git upload-pack failed");
}
// The streaming response headers have already been sent. If Git failed
// before producing any stdout, surface a protocol-visible ERR pkt-line
// instead of silently completing an empty HTTP 200 body.
let _ = tx
.send(Ok(err_pktline_frame(GitService::UploadPack, &msg, None)))
.await;
} else if !status.success() {
record_git_operation(&metrics, "clone", "error");
let stderr_str = String::from_utf8_lossy(&stderr_output);
if is_routine_missing_object_error(&stderr_output) {
debug!(
exit_code = ?status.code(),
stderr = %stderr_str.trim(),
"Git upload-pack stream referenced an unavailable object"
);
} else if is_git_protocol_error(status.code(), &stderr_output) {
warn!(
"Git upload-pack protocol error after streaming stdout: {}",
stderr_str
);
} else {
error!(
exit_status = %status,
stream_outcome = ?pump_result,
stderr = %stderr_str,
"Git upload-pack failed after streaming stdout"
);
}
} else if matches!(pump_result, PumpResult::Eof { .. }) {
debug!("Git upload-pack stream completed successfully");
record_git_operation(&metrics, "clone", "success");
} else {
record_git_operation(&metrics, "clone", "error");
}
}
/// Handle POST /git-receive-pack (push)
///
/// This includes GRASP authorization validation according to GRASP-01:
/// "MUST accept pushes via this service that match the latest repo state announcement
/// on the relay, respecting the recursive maintainer set."
///
/// Also per GRASP-01: "MUST set repository HEAD per repository state announcement
/// as soon as the git data related to that branch has been received."
///
/// Also purgatory GRASP-01: "Accepted repo state announcements, PRs and PR Updates
/// SHOULD be accepted with message "purgatory: won't be served until git data arrives"
/// and kepted in purgatory (not served) until the related git data arrives and
/// otherwise discarded after 30 minutes."
///
/// # Arguments
/// * `repo_path` - Path to the bare git repository
/// * `request_body` - The git pack data from the client
/// * `database` - Database reference for authorization queries
/// * `identifier` - The repository identifier (d tag) for authorization lookup
/// * `owner_pubkey` - Announcement pubkey from the URL path; selects the repository coordinate whose authority is resolved
/// * `git_data_path` - Base path for git repositories (for syncing to other owner repos)
/// * `git_protocol` - Optional Git protocol version (e.g., "version=2")
#[allow(clippy::too_many_arguments)]
pub async fn handle_receive_pack(
repo_path: PathBuf,
request_body: Bytes,
database: SharedDatabase,
relay: LocalRelay,
identifier: &str,
owner_pubkey: &str,
purgatory: Arc<Purgatory>,
git_data_path: &str,
git_protocol: Option<&str>,
repo_lifecycle_guard: Option<LifecycleReadGuard>,
promotion_hooks: Option<Arc<dyn PurgatoryPromotionHooks>>,
metrics: Option<Arc<Metrics>>,
) -> Result<Response<GitResponseBody>, GitError> {
debug!("Handling receive-pack for {:?}", repo_path);
if !repo_path.exists() {
return Err(GitError::RepositoryNotFound);
}
let storage = LocalGitStorage::new(git_data_path);
let family_key = FamilyKey::sha1(identifier).map_err(|e| GitError::Storage(e.to_string()))?;
let family_lease = if storage.is_thin_view(&family_key, &repo_path) {
Some(
storage
.write_lease(&family_key)
.await
.map_err(|e| GitError::Storage(e.to_string()))?,
)
} else {
// Legacy repositories remain self-contained until startup migration.
// This fallback also keeps direct library callers safe: never point a
// ref at a family object database the view cannot read.
None
};
// A background fetch may have applied the client's target since refs were
// advertised and removed its state from purgatory. Normalize only commands
// already satisfied by this view; keep ordinary authorization for changes
// and Git's compare-and-swap protection against a later writer.
let (request_body, push_plan) = match crate::git::list_refs(&repo_path) {
Ok(refs) => {
let refs = refs.into_iter().collect();
let normalized = normalize_applied_ref_updates(&request_body, &refs);
let push_plan = ReceivePackPlan::prepare(&normalized, &refs);
(normalized, push_plan)
}
Err(_) => (request_body, None),
};
let authorization_body = push_plan
.as_ref()
.map(|push| &push.authorization_body)
.unwrap_or(&request_body);
// GRASP Authorization Check
debug!(
"Authorizing push for {} owned by {} via database query",
identifier, owner_pubkey
);
// check push is authorised
let auth_result = match authorize_push(
&database,
identifier,
owner_pubkey,
authorization_body,
&purgatory,
&repo_path,
)
.await
{
Ok(auth_result) => {
if !auth_result.authorized {
warn!("Push rejected for {}: {}", identifier, auth_result.reason);
record_git_operation(&metrics, "push", "error");
return Ok(build_git_protocol_error_response(
GitService::ReceivePack,
&format!("authorisation failed: {}", auth_result.reason),
Some(&request_body),
));
}
info!(
"Push authorized for {} - {} maintainers, {} purgatory events: {}",
identifier,
auth_result.maintainers.len(),
auth_result.purgatory_events.len(),
auth_result.reason
);
auth_result
}
Err(e) => {
warn!("Authorization check failed for {}: {}", identifier, e);
record_git_operation(&metrics, "push", "error");
return Ok(build_git_protocol_error_response(
GitService::ReceivePack,
&format!("authorisation failed: {}", e),
Some(&request_body),
));
}
};
let mut push_plan = if let Some(mut push_plan) = push_plan {
let verify_path = repo_path.clone();
Some(
tokio::task::spawn_blocking(move || {
push_plan.verify(&verify_path);
push_plan
})
.await
.map_err(|error| GitError::Storage(error.to_string()))?,
)
} else {
None
};
if let Some(plan) = push_plan.as_mut() {
if !plan.needs_git() {
let body = plan.response(&[]);
plan.release_locks();
record_git_operation(
&metrics,
"push",
if plan.failed { "error" } else { "success" },
);
return Ok(Response::builder()
.status(StatusCode::OK)
.header(
"content-type",
GitService::ReceivePack.result_content_type(),
)
.header("cache-control", "no-cache")
.body(full_body(body))
.unwrap());
}
}
let forwarded_body = push_plan
.as_ref()
.map(|plan| &plan.forwarded_body)
.unwrap_or(&request_body);
let pushed_refs = parse_pushed_refs(&request_body);
let staged = if family_lease.is_some() {
stage_push(&repo_path, &pushed_refs, &auth_result.unsigned_refs).await?
} else {
false
};
// Ref updates remain in the selected view. Objects of a push backed by
// signed events are installed directly in the shared family inventory;
// a staged push keeps them in the view until promotion.
let mut git = GitSubprocess::spawn_with_object_directory(
GitService::ReceivePack,
&repo_path,
false,
git_protocol,
family_lease
.as_ref()
.filter(|_| !staged)
.map(|lease| lease.family_objects_path.as_path()),
)
.map_err(GitError::ProcessSpawnFailed)?;
let stdin = git.take_stdin();
let forwarded_body = forwarded_body.clone();
let stdout = git.take_stdout().ok_or_else(|| {
GitError::IoError(io::Error::new(
io::ErrorKind::BrokenPipe,
"git receive-pack stdout unavailable",
))
})?;
let stderr = git.take_stderr();
// Do not buffer receive-pack stdout: for large pushes Git may spend tens of
// seconds resolving deltas/checking connectivity after the client finishes
// uploading. Streaming sideband progress keeps libgit2 clients from
// hitting their per-recv timeout during that otherwise-silent window.
let (tx, rx) = mpsc::channel::<Result<Frame<Bytes>, io::Error>>(STREAM_CHANNEL_DEPTH);
let new_oids: HashSet<String> = pushed_refs
.iter()
.filter(|(_, new_oid, _)| new_oid != "0000000000000000000000000000000000000000")
.map(|(_, new_oid, _)| new_oid.clone())
.collect();
let request_body_for_errors = request_body.clone();
let identifier = identifier.to_owned();
let git_data_path = git_data_path.to_owned();
tokio::spawn(async move {
// Own the subprocess and verification locks together before the first
// upload await. Cancellation must stop Git before those locks release.
stream_receive_pack_output(
git,
stdin,
forwarded_body,
stdout,
stderr,
tx,
repo_path,
new_oids,
database,
relay,
identifier,
purgatory,
git_data_path,
request_body_for_errors,
repo_lifecycle_guard,
promotion_hooks,
metrics,
storage,
family_key,
family_lease,
pushed_refs,
staged,
push_plan,
)
.await;
});
Ok(streaming_response(GitService::ReceivePack, rx))
}
#[allow(clippy::too_many_arguments)]
async fn stream_receive_pack_output<S, E, I>(
mut git: GitSubprocess,
stdin: Option<I>,
forwarded_body: Bytes,
stdout: S,
stderr: Option<E>,
tx: mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
repo_path: PathBuf,
new_oids: HashSet<String>,
database: SharedDatabase,
relay: LocalRelay,
identifier: String,
purgatory: Arc<Purgatory>,
git_data_path: String,
request_body: Bytes,
repo_lifecycle_guard: Option<LifecycleReadGuard>,
promotion_hooks: Option<Arc<dyn PurgatoryPromotionHooks>>,
metrics: Option<Arc<Metrics>>,
storage: LocalGitStorage,
family_key: FamilyKey,
family_lease: Option<FamilyWriteLease>,
pushed_refs: Vec<(String, String, String)>,
staged: bool,
mut push_plan: Option<ReceivePackPlan>,
) where
I: tokio::io::AsyncWrite + Unpin + Send + 'static,
S: tokio::io::AsyncRead + Unpin + Send + 'static,
E: tokio::io::AsyncRead + Unpin + Send + 'static,
{
let stderr_task = stderr.map(|stderr| tokio::spawn(read_stderr_to_end(stderr)));
// Drain both output pipes while feeding the pack. Git can produce enough
// diagnostics during unpacking to block before it finishes reading stdin.
let upload = async {
if let Some(mut stdin) = stdin {
tokio::select! {
result = stdin.write_all(&forwarded_body) => result?,
_ = tx.closed() => return Err(io::Error::new(io::ErrorKind::BrokenPipe, "push client disconnected")),
}
}
Ok::<_, io::Error>(())
};
let output = async {
let result = if let Some(plan) = push_plan.as_mut() {
plan.pump(stdout, &tx).await
} else {
pump_receive_pack_stdout_to_channel(stdout, &tx).await
};
match result {
(PumpResult::Eof { sent_stdout }, report) => Ok((sent_stdout, report)),
_ => Err(io::Error::new(
io::ErrorKind::BrokenPipe,
"push output interrupted",
)),
}
};
let (mut sent_stdout, terminal_flush) = match tokio::try_join!(upload, output) {
Ok(((), result)) => result,
Err(error) => {
// Reap Git before releasing prepared ref locks or the family lease.
let _ = git.kill().await;
let _ = git.wait().await;
if let Some(task) = stderr_task {
let _ = task.await;
}
record_git_operation(&metrics, "push", "error");
let _ = tx.send(Err(error)).await;
return;
}
};
let status = match git.wait().await {
Ok(status) => status,
Err(e) => {
let _ = tx.send(Err(e)).await;
drop(repo_lifecycle_guard);
record_git_operation(&metrics, "push", "error");
return;
}
};
if let Some(plan) = push_plan.as_mut() {
plan.release_locks();
}
let stderr_output = match stderr_task {
Some(task) => task.await.unwrap_or_default(),
None => Vec::new(),
};
if !status.success() {
if push_plan.is_some() {
// The merged report is still buffered. A failed process must not
// complete the push response, even if some ref updates succeeded.
record_git_operation(&metrics, "push", "error");
error!(
"Git receive-pack failed with {status}: {}",
String::from_utf8_lossy(&stderr_output)
);
let _ = tx
.send(Err(io::Error::other(format!(
"git receive-pack failed with {status}"
))))
.await;
return;
}
if let Some(flush) = terminal_flush {
sent_stdout = true;
if send_body_bytes(&tx, flush).await.is_err() {
drop(repo_lifecycle_guard);
record_git_operation(&metrics, "push", "error");
return;
}
}
drop(repo_lifecycle_guard);
record_git_operation(&metrics, "push", "error");
let stderr_str = String::from_utf8_lossy(&stderr_output);
if is_git_protocol_error(status.code(), &stderr_output) {
warn!(
"Git receive-pack protocol error (returning ERR pkt-line): {}",
stderr_str
);
if !sent_stdout {
let _ = tx
.send(Ok(err_pktline_frame(
GitService::ReceivePack,
&stderr_str,
Some(&request_body),
)))
.await;
}
} else {
error!("Git receive-pack failed: {}", stderr_str);
if !sent_stdout {
let msg = if stderr_str.trim().is_empty() {
format!("git receive-pack failed with code {:?}", status.code())
} else {
stderr_str.to_string()
};
let _ = tx
.send(Ok(err_pktline_frame(
GitService::ReceivePack,
&msg,
Some(&request_body),
)))
.await;
}
}
return;
}
debug!("Git receive-pack stream completed successfully");
if staged {
settle_staged_push(&repo_path).await;
} else if family_lease.is_some() {
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
}
// Received objects and their retention roots are complete. Release the
// family writer before post-push processing: purgatory promotion and
// archive recovery may need to re-enter this same identifier family.
drop(family_lease);
// Release the repository lifecycle read lock once git-receive-pack itself
// has finished. The lock's purpose is to keep deletion/archive/restore from
// removing or replacing the bare repository while Git is actively reading or
// writing it; Git's own internal locking handles concurrent push/fetch ref
// and object consistency. Post-push purgatory promotion may re-enter
// recovery paths that take the lifecycle write lock, so holding this beyond
// the subprocess boundary would deadlock.
drop(repo_lifecycle_guard);
// A successful receive-pack exit can still contain rejected refs. Do not
// immediately apply the complete pending state over Git's partial result.
// The normal background state-recovery path remains responsible for it.
if push_plan.as_ref().is_some_and(|plan| plan.failed) {
record_git_operation(&metrics, "push", "error");
if let Some(report) = terminal_flush {
let _ = send_body_bytes(&tx, report).await;
}
return;
}
// Git's receive-pack progress has already been streamed verbatim, but its
// final flush remains withheld. Run GRASP follow-up work before releasing
// that client-visible success boundary so push completion implies local
// refs/HEAD and promoted events are ready. Follow-up failures remain
// internal/log-only rather than client-visible push rejections.
//
// A complex promotion may itself run longer than the client's receive
// timeout. Only clients that negotiated side-band-64k can safely receive
// invented progress pkt-lines, so keep those clients alive while leaving
// the response for non-sideband clients byte-for-byte unchanged.
let keepalive_task =
if terminal_flush.is_some() && client_negotiated_sideband_64k(&request_body) {
Some(ReceivePackKeepalive::spawn(
tx.clone(),
POST_PUSH_KEEPALIVE_INTERVAL,
))
} else {
None
};
let processing_result = process_newly_available_git_data(
&repo_path,
&new_oids,
&database,
Some(&relay),
&purgatory,
std::path::Path::new(&git_data_path),
promotion_hooks.as_deref(),
)
.await;
// Stop and join the sender before releasing the flush so no keepalive can
// race behind the terminal protocol boundary.
if let Some(task) = keepalive_task {
task.stop().await;
}
match processing_result {
Ok(result) => {
if result.released_any() {
info!(
"Processed push for {}: {} states released, {} PRs released, {} repos synced",
identifier, result.states_released, result.prs_released, result.repos_synced
);
}
if !result.errors.is_empty() {
for error in &result.errors {
warn!(
"Error during post-push processing for {}: {}",
identifier, error
);
}
}
}
Err(e) => {
warn!(
"Failed to process newly available git data after push to {}: {}",
identifier, e
);
}
}
// The final receive-pack flush is the client-visible success boundary.
// Release it only after promoted events have been saved and subscribers
// notified, so a completed push implies that local GRASP state is ready.
if let Some(flush) = terminal_flush {
if send_body_bytes(&tx, flush).await.is_err() {
record_git_operation(&metrics, "push", "error");
return;
}
}
record_git_operation(&metrics, "push", "success");
}
pub(crate) async fn send_body_bytes(
tx: &mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
bytes: Vec<u8>,
) -> Result<(), mpsc::error::SendError<Result<Frame<Bytes>, io::Error>>> {
tx.send(Ok(Frame::data(Bytes::from(bytes)))).await
}
/// Errors that can occur in Git handlers
///
/// These represent transport/infrastructure failures, not application-level
/// rejections. Application-level rejections (e.g. auth failures) are returned
/// as HTTP 200 with an ERR pkt-line so git clients can display the message.
#[derive(Debug)]
pub enum GitError {
RepositoryNotFound,
ProcessSpawnFailed(std::io::Error),
IoError(std::io::Error),
GitFailed(Option<i32>),
Storage(String),
}
impl std::fmt::Display for GitError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::RepositoryNotFound => write!(f, "repository not found"),
Self::ProcessSpawnFailed(e) => write!(f, "failed to spawn git process: {}", e),
Self::IoError(e) => write!(f, "IO error: {}", e),
Self::GitFailed(code) => write!(f, "git process failed with code: {:?}", code),
Self::Storage(error) => write!(f, "Git storage error: {error}"),
}
}
}
/// Decide where a push to a thin view stores its objects, before Git runs.
///
/// A push is staged in the view when it carries a ref that no signed event
/// names, or when the view already holds staged objects: the client may have
/// omitted objects that only staging holds. Signed tips of a staged push are
/// recorded as owed to the family first. The caller holds the family lease.
pub(crate) async fn stage_push(
repo_path: &std::path::Path,
pushed_refs: &[(String, String, String)],
unsigned_refs: &HashSet<String>,
) -> Result<bool, GitError> {
if unsigned_refs.is_empty() && !super::staging::is_staged(repo_path) {
return Ok(false);
}
let owed: Vec<_> = pushed_refs
.iter()
.filter(|(_, new_oid, name)| {
!unsigned_refs.contains(name) && new_oid.bytes().any(|digit| digit != b'0')
})
.map(|(_, new_oid, name)| super::staging::Tip::new(name, new_oid))
.collect();
let unsigned: Vec<_> = pushed_refs
.iter()
.filter(|(_, _, name)| unsigned_refs.contains(name))
.map(|(_, oid, name)| super::staging::Tip::new(name, oid))
.collect();
let view = repo_path.to_owned();
tokio::task::spawn_blocking(move || super::staging::stage_upload(&view, &owed, &unsigned))
.await
.map_err(|error| GitError::Storage(error.to_string()))?
.map_err(|error| GitError::Storage(format!("cannot stage push: {error:#}")))
}
/// Promote the signed tips of a staged push while the family lease is held.
///
/// A failure leaves the tips owed for staging maintenance to retry. The push
/// has already succeeded, so it is not reported to the client.
pub(crate) async fn settle_staged_push(repo_path: &std::path::Path) {
let view = repo_path.to_owned();
match tokio::task::spawn_blocking(move || super::staging::settle(&view)).await {
Ok(Ok(0)) => {}
Ok(Ok(owed)) => warn!(
repo = %repo_path.display(),
owed,
"Signed history of a staged push is not yet in the family"
),
Ok(Err(error)) => error!(
repo = %repo_path.display(),
error = %format!("{error:#}"),
"Failed to promote signed history of a staged push"
),
Err(error) => error!(repo = %repo_path.display(), %error, "Staged push promotion panicked"),
}
super::staging::request_maintenance(repo_path);
}
pub(crate) fn retain_accepted_tips(
storage: &LocalGitStorage,
family_key: &FamilyKey,
repo_path: &std::path::Path,
pushed_refs: &[(String, String, String)],
) {
const ZERO_OID: &str = "0000000000000000000000000000000000000000";
for (_, new_oid, ref_name) in pushed_refs {
if new_oid == ZERO_OID
|| super::get_ref_commit(repo_path, ref_name).as_deref() != Some(new_oid.as_str())
{
continue;
}
if let Err(error) = storage.retain_tip(family_key, ref_name, new_oid) {
warn!(
family = %family_key.identifier,
%ref_name,
%new_oid,
%error,
"Failed to install family retention root"
);
}
if let Err(error) = storage.advertise_base_tip(family_key, ref_name, new_oid) {
warn!(
family = %family_key.identifier,
%ref_name,
%new_oid,
%error,
"Failed to install family negotiation root"
);
}
}
}
impl std::error::Error for GitError {}
impl GitError {
/// Convert to HTTP status code
pub fn status_code(&self) -> StatusCode {
match self {
Self::RepositoryNotFound => StatusCode::NOT_FOUND,
_ => StatusCode::INTERNAL_SERVER_ERROR,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use http_body_util::BodyExt;
#[test]
fn missing_object_probe_errors_are_routine_client_diagnostics() {
assert!(is_routine_missing_object_error(
b"fatal: git upload-pack: not our ref deadbeef"
));
assert!(is_routine_missing_object_error(
b"error: Server does not allow request for unadvertised object deadbeef"
));
assert!(!is_routine_missing_object_error(
b"fatal: unable to read repository data"
));
}
/// Build a single-update receive-pack request body pkt-line for tests.
///
/// Matches the smart-HTTP format:
/// <old-oid> SP <new-oid> SP <ref-name> NUL <capabilities>\n
fn build_receive_pack_pktline(caps: &str) -> Vec<u8> {
let old = "0".repeat(40);
let new = "1".repeat(40);
let refname = "refs/heads/main";
// payload format with NUL between ref-name and capabilities
let mut payload = Vec::new();
payload.extend_from_slice(format!("{} {} {}", old, new, refname).as_bytes());
payload.push(0);
payload.extend_from_slice(caps.as_bytes());
payload.push(b'\n');
// pkt-line frame: hex length prefix
let total_len = payload.len() + 4;
let mut out = Vec::new();
out.extend_from_slice(format!("{:04x}", total_len).as_bytes());
out.extend_from_slice(&payload);
out
}
async fn response_body_bytes(resp: Response<GitResponseBody>) -> Vec<u8> {
resp.into_body()
.collect()
.await
.expect("collect body")
.to_bytes()
.to_vec()
}
#[test]
fn detects_side_band_64k_in_first_pktline() {
let body = build_receive_pack_pktline("report-status side-band-64k agent=git/2.42");
assert!(client_negotiated_sideband_64k(&body));
}
#[test]
fn detects_side_band_64k_when_only_capability() {
let body = build_receive_pack_pktline("side-band-64k");
assert!(client_negotiated_sideband_64k(&body));
}
#[test]
fn detects_no_sideband_when_only_report_status() {
let body = build_receive_pack_pktline("report-status agent=git/2.42");
assert!(!client_negotiated_sideband_64k(&body));
}
#[test]
fn detects_no_sideband_for_legacy_side_band() {
// The legacy `side-band` (without -64k) is NOT what modern clients
// use for the demuxer; treat it as no-sideband to be conservative.
let body = build_receive_pack_pktline("report-status side-band");
assert!(!client_negotiated_sideband_64k(&body));
}
#[test]
fn handles_short_or_malformed_body() {
assert!(!client_negotiated_sideband_64k(b""));
assert!(!client_negotiated_sideband_64k(b"abc"));
// claimed length 4 = flush
assert!(!client_negotiated_sideband_64k(b"0000"));
// invalid hex length
assert!(!client_negotiated_sideband_64k(b"zzzz1234"));
// length larger than body
assert!(!client_negotiated_sideband_64k(b"0099short"));
}
#[tokio::test]
async fn receive_pack_err_is_sideband_wrapped_when_client_negotiates() {
let body = build_receive_pack_pktline("report-status side-band-64k");
let resp = build_git_protocol_error_response(
GitService::ReceivePack,
"authorisation failed: nope",
Some(&body),
);
assert_eq!(resp.status(), StatusCode::OK);
assert_eq!(
resp.headers()
.get("content-type")
.map(|v| v.to_str().unwrap_or("").to_string()),
Some("application/x-git-receive-pack-result".to_string())
);
let raw = response_body_bytes(resp).await;
// pkt-line: <4-hex-length><payload>
assert!(raw.len() >= 4);
let len =
u16::from_str_radix(std::str::from_utf8(&raw[0..4]).unwrap(), 16).unwrap() as usize;
assert_eq!(raw.len(), len);
// First payload byte MUST be the band id 0x03 (error band)
assert_eq!(raw[4], 0x03);
// Remainder is the ERR pkt-line content
let rest = std::str::from_utf8(&raw[5..]).unwrap();
assert!(
rest.starts_with("ERR authorisation failed: nope"),
"unexpected error payload: {:?}",
rest
);
}
#[tokio::test]
async fn receive_pack_err_is_bare_when_no_sideband() {
let body = build_receive_pack_pktline("report-status");
let resp = build_git_protocol_error_response(
GitService::ReceivePack,
"authorisation failed: nope",
Some(&body),
);
let raw = response_body_bytes(resp).await;
// Without sideband, first payload byte is 'E' (ASCII 0x45) from "ERR".
assert!(raw.len() >= 5);
assert_eq!(raw[4], b'E');
let rest = std::str::from_utf8(&raw[4..]).unwrap();
assert!(rest.starts_with("ERR authorisation failed: nope"));
}
#[tokio::test]
async fn upload_pack_err_is_never_sideband_wrapped() {
// Even if we somehow had a body that mentioned side-band-64k, upload-pack
// responses must remain bare ERR pkt-lines.
let body = build_receive_pack_pktline("side-band-64k");
let resp =
build_git_protocol_error_response(GitService::UploadPack, "not our ref", Some(&body));
let raw = response_body_bytes(resp).await;
assert_eq!(raw[4], b'E');
}
#[tokio::test]
async fn streaming_body_yields_stdout_before_eof() {
use tokio::io::AsyncWriteExt;
use tokio::time::{timeout, Duration};
let (mut writer, reader) = tokio::io::duplex(64);
let (tx, rx) = mpsc::channel::<Result<Frame<Bytes>, io::Error>>(STREAM_CHANNEL_DEPTH);
let pump_task = tokio::spawn(async move { pump_stdout_to_channel(reader, &tx).await });
writer.write_all(b"progress chunk").await.unwrap();
let mut body = streaming_response(GitService::ReceivePack, rx).into_body();
let frame = timeout(Duration::from_secs(1), body.frame())
.await
.expect("stream should yield first stdout chunk before EOF")
.expect("stream should still be open")
.expect("stdout chunk should not be a body error");
let data = frame.into_data().expect("frame should contain data");
assert_eq!(data, Bytes::from_static(b"progress chunk"));
assert!(
!pump_task.is_finished(),
"stdout pump should still be waiting for EOF after yielding first chunk"
);
drop(writer);
assert_eq!(
pump_task.await.unwrap(),
PumpResult::Eof { sent_stdout: true }
);
}
#[tokio::test]
async fn post_push_keepalive_is_sideband_progress() {
use tokio::time::timeout;
let (tx, mut rx) = mpsc::channel::<Result<Frame<Bytes>, io::Error>>(STREAM_CHANNEL_DEPTH);
let keepalive = ReceivePackKeepalive::spawn(tx, Duration::from_millis(10));
let frame = timeout(Duration::from_secs(1), rx.recv())
.await
.expect("keepalive should arrive before timeout")
.expect("keepalive channel should remain open")
.expect("keepalive should not be an HTTP body error");
let data = frame.into_data().expect("keepalive should contain data");
assert_eq!(data.as_ref(), receive_pack_keepalive_pktline());
assert_eq!(
data[4], 0x02,
"keepalive pkt-line payload must use sideband band 2"
);
keepalive.stop().await;
}
#[tokio::test]
async fn receive_pack_err_without_body_is_not_wrapped() {
// No request body → conservative path: do not wrap.
let resp = build_git_protocol_error_response(GitService::ReceivePack, "boom", None);
let raw = response_body_bytes(resp).await;
assert_eq!(raw[4], b'E');
}
}