queue: wire attempt counts into retry backoff

Every network retry hard-coded retry_delay_seconds(0), so backoff was
flat at 1.2-1.8s regardless of attempt; a multi-minute outage dead-
lettered after three rapid tries. The queue now scales the handler's
delay by 2^attempts (cap 300s) before rescheduling.
This commit is contained in:
2026-08-08 20:07:26 +08:00
parent 6849006ad7
commit aa3083792a
+21 -2
View File
@@ -79,6 +79,14 @@ fn recover_update(conn: &rusqlite::Connection) -> rusqlite::Result<()> {
Ok(()) Ok(())
} }
/// Base delay × 2^attempts (attempts = retries already done), capped at 300s.
/// Applied at the queue layer so the attempt count actually reaches the
/// backoff computation; Telegram `RetryAfter` delays get the same treatment
/// (conservatively larger wait, no API change needed).
fn scaled_retry_delay(base: f64, attempts: i32) -> f64 {
(base * 2f64.powi(attempts)).min(300.0)
}
fn ensure_schema(conn: &rusqlite::Connection) -> rusqlite::Result<()> { fn ensure_schema(conn: &rusqlite::Connection) -> rusqlite::Result<()> {
conn.execute_batch( conn.execute_batch(
"CREATE TABLE IF NOT EXISTS tasks (id TEXT PRIMARY KEY, payload TEXT NOT NULL, \ "CREATE TABLE IF NOT EXISTS tasks (id TEXT PRIMARY KEY, payload TEXT NOT NULL, \
@@ -348,12 +356,13 @@ impl QueueWorker {
self.delete_row(&row.id).await; self.delete_row(&row.id).await;
(self.dead_letter)(payload, message).await; (self.dead_letter)(payload, message).await;
} else { } else {
let delay = scaled_retry_delay(delay_seconds, row.attempts);
log::info!( log::info!(
"task {} rescheduled in {delay_seconds:.1}s (attempt {})", "task {} rescheduled in {delay:.1}s (attempt {})",
row.id, row.id,
row.attempts + 1 row.attempts + 1
); );
self.reschedule(&row.id, payload, delay_seconds, row.attempts + 1) self.reschedule(&row.id, payload, delay, row.attempts + 1)
.await; .await;
} }
} }
@@ -400,6 +409,16 @@ mod tests {
use super::*; use super::*;
use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}; use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
#[test]
fn scaled_retry_delay_scales_and_caps() {
assert_eq!(scaled_retry_delay(1.0, 0), 1.0);
assert_eq!(scaled_retry_delay(1.0, 1), 2.0);
assert_eq!(scaled_retry_delay(1.0, 2), 4.0);
assert_eq!(scaled_retry_delay(1.5, 1), 3.0);
assert_eq!(scaled_retry_delay(1.0, 10), 300.0, "capped at 300s");
assert_eq!(scaled_retry_delay(300.0, 0), 300.0);
}
async fn new_queue() -> (PersistentTaskQueue, tempfile::TempDir) { async fn new_queue() -> (PersistentTaskQueue, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap(); let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("queue.db"); let path = dir.path().join("queue.db");