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() {