1use super::*;
2
3pub(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 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 {
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 {
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 *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 _ => {} }
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 {
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 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 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 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 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 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}