|
1 | 1 | //! 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}; |
4 | 4 | use init4_bin_base::perms::tx_cache::{BuilderTxCache, BuilderTxCacheError}; |
5 | 5 | use signet_sim::{ProviderStateSource, SimItemValidity, check_bundle_tx_list}; |
6 | 6 | use signet_tx_cache::{TxCacheError, types::CachedBundle}; |
| 7 | +use std::{ops::ControlFlow, pin::Pin, time::Duration}; |
7 | 8 | use tokio::{ |
8 | | - sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel}, |
| 9 | + sync::{mpsc, watch}, |
9 | 10 | task::JoinHandle, |
10 | | - time::{self, Duration}, |
| 11 | + time, |
11 | 12 | }; |
12 | | -use tracing::{Instrument, debug_span, trace, trace_span, warn}; |
| 13 | +use tracing::{Instrument, debug, debug_span, trace, warn}; |
13 | 14 |
|
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>>; |
16 | 16 |
|
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. |
18 | 23 | #[derive(Debug)] |
19 | 24 | pub struct BundlePoller { |
20 | 25 | /// The builder configuration values. |
21 | 26 | config: &'static BuilderConfig, |
22 | | - |
23 | 27 | /// Client for the tx cache. |
24 | 28 | 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>>, |
28 | 31 | } |
29 | 32 |
|
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. |
37 | 33 | 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 { |
45 | 36 | let config = crate::config(); |
46 | 37 | 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 } |
48 | 39 | } |
49 | 40 |
|
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 |
53 | 45 | } |
54 | 46 |
|
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 | + } |
58 | 74 | } |
59 | 75 |
|
60 | 76 | /// Spawns a tokio task to check the validity of all host transactions in a |
61 | 77 | /// bundle before sending it to the cache task via the outbound channel. |
62 | 78 | /// |
63 | 79 | /// Uses [`check_bundle_tx_list`] from `signet-sim` to validate host tx nonces |
64 | 80 | /// 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 | + ) { |
66 | 85 | let span = debug_span!("check_bundle_nonces", bundle_id = %bundle.id); |
67 | 86 | tokio::spawn(async move { |
68 | 87 | let recovered = match bundle.bundle.try_to_recovered() { |
@@ -114,55 +133,162 @@ impl BundlePoller { |
114 | 133 | }); |
115 | 134 | } |
116 | 135 |
|
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> { |
118 | 169 | 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 | + } |
120 | 190 |
|
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 { |
123 | 250 | if outbound.is_closed() { |
124 | | - span.in_scope(|| trace!("No receivers left, shutting down")); |
| 251 | + debug!("Outbound channel closed, shutting down"); |
125 | 252 | break; |
126 | 253 | } |
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; |
134 | 262 | } |
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; |
138 | 268 | } |
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; |
153 | 283 | } |
154 | 284 | } |
155 | | - |
156 | | - time::sleep(self.poll_duration()).await; |
157 | 285 | } |
158 | 286 | } |
159 | 287 |
|
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(); |
164 | 291 | let jh = tokio::spawn(self.task_future(outbound)); |
165 | | - |
166 | 292 | (inbound, jh) |
167 | 293 | } |
168 | 294 | } |
0 commit comments