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
| Metric | Description |
|---|---|
dlq_replayed_total | Requests successfully re-injected into the pipeline |
dlq_aged_out_total | Requests 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§
- DlqReplay
Scheduler - Dead-letter queue replay scheduler.
- Replay
Entry - A single entry in the
DlqReplaySchedulerinternal queue.