MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / process_retry_queue

Function process_retry_queue

nodedb/src/event/consumer.rs:397–424  ·  view source on GitHub ↗

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>,
)

Source from the content-addressed store, hash-verified

395
396/// Process the retry queue: DLQ exhausted entries and retry ready ones.
397async 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)]
427mod tests {

Callers 1

consumer_loopFunction · 0.85

Calls 6

retry_singleFunction · 0.85
lockMethod · 0.80
drain_dueMethod · 0.45
is_emptyMethod · 0.45
enqueueMethod · 0.45
cloneMethod · 0.45

Tested by

no test coverage detected