feat(storage): route Git traffic through identifier families

Motivation: owner repositories and GRASP-06 contributor routes currently store or copy the same Git objects independently, forcing clients to upload data the relay already has.

Approach: create new repositories as thin ref views, direct receive-pack and proactive fetch writes into the shared identifier family, retain accepted tips, and expose family bases as anonymous receive negotiation haves. Existing legacy repositories deliberately remain self-contained until the startup migration layer lands.

Correctness assumptions: repositories only share objects when their validated identifier and object format match. A process-wide per-family lease serializes object-producing operations, while the terminal receive-pack flush remains behind ref validation and post-push processing.

Excluded scope: remote S3 durability, cache eviction, garbage collection, and conversion of existing repositories are separate stack layers.

Validation: cargo check --all-targets; cargo test --test git_response_streaming; cargo test --test grasp06_pr_hosting; plus repository_creation, git_clone, and storage/purgatory targeted suites from the preceding review.
This commit is contained in:
DanConwayDev
2026-08-17 15:11:27 +00:00
parent 27f9eaaf98
commit f084f7e973
10 changed files with 341 additions and 108 deletions
+2 -2
View File
@@ -711,8 +711,8 @@ Optional endpoint at `/prs/<npub>/<identifier>.git`, gated on `NGIT_GRASP06_ENAB
- [`src/grasp06/endpoint.rs`](../../src/grasp06/endpoint.rs) — URL parsing.
- [`src/grasp06/paths.rs`](../../src/grasp06/paths.rs) — on-disk path conventions under `<git_data_path>/prs/<hex>/<identifier>.git`.
- [`src/grasp06/fetch.rs`](../../src/grasp06/fetch.rs) — empty-repo synthesis for `info/refs` and `git-upload-pack` against repos that don't yet exist on disk.
- [`src/grasp06/receive.rs`](../../src/grasp06/receive.rs) — `git-receive-pack` with init-on-push, strict `refs/nostr/<event-id>` ref-name validation, and per-ref post-push validation against the database and purgatory. Uses a per-`(submitter, identifier)` `PrsPathState` (mutex + `in_flight` counter) so concurrent pushes to the same path proceed in parallel: the mutex is held only for init+register and decrement+end-of-push cleanup, not across `git-receive-pack` itself. The same state is consulted (with `try_lock` in synchronous contexts) by the PR-event policy and the purgatory expiry sweep, which only remove a `/prs/` ref or zero-ref bare repo when `in_flight == 0`.
- [`src/grasp06/fetch.rs`](../../src/grasp06/fetch.rs) — empty thin-view synthesis for `info/refs` and `git-upload-pack` against routes that do not yet exist on disk. Named refs stay empty while an existing identifier family supplies anonymous receive negotiation bases.
- [`src/grasp06/receive.rs`](../../src/grasp06/receive.rs) — `git-receive-pack` with thin-view init-on-push, strict `refs/nostr/<event-id>` ref-name validation, and per-ref post-push validation against the database and purgatory. A per-`(submitter, identifier)` `PrsPathState` protects view creation/removal; a per-family lease serializes object-producing work across related views. The same path state is consulted by PR-event policy and purgatory expiry, which only remove a `/prs/` ref or zero-ref view when `in_flight == 0`.
- [`src/grasp06/policy.rs`](../../src/grasp06/policy.rs) — strict clone-tag URL comparator used by the PR-event acceptance relaxation.
- [`src/grasp06/cleanup.rs`](../../src/grasp06/cleanup.rs) — one-shot startup scan over `<git_data_path>/prs/` that removes zero-ref bare repos left behind by a previous run (crash mid-push, crash mid-cleanup, or shutdown with unresolved scoped placeholders). Runs before the HTTP server starts accepting requests; no locking is needed because nothing else is touching `/prs/` yet.
@@ -55,7 +55,9 @@ Two reasons:
The contributor chose to publish to `/prs/`. Mirroring into any accepted repository on this relay makes the PR visible at the expected location for clients browsing that repo. Not mirroring the reverse direction preserves the maintainer's declared `clone` intent on their own pushes — we do not invent new hosting locations for their events.
Until object-pool dedup lands, the mirror doubles storage for affected refs. We accept this cost as the simplest correct implementation; dedup lands later and makes the mirror effectively free.
Identifier-family object storage makes the mirror a ref-only operation. Owner
and contributor views for the same identifier resolve one shared object
inventory, so accepting or mirroring a PR does not copy its object graph.
## Architecture
@@ -94,7 +96,7 @@ POST /prs/<npub>/<id>.git/git-receive-pack
signer ≠ URL npub → reject ref
identifier (d-tag) ≠ URL identifier → reject ref
commit ≠ event's c tag → delete ref
all match → ref locked; release event from purgatory;
all match → ref locked; retain family tip; release event from purgatory;
mirror to matching standard repos
event not yet seen:
accept ref; create or update PR placeholder in purgatory
@@ -107,7 +109,14 @@ The flow mirrors the existing `refs/nostr/<event-id>` path at the standard endpo
#### On-demand bare repo creation
The first push to `/prs/<submitter>/<identifier>.git` creates the bare repo on disk. A per-`(submitter, identifier)` [`PrsPathState`](../../src/grasp06/receive.rs) (kept in a `DashMap` on the `HttpService`, see [`crate::grasp06::receive::RepoInitLocks`](../../src/grasp06/receive.rs)) provides a `tokio::sync::Mutex` and an `in_flight: AtomicUsize` counter. The mutex is held only briefly — for `git init --bare` plus the `in_flight` increment at the start of a push, and for the `in_flight` decrement plus end-of-push zero-ref cleanup at the end. The pack upload and per-ref validation run *without* the mutex held, so concurrent pushes to the same path proceed in parallel; git's own ref locking handles intra-push concurrency. Off-push cleanup paths (PR-event policy, purgatory expiry) take the same mutex briefly and only `rm -rf` the bare repo when they see `in_flight == 0` and `list_refs` is empty — so no off-push code path can delete the bare repo while a push is in flight, and no push is serialised behind another push to the same identity.
The first push to `/prs/<submitter>/<identifier>.git` creates a thin bare view
whose alternate is the identifier family. A per-`(submitter, identifier)`
[`PrsPathState`](../../src/grasp06/receive.rs) still protects view creation and
cleanup. A separate per-family write lease serializes object-producing work
across all owner and contributor views for that identifier. `git-receive-pack`
writes objects to the family while updating only the selected view's refs.
Off-push cleanup may remove an empty contributor view, but never removes family
objects or retained roots.
### Event acceptance relaxation
@@ -130,9 +139,10 @@ When a PR or PR Update's purgatory entry is released via a `/prs/` push:
1. Save event to DB and remove from purgatory (as today).
2. For each `a` tag in the event of the form `30617:<pubkey>:<d-tag>`:
- Resolve to a local repo path `<git_data_path>/<pubkey-npub>/<d-tag>.git`.
- If that repo has an active (non-purgatory) announcement, copy objects + install `refs/nostr/<event-id>` into it (same mechanism as existing cross-owner sync in [`src/git/sync.rs`](../../src/git/sync.rs)).
- If that repo has an active (non-purgatory) announcement, install `refs/nostr/<event-id>` in its view after checking the commit resolves through the shared family.
The mirror copies the same ref, same commits. No separate object store. Dedup can be added transparently later via git alternates keyed on d-tag.
The mirror installs the same ref to the same commit. Git alternates keyed by
identifier make the objects immediately available without a copy.
The mirror is **one-directional**: pushes to `/<maintainer>/<id>.git` are not mirrored into `/prs/*`. Only the `/prs/` → `<maintainer>/` direction fires, and only when the source repo path is under `prs_base_path`.
@@ -250,9 +260,10 @@ This is the design rationale for the mirror being one-directional and the `/prs/
Deletion of the contributor's own PR event (NIP-09) → the deletion-request branch will decide and implement the ref-lifecycle behaviour. GRASP-06 v1 does not add hooks or no-op scaffolding for this.
### Object-pool deduplication (not yet specified)
### Object-pool deduplication
The mirror copies objects today. When object-pool dedup exists — likely as git alternates keyed on `(identifier, a-tag-coord-set)` — the mirror becomes a pointer operation with near-zero storage cost. No URL or spec change is required.
The identifier family is the shared object inventory. Mirroring installs a view
ref and does not copy objects; no URL or spec change is required.
## Anticipated failure modes and mitigations
+79 -5
View File
@@ -17,6 +17,7 @@ use tokio::time::MissedTickBehavior;
use tracing::{debug, error, info, warn};
use super::protocol::{GitService, PktLine};
use super::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
use super::subprocess::GitSubprocess;
use super::{full_body, GitResponseBody};
@@ -343,7 +344,7 @@ where
/// caller finish GRASP post-push processing before making success visible,
/// without buffering the progress stream that keeps clients alive during
/// expensive pack processing.
async fn pump_receive_pack_stdout_to_channel<R>(
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>>)
@@ -701,9 +702,34 @@ pub async fn handle_receive_pack(
}
};
// Spawn git receive-pack
let mut git = GitSubprocess::spawn(GitService::ReceivePack, &repo_path, false, git_protocol)
.map_err(GitError::ProcessSpawnFailed)?;
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
};
// Ref updates remain in the selected view; receive-pack's quarantine and
// final objects are installed directly in the shared family inventory.
let mut git = GitSubprocess::spawn_with_object_directory(
GitService::ReceivePack,
&repo_path,
false,
git_protocol,
family_lease
.as_ref()
.map(|lease| lease.family_objects_path.as_path()),
)
.map_err(GitError::ProcessSpawnFailed)?;
// Write request to git's stdin
if let Some(mut stdin) = git.take_stdin() {
@@ -754,6 +780,10 @@ pub async fn handle_receive_pack(
repo_lifecycle_guard,
promotion_hooks,
metrics,
storage,
family_key,
family_lease,
pushed_refs,
)
.await;
});
@@ -778,6 +808,10 @@ async fn stream_receive_pack_output<S, E>(
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)>,
) where
S: tokio::io::AsyncRead + Unpin + Send + 'static,
E: tokio::io::AsyncRead + Unpin + Send + 'static,
@@ -863,6 +897,10 @@ async fn stream_receive_pack_output<S, E>(
debug!("Git receive-pack stream completed successfully");
if family_lease.is_some() {
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
}
// 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
@@ -948,7 +986,7 @@ async fn stream_receive_pack_output<S, E>(
record_git_operation(&metrics, "push", "success");
}
async fn send_body_bytes(
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>>> {
@@ -966,6 +1004,7 @@ pub enum GitError {
ProcessSpawnFailed(std::io::Error),
IoError(std::io::Error),
GitFailed(Option<i32>),
Storage(String),
}
impl std::fmt::Display for GitError {
@@ -975,6 +1014,41 @@ impl std::fmt::Display for GitError {
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}"),
}
}
}
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"
);
}
}
}
+31 -9
View File
@@ -8,7 +8,7 @@
use std::fmt;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use std::sync::{Arc, LazyLock};
use anyhow::{anyhow, Context, Result};
use bitcoin_hashes::{sha256, Hash};
@@ -26,6 +26,9 @@ const FAMILY_DIR: &str = "families";
const BASE_REFS_PREFIX: &str = "refs/grasp/bases/";
const RETAINED_REFS_PREFIX: &str = "refs/grasp/retained/";
type FamilyLockId = (PathBuf, FamilyKey);
static FAMILY_LOCKS: LazyLock<DashMap<FamilyLockId, Arc<Mutex<()>>>> = LazyLock::new(DashMap::new);
/// Git object hash algorithms are separate family namespaces.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ObjectFormat {
@@ -100,15 +103,19 @@ impl FamilyKey {
#[derive(Clone)]
pub struct LocalGitStorage {
git_data_path: PathBuf,
family_locks: Arc<DashMap<FamilyKey, Arc<Mutex<()>>>>,
}
impl LocalGitStorage {
pub fn new(git_data_path: impl Into<PathBuf>) -> Self {
Self {
git_data_path: git_data_path.into(),
family_locks: Arc::new(DashMap::new()),
}
let git_data_path = git_data_path.into();
let git_data_path = if git_data_path.is_absolute() {
git_data_path
} else {
std::env::current_dir()
.unwrap_or_else(|_| PathBuf::from("."))
.join(git_data_path)
};
Self { git_data_path }
}
pub fn git_data_path(&self) -> &Path {
@@ -134,11 +141,26 @@ impl LocalGitStorage {
self.family_repo_path(key).join("objects")
}
pub fn family_exists(&self, key: &FamilyKey) -> bool {
self.family_repo_path(key).is_dir()
}
pub fn is_thin_view(&self, key: &FamilyKey, view_path: &Path) -> bool {
let alternate = view_path.join("objects/info/alternates");
let Ok(configured) = std::fs::read_to_string(alternate) else {
return false;
};
let Ok(expected) = std::fs::canonicalize(self.family_objects_path(key)) else {
return false;
};
configured.lines().any(|line| Path::new(line) == expected)
}
/// Serialize mutations to one identifier family.
pub async fn write_lease(&self, key: &FamilyKey) -> Result<FamilyWriteLease> {
let lock = self
.family_locks
.entry(key.clone())
let lock_id = (self.git_data_path.clone(), key.clone());
let lock = FAMILY_LOCKS
.entry(lock_id)
.or_insert_with(|| Arc::new(Mutex::new(())))
.clone();
let guard = lock.lock_owned().await;
+24
View File
@@ -28,6 +28,27 @@ impl GitSubprocess {
repo_path: impl AsRef<Path>,
advertise: bool,
git_protocol: Option<&str>,
) -> std::io::Result<Self> {
Self::spawn_with_object_directory(
service,
repo_path,
advertise,
git_protocol,
None::<&Path>,
)
}
/// Spawn Git with an optional writable object database override.
///
/// Thin views normally read through `objects/info/alternates`. Commands
/// that create objects use this override so receive-pack quarantine output
/// and fetched objects land directly in the identifier family.
pub fn spawn_with_object_directory(
service: GitService,
repo_path: impl AsRef<Path>,
advertise: bool,
git_protocol: Option<&str>,
object_directory: Option<impl AsRef<Path>>,
) -> std::io::Result<Self> {
let repo_path = repo_path.as_ref();
@@ -62,6 +83,9 @@ impl GitSubprocess {
if let Some(protocol) = git_protocol {
cmd.env("GIT_PROTOCOL", protocol);
}
if let Some(object_directory) = object_directory {
cmd.env("GIT_OBJECT_DIRECTORY", object_directory.as_ref());
}
let child = cmd.spawn()?;
+22 -8
View File
@@ -5,10 +5,11 @@
//! > MUST respond to upload-pack requests for any well-formed path as if
//! > serving an empty bare repository.
//!
//! When a real `/prs/<submitter>/<identifier>.git` repo exists on disk
//! we delegate to the standard handlers in [`crate::git::handlers`].
//! Otherwise we synthesise a brand-new empty bare repo in a per-request
//! temporary directory and run the standard upload-pack against that.
//! When a real `/prs/<submitter>/<identifier>.git` repo exists on disk we
//! delegate to the standard handlers in [`crate::git::handlers`]. Otherwise we
//! synthesize an empty thin view in a per-request temporary directory. If the
//! identifier family exists, receive-pack discovery can advertise its bases as
//! anonymous `.have` entries without exposing any named refs.
//!
//! Receive-pack lives in [`crate::grasp06::receive`].
@@ -22,6 +23,7 @@ use tracing::{debug, warn};
use crate::git::handlers::{handle_info_refs, handle_upload_pack, GitError};
use crate::git::protocol::GitService;
use crate::git::storage::{FamilyKey, LocalGitStorage};
use crate::git::GitResponseBody;
use crate::grasp06::endpoint::PrsUrl;
use crate::grasp06::paths::prs_repo_path;
@@ -60,7 +62,9 @@ pub async fn handle_prs_info_refs(
prs.submitter.to_hex(),
prs.identifier
);
let temp = init_empty_bare_repo()?;
let storage = LocalGitStorage::new(git_data_path);
let key = FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
let temp = init_empty_bare_repo(&storage, &key)?;
let repo_path = temp.path().to_path_buf();
let response = handle_info_refs(repo_path, service, git_protocol).await;
// `temp` is dropped here, deleting the directory. Git has already
@@ -98,7 +102,9 @@ pub async fn handle_prs_upload_pack(
prs.submitter.to_hex(),
prs.identifier
);
let temp = init_empty_bare_repo()?;
let storage = LocalGitStorage::new(git_data_path);
let key = FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
let temp = init_empty_bare_repo(&storage, &key)?;
handle_prs_upload_pack_buffered(temp, body, git_protocol).await
}
@@ -169,7 +175,7 @@ async fn handle_prs_upload_pack_buffered(
.unwrap())
}
/// Create a fresh empty bare repo in a temp directory.
/// Create a fresh empty thin view in a temp directory.
///
/// Per-request temp dirs are deliberate. They are cheap on
/// any reasonable filesystem (one `mkdir`, one `git init --bare`) and
@@ -177,7 +183,10 @@ async fn handle_prs_upload_pack_buffered(
/// profiling shows this is a bottleneck we can switch to a shared
/// `<git_data_path>/prs/.empty-template.git`; do not optimise until
/// measurable.
fn init_empty_bare_repo() -> Result<TempDir, GitError> {
fn init_empty_bare_repo(
storage: &LocalGitStorage,
family_key: &FamilyKey,
) -> Result<TempDir, GitError> {
let temp = TempDir::new().map_err(GitError::IoError)?;
let path = temp.path();
let output = Command::new("git")
@@ -194,5 +203,10 @@ fn init_empty_bare_repo() -> Result<TempDir, GitError> {
);
return Err(GitError::GitFailed(output.status.code()));
}
if storage.family_exists(family_key) {
storage
.configure_thin_view(family_key, path)
.map_err(|error| GitError::Storage(error.to_string()))?;
}
Ok(temp)
}
+82 -53
View File
@@ -19,19 +19,21 @@
//! reject the whole push with an `ERR` pkt-line — matching the
//! standard-endpoint UX. Nothing on disk has been touched yet so failed
//! probes leave no state.
//! 3. Acquires the per-path coordination state from [`RepoInitLocks`]
//! briefly: under its mutex it runs the on-demand `git init --bare`
//! 3. Acquires the identifier-family write lease and the per-path coordination
//! state from [`RepoInitLocks`] briefly: under the path mutex it creates an
//! on-demand thin view
//! and increments the `in_flight` counter, then releases the mutex.
//! Steps 4 and 5 run *without* the per-path lock so concurrent pushes
//! to the same `(submitter, identifier)` proceed in parallel — git's
//! own ref locking handles intra-push concurrency, and the
//! `in_flight` counter is what off-push cleanup paths consult to know
//! a push is active.
//! 4. Starts `git-receive-pack`, writes the full request body to the child, and
//! Steps 4 and 5 run *without* the per-path lock. The family write lease
//! serializes object-producing operations for the identifier, while the
//! `in_flight` counter is what off-push cleanup paths consult to know a push
//! is active.
//! 4. Starts `git-receive-pack` with the family as its writable object database,
//! writes the full request body to the child, and
//! immediately returns an HTTP response whose body is backed by a bounded
//! channel. From this point on Hyper can stream stdout to the client while a
//! detached task owns the subprocess, stderr, and all post-push work.
//! 5. In the detached task, each stdout chunk is forwarded as it is read. If Git
//! 5. In the detached task, progress is forwarded while the terminal flush is
//! retained until family roots and GRASP post-processing are complete. If Git
//! exits with a protocol-level error before writing stdout, the task sends a
//! Git `ERR` pkt-line through the same stream so clients still see a normal
//! receive-pack failure. If stdout has already been sent, the task cannot
@@ -76,10 +78,12 @@ use crate::git::authorization::{
};
use crate::git::handlers::{
build_git_protocol_error_response, err_pktline_frame, is_git_protocol_error,
pump_stdout_to_channel, read_stderr_to_end, record_git_operation, streaming_response, GitError,
PumpResult, STREAM_CHANNEL_DEPTH,
pump_receive_pack_stdout_to_channel, read_stderr_to_end, record_git_operation,
retain_accepted_tips, send_body_bytes, streaming_response, GitError, PumpResult,
STREAM_CHANNEL_DEPTH,
};
use crate::git::protocol::GitService;
use crate::git::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
use crate::git::subprocess::GitSubprocess;
use crate::git::sync::process_newly_available_git_data;
use crate::git::{delete_ref, list_refs, GitResponseBody};
@@ -107,9 +111,9 @@ use crate::sync::rejected_index::RejectedEventsIndex;
/// expiry) for the duration of one `delete_ref` + optional
/// `remove_dir_all`.
///
/// `git-receive-pack` itself and per-ref validation run *without* the
/// mutex held, so two pushes to the same path proceed in parallel; git's
/// own ref locking handles intra-push concurrency.
/// `git-receive-pack` itself and per-ref validation run *without* the path
/// mutex held. A separate identifier-family lease serializes their shared
/// object inventory.
///
/// Off-push cleanup paths only `rm -rf` the bare repo when both
/// `in_flight.load() == 0` *and* `list_refs` returns empty while they
@@ -258,20 +262,27 @@ pub async fn handle_prs_receive_pack(
// 3. Acquire the per-path coordination state and, under its mutex,
// initialise the bare repo on demand and register this request as
// in-flight. The mutex is then released — `git-receive-pack` and
// per-ref validation run WITHOUT the lock so concurrent pushes to
// the same `(submitter, identifier)` proceed in parallel. Cleanup
// paths consult `in_flight` (under the same mutex) before
// per-ref validation run WITHOUT the path lock. The family lease above
// serializes object-producing work for the identifier. Cleanup paths
// consult `in_flight` (under the same mutex) before
// deleting the bare repo, so a repo can never vanish mid-receive.
let repo_path = prs_repo_path(
Path::new(git_data_path),
&prs.submitter.to_hex(),
&prs.identifier,
);
let storage = LocalGitStorage::new(git_data_path);
let family_key =
FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
let family_lease = storage
.write_lease(&family_key)
.await
.map_err(|e| GitError::Storage(e.to_string()))?;
let state = path_state(&repo_init_locks, &repo_path);
{
let _g = state.mu.lock().expect("prs path mutex poisoned");
if let Err(e) = ensure_repo_initialised(&repo_path) {
if let Err(e) = ensure_repo_initialised(&storage, &family_key, &repo_path) {
error!(
"/prs/ receive-pack: failed to initialise repo at {}: {}",
repo_path.display(),
@@ -289,17 +300,23 @@ pub async fn handle_prs_receive_pack(
// to the same streaming shape as the standard receive-pack endpoint:
// write stdin here, then hand stdout/stderr and all follow-up state to a
// detached task that feeds the response body channel.
let mut git =
match GitSubprocess::spawn(GitService::ReceivePack, &repo_path, false, git_protocol)
.map_err(GitError::ProcessSpawnFailed)
{
Ok(git) => git,
Err(e) => {
finish_prs_receive_pack(&state, &repo_path);
record_git_operation(&metrics, "push", "error");
return Err(e);
}
};
let use_family = storage.is_thin_view(&family_key, &repo_path);
let mut git = match GitSubprocess::spawn_with_object_directory(
GitService::ReceivePack,
&repo_path,
false,
git_protocol,
use_family.then_some(&family_lease.family_objects_path),
)
.map_err(GitError::ProcessSpawnFailed)
{
Ok(git) => git,
Err(e) => {
finish_prs_receive_pack(&state, &repo_path);
record_git_operation(&metrics, "push", "error");
return Err(e);
}
};
if let Some(mut stdin) = git.take_stdin() {
if let Err(e) = stdin.write_all(&request_body).await {
@@ -350,6 +367,9 @@ pub async fn handle_prs_receive_pack(
prs.identifier.clone(),
domain.to_string(),
metrics,
storage,
family_key,
use_family.then_some(family_lease),
));
Ok(streaming_response(GitService::ReceivePack, rx))
@@ -389,30 +409,18 @@ fn invalid_ref_reason(ref_name: &str) -> Option<String> {
///
/// The caller must hold the per-path mutex from [`PrsPathState`] for
/// `repo_path` before invoking this function.
fn ensure_repo_initialised(repo_path: &Path) -> Result<(), GitError> {
fn ensure_repo_initialised(
storage: &LocalGitStorage,
family_key: &FamilyKey,
repo_path: &Path,
) -> Result<(), GitError> {
if repo_path.exists() {
return Ok(());
}
if let Some(parent) = repo_path.parent() {
std::fs::create_dir_all(parent).map_err(GitError::IoError)?;
}
let output = std::process::Command::new("git")
.args(["init", "--bare", "--initial-branch=main", "--quiet"])
.arg(repo_path)
.output()
.map_err(GitError::ProcessSpawnFailed)?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
error!(
"/prs/ git init --bare failed at {}: {}",
repo_path.display(),
stderr.trim()
);
return Err(GitError::GitFailed(output.status.code()));
}
storage
.create_thin_view(family_key, repo_path)
.map_err(|error| GitError::Storage(error.to_string()))?;
info!(
"/prs/ initialised bare repo at {} on demand",
@@ -482,6 +490,9 @@ async fn stream_prs_receive_pack_output<S, E>(
identifier: String,
domain: String,
metrics: Option<Arc<Metrics>>,
storage: LocalGitStorage,
family_key: FamilyKey,
family_lease: Option<FamilyWriteLease>,
) where
S: tokio::io::AsyncRead + Unpin + Send + 'static,
E: tokio::io::AsyncRead + Unpin + Send + 'static,
@@ -490,7 +501,7 @@ async fn stream_prs_receive_pack_output<S, E>(
// fill its stderr pipe while still producing stdout; draining both prevents
// child-process deadlock and preserves stderr for protocol-error reporting.
let stderr_task = stderr.map(|stderr| tokio::spawn(read_stderr_to_end(stderr)));
let pump_result = pump_stdout_to_channel(stdout, &tx).await;
let (pump_result, terminal_flush) = pump_receive_pack_stdout_to_channel(stdout, &tx).await;
// If the client goes away or stdout read fails, stop Git rather than letting
// it continue writing into a response nobody can receive. Cleanup below will
@@ -514,7 +525,7 @@ async fn stream_prs_receive_pack_output<S, E>(
None => Vec::new(),
};
let sent_stdout = match pump_result {
let mut sent_stdout = match pump_result {
PumpResult::Eof { sent_stdout } => sent_stdout,
PumpResult::ClientDisconnected | PumpResult::ReadError => {
finish_prs_receive_pack(&state, &repo_path);
@@ -524,6 +535,14 @@ async fn stream_prs_receive_pack_output<S, E>(
};
if !status.success() {
if let Some(flush) = terminal_flush {
sent_stdout = true;
if send_body_bytes(&tx, flush).await.is_err() {
finish_prs_receive_pack(&state, &repo_path);
record_git_operation(&metrics, "push", "error");
return;
}
}
record_git_operation(&metrics, "push", "error");
let stderr_str = String::from_utf8_lossy(&stderr_output);
if is_git_protocol_error(status.code(), &stderr_output) {
@@ -568,8 +587,6 @@ async fn stream_prs_receive_pack_output<S, E>(
}
debug!("/prs/ git-receive-pack stream completed successfully");
record_git_operation(&metrics, "push", "success");
// Race safety net. The pre-validation in `handle_prs_receive_pack` was
// performed before `git-receive-pack` ran, so an event with one of the
// pushed ids may have arrived via WebSocket during the receive-pack window.
@@ -595,6 +612,10 @@ async fn stream_prs_receive_pack_output<S, E>(
.await;
}
if family_lease.is_some() {
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
}
finish_prs_receive_pack(&state, &repo_path);
// Drive the standard purgatory-release pipeline so PR events already
@@ -633,6 +654,14 @@ async fn stream_prs_receive_pack_output<S, E>(
);
}
}
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");
}
/// Race safety net for the `/prs/` receive-pack post-push phase.
+7 -19
View File
@@ -411,26 +411,14 @@ impl AnnouncementPolicy {
return Ok(());
}
// Create parent directory (npub directory)
let parent = repo_path
.parent()
.ok_or_else(|| format!("Invalid repository path: {}", repo_path.display()))?;
let storage = crate::git::storage::LocalGitStorage::new(&self.ctx.git_data_path);
let key = crate::git::storage::FamilyKey::sha1(&announcement.identifier)
.map_err(|error| error.to_string())?;
storage
.create_thin_view(&key, &repo_path)
.map_err(|error| error.to_string())?;
std::fs::create_dir_all(parent)
.map_err(|e| format!("Failed to create directory {}: {}", parent.display(), e))?;
// Initialize bare repository using git command
let output = std::process::Command::new("git")
.args(["init", "--bare", repo_path.to_str().unwrap()])
.output()
.map_err(|e| format!("Failed to execute git init: {}", e))?;
if !output.status.success() {
let stderr = String::from_utf8_lossy(&output.stderr);
return Err(format!("git init failed: {}", stderr));
}
tracing::info!("Created bare repository at {}", repo_path.display());
tracing::info!("Created thin repository view at {}", repo_path.display());
Ok(())
}
+46
View File
@@ -229,6 +229,7 @@ use tokio::io::{AsyncRead, AsyncReadExt};
use tokio::process::Command;
use tracing::debug;
use crate::git::storage::{FamilyKey, LocalGitStorage};
use crate::nostr::builder::Nip34WritePolicy;
use crate::nostr::events::RepositoryState;
use crate::nostr::SharedDatabase;
@@ -436,6 +437,7 @@ fn hardened_git_command(
repo_path: &Path,
resolve_pin: Option<&str>,
auth_header: Option<&str>,
object_directory: Option<&Path>,
args: &[String],
) -> Command {
let mut command = Command::new("git");
@@ -469,6 +471,9 @@ fn hardened_git_command(
.stderr(std::process::Stdio::piped())
.kill_on_drop(true)
.current_dir(repo_path);
if let Some(object_directory) = object_directory {
command.env("GIT_OBJECT_DIRECTORY", object_directory);
}
command
}
@@ -974,6 +979,23 @@ impl SyncContext for RealSyncContext {
};
let resolve_pin = resolve_pin_entry(&resolved);
let storage = LocalGitStorage::new(&self.git_data_path);
let family_key =
crate::git::sync::extract_identifier_from_repo_path(repo_path, &self.git_data_path)
.and_then(|identifier| FamilyKey::sha1(identifier).ok());
let family_lease = match family_key.as_ref() {
Some(key) if storage.is_thin_view(key, repo_path) => Some(
storage
.write_lease(key)
.await
.map_err(|error| anyhow::anyhow!("open family for git fetch: {error}"))?,
),
_ => None,
};
let family_objects = family_lease
.as_ref()
.map(|lease| lease.family_objects_path.as_path());
let repo_path = repo_path.to_path_buf();
let url = url.to_string();
let domain = extract_domain(&url).unwrap_or_else(|| "unknown".to_string());
@@ -1028,6 +1050,7 @@ impl SyncContext for RealSyncContext {
&repo_path,
resolve_pin.as_deref(),
fresh_auth_header().as_deref(),
None,
&ls_remote_args,
),
&domain,
@@ -1098,6 +1121,7 @@ impl SyncContext for RealSyncContext {
&repo_path,
resolve_pin.as_deref(),
fresh_auth_header().as_deref(),
family_objects,
&args,
),
&domain,
@@ -1168,6 +1192,7 @@ impl SyncContext for RealSyncContext {
&repo_path,
resolve_pin.as_deref(),
fresh_auth_header().as_deref(),
family_objects,
&args,
),
&domain,
@@ -1219,6 +1244,27 @@ impl SyncContext for RealSyncContext {
.cloned()
.collect();
if let Some(key) = family_key.as_ref().filter(|_| family_lease.is_some()) {
for oid in &fetched {
if let Err(error) = storage.retain_tip(key, "purgatory-fetch", oid) {
tracing::warn!(
identifier = %key.identifier,
%oid,
%error,
"Failed to retain proactively fetched family tip"
);
}
if let Err(error) = storage.advertise_base_tip(key, "purgatory-fetch", oid) {
tracing::warn!(
identifier = %key.identifier,
%oid,
%error,
"Failed to advertise proactively fetched family tip"
);
}
}
}
// One summary line per outbound fetch pass: this is the
// client-side visibility the drop-one-and-retry loop lacked
// (it logged only at debug level, so a relay crawling a remote
+30 -5
View File
@@ -300,7 +300,7 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() {
let _env_lock = PATH_ENV_LOCK.lock().await;
let fake_bin = tempfile::tempdir().expect("fake git bin tempdir");
write_fake_git(fake_bin.path());
write_fake_git_with_terminal_flush(fake_bin.path());
let _path = PathOverride::prepend(fake_bin.path());
let keys = Keys::generate();
@@ -355,7 +355,11 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() {
.expect("first /prs/ stdout chunk should arrive before fake git exits")
.expect("/prs/ body should still be open")
.expect("first /prs/ frame should not be an HTTP body error");
assert_eq!(frame_data(first), Bytes::from_static(b"first-progress\n"));
let mut streamed = frame_data(first).to_vec();
assert!(
b"first-progress\n".starts_with(&streamed),
"the four-byte terminal look-behind may split the first progress chunk"
);
assert!(
repo_path.exists(),
@@ -374,7 +378,15 @@ async fn prs_receive_pack_streams_stdout_before_cleanup_removes_empty_repo() {
.expect("second /prs/ stdout chunk should arrive after fake git wakes")
.expect("/prs/ body should still be open for second chunk")
.expect("second /prs/ frame should not be an HTTP body error");
assert_eq!(frame_data(second), Bytes::from_static(b"second-progress\n"));
streamed.extend_from_slice(&frame_data(second));
let terminal = timeout(Duration::from_secs(1), body.frame())
.await
.expect("/prs/ terminal flush should arrive after cleanup")
.expect("/prs/ body should contain the terminal flush")
.expect("/prs/ terminal frame should not be an HTTP body error");
streamed.extend_from_slice(&frame_data(terminal));
assert_eq!(streamed, b"first-progress\nsecond-progress\n0000");
let eof = timeout(Duration::from_secs(1), body.frame())
.await
@@ -474,15 +486,28 @@ fi
if [ "$1" = "init" ]; then
repo="${@: -1}"
mkdir -p "$repo"
mkdir -p "$repo/objects/info"
exit 0
fi
if [ "$1" = "config" ]; then
exit 0
fi
if [ "$1" = "rev-parse" ] && [ "${2:-}" = "--show-object-format" ]; then
printf 'sha1\n'
exit 0
fi
if [ "$1" = "cat-file" ]; then
exit 1
fi
if [ "$1" = "for-each-ref" ]; then
exit 0
fi
echo "fake git only supports init, for-each-ref, and receive-pack" >&2
echo "unsupported fake git invocation: $*" >&2
exit 1
"#
.replace("__TERMINAL_FLUSH__", terminal_flush),