From dd9c329874c5c11f98cd9acc5488c1e5d792274d Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 28 Jun 2026 00:35:10 +0300 Subject: [PATCH] [ingest] loop wait_for_allow to prevent enabled notification bypass --- src/ingest/firehose.rs | 26 +++++++++++++++++--------- 1 file changed, 17 insertions(+), 9 deletions(-) diff --git a/src/ingest/firehose.rs b/src/ingest/firehose.rs index f332eaa..530a502 100644 --- a/src/ingest/firehose.rs +++ b/src/ingest/firehose.rs @@ -301,18 +301,26 @@ impl FirehoseIngestor { let accounts = self.state.db.get_count(&count_key).await; #[cfg(feature = "firehose-diagnostics")] let throttle_started = Instant::now(); - tokio::select! { - _ = self.throttle.wait_for_allow(accounts, &tier) => { - #[cfg(feature = "firehose-diagnostics")] - self.stats.record_throttle_wait(throttle_started.elapsed()); - } - _ = self.enabled.changed() => { - if !*self.enabled.borrow() { - info!("firehose disabled, disconnecting"); - break Ok(()); + let mut disabled = false; + loop { + tokio::select! { + _ = self.throttle.wait_for_allow(accounts, &tier) => { + #[cfg(feature = "firehose-diagnostics")] + self.stats.record_throttle_wait(throttle_started.elapsed()); + break; + } + res = self.enabled.changed() => { + if res.is_err() || !*self.enabled.borrow() { + disabled = true; + break; + } } } } + if disabled { + info!("firehose disabled, disconnecting"); + break Ok(()); + } } self.handle_message(msg).await; }, -- 2.51.2