fix: fence queue dead-letter side effects

This commit is contained in:
2026-09-24 18:33:21 +08:00
parent 5e265f2413
commit e7855b02fe
+76 -14
View File
@@ -496,8 +496,9 @@ impl QueueWorker {
Ok(value) => value, Ok(value) => value,
Err(e) => { Err(e) => {
log::error!("queue: unparseable payload for {}: {e}", row.id); log::error!("queue: unparseable payload for {}: {e}", row.id);
self.delete_row(&row.id, &row.lease_token).await; if self.delete_row(&row.id, &row.lease_token).await {
(self.dead_letter)(Value::Null, format!("invalid stored payload: {e}")).await; (self.dead_letter)(Value::Null, format!("invalid stored payload: {e}")).await;
}
return; return;
} }
}; };
@@ -546,8 +547,9 @@ impl QueueWorker {
row.id, row.id,
row.attempts + 1 row.attempts + 1
); );
self.delete_row(&row.id, &row.lease_token).await; if self.delete_row(&row.id, &row.lease_token).await {
(self.dead_letter)(payload, message).await; (self.dead_letter)(payload, message).await;
}
} else { } else {
let delay = scaled_retry_delay(delay_seconds, row.attempts); let delay = scaled_retry_delay(delay_seconds, row.attempts);
log::debug!( log::debug!(
@@ -561,11 +563,12 @@ impl QueueWorker {
} }
Err(QueueError::Permanent { message, payload }) => { Err(QueueError::Permanent { message, payload }) => {
log::error!("dead-lettering {} {fields}: {message}", row.id); log::error!("dead-lettering {} {fields}: {message}", row.id);
self.delete_row(&row.id, &row.lease_token).await; if self.delete_row(&row.id, &row.lease_token).await {
(self.dead_letter)(payload, message).await; (self.dead_letter)(payload, message).await;
} }
} }
} }
}
/// Drives the handler to completion, refreshing the row's `locked_until` /// Drives the handler to completion, refreshing the row's `locked_until`
/// every 30 s so the expiry sweep never re-leases a still-running task. /// every 30 s so the expiry sweep never re-leases a still-running task.
@@ -611,7 +614,10 @@ impl QueueWorker {
// future here stops this attempt instead of racing the // future here stops this attempt instead of racing the
// new holder through the same send. // new holder through the same send.
Ok(_) => return Err(LeaseLost), Ok(_) => return Err(LeaseLost),
Err(e) => log::error!("queue lease heartbeat failed: {e}"), Err(e) => {
log::error!("queue lease heartbeat failed: {e}; abandoning attempt");
return Err(LeaseLost);
}
} }
} }
} }
@@ -625,19 +631,18 @@ impl QueueWorker {
/// (a busy/contended DB is the usual cause and clears), and if the DB still /// (a busy/contended DB is the usual cause and clears), and if the DB still
/// refuses, the row is marked `done` — a status neither the lease query /// refuses, the row is marked `done` — a status neither the lease query
/// (`pending`) nor the sweep (`in_progress`) looks at — so a task that /// (`pending`) nor the sweep (`in_progress`) looks at — so a task that
/// already ran can never be re-leased. Both writes failing is logged at /// already ran can never be re-leased. Returns `false` when the token is
/// error level with the row id, since that is the one case where a /// no longer ours; callers must not run dead-letter side effects then.
/// duplicate send stays possible. async fn delete_row(&self, id: &str, lease_token: &str) -> bool {
async fn delete_row(&self, id: &str, lease_token: &str) {
for attempt in 0..TERMINAL_WRITE_ATTEMPTS { for attempt in 0..TERMINAL_WRITE_ATTEMPTS {
match self.try_delete_row(id, lease_token).await { match self.try_delete_row(id, lease_token).await {
Ok(true) => return, Ok(true) => return true,
// The row is not ours any more (re-leased while we worked): // The row is not ours any more (re-leased while we worked):
// leaving it alone *is* the clean outcome — retrying or // leaving it alone *is* the clean outcome — retrying or
// tombstoning here would erase the new holder's work. // tombstoning here would erase the new holder's work.
Ok(false) => { Ok(false) => {
log::warn!("queue: row {id} was re-leased; not deleting it"); log::warn!("queue: row {id} was re-leased; not deleting it");
return; return false;
} }
Err(e) => { Err(e) => {
log::error!("queue delete failed (attempt {}): {e}", attempt + 1); log::error!("queue delete failed (attempt {}): {e}", attempt + 1);
@@ -646,12 +651,21 @@ impl QueueWorker {
} }
} }
match mark_done(&self.pool, id, lease_token).await { match mark_done(&self.pool, id, lease_token).await {
Ok(true) => log::warn!("queue: row {id} marked done instead of deleted"), Ok(true) => {
Ok(false) => log::warn!("queue: row {id} was re-leased; nothing to tombstone"), log::warn!("queue: row {id} marked done instead of deleted");
Err(e) => log::error!( true
}
Ok(false) => {
log::warn!("queue: row {id} was re-leased; nothing to tombstone");
false
}
Err(e) => {
log::error!(
"queue: row {id} could not be deleted or marked done ({e}); \ "queue: row {id} could not be deleted or marked done ({e}); \
the expiry sweep may run this finished task again" the expiry sweep may run this finished task again"
), );
false
}
} }
} }
@@ -871,6 +885,54 @@ mod tests {
queue.stop().await; queue.stop().await;
} }
#[tokio::test]
async fn a_re_leased_row_cannot_dead_letter_the_old_attempt() {
let (queue, _dir) = new_queue().await;
queue
.enqueue(serde_json::json!({"a": 1}), now_f64())
.await
.unwrap();
let id: String = queue
.pool
.with_conn(|conn| conn.query_row("SELECT id FROM tasks", [], |r| r.get(0)))
.await
.unwrap();
set_lease(&queue, &id, "old-holder").await;
let dead_calls = Arc::new(AtomicUsize::new(0));
let worker = QueueWorker {
pool: std::sync::Arc::clone(&queue.pool),
notify: Arc::new(Notify::new()),
stop: Arc::new(AtomicBool::new(false)),
handler: Arc::new(|_payload| {
Box::pin(async {
Err(QueueError::Permanent {
message: "stale failure".into(),
payload: serde_json::json!({"a": 1}),
})
})
}),
dead_letter: Arc::new({
let dead_calls = Arc::clone(&dead_calls);
move |_payload, _message| {
let dead_calls = Arc::clone(&dead_calls);
Box::pin(async move {
dead_calls.fetch_add(1, AtomicOrdering::SeqCst);
})
}
}),
};
set_lease(&queue, &id, "new-holder").await;
worker
.process(LeasedRow {
id,
payload: "{\"a\":1}".into(),
attempts: 0,
lease_token: "old-holder".into(),
})
.await;
assert_eq!(dead_calls.load(AtomicOrdering::SeqCst), 0);
}
#[tokio::test] #[tokio::test]
async fn pending_backlog_counts_only_unleased_rows() { async fn pending_backlog_counts_only_unleased_rows() {
let (queue, _dir) = new_queue().await; let (queue, _dir) = new_queue().await;