Skip to main content

DlqReplayScheduler

Struct DlqReplayScheduler 

Source
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 given Duration.

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 by age_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

Source

pub fn new() -> Self

Create a new empty DlqReplayScheduler.

§Panics

This function does not panic.

Source

pub fn push(&self, entry: ReplayEntry)

Push a new ReplayEntry into the scheduler queue.

§Panics

This function does not panic.

Source

pub fn len(&self) -> usize

Return the number of entries currently waiting for replay.

§Panics

This function does not panic.

Source

pub fn is_empty(&self) -> bool

Return true if no entries are queued.

§Panics

This function does not panic.

Source

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.

Source

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.

Source

pub fn age_out(&self, max_age: Duration)

Remove and discard all entries older than max_age.

Increments dlq_aged_out_total and the aged_out_total atomic for each evicted entry.

§Panics

This function does not panic.

Trait Implementations§

Source§

impl Clone for DlqReplayScheduler

Source§

fn clone(&self) -> DlqReplayScheduler

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Default for DlqReplayScheduler

Source§

fn default() -> Self

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> FromRef<T> for T
where T: Clone,

§

fn from_ref(input: &T) -> T

Converts to this type from a reference to the input type.
§

impl<T> FutureExt for T

§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
§

impl<T> PolicyExt for T
where T: ?Sized,

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more