mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(sync): recover stalled invitation acceptance
A live ngit acceptance reached gitnostr while its five-second purgatory pass was still holding the sync actor across an older, unbounded backlog. The fresh reciprocal announcement missed the inviter event's two-minute hot-cache window; once expired, the cold index knew the event ID but had forgotten which relay could serve it, leaving an empty invitee repository in purgatory indefinitely. Prioritize unseen announcements within a fixed dependency budget, limit each pass to one relay action, and keep deferred actions discoverable on later ticks. Persist and union the relays that supplied rejected events, retain those sources while connecting, and refetch cold dependencies by exact ID before state and Git convergence continue. Cover the production topology with owner-only and invitee-only relays, an expired hot cache, no shared server, no invitee state event, and no client push. Also prove that a fresh acceptance jumps ahead of a 1,000-event attempted backlog and that source hints survive persistence.
This commit is contained in:
@@ -13,6 +13,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fix maintainer invitation acceptance by prioritizing fresh purgatory
|
||||
dependencies within a bounded reconciliation pass and retaining the source
|
||||
relays needed to recover inviter events after the short-lived hot cache
|
||||
expires.
|
||||
- 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.
|
||||
|
||||
+206
-39
@@ -56,6 +56,61 @@ use crate::nostr::builder::Nip34WritePolicy;
|
||||
use crate::nostr::SharedDatabase;
|
||||
use nostr_relay_builder::prelude::LocalRelay;
|
||||
|
||||
const MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK: usize = 32;
|
||||
const MAX_PURGATORY_FILTER_ACTIONS_PER_TICK: usize = 1;
|
||||
|
||||
fn purgatory_dependency_retry_after() -> Duration {
|
||||
if std::env::var("NGIT_TEST").as_deref() == Ok("1") {
|
||||
Duration::from_secs(2)
|
||||
} else {
|
||||
Duration::from_secs(30)
|
||||
}
|
||||
}
|
||||
|
||||
fn dependency_relay_retention() -> Duration {
|
||||
if std::env::var("NGIT_TEST").as_deref() == Ok("1") {
|
||||
Duration::from_secs(10)
|
||||
} else {
|
||||
Duration::from_secs(60)
|
||||
}
|
||||
}
|
||||
|
||||
fn select_purgatory_dependency_events(
|
||||
mut events: Vec<Event>,
|
||||
attempts: &mut HashMap<EventId, Instant>,
|
||||
now: Instant,
|
||||
retry_after: Duration,
|
||||
limit: usize,
|
||||
) -> Vec<Event> {
|
||||
let current_ids: HashSet<EventId> = events.iter().map(|event| event.id).collect();
|
||||
attempts.retain(|event_id, _| current_ids.contains(event_id));
|
||||
|
||||
events.sort_by(|left, right| {
|
||||
let left_is_new = !attempts.contains_key(&left.id);
|
||||
let right_is_new = !attempts.contains_key(&right.id);
|
||||
right_is_new
|
||||
.cmp(&left_is_new)
|
||||
.then_with(|| right.created_at.cmp(&left.created_at))
|
||||
.then_with(|| right.id.cmp(&left.id))
|
||||
});
|
||||
|
||||
let selected: Vec<Event> = events
|
||||
.into_iter()
|
||||
.filter(|event| {
|
||||
attempts
|
||||
.get(&event.id)
|
||||
.is_none_or(|attempted_at| now.duration_since(*attempted_at) >= retry_after)
|
||||
})
|
||||
.take(limit)
|
||||
.collect();
|
||||
|
||||
for event in &selected {
|
||||
attempts.insert(event.id, now);
|
||||
}
|
||||
|
||||
selected
|
||||
}
|
||||
|
||||
/// Return one stable identity for a relay URL throughout all sync indexes.
|
||||
///
|
||||
/// URL parsers represent a root path as `/`, so announcements that alternate
|
||||
@@ -807,6 +862,10 @@ pub struct SyncManager {
|
||||
connections: HashMap<String, RelayConnection>,
|
||||
/// Last exact-ID dependency recovery attempt, used to bound retries.
|
||||
dependency_refetch_attempts: Arc<std::sync::Mutex<HashMap<EventId, Instant>>>,
|
||||
/// Last dependency pass for each purgatory announcement.
|
||||
purgatory_dependency_attempts: HashMap<EventId, Instant>,
|
||||
/// Temporary source relays retained while rejected dependencies are recovered.
|
||||
dependency_relay_deadlines: HashMap<String, Instant>,
|
||||
/// Health tracker for relay connection state
|
||||
health_tracker: Arc<RelayHealthTracker>,
|
||||
/// Counter for generating unique batch IDs
|
||||
@@ -903,6 +962,8 @@ impl SyncManager {
|
||||
rejected_events_index,
|
||||
connections: HashMap::new(),
|
||||
dependency_refetch_attempts: Arc::new(std::sync::Mutex::new(HashMap::new())),
|
||||
purgatory_dependency_attempts: HashMap::new(),
|
||||
dependency_relay_deadlines: HashMap::new(),
|
||||
health_tracker: Arc::new(RelayHealthTracker::new(config)),
|
||||
next_batch_id: 0,
|
||||
next_connect_attempt_token: 0,
|
||||
@@ -2778,9 +2839,7 @@ impl SyncManager {
|
||||
return;
|
||||
}
|
||||
|
||||
// Register any new entries in repo_sync_index as StateOnly
|
||||
let mut new_relay_urls: std::collections::HashSet<String> =
|
||||
std::collections::HashSet::new();
|
||||
// Register any new entries in repo_sync_index as StateOnly.
|
||||
{
|
||||
let mut index = self.repo_sync_index.write().await;
|
||||
for (repo_id, relays) in &announcements {
|
||||
@@ -2798,35 +2857,7 @@ impl SyncManager {
|
||||
// Don't downgrade an already-Full entry
|
||||
// Add any new relay URLs
|
||||
for relay in relays {
|
||||
if entry.relays.insert(relay.clone()) {
|
||||
new_relay_urls.insert(relay.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if !new_relay_urls.is_empty() {
|
||||
// For any relay URLs that are new, compute and send AddFilters actions
|
||||
let all_targets = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
derive_relay_targets(&repo_index)
|
||||
};
|
||||
|
||||
let actions = {
|
||||
let pending_index = self.pending_sync_index.read().await;
|
||||
let relay_index = self.relay_sync_index.read().await;
|
||||
compute_actions(&all_targets, &pending_index, &relay_index)
|
||||
};
|
||||
|
||||
for action in actions {
|
||||
// Only act on relays that have new URLs (avoids redundant work)
|
||||
if new_relay_urls.contains(&action.relay_url) {
|
||||
tracing::info!(
|
||||
relay = %action.relay_url,
|
||||
repos = action.items.repos.len(),
|
||||
"Purgatory sync timer: connecting to new relay from purgatory announcement"
|
||||
);
|
||||
self.handle_new_sync_filters(action).await;
|
||||
entry.relays.insert(relay.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -2835,10 +2866,47 @@ impl SyncManager {
|
||||
// before the reciprocal owner announcement reached this relay. Retry those
|
||||
// dependencies after connecting the owner's relay hints so cold-cache
|
||||
// misses can be fetched directly by event ID.
|
||||
for event in announcement_events {
|
||||
let dependency_events = select_purgatory_dependency_events(
|
||||
announcement_events,
|
||||
&mut self.purgatory_dependency_attempts,
|
||||
Instant::now(),
|
||||
purgatory_dependency_retry_after(),
|
||||
MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK,
|
||||
);
|
||||
for event in dependency_events {
|
||||
self.reprocess_purgatory_announcement_dependencies(&event)
|
||||
.await;
|
||||
}
|
||||
|
||||
// Recompute all outstanding actions on every pass. A bounded pass may
|
||||
// intentionally defer some relays, so restricting this to newly seen
|
||||
// URLs would lose that work on the next tick.
|
||||
let all_targets = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
derive_relay_targets(&repo_index)
|
||||
};
|
||||
let mut actions = {
|
||||
let pending_index = self.pending_sync_index.read().await;
|
||||
let relay_index = self.relay_sync_index.read().await;
|
||||
compute_actions(&all_targets, &pending_index, &relay_index)
|
||||
};
|
||||
actions.sort_by_key(|action| {
|
||||
!self
|
||||
.dependency_relay_deadlines
|
||||
.contains_key(&action.relay_url)
|
||||
});
|
||||
|
||||
for action in actions
|
||||
.into_iter()
|
||||
.take(MAX_PURGATORY_FILTER_ACTIONS_PER_TICK)
|
||||
{
|
||||
tracing::info!(
|
||||
relay = %action.relay_url,
|
||||
repos = action.items.repos.len(),
|
||||
"Purgatory sync timer: processing one bounded relay action"
|
||||
);
|
||||
self.handle_new_sync_filters(action).await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Retry events whose authorization depends on a newly admitted owner announcement.
|
||||
@@ -2847,7 +2915,7 @@ impl SyncManager {
|
||||
/// `process_event_static`, but runs while the owner announcement is still in
|
||||
/// purgatory. Announcement policy already treats purgatory announcements as
|
||||
/// maintainer authority; doing the retry here makes arrival order irrelevant.
|
||||
async fn reprocess_purgatory_announcement_dependencies(&self, event: &Event) {
|
||||
async fn reprocess_purgatory_announcement_dependencies(&mut self, event: &Event) {
|
||||
let announcement =
|
||||
match crate::nostr::events::RepositoryAnnouncement::from_event(event.clone()) {
|
||||
Ok(announcement) => announcement,
|
||||
@@ -2890,6 +2958,11 @@ impl SyncManager {
|
||||
rejected_index::EventType::Announcement,
|
||||
rejected_index::EventType::State,
|
||||
] {
|
||||
dependency_relay_urls.extend(self.rejected_events_index.dependency_relay_hints(
|
||||
&maintainer_pubkey,
|
||||
&announcement.identifier,
|
||||
Some(event_type),
|
||||
));
|
||||
let (event_ids, hot_events) = self.rejected_events_index.dependency_candidates(
|
||||
&maintainer_pubkey,
|
||||
&announcement.identifier,
|
||||
@@ -2944,6 +3017,11 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
// The owner's state may also have arrived before their announcement.
|
||||
dependency_relay_urls.extend(self.rejected_events_index.dependency_relay_hints(
|
||||
&event.pubkey,
|
||||
&announcement.identifier,
|
||||
Some(rejected_index::EventType::State),
|
||||
));
|
||||
let (event_ids, hot_events) = self.rejected_events_index.dependency_candidates(
|
||||
&event.pubkey,
|
||||
&announcement.identifier,
|
||||
@@ -2983,9 +3061,43 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
if !dependency_refetch_ids.is_empty() {
|
||||
let dependency_relay_urls: Vec<String> = dependency_relay_urls.into_iter().collect();
|
||||
let mut canonical_dependency_relays = Vec::new();
|
||||
for relay_url in dependency_relay_urls {
|
||||
let relay_url = match canonical_relay_key(&relay_url) {
|
||||
Ok(relay_url) => relay_url,
|
||||
Err(error) => {
|
||||
tracing::warn!(
|
||||
relay = %relay_url,
|
||||
error = %error,
|
||||
"Ignoring invalid rejected dependency source relay"
|
||||
);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
self.dependency_relay_deadlines.insert(
|
||||
relay_url.clone(),
|
||||
Instant::now() + dependency_relay_retention(),
|
||||
);
|
||||
if !self.connections.contains_key(&relay_url) {
|
||||
self.register_relay(relay_url.clone(), false).await;
|
||||
self.schedule_connect_relay(&relay_url).await;
|
||||
}
|
||||
canonical_dependency_relays.push(relay_url);
|
||||
}
|
||||
|
||||
let connected_dependency_relays: Vec<String> = {
|
||||
let relay_index = self.relay_sync_index.read().await;
|
||||
canonical_dependency_relays
|
||||
.into_iter()
|
||||
.filter(|relay_url| {
|
||||
relay_index
|
||||
.get(relay_url)
|
||||
.is_some_and(|state| state.connection_status.is_live_sync_active())
|
||||
})
|
||||
.collect()
|
||||
};
|
||||
self.spawn_purgatory_dependency_refetch(
|
||||
&dependency_relay_urls,
|
||||
&connected_dependency_relays,
|
||||
dependency_refetch_ids,
|
||||
&announcement.identifier,
|
||||
);
|
||||
@@ -3679,11 +3791,12 @@ impl SyncManager {
|
||||
|
||||
// Use appropriate method based on event kind
|
||||
if event.kind == Kind::RepoState {
|
||||
rejected_events_index.add_state(
|
||||
rejected_events_index.add_state_from_relay(
|
||||
event.clone(),
|
||||
event.pubkey,
|
||||
identifier.to_string(),
|
||||
reason,
|
||||
Some(relay_url.to_string()),
|
||||
);
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
@@ -3692,11 +3805,12 @@ impl SyncManager {
|
||||
"Added rejected state event to two-tier index"
|
||||
);
|
||||
} else {
|
||||
rejected_events_index.add_announcement(
|
||||
rejected_events_index.add_announcement_from_relay(
|
||||
event.clone(),
|
||||
event.pubkey,
|
||||
identifier.to_string(),
|
||||
reason,
|
||||
Some(relay_url.to_string()),
|
||||
);
|
||||
tracing::debug!(
|
||||
event_id = %event.id,
|
||||
@@ -3874,12 +3988,17 @@ impl SyncManager {
|
||||
///
|
||||
/// Bootstrap relays are NEVER disconnected, even if empty.
|
||||
async fn check_disconnects(&mut self) {
|
||||
let desired_relays: HashSet<String> = {
|
||||
let now = Instant::now();
|
||||
self.dependency_relay_deadlines
|
||||
.retain(|_, deadline| *deadline > now);
|
||||
|
||||
let mut desired_relays: HashSet<String> = {
|
||||
let repo_index = self.repo_sync_index.read().await;
|
||||
algorithms::derive_relay_targets(&repo_index)
|
||||
.into_keys()
|
||||
.collect()
|
||||
};
|
||||
desired_relays.extend(self.dependency_relay_deadlines.keys().cloned());
|
||||
|
||||
// Collect relays to disconnect
|
||||
let to_disconnect: Vec<String> = {
|
||||
@@ -4580,6 +4699,54 @@ impl SyncManager {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn purgatory_dependency_budget_prioritizes_a_fresh_announcement() {
|
||||
let keys = Keys::generate();
|
||||
let now = Instant::now();
|
||||
let retry_after = Duration::from_secs(30);
|
||||
let mut attempts = HashMap::new();
|
||||
let mut events = Vec::new();
|
||||
|
||||
for created_at in 1..=1000 {
|
||||
let event = EventBuilder::new(Kind::GitRepoAnnouncement, "")
|
||||
.tag(Tag::identifier(format!("old-{created_at}")))
|
||||
.custom_created_at(Timestamp::from_secs(created_at))
|
||||
.finalize(&keys)
|
||||
.expect("Failed to create old purgatory announcement");
|
||||
attempts.insert(event.id, now);
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
let fresh = EventBuilder::new(Kind::GitRepoAnnouncement, "")
|
||||
.tag(Tag::identifier("fresh"))
|
||||
.custom_created_at(Timestamp::from_secs(1001))
|
||||
.finalize(&keys)
|
||||
.expect("Failed to create fresh purgatory announcement");
|
||||
events.push(fresh.clone());
|
||||
|
||||
let selected = select_purgatory_dependency_events(
|
||||
events,
|
||||
&mut attempts,
|
||||
now,
|
||||
retry_after,
|
||||
MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK,
|
||||
);
|
||||
|
||||
assert_eq!(selected.len(), 1);
|
||||
assert_eq!(selected[0].id, fresh.id);
|
||||
assert!(
|
||||
select_purgatory_dependency_events(
|
||||
vec![fresh],
|
||||
&mut attempts,
|
||||
now,
|
||||
retry_after,
|
||||
MAX_PURGATORY_DEPENDENCY_EVENTS_PER_TICK,
|
||||
)
|
||||
.is_empty(),
|
||||
"an attempted announcement must wait for its bounded retry deadline"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn connect_attempt_tokens_deduplicate_and_reject_stale_results() {
|
||||
let relay = "wss://relay.example";
|
||||
|
||||
+185
-6
@@ -140,6 +140,7 @@ struct HotCacheEntry {
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
relay_hints: HashSet<String>,
|
||||
#[allow(dead_code)] // Used for metrics/debugging in future
|
||||
reason: RejectionReason,
|
||||
cached_at: Instant,
|
||||
@@ -154,6 +155,8 @@ struct SerializableHotCacheEntry {
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
#[serde(default)]
|
||||
relay_hints: HashSet<String>,
|
||||
reason: RejectionReason,
|
||||
/// Duration since saved_at when this entry was cached
|
||||
cached_at_offset_secs: u64,
|
||||
@@ -167,6 +170,7 @@ struct ColdIndexEntry {
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
relay_hints: HashSet<String>,
|
||||
#[allow(dead_code)] // Used for metrics/debugging in future
|
||||
reason: RejectionReason,
|
||||
rejected_at: Instant,
|
||||
@@ -180,6 +184,8 @@ struct SerializableColdIndexEntry {
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
#[serde(default)]
|
||||
relay_hints: HashSet<String>,
|
||||
reason: RejectionReason,
|
||||
/// Duration since saved_at when this entry was rejected
|
||||
rejected_at_offset_secs: u64,
|
||||
@@ -237,6 +243,7 @@ impl HotCache {
|
||||
}
|
||||
|
||||
/// Add event to hot cache
|
||||
#[cfg(test)]
|
||||
fn add(
|
||||
&self,
|
||||
event: Event,
|
||||
@@ -244,12 +251,32 @@ impl HotCache {
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
self.add_with_relay_hints(
|
||||
event,
|
||||
pubkey,
|
||||
identifier,
|
||||
event_type,
|
||||
HashSet::new(),
|
||||
reason,
|
||||
);
|
||||
}
|
||||
|
||||
fn add_with_relay_hints(
|
||||
&self,
|
||||
event: Event,
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
relay_hints: HashSet<String>,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
let entry = HotCacheEntry {
|
||||
event,
|
||||
pubkey,
|
||||
identifier,
|
||||
event_type,
|
||||
relay_hints,
|
||||
reason,
|
||||
cached_at: Instant::now(),
|
||||
};
|
||||
@@ -355,6 +382,7 @@ impl ColdIndex {
|
||||
}
|
||||
|
||||
/// Add metadata to cold index
|
||||
#[cfg(test)]
|
||||
fn add(
|
||||
&self,
|
||||
event_id: EventId,
|
||||
@@ -363,15 +391,41 @@ impl ColdIndex {
|
||||
event_type: EventType,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
self.add_with_relay_hints(
|
||||
event_id,
|
||||
pubkey,
|
||||
identifier,
|
||||
event_type,
|
||||
HashSet::new(),
|
||||
reason,
|
||||
);
|
||||
}
|
||||
|
||||
fn add_with_relay_hints(
|
||||
&self,
|
||||
event_id: EventId,
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
event_type: EventType,
|
||||
relay_hints: HashSet<String>,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
let mut entries = self.entries.write().unwrap();
|
||||
if let Some(entry) = entries.get_mut(&event_id) {
|
||||
entry.relay_hints.extend(relay_hints);
|
||||
return;
|
||||
}
|
||||
|
||||
let entry = ColdIndexEntry {
|
||||
pubkey,
|
||||
identifier,
|
||||
event_type,
|
||||
relay_hints,
|
||||
reason,
|
||||
rejected_at: Instant::now(),
|
||||
};
|
||||
|
||||
self.entries.write().unwrap().insert(event_id, entry);
|
||||
entries.insert(event_id, entry);
|
||||
}
|
||||
|
||||
/// Check if event is in cold index
|
||||
@@ -433,6 +487,29 @@ impl ColdIndex {
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn dependency_relay_hints(
|
||||
&self,
|
||||
maintainer_pubkey: &PublicKey,
|
||||
identifier: &str,
|
||||
event_type: Option<EventType>,
|
||||
) -> HashSet<String> {
|
||||
let entries = self.entries.read().unwrap();
|
||||
let now = Instant::now();
|
||||
|
||||
entries
|
||||
.values()
|
||||
.filter(|entry| {
|
||||
let matches_type = event_type.is_none_or(|et| entry.event_type == et);
|
||||
entry.pubkey == *maintainer_pubkey
|
||||
&& entry.identifier == identifier
|
||||
&& matches_type
|
||||
&& entry.reason != RejectionReason::Other
|
||||
&& now.duration_since(entry.rejected_at) < self.expiry_duration
|
||||
})
|
||||
.flat_map(|entry| entry.relay_hints.iter().cloned())
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn remove(&self, event_id: &EventId) {
|
||||
self.entries.write().unwrap().remove(event_id);
|
||||
}
|
||||
@@ -542,21 +619,36 @@ impl RejectedEventsIndex {
|
||||
identifier: String,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
self.add_announcement_from_relay(event, pubkey, identifier, reason, None);
|
||||
}
|
||||
|
||||
pub fn add_announcement_from_relay(
|
||||
&self,
|
||||
event: Event,
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
reason: RejectionReason,
|
||||
relay_url: Option<String>,
|
||||
) {
|
||||
let relay_hints: HashSet<String> = relay_url.into_iter().collect();
|
||||
|
||||
// Add to hot cache (full event)
|
||||
self.hot_cache.add(
|
||||
self.hot_cache.add_with_relay_hints(
|
||||
event.clone(),
|
||||
pubkey,
|
||||
identifier.clone(),
|
||||
EventType::Announcement,
|
||||
relay_hints.clone(),
|
||||
reason,
|
||||
);
|
||||
|
||||
// Add to cold index (metadata only)
|
||||
self.cold_index.add(
|
||||
self.cold_index.add_with_relay_hints(
|
||||
event.id,
|
||||
pubkey,
|
||||
identifier,
|
||||
EventType::Announcement,
|
||||
relay_hints,
|
||||
reason,
|
||||
);
|
||||
|
||||
@@ -579,18 +671,38 @@ impl RejectedEventsIndex {
|
||||
identifier: String,
|
||||
reason: RejectionReason,
|
||||
) {
|
||||
self.add_state_from_relay(event, pubkey, identifier, reason, None);
|
||||
}
|
||||
|
||||
pub fn add_state_from_relay(
|
||||
&self,
|
||||
event: Event,
|
||||
pubkey: PublicKey,
|
||||
identifier: String,
|
||||
reason: RejectionReason,
|
||||
relay_url: Option<String>,
|
||||
) {
|
||||
let relay_hints: HashSet<String> = relay_url.into_iter().collect();
|
||||
|
||||
// Add to hot cache (full event)
|
||||
self.hot_cache.add(
|
||||
self.hot_cache.add_with_relay_hints(
|
||||
event.clone(),
|
||||
pubkey,
|
||||
identifier.clone(),
|
||||
EventType::State,
|
||||
relay_hints.clone(),
|
||||
reason,
|
||||
);
|
||||
|
||||
// Add to cold index (metadata only)
|
||||
self.cold_index
|
||||
.add(event.id, pubkey, identifier, EventType::State, reason);
|
||||
self.cold_index.add_with_relay_hints(
|
||||
event.id,
|
||||
pubkey,
|
||||
identifier,
|
||||
EventType::State,
|
||||
relay_hints,
|
||||
reason,
|
||||
);
|
||||
|
||||
// Update metrics
|
||||
self.update_metrics_for_type("state");
|
||||
@@ -681,6 +793,16 @@ impl RejectedEventsIndex {
|
||||
)
|
||||
}
|
||||
|
||||
pub fn dependency_relay_hints(
|
||||
&self,
|
||||
pubkey: &PublicKey,
|
||||
identifier: &str,
|
||||
event_type: Option<EventType>,
|
||||
) -> HashSet<String> {
|
||||
self.cold_index
|
||||
.dependency_relay_hints(pubkey, identifier, event_type)
|
||||
}
|
||||
|
||||
/// Remove a successfully processed event from both rejected-event tiers.
|
||||
pub fn remove(&self, event_id: &EventId) {
|
||||
self.hot_cache.remove(event_id);
|
||||
@@ -776,6 +898,7 @@ impl RejectedEventsIndex {
|
||||
pubkey: entry.pubkey,
|
||||
identifier: entry.identifier.clone(),
|
||||
event_type: entry.event_type,
|
||||
relay_hints: entry.relay_hints.clone(),
|
||||
reason: entry.reason,
|
||||
cached_at_offset_secs,
|
||||
};
|
||||
@@ -794,6 +917,7 @@ impl RejectedEventsIndex {
|
||||
pubkey: entry.pubkey,
|
||||
identifier: entry.identifier.clone(),
|
||||
event_type: entry.event_type,
|
||||
relay_hints: entry.relay_hints.clone(),
|
||||
reason: entry.reason,
|
||||
rejected_at_offset_secs,
|
||||
};
|
||||
@@ -869,6 +993,7 @@ impl RejectedEventsIndex {
|
||||
pubkey: serializable_entry.pubkey,
|
||||
identifier: serializable_entry.identifier,
|
||||
event_type: serializable_entry.event_type,
|
||||
relay_hints: serializable_entry.relay_hints,
|
||||
reason: serializable_entry.reason,
|
||||
cached_at,
|
||||
};
|
||||
@@ -889,6 +1014,7 @@ impl RejectedEventsIndex {
|
||||
pubkey: serializable_entry.pubkey,
|
||||
identifier: serializable_entry.identifier,
|
||||
event_type: serializable_entry.event_type,
|
||||
relay_hints: serializable_entry.relay_hints,
|
||||
reason: serializable_entry.reason,
|
||||
rejected_at,
|
||||
};
|
||||
@@ -941,6 +1067,59 @@ mod tests {
|
||||
assert_eq!(retrieved[0].id, event.id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_cold_dependency_relay_hints_are_unioned_and_persisted() {
|
||||
let index = RejectedEventsIndex::new(Duration::from_millis(1), Duration::from_secs(604800));
|
||||
let event = create_test_event().await;
|
||||
let identifier = "relay-hint-recovery".to_string();
|
||||
let reason = RejectionReason::MaintainerNotYetValid;
|
||||
|
||||
index.add_announcement_from_relay(
|
||||
event.clone(),
|
||||
event.pubkey,
|
||||
identifier.clone(),
|
||||
reason,
|
||||
Some("wss://relay-a.example".to_string()),
|
||||
);
|
||||
index.add_announcement_from_relay(
|
||||
event.clone(),
|
||||
event.pubkey,
|
||||
identifier.clone(),
|
||||
reason,
|
||||
Some("wss://relay-b.example".to_string()),
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
index
|
||||
.dependency_relay_hints(&event.pubkey, &identifier, Some(EventType::Announcement),),
|
||||
HashSet::from([
|
||||
"wss://relay-a.example".to_string(),
|
||||
"wss://relay-b.example".to_string(),
|
||||
])
|
||||
);
|
||||
|
||||
let directory = tempfile::tempdir().expect("Failed to create cache directory");
|
||||
let path = directory.path().join("rejected-events-cache.json");
|
||||
index.save_to_disk(&path).expect("Failed to save cache");
|
||||
|
||||
let restored =
|
||||
RejectedEventsIndex::new(Duration::from_millis(1), Duration::from_secs(604800));
|
||||
restored
|
||||
.restore_from_disk(&path)
|
||||
.expect("Failed to restore cache");
|
||||
assert_eq!(
|
||||
restored.dependency_relay_hints(
|
||||
&event.pubkey,
|
||||
&identifier,
|
||||
Some(EventType::Announcement),
|
||||
),
|
||||
HashSet::from([
|
||||
"wss://relay-a.example".to_string(),
|
||||
"wss://relay-b.example".to_string(),
|
||||
])
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_hot_cache_expires_after_duration() {
|
||||
let cache = HotCache::new(Duration::from_millis(50));
|
||||
|
||||
@@ -424,6 +424,170 @@ async fn test_existing_repository_invitation_acceptance_syncs_without_invitee_pu
|
||||
owner_relay.stop().await;
|
||||
}
|
||||
|
||||
/// Acceptance on an invitee-only server must recover an expired owner
|
||||
/// invitation from the owner-only server that originally supplied it.
|
||||
///
|
||||
/// This matches the CLI flow where the owner has already pushed and issued a
|
||||
/// state event, while the invitee publishes only a reciprocal announcement.
|
||||
/// The rejected owner announcement is allowed to leave the one-second hot
|
||||
/// cache before acceptance, so convergence depends on the persisted source
|
||||
/// relay hint and exact-ID recovery rather than a shared GRASP server.
|
||||
#[tokio::test]
|
||||
async fn test_invitee_only_acceptance_recovers_cold_owner_invitation() {
|
||||
use crate::common::{create_state_event, create_test_repo_with_commit, CommitVariant};
|
||||
|
||||
let owner_relay = TestRelay::start_with_sync(None).await;
|
||||
let owner_keys = Keys::generate();
|
||||
let invitee_keys = Keys::generate();
|
||||
let identifier = "cold-owner-invitation-acceptance";
|
||||
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");
|
||||
let owner_npub = owner_keys
|
||||
.public_key()
|
||||
.to_bech32()
|
||||
.expect("Failed to encode owner npub");
|
||||
let owner_servers = [&owner_relay];
|
||||
let (owner_clone_urls, owner_relay_urls) =
|
||||
repository_urls(&owner_keys, &owner_servers, identifier);
|
||||
let invitation = repository_announcement(
|
||||
&owner_keys,
|
||||
&owner_servers,
|
||||
&[invitee_keys.public_key()],
|
||||
identifier,
|
||||
)
|
||||
.finalize(&owner_keys)
|
||||
.expect("Failed to create owner invitation");
|
||||
let owner_state = create_state_event(
|
||||
&owner_keys,
|
||||
identifier,
|
||||
&[("main", &owner_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");
|
||||
|
||||
send_to_relay(&owner_relay, &invitation)
|
||||
.await
|
||||
.expect("Failed to publish owner invitation");
|
||||
send_to_relay(&owner_relay, &owner_state)
|
||||
.await
|
||||
.expect("Failed to publish owner state");
|
||||
crate::common::push_to_relay(
|
||||
owner_git.path(),
|
||||
&owner_relay.domain(),
|
||||
&owner_npub,
|
||||
identifier,
|
||||
)
|
||||
.expect("Failed to push owner repository");
|
||||
assert_exact_event_served(&owner_relay, &invitation, "Owner invitation").await;
|
||||
assert_exact_event_served(&owner_relay, &owner_state, "Owner state").await;
|
||||
|
||||
let invitee_relay =
|
||||
TestRelay::start_with_sync_and_rejected_hot_cache(Some(owner_relay.url().to_string()), 1)
|
||||
.await;
|
||||
let invitation_note = invitation
|
||||
.id
|
||||
.to_bech32()
|
||||
.expect("Failed to encode invitation event ID");
|
||||
assert!(
|
||||
wait_for_log(
|
||||
&invitee_relay.log_path(),
|
||||
&invitation_note,
|
||||
Duration::from_secs(10),
|
||||
)
|
||||
.await,
|
||||
"Invitee-only server should reject and index the owner invitation"
|
||||
);
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
|
||||
let pushes_before_acceptance = [
|
||||
push_counts(&owner_relay).await,
|
||||
push_counts(&invitee_relay).await,
|
||||
];
|
||||
let invitee_servers = [&invitee_relay];
|
||||
let acceptance = repository_announcement(
|
||||
&invitee_keys,
|
||||
&invitee_servers,
|
||||
&[owner_keys.public_key()],
|
||||
identifier,
|
||||
)
|
||||
.custom_created_at(Timestamp::from_secs(invitation.created_at.as_secs() + 1))
|
||||
.finalize(&invitee_keys)
|
||||
.expect("Failed to create invitee acceptance");
|
||||
let (invitee_clone_urls, _) = repository_urls(&invitee_keys, &invitee_servers, identifier);
|
||||
assert_eq!(announcement_clone_urls(&acceptance), invitee_clone_urls);
|
||||
assert!(
|
||||
!announcement_clone_urls(&acceptance)
|
||||
.iter()
|
||||
.any(|url| owner_clone_urls.contains(url)),
|
||||
"Acceptance must not copy the owner's Git endpoint"
|
||||
);
|
||||
|
||||
send_to_relay(&invitee_relay, &acceptance)
|
||||
.await
|
||||
.expect("Failed to publish invitee acceptance");
|
||||
assert!(
|
||||
wait_for_log(
|
||||
&invitee_relay.log_path(),
|
||||
"Fetched purgatory dependencies by exact event ID",
|
||||
Duration::from_secs(20),
|
||||
)
|
||||
.await,
|
||||
"Invitee-only server should recover the expired invitation by exact event ID"
|
||||
);
|
||||
|
||||
assert_exact_event_served(&invitee_relay, &invitation, "Recovered owner invitation").await;
|
||||
assert_exact_event_served(&invitee_relay, &owner_state, "Recovered owner state").await;
|
||||
assert_exact_event_served(&invitee_relay, &acceptance, "Invitee acceptance").await;
|
||||
|
||||
let expected_refs = BTreeMap::from([("refs/heads/main".to_string(), owner_commit.clone())]);
|
||||
assert_remote_refs(
|
||||
&invitee_clone_urls[0],
|
||||
&expected_refs,
|
||||
Duration::from_secs(20),
|
||||
)
|
||||
.await;
|
||||
assert_remote_default_branch(
|
||||
&invitee_clone_urls[0],
|
||||
"refs/heads/main",
|
||||
Duration::from_secs(20),
|
||||
)
|
||||
.await;
|
||||
|
||||
let invitee_state_exists = wait_for_event_on_relay(
|
||||
invitee_relay.url(),
|
||||
Filter::new()
|
||||
.kind(Kind::RepoState)
|
||||
.author(invitee_keys.public_key())
|
||||
.identifier(identifier),
|
||||
Duration::from_secs(1),
|
||||
)
|
||||
.await;
|
||||
assert!(
|
||||
!invitee_state_exists,
|
||||
"Invitation acceptance must not publish an invitee state event"
|
||||
);
|
||||
assert_eq!(
|
||||
[
|
||||
push_counts(&owner_relay).await,
|
||||
push_counts(&invitee_relay).await,
|
||||
],
|
||||
pushes_before_acceptance,
|
||||
"Invitation acceptance must converge without a client Git push"
|
||||
);
|
||||
|
||||
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.
|
||||
///
|
||||
@@ -974,13 +1138,10 @@ async fn test_purgatory_owner_uses_rejected_maintainer_clone_to_sync_git() {
|
||||
.tags(vec![
|
||||
Tag::identifier(identifier),
|
||||
Tag::custom("clone", vec![owner_clone.clone()]),
|
||||
Tag::custom(
|
||||
"relays",
|
||||
vec![
|
||||
source_relay.url().to_string(),
|
||||
target_relay.url().to_string(),
|
||||
],
|
||||
),
|
||||
// The reciprocal announcement intentionally advertises only the
|
||||
// target. Cold dependency recovery must remember that the
|
||||
// rejected maintainer event was actually received from source.
|
||||
Tag::custom("relays", vec![target_relay.url().to_string()]),
|
||||
Tag::custom("maintainers", vec![maintainer_keys.public_key().to_hex()]),
|
||||
])
|
||||
.finalize(&owner_keys)
|
||||
@@ -1215,100 +1376,6 @@ async fn test_maintainer_announcement_reprocessed_immediately() {
|
||||
relay_b.stop().await;
|
||||
}
|
||||
|
||||
/// Test that maintainer announcements NOT in hot cache are still prevented from re-fetching
|
||||
///
|
||||
/// Flow:
|
||||
/// 1. Maintainer announcement arrives → Rejected (added to hot cache + cold index)
|
||||
/// 2. Wait for hot cache to expire (2+ minutes)
|
||||
/// 3. Owner announcement arrives → Invalidates cold index
|
||||
/// 4. Maintainer announcement should NOT be re-fetched (cold index prevents)
|
||||
/// 5. Only owner announcement should be in database
|
||||
///
|
||||
/// This test verifies the cold index prevents repeated downloads after hot cache expiry.
|
||||
/// Note: This test is slow (2+ minutes) so we'll skip it in normal test runs.
|
||||
#[tokio::test]
|
||||
#[ignore] // Skip by default due to 2+ minute duration
|
||||
async fn test_maintainer_announcement_cold_index_prevents_refetch() {
|
||||
let relay = TestRelay::start().await;
|
||||
|
||||
// Create keys
|
||||
let owner_keys = Keys::generate();
|
||||
let maintainer_keys = Keys::generate();
|
||||
|
||||
let identifier = "test-repo-cold";
|
||||
|
||||
// Create client using TestClient helper
|
||||
let client = TestClient::new(relay.url(), maintainer_keys.clone())
|
||||
.await
|
||||
.expect("Failed to connect to relay");
|
||||
|
||||
// Step 1: Send maintainer announcement (will be rejected - doesn't list our relay)
|
||||
let maintainer_announcement =
|
||||
EventBuilder::new(Kind::GitRepoAnnouncement, "Maintainer's repository")
|
||||
.tags(vec![
|
||||
Tag::identifier(identifier),
|
||||
Tag::custom(
|
||||
"clone",
|
||||
vec![format!("https://example.com/{}.git", identifier)],
|
||||
),
|
||||
Tag::custom("relays", vec!["wss://example.com".to_string()]),
|
||||
])
|
||||
.finalize(&maintainer_keys)
|
||||
.unwrap();
|
||||
|
||||
// Send maintainer announcement - expect it to be rejected
|
||||
let _ = client.send_event(&maintainer_announcement).await;
|
||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||
|
||||
// Step 2: Wait for hot cache to expire (default: 120 seconds)
|
||||
println!("⏳ Waiting for hot cache to expire (120 seconds)...");
|
||||
tokio::time::sleep(Duration::from_secs(125)).await;
|
||||
|
||||
// Step 3: Send owner announcement (lists maintainer)
|
||||
let owner_announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "Owner's repository")
|
||||
.tags(vec![
|
||||
Tag::identifier(identifier),
|
||||
Tag::custom(
|
||||
"clone",
|
||||
vec![format!("https://{}/{}.git", relay.domain(), identifier)],
|
||||
),
|
||||
Tag::custom("relays", vec![relay.url().to_string()]),
|
||||
Tag::custom("maintainers", vec![maintainer_keys.public_key().to_hex()]),
|
||||
])
|
||||
.finalize(&owner_keys)
|
||||
.unwrap();
|
||||
|
||||
client.send_event(&owner_announcement).await.unwrap();
|
||||
tokio::time::sleep(Duration::from_millis(500)).await;
|
||||
|
||||
// Step 4: Verify only owner announcement is in database
|
||||
let owner_filter = Filter::new()
|
||||
.kind(Kind::GitRepoAnnouncement)
|
||||
.author(owner_keys.public_key())
|
||||
.identifier(identifier);
|
||||
|
||||
let owner_found =
|
||||
wait_for_event_on_relay(relay.url(), owner_filter, Duration::from_secs(2)).await;
|
||||
assert!(owner_found, "Owner announcement should be accepted");
|
||||
|
||||
let maintainer_filter = Filter::new()
|
||||
.kind(Kind::GitRepoAnnouncement)
|
||||
.author(maintainer_keys.public_key())
|
||||
.identifier(identifier);
|
||||
|
||||
let maintainer_found =
|
||||
wait_for_event_on_relay(relay.url(), maintainer_filter, Duration::from_millis(500)).await;
|
||||
assert!(
|
||||
!maintainer_found,
|
||||
"Maintainer announcement should NOT be re-processed (hot cache expired)"
|
||||
);
|
||||
|
||||
println!("✅ Cold index prevented re-fetch after hot cache expiry");
|
||||
|
||||
client.disconnect().await;
|
||||
relay.stop().await;
|
||||
}
|
||||
|
||||
/// Test that all maintainer announcements are re-processed when the owner announcement
|
||||
/// is promoted from purgatory via a git push.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user