mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
fix(startup): serve before retention catch-up
A production restart spent minutes scanning more than 52,000 retained deletion requests before the SyncManager and its purgatory timer could start. The relay accepted a fresh invitation announcement during that window, but no background worker existed yet to recover it. Keep policy reconciliation in the blocking startup path, then run request and holding retention catch-up immediately from the existing maintenance task after the sync machinery is live. The regression proves cleanup is deferred without losing its startup catch-up behavior.
This commit is contained in:
@@ -63,14 +63,6 @@ impl DeletionRuntime {
|
||||
/// Run deletion-owned startup tasks before the relay begins serving traffic.
|
||||
pub async fn run_startup_tasks(&self) -> Result<()> {
|
||||
log_startup_reconciliation(self.service.run_startup_reconciliation().await?);
|
||||
// Enumeration failure is fatal before serving traffic; individual record
|
||||
// failures are conservatively retained and retried by the timer.
|
||||
let request_stats = self
|
||||
.service
|
||||
.cleanup_expired_requests(Timestamp::now())
|
||||
.await?;
|
||||
log_request_cleanup("startup catch-up", request_stats);
|
||||
self.run_holding_startup_cleanup().await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -84,16 +76,22 @@ impl DeletionRuntime {
|
||||
let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
|
||||
|
||||
let handle = tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval_at(
|
||||
tokio::time::Instant::now() + interval_duration,
|
||||
interval_duration,
|
||||
);
|
||||
let mut first_pass = true;
|
||||
let mut interval =
|
||||
tokio::time::interval_at(tokio::time::Instant::now(), interval_duration);
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = interval.tick() => {
|
||||
// The first pass is the startup catch-up. It deliberately
|
||||
// runs here, after the relay and SyncManager are live:
|
||||
// large retained request sets can take minutes to scan,
|
||||
// but must not prevent fresh purgatory work from starting.
|
||||
match service.cleanup_expired_requests(Timestamp::now()).await {
|
||||
Ok(stats) => log_request_cleanup("periodic pass", stats),
|
||||
Ok(stats) => log_request_cleanup(
|
||||
if first_pass { "startup catch-up" } else { "periodic pass" },
|
||||
stats,
|
||||
),
|
||||
Err(error) => tracing::warn!(error = %error, "Deletion-request cleanup periodic pass failed"),
|
||||
}
|
||||
match holding.cleanup_expired_with_lifecycle(&lifecycle, Timestamp::now(), retention).await {
|
||||
@@ -112,6 +110,7 @@ impl DeletionRuntime {
|
||||
tracing::warn!(error = %e, "Holding cleanup periodic pass failed");
|
||||
}
|
||||
}
|
||||
first_pass = false;
|
||||
}
|
||||
changed = shutdown_rx.changed() => {
|
||||
if changed.is_ok() && *shutdown_rx.borrow() {
|
||||
@@ -134,33 +133,6 @@ impl DeletionRuntime {
|
||||
handle,
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_holding_startup_cleanup(&self) {
|
||||
match self
|
||||
.holding
|
||||
.cleanup_expired_with_lifecycle(
|
||||
&self.lifecycle,
|
||||
Timestamp::now(),
|
||||
self.holding_retention,
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(stats) => {
|
||||
if stats.expired_records > 0 {
|
||||
tracing::info!(
|
||||
examined = stats.metadata_examined,
|
||||
expired = stats.expired_records,
|
||||
metadata_deleted = stats.metadata_deleted,
|
||||
archived_events_deleted = stats.archived_events_deleted,
|
||||
"Holding cleanup startup catch-up completed"
|
||||
);
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "Holding cleanup startup catch-up failed");
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct DeletionCleanupTask {
|
||||
@@ -284,3 +256,101 @@ fn log_whitelist_restore(stats: WhitelistRestoreStats) {
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use nostr_memory::MemoryDatabase;
|
||||
use nostr_relay_builder::prelude::{
|
||||
EventBuilder, EventId, FinalizeEvent, Keys, Kind, NostrDatabase, Tag,
|
||||
};
|
||||
|
||||
use super::*;
|
||||
use crate::grasp06::receive::new_repo_init_locks;
|
||||
use crate::nostr::lifecycle::{
|
||||
DeletionContext, HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle,
|
||||
RequestClassification, Tombstones,
|
||||
};
|
||||
use crate::purgatory::Purgatory;
|
||||
|
||||
#[tokio::test]
|
||||
async fn retained_request_cleanup_starts_after_blocking_startup_finishes() {
|
||||
let database = Arc::new(MemoryDatabase::unbounded());
|
||||
let tombstones = Tombstones::in_memory();
|
||||
let holding = HoldingStore::in_memory();
|
||||
let lifecycle = RepositoryLifecycle::in_memory();
|
||||
let config = crate::config::Config {
|
||||
deletion_request_retention_unused_served_secs: 1,
|
||||
deletion_request_retention_unused_unserved_gating_additional_secs: 1,
|
||||
..crate::config::Config::for_testing()
|
||||
};
|
||||
let service = DeletionService::new(DeletionContext::new(
|
||||
"test.example.com",
|
||||
database.clone(),
|
||||
tombstones.clone(),
|
||||
holding.clone(),
|
||||
lifecycle.clone(),
|
||||
ReplaceableHistoryStore::in_memory(),
|
||||
PathBuf::new(),
|
||||
Arc::new(Purgatory::new(PathBuf::new())),
|
||||
config,
|
||||
new_repo_init_locks(),
|
||||
));
|
||||
let runtime = DeletionRuntime::new(
|
||||
service,
|
||||
holding,
|
||||
Duration::from_secs(60),
|
||||
Duration::from_secs(60),
|
||||
);
|
||||
let request = EventBuilder::new(Kind::EventDeletion, "")
|
||||
.tags(vec![Tag::event(EventId::all_zeros())])
|
||||
.finalize(&Keys::generate())
|
||||
.expect("deletion request should sign");
|
||||
|
||||
tombstones
|
||||
.record_request(
|
||||
&request,
|
||||
Timestamp::from_secs(1),
|
||||
RequestClassification::LocallyActionable,
|
||||
)
|
||||
.await
|
||||
.expect("request lifecycle should be retained");
|
||||
database
|
||||
.save_event(&request)
|
||||
.await
|
||||
.expect("request payload should be served");
|
||||
|
||||
runtime
|
||||
.run_startup_tasks()
|
||||
.await
|
||||
.expect("blocking startup reconciliation should complete");
|
||||
assert!(
|
||||
tombstones
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.expect("lifecycle lookup should succeed")
|
||||
.is_some(),
|
||||
"retention cleanup must not delay relay and SyncManager startup"
|
||||
);
|
||||
|
||||
let cleanup = runtime.spawn_cleanup_task();
|
||||
tokio::time::timeout(Duration::from_secs(2), async {
|
||||
loop {
|
||||
if tombstones
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.expect("lifecycle lookup should succeed")
|
||||
.is_none()
|
||||
{
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.expect("background startup catch-up should run immediately");
|
||||
cleanup.shutdown().await;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user