mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
Historic REQ+EOSE sync subscribed every byte-budgeted filter group of a batch in one loop and left them all awaiting EOSE concurrently, as did the REQ+EOSE fallback, negentropy ID fetches and missing-event retries. On strfry-family relays these REQs share maxSubsPerConnection with negentropy views and live subscriptions; nos.lol (budget 20) answered each gitnostr.com startup with a burst of 'ERROR: too many concurrent REQs' NOTICEs (6 in one second at the 2026-08-05 06:41 startup; 7,739 over the prior three days). A per-connection semaphore (5 permits, shared across clones) now gates every auto-close subscription inside subscribe_filters: the permit is registered against the subscription id on success and released when the connection's own event loop sees that subscription's EOSE or CLOSED frame - before forwarding the notification, so release never depends on downstream channel consumers. This keeps a plain blocking acquire deadlock-free even though the SyncManager actor both creates subscriptions and processes EOSE: bursts pipeline at five in flight, matching the precedent of the actor already stalling inline for negentropy batches. Permits are additionally freed on unsubscribe, disconnect, event-loop termination, and by a 30-second watchdog so a relay that never answers cannot starve later subscriptions. The negentropy semaphore stays separate (different lifetimes); live subscriptions are not gated. Budget: 4 NEG + 5 REQ + 2 margin leaves at least nine slots of the tightest observed budget (20) for live subscriptions. Scheduling is deliberately per-REQ rather than per-core-filter: a paginating filter chain holds no slot between pages, so queued groups interleave breadth-first. Relay-visible concurrency is identical either way and pagination chains have no durable identity across disconnects. Reproduction: a new proxy fixture mimics the strfry limit (rejecting REQs beyond 5 with the production NOTICE, exempting limit:0 live subscriptions, delaying EOSE so REQs provably overlap), and a scenario test syncs 2500 root events into persistent storage, restarts the relay - startup recomputes filters from the full index, the shape that bursts in production - and asserts zero rejections with overlap retained. Unfixed: 3 REQs rejected (opened 91, peak pinned at the limit). Fixed: zero rejections, peak <= 5, sync completes. The TestRelay fixture gains same-port restart support with persistent LMDB storage for this. Deliberately excluded: the unified budget ledger (NIP-11-aware B/M, per-query result caps), raising the 300-ID exact-ID chunks, and any configuration surface for the bound. Validated with the full test suite (nix develop -c cargo test).
325 lines
12 KiB
Rust
325 lines
12 KiB
Rust
//! Concurrency-Limiting REQ Proxy for Sync Tests
|
|
//!
|
|
//! A transparent WebSocket proxy that sits between a syncing relay and its
|
|
//! bootstrap relay and enforces a strfry-style bound on concurrent transient
|
|
//! REQ subscriptions per connection. strfry counts REQ subscriptions against
|
|
//! `maxSubsPerConnection` and answers excess REQ frames with
|
|
//! `NOTICE ERROR: too many concurrent REQs` — the production behaviour
|
|
//! observed from nos.lol (which advertises `max_subscriptions: 20`) during
|
|
//! gitnostr.com startup bursts.
|
|
//!
|
|
//! The proxy models *awaiting-EOSE* occupancy — the quantity the client-side
|
|
//! transient-REQ bound controls: a REQ joins the active set when it is
|
|
//! opened and leaves it when the client sends CLOSE or the backend answers
|
|
//! EOSE, whichever comes first. Live subscriptions (every filter carrying
|
|
//! `limit: 0`) are exempt; they are persistent by design and budgeted
|
|
//! separately.
|
|
//!
|
|
//! The proxy additionally:
|
|
//! - records the peak number of concurrently open transient REQs, so tests
|
|
//! can assert the syncing relay's concurrency bound end to end;
|
|
//! - delays backend EOSE frames for counted REQs by a fixed interval, so
|
|
//! REQs opened together provably overlap instead of racing the loopback
|
|
//! round-trip;
|
|
//! - rejects NEG-OPEN frames as unsupported, so a syncing relay with
|
|
//! negentropy enabled still exercises the REQ+EOSE path.
|
|
//!
|
|
//! # Usage
|
|
//!
|
|
//! ```ignore
|
|
//! let source = TestRelay::start().await;
|
|
//! let proxy = ReqLimitingProxy::start(source.url(), 5).await;
|
|
//! let syncing = TestRelay::start_with_sync_no_negentropy(Some(proxy.url().into())).await;
|
|
//! // ... assert proxy.rejected_count() == 0 && proxy.peak_concurrent() <= 5 ...
|
|
//! ```
|
|
|
|
use std::collections::HashSet;
|
|
use std::sync::atomic::{AtomicUsize, Ordering};
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use futures_util::{SinkExt, StreamExt};
|
|
use tokio::net::TcpListener;
|
|
use tokio::sync::oneshot;
|
|
use tokio_tungstenite::tungstenite::Message;
|
|
|
|
/// Fixed delay applied to backend EOSE frames for counted REQs.
|
|
///
|
|
/// Long enough that a burst of REQ frames sent together is observed before
|
|
/// any subscription can complete; short enough to keep tests fast.
|
|
const REQ_RESPONSE_DELAY: Duration = Duration::from_millis(150);
|
|
|
|
/// WebSocket proxy bounding concurrent transient REQs like a strfry relay.
|
|
pub struct ReqLimitingProxy {
|
|
url: String,
|
|
peak: Arc<AtomicUsize>,
|
|
rejected: Arc<AtomicUsize>,
|
|
opened: Arc<AtomicUsize>,
|
|
shutdown_tx: Option<oneshot::Sender<()>>,
|
|
handle: Option<tokio::task::JoinHandle<()>>,
|
|
}
|
|
|
|
impl ReqLimitingProxy {
|
|
/// Start a proxy on a random loopback port, forwarding to `backend_url`
|
|
/// and rejecting transient REQ frames that would exceed `limit`
|
|
/// concurrent subscriptions on a connection.
|
|
pub async fn start(backend_url: &str, limit: usize) -> Self {
|
|
let listener = TcpListener::bind("127.0.0.1:0")
|
|
.await
|
|
.expect("ReqLimitingProxy failed to bind");
|
|
let port = listener
|
|
.local_addr()
|
|
.expect("ReqLimitingProxy local_addr")
|
|
.port();
|
|
|
|
let peak = Arc::new(AtomicUsize::new(0));
|
|
let rejected = Arc::new(AtomicUsize::new(0));
|
|
let opened = Arc::new(AtomicUsize::new(0));
|
|
let (shutdown_tx, mut shutdown_rx) = oneshot::channel::<()>();
|
|
|
|
let backend_url = backend_url.to_string();
|
|
let accept_peak = peak.clone();
|
|
let accept_rejected = rejected.clone();
|
|
let accept_opened = opened.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 peak = accept_peak.clone();
|
|
let rejected = accept_rejected.clone();
|
|
let opened = accept_opened.clone();
|
|
tokio::spawn(async move {
|
|
if let Err(error) = proxy_connection(
|
|
stream,
|
|
&backend_url,
|
|
limit,
|
|
peak,
|
|
rejected,
|
|
opened,
|
|
)
|
|
.await
|
|
{
|
|
// Disconnects mid-test are expected; log only.
|
|
eprintln!("ReqLimitingProxy connection ended: {error}");
|
|
}
|
|
});
|
|
}
|
|
_ = &mut shutdown_rx => break,
|
|
}
|
|
}
|
|
});
|
|
|
|
Self {
|
|
url: format!("ws://127.0.0.1:{port}"),
|
|
peak,
|
|
rejected,
|
|
opened,
|
|
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
|
|
}
|
|
|
|
/// Highest number of transient REQs observed open at once on any
|
|
/// connection.
|
|
pub fn peak_concurrent(&self) -> usize {
|
|
self.peak.load(Ordering::Relaxed)
|
|
}
|
|
|
|
/// Number of REQ frames rejected for exceeding the limit.
|
|
pub fn rejected_count(&self) -> usize {
|
|
self.rejected.load(Ordering::Relaxed)
|
|
}
|
|
|
|
/// Total transient REQ frames accepted and forwarded to the backend.
|
|
pub fn opened_count(&self) -> usize {
|
|
self.opened.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 ReqLimitingProxy {
|
|
fn drop(&mut self) {
|
|
if let Some(tx) = self.shutdown_tx.take() {
|
|
let _ = tx.send(());
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Forward one client connection to the backend, bounding concurrent
|
|
/// transient REQs and recording the observed peak.
|
|
///
|
|
/// A single loop owns both writers so rejections can be answered directly
|
|
/// to the client without forwarding to the backend.
|
|
async fn proxy_connection(
|
|
client_stream: tokio::net::TcpStream,
|
|
backend_url: &str,
|
|
limit: usize,
|
|
peak: Arc<AtomicUsize>,
|
|
rejected: Arc<AtomicUsize>,
|
|
opened: 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();
|
|
|
|
// Transient REQ subscription ids currently open on this connection.
|
|
let mut active: HashSet<String> = HashSet::new();
|
|
|
|
loop {
|
|
tokio::select! {
|
|
message = client_rx.next() => {
|
|
let Some(message) = message else { break };
|
|
let message = message.map_err(|e| format!("client read: {e}"))?;
|
|
if let Message::Text(text) = &message {
|
|
match parse_client_frame(text.as_str()) {
|
|
Some(ClientFrame::Req(subid)) => {
|
|
if active.len() >= limit {
|
|
rejected.fetch_add(1, Ordering::Relaxed);
|
|
client_tx
|
|
.send(Message::text(
|
|
r#"["NOTICE","ERROR: too many concurrent REQs"]"#
|
|
.to_string(),
|
|
))
|
|
.await
|
|
.map_err(|e| format!("client write: {e}"))?;
|
|
client_tx
|
|
.send(Message::text(format!(
|
|
r#"["CLOSED","{subid}","blocked: too many concurrent REQs"]"#
|
|
)))
|
|
.await
|
|
.map_err(|e| format!("client write: {e}"))?;
|
|
continue;
|
|
}
|
|
active.insert(subid);
|
|
opened.fetch_add(1, Ordering::Relaxed);
|
|
peak.fetch_max(active.len(), Ordering::Relaxed);
|
|
}
|
|
Some(ClientFrame::LiveReq) => {
|
|
// Persistent live subscription (all filters
|
|
// limit:0) - budgeted separately, not counted.
|
|
}
|
|
Some(ClientFrame::Close(subid)) => {
|
|
active.remove(&subid);
|
|
}
|
|
Some(ClientFrame::NegOpen(subid)) => {
|
|
// This proxy models a relay without NIP-77 so
|
|
// historic sync exercises the REQ+EOSE path.
|
|
client_tx
|
|
.send(Message::text(format!(
|
|
r#"["NEG-ERR","{subid}","blocked: negentropy not supported"]"#
|
|
)))
|
|
.await
|
|
.map_err(|e| format!("client write: {e}"))?;
|
|
continue;
|
|
}
|
|
None => {}
|
|
}
|
|
}
|
|
backend_tx
|
|
.send(message)
|
|
.await
|
|
.map_err(|e| format!("backend write: {e}"))?;
|
|
}
|
|
message = backend_rx.next() => {
|
|
let Some(message) = message else { break };
|
|
let message = message.map_err(|e| format!("backend read: {e}"))?;
|
|
if let Message::Text(text) = &message {
|
|
if let Some(subid) = parse_backend_eose(text.as_str()) {
|
|
if active.remove(&subid) {
|
|
// Hold EOSE frames briefly so simultaneously
|
|
// opened REQs provably overlap at the proxy.
|
|
tokio::time::sleep(REQ_RESPONSE_DELAY).await;
|
|
}
|
|
}
|
|
}
|
|
client_tx
|
|
.send(message)
|
|
.await
|
|
.map_err(|e| format!("client write: {e}"))?;
|
|
}
|
|
}
|
|
}
|
|
|
|
Ok(())
|
|
}
|
|
|
|
/// Parsed shape of a client-to-relay frame the proxy cares about.
|
|
enum ClientFrame {
|
|
/// `["REQ", <subid>, <filters...>]` - transient, counted.
|
|
Req(String),
|
|
/// `["REQ", ...]` where every filter carries `limit: 0` - a persistent
|
|
/// live subscription, exempt from the transient bound.
|
|
LiveReq,
|
|
/// `["CLOSE", <subid>]` - ends the subscription.
|
|
Close(String),
|
|
/// `["NEG-OPEN", <subid>, ...]` - rejected as unsupported.
|
|
NegOpen(String),
|
|
}
|
|
|
|
/// Classify a client text frame if the proxy tracks it, else `None`.
|
|
fn parse_client_frame(text: &str) -> Option<ClientFrame> {
|
|
if !(text.starts_with("[\"REQ")
|
|
|| text.starts_with("[\"CLOSE")
|
|
|| text.starts_with("[\"NEG-OPEN"))
|
|
{
|
|
return None;
|
|
}
|
|
let value = serde_json::from_str::<serde_json::Value>(text).ok()?;
|
|
let array = value.as_array()?;
|
|
let kind = array.first()?.as_str()?;
|
|
let subid = array.get(1)?.as_str()?.to_string();
|
|
match kind {
|
|
"REQ" => {
|
|
let filters = &array[2..];
|
|
let live = !filters.is_empty()
|
|
&& filters
|
|
.iter()
|
|
.all(|f| f.get("limit").and_then(|l| l.as_u64()) == Some(0));
|
|
if live {
|
|
Some(ClientFrame::LiveReq)
|
|
} else {
|
|
Some(ClientFrame::Req(subid))
|
|
}
|
|
}
|
|
"CLOSE" => Some(ClientFrame::Close(subid)),
|
|
"NEG-OPEN" => Some(ClientFrame::NegOpen(subid)),
|
|
_ => None,
|
|
}
|
|
}
|
|
|
|
/// Extract the subscription id from a backend `["EOSE", <subid>]` frame.
|
|
fn parse_backend_eose(text: &str) -> Option<String> {
|
|
if !text.starts_with("[\"EOSE") {
|
|
return None;
|
|
}
|
|
let value = serde_json::from_str::<serde_json::Value>(text).ok()?;
|
|
let array = value.as_array()?;
|
|
if array.first()?.as_str()? != "EOSE" {
|
|
return None;
|
|
}
|
|
Some(array.get(1)?.as_str()?.to_string())
|
|
}
|