mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
Production logs after deploying a8964bb to gitnostr.com showed ~41
incomplete negentropy retries and 20 batches completing with partial
results within six minutes, some batches missing hundreds of events.
Negentropy reconciliation identifies event IDs missing locally, but a
relay's exact-ID response can return only a subset (or nothing on the
retry). Batches without repository/root-event metadata - the generic
Layer 1 announcements batch - cannot build a semantic REQ+EOSE
fallback, so handle_eose finalized them "with partial results" and
dropped the missing IDs entirely. Nothing retried them until the next
daily sync up to 25 hours later, leaving repository announcements and
their dependencies absent indefinitely.
Missing IDs from a batch that finalizes incomplete are now registered
in a per-relay recovery index (sync::missing_events), and the existing
sync maintenance timer refetches them over the relay's live connection
with bounded exponential backoff (30s doubling to 15min, one in-flight
attempt per relay, 300 IDs per fetch). Network I/O runs outside the
sync actor lock. Startup remains non-blocking: the batch still
finalizes as failed, the relay transitions to
ConnectedHistoricSyncFailures, and traffic is served while recovery
runs in the background.
Semantics:
- progress clears only the IDs actually recovered and resets backoff;
- duplicate incomplete responses merge into the pending set without
duplicating work;
- IDs satisfied by live sync or user submission are cleared on the
next tick without consuming attempt budget;
- attempts against a disconnected relay are deferred, not counted, so
an unavailable relay neither expires its work nor loops tightly;
- 12 consecutive zero-progress attempts expire the pending IDs with an
explicit warning; the relay stays observably degraded until the
daily sync re-discovers the gap;
- full recovery promotes the relay back to Connected unless an
unrelated batch failure was observed for it;
- nothing persists across restarts: historic sync re-runs from scratch
and re-detects any still-missing events, so incomplete work is never
falsely reported as complete.
Also fixes the retry-subscription-failure path, which confirmed an
incomplete batch without marking it failed (falsely reporting
Connected), and bounds the previously unbounded missing_ids log arrays
to a five-ID sample.
Regression coverage: a new censoring WebSocket proxy fixture sits
between a syncing relay and a real ngit-grasp bootstrap relay,
forwarding NIP-77 frames unchanged while withholding chosen EVENT
frames. The integration test reproduces the full production sequence
(subset response, zero-progress retry, no semantic fallback,
ConnectedHistoricSyncFailures) and proves the withheld event is
recovered and the relay promoted to Connected once the event becomes
available - without a restart and while live sync continues unstarved.
Unit tests cover registration dedupe, partial clears, backoff growth
and cap, explicit expiry, deferral, and health-restoration poisoning.
Full cargo test suite passes.
tokio-tungstenite was added as a dev-dependency for the proxy fixture;
it was already present transitively in Cargo.lock, so no Nix hash
updates are required (crates.io dependency under cargoLock).
Deliberately out of scope: durable persistence of pending recovery
work, retrying missing IDs against other relays, outbound-target
policy changes, and broader logging cleanup.
Confirms the closed issue
nostr:nevent1qqs94up6nnkzjlz4fcy5tesh8yxvr63xqjhg79etmc573fuunjt0qeqpz3mhxue69uhhyetvv9ujumn8d96zuer9wc5tdht6
246 lines
8.3 KiB
Rust
246 lines
8.3 KiB
Rust
//! Censoring WebSocket Proxy for Sync Tests
|
|
//!
|
|
//! A transparent WebSocket proxy that sits between a syncing relay and its
|
|
//! bootstrap relay, forwarding every frame except `["EVENT", ...]` messages
|
|
//! whose event ID is currently withheld.
|
|
//!
|
|
//! This simulates a relay that reports events during NIP-77 negentropy
|
|
//! reconciliation (NEG-* frames pass through untouched, so the backend's
|
|
//! full event set is visible to reconciliation) but fails to deliver some
|
|
//! of those events on exact-ID REQ fetches — the production behaviour
|
|
//! behind incomplete historic-sync batches.
|
|
//!
|
|
//! # Usage
|
|
//!
|
|
//! ```ignore
|
|
//! let source = TestRelay::start().await;
|
|
//! let proxy = CensoringProxy::start(source.url()).await;
|
|
//! proxy.withhold(event.id);
|
|
//! let syncing = TestRelay::start_with_sync(Some(proxy.url().into())).await;
|
|
//! // ... syncing relay never receives `event` ...
|
|
//! proxy.release(event.id);
|
|
//! // ... the next fetch through the proxy can deliver it ...
|
|
//! ```
|
|
|
|
use std::collections::HashSet;
|
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
use std::sync::{Arc, RwLock};
|
|
|
|
use futures_util::{SinkExt, StreamExt};
|
|
use nostr_sdk::prelude::EventId;
|
|
use tokio::net::TcpListener;
|
|
use tokio::sync::oneshot;
|
|
use tokio_tungstenite::tungstenite::Message;
|
|
|
|
/// WebSocket proxy that withholds selected EVENT frames from its backend.
|
|
pub struct CensoringProxy {
|
|
url: String,
|
|
withheld: Arc<RwLock<HashSet<String>>>,
|
|
dropped: Arc<AtomicUsize>,
|
|
shutdown_tx: Option<oneshot::Sender<()>>,
|
|
handle: Option<tokio::task::JoinHandle<()>>,
|
|
}
|
|
|
|
impl CensoringProxy {
|
|
/// Start a proxy on a random loopback port, forwarding to `backend_url`.
|
|
pub async fn start(backend_url: &str) -> Self {
|
|
let listener = TcpListener::bind("127.0.0.1:0")
|
|
.await
|
|
.expect("CensoringProxy failed to bind");
|
|
let port = listener
|
|
.local_addr()
|
|
.expect("CensoringProxy local_addr")
|
|
.port();
|
|
|
|
let withheld: Arc<RwLock<HashSet<String>>> = Arc::new(RwLock::new(HashSet::new()));
|
|
let dropped = Arc::new(AtomicUsize::new(0));
|
|
let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>();
|
|
|
|
let backend_url = backend_url.to_string();
|
|
let accept_withheld = withheld.clone();
|
|
let accept_dropped = dropped.clone();
|
|
|
|
let handle = tokio::spawn(async move {
|
|
loop {
|
|
tokio::select! {
|
|
accepted = listener.accept() => {
|
|
let Ok((stream, _)) = accepted else { break };
|
|
let backend_url = backend_url.clone();
|
|
let withheld = accept_withheld.clone();
|
|
let dropped = accept_dropped.clone();
|
|
tokio::spawn(async move {
|
|
if let Err(error) =
|
|
proxy_connection(stream, &backend_url, withheld, dropped).await
|
|
{
|
|
// Disconnects mid-test are expected; log for debugging only.
|
|
eprintln!("CensoringProxy connection ended: {error}");
|
|
}
|
|
});
|
|
}
|
|
_ = &mut shutdown_rx => break,
|
|
}
|
|
}
|
|
});
|
|
|
|
Self {
|
|
url: format!("ws://127.0.0.1:{port}"),
|
|
withheld,
|
|
dropped,
|
|
shutdown_tx: Some(shutdown_tx),
|
|
handle: Some(handle),
|
|
}
|
|
}
|
|
|
|
/// The ws:// URL the syncing relay should use as its bootstrap relay.
|
|
pub fn url(&self) -> &str {
|
|
&self.url
|
|
}
|
|
|
|
/// Start withholding EVENT frames carrying this event ID.
|
|
pub fn withhold(&self, event_id: EventId) {
|
|
self.withheld
|
|
.write()
|
|
.expect("withheld lock poisoned")
|
|
.insert(event_id.to_hex());
|
|
}
|
|
|
|
/// Stop withholding this event ID. Frames sent by the backend after this
|
|
/// call pass through; already-dropped frames are not replayed.
|
|
pub fn release(&self, event_id: EventId) {
|
|
self.withheld
|
|
.write()
|
|
.expect("withheld lock poisoned")
|
|
.remove(&event_id.to_hex());
|
|
}
|
|
|
|
/// Number of EVENT frames dropped so far.
|
|
pub fn dropped_count(&self) -> usize {
|
|
self.dropped.load(Ordering::Relaxed)
|
|
}
|
|
|
|
/// Stop the proxy.
|
|
pub async fn stop(mut self) {
|
|
if let Some(tx) = self.shutdown_tx.take() {
|
|
let _ = tx.send(());
|
|
}
|
|
if let Some(handle) = self.handle.take() {
|
|
let _ = handle.await;
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Drop for CensoringProxy {
|
|
fn drop(&mut self) {
|
|
if let Some(tx) = self.shutdown_tx.take() {
|
|
let _ = tx.send(());
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Forward one client connection to the backend, censoring backend->client
|
|
/// EVENT frames for withheld IDs.
|
|
async fn proxy_connection(
|
|
client_stream: tokio::net::TcpStream,
|
|
backend_url: &str,
|
|
withheld: Arc<RwLock<HashSet<String>>>,
|
|
dropped: Arc<AtomicUsize>,
|
|
) -> Result<(), String> {
|
|
let client_ws = tokio_tungstenite::accept_async(client_stream)
|
|
.await
|
|
.map_err(|e| format!("client handshake failed: {e}"))?;
|
|
let (backend_ws, _) = tokio_tungstenite::connect_async(backend_url)
|
|
.await
|
|
.map_err(|e| format!("backend connect failed: {e}"))?;
|
|
|
|
let (mut client_tx, mut client_rx) = client_ws.split();
|
|
let (mut backend_tx, mut backend_rx) = backend_ws.split();
|
|
|
|
// Client -> backend: forward untouched.
|
|
let upstream = async {
|
|
while let Some(message) = client_rx.next().await {
|
|
let message = message.map_err(|e| format!("client read: {e}"))?;
|
|
backend_tx
|
|
.send(message)
|
|
.await
|
|
.map_err(|e| format!("backend write: {e}"))?;
|
|
}
|
|
Ok::<(), String>(())
|
|
};
|
|
|
|
// Backend -> client: drop withheld EVENT frames, forward everything else
|
|
// (including NEG-MSG, EOSE, NOTICE, and control frames).
|
|
let downstream = async {
|
|
while let Some(message) = backend_rx.next().await {
|
|
let message = message.map_err(|e| format!("backend read: {e}"))?;
|
|
if let Message::Text(text) = &message {
|
|
let ids = withheld.read().expect("withheld lock poisoned");
|
|
if is_withheld_event_frame(text.as_str(), &ids) {
|
|
drop(ids);
|
|
dropped.fetch_add(1, Ordering::Relaxed);
|
|
continue;
|
|
}
|
|
}
|
|
client_tx
|
|
.send(message)
|
|
.await
|
|
.map_err(|e| format!("client write: {e}"))?;
|
|
}
|
|
Ok::<(), String>(())
|
|
};
|
|
|
|
// Either side ending tears the whole proxied connection down.
|
|
tokio::select! {
|
|
result = upstream => result,
|
|
result = downstream => result,
|
|
}
|
|
}
|
|
|
|
/// True if `text` is a NIP-01 `["EVENT", <sub>, {..}]` frame whose event ID
|
|
/// is in the withheld set.
|
|
fn is_withheld_event_frame(text: &str, withheld: &HashSet<String>) -> bool {
|
|
if withheld.is_empty() {
|
|
return false;
|
|
}
|
|
let Ok(value) = serde_json::from_str::<serde_json::Value>(text) else {
|
|
return false;
|
|
};
|
|
let Some(array) = value.as_array() else {
|
|
return false;
|
|
};
|
|
if array.first().and_then(|v| v.as_str()) != Some("EVENT") {
|
|
return false;
|
|
}
|
|
array
|
|
.get(2)
|
|
.and_then(|event| event.get("id"))
|
|
.and_then(|id| id.as_str())
|
|
.is_some_and(|id| withheld.contains(id))
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn withheld_event_frames_are_detected() {
|
|
let mut withheld = HashSet::new();
|
|
withheld.insert("ab".repeat(32));
|
|
let id = "ab".repeat(32);
|
|
|
|
let event_frame = format!(r#"["EVENT","sub",{{"id":"{id}","kind":1}}]"#);
|
|
assert!(is_withheld_event_frame(&event_frame, &withheld));
|
|
|
|
let other_frame = r#"["EVENT","sub",{"id":"cd","kind":1}]"#;
|
|
assert!(!is_withheld_event_frame(other_frame, &withheld));
|
|
|
|
let eose_frame = r#"["EOSE","sub"]"#;
|
|
assert!(!is_withheld_event_frame(eose_frame, &withheld));
|
|
|
|
let neg_frame = format!(r#"["NEG-MSG","sub","{id}"]"#);
|
|
assert!(
|
|
!is_withheld_event_frame(&neg_frame, &withheld),
|
|
"negentropy frames must pass through so reconciliation still reports the event"
|
|
);
|
|
}
|
|
}
|