mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Merge #4584cec7: Route Git traffic through identifier families
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsytpxwcu3rtwvqntrk64dnekh6xf9czgv8sedtd8rd95fsugc53xcwp9h3l PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 CoverNote: Completes the local identifier-family storage model after the prerequisite storage-primitives PR was merged. - Routes owner and `/prs/` reads, pushes, and proactive fetches through a shared `(object format, identifier)` object family. - Lets related repositories satisfy reachable SHA wants and advertises retained family base refs, avoiding repeat uploads of objects already stored by the server. - Migrates legacy repositories deterministically on launch while retaining rollback backups and preserving incomplete refs and readable objects. - Adds one permanent family integrity/healing engine for packs, object connectivity, view alternates, and ref targets. It fetches exact missing OIDs from clone URLs in accepted announcements through the existing hardened outbound path, rechecks the family, and logs unresolved damage at `ERROR`. - Runs that engine asynchronously after migration and exposes `ngit-grasp integrity-check --identifier <id> [--repair]` through a durable live-process request queue. Migration does not get a separate recovery subsystem: it performs the structural conversion, then hands the resulting family to the ordinary steady-state checker. Unindexed legacy packs remain in the rollback backup. Garbage collection, legacy backup archaeology, and S3 storage remain out of scope. Testing on gitnostr.com: the already-installed storage version makes structural migration a no-op, but the startup integrity pass still runs unconditionally, so this is a valid test of the permanent steady-state path. To prove remote self-healing, use a sacrificial identifier whose accepted announcement lists a second Git server containing the same reachable object; snapshot its family and views, move one verified loose object into quarantine, invoke `integrity-check --repair` or restart, and verify the repair log, restored object, `git fsck`, and a fresh clone. This does not re-test the first legacy-to-family transition; that transition should remain covered by the migration fixtures or a disposable pre-migration data copy. Do not remove the production migration marker to force a rerun.
This commit is contained in:
@@ -713,8 +713,8 @@ Optional endpoint at `/prs/<npub>/<identifier>.git`, gated on `NGIT_GRASP06_ENAB
|
||||
|
||||
- [`src/grasp06/endpoint.rs`](../../src/grasp06/endpoint.rs) — URL parsing.
|
||||
- [`src/grasp06/paths.rs`](../../src/grasp06/paths.rs) — on-disk path conventions under `<git_data_path>/prs/<hex>/<identifier>.git`.
|
||||
- [`src/grasp06/fetch.rs`](../../src/grasp06/fetch.rs) — empty-repo synthesis for `info/refs` and `git-upload-pack` against repos that don't yet exist on disk.
|
||||
- [`src/grasp06/receive.rs`](../../src/grasp06/receive.rs) — `git-receive-pack` with init-on-push, strict `refs/nostr/<event-id>` ref-name validation, and per-ref post-push validation against the database and purgatory. Uses a per-`(submitter, identifier)` `PrsPathState` (mutex + `in_flight` counter) so concurrent pushes to the same path proceed in parallel: the mutex is held only for init+register and decrement+end-of-push cleanup, not across `git-receive-pack` itself. The same state is consulted (with `try_lock` in synchronous contexts) by the PR-event policy and the purgatory expiry sweep, which only remove a `/prs/` ref or zero-ref bare repo when `in_flight == 0`.
|
||||
- [`src/grasp06/fetch.rs`](../../src/grasp06/fetch.rs) — empty thin-view synthesis for `info/refs` and `git-upload-pack` against routes that do not yet exist on disk. Named refs stay empty while an existing identifier family supplies anonymous receive negotiation bases.
|
||||
- [`src/grasp06/receive.rs`](../../src/grasp06/receive.rs) — `git-receive-pack` with thin-view init-on-push, strict `refs/nostr/<event-id>` ref-name validation, and per-ref post-push validation against the database and purgatory. A per-`(submitter, identifier)` `PrsPathState` protects view creation/removal; a per-family lease serializes object-producing work across related views. The same path state is consulted by PR-event policy and purgatory expiry, which only remove a `/prs/` ref or zero-ref view when `in_flight == 0`.
|
||||
- [`src/grasp06/policy.rs`](../../src/grasp06/policy.rs) — strict clone-tag URL comparator used by the PR-event acceptance relaxation.
|
||||
- [`src/grasp06/cleanup.rs`](../../src/grasp06/cleanup.rs) — one-shot startup scan over `<git_data_path>/prs/` that removes zero-ref bare repos left behind by a previous run (crash mid-push, crash mid-cleanup, or shutdown with unresolved scoped placeholders). Runs before the HTTP server starts accepting requests; no locking is needed because nothing else is touching `/prs/` yet.
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -55,7 +55,9 @@ Two reasons:
|
||||
|
||||
The contributor chose to publish to `/prs/`. Mirroring into any accepted repository on this relay makes the PR visible at the expected location for clients browsing that repo. Not mirroring the reverse direction preserves the maintainer's declared `clone` intent on their own pushes — we do not invent new hosting locations for their events.
|
||||
|
||||
Until object-pool dedup lands, the mirror doubles storage for affected refs. We accept this cost as the simplest correct implementation; dedup lands later and makes the mirror effectively free.
|
||||
Identifier-family object storage makes the mirror a ref-only operation. Owner
|
||||
and contributor views for the same identifier resolve one shared object
|
||||
inventory, so accepting or mirroring a PR does not copy its object graph.
|
||||
|
||||
## Architecture
|
||||
|
||||
@@ -94,7 +96,7 @@ POST /prs/<npub>/<id>.git/git-receive-pack
|
||||
signer ≠ URL npub → reject ref
|
||||
identifier (d-tag) ≠ URL identifier → reject ref
|
||||
commit ≠ event's c tag → delete ref
|
||||
all match → ref locked; release event from purgatory;
|
||||
all match → ref locked; retain family tip; release event from purgatory;
|
||||
mirror to matching standard repos
|
||||
event not yet seen:
|
||||
accept ref; create or update PR placeholder in purgatory
|
||||
@@ -107,7 +109,14 @@ The flow mirrors the existing `refs/nostr/<event-id>` path at the standard endpo
|
||||
|
||||
#### On-demand bare repo creation
|
||||
|
||||
The first push to `/prs/<submitter>/<identifier>.git` creates the bare repo on disk. A per-`(submitter, identifier)` [`PrsPathState`](../../src/grasp06/receive.rs) (kept in a `DashMap` on the `HttpService`, see [`crate::grasp06::receive::RepoInitLocks`](../../src/grasp06/receive.rs)) provides a `tokio::sync::Mutex` and an `in_flight: AtomicUsize` counter. The mutex is held only briefly — for `git init --bare` plus the `in_flight` increment at the start of a push, and for the `in_flight` decrement plus end-of-push zero-ref cleanup at the end. The pack upload and per-ref validation run *without* the mutex held, so concurrent pushes to the same path proceed in parallel; git's own ref locking handles intra-push concurrency. Off-push cleanup paths (PR-event policy, purgatory expiry) take the same mutex briefly and only `rm -rf` the bare repo when they see `in_flight == 0` and `list_refs` is empty — so no off-push code path can delete the bare repo while a push is in flight, and no push is serialised behind another push to the same identity.
|
||||
The first push to `/prs/<submitter>/<identifier>.git` creates a thin bare view
|
||||
whose alternate is the identifier family. A per-`(submitter, identifier)`
|
||||
[`PrsPathState`](../../src/grasp06/receive.rs) still protects view creation and
|
||||
cleanup. A separate per-family write lease serializes object-producing work
|
||||
across all owner and contributor views for that identifier. `git-receive-pack`
|
||||
writes objects to the family while updating only the selected view's refs.
|
||||
Off-push cleanup may remove an empty contributor view, but never removes family
|
||||
objects or retained roots.
|
||||
|
||||
### Event acceptance relaxation
|
||||
|
||||
@@ -130,9 +139,10 @@ When a PR or PR Update's purgatory entry is released via a `/prs/` push:
|
||||
1. Save event to DB and remove from purgatory (as today).
|
||||
2. For each `a` tag in the event of the form `30617:<pubkey>:<d-tag>`:
|
||||
- Resolve to a local repo path `<git_data_path>/<pubkey-npub>/<d-tag>.git`.
|
||||
- If that repo has an active (non-purgatory) announcement, copy objects + install `refs/nostr/<event-id>` into it (same mechanism as existing cross-owner sync in [`src/git/sync.rs`](../../src/git/sync.rs)).
|
||||
- If that repo has an active (non-purgatory) announcement, install `refs/nostr/<event-id>` in its view after checking the commit resolves through the shared family.
|
||||
|
||||
The mirror copies the same ref, same commits. No separate object store. Dedup can be added transparently later via git alternates keyed on d-tag.
|
||||
The mirror installs the same ref to the same commit. Git alternates keyed by
|
||||
identifier make the objects immediately available without a copy.
|
||||
|
||||
The mirror is **one-directional**: pushes to `/<maintainer>/<id>.git` are not mirrored into `/prs/*`. Only the `/prs/` → `<maintainer>/` direction fires, and only when the source repo path is under `prs_base_path`.
|
||||
|
||||
@@ -250,9 +260,10 @@ This is the design rationale for the mirror being one-directional and the `/prs/
|
||||
|
||||
Deletion of the contributor's own PR event (NIP-09) → the deletion-request branch will decide and implement the ref-lifecycle behaviour. GRASP-06 v1 does not add hooks or no-op scaffolding for this.
|
||||
|
||||
### Object-pool deduplication (not yet specified)
|
||||
### Object-pool deduplication
|
||||
|
||||
The mirror copies objects today. When object-pool dedup exists — likely as git alternates keyed on `(identifier, a-tag-coord-set)` — the mirror becomes a pointer operation with near-zero storage cost. No URL or spec change is required.
|
||||
The identifier family is the shared object inventory. Mirroring installs a view
|
||||
ref and does not copy objects; no URL or spec change is required.
|
||||
|
||||
## Anticipated failure modes and mitigations
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -283,6 +283,9 @@ git ls-remote https://ngit.example.com/<npub>/<repo>.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
|
||||
|
||||
@@ -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
|
||||
<git-data>/.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.
|
||||
+84
-5
@@ -17,6 +17,7 @@ use tokio::time::MissedTickBehavior;
|
||||
use tracing::{debug, error, info, warn};
|
||||
|
||||
use super::protocol::{GitService, PktLine};
|
||||
use super::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
|
||||
use super::subprocess::GitSubprocess;
|
||||
use super::{full_body, GitResponseBody};
|
||||
|
||||
@@ -343,7 +344,7 @@ where
|
||||
/// caller finish GRASP post-push processing before making success visible,
|
||||
/// without buffering the progress stream that keeps clients alive during
|
||||
/// expensive pack processing.
|
||||
async fn pump_receive_pack_stdout_to_channel<R>(
|
||||
pub(crate) async fn pump_receive_pack_stdout_to_channel<R>(
|
||||
mut stdout: R,
|
||||
tx: &mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
|
||||
) -> (PumpResult, Option<Vec<u8>>)
|
||||
@@ -701,9 +702,34 @@ pub async fn handle_receive_pack(
|
||||
}
|
||||
};
|
||||
|
||||
// Spawn git receive-pack
|
||||
let mut git = GitSubprocess::spawn(GitService::ReceivePack, &repo_path, false, git_protocol)
|
||||
.map_err(GitError::ProcessSpawnFailed)?;
|
||||
let storage = LocalGitStorage::new(git_data_path);
|
||||
let family_key = FamilyKey::sha1(identifier).map_err(|e| GitError::Storage(e.to_string()))?;
|
||||
let family_lease = if storage.is_thin_view(&family_key, &repo_path) {
|
||||
Some(
|
||||
storage
|
||||
.write_lease(&family_key)
|
||||
.await
|
||||
.map_err(|e| GitError::Storage(e.to_string()))?,
|
||||
)
|
||||
} else {
|
||||
// Legacy repositories remain self-contained until startup migration.
|
||||
// This fallback also keeps direct library callers safe: never point a
|
||||
// ref at a family object database the view cannot read.
|
||||
None
|
||||
};
|
||||
|
||||
// Ref updates remain in the selected view; receive-pack's quarantine and
|
||||
// final objects are installed directly in the shared family inventory.
|
||||
let mut git = GitSubprocess::spawn_with_object_directory(
|
||||
GitService::ReceivePack,
|
||||
&repo_path,
|
||||
false,
|
||||
git_protocol,
|
||||
family_lease
|
||||
.as_ref()
|
||||
.map(|lease| lease.family_objects_path.as_path()),
|
||||
)
|
||||
.map_err(GitError::ProcessSpawnFailed)?;
|
||||
|
||||
// Write request to git's stdin
|
||||
if let Some(mut stdin) = git.take_stdin() {
|
||||
@@ -754,6 +780,10 @@ pub async fn handle_receive_pack(
|
||||
repo_lifecycle_guard,
|
||||
promotion_hooks,
|
||||
metrics,
|
||||
storage,
|
||||
family_key,
|
||||
family_lease,
|
||||
pushed_refs,
|
||||
)
|
||||
.await;
|
||||
});
|
||||
@@ -778,6 +808,10 @@ async fn stream_receive_pack_output<S, E>(
|
||||
repo_lifecycle_guard: Option<LifecycleReadGuard>,
|
||||
promotion_hooks: Option<Arc<dyn PurgatoryPromotionHooks>>,
|
||||
metrics: Option<Arc<Metrics>>,
|
||||
storage: LocalGitStorage,
|
||||
family_key: FamilyKey,
|
||||
family_lease: Option<FamilyWriteLease>,
|
||||
pushed_refs: Vec<(String, String, String)>,
|
||||
) where
|
||||
S: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
E: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
@@ -863,6 +897,15 @@ async fn stream_receive_pack_output<S, E>(
|
||||
|
||||
debug!("Git receive-pack stream completed successfully");
|
||||
|
||||
if family_lease.is_some() {
|
||||
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
|
||||
}
|
||||
|
||||
// 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<S, E>(
|
||||
record_git_operation(&metrics, "push", "success");
|
||||
}
|
||||
|
||||
async fn send_body_bytes(
|
||||
pub(crate) async fn send_body_bytes(
|
||||
tx: &mpsc::Sender<Result<Frame<Bytes>, io::Error>>,
|
||||
bytes: Vec<u8>,
|
||||
) -> Result<(), mpsc::error::SendError<Result<Frame<Bytes>, io::Error>>> {
|
||||
@@ -966,6 +1009,7 @@ pub enum GitError {
|
||||
ProcessSpawnFailed(std::io::Error),
|
||||
IoError(std::io::Error),
|
||||
GitFailed(Option<i32>),
|
||||
Storage(String),
|
||||
}
|
||||
|
||||
impl std::fmt::Display for GitError {
|
||||
@@ -975,6 +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"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<PathBuf> {
|
||||
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<PathBuf>,
|
||||
pub refs_checked: usize,
|
||||
pub missing_oids: BTreeSet<String>,
|
||||
pub missing_ref_targets: Vec<MissingRefTarget>,
|
||||
pub invalid_alternates: Vec<PathBuf>,
|
||||
pub pack_errors: Vec<String>,
|
||||
pub fsck_diagnostics: Vec<String>,
|
||||
}
|
||||
|
||||
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<Vec<String>>;
|
||||
|
||||
async fn fetch_missing_oids(
|
||||
&self,
|
||||
target_view: &Path,
|
||||
url: &str,
|
||||
oids: &[String],
|
||||
) -> Result<Vec<String>>;
|
||||
}
|
||||
|
||||
/// 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<String>,
|
||||
pub source_failures: Vec<String>,
|
||||
}
|
||||
|
||||
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<crate::purgatory::sync::RealSyncContext>,
|
||||
) -> 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<S: FamilyRepairSource + ?Sized>(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<S: FamilyRepairSource + ?Sized>(
|
||||
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::<Vec<_>>(),
|
||||
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<S: FamilyRepairSource + ?Sized>(
|
||||
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::<Vec<_>>()
|
||||
.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<Vec<(PathBuf, IntegrityRequest)>> {
|
||||
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::<Vec<_>>();
|
||||
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<S: FamilyRepairSource + ?Sized>(
|
||||
storage: &LocalGitStorage,
|
||||
key: &FamilyKey,
|
||||
source: &S,
|
||||
repair: bool,
|
||||
) -> Result<FamilyRepairOutcome> {
|
||||
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<Vec<String>> {
|
||||
self.accepted_repository_clone_urls(identifier).await
|
||||
}
|
||||
|
||||
async fn fetch_missing_oids(
|
||||
&self,
|
||||
target_view: &Path,
|
||||
url: &str,
|
||||
oids: &[String],
|
||||
) -> Result<Vec<String>> {
|
||||
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<Vec<FamilyKey>> {
|
||||
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<FamilyIntegrityReport> {
|
||||
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<Vec<PathBuf>> {
|
||||
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<PathBuf>) -> 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<PathBuf>) -> 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<Vec<(String, String)>> {
|
||||
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<Vec<String>> {
|
||||
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<String>, Vec<String>)> {
|
||||
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<bool> {
|
||||
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<Vec<String>> {
|
||||
Ok(vec![self.url.clone()])
|
||||
}
|
||||
|
||||
async fn fetch_missing_oids(
|
||||
&self,
|
||||
target_view: &Path,
|
||||
url: &str,
|
||||
oids: &[String],
|
||||
) -> Result<Vec<String>> {
|
||||
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 <test@example.com> 0 +0000\ncommitter Test <test@example.com> 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);
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
@@ -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;
|
||||
|
||||
+31
-9
@@ -8,7 +8,7 @@
|
||||
use std::fmt;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Command;
|
||||
use std::sync::Arc;
|
||||
use std::sync::{Arc, LazyLock};
|
||||
|
||||
use anyhow::{anyhow, Context, Result};
|
||||
use bitcoin_hashes::{sha256, Hash};
|
||||
@@ -26,6 +26,9 @@ const FAMILY_DIR: &str = "families";
|
||||
const BASE_REFS_PREFIX: &str = "refs/grasp/bases/";
|
||||
const RETAINED_REFS_PREFIX: &str = "refs/grasp/retained/";
|
||||
|
||||
type FamilyLockId = (PathBuf, FamilyKey);
|
||||
static FAMILY_LOCKS: LazyLock<DashMap<FamilyLockId, Arc<Mutex<()>>>> = LazyLock::new(DashMap::new);
|
||||
|
||||
/// Git object hash algorithms are separate family namespaces.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
|
||||
pub enum ObjectFormat {
|
||||
@@ -100,15 +103,19 @@ impl FamilyKey {
|
||||
#[derive(Clone)]
|
||||
pub struct LocalGitStorage {
|
||||
git_data_path: PathBuf,
|
||||
family_locks: Arc<DashMap<FamilyKey, Arc<Mutex<()>>>>,
|
||||
}
|
||||
|
||||
impl LocalGitStorage {
|
||||
pub fn new(git_data_path: impl Into<PathBuf>) -> Self {
|
||||
Self {
|
||||
git_data_path: git_data_path.into(),
|
||||
family_locks: Arc::new(DashMap::new()),
|
||||
}
|
||||
let git_data_path = git_data_path.into();
|
||||
let git_data_path = if git_data_path.is_absolute() {
|
||||
git_data_path
|
||||
} else {
|
||||
std::env::current_dir()
|
||||
.unwrap_or_else(|_| PathBuf::from("."))
|
||||
.join(git_data_path)
|
||||
};
|
||||
Self { git_data_path }
|
||||
}
|
||||
|
||||
pub fn git_data_path(&self) -> &Path {
|
||||
@@ -134,11 +141,26 @@ impl LocalGitStorage {
|
||||
self.family_repo_path(key).join("objects")
|
||||
}
|
||||
|
||||
pub fn family_exists(&self, key: &FamilyKey) -> bool {
|
||||
self.family_repo_path(key).is_dir()
|
||||
}
|
||||
|
||||
pub fn is_thin_view(&self, key: &FamilyKey, view_path: &Path) -> bool {
|
||||
let alternate = view_path.join("objects/info/alternates");
|
||||
let Ok(configured) = std::fs::read_to_string(alternate) else {
|
||||
return false;
|
||||
};
|
||||
let Ok(expected) = std::fs::canonicalize(self.family_objects_path(key)) else {
|
||||
return false;
|
||||
};
|
||||
configured.lines().any(|line| Path::new(line) == expected)
|
||||
}
|
||||
|
||||
/// Serialize mutations to one identifier family.
|
||||
pub async fn write_lease(&self, key: &FamilyKey) -> Result<FamilyWriteLease> {
|
||||
let lock = self
|
||||
.family_locks
|
||||
.entry(key.clone())
|
||||
let lock_id = (self.git_data_path.clone(), key.clone());
|
||||
let lock = FAMILY_LOCKS
|
||||
.entry(lock_id)
|
||||
.or_insert_with(|| Arc::new(Mutex::new(())))
|
||||
.clone();
|
||||
let guard = lock.lock_owned().await;
|
||||
|
||||
@@ -28,6 +28,27 @@ impl GitSubprocess {
|
||||
repo_path: impl AsRef<Path>,
|
||||
advertise: bool,
|
||||
git_protocol: Option<&str>,
|
||||
) -> std::io::Result<Self> {
|
||||
Self::spawn_with_object_directory(
|
||||
service,
|
||||
repo_path,
|
||||
advertise,
|
||||
git_protocol,
|
||||
None::<&Path>,
|
||||
)
|
||||
}
|
||||
|
||||
/// Spawn Git with an optional writable object database override.
|
||||
///
|
||||
/// Thin views normally read through `objects/info/alternates`. Commands
|
||||
/// that create objects use this override so receive-pack quarantine output
|
||||
/// and fetched objects land directly in the identifier family.
|
||||
pub fn spawn_with_object_directory(
|
||||
service: GitService,
|
||||
repo_path: impl AsRef<Path>,
|
||||
advertise: bool,
|
||||
git_protocol: Option<&str>,
|
||||
object_directory: Option<impl AsRef<Path>>,
|
||||
) -> std::io::Result<Self> {
|
||||
let repo_path = repo_path.as_ref();
|
||||
|
||||
@@ -62,6 +83,9 @@ impl GitSubprocess {
|
||||
if let Some(protocol) = git_protocol {
|
||||
cmd.env("GIT_PROTOCOL", protocol);
|
||||
}
|
||||
if let Some(object_directory) = object_directory {
|
||||
cmd.env("GIT_OBJECT_DIRECTORY", object_directory.as_ref());
|
||||
}
|
||||
|
||||
let child = cmd.spawn()?;
|
||||
|
||||
|
||||
+22
-8
@@ -5,10 +5,11 @@
|
||||
//! > MUST respond to upload-pack requests for any well-formed path as if
|
||||
//! > serving an empty bare repository.
|
||||
//!
|
||||
//! When a real `/prs/<submitter>/<identifier>.git` repo exists on disk
|
||||
//! we delegate to the standard handlers in [`crate::git::handlers`].
|
||||
//! Otherwise we synthesise a brand-new empty bare repo in a per-request
|
||||
//! temporary directory and run the standard upload-pack against that.
|
||||
//! When a real `/prs/<submitter>/<identifier>.git` repo exists on disk we
|
||||
//! delegate to the standard handlers in [`crate::git::handlers`]. Otherwise we
|
||||
//! synthesize an empty thin view in a per-request temporary directory. If the
|
||||
//! identifier family exists, receive-pack discovery can advertise its bases as
|
||||
//! anonymous `.have` entries without exposing any named refs.
|
||||
//!
|
||||
//! Receive-pack lives in [`crate::grasp06::receive`].
|
||||
|
||||
@@ -22,6 +23,7 @@ use tracing::{debug, warn};
|
||||
|
||||
use crate::git::handlers::{handle_info_refs, handle_upload_pack, GitError};
|
||||
use crate::git::protocol::GitService;
|
||||
use crate::git::storage::{FamilyKey, LocalGitStorage};
|
||||
use crate::git::GitResponseBody;
|
||||
use crate::grasp06::endpoint::PrsUrl;
|
||||
use crate::grasp06::paths::prs_repo_path;
|
||||
@@ -60,7 +62,9 @@ pub async fn handle_prs_info_refs(
|
||||
prs.submitter.to_hex(),
|
||||
prs.identifier
|
||||
);
|
||||
let temp = init_empty_bare_repo()?;
|
||||
let storage = LocalGitStorage::new(git_data_path);
|
||||
let key = FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
|
||||
let temp = init_empty_bare_repo(&storage, &key)?;
|
||||
let repo_path = temp.path().to_path_buf();
|
||||
let response = handle_info_refs(repo_path, service, git_protocol).await;
|
||||
// `temp` is dropped here, deleting the directory. Git has already
|
||||
@@ -98,7 +102,9 @@ pub async fn handle_prs_upload_pack(
|
||||
prs.submitter.to_hex(),
|
||||
prs.identifier
|
||||
);
|
||||
let temp = init_empty_bare_repo()?;
|
||||
let storage = LocalGitStorage::new(git_data_path);
|
||||
let key = FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
|
||||
let temp = init_empty_bare_repo(&storage, &key)?;
|
||||
handle_prs_upload_pack_buffered(temp, body, git_protocol).await
|
||||
}
|
||||
|
||||
@@ -169,7 +175,7 @@ async fn handle_prs_upload_pack_buffered(
|
||||
.unwrap())
|
||||
}
|
||||
|
||||
/// Create a fresh empty bare repo in a temp directory.
|
||||
/// Create a fresh empty thin view in a temp directory.
|
||||
///
|
||||
/// Per-request temp dirs are deliberate. They are cheap on
|
||||
/// any reasonable filesystem (one `mkdir`, one `git init --bare`) and
|
||||
@@ -177,7 +183,10 @@ async fn handle_prs_upload_pack_buffered(
|
||||
/// profiling shows this is a bottleneck we can switch to a shared
|
||||
/// `<git_data_path>/prs/.empty-template.git`; do not optimise until
|
||||
/// measurable.
|
||||
fn init_empty_bare_repo() -> Result<TempDir, GitError> {
|
||||
fn init_empty_bare_repo(
|
||||
storage: &LocalGitStorage,
|
||||
family_key: &FamilyKey,
|
||||
) -> Result<TempDir, GitError> {
|
||||
let temp = TempDir::new().map_err(GitError::IoError)?;
|
||||
let path = temp.path();
|
||||
let output = Command::new("git")
|
||||
@@ -194,5 +203,10 @@ fn init_empty_bare_repo() -> Result<TempDir, GitError> {
|
||||
);
|
||||
return Err(GitError::GitFailed(output.status.code()));
|
||||
}
|
||||
if storage.family_exists(family_key) {
|
||||
storage
|
||||
.configure_thin_view(family_key, path)
|
||||
.map_err(|error| GitError::Storage(error.to_string()))?;
|
||||
}
|
||||
Ok(temp)
|
||||
}
|
||||
|
||||
+82
-53
@@ -19,19 +19,21 @@
|
||||
//! reject the whole push with an `ERR` pkt-line — matching the
|
||||
//! standard-endpoint UX. Nothing on disk has been touched yet so failed
|
||||
//! probes leave no state.
|
||||
//! 3. Acquires the per-path coordination state from [`RepoInitLocks`]
|
||||
//! briefly: under its mutex it runs the on-demand `git init --bare`
|
||||
//! 3. Acquires the identifier-family write lease and the per-path coordination
|
||||
//! state from [`RepoInitLocks`] briefly: under the path mutex it creates an
|
||||
//! on-demand thin view
|
||||
//! and increments the `in_flight` counter, then releases the mutex.
|
||||
//! Steps 4 and 5 run *without* the per-path lock so concurrent pushes
|
||||
//! to the same `(submitter, identifier)` proceed in parallel — git's
|
||||
//! own ref locking handles intra-push concurrency, and the
|
||||
//! `in_flight` counter is what off-push cleanup paths consult to know
|
||||
//! a push is active.
|
||||
//! 4. Starts `git-receive-pack`, writes the full request body to the child, and
|
||||
//! Steps 4 and 5 run *without* the per-path lock. The family write lease
|
||||
//! serializes object-producing operations for the identifier, while the
|
||||
//! `in_flight` counter is what off-push cleanup paths consult to know a push
|
||||
//! is active.
|
||||
//! 4. Starts `git-receive-pack` with the family as its writable object database,
|
||||
//! writes the full request body to the child, and
|
||||
//! immediately returns an HTTP response whose body is backed by a bounded
|
||||
//! channel. From this point on Hyper can stream stdout to the client while a
|
||||
//! detached task owns the subprocess, stderr, and all post-push work.
|
||||
//! 5. In the detached task, each stdout chunk is forwarded as it is read. If Git
|
||||
//! 5. In the detached task, progress is forwarded while the terminal flush is
|
||||
//! retained until family roots and GRASP post-processing are complete. If Git
|
||||
//! exits with a protocol-level error before writing stdout, the task sends a
|
||||
//! Git `ERR` pkt-line through the same stream so clients still see a normal
|
||||
//! receive-pack failure. If stdout has already been sent, the task cannot
|
||||
@@ -76,10 +78,12 @@ use crate::git::authorization::{
|
||||
};
|
||||
use crate::git::handlers::{
|
||||
build_git_protocol_error_response, err_pktline_frame, is_git_protocol_error,
|
||||
pump_stdout_to_channel, read_stderr_to_end, record_git_operation, streaming_response, GitError,
|
||||
PumpResult, STREAM_CHANNEL_DEPTH,
|
||||
pump_receive_pack_stdout_to_channel, read_stderr_to_end, record_git_operation,
|
||||
retain_accepted_tips, send_body_bytes, streaming_response, GitError, PumpResult,
|
||||
STREAM_CHANNEL_DEPTH,
|
||||
};
|
||||
use crate::git::protocol::GitService;
|
||||
use crate::git::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
|
||||
use crate::git::subprocess::GitSubprocess;
|
||||
use crate::git::sync::process_newly_available_git_data;
|
||||
use crate::git::{delete_ref, list_refs, GitResponseBody};
|
||||
@@ -107,9 +111,9 @@ use crate::sync::rejected_index::RejectedEventsIndex;
|
||||
/// expiry) for the duration of one `delete_ref` + optional
|
||||
/// `remove_dir_all`.
|
||||
///
|
||||
/// `git-receive-pack` itself and per-ref validation run *without* the
|
||||
/// mutex held, so two pushes to the same path proceed in parallel; git's
|
||||
/// own ref locking handles intra-push concurrency.
|
||||
/// `git-receive-pack` itself and per-ref validation run *without* the path
|
||||
/// mutex held. A separate identifier-family lease serializes their shared
|
||||
/// object inventory.
|
||||
///
|
||||
/// Off-push cleanup paths only `rm -rf` the bare repo when both
|
||||
/// `in_flight.load() == 0` *and* `list_refs` returns empty while they
|
||||
@@ -258,20 +262,27 @@ pub async fn handle_prs_receive_pack(
|
||||
// 3. Acquire the per-path coordination state and, under its mutex,
|
||||
// initialise the bare repo on demand and register this request as
|
||||
// in-flight. The mutex is then released — `git-receive-pack` and
|
||||
// per-ref validation run WITHOUT the lock so concurrent pushes to
|
||||
// the same `(submitter, identifier)` proceed in parallel. Cleanup
|
||||
// paths consult `in_flight` (under the same mutex) before
|
||||
// per-ref validation run WITHOUT the path lock. The family lease above
|
||||
// serializes object-producing work for the identifier. Cleanup paths
|
||||
// consult `in_flight` (under the same mutex) before
|
||||
// deleting the bare repo, so a repo can never vanish mid-receive.
|
||||
let repo_path = prs_repo_path(
|
||||
Path::new(git_data_path),
|
||||
&prs.submitter.to_hex(),
|
||||
&prs.identifier,
|
||||
);
|
||||
let storage = LocalGitStorage::new(git_data_path);
|
||||
let family_key =
|
||||
FamilyKey::sha1(&prs.identifier).map_err(|e| GitError::Storage(e.to_string()))?;
|
||||
let family_lease = storage
|
||||
.write_lease(&family_key)
|
||||
.await
|
||||
.map_err(|e| GitError::Storage(e.to_string()))?;
|
||||
let state = path_state(&repo_init_locks, &repo_path);
|
||||
|
||||
{
|
||||
let _g = state.mu.lock().expect("prs path mutex poisoned");
|
||||
if let Err(e) = ensure_repo_initialised(&repo_path) {
|
||||
if let Err(e) = ensure_repo_initialised(&storage, &family_key, &repo_path) {
|
||||
error!(
|
||||
"/prs/ receive-pack: failed to initialise repo at {}: {}",
|
||||
repo_path.display(),
|
||||
@@ -289,17 +300,23 @@ pub async fn handle_prs_receive_pack(
|
||||
// to the same streaming shape as the standard receive-pack endpoint:
|
||||
// write stdin here, then hand stdout/stderr and all follow-up state to a
|
||||
// detached task that feeds the response body channel.
|
||||
let mut git =
|
||||
match GitSubprocess::spawn(GitService::ReceivePack, &repo_path, false, git_protocol)
|
||||
.map_err(GitError::ProcessSpawnFailed)
|
||||
{
|
||||
Ok(git) => git,
|
||||
Err(e) => {
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
return Err(e);
|
||||
}
|
||||
};
|
||||
let use_family = storage.is_thin_view(&family_key, &repo_path);
|
||||
let mut git = match GitSubprocess::spawn_with_object_directory(
|
||||
GitService::ReceivePack,
|
||||
&repo_path,
|
||||
false,
|
||||
git_protocol,
|
||||
use_family.then_some(&family_lease.family_objects_path),
|
||||
)
|
||||
.map_err(GitError::ProcessSpawnFailed)
|
||||
{
|
||||
Ok(git) => git,
|
||||
Err(e) => {
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
return Err(e);
|
||||
}
|
||||
};
|
||||
|
||||
if let Some(mut stdin) = git.take_stdin() {
|
||||
if let Err(e) = stdin.write_all(&request_body).await {
|
||||
@@ -350,6 +367,9 @@ pub async fn handle_prs_receive_pack(
|
||||
prs.identifier.clone(),
|
||||
domain.to_string(),
|
||||
metrics,
|
||||
storage,
|
||||
family_key,
|
||||
use_family.then_some(family_lease),
|
||||
));
|
||||
|
||||
Ok(streaming_response(GitService::ReceivePack, rx))
|
||||
@@ -389,30 +409,18 @@ fn invalid_ref_reason(ref_name: &str) -> Option<String> {
|
||||
///
|
||||
/// The caller must hold the per-path mutex from [`PrsPathState`] for
|
||||
/// `repo_path` before invoking this function.
|
||||
fn ensure_repo_initialised(repo_path: &Path) -> Result<(), GitError> {
|
||||
fn ensure_repo_initialised(
|
||||
storage: &LocalGitStorage,
|
||||
family_key: &FamilyKey,
|
||||
repo_path: &Path,
|
||||
) -> Result<(), GitError> {
|
||||
if repo_path.exists() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if let Some(parent) = repo_path.parent() {
|
||||
std::fs::create_dir_all(parent).map_err(GitError::IoError)?;
|
||||
}
|
||||
|
||||
let output = std::process::Command::new("git")
|
||||
.args(["init", "--bare", "--initial-branch=main", "--quiet"])
|
||||
.arg(repo_path)
|
||||
.output()
|
||||
.map_err(GitError::ProcessSpawnFailed)?;
|
||||
|
||||
if !output.status.success() {
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
error!(
|
||||
"/prs/ git init --bare failed at {}: {}",
|
||||
repo_path.display(),
|
||||
stderr.trim()
|
||||
);
|
||||
return Err(GitError::GitFailed(output.status.code()));
|
||||
}
|
||||
storage
|
||||
.create_thin_view(family_key, repo_path)
|
||||
.map_err(|error| GitError::Storage(error.to_string()))?;
|
||||
|
||||
info!(
|
||||
"/prs/ initialised bare repo at {} on demand",
|
||||
@@ -482,6 +490,9 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
identifier: String,
|
||||
domain: String,
|
||||
metrics: Option<Arc<Metrics>>,
|
||||
storage: LocalGitStorage,
|
||||
family_key: FamilyKey,
|
||||
family_lease: Option<FamilyWriteLease>,
|
||||
) where
|
||||
S: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
E: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
@@ -490,7 +501,7 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
// fill its stderr pipe while still producing stdout; draining both prevents
|
||||
// child-process deadlock and preserves stderr for protocol-error reporting.
|
||||
let stderr_task = stderr.map(|stderr| tokio::spawn(read_stderr_to_end(stderr)));
|
||||
let pump_result = pump_stdout_to_channel(stdout, &tx).await;
|
||||
let (pump_result, terminal_flush) = pump_receive_pack_stdout_to_channel(stdout, &tx).await;
|
||||
|
||||
// If the client goes away or stdout read fails, stop Git rather than letting
|
||||
// it continue writing into a response nobody can receive. Cleanup below will
|
||||
@@ -514,7 +525,7 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
None => Vec::new(),
|
||||
};
|
||||
|
||||
let sent_stdout = match pump_result {
|
||||
let mut sent_stdout = match pump_result {
|
||||
PumpResult::Eof { sent_stdout } => sent_stdout,
|
||||
PumpResult::ClientDisconnected | PumpResult::ReadError => {
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
@@ -524,6 +535,14 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
};
|
||||
|
||||
if !status.success() {
|
||||
if let Some(flush) = terminal_flush {
|
||||
sent_stdout = true;
|
||||
if send_body_bytes(&tx, flush).await.is_err() {
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
return;
|
||||
}
|
||||
}
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
let stderr_str = String::from_utf8_lossy(&stderr_output);
|
||||
if is_git_protocol_error(status.code(), &stderr_output) {
|
||||
@@ -568,8 +587,6 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
}
|
||||
|
||||
debug!("/prs/ git-receive-pack stream completed successfully");
|
||||
record_git_operation(&metrics, "push", "success");
|
||||
|
||||
// Race safety net. The pre-validation in `handle_prs_receive_pack` was
|
||||
// performed before `git-receive-pack` ran, so an event with one of the
|
||||
// pushed ids may have arrived via WebSocket during the receive-pack window.
|
||||
@@ -595,6 +612,10 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
.await;
|
||||
}
|
||||
|
||||
if family_lease.is_some() {
|
||||
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
|
||||
}
|
||||
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
|
||||
// Drive the standard purgatory-release pipeline so PR events already
|
||||
@@ -633,6 +654,14 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(flush) = terminal_flush {
|
||||
if send_body_bytes(&tx, flush).await.is_err() {
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
return;
|
||||
}
|
||||
}
|
||||
record_git_operation(&metrics, "push", "success");
|
||||
}
|
||||
|
||||
/// Race safety net for the `/prs/` receive-pack post-push phase.
|
||||
|
||||
+24
-1
@@ -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<String> = 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()?);
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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<bool> {
|
||||
let result: anyhow::Result<bool> = 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;
|
||||
}
|
||||
|
||||
@@ -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(())
|
||||
}
|
||||
|
||||
|
||||
+55
-41
@@ -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;
|
||||
|
||||
@@ -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<NaughtyListTracker> {
|
||||
&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<Vec<String>> {
|
||||
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<FamilyKey> {
|
||||
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()]);
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user