diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index be4f50f..63f0ed5 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -713,8 +713,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/git-family-object-storage.md b/docs/explanation/git-family-object-storage.md index 26da891..600a3c4 100644 --- a/docs/explanation/git-family-object-storage.md +++ b/docs/explanation/git-family-object-storage.md @@ -302,6 +302,42 @@ View lifecycle locks remain responsible for “may this path be removed?” The family lock is responsible for “is this object inventory and manifest update atomic?” Neither lock grants authorization. +## Integrity and healing + +Integrity is defined for an object-format/identifier family, not for one +owner path. A pass checks the family pack set and object graph, then checks +every owner and `/prs/` view with the same identifier for the correct alternate +and for refs whose targets are available in the family. Multiple independent +histories in one identifier are valid, and unreachable objects are not an +error: retained delete-state and rollback data intentionally remain present. + +The relay starts one non-blocking pass after migration, database +initialization, and construction of the hardened outbound Git client. Broken +alternate wiring is repaired locally. For missing OIDs, the pass tries clone +URLs from accepted repository announcements, excluding this service and +applying the same SSRF policy, DNS pinning, credentials, process containment, +and missing-OID behavior used by proactive sync. It then runs the same +integrity check again. An unresolved family produces an `ERROR` log containing +bounded counts and OID/diagnostic samples; it does not make availability depend +on remote servers. + +This is also the migration repair path. Migration remains a deterministic, +offline conversion that preserves every Git-readable object and the exact +legacy refs. Once those paths are thin family views, the ordinary family pass +can heal pre-existing missing objects. Unindexed legacy packs remain in the +migration backup; there is no separate legacy repair subsystem. + +Operators can queue the same identifier-scoped check in the live process: + +```console +ngit-grasp integrity-check --identifier example +ngit-grasp integrity-check --identifier example --repair +``` + +The command writes a durable request beneath `.grasp/integrity-requests/`. +The server consumes it while holding its normal in-process family locks, so a +manual repair cannot race an object-producing request in another view. + ## Security and privacy trade-offs - Sharing is restricted to one validated identifier, object format, and @@ -349,13 +385,16 @@ and physical deletion have an explicit policy. ## Delivery sequence -The implementation is intentionally reviewable as a stack: +The implementation is intentionally reviewable as a local-first stack: 1. this decision and its invariants; 2. family paths, local inventory, thin-view construction, and Git-level tests; 3. standard and `/prs/` handler integration plus ref-only synchronization; -4. opt-in S3 manifests, verified hydration, cache, and the durability fence; -5. launch-time legacy migration, restart recovery, and operator documentation. +4. launch-time legacy migration, restart recovery, and operator documentation; +5. one optional S3 layer containing manifests, verified hydration, cache, + local-family adoption, configuration, and the durability fence. -Each layer retains the public repository layout and can be reviewed against the -invariants above. The S3 layer never changes the default from local storage. +Steps 3 and 4 ship together as the first deployable local-storage milestone. +It can be operated and observed before the S3 layer is considered. Every layer +retains the public repository layout, and S3 never changes the default from +local storage. 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/docs/how-to/README.md b/docs/how-to/README.md index 087ae53..404294d 100644 --- a/docs/how-to/README.md +++ b/docs/how-to/README.md @@ -24,6 +24,19 @@ How-to guides are **recipes** that show you how to solve specific problems or ac ## Available How-To Guides +### [Upgrade Git family storage](upgrade-git-family-storage.md) +**Problem:** Deduplicate existing repository objects during a server upgrade +**Difficulty:** Advanced + +**You'll learn:** +- Prepare capacity and a release rollback point +- Run the automatic crash-safe launch migration +- Verify owner and `/prs/` repository views +- Check or repair one identifier family on demand +- Recover safely from an interrupted launch + +--- + ### [Configure Nix Flakes](nix-flakes.md) **Problem:** Set up reproducible development environment **Difficulty:** Intermediate diff --git a/docs/how-to/deploy.md b/docs/how-to/deploy.md index 41cf889..d0749fe 100644 --- a/docs/how-to/deploy.md +++ b/docs/how-to/deploy.md @@ -283,6 +283,9 @@ git ls-remote https://ngit.example.com//.git - `dataDir` - Base directory for data (default: /var/lib/ngit-grasp-{name}) - `databaseBackend` - "lmdb" | "memory" (default: "lmdb") +See [Upgrade Git family storage](upgrade-git-family-storage.md) before updating +an existing instance to a release that enables identifier-family storage. + ### Identity - `relayName` - Relay name for NIP-11 (default: "{domain} grasp relay") - `relayDescription` - Relay description diff --git a/docs/how-to/upgrade-git-family-storage.md b/docs/how-to/upgrade-git-family-storage.md new file mode 100644 index 0000000..e41a23d --- /dev/null +++ b/docs/how-to/upgrade-git-family-storage.md @@ -0,0 +1,106 @@ +# Upgrade a server to identifier-family Git storage + +This procedure upgrades existing owner repositories and `/prs/` repositories +to ref-only views backed by one local object family per repository identifier. +It requires no new configuration and keeps all durable Git storage local. + +The upgrader runs automatically on every relay launch after configuration has +been validated and before purgatory restoration, background sync, or HTTP +request handling. A failed migration stops startup; restarting resumes from its +fsynced journal. + +After runtime database initialization, a non-blocking integrity pass checks +the resulting identifier families and views. It attempts to fetch missing +objects from clone URLs in accepted repository announcements and emits an +`ERROR` log for any family that remains unhealthy. Network repair never holds +up the listening service or changes whether the structural migration commits. + +## Before deploying + +1. Stop writes to the relay and take a filesystem snapshot of both the Git and + relay data directories. An older binary cannot serve the new thin views, so + this external snapshot is the clean release-level rollback boundary. +2. Check free space on `NGIT_GIT_DATA_PATH`. During the first launch the server + retains every original repository and builds a verified family union. Plan + for at least the current owner and `/prs/` repository footprint again, plus + room for the largest identifier family. Deduplication reduces the final + family size but should not be assumed for preflight capacity planning. +3. Keep Git and relay data on durable storage. +4. Do not configure Git GC for this rollout. Unreachable objects and migration + backups intentionally preserve delete-state rollback material. + +## Perform the upgrade + +Deploy the new release with the existing configuration and start the relay. +For every identifier, launch migration inventories all objects from all +matching owner and `/prs/` paths, including unreachable objects, verifies the +union, then atomically replaces each repository with a thin view. Original +repositories move to: + +```text +/.grasp/migration/backups/ +``` + +Progress and restart state live under `.grasp/migration/journal/`. The global +`.grasp/storage-version` marker is written only after every eligible view has +been reconciled and verified. Do not edit these files while the service is +running. + +Confirm the service reaches its normal listening state and exercise a clone, +fetch, and push for both a normal repository and a `/prs/` route. Retain the +external snapshot and `.grasp/migration/backups/` for the rollback window. The +server never deletes those backups automatically. + +Watch for the terminal startup-pass summary: + +```text +Git identifier-family integrity startup pass completed +``` + +An `unresolved` or `failed` count above zero is accompanied by an `ERROR` log +for each affected identifier. This reports pre-existing missing data without +putting startup into a network-dependent restart loop. + +## Check or repair one identifier on demand + +Queue a read-only check for every object format and owner/`/prs/` view sharing +an identifier: + +```console +ngit-grasp integrity-check \ + --git-data-path /var/lib/ngit-grasp/git \ + --identifier example +``` + +Add `--repair` to repair alternate wiring and try accepted clone servers for +missing OIDs: + +```console +ngit-grasp integrity-check \ + --git-data-path /var/lib/ngit-grasp/git \ + --identifier example \ + --repair +``` + +The command queues a durable request for the running relay rather than opening +or mutating a family from a second process. The worker normally consumes it +within five seconds and writes the result to the service log. A request queued +while the relay is stopped is processed after its next startup integrity pass. +Repeated requests for the same identifier and mode safely coalesce. + +## Failure and rollback + +- A migration failure is fail-closed. Fix the reported filesystem or Git error + and restart; the journal resumes the safe transition. +- A post-migration integrity repair failure is fail-open because the damage + predates conversion or arose after it. Inspect the identifier's `ERROR` log, + repair or update its listed clone sources, and queue `integrity-check + --repair` again. +- To roll the software release back, stop the service and restore the complete + pre-upgrade Git and relay-data snapshot together. Do not point an older + binary at migrated thin views. + +There is deliberately no automatic cleanup step. Backup retirement and +unreachable-object pruning remain deferred until rollback and delete-state +retention have an explicit policy. S3 adoption is a separate optional rollout +and is not required for this local storage model. diff --git a/src/git/handlers.rs b/src/git/handlers.rs index 9b28d15..ee1b4b3 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,15 @@ 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); + } + + // 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 @@ -948,7 +991,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 +1009,7 @@ pub enum GitError { ProcessSpawnFailed(std::io::Error), IoError(std::io::Error), GitFailed(Option), + Storage(String), } impl std::fmt::Display for GitError { @@ -975,6 +1019,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/integrity.rs b/src/git/integrity.rs new file mode 100644 index 0000000..e70cdfa --- /dev/null +++ b/src/git/integrity.rs @@ -0,0 +1,989 @@ +//! Identifier-family Git integrity inspection and repair orchestration. +//! +//! The durable integrity boundary is one object-format/identifier family plus +//! every owner and GRASP-06 view backed by it. Legacy repository migration is +//! only one producer of these families; steady-state and operator-triggered +//! checks use the same report. + +use std::collections::BTreeSet; +use std::path::{Path, PathBuf}; +use std::process::Command; +use std::sync::Arc; +use std::time::Duration; + +use anyhow::{anyhow, Context, Result}; +use async_trait::async_trait; +use bitcoin_hashes::{sha256, Hash}; +use clap::Args; +use nostr_sdk::prelude::{FromBech32, PublicKey}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinHandle; +use tracing::{error, info, warn}; + +use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat}; +use super::validate_repository_identifier; + +const REQUEST_VERSION: u32 = 1; +const REQUEST_POLL_INTERVAL: Duration = Duration::from_secs(5); + +/// Arguments for queueing an identifier-family integrity check in the live +/// relay process. +#[derive(Debug, Args)] +pub struct IntegrityCheckArgs { + /// Repository identifier (`d` tag value). All object formats and views for + /// this identifier are checked together. + #[arg(long)] + pub identifier: String, + + /// Attempt repair using accepted repository clone URLs. + #[arg(long, default_value_t = false)] + pub repair: bool, + + /// Git data path containing the `.grasp` family storage directory. + #[arg(long, env = "NGIT_GIT_DATA_PATH", default_value = "./data/git")] + pub git_data_path: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +struct IntegrityRequest { + version: u32, + identifier: String, + repair: bool, +} + +/// Queue an integrity request for the running relay. Requests are durable and +/// idempotent; if the service is stopped, the next launch consumes them. +pub fn enqueue_manual_check(args: &IntegrityCheckArgs) -> Result { + if !validate_repository_identifier(&args.identifier) { + return Err(anyhow!( + "invalid repository identifier {:?}", + args.identifier + )); + } + let storage = LocalGitStorage::new(&args.git_data_path); + let directory = request_directory(&storage); + std::fs::create_dir_all(&directory) + .with_context(|| format!("create integrity request directory {}", directory.display()))?; + let request = IntegrityRequest { + version: REQUEST_VERSION, + identifier: args.identifier.clone(), + repair: args.repair, + }; + let digest = + sha256::Hash::hash(format!("{}\0{}", request.identifier, request.repair).as_bytes()); + let path = directory.join(format!("{digest}.json")); + crate::atomic_file::write(&path, &serde_json::to_vec_pretty(&request)?) + .with_context(|| format!("write integrity request {}", path.display()))?; + Ok(path) +} + +/// A view ref whose target is absent from the identifier family. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct MissingRefTarget { + pub view: PathBuf, + pub reference: String, + pub oid: String, +} + +/// Integrity state for one identifier family and all views that use it. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct FamilyIntegrityReport { + pub key: FamilyKey, + pub views: Vec, + pub refs_checked: usize, + pub missing_oids: BTreeSet, + pub missing_ref_targets: Vec, + pub invalid_alternates: Vec, + pub pack_errors: Vec, + pub fsck_diagnostics: Vec, +} + +impl FamilyIntegrityReport { + /// Whether the family and every one of its views passed inspection. + pub fn is_healthy(&self) -> bool { + self.missing_oids.is_empty() + && self.missing_ref_targets.is_empty() + && self.invalid_alternates.is_empty() + && self.pack_errors.is_empty() + && self.fsck_diagnostics.is_empty() + } +} + +/// Network operations needed by the family healer. +/// +/// Production delegates to the existing hardened purgatory fetch client; +/// tests can supply deterministic local object sources. +#[async_trait] +pub trait FamilyRepairSource: Send + Sync { + async fn accepted_clone_urls(&self, identifier: &str) -> Result>; + + async fn fetch_missing_oids( + &self, + target_view: &Path, + url: &str, + oids: &[String], + ) -> Result>; +} + +/// Result of an optional repair attempt followed by a fresh integrity check. +#[derive(Debug)] +pub struct FamilyRepairOutcome { + pub initial: FamilyIntegrityReport, + pub final_report: FamilyIntegrityReport, + pub sources_tried: Vec, + pub source_failures: Vec, +} + +impl FamilyRepairOutcome { + pub fn repaired(&self) -> bool { + !self.initial.is_healthy() && self.final_report.is_healthy() + } +} + +#[derive(Default)] +struct PassStats { + checked: usize, + healthy: usize, + repaired: usize, + unresolved: usize, + failed: usize, +} + +/// Start the non-blocking integrity worker. +/// +/// It checks and heals every installed family once after startup, then +/// consumes identifier-scoped requests written by [`enqueue_manual_check`]. +pub fn spawn_integrity_worker( + storage: LocalGitStorage, + source: Arc, +) -> JoinHandle<()> { + tokio::spawn(async move { + run_startup_pass(&storage, source.as_ref()).await; + let first = tokio::time::Instant::now() + REQUEST_POLL_INTERVAL; + let mut interval = tokio::time::interval_at(first, REQUEST_POLL_INTERVAL); + loop { + interval.tick().await; + process_manual_requests(&storage, source.as_ref()).await; + } + }) +} + +async fn run_startup_pass(storage: &LocalGitStorage, source: &S) { + let families = match discover_families(storage) { + Ok(families) => families, + Err(error) => { + error!(%error, "Git identifier-family integrity startup discovery failed"); + return; + } + }; + info!( + families = families.len(), + "Git identifier-family integrity startup pass started" + ); + let mut stats = PassStats::default(); + for key in families { + run_one_family(storage, &key, source, true, "startup", &mut stats).await; + tokio::task::yield_now().await; + } + info!( + checked = stats.checked, + healthy = stats.healthy, + repaired = stats.repaired, + unresolved = stats.unresolved, + failed = stats.failed, + "Git identifier-family integrity startup pass completed" + ); +} + +async fn process_manual_requests( + storage: &LocalGitStorage, + source: &S, +) { + let requests = match load_requests(storage) { + Ok(requests) => requests, + Err(error) => { + error!(%error, "Read manual Git integrity requests"); + return; + } + }; + for (path, request) in requests { + let families = match discover_families(storage) { + Ok(families) => families + .into_iter() + .filter(|key| key.identifier == request.identifier) + .collect::>(), + Err(error) => { + error!( + identifier = %request.identifier, + %error, + "Manual Git integrity family discovery failed" + ); + remove_consumed_request(&path); + continue; + } + }; + if families.is_empty() { + error!( + identifier = %request.identifier, + "Manual Git integrity request matched no identifier family" + ); + remove_consumed_request(&path); + continue; + } + let mut stats = PassStats::default(); + for key in families { + run_one_family(storage, &key, source, request.repair, "manual", &mut stats).await; + } + info!( + identifier = %request.identifier, + repair = request.repair, + checked = stats.checked, + healthy = stats.healthy, + repaired = stats.repaired, + unresolved = stats.unresolved, + failed = stats.failed, + "Manual Git identifier-family integrity request completed" + ); + remove_consumed_request(&path); + } +} + +async fn run_one_family( + storage: &LocalGitStorage, + key: &FamilyKey, + source: &S, + repair: bool, + trigger: &'static str, + stats: &mut PassStats, +) { + stats.checked += 1; + let outcome = match check_and_repair_family(storage, key, source, repair).await { + Ok(outcome) => outcome, + Err(error) => { + stats.failed += 1; + error!( + identifier = %key.identifier, + object_format = %key.object_format, + trigger, + repair, + %error, + "Git identifier-family integrity check failed" + ); + return; + } + }; + if outcome.final_report.is_healthy() { + if outcome.repaired() { + stats.repaired += 1; + info!( + identifier = %key.identifier, + object_format = %key.object_format, + trigger, + sources_tried = outcome.sources_tried.len(), + "Git identifier family repaired" + ); + } else { + stats.healthy += 1; + } + return; + } + + stats.unresolved += 1; + let missing_sample = outcome + .final_report + .missing_oids + .iter() + .take(16) + .cloned() + .collect::>() + .join(","); + let diagnostic = outcome + .final_report + .fsck_diagnostics + .first() + .map(String::as_str) + .unwrap_or(""); + let source_failure = outcome + .source_failures + .first() + .map(String::as_str) + .unwrap_or(""); + error!( + identifier = %key.identifier, + object_format = %key.object_format, + trigger, + repair, + views = outcome.final_report.views.len(), + missing_oids = outcome.final_report.missing_oids.len(), + missing_sample, + missing_ref_targets = outcome.final_report.missing_ref_targets.len(), + invalid_alternates = outcome.final_report.invalid_alternates.len(), + pack_errors = outcome.final_report.pack_errors.len(), + fsck_diagnostics = outcome.final_report.fsck_diagnostics.len(), + diagnostic, + sources_tried = outcome.sources_tried.len(), + source_failures = outcome.source_failures.len(), + source_failure, + "Git identifier family remains unhealthy after integrity check" + ); +} + +fn request_directory(storage: &LocalGitStorage) -> PathBuf { + storage.internal_path().join("integrity-requests") +} + +fn load_requests(storage: &LocalGitStorage) -> Result> { + let directory = request_directory(storage); + let entries = match std::fs::read_dir(&directory) { + Ok(entries) => entries, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(error) => return Err(error.into()), + }; + let mut paths = entries + .filter_map(|entry| entry.ok()) + .filter(|entry| entry.file_type().is_ok_and(|kind| kind.is_file())) + .map(|entry| entry.path()) + .filter(|path| path.extension().and_then(|value| value.to_str()) == Some("json")) + .collect::>(); + paths.sort(); + let mut requests = Vec::new(); + for path in paths { + let request: IntegrityRequest = serde_json::from_slice(&std::fs::read(&path)?) + .with_context(|| format!("parse integrity request {}", path.display()))?; + if request.version != REQUEST_VERSION + || !validate_repository_identifier(&request.identifier) + { + return Err(anyhow!("invalid integrity request {}", path.display())); + } + requests.push((path, request)); + } + Ok(requests) +} + +fn remove_consumed_request(path: &Path) { + if let Err(error) = std::fs::remove_file(path) { + warn!(request = %path.display(), %error, "Remove consumed Git integrity request"); + } +} + +/// Inspect a family, repair local view wiring, fetch missing OIDs from every +/// accepted clone source as needed, then inspect it again. +pub async fn check_and_repair_family( + storage: &LocalGitStorage, + key: &FamilyKey, + source: &S, + repair: bool, +) -> Result { + let lease = storage.write_lease(key).await?; + let initial = inspect_family(storage, key)?; + if initial.is_healthy() || !repair { + return Ok(FamilyRepairOutcome { + final_report: initial.clone(), + initial, + sources_tried: Vec::new(), + source_failures: Vec::new(), + }); + } + + for view in &initial.invalid_alternates { + storage.configure_thin_view(key, view)?; + } + drop(lease); + + let mut current = inspect_family(storage, key)?; + let mut sources_tried = Vec::new(); + let mut source_failures = Vec::new(); + if !current.missing_oids.is_empty() { + let mut urls = source.accepted_clone_urls(&key.identifier).await?; + urls.sort(); + urls.dedup(); + let target = current + .views + .first() + .cloned() + .unwrap_or_else(|| storage.family_repo_path(key)); + for url in urls { + let missing: Vec<_> = current.missing_oids.iter().cloned().collect(); + if missing.is_empty() { + break; + } + sources_tried.push(url.clone()); + if let Err(error) = source.fetch_missing_oids(&target, &url, &missing).await { + source_failures.push(format!("{url}: {error}")); + } + current = inspect_family(storage, key)?; + } + } + + // A view ref that was deliberately preserved while its target was absent + // becomes a useful family root as soon as repair restores that object. + let _lease = storage.write_lease(key).await?; + for target in &initial.missing_ref_targets { + if oid_exists(&storage.family_repo_path(key), &target.oid)? { + let source_ref = format!("{}:{}", target.view.display(), target.reference); + storage.retain_tip(key, &source_ref, &target.oid)?; + storage.advertise_base_tip(key, &source_ref, &target.oid)?; + } + } + let final_report = inspect_family(storage, key)?; + Ok(FamilyRepairOutcome { + initial, + final_report, + sources_tried, + source_failures, + }) +} + +#[async_trait] +impl FamilyRepairSource for crate::purgatory::sync::RealSyncContext { + async fn accepted_clone_urls(&self, identifier: &str) -> Result> { + self.accepted_repository_clone_urls(identifier).await + } + + async fn fetch_missing_oids( + &self, + target_view: &Path, + url: &str, + oids: &[String], + ) -> Result> { + crate::purgatory::sync::SyncContext::fetch_oids_with_role( + self, + target_view, + url, + oids, + crate::purgatory::sync::GitFetchRole::Integrity, + ) + .await + } +} + +/// Discover every installed identifier family in deterministic order. +pub fn discover_families(storage: &LocalGitStorage) -> Result> { + let mut families = Vec::new(); + for object_format in [ObjectFormat::Sha1, ObjectFormat::Sha256] { + let root = storage + .internal_path() + .join("families") + .join(object_format.as_str()); + let entries = match std::fs::read_dir(&root) { + Ok(entries) => entries, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => continue, + Err(error) => return Err(error).with_context(|| format!("read {}", root.display())), + }; + for entry in entries { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + let name = entry.file_name(); + let Some(identifier) = name.to_str().and_then(|name| name.strip_suffix(".git")) else { + continue; + }; + if !validate_repository_identifier(identifier) { + return Err(anyhow!( + "invalid identifier-family path {}", + entry.path().display() + )); + } + families.push(FamilyKey::new(object_format, identifier)?); + } + } + families.sort_by(|left, right| { + (left.object_format.as_str(), left.identifier.as_str()) + .cmp(&(right.object_format.as_str(), right.identifier.as_str())) + }); + Ok(families) +} + +/// Inspect one family and every owner or `/prs/` view backed by its identifier. +pub fn inspect_family(storage: &LocalGitStorage, key: &FamilyKey) -> Result { + let family = storage.family_repo_path(key); + if !family.is_dir() { + return Err(anyhow!("family is missing: {}", family.display())); + } + let actual_format = ObjectFormat::detect(&family)?; + if actual_format != key.object_format { + return Err(anyhow!( + "family {} has object format {}, expected {}", + family.display(), + actual_format, + key.object_format + )); + } + + let views = discover_views(storage, key)?; + let mut missing_oids = BTreeSet::new(); + let mut missing_ref_targets = Vec::new(); + let mut invalid_alternates = Vec::new(); + let mut refs_checked = 0; + + for view in &views { + if !storage.is_thin_view(key, view) { + invalid_alternates.push(view.clone()); + } + for (reference, oid) in list_refs(view)? { + refs_checked += 1; + if !oid_exists(&family, &oid)? { + missing_oids.insert(oid.clone()); + missing_ref_targets.push(MissingRefTarget { + view: view.clone(), + reference, + oid, + }); + } + } + } + + let pack_errors = inspect_packs(&family)?; + let (fsck_missing, fsck_diagnostics) = inspect_connectivity(&family, key.object_format)?; + missing_oids.extend(fsck_missing); + + Ok(FamilyIntegrityReport { + key: key.clone(), + views, + refs_checked, + missing_oids, + missing_ref_targets, + invalid_alternates, + pack_errors, + fsck_diagnostics, + }) +} + +fn discover_views(storage: &LocalGitStorage, key: &FamilyKey) -> Result> { + let mut views = Vec::new(); + let entries = match std::fs::read_dir(storage.git_data_path()) { + Ok(entries) => entries, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(views), + Err(error) => return Err(error.into()), + }; + for entry in entries { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + let name = entry.file_name().to_string_lossy().into_owned(); + if name == "prs" { + discover_pr_views(&entry.path(), key, &mut views)?; + } else if name.starts_with("npub1") && PublicKey::from_bech32(&name).is_ok() { + maybe_add_view(&entry.path(), key, &mut views)?; + } + } + views.sort(); + Ok(views) +} + +fn discover_pr_views(root: &Path, key: &FamilyKey, views: &mut Vec) -> Result<()> { + for entry in std::fs::read_dir(root)? { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + let name = entry.file_name().to_string_lossy().into_owned(); + if name.len() == 64 && name.chars().all(|character| character.is_ascii_hexdigit()) { + maybe_add_view(&entry.path(), key, views)?; + } + } + Ok(()) +} + +fn maybe_add_view(directory: &Path, key: &FamilyKey, views: &mut Vec) -> Result<()> { + let candidate = directory.join(format!("{}.git", key.identifier)); + if !candidate.is_dir() { + return Ok(()); + } + if ObjectFormat::detect(&candidate)? == key.object_format { + views.push(candidate); + } + Ok(()) +} + +fn list_refs(repo: &Path) -> Result> { + let output = Command::new("git") + .args(["for-each-ref", "--format=%(refname)%00%(objectname)"]) + .current_dir(repo) + .output() + .with_context(|| format!("list refs in {}", repo.display()))?; + if !output.status.success() { + return Err(anyhow!( + "list refs in {}: {}", + repo.display(), + String::from_utf8_lossy(&output.stderr).trim() + )); + } + let mut refs = Vec::new(); + for line in output.stdout.split(|byte| *byte == b'\n') { + if line.is_empty() { + continue; + } + let Some(separator) = line.iter().position(|byte| *byte == 0) else { + return Err(anyhow!("unexpected git for-each-ref output")); + }; + refs.push(( + String::from_utf8(line[..separator].to_vec())?, + String::from_utf8(line[separator + 1..].to_vec())?, + )); + } + Ok(refs) +} + +fn inspect_packs(family: &Path) -> Result> { + let pack_dir = family.join("objects/pack"); + let entries = match std::fs::read_dir(&pack_dir) { + Ok(entries) => entries, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(error) => return Err(error.into()), + }; + let mut errors = Vec::new(); + let mut indexes = Vec::new(); + for entry in entries { + let entry = entry?; + if !entry.file_type()?.is_file() { + continue; + } + let path = entry.path(); + match path.extension().and_then(|extension| extension.to_str()) { + Some("pack") if !path.with_extension("idx").is_file() => { + errors.push(format!("pack has no index: {}", path.display())); + } + Some("idx") if !path.with_extension("pack").is_file() => { + errors.push(format!("index has no pack: {}", path.display())); + } + Some("idx") => indexes.push(path), + _ => {} + } + } + indexes.sort(); + for index in indexes { + let status = Command::new("git") + .arg("verify-pack") + .arg(&index) + .stdout(std::process::Stdio::null()) + .stderr(std::process::Stdio::null()) + .status() + .with_context(|| format!("verify pack index {}", index.display()))?; + if !status.success() { + errors.push(format!("pack verification failed: {}", index.display())); + } + } + errors.sort(); + Ok(errors) +} + +fn inspect_connectivity( + family: &Path, + object_format: ObjectFormat, +) -> Result<(BTreeSet, Vec)> { + let output = Command::new("git") + .args(["fsck", "--full", "--no-reflogs", "--no-dangling"]) + .env("LC_ALL", "C") + .current_dir(family) + .output() + .with_context(|| format!("inspect family connectivity in {}", family.display()))?; + if output.status.success() { + return Ok((BTreeSet::new(), Vec::new())); + } + + let diagnostic = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let oid_len = match object_format { + ObjectFormat::Sha1 => 40, + ObjectFormat::Sha256 => 64, + }; + let mut missing = BTreeSet::new(); + let mut diagnostics = Vec::new(); + for line in diagnostic + .lines() + .map(str::trim) + .filter(|line| !line.is_empty()) + { + let fields: Vec<_> = line.split_whitespace().collect(); + if fields.first() == Some(&"missing") { + if let Some(oid) = fields.iter().rev().find(|field| { + field.len() == oid_len && field.chars().all(|ch| ch.is_ascii_hexdigit()) + }) { + missing.insert((*oid).to_owned()); + } + } + if diagnostics.len() < 64 { + diagnostics.push(line.chars().take(2048).collect()); + } + } + if diagnostics.is_empty() { + diagnostics.push(format!("git fsck exited with {}", output.status)); + } + Ok((missing, diagnostics)) +} + +fn oid_exists(repo: &Path, oid: &str) -> Result { + let output = Command::new("git") + .args(["cat-file", "-e", oid]) + .current_dir(repo) + .output() + .with_context(|| format!("check object {oid} in {}", repo.display()))?; + Ok(output.status.success()) +} + +#[cfg(test)] +mod tests { + use std::fs; + use std::io::Write; + use std::process::Stdio; + + use nostr_sdk::prelude::{Keys, ToBech32}; + + use super::*; + + struct LocalRepairSource { + url: String, + family_objects: PathBuf, + } + + #[async_trait] + impl FamilyRepairSource for LocalRepairSource { + async fn accepted_clone_urls(&self, _identifier: &str) -> Result> { + Ok(vec![self.url.clone()]) + } + + async fn fetch_missing_oids( + &self, + target_view: &Path, + url: &str, + oids: &[String], + ) -> Result> { + let mut command = Command::new("git"); + command + .arg("fetch") + .arg("--no-write-fetch-head") + .arg(url) + .args(oids) + .env("GIT_OBJECT_DIRECTORY", &self.family_objects) + .current_dir(target_view); + let output = command.output()?; + if !output.status.success() { + return Err(anyhow!(String::from_utf8_lossy(&output.stderr).into_owned())); + } + Ok(oids.to_vec()) + } + } + + fn git(repo: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .args(args) + .current_dir(repo) + .env("GIT_CONFIG_NOSYSTEM", "1") + .env("HOME", "/nonexistent") + .output() + .unwrap(); + assert!( + output.status.success(), + "git {args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_owned() + } + + fn write_commit(repo: &Path, parent: Option<&str>) -> String { + let tree = git(repo, &["mktree"]); + let mut contents = format!( + "tree {tree}\nauthor Test 0 +0000\ncommitter Test 0 +0000\n\nintegrity fixture\n" + ); + if let Some(parent) = parent { + contents.insert_str( + contents.find("author ").unwrap(), + &format!("parent {parent}\n"), + ); + } + let mut child = Command::new("git") + .args(["hash-object", "-t", "commit", "-w", "--stdin"]) + .current_dir(repo) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child + .stdin + .take() + .unwrap() + .write_all(contents.as_bytes()) + .unwrap(); + let output = child.wait_with_output().unwrap(); + assert!(output.status.success()); + String::from_utf8_lossy(&output.stdout).trim().to_owned() + } + + fn fixture() -> ( + tempfile::TempDir, + LocalGitStorage, + FamilyKey, + PathBuf, + String, + ) { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let key = FamilyKey::sha1("shared").unwrap(); + let family = storage.ensure_family(&key).unwrap(); + let commit = write_commit(&family, None); + storage.retain_tip(&key, "fixture", &commit).unwrap(); + storage + .advertise_base_tip(&key, "fixture", &commit) + .unwrap(); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let view = storage.git_data_path().join(owner).join("shared.git"); + storage.create_thin_view(&key, &view).unwrap(); + git(&view, &["update-ref", "refs/heads/main", &commit]); + (temp, storage, key, view, commit) + } + + #[test] + fn healthy_family_checks_family_and_views() { + let (_temp, storage, key, _view, _commit) = fixture(); + + let report = inspect_family(&storage, &key).unwrap(); + + assert!(report.is_healthy(), "{report:#?}"); + assert_eq!(report.views.len(), 1); + assert_eq!(report.refs_checked, 1); + assert_eq!(discover_families(&storage).unwrap(), vec![key]); + } + + #[test] + fn reports_missing_graph_objects_and_view_ref_targets() { + let (_temp, storage, key, view, commit) = fixture(); + let family = storage.family_repo_path(&key); + let child = write_commit(&family, Some(&commit)); + git(&view, &["update-ref", "refs/heads/main", &child]); + storage.retain_tip(&key, "child", &child).unwrap(); + let object = family.join("objects").join(&commit[..2]).join(&commit[2..]); + fs::remove_file(object).unwrap(); + let broken = "2".repeat(40); + fs::create_dir_all(view.join("refs/heads")).unwrap(); + fs::write(view.join("refs/heads/broken"), format!("{broken}\n")).unwrap(); + + let report = inspect_family(&storage, &key).unwrap(); + + assert!(!report.is_healthy()); + assert!(report.missing_oids.contains(&commit)); + assert!(report.missing_oids.contains(&broken)); + assert!(report + .missing_ref_targets + .iter() + .any(|target| target.reference == "refs/heads/broken" && target.oid == broken)); + assert!(!report.fsck_diagnostics.is_empty()); + } + + #[test] + fn reports_a_view_whose_family_alternate_changed() { + let (_temp, storage, key, view, _commit) = fixture(); + fs::write(view.join("objects/info/alternates"), "/wrong/objects\n").unwrap(); + + let report = inspect_family(&storage, &key).unwrap(); + + assert_eq!(report.invalid_alternates, vec![view]); + assert!(!report.is_healthy()); + } + + #[test] + fn reports_an_unindexed_family_pack() { + let (_temp, storage, key, _view, _commit) = fixture(); + let pack = storage + .family_repo_path(&key) + .join("objects/pack/pack-0000000000000000000000000000000000000000.pack"); + fs::write(pack, b"incomplete").unwrap(); + + let report = inspect_family(&storage, &key).unwrap(); + + assert_eq!(report.pack_errors.len(), 1); + assert!(report.pack_errors[0].contains("pack has no index")); + } + + #[test] + fn manual_requests_are_durable_and_identifier_scoped() { + let temp = tempfile::tempdir().unwrap(); + let args = IntegrityCheckArgs { + identifier: "shared".to_owned(), + repair: true, + git_data_path: temp.path().join("git").to_string_lossy().into_owned(), + }; + + let path = enqueue_manual_check(&args).unwrap(); + let storage = LocalGitStorage::new(&args.git_data_path); + let requests = load_requests(&storage).unwrap(); + + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].0, path); + assert_eq!(requests[0].1.identifier, "shared"); + assert!(requests[0].1.repair); + } + + #[test] + fn manual_requests_reject_unsafe_identifiers() { + let temp = tempfile::tempdir().unwrap(); + let args = IntegrityCheckArgs { + identifier: "../escape".to_owned(), + repair: false, + git_data_path: temp.path().join("git").to_string_lossy().into_owned(), + }; + + assert!(enqueue_manual_check(&args).is_err()); + } + + #[tokio::test] + async fn manual_requests_are_consumed_by_the_family_worker() { + let (temp, storage, _key, _view, _commit) = fixture(); + let args = IntegrityCheckArgs { + identifier: "shared".to_owned(), + repair: false, + git_data_path: storage.git_data_path().to_string_lossy().into_owned(), + }; + let request = enqueue_manual_check(&args).unwrap(); + let source = LocalRepairSource { + url: String::new(), + family_objects: temp.path().join("unused"), + }; + + process_manual_requests(&storage, &source).await; + + assert!(!request.exists()); + } + + #[tokio::test] + async fn repairs_missing_objects_from_an_identifier_clone_source() { + let (temp, storage, key, view, commit) = fixture(); + let family = storage.family_repo_path(&key); + let source_repo = temp.path().join("source.git"); + git( + temp.path(), + &[ + "clone", + "--bare", + family.to_str().unwrap(), + source_repo.to_str().unwrap(), + ], + ); + git( + &source_repo, + &["update-ref", "refs/heads/recovery", &commit], + ); + let object = family.join("objects").join(&commit[..2]).join(&commit[2..]); + fs::remove_file(object).unwrap(); + assert!(!inspect_family(&storage, &key).unwrap().is_healthy()); + let source = LocalRepairSource { + url: source_repo.to_string_lossy().into_owned(), + family_objects: storage.family_objects_path(&key), + }; + + let outcome = check_and_repair_family(&storage, &key, &source, true) + .await + .unwrap(); + + assert!(outcome.repaired(), "{outcome:#?}"); + assert_eq!(outcome.sources_tried, vec![source.url]); + assert!(outcome.source_failures.is_empty()); + assert!(oid_exists(&family, &commit).unwrap()); + assert_eq!(git(&view, &["rev-parse", "refs/heads/main"]), commit); + } +} diff --git a/src/git/migration.rs b/src/git/migration.rs new file mode 100644 index 0000000..684f3c5 --- /dev/null +++ b/src/git/migration.rs @@ -0,0 +1,1336 @@ +//! Crash-safe launch migration from full bare repositories to thin family views. +//! +//! The migration runs before any relay task can access Git storage. Original +//! repositories are moved under `.grasp/migration/backups` and are never +//! deleted automatically. Every object reported by `cat-file +//! --batch-all-objects`, including unreachable objects, is copied into a +//! verified union pack before a view can be replaced. + +use std::collections::{BTreeMap, HashMap, HashSet}; +use std::io::Write; +use std::path::{Path, PathBuf}; +use std::process::{Command, Stdio}; + +use anyhow::{anyhow, Context, Result}; +use bitcoin_hashes::{sha256, Hash}; +use nostr_sdk::prelude::{FromBech32, PublicKey}; +use serde::{Deserialize, Serialize}; +use tracing::warn; + +use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat, STORAGE_VERSION}; +use super::validate_repository_identifier; + +const MIGRATION_VERSION: u32 = 1; + +#[derive(Debug, Default, Clone, PartialEq, Eq)] +pub struct MigrationReport { + pub migrated_views: usize, + pub recovered_views: usize, + pub families_built: usize, + pub already_current: bool, +} + +#[derive(Debug, Clone)] +struct Candidate { + relative_path: PathBuf, + view_path: PathBuf, + source_path: PathBuf, + backup_path: PathBuf, + journal_path: PathBuf, + key: FamilyKey, + recovered: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +struct RefSnapshot { + oid: String, + symbolic_target: Option, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +enum JournalState { + Prepared, + BackedUp, + ViewInstalled, + Complete, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct ViewJournal { + version: u32, + relative_path: PathBuf, + object_format: String, + identifier: String, + refs: BTreeMap, + head: Vec, + state: JournalState, +} + +/// Upgrade all eligible Git repositories under one storage root. +/// +/// This function is safe to call on every launch. The storage-version marker +/// advances only after every legacy path and incomplete journal is reconciled. +pub async fn migrate_on_startup(storage: &LocalGitStorage) -> Result { + std::fs::create_dir_all(storage.git_data_path()).with_context(|| { + format!( + "create Git data directory {}", + storage.git_data_path().display() + ) + })?; + let storage_is_current = validate_storage_version(storage)?; + let mut candidates = scan_active_legacy_views(storage)?; + if storage_is_current && candidates.is_empty() { + return Ok(MigrationReport { + already_current: true, + ..MigrationReport::default() + }); + } + + let journals = load_journals(storage)?; + reconcile_journals(storage, journals, &mut candidates)?; + candidates.sort_by(|left, right| left.relative_path.cmp(&right.relative_path)); + candidates.dedup_by(|left, right| left.relative_path == right.relative_path); + + let mut report = MigrationReport::default(); + let mut by_family: HashMap> = HashMap::new(); + for candidate in candidates { + by_family + .entry(candidate.key.clone()) + .or_default() + .push(candidate); + } + let mut families: Vec<_> = by_family.into_iter().collect(); + families.sort_by(|(left, _), (right, _)| { + left.object_format + .as_str() + .cmp(right.object_format.as_str()) + .then(left.identifier.cmp(&right.identifier)) + }); + + for (key, candidates) in families { + build_family_union(storage, &key, &candidates)?; + report.families_built += 1; + for candidate in candidates { + install_view(storage, &candidate)?; + if candidate.recovered { + report.recovered_views += 1; + } else { + report.migrated_views += 1; + } + } + } + + verify_all_views(storage)?; + write_atomic( + &storage.storage_version_path(), + format!("{STORAGE_VERSION}\n").as_bytes(), + )?; + report.already_current = report.migrated_views == 0 && report.recovered_views == 0; + Ok(report) +} + +/// Enumerate installed family repositories for S3 adoption after migration. +pub fn family_keys(storage: &LocalGitStorage) -> Result> { + let mut keys = Vec::new(); + let families = storage.internal_path().join("families"); + for format in [ObjectFormat::Sha1, ObjectFormat::Sha256] { + let path = families.join(format.as_str()); + let Ok(entries) = std::fs::read_dir(path) else { + continue; + }; + for entry in entries { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + let Some(identifier) = entry + .file_name() + .to_str() + .and_then(|name| name.strip_suffix(".git")) + .map(str::to_owned) + else { + continue; + }; + keys.push(FamilyKey::new(format, identifier)?); + } + } + keys.sort_by(|left, right| { + left.object_format + .as_str() + .cmp(right.object_format.as_str()) + .then(left.identifier.cmp(&right.identifier)) + }); + Ok(keys) +} + +/// Import every object and ref tip from an archived repository into a family. +/// +/// The caller must hold the family's write lease. This deliberately includes +/// unreachable objects so restoring a pre-family archive does not narrow the +/// rollback history retained by the original repository. +pub(crate) fn import_archived_repository( + storage: &LocalGitStorage, + key: &FamilyKey, + source: &Path, +) -> Result<()> { + let source_format = ObjectFormat::detect(source)?; + if source_format != key.object_format { + return Err(anyhow!( + "archived repository {} uses {}, expected {}", + source.display(), + source_format, + key.object_format + )); + } + let family = storage.ensure_family(key)?; + let expected_oids = list_all_objects(source)?; + copy_object_database(source, &family, key.object_format)?; + ensure_oids(&family, &expected_oids)?; + for (name, reference) in snapshot_refs(source)? { + retain_family_tip(storage, key, &family, &name, &reference.oid)?; + } + Ok(()) +} + +fn validate_storage_version(storage: &LocalGitStorage) -> Result { + let path = storage.storage_version_path(); + let Ok(contents) = std::fs::read_to_string(&path) else { + return Ok(false); + }; + let version: u32 = contents + .trim() + .parse() + .with_context(|| format!("parse storage version in {}", path.display()))?; + if version > STORAGE_VERSION { + return Err(anyhow!( + "Git storage version {version} is newer than supported version {STORAGE_VERSION}" + )); + } + Ok(version == STORAGE_VERSION) +} + +fn scan_active_legacy_views(storage: &LocalGitStorage) -> Result> { + let mut candidates = Vec::new(); + let root = storage.git_data_path(); + for owner_entry in std::fs::read_dir(root)? { + let owner_entry = owner_entry?; + if !owner_entry.file_type()?.is_dir() { + continue; + } + let owner = owner_entry.file_name().to_string_lossy().into_owned(); + if owner == "prs" { + scan_prs_tree(storage, &owner_entry.path(), &mut candidates)?; + continue; + } + if !owner.starts_with("npub1") || PublicKey::from_bech32(&owner).is_err() { + continue; + } + scan_repo_directory(storage, &owner_entry.path(), &mut candidates)?; + } + Ok(candidates) +} + +fn scan_prs_tree( + storage: &LocalGitStorage, + prs_path: &Path, + candidates: &mut Vec, +) -> Result<()> { + for submitter in std::fs::read_dir(prs_path)? { + let submitter = submitter?; + if !submitter.file_type()?.is_dir() { + continue; + } + let name = submitter.file_name().to_string_lossy().into_owned(); + if name.len() != 64 || !name.chars().all(|ch| ch.is_ascii_hexdigit()) { + continue; + } + scan_repo_directory(storage, &submitter.path(), candidates)?; + } + Ok(()) +} + +fn scan_repo_directory( + storage: &LocalGitStorage, + directory: &Path, + candidates: &mut Vec, +) -> Result<()> { + for entry in std::fs::read_dir(directory)? { + let entry = entry?; + if !entry.file_type()?.is_dir() { + continue; + } + let file_name = entry.file_name(); + let Some(identifier) = file_name + .to_str() + .and_then(|name| name.strip_suffix(".git")) + else { + continue; + }; + if !validate_repository_identifier(identifier) { + return Err(anyhow!( + "unsafe repository identifier in migration path {}", + entry.path().display() + )); + } + let object_format = ObjectFormat::detect(&entry.path())?; + let key = FamilyKey::new(object_format, identifier)?; + if storage.is_thin_view(&key, &entry.path()) { + continue; + } + if entry.path().join("objects/info/alternates").exists() { + return Err(anyhow!( + "legacy repository {} uses an unknown alternate; refusing lossy migration", + entry.path().display() + )); + } + candidates.push(candidate_for_active(storage, &entry.path(), key)?); + } + Ok(()) +} + +fn candidate_for_active( + storage: &LocalGitStorage, + view_path: &Path, + key: FamilyKey, +) -> Result { + let relative_path = view_path + .strip_prefix(storage.git_data_path()) + .context("repository path escaped Git data root")? + .to_path_buf(); + let backup_path = backup_root(storage).join(&relative_path); + Ok(Candidate { + journal_path: journal_path(storage, &relative_path), + relative_path, + view_path: view_path.to_path_buf(), + source_path: view_path.to_path_buf(), + backup_path, + key, + recovered: false, + }) +} + +fn reconcile_journals( + storage: &LocalGitStorage, + journals: Vec<(PathBuf, ViewJournal)>, + candidates: &mut Vec, +) -> Result<()> { + for (journal_path, mut journal) in journals { + if journal.version != MIGRATION_VERSION { + return Err(anyhow!( + "unsupported migration journal version {} in {}", + journal.version, + journal_path.display() + )); + } + let path_identifier = validate_relative_view_path(&journal.relative_path)?; + if path_identifier != journal.identifier { + return Err(anyhow!( + "migration journal identifier {:?} does not match path {}", + journal.identifier, + journal.relative_path.display() + )); + } + let format = parse_object_format(&journal.object_format)?; + let key = FamilyKey::new(format, journal.identifier.clone())?; + let view_path = storage.git_data_path().join(&journal.relative_path); + let backup_path = backup_root(storage).join(&journal.relative_path); + + if view_path.exists() && storage.is_thin_view(&key, &view_path) { + verify_view(&view_path, &journal.refs, &journal.head)?; + if !backup_path.is_dir() { + return Err(anyhow!( + "migrated view {} has no rollback backup", + view_path.display() + )); + } + if journal.state != JournalState::Complete { + journal.state = JournalState::Complete; + write_journal(&journal_path, &journal)?; + } + candidates.retain(|candidate| candidate.relative_path != journal.relative_path); + continue; + } + + let source_path = if backup_path.is_dir() { + if view_path.exists() { + return Err(anyhow!( + "both legacy view and migration backup exist for {}", + journal.relative_path.display() + )); + } + backup_path.clone() + } else if view_path.is_dir() { + view_path.clone() + } else { + return Err(anyhow!( + "migration journal {} has neither source nor backup repository", + journal_path.display() + )); + }; + candidates.retain(|candidate| candidate.relative_path != journal.relative_path); + candidates.push(Candidate { + relative_path: journal.relative_path, + view_path, + source_path, + backup_path, + journal_path, + key, + recovered: true, + }); + } + Ok(()) +} + +fn build_family_union( + storage: &LocalGitStorage, + key: &FamilyKey, + candidates: &[Candidate], +) -> Result<()> { + let family = storage.ensure_family(key)?; + let staging = migration_root(storage) + .join("family-staging") + .join(key.object_format.as_str()) + .join(format!("{}.git", key.identifier)); + remove_migration_staging(&staging, &migration_root(storage).join("family-staging"))?; + if let Some(parent) = staging.parent() { + std::fs::create_dir_all(parent)?; + } + run_git( + None, + &[ + "init", + "--bare", + "--quiet", + &format!("--object-format={}", key.object_format), + staging.to_string_lossy().as_ref(), + ], + "initialize migration family staging repository", + )?; + + let mut sources = vec![family.clone()]; + sources.extend( + candidates + .iter() + .map(|candidate| candidate.source_path.clone()), + ); + let mut expected_oids = Vec::new(); + for source in &sources { + let format = ObjectFormat::detect(source)?; + if format != key.object_format { + return Err(anyhow!( + "object format mismatch while merging {}", + source.display() + )); + } + expected_oids.extend(list_all_objects(source)?); + copy_object_database(source, &staging, key.object_format)?; + } + expected_oids.sort(); + expected_oids.dedup(); + ensure_oids(&staging, &expected_oids)?; + + if !expected_oids.is_empty() { + let union_dir = migration_root(storage) + .join("union-packs") + .join(key.object_format.as_str()) + .join(&key.identifier); + remove_migration_staging(&union_dir, &migration_root(storage).join("union-packs"))?; + std::fs::create_dir_all(&union_dir)?; + let base = union_dir.join("pack"); + let pack_hash = pack_explicit_objects(&staging, &base, &expected_oids)?; + let pack = union_dir.join(format!("pack-{pack_hash}.pack")); + let index = union_dir.join(format!("pack-{pack_hash}.idx")); + verify_pack(&index)?; + let family_pack_dir = family.join("objects/pack"); + std::fs::create_dir_all(&family_pack_dir)?; + install_if_absent(&pack, &family_pack_dir.join(pack.file_name().unwrap()))?; + install_if_absent(&index, &family_pack_dir.join(index.file_name().unwrap()))?; + fsync_directory(&family_pack_dir)?; + } + + ensure_oids(&family, &expected_oids)?; + for candidate in candidates { + for (name, reference) in snapshot_refs(&candidate.source_path)? { + retain_family_tip(storage, key, &family, &name, &reference.oid)?; + } + } + audit_family_connectivity(&family)?; + remove_migration_staging(&staging, &migration_root(storage).join("family-staging"))?; + Ok(()) +} + +fn install_view(storage: &LocalGitStorage, candidate: &Candidate) -> Result<()> { + let refs = snapshot_refs(&candidate.source_path)?; + let head = std::fs::read(candidate.source_path.join("HEAD")) + .with_context(|| format!("read HEAD from {}", candidate.source_path.display()))?; + + let mut journal = ViewJournal { + version: MIGRATION_VERSION, + relative_path: candidate.relative_path.clone(), + object_format: candidate.key.object_format.as_str().to_owned(), + identifier: candidate.key.identifier.clone(), + refs: refs.clone(), + head: head.clone(), + state: JournalState::Prepared, + }; + write_journal(&candidate.journal_path, &journal)?; + + let staged_view = migration_root(storage) + .join("view-staging") + .join(journal_name(&candidate.relative_path)); + remove_migration_staging(&staged_view, &migration_root(storage).join("view-staging"))?; + storage.create_thin_view(&candidate.key, &staged_view)?; + preserve_view_configuration( + storage, + &candidate.key, + &candidate.source_path, + &staged_view, + )?; + install_refs(&staged_view, &refs)?; + write_atomic(&staged_view.join("HEAD"), &head)?; + verify_view(&staged_view, &refs, &head)?; + + if !candidate.backup_path.exists() { + if candidate.view_path != candidate.source_path || !candidate.view_path.is_dir() { + return Err(anyhow!( + "migration source/view mismatch for {}", + candidate.relative_path.display() + )); + } + if let Some(parent) = candidate.backup_path.parent() { + std::fs::create_dir_all(parent)?; + } + std::fs::rename(&candidate.view_path, &candidate.backup_path).with_context(|| { + format!( + "move legacy repository {} to rollback backup {}", + candidate.view_path.display(), + candidate.backup_path.display() + ) + })?; + fsync_directory(candidate.view_path.parent().context("view has no parent")?)?; + fsync_directory( + candidate + .backup_path + .parent() + .context("backup has no parent")?, + )?; + } + journal.state = JournalState::BackedUp; + write_journal(&candidate.journal_path, &journal)?; + + if !candidate.view_path.exists() { + if let Some(parent) = candidate.view_path.parent() { + std::fs::create_dir_all(parent)?; + } + std::fs::rename(&staged_view, &candidate.view_path) + .with_context(|| format!("install migrated view {}", candidate.view_path.display()))?; + fsync_directory(candidate.view_path.parent().context("view has no parent")?)?; + } + journal.state = JournalState::ViewInstalled; + write_journal(&candidate.journal_path, &journal)?; + verify_view(&candidate.view_path, &refs, &head)?; + journal.state = JournalState::Complete; + write_journal(&candidate.journal_path, &journal)?; + Ok(()) +} + +fn preserve_view_configuration( + storage: &LocalGitStorage, + key: &FamilyKey, + source: &Path, + target: &Path, +) -> Result<()> { + let source_config = source.join("config"); + if source_config.is_file() { + std::fs::copy(&source_config, target.join("config"))?; + } + storage.configure_thin_view(key, target)?; + let source_description = source.join("description"); + if source_description.is_file() { + std::fs::copy(source_description, target.join("description"))?; + } + copy_directory_contents(&source.join("hooks"), &target.join("hooks"))?; + Ok(()) +} + +fn verify_all_views(storage: &LocalGitStorage) -> Result<()> { + if !scan_active_legacy_views(storage)?.is_empty() { + return Err(anyhow!( + "legacy Git repositories remain after storage migration" + )); + } + for (_, journal) in load_journals(storage)? { + if journal.state != JournalState::Complete { + return Err(anyhow!( + "incomplete migration journal remains for {}", + journal.relative_path.display() + )); + } + } + Ok(()) +} + +fn snapshot_refs(repo: &Path) -> Result> { + let output = Command::new("git") + .args([ + "for-each-ref", + "--format=%(refname)%00%(objectname)%00%(symref)", + ]) + .current_dir(repo) + .output()?; + if !output.status.success() { + return Err(anyhow!( + "list refs in {}: {}", + repo.display(), + String::from_utf8_lossy(&output.stderr).trim() + )); + } + let mut refs = BTreeMap::new(); + for line in output.stdout.split(|byte| *byte == b'\n') { + if line.is_empty() { + continue; + } + let fields: Vec<_> = line.split(|byte| *byte == 0).collect(); + if fields.len() != 3 { + return Err(anyhow!("unexpected git for-each-ref output")); + } + let name = String::from_utf8(fields[0].to_vec())?; + let oid = String::from_utf8(fields[1].to_vec())?; + let symbolic_target = if fields[2].is_empty() { + None + } else { + Some(String::from_utf8(fields[2].to_vec())?) + }; + refs.insert( + name, + RefSnapshot { + oid, + symbolic_target, + }, + ); + } + Ok(refs) +} + +fn install_refs(repo: &Path, refs: &BTreeMap) -> Result<()> { + for (name, reference) in refs { + if let Some(target) = &reference.symbolic_target { + run_git( + Some(repo), + &["symbolic-ref", name, target], + "install migrated symbolic ref", + )?; + } else if oid_is_available(repo, &reference.oid)? { + run_git( + Some(repo), + &["update-ref", name, &reference.oid], + "install migrated ref", + )?; + } else { + let relative = Path::new(name); + if !name.starts_with("refs/") + || relative.is_absolute() + || relative + .components() + .any(|component| !matches!(component, std::path::Component::Normal(_))) + { + return Err(anyhow!("unsafe migrated ref name {name:?}")); + } + write_atomic( + &repo.join(relative), + format!("{}\n", reference.oid).as_bytes(), + )?; + warn!( + repository = %repo.display(), + reference = %name, + oid = %reference.oid, + "Preserving a pre-existing ref whose target object is missing" + ); + } + } + Ok(()) +} + +fn verify_view( + view: &Path, + expected_refs: &BTreeMap, + expected_head: &[u8], +) -> Result<()> { + if snapshot_refs(view)? != *expected_refs { + return Err(anyhow!("migrated refs differ in {}", view.display())); + } + if std::fs::read(view.join("HEAD"))? != expected_head { + return Err(anyhow!("migrated HEAD differs in {}", view.display())); + } + for (name, reference) in expected_refs { + if !oid_is_available(view, &reference.oid)? { + warn!( + repository = %view.display(), + reference = %name, + oid = %reference.oid, + "Verified a preserved ref whose target object is missing" + ); + } + } + Ok(()) +} + +fn list_all_objects(repo: &Path) -> Result> { + let output = Command::new("git") + .args([ + "cat-file", + "--batch-all-objects", + "--batch-check=%(objectname)", + ]) + .current_dir(repo) + .output()?; + if !output.status.success() { + return Err(anyhow!( + "enumerate every object in {}: {}", + repo.display(), + String::from_utf8_lossy(&output.stderr).trim() + )); + } + Ok(String::from_utf8(output.stdout)? + .lines() + .map(str::to_owned) + .collect()) +} + +fn copy_object_database(source: &Path, target: &Path, format: ObjectFormat) -> Result<()> { + let source_objects = source.join("objects"); + let target_objects = target.join("objects"); + let oid_len = match format { + ObjectFormat::Sha1 => 40, + ObjectFormat::Sha256 => 64, + }; + for fanout in std::fs::read_dir(&source_objects)? { + let fanout = fanout?; + let name = fanout.file_name().to_string_lossy().into_owned(); + if fanout.file_type()?.is_dir() + && name.len() == 2 + && name.chars().all(|ch| ch.is_ascii_hexdigit()) + { + for object in std::fs::read_dir(fanout.path())? { + let object = object?; + let suffix = object.file_name().to_string_lossy().into_owned(); + if object.file_type()?.is_file() + && suffix.len() == oid_len - 2 + && suffix.chars().all(|ch| ch.is_ascii_hexdigit()) + { + let destination = target_objects.join(&name).join(&suffix); + install_if_absent(&object.path(), &destination)?; + } + } + } + } + let source_packs = source_objects.join("pack"); + let target_packs = target_objects.join("pack"); + if source_packs.is_dir() { + for pack in std::fs::read_dir(source_packs)? { + let pack = pack?; + let name = pack.file_name().to_string_lossy().into_owned(); + if !pack.file_type()?.is_file() || !name.ends_with(".pack") { + continue; + } + let index = pack.path().with_extension("idx"); + if !index.is_file() { + warn!( + pack = %pack.path().display(), + "Skipping Git-invisible pack without an index; the migration backup retains it" + ); + continue; + } + install_if_absent(&pack.path(), &target_packs.join(&name))?; + install_if_absent(&index, &target_packs.join(index.file_name().unwrap()))?; + } + } + Ok(()) +} + +fn pack_explicit_objects(repo: &Path, base: &Path, oids: &[String]) -> Result { + let mut child = Command::new("git") + .arg("pack-objects") + .arg(base) + .current_dir(repo) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn()?; + { + let stdin = child.stdin.as_mut().context("git pack-objects stdin")?; + for oid in oids { + writeln!(stdin, "{oid}")?; + } + } + let output = child.wait_with_output()?; + if !output.status.success() { + return Err(anyhow!( + "build migration union pack: {}", + String::from_utf8_lossy(&output.stderr).trim() + )); + } + Ok(String::from_utf8(output.stdout)?.trim().to_owned()) +} + +fn verify_pack(index: &Path) -> Result<()> { + let status = Command::new("git") + .arg("verify-pack") + .arg(index) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status()?; + if !status.success() { + return Err(anyhow!( + "migration pack does not verify: {}", + index.display() + )); + } + Ok(()) +} + +fn oid_is_available(repo: &Path, oid: &str) -> Result { + let output = Command::new("git") + .args(["cat-file", "-e", oid]) + .current_dir(repo) + .output()?; + Ok(output.status.success()) +} + +fn ensure_oids(repo: &Path, expected: &[String]) -> Result<()> { + if expected.is_empty() { + return Ok(()); + } + let available: HashSet<_> = list_all_objects(repo)?.into_iter().collect(); + if let Some(missing) = expected.iter().find(|oid| !available.contains(*oid)) { + return Err(anyhow!( + "object {missing} is missing from {}", + repo.display() + )); + } + Ok(()) +} + +fn audit_family_connectivity(repo: &Path) -> Result<()> { + let output = Command::new("git") + .args(["fsck", "--full", "--no-reflogs", "--no-dangling"]) + .current_dir(repo) + .output()?; + if output.status.success() { + return Ok(()); + } + + let diagnostic = format!( + "{}{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let diagnostic: String = diagnostic.trim().chars().take(2048).collect(); + warn!( + repository = %repo.display(), + status = %output.status, + diagnostic, + "Migrated family preserves a pre-existing incomplete object graph" + ); + Ok(()) +} + +fn retain_family_tip( + storage: &LocalGitStorage, + key: &FamilyKey, + family: &Path, + name: &str, + oid: &str, +) -> Result<()> { + if oid_is_available(family, oid)? { + storage.retain_tip(key, name, oid)?; + storage.advertise_base_tip(key, name, oid)?; + } else { + warn!( + repository = %family.display(), + reference = %name, + oid = %oid, + "Not advertising a pre-existing ref whose target object is missing" + ); + } + Ok(()) +} + +fn install_if_absent(source: &Path, destination: &Path) -> Result<()> { + if destination.exists() { + return Ok(()); + } + if let Some(parent) = destination.parent() { + std::fs::create_dir_all(parent)?; + } + match std::fs::hard_link(source, destination) { + Ok(()) => Ok(()), + Err(_) => { + std::fs::copy(source, destination)?; + Ok(()) + } + } +} + +fn copy_directory_contents(source: &Path, target: &Path) -> Result<()> { + if !source.is_dir() { + return Ok(()); + } + std::fs::create_dir_all(target)?; + for entry in std::fs::read_dir(source)? { + let entry = entry?; + let destination = target.join(entry.file_name()); + if entry.file_type()?.is_dir() { + copy_directory_contents(&entry.path(), &destination)?; + } else if entry.file_type()?.is_file() { + std::fs::copy(entry.path(), destination)?; + } + } + Ok(()) +} + +fn load_journals(storage: &LocalGitStorage) -> Result> { + let root = migration_root(storage).join("journal"); + let Ok(entries) = std::fs::read_dir(root) else { + return Ok(Vec::new()); + }; + let mut journals = Vec::new(); + for entry in entries { + let entry = entry?; + if !entry.file_type()?.is_file() + || entry.path().extension().and_then(|value| value.to_str()) != Some("json") + { + continue; + } + let journal = serde_json::from_slice(&std::fs::read(entry.path())?) + .with_context(|| format!("parse migration journal {}", entry.path().display()))?; + journals.push((entry.path(), journal)); + } + journals.sort_by(|left, right| left.0.cmp(&right.0)); + Ok(journals) +} + +fn write_journal(path: &Path, journal: &ViewJournal) -> Result<()> { + write_atomic(path, &serde_json::to_vec_pretty(journal)?) +} + +fn write_atomic(path: &Path, bytes: &[u8]) -> Result<()> { + let parent = path + .parent() + .context("atomic migration file has no parent")?; + std::fs::create_dir_all(parent)?; + let pending = parent.join(format!( + ".{}.pending-{}", + path.file_name().unwrap_or_default().to_string_lossy(), + std::process::id() + )); + let mut file = std::fs::OpenOptions::new() + .write(true) + .create(true) + .truncate(true) + .open(&pending)?; + file.write_all(bytes)?; + file.sync_all()?; + std::fs::rename(&pending, path)?; + fsync_directory(parent) +} + +fn fsync_directory(path: &Path) -> Result<()> { + std::fs::File::open(path)?.sync_all()?; + Ok(()) +} + +fn run_git(current_dir: Option<&Path>, args: &[&str], description: &str) -> Result<()> { + let mut command = Command::new("git"); + if let Some(current_dir) = current_dir { + command.current_dir(current_dir); + } + let output = command.args(args).output()?; + if !output.status.success() { + return Err(anyhow!( + "{description}: {}", + String::from_utf8_lossy(&output.stderr).trim() + )); + } + Ok(()) +} + +fn migration_root(storage: &LocalGitStorage) -> PathBuf { + storage.internal_path().join("migration") +} + +fn backup_root(storage: &LocalGitStorage) -> PathBuf { + migration_root(storage).join("backups") +} + +fn journal_path(storage: &LocalGitStorage, relative_path: &Path) -> PathBuf { + migration_root(storage) + .join("journal") + .join(format!("{}.json", journal_name(relative_path))) +} + +fn journal_name(relative_path: &Path) -> String { + sha256::Hash::hash(relative_path.to_string_lossy().as_bytes()).to_string() +} + +fn parse_object_format(value: &str) -> Result { + match value { + "sha1" => Ok(ObjectFormat::Sha1), + "sha256" => Ok(ObjectFormat::Sha256), + _ => Err(anyhow!("unsupported journal object format {value:?}")), + } +} + +fn validate_relative_view_path(path: &Path) -> Result { + if path.is_absolute() + || path + .components() + .any(|component| !matches!(component, std::path::Component::Normal(_))) + { + return Err(anyhow!("unsafe migration view path {}", path.display())); + } + let components: Vec<_> = path + .components() + .filter_map(|component| component.as_os_str().to_str()) + .collect(); + let repo_name = match components.as_slice() { + [owner, repo] if owner.starts_with("npub1") && PublicKey::from_bech32(owner).is_ok() => { + *repo + } + ["prs", submitter, repo] + if submitter.len() == 64 && submitter.chars().all(|ch| ch.is_ascii_hexdigit()) => + { + *repo + } + _ => return Err(anyhow!("unsafe migration view path {}", path.display())), + }; + let identifier = repo_name + .strip_suffix(".git") + .filter(|identifier| validate_repository_identifier(identifier)) + .ok_or_else(|| anyhow!("unsafe migration view path {}", path.display()))?; + Ok(identifier.to_owned()) +} + +fn remove_migration_staging(path: &Path, allowed_root: &Path) -> Result<()> { + if !path.starts_with(allowed_root) || path == allowed_root { + return Err(anyhow!( + "refusing to remove migration path outside staging root: {}", + path.display() + )); + } + match std::fs::remove_dir_all(path) { + Ok(()) => Ok(()), + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(error) => Err(error.into()), + } +} + +#[cfg(test)] +mod tests { + use std::io::Write; + + use nostr_sdk::prelude::{Keys, ToBech32}; + + use super::*; + + fn git(repo: &Path, args: &[&str]) -> String { + let output = Command::new("git") + .args(args) + .current_dir(repo) + .output() + .unwrap(); + assert!( + output.status.success(), + "git {args:?}: {}", + String::from_utf8_lossy(&output.stderr) + ); + String::from_utf8_lossy(&output.stdout).trim().to_owned() + } + + fn legacy_repo( + root: &Path, + owner: &str, + identifier: &str, + contents: &[u8], + ) -> (PathBuf, String) { + let repo = root.join(owner).join(format!("{identifier}.git")); + std::fs::create_dir_all(repo.parent().unwrap()).unwrap(); + assert!(Command::new("git") + .args(["init", "--bare", "--quiet", "--initial-branch=main"]) + .arg(&repo) + .status() + .unwrap() + .success()); + let mut child = Command::new("git") + .args(["hash-object", "-w", "--stdin"]) + .current_dir(&repo) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child.stdin.take().unwrap().write_all(contents).unwrap(); + let output = child.wait_with_output().unwrap(); + assert!(output.status.success()); + let oid = String::from_utf8_lossy(&output.stdout).trim().to_owned(); + git(&repo, &["update-ref", "refs/test/data", &oid]); + (repo, oid) + } + + #[test] + fn batch_object_verification_reports_missing_objects() { + let temp = tempfile::tempdir().unwrap(); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let (repo, oid) = legacy_repo(temp.path(), &owner, "example", b"present\n"); + + ensure_oids(&repo, std::slice::from_ref(&oid)).unwrap(); + + let missing = "0".repeat(40); + let error = ensure_oids(&repo, std::slice::from_ref(&missing)).unwrap_err(); + assert_eq!( + error.to_string(), + format!("object {missing} is missing from {}", repo.display()) + ); + } + + #[tokio::test] + async fn migration_preserves_a_preexisting_incomplete_object_graph() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let (repo, _) = legacy_repo(storage.git_data_path(), &owner, "incomplete", b"present\n"); + let tree = git(&repo, &["mktree"]); + let missing_parent = "1".repeat(40); + let commit = format!( + "tree {tree}\nparent {missing_parent}\nauthor Test 0 +0000\ncommitter Test 0 +0000\n\nincomplete history\n" + ); + let mut child = Command::new("git") + .args(["hash-object", "-t", "commit", "-w", "--stdin"]) + .current_dir(&repo) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child + .stdin + .take() + .unwrap() + .write_all(commit.as_bytes()) + .unwrap(); + let output = child.wait_with_output().unwrap(); + assert!(output.status.success()); + let commit_oid = String::from_utf8_lossy(&output.stdout).trim().to_owned(); + git(&repo, &["update-ref", "refs/heads/main", &commit_oid]); + assert!(!Command::new("git") + .args(["fsck", "--full", "--no-reflogs", "--no-dangling"]) + .current_dir(&repo) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .status() + .unwrap() + .success()); + + let report = migrate_on_startup(&storage).await.unwrap(); + + assert_eq!(report.migrated_views, 1); + let key = FamilyKey::new(ObjectFormat::Sha1, "incomplete").unwrap(); + assert!(oid_is_available(&storage.family_repo_path(&key), &commit_oid).unwrap()); + } + + #[tokio::test] + async fn migration_archives_a_pack_without_an_index() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let (repo, oid) = legacy_repo(storage.git_data_path(), &owner, "orphan-pack", b"present\n"); + let relative = repo.strip_prefix(storage.git_data_path()).unwrap(); + let orphan_name = "pack-0000000000000000000000000000000000000000.pack"; + let orphan = repo.join("objects/pack").join(orphan_name); + std::fs::create_dir_all(orphan.parent().unwrap()).unwrap(); + std::fs::write(&orphan, b"not a readable Git pack\n").unwrap(); + + let report = migrate_on_startup(&storage).await.unwrap(); + + assert_eq!(report.migrated_views, 1); + let archived = backup_root(&storage) + .join(relative) + .join("objects/pack") + .join(orphan_name); + assert_eq!( + std::fs::read(archived).unwrap(), + b"not a readable Git pack\n" + ); + let key = FamilyKey::new(ObjectFormat::Sha1, "orphan-pack").unwrap(); + let family = storage.family_repo_path(&key); + assert!(!family.join("objects/pack").join(orphan_name).exists()); + assert!(oid_is_available(&family, &oid).unwrap()); + } + + #[tokio::test] + async fn migration_preserves_a_ref_with_a_missing_target() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let (repo, _) = legacy_repo(storage.git_data_path(), &owner, "missing-ref", b"present\n"); + let missing = "2".repeat(40); + let broken_ref = repo.join("refs/heads/broken"); + std::fs::create_dir_all(broken_ref.parent().unwrap()).unwrap(); + std::fs::write(&broken_ref, format!("{missing}\n")).unwrap(); + assert_eq!( + snapshot_refs(&repo).unwrap()["refs/heads/broken"].oid, + missing + ); + + let report = migrate_on_startup(&storage).await.unwrap(); + + assert_eq!(report.migrated_views, 1); + assert_eq!( + std::fs::read_to_string(&broken_ref).unwrap(), + format!("{missing}\n") + ); + assert_eq!( + snapshot_refs(&repo).unwrap()["refs/heads/broken"].oid, + missing + ); + let key = FamilyKey::new(ObjectFormat::Sha1, "missing-ref").unwrap(); + let family = storage.family_repo_path(&key); + assert!(!oid_is_available(&family, &missing).unwrap()); + assert!(!git( + &family, + &["for-each-ref", "--format=%(objectname)", "refs/grasp/"] + ) + .lines() + .any(|oid| oid == missing)); + } + + #[tokio::test] + async fn migration_deduplicates_views_and_preserves_unreachable_objects() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let owner_one = Keys::generate().public_key().to_bech32().unwrap(); + let owner_two = Keys::generate().public_key().to_bech32().unwrap(); + let (first, first_oid) = legacy_repo( + storage.git_data_path(), + &owner_one, + "example", + b"shared data\n", + ); + let (second, second_oid) = legacy_repo( + storage.git_data_path(), + &owner_two, + "example", + b"other data\n", + ); + let mut child = Command::new("git") + .args(["hash-object", "-w", "--stdin"]) + .current_dir(&first) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child + .stdin + .take() + .unwrap() + .write_all(b"unreachable rollback data\n") + .unwrap(); + let output = child.wait_with_output().unwrap(); + let unreachable = String::from_utf8_lossy(&output.stdout).trim().to_owned(); + + let report = migrate_on_startup(&storage).await.unwrap(); + assert_eq!(report.migrated_views, 2); + assert_eq!(report.families_built, 1); + let key = FamilyKey::sha1("example").unwrap(); + assert!(storage.is_thin_view(&key, &first)); + assert!(storage.is_thin_view(&key, &second)); + assert_eq!( + super::super::get_ref_commit(&first, "refs/test/data"), + Some(first_oid) + ); + assert_eq!( + super::super::get_ref_commit(&second, "refs/test/data"), + Some(second_oid) + ); + assert!(oid_is_available(&storage.family_repo_path(&key), &unreachable).unwrap()); + assert!(backup_root(&storage) + .join(&owner_one) + .join("example.git") + .is_dir()); + assert!(backup_root(&storage) + .join(&owner_two) + .join("example.git") + .is_dir()); + + let second_report = migrate_on_startup(&storage).await.unwrap(); + assert!(second_report.already_current); + assert_eq!(second_report.migrated_views, 0); + } + + #[tokio::test] + async fn migration_resumes_after_legacy_repo_was_renamed_to_backup() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let owner = Keys::generate().public_key().to_bech32().unwrap(); + let (view, oid) = legacy_repo(storage.git_data_path(), &owner, "recovery", b"recover me\n"); + let candidate = + candidate_for_active(&storage, &view, FamilyKey::sha1("recovery").unwrap()).unwrap(); + let refs = snapshot_refs(&view).unwrap(); + let head = std::fs::read(view.join("HEAD")).unwrap(); + write_journal( + &candidate.journal_path, + &ViewJournal { + version: MIGRATION_VERSION, + relative_path: candidate.relative_path.clone(), + object_format: "sha1".to_owned(), + identifier: "recovery".to_owned(), + refs, + head, + state: JournalState::Prepared, + }, + ) + .unwrap(); + std::fs::create_dir_all(candidate.backup_path.parent().unwrap()).unwrap(); + std::fs::rename(&view, &candidate.backup_path).unwrap(); + + let report = migrate_on_startup(&storage).await.unwrap(); + assert_eq!(report.recovered_views, 1); + assert_eq!( + super::super::get_ref_commit(&view, "refs/test/data"), + Some(oid) + ); + assert!(storage.is_thin_view(&FamilyKey::sha1("recovery").unwrap(), &view)); + } + + #[tokio::test] + async fn archived_repository_import_preserves_unreachable_objects() { + let temp = tempfile::tempdir().unwrap(); + let storage = LocalGitStorage::new(temp.path().join("git")); + let (archive, referenced) = + legacy_repo(temp.path(), "archive", "restored", b"referenced\n"); + let mut child = Command::new("git") + .args(["hash-object", "-w", "--stdin"]) + .current_dir(&archive) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .spawn() + .unwrap(); + child + .stdin + .take() + .unwrap() + .write_all(b"unreachable archive data\n") + .unwrap(); + let output = child.wait_with_output().unwrap(); + assert!(output.status.success()); + let unreachable = String::from_utf8_lossy(&output.stdout).trim().to_owned(); + + let key = FamilyKey::sha1("restored").unwrap(); + let _lease = storage.write_lease(&key).await.unwrap(); + import_archived_repository(&storage, &key, &archive).unwrap(); + + let family = storage.family_repo_path(&key); + assert!(oid_is_available(&family, &referenced).unwrap()); + assert!(oid_is_available(&family, &unreachable).unwrap()); + assert!(!git(&family, &["for-each-ref", "refs/grasp/retained/"]).is_empty()); + } +} diff --git a/src/git/mod.rs b/src/git/mod.rs index 4444fae..8c09327 100644 --- a/src/git/mod.rs +++ b/src/git/mod.rs @@ -19,6 +19,8 @@ pub mod authorization; pub mod handlers; +pub mod integrity; +pub mod migration; pub mod process; pub mod protocol; pub mod storage; 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/main.rs b/src/main.rs index 345ef4f..653caf9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -36,6 +36,9 @@ enum Cli { /// /// This is an operator/admin maintenance command and is idempotent. HoldingEject(nostr::lifecycle::HoldingEjectArgs), + + /// Queue an identifier-family integrity check in the running relay. + IntegrityCheck(ngit_grasp::git::integrity::IntegrityCheckArgs), } #[tokio::main] @@ -47,7 +50,13 @@ async fn main() -> Result<()> { // If not, prepend the implicit "serve" subcommand so that clap routes to Cli::Serve // and all relay flags are parsed normally (preserving backward compatibility). let mut args: Vec = std::env::args().collect(); - let known_subcommands = ["serve", "cleanup-empty-repos", "holding-eject", "help"]; + let known_subcommands = [ + "serve", + "cleanup-empty-repos", + "holding-eject", + "integrity-check", + "help", + ]; let has_subcommand = args.get(1).is_some_and(|a| { known_subcommands.contains(&a.as_str()) || matches!(a.as_str(), "-h" | "--help" | "-V" | "--version") @@ -59,6 +68,20 @@ async fn main() -> Result<()> { match Cli::parse_from(args) { Cli::CleanupEmptyRepos(cleanup_args) => cleanup_empty_repos::run(&cleanup_args).await, Cli::HoldingEject(eject_args) => nostr::lifecycle::run_holding_eject(eject_args).await, + Cli::IntegrityCheck(integrity_args) => { + let path = ngit_grasp::git::integrity::enqueue_manual_check(&integrity_args)?; + println!( + "Queued {} for identifier '{}' at {}", + if integrity_args.repair { + "integrity check and repair" + } else { + "integrity check" + }, + integrity_args.identifier, + path.display() + ); + Ok(()) + } Cli::Serve(config) => { let mut config = *config; config.relay_owner_nsec = Some(Config::load_relay_owner_key()?); diff --git a/src/nostr/lifecycle/deletion/purgatory.rs b/src/nostr/lifecycle/deletion/purgatory.rs index d87c4cb..7540f2a 100644 --- a/src/nostr/lifecycle/deletion/purgatory.rs +++ b/src/nostr/lifecycle/deletion/purgatory.rs @@ -1,5 +1,4 @@ use std::collections::HashSet; -use std::process::Command; use nostr_sdk::prelude::{Event, Kind, Timestamp}; @@ -57,7 +56,7 @@ impl DeletionPolicy { } } - if let Err(e) = Self::replace_live_repo_with_empty_purgatory_repo(&repo_path) { + if let Err(e) = self.replace_live_repo_with_empty_purgatory_repo(&repo_path, &identifier) { tracing::warn!( owner = %owner.to_hex(), identifier = %identifier, @@ -84,7 +83,9 @@ impl DeletionPolicy { } fn replace_live_repo_with_empty_purgatory_repo( + &self, repo_path: &std::path::Path, + identifier: &str, ) -> anyhow::Result<()> { if repo_path.exists() { std::fs::remove_dir_all(repo_path).map_err(|e| { @@ -96,37 +97,9 @@ impl DeletionPolicy { })?; } - let parent = repo_path - .parent() - .ok_or_else(|| anyhow::anyhow!("invalid repository path {}", repo_path.display()))?; - std::fs::create_dir_all(parent).map_err(|e| { - anyhow::anyhow!( - "failed to create repository parent {}: {}", - parent.display(), - e - ) - })?; - - let output = Command::new("git") - .args([ - "init", - "--bare", - repo_path.to_str().ok_or_else(|| { - anyhow::anyhow!("non-utf8 repository path {}", repo_path.display()) - })?, - ]) - .output() - .map_err(|e| anyhow::anyhow!("failed to execute git init: {}", e))?; - - if !output.status.success() { - return Err(anyhow::anyhow!( - "git init failed for {}: {}", - repo_path.display(), - String::from_utf8_lossy(&output.stderr) - )); - } - - Ok(()) + let storage = crate::git::storage::LocalGitStorage::new(&self.ctx.git_data_path); + let family = crate::git::storage::FamilyKey::sha1(identifier)?; + storage.create_thin_view(&family, repo_path) } /// Remove any purgatory entries targeted by this deletion event. diff --git a/src/nostr/lifecycle/deletion/recovery.rs b/src/nostr/lifecycle/deletion/recovery.rs index 121b636..e140ecd 100644 --- a/src/nostr/lifecycle/deletion/recovery.rs +++ b/src/nostr/lifecycle/deletion/recovery.rs @@ -18,6 +18,7 @@ use nostr_sdk::prelude::{Event, Kind, PublicKey, Timestamp}; use tar::Archive as TarArchive; use super::DeletionService; +use crate::git::storage::{FamilyKey, LocalGitStorage, ObjectFormat}; use crate::nostr::lifecycle::RecoveryMetadataRecord; // --------------------------------------------------------------------------- @@ -47,7 +48,7 @@ impl DeletionService { // --------------------------------------------------------------------------- impl DeletionService { - pub(crate) fn restore_git_archive_to_repo( + pub(crate) async fn restore_git_archive_to_repo( &self, archive_path: &Path, owner_path_component: &str, @@ -82,7 +83,7 @@ impl DeletionService { ) })?; - let result = (|| -> anyhow::Result { + let result: anyhow::Result = async { let archive_file = File::open(archive_path).map_err(|e| { anyhow::anyhow!("failed to open archive {}: {}", archive_path.display(), e) })?; @@ -106,10 +107,36 @@ impl DeletionService { )); } + let storage = LocalGitStorage::new(&self.ctx.git_data_path); + let key = FamilyKey::new(ObjectFormat::detect(&extracted_repo)?, identifier)?; + let _family_lease = storage.write_lease(&key).await?; + crate::git::migration::import_archived_repository(&storage, &key, &extracted_repo)?; + if target_repo.is_dir() { + if !storage.is_thin_view(&key, &target_repo) { + return Err(anyhow::anyhow!( + "active repository {} is not an identifier-family view; restart to migrate it before recovery", + target_repo.display() + )); + } return Self::restore_archived_nostr_refs(&extracted_repo, &target_repo); } + std::fs::remove_dir_all(extracted_repo.join("objects")).map_err(|e| { + anyhow::anyhow!( + "failed to remove imported object database from {}: {}", + extracted_repo.display(), + e + ) + })?; + std::fs::create_dir_all(extracted_repo.join("objects")).map_err(|e| { + anyhow::anyhow!( + "failed to create thin object directory in {}: {}", + extracted_repo.display(), + e + ) + })?; + storage.configure_thin_view(&key, &extracted_repo)?; std::fs::rename(&extracted_repo, &target_repo) .map_err(|e| { anyhow::anyhow!( @@ -120,7 +147,8 @@ impl DeletionService { ) }) .map(|_| true) - })(); + } + .await; let _ = std::fs::remove_dir_all(&staging_dir); result @@ -551,11 +579,10 @@ impl DeletionService { ); git_ready = false; } else { - match self.restore_git_archive_to_repo( - &abs_path, - &owner_component, - identifier, - ) { + match self + .restore_git_archive_to_repo(&abs_path, &owner_component, identifier) + .await + { Ok(restored) => { archive_restored = restored; } diff --git a/src/nostr/policy/announcement.rs b/src/nostr/policy/announcement.rs index 60961bc..7e33148 100644 --- a/src/nostr/policy/announcement.rs +++ b/src/nostr/policy/announcement.rs @@ -412,26 +412,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/mod.rs b/src/purgatory/mod.rs index 48baa5a..96c2334 100644 --- a/src/purgatory/mod.rs +++ b/src/purgatory/mod.rs @@ -964,56 +964,32 @@ impl Purgatory { } } - // If the entry was soft-expired, recreate the bare repo outside the + // If the entry was soft-expired, recreate the thin view outside the // mutable borrow so we don't hold the DashMap lock during I/O. if let Some((repo_path, was_soft_expired)) = revival_info { if was_soft_expired { if !repo_path.exists() { - match std::fs::create_dir_all(&repo_path) { - Ok(()) => { - // Initialise as a bare git repository - let status = std::process::Command::new("git") - .args(["init", "--bare"]) - .arg(&repo_path) - .status(); - match status { - Ok(s) if s.success() => { - tracing::info!( - path = %repo_path.display(), - owner = %owner, - identifier = %identifier, - "Recreated bare repository for revived soft-expired announcement" - ); - } - Ok(s) => { - tracing::warn!( - path = %repo_path.display(), - exit_code = ?s.code(), - "git init --bare failed when reviving soft-expired announcement" - ); - } - Err(e) => { - tracing::warn!( - path = %repo_path.display(), - error = %e, - "Failed to run git init --bare when reviving soft-expired announcement" - ); - } - } - } - Err(e) => { - tracing::warn!( - path = %repo_path.display(), - error = %e, - "Failed to create directory when reviving soft-expired announcement" - ); - } + let storage = crate::git::storage::LocalGitStorage::new(&self._git_data_path); + let result = crate::git::storage::FamilyKey::sha1(identifier) + .and_then(|family| storage.create_thin_view(&family, &repo_path)); + match result { + Ok(()) => tracing::info!( + path = %repo_path.display(), + owner = %owner, + identifier = %identifier, + "Recreated thin repository view for revived soft-expired announcement" + ), + Err(e) => tracing::warn!( + path = %repo_path.display(), + error = %e, + "Failed to recreate thin repository view for revived soft-expired announcement" + ), } } tracing::info!( owner = %owner, identifier = %identifier, - "Revived soft-expired announcement (bare repo recreated, expiry extended)" + "Revived soft-expired announcement (thin view recreated, expiry extended)" ); } } @@ -2262,6 +2238,44 @@ fn cleanup_defers_expiry_only_while_repository_sync_is_active() { assert_eq!(purgatory.expired_count(), 2); } +#[test] +fn extending_soft_expired_announcement_recreates_a_thin_view() { + let directory = tempfile::tempdir().unwrap(); + let git_data_path = directory.path().join("git"); + let repo_path = git_data_path.join("owner").join("revived.git"); + let purgatory = Purgatory::new(&git_data_path); + let keys = Keys::generate(); + let identifier = "revived"; + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .finalize(&keys) + .unwrap(); + + purgatory.add_announcement( + announcement, + identifier.to_string(), + keys.public_key(), + repo_path.clone(), + HashSet::new(), + ); + purgatory + .announcement_purgatory + .get_mut(&(keys.public_key(), identifier.to_string())) + .unwrap() + .soft_expired = true; + + purgatory.extend_announcement_expiry(&keys.public_key(), identifier, Duration::from_secs(60)); + + let storage = crate::git::storage::LocalGitStorage::new(&git_data_path); + let family = crate::git::storage::FamilyKey::sha1(identifier).unwrap(); + assert!(storage.is_thin_view(&family, &repo_path)); + assert!( + !purgatory + .find_announcement(&keys.public_key(), identifier) + .unwrap() + .soft_expired + ); +} + #[test] fn test_cleanup_mixed_expired_and_fresh() { use std::time::Duration; diff --git a/src/purgatory/sync/context.rs b/src/purgatory/sync/context.rs index 07f9835..9cebcea 100644 --- a/src/purgatory/sync/context.rs +++ b/src/purgatory/sync/context.rs @@ -21,6 +21,7 @@ use std::time::{Duration, Instant}; pub enum GitFetchRole { Primary, Hedge, + Integrity, } impl GitFetchRole { @@ -28,6 +29,7 @@ impl GitFetchRole { match self { Self::Primary => "primary", Self::Hedge => "hedge", + Self::Integrity => "integrity", } } } @@ -229,6 +231,7 @@ use tokio::io::{AsyncRead, AsyncReadExt}; use tokio::process::Command; use tracing::debug; +use crate::git::storage::{FamilyKey, LocalGitStorage, ObjectFormat}; use crate::nostr::builder::Nip34WritePolicy; use crate::nostr::events::RepositoryState; use crate::nostr::SharedDatabase; @@ -330,6 +333,33 @@ impl RealSyncContext { pub fn git_naughty_list(&self) -> &Arc { &self.git_naughty_list } + + /// Clone URLs from accepted repository announcements for one identifier. + /// + /// Integrity repair deliberately excludes purgatory URLs: it heals + /// durable families only from repository sources already accepted into + /// the relay database. The normal outbound policy is still applied at the + /// point of each fetch. + pub async fn accepted_repository_clone_urls(&self, identifier: &str) -> Result> { + let data = crate::git::authorization::fetch_repository_data_excluding_purgatory( + &self.database, + identifier, + ) + .await?; + let mut urls: Vec<_> = data + .announcements + .into_iter() + .flat_map(|announcement| announcement.clone_urls) + .filter(|url| { + self.our_domain_value + .as_deref() + .is_none_or(|domain| !crate::outbound::url_matches_service_domain(url, domain)) + }) + .collect(); + urls.sort(); + urls.dedup(); + Ok(urls) + } } const MISS_MEMO_TTL: Duration = Duration::from_secs(30 * 60); @@ -436,6 +466,8 @@ fn hardened_git_command( repo_path: &Path, resolve_pin: Option<&str>, auth_header: Option<&str>, + object_directory: Option<&Path>, + negotiation_algorithm: Option<&str>, args: &[String], ) -> Command { let mut command = Command::new("git"); @@ -446,6 +478,11 @@ fn hardened_git_command( .arg("http.proxy=") .arg("-c") .arg("credential.helper="); + if let Some(algorithm) = negotiation_algorithm { + command + .arg("-c") + .arg(format!("fetch.negotiationAlgorithm={algorithm}")); + } if let Some(pin) = resolve_pin { command.arg("-c").arg(format!("http.curloptResolve={pin}")); } @@ -469,6 +506,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 } @@ -799,6 +839,15 @@ fn is_object_missing_error(stderr: &str) -> bool { || stderr.contains("Server does not allow request for unadvertised object") } +/// Integrity repair starts from a repository whose refs may name objects that +/// are already missing. Normal fetch negotiation walks those refs as local +/// "have" tips, which can make Git request no pack and then fail its own +/// connectivity check. Repair fetches therefore advertise no local haves and +/// ask the source for the complete requested object closure. +fn fetch_negotiation_algorithm(role: GitFetchRole) -> Option<&'static str> { + (role == GitFetchRole::Integrity).then_some("noop") +} + /// Record a remote failure with the naughty list tracker (only error /// categories the tracker classifies as persistent are recorded). fn record_remote_failure( @@ -974,6 +1023,21 @@ impl SyncContext for RealSyncContext { }; let resolve_pin = resolve_pin_entry(&resolved); + let storage = LocalGitStorage::new(&self.git_data_path); + let family_key = family_key_for_view(repo_path, &self.git_data_path); + 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 +1092,8 @@ impl SyncContext for RealSyncContext { &repo_path, resolve_pin.as_deref(), fresh_auth_header().as_deref(), + None, + None, &ls_remote_args, ), &domain, @@ -1098,6 +1164,8 @@ impl SyncContext for RealSyncContext { &repo_path, resolve_pin.as_deref(), fresh_auth_header().as_deref(), + family_objects, + fetch_negotiation_algorithm(role), &args, ), &domain, @@ -1168,6 +1236,8 @@ impl SyncContext for RealSyncContext { &repo_path, resolve_pin.as_deref(), fresh_auth_header().as_deref(), + family_objects, + fetch_negotiation_algorithm(role), &args, ), &domain, @@ -1219,6 +1289,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 @@ -1310,10 +1401,28 @@ impl SyncContext for RealSyncContext { } } +fn family_key_for_view(repo_path: &Path, git_data_path: &Path) -> Option { + let identifier = crate::git::sync::extract_identifier_from_repo_path(repo_path, git_data_path)?; + let object_format = ObjectFormat::detect(repo_path).ok()?; + FamilyKey::new(object_format, identifier).ok() +} + #[cfg(test)] mod fetch_helper_tests { use super::*; + #[test] + fn family_key_for_view_preserves_the_repository_object_format() { + let temp = tempfile::tempdir().unwrap(); + let git_data_path = temp.path().join("git"); + let storage = LocalGitStorage::new(&git_data_path); + let expected = FamilyKey::new(ObjectFormat::Sha256, "shared").unwrap(); + let view = git_data_path.join("owner").join("shared.git"); + storage.create_thin_view(&expected, &view).unwrap(); + + assert_eq!(family_key_for_view(&view, &git_data_path), Some(expected)); + } + /// Hermetic git for these tests, immune to ambient git configuration /// such as hooks, signing, and templates. fn fixture_git() -> std::process::Command { @@ -1665,6 +1774,17 @@ mod fetch_helper_tests { )); } + #[test] + fn integrity_fetches_disable_local_have_negotiation() { + assert_eq!(fetch_negotiation_algorithm(GitFetchRole::Primary), None); + assert_eq!(fetch_negotiation_algorithm(GitFetchRole::Hedge), None); + assert_eq!( + fetch_negotiation_algorithm(GitFetchRole::Integrity), + Some("noop") + ); + assert_eq!(GitFetchRole::Integrity.as_str(), "integrity"); + } + #[test] fn miss_memo_is_stable_order_independent_invalidated_and_bounded() { let first = HashSet::from(["b".to_string(), "a".to_string()]); diff --git a/src/purgatory/sync/mod.rs b/src/purgatory/sync/mod.rs index 022a556..e46b29a 100644 --- a/src/purgatory/sync/mod.rs +++ b/src/purgatory/sync/mod.rs @@ -13,7 +13,7 @@ mod r#loop; mod queue; mod throttle; -pub use context::{ProcessResult, RealSyncContext, SyncContext}; +pub use context::{GitFetchRole, ProcessResult, RealSyncContext, SyncContext}; pub use functions::{ get_throttled_domains_with_untried_urls, sync_identifier, sync_identifier_from_url, sync_identifier_next_url, ThrottledDomainInfo, diff --git a/src/server.rs b/src/server.rs index 48bff84..5a51ff6 100644 --- a/src/server.rs +++ b/src/server.rs @@ -125,6 +125,23 @@ impl RelayServer { } info!("Database backend: {}", config.database_backend); + // Upgrade Git storage before any runtime component can inspect or + // mutate repositories. Migration is resumable and retains original + // repositories as operator-managed rollback backups. + let git_storage = git::storage::LocalGitStorage::new(config.effective_git_data_path()); + let migration = git::migration::migrate_on_startup(&git_storage) + .await + .context("upgrade Git repositories to identifier-family storage")?; + if !migration.already_current { + info!( + migrated_views = migration.migrated_views, + recovered_views = migration.recovered_views, + families_built = migration.families_built, + "Git identifier-family storage migration completed" + ); + } + info!("Git object backend: local identifier families"); + // Initialize metrics if enabled let metrics = if config.metrics_enabled { info!("Metrics enabled on /metrics endpoint"); @@ -422,6 +439,17 @@ impl RelayServer { outbound_credential_keys, )); + // Check the permanent identifier-family model after migration and + // database initialization. The pass is intentionally non-blocking: + // remote repair must not make availability depend on listed clone + // servers. It also owns filesystem-queued operator requests so all + // repair work remains inside the process holding family write leases. + background_tasks.push(git::integrity::spawn_integrity_worker( + git_storage, + sync_ctx.clone(), + )); + info!("Git identifier-family integrity worker started"); + // Create throttle manager for rate limiting remote git servers // Default: 5 concurrent requests per domain, 60 requests per minute per domain let throttle_manager = Arc::new(ThrottleManager::new(5, 60)); diff --git a/tests/git_response_streaming.rs b/tests/git_response_streaming.rs index 5f64ceb..f669493 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,18 @@ 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)); + + // Family ref retention happens behind this boundary. Keep the deadline + // bounded but allow finalization the same scheduling headroom as the + // subprocess wake above. + let terminal = timeout(Duration::from_secs(3), 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 +489,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), diff --git a/tests/lifecycle/nip09_blacklist_ops.rs b/tests/lifecycle/nip09_blacklist_ops.rs index 6365933..3fa8011 100644 --- a/tests/lifecycle/nip09_blacklist_ops.rs +++ b/tests/lifecycle/nip09_blacklist_ops.rs @@ -1,6 +1,7 @@ //! Blacklist operations tests: startup parity, disrespector interaction, //! manual ejection, and operational metrics wiring. +use std::process::Command; use std::sync::Arc; use std::time::Duration; @@ -119,6 +120,15 @@ fn base_config() -> Config { Config::parse_from(["ngit-grasp-test", "--domain", "test.example.com"]) } +fn init_bare_repo(path: &std::path::Path) { + let status = Command::new("git") + .args(["init", "--bare", "--quiet"]) + .arg(path) + .status() + .expect("start git init --bare"); + assert!(status.success(), "initialize bare repository"); +} + #[tokio::test] async fn startup_blacklist_scan_deletes_matching_repositories_via_holding_archive_path() { let relay_dir = tempfile::tempdir().expect("relay tempdir"); @@ -141,7 +151,7 @@ async fn startup_blacklist_scan_deletes_matching_repositories_via_holding_archiv let owner_npub = owner.public_key().to_bech32().expect("npub"); let owner_dir = owner_npub.clone(); let repo_path = git_dir.path().join(&owner_dir).join("blacklisted-repo.git"); - std::fs::create_dir_all(repo_path.join("refs")).expect("create bare repo dir"); + init_bare_repo(&repo_path); let mut config = base_config(); config.repository_blacklist = owner_npub; @@ -394,7 +404,7 @@ async fn startup_whitelist_restore_recovers_now_whitelisted_scope_and_is_idempot .path() .join(owner_npub.clone()) .join("whitelist-restore-repo.git"); - std::fs::create_dir_all(repo_path.join("refs")).expect("create bare repo dir"); + init_bare_repo(&repo_path); let deletion_policy = make_policy( Config { @@ -626,7 +636,7 @@ async fn startup_blacklist_restore_recovers_unblacklisted_scope_and_is_idempoten .path() .join(owner_npub.clone()) .join("blacklist-restore-repo.git"); - std::fs::create_dir_all(repo_path.join("refs")).expect("create bare repo dir"); + init_bare_repo(&repo_path); let deletion_policy = make_policy( Config {