From e998860e0938d33ad5bee94674335ed6aae94efe Mon Sep 17 00:00:00 2001 From: Bendik Aagaard Lynghaug Date: Sun, 13 Sep 2026 12:16:29 +0200 Subject: [PATCH] Archive consumer: capped exponential backoff (resets after a healthy run) Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx --- src/server/store.rs | 21 ++++++++++++++++++--- 1 file changed, 18 insertions(+), 3 deletions(-) diff --git a/src/server/store.rs b/src/server/store.rs index c226b36..11d820e 100644 --- a/src/server/store.rs +++ b/src/server/store.rs @@ -2,7 +2,7 @@ //! into Postgres, so history survives restarts and includes messages //! published by any client on the bus (not just this app). -use std::time::Duration; +use std::time::{Duration, Instant}; use async_nats::jetstream; use futures::StreamExt; @@ -38,11 +38,26 @@ pub async fn init_schema(pool: &PgPool) -> anyhow::Result<()> { /// Runs forever; (re)creates the stream/consumer and retries on any failure, /// so a NATS or Postgres outage never takes the chat server down. pub async fn run_consumer(nats: async_nats::Client, pool: PgPool) { + const MIN_BACKOFF: Duration = Duration::from_secs(5); + const MAX_BACKOFF: Duration = Duration::from_secs(60); + let mut backoff = MIN_BACKOFF; loop { + let started = Instant::now(); if let Err(err) = consume(&nats, &pool).await { - tracing::error!("archive consumer failed: {err:#}; retrying in 5s"); + // A failure after a long healthy run is a fresh incident, not an + // escalating one - reset the backoff so we retry promptly. + if started.elapsed() >= MAX_BACKOFF { + backoff = MIN_BACKOFF; + } + tracing::error!( + "archive consumer failed after {:?}: {err:#}; retrying in {}s", + started.elapsed(), + backoff.as_secs() + ); + tokio::time::sleep(backoff).await; + // Cap the backoff so a persistent outage doesn't hammer NATS/PG. + backoff = (backoff * 2).min(MAX_BACKOFF); } - tokio::time::sleep(Duration::from_secs(5)).await; } }