diff --git a/CHANGELOG.md b/CHANGELOG.md index 47bd5f0..bd15449 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,6 +20,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Keep purgatory announcements, state events, and PR events alive while a + concrete background Git sync for their repository is actively running. A + large healthy clone can exceed the nominal 30-minute purgatory window; + queued and backoff-only work remains eligible for normal expiry so broken + repositories cannot extend retention indefinitely. - Preserve purgatory and rejected-event recovery across abrupt termination. Their snapshots now remain after restore and are atomically replaced every 60 seconds, bounding crash loss to one interval instead of consuming the diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index b4f1444..3737c64 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -307,6 +307,8 @@ pub struct Purgatory { 5. **Automatic Expiry**: 30-minute default expiry, extensible during processing - Background cleanup task runs every 60 seconds - Removes expired entries from all stores + - Defers matching expiry while a concrete background Git sync is actively + running; queued/backoff work receives no extension #### Data Types diff --git a/docs/explanation/purgatory-design.md b/docs/explanation/purgatory-design.md index bc98052..a801e85 100644 --- a/docs/explanation/purgatory-design.md +++ b/docs/explanation/purgatory-design.md @@ -348,6 +348,8 @@ fetch_repository_data_with_purgatory() checks DB + purgatory The protocol specifies 30-minute expiry for announcements. We implement a two-phase soft expiry: **Phase 1 — Initial 30-minute expiry (`soft_expired == false`):** +- Defer the expiry sweep while a concrete background Git sync for the + repository is actively running - Delete the bare git repo (frees disk space, respects protocol expiry) - Set `soft_expired = true` - Extend `expires_at` by 24 hours (`SOFT_EXPIRY_EXTENDED`) @@ -378,6 +380,14 @@ The 30-minute purgatory timer is reset (extended) in three scenarios: All three call `purgatory.extend_announcement_expiry(owner, identifier, 1800s)`. +Background Git synchronization does not reset the 30-minute deadline. Instead, +the expiry sweep skips matching announcements, state events and PR events only +while `sync_identifier` is actively running. This lets a large but healthy +clone finish without allowing an unavailable repository to live forever: +queue entries waiting for their next attempt or in backoff remain eligible for +normal expiry. When active work finishes, it either promotes the satisfied +events or the next sweep expires any still-unsatisfied entries. + ### Purgatory Lifecycle ``` diff --git a/src/purgatory/mod.rs b/src/purgatory/mod.rs index d2d9932..48baa5a 100644 --- a/src/purgatory/mod.rs +++ b/src/purgatory/mod.rs @@ -344,6 +344,35 @@ impl Purgatory { false } + /// Whether a concrete background Git sync is currently running for an + /// identifier. + /// + /// Queue membership alone is deliberately insufficient: entries waiting + /// for their next backoff attempt must still expire normally when remote + /// Git data remains unavailable. Only the interval in which + /// `sync_identifier` is actively working protects matching purgatory + /// entries from the expiry sweep. + fn identifier_sync_in_progress(&self, identifier: &str) -> bool { + self.sync_queue + .get(identifier) + .is_some_and(|entry| entry.in_progress) + } + + /// Whether any repository referenced by an event has active Git sync. + fn event_sync_in_progress(&self, event: &Event) -> bool { + event.tags.iter().any(|tag| { + let values = tag.clone().to_vec(); + if values.len() < 2 || values[0] != "a" || !values[1].starts_with("30617:") { + return false; + } + + values[1] + .splitn(3, ':') + .nth(2) + .is_some_and(|identifier| self.identifier_sync_in_progress(identifier)) + }) + } + /// Get a reference to the sync queue (for the sync loop). pub fn sync_queue(&self) -> &Arc> { &self.sync_queue @@ -1143,7 +1172,9 @@ impl Purgatory { let expired_announcements: Vec<(PublicKey, String, PathBuf, EventId, bool)> = self .announcement_purgatory .iter() - .filter(|entry| entry.value().expires_at <= now) + .filter(|entry| { + entry.value().expires_at <= now && !self.identifier_sync_in_progress(&entry.key().1) + }) .map(|entry| { let key = entry.key(); let v = entry.value(); @@ -1224,6 +1255,10 @@ impl Purgatory { // Remove expired state events and mark them as expired self.state_events.retain(|identifier, entries| { + if self.identifier_sync_in_progress(identifier) { + return true; + } + let original_len = entries.len(); // Log and collect expired entries before removing @@ -1267,7 +1302,14 @@ impl Purgatory { let expired_prs: Vec<_> = self .pr_events .iter() - .filter(|entry| entry.value().expires_at <= now) + .filter(|entry| { + let value = entry.value(); + value.expires_at <= now + && !value + .event + .as_ref() + .is_some_and(|event| self.event_sync_in_progress(event)) + }) .map(|entry| { let pr_entry = entry.value(); let event_id_str = entry.key().clone(); @@ -2146,6 +2188,80 @@ fn test_cleanup_preserves_non_expired_entries() { assert_eq!(pr_count, 1); } +#[test] +fn cleanup_defers_expiry_only_while_repository_sync_is_active() { + let purgatory = Purgatory::new(PathBuf::new()); + let keys = Keys::generate(); + let identifier = "large-repository"; + + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .finalize(&keys) + .unwrap(); + purgatory.add_announcement( + announcement, + identifier.to_string(), + keys.public_key(), + PathBuf::from("/path/that/does/not/exist"), + HashSet::new(), + ); + + let state = EventBuilder::new(Kind::RepoState, "") + .finalize(&keys) + .unwrap(); + purgatory.add_state(state, identifier.to_string(), keys.public_key(), true); + + let pr = EventBuilder::new(Kind::GitPatch, "") + .tag(Tag::custom( + "a", + vec![format!("30617:{}:{identifier}", keys.public_key().to_hex())], + )) + .finalize(&keys) + .unwrap(); + purgatory.add_pr(pr, "pr-event".to_string(), "commit".to_string(), true); + + let expired_at = Instant::now() - Duration::from_secs(1); + purgatory + .announcement_purgatory + .get_mut(&(keys.public_key(), identifier.to_string())) + .unwrap() + .expires_at = expired_at; + purgatory.state_events.get_mut(identifier).unwrap()[0].expires_at = expired_at; + purgatory.pr_events.get_mut("pr-event").unwrap().expires_at = expired_at; + purgatory + .sync_queue + .get_mut(identifier) + .unwrap() + .in_progress = true; + + let removed = purgatory.cleanup(); + assert_eq!(removed, (0, 0, 0)); + assert!( + !purgatory + .find_announcement(&keys.public_key(), identifier) + .unwrap() + .soft_expired + ); + assert_eq!(purgatory.count(), (1, 1, 1)); + assert_eq!(purgatory.expired_count(), 0); + + // An entry merely waiting in the queue or in backoff is not active work. + // Once the concrete Git operation ends, an unsatisfied event may expire. + purgatory + .sync_queue + .get_mut(identifier) + .unwrap() + .in_progress = false; + let removed = purgatory.cleanup(); + assert_eq!(removed, (0, 1, 1)); + assert!( + purgatory + .find_announcement(&keys.public_key(), identifier) + .unwrap() + .soft_expired + ); + assert_eq!(purgatory.expired_count(), 2); +} + #[test] fn test_cleanup_mixed_expired_and_fresh() { use std::time::Duration;