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
impl ResponseBuffer
Sourcepub fn new(hold_duration: Duration) -> Self
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.
Sourcepub fn hold_duration(&self) -> Duration
pub fn hold_duration(&self) -> Duration
The current effective hold duration.
Sourcepub fn base_hold_duration(&self) -> Duration
pub fn base_hold_duration(&self) -> Duration
The base hold duration as configured at startup.
Sourcepub fn set_hold_duration(&self, duration: Duration)
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.
Sourcepub fn set_response_sla(&self, sla: Duration)
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.
Sourcepub fn response_sla(&self) -> Option<Duration>
pub fn response_sla(&self) -> Option<Duration>
Current response SLA duration, or None if not configured.
Sourcepub async fn push(&self, entry: BufferedResponse)
pub async fn push(&self, entry: BufferedResponse)
Add a completed response to the buffer.
Sourcepub async fn push_with_deadline(
&self,
entry: BufferedResponse,
task_received: Instant,
)
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.
Sourcepub async fn drain_ready(&self) -> Vec<BufferedResponse>
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.
Sourcepub async fn list(&self) -> Vec<BufferEntrySummary>
pub async fn list(&self) -> Vec<BufferEntrySummary>
Return summaries of all buffered entries (for the dashboard UI).
Sourcepub async fn release(&self, id: &str) -> Option<BufferedResponse>
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.
Sourcepub async fn reject(&self, id: &str) -> Option<BufferedResponse>
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.
Sourcepub async fn drain_stale(&self, current_job_id: &str) -> Vec<BufferedResponse>
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.
Sourcepub fn pause(&self)
pub fn pause(&self)
Pause the buffer: drain_ready() will return empty and the worker
should also stop pulling new NATS tasks.
Sourcepub fn set_auto_approve(&self, enabled: bool)
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.
Sourcepub fn is_auto_approve(&self) -> bool
pub fn is_auto_approve(&self) -> bool
Whether auto-approve mode is currently enabled.
Sourcepub fn set_auto_approve_threshold(&self, threshold: f32)
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.
Sourcepub fn auto_approve_threshold(&self) -> f32
pub fn auto_approve_threshold(&self) -> f32
Current auto-approve divergence threshold (0.0 to 1.0).
Sourcepub async fn auto_release_if_eligible(&self, divergence: Option<f32>) -> usize
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.
Sourcepub async fn get_detail(&self, id: &str) -> Option<BufferEntryDetail>
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).
Sourcepub async fn update_payload(&self, id: &str, new_payload: Vec<u8>) -> bool
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.
Sourcepub async fn add_comment(
&self,
id: &str,
annotation: OperatorAnnotation,
) -> bool
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.
Sourcepub async fn mark_for_release(&self, id: &str) -> bool
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.
Sourcepub async fn force_release(&self, id: &str) -> bool
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.
Sourcepub async fn stop(&self, id: &str) -> bool
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.
Sourcepub async fn unstop(&self, id: &str) -> bool
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.
Sourcepub async fn update_payload_with_annotation(
&self,
id: &str,
new_payload: Vec<u8>,
annotation: OperatorAnnotation,
) -> bool
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§
impl !Freeze for ResponseBuffer
impl !RefUnwindSafe for ResponseBuffer
impl !UnwindSafe for ResponseBuffer
impl Send for ResponseBuffer
impl Sync for ResponseBuffer
impl Unpin for ResponseBuffer
impl UnsafeUnpin for ResponseBuffer
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
impl<A, B, T> HttpServerConnExec<A, B> for Twhere
B: Body,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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 moreSource§impl<T, U> OverflowingInto<U> for Twhere
U: OverflowingFrom<T>,
impl<T, U> OverflowingInto<U> for Twhere
U: OverflowingFrom<T>,
fn overflowing_into(self) -> (U, bool)
Source§impl<D> OwoColorize for D
impl<D> OwoColorize for D
Source§fn fg<C>(&self) -> FgColorDisplay<'_, C, Self>where
C: Color,
fn fg<C>(&self) -> FgColorDisplay<'_, C, Self>where
C: Color,
Source§fn bg<C>(&self) -> BgColorDisplay<'_, C, Self>where
C: Color,
fn bg<C>(&self) -> BgColorDisplay<'_, C, Self>where
C: Color,
Source§fn black(&self) -> FgColorDisplay<'_, Black, Self>
fn black(&self) -> FgColorDisplay<'_, Black, Self>
Source§fn on_black(&self) -> BgColorDisplay<'_, Black, Self>
fn on_black(&self) -> BgColorDisplay<'_, Black, Self>
Source§fn red(&self) -> FgColorDisplay<'_, Red, Self>
fn red(&self) -> FgColorDisplay<'_, Red, Self>
Source§fn on_red(&self) -> BgColorDisplay<'_, Red, Self>
fn on_red(&self) -> BgColorDisplay<'_, Red, Self>
Source§fn green(&self) -> FgColorDisplay<'_, Green, Self>
fn green(&self) -> FgColorDisplay<'_, Green, Self>
Source§fn on_green(&self) -> BgColorDisplay<'_, Green, Self>
fn on_green(&self) -> BgColorDisplay<'_, Green, Self>
Source§fn yellow(&self) -> FgColorDisplay<'_, Yellow, Self>
fn yellow(&self) -> FgColorDisplay<'_, Yellow, Self>
Source§fn on_yellow(&self) -> BgColorDisplay<'_, Yellow, Self>
fn on_yellow(&self) -> BgColorDisplay<'_, Yellow, Self>
Source§fn blue(&self) -> FgColorDisplay<'_, Blue, Self>
fn blue(&self) -> FgColorDisplay<'_, Blue, Self>
Source§fn on_blue(&self) -> BgColorDisplay<'_, Blue, Self>
fn on_blue(&self) -> BgColorDisplay<'_, Blue, Self>
Source§fn magenta(&self) -> FgColorDisplay<'_, Magenta, Self>
fn magenta(&self) -> FgColorDisplay<'_, Magenta, Self>
Source§fn on_magenta(&self) -> BgColorDisplay<'_, Magenta, Self>
fn on_magenta(&self) -> BgColorDisplay<'_, Magenta, Self>
Source§fn purple(&self) -> FgColorDisplay<'_, Magenta, Self>
fn purple(&self) -> FgColorDisplay<'_, Magenta, Self>
Source§fn on_purple(&self) -> BgColorDisplay<'_, Magenta, Self>
fn on_purple(&self) -> BgColorDisplay<'_, Magenta, Self>
Source§fn cyan(&self) -> FgColorDisplay<'_, Cyan, Self>
fn cyan(&self) -> FgColorDisplay<'_, Cyan, Self>
Source§fn on_cyan(&self) -> BgColorDisplay<'_, Cyan, Self>
fn on_cyan(&self) -> BgColorDisplay<'_, Cyan, Self>
Source§fn white(&self) -> FgColorDisplay<'_, White, Self>
fn white(&self) -> FgColorDisplay<'_, White, Self>
Source§fn on_white(&self) -> BgColorDisplay<'_, White, Self>
fn on_white(&self) -> BgColorDisplay<'_, White, Self>
Source§fn default_color(&self) -> FgColorDisplay<'_, Default, Self>
fn default_color(&self) -> FgColorDisplay<'_, Default, Self>
Source§fn on_default_color(&self) -> BgColorDisplay<'_, Default, Self>
fn on_default_color(&self) -> BgColorDisplay<'_, Default, Self>
Source§fn bright_black(&self) -> FgColorDisplay<'_, BrightBlack, Self>
fn bright_black(&self) -> FgColorDisplay<'_, BrightBlack, Self>
Source§fn on_bright_black(&self) -> BgColorDisplay<'_, BrightBlack, Self>
fn on_bright_black(&self) -> BgColorDisplay<'_, BrightBlack, Self>
Source§fn bright_red(&self) -> FgColorDisplay<'_, BrightRed, Self>
fn bright_red(&self) -> FgColorDisplay<'_, BrightRed, Self>
Source§fn on_bright_red(&self) -> BgColorDisplay<'_, BrightRed, Self>
fn on_bright_red(&self) -> BgColorDisplay<'_, BrightRed, Self>
Source§fn bright_green(&self) -> FgColorDisplay<'_, BrightGreen, Self>
fn bright_green(&self) -> FgColorDisplay<'_, BrightGreen, Self>
Source§fn on_bright_green(&self) -> BgColorDisplay<'_, BrightGreen, Self>
fn on_bright_green(&self) -> BgColorDisplay<'_, BrightGreen, Self>
Source§fn bright_yellow(&self) -> FgColorDisplay<'_, BrightYellow, Self>
fn bright_yellow(&self) -> FgColorDisplay<'_, BrightYellow, Self>
Source§fn on_bright_yellow(&self) -> BgColorDisplay<'_, BrightYellow, Self>
fn on_bright_yellow(&self) -> BgColorDisplay<'_, BrightYellow, Self>
Source§fn bright_blue(&self) -> FgColorDisplay<'_, BrightBlue, Self>
fn bright_blue(&self) -> FgColorDisplay<'_, BrightBlue, Self>
Source§fn on_bright_blue(&self) -> BgColorDisplay<'_, BrightBlue, Self>
fn on_bright_blue(&self) -> BgColorDisplay<'_, BrightBlue, Self>
Source§fn bright_magenta(&self) -> FgColorDisplay<'_, BrightMagenta, Self>
fn bright_magenta(&self) -> FgColorDisplay<'_, BrightMagenta, Self>
Source§fn on_bright_magenta(&self) -> BgColorDisplay<'_, BrightMagenta, Self>
fn on_bright_magenta(&self) -> BgColorDisplay<'_, BrightMagenta, Self>
Source§fn bright_purple(&self) -> FgColorDisplay<'_, BrightMagenta, Self>
fn bright_purple(&self) -> FgColorDisplay<'_, BrightMagenta, Self>
Source§fn on_bright_purple(&self) -> BgColorDisplay<'_, BrightMagenta, Self>
fn on_bright_purple(&self) -> BgColorDisplay<'_, BrightMagenta, Self>
Source§fn bright_cyan(&self) -> FgColorDisplay<'_, BrightCyan, Self>
fn bright_cyan(&self) -> FgColorDisplay<'_, BrightCyan, Self>
Source§fn on_bright_cyan(&self) -> BgColorDisplay<'_, BrightCyan, Self>
fn on_bright_cyan(&self) -> BgColorDisplay<'_, BrightCyan, Self>
Source§fn bright_white(&self) -> FgColorDisplay<'_, BrightWhite, Self>
fn bright_white(&self) -> FgColorDisplay<'_, BrightWhite, Self>
Source§fn on_bright_white(&self) -> BgColorDisplay<'_, BrightWhite, Self>
fn on_bright_white(&self) -> BgColorDisplay<'_, BrightWhite, Self>
Source§fn bold(&self) -> BoldDisplay<'_, Self>
fn bold(&self) -> BoldDisplay<'_, Self>
Source§fn dimmed(&self) -> DimDisplay<'_, Self>
fn dimmed(&self) -> DimDisplay<'_, Self>
Source§fn italic(&self) -> ItalicDisplay<'_, Self>
fn italic(&self) -> ItalicDisplay<'_, Self>
Source§fn underline(&self) -> UnderlineDisplay<'_, Self>
fn underline(&self) -> UnderlineDisplay<'_, Self>
Source§fn blink(&self) -> BlinkDisplay<'_, Self>
fn blink(&self) -> BlinkDisplay<'_, Self>
Source§fn blink_fast(&self) -> BlinkFastDisplay<'_, Self>
fn blink_fast(&self) -> BlinkFastDisplay<'_, Self>
Source§fn reversed(&self) -> ReversedDisplay<'_, Self>
fn reversed(&self) -> ReversedDisplay<'_, Self>
Source§fn strikethrough(&self) -> StrikeThroughDisplay<'_, Self>
fn strikethrough(&self) -> StrikeThroughDisplay<'_, Self>
Source§fn color<Color>(&self, color: Color) -> FgDynColorDisplay<'_, Color, Self>where
Color: DynColor,
fn color<Color>(&self, color: Color) -> FgDynColorDisplay<'_, Color, Self>where
Color: DynColor,
OwoColorize::fg or
a color-specific method, such as OwoColorize::green, Read moreSource§fn on_color<Color>(&self, color: Color) -> BgDynColorDisplay<'_, Color, Self>where
Color: DynColor,
fn on_color<Color>(&self, color: Color) -> BgDynColorDisplay<'_, Color, Self>where
Color: DynColor,
OwoColorize::bg or
a color-specific method, such as OwoColorize::on_yellow, Read more