Sitelet https://github.com/init4tech/builder/commit/22cd9b00f68533f489a91e6088c5d9eeae37a5bd
Skip to content

Commit 22cd9b0

Browse files
authored
feat: replace BundlePoller polling with SSE streaming (#269)
## Description Mirrors #259 (tx-poller SSE) for bundles. Replaces the 1s polling loop in `BundlePoller` with an SSE subscription to `/bundles/feed` via `BuilderTxCache::subscribe_bundles` (newly exposed by the SDK, matched by the tx-pool-webservice endpoint). The structure matches `TxPoller` exactly: - Initial `full_fetch` paginates through `stream_bundles` to seed the cache on startup. - `subscribe` opens the SSE stream and yields `CachedBundle`s in real time. - On stream error / stream end, reconnects with exponential backoff (1s → 30s cap), racing the sleep against the block-env watcher so an env change wins and triggers a fresh full refetch immediately. - On block-env change, the full_fetch runs again to cover anything missed during disconnection. - `NotOurSlot` is logged at trace level (expected when the builder is not slot-permissioned) in both the full-fetch and subscribe paths — avoids spurious warn-level noise. The public `check_bundle_cache()` wrapper is dropped; the integration test now uses `BuilderTxCache::stream_bundles` directly, matching the style of `tx_poller_test.rs`. ## Related Issue Stacked on #259. ## Testing - [x] `make fmt` passes - [x] `make clippy` passes - [x] `make test` passes - [x] `cargo doc --no-deps` passes with `-D warnings` - [x] `cargo test --features test-utils --no-run` builds all integration tests - [ ] Exercise against a live tx-pool-webservice to verify SSE reconnect and env-driven refetch behavior end-to-end
1 parent 390c4eb commit 22cd9b0

4 files changed

Lines changed: 212 additions & 74 deletions

File tree

‎src/metrics.rs‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,10 @@ const TXS_FETCHED: &str = "signet.builder.cache.txs_fetched";
3434
const TXS_FETCHED_HELP: &str = "Transactions fetched per poll cycle.";
3535

3636
const SSE_RECONNECT_ATTEMPTS: &str = "signet.builder.cache.sse_reconnect_attempts";
37-
const SSE_RECONNECT_ATTEMPTS_HELP: &str = "SSE transaction stream reconnect attempts.";
37+
const SSE_RECONNECT_ATTEMPTS_HELP: &str = "SSE stream reconnect attempts.";
38+
39+
const SSE_SUBSCRIBE_ERRORS: &str = "signet.builder.cache.sse_subscribe_errors";
40+
const SSE_SUBSCRIBE_ERRORS_HELP: &str = "SSE stream subscription failures.";
3841

3942
const BUNDLE_POLL_COUNT: &str = "signet.builder.cache.bundle_poll_count";
4043
const BUNDLE_POLL_COUNT_HELP: &str = "Bundle cache poll attempts.";
@@ -152,6 +155,7 @@ static DESCRIPTIONS: LazyLock<()> = LazyLock::new(|| {
152155
describe_counter!(TX_POLL_ERRORS, TX_POLL_ERRORS_HELP);
153156
describe_histogram!(TXS_FETCHED, TXS_FETCHED_HELP);
154157
describe_counter!(SSE_RECONNECT_ATTEMPTS, SSE_RECONNECT_ATTEMPTS_HELP);
158+
describe_counter!(SSE_SUBSCRIBE_ERRORS, SSE_SUBSCRIBE_ERRORS_HELP);
155159
describe_counter!(BUNDLE_POLL_COUNT, BUNDLE_POLL_COUNT_HELP);
156160
describe_counter!(BUNDLE_POLL_ERRORS, BUNDLE_POLL_ERRORS_HELP);
157161
describe_histogram!(BUNDLES_FETCHED, BUNDLES_FETCHED_HELP);
@@ -243,6 +247,11 @@ pub(crate) fn inc_sse_reconnect_attempts() {
243247
counter!(SSE_RECONNECT_ATTEMPTS).increment(1);
244248
}
245249

250+
/// Increment the SSE subscribe error counter.
251+
pub(crate) fn inc_sse_subscribe_errors() {
252+
counter!(SSE_SUBSCRIBE_ERRORS).increment(1);
253+
}
254+
246255
/// Increment the bundle poll attempt counter.
247256
pub(crate) fn inc_bundle_poll_count() {
248257
counter!(BUNDLE_POLL_COUNT).increment(1);

‎src/tasks/cache/bundle.rs‎

Lines changed: 196 additions & 70 deletions
Original file line numberDiff line numberDiff line change
@@ -1,68 +1,87 @@
11
//! Bundler service responsible for fetching bundles and sending them to the simulator.
2-
use crate::config::BuilderConfig;
3-
use futures_util::{TryFutureExt, TryStreamExt};
2+
use crate::{config::BuilderConfig, tasks::env::SimEnv};
3+
use futures_util::{Stream, StreamExt, TryFutureExt, TryStreamExt};
44
use init4_bin_base::perms::tx_cache::{BuilderTxCache, BuilderTxCacheError};
55
use signet_sim::{ProviderStateSource, SimItemValidity, check_bundle_tx_list};
66
use signet_tx_cache::{TxCacheError, types::CachedBundle};
7+
use std::{ops::ControlFlow, pin::Pin, time::Duration};
78
use tokio::{
8-
sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel},
9+
sync::{mpsc, watch},
910
task::JoinHandle,
10-
time::{self, Duration},
11+
time,
1112
};
12-
use tracing::{Instrument, debug_span, trace, trace_span, warn};
13+
use tracing::{Instrument, debug, debug_span, trace, warn};
1314

14-
/// Poll interval for the bundle poller in milliseconds.
15-
const POLL_INTERVAL_MS: u64 = 1000;
15+
type SseStream = Pin<Box<dyn Stream<Item = Result<CachedBundle, BuilderTxCacheError>> + Send>>;
1616

17-
/// The BundlePoller polls the tx-pool for bundles.
17+
const INITIAL_RECONNECT_BACKOFF: Duration = Duration::from_secs(1);
18+
const MAX_RECONNECT_BACKOFF: Duration = Duration::from_secs(30);
19+
20+
/// The BundlePoller fetches bundles from the tx-pool on startup and on each
21+
/// block environment change, and subscribes to an SSE stream for real-time
22+
/// delivery of new bundles in between.
1823
#[derive(Debug)]
1924
pub struct BundlePoller {
2025
/// The builder configuration values.
2126
config: &'static BuilderConfig,
22-
2327
/// Client for the tx cache.
2428
tx_cache: BuilderTxCache,
25-
26-
/// Defines the interval at which the bundler polls the tx-pool for bundles.
27-
poll_interval_ms: u64,
29+
/// Receiver for block environment updates, used to trigger refetches.
30+
envs: watch::Receiver<Option<SimEnv>>,
2831
}
2932

30-
impl Default for BundlePoller {
31-
fn default() -> Self {
32-
Self::new()
33-
}
34-
}
35-
36-
/// Implements a poller for the block builder to pull bundles from the tx-pool.
3733
impl BundlePoller {
38-
/// Creates a new BundlePoller from the provided builder config.
39-
pub fn new() -> Self {
40-
Self::new_with_poll_interval_ms(POLL_INTERVAL_MS)
41-
}
42-
43-
/// Creates a new BundlePoller from the provided builder config and with the specified poll interval in ms.
44-
pub fn new_with_poll_interval_ms(poll_interval_ms: u64) -> Self {
34+
/// Returns a new [`BundlePoller`] with the given block environment receiver.
35+
pub fn new(envs: watch::Receiver<Option<SimEnv>>) -> Self {
4536
let config = crate::config();
4637
let tx_cache = BuilderTxCache::new(config.tx_pool_url.clone(), config.oauth_token());
47-
Self { config, tx_cache, poll_interval_ms }
38+
Self { config, tx_cache, envs }
4839
}
4940

50-
/// Returns the poll duration as a [`Duration`].
51-
const fn poll_duration(&self) -> Duration {
52-
Duration::from_millis(self.poll_interval_ms)
41+
/// Pulls every bundle currently in the cache, paginating until the stream
42+
/// is exhausted. Pure fetch — no metrics, no forwarding.
43+
async fn check_bundle_cache(&self) -> Result<Vec<CachedBundle>, BuilderTxCacheError> {
44+
self.tx_cache.stream_bundles().try_collect().await
5345
}
5446

55-
/// Fetches all bundles from the tx-cache, paginating through all available pages.
56-
pub async fn check_bundle_cache(&self) -> Result<Vec<CachedBundle>, BuilderTxCacheError> {
57-
self.tx_cache.stream_bundles().try_collect().await
47+
/// Fetches all bundles from the cache and forwards each to the outbound
48+
/// channel. Records poll metrics around the fetch.
49+
async fn fetch_and_forward(&self, outbound: &mpsc::UnboundedSender<CachedBundle>) {
50+
crate::metrics::inc_bundle_poll_count();
51+
// NotOurSlot is expected whenever the builder isn't slot-permissioned;
52+
// don't bump the error counter or warn.
53+
let Ok(bundles) = self
54+
.check_bundle_cache()
55+
.inspect_err(|error| match error {
56+
BuilderTxCacheError::TxCache(TxCacheError::NotOurSlot) => {
57+
trace!("Not our slot to fetch bundles");
58+
}
59+
_ => {
60+
crate::metrics::inc_bundle_poll_errors();
61+
warn!(%error, "Failed to fetch bundles from tx-cache");
62+
}
63+
})
64+
.await
65+
else {
66+
return;
67+
};
68+
69+
crate::metrics::record_bundles_fetched(bundles.len());
70+
trace!(count = bundles.len(), "found bundles");
71+
for bundle in bundles {
72+
Self::spawn_check_bundle_nonces(bundle, outbound.clone());
73+
}
5874
}
5975

6076
/// Spawns a tokio task to check the validity of all host transactions in a
6177
/// bundle before sending it to the cache task via the outbound channel.
6278
///
6379
/// Uses [`check_bundle_tx_list`] from `signet-sim` to validate host tx nonces
6480
/// and balance against the host chain. Drops bundles that are not currently valid.
65-
fn spawn_check_bundle_nonces(bundle: CachedBundle, outbound: UnboundedSender<CachedBundle>) {
81+
fn spawn_check_bundle_nonces(
82+
bundle: CachedBundle,
83+
outbound: mpsc::UnboundedSender<CachedBundle>,
84+
) {
6685
let span = debug_span!("check_bundle_nonces", bundle_id = %bundle.id);
6786
tokio::spawn(async move {
6887
let recovered = match bundle.bundle.try_to_recovered() {
@@ -114,55 +133,162 @@ impl BundlePoller {
114133
});
115134
}
116135

117-
async fn task_future(self, outbound: UnboundedSender<CachedBundle>) {
136+
/// Returns `None` on connection failure; the caller is responsible for
137+
/// scheduling a retry. Avoids the empty-stream sentinel pattern that
138+
/// would double-log "stream ended" on a failure that never opened.
139+
async fn subscribe(&self) -> Option<SseStream> {
140+
self.tx_cache
141+
.subscribe_bundles()
142+
.await
143+
.inspect(
144+
|_| debug!(url = %self.config.tx_pool_url, "SSE bundle subscription established"),
145+
)
146+
.inspect_err(|error| match error {
147+
BuilderTxCacheError::TxCache(TxCacheError::NotOurSlot) => {
148+
trace!("Not our slot to subscribe to bundles");
149+
}
150+
_ => {
151+
crate::metrics::inc_sse_subscribe_errors();
152+
warn!(%error, "Failed to open SSE bundle subscription");
153+
}
154+
})
155+
.ok()
156+
.map(|s| Box::pin(s) as SseStream)
157+
}
158+
159+
/// Loops with exponential backoff until either a fresh SSE stream is
160+
/// established (returned as `Some`) or the outbound channel is closed
161+
/// (returned as `None`, signalling the task should shut down). Runs a
162+
/// full refetch alongside each subscribe attempt to cover items missed
163+
/// while disconnected.
164+
async fn reconnect(
165+
&mut self,
166+
outbound: &mpsc::UnboundedSender<CachedBundle>,
167+
backoff: &mut Duration,
168+
) -> Option<SseStream> {
118169
loop {
119-
let span = trace_span!("BundlePoller::loop", url = %self.config.tx_pool_url);
170+
if outbound.is_closed() {
171+
return None;
172+
}
173+
crate::metrics::inc_sse_reconnect_attempts();
174+
tokio::select! {
175+
// Biased: a block env change wins over the backoff sleep. An env
176+
// change triggers a full refetch below anyway, which supersedes the
177+
// sleep-then-reconnect path — so there's no point waiting out the
178+
// backoff.
179+
biased;
180+
_ = self.envs.changed() => {}
181+
_ = time::sleep(*backoff) => {}
182+
}
183+
*backoff = (*backoff * 2).min(MAX_RECONNECT_BACKOFF);
184+
let (_, stream) = tokio::join!(self.fetch_and_forward(outbound), self.subscribe());
185+
if let Some(stream) = stream {
186+
return Some(stream);
187+
}
188+
}
189+
}
120190

121-
// Check this here to avoid making the web request if we know
122-
// we don't need the results.
191+
/// Reconnects and swaps in the fresh stream, or returns `Break` if the
192+
/// outbound channel closed during the reconnect loop.
193+
async fn try_reconnect(
194+
&mut self,
195+
outbound: &mpsc::UnboundedSender<CachedBundle>,
196+
backoff: &mut Duration,
197+
stream: &mut SseStream,
198+
) -> ControlFlow<()> {
199+
match self.reconnect(outbound, backoff).await {
200+
Some(s) => {
201+
*stream = s;
202+
ControlFlow::Continue(())
203+
}
204+
None => ControlFlow::Break(()),
205+
}
206+
}
207+
208+
/// Returns `Break` when the outbound channel has closed and the task
209+
/// should shut down.
210+
async fn handle_sse_item(
211+
&mut self,
212+
item: Option<Result<CachedBundle, BuilderTxCacheError>>,
213+
outbound: &mpsc::UnboundedSender<CachedBundle>,
214+
backoff: &mut Duration,
215+
stream: &mut SseStream,
216+
) -> ControlFlow<()> {
217+
match item {
218+
Some(Ok(bundle)) => {
219+
*backoff = INITIAL_RECONNECT_BACKOFF;
220+
if outbound.is_closed() {
221+
trace!("No receivers left, shutting down");
222+
return ControlFlow::Break(());
223+
}
224+
Self::spawn_check_bundle_nonces(bundle, outbound.clone());
225+
ControlFlow::Continue(())
226+
}
227+
Some(Err(error)) => {
228+
warn!(%error, "SSE bundle stream interrupted, reconnecting");
229+
self.try_reconnect(outbound, backoff, stream).await
230+
}
231+
None => {
232+
warn!("SSE bundle stream ended, reconnecting");
233+
self.try_reconnect(outbound, backoff, stream).await
234+
}
235+
}
236+
}
237+
238+
async fn task_future(mut self, outbound: mpsc::UnboundedSender<CachedBundle>) {
239+
let (_, sub) = tokio::join!(self.fetch_and_forward(&outbound), self.subscribe());
240+
let mut backoff = INITIAL_RECONNECT_BACKOFF;
241+
let mut sse_stream = match sub {
242+
Some(s) => s,
243+
None => match self.reconnect(&outbound, &mut backoff).await {
244+
Some(s) => s,
245+
None => return,
246+
},
247+
};
248+
249+
loop {
123250
if outbound.is_closed() {
124-
span.in_scope(|| trace!("No receivers left, shutting down"));
251+
debug!("Outbound channel closed, shutting down");
125252
break;
126253
}
127-
128-
crate::metrics::inc_bundle_poll_count();
129-
let Ok(bundles) = self
130-
.check_bundle_cache()
131-
.inspect_err(|error| match error {
132-
BuilderTxCacheError::TxCache(TxCacheError::NotOurSlot) => {
133-
trace!("Not our slot to fetch bundles");
254+
tokio::select! {
255+
item = sse_stream.next() => {
256+
if self
257+
.handle_sse_item(item, &outbound, &mut backoff, &mut sse_stream)
258+
.await
259+
.is_break()
260+
{
261+
break;
134262
}
135-
_ => {
136-
crate::metrics::inc_bundle_poll_errors();
137-
warn!(%error, "Failed to fetch bundles from tx-cache");
263+
}
264+
res = self.envs.changed() => {
265+
if res.is_err() {
266+
debug!("Block env channel closed, shutting down");
267+
break;
138268
}
139-
})
140-
.instrument(span.clone())
141-
.await
142-
else {
143-
time::sleep(self.poll_duration()).await;
144-
continue;
145-
};
146-
147-
{
148-
let _guard = span.entered();
149-
crate::metrics::record_bundles_fetched(bundles.len());
150-
trace!(count = bundles.len(), "fetched bundles from tx-cache");
151-
for bundle in bundles {
152-
Self::spawn_check_bundle_nonces(bundle, outbound.clone());
269+
// Run the refetch under the BlockConstruction span built by
270+
// EnvTask, so its sim.ru.number / sim.host.number fields
271+
// attach to anything the refetch logs.
272+
let span = self
273+
.envs
274+
.borrow()
275+
.as_ref()
276+
.map_or_else(tracing::Span::none, |env| env.clone_span());
277+
async {
278+
debug!("Block env changed, refetching all bundles");
279+
self.fetch_and_forward(&outbound).await;
280+
}
281+
.instrument(span)
282+
.await;
153283
}
154284
}
155-
156-
time::sleep(self.poll_duration()).await;
157285
}
158286
}
159287

160-
/// Spawns a task that sends bundles it finds to its channel sender.
161-
pub fn spawn(self) -> (UnboundedReceiver<CachedBundle>, JoinHandle<()>) {
162-
let (outbound, inbound) = unbounded_channel();
163-
288+
/// Spawns the task future and returns a receiver for bundles it finds.
289+
pub fn spawn(self) -> (mpsc::UnboundedReceiver<CachedBundle>, JoinHandle<()>) {
290+
let (outbound, inbound) = mpsc::unbounded_channel();
164291
let jh = tokio::spawn(self.task_future(outbound));
165-
166292
(inbound, jh)
167293
}
168294
}

‎src/tasks/cache/system.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,7 @@ impl CacheTasks {
2727
let (tx_receiver, tx_poller) = tx_poller.spawn();
2828

2929
// Bundle Poller pulls bundles from the cache
30-
let bundle_poller = BundlePoller::new();
30+
let bundle_poller = BundlePoller::new(self.block_env.clone());
3131
let (bundle_receiver, bundle_poller) = bundle_poller.spawn();
3232

3333
// Set up the cache task

‎tests/bundle_poller_test.rs‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,15 +2,18 @@
22

33
use builder::test_utils::{setup_logging, setup_test_config};
44
use eyre::Result;
5+
use futures_util::TryStreamExt;
6+
use init4_bin_base::perms::tx_cache::BuilderTxCache;
57

68
#[tokio::test]
79
async fn test_bundle_poller_roundtrip() -> Result<()> {
810
setup_logging();
911
setup_test_config();
1012

11-
let bundle_poller = builder::tasks::cache::BundlePoller::new();
13+
let config = builder::config();
14+
let tx_cache = BuilderTxCache::new(config.tx_pool_url.clone(), config.oauth_token());
1215

13-
let _ = bundle_poller.check_bundle_cache().await?;
16+
let _bundles: Vec<_> = tx_cache.stream_bundles().try_collect().await?;
1417

1518
Ok(())
1619
}

0 commit comments

Comments
 (0)