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) 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 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 {
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 {
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 *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 _ => {} }
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 {
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 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 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 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 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 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 #[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}