use std::time::{Duration, Instant};
use super::{
markdown::CodeFenceState, stream::AppendOnlyStream, StreamKind, StreamUi,
STREAM_PREVIEW_MIN_CHARS, STREAM_UI_TICK,
};
const TARGET_LEAD: Duration = Duration::from_millis(100);
const MAX_RESERVE_CHARS: usize = 2048;
#[derive(Debug, Default)]
pub(super) struct StreamPacer {
last_release: Option<Instant>,
carry: f64,
}
impl StreamPacer {
pub(super) fn reset(&mut self) {
*self = Self::default();
}
fn note_refill(&mut self, now: Instant) {
self.last_release = Some(now);
self.carry = 0.0;
}
fn release_allowance(&mut self, now: Instant, reserve_chars: usize) -> usize {
if reserve_chars == 0 {
self.last_release = Some(now);
self.carry = 0.0;
return 0;
}
if reserve_chars > MAX_RESERVE_CHARS {
self.last_release = Some(now);
self.carry = 0.0;
return reserve_chars;
}
let Some(last_release) = self.last_release else {
self.last_release = Some(now);
return 0;
};
let elapsed = now.saturating_duration_since(last_release).as_secs_f64();
if elapsed <= 0.0 {
return 0;
}
self.last_release = Some(now);
let released = reserve_chars as f64 * elapsed / TARGET_LEAD.as_secs_f64() + self.carry;
let whole = released.floor();
self.carry = released - whole;
(whole as usize).min(reserve_chars)
}
}
impl StreamUi {
pub(super) fn stream(&self, kind: StreamKind) -> &AppendOnlyStream {
match kind {
StreamKind::Assistant => &self.assistant_stream,
StreamKind::Reasoning => &self.reasoning_stream,
}
}
pub(super) fn stream_mut(&mut self, kind: StreamKind) -> &mut AppendOnlyStream {
match kind {
StreamKind::Assistant => &mut self.assistant_stream,
StreamKind::Reasoning => &mut self.reasoning_stream,
}
}
pub(super) fn code_fence(&self, kind: StreamKind) -> &CodeFenceState {
match kind {
StreamKind::Assistant => &self.assistant_stream_code_fence,
StreamKind::Reasoning => &self.reasoning_stream_code_fence,
}
}
fn code_fence_mut(&mut self, kind: StreamKind) -> &mut CodeFenceState {
match kind {
StreamKind::Assistant => &mut self.assistant_stream_code_fence,
StreamKind::Reasoning => &mut self.reasoning_stream_code_fence,
}
}
pub(super) fn advance_code_fence(&mut self, kind: StreamKind, text: &str) {
super::markdown::update_code_block_state(text, self.code_fence_mut(kind));
self.invalidate_preview_cache();
}
pub(super) fn push_delta(&mut self, kind: StreamKind, text: &str, now: Instant) {
if text.is_empty() {
self.schedule_tick(kind, now);
return;
}
let was_empty = self.hold.is_empty();
self.hold.push_str(text);
if was_empty {
self.pacer.note_refill(now);
}
self.release_into(kind, now);
self.schedule_tick(kind, now);
}
pub(super) fn on_tick(&mut self, now: Instant) -> bool {
if self
.stream_tick_deadline
.is_none_or(|deadline| now < deadline)
{
return false;
}
self.stream_tick_deadline = None;
let Some(kind) = self.current_stream_kind else {
return false;
};
let released = self.release_into(kind, now);
if !self.hold.is_empty() {
self.schedule_tick(kind, now);
}
released
}
pub(super) fn flush_hold(&mut self, kind: StreamKind) {
if self.hold.is_empty() {
return;
}
let text = std::mem::take(&mut self.hold);
self.stream_mut(kind).push_delta(&text);
self.pacer.reset();
}
pub(super) fn discard_hold(&mut self) {
self.hold.clear();
self.pacer.reset();
}
pub(super) fn schedule_tick(&mut self, kind: StreamKind, now: Instant) {
let pending_chars = self.stream(kind).pending_text().chars().count();
let needs_tick = !self.hold.is_empty() || pending_chars >= STREAM_PREVIEW_MIN_CHARS;
if !needs_tick {
self.stream_tick_deadline = None;
} else if self.stream_tick_deadline.is_none() {
self.stream_tick_deadline = Some(now + STREAM_UI_TICK);
}
}
pub(super) fn clear_tick_deadline(&mut self) {
self.stream_tick_deadline = None;
}
fn release_into(&mut self, kind: StreamKind, now: Instant) -> bool {
let reserve = self.hold.chars().count();
let chars = self.pacer.release_allowance(now, reserve);
if chars == 0 {
return false;
}
let byte_end = self
.hold
.char_indices()
.nth(chars)
.map_or(self.hold.len(), |(byte_index, _)| byte_index);
let released: String = self.hold.drain(..byte_end).collect();
self.stream_mut(kind).push_delta(&released);
true
}
}
#[cfg(test)]
#[path = "stream_pace_tests.rs"]
mod tests;