Skip to main content

beam_worker/worker_runtime/
run_loop.rs

1use super::*;
2
3/// Returns `Some(("TERM", "xterm-256color"))` when the CLI's spec sets
4/// `inject_term_xterm` (codex/traex) and the inherited `TERM` is missing,
5/// empty, or `"dumb"`.
6///
7/// For any other CLI, or when `TERM` is already set to a non-empty value other
8/// than `"dumb"`, returns `None` — the environment is never overwritten.
9///
10/// This is a pure function so it can be tested deterministically without
11/// touching the real process environment.
12pub(crate) fn maybe_inject_term(
13    cli_id: &str,
14    current_term: Option<&str>,
15) -> Option<(String, String)> {
16    let inject = beam_core::cli_specs::cli_spec(cli_id)
17        .map(|spec| spec.inject_term_xterm)
18        .unwrap_or(false);
19    if !inject {
20        return None;
21    }
22    match current_term {
23        None | Some("") | Some("dumb") => Some(("TERM".to_string(), "xterm-256color".to_string())),
24        _ => None,
25    }
26}
27
28/// Consecutive dead `is_alive` samples required before emitting `CliExit`.
29/// A single false negative (empty `list-panes`, transient pane hide) must
30/// not tear the worker down.
31pub(crate) const CLI_DEAD_CONFIRM_ATTEMPTS: usize = 3;
32
33/// Count consecutive dead liveness samples. Returns true when the worker
34/// should emit `CliExit`.
35pub(crate) fn should_emit_cli_exit(consecutive_dead: &mut usize, alive: bool) -> bool {
36    if alive {
37        *consecutive_dead = 0;
38        return false;
39    }
40    *consecutive_dead = consecutive_dead.saturating_add(1);
41    *consecutive_dead >= CLI_DEAD_CONFIRM_ATTEMPTS
42}
43
44pub(crate) async fn send_message(
45    stdout: &Arc<Mutex<tokio::io::Stdout>>,
46    msg: &WorkerToDaemon,
47) -> Result<()> {
48    let mut out = stdout.lock().await;
49    out.write_all(serde_json::to_string(msg)?.as_bytes())
50        .await?;
51    out.write_all(b"\n").await?;
52    out.flush().await?;
53    Ok(())
54}
55
56pub async fn run(init: InitConfig) -> Result<()> {
57    let stdout = Arc::new(Mutex::new(tokio::io::stdout()));
58    let session_name = format!("beam-{}", &init.session_id[..8.min(init.session_id.len())]);
59    let paths = BeamPaths::discover()?;
60    let adapter = Arc::new(Mutex::new(CliAdapter::from_init(&init)?));
61    // Backend selection: Herdr adopted sessions observe a user pane; Herdr
62    // managed sessions spawn into a labeled workspace; everything else keeps
63    // the zellij path unchanged.
64    let (backend_impl, attach_context, herdr_handle) =
65        crate::backend::select::select_backend(&init, &session_name);
66    let spawn_spec = adapter.lock().await.build_spawn_spec(&init);
67    let extra_env = maybe_inject_term(&init.cli_id, std::env::var("TERM").ok().as_deref())
68        .into_iter()
69        .collect();
70    let exe_path = std::env::current_exe()
71        .ok()
72        .map(|p| p.display().to_string())
73        .unwrap_or_else(|| "beam".to_string());
74    let bindir = std::env::current_exe()
75        .ok()
76        .and_then(|p| p.parent().map(|v| v.display().to_string()))
77        .unwrap_or_default();
78    let args = if init.adopted_from.is_some() {
79        (spawn_spec.bin, spawn_spec.args)
80    } else {
81        if super::launch::normalize_cgroup_slice(init.cgroup_slice.as_deref()).is_some()
82            && super::launch::current_launch_platform() == super::launch::LaunchPlatform::Linux
83            && !super::launch::systemd_run_available()
84        {
85            anyhow::bail!(
86                "cgroupSlice requires systemd-run on Linux (user systemd); install systemd or unset cgroupSlice"
87            );
88        }
89        super::launch::build_launch_spec(
90            super::launch::current_launch_platform(),
91            super::launch::LaunchInput {
92                cgroup_slice: init.cgroup_slice.clone(),
93                session_id: init.session_id.clone(),
94                beam_home: paths.root().display().to_string(),
95                beam_bin: exe_path,
96                path_prepend: bindir,
97                extra_env,
98                spec: spawn_spec,
99            },
100        )?
101    };
102    let spawn_opts = SpawnOpts {
103        cwd: init.working_dir.clone(),
104        cols: DEFAULT_TERMINAL_COLS,
105        rows: DEFAULT_TERMINAL_ROWS,
106        env: Vec::new(),
107    };
108    backend_impl
109        .spawn(&args.0, &args.1, spawn_opts)
110        .await
111        .with_context(|| format!("failed to {} session {}", attach_context, init.session_id))?;
112    // Shared handle: the backend synchronizes internally per operation, so
113    // no outer Mutex is needed and a long write_input() never blocks screen
114    // capture, terminal keys, or the screenshot coordinator.
115    let backend = backend_impl;
116    let (herdr_workspace_id, herdr_pane_id) = herdr_handle
117        .as_ref()
118        .and_then(|h| h.herdr_ids())
119        .map(|ids| (Some(ids.workspace_id), Some(ids.pane_id)))
120        .unwrap_or((None, None));
121    let mut cli_pid_marker = None;
122    let child_pid = backend.child_pid().await?;
123    adapter.lock().await.on_spawned(child_pid);
124    if let Some(pid) = child_pid {
125        tokio::fs::create_dir_all(paths.cli_pid_markers_dir()).await?;
126        let marker = paths.cli_pid_markers_dir().join(pid.to_string());
127        tokio::fs::write(&marker, init.session_id.as_bytes()).await?;
128        cli_pid_marker = Some(marker);
129    }
130    // TUI CLIs drop keystrokes typed before their input UI is up. Wait for
131    // the CLI's ready marker (kimi's welcome screen, or a generic "welcome"
132    // for TUIs without a known one) before signaling Ready, so the initial
133    // prompt and the first stdin message land on a live input field. Adopted
134    // sessions attach to an already-running CLI; there is nothing to wait for.
135    if init.adopted_from.is_none()
136        && let Some(marker) = crate::adapters::tui_ready_marker(&init.cli_id)
137    {
138        let ready = wait_for_tui_ready(backend.as_ref(), marker).await;
139        if ready {
140            info!(session = %init.session_id, adapter = %init.cli_id, marker, "CLI TUI ready marker observed");
141        } else {
142            warn!(session = %init.session_id, adapter = %init.cli_id, marker, "CLI TUI ready marker not observed within {}s; typing input anyway", TUI_READY_TIMEOUT.as_secs());
143        }
144    }
145    let latest_screen = Arc::new(RwLock::new(String::new()));
146    let latest_raw_screen = Arc::new(RwLock::new(String::new()));
147    let display_mode = Arc::new(RwLock::new(DisplayMode::Hidden));
148    let analyzer_runtime = Arc::new(RwLock::new(AnalyzerRuntime::default()));
149    let usage_limit_tracker = Arc::new(Mutex::new(UsageLimitTracker::default()));
150    let current_turn_id = Arc::new(RwLock::new(String::new()));
151    let (updates, _) = broadcast::channel::<String>(256);
152
153    send_message(
154        &stdout,
155        &WorkerToDaemon::Ready {
156            zellij_session: session_name.clone(),
157            backend_kind: init.backend_kind,
158            herdr_workspace_id,
159            herdr_pane_id,
160        },
161    )
162    .await?;
163
164    let sample_backend = backend.clone();
165    let sample_screen = latest_screen.clone();
166    let sample_raw_screen = latest_raw_screen.clone();
167    let sample_updates = updates.clone();
168    let sample_stdout = stdout.clone();
169    let sample_adapter = adapter.clone();
170    let sample_display_mode = display_mode.clone();
171    let sample_usage_limit_tracker = usage_limit_tracker.clone();
172    let sample_current_turn_id = current_turn_id.clone();
173    let last_uploaded_hash = Arc::new(Mutex::new(None::<String>));
174    let last_broadcast_hash: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
175    let sample_last_broadcast_hash = last_broadcast_hash.clone();
176    let sample_analyzer_runtime = analyzer_runtime.clone();
177    // Channel for trigger events (TurnStarted, PaneUpdate, …) to the screenshot coordinator.
178    let (trigger_tx, trigger_rx) = mpsc::channel::<Trigger>(16);
179    let screen_capture_task = tokio::spawn(async move {
180        let mut last_emitted_status = ScreenStatus::Starting;
181        let mut last_emitted_usage_limit: Option<CliUsageLimitState> = None;
182        let mut consecutive_dead = 0usize;
183        loop {
184            let (screen, alive) = {
185                let screen = sample_backend.capture_viewport().await.unwrap_or_default();
186                // Unknown/error is alive: a wedged probe must not emit CliExit.
187                let alive = sample_backend.is_alive().await.unwrap_or(true);
188                (screen, alive)
189            };
190
191            let hash_changed;
192
193            {
194                *sample_raw_screen.write().await = screen.clone();
195                let mode = *sample_display_mode.read().await;
196                let rendered = render_screen_for_display_mode(&screen, mode);
197                let now_ms = now_ms();
198                let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
199                let base_status = if analyzing {
200                    ScreenStatus::Analyzing
201                } else {
202                    ScreenStatus::Working
203                };
204                let (status, usage_limit) =
205                    sample_usage_limit_tracker
206                        .lock()
207                        .await
208                        .classify(&screen, base_status, now_ms);
209                let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
210                {
211                    let guard = sample_last_broadcast_hash.lock().await;
212                    hash_changed = guard.as_deref() != Some(&rendered_hash);
213                }
214
215                if hash_changed
216                    || last_emitted_status != status
217                    || last_emitted_usage_limit != usage_limit
218                {
219                    *sample_last_broadcast_hash.lock().await = Some(rendered_hash.clone());
220                    let mut current = sample_screen.write().await;
221                    *current = rendered.clone();
222                    let _ = sample_updates.send(rendered.clone());
223                    let _ = send_message(
224                        &sample_stdout,
225                        &WorkerToDaemon::ScreenUpdate {
226                            content: rendered.clone(),
227                            status,
228                            usage_limit: usage_limit.clone(),
229                        },
230                    )
231                    .await;
232                    last_emitted_status = status;
233                    last_emitted_usage_limit = usage_limit.clone();
234                }
235            }
236
237            if let Ok(poll) = sample_adapter.lock().await.poll() {
238                if let Some(cli_session_id) = poll.cli_session_id {
239                    let _ = send_message(
240                        &sample_stdout,
241                        &WorkerToDaemon::CliSessionId { cli_session_id },
242                    )
243                    .await;
244                }
245                if let Some((user_text, assistant_text)) = poll.adopt_preamble {
246                    let _ = send_message(
247                        &sample_stdout,
248                        &WorkerToDaemon::AdoptPreamble {
249                            user_text,
250                            assistant_text,
251                        },
252                    )
253                    .await;
254                }
255                if let Some(content) = poll.final_output {
256                    let turn_id = sample_current_turn_id.read().await.clone();
257                    let _ = send_message(
258                        &sample_stdout,
259                        &WorkerToDaemon::FinalOutput {
260                            content,
261                            turn_id,
262                            kind: poll.final_output_kind,
263                            user_text: poll.final_output_user_text,
264                        },
265                    )
266                    .await;
267                }
268                if poll.prompt_ready {
269                    let _ = send_message(&sample_stdout, &WorkerToDaemon::PromptReady).await;
270                    let rendered = sample_screen.read().await.clone();
271                    let raw = sample_raw_screen.read().await.clone();
272                    let now_ms = now_ms();
273                    let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
274                    let base_status = if analyzing {
275                        ScreenStatus::Analyzing
276                    } else {
277                        ScreenStatus::Idle
278                    };
279                    let (status, usage_limit) =
280                        sample_usage_limit_tracker
281                            .lock()
282                            .await
283                            .classify(&raw, base_status, now_ms);
284                    let _ = send_message(
285                        &sample_stdout,
286                        &WorkerToDaemon::ScreenUpdate {
287                            content: rendered,
288                            status,
289                            usage_limit,
290                        },
291                    )
292                    .await;
293                }
294            }
295
296            if should_emit_cli_exit(&mut consecutive_dead, alive) {
297                warn!(
298                    consecutive_dead,
299                    "cli liveness dead confirmed, emitting CliExit"
300                );
301                let _ = send_message(
302                    &sample_stdout,
303                    &WorkerToDaemon::CliExit {
304                        code: Some(0),
305                        signal: None,
306                    },
307                )
308                .await;
309                break;
310            }
311            if !alive {
312                warn!(
313                    consecutive_dead,
314                    confirm_after = CLI_DEAD_CONFIRM_ATTEMPTS,
315                    "cli liveness dead, waiting for confirmation"
316                );
317            }
318
319            tokio::time::sleep(Duration::from_millis(5000)).await;
320        }
321    });
322    let mut worker_joins = tokio::task::JoinSet::new();
323    worker_joins.spawn(async move {
324        let _ = screen_capture_task.await;
325    });
326
327    if init.cli_id == "grok" {
328        let grok_raw_screen = latest_raw_screen.clone();
329        let grok_runtime_state = analyzer_runtime.clone();
330        let grok_stdout = stdout.clone();
331        worker_joins.spawn(async move {
332            loop {
333                tokio::time::sleep(Duration::from_millis(800)).await;
334                let snapshot = grok_raw_screen.read().await.clone();
335                if snapshot.is_empty() {
336                    continue;
337                }
338                let options =
339                    crate::worker_runtime::grok_prompts::detect_grok_plan_approval(&snapshot);
340                let mut runtime = grok_runtime_state.write().await;
341                match options {
342                    Some(options) if !runtime.prompt_active => {
343                        runtime.prompt_active = true;
344                        let _ = send_message(
345                            &grok_stdout,
346                            &WorkerToDaemon::TuiPrompt {
347                                description: "Grok plan approval".to_string(),
348                                options,
349                                multi_select: false,
350                            },
351                        )
352                        .await;
353                    }
354                    None if runtime.prompt_active => {
355                        runtime.prompt_active = false;
356                        let _ = send_message(
357                            &grok_stdout,
358                            &WorkerToDaemon::TuiPromptResolved {
359                                selected_text: None,
360                            },
361                        )
362                        .await;
363                    }
364                    _ => {}
365                }
366            }
367        });
368    }
369
370    if screen_analyzer_enabled(&init.screen_analyzer) {
371        let analyzer_cfg = init.screen_analyzer.clone();
372        let analyzer_raw_screen = latest_raw_screen.clone();
373        let analyzer_runtime_state = analyzer_runtime.clone();
374        let analyzer_stdout = stdout.clone();
375        let analyzer_cli_id = init.cli_id.clone();
376        let analyzer_task = tokio::spawn(async move {
377            let client = Client::new();
378            loop {
379                tokio::time::sleep(Duration::from_millis(analyzer_cfg.interval_ms)).await;
380                let snapshot = analyzer_raw_screen.read().await.clone();
381                if snapshot.is_empty() {
382                    continue;
383                }
384                let truncated = if snapshot.len() > analyzer_cfg.snapshot_max_chars {
385                    snapshot[snapshot.len() - analyzer_cfg.snapshot_max_chars..].to_string()
386                } else {
387                    snapshot
388                };
389                let now = now_ms();
390                {
391                    let mut runtime = analyzer_runtime_state.write().await;
392                    if truncated == runtime.last_snapshot {
393                        runtime.stable_count = runtime.stable_count.saturating_add(1);
394                    } else {
395                        runtime.stable_count = 1;
396                        runtime.last_snapshot = truncated.clone();
397                        if runtime.waiting_for_content_change {
398                            runtime.waiting_for_content_change = false;
399                        }
400                    }
401                    if runtime.stable_count < analyzer_cfg.stable_count {
402                        continue;
403                    }
404                    if runtime.waiting_for_content_change
405                        && truncated == runtime.last_analyzed_snapshot
406                    {
407                        continue;
408                    }
409                    if runtime.cooldown_until_ms > now {
410                        continue;
411                    }
412                    runtime.is_analyzing = true;
413                    runtime.last_analyzed_snapshot = truncated.clone();
414                }
415
416                if analyzer_cli_id == "grok"
417                    && crate::worker_runtime::grok_prompts::detect_grok_plan_approval(&truncated)
418                        .is_some()
419                {
420                    let mut runtime = analyzer_runtime_state.write().await;
421                    runtime.is_analyzing = false;
422                    continue;
423                }
424
425                let result = call_screen_analyzer(&client, &analyzer_cfg, &truncated).await;
426
427                let mut runtime = analyzer_runtime_state.write().await;
428                runtime.is_analyzing = false;
429                match result {
430                    Ok(analysis) => {
431                        apply_screen_analyzer_result(
432                            &mut runtime,
433                            &analysis.check_again_when,
434                            now_ms(),
435                        );
436                        if analysis.needs_interaction && !analysis.options.is_empty() {
437                            if !runtime.prompt_active {
438                                runtime.prompt_active = true;
439                                let _ = send_message(
440                                    &analyzer_stdout,
441                                    &WorkerToDaemon::TuiPrompt {
442                                        description: analysis.description.clone().unwrap_or_else(
443                                            || "CLI needs your selection".to_string(),
444                                        ),
445                                        options: analysis.options.clone(),
446                                        multi_select: analysis.multi_select,
447                                    },
448                                )
449                                .await;
450                            }
451                        } else if runtime.prompt_active {
452                            runtime.prompt_active = false;
453                            let _ = send_message(
454                                &analyzer_stdout,
455                                &WorkerToDaemon::TuiPromptResolved {
456                                    selected_text: None,
457                                },
458                            )
459                            .await;
460                        }
461                    }
462                    Err(_) => {
463                        runtime.waiting_for_content_change = true;
464                        runtime.cooldown_until_ms = 0;
465                    }
466                }
467            }
468        });
469        worker_joins.spawn(async move {
470            let _ = analyzer_task.await;
471        });
472    }
473
474    // Coordinator task — owns the trigger receiver and the 5-second fallback.
475    {
476        let coord_backend = backend.clone();
477        let coord_stdout = stdout.clone();
478        let coord_session_id = init.session_id.clone();
479        let coord_app_id = init.lark_app_id.clone();
480        let coord_app_secret = init.lark_app_secret.clone();
481        let coord_display_mode = display_mode.clone();
482        let coord_analyzer_runtime = analyzer_runtime.clone();
483        let coord_usage_limit_tracker = usage_limit_tracker.clone();
484        let coord_last_uploaded_hash = last_uploaded_hash.clone();
485        let coord_latest_raw_screen = latest_raw_screen.clone();
486        let coord_rx = trigger_rx;
487        worker_joins.spawn(async move {
488            coordinator_loop(
489                coord_backend,
490                coord_stdout,
491                coord_session_id,
492                coord_app_id,
493                coord_app_secret,
494                coord_display_mode,
495                coord_analyzer_runtime,
496                coord_usage_limit_tracker,
497                coord_last_uploaded_hash,
498                coord_latest_raw_screen,
499                coord_rx,
500            )
501            .await;
502        });
503    }
504    // Subscribe task: forward backend pane-update notifications to the screenshot coordinator.
505    // Also caches the full viewport ANSI chunk so the coordinator can capture
506    // without waiting on a slow backend call inside write_input().
507    {
508        let sub_backend = backend.clone();
509        let sub_trigger_tx = trigger_tx.clone();
510        let sub_latest_raw_screen = latest_raw_screen.clone();
511        worker_joins.spawn(async move {
512            let mut rx = sub_backend.subscribe();
513            loop {
514                match rx.recv().await {
515                    Ok(chunk) => {
516                        // latest wins: the chunk is the full viewport (not incremental)
517                        *sub_latest_raw_screen.write().await = chunk;
518                        if let Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) =
519                            sub_trigger_tx.try_send(Trigger::PaneUpdate)
520                        {
521                            break;
522                        } // Ok or Full → discard, keep listening
523                    }
524                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
525                    Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
526                }
527            }
528        });
529    }
530    // Herdr agent state polling: emit MuxAgentState for side effects only
531    // (blocked → attention). v1 uses a bounded `agent get` poll.
532    if init.backend_kind == beam_core::BackendKind::Herdr
533        && let Some(handle) = herdr_handle.as_ref()
534        && let Some(pane_id) = handle.herdr_ids().map(|ids| ids.pane_id)
535    {
536        let agent_stdout = stdout.clone();
537        let agent_pane_id = pane_id.clone();
538        worker_joins.spawn(async move {
539            loop {
540                tokio::time::sleep(Duration::from_secs(2)).await;
541                let state = crate::backend::herdr::cli::agent_get(&agent_pane_id)
542                    .await
543                    .ok()
544                    .and_then(|payload| crate::backend::herdr::cli::parse_agent_state(&payload));
545                let Some(state) = state else {
546                    continue;
547                };
548                if send_message(
549                    &agent_stdout,
550                    &WorkerToDaemon::MuxAgentState {
551                        state,
552                        agent_name: None,
553                        pane_id: agent_pane_id.clone(),
554                        message: None,
555                    },
556                )
557                .await
558                .is_err()
559                {
560                    break;
561                }
562            }
563        });
564    }
565    if !init.prompt.is_empty() && !crate::adapters::passes_initial_prompt_via_args(&init.cli_id) {
566        usage_limit_tracker.lock().await.begin_turn(
567            "",
568            SystemTime::now()
569                .duration_since(UNIX_EPOCH)
570                .unwrap_or_default()
571                .as_millis() as u64,
572        );
573        *current_turn_id.write().await = init
574            .prompt_turn_id
575            .clone()
576            .unwrap_or_else(|| Uuid::new_v4().to_string());
577        *last_uploaded_hash.lock().await = None;
578        let submit = adapter
579            .lock()
580            .await
581            .write_input(backend.as_ref(), &init.prompt)
582            .await?;
583        if let Some(cli_session_id) = submit.cli_session_id {
584            send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
585        }
586        if !submit.submitted {
587            let message = submit
588                .failure_reason
589                .unwrap_or_else(|| "CLI submit could not be confirmed".to_string());
590            send_message(&stdout, &WorkerToDaemon::UserNotify { message }).await?;
591        }
592    }
593
594    // Init-time transcript source resolution for adapters with resolvable
595    // sources (opencode): resolve cli_session_id before the message loop.
596    {
597        let resolution = {
598            let mut adapter_guard = adapter.lock().await;
599            let resolution = adapter_guard
600                .resolve_transcript_source(backend.as_ref())
601                .await;
602            drop(adapter_guard);
603            resolution
604        };
605        if let Some(resolution) = resolution {
606            let outcome = resolution.unwrap_or_else(|err| {
607                warn!(session = %init.session_id, adapter = %init.cli_id, "resolve_transcript_source error: {:?}", err);
608                ResolveOutcome::NotFound {
609                    reason: format!("transcript source resolution failed: {}", err),
610                }
611            });
612            match outcome {
613                ResolveOutcome::Found(source) => {
614                    send_message(
615                        &stdout,
616                        &WorkerToDaemon::CliSessionId {
617                            cli_session_id: source.session_id.clone(),
618                        },
619                    )
620                    .await?;
621                    // Also set in adapter state for subsequent poll/write_input.
622                    adapter
623                        .lock()
624                        .await
625                        .set_transcript_source(&source.session_id);
626                    info!(
627                        session = %init.session_id, adapter = %init.cli_id,
628                        transcript_session = %source.session_id,
629                        "transcript source resolved automatically"
630                    );
631                }
632                ResolveOutcome::Ambiguous { candidates, .. } => {
633                    info!(
634                        session = %init.session_id, adapter = %init.cli_id,
635                        candidate_count = candidates.len(),
636                        "transcript source ambiguous, requesting user choice"
637                    );
638                    let turn_id = Uuid::new_v4().to_string();
639                    let choices: Vec<TranscriptChoice> = candidates
640                        .iter()
641                        .map(|c| TranscriptChoice {
642                            session_id: c.session_id.clone(),
643                            label: format!("{} ({})", c.session_id, c.db_path.display()),
644                        })
645                        .collect();
646                    send_message(
647                        &stdout,
648                        &WorkerToDaemon::TranscriptChoices {
649                            candidates: choices,
650                            turn_id: turn_id.clone(),
651                        },
652                    )
653                    .await?;
654                }
655                ResolveOutcome::NotFound { reason } => {
656                    info!(session = %init.session_id, adapter = %init.cli_id, "transcript source not found: {}", reason);
657                    send_message(&stdout, &WorkerToDaemon::UserNotify { message: reason }).await?;
658                }
659            }
660        }
661    }
662
663    // Shared "currently processing" stamp: Some((description, start_ms))
664    // while the message loop is busy handling one daemon message. Drives
665    // both the processing watchdog (WARN past the threshold) and the
666    // heartbeat IPC, so the daemon can tell "worker dead" apart from
667    // "worker stuck on a message".
668    let processing_since: Arc<StdMutex<Option<(String, u64)>>> = Arc::new(StdMutex::new(None));
669
670    // stdin reader runs on a dedicated OS thread: while the message loop is
671    // busy handling one message, subsequent daemon messages keep being
672    // drained from the pipe into this channel instead of piling up in the
673    // kernel pipe buffer.
674    let (daemon_msg_tx, mut daemon_msg_rx) = mpsc::channel::<DaemonToWorker>(32);
675    std::thread::spawn(move || {
676        use std::io::BufRead as _;
677        let stdin = std::io::stdin();
678        for line in stdin.lock().lines() {
679            let line = match line {
680                Ok(line) => line,
681                Err(_) => break,
682            };
683            if line.trim().is_empty() {
684                continue;
685            }
686            match serde_json::from_str::<DaemonToWorker>(&line) {
687                Ok(msg) => {
688                    if daemon_msg_tx.blocking_send(msg).is_err() {
689                        break;
690                    }
691                }
692                Err(err) => warn!("failed to parse daemon message: {}", err),
693            }
694        }
695    });
696
697    // Processing watchdog: WARN when one daemon message is being handled
698    // longer than the threshold (e.g. a wedged write_input), so the stuck
699    // message is directly locatable in the log.
700    const MESSAGE_PROCESSING_WARN_THRESHOLD: Duration = Duration::from_secs(120);
701    {
702        let watchdog_processing = processing_since.clone();
703        let watchdog_session_id = init.session_id.clone();
704        worker_joins.spawn(async move {
705            loop {
706                tokio::time::sleep(Duration::from_secs(15)).await;
707                let stuck = {
708                    let guard = watchdog_processing.lock().unwrap();
709                    guard.as_ref().and_then(|(desc, start_ms)| {
710                        let elapsed = now_ms().saturating_sub(*start_ms);
711                        (elapsed > MESSAGE_PROCESSING_WARN_THRESHOLD.as_millis() as u64)
712                            .then(|| (desc.clone(), elapsed))
713                    })
714                };
715                if let Some((desc, elapsed_ms)) = stuck {
716                    warn!(
717                        session = %watchdog_session_id,
718                        message = %desc,
719                        elapsed_ms,
720                        "daemon message processing exceeded {}s threshold",
721                        MESSAGE_PROCESSING_WARN_THRESHOLD.as_secs()
722                    );
723                }
724            }
725        });
726    }
727
728    // Heartbeat: independent of the message loop; lets the daemon
729    // distinguish a dead worker (no heartbeat) from a stuck one (heartbeat
730    // keeps coming with processing_since_ms set).
731    {
732        let hb_stdout = stdout.clone();
733        let hb_processing = processing_since.clone();
734        worker_joins.spawn(async move {
735            loop {
736                tokio::time::sleep(Duration::from_secs(10)).await;
737                let since = hb_processing
738                    .lock()
739                    .unwrap()
740                    .as_ref()
741                    .map(|(_, start)| *start);
742                if send_message(
743                    &hb_stdout,
744                    &WorkerToDaemon::Heartbeat {
745                        processing_since_ms: since,
746                    },
747                )
748                .await
749                .is_err()
750                {
751                    break;
752                }
753            }
754        });
755    }
756
757    while let Some(msg) = daemon_msg_rx.recv().await {
758        *processing_since.lock().unwrap() = Some((daemon_message_desc(&msg), now_ms()));
759        match msg {
760            DaemonToWorker::Message { content, turn_id } => {
761                info!(session = %init.session_id, %turn_id, "Message received, sending TurnStarted to coordinator");
762                handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
763                let snapshot = latest_raw_screen.read().await.clone();
764                usage_limit_tracker.lock().await.begin_turn(
765                    &snapshot,
766                    SystemTime::now()
767                        .duration_since(UNIX_EPOCH)
768                        .unwrap_or_default()
769                        .as_millis() as u64,
770                );
771                *current_turn_id.write().await = turn_id.clone();
772                *last_uploaded_hash.lock().await = None;
773                // Notify coordinator: a new turn has started (best-effort).
774                if trigger_tx
775                    .send(Trigger::TurnStarted { turn_id })
776                    .await
777                    .is_err()
778                {
779                    warn!("coordinator channel closed, TurnStarted not sent");
780                }
781                let submit = adapter
782                    .lock()
783                    .await
784                    .write_input(backend.as_ref(), &content)
785                    .await?;
786                if let Some(cli_session_id) = submit.cli_session_id {
787                    send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
788                }
789                if !submit.submitted {
790                    let message = submit
791                        .failure_reason
792                        .unwrap_or_else(|| "CLI submit could not be confirmed".to_string());
793                    send_message(&stdout, &WorkerToDaemon::UserNotify { message }).await?;
794                }
795            }
796            DaemonToWorker::RawInput { content, turn_id } => {
797                info!(session = %init.session_id, %turn_id, "RawInput received, sending TurnStarted to coordinator");
798                handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
799                let snapshot = latest_raw_screen.read().await.clone();
800                usage_limit_tracker.lock().await.begin_turn(
801                    &snapshot,
802                    SystemTime::now()
803                        .duration_since(UNIX_EPOCH)
804                        .unwrap_or_default()
805                        .as_millis() as u64,
806                );
807                *current_turn_id.write().await = turn_id.clone();
808                *last_uploaded_hash.lock().await = None;
809                // Notify coordinator: a new turn has started (best-effort).
810                if trigger_tx
811                    .send(Trigger::TurnStarted { turn_id })
812                    .await
813                    .is_err()
814                {
815                    warn!("coordinator channel closed, TurnStarted not sent");
816                }
817                backend.raw_input(&content).await?;
818            }
819            DaemonToWorker::Close => {
820                backend.destroy_session().await?;
821                break;
822            }
823            DaemonToWorker::Restart => {
824                backend.destroy_session().await?;
825                break;
826            }
827            DaemonToWorker::RefreshScreen => {
828                let screen = backend.capture_viewport().await?;
829                *latest_raw_screen.write().await = screen.clone();
830                let mode = *display_mode.read().await;
831                let rendered = render_screen_for_display_mode(&screen, mode);
832                let now_ms = now_ms();
833                let analyzing = analyzer_runtime.read().await.is_analyzing;
834                let base_status = if analyzing {
835                    ScreenStatus::Analyzing
836                } else {
837                    ScreenStatus::Working
838                };
839                let (status, usage_limit) =
840                    usage_limit_tracker
841                        .lock()
842                        .await
843                        .classify(&screen, base_status, now_ms);
844                *latest_screen.write().await = rendered.clone();
845                let _ = updates.send(rendered.clone());
846                send_message(
847                    &stdout,
848                    &WorkerToDaemon::ScreenUpdate {
849                        content: rendered.clone(),
850                        status,
851                        usage_limit: usage_limit.clone(),
852                    },
853                )
854                .await?;
855                let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
856                *last_broadcast_hash.lock().await = Some(rendered_hash);
857                // Notify coordinator: refresh request (best-effort).
858                if trigger_tx.send(Trigger::Refresh).await.is_err() {
859                    warn!("coordinator channel closed, Refresh not sent");
860                }
861            }
862            DaemonToWorker::SetDisplayMode { mode } => {
863                *display_mode.write().await = mode;
864                let raw = latest_raw_screen.read().await.clone();
865                let rendered = render_screen_for_display_mode(&raw, mode);
866                let now_ms = now_ms();
867                let analyzing = analyzer_runtime.read().await.is_analyzing;
868                let base_status = if analyzing {
869                    ScreenStatus::Analyzing
870                } else {
871                    ScreenStatus::Working
872                };
873                let (status, usage_limit) =
874                    usage_limit_tracker
875                        .lock()
876                        .await
877                        .classify(&raw, base_status, now_ms);
878                *latest_screen.write().await = rendered.clone();
879                let _ = updates.send(rendered.clone());
880                send_message(
881                    &stdout,
882                    &WorkerToDaemon::ScreenUpdate {
883                        content: rendered,
884                        status,
885                        usage_limit: usage_limit.clone(),
886                    },
887                )
888                .await?;
889                // Notify coordinator: display mode changed (best-effort).
890                if trigger_tx
891                    .send(Trigger::SetDisplayMode(mode))
892                    .await
893                    .is_err()
894                {
895                    warn!("coordinator channel closed, SetDisplayMode not sent");
896                }
897            }
898            DaemonToWorker::TermAction { key } => {
899                let keys = term_action_keys(key);
900                backend.send_special_keys(&keys).await?;
901            }
902            DaemonToWorker::SpecialKeys { keys } => {
903                backend.send_special_keys(&keys).await?;
904            }
905            DaemonToWorker::TuiKeys { keys, is_final } => {
906                handle_tui_keys(&backend, &analyzer_runtime, &keys, is_final).await?;
907            }
908            DaemonToWorker::TuiTextInput { keys, text } => {
909                handle_tui_text_input(&backend, &adapter, &analyzer_runtime, &keys, &text).await?;
910            }
911            DaemonToWorker::SetTranscriptSource { cli_session_id } => {
912                let applied = adapter.lock().await.set_transcript_source(&cli_session_id);
913                if applied {
914                    info!(session = %init.session_id, adapter = %init.cli_id, transcript_session = %cli_session_id, "transcript source set by user");
915                    send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
916                }
917            }
918            DaemonToWorker::Init(_) => {}
919        }
920        *processing_since.lock().unwrap() = None;
921    }
922
923    worker_joins.abort_all();
924    while worker_joins.join_next().await.is_some() {}
925    let _ = backend.kill().await;
926    if let Some(marker) = cli_pid_marker {
927        let _ = tokio::fs::remove_file(marker).await;
928    }
929    info!("worker exiting");
930    Ok(())
931}
932
933/// Short human-readable label for a daemon message, used by the processing
934/// watchdog and heartbeat stamp (avoids dumping full message contents).
935fn daemon_message_desc(msg: &DaemonToWorker) -> String {
936    match msg {
937        DaemonToWorker::Message { turn_id, .. } => format!("Message(turn_id={})", turn_id),
938        DaemonToWorker::RawInput { turn_id, .. } => format!("RawInput(turn_id={})", turn_id),
939        DaemonToWorker::Close => "Close".to_string(),
940        DaemonToWorker::Restart => "Restart".to_string(),
941        DaemonToWorker::SetDisplayMode { .. } => "SetDisplayMode".to_string(),
942        DaemonToWorker::TermAction { .. } => "TermAction".to_string(),
943        DaemonToWorker::SpecialKeys { .. } => "SpecialKeys".to_string(),
944        DaemonToWorker::TuiKeys { .. } => "TuiKeys".to_string(),
945        DaemonToWorker::TuiTextInput { .. } => "TuiTextInput".to_string(),
946        DaemonToWorker::RefreshScreen => "RefreshScreen".to_string(),
947        DaemonToWorker::SetTranscriptSource { .. } => "SetTranscriptSource".to_string(),
948        DaemonToWorker::Init(_) => "Init".to_string(),
949    }
950}
951
952#[cfg(test)]
953mod tests {
954    use super::*;
955
956    #[test]
957    fn daemon_message_desc_labels_variants_without_content() {
958        let msg = DaemonToWorker::Message {
959            content: "secret prompt".to_string(),
960            turn_id: "t-1".to_string(),
961        };
962        let desc = daemon_message_desc(&msg);
963        assert!(desc.contains("Message"));
964        assert!(desc.contains("t-1"));
965        assert!(!desc.contains("secret prompt"));
966        assert_eq!(daemon_message_desc(&DaemonToWorker::Close), "Close");
967        assert_eq!(
968            daemon_message_desc(&DaemonToWorker::RefreshScreen),
969            "RefreshScreen"
970        );
971    }
972
973    #[test]
974    fn cli_exit_requires_consecutive_dead_samples() {
975        assert_eq!(CLI_DEAD_CONFIRM_ATTEMPTS, 3);
976        let mut dead = 0;
977        assert!(!should_emit_cli_exit(&mut dead, false));
978        assert_eq!(dead, 1);
979        assert!(!should_emit_cli_exit(&mut dead, false));
980        assert_eq!(dead, 2);
981        assert!(should_emit_cli_exit(&mut dead, false));
982        assert_eq!(dead, 3);
983        assert!(should_emit_cli_exit(&mut dead, false));
984        assert_eq!(dead, 4);
985    }
986
987    #[test]
988    fn cli_exit_counter_resets_on_alive() {
989        let mut dead = 0;
990        assert!(!should_emit_cli_exit(&mut dead, false));
991        assert!(!should_emit_cli_exit(&mut dead, false));
992        assert!(!should_emit_cli_exit(&mut dead, true));
993        assert_eq!(dead, 0);
994        assert!(!should_emit_cli_exit(&mut dead, false));
995        assert_eq!(dead, 1);
996    }
997}