| 1 | //! Task-panel and shell projection: task-panel refresh, shell live-output |
| 2 | //! reconciliation, detached-job projection, and RLM task entries |
| 3 | //! (TUI_MODULARIZATION.md slice 5). Pure projection — no dispatch here. |
| 4 | |
| 5 | use super::*; |
| 6 | use crate::tui::automation_panel::AutomationScan; |
| 7 | use crate::tui::background_finished::{FinishedOutcome, FinishedWork, MAX_FINISHED_SHELLS}; |
| 8 | use codewhale_localization::{Locale, MessageId, tr}; |
| 9 | |
| 10 | pub(super) async fn refresh_active_task_panel( |
| 11 | app: &mut App, |
| 12 | task_manager: &SharedTaskManager, |
| 13 | ) -> bool { |
| 14 | let namespace_changed = app.task_panel_session_id != app.current_session_id; |
| 15 | if namespace_changed { |
| 16 | app.task_panel.clear(); |
| 17 | app.task_panel_session_id = app.current_session_id.clone(); |
| 18 | app.task_panel_unavailable = false; |
| 19 | } |
| 20 | let tasks = match app.current_session_id.as_deref() { |
| 21 | Some(session_id) => match task_manager |
| 22 | .list_tasks_for_owner(None, None, session_id) |
| 23 | .await |
| 24 | { |
| 25 | Ok(tasks) => tasks, |
| 26 | Err(error) => { |
| 27 | let changed = namespace_changed || !app.task_panel_unavailable; |
| 28 | if !app.task_panel_unavailable { |
| 29 | app.push_status_toast( |
| 30 | codewhale_localization::tr( |
| 31 | app.ui_locale, |
| 32 | codewhale_localization::MessageId::TaskInventoryUnavailable, |
| 33 | ) |
| 34 | .to_string(), |
| 35 | crate::tui::app::StatusToastLevel::Warning, |
| 36 | Some(8_000), |
| 37 | ); |
| 38 | tracing::warn!(%error, "Task inventory unavailable; preserving scoped snapshot"); |
| 39 | } |
| 40 | app.task_panel_unavailable = true; |
| 41 | return changed; |
| 42 | } |
| 43 | }, |
| 44 | None => Vec::new(), |
| 45 | }; |
| 46 | let was_unavailable = std::mem::replace(&mut app.task_panel_unavailable, false); |
| 47 | // #6565: a durable task that failed or was cancelled is as much news as |
| 48 | // one that completed; each lands in the batched notice by its summary. |
| 49 | let newly_finished_tasks = tasks |
| 50 | .iter() |
| 51 | .filter_map(|task| { |
| 52 | let (outcome, word) = match task.status { |
| 53 | TaskStatus::Completed => (FinishedOutcome::Done, MessageId::BackgroundOutcomeDone), |
| 54 | TaskStatus::Failed => (FinishedOutcome::Failed, MessageId::BackgroundOutcomeFailed), |
| 55 | TaskStatus::Canceled => ( |
| 56 | FinishedOutcome::Stopped, |
| 57 | MessageId::BackgroundOutcomeCancelled, |
| 58 | ), |
| 59 | TaskStatus::Queued | TaskStatus::Running => return None, |
| 60 | }; |
| 61 | // Restored history is not a new completion. A task that finishes |
| 62 | // between polls still counts, without ever needing a live row. |
| 63 | if !task.ended_at.is_some_and(|at| at >= app.session_started_at) |
| 64 | || !app.notified_task_ids.insert(task.id.clone()) |
| 65 | { |
| 66 | return None; |
| 67 | } |
| 68 | let word = tr(app.ui_locale, word); |
| 69 | let summary = match task.duration_ms { |
| 70 | Some(ms) => format!("{word} · {}", crate::agent_roster::format_duration(ms)), |
| 71 | None => word.to_string(), |
| 72 | }; |
| 73 | let summary = match task.error.as_deref().map(str::trim) { |
| 74 | Some(error) if !error.is_empty() && outcome != FinishedOutcome::Done => { |
| 75 | format!("{summary} · {}", bound_agent_activity_text(error)) |
| 76 | } |
| 77 | _ => summary, |
| 78 | }; |
| 79 | Some( |
| 80 | FinishedWork::task( |
| 81 | &task.prompt_summary, |
| 82 | outcome, |
| 83 | summary, |
| 84 | std::time::Duration::from_millis(task.duration_ms.unwrap_or_default()), |
| 85 | ) |
| 86 | .in_turn(app.is_loading), |
| 87 | ) |
| 88 | }) |
| 89 | .collect::<Vec<_>>(); |
| 90 | let durable_background_completed = newly_finished_tasks |
| 91 | .iter() |
| 92 | .any(|task| task.outcome == FinishedOutcome::Done); |
| 93 | app.background_finished.extend(newly_finished_tasks); |
| 94 | let mut lifecycle_changed = false; |
| 95 | if let (Some(work), Some(session_id)) = ( |
| 96 | app.runtime_services.work.as_ref(), |
| 97 | app.current_session_id.as_deref(), |
| 98 | ) { |
| 99 | for task in &tasks { |
| 100 | if !task.execution_binding_known { |
| 101 | continue; |
| 102 | } |
| 103 | let external = format!("task:{}", task.id); |
| 104 | if !work.has_operation_binding(Some(session_id), &external) { |
| 105 | continue; |
| 106 | } |
| 107 | match work.reconcile_operation( |
| 108 | session_id, |
| 109 | task_owner_snapshot( |
| 110 | &task.id, |
| 111 | task.status, |
| 112 | task.lifecycle_seq, |
| 113 | task.created_at, |
| 114 | task.started_at, |
| 115 | task.ended_at, |
| 116 | ), |
| 117 | ) { |
| 118 | Ok(changed) => lifecycle_changed |= changed, |
| 119 | Err(err) => { |
| 120 | tracing::warn!(task_id = %task.id, error = %err, "failed to reconcile durable task lifecycle"); |
| 121 | } |
| 122 | } |
| 123 | } |
| 124 | } |
| 125 | if lifecycle_changed && let Err(err) = persist_pending_work_checkpoint(app).await { |
| 126 | tracing::warn!(error = %err, "durable task lifecycle checkpoint remains pending"); |
| 127 | } |
| 128 | let session_started_at = app.session_started_at; |
| 129 | let mut entries: Vec<TaskPanelEntry> = |
| 130 | select_work_sidebar_tasks(tasks, session_started_at, app.current_session_id.as_deref()) |
| 131 | .into_iter() |
| 132 | .map(|summary| { |
| 133 | let unverified = !summary.execution_binding_known |
| 134 | && matches!(summary.status, TaskStatus::Queued | TaskStatus::Running); |
| 135 | let mut entry = task_summary_to_panel_entry(summary); |
| 136 | if unverified { |
| 137 | entry.stale = true; |
| 138 | entry.role = Some( |
| 139 | codewhale_localization::tr( |
| 140 | app.ui_locale, |
| 141 | codewhale_localization::MessageId::TaskOwnershipUnverified, |
| 142 | ) |
| 143 | .to_string(), |
| 144 | ); |
| 145 | } |
| 146 | entry |
| 147 | }) |
| 148 | .collect(); |
| 149 | |
| 150 | entries.extend(active_rlm_task_entries(app)); |
| 151 | |
| 152 | // #3804: this is a render-only read of shell jobs and must not block the |
| 153 | // async UI loop on the shell manager's std::sync Mutex. Use try_lock; on |
| 154 | // contention, fail closed (no shell rows this frame) rather than show a |
| 155 | // snapshot that may belong to a replaced session. Shell ownership, |
| 156 | // cancellation, approval state, and output capture never depend on this |
| 157 | // refresh succeeding. |
| 158 | let jobs = match app.runtime_services.shell_manager.as_ref() { |
| 159 | Some(shell_mgr) => match shell_mgr.try_lock() { |
| 160 | Ok(mut mgr) => Some( |
| 161 | mgr.list_jobs_for_session(app.current_session_id.as_deref().unwrap_or_default()), |
| 162 | ), |
| 163 | Err(_) => None, |
| 164 | }, |
| 165 | None => None, |
| 166 | }; |
| 167 | let jobs = jobs.unwrap_or_default(); |
| 168 | let shell_background_completed = project_shell_jobs(app, &mut entries, &jobs); |
| 169 | |
| 170 | // Report whether anything visible changed so the idle tick can skip the |
| 171 | // redraw: an unconditional 2.5 s repaint kept the app from ever going |
| 172 | // quiescent (#3757). |
| 173 | let changed = namespace_changed |
| 174 | || was_unavailable |
| 175 | || lifecycle_changed |
| 176 | || shell_background_completed |
| 177 | || app.task_panel != entries; |
| 178 | app.task_panel = entries; |
| 179 | let tip_shown = (durable_background_completed || shell_background_completed) |
| 180 | && app.maybe_show_behavioral_tip( |
| 181 | crate::tui::behavioral_tips::BehavioralTip::BackgroundJobReceipt, |
| 182 | ); |
| 183 | changed || tip_shown |
| 184 | } |
| 185 | |
| 186 | /// Fold one session-scoped snapshot. An empty snapshot may be a lock miss; |
| 187 | /// neither completion deduplication nor retained IDs depend on visible rows. |
| 188 | pub(super) fn project_shell_jobs( |
| 189 | app: &mut App, |
| 190 | entries: &mut Vec<TaskPanelEntry>, |
| 191 | jobs: &[crate::tools::shell::ShellJobSnapshot], |
| 192 | ) -> bool { |
| 193 | // #6565: every terminal status is a completion a person should hear |
| 194 | // about (a failed, killed or timed-out shell used to vanish silently), |
| 195 | // and a finished shell stays listed, muted, with how it ended. |
| 196 | let finished_now = newly_terminal(&mut app.notified_shell_ids, jobs, app.session_started_at); |
| 197 | for job in &finished_now { |
| 198 | let retained = app |
| 199 | .finished_shell_ids |
| 200 | .entry(app.current_session_id.clone().unwrap_or_default()) |
| 201 | .or_default(); |
| 202 | retained.push_back(job.id.clone()); |
| 203 | while retained.len() > MAX_FINISHED_SHELLS { |
| 204 | retained.pop_front(); |
| 205 | } |
| 206 | let (outcome, summary) = shell_outcome(app.ui_locale, job); |
| 207 | let toast = format!("shell · {} · {summary}", job.command.trim()); |
| 208 | let level = if outcome == FinishedOutcome::Done { |
| 209 | crate::tui::app::StatusToastLevel::Info |
| 210 | } else { |
| 211 | crate::tui::app::StatusToastLevel::Warning |
| 212 | }; |
| 213 | app.push_status_toast(bound_agent_activity_text(&toast), level, Some(6_000)); |
| 214 | let parent_busy = app.is_loading; |
| 215 | app.background_finished.push( |
| 216 | FinishedWork::shell( |
| 217 | &job.command, |
| 218 | outcome, |
| 219 | summary, |
| 220 | std::time::Duration::from_millis(job.elapsed_ms), |
| 221 | ) |
| 222 | .in_turn(parent_busy), |
| 223 | ); |
| 224 | } |
| 225 | let finished_order = app |
| 226 | .finished_shell_ids |
| 227 | .get(app.current_session_id.as_deref().unwrap_or_default()) |
| 228 | .into_iter() |
| 229 | .flatten() |
| 230 | .cloned() |
| 231 | .collect::<Vec<_>>(); |
| 232 | let shell_entry = |job: &crate::tools::shell::ShellJobSnapshot| TaskPanelEntry { |
| 233 | id: job.id.clone(), |
| 234 | status: shell_status_token(&job.status).to_string(), |
| 235 | prompt_summary: format!("shell: {}", job.command), |
| 236 | duration_ms: Some(job.elapsed_ms), |
| 237 | kind: TaskPanelEntryKind::Shell, |
| 238 | stale: job.stale, |
| 239 | elapsed_since_output_ms: job.elapsed_since_output_ms, |
| 240 | owner_agent_id: job.owner_agent_id.clone(), |
| 241 | owner_agent_name: job.owner_agent_name.clone(), |
| 242 | current_tool: None, |
| 243 | role: None, |
| 244 | files_touched: 0, |
| 245 | exit_code: job.exit_code, |
| 246 | }; |
| 247 | entries.extend( |
| 248 | jobs.iter() |
| 249 | .filter(|job| matches!(job.status, crate::tools::shell::ShellStatus::Running)) |
| 250 | .map(shell_entry), |
| 251 | ); |
| 252 | // Finished shells in the order they finished, so the list is stable and |
| 253 | // an idle tick with nothing new changes nothing (#3757). |
| 254 | entries.extend( |
| 255 | finished_order |
| 256 | .iter() |
| 257 | .filter_map(|id| jobs.iter().find(|job| &job.id == id)) |
| 258 | .filter(|job| { |
| 259 | job.background && !matches!(job.status, crate::tools::shell::ShellStatus::Running) |
| 260 | }) |
| 261 | .map(shell_entry), |
| 262 | ); |
| 263 | !finished_now.is_empty() |
| 264 | } |
| 265 | |
| 266 | /// The wire token a shell status shows as in the task panel. |
| 267 | pub(crate) fn shell_status_token(status: &crate::tools::shell::ShellStatus) -> &'static str { |
| 268 | use crate::tools::shell::ShellStatus; |
| 269 | match status { |
| 270 | ShellStatus::Running => "running", |
| 271 | ShellStatus::Completed => "completed", |
| 272 | ShellStatus::Failed => "failed", |
| 273 | ShellStatus::Killed => "killed", |
| 274 | ShellStatus::TimedOut => "timed_out", |
| 275 | } |
| 276 | } |
| 277 | |
| 278 | /// How a finished shell ended, in the words its row and notice use: |
| 279 | /// `exit 0 · 12s`, `failed · exit 2`, `killed`, `timed out`. |
| 280 | pub(crate) fn shell_outcome( |
| 281 | locale: Locale, |
| 282 | job: &crate::tools::shell::ShellJobSnapshot, |
| 283 | ) -> (FinishedOutcome, String) { |
| 284 | use crate::tools::shell::ShellStatus; |
| 285 | let exit = job |
| 286 | .exit_code |
| 287 | .map(|code| tr(locale, MessageId::BackgroundExitCode).replace("{code}", &code.to_string())); |
| 288 | let took = crate::agent_roster::format_duration(job.elapsed_ms); |
| 289 | match job.status { |
| 290 | ShellStatus::Completed | ShellStatus::Running => ( |
| 291 | FinishedOutcome::Done, |
| 292 | format!( |
| 293 | "{} · {took}", |
| 294 | exit.unwrap_or_else(|| tr(locale, MessageId::BackgroundOutcomeDone).into_owned()) |
| 295 | ), |
| 296 | ), |
| 297 | ShellStatus::Failed => ( |
| 298 | FinishedOutcome::Failed, |
| 299 | match exit { |
| 300 | Some(exit) => format!( |
| 301 | "{} · {exit}", |
| 302 | tr(locale, MessageId::BackgroundOutcomeFailed) |
| 303 | ), |
| 304 | None => tr(locale, MessageId::BackgroundOutcomeFailed).into_owned(), |
| 305 | }, |
| 306 | ), |
| 307 | ShellStatus::Killed => ( |
| 308 | FinishedOutcome::Stopped, |
| 309 | tr(locale, MessageId::BackgroundOutcomeKilled).into_owned(), |
| 310 | ), |
| 311 | ShellStatus::TimedOut => ( |
| 312 | FinishedOutcome::Stopped, |
| 313 | tr(locale, MessageId::BackgroundOutcomeTimedOut).into_owned(), |
| 314 | ), |
| 315 | } |
| 316 | } |
| 317 | |
| 318 | /// Every fresh background completion once, even when never observed running. |
| 319 | /// A missed lock, session switch or evicted row never resets this ledger. |
| 320 | /// The session-start cutoff suppresses restored history after a TUI restart. |
| 321 | /// Finish time is the manager's first terminal observation, not OS exit time; |
| 322 | /// legacy snapshots without it are history, never new completion receipts. |
| 323 | pub(super) fn newly_terminal<'a>( |
| 324 | notified_ids: &mut HashSet<String>, |
| 325 | jobs: &'a [crate::tools::shell::ShellJobSnapshot], |
| 326 | session_started_at: chrono::DateTime<chrono::Utc>, |
| 327 | ) -> Vec<&'a crate::tools::shell::ShellJobSnapshot> { |
| 328 | jobs.iter() |
| 329 | .filter(|job| job.background && job.finished_at.is_some_and(|at| at >= session_started_at)) |
| 330 | .filter(|job| !matches!(job.status, crate::tools::shell::ShellStatus::Running)) |
| 331 | .filter(|job| notified_ids.insert(job.id.clone())) |
| 332 | .collect() |
| 333 | } |
| 334 | |
| 335 | /// Newest runs scanned per automation when refreshing the automation |
| 336 | /// projection. Live runs sit at the head of the newest-first listing, and a |
| 337 | /// failure once seen is held unacknowledged by the projection until the |
| 338 | /// operator engages the automation surface, so the band never needs a full |
| 339 | /// run-history scan on the render cadence. |
| 340 | const AUTOMATION_PANEL_RUN_SCAN: usize = 25; |
| 341 | |
| 342 | /// Refresh the activity band's scheduled-work projection |
| 343 | /// (AUTOMATION-VISIBILITY-SPEC §2.1) from the durable automation store. |
| 344 | /// |
| 345 | /// The store is files on disk — every definition plus up to |
| 346 | /// `AUTOMATION_PANEL_RUN_SCAN` run files per definition — so the scan never |
| 347 | /// runs on the async UI loop: it is taken on a blocking thread, and this |
| 348 | /// tick folds whatever scan has finished, then starts the next one. At most |
| 349 | /// one scan is in flight; a slow disk costs the band latency, never the |
| 350 | /// frame. Returns whether the visible state changed, so the idle tick can |
| 351 | /// skip the redraw (#3757). |
| 352 | /// |
| 353 | /// `start_next` says whether the next scan may begin once a finished one is |
| 354 | /// folded. The caller passes [`automation_scan_is_due`], so a quiet session |
| 355 | /// with nothing scheduled does not re-read the store every 2.5 s (#6728); a |
| 356 | /// scan already in flight is still folded on every tick. |
| 357 | pub(super) async fn refresh_automation_panel(app: &mut App, start_next: bool) -> bool { |
| 358 | let mut changed = false; |
| 359 | if let Some(scan) = app.automation_scan.take() { |
| 360 | if scan.is_finished() { |
| 361 | match scan.await { |
| 362 | Ok(scan) => changed = fold_automation_scan(app, &scan), |
| 363 | Err(err) => { |
| 364 | tracing::warn!(error = %err, "automation panel scan task failed"); |
| 365 | } |
| 366 | } |
| 367 | } else { |
| 368 | app.automation_scan = Some(scan); |
| 369 | return false; |
| 370 | } |
| 371 | } |
| 372 | if start_next { |
| 373 | app.automation_scan = start_automation_scan(app, false); |
| 374 | } |
| 375 | changed |
| 376 | } |
| 377 | |
| 378 | /// The 2.5 s cadence the activity band is scanned at whenever it has |
| 379 | /// something to show or the user is around. |
| 380 | pub(super) const AUTOMATION_SCAN_BUSY_INTERVAL: std::time::Duration = |
| 381 | std::time::Duration::from_millis(2_500); |
| 382 | |
| 383 | /// The cadence for a quiet session whose band shows nothing: a new |
| 384 | /// automation created from outside this process lights the band within this |
| 385 | /// long. Anything created from inside it (a command, a tool call, a click) is |
| 386 | /// input or an engine event, which restores the busy cadence at once. |
| 387 | pub(super) const AUTOMATION_SCAN_QUIET_INTERVAL: std::time::Duration = |
| 388 | std::time::Duration::from_secs(15); |
| 389 | |
| 390 | /// How long to wait between automation scans. Pure: busy unless the UI is |
| 391 | /// quiet *and* the band has nothing to track *and* no automation surface is |
| 392 | /// open. A scan already in flight does not matter here: it is folded on every |
| 393 | /// 2.5 s tick whatever this says, and only the *start* of the next one waits. |
| 394 | pub(super) fn automation_scan_interval(app: &App, ui_quiet: bool) -> std::time::Duration { |
| 395 | let panel = &app.automation_panel; |
| 396 | let band_tracks_something = panel.active_automations > 0 |
| 397 | || panel.live_runs > 0 |
| 398 | || !panel.live_run_owners().is_empty() |
| 399 | || panel.has_unacknowledged_failure(); |
| 400 | let surface_open = app.view_stack.top_kind() == Some(crate::tui::views::ModalKind::Automations); |
| 401 | if ui_quiet && !band_tracks_something && !surface_open { |
| 402 | AUTOMATION_SCAN_QUIET_INTERVAL |
| 403 | } else { |
| 404 | AUTOMATION_SCAN_BUSY_INTERVAL |
| 405 | } |
| 406 | } |
| 407 | |
| 408 | /// Whether the next automation scan may start, `since_last_scan` after the |
| 409 | /// previous one. |
| 410 | pub(super) fn automation_scan_is_due( |
| 411 | app: &App, |
| 412 | ui_quiet: bool, |
| 413 | since_last_scan: std::time::Duration, |
| 414 | ) -> bool { |
| 415 | since_last_scan >= automation_scan_interval(app, ui_quiet) |
| 416 | } |
| 417 | |
| 418 | /// Startup variant: take one scan and wait for it, so the first frame |
| 419 | /// already carries the band count (the task panel gets the same courtesy). |
| 420 | /// The startup pass is FULL — every run file — so a long-running task that |
| 421 | /// already sits behind more than a window of newer runs is visible from the |
| 422 | /// first frame. The manager lock is retried briefly: the scheduler tick |
| 423 | /// holds it only while persisting, and giving up here would blank the band |
| 424 | /// on a contended startup. |
| 425 | pub(super) async fn refresh_automation_panel_blocking(app: &mut App) -> bool { |
| 426 | let scan = 'retry: { |
| 427 | for _ in 0..STARTUP_SCAN_LOCK_RETRIES { |
| 428 | if let Some(scan) = start_automation_scan(app, true) { |
| 429 | break 'retry Some(scan); |
| 430 | } |
| 431 | tokio::time::sleep(STARTUP_SCAN_LOCK_RETRY_DELAY).await; |
| 432 | } |
| 433 | None |
| 434 | }; |
| 435 | let Some(scan) = scan else { |
| 436 | return false; |
| 437 | }; |
| 438 | match scan.await { |
| 439 | Ok(scan) => fold_automation_scan(app, &scan), |
| 440 | Err(err) => { |
| 441 | tracing::warn!(error = %err, "automation panel scan task failed"); |
| 442 | false |
| 443 | } |
| 444 | } |
| 445 | } |
| 446 | |
| 447 | /// Startup lock-retry budget: the scheduler's persist phase is short; ten |
| 448 | /// 50 ms attempts covers it without parking startup on a stuck lock. |
| 449 | const STARTUP_SCAN_LOCK_RETRIES: usize = 10; |
| 450 | const STARTUP_SCAN_LOCK_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(50); |
| 451 | |
| 452 | /// Start one store scan on a blocking thread. The manager is cloned out |
| 453 | /// from under its tokio Mutex (`try_lock`, same rule as the shell snapshot: |
| 454 | /// the scheduler tick holds that lock while persisting, and the UI loop |
| 455 | /// must not park behind it; on contention the next tick retries) so the |
| 456 | /// scan holds no lock while it reads — the store's writes are atomic |
| 457 | /// renames, so a concurrent read sees a whole file either way. |
| 458 | /// |
| 459 | /// `full` reads every run file (startup only). The cadence scan reads the |
| 460 | /// newest `AUTOMATION_PANEL_RUN_SCAN` runs per automation PLUS the runs |
| 461 | /// this session already watched go live, wherever they sit in history — a |
| 462 | /// frequent automation can stack newer runs behind a long-running task, |
| 463 | /// and neither the live count nor the settle receipt may depend on the |
| 464 | /// task staying inside the newest window. |
| 465 | fn start_automation_scan(app: &App, full: bool) -> Option<tokio::task::JoinHandle<AutomationScan>> { |
| 466 | let automations = app.runtime_services.automations.as_ref()?; |
| 467 | let manager = automations.try_lock().ok()?.clone(); |
| 468 | let live_owners = app.automation_panel.live_run_owners(); |
| 469 | Some(tokio::task::spawn_blocking(move || { |
| 470 | let records = match manager.list_automations() { |
| 471 | Ok(records) => records, |
| 472 | Err(err) => { |
| 473 | tracing::warn!(error = %err, "automation panel refresh could not list automations"); |
| 474 | return AutomationScan::default(); |
| 475 | } |
| 476 | }; |
| 477 | let mut runs = Vec::new(); |
| 478 | for record in &records { |
| 479 | let limit = if full { |
| 480 | None |
| 481 | } else { |
| 482 | Some(AUTOMATION_PANEL_RUN_SCAN) |
| 483 | }; |
| 484 | match manager.list_runs(&record.id, limit) { |
| 485 | Ok(recent) => runs.extend(recent), |
| 486 | Err(err) => { |
| 487 | tracing::warn!(automation_id = %record.id, error = %err, "automation panel refresh could not list runs"); |
| 488 | } |
| 489 | } |
| 490 | if !full { |
| 491 | let wanted: std::collections::BTreeSet<String> = live_owners |
| 492 | .iter() |
| 493 | .filter(|(_, owner)| owner.as_str() == record.id.as_str()) |
| 494 | .map(|(run_id, _)| run_id.clone()) |
| 495 | .collect(); |
| 496 | if !wanted.is_empty() { |
| 497 | match manager.get_runs_by_ids(&record.id, &wanted) { |
| 498 | Ok(found) => runs.extend(found), |
| 499 | Err(err) => { |
| 500 | tracing::warn!(automation_id = %record.id, error = %err, "automation panel refresh could not re-read live runs"); |
| 501 | } |
| 502 | } |
| 503 | } |
| 504 | } |
| 505 | } |
| 506 | // A re-read live run may also sit inside the newest window; the |
| 507 | // fold counts by id, so dedupe before handing the scan over. |
| 508 | let mut seen = std::collections::BTreeSet::new(); |
| 509 | runs.retain(|run| seen.insert(run.id.clone())); |
| 510 | AutomationScan { records, runs } |
| 511 | })) |
| 512 | } |
| 513 | |
| 514 | /// Fold a finished scan into the projection and post the typed receipt |
| 515 | /// (spec §2.2) for every run this session watched go live and settle — the |
| 516 | /// transcript learns about background work from the same scan that |
| 517 | /// repaints the band. |
| 518 | fn fold_automation_scan(app: &mut App, scan: &AutomationScan) -> bool { |
| 519 | let session_started_at = app.session_started_at; |
| 520 | let delta = app |
| 521 | .automation_panel |
| 522 | .fold_scan(&scan.records, &scan.runs, session_started_at); |
| 523 | let locale = app.ui_locale; |
| 524 | for run in &delta.settled { |
| 525 | app.add_message(crate::tui::automation_routing::settled_run_receipt( |
| 526 | locale, run, |
| 527 | )); |
| 528 | } |
| 529 | delta.changed || !delta.settled.is_empty() |
| 530 | } |
| 531 | |
| 532 | pub(super) fn refresh_shell_exec_live_output(app: &mut App) -> bool { |
| 533 | let Some(shell_mgr) = app.runtime_services.shell_manager.as_ref().cloned() else { |
| 534 | return false; |
| 535 | }; |
| 536 | // #3804: render-only read — try_lock so a contended shell Mutex can never |
| 537 | // block the async UI loop; skip this frame's live-output update on |
| 538 | // contention (the next refresh picks it up). |
| 539 | let jobs = { |
| 540 | let Ok(mut mgr) = shell_mgr.try_lock() else { |
| 541 | return false; |
| 542 | }; |
| 543 | mgr.list_jobs_for_session(app.current_session_id.as_deref().unwrap_or_default()) |
| 544 | .into_iter() |
| 545 | .map(|job| (job.id.clone(), job)) |
| 546 | .collect::<std::collections::HashMap<_, _>>() |
| 547 | }; |
| 548 | let mut changed = false; |
| 549 | for index in 0..app.virtual_cell_count() { |
| 550 | let Some(ShellExecLiveUpdate { |
| 551 | task_id, |
| 552 | status: next_status, |
| 553 | output: next_live, |
| 554 | duration_ms: next_duration, |
| 555 | finalized, |
| 556 | stale_elapsed_since_output_ms, |
| 557 | }) = shell_exec_live_update(app, index, &jobs) |
| 558 | else { |
| 559 | continue; |
| 560 | }; |
| 561 | let Some(HistoryCell::Tool(ToolCell::Exec(exec))) = app.cell_at_virtual_index_mut(index) |
| 562 | else { |
| 563 | continue; |
| 564 | }; |
| 565 | if exec.output.is_some() || exec.shell_task_id.as_deref() != Some(task_id.as_str()) { |
| 566 | continue; |
| 567 | } |
| 568 | exec.status = next_status; |
| 569 | exec.duration_ms = Some(next_duration); |
| 570 | exec.stale_elapsed_since_output_ms = stale_elapsed_since_output_ms; |
| 571 | if finalized { |
| 572 | exec.output = next_live; |
| 573 | exec.output_summary = exec |
| 574 | .output |
| 575 | .as_deref() |
| 576 | .map(crate::tui::history::summarize_tool_output); |
| 577 | exec.live_output = None; |
| 578 | exec.stale_elapsed_since_output_ms = None; |
| 579 | } else { |
| 580 | exec.live_output = next_live; |
| 581 | } |
| 582 | changed = true; |
| 583 | } |
| 584 | changed |
| 585 | } |
| 586 | |
| 587 | pub(super) struct ShellExecLiveUpdate { |
| 588 | pub(super) task_id: String, |
| 589 | pub(super) status: ToolStatus, |
| 590 | pub(super) output: Option<String>, |
| 591 | pub(super) duration_ms: u64, |
| 592 | pub(super) finalized: bool, |
| 593 | pub(super) stale_elapsed_since_output_ms: Option<u64>, |
| 594 | } |
| 595 | |
| 596 | pub(super) fn shell_exec_live_update( |
| 597 | app: &App, |
| 598 | index: usize, |
| 599 | jobs: &std::collections::HashMap<String, ShellJobSnapshot>, |
| 600 | ) -> Option<ShellExecLiveUpdate> { |
| 601 | let HistoryCell::Tool(ToolCell::Exec(exec)) = app.cell_at_virtual_index(index)? else { |
| 602 | return None; |
| 603 | }; |
| 604 | if exec.output.is_some() { |
| 605 | return None; |
| 606 | } |
| 607 | let task_id = exec.shell_task_id.as_deref()?; |
| 608 | let Some(job) = jobs.get(task_id) else { |
| 609 | return Some(ShellExecLiveUpdate { |
| 610 | task_id: task_id.to_string(), |
| 611 | status: ToolStatus::Failed, |
| 612 | output: detached_shell_job_output(task_id, exec), |
| 613 | duration_ms: exec.duration_ms.unwrap_or_default(), |
| 614 | finalized: true, |
| 615 | stale_elapsed_since_output_ms: None, |
| 616 | }); |
| 617 | }; |
| 618 | let next_status = shell_job_tool_status(&job.status); |
| 619 | let next_live = shell_job_live_output(job).or_else(|| exec.live_output.clone()); |
| 620 | let finalized = !matches!(job.status, ShellStatus::Running); |
| 621 | let stale_elapsed_since_output_ms = if matches!(job.status, ShellStatus::Running) && job.stale { |
| 622 | Some(job.elapsed_since_output_ms.unwrap_or(0)) |
| 623 | } else { |
| 624 | None |
| 625 | }; |
| 626 | if exec.status == next_status |
| 627 | && exec.live_output == next_live |
| 628 | && exec.duration_ms == Some(job.elapsed_ms) |
| 629 | && exec.stale_elapsed_since_output_ms == stale_elapsed_since_output_ms |
| 630 | { |
| 631 | return None; |
| 632 | } |
| 633 | Some(ShellExecLiveUpdate { |
| 634 | task_id: task_id.to_string(), |
| 635 | status: next_status, |
| 636 | output: next_live, |
| 637 | duration_ms: job.elapsed_ms, |
| 638 | finalized, |
| 639 | stale_elapsed_since_output_ms, |
| 640 | }) |
| 641 | } |
| 642 | |
| 643 | pub(super) fn detached_shell_job_output(task_id: &str, exec: &ExecCell) -> Option<String> { |
| 644 | let mut output = exec.live_output.clone().unwrap_or_default(); |
| 645 | if !output.trim().is_empty() { |
| 646 | output.push_str("\n\n"); |
| 647 | } |
| 648 | output.push_str(&format!( |
| 649 | "Shell job `{task_id}` is no longer attached to this TUI session." |
| 650 | )); |
| 651 | Some(output) |
| 652 | } |
| 653 | |
| 654 | pub(super) fn shell_job_tool_status(status: &ShellStatus) -> ToolStatus { |
| 655 | match status { |
| 656 | ShellStatus::Running => ToolStatus::Running, |
| 657 | ShellStatus::Completed => ToolStatus::Success, |
| 658 | ShellStatus::Failed | ShellStatus::Killed | ShellStatus::TimedOut => ToolStatus::Failed, |
| 659 | } |
| 660 | } |
| 661 | |
| 662 | pub(super) fn shell_job_live_output(job: &ShellJobSnapshot) -> Option<String> { |
| 663 | match (job.stdout_tail.is_empty(), job.stderr_tail.is_empty()) { |
| 664 | (true, true) => None, |
| 665 | (false, true) => Some(job.stdout_tail.clone()), |
| 666 | (true, false) => Some(format!("STDERR:\n{}", job.stderr_tail)), |
| 667 | (false, false) => Some(format!( |
| 668 | "{}\n\nSTDERR:\n{}", |
| 669 | job.stdout_tail, job.stderr_tail |
| 670 | )), |
| 671 | } |
| 672 | } |
| 673 | |
| 674 | pub(super) fn active_rlm_task_entries(app: &App) -> Vec<TaskPanelEntry> { |
| 675 | let Some(active) = app.active_cell.as_ref() else { |
| 676 | return Vec::new(); |
| 677 | }; |
| 678 | let duration_ms = app |
| 679 | .turn_started_at |
| 680 | .map(|started| u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX)); |
| 681 | active |
| 682 | .entries() |
| 683 | .iter() |
| 684 | .enumerate() |
| 685 | .filter_map(|(idx, entry)| { |
| 686 | let HistoryCell::Tool(ToolCell::Generic(generic)) = entry else { |
| 687 | return None; |
| 688 | }; |
| 689 | if !matches!( |
| 690 | generic.name.as_str(), |
| 691 | "rlm_open" | "rlm_eval" | "rlm_configure" | "rlm_close" | "rlm" |
| 692 | ) || generic.status != ToolStatus::Running |
| 693 | { |
| 694 | return None; |
| 695 | } |
| 696 | let summary = generic |
| 697 | .input_summary |
| 698 | .as_deref() |
| 699 | .filter(|summary| !summary.trim().is_empty()) |
| 700 | .unwrap_or("running chunked analysis"); |
| 701 | Some(TaskPanelEntry { |
| 702 | exit_code: None, |
| 703 | id: format!("rlm-{}", idx + 1), |
| 704 | status: "running".to_string(), |
| 705 | prompt_summary: format!("RLM: {summary}"), |
| 706 | duration_ms, |
| 707 | kind: TaskPanelEntryKind::Background, |
| 708 | stale: false, |
| 709 | elapsed_since_output_ms: None, |
| 710 | owner_agent_id: None, |
| 711 | owner_agent_name: None, |
| 712 | current_tool: None, |
| 713 | role: None, |
| 714 | files_touched: 0, |
| 715 | }) |
| 716 | }) |
| 717 | .collect() |
| 718 | } |
| 719 | |
| 720 | #[cfg(test)] |
| 721 | mod automation_scan_gate_tests { |
| 722 | use super::*; |
| 723 | use std::time::Duration; |
| 724 | |
| 725 | fn app() -> App { |
| 726 | crate::test_support::test_app_with_options(crate::tui::app::TuiOptions { |
| 727 | skip_onboarding: true, |
| 728 | ..crate::test_support::test_tui_options(std::path::PathBuf::from(".")) |
| 729 | }) |
| 730 | } |
| 731 | |
| 732 | #[test] |
| 733 | fn a_quiet_session_with_nothing_scheduled_rescans_on_the_long_interval() { |
| 734 | let app = app(); |
| 735 | assert_eq!( |
| 736 | automation_scan_interval(&app, true), |
| 737 | AUTOMATION_SCAN_QUIET_INTERVAL |
| 738 | ); |
| 739 | assert!(!automation_scan_is_due(&app, true, Duration::from_secs(14))); |
| 740 | assert!(automation_scan_is_due(&app, true, Duration::from_secs(15))); |
| 741 | } |
| 742 | |
| 743 | #[test] |
| 744 | fn a_session_in_use_keeps_the_busy_cadence() { |
| 745 | let app = app(); |
| 746 | assert_eq!( |
| 747 | automation_scan_interval(&app, false), |
| 748 | AUTOMATION_SCAN_BUSY_INTERVAL |
| 749 | ); |
| 750 | assert!(!automation_scan_is_due( |
| 751 | &app, |
| 752 | false, |
| 753 | Duration::from_millis(2_499) |
| 754 | )); |
| 755 | assert!(automation_scan_is_due( |
| 756 | &app, |
| 757 | false, |
| 758 | Duration::from_millis(2_500) |
| 759 | )); |
| 760 | } |
| 761 | |
| 762 | #[test] |
| 763 | fn scheduled_work_in_the_band_keeps_the_busy_cadence_even_when_quiet() { |
| 764 | let mut app = app(); |
| 765 | app.automation_panel.active_automations = 1; |
| 766 | assert_eq!( |
| 767 | automation_scan_interval(&app, true), |
| 768 | AUTOMATION_SCAN_BUSY_INTERVAL, |
| 769 | "an active automation can start a run at any moment" |
| 770 | ); |
| 771 | |
| 772 | let mut app = self::app(); |
| 773 | app.automation_panel.live_runs = 1; |
| 774 | assert_eq!( |
| 775 | automation_scan_interval(&app, true), |
| 776 | AUTOMATION_SCAN_BUSY_INTERVAL, |
| 777 | "a live run settles into a receipt the scan produces" |
| 778 | ); |
| 779 | } |
| 780 | |
| 781 | #[test] |
| 782 | fn a_finished_scan_is_folded_without_starting_the_next_one() { |
| 783 | let rt = tokio::runtime::Builder::new_current_thread() |
| 784 | .enable_all() |
| 785 | .build() |
| 786 | .expect("runtime"); |
| 787 | let mut app = app(); |
| 788 | rt.block_on(async { |
| 789 | let scan = |
| 790 | tokio::spawn(async { crate::tui::automation_panel::AutomationScan::default() }); |
| 791 | while !scan.is_finished() { |
| 792 | tokio::task::yield_now().await; |
| 793 | } |
| 794 | app.automation_scan = Some(scan); |
| 795 | // Not due: the scan is folded and nothing replaces it, so a |
| 796 | // quiet session does not chain scan after scan. |
| 797 | let changed = refresh_automation_panel(&mut app, false).await; |
| 798 | assert!(!changed, "an empty store changes nothing visible"); |
| 799 | }); |
| 800 | assert!(app.automation_scan.is_none()); |
| 801 | } |
| 802 | } |
| 803 |