From 2bdc11d7372b2ca389eca9156558146077d75386 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 6 Aug 2026 21:48:59 +0000 Subject: [PATCH] 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. --- tests/common/req_limiting_proxy.rs | 120 ++++++++++++++++++++++------- 1 file changed, 94 insertions(+), 26 deletions(-) diff --git a/tests/common/req_limiting_proxy.rs b/tests/common/req_limiting_proxy.rs index 3a3d0ae..22eaa35 100644 --- a/tests/common/req_limiting_proxy.rs +++ b/tests/common/req_limiting_proxy.rs @@ -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, 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, 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", , ]` - 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::>(); + + 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::>(); + + 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() {