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