mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat: clean up deletion request lifecycle
This commit is contained in:
@@ -12,9 +12,10 @@ and purgatory transitions.
|
||||
|
||||
> **Status:** Lifecycle metadata, NIP-09 and NIP-62 lifecycle admission
|
||||
> (including disrespector read-only would-have-deleted classification),
|
||||
> deterministic admission-use attribution, and measured destructive outcomes are
|
||||
> implemented. Request cleanup/expiry, migration, target-set deduplication
|
||||
> removal, and final policy-mode reconciliation remain future stages.
|
||||
> deterministic admission-use attribution, measured destructive outcomes,
|
||||
> startup reconciliation, and periodic request cleanup/permanent expiry are
|
||||
> implemented. Target-set deduplication removal, repair-command retirement, and
|
||||
> final rollout coverage remain future work.
|
||||
|
||||
### Production motivation
|
||||
|
||||
@@ -271,9 +272,16 @@ mode reapplies normal gating and the 9-month served plus 3-month unserved expiry
|
||||
schedule using retained lifecycle timestamps. Startup reconciliation must apply
|
||||
these transitions before the relay begins serving traffic.
|
||||
|
||||
Periodic request cleanup, permanent tombstone expiry, removal of target-set
|
||||
deduplication and the `repair-deletion-requests` command, and final rollout
|
||||
coverage remain future work.
|
||||
Startup catch-up and periodic cleanup use the holding-cleanup cadence. They
|
||||
apply the same timestamp-derived half-open lifecycle boundaries, remove expired
|
||||
served copies from Main before Tombstones, and permanently remove payload plus
|
||||
lifecycle metadata only after the additional gating period. A shared lifecycle
|
||||
transition lock prevents cleanup from acting on a stale record while admission
|
||||
is promoting and marking it used.
|
||||
|
||||
Removal of target-set deduplication, retirement of the
|
||||
`repair-deletion-requests` command, and final rollout coverage remain future
|
||||
work.
|
||||
|
||||
## Core consistency invariant
|
||||
|
||||
|
||||
@@ -113,6 +113,30 @@ lazy_static! {
|
||||
.expect("register holding cleanup last-run metric");
|
||||
metric
|
||||
};
|
||||
static ref DELETION_REQUEST_CLEANUP_RUNS_TOTAL: Counter = {
|
||||
let metric = Counter::with_opts(Opts::new(
|
||||
"ngit_deletion_request_cleanup_runs_total",
|
||||
"Number of deletion-request cleanup passes run",
|
||||
)).expect("build deletion request cleanup runs metric");
|
||||
REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup runs metric");
|
||||
metric
|
||||
};
|
||||
static ref DELETION_REQUEST_CLEANUP_REMOVED_TOTAL: CounterVec = {
|
||||
let metric = CounterVec::new(Opts::new(
|
||||
"ngit_deletion_request_cleanup_removed_total",
|
||||
"Deletion-request cleanup removals by storage type",
|
||||
), &["type"]).expect("build deletion request cleanup removed metric");
|
||||
REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup removed metric");
|
||||
metric
|
||||
};
|
||||
static ref DELETION_REQUEST_CLEANUP_OUTCOMES_TOTAL: CounterVec = {
|
||||
let metric = CounterVec::new(Opts::new(
|
||||
"ngit_deletion_request_cleanup_outcomes_total",
|
||||
"Deletion-request cleanup failures and concurrent skips",
|
||||
), &["outcome"]).expect("build deletion request cleanup outcomes metric");
|
||||
REGISTRY.register(Box::new(metric.clone())).expect("register deletion request cleanup outcomes metric");
|
||||
metric
|
||||
};
|
||||
static ref RECOVERY_TOTAL: CounterVec = {
|
||||
let metric = CounterVec::new(
|
||||
Opts::new(
|
||||
@@ -247,6 +271,30 @@ pub fn record_holding_cleanup_run(
|
||||
.set(archive_files_deleted as f64);
|
||||
}
|
||||
|
||||
pub fn record_deletion_request_cleanup_run(
|
||||
main_removed: usize,
|
||||
tombstone_removed: usize,
|
||||
metadata_removed: usize,
|
||||
failures: usize,
|
||||
skipped: usize,
|
||||
) {
|
||||
DELETION_REQUEST_CLEANUP_RUNS_TOTAL.inc();
|
||||
for (kind, value) in [
|
||||
("main", main_removed),
|
||||
("tombstone", tombstone_removed),
|
||||
("metadata", metadata_removed),
|
||||
] {
|
||||
DELETION_REQUEST_CLEANUP_REMOVED_TOTAL
|
||||
.with_label_values(&[kind])
|
||||
.inc_by(value as f64);
|
||||
}
|
||||
for (outcome, value) in [("failure", failures), ("stale_or_concurrent_skip", skipped)] {
|
||||
DELETION_REQUEST_CLEANUP_OUTCOMES_TOTAL
|
||||
.with_label_values(&[outcome])
|
||||
.inc_by(value as f64);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn record_recovery_attempt() {
|
||||
RECOVERY_TOTAL.with_label_values(&["attempted"]).inc();
|
||||
}
|
||||
|
||||
@@ -0,0 +1,457 @@
|
||||
use anyhow::Result;
|
||||
use nostr_relay_builder::prelude::{Filter, Kind, Timestamp};
|
||||
|
||||
use crate::nostr::lifecycle::RequestLifecycleRecord;
|
||||
|
||||
use super::DeletionService;
|
||||
|
||||
/// Result counters for one bounded deletion-request retention pass.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct RequestCleanupStats {
|
||||
pub canonical_records_examined: usize,
|
||||
pub main_payloads_removed: usize,
|
||||
pub tombstone_payloads_removed: usize,
|
||||
pub metadata_rows_removed: usize,
|
||||
pub indefinitely_retained_disrespector_requests: usize,
|
||||
pub stale_or_concurrent_records_skipped: usize,
|
||||
pub failures: usize,
|
||||
}
|
||||
|
||||
impl DeletionService {
|
||||
/// Reconcile served copies and permanently expire elapsed request lifecycles.
|
||||
/// The supplied time makes every retention boundary deterministic in tests.
|
||||
pub(crate) async fn cleanup_expired_requests(
|
||||
&self,
|
||||
now: Timestamp,
|
||||
) -> Result<RequestCleanupStats> {
|
||||
let records = self.ctx.tombstones().lifecycle_records_result().await?;
|
||||
let mut stats = RequestCleanupStats {
|
||||
canonical_records_examined: records.len(),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
for selected in records {
|
||||
if let Err(error) = self.cleanup_one_request(&selected, now, &mut stats).await {
|
||||
stats.failures += 1;
|
||||
tracing::warn!(request_id = %selected.request.id, error = %error, "Deletion-request cleanup failed for record");
|
||||
}
|
||||
}
|
||||
|
||||
match self.ctx.tombstones().remove_orphan_metadata().await {
|
||||
Ok(removed) => stats.metadata_rows_removed += removed,
|
||||
Err(error) => {
|
||||
stats.failures += 1;
|
||||
tracing::warn!(error = %error, "Deletion-request cleanup failed to remove orphan metadata");
|
||||
}
|
||||
}
|
||||
Ok(stats)
|
||||
}
|
||||
|
||||
async fn cleanup_one_request(
|
||||
&self,
|
||||
selected: &RequestLifecycleRecord,
|
||||
now: Timestamp,
|
||||
stats: &mut RequestCleanupStats,
|
||||
) -> Result<()> {
|
||||
let tombstones = self.ctx.tombstones();
|
||||
let _guard = tombstones.lock_lifecycle().await;
|
||||
let Some(current) = tombstones
|
||||
.lifecycle_for_request_result(&selected.request.id)
|
||||
.await?
|
||||
else {
|
||||
stats.stale_or_concurrent_records_skipped += 1;
|
||||
return Ok(());
|
||||
};
|
||||
// A mark-used/reclassification replacement between enumeration and this
|
||||
// lock acquisition means this pass must not act on its stale decision.
|
||||
if current != *selected {
|
||||
stats.stale_or_concurrent_records_skipped += 1;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let targeting = current.request.kind == Kind::EventDeletion
|
||||
|| self.vanish_targets_this_relay(¤t.request);
|
||||
if self.ctx.config.deletion_request_disrespector
|
||||
&& targeting
|
||||
&& current.last_used_at.is_some()
|
||||
{
|
||||
stats.indefinitely_retained_disrespector_requests += 1;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if self.request_is_served(¤t, now)? {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
// Main must be removed before Tombstones. Holding the same lifecycle
|
||||
// lock as gate attribution prevents stale cleanup from unserving a
|
||||
// request that an admission just promoted and marked used.
|
||||
if self
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(¤t.request.id)
|
||||
.await?
|
||||
.is_some()
|
||||
{
|
||||
self.ctx
|
||||
.database()
|
||||
.delete(Filter::new().ids(vec![current.request.id]))
|
||||
.await?;
|
||||
stats.main_payloads_removed += 1;
|
||||
}
|
||||
|
||||
if !self.request_is_expired(¤t, now)? {
|
||||
return Ok(());
|
||||
}
|
||||
let removed = tombstones
|
||||
.permanently_delete_request_locked(¤t.request.id)
|
||||
.await?;
|
||||
stats.tombstone_payloads_removed += removed.payloads_deleted;
|
||||
stats.metadata_rows_removed += removed.metadata_deleted;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) fn request_is_served(
|
||||
&self,
|
||||
record: &RequestLifecycleRecord,
|
||||
now: Timestamp,
|
||||
) -> Result<bool> {
|
||||
let (anchor, duration) = match record.last_used_at {
|
||||
Some(used) => (
|
||||
used,
|
||||
self.ctx
|
||||
.config
|
||||
.deletion_request_retention_used_served_after_last_used(),
|
||||
),
|
||||
None => (
|
||||
record.first_seen_at,
|
||||
self.ctx.config.deletion_request_retention_unused_served(),
|
||||
),
|
||||
};
|
||||
let deadline = anchor
|
||||
.as_secs()
|
||||
.checked_add(duration.as_secs())
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!("deletion-request served retention deadline overflow")
|
||||
})?;
|
||||
Ok(now.as_secs() < deadline)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use nostr_relay_builder::prelude::{Event, EventBuilder, EventId, FinalizeEvent, Keys, Tag};
|
||||
|
||||
use super::*;
|
||||
use crate::grasp06::receive::new_repo_init_locks;
|
||||
use crate::nostr::lifecycle::{
|
||||
HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle, RequestClassification,
|
||||
Tombstones,
|
||||
};
|
||||
use crate::purgatory::Purgatory;
|
||||
|
||||
fn service(disrespector: bool) -> DeletionService {
|
||||
let db = Arc::new(nostr_memory::MemoryDatabase::unbounded());
|
||||
let config = crate::config::Config {
|
||||
deletion_request_disrespector: disrespector,
|
||||
deletion_request_retention_unused_served_secs: 10,
|
||||
deletion_request_retention_unused_unserved_gating_additional_secs: 5,
|
||||
deletion_request_retention_used_served_after_last_used_secs: 20,
|
||||
deletion_request_retention_used_unserved_gating_additional_secs: 5,
|
||||
..crate::config::Config::for_testing()
|
||||
};
|
||||
DeletionService::new(super::super::DeletionContext::new(
|
||||
"test.example.com",
|
||||
db,
|
||||
Tombstones::in_memory(),
|
||||
HoldingStore::in_memory(),
|
||||
RepositoryLifecycle::in_memory(),
|
||||
ReplaceableHistoryStore::in_memory(),
|
||||
PathBuf::new(),
|
||||
Arc::new(Purgatory::new(PathBuf::new())),
|
||||
config,
|
||||
new_repo_init_locks(),
|
||||
))
|
||||
}
|
||||
|
||||
fn deletion() -> Event {
|
||||
EventBuilder::new(Kind::EventDeletion, "")
|
||||
.tags(vec![Tag::event(EventId::all_zeros())])
|
||||
.finalize(&Keys::generate())
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn non_targeting_vanish() -> Event {
|
||||
EventBuilder::new(Kind::RequestToVanish, "")
|
||||
.tags(vec![Tag::custom(
|
||||
"relay",
|
||||
vec!["wss://elsewhere.example".to_owned()],
|
||||
)])
|
||||
.finalize(&Keys::generate())
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn retained(service: &DeletionService, request: &Event, first_seen: u64) {
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.record_request(
|
||||
request,
|
||||
Timestamp::from_secs(first_seen),
|
||||
RequestClassification::LocallyActionable,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
service.ctx.database().save_event(request).await.unwrap();
|
||||
}
|
||||
|
||||
async fn main_has(service: &DeletionService, request: &Event) -> bool {
|
||||
service
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn unused_boundaries_keep_gate_then_permanently_expire() {
|
||||
let service = service(false);
|
||||
let request = deletion();
|
||||
retained(&service, &request, 0).await;
|
||||
assert_eq!(
|
||||
service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(9))
|
||||
.await
|
||||
.unwrap()
|
||||
.main_payloads_removed,
|
||||
0
|
||||
);
|
||||
assert!(main_has(&service, &request).await);
|
||||
|
||||
let stats = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(10))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.main_payloads_removed, 1);
|
||||
assert!(!main_has(&service, &request).await);
|
||||
assert!(service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
assert_eq!(
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.event_deletion_candidates(&EventId::all_zeros(), &request.pubkey)
|
||||
.await
|
||||
.unwrap()
|
||||
.len(),
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_records_result()
|
||||
.await
|
||||
.unwrap()
|
||||
.len(),
|
||||
1
|
||||
);
|
||||
|
||||
let stats = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(15))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.tombstone_payloads_removed, 1);
|
||||
assert!(service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn used_boundaries_and_reuse_reset_the_lifecycle() {
|
||||
let service = service(false);
|
||||
let request = deletion();
|
||||
retained(&service, &request, 0).await;
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&request.id, Timestamp::from_secs(20))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(39))
|
||||
.await
|
||||
.unwrap()
|
||||
.main_payloads_removed
|
||||
== 0
|
||||
);
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&request.id, Timestamp::from_secs(39))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(
|
||||
service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(40))
|
||||
.await
|
||||
.unwrap()
|
||||
.main_payloads_removed
|
||||
== 0
|
||||
);
|
||||
let stats = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(59))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.main_payloads_removed, 1);
|
||||
assert!(service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
let stats = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(64))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.tombstone_payloads_removed, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn disrespector_only_retains_used_targeting_requests_indefinitely() {
|
||||
let service = service(true);
|
||||
let used = deletion();
|
||||
retained(&service, &used, 0).await;
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&used.id, Timestamp::from_secs(1))
|
||||
.await
|
||||
.unwrap();
|
||||
let stats = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(10_000))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.indefinitely_retained_disrespector_requests, 1);
|
||||
assert!(main_has(&service, &used).await);
|
||||
|
||||
let unused = deletion();
|
||||
retained(&service, &unused, 0).await;
|
||||
service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(15))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&unused.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn non_targeting_nip62_expires_and_cleanup_is_idempotent() {
|
||||
let service = service(false);
|
||||
let request = non_targeting_vanish();
|
||||
retained(&service, &request, 0).await;
|
||||
assert!(
|
||||
service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(15))
|
||||
.await
|
||||
.unwrap()
|
||||
.tombstone_payloads_removed
|
||||
== 1
|
||||
);
|
||||
let rerun = service
|
||||
.cleanup_expired_requests(Timestamp::from_secs(15))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
rerun.main_payloads_removed
|
||||
+ rerun.tombstone_payloads_removed
|
||||
+ rerun.metadata_rows_removed,
|
||||
0
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn stale_cleanup_snapshot_cannot_unserve_a_marked_used_request() {
|
||||
let service = service(false);
|
||||
let request = deletion();
|
||||
retained(&service, &request, 0).await;
|
||||
let selected = service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
// This models admission's metadata update and Main promotion happening
|
||||
// after enumeration but before cleanup acquires the shared lock.
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&request.id, Timestamp::from_secs(14))
|
||||
.await
|
||||
.unwrap();
|
||||
let mut stats = RequestCleanupStats::default();
|
||||
service
|
||||
.cleanup_one_request(&selected, Timestamp::from_secs(15), &mut stats)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.stale_or_concurrent_records_skipped, 1);
|
||||
assert!(main_has(&service, &request).await);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn startup_reconciliation_and_periodic_cleanup_share_serving_boundary() {
|
||||
let startup = service(false);
|
||||
let periodic = service(false);
|
||||
let startup_request = deletion();
|
||||
let periodic_request = deletion();
|
||||
retained(&startup, &startup_request, 0).await;
|
||||
retained(&periodic, &periodic_request, 0).await;
|
||||
|
||||
startup
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(10))
|
||||
.await
|
||||
.unwrap();
|
||||
periodic
|
||||
.cleanup_expired_requests(Timestamp::from_secs(10))
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(!main_has(&startup, &startup_request).await);
|
||||
assert!(!main_has(&periodic, &periodic_request).await);
|
||||
assert!(startup
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&startup_request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
assert!(periodic
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&periodic_request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,6 @@
|
||||
mod archival;
|
||||
mod cascade;
|
||||
mod cleanup;
|
||||
mod context;
|
||||
mod policy;
|
||||
mod pr_refs;
|
||||
@@ -42,6 +43,7 @@ impl DeletionOutcome {
|
||||
}
|
||||
}
|
||||
|
||||
pub use cleanup::RequestCleanupStats;
|
||||
pub use context::DeletionContext;
|
||||
pub use policy::DeletionPolicy;
|
||||
pub use runtime::{run_holding_eject, DeletionCleanupTask, DeletionRuntime, HoldingEjectArgs};
|
||||
|
||||
@@ -334,16 +334,15 @@ impl DeletionPolicy {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
Kind::RepoState => {
|
||||
Kind::RepoState
|
||||
if self
|
||||
.ctx
|
||||
.purgatory
|
||||
.find_state(&identifier)
|
||||
.into_iter()
|
||||
.any(|entry| entry.author == owner && covered(entry.event.created_at))
|
||||
{
|
||||
return true;
|
||||
}
|
||||
.any(|entry| entry.author == owner && covered(entry.event.created_at)) =>
|
||||
{
|
||||
return true;
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
|
||||
@@ -14,7 +14,7 @@ use super::startup::{
|
||||
BlacklistParityStats, BlacklistRestoreStats, StartupReconciliationStats, WhitelistParityStats,
|
||||
WhitelistRestoreStats,
|
||||
};
|
||||
use super::DeletionService;
|
||||
use super::{DeletionService, RequestCleanupStats};
|
||||
|
||||
#[derive(Debug, Args)]
|
||||
pub struct HoldingEjectArgs {
|
||||
@@ -63,6 +63,13 @@ 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(())
|
||||
}
|
||||
@@ -71,6 +78,7 @@ impl DeletionRuntime {
|
||||
pub fn spawn_cleanup_task(&self) -> DeletionCleanupTask {
|
||||
let holding = self.holding.clone();
|
||||
let lifecycle = self.lifecycle.clone();
|
||||
let service = self.service.clone();
|
||||
let retention = self.holding_retention;
|
||||
let interval_duration = self.holding_cleanup_interval;
|
||||
let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
|
||||
@@ -84,6 +92,10 @@ impl DeletionRuntime {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = interval.tick() => {
|
||||
match service.cleanup_expired_requests(Timestamp::now()).await {
|
||||
Ok(stats) => log_request_cleanup("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 {
|
||||
Ok(stats) => {
|
||||
if stats.expired_records > 0 {
|
||||
@@ -103,7 +115,7 @@ impl DeletionRuntime {
|
||||
}
|
||||
changed = shutdown_rx.changed() => {
|
||||
if changed.is_ok() && *shutdown_rx.borrow() {
|
||||
tracing::info!("Holding cleanup task received shutdown signal");
|
||||
tracing::info!("Deletion lifecycle maintenance task received shutdown signal");
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -114,7 +126,7 @@ impl DeletionRuntime {
|
||||
tracing::info!(
|
||||
retention_secs = retention.as_secs(),
|
||||
interval_secs = interval_duration.as_secs(),
|
||||
"Holding cleanup task started"
|
||||
"Deletion lifecycle maintenance task started"
|
||||
);
|
||||
|
||||
DeletionCleanupTask {
|
||||
@@ -160,7 +172,7 @@ impl DeletionCleanupTask {
|
||||
pub async fn shutdown(self) {
|
||||
let _ = self.shutdown_tx.send(true);
|
||||
if let Err(e) = self.handle.await {
|
||||
tracing::warn!(error = %e, "Holding cleanup task join failed during shutdown");
|
||||
tracing::warn!(error = %e, "Deletion lifecycle maintenance task join failed during shutdown");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -200,6 +212,27 @@ fn log_startup_reconciliation(stats: StartupReconciliationStats) {
|
||||
log_whitelist_restore(stats.whitelist_restore);
|
||||
}
|
||||
|
||||
fn log_request_cleanup(phase: &str, stats: RequestCleanupStats) {
|
||||
crate::metrics::record_deletion_request_cleanup_run(
|
||||
stats.main_payloads_removed,
|
||||
stats.tombstone_payloads_removed,
|
||||
stats.metadata_rows_removed,
|
||||
stats.failures,
|
||||
stats.stale_or_concurrent_records_skipped,
|
||||
);
|
||||
tracing::info!(
|
||||
phase,
|
||||
examined = stats.canonical_records_examined,
|
||||
main_removed = stats.main_payloads_removed,
|
||||
tombstone_removed = stats.tombstone_payloads_removed,
|
||||
metadata_removed = stats.metadata_rows_removed,
|
||||
indefinitely_retained = stats.indefinitely_retained_disrespector_requests,
|
||||
skipped = stats.stale_or_concurrent_records_skipped,
|
||||
failures = stats.failures,
|
||||
"Deletion-request cleanup completed"
|
||||
);
|
||||
}
|
||||
|
||||
fn log_blacklist_parity(stats: BlacklistParityStats) {
|
||||
if stats.scanned_announcements > 0 || stats.matched_announcements > 0 {
|
||||
tracing::info!(
|
||||
|
||||
@@ -93,13 +93,44 @@ impl DeletionService {
|
||||
.into_values()
|
||||
.min_by(|(left, _), (right, _)| Self::winner_order(left, right))?;
|
||||
|
||||
// Serialize promotion plus last-used attribution with request cleanup.
|
||||
// Cleanup always re-reads metadata under this same lock, so it cannot
|
||||
// remove the Main copy selected by an admission using an old snapshot.
|
||||
let _lifecycle_guard = tombstones.lock_lifecycle().await;
|
||||
let Some(winner) = (match tombstones
|
||||
.lifecycle_for_request_result(&winner.request.id)
|
||||
.await
|
||||
{
|
||||
Ok(record) => record,
|
||||
Err(error) => {
|
||||
return Some(Self::gate_error(
|
||||
"re-reading deletion-request lifecycle",
|
||||
error,
|
||||
))
|
||||
}
|
||||
}) else {
|
||||
return Some(Self::gate_error(
|
||||
"re-reading deletion-request lifecycle",
|
||||
anyhow::anyhow!("request {} disappeared", winner.request.id),
|
||||
));
|
||||
};
|
||||
match self.request_is_expired(&winner, Timestamp::now()) {
|
||||
Ok(true) => return None,
|
||||
Ok(false) => {}
|
||||
Err(error) => {
|
||||
return Some(Self::gate_error(
|
||||
"evaluating deletion-request retention",
|
||||
error,
|
||||
))
|
||||
}
|
||||
}
|
||||
let previous_last_used_at = winner.last_used_at;
|
||||
if let Err(error) = self.ctx.database().save_event(&winner.request).await {
|
||||
return Some(Self::gate_error("promoting deletion request", error));
|
||||
}
|
||||
let used_at = Timestamp::now();
|
||||
let updated = match tombstones
|
||||
.mark_request_used(&winner.request.id, used_at)
|
||||
.mark_request_used_locked(&winner.request.id, used_at)
|
||||
.await
|
||||
{
|
||||
Ok(Some(record)) if record.last_used_at >= Some(used_at) => record,
|
||||
@@ -149,7 +180,11 @@ impl DeletionService {
|
||||
})
|
||||
}
|
||||
|
||||
fn request_is_expired(&self, record: &RequestLifecycleRecord, now: Timestamp) -> Result<bool> {
|
||||
pub(crate) fn request_is_expired(
|
||||
&self,
|
||||
record: &RequestLifecycleRecord,
|
||||
now: Timestamp,
|
||||
) -> Result<bool> {
|
||||
let (anchor, served, additional) = match record.last_used_at {
|
||||
Some(last_used_at) => (
|
||||
last_used_at,
|
||||
|
||||
@@ -5,7 +5,7 @@ use anyhow::Result;
|
||||
use nostr_relay_builder::prelude::{Event, Filter, Kind, PublicKey, Timestamp};
|
||||
|
||||
use crate::nostr::events::RepositoryAnnouncement;
|
||||
use crate::nostr::lifecycle::{DeletionSource, RequestClassification, RequestLifecycleRecord};
|
||||
use crate::nostr::lifecycle::{DeletionSource, RequestClassification};
|
||||
|
||||
use super::DeletionService;
|
||||
|
||||
@@ -248,28 +248,6 @@ impl DeletionService {
|
||||
}
|
||||
}
|
||||
|
||||
fn request_is_served(&self, record: &RequestLifecycleRecord, now: Timestamp) -> Result<bool> {
|
||||
let (anchor, duration) = match record.last_used_at {
|
||||
Some(used) => (
|
||||
used,
|
||||
self.ctx
|
||||
.config
|
||||
.deletion_request_retention_used_served_after_last_used(),
|
||||
),
|
||||
None => (
|
||||
record.first_seen_at,
|
||||
self.ctx.config.deletion_request_retention_unused_served(),
|
||||
),
|
||||
};
|
||||
let deadline = anchor
|
||||
.as_secs()
|
||||
.checked_add(duration.as_secs())
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!("deletion-request served retention deadline overflow")
|
||||
})?;
|
||||
Ok(now.as_secs() < deadline)
|
||||
}
|
||||
|
||||
async fn promote_request(
|
||||
&self,
|
||||
request: &Event,
|
||||
|
||||
@@ -45,6 +45,8 @@ use std::collections::HashSet;
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::{Mutex, MutexGuard};
|
||||
|
||||
use nostr_lmdb::NostrLmdb;
|
||||
use nostr_memory::MemoryDatabase;
|
||||
use nostr_relay_builder::prelude::{
|
||||
@@ -101,6 +103,13 @@ pub struct RequestLifecycleRecord {
|
||||
pub classification: RequestClassification,
|
||||
}
|
||||
|
||||
/// Counts returned after permanently removing one request lifecycle.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct RequestLifecycleDeletionStats {
|
||||
pub payloads_deleted: usize,
|
||||
pub metadata_deleted: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
struct DeletionTargets {
|
||||
e_ids: HashSet<EventId>,
|
||||
@@ -141,6 +150,11 @@ impl DeletionTargets {
|
||||
pub struct Tombstones {
|
||||
db: Arc<dyn NostrDatabase>,
|
||||
metadata_signer: Keys,
|
||||
/// Serializes lifecycle metadata changes with Main-database transitions
|
||||
/// performed by `DeletionService`. This is intentionally one simple lock:
|
||||
/// request lifecycles are a small control plane and correctness matters more
|
||||
/// than parallel metadata writes.
|
||||
lifecycle_lock: Arc<Mutex<()>>,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for Tombstones {
|
||||
@@ -179,6 +193,7 @@ impl Tombstones {
|
||||
Ok(Self {
|
||||
db: Arc::new(db),
|
||||
metadata_signer: Keys::generate(),
|
||||
lifecycle_lock: Arc::new(Mutex::new(())),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -193,9 +208,17 @@ impl Tombstones {
|
||||
.build(),
|
||||
),
|
||||
metadata_signer: Keys::generate(),
|
||||
lifecycle_lock: Arc::new(Mutex::new(())),
|
||||
}
|
||||
}
|
||||
|
||||
/// Acquire the lifecycle transition lock. Callers that also transition the
|
||||
/// Main database must hold this across the Main operation and the matching
|
||||
/// metadata update/removal.
|
||||
pub(crate) async fn lock_lifecycle(&self) -> MutexGuard<'_, ()> {
|
||||
self.lifecycle_lock.lock().await
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn save_request_payload_without_metadata(
|
||||
&self,
|
||||
@@ -215,6 +238,17 @@ impl Tombstones {
|
||||
event: &Event,
|
||||
first_seen_at: Timestamp,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<RequestLifecycleRecord> {
|
||||
let _guard = self.lock_lifecycle().await;
|
||||
self.record_request_locked(event, first_seen_at, classification)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn record_request_locked(
|
||||
&self,
|
||||
event: &Event,
|
||||
first_seen_at: Timestamp,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<RequestLifecycleRecord> {
|
||||
if !Self::is_request_kind(event.kind) {
|
||||
anyhow::bail!("Tombstone lifecycle accepts only kind-5 or kind-62 requests");
|
||||
@@ -339,6 +373,15 @@ impl Tombstones {
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
used_at: Timestamp,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
let _guard = self.lock_lifecycle().await;
|
||||
self.mark_request_used_locked(request_id, used_at).await
|
||||
}
|
||||
|
||||
pub(crate) async fn mark_request_used_locked(
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
used_at: Timestamp,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
let Some(existing) = self.lifecycle_for_request_result(request_id).await? else {
|
||||
return Ok(None);
|
||||
@@ -364,6 +407,16 @@ impl Tombstones {
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
let _guard = self.lock_lifecycle().await;
|
||||
self.reclassify_request_locked(request_id, classification)
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn reclassify_request_locked(
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
let Some(existing) = self.lifecycle_for_request_result(request_id).await? else {
|
||||
return Ok(None);
|
||||
@@ -402,6 +455,86 @@ impl Tombstones {
|
||||
records
|
||||
}
|
||||
|
||||
/// Fallible enumeration for destructive maintenance. Unlike
|
||||
/// `lifecycle_records`, this never turns a database failure into an
|
||||
/// apparently empty result. Private metadata is not returned.
|
||||
pub async fn lifecycle_records_result(&self) -> anyhow::Result<Vec<RequestLifecycleRecord>> {
|
||||
let requests = self.request_payloads().await?;
|
||||
let mut records = Vec::new();
|
||||
for request in requests {
|
||||
if let Some(record) = self.lifecycle_for_request_result(&request.id).await? {
|
||||
records.push(record);
|
||||
}
|
||||
}
|
||||
records.sort_by_key(|record| record.request.id.to_hex());
|
||||
Ok(records)
|
||||
}
|
||||
|
||||
/// Remove one payload and every metadata row linked to it. This is
|
||||
/// idempotent when either side was already removed. The caller must hold
|
||||
/// `lock_lifecycle` and has already completed any required Main removal.
|
||||
pub(crate) async fn permanently_delete_request_locked(
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
) -> anyhow::Result<RequestLifecycleDeletionStats> {
|
||||
let payloads_deleted = usize::from(self.request_payload(request_id).await?.is_some());
|
||||
let metadata = self.metadata_events_for_request_result(request_id).await?;
|
||||
if payloads_deleted > 0 {
|
||||
self.db
|
||||
.delete(Filter::new().ids(vec![*request_id]))
|
||||
.await
|
||||
.map_err(|error| {
|
||||
anyhow::anyhow!("Failed to remove tombstone request {request_id}: {error}")
|
||||
})?;
|
||||
}
|
||||
let metadata_deleted = metadata.len();
|
||||
if metadata_deleted > 0 {
|
||||
self.db
|
||||
.delete(Filter::new().ids(metadata.into_iter().map(|event| event.id)))
|
||||
.await
|
||||
.map_err(|error| {
|
||||
anyhow::anyhow!("Failed to remove lifecycle metadata for {request_id}: {error}")
|
||||
})?;
|
||||
}
|
||||
Ok(RequestLifecycleDeletionStats {
|
||||
payloads_deleted,
|
||||
metadata_deleted,
|
||||
})
|
||||
}
|
||||
|
||||
/// Remove metadata whose unambiguous `e` link has no retained request.
|
||||
/// This is deliberately separate from canonical enumeration: malformed
|
||||
/// metadata remains invisible to queries but is not guessed at.
|
||||
pub async fn remove_orphan_metadata(&self) -> anyhow::Result<usize> {
|
||||
let _guard = self.lock_lifecycle().await;
|
||||
let metadata = self
|
||||
.db
|
||||
.query(Filter::new().kind(Kind::from(TOMBSTONE_REQUEST_METADATA_KIND)))
|
||||
.await
|
||||
.map_err(|error| anyhow::anyhow!("Failed to enumerate tombstone metadata: {error}"))?;
|
||||
let mut stale = Vec::new();
|
||||
for event in metadata {
|
||||
let Some(request_id) =
|
||||
Self::tag_value(&event, "e").and_then(|value| EventId::from_hex(value).ok())
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
if self.request_payload(&request_id).await?.is_none() {
|
||||
stale.push(event.id);
|
||||
}
|
||||
}
|
||||
let removed = stale.len();
|
||||
if removed > 0 {
|
||||
self.db
|
||||
.delete(Filter::new().ids(stale))
|
||||
.await
|
||||
.map_err(|error| {
|
||||
anyhow::anyhow!("Failed to remove orphan lifecycle metadata: {error}")
|
||||
})?;
|
||||
}
|
||||
Ok(removed)
|
||||
}
|
||||
|
||||
/// List the original kind-5 and kind-62 payloads retained by this store.
|
||||
/// This deliberately does not require lifecycle metadata, and never exposes
|
||||
/// relay-private metadata events to callers.
|
||||
|
||||
Reference in New Issue
Block a user