返回 CodeWhale
streaming_thinking.rs
根目录 / crates / tui / src / tui / streaming_thinking.rs
1 //! Streaming-thinking lifecycle for the active cell.
2 //!
3 //! DeepSeek V4 emits `reasoning_content` chunks before final answers.
4 //! These get rendered as a "Thinking" entry inside the per-turn active
5 //! cell. This module is the single source of truth for:
6 //!
7 //! - creating a streaming thinking entry on first chunk
8 //! - appending chunks to the live entry
9 //! - showing a localized placeholder while a translation is in-flight
10 //! (and animating its elapsed/spinner suffix)
11 //! - replacing the placeholder when the translation arrives
12 //! - finalizing the entry (stopping the spinner, stamping duration)
13 //! when a thinking block ends
14 //! - stashing the reasoning buffer onto `app.last_reasoning` so the
15 //! summary survives compaction
16
17 use std::time::Duration;
18 use std::time::Instant;
19
20 use crate::tui::active_cell::ActiveCell;
21 use crate::tui::app::App;
22 use crate::tui::history::HistoryCell;
23
24 /// Debounce window for active-cell revision bumps while a thinking block is
25 /// streaming (#1620). Reasoning deltas arrive far faster than the eye can
26 /// follow, and each revision bump invalidates the active cell's wrap cache,
27 /// forcing a full re-wrap of the live tail. Coalescing intermediate bumps to
28 /// one per window keeps the perceived stream smooth without re-wrapping per
29 /// character. ~100ms ≈ 10 intermediate repaints/sec, well below the 120 FPS
30 /// frame cap (see `frame_rate_limiter`) yet imperceptible as lag.
31 ///
32 /// Correctness: this only skips *intermediate* repaints. Appended content is
33 /// never dropped — it lands in the cell immediately — and finalize always
34 /// forces a bump so the final reasoning text is fully rendered.
35 const THINKING_REVISION_THROTTLE: Duration = Duration::from_millis(100);
36
37 /// Bump the active-cell revision for a streaming thinking mutation, but at
38 /// most once per [`THINKING_REVISION_THROTTLE`] window. Returns whether a bump
39 /// was actually emitted. Skipped bumps coalesce into the next one (or into the
40 /// forced finalize bump), so no content is ever lost — only redundant
41 /// intermediate re-wraps are dropped.
42 fn bump_thinking_revision_throttled(app: &mut App, now: Instant) -> bool {
43 let due = app
44 .thinking_revision_last_bump_at
45 .is_none_or(|last| now.saturating_duration_since(last) >= THINKING_REVISION_THROTTLE);
46 if due {
47 app.thinking_revision_last_bump_at = Some(now);
48 app.bump_active_cell_revision();
49 }
50 due
51 }
52
53 /// Ensure an in-flight Thinking entry exists in `active_cell` and return its
54 /// entry index. If no thinking entry is currently streaming, push a fresh one.
55 /// P2.3: thinking shares the active cell with subsequent tool calls so the
56 /// pair render as one logical "Working…" block.
57 pub(super) fn ensure_active_entry(app: &mut App) -> usize {
58 if let Some(idx) = app.streaming_thinking_active_entry {
59 return idx;
60 }
61 if app.active_cell.is_none() {
62 app.active_cell = Some(ActiveCell::new());
63 }
64 let active = app.active_cell.as_mut().expect("active_cell just ensured");
65 let entry_idx = active.push_thinking(HistoryCell::Thinking {
66 content: String::new(),
67 streaming: true,
68 duration_secs: None,
69 });
70 app.streaming_thinking_active_entry = Some(entry_idx);
71 app.bump_active_cell_revision();
72 entry_idx
73 }
74
75 /// Append text to a streaming Thinking entry inside `active_cell`. The text is
76 /// committed to the cell immediately; the active-cell revision bump that
77 /// triggers a re-wrap of the live tail is debounced to at most one per
78 /// [`THINKING_REVISION_THROTTLE`] window (#1620). Skipped bumps coalesce into
79 /// the next append or the forced finalize bump, so no content is ever lost.
80 pub(super) fn append(app: &mut App, entry_idx: usize, text: &str) {
81 append_at(app, entry_idx, text, Instant::now());
82 }
83
84 /// `append` with an injectable clock so the debounce can be tested
85 /// deterministically.
86 fn append_at(app: &mut App, entry_idx: usize, text: &str, now: Instant) {
87 if text.is_empty() {
88 return;
89 }
90 let mutated = if let Some(active) = app.active_cell.as_mut()
91 && let Some(HistoryCell::Thinking { content, .. }) = active.entry_mut(entry_idx)
92 {
93 content.push_str(text);
94 true
95 } else {
96 false
97 };
98 if mutated {
99 bump_thinking_revision_throttled(app, now);
100 }
101 }
102
103 /// Build the spinner-decorated placeholder shown in the thinking entry
104 /// while a translation is in flight (`Thinking… (1.2s |)`).
105 fn translation_placeholder_spinner_frame(app: &App, elapsed: f32) -> &'static str {
106 let elapsed = std::time::Duration::try_from_secs_f32(elapsed.max(0.0))
107 .unwrap_or(std::time::Duration::MAX);
108 let animated_frame =
109 codewhale_ratatui::spin::frame(elapsed, codewhale_ratatui::MotionMode::Full, true);
110 app.motion_policy().spinner_glyph(animated_frame, true)
111 }
112
113 pub(super) fn translation_placeholder_frame(app: &App) -> String {
114 let base = codewhale_localization::thinking_translation_placeholder(app.ui_locale);
115 let elapsed = app
116 .thinking_started_at
117 .or(app.turn_started_at)
118 .map(|started| started.elapsed().as_secs_f32())
119 .unwrap_or_default();
120 let frame = translation_placeholder_spinner_frame(app, elapsed);
121 format!("{base} ({elapsed:.1}s {frame})")
122 }
123
124 /// If the given entry is empty or still showing the translation
125 /// placeholder prefix, replace it with the latest animated frame.
126 pub(super) fn set_placeholder(app: &mut App, entry_idx: usize) {
127 let base = codewhale_localization::thinking_translation_placeholder(app.ui_locale);
128 let next = translation_placeholder_frame(app);
129 let mutated = if let Some(active) = app.active_cell.as_mut()
130 && let Some(HistoryCell::Thinking { content, .. }) = active.entry_mut(entry_idx)
131 && (content.is_empty() || content.starts_with(base))
132 {
133 if *content != next {
134 *content = next;
135 true
136 } else {
137 false
138 }
139 } else {
140 false
141 };
142 if mutated {
143 app.bump_active_cell_revision();
144 }
145 }
146
147 /// Advance the spinner suffix on every existing translation placeholder
148 /// in `active_cell`. Returns true when at least one cell was updated so
149 /// the dispatch loop can schedule another tick.
150 pub(super) fn animate_pending_translation(app: &mut App, translation_pending: bool) -> bool {
151 if !app.translation_enabled {
152 return false;
153 }
154 let thinking_streaming = app.streaming_thinking_active_entry.is_some();
155 if !translation_pending && !thinking_streaming {
156 return false;
157 }
158 let base = codewhale_localization::thinking_translation_placeholder(app.ui_locale);
159 let next = translation_placeholder_frame(app);
160
161 if let Some(active) = app.active_cell.as_mut() {
162 for idx in (0..active.entry_count()).rev() {
163 if let Some(HistoryCell::Thinking { content, .. }) = active.entry_mut(idx)
164 && content.starts_with(base)
165 && *content != next
166 {
167 *content = next.clone();
168 app.bump_active_cell_revision();
169 return true;
170 }
171 }
172 }
173 false
174 }
175
176 /// Replace a translation placeholder with the finished translated text.
177 /// Searches the active cell first, then the finalized history (covers
178 /// the case where the translation lands after the thinking block was
179 /// already moved into history).
180 pub(super) fn replace_pending_translation(
181 app: &mut App,
182 placeholder: &str,
183 translated_text: String,
184 ) {
185 if let Some(active) = app.active_cell.as_mut() {
186 for idx in (0..active.entry_count()).rev() {
187 if let Some(HistoryCell::Thinking { content, .. }) = active.entry_mut(idx)
188 && content.starts_with(placeholder)
189 {
190 *content = translated_text;
191 app.bump_active_cell_revision();
192 return;
193 }
194 }
195 }
196
197 for idx in (0..app.history.len()).rev() {
198 if let Some(HistoryCell::Thinking { content, .. }) = app.history.get_mut(idx)
199 && content.starts_with(placeholder)
200 {
201 *content = translated_text;
202 app.bump_history_cell(idx);
203 return;
204 }
205 }
206 }
207
208 /// Start a new streaming thinking block. If another thinking block is still
209 /// active, first drain its pending UI tail so a late block boundary cannot
210 /// discard content buffered inside `StreamingState`.
211 pub(super) fn start_block(app: &mut App) -> bool {
212 let finalized_previous = if app.streaming_thinking_active_entry.is_some() {
213 let finalized = finalize_current(app);
214 stash_reasoning_buffer_into_last_reasoning(app);
215 finalized
216 } else {
217 false
218 };
219
220 app.reasoning_buffer.clear();
221 app.reasoning_header = None;
222 app.thinking_started_at = Some(Instant::now());
223 app.streaming_state.reset();
224 app.streaming_state.start_thinking(0);
225 let _ = ensure_active_entry(app);
226 finalized_previous
227 }
228
229 /// Finalize the currently-streaming thinking entry: drain the pending
230 /// state buffer, compute elapsed duration, stop the spinner.
231 pub(super) fn finalize_current(app: &mut App) -> bool {
232 let duration = app
233 .thinking_started_at
234 .take()
235 .map(|t| t.elapsed().as_secs_f32());
236 let remaining = app.streaming_state.finalize_block_text(0);
237 finalize_active_entry(app, duration, &remaining)
238 }
239
240 /// Move the in-flight reasoning buffer onto `app.last_reasoning` so the
241 /// summary survives compaction or transcript trimming.
242 pub(super) fn stash_reasoning_buffer_into_last_reasoning(app: &mut App) {
243 if app.reasoning_buffer.is_empty() {
244 return;
245 }
246
247 if let Some(existing) = app.last_reasoning.as_mut()
248 && !existing.is_empty()
249 {
250 if !existing.ends_with('\n') {
251 existing.push('\n');
252 }
253 existing.push_str(&app.reasoning_buffer);
254 } else {
255 app.last_reasoning = Some(app.reasoning_buffer.clone());
256 }
257 app.reasoning_buffer.clear();
258 }
259
260 /// Finalize the in-flight thinking entry in `active_cell`: append the
261 /// collector's remaining buffered text, stop the spinner, and stamp the
262 /// duration. Returns `true` when a thinking entry was finalized (so the
263 /// dispatch loop knows the transcript was touched). No-op if no thinking
264 /// entry is currently streaming.
265 pub(super) fn finalize_active_entry(app: &mut App, duration: Option<f32>, remaining: &str) -> bool {
266 let Some(entry_idx) = app.streaming_thinking_active_entry.take() else {
267 return false;
268 };
269 if !remaining.is_empty() {
270 append(app, entry_idx, remaining);
271 }
272 if let Some(active) = app.active_cell.as_mut()
273 && let Some(HistoryCell::Thinking {
274 streaming,
275 duration_secs,
276 ..
277 }) = active.entry_mut(entry_idx)
278 {
279 *streaming = false;
280 *duration_secs = duration;
281 }
282 // Red line (#1620): finalize must force a bump so the final reasoning text
283 // is fully rendered even if the last appended chunk was throttled. Reset
284 // the debounce window so the next thinking block's first chunk renders
285 // immediately rather than being coalesced into a stale window.
286 app.thinking_revision_last_bump_at = None;
287 app.bump_active_cell_revision();
288 true
289 }
290
291 #[cfg(test)]
292 mod tests {
293 use super::*;
294 use crate::config::Config;
295 use crate::tui::app::{App, TuiOptions};
296 use std::path::PathBuf;
297
298 fn test_app() -> App {
299 let options = TuiOptions {
300 start_in_agent_mode: true,
301 skip_onboarding: false,
302 ..crate::test_support::test_tui_options(PathBuf::from("."))
303 };
304 App::new(options, &Config::default())
305 }
306
307 fn thinking_content(app: &App, entry_idx: usize) -> String {
308 match app
309 .active_cell
310 .as_ref()
311 .and_then(|active| active.entries().get(entry_idx))
312 {
313 Some(HistoryCell::Thinking { content, .. }) => content.clone(),
314 other => panic!("expected a Thinking entry at {entry_idx}, got {other:?}"),
315 }
316 }
317
318 #[test]
319 fn translation_placeholder_spinner_uses_full_motion_only() {
320 let mut app = test_app();
321 app.low_motion = false;
322 app.fancy_animations = true;
323 assert_eq!(translation_placeholder_spinner_frame(&app, 0.0), ">");
324 assert_eq!(translation_placeholder_spinner_frame(&app, 0.6), "\\");
325
326 app.low_motion = true;
327 assert_eq!(translation_placeholder_spinner_frame(&app, 0.0), "●");
328 assert_eq!(translation_placeholder_spinner_frame(&app, 0.6), "●");
329
330 app.low_motion = false;
331 app.fancy_animations = false;
332 assert_eq!(translation_placeholder_spinner_frame(&app, 0.0), "›");
333 assert_eq!(translation_placeholder_spinner_frame(&app, 0.6), "›");
334 }
335
336 /// #1620: a burst of reasoning chunks inside one throttle window must
337 /// coalesce to a single active-cell revision bump (so the renderer
338 /// re-wraps the live tail ~10x/sec instead of once per character), while
339 /// every byte of content is preserved and finalize forces a final bump.
340 #[test]
341 fn issue_1620_throttles_thinking_bumps_without_losing_content() {
342 let mut app = test_app();
343 let entry = ensure_active_entry(&mut app);
344 // `ensure_active_entry` bumped once on creation; start the measurement
345 // from a clean throttle window so the first append renders immediately.
346 app.thinking_revision_last_bump_at = None;
347 let rev_before = app.active_cell_revision;
348
349 let t0 = Instant::now();
350 let chunks = [
351 "Hel", "lo, ", "this", " is", " a", " lo", "ng", " re", "ason", "ing",
352 ];
353 // All ten chunks land within a single 100ms window (5ms apart).
354 for (i, chunk) in chunks.iter().enumerate() {
355 append_at(
356 &mut app,
357 entry,
358 chunk,
359 t0 + Duration::from_millis(i as u64 * 5),
360 );
361 }
362 assert_eq!(
363 app.active_cell_revision.wrapping_sub(rev_before),
364 1,
365 "rapid chunks within one throttle window must coalesce to one bump"
366 );
367
368 // A chunk after the window expires is allowed to bump again.
369 append_at(
370 &mut app,
371 entry,
372 " stream",
373 t0 + THINKING_REVISION_THROTTLE + Duration::from_millis(10),
374 );
375 assert_eq!(
376 app.active_cell_revision.wrapping_sub(rev_before),
377 2,
378 "a chunk past the throttle window should bump once more"
379 );
380
381 // No content was dropped despite the skipped intermediate bumps.
382 let expected = format!("{} stream", chunks.concat());
383 assert_eq!(thinking_content(&app, entry), expected);
384
385 // Red line: finalize forces exactly one bump and flushes the tail.
386 let rev_pre_final = app.active_cell_revision;
387 let finalized = finalize_active_entry(&mut app, Some(1.5), " [end]");
388 assert!(finalized, "finalize should report it finalized an entry");
389 assert_eq!(
390 app.active_cell_revision,
391 rev_pre_final.wrapping_add(1),
392 "finalize must always force exactly one revision bump"
393 );
394 assert_eq!(
395 thinking_content(&app, entry),
396 format!("{expected} [end]"),
397 "finalize must not drop the trailing reasoning text"
398 );
399 assert!(
400 app.thinking_revision_last_bump_at.is_none(),
401 "finalize should reset the throttle window for the next block"
402 );
403 }
404 }
405
405 lines RUST