Archive consumer: capped exponential backoff (resets after a healthy run)
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
7361a39618
commit
e998860e09
+18
-3
@@ -2,7 +2,7 @@
|
|||||||
//! into Postgres, so history survives restarts and includes messages
|
//! into Postgres, so history survives restarts and includes messages
|
||||||
//! published by any client on the bus (not just this app).
|
//! 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 async_nats::jetstream;
|
||||||
use futures::StreamExt;
|
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,
|
/// Runs forever; (re)creates the stream/consumer and retries on any failure,
|
||||||
/// so a NATS or Postgres outage never takes the chat server down.
|
/// so a NATS or Postgres outage never takes the chat server down.
|
||||||
pub async fn run_consumer(nats: async_nats::Client, pool: PgPool) {
|
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 {
|
loop {
|
||||||
|
let started = Instant::now();
|
||||||
if let Err(err) = consume(&nats, &pool).await {
|
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;
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user