Skip to main content

ResponseBuffer

Struct ResponseBuffer 

Source
pub struct ResponseBuffer { /* private fields */ }
Expand description

Thread-safe response buffer with pause/resume and manual release/reject.

Hold duration is dynamically adjustable via [set_hold_duration] to support adaptive SLA — low-scoring agents can have their hold increased to give operators more review time.

Implementations§

Source§

impl ResponseBuffer

Source

pub fn new(hold_duration: Duration) -> Self

Create a new buffer with the given hold duration.

The hold duration initializes the response SLA for deadline-based release via [push_with_deadline]. The dashboard can later override the SLA via [set_response_sla].

Duration::ZERO means responses drain immediately (pass-through mode). Any non-zero duration is used as-is — callers control the review window.

Source

pub fn hold_duration(&self) -> Duration

The current effective hold duration.

Source

pub fn base_hold_duration(&self) -> Duration

The base hold duration as configured at startup.

Source

pub fn set_hold_duration(&self, duration: Duration)

Dynamically update the effective hold duration.

Used by the adaptive SLA system to slow down or speed up the buffer based on agent scores.

Source

pub fn set_response_sla(&self, sla: Duration)

Set the response SLA (used by [push_with_deadline] to compute release time).

Duration::ZERO disables the deadline (pass-through). Any non-zero value is used as-is.

Source

pub fn response_sla(&self) -> Option<Duration>

Current response SLA duration, or None if not configured.

Source

pub async fn push(&self, entry: BufferedResponse)

Add a completed response to the buffer.

Source

pub async fn push_with_deadline( &self, entry: BufferedResponse, task_received: Instant, )

Add a completed response, computing release_at from the SLA deadline.

release_at = task_received + response_sla

The buffer holds the response for the full SLA duration. The rainfall animation on the dashboard runs for exactly this duration — card reaching the bottom = SLA expired = auto-release.

The orchestrator already sets round timeouts that respect all agent SLAs, so no reserve is needed at the buffer level.

Falls back to the caller-provided release_at if no SLA is configured.

Source

pub async fn drain_ready(&self) -> Vec<BufferedResponse>

Drain entries whose release_at has passed, unless paused.

Returns the drained entries (caller is responsible for publishing + ack). When paused, always returns an empty vec.

Source

pub async fn list(&self) -> Vec<BufferEntrySummary>

Return summaries of all buffered entries (for the dashboard UI).

Source

pub async fn release(&self, id: &str) -> Option<BufferedResponse>

Manually release a specific entry by ID, regardless of hold duration.

Returns the entry if found, or None if not in the buffer.

Source

pub async fn reject(&self, id: &str) -> Option<BufferedResponse>

Reject (discard) a specific entry by ID.

Returns the entry if found (caller should ack the message without publishing), or None if not in the buffer.

Source

pub async fn drain_stale(&self, current_job_id: &str) -> Vec<BufferedResponse>

Remove all entries whose job_id does NOT match the given current job.

Returns the removed entries (caller is responsible for ack-ing them). This prevents stale responses from previous deliberations from lingering in the operator review queue.

Source

pub fn pause(&self)

Pause the buffer: drain_ready() will return empty and the worker should also stop pulling new NATS tasks.

Source

pub fn resume(&self)

Resume the buffer: drain_ready() resumes normal operation.

Source

pub fn is_paused(&self) -> bool

Whether the buffer is currently paused.

Source

pub fn set_auto_approve(&self, enabled: bool)

Enable or disable auto-approve mode.

When enabled, entries whose agent divergence is below the configured threshold are auto-released immediately instead of waiting for the hold timer or manual operator action.

Source

pub fn is_auto_approve(&self) -> bool

Whether auto-approve mode is currently enabled.

Source

pub fn set_auto_approve_threshold(&self, threshold: f32)

Set the divergence threshold for auto-approve (0.0 to 1.0).

Values are clamped to [0.0, 1.0]. Stored internally as thousandths.

Source

pub fn auto_approve_threshold(&self) -> f32

Current auto-approve divergence threshold (0.0 to 1.0).

Source

pub async fn auto_release_if_eligible(&self, divergence: Option<f32>) -> usize

When auto-approve is enabled and the agent’s divergence is at or below the threshold, mark all non-stopped pending entries for immediate release. The gate is strictly div > threshold → block, so a threshold of 1.0 (the default) releases every entry because compute_divergence clamps to [0.0, 1.0].

When divergence is None (no scores yet), the operator’s explicit opt-in to auto-approve takes precedence — entries are released. The threshold only blocks release when we have divergence data strictly exceeding the threshold.

Returns the number of entries marked for auto-release.

Source

pub async fn len(&self) -> usize

Number of entries currently in the buffer.

Source

pub async fn is_empty(&self) -> bool

Whether the buffer is empty.

Source

pub async fn get_detail(&self, id: &str) -> Option<BufferEntryDetail>

Return the full detail of a specific buffer entry, including the deserialized response payload (for operator inspection/editing).

Source

pub async fn update_payload(&self, id: &str, new_payload: Vec<u8>) -> bool

Update the payload of a specific buffer entry (operator edit).

Returns true if the entry was found and updated, false otherwise.

Source

pub async fn add_comment( &self, id: &str, annotation: OperatorAnnotation, ) -> bool

Add an operator comment to a buffer entry without modifying the payload.

Returns true if the entry was found and annotated, false otherwise.

Source

pub async fn mark_for_release(&self, id: &str) -> bool

Mark a buffer entry for immediate release by setting its release_at to now.

The entry stays in the buffer — the worker’s drain_buffer() loop will pick it up on the next cycle (≤500ms) and handle the NATS publish. This avoids needing a NATS client in the status server.

Note: The stopped flag is preserved. Stopped entries must be explicitly unstopped (or use [force_release]) before they can drain.

Returns true if the entry was found and marked, false otherwise.

Source

pub async fn force_release(&self, id: &str) -> bool

Atomically unstop and mark a buffer entry for immediate release.

Combines [unstop] + [mark_for_release] in a single lock acquisition, eliminating the race window where drain_ready() could observe the entry as unstopped with a stale (already-passed) release_at.

Returns true if the entry was found, false otherwise.

Source

pub async fn stop(&self, id: &str) -> bool

Stop (reversibly reject) a buffer entry.

Stopped entries remain in the buffer but are skipped by [drain_ready] — they won’t auto-release. The operator can later call [unstop] to make the entry eligible for release again.

Returns true if the entry was found and stopped, false otherwise.

Source

pub async fn unstop(&self, id: &str) -> bool

Un-stop a previously stopped buffer entry.

The entry becomes eligible for [drain_ready] again. If its release_at has already passed, it will drain on the next cycle.

Returns true if the entry was found and un-stopped, false otherwise.

Source

pub async fn update_payload_with_annotation( &self, id: &str, new_payload: Vec<u8>, annotation: OperatorAnnotation, ) -> bool

Update the payload of a buffer entry AND record an edit annotation.

Marks the entry as edited = true and appends the annotation. Returns true if the entry was found and updated, false otherwise.

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
Source§

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

Source§

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

Source§

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

Source§

fn exact_from(value: T) -> U

Source§

impl<T, U> ExactInto<U> for T
where U: ExactFrom<T>,

Source§

fn exact_into(self) -> U

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

Source§

impl<T> Instrument for T

Source§

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

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

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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T, U> OverflowingInto<U> for T
where U: OverflowingFrom<T>,

Source§

impl<D> OwoColorize for D

Source§

fn fg<C>(&self) -> FgColorDisplay<'_, C, Self>
where C: Color,

Set the foreground color generically Read more
Source§

fn bg<C>(&self) -> BgColorDisplay<'_, C, Self>
where C: Color,

Set the background color generically. Read more
Source§

fn black(&self) -> FgColorDisplay<'_, Black, Self>

Change the foreground color to black
Source§

fn on_black(&self) -> BgColorDisplay<'_, Black, Self>

Change the background color to black
Source§

fn red(&self) -> FgColorDisplay<'_, Red, Self>

Change the foreground color to red
Source§

fn on_red(&self) -> BgColorDisplay<'_, Red, Self>

Change the background color to red
Source§

fn green(&self) -> FgColorDisplay<'_, Green, Self>

Change the foreground color to green
Source§

fn on_green(&self) -> BgColorDisplay<'_, Green, Self>

Change the background color to green
Source§

fn yellow(&self) -> FgColorDisplay<'_, Yellow, Self>

Change the foreground color to yellow
Source§

fn on_yellow(&self) -> BgColorDisplay<'_, Yellow, Self>

Change the background color to yellow
Source§

fn blue(&self) -> FgColorDisplay<'_, Blue, Self>

Change the foreground color to blue
Source§

fn on_blue(&self) -> BgColorDisplay<'_, Blue, Self>

Change the background color to blue
Source§

fn magenta(&self) -> FgColorDisplay<'_, Magenta, Self>

Change the foreground color to magenta
Source§

fn on_magenta(&self) -> BgColorDisplay<'_, Magenta, Self>

Change the background color to magenta
Source§

fn purple(&self) -> FgColorDisplay<'_, Magenta, Self>

Change the foreground color to purple
Source§

fn on_purple(&self) -> BgColorDisplay<'_, Magenta, Self>

Change the background color to purple
Source§

fn cyan(&self) -> FgColorDisplay<'_, Cyan, Self>

Change the foreground color to cyan
Source§

fn on_cyan(&self) -> BgColorDisplay<'_, Cyan, Self>

Change the background color to cyan
Source§

fn white(&self) -> FgColorDisplay<'_, White, Self>

Change the foreground color to white
Source§

fn on_white(&self) -> BgColorDisplay<'_, White, Self>

Change the background color to white
Source§

fn default_color(&self) -> FgColorDisplay<'_, Default, Self>

Change the foreground color to the terminal default
Source§

fn on_default_color(&self) -> BgColorDisplay<'_, Default, Self>

Change the background color to the terminal default
Source§

fn bright_black(&self) -> FgColorDisplay<'_, BrightBlack, Self>

Change the foreground color to bright black
Source§

fn on_bright_black(&self) -> BgColorDisplay<'_, BrightBlack, Self>

Change the background color to bright black
Source§

fn bright_red(&self) -> FgColorDisplay<'_, BrightRed, Self>

Change the foreground color to bright red
Source§

fn on_bright_red(&self) -> BgColorDisplay<'_, BrightRed, Self>

Change the background color to bright red
Source§

fn bright_green(&self) -> FgColorDisplay<'_, BrightGreen, Self>

Change the foreground color to bright green
Source§

fn on_bright_green(&self) -> BgColorDisplay<'_, BrightGreen, Self>

Change the background color to bright green
Source§

fn bright_yellow(&self) -> FgColorDisplay<'_, BrightYellow, Self>

Change the foreground color to bright yellow
Source§

fn on_bright_yellow(&self) -> BgColorDisplay<'_, BrightYellow, Self>

Change the background color to bright yellow
Source§

fn bright_blue(&self) -> FgColorDisplay<'_, BrightBlue, Self>

Change the foreground color to bright blue
Source§

fn on_bright_blue(&self) -> BgColorDisplay<'_, BrightBlue, Self>

Change the background color to bright blue
Source§

fn bright_magenta(&self) -> FgColorDisplay<'_, BrightMagenta, Self>

Change the foreground color to bright magenta
Source§

fn on_bright_magenta(&self) -> BgColorDisplay<'_, BrightMagenta, Self>

Change the background color to bright magenta
Source§

fn bright_purple(&self) -> FgColorDisplay<'_, BrightMagenta, Self>

Change the foreground color to bright purple
Source§

fn on_bright_purple(&self) -> BgColorDisplay<'_, BrightMagenta, Self>

Change the background color to bright purple
Source§

fn bright_cyan(&self) -> FgColorDisplay<'_, BrightCyan, Self>

Change the foreground color to bright cyan
Source§

fn on_bright_cyan(&self) -> BgColorDisplay<'_, BrightCyan, Self>

Change the background color to bright cyan
Source§

fn bright_white(&self) -> FgColorDisplay<'_, BrightWhite, Self>

Change the foreground color to bright white
Source§

fn on_bright_white(&self) -> BgColorDisplay<'_, BrightWhite, Self>

Change the background color to bright white
Source§

fn bold(&self) -> BoldDisplay<'_, Self>

Make the text bold
Source§

fn dimmed(&self) -> DimDisplay<'_, Self>

Make the text dim
Source§

fn italic(&self) -> ItalicDisplay<'_, Self>

Make the text italicized
Source§

fn underline(&self) -> UnderlineDisplay<'_, Self>

Make the text underlined
Make the text blink
Make the text blink (but fast!)
Source§

fn reversed(&self) -> ReversedDisplay<'_, Self>

Swap the foreground and background colors
Source§

fn hidden(&self) -> HiddenDisplay<'_, Self>

Hide the text
Source§

fn strikethrough(&self) -> StrikeThroughDisplay<'_, Self>

Cross out the text
Source§

fn color<Color>(&self, color: Color) -> FgDynColorDisplay<'_, Color, Self>
where Color: DynColor,

Set the foreground color at runtime. Only use if you do not know which color will be used at compile-time. If the color is constant, use either OwoColorize::fg or a color-specific method, such as OwoColorize::green, Read more
Source§

fn on_color<Color>(&self, color: Color) -> BgDynColorDisplay<'_, Color, Self>
where Color: DynColor,

Set the background color at runtime. Only use if you do not know what color to use at compile-time. If the color is constant, use either OwoColorize::bg or a color-specific method, such as OwoColorize::on_yellow, Read more
Source§

fn fg_rgb<const R: u8, const G: u8, const B: u8>( &self, ) -> FgColorDisplay<'_, CustomColor<R, G, B>, Self>

Set the foreground color to a specific RGB value.
Source§

fn bg_rgb<const R: u8, const G: u8, const B: u8>( &self, ) -> BgColorDisplay<'_, CustomColor<R, G, B>, Self>

Set the background color to a specific RGB value.
Source§

fn truecolor(&self, r: u8, g: u8, b: u8) -> FgDynColorDisplay<'_, Rgb, Self>

Sets the foreground color to an RGB value.
Source§

fn on_truecolor(&self, r: u8, g: u8, b: u8) -> BgDynColorDisplay<'_, Rgb, Self>

Sets the background color to an RGB value.
Source§

fn style(&self, style: Style) -> Styled<&Self>

Apply a runtime-determined style
Source§

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

Source§

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
Source§

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
Source§

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

Source§

impl<T, U> RoundingInto<U> for T
where U: RoundingFrom<T>,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> SaturatingInto<U> for T
where U: SaturatingFrom<T>,

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.
Source§

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

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

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
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

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

impl<T, U> WrappingInto<U> for T
where U: WrappingFrom<T>,

Source§

fn wrapping_into(self) -> U