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

Commit a54ea30

Browse files
Evalirclaude
andcommitted
refactor: apply fraser's PR #259 review feedback
- Move struct doc to reflect SSE behavior (was still describing the old polling implementation); drop the redundant impl-block doc. - Run initial fetch + SSE subscribe concurrently in task_future via tokio::join!, mirroring the reconnect path. - Bump "Block env changed" log from trace to debug — env changes are infrequent and worth seeing in normal debug output. - Add a sse_reconnect_attempts counter; increment once per reconnect call. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
1 parent 9bc0ff4 commit a54ea30

2 files changed

Lines changed: 19 additions & 11 deletions

File tree

‎src/metrics.rs‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ const TX_POLL_ERRORS_HELP: &str = "Transaction cache poll errors.";
3333
const TXS_FETCHED: &str = "signet.builder.cache.txs_fetched";
3434
const TXS_FETCHED_HELP: &str = "Transactions fetched per poll cycle.";
3535

36+
const SSE_RECONNECT_ATTEMPTS: &str = "signet.builder.cache.sse_reconnect_attempts";
37+
const SSE_RECONNECT_ATTEMPTS_HELP: &str = "SSE transaction stream reconnect attempts.";
38+
3639
const BUNDLE_POLL_COUNT: &str = "signet.builder.cache.bundle_poll_count";
3740
const BUNDLE_POLL_COUNT_HELP: &str = "Bundle cache poll attempts.";
3841

@@ -148,6 +151,7 @@ static DESCRIPTIONS: LazyLock<()> = LazyLock::new(|| {
148151
describe_counter!(TX_POLL_COUNT, TX_POLL_COUNT_HELP);
149152
describe_counter!(TX_POLL_ERRORS, TX_POLL_ERRORS_HELP);
150153
describe_histogram!(TXS_FETCHED, TXS_FETCHED_HELP);
154+
describe_counter!(SSE_RECONNECT_ATTEMPTS, SSE_RECONNECT_ATTEMPTS_HELP);
151155
describe_counter!(BUNDLE_POLL_COUNT, BUNDLE_POLL_COUNT_HELP);
152156
describe_counter!(BUNDLE_POLL_ERRORS, BUNDLE_POLL_ERRORS_HELP);
153157
describe_histogram!(BUNDLES_FETCHED, BUNDLES_FETCHED_HELP);
@@ -234,6 +238,11 @@ pub(crate) fn record_txs_fetched(count: usize) {
234238
histogram!(TXS_FETCHED).record(count as f64);
235239
}
236240

241+
/// Increment the SSE reconnect attempts counter.
242+
pub(crate) fn inc_sse_reconnect_attempts() {
243+
counter!(SSE_RECONNECT_ATTEMPTS).increment(1);
244+
}
245+
237246
/// Increment the bundle poll attempt counter.
238247
pub(crate) fn inc_bundle_poll_count() {
239248
counter!(BUNDLE_POLL_COUNT).increment(1);

‎src/tasks/cache/tx.rs‎

Lines changed: 10 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,9 @@ type SseStream = Pin<Box<dyn Stream<Item = Result<TxEnvelope, TxCacheError>> + S
1919
const INITIAL_RECONNECT_BACKOFF: Duration = Duration::from_secs(1);
2020
const MAX_RECONNECT_BACKOFF: Duration = Duration::from_secs(30);
2121

22-
/// Implements a poller for the block builder to pull transactions from the
23-
/// transaction pool.
22+
/// Fetches transactions from the transaction pool on startup and on each
23+
/// block environment change, and subscribes to an SSE stream for real-time
24+
/// delivery of new transactions in between.
2425
#[derive(Debug)]
2526
pub struct TxPoller {
2627
/// Config values from the Builder.
@@ -31,9 +32,6 @@ pub struct TxPoller {
3132
envs: watch::Receiver<Option<SimEnv>>,
3233
}
3334

34-
/// [`TxPoller`] fetches transactions from the transaction pool on startup
35-
/// and on each block environment change, and subscribes to an SSE stream
36-
/// for real-time delivery of new transactions in between.
3735
impl TxPoller {
3836
/// Returns a new [`TxPoller`] with the given block environment receiver.
3937
pub fn new(envs: watch::Receiver<Option<SimEnv>>) -> Self {
@@ -137,6 +135,7 @@ impl TxPoller {
137135
outbound: &mpsc::UnboundedSender<ReceivedTx>,
138136
backoff: &mut Duration,
139137
) -> SseStream {
138+
crate::metrics::inc_sse_reconnect_attempts();
140139
tokio::select! {
141140
// Biased: a block env change wins over the backoff sleep. An env
142141
// change triggers a full refetch below anyway, which supersedes the
@@ -184,11 +183,11 @@ impl TxPoller {
184183
}
185184

186185
async fn task_future(mut self, outbound: mpsc::UnboundedSender<ReceivedTx>) {
187-
// Initial full fetch of all transactions currently in the cache.
188-
self.fetch_and_dispatch(&outbound).await;
189-
190-
// Open the SSE stream for real-time delivery of new transactions.
191-
let mut sse_stream = self.subscribe().await;
186+
// Initial full fetch of all currently-cached transactions, plus SSE
187+
// subscription for real-time delivery, run concurrently — symmetric
188+
// with the reconnect path.
189+
let (_, mut sse_stream) =
190+
tokio::join!(self.fetch_and_dispatch(&outbound), self.subscribe());
192191
let mut backoff = INITIAL_RECONNECT_BACKOFF;
193192

194193
loop {
@@ -207,7 +206,7 @@ impl TxPoller {
207206
debug!("Block env channel closed, shutting down");
208207
break;
209208
}
210-
trace!("Block env changed, refetching all transactions");
209+
debug!("Block env changed, refetching all transactions");
211210
self.fetch_and_dispatch(&outbound).await;
212211
}
213212
}

0 commit comments

Comments
 (0)