Process the retry queue: DLQ exhausted entries and retry ready ones.
(
retry_queue: &mut TriggerRetryQueue,
trigger_dlq: &Arc<std::sync::Mutex<TriggerDlq>>,
shared_state: &Arc<SharedState>,
)
| 395 | |
| 396 | /// Process the retry queue: DLQ exhausted entries and retry ready ones. |
| 397 | async fn process_retry_queue( |
| 398 | retry_queue: &mut TriggerRetryQueue, |
| 399 | trigger_dlq: &Arc<std::sync::Mutex<TriggerDlq>>, |
| 400 | shared_state: &Arc<SharedState>, |
| 401 | ) { |
| 402 | let (ready, exhausted) = retry_queue.drain_due(); |
| 403 | if !exhausted.is_empty() { |
| 404 | let mut dlq = trigger_dlq.lock().unwrap_or_else(|p| p.into_inner()); |
| 405 | for entry in &exhausted { |
| 406 | let _ = dlq.enqueue(super::trigger::dlq::DlqEnqueueParams { |
| 407 | tenant_id: entry.tenant_id, |
| 408 | source_collection: entry.collection.clone(), |
| 409 | row_id: entry.row_id.clone(), |
| 410 | operation: entry.operation.clone(), |
| 411 | trigger_name: entry.trigger_name.clone(), |
| 412 | error: entry.last_error.clone(), |
| 413 | retry_count: entry.attempts, |
| 414 | source_lsn: entry.source_lsn, |
| 415 | source_sequence: entry.source_sequence, |
| 416 | }); |
| 417 | } |
| 418 | // dlq MutexGuard dropped before any await. |
| 419 | } |
| 420 | |
| 421 | for entry in ready { |
| 422 | super::trigger::dispatcher::retry_single(&entry, shared_state, retry_queue).await; |
| 423 | } |
| 424 | } |
| 425 | |
| 426 | #[cfg(test)] |
| 427 | mod tests { |
no test coverage detected