mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix: batch deletion lifecycle startup migration
This commit is contained in:
@@ -13,6 +13,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
|
||||
### Fixed
|
||||
|
||||
- Batch the one-time deletion-request lifecycle migration so large production databases do not remain unavailable while LMDB commits every historical request in separate transactions.
|
||||
- Prevent invitation syncing from dropping a source relay while its initial repository history is still being downloaded.
|
||||
- Fix maintainership invitation syncing by retrying the source maintainer announcement once reciprocal membership is known, allowing GRASP to discover and copy the existing repository.
|
||||
- Purgatory promotion now applies NIP-01's lowest-event-ID tie-break for same-second repository state replacements.
|
||||
|
||||
@@ -288,6 +288,10 @@ 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`.
|
||||
Missing lifecycle payloads and metadata are queued in bounded batches so LMDB
|
||||
can commit many migration writes in one transaction. The same startup snapshot
|
||||
tracks which requests are already served, avoiding a separate main-database
|
||||
lookup for every historical request while traffic is still paused.
|
||||
|
||||
The same startup pass reclassifies retained requests using current relay
|
||||
configuration and reconciles served copies against the lifecycle deadline.
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use nostr::nips::nip19::ToBech32;
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
|
||||
use anyhow::Result;
|
||||
use nostr_relay_builder::prelude::{Event, Filter, Kind, PublicKey, Timestamp};
|
||||
@@ -115,6 +115,7 @@ impl DeletionService {
|
||||
now: Timestamp,
|
||||
) -> Result<RequestLifecycleStartupStats> {
|
||||
let main = self.request_payloads_from_main().await?;
|
||||
let main_request_ids: HashSet<_> = main.iter().map(|request| request.id).collect();
|
||||
let tombstones = self.ctx.tombstones();
|
||||
// This bulk query resolves all payload/metadata pairs in one pass. Do
|
||||
// not call lifecycle_for_request_result once per request here: startup
|
||||
@@ -139,26 +140,34 @@ impl DeletionService {
|
||||
requests.entry(request.id).or_insert(request);
|
||||
}
|
||||
stats.unique_requests_discovered = requests.len();
|
||||
stats.valid_metadata_existing = existing_by_request.len();
|
||||
|
||||
let missing: Vec<_> = requests
|
||||
.values()
|
||||
.filter(|request| !existing_by_request.contains_key(&request.id))
|
||||
.map(|request| (request.clone(), self.current_classification(request)))
|
||||
.collect();
|
||||
if !missing.is_empty() {
|
||||
tracing::info!(
|
||||
requests = missing.len(),
|
||||
existing_lifecycles = stats.valid_metadata_existing,
|
||||
"Batch-migrating historical deletion-request lifecycles"
|
||||
);
|
||||
}
|
||||
let migrated = tombstones
|
||||
.record_requests_for_startup_migration(missing, now)
|
||||
.await?;
|
||||
stats.requests_migrated = migrated.len();
|
||||
existing_by_request.extend(
|
||||
migrated
|
||||
.into_iter()
|
||||
.map(|record| (record.request.id, record)),
|
||||
);
|
||||
|
||||
for request in requests.into_values() {
|
||||
let existing = existing_by_request.remove(&request.id);
|
||||
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 mut record = existing_by_request.remove(&request.id).ok_or_else(|| {
|
||||
anyhow::anyhow!("lifecycle missing after startup migration: {}", request.id)
|
||||
})?;
|
||||
|
||||
let desired = self.current_classification(&record.request);
|
||||
if record.classification != desired {
|
||||
@@ -207,10 +216,19 @@ impl DeletionService {
|
||||
&& targeting
|
||||
&& record.last_used_at.is_some();
|
||||
if serve_indefinitely || self.request_is_served(&record, now)? {
|
||||
self.promote_request(&record.request, &mut stats).await?;
|
||||
self.promote_request(
|
||||
&record.request,
|
||||
main_request_ids.contains(&record.request.id),
|
||||
&mut stats,
|
||||
)
|
||||
.await?;
|
||||
} else {
|
||||
self.remove_request_from_main(&record.request, &mut stats)
|
||||
.await?;
|
||||
self.remove_request_from_main(
|
||||
&record.request,
|
||||
main_request_ids.contains(&record.request.id),
|
||||
&mut stats,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
tracing::info!(
|
||||
@@ -260,15 +278,10 @@ impl DeletionService {
|
||||
async fn promote_request(
|
||||
&self,
|
||||
request: &Event,
|
||||
exists_in_main: bool,
|
||||
stats: &mut RequestLifecycleStartupStats,
|
||||
) -> Result<()> {
|
||||
let exists = self
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(&request.id)
|
||||
.await?
|
||||
.is_some();
|
||||
if !exists {
|
||||
if !exists_in_main {
|
||||
self.ctx.database().save_event(request).await?;
|
||||
stats.main_requests_promoted += 1;
|
||||
}
|
||||
@@ -278,15 +291,10 @@ impl DeletionService {
|
||||
async fn remove_request_from_main(
|
||||
&self,
|
||||
request: &Event,
|
||||
exists_in_main: bool,
|
||||
stats: &mut RequestLifecycleStartupStats,
|
||||
) -> Result<()> {
|
||||
let exists = self
|
||||
.ctx
|
||||
.database()
|
||||
.event_by_id(&request.id)
|
||||
.await?
|
||||
.is_some();
|
||||
if exists {
|
||||
if exists_in_main {
|
||||
self.ctx
|
||||
.database()
|
||||
.delete(Filter::new().ids(vec![request.id]))
|
||||
@@ -706,7 +714,7 @@ mod tests {
|
||||
};
|
||||
use crate::purgatory::Purgatory;
|
||||
|
||||
fn service(disrespector: bool) -> DeletionService {
|
||||
fn service_with_tombstones(disrespector: bool, tombstones: Tombstones) -> DeletionService {
|
||||
let db = Arc::new(nostr_memory::MemoryDatabase::unbounded());
|
||||
let config = crate::config::Config {
|
||||
deletion_request_disrespector: disrespector,
|
||||
@@ -718,7 +726,7 @@ mod tests {
|
||||
DeletionService::new(DeletionContext::new(
|
||||
"test.example.com",
|
||||
db,
|
||||
Tombstones::in_memory(),
|
||||
tombstones,
|
||||
HoldingStore::in_memory(),
|
||||
RepositoryLifecycle::in_memory(),
|
||||
ReplaceableHistoryStore::in_memory(),
|
||||
@@ -729,6 +737,10 @@ mod tests {
|
||||
))
|
||||
}
|
||||
|
||||
fn service(disrespector: bool) -> DeletionService {
|
||||
service_with_tombstones(disrespector, Tombstones::in_memory())
|
||||
}
|
||||
|
||||
fn deletion(keys: &Keys) -> Event {
|
||||
EventBuilder::new(Kind::EventDeletion, "")
|
||||
.finalize(keys)
|
||||
@@ -804,6 +816,43 @@ mod tests {
|
||||
.is_some());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn startup_batch_migrates_thousands_of_main_only_requests() {
|
||||
const REQUEST_COUNT: usize = 2_000;
|
||||
|
||||
let tempdir = tempfile::tempdir().unwrap();
|
||||
let tombstones = Tombstones::open_lmdb(tempdir.path()).await.unwrap();
|
||||
let service = service_with_tombstones(false, tombstones);
|
||||
let keys = Keys::generate();
|
||||
for sequence in 0..REQUEST_COUNT {
|
||||
let request = EventBuilder::new(Kind::EventDeletion, sequence.to_string())
|
||||
.finalize(&keys)
|
||||
.unwrap();
|
||||
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.main_payloads_scanned, REQUEST_COUNT);
|
||||
assert_eq!(stats.unique_requests_discovered, REQUEST_COUNT);
|
||||
assert_eq!(stats.valid_metadata_existing, 0);
|
||||
assert_eq!(stats.requests_migrated, REQUEST_COUNT);
|
||||
assert_eq!(stats.main_requests_promoted, 0);
|
||||
assert_eq!(
|
||||
service
|
||||
.ctx
|
||||
.tombstones()
|
||||
.lifecycle_records_result()
|
||||
.await
|
||||
.unwrap()
|
||||
.len(),
|
||||
REQUEST_COUNT
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn startup_reconciles_thousands_of_retained_lifecycles() {
|
||||
// Keep this large enough that a per-request full payload scan is a
|
||||
|
||||
@@ -45,6 +45,7 @@ use std::collections::HashMap;
|
||||
use std::path::Path;
|
||||
use std::sync::Arc;
|
||||
|
||||
use futures_util::future::join_all;
|
||||
use tokio::sync::{Mutex, MutexGuard};
|
||||
|
||||
use nostr_lmdb::NostrLmdb;
|
||||
@@ -69,6 +70,10 @@ pub const TOMBSTONE_REQUEST_CLASSIFICATION_TAG: &str = "tombstone-request-classi
|
||||
/// Bound the number of request IDs in a metadata `#e` query so a target with
|
||||
/// many deletion requests cannot create an unbounded database filter.
|
||||
const METADATA_REQUEST_ID_QUERY_CHUNK_SIZE: usize = 256;
|
||||
/// Number of request lifecycles queued together during the one-time migration
|
||||
/// from the main relay database. `nostr-lmdb` drains concurrent writes into a
|
||||
/// single transaction, avoiding one fsync-heavy transaction per event.
|
||||
const STARTUP_MIGRATION_CHUNK_SIZE: usize = 512;
|
||||
|
||||
/// Relay policy classification retained with a deletion request's lifecycle.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
@@ -281,6 +286,63 @@ impl Tombstones {
|
||||
})
|
||||
}
|
||||
|
||||
/// Persist historical request lifecycles in write batches during startup.
|
||||
///
|
||||
/// The caller has already enumerated canonical lifecycle metadata and only
|
||||
/// passes requests that lack a valid record. Startup remains unavailable
|
||||
/// until this completes, so a partial batch failure is safe to retry: a
|
||||
/// payload without metadata is not a valid lifecycle and will be included
|
||||
/// in the next migration pass.
|
||||
pub(crate) async fn record_requests_for_startup_migration(
|
||||
&self,
|
||||
requests: Vec<(Event, RequestClassification)>,
|
||||
first_seen_at: Timestamp,
|
||||
) -> anyhow::Result<Vec<RequestLifecycleRecord>> {
|
||||
let _guard = self.lock_lifecycle().await;
|
||||
let mut records = Vec::with_capacity(requests.len());
|
||||
|
||||
for chunk in requests.chunks(STARTUP_MIGRATION_CHUNK_SIZE) {
|
||||
let mut metadata = Vec::with_capacity(chunk.len());
|
||||
for (request, classification) in chunk {
|
||||
if !Self::is_request_kind(request.kind) {
|
||||
anyhow::bail!("Tombstone lifecycle accepts only kind-5 or kind-62 requests");
|
||||
}
|
||||
metadata.push(self.build_metadata(
|
||||
request.id,
|
||||
first_seen_at,
|
||||
None,
|
||||
*classification,
|
||||
)?);
|
||||
}
|
||||
|
||||
// Poll every save future together. The LMDB backend's ingester can
|
||||
// then drain these operations into a small number of transactions
|
||||
// instead of committing each payload and metadata event separately.
|
||||
let writes = chunk
|
||||
.iter()
|
||||
.map(|(request, _)| request)
|
||||
.chain(metadata.iter())
|
||||
.map(|event| self.db.save_event(event));
|
||||
for result in join_all(writes).await {
|
||||
result.map_err(|error| {
|
||||
anyhow::anyhow!("Failed to batch-migrate deletion-request lifecycle: {error}")
|
||||
})?;
|
||||
}
|
||||
|
||||
records.extend(chunk.iter().zip(metadata).map(
|
||||
|((request, classification), metadata)| RequestLifecycleRecord {
|
||||
metadata_event_id: metadata.id,
|
||||
request: request.clone(),
|
||||
first_seen_at,
|
||||
last_used_at: None,
|
||||
classification: *classification,
|
||||
},
|
||||
));
|
||||
}
|
||||
|
||||
Ok(records)
|
||||
}
|
||||
|
||||
/// Return the canonical lifecycle record, if both a valid payload and valid
|
||||
/// metadata exist. Multiple metadata rows are compacted lazily by updates.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user