From 67c7b64999b8c1ccb02beda7b5add78c645a0102 Mon Sep 17 00:00:00 2001 From: jnoorchashm37 Date: Tue, 25 Aug 2026 12:47:41 -0400 Subject: [PATCH] with updated --- bin/angstrom/src/components.rs | 2 +- crates/rpc/src/orders.rs | 14 +++++++++++++- 2 files changed, 14 insertions(+), 2 deletions(-) diff --git a/bin/angstrom/src/components.rs b/bin/angstrom/src/components.rs index 0fe1092da..bfa64f460 100644 --- a/bin/angstrom/src/components.rs +++ b/bin/angstrom/src/components.rs @@ -141,7 +141,7 @@ impl StromHandles { pub fn initialize_strom_handles() -> StromHandles { let (eth_tx, eth_rx) = channel(100); let (matching_tx, matching_rx) = channel(100); - let (pool_manager_tx, _) = tokio::sync::broadcast::channel(100); + let (pool_manager_tx, _) = tokio::sync::broadcast::channel(10000); let (pool_tx, pool_rx) = reth_metrics::common::mpsc::metered_unbounded_channel("orderpool"); let (orderpool_tx, orderpool_rx) = unbounded_channel(); let (validator_tx, validator_rx) = unbounded_channel(); diff --git a/crates/rpc/src/orders.rs b/crates/rpc/src/orders.rs index 54ba88231..be5234a6b 100644 --- a/crates/rpc/src/orders.rs +++ b/crates/rpc/src/orders.rs @@ -16,6 +16,7 @@ use futures::StreamExt; use jsonrpsee::{PendingSubscriptionSink, SubscriptionMessage, core::RpcResult}; use order_pool::{OrderPoolHandle, PoolManagerUpdate}; use reth_tasks::TaskExecutor; +use tokio_stream::wrappers::errors::BroadcastStreamRecvError; use validation::order::OrderValidatorHandle; pub struct OrderApi { @@ -143,11 +144,22 @@ where .map(move |update| update.map(|value| value.filter_out_order(&kind, &filter))); self.task_executor.spawn_task(async move { - while let Some(Ok(order)) = subscription.next().await { + while let Some(update) = subscription.next().await { if sink.is_closed() { break; } + let order = match update { + Ok(order) => order, + Err(BroadcastStreamRecvError::Lagged(dropped)) => { + tracing::warn!( + dropped, + "order subscription lagged; continuing with retained updates" + ); + continue; + } + }; + if let Some(result) = order { match SubscriptionMessage::new( sink.method_name(),