Skip to main content

Module dlq_replay

Module dlq_replay 

Source
Expand description

Dead-Letter Queue Replay Scheduler

DlqReplayScheduler wraps the pipeline’s DeadLetterQueue and adds controlled re-injection of shed requests back into the pipeline with exponential backoff between each replay, optional per-session filtering, and age-based eviction.

§Design

The scheduler holds its own internal queue of ReplayEntry items. A “dropped” request from the DeadLetterQueue can be converted into a ReplayEntry and pushed here. The scheduler does not automatically drain the main DeadLetterQueue — callers feed it explicitly, which keeps the concern separation clean and lets callers control replay policy.

§Prometheus Counters

MetricDescription
dlq_replayed_totalRequests successfully re-injected into the pipeline
dlq_aged_out_totalRequests evicted by DlqReplayScheduler::age_out

§Example

use std::collections::HashMap;
use std::time::Duration;
use tokio::sync::mpsc;
use tokio_prompt_orchestrator::{PromptRequest, SessionId};
use tokio_prompt_orchestrator::enhanced::dlq_replay::{DlqReplayScheduler, ReplayEntry};

let (tx, mut rx) = mpsc::channel::<PromptRequest>(16);
let scheduler = DlqReplayScheduler::new();

let req = PromptRequest {
    session: SessionId::new("session-1"),
    request_id: "req-42".to_string(),
    input: "retry me".to_string(),
    meta: HashMap::new(),
    deadline: None,
};

scheduler.push(ReplayEntry::new(req));
scheduler.replay_all(&tx).await;

let replayed = rx.recv().await.expect("replayed request");
assert_eq!(replayed.request_id, "req-42");

Structs§

DlqReplayScheduler
Dead-letter queue replay scheduler.
ReplayEntry
A single entry in the DlqReplayScheduler internal queue.