From 0aeeea433f50c56e8de562ec4f1c96fc6c31b275 Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Wed, 15 Apr 2026 19:58:42 +0300 Subject: [PATCH] [firehose] jitter connection startup when we are spawning persisted / configged firehoses --- src/control/firehose.rs | 16 +++++++++++++++- src/control/mod.rs | 4 ++-- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/src/control/firehose.rs b/src/control/firehose.rs index 88a3cfe..894f37b 100644 --- a/src/control/firehose.rs +++ b/src/control/firehose.rs @@ -1,7 +1,9 @@ use std::sync::Arc; use std::sync::atomic::{AtomicUsize, Ordering}; +use std::time::Duration; use miette::{IntoDiagnostic, Result}; +use rand::RngExt; use tokio_util::sync::CancellationToken; use tracing::{error, info}; use url::Url; @@ -67,6 +69,7 @@ impl FirehoseHandle { &self, source: &FirehoseSource, shared: &FirehoseShared, + delay_startup: bool, ) -> Result<()> { use std::sync::atomic::AtomicI64; let state = &self.state; @@ -100,6 +103,17 @@ impl FirehoseHandle { let tasks = self.tasks.clone(); let token = cancel.clone(); async move { + // jitter connection start so we dont cause thundering herd problems + if delay_startup { + let jitter_ms = rand::rng().random_range(0u64..2000); + tokio::select! { + _ = tokio::time::sleep(Duration::from_millis(jitter_ms)) => {} + _ = token.cancelled() => { + info!(relay = %relay_url, "firehose ingestor cancelled"); + return; + } + } + } tokio::select! { res = ingestor.run() => { // only remove our own entry because an upsert could replace us @@ -185,7 +199,7 @@ impl FirehoseHandle { let _ = self.persisted.insert_async(url.clone()).await; - self.spawn_firehose_ingestor(&FirehoseSource { url, is_pds }, shared) + self.spawn_firehose_ingestor(&FirehoseSource { url, is_pds }, shared, false) .await?; Ok(()) diff --git a/src/control/mod.rs b/src/control/mod.rs index dff981b..2dedf61 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -409,7 +409,7 @@ impl Hydrant { ); for source in &relay_hosts { firehose - .spawn_firehose_ingestor(source, fire_shared) + .spawn_firehose_ingestor(source, fire_shared, true) .await?; } } @@ -427,7 +427,7 @@ impl Hydrant { continue; } firehose - .spawn_firehose_ingestor(source, fire_shared) + .spawn_firehose_ingestor(source, fire_shared, true) .await?; } -- 2.51.2