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
28pub(crate) async fn prepare_wrapper(
29    init: &InitConfig,
30    paths: &BeamPaths,
31) -> Result<std::path::PathBuf> {
32    tokio::fs::create_dir_all(paths.run_dir()).await?;
33    let wrapper = paths.worker_wrapper_sh(&init.session_id);
34    let exe_path = std::env::current_exe()
35        .ok()
36        .map(|p| p.display().to_string())
37        .unwrap_or_else(|| "beam".to_string());
38    let content = format!(
39        "#!/bin/sh\ncd {cwd}\nexport BEAM_SESSION_ID={sid}\nexport BEAM_HOME={home}\nexport BEAM_BIN={exe}\nif [ -n \"$PATH\" ]; then\n  export PATH={bindir}:$PATH\nelse\n  export PATH={bindir}\nfi\nexec \"$@\"\n",
40        cwd = shell_quote(&init.working_dir),
41        sid = shell_quote(&init.session_id),
42        home = shell_quote(&paths.root().display().to_string()),
43        exe = shell_quote(&exe_path),
44        bindir = shell_quote(
45            &std::env::current_exe()
46                .ok()
47                .and_then(|p| p.parent().map(|v| v.display().to_string()))
48                .unwrap_or_default()
49        ),
50    );
51    tokio::fs::write(&wrapper, content).await?;
52    #[cfg(unix)]
53    {
54        use std::os::unix::fs::PermissionsExt;
55        let perms = std::fs::Permissions::from_mode(0o755);
56        tokio::fs::set_permissions(&wrapper, perms).await?;
57    }
58    Ok(wrapper)
59}
60
61pub(crate) async fn send_message(
62    stdout: &Arc<Mutex<tokio::io::Stdout>>,
63    msg: &WorkerToDaemon,
64) -> Result<()> {
65    let mut out = stdout.lock().await;
66    out.write_all(serde_json::to_string(msg)?.as_bytes())
67        .await?;
68    out.write_all(b"\n").await?;
69    out.flush().await?;
70    Ok(())
71}
72
73pub async fn run(init: InitConfig) -> Result<()> {
74    let stdout = Arc::new(Mutex::new(tokio::io::stdout()));
75    let session_name = format!("beam-{}", &init.session_id[..8.min(init.session_id.len())]);
76    let paths = BeamPaths::discover()?;
77    let adapter = Arc::new(Mutex::new(CliAdapter::from_init(&init)?));
78    let wrapper = if init.resume || init.adopted_from.is_some() {
79        None
80    } else {
81        Some(prepare_wrapper(&init, &paths).await?)
82    };
83    let (mut backend_impl, attach_context): (Box<dyn SessionBackend>, &'static str) =
84        if let Some(adopted) = init.adopted_from.as_ref() {
85            if let Some(pane_id) = adopted.zellij_pane_id.clone() {
86                let session = adopted.zellij_session.clone().unwrap_or_else(|| {
87                    format!("beam-{}", &init.session_id[..8.min(init.session_id.len())])
88                });
89                let observe = ZellijObserveBackend::new(
90                    session,
91                    pane_id,
92                    u32::try_from(adopted.original_cli_pid).ok(),
93                );
94                (Box::new(observe), "observe")
95            } else {
96                let zellij = ZellijBackend::new(session_name.clone());
97                (Box::new(zellij), "spawn")
98            }
99        } else {
100            let zellij = ZellijBackend::new(session_name.clone());
101            (Box::new(zellij), "spawn")
102        };
103    let spawn_spec = adapter.lock().await.build_spawn_spec(&init);
104    let args = if let Some(wrapper) = wrapper {
105        let mut args = Vec::with_capacity(2 + init.cli_args.len());
106        args.push(wrapper.display().to_string());
107        args.push(spawn_spec.bin.clone());
108        args.extend(spawn_spec.args.clone());
109        ("/bin/sh".to_string(), args)
110    } else {
111        (spawn_spec.bin, spawn_spec.args)
112    };
113    let mut env = Vec::new();
114    if let Some((k, v)) = maybe_inject_term(&init.cli_id, std::env::var("TERM").ok().as_deref()) {
115        env.push((k, v));
116    }
117    let spawn_opts = SpawnOpts {
118        cwd: init.working_dir.clone(),
119        cols: DEFAULT_TERMINAL_COLS,
120        rows: DEFAULT_TERMINAL_ROWS,
121        env,
122    };
123    backend_impl
124        .spawn(&args.0, &args.1, spawn_opts)
125        .await
126        .with_context(|| format!("failed to {} session {}", attach_context, init.session_id))?;
127    let backend: Arc<Mutex<Box<dyn SessionBackend>>> = Arc::new(Mutex::new(backend_impl));
128    let mut cli_pid_marker = None;
129    let child_pid = backend.lock().await.child_pid().await?;
130    adapter.lock().await.on_spawned(child_pid);
131    if let Some(pid) = child_pid {
132        tokio::fs::create_dir_all(paths.cli_pid_markers_dir()).await?;
133        let marker = paths.cli_pid_markers_dir().join(pid.to_string());
134        tokio::fs::write(&marker, init.session_id.as_bytes()).await?;
135        cli_pid_marker = Some(marker);
136    }
137    let latest_screen = Arc::new(RwLock::new(String::new()));
138    let latest_raw_screen = Arc::new(RwLock::new(String::new()));
139    let display_mode = Arc::new(RwLock::new(DisplayMode::Hidden));
140    let analyzer_runtime = Arc::new(RwLock::new(AnalyzerRuntime::default()));
141    let usage_limit_tracker = Arc::new(Mutex::new(UsageLimitTracker::default()));
142    let current_turn_id = Arc::new(RwLock::new(String::new()));
143    let (updates, _) = broadcast::channel::<String>(256);
144
145    send_message(
146        &stdout,
147        &WorkerToDaemon::Ready {
148            zellij_session: session_name.clone(),
149        },
150    )
151    .await?;
152
153    let sample_backend = backend.clone();
154    let sample_screen = latest_screen.clone();
155    let sample_raw_screen = latest_raw_screen.clone();
156    let sample_updates = updates.clone();
157    let sample_stdout = stdout.clone();
158    let sample_adapter = adapter.clone();
159    let sample_display_mode = display_mode.clone();
160    let sample_usage_limit_tracker = usage_limit_tracker.clone();
161    let sample_current_turn_id = current_turn_id.clone();
162    let last_uploaded_hash = Arc::new(Mutex::new(None::<String>));
163    let last_broadcast_hash: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
164    let sample_last_broadcast_hash = last_broadcast_hash.clone();
165    let sample_analyzer_runtime = analyzer_runtime.clone();
166    // Channel for trigger events (TurnStarted, PaneUpdate, …) to the screenshot coordinator.
167    let (trigger_tx, trigger_rx) = mpsc::channel::<Trigger>(16);
168    let screen_capture_task = tokio::spawn(async move {
169        let mut last_emitted_status = ScreenStatus::Starting;
170        let mut last_emitted_usage_limit: Option<CliUsageLimitState> = None;
171        loop {
172            let (screen, alive) = {
173                let guard = sample_backend.lock().await;
174                let screen = guard.capture_viewport().await.unwrap_or_default();
175                let alive = guard.is_alive().await.unwrap_or(false);
176                (screen, alive)
177            };
178
179            let hash_changed;
180
181            {
182                *sample_raw_screen.write().await = screen.clone();
183                let mode = *sample_display_mode.read().await;
184                let rendered = render_screen_for_display_mode(&screen, mode);
185                let now_ms = now_ms();
186                let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
187                let base_status = if analyzing {
188                    ScreenStatus::Analyzing
189                } else {
190                    ScreenStatus::Working
191                };
192                let (status, usage_limit) =
193                    sample_usage_limit_tracker
194                        .lock()
195                        .await
196                        .classify(&screen, base_status, now_ms);
197                let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
198                {
199                    let guard = sample_last_broadcast_hash.lock().await;
200                    hash_changed = guard.as_deref() != Some(&rendered_hash);
201                }
202
203                if hash_changed
204                    || last_emitted_status != status
205                    || last_emitted_usage_limit != usage_limit
206                {
207                    *sample_last_broadcast_hash.lock().await = Some(rendered_hash.clone());
208                    let mut current = sample_screen.write().await;
209                    *current = rendered.clone();
210                    let _ = sample_updates.send(rendered.clone());
211                    let _ = send_message(
212                        &sample_stdout,
213                        &WorkerToDaemon::ScreenUpdate {
214                            content: rendered.clone(),
215                            status,
216                            usage_limit: usage_limit.clone(),
217                        },
218                    )
219                    .await;
220                    last_emitted_status = status;
221                    last_emitted_usage_limit = usage_limit.clone();
222                }
223            }
224
225            if let Ok(poll) = sample_adapter.lock().await.poll() {
226                if let Some(cli_session_id) = poll.cli_session_id {
227                    let _ = send_message(
228                        &sample_stdout,
229                        &WorkerToDaemon::CliSessionId { cli_session_id },
230                    )
231                    .await;
232                }
233                if let Some((user_text, assistant_text)) = poll.adopt_preamble {
234                    let _ = send_message(
235                        &sample_stdout,
236                        &WorkerToDaemon::AdoptPreamble {
237                            user_text,
238                            assistant_text,
239                        },
240                    )
241                    .await;
242                }
243                if let Some(content) = poll.final_output {
244                    let turn_id = sample_current_turn_id.read().await.clone();
245                    let _ = send_message(
246                        &sample_stdout,
247                        &WorkerToDaemon::FinalOutput {
248                            content,
249                            turn_id,
250                            kind: poll.final_output_kind,
251                            user_text: poll.final_output_user_text,
252                        },
253                    )
254                    .await;
255                }
256                if poll.prompt_ready {
257                    let _ = send_message(&sample_stdout, &WorkerToDaemon::PromptReady).await;
258                    let rendered = sample_screen.read().await.clone();
259                    let raw = sample_raw_screen.read().await.clone();
260                    let now_ms = now_ms();
261                    let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
262                    let base_status = if analyzing {
263                        ScreenStatus::Analyzing
264                    } else {
265                        ScreenStatus::Idle
266                    };
267                    let (status, usage_limit) =
268                        sample_usage_limit_tracker
269                            .lock()
270                            .await
271                            .classify(&raw, base_status, now_ms);
272                    let _ = send_message(
273                        &sample_stdout,
274                        &WorkerToDaemon::ScreenUpdate {
275                            content: rendered,
276                            status,
277                            usage_limit,
278                        },
279                    )
280                    .await;
281                }
282            }
283
284            if !alive {
285                let _ = send_message(
286                    &sample_stdout,
287                    &WorkerToDaemon::CliExit {
288                        code: Some(0),
289                        signal: None,
290                    },
291                )
292                .await;
293                break;
294            }
295
296            tokio::time::sleep(Duration::from_millis(5000)).await;
297        }
298    });
299    let mut worker_joins = tokio::task::JoinSet::new();
300    worker_joins.spawn(async move {
301        let _ = screen_capture_task.await;
302    });
303
304    if screen_analyzer_enabled(&init.screen_analyzer) {
305        let analyzer_cfg = init.screen_analyzer.clone();
306        let analyzer_raw_screen = latest_raw_screen.clone();
307        let analyzer_runtime_state = analyzer_runtime.clone();
308        let analyzer_stdout = stdout.clone();
309        let analyzer_task = tokio::spawn(async move {
310            let client = Client::new();
311            loop {
312                tokio::time::sleep(Duration::from_millis(analyzer_cfg.interval_ms)).await;
313                let snapshot = analyzer_raw_screen.read().await.clone();
314                if snapshot.is_empty() {
315                    continue;
316                }
317                let truncated = if snapshot.len() > analyzer_cfg.snapshot_max_chars {
318                    snapshot[snapshot.len() - analyzer_cfg.snapshot_max_chars..].to_string()
319                } else {
320                    snapshot
321                };
322                let now = now_ms();
323                {
324                    let mut runtime = analyzer_runtime_state.write().await;
325                    if truncated == runtime.last_snapshot {
326                        runtime.stable_count = runtime.stable_count.saturating_add(1);
327                    } else {
328                        runtime.stable_count = 1;
329                        runtime.last_snapshot = truncated.clone();
330                        if runtime.waiting_for_content_change {
331                            runtime.waiting_for_content_change = false;
332                        }
333                    }
334                    if runtime.stable_count < analyzer_cfg.stable_count {
335                        continue;
336                    }
337                    if runtime.waiting_for_content_change
338                        && truncated == runtime.last_analyzed_snapshot
339                    {
340                        continue;
341                    }
342                    if runtime.cooldown_until_ms > now {
343                        continue;
344                    }
345                    runtime.is_analyzing = true;
346                    runtime.last_analyzed_snapshot = truncated.clone();
347                }
348
349                let result = call_screen_analyzer(&client, &analyzer_cfg, &truncated).await;
350
351                let mut runtime = analyzer_runtime_state.write().await;
352                runtime.is_analyzing = false;
353                match result {
354                    Ok(analysis) => {
355                        apply_screen_analyzer_result(
356                            &mut runtime,
357                            &analysis.check_again_when,
358                            now_ms(),
359                        );
360                        if analysis.needs_interaction && !analysis.options.is_empty() {
361                            if !runtime.prompt_active {
362                                runtime.prompt_active = true;
363                                let _ = send_message(
364                                    &analyzer_stdout,
365                                    &WorkerToDaemon::TuiPrompt {
366                                        description: analysis.description.clone().unwrap_or_else(
367                                            || "CLI needs your selection".to_string(),
368                                        ),
369                                        options: analysis.options.clone(),
370                                        multi_select: analysis.multi_select,
371                                    },
372                                )
373                                .await;
374                            }
375                        } else if runtime.prompt_active {
376                            runtime.prompt_active = false;
377                            let _ = send_message(
378                                &analyzer_stdout,
379                                &WorkerToDaemon::TuiPromptResolved {
380                                    selected_text: None,
381                                },
382                            )
383                            .await;
384                        }
385                    }
386                    Err(_) => {
387                        runtime.waiting_for_content_change = true;
388                        runtime.cooldown_until_ms = 0;
389                    }
390                }
391            }
392        });
393        worker_joins.spawn(async move {
394            let _ = analyzer_task.await;
395        });
396    }
397
398    // Coordinator task — owns the trigger receiver and the 5-second fallback.
399    {
400        let coord_backend = backend.clone();
401        let coord_stdout = stdout.clone();
402        let coord_session_id = init.session_id.clone();
403        let coord_app_id = init.lark_app_id.clone();
404        let coord_app_secret = init.lark_app_secret.clone();
405        let coord_display_mode = display_mode.clone();
406        let coord_analyzer_runtime = analyzer_runtime.clone();
407        let coord_usage_limit_tracker = usage_limit_tracker.clone();
408        let coord_last_uploaded_hash = last_uploaded_hash.clone();
409        let coord_latest_raw_screen = latest_raw_screen.clone();
410        let coord_rx = trigger_rx;
411        worker_joins.spawn(async move {
412            coordinator_loop(
413                coord_backend,
414                coord_stdout,
415                coord_session_id,
416                coord_app_id,
417                coord_app_secret,
418                coord_display_mode,
419                coord_analyzer_runtime,
420                coord_usage_limit_tracker,
421                coord_last_uploaded_hash,
422                coord_latest_raw_screen,
423                coord_rx,
424            )
425            .await;
426        });
427    }
428    // Subscribe task: forward backend pane-update notifications to the screenshot coordinator.
429    // Also caches the full viewport ANSI chunk so the coordinator can capture
430    // without waiting for the backend Mutex held by write_input().
431    {
432        let sub_backend = backend.clone();
433        let sub_trigger_tx = trigger_tx.clone();
434        let sub_latest_raw_screen = latest_raw_screen.clone();
435        worker_joins.spawn(async move {
436            let mut rx = sub_backend.lock().await.subscribe();
437            loop {
438                match rx.recv().await {
439                    Ok(chunk) => {
440                        // latest wins: the chunk is the full viewport (not incremental)
441                        *sub_latest_raw_screen.write().await = chunk;
442                        match sub_trigger_tx.try_send(Trigger::PaneUpdate) {
443                            Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => break,
444                            _ => {} // Ok or Full → discard, keep listening
445                        }
446                    }
447                    Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
448                    Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
449                }
450            }
451        });
452    }
453    if !init.prompt.is_empty() && !crate::adapters::passes_initial_prompt_via_args(&init.cli_id) {
454        usage_limit_tracker.lock().await.begin_turn(
455            "",
456            SystemTime::now()
457                .duration_since(UNIX_EPOCH)
458                .unwrap_or_default()
459                .as_millis() as u64,
460        );
461        *current_turn_id.write().await = init
462            .prompt_turn_id
463            .clone()
464            .unwrap_or_else(|| Uuid::new_v4().to_string());
465        *last_uploaded_hash.lock().await = None;
466        let guard = backend.lock().await;
467        let submit = adapter
468            .lock()
469            .await
470            .write_input(guard.as_ref(), &init.prompt)
471            .await?;
472        if let Some(cli_session_id) = submit.cli_session_id {
473            send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
474        }
475    }
476
477    // Init-time transcript source resolution for adapters with resolvable
478    // sources (opencode): resolve cli_session_id before the message loop.
479    {
480        let resolution = {
481            let mut adapter_guard = adapter.lock().await;
482            let backend_guard = backend.lock().await;
483            let resolution = adapter_guard
484                .resolve_transcript_source(backend_guard.as_ref())
485                .await;
486            drop(backend_guard);
487            drop(adapter_guard);
488            resolution
489        };
490        if let Some(resolution) = resolution {
491            let outcome = resolution.unwrap_or_else(|err| {
492                warn!(session = %init.session_id, adapter = %init.cli_id, "resolve_transcript_source error: {:?}", err);
493                ResolveOutcome::NotFound {
494                    reason: format!("transcript source resolution failed: {}", err),
495                }
496            });
497            match outcome {
498                ResolveOutcome::Found(source) => {
499                    send_message(
500                        &stdout,
501                        &WorkerToDaemon::CliSessionId {
502                            cli_session_id: source.session_id.clone(),
503                        },
504                    )
505                    .await?;
506                    // Also set in adapter state for subsequent poll/write_input.
507                    adapter
508                        .lock()
509                        .await
510                        .set_transcript_source(&source.session_id);
511                    info!(
512                        session = %init.session_id, adapter = %init.cli_id,
513                        transcript_session = %source.session_id,
514                        "transcript source resolved automatically"
515                    );
516                }
517                ResolveOutcome::Ambiguous { candidates, .. } => {
518                    info!(
519                        session = %init.session_id, adapter = %init.cli_id,
520                        candidate_count = candidates.len(),
521                        "transcript source ambiguous, requesting user choice"
522                    );
523                    let turn_id = Uuid::new_v4().to_string();
524                    let choices: Vec<TranscriptChoice> = candidates
525                        .iter()
526                        .map(|c| TranscriptChoice {
527                            session_id: c.session_id.clone(),
528                            label: format!("{} ({})", c.session_id, c.db_path.display()),
529                        })
530                        .collect();
531                    send_message(
532                        &stdout,
533                        &WorkerToDaemon::TranscriptChoices {
534                            candidates: choices,
535                            turn_id: turn_id.clone(),
536                        },
537                    )
538                    .await?;
539                }
540                ResolveOutcome::NotFound { reason } => {
541                    info!(session = %init.session_id, adapter = %init.cli_id, "transcript source not found: {}", reason);
542                    send_message(&stdout, &WorkerToDaemon::UserNotify { message: reason })
543                        .await?;
544                }
545            }
546        }
547    }
548
549    let stdin = BufReader::new(tokio::io::stdin());
550    let mut lines = stdin.lines();
551    loop {
552        let line = match lines.next_line().await {
553            Ok(Some(line)) => line,
554            Ok(None) => break,
555            Err(_) => break,
556        };
557        if line.trim().is_empty() {
558            continue;
559        }
560        let msg: DaemonToWorker = serde_json::from_str(&line)?;
561        match msg {
562            DaemonToWorker::Message { content, turn_id } => {
563                info!(session = %init.session_id, %turn_id, "Message received, sending TurnStarted to coordinator");
564                handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
565                let snapshot = latest_raw_screen.read().await.clone();
566                usage_limit_tracker.lock().await.begin_turn(
567                    &snapshot,
568                    SystemTime::now()
569                        .duration_since(UNIX_EPOCH)
570                        .unwrap_or_default()
571                        .as_millis() as u64,
572                );
573                *current_turn_id.write().await = turn_id.clone();
574                *last_uploaded_hash.lock().await = None;
575                // Notify coordinator: a new turn has started (best-effort).
576                if trigger_tx
577                    .send(Trigger::TurnStarted { turn_id })
578                    .await
579                    .is_err()
580                {
581                    warn!("coordinator channel closed, TurnStarted not sent");
582                }
583                let guard = backend.lock().await;
584                let submit = adapter
585                    .lock()
586                    .await
587                    .write_input(guard.as_ref(), &content)
588                    .await?;
589                if let Some(cli_session_id) = submit.cli_session_id {
590                    send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
591                }
592                if !submit.submitted {
593                    let message = submit
594                        .failure_reason
595                        .unwrap_or_else(|| "CLI submit could not be confirmed".to_string());
596                    send_message(&stdout, &WorkerToDaemon::UserNotify { message }).await?;
597                }
598            }
599            DaemonToWorker::RawInput { content, turn_id } => {
600                info!(session = %init.session_id, %turn_id, "RawInput received, sending TurnStarted to coordinator");
601                handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
602                let snapshot = latest_raw_screen.read().await.clone();
603                usage_limit_tracker.lock().await.begin_turn(
604                    &snapshot,
605                    SystemTime::now()
606                        .duration_since(UNIX_EPOCH)
607                        .unwrap_or_default()
608                        .as_millis() as u64,
609                );
610                *current_turn_id.write().await = turn_id.clone();
611                *last_uploaded_hash.lock().await = None;
612                // Notify coordinator: a new turn has started (best-effort).
613                if trigger_tx
614                    .send(Trigger::TurnStarted { turn_id })
615                    .await
616                    .is_err()
617                {
618                    warn!("coordinator channel closed, TurnStarted not sent");
619                }
620                let guard = backend.lock().await;
621                guard.raw_input(&content).await?;
622            }
623            DaemonToWorker::Close => {
624                let mut guard = backend.lock().await;
625                guard.destroy_session().await?;
626                break;
627            }
628            DaemonToWorker::Restart => {
629                let mut guard = backend.lock().await;
630                guard.destroy_session().await?;
631                break;
632            }
633            DaemonToWorker::RefreshScreen => {
634                let guard = backend.lock().await;
635                let screen = guard.capture_viewport().await?;
636                *latest_raw_screen.write().await = screen.clone();
637                let mode = *display_mode.read().await;
638                let rendered = render_screen_for_display_mode(&screen, mode);
639                let now_ms = now_ms();
640                let analyzing = analyzer_runtime.read().await.is_analyzing;
641                let base_status = if analyzing {
642                    ScreenStatus::Analyzing
643                } else {
644                    ScreenStatus::Working
645                };
646                let (status, usage_limit) =
647                    usage_limit_tracker
648                        .lock()
649                        .await
650                        .classify(&screen, base_status, now_ms);
651                *latest_screen.write().await = rendered.clone();
652                let _ = updates.send(rendered.clone());
653                send_message(
654                    &stdout,
655                    &WorkerToDaemon::ScreenUpdate {
656                        content: rendered.clone(),
657                        status,
658                        usage_limit: usage_limit.clone(),
659                    },
660                )
661                .await?;
662                let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
663                *last_broadcast_hash.lock().await = Some(rendered_hash);
664                // Notify coordinator: refresh request (best-effort).
665                if trigger_tx.send(Trigger::Refresh).await.is_err() {
666                    warn!("coordinator channel closed, Refresh not sent");
667                }
668            }
669            DaemonToWorker::SetDisplayMode { mode } => {
670                *display_mode.write().await = mode;
671                let raw = latest_raw_screen.read().await.clone();
672                let rendered = render_screen_for_display_mode(&raw, mode);
673                let now_ms = now_ms();
674                let analyzing = analyzer_runtime.read().await.is_analyzing;
675                let base_status = if analyzing {
676                    ScreenStatus::Analyzing
677                } else {
678                    ScreenStatus::Working
679                };
680                let (status, usage_limit) =
681                    usage_limit_tracker
682                        .lock()
683                        .await
684                        .classify(&raw, base_status, now_ms);
685                *latest_screen.write().await = rendered.clone();
686                let _ = updates.send(rendered.clone());
687                send_message(
688                    &stdout,
689                    &WorkerToDaemon::ScreenUpdate {
690                        content: rendered,
691                        status,
692                        usage_limit: usage_limit.clone(),
693                    },
694                )
695                .await?;
696                // Notify coordinator: display mode changed (best-effort).
697                if trigger_tx
698                    .send(Trigger::SetDisplayMode(mode))
699                    .await
700                    .is_err()
701                {
702                    warn!("coordinator channel closed, SetDisplayMode not sent");
703                }
704            }
705            DaemonToWorker::TermAction { key } => {
706                let keys = term_action_keys(key);
707                let guard = backend.lock().await;
708                guard.send_special_keys(&keys).await?;
709            }
710            DaemonToWorker::SpecialKeys { keys } => {
711                let guard = backend.lock().await;
712                guard.send_special_keys(&keys).await?;
713            }
714            DaemonToWorker::TuiKeys { keys, is_final } => {
715                handle_tui_keys(&backend, &analyzer_runtime, &keys, is_final).await?;
716            }
717            DaemonToWorker::TuiTextInput { keys, text } => {
718                handle_tui_text_input(&backend, &adapter, &analyzer_runtime, &keys, &text).await?;
719            }
720            DaemonToWorker::SetTranscriptSource { cli_session_id } => {
721                let applied = adapter.lock().await.set_transcript_source(&cli_session_id);
722                if applied {
723                    info!(session = %init.session_id, adapter = %init.cli_id, transcript_session = %cli_session_id, "transcript source set by user");
724                    send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
725                }
726            }
727            DaemonToWorker::Init(_) => {}
728        }
729    }
730
731    worker_joins.abort_all();
732    while worker_joins.join_next().await.is_some() {}
733    {
734        let mut guard = backend.lock().await;
735        let _ = guard.kill().await;
736    }
737    if let Some(marker) = cli_pid_marker {
738        let _ = tokio::fs::remove_file(marker).await;
739    }
740    info!("worker exiting");
741    Ok(())
742}