fix(sync): account for hydration outcomes

Historic missing-event recovery treated relay delivery as success even when admission or persistence left the event unservable. That could promote a relay to healthy while repository coverage was still incomplete and gave operators no way to distinguish transport, policy, and storage gaps.

Classify sync processing into saved, duplicate, purgatory, tombstone, bounded rejection, and persistence outcomes. Keep dependency-sensitive and transient failures pending during exact-ID recovery, export phase-scoped hydration counters, and document the operational queries.

Permanent invalid/blocked events and events already owned by the rejected dependency index remain terminally accounted. Recursive frontier expansion and participant relay discovery are deliberately excluded for stacked follow-up changes.

Validated with 750 library tests, the censored-relay historic recovery scenario, all sync-metrics integration tests, formatting, and clippy across all targets with warnings denied.
This commit is contained in:
DanConwayDev
2026-08-13 22:35:52 +00:00
parent 79f1d3b993
commit 39e1f692c0
4 changed files with 303 additions and 26 deletions
+10 -3
View File
@@ -898,9 +898,12 @@ maintenance timer:
(30s base doubling up to 15min; sub-second in `NGIT_TEST`), one in-flight
attempt per relay, at most 300 IDs per fetch. One persistently incomplete
relay cannot starve other relays or later batches.
- **Progress-aware**: a successful attempt clears only the IDs actually
recovered and resets the backoff; duplicate incomplete responses merge into
the existing pending set. IDs that arrive by other means (live sync, user
- **Outcome-aware**: a successful attempt clears only IDs with a durable
terminal explanation: saved, already stored, in purgatory, validly
tombstoned, blocked, or permanently invalid. Restricted, policy-error,
unknown, and persistence-error outcomes remain pending because a dependency
or transient fault may clear. Duplicate incomplete responses merge into the
existing pending set. IDs that arrive by other means (live sync, user
submission) are cleared on the next tick without consuming attempt budget.
- **Explicit expiry**: after 12 consecutive zero-progress attempts the relay's
pending IDs are dropped with a warning, and the relay stays in
@@ -1182,6 +1185,10 @@ The [`SyncMetrics`](src/sync/metrics.rs:18) module provides comprehensive monito
### Event Metrics
- `ngit_sync_events_synced_total`: Total events synced (newly saved events only, not duplicates or rejected)
- `ngit_sync_hydration_events_total{relay,phase,outcome}`: Remote hydration
work by fixed phase and outcome. Outcomes distinguish requests, deliveries, source
non-delivery, saved/duplicate/purgatory/tombstone states, bounded policy
rejection classes, and persistence failures.
### Summary Metrics
+9
View File
@@ -121,6 +121,7 @@ When GRASP-02 proactive sync is implemented, the following metrics will be added
| `ngit_sync_policy_refusals_total` | Counter | relay, category | Subscription policy refusals using bounded categories; raw reasons remain in logs |
| `ngit_sync_relay_failures` | Gauge | relay | Current consecutive failure count |
| `ngit_sync_events_synced_total` | Counter | - | Events synced (newly saved events only) |
| `ngit_sync_hydration_events_total` | Counter | relay, phase, outcome | Remote hydration requests, deliveries, and bounded persistence/admission outcomes; phase is `stream` or `recovery` |
| `ngit_sync_relays_tracked_total` | Gauge | - | Total relays discovered |
| `ngit_sync_relays_connected_total` | Gauge | - | Currently connected relay count |
| `ngit_sync_relays_dead_total` | Gauge | - | Relays marked as dead |
@@ -181,6 +182,14 @@ sum(rate(ngit_sync_connection_attempts_total{result="success"}[1h]))
# Event sync rate (newly saved events)
rate(ngit_sync_events_synced_total[5m])
# Exact-ID responses that a source did not deliver
sum by (relay) (rate(ngit_sync_hydration_events_total{phase="recovery",outcome="not_delivered"}[15m]))
# Delivered events that were not made servable
sum by (relay, outcome) (
rate(ngit_sync_hydration_events_total{outcome=~"purgatory|rejected_.*|persistence_error"}[15m])
)
# Relays with high failure counts (potential issues)
topk(10, ngit_sync_relay_failures)
+43
View File
@@ -34,6 +34,8 @@ pub struct SyncMetrics {
// === Event metrics ===
/// Total events synced (newly saved events only)
events_synced_total: IntCounter,
/// Historic/live hydration deliveries and terminal processing outcomes.
hydration_events_total: IntCounterVec,
// === Summary metrics ===
/// Total relays discovered and tracked
@@ -127,6 +129,15 @@ impl SyncMetrics {
))?;
registry.register(Box::new(events_synced_total.clone()))?;
let hydration_events_total = IntCounterVec::new(
Opts::new(
"ngit_sync_hydration_events_total",
"Hydration events by relay, phase, and bounded delivery or persistence outcome",
),
&["relay", "phase", "outcome"],
)?;
registry.register(Box::new(hydration_events_total.clone()))?;
// Summary metrics
let relays_tracked_total = IntGauge::with_opts(Opts::new(
"ngit_sync_relays_tracked_total",
@@ -236,6 +247,7 @@ impl SyncMetrics {
relay_failures,
policy_refusals_total,
events_synced_total,
hydration_events_total,
relays_tracked_total,
relays_connected_total,
relays_dead_total,
@@ -402,6 +414,18 @@ impl SyncMetrics {
self.events_synced_total.inc();
}
/// Record hydration work using a fixed outcome vocabulary supplied by the
/// sync manager. `count` allows exact-ID attempts to account for a batch
/// without performing one Prometheus update per absent response.
pub fn record_hydration_events(&self, relay: &str, phase: &str, outcome: &str, count: usize) {
if count == 0 {
return;
}
self.hydration_events_total
.with_label_values(&[relay, phase, outcome])
.inc_by(count as u64);
}
// === Summary Recording Methods ===
/// Set the total tracked relay count.
@@ -643,6 +667,25 @@ mod tests {
metrics.record_synced_event();
metrics.record_synced_event();
metrics.record_synced_event();
metrics.record_hydration_events("wss://relay.example", "recovery", "requested", 3);
metrics.record_hydration_events("wss://relay.example", "recovery", "saved", 2);
metrics.record_hydration_events("wss://relay.example", "recovery", "not_delivered", 1);
let hydration = registry
.gather()
.into_iter()
.find(|family| family.name() == "ngit_sync_hydration_events_total")
.expect("hydration outcome metric must be exported");
assert_eq!(hydration.get_metric().len(), 3);
assert_eq!(
hydration
.get_metric()
.iter()
.map(|metric| metric.get_counter().value() as u64)
.sum::<u64>(),
6
);
}
#[test]
+241 -23
View File
@@ -432,8 +432,64 @@ pub enum ProcessResult {
Duplicate,
/// Event added to Purgatory
Purgatory,
/// Event rejected by write policy
Rejected,
/// Event is absent by a valid persisted deletion or vanish request.
Tombstoned,
/// Event rejected by write policy.
Rejected(PolicyRejection),
/// The database could not be read or an accepted event could not be saved.
PersistenceError,
}
/// Bounded admission-rejection classes used by hydration accounting.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PolicyRejection {
Blocked,
Invalid,
Restricted,
Error,
Other,
PreviouslyRejected,
}
impl ProcessResult {
fn hydration_outcome(self) -> &'static str {
match self {
Self::Saved => "saved",
Self::Duplicate => "duplicate",
Self::Purgatory => "purgatory",
Self::Tombstoned => "tombstoned",
Self::Rejected(PolicyRejection::Blocked) => "rejected_blocked",
Self::Rejected(PolicyRejection::Invalid) => "rejected_invalid",
Self::Rejected(PolicyRejection::Restricted) => "rejected_restricted",
Self::Rejected(PolicyRejection::Error) => "rejected_error",
Self::Rejected(PolicyRejection::Other) => "rejected_other",
Self::Rejected(PolicyRejection::PreviouslyRejected) => "rejected_cached",
Self::PersistenceError => "persistence_error",
}
}
fn is_rejected(self) -> bool {
matches!(self, Self::Rejected(_))
}
/// Whether this result gives a durable explanation for the requested ID.
///
/// Restricted, server-error, and unknown policy failures can become valid
/// after dependencies arrive or a transient fault clears, so exact-ID
/// recovery keeps them pending. Invalid/blocked results are permanent for
/// the event, while purgatory and tombstones are explicit terminal states.
fn is_terminally_accounted(self) -> bool {
matches!(
self,
Self::Saved
| Self::Duplicate
| Self::Purgatory
| Self::Tombstoned
| Self::Rejected(PolicyRejection::Blocked)
| Self::Rejected(PolicyRejection::Invalid)
| Self::Rejected(PolicyRejection::PreviouslyRejected)
)
}
}
/// A low-volume summary of the per-relay data lane.
@@ -449,7 +505,9 @@ struct EventPipelineWindow {
saved: u64,
duplicate: u64,
purgatory: u64,
tombstoned: u64,
rejected: u64,
persistence_error: u64,
queue_delay: std::time::Duration,
max_queue_delay: std::time::Duration,
processing_time: std::time::Duration,
@@ -464,7 +522,9 @@ impl Default for EventPipelineWindow {
saved: 0,
duplicate: 0,
purgatory: 0,
tombstoned: 0,
rejected: 0,
persistence_error: 0,
queue_delay: std::time::Duration::ZERO,
max_queue_delay: std::time::Duration::ZERO,
processing_time: std::time::Duration::ZERO,
@@ -487,7 +547,9 @@ impl EventPipelineWindow {
ProcessResult::Saved => self.saved += 1,
ProcessResult::Duplicate => self.duplicate += 1,
ProcessResult::Purgatory => self.purgatory += 1,
ProcessResult::Rejected => self.rejected += 1,
ProcessResult::Tombstoned => self.tombstoned += 1,
ProcessResult::Rejected(_) => self.rejected += 1,
ProcessResult::PersistenceError => self.persistence_error += 1,
}
self.queue_delay += queue_delay;
self.max_queue_delay = self.max_queue_delay.max(queue_delay);
@@ -509,7 +571,9 @@ impl EventPipelineWindow {
saved = self.saved,
duplicate = self.duplicate,
purgatory = self.purgatory,
tombstoned = self.tombstoned,
rejected = self.rejected,
persistence_error = self.persistence_error,
events_per_second = delivered / elapsed.as_secs_f64(),
average_queue_delay_ms = self.queue_delay.as_secs_f64() * 1000.0 / delivered,
max_queue_delay_ms = self.max_queue_delay.as_secs_f64() * 1000.0,
@@ -2870,9 +2934,9 @@ impl SyncManager {
/// One bounded exact-ID recovery fetch against a single relay.
///
/// Runs outside the sync actor lock. Every event the relay returns is
/// passed through the normal write policy; an ID counts as recovered once
/// the relay has delivered the event, regardless of the policy verdict
/// (rejected events have their own re-processing machinery).
/// passed through the normal write policy. An ID counts as recovered only
/// once it has a durable terminal outcome; transient persistence and
/// dependency-sensitive policy failures remain pending.
#[allow(clippy::too_many_arguments)]
async fn run_missing_event_recovery_attempt(
relay_url: String,
@@ -2893,6 +2957,9 @@ impl SyncManager {
requested = requested.len(),
"Retrying events missing from an incomplete historic sync batch"
);
if let Some(metrics) = metrics.as_ref() {
metrics.record_hydration_events(&relay_url, "recovery", "requested", requested.len());
}
let events = match connection
.fetch_events(
@@ -2914,11 +2981,18 @@ impl SyncManager {
}
};
let mut delivered = HashSet::new();
let mut recovered = HashSet::new();
for event in events {
if !requested.contains(&event.id) {
continue;
}
if !delivered.insert(event.id) {
continue;
}
if let Some(metrics) = metrics.as_ref() {
metrics.record_hydration_events(&relay_url, "recovery", "delivered", 1);
}
// Events already tracked as rejected count as recovered without
// re-processing: re-validation is owned by the rejected-index
// re-processing machinery, and unrecoverable IDs never revalidate.
@@ -2928,6 +3002,9 @@ impl SyncManager {
event_id = %event.id,
"Recovered missing event already tracked as rejected, skipping re-processing"
);
if let Some(metrics) = metrics.as_ref() {
metrics.record_hydration_events(&relay_url, "recovery", "rejected_cached", 1);
}
recovered.insert(event.id);
continue;
}
@@ -2941,14 +3018,38 @@ impl SyncManager {
crate::nostr::persistence::SaveContext::RelaySync,
)
.await;
if result == ProcessResult::Rejected {
if let Some(metrics) = metrics.as_ref() {
metrics.record_hydration_events(
&relay_url,
"recovery",
result.hydration_outcome(),
1,
);
if result == ProcessResult::Saved {
metrics.record_synced_event();
}
}
if result.is_rejected() || result == ProcessResult::PersistenceError {
tracing::debug!(
relay = %relay_url,
event_id = %event.id,
"Recovered missing event was rejected by the write policy"
outcome = result.hydration_outcome(),
terminal = result.is_terminally_accounted(),
"Recovered missing event was not persisted"
);
}
recovered.insert(event.id);
if result.is_terminally_accounted() {
recovered.insert(event.id);
}
}
if let Some(metrics) = metrics.as_ref() {
metrics.record_hydration_events(
&relay_url,
"recovery",
"not_delivered",
requested.len().saturating_sub(delivered.len()),
);
}
let outcome = missing_event_recovery.lock().unwrap().complete_attempt(
@@ -3832,10 +3933,24 @@ impl SyncManager {
"Skipping previously rejected announcement event"
);
pipeline_window.record(
ProcessResult::Rejected,
ProcessResult::Rejected(PolicyRejection::PreviouslyRejected),
queue_delay,
processing_started.elapsed(),
);
if let Some(ref metrics) = metrics_clone {
metrics.record_hydration_events(
&relay_url_clone,
"stream",
"delivered",
1,
);
metrics.record_hydration_events(
&relay_url_clone,
"stream",
"rejected_cached",
1,
);
}
pipeline_window.report_if_due(&relay_url_clone, event_rx.len());
continue;
}
@@ -3850,6 +3965,20 @@ impl SyncManager {
crate::nostr::persistence::SaveContext::RelaySync,
)
.await;
if let Some(ref metrics) = metrics_clone {
metrics.record_hydration_events(
&relay_url_clone,
"stream",
"delivered",
1,
);
metrics.record_hydration_events(
&relay_url_clone,
"stream",
result.hydration_outcome(),
1,
);
}
// Only record metric when event is actually saved
if result == ProcessResult::Saved {
if let Some(ref metrics) = metrics_clone {
@@ -3900,7 +4029,7 @@ impl SyncManager {
// Track received event IDs for negentropy batches. Unlike REQ+EOSE
// pagination above, negentropy completion is concerned with events that
// were actually saved or already present locally.
if result == ProcessResult::Saved || result == ProcessResult::Duplicate {
if result.is_terminally_accounted() {
let mut pending = pending_sync_index.write().await;
if let Some(batches) = pending.get_mut(&relay_url_clone) {
for batch in batches.iter_mut() {
@@ -5880,7 +6009,7 @@ impl SyncManager {
)
.await;
if result != ProcessResult::Rejected {
if result.is_terminally_accounted() {
rejected_events_index.remove(&event.id);
dependency_refetch_attempts
.lock()
@@ -6161,7 +6290,15 @@ impl SyncManager {
write_policy.purgatory().enqueue_sync_immediate(&identifier);
}
}
ProcessResult::Rejected => {
ProcessResult::Tombstoned => {
stats.duplicate += 1;
tracing::debug!(
event_id = %event.id,
"{} remains absent under deletion semantics",
context
);
}
ProcessResult::Rejected(_) | ProcessResult::PersistenceError => {
stats.rejected += 1;
tracing::warn!(
event_id = %event.id,
@@ -6173,7 +6310,7 @@ impl SyncManager {
}
}
if reprocess_result != ProcessResult::Rejected {
if reprocess_result.is_terminally_accounted() {
rejected_events_index.remove(&event.id);
}
}
@@ -6190,6 +6327,35 @@ impl SyncManager {
/// - Broadcast to WebSocket subscribers via notify_event (enables recursive relay discovery)
///
/// Returns `ProcessResult` to indicate whether the event was saved, duplicate, or rejected.
fn classify_policy_result(
prefix: &nostr_sdk::prelude::MachineReadablePrefix,
message: &str,
status: bool,
) -> ProcessResult {
use nostr_sdk::prelude::MachineReadablePrefix;
if status {
return if prefix.as_str() == "purgatory" || message.contains("purgatory") {
ProcessResult::Purgatory
} else {
ProcessResult::Duplicate
};
}
if message == "this event is deleted" || message == "this pubkey has requested to vanish" {
return ProcessResult::Tombstoned;
}
let rejection = match prefix {
MachineReadablePrefix::Blocked => PolicyRejection::Blocked,
MachineReadablePrefix::Invalid => PolicyRejection::Invalid,
MachineReadablePrefix::Restricted => PolicyRejection::Restricted,
MachineReadablePrefix::Error => PolicyRejection::Error,
_ => PolicyRejection::Other,
};
ProcessResult::Rejected(rejection)
}
async fn process_event_static(
event: &Event,
relay_url: &str,
@@ -6209,7 +6375,7 @@ impl SyncManager {
}
Err(e) => {
tracing::warn!(event_id = %event.id, error = %e, "Database error checking event");
return ProcessResult::Rejected;
return ProcessResult::PersistenceError;
}
Ok(None) => {} // Continue processing
}
@@ -6229,7 +6395,7 @@ impl SyncManager {
error = %e,
"Failed to save synced event"
);
return ProcessResult::Rejected;
return ProcessResult::PersistenceError;
}
// Broadcast to WebSocket subscribers (enables recursive relay discovery)
@@ -6405,23 +6571,31 @@ impl SyncManager {
ProcessResult::Saved
}
WritePolicyResult::Reject {
message, status, ..
prefix,
message,
status,
} => {
if status {
let process_result = Self::classify_policy_result(&prefix, &message, status);
if matches!(
process_result,
ProcessResult::Purgatory | ProcessResult::Duplicate
) {
tracing::debug!(
event_id = %event.id,
kind = %event.kind.as_u16(),
outcome = process_result.hydration_outcome(),
reason = %message,
"Event added to purgatory"
"Event accepted without main-database persistence"
);
// Note: git data sync for state events is triggered by the policy
// layer when adding to purgatory (via start_state_sync)
ProcessResult::Purgatory
process_result
} else {
tracing::debug!(
event_id = %event.id,
relay = %relay_url,
kind = %event.kind.as_u16(),
outcome = process_result.hydration_outcome(),
reason = %message,
"Event rejected by write policy"
);
@@ -6498,7 +6672,7 @@ impl SyncManager {
}
}
ProcessResult::Rejected
process_result
}
}
}
@@ -7899,6 +8073,50 @@ mod tests {
);
}
#[test]
fn hydration_outcomes_separate_terminal_and_retryable_failures() {
use nostr::message::relay::SingleWord;
use nostr_sdk::prelude::MachineReadablePrefix;
let purgatory = MachineReadablePrefix::Custom(
SingleWord::from_static("purgatory").expect("valid custom prefix"),
);
assert_eq!(
SyncManager::classify_policy_result(
&purgatory,
"won't be served until git data arrives",
true,
),
ProcessResult::Purgatory
);
assert_eq!(
SyncManager::classify_policy_result(
&MachineReadablePrefix::Invalid,
"this event is deleted",
false,
),
ProcessResult::Tombstoned
);
let invalid = SyncManager::classify_policy_result(
&MachineReadablePrefix::Invalid,
"malformed event",
false,
);
let restricted = SyncManager::classify_policy_result(
&MachineReadablePrefix::Restricted,
"dependency not accepted yet",
false,
);
let persistence_error = ProcessResult::PersistenceError;
assert_eq!(invalid.hydration_outcome(), "rejected_invalid");
assert!(invalid.is_terminally_accounted());
assert_eq!(restricted.hydration_outcome(), "rejected_restricted");
assert!(!restricted.is_terminally_accounted());
assert!(!persistence_error.is_terminally_accounted());
}
#[test]
fn lifecycle_inbox_accepts_production_sized_terminal_burst() {
let (tx, mut rx) = lifecycle_notification_channel::<EoseNotification>();
@@ -8460,7 +8678,7 @@ mod tests {
.custom_created_at(Timestamp::from_secs(10))
.finalize(&keys)
.expect("build rejected event"),
ProcessResult::Rejected,
ProcessResult::Rejected(PolicyRejection::Restricted),
),
];
@@ -8469,7 +8687,7 @@ mod tests {
pagination.record_event(&event);
assert!(matches!(
policy_result,
ProcessResult::Purgatory | ProcessResult::Rejected
ProcessResult::Purgatory | ProcessResult::Rejected(_)
));
}