Sitelet https://github.com/init4tech/builder/pull/259/commits/6dfe101ca8c96d0bfb38ba6576ea2184ec3d12a7
Skip to content
Prev Previous commit
Next Next commit
refactor: hold reconnect backoff on TxPoller
Move the SSE reconnect backoff from a `&mut Duration` plumbed
through `reconnect` and `handle_sse_item` to a field on `TxPoller`.
The state was already specific to the running task; carrying it
on self drops two parameters and one local in `task_future`.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
  • Loading branch information
Evalir and claude committed Apr 29, 2026
commit 6dfe101ca8c96d0bfb38ba6576ea2184ec3d12a7
26 changes: 12 additions & 14 deletions src/tasks/cache/tx.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,14 +30,18 @@ pub struct TxPoller {
tx_cache: TxCache,
/// Receiver for block environment updates, used to trigger refetches.
envs: watch::Receiver<Option<SimEnv>>,
/// SSE reconnect backoff. Doubles on each reconnect (capped at
/// `MAX_RECONNECT_BACKOFF`) and resets to `INITIAL_RECONNECT_BACKOFF`
/// on each successfully received tx.
backoff: Duration,
}

impl TxPoller {
/// Returns a new [`TxPoller`] with the given block environment receiver.
pub fn new(envs: watch::Receiver<Option<SimEnv>>) -> Self {
let config = crate::config();
let tx_cache = TxCache::new(config.tx_pool_url.clone());
Self { config, tx_cache, envs }
Self { config, tx_cache, envs, backoff: INITIAL_RECONNECT_BACKOFF }
}

/// Spawn a tokio task to check the nonce of a transaction before sending
Expand Down Expand Up @@ -130,11 +134,7 @@ impl TxPoller {

/// Reconnects the SSE stream with backoff. Performs a full refetch to
/// cover any items missed while disconnected.
async fn reconnect(
&mut self,
outbound: &mpsc::UnboundedSender<ReceivedTx>,
backoff: &mut Duration,
) -> SseStream {
async fn reconnect(&mut self, outbound: &mpsc::UnboundedSender<ReceivedTx>) -> SseStream {
crate::metrics::inc_sse_reconnect_attempts();
tokio::select! {
// Biased: a block env change wins over the backoff sleep. An env
Expand All @@ -143,9 +143,9 @@ impl TxPoller {
// backoff.
biased;
_ = self.envs.changed() => {}
_ = time::sleep(*backoff) => {}
_ = time::sleep(self.backoff) => {}
}
*backoff = (*backoff * 2).min(MAX_RECONNECT_BACKOFF);
self.backoff = (self.backoff * 2).min(MAX_RECONNECT_BACKOFF);
let (_, stream) = tokio::join!(self.fetch_and_dispatch(outbound), self.subscribe());
stream
}
Expand All @@ -158,12 +158,11 @@ impl TxPoller {
&mut self,
item: Option<Result<TxEnvelope, TxCacheError>>,
outbound: &mpsc::UnboundedSender<ReceivedTx>,
backoff: &mut Duration,
stream: &mut SseStream,
) -> ControlFlow<()> {
match item {
Some(Ok(tx)) => {
*backoff = INITIAL_RECONNECT_BACKOFF;
self.backoff = INITIAL_RECONNECT_BACKOFF;
if outbound.is_closed() {
trace!("No receivers left, shutting down");
return ControlFlow::Break(());
Expand All @@ -172,11 +171,11 @@ impl TxPoller {
}
Some(Err(error)) => {
warn!(%error, "SSE transaction stream error, reconnecting");
*stream = self.reconnect(outbound, backoff).await;
*stream = self.reconnect(outbound).await;
}
None => {
warn!("SSE transaction stream ended, reconnecting");
*stream = self.reconnect(outbound, backoff).await;
*stream = self.reconnect(outbound).await;
}
}
ControlFlow::Continue(())
Expand All @@ -188,13 +187,12 @@ impl TxPoller {
// with the reconnect path.
let (_, mut sse_stream) =
tokio::join!(self.fetch_and_dispatch(&outbound), self.subscribe());
let mut backoff = INITIAL_RECONNECT_BACKOFF;

loop {
tokio::select! {
item = sse_stream.next() => {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

i think this arm can be collapsed a bit by adding stream combinators that transform the sse_stream into a future that only ever returns on ControlFlow::Break?

sketch:

let stream = sse_stream
   .then(|item| self.handle_sse_item(...))
   .skip_while(|flow| flow.is_continue)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

passing the &mut sse_stream as an arg break this tho

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

but that's a bit of a code smell anyway, so it might be worth re-flowing to remove it by adding more information in the ControlFLow

if self
.handle_sse_item(item, &outbound, &mut backoff, &mut sse_stream)
.handle_sse_item(item, &outbound, &mut sse_stream)
.await
.is_break()
{
Expand Down