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