pub struct DlqReplayScheduler {
pub replayed_total: Arc<AtomicU64>,
pub aged_out_total: Arc<AtomicU64>,
/* private fields */
}Expand description
Dead-letter queue replay scheduler.
Wraps the pipeline’s dead-letter concept and adds:
replay_all— re-inject every eligible entry with exponential backoff between sends.replay_by_session— replay only entries belonging to a specific session.age_out— drop entries older than a givenDuration.
DlqReplayScheduler is Clone + Send + Sync. All clones share the same
internal queue.
§Prometheus counters
dlq_replayed_total— incremented on each successful re-injection.dlq_aged_out_total— incremented on each eviction byage_out.
Fields§
§replayed_total: Arc<AtomicU64>Total entries replayed across the lifetime of this scheduler.
aged_out_total: Arc<AtomicU64>Total entries aged out across the lifetime of this scheduler.
Implementations§
Source§impl DlqReplayScheduler
impl DlqReplayScheduler
Sourcepub fn push(&self, entry: ReplayEntry)
pub fn push(&self, entry: ReplayEntry)
Sourcepub async fn replay_all(&self, sender: &Sender<PromptRequest>)
pub async fn replay_all(&self, sender: &Sender<PromptRequest>)
Re-inject all eligible entries back into the pipeline.
An entry is eligible if ReplayEntry::is_ready returns true
(i.e. its exponential-backoff window has elapsed).
Between consecutive sends the scheduler applies exponential backoff
(from the entry’s own attempts counter), sleeping the current task
for the computed delay. Entries that are sent successfully are removed
from the queue. Entries whose Sender is closed are left in the queue
and a warn! is emitted.
§Panics
This function does not panic.
Sourcepub async fn replay_by_session(
&self,
session_id: &SessionId,
sender: &Sender<PromptRequest>,
)
pub async fn replay_by_session( &self, session_id: &SessionId, sender: &Sender<PromptRequest>, )
Re-inject all entries belonging to session_id back into the pipeline.
Behaviour is identical to replay_all but only
entries whose request.session matches session_id are attempted.
§Panics
This function does not panic.
Trait Implementations§
Source§impl Clone for DlqReplayScheduler
impl Clone for DlqReplayScheduler
Source§fn clone(&self) -> DlqReplayScheduler
fn clone(&self) -> DlqReplayScheduler
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreAuto Trait Implementations§
impl Freeze for DlqReplayScheduler
impl RefUnwindSafe for DlqReplayScheduler
impl Send for DlqReplayScheduler
impl Sync for DlqReplayScheduler
impl Unpin for DlqReplayScheduler
impl UnsafeUnpin for DlqReplayScheduler
impl UnwindSafe for DlqReplayScheduler
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request