Merge #0caf6d9a: Fix maintainer invitation acceptance and sync scenarios

nostr:nevent1qqsqetmdntx5lj2ad969l67fkj8nx7cglgvrlznkf47zr7xc8usqtpcpz3mhxue69uhhyetvv9ujumn8d96zuer9wcuj0sjx

PR-Author: DanConwayDev's Agent
nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0

CoverNote:

Fixes maintainer invitation acceptance and synchronization across owner-only, invitee-only, and shared GRASP servers, including repositories with conflicting existing state.

Adds integration tests covering invitation acceptance, one-way maintainer authorization, announcement and state propagation, refs and HEAD convergence, and synchronization without an invitee push. Only some of the underlying fixes were extracted from the `pr/proactive-sync` branch; the remaining fixes were developed here as these scenarios were tested.
This commit is contained in:
DanConwayDev
2026-07-27 11:23:35 +01:00
14 changed files with 2290 additions and 289 deletions
+25
View File
@@ -13,6 +13,31 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Fix invitation acceptance sync by deferring subscription consolidation until
in-flight relay batches finish, keeping the sync actor available to process
EOSE messages and the five-second purgatory reconciliation pass.
- Fix invitation acceptance sync when a listed source relay is initially
unavailable or empty by retaining and retrying desired GRASP-02 work until it
is confirmed. Root-slash URL variants now share one relay lifecycle instead
of multiplying connections and subscription batches.
- Fix invitation acceptance sync across large relay lists by moving websocket
handshakes out of the sync actor, limiting them to eight concurrent attempts,
and waiting for each relay to be connected before starting subscriptions.
- Fix invitation dependency recovery by targeting retained event IDs at their
associated relay and falling back from NIP-77 after a 15-second total
deadline, so a missing or endlessly active exchange cannot stall the actor.
- Fix large invitation relay sets repeatedly rebuilding their subscriptions by
applying the 70-filter consolidation threshold only to fragmentation above
the relay's irreducible desired live-filter baseline.
- Keep invitation relay connection ownership inside the bounded scheduler by
disabling the SDK's independent auto-reconnect loop and cancelling queued or
active connection workers when the sync manager shuts down.
- Fix invitation acceptance on a shared GRASP server by reapplying an existing
owner state event to the invitee's newly created repository, so GRASP-02 can
align it without an invitee state event or Git push.
- Fix invitation acceptance when the invitee already owns the same repository
identifier by routing maintainer-changing replacements through purgatory and
reprocessing the owner dependencies that align the existing Git repository.
- Sideband-aware Git clients now receive periodic progress while GRASP performs post-push purgatory promotion and cross-owner repository alignment, preventing the client I/O timeout from expiring during unusually complex finalization.
- Smart HTTP pushes now expose the terminal receive-pack flush only after GRASP has finished promoting the matching repository announcement and state from purgatory. Git progress remains streamed while large packs are resolved and checked, but an immediately following clone or proposal push can now rely on a completed push being queryable on the relay.
- Batch the one-time deletion-request lifecycle migration so large production databases do not remain unavailable while LMDB commits every historical request in separate transactions.
+7 -1
View File
@@ -514,9 +514,15 @@ The ngit-grasp relay implements **Proactive Sync of Nostr Events**, which synchr
- **Three-way diff** (`compute_actions`) determines new subscriptions needed
- **Smart reconnection** - uses `since` filter for quick reconnects (<15 min), fresh sync otherwise
- **Health tracking** with exponential backoff for failing relays
- **Bounded connection workers** keep slow DNS and websocket handshakes out of
the sync actor while limiting network pressure to eight concurrent attempts
- **Daily sync** with random 23-25h timer to detect state drift
- **Filter consolidation** when count exceeds 70 to prevent subscription explosion
- **Filter consolidation** when incremental fragmentation exceeds the desired
live-filter baseline by 70; rebuilds are deferred until in-flight batches
drain so EOSE and purgatory work remain responsive
- **Rejected events index** - prevents wasteful broad re-fetching while retaining exact IDs for dependency recovery
- **Desired-source retention** keeps listed GRASP-02 relays retryable until
repository work is actually confirmed, including StateOnly invitation sync
**Architecture:**
+31 -6
View File
@@ -25,8 +25,14 @@ Key Architectural Points:
- **Discovery management**: The nature of discovery inherently leads to a drip feed of root_events (e.g., Repo Announcements, Issues, Patches and PRs) that require additional subscriptions. Without careful management this can lead to large numbers of subscriptions and potentially rate limiting. Mitigation strategies:
- Self-subscriber waits for 5s to batch updates before creating new filters / subscriptions, allowing time for most events to be received from outstanding subscriptions from connected relays
- PendingBatch tracks each new set of filters that may require pagination until they are complete
- Avoid long awaits - recompute desired filters when connection is established to ensure filters are as consolidated as possible
- Consolidation function ensures number of live_sync subscriptions don't reach rate-limiting limits (threshold: 70 filters)
- Websocket handshakes run in at most eight bounded workers outside the sync
actor; only the actor applies their results, and subscriptions start only
after the relay reports `Connected`. The sync manager owns retry/backoff
rather than the SDK, and shutdown cancels queued or active workers
- Recompute desired filters when connection is established to ensure filters are as consolidated as possible
- Consolidation bounds incremental fragmentation to 70 subscriptions above
the irreducible desired live-filter baseline, so large desired sets remain
stable after rebuilding
- **Quick Reconnect** (< 15mins) - doesn't do a full reconciliation vs fresh start (longer disconnect or relaunch binary)
- **Background timers** handle relay connection health and metrics, handling reconnects after backoff and recovery after rate-limiting
@@ -334,8 +340,11 @@ The sync system uses three background tasks that run continuously:
**Actions**:
1. **Disconnect checking**: Calls `check_disconnects()` to remove relays with no repos/events (except bootstrap)
2. **Retry disconnected**: Calls `retry_disconnected_relays()` to attempt reconnection per health tracker backoff
1. **Disconnect checking**: Calls `check_disconnects()` to remove relays with
neither confirmed nor desired repository work (except bootstrap)
2. **Retry disconnected**: Calls `retry_disconnected_relays()` to attempt
reconnection per health tracker backoff while either confirmed or desired
`RepoSyncIndex` work remains
3. **Rate limit recovery**: Calls `check_rate_limit_recovery()` to clear expired rate limits
4. **Metrics update**: Updates Prometheus metrics with current health states
@@ -691,7 +700,10 @@ fn compute_actions(
- **`fresh_start()`**: Full sync - clears all state, L1 historic (with negentropy if available), then L2+L3 via recompute
- **`quick_reconnect()`**: Incremental sync - preserves confirmed state, L1 historic with `since`, L2+L3 rebuild with `since`, then recompute for new items
- **`daily_sync()`**: Wrapper around `fresh_start()` without disconnect metrics
- **`consolidate()`**: Reduces filter count - clears pending, unsubscribes all, rebuilds live subscriptions only, then recompute for new items
- **`consolidate()`**: Reduces filter count after in-flight historic batches have
drained. If a relay crosses the threshold while batches are pending,
consolidation is queued and batch completion wakes the sync actor to rebuild
subscriptions. The actor never polls for EOSE while holding its own lock.
### Sync Primitives
@@ -770,6 +782,11 @@ If negentropy fails (relay doesn't support NIP-77, network error, etc.):
2. The sync falls back to traditional REQ+EOSE
3. No error is raised - fallback is automatic
Each dry-run diff also has a 15-second total wall-clock deadline. This is
separate from the SDK's initial-response and idle timers: a relay that keeps an
exchange active without completing it cannot hold the sync actor indefinitely.
Reaching the deadline uses the same unsupported-relay fallback path.
### Integration with Rejected Events Index
The rejected events index prevents wasteful re-fetching during negentropy sync by excluding rejected event IDs from the reconciliation process:
@@ -803,7 +820,10 @@ needs the inviter's Git data. In that flow:
1. Dependency-resolvable entries remain in the cold index until processing succeeds.
2. If the full event is still in the hot cache, it is re-processed immediately.
3. If the full event expired, the sync manager requests only the retained event IDs from connected relays across the recursive maintainer chain.
3. If the full event expired, the sync manager requests only the retained event
IDs from connected relays across the recursive maintainer chain. Each
request explicitly targets its associated relay connection instead of the
SDK's automatic pool-wide relay selection.
4. Relay requests run in parallel as bounded background work; announcements are processed before state events that may depend on them.
5. Successful or duplicate events are removed from both tiers. Failed and empty requests retain their IDs for a throttled retry.
@@ -970,8 +990,13 @@ RateLimited -> previous state: After 65-second cooldown expires
### Special Behaviors
- **Bootstrap relays**: Never disconnected by cleanup system, even if empty
- **Desired GRASP-02 sources**: Remain registered and retryable before their
first successful historic batch; an initially empty or unavailable source
cannot make a purgatory invitation permanently lose its sync path
- **Rate limiting**: Distinct from connection failures - triggered by relay NOTICE messages
- **Connection timeout**: Set to `base_backoff_secs` to ensure retry timing works correctly
- **Connection concurrency**: At most eight DNS/websocket attempts run at once;
queued attempts do not start health backoff until a worker slot is available
---
+8 -1
View File
@@ -280,7 +280,14 @@ pub async fn fetch_repository_data_with_purgatory(
for entry in purgatory_announcements {
if let Ok(announcement) = RepositoryAnnouncement::from_event(entry.event) {
repo_data.announcements.push(announcement);
if let Some(current) = repo_data.announcements.iter_mut().find(|current| {
current.event.pubkey == announcement.event.pubkey
&& current.identifier == announcement.identifier
}) {
*current = announcement;
} else {
repo_data.announcements.push(announcement);
}
}
}
+18 -6
View File
@@ -981,16 +981,19 @@ pub async fn process_newly_available_git_data(
Ok(result)
}
/// Process purgatory promotions for owner repos populated by state-event OID copies.
/// Process purgatory promotions for owner repos aligned by state processing.
///
/// [`crate::git::process::process_state_with_git_data`] copies OIDs from the
/// source repo to other authorized owner repos for the same identifier. Those
/// target repos may now satisfy purgatory announcements even though no direct
/// push or fetch happened for the target repo path. This helper runs the normal
/// newly-available-git-data pipeline for each populated target repo so promotion
/// behavior stays centralized in this module.
/// push or fetch happened for the target repo path. The source repo can also
/// have a maintainer-topology replacement waiting in purgatory while its
/// existing Git data is re-evaluated. This helper runs the normal
/// newly-available-git-data pipeline only for repos authorized by the preferred
/// state and containing every required OID, so an unrelated or failed copy
/// cannot promote an empty announcement.
pub async fn process_repos_populated_by_state_copy(
source_repo_path: &Path,
state: &RepositoryState,
repo_data: &RepositoryData,
database: &SharedDatabase,
local_relay: Option<&nostr_relay_builder::LocalRelay>,
@@ -1000,10 +1003,19 @@ pub async fn process_repos_populated_by_state_copy(
) -> ProcessResult {
let empty_oids: HashSet<String> = HashSet::new();
let mut aggregate = ProcessResult::default();
let maintainers_by_owner = collect_authorized_maintainers(&repo_data.announcements);
for announcement in &repo_data.announcements {
let target_repo_path = git_data_path.join(announcement.repo_path());
if target_repo_path == source_repo_path {
let owner = announcement.event.pubkey.to_hex();
let Some(maintainers) = maintainers_by_owner.get(&owner) else {
continue;
};
if !maintainers.contains(&state.event.pubkey.to_hex())
|| !is_latest_authorized_state(state, maintainers, &repo_data.states)
|| !can_apply_state(&state.event, &target_repo_path)
{
continue;
}
+80 -1
View File
@@ -7,7 +7,7 @@
use std::net::SocketAddr;
use std::num::NonZeroUsize;
use std::path::Path;
use std::sync::Arc;
use std::sync::{Arc, RwLock};
use anyhow::Result;
use nostr::nips::nip19::ToBech32;
@@ -29,6 +29,7 @@ use crate::nostr::policy::{
};
use crate::nostr::SharedDatabase;
use crate::purgatory::promotion_hooks::NostrPurgatoryPromotionHooks;
use crate::sync::rejected_index::RejectedEventsIndex;
/// NIP-34 Write Policy — admission and routing for GRASP-01 events
///
@@ -53,6 +54,7 @@ pub struct Nip34WritePolicy {
pr_event_policy: PrEventPolicy,
related_event_policy: RelatedEventPolicy,
deletion: DeletionService,
rejected_events_index: Arc<RwLock<Option<Arc<RejectedEventsIndex>>>>,
}
impl std::fmt::Debug for Nip34WritePolicy {
@@ -100,6 +102,7 @@ impl Nip34WritePolicy {
pr_event_policy: PrEventPolicy::new(ctx.clone(), repo_init_locks),
related_event_policy: RelatedEventPolicy::new(ctx.clone()),
deletion: DeletionService::new(deletion_ctx),
rejected_events_index: Arc::new(RwLock::new(None)),
ctx,
}
}
@@ -144,6 +147,26 @@ impl Nip34WritePolicy {
self.ctx.set_local_relay(relay);
}
/// Make sync's rejected-event dependencies available to direct-write
/// purgatory promotions.
pub fn set_rejected_events_index(&self, index: Arc<RejectedEventsIndex>) {
*self
.rejected_events_index
.write()
.expect("rejected events index lock poisoned") = Some(index);
}
pub(crate) fn rejected_events_index(&self) -> Option<Arc<RejectedEventsIndex>> {
self.rejected_events_index
.read()
.expect("rejected events index lock poisoned")
.clone()
}
pub(crate) fn local_relay(&self) -> Option<nostr_relay_builder::LocalRelay> {
self.ctx.get_local_relay()
}
/// Access the deletion subsystem directly.
///
/// Exposes startup passes and recovery without requiring delegation wrappers
@@ -310,6 +333,13 @@ impl Nip34WritePolicy {
// New announcement - add to purgatory
match self.announcement_policy.add_to_purgatory(event) {
Ok(()) => {
// The server may already have the repository's owner-authored state
// and Git objects. Reapply that stored state now so a newly accepted
// maintainer announcement can be promoted without waiting for the
// same replaceable state event to arrive over sync again.
self.reconcile_stored_state_events_for_identifier(&announcement.identifier)
.await;
tracing::info!(
"Accepted announcement to purgatory: {} (waiting for git data)",
event_id_str
@@ -650,6 +680,55 @@ impl Nip34WritePolicy {
}
}
/// Reapply stored state after a new maintainer repository enters purgatory.
///
/// A shared GRASP server can already hold the owner's latest state and Git
/// objects when an invitee publishes their reciprocal announcement. Since
/// the state is replaceable and already stored, sync may not deliver it
/// again. Processing stored candidates newest-first lets each maintainer
/// chain apply its preferred state before considering older candidates,
/// without requiring an invitee-authored state event or push.
async fn reconcile_stored_state_events_for_identifier(&self, identifier: &str) {
let filter = Filter::new().kind(Kind::RepoState).custom_tag(
SingleLetterTag::lowercase(Alphabet::D),
identifier.to_string(),
);
let mut states: Vec<Event> = match self.ctx.database.query(filter).await {
Ok(events) => events.into_iter().collect(),
Err(error) => {
tracing::warn!(
identifier = %identifier,
error = %error,
"Failed to query stored state for new maintainer repository"
);
return;
}
};
states.sort_by(|left, right| {
right
.created_at
.cmp(&left.created_at)
.then_with(|| left.id.cmp(&right.id))
});
for state in states {
let promotion_hooks = NostrPurgatoryPromotionHooks::recovery_only(self);
if let Err(error) = self
.state_policy
.process_state_event(&state, true, Some(&promotion_hooks))
.await
{
tracing::warn!(
identifier = %identifier,
event_id = %state.id,
error = %error,
"Failed to reapply stored state to new maintainer repository"
);
}
}
}
/// Handle events that must reference accepted repositories or events
async fn handle_related_event(&self, event: &Event, event_type: &str) -> WritePolicyResult {
let event_id_str = event.id.to_bech32().unwrap_or_else(|_| event.id.to_hex());
+52 -16
View File
@@ -78,16 +78,28 @@ impl AnnouncementPolicy {
// Owner de-list replacement for an already served repo:
// do NOT allow maintainer-exception acceptance to bypass
// de-list enforcement. Let the caller treat this as a
// reject path so runtime deletion parity can run.
// de-list enforcement when the previous announcement
// actually listed this service. An external maintainer
// announcement never listed this service, so its
// replacement must remain eligible for the maintainer
// exception instead of being mistaken for a de-list.
let lists_service = announcement.lists_service(&self.config.domain);
if !lists_service {
match self
.has_db_announcement(&event.pubkey, &announcement.identifier)
.db_announcement(&event.pubkey, &announcement.identifier)
.await
{
Ok(true) => return AnnouncementResult::Reject(reason),
Ok(false) => {}
Ok(Some(current_event)) => {
let current =
match RepositoryAnnouncement::from_event(current_event) {
Ok(current) => current,
Err(_) => return AnnouncementResult::Reject(reason),
};
if current.lists_service(&self.config.domain) {
return AnnouncementResult::Reject(reason);
}
}
Ok(None) => {}
Err(_) => return AnnouncementResult::Reject(reason),
}
}
@@ -114,11 +126,11 @@ impl AnnouncementPolicy {
// Parse announcement to check for existing active announcement
match RepositoryAnnouncement::from_event(event.clone()) {
Ok(announcement) => {
let in_db = match self
.has_db_announcement(&event.pubkey, &announcement.identifier)
let db_announcement = match self
.db_announcement(&event.pubkey, &announcement.identifier)
.await
{
Ok(v) => v,
Ok(event) => event,
Err(e) => {
tracing::warn!(
error = %e,
@@ -131,7 +143,30 @@ impl AnnouncementPolicy {
}
};
if in_db {
if let Some(current_event) = db_announcement {
let current = match RepositoryAnnouncement::from_event(current_event) {
Ok(current) => current,
Err(e) => {
return AnnouncementResult::Reject(format!(
"Failed to parse stored announcement: {e}"
))
}
};
let current_maintainers: HashSet<&str> =
current.maintainers.iter().map(String::as_str).collect();
let incoming_maintainers: HashSet<&str> = announcement
.maintainers
.iter()
.map(String::as_str)
.collect();
if current_maintainers != incoming_maintainers {
tracing::debug!(
identifier = %announcement.identifier,
"Replacement announcement changes maintainers - routing through purgatory"
);
return AnnouncementResult::AcceptPurgatory;
}
// Replacement announcement with DB entry - accept immediately
tracing::debug!(
identifier = %announcement.identifier,
@@ -272,15 +307,12 @@ impl AnnouncementPolicy {
);
}
/// Check if there's an announcement in the database for this (pubkey, identifier).
///
/// Only checks the database (promoted announcements). For purgatory checks use
/// `purgatory.has_purgatory_announcement()` directly.
async fn has_db_announcement(
/// Return the NIP-01-preferred served announcement for this coordinate.
async fn db_announcement(
&self,
pubkey: &PublicKey,
identifier: &str,
) -> Result<bool, String> {
) -> Result<Option<Event>, String> {
let filter = Filter::new()
.kind(Kind::GitRepoAnnouncement)
.author(*pubkey)
@@ -294,7 +326,11 @@ impl AnnouncementPolicy {
Err(e) => return Err(format!("Database query failed: {}", e)),
};
Ok(!events.is_empty())
Ok(events.into_iter().max_by(|left, right| {
left.created_at
.cmp(&right.created_at)
.then_with(|| right.id.cmp(&left.id))
}))
}
/// Add an announcement to purgatory
+32 -5
View File
@@ -172,11 +172,32 @@ impl StatePolicy {
}
}
// Duplicate check in db
if db_repo_data.states.iter().any(|e| e.event.id.eq(&event.id)) {
// A state already stored in the database may still need to be applied to
// a newly-created maintainer repository. This happens when an invitee
// publishes only their reciprocal announcement to a GRASP server that
// already serves the owner's state and Git data. Reconcile that state
// while the authorized invitee announcement is in purgatory instead of
// returning early as an ordinary duplicate.
let state_already_in_db = db_repo_data.states.iter().any(|e| e.event.id.eq(&event.id));
let has_authorized_purgatory_announcement = authorized_owners.iter().any(|owner_hex| {
nostr_sdk::prelude::PublicKey::from_hex(owner_hex).is_ok_and(|owner| {
self.ctx
.purgatory
.has_purgatory_announcement(&owner, &state.identifier)
})
});
if state_already_in_db && !has_authorized_purgatory_announcement {
tracing::debug!("processed state event duplicate (in db): {}", event.id);
return Ok(duplicate("already have this event"));
}
if state_already_in_db {
tracing::info!(
event_id = %event.id,
identifier = %state.identifier,
"Reapplying stored state to a new maintainer repository"
);
}
// Check if git data is available
if let Some(repo_with_git_data) =
@@ -220,7 +241,7 @@ impl StatePolicy {
// for repos that became populated by the copy.
let local_relay = self.ctx.get_local_relay();
let promotion_result = crate::git::sync::process_repos_populated_by_state_copy(
&repo_with_git_data,
&state,
&db_repo_data,
&self.ctx.database,
local_relay.as_ref(),
@@ -239,8 +260,14 @@ impl StatePolicy {
);
}
// Event will be saved and broadcast by relay builder
Ok(WritePolicyResult::Accept)
if state_already_in_db {
Ok(duplicate(
"already stored; reapplied to pending maintainer repository",
))
} else {
// Event will be saved and broadcast by relay builder
Ok(WritePolicyResult::Accept)
}
} else {
// Only reject expired events if they're from sync (not user-submitted)
// User-submitted events should be allowed to retry in case git data became available
+85 -10
View File
@@ -29,13 +29,13 @@ pub struct NostrPurgatoryPromotionHooks {
}
impl NostrPurgatoryPromotionHooks {
/// Run deletion recovery and accepted-event persistence hooks only.
/// Run deletion recovery and any configured rejected-event dependency hooks.
pub fn recovery_only(write_policy: &Nip34WritePolicy) -> Self {
Self {
deletion: write_policy.deletion().clone(),
write_policy: Some(Arc::new(write_policy.clone())),
rejected_events_index: None,
local_relay: None,
rejected_events_index: write_policy.rejected_events_index(),
local_relay: write_policy.local_relay(),
}
}
@@ -52,6 +52,68 @@ impl NostrPurgatoryPromotionHooks {
local_relay,
}
}
async fn reprocess_state_dependencies(
&self,
write_policy: &Nip34WritePolicy,
rejected_events_index: &RejectedEventsIndex,
relay: &LocalRelay,
author: &PublicKey,
identifier: &str,
) {
let (event_ids, hot_events) =
rejected_events_index.dependency_candidates(author, identifier, Some(EventType::State));
if !event_ids.is_empty() {
info!(
pubkey = %author,
identifier = %identifier,
dependency_event_count = event_ids.len(),
hot_cache_events = hot_events.len(),
"Found rejected state dependencies after purgatory announcement promotion"
);
}
let dummy_addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0);
for state in hot_events {
info!(
event_id = %state.id,
pubkey = %author,
identifier = %identifier,
"Re-processing state event after purgatory announcement promotion"
);
match write_policy.admit_event(&state, &dummy_addr).await {
WritePolicyResult::Accept => {
match write_policy
.save_accepted_event(&state, SaveContext::HotCacheReprocess)
.await
{
Ok(_) => {
rejected_events_index.remove(&state.id);
relay.notify_event(state.clone());
}
Err(e) => {
warn!(
event_id = %state.id,
error = %e,
"Failed to save re-processed state event"
);
}
}
}
WritePolicyResult::Reject { status: true, .. } => {
rejected_events_index.remove(&state.id);
write_policy.purgatory().enqueue_sync_immediate(identifier);
}
_ => {
warn!(
event_id = %state.id,
"State event still rejected after announcement promotion"
);
}
}
}
}
}
#[async_trait]
@@ -88,15 +150,11 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks {
return;
};
if announcement.maintainers.is_empty() {
return;
}
debug!(
identifier = %announcement.identifier,
event_id = %event.id,
maintainer_count = announcement.maintainers.len(),
"Owner announcement promoted via git push, checking hot cache for rejected maintainer announcements"
"Owner announcement promoted from purgatory, checking hot cache for rejected maintainer announcements"
);
for maintainer_hex in &announcement.maintainers {
@@ -114,7 +172,7 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks {
identifier = %announcement.identifier,
dependency_event_count = event_ids.len(),
hot_cache_events = hot_events.len(),
"Found rejected maintainer dependencies after git push promotion"
"Found rejected maintainer dependencies after purgatory promotion"
);
}
@@ -124,7 +182,7 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks {
event_id = %hot_event.id,
maintainer = %maintainer_hex,
identifier = %announcement.identifier,
"Re-processing maintainer announcement from hot cache after git push promotion"
"Re-processing maintainer announcement from hot cache after purgatory promotion"
);
match write_policy.admit_event(&hot_event, &dummy_addr).await {
@@ -136,6 +194,14 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks {
Ok(_) => {
rejected_events_index.remove(&hot_event.id);
relay.notify_event(hot_event.clone());
self.reprocess_state_dependencies(
write_policy,
rejected_events_index,
relay,
&hot_event.pubkey,
&announcement.identifier,
)
.await;
info!(
event_id = %hot_event.id,
"Maintainer announcement accepted and saved on re-processing"
@@ -168,5 +234,14 @@ impl PurgatoryPromotionHooks for NostrPurgatoryPromotionHooks {
}
}
}
self.reprocess_state_dependencies(
write_policy,
rejected_events_index,
relay,
&event.pubkey,
&announcement.identifier,
)
.await;
}
}
+3
View File
@@ -242,6 +242,9 @@ impl RelayServer {
// Get a reference to the rejected events index for shutdown persistence
// and for the HTTP server's git push path (hot-cache re-processing)
let rejected_events_index = sync_manager.rejected_events_index();
relay_runtime
.write_policy
.set_rejected_events_index(rejected_events_index.clone());
let mut background_tasks = Vec::new();
+36 -1
View File
@@ -67,7 +67,9 @@ pub fn derive_relay_targets(
for (repo_id, needs) in repo_index {
for relay_url in &needs.relays {
let entry = relay_targets.entry(relay_url.clone()).or_default();
let relay_key =
super::canonical_relay_key(relay_url).unwrap_or_else(|_| relay_url.clone());
let entry = relay_targets.entry(relay_key).or_default();
match needs.sync_level {
super::SyncLevel::Full => {
@@ -269,6 +271,39 @@ mod tests {
assert_eq!(relay_needs.repos.len(), 3);
}
#[test]
fn test_derive_relay_targets_deduplicates_root_slash_variants() {
let repo_index = HashMap::from([
(
"repo1".to_string(),
ModRepoSyncNeeds {
relays: HashSet::from(["wss://relay1.com".to_string()]),
root_events: HashSet::new(),
sync_level: Default::default(),
},
),
(
"repo2".to_string(),
ModRepoSyncNeeds {
relays: HashSet::from(["wss://relay1.com/".to_string()]),
root_events: HashSet::new(),
sync_level: Default::default(),
},
),
]);
let targets = derive_relay_targets(&repo_index);
assert_eq!(targets.len(), 1);
assert_eq!(
targets
.get("wss://relay1.com")
.expect("root slash variants must share one key")
.repos,
HashSet::from(["repo1".to_string(), "repo2".to_string()])
);
}
#[test]
fn test_derive_relay_targets_repo_across_multiple_relays() {
let mut repo_index = HashMap::new();
+826 -194
View File
File diff suppressed because it is too large Load Diff
+261 -47
View File
@@ -17,11 +17,44 @@
use futures_util::StreamExt;
use nostr_sdk::prelude::*;
use std::future::Future;
use std::time::Duration;
use tokio::sync::mpsc;
use crate::nostr::SharedDatabase;
/// Maximum wall-clock time for one dry-run NIP-77 reconciliation.
///
/// The SDK has initial-response and idle timers, but a relay can keep an
/// exchange alive indefinitely by continuing to send reconciliation messages.
const NEGENTROPY_DIFF_TIMEOUT: Duration = Duration::from_secs(15);
/// Interval between relay-status checks while a connection attempt is in flight.
///
/// nostr-sdk may return from `try_connect_relay` while another task still has
/// the relay in `Connecting`; subscriptions must wait for actual readiness.
const CONNECTION_STATUS_POLL_INTERVAL: Duration = Duration::from_millis(25);
async fn wait_for_connected_status<F>(mut relay_status: F) -> Result<(), RelayStatus>
where
F: FnMut() -> RelayStatus,
{
loop {
let status = relay_status();
match status {
RelayStatus::Connected => return Ok(()),
RelayStatus::Disconnected
| RelayStatus::Terminated
| RelayStatus::Banned
| RelayStatus::Sleeping
| RelayStatus::Shutdown => return Err(status),
RelayStatus::Initialized | RelayStatus::Pending | RelayStatus::Connecting => {
tokio::time::sleep(CONNECTION_STATUS_POLL_INTERVAL).await;
}
}
}
}
/// Events from a relay connection
#[derive(Debug)]
pub enum RelayEvent {
@@ -150,7 +183,7 @@ impl RelayConnection {
/// This method:
/// 1. Adds the relay to the client
/// 2. Establishes the WebSocket connection
/// 3. Verifies connection was established
/// 3. Waits for the relay status to report `Connected`
///
/// Subscriptions are handled separately via handle_connect_or_reconnect.
///
@@ -163,34 +196,53 @@ impl RelayConnection {
/// * `Ok(())` - Connection established successfully
/// * `Err(String)` with error description on failure
pub async fn connect(&self, connection_timeout_secs: u64) -> Result<(), String> {
// Add relay to client
self.client
.add_relay(&self.url)
.await
.map_err(|e| format!("Failed to add relay {}: {}", self.url, e))?;
// Establish connection using try_connect_relay for immediate failure detection
//
// Key difference from client.connect():
// - try_connect_relay: Single attempt with timeout, returns Err on failure,
// does NOT spawn background retry task (we control retries via HealthTracker)
// - connect(): Spawns background task, returns immediately, auto-retries forever
//
// Using try_connect_relay gives us:
// 1. Immediate error return on connection failure
// 2. Configurable timeout (set to base_backoff_secs to ensure retry timing works)
// 3. No conflicting retry logic (we use HealthTracker for backoff)
// 4. Cleaner error messages for metrics recording
//
// See: nostr-sdk-0.44 Client::try_connect_relay documentation
self.client
.try_connect_relay(
&self.url,
std::time::Duration::from_secs(connection_timeout_secs),
let connection_timeout = Duration::from_secs(connection_timeout_secs);
let connection_deadline = tokio::time::Instant::now() + connection_timeout;
let timeout_error = || {
format!(
"Timed out connecting to relay {} after {} seconds",
self.url, connection_timeout_secs
)
};
let relay = tokio::time::timeout_at(connection_deadline, async {
self.client
.add_relay(&self.url)
.reconnect(false)
.await
.map_err(|e| format!("Failed to add relay {}: {}", self.url, e))?;
self.client
.relay(&self.url)
.await
.map_err(|e| format!("Failed to get relay {}: {}", self.url, e))?
.ok_or_else(|| format!("Relay {} was not added to the client", self.url))
})
.await
.map_err(|_| timeout_error())??;
// Use one deadline for both the SDK attempt and readiness verification.
// Cancelling the SDK call with an outer timeout can leave its relay
// status stranded at Connecting, so pass it the remaining budget.
let remaining = connection_deadline.saturating_duration_since(tokio::time::Instant::now());
self.client
.try_connect_relay(&self.url, remaining)
.await
.map_err(|e| format!("Failed to connect to relay {}: {}", self.url, e))?;
tokio::time::timeout_at(
connection_deadline,
wait_for_connected_status(|| relay.status()),
)
.await
.map_err(|_| timeout_error())?
.map_err(|status| {
format!(
"Relay {} entered terminal status {} before connecting",
self.url, status
)
})?;
tracing::info!(url = %self.url, "Connected to relay");
Ok(())
}
@@ -410,7 +462,24 @@ impl RelayConnection {
filter: Filter,
timeout: Duration,
) -> Result<Vec<Event>, String> {
self.client
let relay = self
.client
.relay(&self.url)
.await
.map_err(|error| {
format!(
"Failed to fetch events from {}: failed to resolve relay: {}",
self.url, error
)
})?
.ok_or_else(|| {
format!(
"Failed to fetch events from {}: relay is not registered",
self.url
)
})?;
relay
.fetch_events(filter)
.timeout(timeout)
.await
@@ -537,7 +606,32 @@ impl RelayConnection {
) -> Result<nostr_sdk::client::SyncSummary, String> {
// Use dry_run to only identify differences without downloading events
let sync_opts = SyncOptions::default().dry_run();
let client = self.client.clone();
let sync_task = async move {
match client.sync(filter).opts(sync_opts).await {
Ok(output) => {
if !output.failed.is_empty() {
Err(format!("Negentropy diff had failures: {:?}", output.failed))
} else {
Ok(output.value)
}
}
Err(error) => Err(format!("Negentropy diff failed: {}", error)),
}
};
self.run_negentropy_diff_with_timeout(sync_task, NEGENTROPY_DIFF_TIMEOUT)
.await
}
async fn run_negentropy_diff_with_timeout<F>(
&self,
sync_task: F,
timeout_duration: Duration,
) -> Result<nostr_sdk::client::SyncSummary, String>
where
F: Future<Output = Result<nostr_sdk::client::SyncSummary, String>>,
{
// Clone the atomic for the polling task
let nip77_status = self.nip77_supported.clone();
let url = self.url.clone();
@@ -557,29 +651,23 @@ impl RelayConnection {
}
};
// Race the sync operation against the polling task
let sync_task = self.client.sync(filter.clone()).opts(sync_opts);
let result = tokio::select! {
poll_result = poll_task => {
// Polling detected NIP-77 not supported
poll_result
}
sync_result = sync_task => {
// Sync completed (or failed) first
match sync_result {
Ok(output) => {
// Check for any failures
// Note: Timeouts are common for relays without NIP-77 support
if !output.failed.is_empty() {
Err(format!("Negentropy diff had failures: {:?}", output.failed))
} else {
Ok(output.value)
}
}
Err(e) => Err(format!("Negentropy diff failed: {}", e))
let result = match tokio::time::timeout(timeout_duration, async {
tokio::select! {
poll_result = poll_task => {
poll_result
}
sync_result = sync_task => {
sync_result
}
}
})
.await
{
Ok(result) => result,
Err(_) => Err(format!(
"Negentropy diff timed out after {:.3}s",
timeout_duration.as_secs_f64()
)),
};
match result {
@@ -620,6 +708,132 @@ impl RelayConnection {
#[cfg(test)]
mod tests {
use super::*;
use nostr_relay_builder::prelude::LocalRelayBuilder;
use std::future::pending;
#[tokio::test]
async fn connection_readiness_waits_until_status_is_connected() {
let mut statuses = [
RelayStatus::Connecting,
RelayStatus::Connecting,
RelayStatus::Connected,
]
.into_iter();
wait_for_connected_status(|| statuses.next().expect("status sequence exhausted"))
.await
.expect("delayed connection should become ready");
}
#[tokio::test]
async fn connection_readiness_remains_bounded_when_connecting_never_completes() {
let result = tokio::time::timeout(
Duration::from_millis(60),
wait_for_connected_status(|| RelayStatus::Connecting),
)
.await;
assert!(
result.is_err(),
"a relay stuck in Connecting must not be reported as ready"
);
}
#[tokio::test]
async fn connection_readiness_fails_on_terminal_status() {
let mut statuses = [RelayStatus::Connecting, RelayStatus::Disconnected].into_iter();
assert_eq!(
wait_for_connected_status(|| statuses.next().expect("status sequence exhausted")).await,
Err(RelayStatus::Disconnected)
);
}
#[tokio::test]
async fn hung_negentropy_diff_times_out_and_disables_future_attempts() {
let connection = RelayConnection::new("ws://127.0.0.1:1".to_string(), Keys::generate());
assert!(connection.supports_negentropy().await);
let result = connection
.run_negentropy_diff_with_timeout(
pending::<Result<nostr_sdk::client::SyncSummary, String>>(),
Duration::from_millis(20),
)
.await;
assert!(result
.expect_err("a hung NIP-77 exchange must reach its local deadline")
.contains("timed out"));
assert!(
!connection.supports_negentropy().await,
"a timed-out relay must use REQ+EOSE for future historic batches"
);
}
#[tokio::test]
async fn fetch_events_targets_the_connections_exact_relay() {
let configured = LocalRelayBuilder::default().build();
configured.run().await.expect("start configured relay");
let other = LocalRelayBuilder::default().build();
other.run().await.expect("start other relay");
let expected = EventBuilder::text_note("only on the other relay")
.finalize(&Keys::generate())
.expect("build event");
other
.add_event(expected.clone())
.await
.expect("seed other relay");
let connection = RelayConnection::new(configured.url().await.to_string(), Keys::generate());
connection
.connect(3)
.await
.expect("connect configured relay");
let other_url = other.url().await;
connection
.client
.add_relay(other_url.clone())
.await
.expect("register other relay");
connection
.client
.try_connect_relay(other_url, Duration::from_secs(3))
.await
.expect("connect other relay");
let fetched = connection
.fetch_events(Filter::new().id(expected.id), Duration::from_secs(2))
.await
.expect("fetch from configured relay");
assert!(
fetched.is_empty(),
"a relay-specific fetch must not return an event from another client relay"
);
connection.disconnect().await;
configured.shutdown();
other.shutdown();
}
#[tokio::test]
async fn fetch_events_reports_an_unregistered_exact_relay() {
let connection = RelayConnection::new("ws://127.0.0.1:1".to_string(), Keys::generate());
let error = connection
.fetch_events(
Filter::new().kind(Kind::TextNote),
Duration::from_millis(50),
)
.await
.expect_err("an unregistered relay must fail explicitly");
assert!(error.contains("relay is not registered"), "{error}");
assert!(
!error.contains("relay/s not specified"),
"exact-relay fetch must not use an empty automatic target: {error}"
);
}
#[test]
fn test_normalize_url_with_wss_scheme() {
+826 -1
View File
@@ -6,7 +6,7 @@
//! ## Test design
//!
//! Announcements require git data before they are released from purgatory and
//! served to other relays. The tests exercise both dependency-recovery paths:
//! served to other relays. The tests exercise dependency and source recovery:
//!
//! relay_b syncs maintainer announcement from relay_a
//! → write policy rejects it (no owner announcement in DB yet)
@@ -15,6 +15,12 @@
//! → hot copy is re-processed immediately when available
//! → expired hot copy is fetched by its retained cold-index event ID
//! → maintainer announcement supplies the source clone and git data syncs
//! an existing owner publishes an invitation for a new maintainer
//! → the invitee publishes only their reciprocal announcement
//! → GRASP-02 distributes both announcements across owner-only, invitee-only,
//! and shared servers
//! → the owner's state and Git data align the invitee repositories without
//! an invitee state event or Git push
//!
//! To guarantee the maintainer announcements arrive at relay_b *before* the owner
//! git push, relay_b is started with relay_a as its bootstrap relay. That way
@@ -23,6 +29,7 @@
//! announcements are in relay_a's DB), wait briefly for the sync round-trip, then
//! send the owner announcement + git push.
use std::collections::BTreeMap;
use std::path::Path;
use std::process::Command;
use std::time::Duration;
@@ -44,6 +51,824 @@ async fn wait_for_log(path: &Path, needle: &str, timeout: Duration) -> bool {
}
}
fn repository_urls(
keys: &Keys,
relays: &[&TestRelay],
identifier: &str,
) -> (Vec<String>, Vec<String>) {
let npub = keys
.public_key()
.to_bech32()
.expect("Failed to encode repository owner npub");
let clone_urls = relays
.iter()
.map(|relay| format!("http://{}/{}/{}.git", relay.domain(), npub, identifier))
.collect();
let relay_urls = relays.iter().map(|relay| relay.url().to_string()).collect();
(clone_urls, relay_urls)
}
fn repository_announcement(
keys: &Keys,
relays: &[&TestRelay],
maintainers: &[PublicKey],
identifier: &str,
) -> EventBuilder {
let (clone_urls, relay_urls) = repository_urls(keys, relays, identifier);
let mut tags = vec![
Tag::identifier(identifier),
Tag::custom("clone", clone_urls),
Tag::custom("relays", relay_urls),
];
if !maintainers.is_empty() {
tags.push(Tag::custom(
"maintainers",
maintainers
.iter()
.map(PublicKey::to_hex)
.collect::<Vec<_>>(),
));
}
EventBuilder::new(Kind::GitRepoAnnouncement, "Repository announcement").tags(tags)
}
fn repository_state_with_default_branch(
clone_urls: &[String],
relay_urls: &[String],
identifier: &str,
default_branch: &str,
commit: &str,
) -> EventBuilder {
EventBuilder::new(Kind::RepoState, "").tags(vec![
Tag::identifier(identifier),
Tag::custom("clone", clone_urls.to_vec()),
Tag::custom("relays", relay_urls.to_vec()),
Tag::custom(
format!("refs/heads/{default_branch}"),
vec![commit.to_string()],
),
Tag::custom("HEAD", vec![format!("refs/heads/{default_branch}")]),
])
}
async fn assert_exact_event_served(relay: &TestRelay, event: &Event, description: &str) {
let found = wait_for_event_on_relay(
relay.url(),
Filter::new().id(event.id),
Duration::from_secs(20),
)
.await;
assert!(
found,
"{description} should be served by {}",
relay.domain()
);
}
async fn push_counts(relay: &TestRelay) -> (u64, u64) {
let text = fetch_metrics(relay.url())
.await
.expect("Failed to fetch relay metrics");
let metrics = ParsedMetrics::parse(&text);
let success = metrics
.counter(
"ngit_git_operations_total",
&[("operation", "push"), ("status", "success")],
)
.unwrap_or(0);
let error = metrics
.counter(
"ngit_git_operations_total",
&[("operation", "push"), ("status", "error")],
)
.unwrap_or(0);
(success, error)
}
fn list_remote_refs(clone_url: &str) -> Result<BTreeMap<String, String>, String> {
let output = Command::new("git")
.args(["ls-remote", "--refs", clone_url])
.output()
.map_err(|error| format!("Failed to list refs from {clone_url}: {error}"))?;
if !output.status.success() {
return Err(format!(
"Failed to list refs from {clone_url}: {}",
String::from_utf8_lossy(&output.stderr)
));
}
String::from_utf8(output.stdout)
.map_err(|error| format!("Invalid UTF-8 from {clone_url}: {error}"))?
.lines()
.map(|line| {
let (commit, name) = line
.split_once('\t')
.ok_or_else(|| format!("Invalid ls-remote line from {clone_url}: {line}"))?;
Ok((name.to_string(), commit.to_string()))
})
.collect()
}
fn announcement_clone_urls(event: &Event) -> Vec<String> {
event
.tags
.iter()
.find(|tag| tag.kind() == "clone")
.expect("Repository announcement should have a clone tag")
.clone()
.to_vec()
.into_iter()
.skip(1)
.collect()
}
async fn assert_remote_refs(
clone_url: &str,
expected: &BTreeMap<String, String>,
timeout: Duration,
) {
let deadline = tokio::time::Instant::now() + timeout;
loop {
let actual = list_remote_refs(clone_url);
if actual.as_ref().is_ok_and(|refs| refs == expected) {
return;
}
assert!(
tokio::time::Instant::now() < deadline,
"Announced Git endpoint should exactly match the owner state: {clone_url}\n\
expected: {expected:?}\nactual: {actual:?}"
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
fn rename_default_branch(repository: &Path, branch: &str) {
let output = Command::new("git")
.args(["branch", "-m", branch])
.current_dir(repository)
.output()
.expect("Failed to rename test repository default branch");
assert!(
output.status.success(),
"Failed to rename test repository default branch: {}",
String::from_utf8_lossy(&output.stderr)
);
}
fn remote_default_branch(clone_url: &str) -> Result<String, String> {
let output = Command::new("git")
.args(["ls-remote", "--symref", clone_url, "HEAD"])
.output()
.map_err(|error| format!("Failed to inspect HEAD from {clone_url}: {error}"))?;
if !output.status.success() {
return Err(format!(
"Failed to inspect HEAD from {clone_url}: {}",
String::from_utf8_lossy(&output.stderr)
));
}
String::from_utf8(output.stdout)
.map_err(|error| format!("Invalid UTF-8 from {clone_url}: {error}"))?
.lines()
.find_map(|line| {
line.strip_prefix("ref: ")
.and_then(|line| line.strip_suffix("\tHEAD"))
.map(str::to_string)
})
.ok_or_else(|| format!("No symbolic HEAD returned by {clone_url}"))
}
async fn assert_remote_default_branch(clone_url: &str, expected: &str, timeout: Duration) {
let deadline = tokio::time::Instant::now() + timeout;
loop {
let actual = remote_default_branch(clone_url);
if actual.as_deref() == Ok(expected) {
return;
}
assert!(
tokio::time::Instant::now() < deadline,
"Announced Git endpoint should use {expected} as its default branch: {clone_url}\n\
actual: {actual:?}"
);
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
/// Exercise the production invitation flow for a repository that already has
/// owner state and Git data.
///
/// The owner selects an owner-only server and a server shared with the invitee.
/// The invitee selects an invitee-only server and that shared server. Acceptance
/// publishes only the invitee's reciprocal announcement: no invitee state event
/// is created and the invitee never pushes Git. GRASP-02 must distribute the
/// signed announcements and use the owner's state and Git data to create both
/// invitee repositories.
#[tokio::test]
async fn test_existing_repository_invitation_acceptance_syncs_without_invitee_push() {
use crate::common::{create_state_event, create_test_repo_with_commit, CommitVariant};
let owner_relay = TestRelay::start_with_sync(None).await;
let invitee_relay = TestRelay::start_with_sync(None).await;
let shared_relay = TestRelay::start_with_sync(None).await;
let owner_keys = Keys::generate();
let invitee_keys = Keys::generate();
let identifier = "existing-repository-invitation-acceptance";
let owner_git = tempfile::tempdir().expect("Failed to create owner repository directory");
let commit = create_test_repo_with_commit(owner_git.path(), CommitVariant::StateTest)
.expect("Failed to create owner repository commit");
let owner_npub = owner_keys
.public_key()
.to_bech32()
.expect("Failed to encode owner npub");
let owner_servers = [&owner_relay, &shared_relay];
let (owner_clone_urls, owner_relay_urls) =
repository_urls(&owner_keys, &owner_servers, identifier);
let initial_announcement =
repository_announcement(&owner_keys, &owner_servers, &[], identifier)
.finalize(&owner_keys)
.expect("Failed to create initial owner announcement");
let owner_state = create_state_event(
&owner_keys,
identifier,
&[("main", &commit)],
&[],
&owner_clone_urls
.iter()
.map(String::as_str)
.collect::<Vec<_>>(),
&owner_relay_urls
.iter()
.map(String::as_str)
.collect::<Vec<_>>(),
)
.expect("Failed to create owner state event");
for relay in owner_servers {
send_to_relay(relay, &initial_announcement)
.await
.expect("Failed to publish initial owner announcement");
send_to_relay(relay, &owner_state)
.await
.expect("Failed to publish owner state");
crate::common::push_to_relay(owner_git.path(), &relay.domain(), &owner_npub, identifier)
.expect("Failed to push existing owner repository");
assert_exact_event_served(relay, &initial_announcement, "Initial owner announcement").await;
assert_exact_event_served(relay, &owner_state, "Owner state event").await;
let owner_ref_aligned = crate::common::check_ref_at_commit(
&relay.domain(),
&owner_npub,
identifier,
"refs/heads/main",
&commit,
)
.await
.expect("Failed to inspect owner repository ref");
assert!(
owner_ref_aligned,
"Owner repository should be aligned on {}",
relay.domain()
);
}
let invitation = repository_announcement(
&owner_keys,
&owner_servers,
&[invitee_keys.public_key()],
identifier,
)
.custom_created_at(Timestamp::from_secs(
initial_announcement.created_at.as_secs() + 1,
))
.finalize(&owner_keys)
.expect("Failed to create owner invitation announcement");
for relay in owner_servers {
send_to_relay(relay, &invitation)
.await
.expect("Failed to publish owner invitation");
assert_exact_event_served(relay, &invitation, "Owner invitation announcement").await;
}
let pushes_before_acceptance = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
let invitee_servers = [&invitee_relay, &shared_relay];
let acceptance = repository_announcement(
&invitee_keys,
&invitee_servers,
&[owner_keys.public_key()],
identifier,
)
.finalize(&invitee_keys)
.expect("Failed to create invitee acceptance announcement");
let (invitee_clone_urls, _) = repository_urls(&invitee_keys, &invitee_servers, identifier);
for relay in invitee_servers {
send_to_relay(relay, &acceptance)
.await
.expect("Failed to publish invitee acceptance announcement");
}
for relay in [&owner_relay, &invitee_relay, &shared_relay] {
assert_exact_event_served(relay, &invitation, "Owner invitation announcement").await;
assert_exact_event_served(relay, &acceptance, "Invitee acceptance announcement").await;
assert_exact_event_served(relay, &owner_state, "Owner state event").await;
}
for relay in [&owner_relay, &invitee_relay, &shared_relay] {
let invitee_state_exists = wait_for_event_on_relay(
relay.url(),
Filter::new()
.kind(Kind::RepoState)
.author(invitee_keys.public_key())
.identifier(identifier),
Duration::from_secs(1),
)
.await;
assert!(
!invitee_state_exists,
"Invitee must not issue a state event served by {}",
relay.domain()
);
}
let pushes_after_acceptance = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
assert_eq!(
pushes_after_acceptance, pushes_before_acceptance,
"Invitation acceptance must not perform or attempt a client Git push"
);
let expected_refs = BTreeMap::from([("refs/heads/main".to_string(), commit.clone())]);
let announced_owner_clones = announcement_clone_urls(&invitation);
let announced_invitee_clones = announcement_clone_urls(&acceptance);
assert_eq!(announced_owner_clones, owner_clone_urls);
assert_eq!(announced_invitee_clones, invitee_clone_urls);
assert_eq!(
announced_owner_clones.len() + announced_invitee_clones.len(),
4,
"Owner and invitee announcements should advertise two Git endpoints each"
);
for clone_url in announced_owner_clones
.iter()
.chain(announced_invitee_clones.iter())
{
assert_remote_refs(clone_url, &expected_refs, Duration::from_secs(20)).await;
}
shared_relay.stop().await;
invitee_relay.stop().await;
owner_relay.stop().await;
}
/// Listing a maintainer immediately authorizes that maintainer's state for the
/// inviting owner's repository; reciprocal acceptance is not required.
///
/// The owner first publishes an older state on an owner-only and shared server.
/// The invitee later publishes a newer, different state on an invitee-only and
/// shared server. When the owner lists the invitee, GRASP-02 must apply the
/// invitee's newer state to both owner endpoints without changing the invitee's
/// announcement or requiring another Git push.
#[tokio::test]
async fn test_invitation_applies_newer_invitee_state_to_owner_before_acceptance() {
use crate::common::{create_test_repo_with_commit, CommitVariant};
let owner_relay = TestRelay::start_with_sync(None).await;
let invitee_relay = TestRelay::start_with_sync(None).await;
let shared_relay = TestRelay::start_with_sync(None).await;
let owner_keys = Keys::generate();
let invitee_keys = Keys::generate();
let identifier = "one-way-invitation-newer-maintainer-state";
let owner_branch = "owner-default";
let invitee_branch = "invitee-default";
let owner_servers = [&owner_relay, &shared_relay];
let invitee_servers = [&invitee_relay, &shared_relay];
let owner_git = tempfile::tempdir().expect("Failed to create owner repository directory");
let owner_commit = create_test_repo_with_commit(owner_git.path(), CommitVariant::StateTest)
.expect("Failed to create owner repository commit");
rename_default_branch(owner_git.path(), owner_branch);
let owner_npub = owner_keys
.public_key()
.to_bech32()
.expect("Failed to encode owner npub");
let (owner_clone_urls, owner_relay_urls) =
repository_urls(&owner_keys, &owner_servers, identifier);
let initial_owner_announcement =
repository_announcement(&owner_keys, &owner_servers, &[], identifier)
.finalize(&owner_keys)
.expect("Failed to create initial owner announcement");
let owner_state = repository_state_with_default_branch(
&owner_clone_urls,
&owner_relay_urls,
identifier,
owner_branch,
&owner_commit,
)
.finalize(&owner_keys)
.expect("Failed to create owner state");
let owner_refs = BTreeMap::from([(format!("refs/heads/{owner_branch}"), owner_commit)]);
for relay in owner_servers {
send_to_relay(relay, &initial_owner_announcement)
.await
.expect("Failed to publish initial owner announcement");
send_to_relay(relay, &owner_state)
.await
.expect("Failed to publish owner state");
crate::common::push_to_relay(owner_git.path(), &relay.domain(), &owner_npub, identifier)
.expect("Failed to push owner repository");
assert_exact_event_served(
relay,
&initial_owner_announcement,
"Initial owner announcement",
)
.await;
assert_exact_event_served(relay, &owner_state, "Owner state event").await;
}
for clone_url in &owner_clone_urls {
assert_remote_refs(clone_url, &owner_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{owner_branch}"),
Duration::from_secs(20),
)
.await;
}
tokio::time::sleep(Duration::from_secs(1)).await;
let invitee_git = tempfile::tempdir().expect("Failed to create invitee repository directory");
let invitee_commit = create_test_repo_with_commit(invitee_git.path(), CommitVariant::PrTest)
.expect("Failed to create invitee repository commit");
rename_default_branch(invitee_git.path(), invitee_branch);
let invitee_npub = invitee_keys
.public_key()
.to_bech32()
.expect("Failed to encode invitee npub");
let (invitee_clone_urls, invitee_relay_urls) =
repository_urls(&invitee_keys, &invitee_servers, identifier);
let invitee_created_at = Timestamp::from_secs(owner_state.created_at.as_secs() + 1);
let invitee_announcement =
repository_announcement(&invitee_keys, &invitee_servers, &[], identifier)
.custom_created_at(invitee_created_at)
.finalize(&invitee_keys)
.expect("Failed to create invitee announcement");
let invitee_state = repository_state_with_default_branch(
&invitee_clone_urls,
&invitee_relay_urls,
identifier,
invitee_branch,
&invitee_commit,
)
.custom_created_at(invitee_created_at)
.finalize(&invitee_keys)
.expect("Failed to create invitee state");
assert!(
invitee_state.created_at > owner_state.created_at,
"Invitee state must be newer than the owner's existing state"
);
let invitee_refs = BTreeMap::from([(
format!("refs/heads/{invitee_branch}"),
invitee_commit.clone(),
)]);
for relay in invitee_servers {
send_to_relay(relay, &invitee_announcement)
.await
.expect("Failed to publish invitee announcement");
send_to_relay(relay, &invitee_state)
.await
.expect("Failed to publish invitee state");
crate::common::push_to_relay(
invitee_git.path(),
&relay.domain(),
&invitee_npub,
identifier,
)
.expect("Failed to push invitee repository");
assert_exact_event_served(relay, &invitee_announcement, "Invitee announcement").await;
assert_exact_event_served(relay, &invitee_state, "Invitee state event").await;
}
for clone_url in &invitee_clone_urls {
assert_remote_refs(clone_url, &invitee_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{invitee_branch}"),
Duration::from_secs(20),
)
.await;
}
for clone_url in &owner_clone_urls {
assert_remote_refs(clone_url, &owner_refs, Duration::from_secs(5)).await;
}
let pushes_before_invitation = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
let invitation = repository_announcement(
&owner_keys,
&owner_servers,
&[invitee_keys.public_key()],
identifier,
)
.custom_created_at(Timestamp::from_secs(invitee_created_at.as_secs() + 1))
.finalize(&owner_keys)
.expect("Failed to create owner invitation");
for relay in owner_servers {
send_to_relay(relay, &invitation)
.await
.expect("Failed to publish owner invitation");
}
for relay in owner_servers {
assert_exact_event_served(relay, &invitation, "Owner invitation announcement").await;
assert_exact_event_served(
relay,
&invitee_announcement,
"Invited maintainer announcement",
)
.await;
assert_exact_event_served(relay, &invitee_state, "Newer invited maintainer state").await;
}
for clone_url in &owner_clone_urls {
assert_remote_refs(clone_url, &invitee_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{invitee_branch}"),
Duration::from_secs(20),
)
.await;
}
for clone_url in &invitee_clone_urls {
assert_remote_refs(clone_url, &invitee_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{invitee_branch}"),
Duration::from_secs(20),
)
.await;
}
let pushes_after_invitation = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
assert_eq!(
pushes_after_invitation, pushes_before_invitation,
"One-way maintainer authorization must sync the owner without a client Git push"
);
shared_relay.stop().await;
invitee_relay.stop().await;
owner_relay.stop().await;
}
/// An invitation must not overwrite an invitee's unrelated repository merely
/// because the inviter's state is newer.
///
/// Both users first create a repository with the same identifier but different
/// default branches. The invitee's repository remains on its own state while
/// the one-way invitation syncs. Once the invitee publishes a reciprocal
/// announcement, the newer owner state becomes authoritative for the combined
/// maintainer set and replaces the invitee repository's refs without a new push.
#[tokio::test]
async fn test_acceptance_replaces_existing_invitee_repository_with_newer_owner_state() {
use crate::common::{create_test_repo_with_commit, CommitVariant};
let owner_relay = TestRelay::start_with_sync(None).await;
let invitee_relay = TestRelay::start_with_sync(None).await;
let shared_relay = TestRelay::start_with_sync(None).await;
let owner_keys = Keys::generate();
let invitee_keys = Keys::generate();
let identifier = "existing-invitee-repository-acceptance";
let invitee_branch = "invitee-default";
let owner_branch = "owner-default";
let invitee_servers = [&invitee_relay, &shared_relay];
let owner_servers = [&owner_relay, &shared_relay];
let invitee_git = tempfile::tempdir().expect("Failed to create invitee repository directory");
let invitee_commit = create_test_repo_with_commit(invitee_git.path(), CommitVariant::StateTest)
.expect("Failed to create invitee repository commit");
rename_default_branch(invitee_git.path(), invitee_branch);
let invitee_npub = invitee_keys
.public_key()
.to_bech32()
.expect("Failed to encode invitee npub");
let (invitee_clone_urls, invitee_relay_urls) =
repository_urls(&invitee_keys, &invitee_servers, identifier);
let initial_invitee_announcement =
repository_announcement(&invitee_keys, &invitee_servers, &[], identifier)
.finalize(&invitee_keys)
.expect("Failed to create initial invitee announcement");
let invitee_state = repository_state_with_default_branch(
&invitee_clone_urls,
&invitee_relay_urls,
identifier,
invitee_branch,
&invitee_commit,
)
.finalize(&invitee_keys)
.expect("Failed to create invitee state");
let invitee_refs = BTreeMap::from([(
format!("refs/heads/{invitee_branch}"),
invitee_commit.clone(),
)]);
for relay in invitee_servers {
send_to_relay(relay, &initial_invitee_announcement)
.await
.expect("Failed to publish initial invitee announcement");
send_to_relay(relay, &invitee_state)
.await
.expect("Failed to publish invitee state");
crate::common::push_to_relay(
invitee_git.path(),
&relay.domain(),
&invitee_npub,
identifier,
)
.expect("Failed to push existing invitee repository");
assert_exact_event_served(
relay,
&initial_invitee_announcement,
"Initial invitee announcement",
)
.await;
assert_exact_event_served(relay, &invitee_state, "Invitee state event").await;
}
for clone_url in &invitee_clone_urls {
assert_remote_refs(clone_url, &invitee_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{invitee_branch}"),
Duration::from_secs(20),
)
.await;
}
tokio::time::sleep(Duration::from_secs(1)).await;
let owner_git = tempfile::tempdir().expect("Failed to create owner repository directory");
let owner_commit = create_test_repo_with_commit(owner_git.path(), CommitVariant::PrTest)
.expect("Failed to create owner repository commit");
rename_default_branch(owner_git.path(), owner_branch);
let owner_npub = owner_keys
.public_key()
.to_bech32()
.expect("Failed to encode owner npub");
let (owner_clone_urls, owner_relay_urls) =
repository_urls(&owner_keys, &owner_servers, identifier);
let owner_created_at = Timestamp::from_secs(invitee_state.created_at.as_secs() + 1);
let initial_owner_announcement =
repository_announcement(&owner_keys, &owner_servers, &[], identifier)
.custom_created_at(owner_created_at)
.finalize(&owner_keys)
.expect("Failed to create initial owner announcement");
let owner_state = repository_state_with_default_branch(
&owner_clone_urls,
&owner_relay_urls,
identifier,
owner_branch,
&owner_commit,
)
.custom_created_at(owner_created_at)
.finalize(&owner_keys)
.expect("Failed to create owner state");
assert!(
owner_state.created_at > invitee_state.created_at,
"Owner state must be newer than the invitee's existing state"
);
for relay in owner_servers {
send_to_relay(relay, &initial_owner_announcement)
.await
.expect("Failed to publish initial owner announcement");
send_to_relay(relay, &owner_state)
.await
.expect("Failed to publish owner state");
crate::common::push_to_relay(owner_git.path(), &relay.domain(), &owner_npub, identifier)
.expect("Failed to push owner repository");
assert_exact_event_served(
relay,
&initial_owner_announcement,
"Initial owner announcement",
)
.await;
assert_exact_event_served(relay, &owner_state, "Owner state event").await;
}
let invitation = repository_announcement(
&owner_keys,
&owner_servers,
&[invitee_keys.public_key()],
identifier,
)
.custom_created_at(Timestamp::from_secs(owner_created_at.as_secs() + 1))
.finalize(&owner_keys)
.expect("Failed to create owner invitation");
for relay in owner_servers {
send_to_relay(relay, &invitation)
.await
.expect("Failed to publish owner invitation");
assert_exact_event_served(relay, &invitation, "Owner invitation announcement").await;
}
let invitation_note = invitation
.id
.to_bech32()
.expect("Failed to encode owner invitation note ID");
assert!(
wait_for_log(
&invitee_relay.log_path(),
&invitation_note,
Duration::from_secs(20),
)
.await,
"Invitee-only server should process the unilateral invitation before acceptance"
);
tokio::time::sleep(Duration::from_millis(500)).await;
for relay in invitee_servers {
assert_exact_event_served(relay, &invitee_state, "Invitee state before acceptance").await;
}
for clone_url in &invitee_clone_urls {
assert_remote_refs(clone_url, &invitee_refs, Duration::from_secs(5)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{invitee_branch}"),
Duration::from_secs(5),
)
.await;
}
let pushes_before_acceptance = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
let acceptance = repository_announcement(
&invitee_keys,
&invitee_servers,
&[owner_keys.public_key()],
identifier,
)
.custom_created_at(Timestamp::from_secs(owner_created_at.as_secs() + 2))
.finalize(&invitee_keys)
.expect("Failed to create invitee acceptance");
for relay in invitee_servers {
send_to_relay(relay, &acceptance)
.await
.expect("Failed to publish invitee acceptance");
}
for relay in [&owner_relay, &invitee_relay, &shared_relay] {
assert_exact_event_served(relay, &invitation, "Owner invitation announcement").await;
assert_exact_event_served(relay, &acceptance, "Invitee acceptance announcement").await;
assert_exact_event_served(relay, &owner_state, "Newer owner state event").await;
}
let owner_refs = BTreeMap::from([(format!("refs/heads/{owner_branch}"), owner_commit.clone())]);
for clone_url in &invitee_clone_urls {
assert_remote_refs(clone_url, &owner_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{owner_branch}"),
Duration::from_secs(20),
)
.await;
}
for clone_url in &owner_clone_urls {
assert_remote_refs(clone_url, &owner_refs, Duration::from_secs(20)).await;
assert_remote_default_branch(
clone_url,
&format!("refs/heads/{owner_branch}"),
Duration::from_secs(20),
)
.await;
}
let pushes_after_acceptance = [
push_counts(&owner_relay).await,
push_counts(&invitee_relay).await,
push_counts(&shared_relay).await,
];
assert_eq!(
pushes_after_acceptance, pushes_before_acceptance,
"Invitation acceptance must converge existing repositories without a client Git push"
);
shared_relay.stop().await;
invitee_relay.stop().await;
owner_relay.stop().await;
}
/// A reciprocal owner announcement in purgatory must be enough to unlock an
/// earlier rejected maintainer announcement and use that maintainer's clone URL.
///