feat(deletion): add blacklist parity and observability

This commit is contained in:
DanConwayDev
2026-06-17 14:24:24 +00:00
parent ccdc288946
commit 2f579e8e73
9 changed files with 985 additions and 50 deletions
+34 -26
View File
@@ -1,10 +1,10 @@
# Deletion Request Support (NIP-09)
**Status:** ✅ **PARTIALLY IMPLEMENTED (Phases 1–5 live; Phase 6 deferred)**
**Status:** ✅ **PARTIALLY IMPLEMENTED (core deletion lifecycle and ops visibility live; runtime hardening deferred)**
This document now reflects both the implemented single-node deletion path
(main DB + tombstones + holding DB + git archive lifecycle + recovery) and the
still-planned Phase 6 operational hardening.
still-planned runtime operational hardening.
---
@@ -173,7 +173,7 @@ The deletion system uses three separate data stores:
### Holding Database Operations
**Automatic Operations:**
- **Move to holding:** NIP-09 deletions, blacklist deletions
- **Move to holding:** NIP-09 deletions, startup blacklist parity deletions
- **Automatic recovery:** Re-publishing after NIP-09 deletion (within retention window)
- **Expiry cleanup:** Daily background task removes entries older than retention period
@@ -181,7 +181,7 @@ The deletion system uses three separate data stores:
- **Manual ejection:** Operator force-deletes before retention expires
- Use case: Large repos consuming excessive storage
- Use case: Confirmed malware requiring immediate permanent deletion
- Mechanism: CLI command or admin tool (design in Phase 6)
- Mechanism: `ngit-grasp holding-eject --owner <npub|hex> --identifier <id>`
- Logged for audit trail
- **Manual restoration:** Operator restores blacklisted repo after removal from blacklist
- Future: May support automatic restoration
@@ -259,7 +259,7 @@ When `deletion_request_disrespector = true`:
Result: Archival relay preserves all content
```
**Important:** Disrespector mode ONLY affects NIP-09 user-initiated deletions. It does NOT prevent blacklist-triggered deletions.
**Important:** Disrespector mode ONLY affects NIP-09 user-initiated deletions. It does NOT prevent blacklist-triggered deletions (including startup parity).
**Rationale:**
- NIP-09 deletions are user agency decisions (left-pad protection needed)
@@ -595,7 +595,7 @@ When implementation is complete, the following documentation will be updated:
### Completed in the current implementation
The following are implemented now (recovery still pending):
The following are implemented now:
- ngit-grasp-owned NIP-09/NIP-62 handling (backend auto-processing disabled)
- tombstone persistence and deletion re-submission gate
@@ -671,21 +671,32 @@ using served PR/PR-update fixtures (event + git-ref promotion).
- Recovery cleanup removes recovered holding metadata/payload and archive
tracking artifacts.
**Phase 6: Ops hardening and policy edge-cases** 🔄
- Blacklist-triggered deletion parity and blacklist/disrespector interaction.
- Manual ejection flow and safety/audit semantics.
- Observability: metrics for holding DB size/count, cleanup, recoveries, ejections.
**Ops parity + observability** ✅
- Startup blacklist parity pass is implemented via
`Nip34WritePolicy::run_startup_blacklist_parity_pass`.
- Matching stored announcements are deleted through the same cascade +
holding/archive pipeline as NIP-09, with holding metadata source set to
`blacklist`.
- `deletion_request_disrespector` does not block blacklist-triggered moderation
deletion paths.
- Manual operator ejection is available via:
- `ngit-grasp holding-eject --owner <npub|hex> --identifier <id>`
- Metrics are wired for blacklist deletions, holding cleanup runs/deletes,
recovery outcomes, and manual ejections.
**Remaining hardening** 🔄
- Dynamic blacklist update handling (runtime watcher/reconciliation).
- Expanded operator controls (restore flows, confirmation UX).
- Concurrency/race analysis, max-depth/scale limits, lock strategy finalization.
### Deferred items after Phase 5
The remaining intentionally deferred items are:
1. **Blacklist runtime deletion parity (Phase 6):** there is currently no active
runtime blacklist-triggered deletion path in the codebase; when introduced,
it must route through holding with `DeletionSource::Blacklist`.
2. **Manual/operator controls (Phase 6):** explicit ejection/restore workflows,
safety confirmation UX, and policy knobs beyond automatic owner re-announcement.
1. **Runtime blacklist update handling:** blacklist reconciliation is
startup-only today.
2. **Extended manual/operator controls:** restore workflows,
safety confirmation UX, and policy knobs beyond the current ejection command.
> ✅ Resolved prerequisite from earlier plan: rust-nostr backend auto-processing
> behavior is understood; ngit-grasp runs with `process_nip09(false)` and
@@ -727,17 +738,14 @@ The remaining intentionally deferred items are:
## Monitoring & Metrics
**Prometheus Metrics (Planned):**
- `ngit_deletion_requests_total` - Count of NIP-09 deletion requests received
- `ngit_deletion_requests_processed` - Count actually processed (disrespector mode = 0)
- `ngit_blacklist_deletions_total` - Count of blacklist-triggered deletions
- `ngit_holding_database_events` - Current event count in holding DB
- `ngit_holding_database_size_bytes` - Holding DB disk usage
- `ngit_archive_files_total` - Count of archive tar.gz files
- `ngit_archive_size_bytes` - Total archive disk usage
- `ngit_recoveries_total` - Count of successful automatic recoveries
- `ngit_permanent_deletions_total` - Count of events permanently deleted (post-retention)
- `ngit_manual_ejections_total` - Count of operator-initiated ejections from holding area
**Prometheus Metrics (implemented for current operational paths):**
- `ngit_blacklist_deletions_total{phase,result}`
- `ngit_holding_cleanup_runs_total`
- `ngit_holding_cleanup_deleted_total{type}`
- `ngit_holding_cleanup_last_run_deleted{type}`
- `ngit_recovery_total{result}`
- `ngit_manual_ejections_total`
- `ngit_manual_ejection_deleted_total{type}`
## Testing Strategy
+15 -1
View File
@@ -68,6 +68,20 @@ When an IP exceeds the abuse threshold, a warning is logged but the IP is never
See [Prometheus Setup Guide](../how-to/prometheus-setup.md) for NixOS configuration and Grafana dashboard provisioning.
## Deletion Lifecycle Metrics
The deletion/recovery operational paths expose these additional metrics:
| Metric | Type | Labels | Description |
|--------|------|--------|-------------|
| `ngit_blacklist_deletions_total` | Counter | `phase`, `result` | Startup blacklist parity deletion attempts/success/failure |
| `ngit_holding_cleanup_runs_total` | Counter | - | Number of holding cleanup passes run |
| `ngit_holding_cleanup_deleted_total` | Counter | `type` | Total deleted objects by cleanup (`metadata`, `payload`, `archive_file`) |
| `ngit_holding_cleanup_last_run_deleted` | Gauge | `type` | Deleted object counts for most recent cleanup pass |
| `ngit_recovery_total` | Counter | `result` | Recovery attempts and outcomes (`attempted`, `succeeded`, `failed`, `partial`) |
| `ngit_manual_ejections_total` | Counter | - | Number of operator manual ejection operations |
| `ngit_manual_ejection_deleted_total` | Counter | `type` | Objects removed by manual ejection (`metadata`, `payload`, `archive_file`) |
## Future: Load-Based Sync Scheduling (GRASP-02)
The metrics infrastructure enables future load-based scheduling for GRASP-02 sync jobs:
@@ -298,4 +312,4 @@ When a maintainer announcement arrives before the owner announcement:
**Negentropy Sync Efficiency:**
During sync, cold index IDs are excluded from "missing events" calculation, preventing wasteful re-download of events that will be rejected again.
See [work/rejected-events-index-summary.md](../../work/rejected-events-index-summary.md) for complete implementation details.
See [work/rejected-events-index-summary.md](../../work/rejected-events-index-summary.md) for complete implementation details.
+69 -1
View File
@@ -3,6 +3,7 @@ use std::{path::PathBuf, sync::Arc};
use anyhow::Result;
use clap::Parser;
use nostr_sdk::prelude::PublicKey;
use tokio::signal;
use tokio::sync::watch;
use tracing::{error, info, warn};
@@ -35,6 +36,30 @@ enum Cli {
/// Runs in dry-run mode by default. Pass --execute to make changes.
/// Stop the relay service before running with --execute.
CleanupEmptyRepos(cleanup_empty_repos::CleanupArgs),
/// Permanently eject deleted repository data from holding/archive stores.
///
/// This is an operator/admin maintenance command and is idempotent.
HoldingEject(HoldingEjectArgs),
}
#[derive(Debug, clap::Args)]
struct HoldingEjectArgs {
/// Owner pubkey (npub or hex) for the repository scope to eject.
#[arg(long)]
owner: String,
/// Repository identifier (`d` tag value).
#[arg(long)]
identifier: String,
/// Relay data path containing the holding LMDB directory.
#[arg(long, env = "NGIT_RELAY_DATA_PATH", default_value = "./data/relay")]
relay_data_path: String,
/// Git data path containing the `.archive` subtree.
#[arg(long, env = "NGIT_GIT_DATA_PATH", default_value = "./data/git")]
git_data_path: String,
}
#[tokio::main]
@@ -46,7 +71,7 @@ async fn main() -> Result<()> {
// If not, prepend the implicit "serve" subcommand so that clap routes to Cli::Serve
// and all relay flags are parsed normally (preserving backward compatibility).
let mut args: Vec<String> = std::env::args().collect();
let known_subcommands = ["serve", "cleanup-empty-repos", "help"];
let known_subcommands = ["serve", "cleanup-empty-repos", "holding-eject", "help"];
let has_subcommand = args.get(1).is_some_and(|a| {
known_subcommands.contains(&a.as_str())
|| matches!(a.as_str(), "-h" | "--help" | "-V" | "--version")
@@ -57,6 +82,7 @@ async fn main() -> Result<()> {
match Cli::parse_from(args) {
Cli::CleanupEmptyRepos(cleanup_args) => cleanup_empty_repos::run(&cleanup_args).await,
Cli::HoldingEject(eject_args) => run_holding_eject(eject_args).await,
Cli::Serve(config) => {
let mut config = *config;
// Finish initialising the Config (load relay owner key if not provided).
@@ -71,6 +97,33 @@ async fn main() -> Result<()> {
}
}
async fn run_holding_eject(args: HoldingEjectArgs) -> Result<()> {
let owner_hex = PublicKey::parse(&args.owner)
.map(|pk| pk.to_hex())
.map_err(|e| anyhow::anyhow!("invalid --owner pubkey '{}': {}", args.owner, e))?;
let holding = nostr::holding::HoldingStore::open_lmdb(
std::path::Path::new(&args.relay_data_path),
std::path::Path::new(&args.git_data_path),
)
.await?;
let stats = holding
.manual_eject_repository(&owner_hex, &args.identifier)
.await?;
println!(
"Holding ejection complete: owner={} identifier={} metadata_deleted={} payload_deleted={} archive_files_deleted={}",
owner_hex,
args.identifier,
stats.metadata_deleted,
stats.archived_events_deleted,
stats.archive_files_deleted
);
Ok(())
}
async fn run_relay(config: Config) -> Result<()> {
// Initialize tracing with configured log level
let subscriber = FmtSubscriber::builder()
@@ -177,6 +230,21 @@ async fn run_relay(config: Config) -> Result<()> {
.write_policy
.set_local_relay(relay_with_db.relay.clone());
let blacklist_stats = relay_with_db
.write_policy
.run_startup_blacklist_parity_pass()
.await;
if blacklist_stats.scanned_announcements > 0 || blacklist_stats.matched_announcements > 0 {
info!(
scanned = blacklist_stats.scanned_announcements,
matched = blacklist_stats.matched_announcements,
attempted = blacklist_stats.attempted_deletions,
succeeded = blacklist_stats.successful_deletions,
failed = blacklist_stats.failed_deletions,
"Startup blacklist parity pass completed"
);
}
// Wire the GRASP-06 `/prs/` filesystem cleanup context into
// purgatory so the standard expiry sweep can delete dangling
// refs/nostr/<event-id> refs (and zero-ref bare repos) when a
+170
View File
@@ -32,6 +32,176 @@ use connection::ConnectionTracker;
lazy_static! {
/// Global Prometheus registry for ngit-grasp metrics
pub static ref REGISTRY: Registry = Registry::new();
static ref BLACKLIST_DELETIONS_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new(
"ngit_blacklist_deletions_total",
"Blacklist-triggered deletion attempts by phase and result",
),
&["phase", "result"],
)
.expect("build blacklist deletion metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register blacklist deletion metric");
metric
};
static ref HOLDING_CLEANUP_RUNS_TOTAL: Counter = {
let metric = Counter::with_opts(Opts::new(
"ngit_holding_cleanup_runs_total",
"Number of holding cleanup passes run",
))
.expect("build holding cleanup runs metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register holding cleanup runs metric");
metric
};
static ref HOLDING_CLEANUP_DELETED_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new(
"ngit_holding_cleanup_deleted_total",
"Holding cleanup deleted objects by type",
),
&["type"],
)
.expect("build holding cleanup deleted metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register holding cleanup deleted metric");
metric
};
static ref HOLDING_CLEANUP_LAST_RUN_DELETED: GaugeVec = {
let metric = GaugeVec::new(
Opts::new(
"ngit_holding_cleanup_last_run_deleted",
"Objects deleted in the most recent holding cleanup pass by type",
),
&["type"],
)
.expect("build holding cleanup last-run metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register holding cleanup last-run metric");
metric
};
static ref RECOVERY_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new(
"ngit_recovery_total",
"Recovery attempts by outcome",
),
&["result"],
)
.expect("build recovery metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register recovery metric");
metric
};
static ref MANUAL_EJECTIONS_TOTAL: Counter = {
let metric = Counter::with_opts(Opts::new(
"ngit_manual_ejections_total",
"Number of manual holding ejections",
))
.expect("build manual ejections metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register manual ejections metric");
metric
};
static ref MANUAL_EJECTION_DELETED_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new(
"ngit_manual_ejection_deleted_total",
"Objects deleted by manual ejection operations",
),
&["type"],
)
.expect("build manual ejection deleted metric");
REGISTRY
.register(Box::new(metric.clone()))
.expect("register manual ejection deleted metric");
metric
};
}
pub fn record_blacklist_deletion_attempt(phase: &str) {
BLACKLIST_DELETIONS_TOTAL
.with_label_values(&[phase, "attempted"])
.inc();
}
pub fn record_blacklist_deletion_success(phase: &str) {
BLACKLIST_DELETIONS_TOTAL
.with_label_values(&[phase, "succeeded"])
.inc();
}
pub fn record_blacklist_deletion_failure(phase: &str) {
BLACKLIST_DELETIONS_TOTAL
.with_label_values(&[phase, "failed"])
.inc();
}
pub fn record_holding_cleanup_run(
metadata_deleted: usize,
archived_events_deleted: usize,
archive_files_deleted: usize,
) {
HOLDING_CLEANUP_RUNS_TOTAL.inc();
HOLDING_CLEANUP_DELETED_TOTAL
.with_label_values(&["metadata"])
.inc_by(metadata_deleted as f64);
HOLDING_CLEANUP_DELETED_TOTAL
.with_label_values(&["payload"])
.inc_by(archived_events_deleted as f64);
HOLDING_CLEANUP_DELETED_TOTAL
.with_label_values(&["archive_file"])
.inc_by(archive_files_deleted as f64);
HOLDING_CLEANUP_LAST_RUN_DELETED
.with_label_values(&["metadata"])
.set(metadata_deleted as f64);
HOLDING_CLEANUP_LAST_RUN_DELETED
.with_label_values(&["payload"])
.set(archived_events_deleted as f64);
HOLDING_CLEANUP_LAST_RUN_DELETED
.with_label_values(&["archive_file"])
.set(archive_files_deleted as f64);
}
pub fn record_recovery_attempt() {
RECOVERY_TOTAL.with_label_values(&["attempted"]).inc();
}
pub fn record_recovery_success() {
RECOVERY_TOTAL.with_label_values(&["succeeded"]).inc();
}
pub fn record_recovery_failed() {
RECOVERY_TOTAL.with_label_values(&["failed"]).inc();
}
pub fn record_recovery_partial() {
RECOVERY_TOTAL.with_label_values(&["partial"]).inc();
}
pub fn record_manual_ejection(
metadata_deleted: usize,
archived_events_deleted: usize,
archive_files_deleted: usize,
) {
MANUAL_EJECTIONS_TOTAL.inc();
MANUAL_EJECTION_DELETED_TOTAL
.with_label_values(&["metadata"])
.inc_by(metadata_deleted as f64);
MANUAL_EJECTION_DELETED_TOTAL
.with_label_values(&["payload"])
.inc_by(archived_events_deleted as f64);
MANUAL_EJECTION_DELETED_TOTAL
.with_label_values(&["archive_file"])
.inc_by(archive_files_deleted as f64);
}
/// Central metrics collection for ngit-grasp relay.
+109
View File
@@ -48,6 +48,15 @@ pub struct Nip34WritePolicy {
deletion_policy: DeletionPolicy,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct BlacklistParityStats {
pub scanned_announcements: usize,
pub matched_announcements: usize,
pub attempted_deletions: usize,
pub successful_deletions: usize,
pub failed_deletions: usize,
}
impl std::fmt::Debug for Nip34WritePolicy {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Nip34WritePolicy")
@@ -111,6 +120,95 @@ impl Nip34WritePolicy {
&self.ctx.holding
}
/// Startup-only blacklist parity pass.
///
/// Scans already-stored kind-30617 announcements, identifies entries
/// matching `repository_blacklist`, and applies the same deletion pipeline as
/// NIP-09 (cascade + holding/archive) but with `DeletionSource::Blacklist`.
pub async fn run_startup_blacklist_parity_pass(&self) -> BlacklistParityStats {
let blacklist = self.ctx.config.blacklist_config();
if !blacklist.enabled() {
return BlacklistParityStats::default();
}
let announcements = match self
.ctx
.database
.query(Filter::new().kind(Kind::GitRepoAnnouncement))
.await
{
Ok(events) => events,
Err(e) => {
tracing::error!(error = %e, "Blacklist startup parity scan failed to query announcements");
return BlacklistParityStats::default();
}
};
let mut stats = BlacklistParityStats {
scanned_announcements: announcements.len(),
..BlacklistParityStats::default()
};
for announcement in announcements {
let parsed = match RepositoryAnnouncement::from_event(announcement.clone()) {
Ok(p) => p,
Err(e) => {
tracing::warn!(
event_id = %announcement.id.to_hex(),
error = %e,
"Blacklist startup parity: skipping unparseable announcement"
);
continue;
}
};
let npub = announcement
.pubkey
.to_bech32()
.expect("public key to bech32 should be infallible");
let Some(reason) = blacklist.check(&npub, &parsed.identifier) else {
continue;
};
stats.matched_announcements += 1;
stats.attempted_deletions += 1;
crate::metrics::record_blacklist_deletion_attempt("startup");
match self
.deletion_policy
.apply_blacklist_deletion_for_announcement(&announcement)
.await
{
Ok(()) => {
stats.successful_deletions += 1;
crate::metrics::record_blacklist_deletion_success("startup");
tracing::info!(
event_id = %announcement.id.to_hex(),
owner = %announcement.pubkey.to_hex(),
identifier = %parsed.identifier,
reason = %reason,
"Blacklist startup parity: deleted stored repository"
);
}
Err(e) => {
stats.failed_deletions += 1;
crate::metrics::record_blacklist_deletion_failure("startup");
tracing::error!(
event_id = %announcement.id.to_hex(),
owner = %announcement.pubkey.to_hex(),
identifier = %parsed.identifier,
reason = %reason,
error = %e,
"Blacklist startup parity: deletion failed"
);
}
}
}
stats
}
/// Set the local relay for purgatory notifications.
///
/// This must be called after the relay is created since the relay depends
@@ -1235,6 +1333,8 @@ impl Nip34WritePolicy {
return false;
}
crate::metrics::record_recovery_attempt();
let mut archive_paths = BTreeSet::new();
for record in &records {
if let Some(path) = &record.archive_relative_path {
@@ -1286,6 +1386,7 @@ impl Nip34WritePolicy {
archive = %abs_path.display(),
"Recovery failed restoring git archive; aborting recovery"
);
crate::metrics::record_recovery_failed();
return false;
}
}
@@ -1311,6 +1412,7 @@ impl Nip34WritePolicy {
identifier = %identifier,
"Recovery found only non-restorable state"
);
crate::metrics::record_recovery_failed();
return false;
}
@@ -1370,6 +1472,13 @@ impl Nip34WritePolicy {
"Recovery workflow completed for re-announcement"
);
let partial = !blocked_records.is_empty() || !git_ready;
if partial {
crate::metrics::record_recovery_partial();
} else {
crate::metrics::record_recovery_success();
}
true
}
}
+112
View File
@@ -7,6 +7,7 @@
use std::path::{Component, Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use std::{collections::BTreeSet, io::ErrorKind};
use nostr_lmdb::NostrLmdb;
use nostr_memory::MemoryDatabase;
@@ -91,6 +92,14 @@ pub struct CleanupStats {
pub archive_files_deleted: usize,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct ManualEjectionStats {
pub records_matched: usize,
pub metadata_deleted: usize,
pub archived_events_deleted: usize,
pub archive_files_deleted: usize,
}
impl std::fmt::Debug for HoldingStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HoldingStore").finish_non_exhaustive()
@@ -617,6 +626,109 @@ impl HoldingStore {
}
}
crate::metrics::record_holding_cleanup_run(
stats.metadata_deleted,
stats.archived_events_deleted,
stats.archive_files_deleted,
);
Ok(stats)
}
/// Permanently eject all holding/archive artifacts for one
/// `(owner_pubkey_hex, identifier)` scope.
///
/// This is an idempotent operator operation intended for manual moderation.
/// Missing metadata/payload/archive files are tolerated.
pub async fn manual_eject_repository(
&self,
owner_pubkey_hex: &str,
identifier: &str,
) -> anyhow::Result<ManualEjectionStats> {
let filter = Filter::new().kind(Kind::from(HOLDING_METADATA_KIND));
let metadata_events = self.db.query(filter).await.map_err(|e| {
anyhow::anyhow!("Failed to query holding metadata for manual ejection: {e}")
})?;
let mut records = Vec::new();
let mut archive_paths = BTreeSet::new();
for metadata_event in metadata_events {
if !Self::metadata_matches_recovery_scope(&metadata_event, owner_pubkey_hex, identifier)
{
continue;
}
let archive_relative_path = Self::parse_archive_relative_path(&metadata_event);
if let Some(path) = &archive_relative_path {
archive_paths.insert(path.clone());
}
records.push(RecoveryMetadataRecord {
metadata_event_id: metadata_event.id,
archived_event_id: Self::parse_archived_event_id(&metadata_event),
archive_relative_path,
});
}
let mut stats = ManualEjectionStats {
records_matched: records.len(),
..ManualEjectionStats::default()
};
for record in &records {
let (metadata_deleted, payload_deleted) = self.delete_recovery_record(record).await?;
if metadata_deleted {
stats.metadata_deleted += 1;
}
if payload_deleted {
stats.archived_events_deleted += 1;
}
}
for relative_path in archive_paths {
let Some(absolute_path) = self.resolve_archive_path(&relative_path) else {
tracing::warn!(
owner = %owner_pubkey_hex,
identifier = %identifier,
path = %relative_path,
"Manual ejection skipped invalid archive path"
);
continue;
};
match std::fs::remove_file(&absolute_path) {
Ok(()) => {
stats.archive_files_deleted += 1;
}
Err(e) if e.kind() == ErrorKind::NotFound => {}
Err(e) => {
tracing::warn!(
owner = %owner_pubkey_hex,
identifier = %identifier,
path = %absolute_path.display(),
error = %e,
"Manual ejection failed deleting archive file"
);
}
}
}
tracing::info!(
owner = %owner_pubkey_hex,
identifier = %identifier,
records_matched = stats.records_matched,
metadata_deleted = stats.metadata_deleted,
payload_deleted = stats.archived_events_deleted,
archive_files_deleted = stats.archive_files_deleted,
"Manual holding ejection completed"
);
crate::metrics::record_manual_ejection(
stats.metadata_deleted,
stats.archived_events_deleted,
stats.archive_files_deleted,
);
Ok(stats)
}
}
+50 -9
View File
@@ -103,8 +103,7 @@ impl DeletionPolicy {
/// main database) so clients see an OK and the request is preserved, but it
/// is NOT acted upon — no purgatory eviction, no main-DB deletion, and no
/// tombstone recording. The targeted events therefore remain fully
/// accessible. This only affects NIP-09 user-initiated deletions. Runtime
/// blacklist-triggered deletion flow is currently deferred.
/// accessible. This only affects NIP-09 user-initiated deletions.
pub async fn handle(&self, event: &Event) -> WritePolicyResult {
// Archival mode: store the deletion request but do not process it.
if self.ctx.config.deletion_request_disrespector {
@@ -210,8 +209,13 @@ impl DeletionPolicy {
continue;
}
if Self::is_announcement_coordinate_for_author(&v[1], &event.pubkey) {
self.cascade_delete_announcement(&event.pubkey, &v[1], event.created_at)
.await;
self.cascade_delete_announcement(
&event.pubkey,
&v[1],
event.created_at,
DeletionSource::Nip09,
)
.await;
} else {
self.delete_coordinate_from_main_db(
&event.pubkey,
@@ -266,6 +270,7 @@ impl DeletionPolicy {
author: &PublicKey,
announcement_addr: &str,
deletion_created_at: Timestamp,
source: DeletionSource,
) {
let mut moved_ids = HashSet::new();
@@ -288,7 +293,7 @@ impl DeletionPolicy {
announcement_addr,
deletion_created_at,
&mut moved_ids,
DeletionSource::Nip09,
source,
)
.await;
return;
@@ -319,7 +324,7 @@ impl DeletionPolicy {
announcement_addr,
deletion_created_at,
&mut moved_ids,
DeletionSource::Nip09,
source,
)
.await;
return;
@@ -389,7 +394,7 @@ impl DeletionPolicy {
let filter = Filter::new().ids(orphan_ids);
let metadata = HoldingMetadata {
deleted_at: deletion_created_at,
source: DeletionSource::Nip09,
source,
coordinate: Some(announcement_addr.to_string()),
identifier: Some(identifier.to_string()),
owner_pubkey: Some(author.to_hex()),
@@ -417,7 +422,7 @@ impl DeletionPolicy {
announcement_addr,
deletion_created_at,
&mut moved_ids,
DeletionSource::Nip09,
source,
)
.await;
@@ -429,11 +434,47 @@ impl DeletionPolicy {
identifier,
deletion_created_at,
&mut moved_ids,
DeletionSource::Nip09,
source,
)
.await;
}
/// Apply operator-driven blacklist deletion for a stored kind-30617
/// announcement.
///
/// This path intentionally bypasses `deletion_request_disrespector` and uses
/// the same cascade + holding/archive flow as user NIP-09 deletions, but
/// marks holding metadata with `DeletionSource::Blacklist`.
pub async fn apply_blacklist_deletion_for_announcement(
&self,
announcement: &Event,
) -> anyhow::Result<()> {
if announcement.kind != Kind::GitRepoAnnouncement {
return Err(anyhow::anyhow!(
"blacklist deletion requires kind 30617 announcement, got {}",
announcement.kind.as_u16()
));
}
let Some(identifier) = identifier_from_event(announcement) else {
return Err(anyhow::anyhow!(
"announcement {} missing identifier tag",
announcement.id.to_hex()
));
};
let coordinate = format!("30617:{}:{}", announcement.pubkey.to_hex(), identifier);
self.cascade_delete_announcement(
&announcement.pubkey,
&coordinate,
Timestamp::now(),
DeletionSource::Blacklist,
)
.await;
Ok(())
}
/// Delete kind-30618 state events for `identifier` when the repository is
/// now unanchored (no surviving kind-30617 announcement with that `d` tag).
///
+64 -13
View File
@@ -53,11 +53,11 @@ pub struct TestRelay {
port: u16,
/// Temporary directory for git repositories
/// Kept alive for the lifetime of the relay
_git_data_dir: tempfile::TempDir,
_git_data_dir: Option<tempfile::TempDir>,
/// Path to git data directory (for test assertions)
git_data_path: PathBuf,
/// Temporary directory for relay data (LMDB side stores)
_relay_data_dir: tempfile::TempDir,
_relay_data_dir: Option<tempfile::TempDir>,
/// Path to relay data directory (for test assertions)
relay_data_path: PathBuf,
}
@@ -74,6 +74,9 @@ struct RelayOptions {
grasp06_enable: bool,
deletion_request_disrespector: bool,
lmdb_backend: bool,
repository_blacklist: Option<String>,
git_data_path: Option<PathBuf>,
relay_data_path: Option<PathBuf>,
}
impl TestRelay {
@@ -263,6 +266,40 @@ impl TestRelay {
.await
}
/// Start a relay with LMDB backend and repository blacklist.
pub async fn start_with_lmdb_blacklist(repository_blacklist: impl Into<String>) -> Self {
Self::start_internal(
port::reserve_port(),
RelayOptions {
lmdb_backend: true,
repository_blacklist: Some(repository_blacklist.into()),
..RelayOptions::default()
},
)
.await
}
/// Start a relay on fresh port using pre-existing LMDB/git directories.
pub async fn start_with_existing_lmdb_paths(
git_data_path: PathBuf,
relay_data_path: PathBuf,
repository_blacklist: Option<String>,
deletion_request_disrespector: bool,
) -> Self {
Self::start_internal(
port::reserve_port(),
RelayOptions {
lmdb_backend: true,
repository_blacklist,
git_data_path: Some(git_data_path),
relay_data_path: Some(relay_data_path),
deletion_request_disrespector,
..RelayOptions::default()
},
)
.await
}
/// Start a relay with every configurable option, on a pre-reserved port.
///
/// Prefer the narrower constructors above — this exists so the option
@@ -339,11 +376,25 @@ impl TestRelay {
let bind_address = format!("127.0.0.1:{}", port);
let url = format!("ws://127.0.0.1:{}", port);
// Create temporary directory for git repositories
let git_data_dir =
tempfile::tempdir().expect("Failed to create temporary git data directory");
let relay_data_dir =
tempfile::tempdir().expect("Failed to create temporary relay data directory");
// Create temporary directories unless caller provided explicit paths.
let (git_data_dir, git_data_path) = if let Some(path) = options.git_data_path.clone() {
std::fs::create_dir_all(&path).expect("Failed to create provided git data directory");
(None, path)
} else {
let dir = tempfile::tempdir().expect("Failed to create temporary git data directory");
let path = dir.path().to_path_buf();
(Some(dir), path)
};
let (relay_data_dir, relay_data_path) = if let Some(path) = options.relay_data_path.clone()
{
std::fs::create_dir_all(&path).expect("Failed to create provided relay data directory");
(None, path)
} else {
let dir = tempfile::tempdir().expect("Failed to create temporary relay data directory");
let path = dir.path().to_path_buf();
(Some(dir), path)
};
// Use the built binary directly (faster than cargo run)
let binary_path = std::env::current_exe()
@@ -367,8 +418,8 @@ impl TestRelay {
let mut cmd = Command::new(&binary_path);
cmd.env("NGIT_BIND_ADDRESS", &bind_address)
.env("NGIT_DOMAIN", &bind_address) // Set domain to match bind address
.env("NGIT_GIT_DATA_PATH", git_data_dir.path())
.env("NGIT_RELAY_DATA_PATH", relay_data_dir.path())
.env("NGIT_GIT_DATA_PATH", &git_data_path)
.env("NGIT_RELAY_DATA_PATH", &relay_data_path)
.env("NGIT_OWNER_NPUB", &test_npub)
.env("NGIT_TEST", "1") // Enable test mode: fast timers (200ms batch window, 200ms purgatory sync)
.env("NGIT_SYNC_STARTUP_DELAY_SECS", "0") // No startup delay for faster tests
@@ -427,6 +478,10 @@ impl TestRelay {
cmd.env("NGIT_DELETION_REQUEST_DISRESPECTOR", "true");
}
if let Some(ref blacklist) = options.repository_blacklist {
cmd.env("NGIT_REPOSITORY_BLACKLIST", blacklist);
}
// Release the port reservation immediately before spawning the
// subprocess that will bind it. Holding the reservation through
// env-var setup above is what keeps any concurrent
@@ -435,10 +490,6 @@ impl TestRelay {
let process = cmd.spawn().expect("Failed to start relay process");
// Store git data path for test assertions
let git_data_path = git_data_dir.path().to_path_buf();
let relay_data_path = relay_data_dir.path().to_path_buf();
let mut relay = Self {
process,
url,
+362
View File
@@ -0,0 +1,362 @@
//! Blacklist operations tests: startup parity, disrespector interaction,
//! manual ejection, and operational metrics wiring.
use std::sync::Arc;
use std::time::Duration;
use clap::Parser;
use ngit_grasp::config::Config;
use ngit_grasp::grasp06::receive::new_repo_init_locks;
use ngit_grasp::metrics::{self, Metrics};
use ngit_grasp::nostr::builder::{Nip34WritePolicy, SharedDatabase};
use ngit_grasp::nostr::holding::{
DeletionSource, GitArchiveMetadata, HoldingMetadata, HoldingStore, HOLDING_ARCHIVE_PATH_TAG,
};
use ngit_grasp::nostr::tombstones::Tombstones;
use ngit_grasp::purgatory::Purgatory;
use nostr_sdk::prelude::*;
fn metadata_has_tag(event: &Event, key: &str, value: Option<&str>) -> bool {
event.tags.iter().any(|tag| {
let v = tag.as_slice();
v.len() >= 2 && v[0] == key && value.map(|expected| v[1] == expected).unwrap_or(true)
})
}
fn repo_identifier(event: &Event) -> String {
event
.tags
.iter()
.find_map(|tag| {
let v = tag.as_slice();
(v.len() >= 2 && v[0] == "d").then(|| v[1].clone())
})
.expect("announcement identifier tag")
}
fn make_announcement(keys: &Keys, identifier: &str) -> Event {
EventBuilder::new(Kind::GitRepoAnnouncement, "")
.tags(vec![
Tag::identifier(identifier),
Tag::custom("clone", vec!["https://example.com/repo.git"]),
Tag::custom("relays", vec!["wss://example.com"]),
])
.finalize(keys)
.expect("build announcement")
}
fn make_issue(keys: &Keys, announcement: &Event) -> Event {
let coordinate = format!(
"30617:{}:{}",
announcement.pubkey.to_hex(),
repo_identifier(announcement)
);
EventBuilder::new(Kind::from(1621), "issue")
.tags(vec![Tag::custom("a", vec![coordinate])])
.finalize(keys)
.expect("build issue")
}
fn make_policy(
config: Config,
database: SharedDatabase,
holding: HoldingStore,
git_data_path: &std::path::Path,
) -> Nip34WritePolicy {
let purgatory = Arc::new(Purgatory::new(git_data_path.to_path_buf()));
Nip34WritePolicy::new(
database,
Tombstones::in_memory(),
holding,
git_data_path.to_path_buf(),
purgatory,
config,
new_repo_init_locks(),
)
}
fn base_config() -> Config {
Config::parse_from(["ngit-grasp-test", "--domain", "test.example.com"])
}
#[tokio::test]
async fn startup_blacklist_scan_deletes_matching_repositories_via_holding_archive_path() {
let relay_dir = tempfile::tempdir().expect("relay tempdir");
let git_dir = tempfile::tempdir().expect("git tempdir");
let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded());
let holding = HoldingStore::open_lmdb(relay_dir.path(), git_dir.path())
.await
.expect("open holding lmdb");
let owner = Keys::generate();
let announcement = make_announcement(&owner, "blacklisted-repo");
let issue = make_issue(&owner, &announcement);
db.save_event(&announcement)
.await
.expect("save announcement");
db.save_event(&issue).await.expect("save issue");
let owner_npub = owner.public_key().to_bech32().expect("npub");
let owner_dir = owner_npub.clone();
let repo_path = git_dir.path().join(&owner_dir).join("blacklisted-repo.git");
std::fs::create_dir_all(repo_path.join("refs")).expect("create bare repo dir");
let mut config = base_config();
config.repository_blacklist = owner_npub;
let policy = make_policy(config, db.clone(), holding.clone(), git_dir.path());
let stats = policy.run_startup_blacklist_parity_pass().await;
assert_eq!(stats.matched_announcements, 1);
assert_eq!(stats.successful_deletions, 1);
assert!(
db.event_by_id(&announcement.id)
.await
.expect("query announcement")
.is_none(),
"matching announcement must be deleted from main DB"
);
assert!(
db.event_by_id(&issue.id)
.await
.expect("query issue")
.is_none(),
"dependent events must be cascade-deleted from main DB"
);
assert!(holding.has_event(&announcement.id).await);
assert!(holding.has_event(&issue.id).await);
let announcement_meta = holding.metadata_for_event(&announcement.id).await;
assert!(
announcement_meta
.iter()
.any(|m| metadata_has_tag(m, "holding-source", Some("blacklist"))),
"holding metadata must mark blacklist deletion source"
);
let archive_rel = announcement_meta
.iter()
.flat_map(|m| m.tags.iter())
.find_map(|tag| {
let v = tag.as_slice();
(v.len() >= 2 && v[0] == HOLDING_ARCHIVE_PATH_TAG).then(|| v[1].clone())
})
.expect("blacklist deletion should record archive path");
let archive_abs = holding
.archive_absolute_path(&archive_rel)
.expect("resolve archive path");
assert!(
archive_abs.exists(),
"archive artifact must exist after deletion"
);
}
#[tokio::test]
async fn blacklist_startup_deletion_runs_even_when_disrespector_is_enabled() {
let relay_dir = tempfile::tempdir().expect("relay tempdir");
let git_dir = tempfile::tempdir().expect("git tempdir");
let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded());
let holding = HoldingStore::open_lmdb(relay_dir.path(), git_dir.path())
.await
.expect("open holding lmdb");
let owner = Keys::generate();
let announcement = make_announcement(&owner, "repo-disrespector-blacklist");
db.save_event(&announcement)
.await
.expect("save announcement");
let owner_npub = owner.public_key().to_bech32().expect("npub");
std::fs::create_dir_all(
git_dir
.path()
.join(owner_npub.clone())
.join("repo-disrespector-blacklist.git"),
)
.expect("create bare repo dir");
let config = Config {
repository_blacklist: owner_npub,
deletion_request_disrespector: true,
..base_config()
};
let policy = make_policy(config, db.clone(), holding, git_dir.path());
let stats = policy.run_startup_blacklist_parity_pass().await;
assert_eq!(stats.successful_deletions, 1);
assert!(
db.event_by_id(&announcement.id)
.await
.expect("query announcement")
.is_none(),
"disrespector mode must not block blacklist deletion path"
);
}
#[tokio::test]
async fn startup_blacklist_scan_leaves_non_matching_repositories_untouched() {
let relay_dir = tempfile::tempdir().expect("relay tempdir");
let git_dir = tempfile::tempdir().expect("git tempdir");
let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded());
let holding = HoldingStore::open_lmdb(relay_dir.path(), git_dir.path())
.await
.expect("open holding lmdb");
let owner = Keys::generate();
let announcement = make_announcement(&owner, "safe-repo");
db.save_event(&announcement)
.await
.expect("save announcement");
let mut config = base_config();
config.repository_blacklist = "different-repo".to_string();
let policy = make_policy(config, db.clone(), holding.clone(), git_dir.path());
let stats = policy.run_startup_blacklist_parity_pass().await;
assert_eq!(stats.matched_announcements, 0);
assert!(
db.event_by_id(&announcement.id)
.await
.expect("query announcement")
.is_some(),
"non-matching repository must remain in main DB"
);
assert!(holding
.metadata_for_event(&announcement.id)
.await
.is_empty());
}
#[tokio::test]
async fn manual_ejection_removes_holding_and_archive_idempotently() {
let relay_dir = tempfile::tempdir().expect("relay tempdir");
let git_dir = tempfile::tempdir().expect("git tempdir");
let holding = HoldingStore::open_lmdb(relay_dir.path(), git_dir.path())
.await
.expect("open holding lmdb");
let owner = Keys::generate();
let owner_hex = owner.public_key().to_hex();
let owner_npub = owner.public_key().to_bech32().expect("owner npub");
let event = EventBuilder::new(Kind::TextNote, "payload")
.finalize(&owner)
.expect("build event");
let archive_rel = format!("{}/{}-{}.tar.gz", owner_npub, "manual-repo", 1234);
let archive_abs = git_dir.path().join(".archive").join(&archive_rel);
std::fs::create_dir_all(archive_abs.parent().expect("archive parent"))
.expect("create archive dir");
std::fs::write(&archive_abs, b"archive-bytes").expect("write archive file");
holding
.archive_event(
&event,
&HoldingMetadata {
deleted_at: Timestamp::from_secs(1234),
source: DeletionSource::Blacklist,
coordinate: Some(format!("30617:{}:manual-repo", owner_hex)),
identifier: Some("manual-repo".to_string()),
owner_pubkey: Some(owner_hex.clone()),
git_archive: Some(GitArchiveMetadata {
relative_path: archive_rel,
created_at: Timestamp::from_secs(1234),
}),
},
)
.await
.expect("archive event with metadata");
let first = holding
.manual_eject_repository(&owner_hex, "manual-repo")
.await
.expect("first ejection");
assert_eq!(first.records_matched, 1);
assert_eq!(first.metadata_deleted, 1);
assert_eq!(first.archived_events_deleted, 1);
assert_eq!(first.archive_files_deleted, 1);
let second = holding
.manual_eject_repository(&owner_hex, "manual-repo")
.await
.expect("second ejection");
assert_eq!(second.records_matched, 0);
assert_eq!(second.metadata_deleted, 0);
assert_eq!(second.archived_events_deleted, 0);
assert_eq!(second.archive_files_deleted, 0);
}
#[tokio::test]
async fn metrics_increment_on_blacklist_cleanup_recovery_and_manual_ejection_paths() {
let relay_dir = tempfile::tempdir().expect("relay tempdir");
let git_dir = tempfile::tempdir().expect("git tempdir");
// Blacklist deletion path
let db: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded());
let holding = HoldingStore::open_lmdb(relay_dir.path(), git_dir.path())
.await
.expect("open holding lmdb");
let owner = Keys::generate();
let owner_npub = owner.public_key().to_bech32().expect("owner npub");
let announcement = make_announcement(&owner, "metrics-blacklist");
db.save_event(&announcement)
.await
.expect("save announcement");
std::fs::create_dir_all(
git_dir
.path()
.join(owner_npub.clone())
.join("metrics-blacklist.git"),
)
.expect("create bare repo dir");
let config = Config {
repository_blacklist: owner_npub,
..base_config()
};
let policy = make_policy(config, db.clone(), holding.clone(), git_dir.path());
let _ = policy.run_startup_blacklist_parity_pass().await;
// Holding cleanup path
let stale = EventBuilder::new(Kind::TextNote, "stale")
.finalize(&Keys::generate())
.expect("build stale event");
holding
.archive_event(
&stale,
&HoldingMetadata {
deleted_at: Timestamp::from_secs(1),
source: DeletionSource::Nip09,
coordinate: None,
identifier: None,
owner_pubkey: None,
git_archive: None,
},
)
.await
.expect("archive stale event");
let _ = holding
.cleanup_expired(Timestamp::from_secs(500), Duration::from_secs(10))
.await
.expect("cleanup");
// Recovery path metrics wiring (called by recovery workflow).
metrics::record_recovery_attempt();
metrics::record_recovery_partial();
// Manual ejection path
let _ = holding
.manual_eject_repository(&owner.public_key().to_hex(), "metrics-blacklist")
.await
.expect("manual ejection");
let rendered = Metrics::new(10, None).render();
assert!(rendered.contains("ngit_blacklist_deletions_total"));
assert!(rendered.contains("ngit_holding_cleanup_runs_total"));
assert!(rendered.contains("ngit_recovery_total"));
assert!(rendered.contains("ngit_manual_ejections_total"));
}