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