Sitelet https://github.com/init4tech/node-components/commit/54c0d97cf24b5265b2403fb909bb8ceed2650347
Skip to content

Commit 54c0d97

Browse files
prestwichclaude
andcommitted
feat: restore permit-based subscription notification batching
Replace the one-notification-per-iteration drain pattern with permit_many batching for better throughput and fairness. Bumps ajj to 0.6.2 for permit_many/notification_capacity support. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 86d65ce commit 54c0d97

2 files changed

Lines changed: 66 additions & 60 deletions

File tree

‎crates/rpc-storage/src/interest/subs.rs‎

Lines changed: 33 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,9 @@ use ajj::HandlerCtx;
55
use alloy::{primitives::U64, rpc::types::Log};
66
use dashmap::DashMap;
77
use std::{
8+
cmp::min,
89
collections::VecDeque,
10+
future::pending,
911
sync::{
1012
Arc, Weak,
1113
atomic::{AtomicU64, Ordering},
@@ -14,7 +16,7 @@ use std::{
1416
};
1517
use tokio::sync::broadcast::{self, error::RecvError};
1618
use tokio_util::sync::{CancellationToken, WaitForCancellationFutureOwned};
17-
use tracing::{debug, debug_span, enabled, trace};
19+
use tracing::{Instrument, debug, debug_span, enabled, trace};
1820

1921
/// Either type for subscription outputs.
2022
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
@@ -232,38 +234,23 @@ impl SubscriptionTask {
232234
span.record("filter", format!("{filter:?}"));
233235
}
234236

235-
// Drain one buffered item per iteration, checking for
236-
// cancellation between each send.
237-
if let Some(item) = notif_buffer.pop_front() {
238-
let notification = SubscriptionNotification {
239-
jsonrpc: "2.0",
240-
method: "eth_subscription",
241-
params: SubscriptionParams { result: &item, subscription: id },
242-
};
243-
244-
let _guard = span.enter();
245-
tokio::select! {
246-
biased;
247-
_ = &mut ajj_cancel => {
248-
trace!("subscription cancelled by client disconnect");
249-
token.cancel();
250-
break;
251-
}
252-
_ = token.cancelled() => {
253-
trace!("subscription cancelled by user");
254-
break;
255-
}
256-
result = ajj_ctx.notify(&notification) => {
257-
if result.is_err() {
258-
trace!("channel to client closed");
259-
break;
260-
}
261-
}
237+
// NB: reserve half the capacity to avoid blocking other
238+
// usage. This is a heuristic and can be adjusted as needed.
239+
let guard = span.enter();
240+
let permit_fut = async {
241+
if !notif_buffer.is_empty() {
242+
ajj_ctx
243+
.permit_many(min(ajj_ctx.notification_capacity() / 2, notif_buffer.len()))
244+
.await
245+
} else {
246+
pending().await
262247
}
263-
continue;
264248
}
249+
.in_current_span();
250+
drop(guard);
265251

266-
// Buffer empty — wait for incoming broadcast notifications.
252+
// NB: biased select ensures we check cancellation before
253+
// processing new notifications.
267254
let _guard = span.enter();
268255
tokio::select! {
269256
biased;
@@ -276,6 +263,22 @@ impl SubscriptionTask {
276263
trace!("subscription cancelled by user");
277264
break;
278265
}
266+
permits = permit_fut => {
267+
let Some(permits) = permits else {
268+
trace!("channel to client closed");
269+
break
270+
};
271+
272+
for permit in permits {
273+
let Some(item) = notif_buffer.pop_front() else { break };
274+
let notification = SubscriptionNotification {
275+
jsonrpc: "2.0",
276+
method: "eth_subscription",
277+
params: SubscriptionParams { result: &item, subscription: id },
278+
};
279+
let _ = permit.send(&notification);
280+
}
281+
}
279282
notif_res = notifs.recv() => {
280283
let notif = match notif_res {
281284
Ok(notif) => notif,

‎crates/rpc/src/interest/subs.rs‎

Lines changed: 33 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,9 @@ use reth::{
88
};
99
use signet_node_types::Pnt;
1010
use std::{
11+
cmp::min,
1112
collections::VecDeque,
13+
future::pending,
1214
sync::{
1315
Arc, Weak,
1416
atomic::{AtomicU64, Ordering},
@@ -17,7 +19,7 @@ use std::{
1719
};
1820
use tokio::sync::broadcast::error::RecvError;
1921
use tokio_util::sync::{CancellationToken, WaitForCancellationFutureOwned};
20-
use tracing::{debug, debug_span, enabled, trace};
22+
use tracing::{Instrument, debug, debug_span, enabled, trace};
2123

2224
/// Either type for subscription outputs.
2325
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
@@ -237,38 +239,23 @@ impl SubscriptionTask {
237239
span.record("filter", format!("{filter:?}"));
238240
}
239241

240-
// Drain one buffered item per iteration, checking for
241-
// cancellation between each send.
242-
if let Some(item) = notif_buffer.pop_front() {
243-
let notification = SubscriptionNotification {
244-
jsonrpc: "2.0",
245-
method: "eth_subscription",
246-
params: SubscriptionParams { result: &item, subscription: id },
247-
};
248-
249-
let _guard = span.enter();
250-
tokio::select! {
251-
biased;
252-
_ = &mut ajj_cancel => {
253-
trace!("subscription cancelled by client disconnect");
254-
token.cancel();
255-
break;
256-
}
257-
_ = token.cancelled() => {
258-
trace!("subscription cancelled by user");
259-
break;
260-
}
261-
result = ajj_ctx.notify(&notification) => {
262-
if result.is_err() {
263-
trace!("channel to client closed");
264-
break;
265-
}
266-
}
242+
// NB: reserve half the capacity to avoid blocking other
243+
// usage. This is a heuristic and can be adjusted as needed.
244+
let guard = span.enter();
245+
let permit_fut = async {
246+
if !notif_buffer.is_empty() {
247+
ajj_ctx
248+
.permit_many(min(ajj_ctx.notification_capacity() / 2, notif_buffer.len()))
249+
.await
250+
} else {
251+
pending().await
267252
}
268-
continue;
269253
}
254+
.in_current_span();
255+
drop(guard);
270256

271-
// Buffer empty — wait for incoming notifications.
257+
// NB: biased select ensures we check cancellation before
258+
// processing new notifications.
272259
let _guard = span.enter();
273260
tokio::select! {
274261
biased;
@@ -281,6 +268,22 @@ impl SubscriptionTask {
281268
trace!("subscription cancelled by user");
282269
break;
283270
}
271+
permits = permit_fut => {
272+
let Some(permits) = permits else {
273+
trace!("channel to client closed");
274+
break
275+
};
276+
277+
for permit in permits {
278+
let Some(item) = notif_buffer.pop_front() else { break };
279+
let notification = SubscriptionNotification {
280+
jsonrpc: "2.0",
281+
method: "eth_subscription",
282+
params: SubscriptionParams { result: &item, subscription: id },
283+
};
284+
let _ = permit.send(&notification);
285+
}
286+
}
284287
notif_res = notifs.recv() => {
285288
let notif = match notif_res {
286289
Ok(notif) => notif,

0 commit comments

Comments
 (0)