Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 67 additions & 5 deletions crates/buzz-acp/src/relay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1234,6 +1234,19 @@ impl BgState {
while let Some(event) = self.observer_in_flight.pop_back() {
self.gated_observer_pending.push_front(event);
}
self.trim_gated_observer_pending();
}

/// Restore one observer write the relay refused with `OK false
/// rate-limited:` — correlated by id, so only that frame is retried.
fn requeue_observer_frame(&mut self, event_id: &str) {
if let Some(event) = self.take_observer_in_flight(event_id) {
self.gated_observer_pending.push_front(event);
self.trim_gated_observer_pending();
}
}

fn trim_gated_observer_pending(&mut self) {
while self.gated_observer_pending.len() > GATED_OBSERVER_QUEUE_CAP {
self.gated_observer_pending.pop_front();
self.gated_observer_dropped += 1;
Expand All @@ -1253,13 +1266,15 @@ impl BgState {
}

fn acknowledge_observer_frame(&mut self, event_id: &str) {
if let Some(index) = self
self.take_observer_in_flight(event_id);
}

fn take_observer_in_flight(&mut self, event_id: &str) -> Option<Box<Event>> {
let index = self
.observer_in_flight
.iter()
.position(|event| event.id.to_hex() == event_id)
{
self.observer_in_flight.remove(index);
}
.position(|event| event.id.to_hex() == event_id)?;
self.observer_in_flight.remove(index)
}
}

Expand Down Expand Up @@ -2392,6 +2407,22 @@ async fn handle_ws_message(
warn!("mid-session AUTH rejected (event {event_id}): {message} — triggering reconnect");
return false;
}
if !accepted && message.starts_with("rate-limited:") {
// The relay refused this EVENT for back-pressure. Arm the
// gate and put the frame (if it was a durable observer
// write) back at the head of the paced drain.
let secs = parse_rate_limit_retry_secs(&message).unwrap_or(0);
let deadline = state.set_rate_limit_gate(secs);
state.requeue_observer_frame(&event_id);
warn!(
"event {event_id} rate-limited — gate armed until ~{:.1}s from now",
deadline
.checked_duration_since(tokio::time::Instant::now())
.unwrap_or_default()
.as_secs_f64()
);
return true;
}
state.acknowledge_observer_frame(&event_id);
debug!("OK for event {event_id}: accepted={accepted} message={message}");
}
Expand Down Expand Up @@ -6027,6 +6058,37 @@ mod tests {
assert!(state.observer_in_flight.is_empty());
}

#[test]
fn rate_limited_ok_requeues_only_the_refused_frame() {
let mut state = BgState::new();
let keys = Keys::generate();
let still_in_flight = make_observer_frame(&keys);
let refused = make_observer_frame(&keys);
let later = make_observer_frame(&keys);

state.track_observer_in_flight(Box::new(still_in_flight.clone()));
state.track_observer_in_flight(Box::new(refused.clone()));
state.park_gated_observer_frame(Box::new(later.clone()));
state.requeue_observer_frame(&refused.id.to_hex());

let pending: Vec<_> = state
.gated_observer_pending
.iter()
.map(|event| event.id)
.collect();
assert_eq!(
pending,
[refused.id, later.id],
"refused frame drains first"
);
let in_flight: Vec<_> = state.observer_in_flight.iter().map(|e| e.id).collect();
assert_eq!(
in_flight,
[still_in_flight.id],
"unrefused frame keeps waiting for OK"
);
}

/// The parked-frame queue is bounded: overflow evicts the oldest frame and
/// counts it; the drain resets the counter after logging the summary.
#[tokio::test]
Expand Down
77 changes: 53 additions & 24 deletions crates/buzz-relay/src/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use uuid::Uuid;

use buzz_auth::{generate_challenge, AuthContext, LimitType};
use buzz_core::tenant::TenantContext;
use nostr::Filter;
use nostr::{Event, Filter};

use crate::handlers;
use crate::protocol::{ClientMessage, RelayMessage};
Expand Down Expand Up @@ -571,7 +571,8 @@ async fn handle_text_message(text: String, conn: Arc<ConnectionState>, state: Ar
let permit = match state.handler_semaphore.clone().try_acquire_owned() {
Ok(p) => p,
Err(_) => {
conn.send(RelayMessage::notice(
conn.send(request_rejection_message(
RejectionTarget::Event(&event),
"rate-limited: too many concurrent requests",
));
return;
Expand Down Expand Up @@ -600,7 +601,7 @@ async fn handle_text_message(text: String, conn: Arc<ConnectionState>, state: Ar
Ok(p) => p,
Err(_) => {
conn.send(request_rejection_message(
Some(&sub_id),
RejectionTarget::Subscription(&sub_id),
"rate-limited: too many concurrent requests",
));
return;
Expand Down Expand Up @@ -642,10 +643,24 @@ async fn handle_text_message(text: String, conn: Arc<ConnectionState>, state: Ar
}
}

fn request_rejection_message(sub_id: Option<&str>, reason: &str) -> String {
match sub_id {
Some(sub_id) => RelayMessage::closed(sub_id, reason),
None => RelayMessage::notice(reason),
/// Which client-side frame a rejection is correlated with.
///
/// A publisher waits on `OK <event_id>` and a subscriber waits on
/// `CLOSED <sub_id>`; an uncorrelated `NOTICE` leaves the client hanging
/// until its own timeout fires.
#[derive(Debug, Clone, Copy)]
enum RejectionTarget<'a> {
Event(&'a Event),
Subscription(&'a str),
/// COUNT has no correlated rejection frame; fall back to NOTICE.
Uncorrelated,
}

fn request_rejection_message(target: RejectionTarget<'_>, reason: &str) -> String {
match target {
RejectionTarget::Event(event) => RelayMessage::ok(&event.id.to_hex(), false, reason),
RejectionTarget::Subscription(sub_id) => RelayMessage::closed(sub_id, reason),
RejectionTarget::Uncorrelated => RelayMessage::notice(reason),
}
}

Expand Down Expand Up @@ -679,11 +694,12 @@ async fn enforce_ws_admission(
ws_limit,
)
.await;
let sub_id = match msg {
ClientMessage::Req { sub_id, .. } => Some(sub_id.as_str()),
_ => None,
let target = match msg {
ClientMessage::Event(event) => RejectionTarget::Event(event),
ClientMessage::Req { sub_id, .. } => RejectionTarget::Subscription(sub_id),
_ => RejectionTarget::Uncorrelated,
};
if !send_admission_result(conn, ws_result, sub_id) {
if !send_admission_result(conn, ws_result, target) {
return false;
}

Expand All @@ -702,7 +718,7 @@ async fn enforce_ws_admission(
message_limit,
)
.await;
if !send_admission_result(conn, message_result, None) {
if !send_admission_result(conn, message_result, target) {
return false;
}
}
Expand All @@ -713,22 +729,22 @@ async fn enforce_ws_admission(
fn send_admission_result(
conn: &ConnectionState,
result: Result<(), crate::admission::AdmissionError>,
sub_id: Option<&str>,
target: RejectionTarget<'_>,
) -> bool {
match result {
Ok(()) => true,
Err(crate::admission::AdmissionError::Exceeded { reset_in_secs }) => {
metrics::counter!("buzz_admission_rejections_total", "transport" => "websocket", "reason" => "quota").increment(1);
conn.send(request_rejection_message(
sub_id,
target,
&format!("rate-limited: quota exceeded; retry in {reset_in_secs}s"),
));
false
}
Err(crate::admission::AdmissionError::Unavailable) => {
metrics::counter!("buzz_admission_rejections_total", "transport" => "websocket", "reason" => "unavailable").increment(1);
conn.send(request_rejection_message(
sub_id,
target,
"rate-limited: shared admission unavailable",
));
false
Expand Down Expand Up @@ -835,16 +851,29 @@ mod tests {
}

#[test]
fn req_rejections_are_subscription_scoped() {
fn rejections_are_correlated_with_the_rejected_frame() {
let reason = "rate-limited: too many concurrent requests";
let closed: serde_json::Value =
serde_json::from_str(&request_rejection_message(Some("history-123"), reason))
.expect("parse CLOSED");
assert_eq!(closed, serde_json::json!(["CLOSED", "history-123", reason]));

let notice: serde_json::Value =
serde_json::from_str(&request_rejection_message(None, reason)).expect("parse NOTICE");
assert_eq!(notice, serde_json::json!(["NOTICE", reason]));
let parse = |target| -> serde_json::Value {
serde_json::from_str(&request_rejection_message(target, reason)).expect("parse frame")
};

// A rejected EVENT must settle the publisher's pending OK, not leave
// it waiting on the publish timeout.
let event = nostr::EventBuilder::text_note("hi")
.sign_with_keys(&nostr::Keys::generate())
.expect("sign event");
assert_eq!(
parse(RejectionTarget::Event(&event)),
serde_json::json!(["OK", event.id.to_hex(), false, reason])
);
assert_eq!(
parse(RejectionTarget::Subscription("history-123")),
serde_json::json!(["CLOSED", "history-123", reason])
);
assert_eq!(
parse(RejectionTarget::Uncorrelated),
serde_json::json!(["NOTICE", reason])
);
}

#[tokio::test]
Expand Down
Loading
Loading