返回 CodeWhale
turn_heartbeat.rs
根目录 / crates / tui / src / core / engine / turn_heartbeat.rs
1 //! Turn-phase heartbeat and stall self-report (#6184).
2 //!
3 //! A turn that stops producing output used to leave no trace: no log line,
4 //! nothing in `crashes/`, and a UI that could not tell a quiet model wait from
5 //! a wedged engine. The engine now publishes *where* a turn is (its phase), a
6 //! monotonic last-progress stamp, and the bound the current phase may stay
7 //! silent for. A watchdog task, independent of the turn future, turns an
8 //! overdue bounded phase into a log line, a stall record under `crashes/`, and
9 //! a status event naming the phase.
10 //!
11 //! Phases that are owned by their own bound elsewhere — a tool batch (per-tool
12 //! timeouts plus the UI tool-hang watchdog), a compaction pass, a human
13 //! approval — are declared *parked* (`bound = None`) and never reported here.
14 //!
15 //! The heartbeat is shared with the UI through `EngineHandle`, so the UI's own
16 //! watchdog reads the engine's liveness directly instead of inferring it from
17 //! the stream-chunk timeout.
18
19 use std::fmt;
20 use std::path::PathBuf;
21 use std::sync::{Arc, Mutex};
22 use std::time::Duration;
23
24 use tokio::sync::mpsc;
25 use tokio::time::Instant;
26
27 use crate::core::events::Event;
28
29 /// How often the watchdog samples the heartbeat.
30 pub(crate) const STALL_WATCHDOG_TICK: Duration = Duration::from_secs(5);
31 /// Bound for the engine's own between-request work (context assembly, hooks,
32 /// MCP refresh, post-stream bookkeeping). None of it waits on a provider, so a
33 /// few minutes of silence here is a wedge, not a slow model.
34 pub(crate) const PREPARING_PHASE_BOUND: Duration = Duration::from_secs(180);
35 /// Grace added on top of a wait's own timeout. The inner timeout should fire
36 /// first; the heartbeat only reports when that timeout itself failed to.
37 pub(crate) const STALL_BOUND_GRACE: Duration = Duration::from_secs(30);
38
39 /// Where the active turn currently is.
40 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
41 pub(crate) enum TurnPhase {
42 Idle,
43 /// Engine-local work between provider requests.
44 Preparing,
45 /// Request sent; waiting for the stream to open and produce its first event.
46 AwaitingModel,
47 /// Stream open; waiting on the next event.
48 Streaming,
49 /// Planning or executing a tool batch (parked: per-tool bounds own it).
50 Tools,
51 /// Automatic compaction pass (parked: the pass owns its bound).
52 Compacting,
53 }
54
55 impl TurnPhase {
56 #[must_use]
57 pub(crate) const fn label(self) -> &'static str {
58 match self {
59 Self::Idle => "idle",
60 Self::Preparing => "preparing the next request",
61 Self::AwaitingModel => "waiting for the model's first response",
62 Self::Streaming => "streaming the model response",
63 Self::Tools => "running tools",
64 Self::Compacting => "compacting context",
65 }
66 }
67 }
68
69 impl fmt::Display for TurnPhase {
70 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
71 f.write_str(self.label())
72 }
73 }
74
75 /// One detected stall episode.
76 #[derive(Debug, Clone, PartialEq, Eq)]
77 pub(crate) struct StallReport {
78 /// Which watchdog saw it (`engine`, `ui`, `client`).
79 pub source: &'static str,
80 pub phase: String,
81 pub detail: Option<String>,
82 pub turn_id: Option<String>,
83 /// Provider response/request id when the stream reported one, else the
84 /// route label the request went to.
85 pub provider_request: Option<String>,
86 pub since_progress: Duration,
87 pub bound: Option<Duration>,
88 }
89
90 impl StallReport {
91 /// One user-facing line: where it stalled and what to do.
92 #[must_use]
93 pub(crate) fn status_line(&self) -> String {
94 let mut line = format!(
95 "Turn stalled {} — no progress for {}s",
96 self.phase,
97 self.since_progress.as_secs()
98 );
99 if let Some(detail) = self.detail.as_deref().filter(|d| !d.is_empty()) {
100 line.push_str(&format!(" ({detail})"));
101 }
102 line.push_str(". Press Esc to cancel and retry.");
103 line
104 }
105
106 fn record_body(&self) -> String {
107 let timestamp = chrono::Utc::now().to_rfc3339();
108 let bound = self.bound.map_or_else(
109 || "none (parked)".to_string(),
110 |b| format!("{}s", b.as_secs()),
111 );
112 format!(
113 "Kind: turn-stall\nSource: {source}\nTimestamp: {timestamp}\nPhase: {phase}\n\
114 Detail: {detail}\nTurn: {turn}\nProvider request: {request}\n\
115 No progress for: {since}s\nPhase bound: {bound}\n",
116 source = self.source,
117 phase = self.phase,
118 detail = self.detail.as_deref().unwrap_or("-"),
119 turn = self.turn_id.as_deref().unwrap_or("-"),
120 request = self.provider_request.as_deref().unwrap_or("-"),
121 since = self.since_progress.as_secs(),
122 )
123 }
124 }
125
126 #[cfg(test)]
127 thread_local! {
128 static TEST_STALL_RECORD_DIR: std::cell::RefCell<Option<PathBuf>> =
129 const { std::cell::RefCell::new(None) };
130 }
131
132 /// Route stall records for the current test thread into `dir` (tests never
133 /// write into the real `~/.codewhale/crashes`).
134 #[cfg(test)]
135 pub(crate) fn set_test_stall_record_dir(dir: Option<PathBuf>) {
136 TEST_STALL_RECORD_DIR.with(|slot| *slot.borrow_mut() = dir);
137 }
138
139 /// The selected profile's crash directory, shared with panic dumps.
140 fn stall_record_dir() -> Option<PathBuf> {
141 #[cfg(test)]
142 {
143 TEST_STALL_RECORD_DIR.with(|slot| slot.borrow().clone())
144 }
145 #[cfg(not(test))]
146 {
147 codewhale_config::codewhale_home()
148 .ok()
149 .map(|home| home.join("crashes"))
150 }
151 }
152
153 /// Log a stall and write its record to `crashes/<timestamp>-turn-stall-<source>.log`.
154 /// Best effort; returns the record path when a record directory exists. The
155 /// write runs on its own short-lived thread so no caller (engine task or UI
156 /// event loop) blocks a runtime worker on disk I/O (#6149).
157 pub(crate) fn report_stall(report: &StallReport) -> Option<PathBuf> {
158 let path = stall_record_dir().map(|dir| {
159 let stamp = chrono::Utc::now().format("%Y%m%dT%H%M%S%.3fZ");
160 dir.join(format!("{stamp}-turn-stall-{}.log", report.source))
161 });
162 if let Some(path) = path.clone() {
163 let body = report.record_body();
164 let writer = std::thread::spawn(move || {
165 if let Some(dir) = path.parent() {
166 let _ = std::fs::create_dir_all(dir);
167 }
168 let _ = std::fs::write(&path, body);
169 });
170 // Tests read the record right after reporting.
171 #[cfg(test)]
172 let _ = writer.join();
173 #[cfg(not(test))]
174 drop(writer);
175 }
176 let message = format!(
177 "turn stall ({source}): phase={phase} since_progress={since}s bound={bound:?} turn={turn} request={request} detail={detail} record={record}",
178 source = report.source,
179 phase = report.phase,
180 since = report.since_progress.as_secs(),
181 bound = report.bound.map(|b| b.as_secs()),
182 turn = report.turn_id.as_deref().unwrap_or("-"),
183 request = report.provider_request.as_deref().unwrap_or("-"),
184 detail = report.detail.as_deref().unwrap_or("-"),
185 record = path
186 .as_deref()
187 .map_or_else(|| "-".to_string(), |p| p.display().to_string()),
188 );
189 tracing::warn!(target: "turn_stall", "{message}");
190 crate::logging::warn(&message);
191 path
192 }
193
194 #[derive(Debug, Clone)]
195 struct HeartbeatState {
196 phase: TurnPhase,
197 detail: Option<String>,
198 bound: Option<Duration>,
199 last_progress: Instant,
200 turn_id: Option<String>,
201 provider_request: Option<String>,
202 /// Bumped on every phase change or progress touch; a stall is reported at
203 /// most once per value.
204 progress_seq: u64,
205 reported_seq: Option<u64>,
206 /// Latest report, cleared by the next progress.
207 stall: Option<StallReport>,
208 }
209
210 /// Point-in-time view for the UI watchdog.
211 #[derive(Debug, Clone, PartialEq, Eq)]
212 pub(crate) struct HeartbeatSnapshot {
213 pub phase: TurnPhase,
214 pub since_progress: Duration,
215 pub bound: Option<Duration>,
216 pub stall: Option<StallReport>,
217 }
218
219 impl HeartbeatSnapshot {
220 /// The engine is inside a wait it bounds itself and has not reported as
221 /// overdue. The UI must not second-guess it. Parked phases (tools,
222 /// compaction) stay under the UI's own tool-hang and turn watchdogs.
223 #[must_use]
224 pub(crate) fn engine_owns_live_wait(&self) -> bool {
225 self.phase != TurnPhase::Idle && self.bound.is_some() && self.stall.is_none()
226 }
227 }
228
229 /// Shared turn-phase heartbeat. Cheap to update from the turn loop.
230 #[derive(Debug)]
231 pub(crate) struct TurnHeartbeat {
232 state: Mutex<HeartbeatState>,
233 }
234
235 impl Default for TurnHeartbeat {
236 fn default() -> Self {
237 Self {
238 state: Mutex::new(HeartbeatState {
239 phase: TurnPhase::Idle,
240 detail: None,
241 bound: None,
242 last_progress: Instant::now(),
243 turn_id: None,
244 provider_request: None,
245 progress_seq: 0,
246 reported_seq: None,
247 stall: None,
248 }),
249 }
250 }
251 }
252
253 impl TurnHeartbeat {
254 #[must_use]
255 pub(crate) fn new() -> Arc<Self> {
256 Arc::new(Self::default())
257 }
258
259 fn with_state<R>(&self, f: impl FnOnce(&mut HeartbeatState) -> R) -> R {
260 let mut guard = self
261 .state
262 .lock()
263 .unwrap_or_else(std::sync::PoisonError::into_inner);
264 f(&mut guard)
265 }
266
267 /// A new turn starts in [`TurnPhase::Preparing`].
268 pub(crate) fn begin_turn(&self, turn_id: &str) {
269 self.with_state(|state| {
270 state.turn_id = Some(turn_id.to_string());
271 state.provider_request = None;
272 });
273 self.enter(TurnPhase::Preparing, None, Some(PREPARING_PHASE_BOUND));
274 }
275
276 /// Enter `phase`. `bound = None` declares a parked wait that this watchdog
277 /// never reports.
278 pub(crate) fn enter(&self, phase: TurnPhase, detail: Option<String>, bound: Option<Duration>) {
279 self.with_state(|state| {
280 state.phase = phase;
281 state.detail = detail;
282 state.bound = bound;
283 state.last_progress = Instant::now();
284 state.progress_seq = state.progress_seq.wrapping_add(1);
285 state.stall = None;
286 });
287 }
288
289 /// Record progress inside the current phase.
290 pub(crate) fn touch(&self) {
291 self.with_state(|state| {
292 state.last_progress = Instant::now();
293 state.progress_seq = state.progress_seq.wrapping_add(1);
294 state.stall = None;
295 });
296 }
297
298 /// Stream progress: the first event of a request moves the phase from
299 /// awaiting-model to streaming with the inter-chunk bound; later events
300 /// only touch.
301 pub(crate) fn stream_progress(&self, streaming_bound: Duration) {
302 let entering = self.with_state(|state| state.phase != TurnPhase::Streaming);
303 if entering {
304 let detail = self.with_state(|state| state.detail.clone());
305 self.enter(TurnPhase::Streaming, detail, Some(streaming_bound));
306 } else {
307 self.touch();
308 }
309 }
310
311 /// Remember the provider's id for the in-flight response.
312 pub(crate) fn set_provider_request(&self, id: impl Into<String>) {
313 let id = id.into();
314 if id.is_empty() {
315 return;
316 }
317 self.with_state(|state| state.provider_request = Some(id));
318 }
319
320 pub(crate) fn idle(&self) {
321 self.enter(TurnPhase::Idle, None, None);
322 self.with_state(|state| state.turn_id = None);
323 }
324
325 #[must_use]
326 pub(crate) fn snapshot_at(&self, now: Instant) -> HeartbeatSnapshot {
327 self.with_state(|state| HeartbeatSnapshot {
328 phase: state.phase,
329 since_progress: now.saturating_duration_since(state.last_progress),
330 bound: state.bound,
331 stall: state.stall.clone(),
332 })
333 }
334
335 #[must_use]
336 pub(crate) fn snapshot(&self) -> HeartbeatSnapshot {
337 self.snapshot_at(Instant::now())
338 }
339
340 /// Return a report the first time the current bounded phase is overdue.
341 pub(crate) fn detect_stall_at(&self, now: Instant) -> Option<StallReport> {
342 self.with_state(|state| {
343 if state.phase == TurnPhase::Idle || state.reported_seq == Some(state.progress_seq) {
344 return None;
345 }
346 let bound = state.bound?;
347 let since_progress = now.saturating_duration_since(state.last_progress);
348 if since_progress <= bound {
349 return None;
350 }
351 state.reported_seq = Some(state.progress_seq);
352 let report = StallReport {
353 source: "engine",
354 phase: format!("while {}", state.phase.label()),
355 detail: state.detail.clone(),
356 turn_id: state.turn_id.clone(),
357 provider_request: state.provider_request.clone(),
358 since_progress,
359 bound: Some(bound),
360 };
361 state.stall = Some(report.clone());
362 Some(report)
363 })
364 }
365 }
366
367 /// Aborts the wrapped task when dropped (the watchdog must not outlive the
368 /// engine that owns its heartbeat).
369 pub(crate) struct AbortOnDrop(pub(crate) tokio::task::JoinHandle<()>);
370
371 impl Drop for AbortOnDrop {
372 fn drop(&mut self) {
373 self.0.abort();
374 }
375 }
376
377 /// Supervise `heartbeat` until the event channel closes: every overdue bounded
378 /// phase yields one log line, one stall record, and one status event.
379 pub(crate) fn spawn_turn_stall_watchdog(
380 heartbeat: Arc<TurnHeartbeat>,
381 tx_event: mpsc::Sender<Event>,
382 ) -> tokio::task::JoinHandle<()> {
383 spawn_turn_stall_watchdog_every(heartbeat, tx_event, STALL_WATCHDOG_TICK)
384 }
385
386 fn spawn_turn_stall_watchdog_every(
387 heartbeat: Arc<TurnHeartbeat>,
388 tx_event: mpsc::Sender<Event>,
389 tick: Duration,
390 ) -> tokio::task::JoinHandle<()> {
391 #[cfg(test)]
392 let test_dir = TEST_STALL_RECORD_DIR.with(|slot| slot.borrow().clone());
393 tokio::spawn(async move {
394 #[cfg(test)]
395 set_test_stall_record_dir(test_dir);
396 let mut ticker = tokio::time::interval(tick);
397 ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
398 loop {
399 ticker.tick().await;
400 if tx_event.is_closed() {
401 break;
402 }
403 if let Some(report) = heartbeat.detect_stall_at(Instant::now()) {
404 let record = report_stall(&report);
405 let mut line = report.status_line();
406 if let Some(path) = record {
407 line.push_str(&format!(" Stall record: {}", path.display()));
408 }
409 // Never block the watchdog on a full mailbox: a wedged
410 // consumer is exactly the case it exists to survive.
411 let _ = tx_event.try_send(Event::status(line));
412 }
413 }
414 })
415 }
416
417 #[cfg(test)]
418 mod tests {
419 use super::*;
420
421 fn later(secs: u64) -> Instant {
422 Instant::now() + Duration::from_secs(secs)
423 }
424
425 #[test]
426 fn stall_bounded_phase_reports_once_per_episode() {
427 let heartbeat = TurnHeartbeat::new();
428 heartbeat.begin_turn("turn_1");
429 heartbeat.enter(
430 TurnPhase::AwaitingModel,
431 Some("deepseek/deepseek-v4-pro".into()),
432 Some(Duration::from_secs(60)),
433 );
434 assert!(heartbeat.detect_stall_at(Instant::now()).is_none());
435 let report = heartbeat
436 .detect_stall_at(later(61))
437 .expect("overdue bounded phase must report");
438 assert_eq!(report.turn_id.as_deref(), Some("turn_1"));
439 assert!(report.phase.contains("first response"), "{}", report.phase);
440 assert!(
441 heartbeat.detect_stall_at(later(120)).is_none(),
442 "once per episode"
443 );
444 assert!(!heartbeat.snapshot().engine_owns_live_wait());
445 heartbeat.touch();
446 assert!(
447 heartbeat.snapshot().engine_owns_live_wait(),
448 "progress clears the stall"
449 );
450 }
451
452 #[test]
453 fn stall_parked_phase_is_never_reported() {
454 let heartbeat = TurnHeartbeat::new();
455 heartbeat.begin_turn("turn_1");
456 heartbeat.enter(TurnPhase::Tools, Some("exec_shell".into()), None);
457 assert!(heartbeat.detect_stall_at(later(24 * 60 * 60)).is_none());
458 assert!(
459 !heartbeat.snapshot().engine_owns_live_wait(),
460 "parked phases stay under the UI's own watchdogs"
461 );
462 heartbeat.idle();
463 assert!(heartbeat.detect_stall_at(later(24 * 60 * 60)).is_none());
464 }
465
466 #[test]
467 fn stall_first_stream_event_switches_to_inter_chunk_bound() {
468 let heartbeat = TurnHeartbeat::new();
469 heartbeat.begin_turn("turn_1");
470 heartbeat.enter(
471 TurnPhase::AwaitingModel,
472 None,
473 Some(Duration::from_secs(10)),
474 );
475 heartbeat.stream_progress(Duration::from_secs(100));
476 let snapshot = heartbeat.snapshot();
477 assert_eq!(snapshot.phase, TurnPhase::Streaming);
478 assert_eq!(snapshot.bound, Some(Duration::from_secs(100)));
479 assert!(heartbeat.detect_stall_at(later(50)).is_none());
480 assert!(heartbeat.detect_stall_at(later(101)).is_some());
481 }
482
483 /// Fault injection: an inter-chunk wait that never ends produces a log
484 /// line, a `crashes/` stall record, and a status event within the bound
485 /// plus one watchdog tick.
486 #[tokio::test]
487 async fn stall_watchdog_writes_record_and_status_within_bound() {
488 let dir = tempfile::tempdir().expect("tempdir");
489 set_test_stall_record_dir(Some(dir.path().to_path_buf()));
490 let heartbeat = TurnHeartbeat::new();
491 heartbeat.begin_turn("turn_wedged");
492 let bound = Duration::from_millis(200);
493 let tick = Duration::from_millis(20);
494 heartbeat.enter(TurnPhase::Streaming, Some("mock/model".into()), Some(bound));
495 heartbeat.set_provider_request("resp_123");
496 let (tx, mut rx) = mpsc::channel(4);
497 let started = std::time::Instant::now();
498 let watchdog = spawn_turn_stall_watchdog_every(Arc::clone(&heartbeat), tx, tick);
499
500 let event = tokio::time::timeout(Duration::from_secs(10), rx.recv())
501 .await
502 .expect("stall status")
503 .expect("event");
504 let elapsed = started.elapsed();
505 assert!(elapsed >= bound, "not before the bound: {elapsed:?}");
506 let Event::Status { message } = event else {
507 panic!("expected status event");
508 };
509 assert!(
510 message.contains("Turn stalled while streaming"),
511 "{message}"
512 );
513 assert!(message.contains("Esc to cancel and retry"), "{message}");
514 assert!(message.contains("Stall record:"), "{message}");
515
516 let records: Vec<_> = std::fs::read_dir(dir.path())
517 .expect("record dir")
518 .flatten()
519 .map(|entry| std::fs::read_to_string(entry.path()).expect("record"))
520 .collect();
521 assert_eq!(records.len(), 1, "exactly one record per stall episode");
522 assert!(records[0].contains("Kind: turn-stall"));
523 assert!(records[0].contains("Turn: turn_wedged"));
524 assert!(records[0].contains("Provider request: resp_123"));
525 watchdog.abort();
526 set_test_stall_record_dir(None);
527 }
528 }
529
529 lines RUST