Skip to main content

beam_worker/worker_runtime/
run_loop.rs

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