mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(git): stage unsigned uploads until a signed event accepts them
A push to refs/nostr/<event-id> is accepted before its PR event is known, and its objects went straight into the identifier family. The family is never garbage-collected, so anyone could fill permanent storage without signing anything, and an upload whose event never arrived stayed forever. Approach Route by what push authorization already decides. A ref named by a signed State or PR event, accepted or in purgatory, is received into the family as before. A refs/nostr ref with no event, or only a placeholder, is received into the view's own object directory. Pre-validation now reports whether a match was signed, because a placeholder match was indistinguishable from a stored event. Accepting a PR event first promotes its tip: fetch the history from the view into the family, check the family alone holds it down to existing retained roots, then install the usual retained and base roots. On failure the event is rejected and its placeholder kept, so the upload expires normally and the event can be sent again. While a view holds staged objects every push to it is staged, because the view advertises pending refs and a client may omit objects only staging holds. Signed tips of such a push are recorded as owed before Git runs and promoted when it finishes. Compaction refuses to run while anything is owed: rollback after a State deletion needs history no ref names, so a missing ref never proves history is disposable. Objects a view holds before it is first staged are moved into the family. Staging is reclaimed with git repack -a -d -l and git prune. Git can install a ref whose parent a concurrent repack removed (see tests/git_cruft_concurrency.rs), so compaction takes the family write lease that every push already holds. It waits at most 250ms and retries with backoff. One worker handles requests from pushes, promotions and ref deletions, and reads the persistent registry at startup. The /prs/ handler now releases the family lease before post-push processing, as the standard handler does. Promotion of a waiting PR event re-enters the family and would otherwise deadlock. Assumptions - Views and their family are on one filesystem; moving pre-existing objects uses hard links. - Retained roots are complete. A damaged root fails promotion and is left to the integrity pass. - One server process per storage root, as the family lease already assumes. Excluded - Storage quotas. Staging bounds how long an unsigned upload is kept, not its size. - A pack from a signed push is stored whole. A signer can make any object reachable from their own tip, so filtering it would protect nothing. - Fetches take no lease. A fetch of a pending ref that expires while being served may fail. - Archives store a view as it is; restoring one imports staged objects. - State acceptance, rollback and purgatory sync are unchanged. Validation - nix develop -c cargo test --lib git::staging: 13 passed, covering reclamation, a pending sibling keeping its history, promotion, refusal of incomplete history, owed history surviving ref deletion, restart recovery and a busy family. - nix develop -c cargo test --test pending_upload_staging: 4 end-to-end tests passed, including an unsigned upload across a relay crash and a signed push that omits objects only staging holds. - cargo clippy --workspace --all-targets -- -D warnings and cargo fmt --check were clean. Assisted-by: Claude Fable 5.1
This commit is contained in:
@@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### Fixed
|
||||
|
||||
- Keep Git uploads that no signed event names out of permanent storage. A push
|
||||
to `refs/nostr/<event-id>` whose PR event is not yet known is staged in the
|
||||
repository that received it, promoted into shared storage when the event is
|
||||
accepted, and reclaimed after its pending ref expires. Pushes backed by a
|
||||
signed State or PR event are stored as before.
|
||||
|
||||
- Accept a GRASP-06 `/prs/` PR event whose pushed ref survived a crash that
|
||||
lost its purgatory placeholder, matching it by the event's service-local
|
||||
clone URL and exact `refs/nostr/<event-id>` ref.
|
||||
|
||||
@@ -532,6 +532,10 @@ See [`types.rs`](../../src/purgatory/types.rs) for complete definitions:
|
||||
keep the placeholder for retry. `/prs/` copies retain their separate cleanup.
|
||||
Old unscoped placeholders cannot safely identify a repository for online cleanup;
|
||||
the authorization-integrity checker preserves unmatched refs for inspection.
|
||||
- Spawns the staging maintenance worker. Uploads that no signed event names
|
||||
are staged in their view; the worker promotes history owed to the family and
|
||||
reclaims staged objects once their pending ref is gone. See
|
||||
[Staging of unsigned uploads](git-family-object-storage.md#staging-of-unsigned-uploads).
|
||||
|
||||
#### Thread Safety
|
||||
|
||||
|
||||
@@ -158,19 +158,144 @@ authorization remain view-specific.
|
||||
|
||||
For local writes, receive-pack writes its quarantine and final objects into the
|
||||
family object directory while it updates refs in the selected view. Other Git
|
||||
commands that can create objects, including proactive fetch and archive
|
||||
restore, must use the same family object directory. Read-only commands can use
|
||||
the view normally because its alternate resolves the inventory.
|
||||
commands that can create objects, including proactive fetch, copies between
|
||||
views and archive restore, must use the same family object directory. Read-only
|
||||
commands can use the view normally because its alternate resolves the
|
||||
inventory. The one exception is an upload that no signed event names yet,
|
||||
described next.
|
||||
|
||||
## Staging of unsigned uploads
|
||||
|
||||
**Added:** 2026-09-29
|
||||
|
||||
A push to `refs/nostr/<event-id>` is accepted before its PR event is known.
|
||||
Until this change its objects went straight into the family, which is never
|
||||
garbage-collected, so anyone could fill permanent storage without signing
|
||||
anything. Such an upload is now **staged**: receive-pack writes its objects to
|
||||
the view's own object directory, and they reach the family only when a signed
|
||||
event accepts them.
|
||||
|
||||
### What earns family storage
|
||||
|
||||
A signed event earns family storage at push time. Push authorization already
|
||||
decides, for each pushed ref, whether a signed event names it:
|
||||
|
||||
| Pushed ref | Objects go to |
|
||||
| ------------------------------------------------------------ | ------------- |
|
||||
| Branch or tag named by a State, accepted or in purgatory | family |
|
||||
| `refs/nostr/<id>` whose PR event is accepted or in purgatory | family |
|
||||
| `refs/nostr/<id>` with no event, or only a placeholder | view staging |
|
||||
|
||||
A push that carries any unsigned ref is staged as a whole, because one push is
|
||||
one pack. While a view holds staged objects, every push to it is staged too:
|
||||
the view advertises its pending refs, so a client may omit objects that only
|
||||
staging holds, and receive-pack could not check such a push against the family
|
||||
alone.
|
||||
|
||||
Trade-off: a pack from a signed push is stored whole, including any object in
|
||||
it that the signed tip does not reach. This is accepted because a signer can
|
||||
make any object reachable from their own tip. Staging defends against uploads
|
||||
nobody signed for, not against a signer's choice of content. Purgatory sync
|
||||
fetches and integrity repair fetch history for signed events and keep writing
|
||||
to the family.
|
||||
|
||||
### Promotion
|
||||
|
||||
Promotion copies a tip's history from the view into the family with
|
||||
`git fetch`, checks that the family alone holds it, and installs the usual
|
||||
retained and base roots. The check walks from the tip down to existing retained
|
||||
roots. Those were complete when installed and the family is append-only, so the
|
||||
cost follows the new history instead of the whole repository. Detecting damage
|
||||
behind a root remains the integrity pass's job; a damaged root fails the check
|
||||
and therefore the promotion.
|
||||
|
||||
Promotion happens in two places:
|
||||
|
||||
- **PR acceptance.** A PR or PR Update event is accepted only after its tip is
|
||||
promoted. On failure the event is rejected and its placeholder is kept, so
|
||||
the upload expires normally and the client can send the event again.
|
||||
- **Signed push into a staged view.** Its signed tips are promoted when
|
||||
receive-pack finishes, while the push still holds the family lease.
|
||||
|
||||
State events never wait on promotion. Their history is in the family when the
|
||||
push that carries it completes.
|
||||
|
||||
### Owed history
|
||||
|
||||
Rollback after a State deletion depends on history that no current ref names.
|
||||
The absence of a ref therefore never proves that staged history is disposable.
|
||||
|
||||
Before receive-pack runs, the signed tips of a staged push are recorded as
|
||||
**owed** in the view's registry record. A tip stays owed until its history is
|
||||
complete in the family alone. If promotion at the end of the push fails, or
|
||||
the server stops before it runs, the tip is still owed after restart and
|
||||
maintenance promotes it by object ID, whether or not a ref still names it. An
|
||||
owed tip whose object is in neither store belongs to a push that never
|
||||
completed and is dropped. An empty `/prs/` view is not removed while it owes
|
||||
history.
|
||||
|
||||
A view that already holds objects when it is first staged has them moved into
|
||||
the family first. They predate staging and cannot be told apart from accepted
|
||||
history.
|
||||
|
||||
### Compaction
|
||||
|
||||
Staging is reclaimed by Git itself. For one view, maintenance:
|
||||
|
||||
1. promotes owed tips, and stops if any remain owed;
|
||||
2. runs `git repack -a -d -l` and `git prune --expire=now`, which keep the
|
||||
history live refs need and the family lacks, and drop everything else;
|
||||
3. removes loose objects the family also holds;
|
||||
4. removes the registry record if nothing is left, and the view itself if it
|
||||
is an empty `/prs/` view.
|
||||
|
||||
Every live ref is a root, so a pending upload keeps exactly its own history
|
||||
and an abandoned one disappears once its ref expires.
|
||||
|
||||
Compaction must not overlap a push to the same view. Git can report a
|
||||
successful push and install a ref whose parent a concurrent repack has just
|
||||
removed; `tests/git_cruft_concurrency.rs` reproduces this. A grace period
|
||||
cannot prevent it, because the push may start after the repack has chosen
|
||||
what to keep. Compaction therefore takes the family write lease, which every
|
||||
push already holds from before receive-pack starts until its refs and roots
|
||||
are installed. The lease is fair: maintenance queues behind active writers for
|
||||
at most 250 ms, then gives way and retries with backoff from one second to five
|
||||
minutes, so it cannot hold up the pushes behind it.
|
||||
|
||||
Fetches take no lease. A fetch of a pending ref that expires while it is being
|
||||
served may fail. That ref named an upload nobody signed for and is being
|
||||
removed deliberately.
|
||||
|
||||
### Registry and maintenance worker
|
||||
|
||||
`.grasp/staging/<digest>.json` records each staged view and its owed tips. It
|
||||
is written and fsynced before receive-pack can write an object. One worker
|
||||
processes maintenance requests, which come from staged pushes, promotions and
|
||||
ref deletions, including placeholder expiry. At startup it reads the registry,
|
||||
so interrupted work resumes. A view whose remaining staged objects belong to
|
||||
pending refs is parked until the next request rather than polled.
|
||||
|
||||
### Limits
|
||||
|
||||
- Staging is not a storage quota. It bounds how long an unsigned upload is
|
||||
kept, not how much one may upload.
|
||||
- Archiving a repository stores its view directory as it is, staged objects
|
||||
included. Restoring the archive imports them into the family.
|
||||
- Staged history exists in one view only. Other views cannot read it until it
|
||||
is promoted.
|
||||
|
||||
## Durability and the success fence
|
||||
|
||||
The invariant is:
|
||||
|
||||
> When a client observes a successful push, every Git object needed by the
|
||||
> accepted ref updates is durable in the configured family backend.
|
||||
> accepted ref updates is durable in the configured family backend, or in the
|
||||
> view's staging for an upload that no signed event names yet.
|
||||
|
||||
The local backend satisfies the fence when Git has atomically installed the
|
||||
objects in the family object directory and receive-pack has completed.
|
||||
objects in the family object directory, or in view staging, and receive-pack
|
||||
has completed. Accepting a PR event additionally requires its history to be
|
||||
complete in the family.
|
||||
|
||||
The S3 backend cannot release receive-pack's terminal success immediately.
|
||||
The handler must retain the final protocol status until it has:
|
||||
@@ -224,8 +349,12 @@ Consequences:
|
||||
|
||||
- disable automatic Git maintenance and pruning for family repositories;
|
||||
- create retained roots for every accepted branch, tag, and PR tip;
|
||||
- unsigned uploads stay outside this inventory until a signed event accepts
|
||||
them;
|
||||
- do not use S3 lifecycle deletion on family objects;
|
||||
- never compact by packing only the current visible ref closure;
|
||||
- never compact the family by packing only the current visible ref closure
|
||||
(view staging is compacted this way because it is not part of the
|
||||
inventory, and only once nothing in it is owed to the family);
|
||||
- if pack-count compaction becomes necessary, repack the union of every object
|
||||
in all selected packs, publish the replacement in addition to the old
|
||||
immutable packs, and leave physical deletion to future GC work.
|
||||
@@ -243,8 +372,10 @@ view whose alternate is the identifier family. Therefore:
|
||||
- receive-pack may advertise family base tips as anonymous `.have` lines;
|
||||
- the first contributor push sends only objects the family does not have;
|
||||
- pushes remain restricted to `refs/nostr/<event-id>`;
|
||||
- a push whose PR event is not yet known is staged in the view;
|
||||
- placeholder and expiry cleanup removes view refs or an empty view, never
|
||||
family objects or retained roots.
|
||||
family objects or retained roots. An empty view that still owes history to
|
||||
the family is kept.
|
||||
|
||||
Mirroring a PR into an owner view becomes a ref update after an object
|
||||
availability check. It no longer copies the object graph.
|
||||
@@ -384,9 +515,10 @@ local objects exist.
|
||||
|
||||
Operations that mutate one family are serialized by a per-family lock. This
|
||||
covers pushes to different owner or `/prs/` views with the same identifier,
|
||||
proactive fetches, archive restores, retained/base ref updates, and S3 manifest
|
||||
publication. Read requests take a hydrated family lease; they do not hold the
|
||||
writer lock after their pack set has been pinned.
|
||||
proactive fetches, archive restores, retained/base ref updates, promotion and
|
||||
compaction of view staging, and S3 manifest publication. Read requests take a
|
||||
hydrated family lease; they do not hold the writer lock after their pack set
|
||||
has been pinned.
|
||||
|
||||
View lifecycle locks remain responsible for “may this path be removed?” The
|
||||
family lock is responsible for “is this object inventory and manifest update
|
||||
@@ -397,8 +529,9 @@ atomic?” Neither lock grants authorization.
|
||||
Storage integrity is defined for an object-format/identifier family, not for
|
||||
one owner path. The storage 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
|
||||
correct alternate and for refs whose targets are available in the family. A
|
||||
ref of a staged view whose history is complete in that view's staging is a
|
||||
pending upload, not a fault. Multiple independent histories in one identifier are valid, and unreachable
|
||||
objects are not an error: retained delete-state and rollback data intentionally
|
||||
remain present.
|
||||
|
||||
|
||||
@@ -614,6 +614,12 @@ on failure for retry. A database check protects accepted events whose placeholde
|
||||
survived in an older checkpoint. Legacy records without destinations remain readable,
|
||||
but cannot safely drive scoped online ref deletion.
|
||||
|
||||
The objects of a git-first push are staged in the view that received them, not
|
||||
stored in the identifier family. Expiring the ref lets staging maintenance
|
||||
reclaim them. Accepting the PR event first promotes its history into the
|
||||
family; the placeholder is released only after that succeeds. See
|
||||
[Staging of unsigned uploads](git-family-object-storage.md#staging-of-unsigned-uploads).
|
||||
|
||||
## Background Sync
|
||||
|
||||
Purgatory includes a background sync system that fetches git data from remote servers when events arrive before git data.
|
||||
|
||||
@@ -76,6 +76,7 @@ pub async fn authorize_push(
|
||||
|
||||
// Collect all purgatory events that authorize this push
|
||||
let mut purgatory_events = Vec::new();
|
||||
let mut unsigned_refs = HashSet::new();
|
||||
|
||||
// Handle refs/nostr/ refs - validate and collect PR/PR-update events from purgatory
|
||||
if !nostr_refs.is_empty() {
|
||||
@@ -101,7 +102,11 @@ pub async fn authorize_push(
|
||||
}
|
||||
NostrRefPreValidation::Authorized {
|
||||
event_from_purgatory,
|
||||
signed,
|
||||
} => {
|
||||
if !signed {
|
||||
unsigned_refs.insert(ref_name.clone());
|
||||
}
|
||||
if let Some(event) = event_from_purgatory {
|
||||
purgatory_events.push(event);
|
||||
false
|
||||
@@ -109,7 +114,10 @@ pub async fn authorize_push(
|
||||
purgatory.find_pr_placeholder(event_id).is_some()
|
||||
}
|
||||
}
|
||||
NostrRefPreValidation::Unknown => true,
|
||||
NostrRefPreValidation::Unknown => {
|
||||
unsigned_refs.insert(ref_name.clone());
|
||||
true
|
||||
}
|
||||
};
|
||||
if track_placeholder {
|
||||
// Remember every destination, including subsequent pushes of the
|
||||
@@ -183,6 +191,7 @@ pub async fn authorize_push(
|
||||
state: auth_result.state,
|
||||
maintainers: auth_result.maintainers,
|
||||
purgatory_events,
|
||||
unsigned_refs,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -193,6 +202,7 @@ pub async fn authorize_push(
|
||||
state: None,
|
||||
maintainers: vec![],
|
||||
purgatory_events,
|
||||
unsigned_refs,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -668,6 +678,7 @@ pub async fn get_state_authorization_for_selected_repo(
|
||||
state: None,
|
||||
maintainers: authorized.into_iter().collect(),
|
||||
purgatory_events: vec![],
|
||||
unsigned_refs: HashSet::new(),
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -739,6 +750,7 @@ pub async fn get_state_authorization_for_selected_repo(
|
||||
state: Some(state),
|
||||
maintainers: authorized.into_iter().collect(),
|
||||
purgatory_events: vec![latest_authorized.clone()],
|
||||
unsigned_refs: HashSet::new(),
|
||||
});
|
||||
} else {
|
||||
warn!(
|
||||
@@ -875,6 +887,9 @@ pub struct AuthorizationResult {
|
||||
pub maintainers: Vec<String>,
|
||||
/// Events from purgatory that authorized this push (state, PR, PR-update events)
|
||||
pub purgatory_events: Vec<Event>,
|
||||
/// Pushed `refs/nostr/` refs that no signed event names yet. Their objects
|
||||
/// are received into view staging instead of the family.
|
||||
pub unsigned_refs: HashSet<String>,
|
||||
}
|
||||
|
||||
impl AuthorizationResult {
|
||||
@@ -886,6 +901,7 @@ impl AuthorizationResult {
|
||||
state: None,
|
||||
maintainers: vec![],
|
||||
purgatory_events: vec![],
|
||||
unsigned_refs: HashSet::new(),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1248,7 +1264,13 @@ pub enum NostrRefPreValidation {
|
||||
/// `authorize_push`) can collect it into `purgatory_events`. `None`
|
||||
/// means the match came from the DB or from a placeholder-only
|
||||
/// purgatory entry — nothing to collect.
|
||||
Authorized { event_from_purgatory: Option<Event> },
|
||||
///
|
||||
/// `signed` is false for a placeholder-only match: no signed event names
|
||||
/// this commit yet, so its objects have not earned family storage.
|
||||
Authorized {
|
||||
event_from_purgatory: Option<Event>,
|
||||
signed: bool,
|
||||
},
|
||||
/// No event with that id is known to the relay yet. The caller may
|
||||
/// create a placeholder (with or without a `/prs/` scope) so the
|
||||
/// purgatory sweep can clean up the ref if the event never arrives.
|
||||
@@ -1318,6 +1340,7 @@ pub async fn pre_validate_refs_nostr_push(
|
||||
}
|
||||
return NostrRefPreValidation::Authorized {
|
||||
event_from_purgatory: None,
|
||||
signed: true,
|
||||
};
|
||||
}
|
||||
Ok(None) => {}
|
||||
@@ -1344,6 +1367,7 @@ pub async fn pre_validate_refs_nostr_push(
|
||||
}
|
||||
return NostrRefPreValidation::Authorized {
|
||||
event_from_purgatory: Some(event),
|
||||
signed: true,
|
||||
};
|
||||
}
|
||||
None => {
|
||||
@@ -1367,6 +1391,7 @@ pub async fn pre_validate_refs_nostr_push(
|
||||
}
|
||||
return NostrRefPreValidation::Authorized {
|
||||
event_from_purgatory: None,
|
||||
signed: false,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
+68
-5
@@ -710,7 +710,7 @@ pub async fn handle_receive_pack(
|
||||
);
|
||||
|
||||
// check push is authorised
|
||||
let _auth_result = match authorize_push(
|
||||
let auth_result = match authorize_push(
|
||||
&database,
|
||||
identifier,
|
||||
owner_pubkey,
|
||||
@@ -788,8 +788,16 @@ pub async fn handle_receive_pack(
|
||||
.map(|plan| &plan.forwarded_body)
|
||||
.unwrap_or(&request_body);
|
||||
|
||||
// Ref updates remain in the selected view; receive-pack's quarantine and
|
||||
// final objects are installed directly in the shared family inventory.
|
||||
let pushed_refs = parse_pushed_refs(&request_body);
|
||||
let staged = if family_lease.is_some() {
|
||||
stage_push(&repo_path, &pushed_refs, &auth_result.unsigned_refs).await?
|
||||
} else {
|
||||
false
|
||||
};
|
||||
|
||||
// Ref updates remain in the selected view. Objects of a push backed by
|
||||
// signed events are installed directly in the shared family inventory;
|
||||
// a staged push keeps them in the view until promotion.
|
||||
let mut git = GitSubprocess::spawn_with_object_directory(
|
||||
GitService::ReceivePack,
|
||||
&repo_path,
|
||||
@@ -797,6 +805,7 @@ pub async fn handle_receive_pack(
|
||||
git_protocol,
|
||||
family_lease
|
||||
.as_ref()
|
||||
.filter(|_| !staged)
|
||||
.map(|lease| lease.family_objects_path.as_path()),
|
||||
)
|
||||
.map_err(GitError::ProcessSpawnFailed)?;
|
||||
@@ -817,7 +826,6 @@ pub async fn handle_receive_pack(
|
||||
// uploading. Streaming sideband progress keeps libgit2 clients from
|
||||
// hitting their per-recv timeout during that otherwise-silent window.
|
||||
let (tx, rx) = mpsc::channel::<Result<Frame<Bytes>, io::Error>>(STREAM_CHANNEL_DEPTH);
|
||||
let pushed_refs = parse_pushed_refs(&request_body);
|
||||
let new_oids: HashSet<String> = pushed_refs
|
||||
.iter()
|
||||
.filter(|(_, new_oid, _)| new_oid != "0000000000000000000000000000000000000000")
|
||||
@@ -852,6 +860,7 @@ pub async fn handle_receive_pack(
|
||||
family_key,
|
||||
family_lease,
|
||||
pushed_refs,
|
||||
staged,
|
||||
push_plan,
|
||||
)
|
||||
.await;
|
||||
@@ -883,6 +892,7 @@ async fn stream_receive_pack_output<S, E, I>(
|
||||
family_key: FamilyKey,
|
||||
family_lease: Option<FamilyWriteLease>,
|
||||
pushed_refs: Vec<(String, String, String)>,
|
||||
staged: bool,
|
||||
mut push_plan: Option<ReceivePackPlan>,
|
||||
) where
|
||||
I: tokio::io::AsyncWrite + Unpin + Send + 'static,
|
||||
@@ -1015,7 +1025,9 @@ async fn stream_receive_pack_output<S, E, I>(
|
||||
|
||||
debug!("Git receive-pack stream completed successfully");
|
||||
|
||||
if family_lease.is_some() {
|
||||
if staged {
|
||||
settle_staged_push(&repo_path).await;
|
||||
} else if family_lease.is_some() {
|
||||
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
|
||||
}
|
||||
|
||||
@@ -1153,6 +1165,57 @@ impl std::fmt::Display for GitError {
|
||||
}
|
||||
}
|
||||
|
||||
/// Decide where a push to a thin view stores its objects, before Git runs.
|
||||
///
|
||||
/// A push is staged in the view when it carries a ref that no signed event
|
||||
/// names, or when the view already holds staged objects: the client may have
|
||||
/// omitted objects that only staging holds. Signed tips of a staged push are
|
||||
/// recorded as owed to the family first. The caller holds the family lease.
|
||||
pub(crate) async fn stage_push(
|
||||
repo_path: &std::path::Path,
|
||||
pushed_refs: &[(String, String, String)],
|
||||
unsigned_refs: &HashSet<String>,
|
||||
) -> Result<bool, GitError> {
|
||||
if unsigned_refs.is_empty() && !super::staging::is_staged(repo_path) {
|
||||
return Ok(false);
|
||||
}
|
||||
let owed: Vec<_> = pushed_refs
|
||||
.iter()
|
||||
.filter(|(_, new_oid, name)| {
|
||||
!unsigned_refs.contains(name) && new_oid.bytes().any(|digit| digit != b'0')
|
||||
})
|
||||
.map(|(_, new_oid, name)| super::staging::Tip::new(name, new_oid))
|
||||
.collect();
|
||||
let view = repo_path.to_owned();
|
||||
tokio::task::spawn_blocking(move || super::staging::stage(&view, &owed))
|
||||
.await
|
||||
.map_err(|error| GitError::Storage(error.to_string()))?
|
||||
.map_err(|error| GitError::Storage(format!("cannot stage push: {error:#}")))
|
||||
}
|
||||
|
||||
/// Promote the signed tips of a staged push while the family lease is held.
|
||||
///
|
||||
/// A failure leaves the tips owed for staging maintenance to retry. The push
|
||||
/// has already succeeded, so it is not reported to the client.
|
||||
pub(crate) async fn settle_staged_push(repo_path: &std::path::Path) {
|
||||
let view = repo_path.to_owned();
|
||||
match tokio::task::spawn_blocking(move || super::staging::settle(&view)).await {
|
||||
Ok(Ok(0)) => {}
|
||||
Ok(Ok(owed)) => warn!(
|
||||
repo = %repo_path.display(),
|
||||
owed,
|
||||
"Signed history of a staged push is not yet in the family"
|
||||
),
|
||||
Ok(Err(error)) => error!(
|
||||
repo = %repo_path.display(),
|
||||
error = %format!("{error:#}"),
|
||||
"Failed to promote signed history of a staged push"
|
||||
),
|
||||
Err(error) => error!(repo = %repo_path.display(), %error, "Staged push promotion panicked"),
|
||||
}
|
||||
super::staging::request_maintenance(repo_path);
|
||||
}
|
||||
|
||||
pub(crate) fn retain_accepted_tips(
|
||||
storage: &LocalGitStorage,
|
||||
family_key: &FamilyKey,
|
||||
|
||||
@@ -629,7 +629,9 @@ pub fn inspect_family(storage: &LocalGitStorage, key: &FamilyKey) -> Result<Fami
|
||||
}
|
||||
for (reference, oid) in list_refs(view)? {
|
||||
refs_checked += 1;
|
||||
if !oid_exists(&family, &oid)? {
|
||||
// A pending upload is complete in its view's staging and is not
|
||||
// family history until its event is accepted.
|
||||
if !oid_exists(&family, &oid)? && !super::staging::holds_history(view, &oid) {
|
||||
missing_oids.insert(oid.clone());
|
||||
missing_ref_targets.push(MissingRefTarget {
|
||||
view: view.clone(),
|
||||
|
||||
@@ -25,6 +25,7 @@ pub mod migration;
|
||||
pub mod process;
|
||||
pub mod protocol;
|
||||
mod receive_pack_plan;
|
||||
pub mod staging;
|
||||
pub mod storage;
|
||||
pub mod subprocess;
|
||||
pub mod sync;
|
||||
@@ -254,6 +255,7 @@ pub fn delete_ref(repo_path: &Path, ref_name: &str) -> Result<(), String> {
|
||||
}
|
||||
|
||||
info!("Deleted ref {} from {}", ref_name, repo_path.display());
|
||||
staging::request_maintenance(repo_path);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,608 @@
|
||||
//! Unsigned uploads stay in their view until a signed event earns family storage.
|
||||
//!
|
||||
//! A push naming a `refs/nostr/<event-id>` for which no signed event is known
|
||||
//! writes its objects to the view's own object directory instead of the
|
||||
//! identifier family. The family is never garbage-collected, so this keeps
|
||||
//! abandoned uploads reclaimable. Accepting the PR event promotes the tip's
|
||||
//! history into the family first.
|
||||
//!
|
||||
//! While a view holds staged objects every push to it is received into the
|
||||
//! view, because a client may omit objects that only staging holds. Signed
|
||||
//! tips of such pushes are recorded as *owed* before Git runs and stay owed
|
||||
//! until their history is complete in the family alone. Compaction refuses to
|
||||
//! run while anything is owed: the absence of a ref never proves that history
|
||||
//! is disposable, since rollback depends on history no current ref names.
|
||||
//!
|
||||
//! Every function that reads or changes staged objects or the registry record
|
||||
//! requires the caller to hold the family write lease. Pushes already hold it
|
||||
//! for their whole duration, so it also excludes them from compaction.
|
||||
use std::fs::File;
|
||||
use std::io::Write;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::{Command, Output, Stdio};
|
||||
|
||||
use anyhow::{bail, ensure, Context, Result};
|
||||
use bitcoin_hashes::{sha256, Hash};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat};
|
||||
|
||||
mod worker;
|
||||
pub use worker::{request_maintenance, run_worker};
|
||||
|
||||
const REGISTRY_DIR: &str = "staging";
|
||||
const RETAINED_GLOB: &str = "--glob=refs/grasp/retained/*";
|
||||
|
||||
/// A signed ref tip whose history the family must hold.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct Tip {
|
||||
pub reference: String,
|
||||
pub oid: String,
|
||||
}
|
||||
|
||||
impl Tip {
|
||||
pub fn new(reference: impl Into<String>, oid: impl Into<String>) -> Self {
|
||||
Self {
|
||||
reference: reference.into(),
|
||||
oid: oid.into(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Persistent registration of one view holding staged objects.
|
||||
#[derive(Debug, Default, Serialize, Deserialize)]
|
||||
struct Record {
|
||||
/// View path relative to the Git data root.
|
||||
view: PathBuf,
|
||||
#[serde(default)]
|
||||
owed: Vec<Tip>,
|
||||
}
|
||||
|
||||
/// Outcome of one compaction.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
pub enum Compaction {
|
||||
/// No staged objects remain and the view is no longer registered.
|
||||
Done,
|
||||
/// Staging holds only history that live refs need and the family lacks.
|
||||
Pending,
|
||||
}
|
||||
|
||||
/// A thin view resolved to its family and registry record.
|
||||
struct View {
|
||||
path: PathBuf,
|
||||
storage: LocalGitStorage,
|
||||
key: FamilyKey,
|
||||
record: PathBuf,
|
||||
}
|
||||
|
||||
impl View {
|
||||
/// `None` for a legacy repository that has no family alternate.
|
||||
fn resolve(view: &Path) -> Result<Option<Self>> {
|
||||
let path = std::fs::canonicalize(view)
|
||||
.with_context(|| format!("resolve view {}", view.display()))?;
|
||||
let text = match std::fs::read_to_string(path.join("objects/info/alternates")) {
|
||||
Ok(text) => text,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
let objects = PathBuf::from(text.lines().next().context("empty family alternate")?);
|
||||
let repo = objects.parent().context("family repository")?;
|
||||
let root = repo.ancestors().nth(4).context("family storage root")?;
|
||||
let identifier = repo
|
||||
.file_name()
|
||||
.and_then(|name| name.to_str())
|
||||
.and_then(|name| name.strip_suffix(".git"))
|
||||
.context("family identifier")?;
|
||||
// The family path names its object format; no Git process is needed.
|
||||
let format = match repo
|
||||
.parent()
|
||||
.and_then(Path::file_name)
|
||||
.and_then(|name| name.to_str())
|
||||
{
|
||||
Some("sha1") => ObjectFormat::Sha1,
|
||||
Some("sha256") => ObjectFormat::Sha256,
|
||||
_ => bail!("unsupported alternate for staging: {}", path.display()),
|
||||
};
|
||||
let storage = LocalGitStorage::new(root);
|
||||
let key = FamilyKey::new(format, identifier)?;
|
||||
ensure!(
|
||||
storage.family_objects_path(&key) == objects,
|
||||
"unsupported alternate for staging: {}",
|
||||
path.display()
|
||||
);
|
||||
let root = std::fs::canonicalize(storage.git_data_path())?;
|
||||
let relative = path
|
||||
.strip_prefix(&root)
|
||||
.context("view outside Git storage")?;
|
||||
let digest = sha256::Hash::hash(relative.as_os_str().as_encoded_bytes());
|
||||
let record = registry_dir(&storage).join(format!("{digest}.json"));
|
||||
Ok(Some(Self {
|
||||
path,
|
||||
storage,
|
||||
key,
|
||||
record,
|
||||
}))
|
||||
}
|
||||
|
||||
fn family(&self) -> PathBuf {
|
||||
self.storage.family_repo_path(&self.key)
|
||||
}
|
||||
|
||||
fn relative(&self) -> Result<PathBuf> {
|
||||
let root = std::fs::canonicalize(self.storage.git_data_path())?;
|
||||
Ok(self.path.strip_prefix(root)?.to_owned())
|
||||
}
|
||||
|
||||
fn load(&self) -> Result<Option<Record>> {
|
||||
match std::fs::read(&self.record) {
|
||||
Ok(bytes) => Ok(Some(serde_json::from_slice(&bytes).with_context(|| {
|
||||
format!("parse staging record {}", self.record.display())
|
||||
})?)),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
|
||||
Err(error) => Err(error.into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn save(&self, record: &Record) -> Result<()> {
|
||||
let directory = self.record.parent().context("staging registry")?;
|
||||
std::fs::create_dir_all(directory)?;
|
||||
let mut file = tempfile::NamedTempFile::new_in(directory)?;
|
||||
file.write_all(&serde_json::to_vec(record)?)?;
|
||||
file.as_file().sync_all()?;
|
||||
file.persist(&self.record)?;
|
||||
File::open(directory)?.sync_all()?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn unregister(&self) -> Result<()> {
|
||||
match std::fs::remove_file(&self.record) {
|
||||
Ok(()) => {}
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
|
||||
Err(error) => return Err(error.into()),
|
||||
}
|
||||
File::open(self.record.parent().context("staging registry")?)?.sync_all()?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
fn registry_dir(storage: &LocalGitStorage) -> PathBuf {
|
||||
storage.internal_path().join(REGISTRY_DIR)
|
||||
}
|
||||
|
||||
fn valid_oid(oid: &str) -> bool {
|
||||
matches!(oid.len(), 40 | 64) && oid.bytes().all(|c| c.is_ascii_hexdigit())
|
||||
}
|
||||
|
||||
fn git(repo: &Path, args: &[&str], input: &[u8]) -> Result<Output> {
|
||||
spawn_git(repo, args, input, Stdio::null())
|
||||
}
|
||||
|
||||
/// Run Git and keep its output. Only for commands that print one short line
|
||||
/// per input line, which the reader thread drains while input is written.
|
||||
fn git_output(repo: &Path, args: &[&str], input: &[u8]) -> Result<String> {
|
||||
let output = spawn_git(repo, args, input, Stdio::piped())?;
|
||||
ensure!(
|
||||
output.status.success(),
|
||||
"git {} failed in {}: {}",
|
||||
args.first().copied().unwrap_or_default(),
|
||||
repo.display(),
|
||||
String::from_utf8_lossy(&output.stderr).trim()
|
||||
);
|
||||
Ok(String::from_utf8(output.stdout)?)
|
||||
}
|
||||
|
||||
fn spawn_git(repo: &Path, args: &[&str], input: &[u8], stdout: Stdio) -> Result<Output> {
|
||||
let mut child = Command::new("git")
|
||||
.current_dir(repo)
|
||||
.args(args)
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(stdout)
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()
|
||||
.with_context(|| format!("spawn git {}", args.first().copied().unwrap_or_default()))?;
|
||||
let mut stdin = child.stdin.take().expect("piped stdin");
|
||||
let input = input.to_owned();
|
||||
let writer = std::thread::spawn(move || stdin.write_all(&input));
|
||||
let output = child.wait_with_output()?;
|
||||
match writer.join() {
|
||||
Ok(Err(error)) if error.kind() != std::io::ErrorKind::BrokenPipe => {
|
||||
return Err(error.into())
|
||||
}
|
||||
Ok(_) => {}
|
||||
Err(_) => bail!("git input writer panicked"),
|
||||
}
|
||||
Ok(output)
|
||||
}
|
||||
|
||||
fn run(repo: &Path, args: &[&str]) -> Result<()> {
|
||||
let output = git(repo, args, b"")?;
|
||||
ensure!(
|
||||
output.status.success(),
|
||||
"git {} failed in {}: {}",
|
||||
args.first().copied().unwrap_or_default(),
|
||||
repo.display(),
|
||||
String::from_utf8_lossy(&output.stderr).trim()
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn object_exists(repo: &Path, oid: &str) -> Result<bool> {
|
||||
Ok(git(repo, &["cat-file", "-e", oid], b"")?.status.success())
|
||||
}
|
||||
|
||||
/// Whether `oids` and their history are present in the family alone.
|
||||
///
|
||||
/// The walk stops at retained roots. Those were complete when installed and
|
||||
/// the family is append-only, so the cost follows the new history rather than
|
||||
/// the whole repository. Detecting damage behind a root is the integrity
|
||||
/// worker's job; a damaged root fails this check and therefore the promotion.
|
||||
fn complete_in_family(family: &Path, oids: &[&str]) -> Result<bool> {
|
||||
let mut input = Vec::new();
|
||||
for oid in oids {
|
||||
input.extend_from_slice(oid.as_bytes());
|
||||
input.push(b'\n');
|
||||
}
|
||||
input.extend_from_slice(format!("--not\n{RETAINED_GLOB}\n").as_bytes());
|
||||
let output = git(
|
||||
family,
|
||||
&["rev-list", "--objects", "--missing=error", "--stdin"],
|
||||
&input,
|
||||
)?;
|
||||
Ok(output.status.success())
|
||||
}
|
||||
|
||||
/// Whether the view is registered as holding staged objects.
|
||||
///
|
||||
/// Reads one directory entry; safe without the family lease as a hint. A
|
||||
/// decision that depends on it must be made while holding the lease.
|
||||
pub fn is_staged(view: &Path) -> bool {
|
||||
View::resolve(view)
|
||||
.ok()
|
||||
.flatten()
|
||||
.is_some_and(|view| view.record.is_file())
|
||||
}
|
||||
|
||||
/// Whether a staged view holds `oid` with all history the family lacks.
|
||||
///
|
||||
/// The walk stops at the family base tips the view advertises, so its cost
|
||||
/// follows the staged history.
|
||||
pub fn holds_history(view: &Path, oid: &str) -> bool {
|
||||
valid_oid(oid)
|
||||
&& is_staged(view)
|
||||
&& git(
|
||||
view,
|
||||
&[
|
||||
"rev-list",
|
||||
"--objects",
|
||||
"--missing=error",
|
||||
oid,
|
||||
"--not",
|
||||
"--alternate-refs",
|
||||
],
|
||||
b"",
|
||||
)
|
||||
.is_ok_and(|output| output.status.success())
|
||||
}
|
||||
|
||||
/// Whether signed history is still owed to the family from this view.
|
||||
///
|
||||
/// A view that owes history must not be removed, even with no refs left.
|
||||
pub fn owes_history(view: &Path) -> bool {
|
||||
match View::resolve(view).and_then(|view| view.map(|view| view.load()).transpose()) {
|
||||
Ok(record) => record
|
||||
.flatten()
|
||||
.is_some_and(|record| !record.owed.is_empty()),
|
||||
// Unreadable ownership metadata is never permission to delete.
|
||||
Err(_) => true,
|
||||
}
|
||||
}
|
||||
|
||||
/// Register the view and record `owed` tips before Git can write any object.
|
||||
///
|
||||
/// Returns `false` for a legacy repository without a family, which keeps its
|
||||
/// own objects and is never staged.
|
||||
pub fn stage(view: &Path, owed: &[Tip]) -> Result<bool> {
|
||||
let Some(view) = View::resolve(view)? else {
|
||||
return Ok(false);
|
||||
};
|
||||
ensure!(
|
||||
owed.iter().all(|tip| valid_oid(&tip.oid)),
|
||||
"invalid owed object ID"
|
||||
);
|
||||
let mut record = match view.load()? {
|
||||
Some(record) => record,
|
||||
None => {
|
||||
// Whatever the view holds before it is first staged predates
|
||||
// staging and cannot be told apart from accepted history.
|
||||
absorb_local_objects(&view)?;
|
||||
Record::default()
|
||||
}
|
||||
};
|
||||
record.view = view.relative()?;
|
||||
let before = record.owed.len();
|
||||
for tip in owed {
|
||||
if !record.owed.contains(tip) {
|
||||
record.owed.push(tip.clone());
|
||||
}
|
||||
}
|
||||
if before != record.owed.len() || !view.record.is_file() {
|
||||
view.save(&record)?;
|
||||
}
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
/// Move every object file of the view into the family.
|
||||
///
|
||||
/// Files are linked into the family before they are removed from the view, so
|
||||
/// a concurrent reader of the view always finds each object in one of them.
|
||||
fn absorb_local_objects(view: &View) -> Result<()> {
|
||||
let source = view.path.join("objects");
|
||||
let target = view.storage.family_objects_path(&view.key);
|
||||
let mut moved = Vec::new();
|
||||
for entry in std::fs::read_dir(&source)? {
|
||||
let entry = entry?;
|
||||
let name = entry.file_name();
|
||||
if name == "info" || !entry.file_type()?.is_dir() {
|
||||
continue;
|
||||
}
|
||||
let directory = target.join(&name);
|
||||
let mut files: Vec<_> = std::fs::read_dir(entry.path())?
|
||||
.map(|file| file.map(|file| file.path()))
|
||||
.collect::<std::io::Result<_>>()?;
|
||||
// Git finds a pack through its index, so the index arrives last.
|
||||
files.sort_by_key(|file| file.extension().is_some_and(|ext| ext == "idx"));
|
||||
for file in files {
|
||||
if !file.is_file() {
|
||||
continue;
|
||||
}
|
||||
std::fs::create_dir_all(&directory)?;
|
||||
let destination = directory.join(file.file_name().context("object file name")?);
|
||||
match std::fs::hard_link(&file, &destination) {
|
||||
Ok(()) => {}
|
||||
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
|
||||
Err(error) => return Err(error.into()),
|
||||
}
|
||||
moved.push(file);
|
||||
}
|
||||
File::open(&directory)?.sync_all()?;
|
||||
}
|
||||
if moved.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
File::open(&target)?.sync_all()?;
|
||||
tracing::info!(
|
||||
view = %view.path.display(),
|
||||
files = moved.len(),
|
||||
"Moved existing view objects into the family before staging"
|
||||
);
|
||||
// Remove indexes first, the reverse of the order they were installed in.
|
||||
moved.sort_by_key(|file| file.extension().is_none_or(|ext| ext != "idx"));
|
||||
for file in moved {
|
||||
std::fs::remove_file(file)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Copy the history of `tips` into the family and retain it there.
|
||||
///
|
||||
/// Succeeds only when every tip is complete in the family alone, which is the
|
||||
/// condition for accepting the event that names it.
|
||||
pub fn promote(view: &Path, tips: &[Tip]) -> Result<()> {
|
||||
let Some(view) = View::resolve(view)? else {
|
||||
return Ok(());
|
||||
};
|
||||
promote_resolved(&view, tips)
|
||||
}
|
||||
|
||||
fn promote_resolved(view: &View, tips: &[Tip]) -> Result<()> {
|
||||
if tips.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
ensure!(
|
||||
tips.iter().all(|tip| valid_oid(&tip.oid)),
|
||||
"invalid staged object ID"
|
||||
);
|
||||
let family = view.family();
|
||||
let mut oids: Vec<&str> = tips.iter().map(|tip| tip.oid.as_str()).collect();
|
||||
oids.sort_unstable();
|
||||
oids.dedup();
|
||||
|
||||
if !complete_in_family(&family, &oids)? {
|
||||
let source = view.path.to_str().context("non-UTF-8 view path")?;
|
||||
let mut args = vec![
|
||||
"-c",
|
||||
"uploadpack.allowAnySHA1InWant=true",
|
||||
"fetch",
|
||||
"--no-tags",
|
||||
"--no-write-fetch-head",
|
||||
"--no-auto-maintenance",
|
||||
"--no-recurse-submodules",
|
||||
source,
|
||||
];
|
||||
args.extend_from_slice(&oids);
|
||||
run(&family, &args).context("copy staged history into the family")?;
|
||||
ensure!(
|
||||
complete_in_family(&family, &oids)?,
|
||||
"promoted history is incomplete in the family"
|
||||
);
|
||||
}
|
||||
for tip in tips {
|
||||
view.storage
|
||||
.retain_tip(&view.key, &tip.reference, &tip.oid)?;
|
||||
view.storage
|
||||
.advertise_base_tip(&view.key, &tip.reference, &tip.oid)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Promote every owed tip and return how many remain owed.
|
||||
///
|
||||
/// A tip whose object is in neither store belongs to a push that never
|
||||
/// completed and is dropped. Failures leave their tips owed for a retry.
|
||||
pub fn settle(view: &Path) -> Result<usize> {
|
||||
let Some(view) = View::resolve(view)? else {
|
||||
return Ok(0);
|
||||
};
|
||||
settle_resolved(&view)
|
||||
}
|
||||
|
||||
fn settle_resolved(view: &View) -> Result<usize> {
|
||||
let Some(mut record) = view.load()? else {
|
||||
return Ok(0);
|
||||
};
|
||||
if record.owed.is_empty() {
|
||||
return Ok(0);
|
||||
}
|
||||
let mut present = Vec::new();
|
||||
for tip in std::mem::take(&mut record.owed) {
|
||||
if object_exists(&view.path, &tip.oid)? {
|
||||
present.push(tip);
|
||||
} else {
|
||||
tracing::debug!(
|
||||
view = %view.path.display(),
|
||||
reference = %tip.reference,
|
||||
oid = %tip.oid,
|
||||
"Dropping owed tip of a push that did not complete"
|
||||
);
|
||||
}
|
||||
}
|
||||
// One fetch covers the usual case. After a failure retry each tip, so one
|
||||
// unavailable tip cannot keep independent history out of the family.
|
||||
if let Err(error) = promote_resolved(view, &present) {
|
||||
let batch = present.len() > 1;
|
||||
for tip in present {
|
||||
let failed = if batch {
|
||||
promote_resolved(view, std::slice::from_ref(&tip)).err()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if !batch || failed.is_some() {
|
||||
tracing::warn!(
|
||||
view = %view.path.display(),
|
||||
reference = %tip.reference,
|
||||
oid = %tip.oid,
|
||||
error = %failed.as_ref().unwrap_or(&error),
|
||||
"Signed history remains staged; promotion will retry"
|
||||
);
|
||||
record.owed.push(tip);
|
||||
}
|
||||
}
|
||||
}
|
||||
view.save(&record)?;
|
||||
Ok(record.owed.len())
|
||||
}
|
||||
|
||||
/// Shrink staging to the history that live refs need and the family lacks.
|
||||
///
|
||||
/// Refuses to run while any tip is owed. Every live ref is a root, so a
|
||||
/// pending upload keeps exactly its own history.
|
||||
pub fn compact(view: &Path) -> Result<Compaction> {
|
||||
let Some(view) = View::resolve(view)? else {
|
||||
return Ok(Compaction::Done);
|
||||
};
|
||||
if !view.record.is_file() {
|
||||
return Ok(Compaction::Done);
|
||||
}
|
||||
let owed = settle_resolved(&view)?;
|
||||
if owed > 0 {
|
||||
bail!("{owed} signed tips are still owed to the family");
|
||||
}
|
||||
run(
|
||||
&view.path,
|
||||
&["repack", "-a", "-d", "-l", "-q", "--no-write-bitmap-index"],
|
||||
)?;
|
||||
run(&view.path, &["prune", "--expire=now"])?;
|
||||
remove_loose_copies(&view)?;
|
||||
if has_local_objects(&view.path)? {
|
||||
return Ok(Compaction::Pending);
|
||||
}
|
||||
view.unregister()?;
|
||||
Ok(Compaction::Done)
|
||||
}
|
||||
|
||||
/// Remove loose objects that the family also holds.
|
||||
///
|
||||
/// Repacking leaves them behind: they are reachable, so pruning keeps them,
|
||||
/// and `--local` keeps them out of the new pack.
|
||||
fn remove_loose_copies(view: &View) -> Result<()> {
|
||||
let objects = view.path.join("objects");
|
||||
let mut loose = Vec::new();
|
||||
for entry in std::fs::read_dir(&objects)? {
|
||||
let entry = entry?;
|
||||
let prefix = entry.file_name();
|
||||
let Some(prefix) = prefix.to_str().filter(|name| name.len() == 2) else {
|
||||
continue;
|
||||
};
|
||||
for object in std::fs::read_dir(entry.path())? {
|
||||
let object = object?;
|
||||
if let Some(rest) = object.file_name().to_str() {
|
||||
let oid = format!("{prefix}{rest}");
|
||||
if valid_oid(&oid) {
|
||||
loose.push((oid, object.path()));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if loose.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
let mut input = Vec::new();
|
||||
for (oid, _) in &loose {
|
||||
input.extend_from_slice(oid.as_bytes());
|
||||
input.push(b'\n');
|
||||
}
|
||||
let report = git_output(
|
||||
&view.family(),
|
||||
&["cat-file", "--batch-check=%(objectname)"],
|
||||
&input,
|
||||
)?;
|
||||
ensure!(
|
||||
report.lines().count() == loose.len(),
|
||||
"unexpected object report from the family"
|
||||
);
|
||||
for ((oid, path), line) in loose.iter().zip(report.lines()) {
|
||||
if line == oid {
|
||||
std::fs::remove_file(path)?;
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn has_local_objects(view: &Path) -> Result<bool> {
|
||||
for entry in std::fs::read_dir(view.join("objects"))? {
|
||||
let entry = entry?;
|
||||
let name = entry.file_name();
|
||||
if name == "info" || !entry.file_type()?.is_dir() {
|
||||
continue;
|
||||
}
|
||||
if std::fs::read_dir(entry.path())?.next().is_some() {
|
||||
return Ok(true);
|
||||
}
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
|
||||
/// Promote an accepted tip out of staging, taking the family lease.
|
||||
///
|
||||
/// A view without staged objects received its history into the family, so
|
||||
/// there is nothing to copy and no lease is taken.
|
||||
pub async fn promote_accepted(view: &Path, tip: Tip) -> Result<()> {
|
||||
if !is_staged(view) {
|
||||
return Ok(());
|
||||
}
|
||||
let Some(resolved) = View::resolve(view)? else {
|
||||
return Ok(());
|
||||
};
|
||||
let lease = resolved.storage.write_lease(&resolved.key).await?;
|
||||
let path = resolved.path.clone();
|
||||
tokio::task::spawn_blocking(move || {
|
||||
let _lease = lease;
|
||||
promote(&path, std::slice::from_ref(&tip))
|
||||
})
|
||||
.await??;
|
||||
request_maintenance(view);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
@@ -0,0 +1,307 @@
|
||||
use super::*;
|
||||
use std::process::Command;
|
||||
|
||||
pub(super) struct Fixture {
|
||||
_root: tempfile::TempDir,
|
||||
pub storage: LocalGitStorage,
|
||||
pub key: FamilyKey,
|
||||
pub view: PathBuf,
|
||||
}
|
||||
|
||||
impl Fixture {
|
||||
pub fn new() -> Self {
|
||||
Self::with_view("owner")
|
||||
}
|
||||
|
||||
/// A view at `<parent>/repo.git` below the Git data root.
|
||||
pub fn with_view(parent: &str) -> Self {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
let storage = LocalGitStorage::new(root.path().canonicalize().unwrap());
|
||||
let key = FamilyKey::sha1("repo").unwrap();
|
||||
storage.ensure_family(&key).unwrap();
|
||||
let view = storage.git_data_path().join(parent).join("repo.git");
|
||||
std::fs::create_dir_all(view.parent().unwrap()).unwrap();
|
||||
storage.create_thin_view(&key, &view).unwrap();
|
||||
Self {
|
||||
_root: root,
|
||||
storage,
|
||||
key,
|
||||
view,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn family(&self) -> PathBuf {
|
||||
self.storage.family_repo_path(&self.key)
|
||||
}
|
||||
|
||||
/// Whether the family holds `oid` and its whole history by itself.
|
||||
fn family_holds(&self, oid: &str) -> bool {
|
||||
Command::new("git")
|
||||
.current_dir(self.family())
|
||||
.args(["rev-list", "--objects", "--missing=error", oid])
|
||||
.output()
|
||||
.unwrap()
|
||||
.status
|
||||
.success()
|
||||
}
|
||||
}
|
||||
|
||||
fn output(repo: &Path, args: &[&str], input: &[u8]) -> String {
|
||||
let mut child = Command::new("git")
|
||||
.current_dir(repo)
|
||||
.args(args)
|
||||
.env("GIT_AUTHOR_NAME", "Test")
|
||||
.env("GIT_AUTHOR_EMAIL", "test@example.com")
|
||||
.env("GIT_COMMITTER_NAME", "Test")
|
||||
.env("GIT_COMMITTER_EMAIL", "test@example.com")
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::piped())
|
||||
.spawn()
|
||||
.unwrap();
|
||||
child.stdin.take().unwrap().write_all(input).unwrap();
|
||||
let output = child.wait_with_output().unwrap();
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"git {args:?}: {}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
String::from_utf8(output.stdout).unwrap().trim().to_owned()
|
||||
}
|
||||
|
||||
/// Write a commit with one unique blob into `repo`'s own object directory.
|
||||
pub(super) fn commit(repo: &Path, content: &str, parent: Option<&str>) -> String {
|
||||
let blob = output(repo, &["hash-object", "-w", "--stdin"], content.as_bytes());
|
||||
let tree = output(
|
||||
repo,
|
||||
&["mktree"],
|
||||
format!("100644 blob {blob}\tfile\n").as_bytes(),
|
||||
);
|
||||
let mut args = vec!["commit-tree", tree.as_str(), "-m", content];
|
||||
if let Some(parent) = parent {
|
||||
args.extend(["-p", parent]);
|
||||
}
|
||||
output(repo, &args, b"")
|
||||
}
|
||||
|
||||
pub(super) fn reference(repo: &Path, name: &str, oid: &str) {
|
||||
output(repo, &["update-ref", name, oid], b"");
|
||||
}
|
||||
|
||||
pub(super) fn delete_reference(repo: &Path, name: &str) {
|
||||
output(repo, &["update-ref", "-d", name], b"");
|
||||
}
|
||||
|
||||
fn family_objects(fixture: &Fixture) -> String {
|
||||
output(
|
||||
&fixture.family(),
|
||||
&["cat-file", "--batch-all-objects", "--batch-check"],
|
||||
b"",
|
||||
)
|
||||
}
|
||||
|
||||
pub(super) const PENDING: &str =
|
||||
"refs/nostr/1111111111111111111111111111111111111111111111111111111111111111";
|
||||
pub(super) const ABANDONED: &str =
|
||||
"refs/nostr/2222222222222222222222222222222222222222222222222222222222222222";
|
||||
|
||||
#[test]
|
||||
fn abandoned_upload_is_reclaimed_without_touching_the_family() {
|
||||
let fixture = Fixture::new();
|
||||
let before = family_objects(&fixture);
|
||||
assert!(stage(&fixture.view, &[]).unwrap());
|
||||
let tip = commit(&fixture.view, "abandoned", None);
|
||||
reference(&fixture.view, ABANDONED, &tip);
|
||||
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Pending);
|
||||
assert!(object_exists(&fixture.view, &tip).unwrap());
|
||||
assert!(is_staged(&fixture.view));
|
||||
|
||||
delete_reference(&fixture.view, ABANDONED);
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Done);
|
||||
assert!(!object_exists(&fixture.view, &tip).unwrap());
|
||||
assert!(!is_staged(&fixture.view));
|
||||
assert_eq!(family_objects(&fixture), before);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn pending_upload_keeps_its_history_when_a_sibling_is_reclaimed() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let parent = commit(&fixture.view, "pending parent", None);
|
||||
let pending = commit(&fixture.view, "pending", Some(&parent));
|
||||
let abandoned = commit(&fixture.view, "abandoned", None);
|
||||
reference(&fixture.view, PENDING, &pending);
|
||||
reference(&fixture.view, ABANDONED, &abandoned);
|
||||
// Pack both uploads together so reclaiming one must rewrite the pack.
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Pending);
|
||||
|
||||
delete_reference(&fixture.view, ABANDONED);
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Pending);
|
||||
|
||||
assert!(!object_exists(&fixture.view, &abandoned).unwrap());
|
||||
output(
|
||||
&fixture.view,
|
||||
&["rev-list", "--objects", "--missing=error", &pending],
|
||||
b"",
|
||||
);
|
||||
assert!(!fixture.family_holds(&pending));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn promotion_makes_history_complete_in_the_family_alone() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let parent = commit(&fixture.view, "parent", None);
|
||||
let tip = commit(&fixture.view, "tip", Some(&parent));
|
||||
let unrelated = commit(&fixture.view, "unrelated", None);
|
||||
reference(&fixture.view, PENDING, &tip);
|
||||
reference(&fixture.view, ABANDONED, &unrelated);
|
||||
|
||||
promote(&fixture.view, &[Tip::new(PENDING, &tip)]).unwrap();
|
||||
|
||||
assert!(fixture.family_holds(&tip));
|
||||
assert!(
|
||||
!object_exists(&fixture.family(), &unrelated).unwrap(),
|
||||
"promotion copies only the accepted history"
|
||||
);
|
||||
delete_reference(&fixture.view, ABANDONED);
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Done);
|
||||
// The view still serves the promoted ref, now through the family.
|
||||
output(
|
||||
&fixture.view,
|
||||
&["rev-list", "--objects", "--missing=error", &tip],
|
||||
b"",
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn promotion_refuses_incomplete_history() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let parent = commit(&fixture.view, "parent", None);
|
||||
let tip = commit(&fixture.view, "tip", Some(&parent));
|
||||
reference(&fixture.view, PENDING, &tip);
|
||||
std::fs::remove_file(
|
||||
fixture
|
||||
.view
|
||||
.join("objects")
|
||||
.join(&parent[..2])
|
||||
.join(&parent[2..]),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert!(promote(&fixture.view, &[Tip::new(PENDING, &tip)]).is_err());
|
||||
assert!(!fixture.family_holds(&tip));
|
||||
let retained = output(
|
||||
&fixture.family(),
|
||||
&["for-each-ref", "refs/grasp/retained/"],
|
||||
b"",
|
||||
);
|
||||
assert_eq!(retained, "", "no retention root for incomplete history");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn owed_history_survives_ref_deletion_until_the_family_holds_it() {
|
||||
let fixture = Fixture::new();
|
||||
let tip_oid = {
|
||||
// Record the signed tip before its objects exist, as a push does.
|
||||
let probe = Fixture::new();
|
||||
commit(&probe.view, "signed", None)
|
||||
};
|
||||
let owed = Tip::new("refs/heads/main", &tip_oid);
|
||||
stage(&fixture.view, std::slice::from_ref(&owed)).unwrap();
|
||||
assert_eq!(commit(&fixture.view, "signed", None), tip_oid);
|
||||
reference(&fixture.view, "refs/heads/main", &tip_oid);
|
||||
|
||||
// Promotion cannot install its retention root while this lock exists.
|
||||
let digest = sha256::Hash::hash(format!("refs/heads/main\0{tip_oid}").as_bytes());
|
||||
let lock = fixture
|
||||
.family()
|
||||
.join(format!("refs/grasp/retained/{digest}.lock"));
|
||||
std::fs::create_dir_all(lock.parent().unwrap()).unwrap();
|
||||
std::fs::write(&lock, b"").unwrap();
|
||||
assert_eq!(settle(&fixture.view).unwrap(), 1);
|
||||
|
||||
// A State deletion removes the ref. Rollback still needs the history.
|
||||
delete_reference(&fixture.view, "refs/heads/main");
|
||||
assert!(
|
||||
compact(&fixture.view).is_err(),
|
||||
"owed tips block compaction"
|
||||
);
|
||||
assert!(object_exists(&fixture.view, &tip_oid).unwrap());
|
||||
assert!(is_staged(&fixture.view));
|
||||
|
||||
std::fs::remove_file(&lock).unwrap();
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Done);
|
||||
assert!(fixture.family_holds(&tip_oid));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn owed_tip_of_a_push_that_never_completed_is_dropped() {
|
||||
let fixture = Fixture::new();
|
||||
let absent = "0123456789012345678901234567890123456789";
|
||||
stage(&fixture.view, &[Tip::new("refs/heads/main", absent)]).unwrap();
|
||||
|
||||
assert_eq!(settle(&fixture.view).unwrap(), 0);
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Done);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn one_unavailable_owed_tip_does_not_hold_back_another() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let parent = commit(&fixture.view, "broken parent", None);
|
||||
let broken = commit(&fixture.view, "broken", Some(&parent));
|
||||
let sound = commit(&fixture.view, "sound", None);
|
||||
stage(
|
||||
&fixture.view,
|
||||
&[
|
||||
Tip::new("refs/heads/broken", &broken),
|
||||
Tip::new("refs/heads/sound", &sound),
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
std::fs::remove_file(
|
||||
fixture
|
||||
.view
|
||||
.join("objects")
|
||||
.join(&parent[..2])
|
||||
.join(&parent[2..]),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(settle(&fixture.view).unwrap(), 1);
|
||||
assert!(fixture.family_holds(&sound));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legacy_repository_without_a_family_is_never_staged() {
|
||||
let root = tempfile::tempdir().unwrap();
|
||||
output(root.path(), &["init", "--bare", "legacy.git"], b"");
|
||||
let legacy = root.path().join("legacy.git");
|
||||
|
||||
assert!(!stage(&legacy, &[]).unwrap());
|
||||
assert!(!is_staged(&legacy));
|
||||
assert_eq!(compact(&legacy).unwrap(), Compaction::Done);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn objects_held_before_first_staging_move_to_the_family() {
|
||||
let fixture = Fixture::new();
|
||||
let packed = commit(&fixture.view, "packed before staging", None);
|
||||
reference(&fixture.view, "refs/heads/packed", &packed);
|
||||
output(&fixture.view, &["repack", "-a", "-d", "-l", "-q"], b"");
|
||||
delete_reference(&fixture.view, "refs/heads/packed");
|
||||
let loose = commit(&fixture.view, "loose before staging", None);
|
||||
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
|
||||
assert!(fixture.family_holds(&packed));
|
||||
assert!(fixture.family_holds(&loose));
|
||||
assert!(!has_local_objects(&fixture.view).unwrap());
|
||||
// Neither commit has a ref, yet compaction cannot reach them any more.
|
||||
assert_eq!(compact(&fixture.view).unwrap(), Compaction::Done);
|
||||
assert!(fixture.family_holds(&packed));
|
||||
assert!(fixture.family_holds(&loose));
|
||||
}
|
||||
@@ -0,0 +1,266 @@
|
||||
//! Event-driven staging maintenance.
|
||||
//!
|
||||
//! Pushes, promotions and ref deletions request maintenance for their view.
|
||||
//! The persistent registry restores those requests after a restart. A view
|
||||
//! whose remaining staged objects belong to pending refs is parked until the
|
||||
//! next request rather than polled.
|
||||
use std::collections::HashMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::{LazyLock, Mutex};
|
||||
use std::time::Duration;
|
||||
|
||||
use anyhow::{Context, Result};
|
||||
use tokio::sync::Notify;
|
||||
use tokio::time::Instant;
|
||||
|
||||
use super::{compact, is_staged, registry_dir, Compaction, Record, View};
|
||||
use crate::git::storage::LocalGitStorage;
|
||||
use crate::grasp06::receive::{remove_prs_repo_if_empty, RepoInitLocks};
|
||||
|
||||
/// How long maintenance queues behind active writers of the family before it
|
||||
/// gives way. The lease is fair, so waiting longer would hold up later pushes.
|
||||
const LEASE_WAIT: Duration = Duration::from_millis(250);
|
||||
const FIRST_RETRY: Duration = Duration::from_secs(1);
|
||||
const LAST_RETRY: Duration = Duration::from_secs(300);
|
||||
|
||||
struct Work {
|
||||
due: Instant,
|
||||
/// Delay before the attempt after next; doubles on each failure.
|
||||
retry: Duration,
|
||||
}
|
||||
|
||||
static QUEUE: LazyLock<Mutex<HashMap<PathBuf, Work>>> = LazyLock::new(Mutex::default);
|
||||
static WAKE: Notify = Notify::const_new();
|
||||
|
||||
fn queue() -> std::sync::MutexGuard<'static, HashMap<PathBuf, Work>> {
|
||||
QUEUE
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
}
|
||||
|
||||
/// Ask the worker to settle and compact a view. Cheap for unstaged views.
|
||||
pub fn request_maintenance(view: &Path) {
|
||||
if !is_staged(view) {
|
||||
return;
|
||||
}
|
||||
let Ok(view) = std::fs::canonicalize(view) else {
|
||||
return;
|
||||
};
|
||||
queue().insert(
|
||||
view,
|
||||
Work {
|
||||
due: Instant::now(),
|
||||
retry: FIRST_RETRY,
|
||||
},
|
||||
);
|
||||
WAKE.notify_waiters();
|
||||
}
|
||||
|
||||
/// Queue every registered view and drop records of views that no longer exist.
|
||||
fn recover(storage: &LocalGitStorage) -> Result<()> {
|
||||
let entries = match std::fs::read_dir(registry_dir(storage)) {
|
||||
Ok(entries) => entries,
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
|
||||
Err(error) => return Err(error.into()),
|
||||
};
|
||||
for entry in entries {
|
||||
let path = entry?.path();
|
||||
if path.extension().is_none_or(|extension| extension != "json") {
|
||||
continue;
|
||||
}
|
||||
let result = (|| -> Result<()> {
|
||||
let record: Record = serde_json::from_slice(&std::fs::read(&path)?)?;
|
||||
anyhow::ensure!(
|
||||
!record.view.as_os_str().is_empty()
|
||||
&& record
|
||||
.view
|
||||
.components()
|
||||
.all(|part| matches!(part, std::path::Component::Normal(_))),
|
||||
"invalid view path"
|
||||
);
|
||||
let view = storage.git_data_path().join(&record.view);
|
||||
if view.try_exists()? {
|
||||
request_maintenance(&view);
|
||||
} else {
|
||||
std::fs::remove_file(&path)?;
|
||||
}
|
||||
Ok(())
|
||||
})();
|
||||
if let Err(error) = result {
|
||||
tracing::warn!(record = %path.display(), %error, "Cannot recover staged Git view");
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// One maintenance attempt. An error means the attempt should be retried.
|
||||
pub(super) async fn maintain(view: &Path, prs_locks: &RepoInitLocks) -> Result<Compaction> {
|
||||
let Some(resolved) = View::resolve(view)? else {
|
||||
return Ok(Compaction::Done);
|
||||
};
|
||||
let lease = tokio::time::timeout(LEASE_WAIT, resolved.storage.write_lease(&resolved.key))
|
||||
.await
|
||||
.context("family is busy")??;
|
||||
let path = resolved.path.clone();
|
||||
let outcome = tokio::task::spawn_blocking(move || {
|
||||
let _lease = lease;
|
||||
compact(&path)
|
||||
})
|
||||
.await??;
|
||||
if outcome == Compaction::Done {
|
||||
remove_prs_repo_if_empty(prs_locks, resolved.storage.git_data_path(), &resolved.path);
|
||||
}
|
||||
Ok(outcome)
|
||||
}
|
||||
|
||||
/// Run staging maintenance until the task is aborted.
|
||||
pub async fn run_worker(storage: LocalGitStorage, prs_locks: RepoInitLocks) {
|
||||
let root = match std::fs::canonicalize(storage.git_data_path()) {
|
||||
Ok(root) => root,
|
||||
Err(error) => {
|
||||
tracing::warn!(%error, "Cannot start staging maintenance");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let recovery = storage.clone();
|
||||
match tokio::task::spawn_blocking(move || recover(&recovery)).await {
|
||||
Ok(Ok(())) => {}
|
||||
result => tracing::warn!(?result, "Staging registry recovery failed"),
|
||||
}
|
||||
loop {
|
||||
// Register for wake-ups before reading the queue, so a request made
|
||||
// between the read and the wait is not missed.
|
||||
let woken = WAKE.notified();
|
||||
tokio::pin!(woken);
|
||||
woken.as_mut().enable();
|
||||
let now = Instant::now();
|
||||
let next = {
|
||||
let mut queue = queue();
|
||||
let next = queue
|
||||
.iter()
|
||||
.filter(|(view, _)| view.starts_with(&root))
|
||||
.min_by_key(|(_, work)| work.due)
|
||||
.map(|(view, work)| (view.clone(), work.due, work.retry));
|
||||
if let Some((view, due, _)) = &next {
|
||||
if *due <= now {
|
||||
queue.remove(view);
|
||||
}
|
||||
}
|
||||
next
|
||||
};
|
||||
match next {
|
||||
Some((view, due, retry)) if due <= now => {
|
||||
if !view.exists() {
|
||||
continue;
|
||||
}
|
||||
if let Err(error) = maintain(&view, &prs_locks).await {
|
||||
tracing::debug!(view = %view.display(), %error, "Staging maintenance will retry");
|
||||
// A request that arrived meanwhile is newer than this failure.
|
||||
queue().entry(view).or_insert(Work {
|
||||
due: Instant::now() + retry,
|
||||
retry: (retry * 2).min(LAST_RETRY),
|
||||
});
|
||||
}
|
||||
}
|
||||
Some((_, due, _)) => {
|
||||
tokio::select! {
|
||||
_ = tokio::time::sleep_until(due) => {}
|
||||
_ = woken => {}
|
||||
}
|
||||
}
|
||||
None => woken.await,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::super::tests::{commit, delete_reference, reference, Fixture, ABANDONED, PENDING};
|
||||
use super::super::{object_exists, stage};
|
||||
use super::*;
|
||||
use crate::grasp06::receive::new_repo_init_locks;
|
||||
|
||||
const DEADLINE: Duration = Duration::from_secs(10);
|
||||
|
||||
async fn wait_until(description: &str, condition: impl Fn() -> bool) {
|
||||
tokio::time::timeout(DEADLINE, async {
|
||||
while !condition() {
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn ref_deletion_wakes_the_worker_to_reclaim_an_abandoned_upload() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let pending = commit(&fixture.view, "pending", None);
|
||||
let abandoned = commit(&fixture.view, "abandoned", None);
|
||||
reference(&fixture.view, PENDING, &pending);
|
||||
reference(&fixture.view, ABANDONED, &abandoned);
|
||||
let worker = tokio::spawn(run_worker(fixture.storage.clone(), new_repo_init_locks()));
|
||||
|
||||
crate::git::delete_ref(&fixture.view, ABANDONED).unwrap();
|
||||
|
||||
let view = fixture.view.clone();
|
||||
wait_until("the abandoned upload to be reclaimed", || {
|
||||
!object_exists(&view, &abandoned).unwrap()
|
||||
})
|
||||
.await;
|
||||
assert!(object_exists(&fixture.view, &pending).unwrap());
|
||||
assert!(is_staged(&fixture.view), "a pending upload stays staged");
|
||||
worker.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn restart_recovers_a_view_registered_before_it() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let abandoned = commit(&fixture.view, "abandoned", None);
|
||||
reference(&fixture.view, ABANDONED, &abandoned);
|
||||
// The ref expired without a maintenance request, as after a crash.
|
||||
delete_reference(&fixture.view, ABANDONED);
|
||||
|
||||
let worker = tokio::spawn(run_worker(fixture.storage.clone(), new_repo_init_locks()));
|
||||
|
||||
let view = fixture.view.clone();
|
||||
wait_until("the recovered view to be reclaimed", || !is_staged(&view)).await;
|
||||
assert!(!object_exists(&fixture.view, &abandoned).unwrap());
|
||||
worker.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintenance_gives_way_to_a_busy_family_and_succeeds_later() {
|
||||
let fixture = Fixture::new();
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
let abandoned = commit(&fixture.view, "abandoned", None);
|
||||
let locks = new_repo_init_locks();
|
||||
|
||||
// A push holds the family lease for as long as it runs.
|
||||
let push = fixture.storage.write_lease(&fixture.key).await.unwrap();
|
||||
assert!(maintain(&fixture.view, &locks).await.is_err());
|
||||
assert!(object_exists(&fixture.view, &abandoned).unwrap());
|
||||
|
||||
drop(push);
|
||||
assert_eq!(
|
||||
maintain(&fixture.view, &locks).await.unwrap(),
|
||||
Compaction::Done
|
||||
);
|
||||
assert!(!object_exists(&fixture.view, &abandoned).unwrap());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn empty_prs_repository_is_removed_once_its_staging_is_reclaimed() {
|
||||
let submitter = "a".repeat(64);
|
||||
let fixture = Fixture::with_view(&format!("prs/{submitter}"));
|
||||
stage(&fixture.view, &[]).unwrap();
|
||||
commit(&fixture.view, "abandoned", None);
|
||||
|
||||
let outcome = maintain(&fixture.view, &new_repo_init_locks()).await;
|
||||
|
||||
assert_eq!(outcome.unwrap(), Compaction::Done);
|
||||
assert!(!fixture.view.exists());
|
||||
}
|
||||
}
|
||||
@@ -1509,6 +1509,16 @@ async fn process_purgatory_pr_events(
|
||||
let owner_pubkey =
|
||||
extract_owner_from_repo_path(source_repo_path, git_data_path).unwrap_or_default();
|
||||
|
||||
// The event may only be accepted once the family holds its history.
|
||||
let tip = crate::git::staging::Tip::new(format!("refs/nostr/{}", event.id), &entry.commit);
|
||||
if let Err(error) = crate::git::staging::promote_accepted(source_repo_path, tip).await {
|
||||
result.errors.push(format!(
|
||||
"PR {} history could not be stored durably: {error:#}",
|
||||
event.id
|
||||
));
|
||||
continue;
|
||||
}
|
||||
|
||||
// Use unified processing function
|
||||
let process_result = crate::git::process::process_pr_with_git_data(
|
||||
event,
|
||||
|
||||
@@ -110,6 +110,8 @@ pub fn scan_on_startup(git_data_path: &Path) -> (usize, usize) {
|
||||
}
|
||||
|
||||
match list_refs(&repo_path) {
|
||||
// Signed history still owed to the family outlives its refs.
|
||||
Ok(refs) if refs.is_empty() && crate::git::staging::owes_history(&repo_path) => {}
|
||||
Ok(refs) if refs.is_empty() => {
|
||||
if let Err(e) = std::fs::remove_dir_all(&repo_path) {
|
||||
warn!(
|
||||
|
||||
+31
-5
@@ -79,8 +79,8 @@ use crate::git::authorization::{
|
||||
use crate::git::handlers::{
|
||||
build_git_protocol_error_response, err_pktline_frame, is_git_protocol_error,
|
||||
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,
|
||||
retain_accepted_tips, send_body_bytes, settle_staged_push, stage_push, streaming_response,
|
||||
GitError, PumpResult, STREAM_CHANNEL_DEPTH,
|
||||
};
|
||||
use crate::git::protocol::GitService;
|
||||
use crate::git::storage::{FamilyKey, FamilyWriteLease, LocalGitStorage};
|
||||
@@ -231,6 +231,7 @@ pub async fn handle_prs_receive_pack(
|
||||
identifier: &prs.identifier,
|
||||
domain,
|
||||
};
|
||||
let mut unsigned_refs = HashSet::new();
|
||||
for (_, new_oid, ref_name) in &pushed_refs {
|
||||
match pre_validate_refs_nostr_push(
|
||||
&database,
|
||||
@@ -255,7 +256,11 @@ pub async fn handle_prs_receive_pack(
|
||||
Some(&request_body),
|
||||
));
|
||||
}
|
||||
NostrRefPreValidation::Authorized { .. } | NostrRefPreValidation::Unknown => {}
|
||||
NostrRefPreValidation::Authorized { signed: true, .. } => {}
|
||||
NostrRefPreValidation::Authorized { signed: false, .. }
|
||||
| NostrRefPreValidation::Unknown => {
|
||||
unsigned_refs.insert(ref_name.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -301,12 +306,24 @@ pub async fn handle_prs_receive_pack(
|
||||
// write stdin here, then hand stdout/stderr and all follow-up state to a
|
||||
// detached task that feeds the response body channel.
|
||||
let use_family = storage.is_thin_view(&family_key, &repo_path);
|
||||
let staged = if use_family {
|
||||
match stage_push(&repo_path, &pushed_refs, &unsigned_refs).await {
|
||||
Ok(staged) => staged,
|
||||
Err(e) => {
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
record_git_operation(&metrics, "push", "error");
|
||||
return Err(e);
|
||||
}
|
||||
}
|
||||
} else {
|
||||
false
|
||||
};
|
||||
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),
|
||||
(use_family && !staged).then_some(&family_lease.family_objects_path),
|
||||
)
|
||||
.map_err(GitError::ProcessSpawnFailed)
|
||||
{
|
||||
@@ -370,6 +387,7 @@ pub async fn handle_prs_receive_pack(
|
||||
storage,
|
||||
family_key,
|
||||
use_family.then_some(family_lease),
|
||||
staged,
|
||||
));
|
||||
|
||||
Ok(streaming_response(GitService::ReceivePack, rx))
|
||||
@@ -447,6 +465,7 @@ fn finish_prs_receive_pack(state: &PrsPathState, repo_path: &Path) {
|
||||
pub(crate) fn remove_idle_empty_repo(state: &PrsPathState, repo_path: &Path) -> bool {
|
||||
if state.in_flight.load(Ordering::Relaxed) != 0
|
||||
|| !matches!(list_refs(repo_path), Ok(refs) if refs.is_empty())
|
||||
|| crate::git::staging::owes_history(repo_path)
|
||||
{
|
||||
return false;
|
||||
}
|
||||
@@ -518,6 +537,7 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
storage: LocalGitStorage,
|
||||
family_key: FamilyKey,
|
||||
family_lease: Option<FamilyWriteLease>,
|
||||
staged: bool,
|
||||
) where
|
||||
S: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
E: tokio::io::AsyncRead + Unpin + Send + 'static,
|
||||
@@ -637,10 +657,16 @@ async fn stream_prs_receive_pack_output<S, E>(
|
||||
.await;
|
||||
}
|
||||
|
||||
if family_lease.is_some() {
|
||||
if staged {
|
||||
settle_staged_push(&repo_path).await;
|
||||
} else if family_lease.is_some() {
|
||||
retain_accepted_tips(&storage, &family_key, &repo_path, &pushed_refs);
|
||||
}
|
||||
|
||||
// Release the family writer before post-push processing: promoting a
|
||||
// waiting event may need to re-enter this same identifier family.
|
||||
drop(family_lease);
|
||||
|
||||
finish_prs_receive_pack(&state, &repo_path);
|
||||
|
||||
// Drive the standard purgatory-release pipeline so PR events already
|
||||
|
||||
@@ -65,6 +65,9 @@ impl PrEventPolicy {
|
||||
}
|
||||
};
|
||||
|
||||
// A standard-endpoint placeholder is released only once the event's
|
||||
// history is durable or no local upload serves the event.
|
||||
let mut standard_placeholder = false;
|
||||
// A crash may happen after Git installs the ref but before the
|
||||
// purgatory checkpoint. Recover the /prs/ scope only from a signed,
|
||||
// service-local clone URL for this signer and an exact event-id/OID ref.
|
||||
@@ -223,6 +226,11 @@ impl PrEventPolicy {
|
||||
)
|
||||
.await?;
|
||||
|
||||
// The placeholder stays until the history is durable, so a
|
||||
// failure here leaves the upload to expire normally.
|
||||
self.promote_staged_history(&prs_repo, &event_id, &commit)
|
||||
.await?;
|
||||
|
||||
let process_result = crate::git::process::process_pr_with_git_data(
|
||||
event,
|
||||
&commit,
|
||||
@@ -279,7 +287,7 @@ impl PrEventPolicy {
|
||||
event_id,
|
||||
commit
|
||||
);
|
||||
self.ctx.purgatory.remove_pr(&event_id);
|
||||
standard_placeholder = true;
|
||||
} else {
|
||||
// Standard endpoint placeholder, mismatched commit — original
|
||||
// behaviour: incoming event supersedes.
|
||||
@@ -289,7 +297,7 @@ impl PrEventPolicy {
|
||||
commit,
|
||||
placeholder_commit
|
||||
);
|
||||
self.ctx.purgatory.remove_pr(&event_id);
|
||||
standard_placeholder = true;
|
||||
// Delete incorrect git data (refs/nostr/<event-id>) will be handled below
|
||||
}
|
||||
}
|
||||
@@ -297,6 +305,9 @@ impl PrEventPolicy {
|
||||
let repo_paths = self.find_relevant_repo_paths(event).await?;
|
||||
|
||||
if repo_paths.is_empty() {
|
||||
if standard_placeholder {
|
||||
self.ctx.purgatory.remove_pr(&event_id);
|
||||
}
|
||||
tracing::debug!("No repository paths found for PR event {}", event_id);
|
||||
return Ok(false);
|
||||
}
|
||||
@@ -357,6 +368,14 @@ impl PrEventPolicy {
|
||||
)
|
||||
.unwrap_or_default();
|
||||
|
||||
// The placeholder stays until the history is durable, so a
|
||||
// failure here leaves the upload to expire normally.
|
||||
self.promote_staged_history(&source_repo, &event_id, &commit)
|
||||
.await?;
|
||||
if standard_placeholder {
|
||||
self.ctx.purgatory.remove_pr(&event_id);
|
||||
}
|
||||
|
||||
// Use unified processing function
|
||||
let result = crate::git::process::process_pr_with_git_data(
|
||||
event,
|
||||
@@ -388,6 +407,9 @@ impl PrEventPolicy {
|
||||
|
||||
Ok(true)
|
||||
} else {
|
||||
if standard_placeholder {
|
||||
self.ctx.purgatory.remove_pr(&event_id);
|
||||
}
|
||||
tracing::debug!(
|
||||
"No git data found for PR event {} with commit {}",
|
||||
event_id,
|
||||
@@ -397,6 +419,31 @@ impl PrEventPolicy {
|
||||
}
|
||||
}
|
||||
|
||||
/// Make the history of an accepted PR tip durable in the family.
|
||||
///
|
||||
/// An upload received before its event was known is staged in its view.
|
||||
/// The event may only be accepted once the family alone holds the history.
|
||||
async fn promote_staged_history(
|
||||
&self,
|
||||
repo: &std::path::Path,
|
||||
event_id: &str,
|
||||
commit: &str,
|
||||
) -> Result<()> {
|
||||
let tip = git::staging::Tip::new(format!("refs/nostr/{event_id}"), commit);
|
||||
git::staging::promote_accepted(repo, tip)
|
||||
.await
|
||||
.map_err(|error| {
|
||||
tracing::error!(
|
||||
event_id = %event_id,
|
||||
commit = %commit,
|
||||
repo = %repo.display(),
|
||||
error = %format!("{error:#}"),
|
||||
"Cannot store PR history durably; event not accepted"
|
||||
);
|
||||
anyhow::anyhow!("PR history could not be stored durably")
|
||||
})
|
||||
}
|
||||
|
||||
async fn find_relevant_repo_paths(&self, event: &Event) -> Result<Vec<std::path::PathBuf>> {
|
||||
// Extract ALL `a` tags (repository references) from the PR event
|
||||
let repo_refs: Vec<String> = event
|
||||
|
||||
@@ -654,6 +654,9 @@ impl Purgatory {
|
||||
tracing::warn!(repo = %path.display(), %reference,
|
||||
"Failed to expire normal PR ref; retaining placeholder for retry");
|
||||
ok = false;
|
||||
} else {
|
||||
// The abandoned upload's staged objects can now be reclaimed.
|
||||
crate::git::staging::request_maintenance(&path);
|
||||
}
|
||||
}
|
||||
if ok && entry.prs_scope.is_some() {
|
||||
|
||||
@@ -386,6 +386,11 @@ impl RelayServer {
|
||||
"Crash-safe sync-state checkpoint task started"
|
||||
);
|
||||
|
||||
background_tasks.push(tokio::spawn(git::staging::run_worker(
|
||||
git::storage::LocalGitStorage::new(config.effective_git_data_path()),
|
||||
repo_init_locks.clone(),
|
||||
)));
|
||||
|
||||
// Spawn background cleanup task for purgatory entries (60s interval)
|
||||
let cleanup_purgatory = purgatory.clone();
|
||||
let cleanup_lifecycle = relay_runtime.lifecycle.clone();
|
||||
|
||||
@@ -0,0 +1,288 @@
|
||||
//! Unsigned uploads are staged in their view and earn family storage only
|
||||
//! when a signed event names them.
|
||||
//!
|
||||
//! ```bash
|
||||
//! cargo test --test pending_upload_staging
|
||||
//! ```
|
||||
|
||||
mod common;
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::time::Duration;
|
||||
|
||||
use common::{CommitVariant, TestRelay};
|
||||
use grasp_audit::{AuditClient, AuditConfig};
|
||||
use ngit_grasp::git::staging;
|
||||
use nostr_sdk::prelude::*;
|
||||
|
||||
const DEADLINE: Duration = Duration::from_secs(20);
|
||||
|
||||
fn git(repo: &Path, args: &[&str]) -> std::process::Output {
|
||||
grasp_audit::git_command()
|
||||
.current_dir(repo)
|
||||
.args(args)
|
||||
.output()
|
||||
.expect("spawn git")
|
||||
}
|
||||
|
||||
/// Every object in the family, which has no alternates of its own.
|
||||
fn family_objects(family: &Path) -> String {
|
||||
let output = git(
|
||||
family,
|
||||
&[
|
||||
"cat-file",
|
||||
"--batch-all-objects",
|
||||
"--batch-check=%(objectname)",
|
||||
],
|
||||
);
|
||||
assert!(output.status.success());
|
||||
String::from_utf8(output.stdout).unwrap()
|
||||
}
|
||||
|
||||
/// Object files in the view's own object directory, which is its staging.
|
||||
fn staged_files(view: &Path) -> Vec<PathBuf> {
|
||||
let mut files = Vec::new();
|
||||
for entry in std::fs::read_dir(view.join("objects")).unwrap() {
|
||||
let entry = entry.unwrap();
|
||||
if entry.file_name() == "info" || !entry.file_type().unwrap().is_dir() {
|
||||
continue;
|
||||
}
|
||||
for file in std::fs::read_dir(entry.path()).unwrap() {
|
||||
files.push(file.unwrap().path());
|
||||
}
|
||||
}
|
||||
files
|
||||
}
|
||||
|
||||
/// Whether `repo` holds `tip` and its whole history without any other store.
|
||||
fn holds_history(repo: &Path, tip: &str) -> bool {
|
||||
git(repo, &["rev-list", "--objects", "--missing=error", tip])
|
||||
.status
|
||||
.success()
|
||||
}
|
||||
|
||||
struct Served {
|
||||
_persistent: tempfile::TempDir,
|
||||
relay: TestRelay,
|
||||
client: AuditClient,
|
||||
announcement: Event,
|
||||
state: Event,
|
||||
identifier: String,
|
||||
npub: String,
|
||||
view: PathBuf,
|
||||
family: PathBuf,
|
||||
}
|
||||
|
||||
impl Served {
|
||||
/// A relay serving one repository whose maintainer push has completed.
|
||||
async fn start(name: &str) -> Self {
|
||||
let persistent = tempfile::tempdir().unwrap();
|
||||
let git_data = persistent.path().join("git");
|
||||
let relay = TestRelay::start_with_existing_lmdb_paths(
|
||||
git_data.clone(),
|
||||
persistent.path().join("relay"),
|
||||
None,
|
||||
false,
|
||||
)
|
||||
.await;
|
||||
let client = AuditClient::new(relay.url(), AuditConfig::isolated())
|
||||
.await
|
||||
.unwrap();
|
||||
let (announcement, identifier, state) =
|
||||
common::publish_served_audit_repo_with_state(&client, name).await;
|
||||
let npub = client.public_key().to_bech32().unwrap();
|
||||
let view = git_data.join(&npub).join(format!("{identifier}.git"));
|
||||
let family = git_data
|
||||
.join(".grasp/families/sha1")
|
||||
.join(format!("{identifier}.git"));
|
||||
Self {
|
||||
_persistent: persistent,
|
||||
relay,
|
||||
client,
|
||||
announcement,
|
||||
state,
|
||||
identifier,
|
||||
npub,
|
||||
view,
|
||||
family,
|
||||
}
|
||||
}
|
||||
|
||||
fn pr_event(&self, tip: &str, title: &str) -> Event {
|
||||
common::create_pr_event(
|
||||
self.client.keys(),
|
||||
&common::announcement_coordinate(&self.announcement, &self.identifier),
|
||||
tip,
|
||||
title,
|
||||
)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn push(&self, local: &Path, tip: &str, event: &Event) {
|
||||
common::push_ref_to_relay(
|
||||
local,
|
||||
&self.relay.domain(),
|
||||
&self.npub,
|
||||
&self.identifier,
|
||||
tip,
|
||||
&format!("refs/nostr/{}", event.id),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
async fn accept(&self, event: &Event) {
|
||||
let client = AuditClient::new(self.relay.url(), AuditConfig::isolated())
|
||||
.await
|
||||
.unwrap();
|
||||
client.send_event(event.clone()).await.unwrap();
|
||||
common::wait_for_event_served(self.relay.url(), &event.id, DEADLINE)
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
async fn wait_until_unstaged(&self) {
|
||||
common::wait_for("staging to be reclaimed", DEADLINE, || async {
|
||||
!staging::is_staged(&self.view) && staged_files(&self.view).is_empty()
|
||||
})
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn signed_push_is_received_into_the_family_without_staging() {
|
||||
let served = Served::start("staging-signed-push").await;
|
||||
|
||||
assert!(!staging::is_staged(&served.view));
|
||||
assert_eq!(staged_files(&served.view), Vec::<PathBuf>::new());
|
||||
let main = ngit_grasp::git::get_ref_commit(&served.view, "refs/heads/main").unwrap();
|
||||
assert!(holds_history(&served.family, &main));
|
||||
|
||||
served.relay.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unsigned_upload_is_staged_across_a_crash_and_promoted_on_acceptance() {
|
||||
let mut served = Served::start("staging-unsigned-upload").await;
|
||||
let family_before = family_objects(&served.family);
|
||||
let local = tempfile::tempdir().unwrap();
|
||||
let tip = common::create_test_repo_with_commit(local.path(), CommitVariant::PrTest).unwrap();
|
||||
let event = served.pr_event(&tip, "staged until signed");
|
||||
|
||||
served.push(local.path(), &tip, &event);
|
||||
|
||||
assert!(staging::is_staged(&served.view));
|
||||
assert!(ngit_grasp::git::oid_exists(&served.view, &tip));
|
||||
assert_eq!(
|
||||
family_objects(&served.family),
|
||||
family_before,
|
||||
"an unsigned upload reached the family"
|
||||
);
|
||||
|
||||
// Killing the relay loses any purgatory placeholder not yet checkpointed;
|
||||
// the staged ref and objects must not depend on it.
|
||||
served.relay = served.relay.restart().await;
|
||||
assert!(staging::is_staged(&served.view));
|
||||
assert!(holds_history(&served.view, &tip));
|
||||
assert_eq!(family_objects(&served.family), family_before);
|
||||
|
||||
served.accept(&event).await;
|
||||
|
||||
assert!(
|
||||
holds_history(&served.family, &tip),
|
||||
"an accepted PR must be complete in the family alone"
|
||||
);
|
||||
served.wait_until_unstaged().await;
|
||||
assert!(holds_history(&served.view, &tip));
|
||||
|
||||
served.relay.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn signed_push_onto_staged_history_is_complete_in_the_family() {
|
||||
let served = Served::start("staging-signed-onto-staged").await;
|
||||
let local = tempfile::tempdir().unwrap();
|
||||
let pending_tip =
|
||||
common::create_test_repo_with_commit(local.path(), CommitVariant::PrTest).unwrap();
|
||||
let pending = served.pr_event(&pending_tip, "never signed");
|
||||
served.push(local.path(), &pending_tip, &pending);
|
||||
assert!(staging::is_staged(&served.view));
|
||||
|
||||
// The view advertises the pending tip, so this push omits its objects.
|
||||
let signed_tip = common::add_commit_to_repo(local.path(), CommitVariant::SecondCommit).unwrap();
|
||||
let signed = served.pr_event(&signed_tip, "signed before its push");
|
||||
let client = AuditClient::new(served.relay.url(), AuditConfig::isolated())
|
||||
.await
|
||||
.unwrap();
|
||||
client
|
||||
.send_event_and_note_purgatory(signed.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
served.push(local.path(), &signed_tip, &signed);
|
||||
|
||||
common::wait_for_event_served(served.relay.url(), &signed.id, DEADLINE)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
holds_history(&served.family, &signed_tip),
|
||||
"signed history depends on objects that only staging holds"
|
||||
);
|
||||
// The pending ref is untouched. Its history is now an ancestor of signed
|
||||
// history, so the family holds it and nothing is left to stage.
|
||||
assert!(holds_history(&served.view, &pending_tip));
|
||||
served.wait_until_unstaged().await;
|
||||
|
||||
served.relay.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn maintainer_push_into_a_staged_view_reaches_the_family() {
|
||||
let served = Served::start("staging-maintainer-push").await;
|
||||
let local = tempfile::tempdir().unwrap();
|
||||
let pending_tip =
|
||||
common::create_test_repo_with_commit(local.path(), CommitVariant::PrTest).unwrap();
|
||||
let pending = served.pr_event(&pending_tip, "never signed");
|
||||
served.push(local.path(), &pending_tip, &pending);
|
||||
assert!(staging::is_staged(&served.view));
|
||||
|
||||
let clone =
|
||||
grasp_audit::clone_repo(&served.relay.domain(), &served.npub, &served.identifier).unwrap();
|
||||
std::fs::write(clone.join("next.txt"), "next state").unwrap();
|
||||
assert!(git(&clone, &["add", "."]).status.success());
|
||||
assert!(git(&clone, &["commit", "-m", "next state"])
|
||||
.status
|
||||
.success());
|
||||
let next = String::from_utf8(git(&clone, &["rev-parse", "HEAD"]).stdout)
|
||||
.unwrap()
|
||||
.trim()
|
||||
.to_owned();
|
||||
let state = common::event_ordering::event_after(
|
||||
common::create_state_event(
|
||||
served.client.keys(),
|
||||
&served.identifier,
|
||||
&[("main", &next)],
|
||||
&[],
|
||||
&[],
|
||||
&[],
|
||||
)
|
||||
.unwrap(),
|
||||
served.client.keys(),
|
||||
served.state.created_at,
|
||||
);
|
||||
served
|
||||
.client
|
||||
.send_event_and_note_purgatory(state.clone())
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(grasp_audit::try_push(&clone).unwrap());
|
||||
let _ = std::fs::remove_dir_all(&clone);
|
||||
|
||||
common::wait_for_event_served(served.relay.url(), &state.id, DEADLINE)
|
||||
.await
|
||||
.unwrap();
|
||||
// Rollback after a State deletion depends on this history, whatever
|
||||
// becomes of the view's refs.
|
||||
assert!(holds_history(&served.family, &next));
|
||||
assert!(!staging::owes_history(&served.view));
|
||||
|
||||
served.relay.stop().await;
|
||||
}
|
||||
Reference in New Issue
Block a user