test(sync): model relay-side REQ occupancy exactly

The fifth whole-file subscription-budget validation run reported one excess REQ at a recorded peak of five. Its aggregate-only proxy evidence could not identify the rejected request, and static inspection found no subscription-opening path outside the ledger.

Make the proxy follow NIP-01 subscription semantics: a repeated REQ with the same id replaces rather than adds occupancy, EOSE retains the relay-side subscription, and only client CLOSE or backend CLOSED frees the slot. Rejections now log the rejected frame and active ids so any future distinct sixth request is attributable.

This is a test-model correction only; production scheduling and ledger behavior are unchanged. It deliberately does not classify the uninstrumented run-five rejection as a proven implementation escape.

Validated with four focused proxy tests and the standalone startup historic REQ concurrency scenario under the stricter CLOSE/CLOSED accounting model.
This commit is contained in:
DanConwayDev
2026-08-06 21:48:59 +00:00
parent 15369ca888
commit 2bdc11d737
+94 -26
View File
@@ -196,26 +196,36 @@ async fn proxy_connection(
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;
match admit_req(&mut active, subid.clone(), limit) {
ReqAdmission::Replacement => {
// NIP-01 replaces an existing subscription with
// the same id; relay-side occupancy is unchanged.
}
ReqAdmission::Rejected => {
eprintln!(
"ReqLimitingProxy rejecting REQ {subid}; active={active:?}; frame={text}"
);
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;
}
ReqAdmission::Opened => {
opened.fetch_add(1, Ordering::Relaxed);
peak.fetch_max(active.len(), Ordering::Relaxed);
}
}
active.insert(subid);
opened.fetch_add(1, Ordering::Relaxed);
peak.fetch_max(active.len(), Ordering::Relaxed);
}
Some(ClientFrame::LiveReq) => {
// Persistent live subscription (all filters
@@ -248,12 +258,14 @@ async fn proxy_connection(
let message = message.map_err(|e| format!("backend read: {e}"))?;
if let Message::Text(text) = &message {
if let Some((subid, is_eose)) = parse_backend_subscription_end(text.as_str()) {
if active.remove(&subid) {
if is_eose {
// Hold EOSE frames briefly so simultaneously
// opened REQs provably overlap at the proxy.
tokio::time::sleep(REQ_RESPONSE_DELAY).await;
}
if is_eose && active.contains(&subid) {
// Hold EOSE frames briefly so simultaneously opened
// REQs provably overlap at the proxy. EOSE ends the
// stored-event phase but does not close the relay-side
// subscription; wait for client CLOSE or backend CLOSED.
tokio::time::sleep(REQ_RESPONSE_DELAY).await;
} else if !is_eose {
apply_backend_end(&mut active, &subid, is_eose);
}
}
}
@@ -268,6 +280,28 @@ async fn proxy_connection(
Ok(())
}
#[derive(Debug, PartialEq, Eq)]
enum ReqAdmission {
Replacement,
Opened,
Rejected,
}
fn admit_req(active: &mut HashSet<String>, subid: String, limit: usize) -> ReqAdmission {
if active.contains(&subid) {
ReqAdmission::Replacement
} else if active.len() >= limit {
ReqAdmission::Rejected
} else {
active.insert(subid);
ReqAdmission::Opened
}
}
fn apply_backend_end(active: &mut HashSet<String>, subid: &str, is_eose: bool) -> bool {
!is_eose && active.remove(subid)
}
/// Parsed shape of a client-to-relay frame the proxy cares about.
enum ClientFrame {
/// `["REQ", <subid>, <filters...>]` - transient, counted.
@@ -325,7 +359,41 @@ fn parse_backend_subscription_end(text: &str) -> Option<(String, bool)> {
#[cfg(test)]
mod tests {
use super::parse_backend_subscription_end;
use std::collections::HashSet;
use super::{admit_req, apply_backend_end, parse_backend_subscription_end, ReqAdmission};
#[test]
fn same_id_req_replaces_without_consuming_another_slot() {
let mut active = (0..5).map(|id| id.to_string()).collect::<HashSet<_>>();
assert_eq!(
admit_req(&mut active, "4".to_string(), 5),
ReqAdmission::Replacement
);
assert_eq!(active.len(), 5);
}
#[test]
fn distinct_req_is_rejected_at_the_limit() {
let mut active = (0..5).map(|id| id.to_string()).collect::<HashSet<_>>();
assert_eq!(
admit_req(&mut active, "5".to_string(), 5),
ReqAdmission::Rejected
);
assert_eq!(active.len(), 5);
}
#[test]
fn eose_retains_occupancy_but_closed_releases_it() {
let mut active = HashSet::from(["historic-id".to_string()]);
assert!(!apply_backend_end(&mut active, "historic-id", true));
assert!(active.contains("historic-id"));
assert!(apply_backend_end(&mut active, "historic-id", false));
assert!(active.is_empty());
}
#[test]
fn backend_closed_definitively_releases_proxy_occupancy() {