diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 86b9137..6547226 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -711,8 +711,8 @@ Optional endpoint at `/prs//.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 `/prs//.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/` 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/` 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 `/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. diff --git a/docs/explanation/grasp-06-contributor-pr-submission.md b/docs/explanation/grasp-06-contributor-pr-submission.md index fdb1a3b..b7228b4 100644 --- a/docs/explanation/grasp-06-contributor-pr-submission.md +++ b/docs/explanation/grasp-06-contributor-pr-submission.md @@ -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//.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/` path at the standard endpo #### On-demand bare repo creation -The first push to `/prs//.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//.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::`: - Resolve to a local repo path `//.git`. - - If that repo has an active (non-purgatory) announcement, copy objects + install `refs/nostr/` 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/` 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 `//.git` are not mirrored into `/prs/*`. Only the `/prs/` → `/` 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 diff --git a/src/git/handlers.rs b/src/git/handlers.rs index 9b28d15..4fae99b 100644 --- a/src/git/handlers.rs +++ b/src/git/handlers.rs @@ -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( +pub(crate) async fn pump_receive_pack_stdout_to_channel( mut stdout: R, tx: &mpsc::Sender, io::Error>>, ) -> (PumpResult, Option>) @@ -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( repo_lifecycle_guard: Option, promotion_hooks: Option>, metrics: Option>, + storage: LocalGitStorage, + family_key: FamilyKey, + family_lease: Option, + 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( 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( record_git_operation(&metrics, "push", "success"); } -async fn send_body_bytes( +pub(crate) async fn send_body_bytes( tx: &mpsc::Sender, io::Error>>, bytes: Vec, ) -> Result<(), mpsc::error::SendError, io::Error>>> { @@ -966,6 +1004,7 @@ pub enum GitError { ProcessSpawnFailed(std::io::Error), IoError(std::io::Error), GitFailed(Option), + 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" + ); } } } diff --git a/src/git/storage.rs b/src/git/storage.rs index 14eb647..06013f7 100644 --- a/src/git/storage.rs +++ b/src/git/storage.rs @@ -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>>> = 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>>>, } impl LocalGitStorage { pub fn new(git_data_path: impl Into) -> 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 { - 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; diff --git a/src/git/subprocess.rs b/src/git/subprocess.rs index 37fa382..1dc9c4e 100644 --- a/src/git/subprocess.rs +++ b/src/git/subprocess.rs @@ -28,6 +28,27 @@ impl GitSubprocess { repo_path: impl AsRef, advertise: bool, git_protocol: Option<&str>, + ) -> std::io::Result { + 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, + advertise: bool, + git_protocol: Option<&str>, + object_directory: Option>, ) -> std::io::Result { 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()?; diff --git a/src/grasp06/fetch.rs b/src/grasp06/fetch.rs index e487b40..b7e336f 100644 --- a/src/grasp06/fetch.rs +++ b/src/grasp06/fetch.rs @@ -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//.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//.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 /// `/prs/.empty-template.git`; do not optimise until /// measurable. -fn init_empty_bare_repo() -> Result { +fn init_empty_bare_repo( + storage: &LocalGitStorage, + family_key: &FamilyKey, +) -> Result { 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 { ); 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) } diff --git a/src/grasp06/receive.rs b/src/grasp06/receive.rs index 9987383..bea8c49 100644 --- a/src/grasp06/receive.rs +++ b/src/grasp06/receive.rs @@ -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 { /// /// 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( identifier: String, domain: String, metrics: Option>, + storage: LocalGitStorage, + family_key: FamilyKey, + family_lease: Option, ) 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( // 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( 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( }; 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( } 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( .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( ); } } + + 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. diff --git a/src/nostr/policy/announcement.rs b/src/nostr/policy/announcement.rs index a311354..7a9af7b 100644 --- a/src/nostr/policy/announcement.rs +++ b/src/nostr/policy/announcement.rs @@ -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(()) } diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index 07f9835..4fc0c13 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -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 diff --git a/tests/git_response_streaming.rs b/tests/git_response_streaming.rs index 5f64ceb..f98072d 100644 --- a/tests/git_response_streaming.rs +++ b/tests/git_response_streaming.rs @@ -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),