mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat: reconcile deletion requests at startup
This commit is contained in:
@@ -244,6 +244,25 @@ request a fresh 30-day probation window. This deliberately favors preservation
|
||||
over immediate production cleanup; normal cleanup moves requests that remain
|
||||
unused out of the served main database after that window.
|
||||
|
||||
Startup performs this migration before blacklist/whitelist reconciliation and
|
||||
before the relay starts serving traffic. It discovers signed kind-5 and kind-62
|
||||
payloads from both the served and tombstone databases, including tombstone
|
||||
payloads whose metadata is missing or malformed, and deduplicates by signed
|
||||
event ID. Valid existing metadata is preserved; missing metadata is written
|
||||
with one startup timestamp and the original payload is promoted to the served
|
||||
database for its fresh probation. The pass is restart-safe: replacement
|
||||
metadata preserves the earliest `first_seen_at` and greatest `last_used_at`.
|
||||
|
||||
The same startup pass reclassifies retained requests using current relay
|
||||
configuration and reconciles served copies at the exact lifecycle deadline.
|
||||
In disrespector mode it read-only evaluates unused targeting requests against
|
||||
main and purgatory data, marking only requests that would currently have an
|
||||
effect as used; it never changes target data during that evaluation. Used
|
||||
targeting requests remain served in that mode, while non-targeting NIP-62 and
|
||||
unused no-op requests retain their ordinary bounded served schedule. Critical
|
||||
query, metadata, promotion, or removal failures fail startup rather than
|
||||
allowing traffic to begin with an incomplete lifecycle reconciliation.
|
||||
|
||||
The relay's current deletion-disrespector configuration governs retained
|
||||
requests; receipt-time mode is not permanent metadata. Switching into
|
||||
disrespector mode removes their local gates and keeps requests that are used or
|
||||
@@ -252,6 +271,10 @@ 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.
|
||||
|
||||
## Core consistency invariant
|
||||
|
||||
ngit-grasp enforces the following invariant across admission, serving, deletion,
|
||||
|
||||
@@ -47,8 +47,8 @@ pub use policy::DeletionPolicy;
|
||||
pub use runtime::{run_holding_eject, DeletionCleanupTask, DeletionRuntime, HoldingEjectArgs};
|
||||
pub use service::DeletionService;
|
||||
pub use startup::{
|
||||
BlacklistParityStats, BlacklistRestoreStats, StartupReconciliationStats, WhitelistParityStats,
|
||||
WhitelistRestoreStats,
|
||||
BlacklistParityStats, BlacklistRestoreStats, RequestLifecycleStartupStats,
|
||||
StartupReconciliationStats, WhitelistParityStats, WhitelistRestoreStats,
|
||||
};
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -253,7 +253,7 @@ impl DeletionPolicy {
|
||||
/// Read-only normal-policy target evaluation for disrespector mode.
|
||||
/// This intentionally only queries the main database and purgatory snapshots;
|
||||
/// it must not call deletion, holding, archive, rollback, or cascade helpers.
|
||||
async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result<bool> {
|
||||
pub(super) async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result<bool> {
|
||||
for id in Self::e_tag_ids(event) {
|
||||
if self.ctx.database.event_by_id(&id).await?.is_some() {
|
||||
return Ok(true);
|
||||
|
||||
@@ -61,9 +61,10 @@ impl DeletionRuntime {
|
||||
}
|
||||
|
||||
/// Run deletion-owned startup tasks before the relay begins serving traffic.
|
||||
pub async fn run_startup_tasks(&self) {
|
||||
log_startup_reconciliation(self.service.run_startup_reconciliation().await);
|
||||
pub async fn run_startup_tasks(&self) -> Result<()> {
|
||||
log_startup_reconciliation(self.service.run_startup_reconciliation().await?);
|
||||
self.run_holding_startup_cleanup().await;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Spawn deletion-owned background maintenance tasks.
|
||||
|
||||
@@ -408,7 +408,7 @@ impl DeletionService {
|
||||
DeletionAdmissionHooks { deletion: self }
|
||||
}
|
||||
|
||||
fn vanish_targets_this_relay(&self, event: &Event) -> bool {
|
||||
pub(crate) fn vanish_targets_this_relay(&self, event: &Event) -> bool {
|
||||
let relay_urls = Self::relay_url_candidates(self.ctx.domain());
|
||||
|
||||
if relay_urls.is_empty() {
|
||||
@@ -422,6 +422,17 @@ impl DeletionService {
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn would_delete_stored_target(&self, event: &Event) -> anyhow::Result<bool> {
|
||||
self.policy.would_delete_stored_target(event).await
|
||||
}
|
||||
|
||||
pub(crate) async fn would_nip62_vanish_stored_data(
|
||||
&self,
|
||||
event: &Event,
|
||||
) -> anyhow::Result<bool> {
|
||||
self.policy.would_nip62_vanish_stored_data(event).await
|
||||
}
|
||||
|
||||
fn relay_url_candidates(domain: &str) -> Vec<RelayUrl> {
|
||||
let domain = domain.trim().trim_end_matches('/');
|
||||
if domain.is_empty() {
|
||||
|
||||
@@ -1,8 +1,11 @@
|
||||
use nostr::nips::nip19::ToBech32;
|
||||
use nostr_relay_builder::prelude::{Filter, Kind, PublicKey, Timestamp};
|
||||
use std::collections::HashMap;
|
||||
|
||||
use anyhow::Result;
|
||||
use nostr_relay_builder::prelude::{Event, Filter, Kind, PublicKey, Timestamp};
|
||||
|
||||
use crate::nostr::events::RepositoryAnnouncement;
|
||||
use crate::nostr::lifecycle::DeletionSource;
|
||||
use crate::nostr::lifecycle::{DeletionSource, RequestClassification, RequestLifecycleRecord};
|
||||
|
||||
use super::DeletionService;
|
||||
|
||||
@@ -48,12 +51,30 @@ pub struct WhitelistRestoreStats {
|
||||
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct StartupReconciliationStats {
|
||||
pub request_lifecycle: RequestLifecycleStartupStats,
|
||||
pub blacklist_parity: BlacklistParityStats,
|
||||
pub blacklist_restore: BlacklistRestoreStats,
|
||||
pub whitelist_parity: WhitelistParityStats,
|
||||
pub whitelist_restore: WhitelistRestoreStats,
|
||||
}
|
||||
|
||||
/// Startup migration and policy reconciliation counters for retained deletion
|
||||
/// requests. A failure is returned to the caller rather than represented as a
|
||||
/// successful zero-count pass.
|
||||
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
|
||||
pub struct RequestLifecycleStartupStats {
|
||||
pub main_payloads_scanned: usize,
|
||||
pub tombstone_payloads_scanned: usize,
|
||||
pub unique_requests_discovered: usize,
|
||||
pub requests_migrated: usize,
|
||||
pub valid_metadata_existing: usize,
|
||||
pub requests_reclassified: usize,
|
||||
pub disrespector_requests_newly_used: usize,
|
||||
pub main_requests_promoted: usize,
|
||||
pub main_requests_removed: usize,
|
||||
pub failures: usize,
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// DeletionService startup pass implementations
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -63,22 +84,229 @@ impl DeletionService {
|
||||
///
|
||||
/// Ordering is intentionally owned here so relay startup only depends on a
|
||||
/// single facade while deletion internals decide sequencing:
|
||||
/// 1) blacklist parity delete
|
||||
/// 2) blacklist auto-restore
|
||||
/// 3) whitelist parity delete
|
||||
/// 4) whitelist restore
|
||||
pub async fn run_startup_reconciliation(&self) -> StartupReconciliationStats {
|
||||
/// 1) deletion request lifecycle migration and reconciliation
|
||||
/// 2) blacklist parity delete
|
||||
/// 3) blacklist auto-restore
|
||||
/// 4) whitelist parity delete
|
||||
/// 5) whitelist restore
|
||||
pub async fn run_startup_reconciliation(&self) -> Result<StartupReconciliationStats> {
|
||||
let request_lifecycle = self
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::now())
|
||||
.await?;
|
||||
let blacklist_parity = self.run_startup_blacklist_parity_pass().await;
|
||||
let blacklist_restore = self.run_startup_blacklist_restore_pass().await;
|
||||
let whitelist_parity = self.run_startup_whitelist_parity_pass().await;
|
||||
let whitelist_restore = self.run_startup_whitelist_restore_pass().await;
|
||||
|
||||
StartupReconciliationStats {
|
||||
Ok(StartupReconciliationStats {
|
||||
request_lifecycle,
|
||||
blacklist_parity,
|
||||
blacklist_restore,
|
||||
whitelist_parity,
|
||||
whitelist_restore,
|
||||
})
|
||||
}
|
||||
|
||||
/// Migrate historical request payloads and reconcile their current-policy
|
||||
/// serving state. The timestamp is supplied for deterministic tests and is
|
||||
/// captured once by the production startup pass.
|
||||
pub(crate) async fn run_request_lifecycle_startup_reconciliation(
|
||||
&self,
|
||||
now: Timestamp,
|
||||
) -> Result<RequestLifecycleStartupStats> {
|
||||
let main = self.request_payloads_from_main().await?;
|
||||
let tombstone_payloads = self.ctx.tombstones().request_payloads().await?;
|
||||
let mut stats = RequestLifecycleStartupStats {
|
||||
main_payloads_scanned: main.len(),
|
||||
tombstone_payloads_scanned: tombstone_payloads.len(),
|
||||
..Default::default()
|
||||
};
|
||||
let mut requests: HashMap<_, Event> = HashMap::new();
|
||||
for request in main.into_iter().chain(tombstone_payloads) {
|
||||
requests.entry(request.id).or_insert(request);
|
||||
}
|
||||
stats.unique_requests_discovered = requests.len();
|
||||
|
||||
for request in requests.into_values() {
|
||||
let existing = self
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await?;
|
||||
let mut record = match existing {
|
||||
Some(record) => {
|
||||
stats.valid_metadata_existing += 1;
|
||||
record
|
||||
}
|
||||
None => {
|
||||
let classification = self.current_classification(&request);
|
||||
let record = self
|
||||
.ctx
|
||||
.tombstones()
|
||||
.record_request(&request, now, classification)
|
||||
.await?;
|
||||
self.promote_request(&request, &mut stats).await?;
|
||||
stats.requests_migrated += 1;
|
||||
record
|
||||
}
|
||||
};
|
||||
|
||||
let desired = self.current_classification(&record.request);
|
||||
if record.classification != desired {
|
||||
record = self
|
||||
.ctx
|
||||
.tombstones()
|
||||
.reclassify_request(&record.request.id, desired)
|
||||
.await?
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!(
|
||||
"lifecycle disappeared during reclassification: {}",
|
||||
record.request.id
|
||||
)
|
||||
})?;
|
||||
stats.requests_reclassified += 1;
|
||||
}
|
||||
|
||||
let targeting = record.request.kind == Kind::EventDeletion
|
||||
|| self.vanish_targets_this_relay(&record.request);
|
||||
if self.ctx.config.deletion_request_disrespector
|
||||
&& targeting
|
||||
&& record.last_used_at.is_none()
|
||||
{
|
||||
let would_use = if record.request.kind == Kind::EventDeletion {
|
||||
self.would_delete_stored_target(&record.request).await?
|
||||
} else {
|
||||
self.would_nip62_vanish_stored_data(&record.request).await?
|
||||
};
|
||||
if would_use {
|
||||
record = self
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&record.request.id, now)
|
||||
.await?
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!(
|
||||
"lifecycle disappeared while marking used: {}",
|
||||
record.request.id
|
||||
)
|
||||
})?;
|
||||
stats.disrespector_requests_newly_used += 1;
|
||||
}
|
||||
}
|
||||
|
||||
let serve_indefinitely = self.ctx.config.deletion_request_disrespector
|
||||
&& targeting
|
||||
&& record.last_used_at.is_some();
|
||||
if serve_indefinitely || self.request_is_served(&record, now)? {
|
||||
self.promote_request(&record.request, &mut stats).await?;
|
||||
} else {
|
||||
self.remove_request_from_main(&record.request, &mut stats)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
tracing::info!(
|
||||
main_payloads_scanned = stats.main_payloads_scanned,
|
||||
tombstone_payloads_scanned = stats.tombstone_payloads_scanned,
|
||||
unique_requests = stats.unique_requests_discovered,
|
||||
migrated = stats.requests_migrated,
|
||||
existing_metadata = stats.valid_metadata_existing,
|
||||
reclassified = stats.requests_reclassified,
|
||||
disrespector_newly_used = stats.disrespector_requests_newly_used,
|
||||
promoted = stats.main_requests_promoted,
|
||||
removed = stats.main_requests_removed,
|
||||
failures = stats.failures,
|
||||
"Deletion-request startup lifecycle reconciliation completed"
|
||||
);
|
||||
Ok(stats)
|
||||
}
|
||||
|
||||
async fn request_payloads_from_main(&self) -> Result<Vec<Event>> {
|
||||
let mut requests = Vec::new();
|
||||
for kind in [Kind::EventDeletion, Kind::RequestToVanish] {
|
||||
requests.extend(
|
||||
self.ctx
|
||||
.database()
|
||||
.query(Filter::new().kind(kind))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
anyhow::anyhow!(
|
||||
"Failed to list main database deletion-request payloads: {e}"
|
||||
)
|
||||
})?,
|
||||
);
|
||||
}
|
||||
Ok(requests)
|
||||
}
|
||||
|
||||
fn current_classification(&self, request: &Event) -> RequestClassification {
|
||||
if request.kind == Kind::RequestToVanish && !self.vanish_targets_this_relay(request) {
|
||||
RequestClassification::NonTargetingNip62
|
||||
} else if self.ctx.config.deletion_request_disrespector {
|
||||
RequestClassification::Disrespector
|
||||
} else {
|
||||
RequestClassification::LocallyActionable
|
||||
}
|
||||
}
|
||||
|
||||
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,
|
||||
stats: &mut RequestLifecycleStartupStats,
|
||||
) -> Result<()> {
|
||||
let exists = self
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(&request.id)
|
||||
.await?
|
||||
.is_some();
|
||||
if !exists {
|
||||
self.ctx.database().save_event(request).await?;
|
||||
stats.main_requests_promoted += 1;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn remove_request_from_main(
|
||||
&self,
|
||||
request: &Event,
|
||||
stats: &mut RequestLifecycleStartupStats,
|
||||
) -> Result<()> {
|
||||
let exists = self
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(&request.id)
|
||||
.await?
|
||||
.is_some();
|
||||
if exists {
|
||||
self.ctx
|
||||
.database()
|
||||
.delete(Filter::new().ids(vec![request.id]))
|
||||
.await?;
|
||||
stats.main_requests_removed += 1;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Startup-only blacklist parity pass.
|
||||
@@ -475,3 +703,160 @@ impl DeletionService {
|
||||
stats
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
use std::sync::Arc;
|
||||
|
||||
use nostr_relay_builder::prelude::{EventBuilder, FinalizeEvent, Keys};
|
||||
|
||||
use super::super::DeletionContext;
|
||||
use super::*;
|
||||
use crate::grasp06::receive::new_repo_init_locks;
|
||||
use crate::nostr::lifecycle::{
|
||||
HoldingStore, ReplaceableHistoryStore, RepositoryLifecycle, 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_used_served_after_last_used_secs: 20,
|
||||
..crate::config::Config::for_testing()
|
||||
};
|
||||
DeletionService::new(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(keys: &Keys) -> Event {
|
||||
EventBuilder::new(Kind::EventDeletion, "")
|
||||
.finalize(keys)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn migrates_main_only_request_with_fresh_probation_idempotently() {
|
||||
let service = service(false);
|
||||
let request = deletion(&Keys::generate());
|
||||
service.ctx.database.save_event(&request).await.unwrap();
|
||||
|
||||
let stats = service
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(100))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.requests_migrated, 1);
|
||||
assert_eq!(stats.main_requests_promoted, 0);
|
||||
let record = service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(record.first_seen_at, Timestamp::from_secs(100));
|
||||
assert_eq!(
|
||||
record.classification,
|
||||
RequestClassification::LocallyActionable
|
||||
);
|
||||
|
||||
let rerun = service
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(200))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(rerun.requests_migrated, 0);
|
||||
assert_eq!(
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap()
|
||||
.first_seen_at,
|
||||
Timestamp::from_secs(100)
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn tombstone_only_request_is_migrated_and_promoted() {
|
||||
let service = service(false);
|
||||
let request = deletion(&Keys::generate());
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.save_request_payload_without_metadata(&request)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let stats = service
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(100))
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(stats.requests_migrated, 1);
|
||||
assert_eq!(stats.main_requests_promoted, 1);
|
||||
assert!(service
|
||||
.ctx
|
||||
.database
|
||||
.event_by_id(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reclassifies_used_request_without_erasing_timestamps() {
|
||||
let service = service(true);
|
||||
let request = deletion(&Keys::generate());
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.record_request(
|
||||
&request,
|
||||
Timestamp::from_secs(10),
|
||||
RequestClassification::LocallyActionable,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.mark_request_used(&request.id, Timestamp::from_secs(20))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
service
|
||||
.run_request_lifecycle_startup_reconciliation(Timestamp::from_secs(30))
|
||||
.await
|
||||
.unwrap();
|
||||
let record = service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_for_request_result(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert_eq!(record.first_seen_at, Timestamp::from_secs(10));
|
||||
assert_eq!(record.last_used_at, Some(Timestamp::from_secs(20)));
|
||||
assert_eq!(record.classification, RequestClassification::Disrespector);
|
||||
assert!(service
|
||||
.ctx
|
||||
.database
|
||||
.event_by_id(&request.id)
|
||||
.await
|
||||
.unwrap()
|
||||
.is_some());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -196,6 +196,18 @@ impl Tombstones {
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn save_request_payload_without_metadata(
|
||||
&self,
|
||||
event: &Event,
|
||||
) -> anyhow::Result<()> {
|
||||
self.db
|
||||
.save_event(event)
|
||||
.await
|
||||
.map(|_| ())
|
||||
.map_err(|error| anyhow::anyhow!("Failed to save test request payload: {error}"))
|
||||
}
|
||||
|
||||
/// Persist a signed deletion/vanish request and relay-generated lifecycle
|
||||
/// metadata. Exact replays retain their original `first_seen_at`.
|
||||
pub async fn record_request(
|
||||
@@ -208,7 +220,7 @@ impl Tombstones {
|
||||
anyhow::bail!("Tombstone lifecycle accepts only kind-5 or kind-62 requests");
|
||||
}
|
||||
|
||||
if let Some(record) = self.lifecycle_for_request(&event.id).await {
|
||||
if let Some(record) = self.lifecycle_for_request_result(&event.id).await? {
|
||||
return Ok(record);
|
||||
}
|
||||
|
||||
@@ -328,7 +340,7 @@ impl Tombstones {
|
||||
request_id: &EventId,
|
||||
used_at: Timestamp,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
let Some(existing) = self.lifecycle_for_request(request_id).await else {
|
||||
let Some(existing) = self.lifecycle_for_request_result(request_id).await? else {
|
||||
return Ok(None);
|
||||
};
|
||||
let last_used_at = existing.last_used_at.max(Some(used_at));
|
||||
@@ -336,37 +348,36 @@ impl Tombstones {
|
||||
return Ok(Some(existing));
|
||||
}
|
||||
|
||||
let replacement = self.build_metadata(
|
||||
*request_id,
|
||||
self.replace_lifecycle_metadata(
|
||||
request_id,
|
||||
existing.first_seen_at,
|
||||
last_used_at,
|
||||
existing.classification,
|
||||
)?;
|
||||
self.db.save_event(&replacement).await.map_err(|e| {
|
||||
anyhow::anyhow!("Failed to update lifecycle metadata for {request_id}: {e}")
|
||||
})?;
|
||||
let updated = self
|
||||
.lifecycle_for_request(request_id)
|
||||
.await
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!("Updated lifecycle metadata for {request_id} was not readable")
|
||||
})?;
|
||||
let stale_ids: Vec<_> = self
|
||||
.metadata_events_for_request(request_id)
|
||||
.await
|
||||
.into_iter()
|
||||
.map(|event| event.id)
|
||||
.filter(|id| *id != updated.metadata_event_id)
|
||||
.collect();
|
||||
if !stale_ids.is_empty() {
|
||||
self.db
|
||||
.delete(Filter::new().ids(stale_ids))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
anyhow::anyhow!("Failed to compact lifecycle metadata for {request_id}: {e}")
|
||||
})?;
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Update a request's policy classification without changing its lifecycle
|
||||
/// timestamps. The replacement is durable before stale rows are compacted,
|
||||
/// so an interrupted update always leaves a readable lifecycle record.
|
||||
pub async fn reclassify_request(
|
||||
&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);
|
||||
};
|
||||
if existing.classification == classification {
|
||||
return Ok(Some(existing));
|
||||
}
|
||||
Ok(self.lifecycle_for_request(request_id).await)
|
||||
self.replace_lifecycle_metadata(
|
||||
request_id,
|
||||
existing.first_seen_at,
|
||||
existing.last_used_at,
|
||||
classification,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// List canonical request lifecycles for migration and retention work.
|
||||
@@ -391,10 +402,89 @@ impl Tombstones {
|
||||
records
|
||||
}
|
||||
|
||||
/// 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.
|
||||
pub async fn request_payloads(&self) -> anyhow::Result<Vec<Event>> {
|
||||
let mut requests = Vec::new();
|
||||
for kind in [Kind::EventDeletion, Kind::RequestToVanish] {
|
||||
let events = self
|
||||
.db
|
||||
.query(Filter::new().kind(kind))
|
||||
.await
|
||||
.map_err(|error| {
|
||||
anyhow::anyhow!("Failed to list tombstone request payloads: {error}")
|
||||
})?;
|
||||
requests.extend(events);
|
||||
}
|
||||
requests.sort_by_key(|event| event.id.to_hex());
|
||||
requests.dedup_by_key(|event| event.id);
|
||||
Ok(requests)
|
||||
}
|
||||
|
||||
fn is_request_kind(kind: Kind) -> bool {
|
||||
kind == Kind::EventDeletion || kind == Kind::RequestToVanish
|
||||
}
|
||||
|
||||
async fn replace_lifecycle_metadata(
|
||||
&self,
|
||||
request_id: &EventId,
|
||||
first_seen_at: Timestamp,
|
||||
last_used_at: Option<Timestamp>,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<Option<RequestLifecycleRecord>> {
|
||||
// Metadata timestamps only order relay-private replacement rows. Ensure
|
||||
// this row is newer even when multiple updates happen in one wall-clock
|
||||
// second, otherwise canonical selection could retain the old class.
|
||||
let newest_existing = self
|
||||
.metadata_events_for_request_result(request_id)
|
||||
.await?
|
||||
.iter()
|
||||
.map(|event| event.created_at)
|
||||
.max();
|
||||
let replacement_created_at = match newest_existing {
|
||||
Some(existing) if existing >= Timestamp::now() => Timestamp::from_secs(
|
||||
existing
|
||||
.as_secs()
|
||||
.checked_add(1)
|
||||
.ok_or_else(|| anyhow::anyhow!("lifecycle metadata timestamp overflow"))?,
|
||||
),
|
||||
_ => Timestamp::now(),
|
||||
};
|
||||
let replacement = self.build_metadata_at(
|
||||
*request_id,
|
||||
first_seen_at,
|
||||
last_used_at,
|
||||
classification,
|
||||
replacement_created_at,
|
||||
)?;
|
||||
self.db.save_event(&replacement).await.map_err(|e| {
|
||||
anyhow::anyhow!("Failed to update lifecycle metadata for {request_id}: {e}")
|
||||
})?;
|
||||
let updated = self
|
||||
.lifecycle_for_request_result(request_id)
|
||||
.await?
|
||||
.ok_or_else(|| {
|
||||
anyhow::anyhow!("Updated lifecycle metadata for {request_id} was not readable")
|
||||
})?;
|
||||
let stale_ids: Vec<_> = self
|
||||
.metadata_events_for_request_result(request_id)
|
||||
.await?
|
||||
.into_iter()
|
||||
.map(|event| event.id)
|
||||
.filter(|id| *id != updated.metadata_event_id)
|
||||
.collect();
|
||||
if !stale_ids.is_empty() {
|
||||
self.db
|
||||
.delete(Filter::new().ids(stale_ids))
|
||||
.await
|
||||
.map_err(|e| {
|
||||
anyhow::anyhow!("Failed to compact lifecycle metadata for {request_id}: {e}")
|
||||
})?;
|
||||
}
|
||||
self.lifecycle_for_request_result(request_id).await
|
||||
}
|
||||
|
||||
async fn request_payload(&self, request_id: &EventId) -> anyhow::Result<Option<Event>> {
|
||||
for kind in [Kind::EventDeletion, Kind::RequestToVanish] {
|
||||
let events = self.db.query(Filter::new().kind(kind)).await.map_err(|e| {
|
||||
@@ -413,6 +503,23 @@ impl Tombstones {
|
||||
first_seen_at: Timestamp,
|
||||
last_used_at: Option<Timestamp>,
|
||||
classification: RequestClassification,
|
||||
) -> anyhow::Result<Event> {
|
||||
self.build_metadata_at(
|
||||
request_id,
|
||||
first_seen_at,
|
||||
last_used_at,
|
||||
classification,
|
||||
Timestamp::now(),
|
||||
)
|
||||
}
|
||||
|
||||
fn build_metadata_at(
|
||||
&self,
|
||||
request_id: EventId,
|
||||
first_seen_at: Timestamp,
|
||||
last_used_at: Option<Timestamp>,
|
||||
classification: RequestClassification,
|
||||
created_at: Timestamp,
|
||||
) -> anyhow::Result<Event> {
|
||||
let mut tags = vec![
|
||||
Tag::event(request_id),
|
||||
@@ -433,10 +540,12 @@ impl Tombstones {
|
||||
}
|
||||
EventBuilder::new(Kind::from(TOMBSTONE_REQUEST_METADATA_KIND), "")
|
||||
.tags(tags)
|
||||
.custom_created_at(created_at)
|
||||
.finalize(&self.metadata_signer)
|
||||
.map_err(|e| anyhow::anyhow!("Failed to build tombstone lifecycle metadata: {e}"))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn metadata_events_for_request(&self, request_id: &EventId) -> Vec<Event> {
|
||||
match self.metadata_events_for_request_result(request_id).await {
|
||||
Ok(events) => events,
|
||||
|
||||
+4
-1
@@ -188,7 +188,10 @@ impl RelayServer {
|
||||
.set_local_relay(relay_runtime.relay.clone());
|
||||
|
||||
let deletion_runtime = relay_runtime.deletion_runtime.clone();
|
||||
deletion_runtime.run_startup_tasks().await;
|
||||
deletion_runtime
|
||||
.run_startup_tasks()
|
||||
.await
|
||||
.context("failed deletion lifecycle startup reconciliation")?;
|
||||
|
||||
// Wire the GRASP-06 `/prs/` filesystem cleanup context into
|
||||
// purgatory so the standard expiry sweep can delete dangling
|
||||
|
||||
Reference in New Issue
Block a user