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) const CLI_DEAD_CONFIRM_ATTEMPTS: usize = 3;
32
33pub(crate) fn should_emit_cli_exit(consecutive_dead: &mut usize, alive: bool) -> bool {
36 if alive {
37 *consecutive_dead = 0;
38 return false;
39 }
40 *consecutive_dead = consecutive_dead.saturating_add(1);
41 *consecutive_dead >= CLI_DEAD_CONFIRM_ATTEMPTS
42}
43
44pub(crate) async fn send_message(
45 stdout: &Arc<Mutex<tokio::io::Stdout>>,
46 msg: &WorkerToDaemon,
47) -> Result<()> {
48 let mut out = stdout.lock().await;
49 out.write_all(serde_json::to_string(msg)?.as_bytes())
50 .await?;
51 out.write_all(b"\n").await?;
52 out.flush().await?;
53 Ok(())
54}
55
56pub async fn run(init: InitConfig) -> Result<()> {
57 let stdout = Arc::new(Mutex::new(tokio::io::stdout()));
58 let session_name = format!("beam-{}", &init.session_id[..8.min(init.session_id.len())]);
59 let paths = BeamPaths::discover()?;
60 let adapter = Arc::new(Mutex::new(CliAdapter::from_init(&init)?));
61 let (backend_impl, attach_context, herdr_handle) =
65 crate::backend::select::select_backend(&init, &session_name);
66 let spawn_spec = adapter.lock().await.build_spawn_spec(&init);
67 let extra_env = maybe_inject_term(&init.cli_id, std::env::var("TERM").ok().as_deref())
68 .into_iter()
69 .collect();
70 let exe_path = std::env::current_exe()
71 .ok()
72 .map(|p| p.display().to_string())
73 .unwrap_or_else(|| "beam".to_string());
74 let bindir = std::env::current_exe()
75 .ok()
76 .and_then(|p| p.parent().map(|v| v.display().to_string()))
77 .unwrap_or_default();
78 let args = if init.adopted_from.is_some() {
79 (spawn_spec.bin, spawn_spec.args)
80 } else {
81 if super::launch::normalize_cgroup_slice(init.cgroup_slice.as_deref()).is_some()
82 && super::launch::current_launch_platform() == super::launch::LaunchPlatform::Linux
83 && !super::launch::systemd_run_available()
84 {
85 anyhow::bail!(
86 "cgroupSlice requires systemd-run on Linux (user systemd); install systemd or unset cgroupSlice"
87 );
88 }
89 super::launch::build_launch_spec(
90 super::launch::current_launch_platform(),
91 super::launch::LaunchInput {
92 cgroup_slice: init.cgroup_slice.clone(),
93 session_id: init.session_id.clone(),
94 beam_home: paths.root().display().to_string(),
95 beam_bin: exe_path,
96 path_prepend: bindir,
97 extra_env,
98 spec: spawn_spec,
99 },
100 )?
101 };
102 let spawn_opts = SpawnOpts {
103 cwd: init.working_dir.clone(),
104 cols: DEFAULT_TERMINAL_COLS,
105 rows: DEFAULT_TERMINAL_ROWS,
106 env: Vec::new(),
107 };
108 backend_impl
109 .spawn(&args.0, &args.1, spawn_opts)
110 .await
111 .with_context(|| format!("failed to {} session {}", attach_context, init.session_id))?;
112 let backend = backend_impl;
116 let (herdr_workspace_id, herdr_pane_id) = crate::backend::select::ready_herdr_identity(
117 herdr_handle
118 .as_ref()
119 .and_then(|h| h.herdr_ids())
120 .map(|ids| ids.workspace_pane()),
121 init.adopted_from.as_ref(),
122 );
123 let mut cli_pid_marker = None;
124 let child_pid = backend.child_pid().await?;
125 adapter.lock().await.on_spawned(child_pid);
126 if let Some(pid) = child_pid {
127 tokio::fs::create_dir_all(paths.cli_pid_markers_dir()).await?;
128 let marker = paths.cli_pid_markers_dir().join(pid.to_string());
129 tokio::fs::write(&marker, init.session_id.as_bytes()).await?;
130 cli_pid_marker = Some(marker);
131 }
132 let ready_probe = crate::adapters::ready_probe(&init.cli_id);
138 if init.adopted_from.is_none() && ready_probe != ReadyProbe::None {
139 let ready = wait_for_ready(backend.as_ref(), ready_probe).await;
140 if ready {
141 info!(session = %init.session_id, adapter = %init.cli_id, probe = ?ready_probe, "CLI TUI ready signal observed");
142 } else {
143 warn!(session = %init.session_id, adapter = %init.cli_id, probe = ?ready_probe, "CLI TUI ready signal not observed within {}s; typing input anyway", TUI_READY_TIMEOUT.as_secs());
144 }
145 }
146 let latest_screen = Arc::new(RwLock::new(String::new()));
147 let latest_raw_screen = Arc::new(RwLock::new(String::new()));
148 let display_mode = Arc::new(RwLock::new(DisplayMode::Hidden));
149 let analyzer_runtime = Arc::new(RwLock::new(AnalyzerRuntime::default()));
150 let usage_limit_tracker = Arc::new(Mutex::new(UsageLimitTracker::default()));
151 let current_turn_id = Arc::new(RwLock::new(String::new()));
152 let (updates, _) = broadcast::channel::<String>(256);
153
154 send_message(
155 &stdout,
156 &WorkerToDaemon::Ready {
157 zellij_session: session_name.clone(),
158 backend_kind: init.backend_kind,
159 herdr_workspace_id,
160 herdr_pane_id,
161 },
162 )
163 .await?;
164
165 let sample_backend = backend.clone();
166 let sample_screen = latest_screen.clone();
167 let sample_raw_screen = latest_raw_screen.clone();
168 let sample_updates = updates.clone();
169 let sample_stdout = stdout.clone();
170 let sample_adapter = adapter.clone();
171 let sample_display_mode = display_mode.clone();
172 let sample_usage_limit_tracker = usage_limit_tracker.clone();
173 let sample_current_turn_id = current_turn_id.clone();
174 let last_uploaded_hash = Arc::new(Mutex::new(None::<String>));
175 let last_broadcast_hash: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(None));
176 let sample_last_broadcast_hash = last_broadcast_hash.clone();
177 let sample_analyzer_runtime = analyzer_runtime.clone();
178 let (trigger_tx, trigger_rx) = mpsc::channel::<Trigger>(16);
180 let screen_capture_task = tokio::spawn(async move {
181 let mut last_emitted_status = ScreenStatus::Starting;
182 let mut last_emitted_usage_limit: Option<CliUsageLimitState> = None;
183 let mut consecutive_dead = 0usize;
184 loop {
185 let (screen, alive) = {
186 let screen = sample_backend.capture_viewport().await.unwrap_or_default();
187 let alive = sample_backend.is_alive().await.unwrap_or(true);
189 (screen, alive)
190 };
191
192 let hash_changed;
193
194 {
195 *sample_raw_screen.write().await = screen.clone();
196 let mode = *sample_display_mode.read().await;
197 let rendered = render_screen_for_display_mode(&screen, mode);
198 let now_ms = now_ms();
199 let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
200 let base_status = if analyzing {
201 ScreenStatus::Analyzing
202 } else {
203 ScreenStatus::Working
204 };
205 let (status, usage_limit) =
206 sample_usage_limit_tracker
207 .lock()
208 .await
209 .classify(&screen, base_status, now_ms);
210 let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
211 {
212 let guard = sample_last_broadcast_hash.lock().await;
213 hash_changed = guard.as_deref() != Some(&rendered_hash);
214 }
215
216 if hash_changed
217 || last_emitted_status != status
218 || last_emitted_usage_limit != usage_limit
219 {
220 *sample_last_broadcast_hash.lock().await = Some(rendered_hash.clone());
221 let mut current = sample_screen.write().await;
222 *current = rendered.clone();
223 let _ = sample_updates.send(rendered.clone());
224 let _ = send_message(
225 &sample_stdout,
226 &WorkerToDaemon::ScreenUpdate {
227 content: rendered.clone(),
228 status,
229 usage_limit: usage_limit.clone(),
230 },
231 )
232 .await;
233 last_emitted_status = status;
234 last_emitted_usage_limit = usage_limit.clone();
235 }
236 }
237
238 if let Ok(poll) = sample_adapter.lock().await.poll() {
239 if let Some(cli_session_id) = poll.cli_session_id {
240 let _ = send_message(
241 &sample_stdout,
242 &WorkerToDaemon::CliSessionId { cli_session_id },
243 )
244 .await;
245 }
246 if let Some((user_text, assistant_text)) = poll.adopt_preamble {
247 let _ = send_message(
248 &sample_stdout,
249 &WorkerToDaemon::AdoptPreamble {
250 user_text,
251 assistant_text,
252 },
253 )
254 .await;
255 }
256 if let Some(content) = poll.final_output {
257 let turn_id = sample_current_turn_id.read().await.clone();
258 let _ = send_message(
259 &sample_stdout,
260 &WorkerToDaemon::FinalOutput {
261 content,
262 turn_id,
263 kind: poll.final_output_kind,
264 user_text: poll.final_output_user_text,
265 },
266 )
267 .await;
268 }
269 if poll.prompt_ready {
270 let _ = send_message(&sample_stdout, &WorkerToDaemon::PromptReady).await;
271 let rendered = sample_screen.read().await.clone();
272 let raw = sample_raw_screen.read().await.clone();
273 let now_ms = now_ms();
274 let analyzing = sample_analyzer_runtime.read().await.is_analyzing;
275 let base_status = if analyzing {
276 ScreenStatus::Analyzing
277 } else {
278 ScreenStatus::Idle
279 };
280 let (status, usage_limit) =
281 sample_usage_limit_tracker
282 .lock()
283 .await
284 .classify(&raw, base_status, now_ms);
285 let _ = send_message(
286 &sample_stdout,
287 &WorkerToDaemon::ScreenUpdate {
288 content: rendered,
289 status,
290 usage_limit,
291 },
292 )
293 .await;
294 }
295 }
296
297 if should_emit_cli_exit(&mut consecutive_dead, alive) {
298 warn!(
299 consecutive_dead,
300 "cli liveness dead confirmed, emitting CliExit"
301 );
302 let _ = send_message(
303 &sample_stdout,
304 &WorkerToDaemon::CliExit {
305 code: Some(0),
306 signal: None,
307 },
308 )
309 .await;
310 break;
311 }
312 if !alive {
313 warn!(
314 consecutive_dead,
315 confirm_after = CLI_DEAD_CONFIRM_ATTEMPTS,
316 "cli liveness dead, waiting for confirmation"
317 );
318 }
319
320 tokio::time::sleep(Duration::from_millis(5000)).await;
321 }
322 });
323 let mut worker_joins = tokio::task::JoinSet::new();
324 worker_joins.spawn(async move {
325 let _ = screen_capture_task.await;
326 });
327
328 if init.cli_id == "grok" {
329 let grok_raw_screen = latest_raw_screen.clone();
330 let grok_runtime_state = analyzer_runtime.clone();
331 let grok_stdout = stdout.clone();
332 worker_joins.spawn(async move {
333 loop {
334 tokio::time::sleep(Duration::from_millis(800)).await;
335 let snapshot = grok_raw_screen.read().await.clone();
336 if snapshot.is_empty() {
337 continue;
338 }
339 let options =
340 crate::worker_runtime::grok_prompts::detect_grok_plan_approval(&snapshot);
341 let mut runtime = grok_runtime_state.write().await;
342 match options {
343 Some(options) if !runtime.prompt_active => {
344 runtime.prompt_active = true;
345 let _ = send_message(
346 &grok_stdout,
347 &WorkerToDaemon::TuiPrompt {
348 description: "Grok plan approval".to_string(),
349 options,
350 multi_select: false,
351 },
352 )
353 .await;
354 }
355 None if runtime.prompt_active => {
356 runtime.prompt_active = false;
357 let _ = send_message(
358 &grok_stdout,
359 &WorkerToDaemon::TuiPromptResolved {
360 selected_text: None,
361 },
362 )
363 .await;
364 }
365 _ => {}
366 }
367 }
368 });
369 }
370
371 if screen_analyzer_enabled(&init.screen_analyzer) {
372 let analyzer_cfg = init.screen_analyzer.clone();
373 let analyzer_raw_screen = latest_raw_screen.clone();
374 let analyzer_runtime_state = analyzer_runtime.clone();
375 let analyzer_stdout = stdout.clone();
376 let analyzer_cli_id = init.cli_id.clone();
377 let analyzer_task = tokio::spawn(async move {
378 let client = Client::new();
379 loop {
380 tokio::time::sleep(Duration::from_millis(analyzer_cfg.interval_ms)).await;
381 let snapshot = analyzer_raw_screen.read().await.clone();
382 if snapshot.is_empty() {
383 continue;
384 }
385 let truncated = if snapshot.len() > analyzer_cfg.snapshot_max_chars {
386 snapshot[snapshot.len() - analyzer_cfg.snapshot_max_chars..].to_string()
387 } else {
388 snapshot
389 };
390 let now = now_ms();
391 {
392 let mut runtime = analyzer_runtime_state.write().await;
393 if truncated == runtime.last_snapshot {
394 runtime.stable_count = runtime.stable_count.saturating_add(1);
395 } else {
396 runtime.stable_count = 1;
397 runtime.last_snapshot = truncated.clone();
398 if runtime.waiting_for_content_change {
399 runtime.waiting_for_content_change = false;
400 }
401 }
402 if runtime.stable_count < analyzer_cfg.stable_count {
403 continue;
404 }
405 if runtime.waiting_for_content_change
406 && truncated == runtime.last_analyzed_snapshot
407 {
408 continue;
409 }
410 if runtime.cooldown_until_ms > now {
411 continue;
412 }
413 runtime.is_analyzing = true;
414 runtime.last_analyzed_snapshot = truncated.clone();
415 }
416
417 if analyzer_cli_id == "grok"
418 && crate::worker_runtime::grok_prompts::detect_grok_plan_approval(&truncated)
419 .is_some()
420 {
421 let mut runtime = analyzer_runtime_state.write().await;
422 runtime.is_analyzing = false;
423 continue;
424 }
425
426 let result = call_screen_analyzer(&client, &analyzer_cfg, &truncated).await;
427
428 let mut runtime = analyzer_runtime_state.write().await;
429 runtime.is_analyzing = false;
430 match result {
431 Ok(analysis) => {
432 apply_screen_analyzer_result(
433 &mut runtime,
434 &analysis.check_again_when,
435 now_ms(),
436 );
437 if analysis.needs_interaction && !analysis.options.is_empty() {
438 if !runtime.prompt_active {
439 runtime.prompt_active = true;
440 let _ = send_message(
441 &analyzer_stdout,
442 &WorkerToDaemon::TuiPrompt {
443 description: analysis.description.clone().unwrap_or_else(
444 || "CLI needs your selection".to_string(),
445 ),
446 options: analysis.options.clone(),
447 multi_select: analysis.multi_select,
448 },
449 )
450 .await;
451 }
452 } else if runtime.prompt_active {
453 runtime.prompt_active = false;
454 let _ = send_message(
455 &analyzer_stdout,
456 &WorkerToDaemon::TuiPromptResolved {
457 selected_text: None,
458 },
459 )
460 .await;
461 }
462 }
463 Err(_) => {
464 runtime.waiting_for_content_change = true;
465 runtime.cooldown_until_ms = 0;
466 }
467 }
468 }
469 });
470 worker_joins.spawn(async move {
471 let _ = analyzer_task.await;
472 });
473 }
474
475 {
477 let coord_backend = backend.clone();
478 let coord_stdout = stdout.clone();
479 let coord_session_id = init.session_id.clone();
480 let coord_app_id = init.lark_app_id.clone();
481 let coord_app_secret = init.lark_app_secret.clone();
482 let coord_display_mode = display_mode.clone();
483 let coord_analyzer_runtime = analyzer_runtime.clone();
484 let coord_usage_limit_tracker = usage_limit_tracker.clone();
485 let coord_last_uploaded_hash = last_uploaded_hash.clone();
486 let coord_latest_raw_screen = latest_raw_screen.clone();
487 let coord_rx = trigger_rx;
488 worker_joins.spawn(async move {
489 coordinator_loop(
490 coord_backend,
491 coord_stdout,
492 coord_session_id,
493 coord_app_id,
494 coord_app_secret,
495 coord_display_mode,
496 coord_analyzer_runtime,
497 coord_usage_limit_tracker,
498 coord_last_uploaded_hash,
499 coord_latest_raw_screen,
500 coord_rx,
501 )
502 .await;
503 });
504 }
505 {
509 let sub_backend = backend.clone();
510 let sub_trigger_tx = trigger_tx.clone();
511 let sub_latest_raw_screen = latest_raw_screen.clone();
512 worker_joins.spawn(async move {
513 let mut rx = sub_backend.subscribe();
514 loop {
515 match rx.recv().await {
516 Ok(chunk) => {
517 *sub_latest_raw_screen.write().await = chunk;
519 if let Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) =
520 sub_trigger_tx.try_send(Trigger::PaneUpdate)
521 {
522 break;
523 } }
525 Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
526 Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => continue,
527 }
528 }
529 });
530 }
531 if init.backend_kind == beam_core::BackendKind::Herdr
534 && let Some(handle) = herdr_handle.as_ref()
535 && let Some(pane_id) = handle.herdr_ids().map(|ids| ids.pane_id)
536 {
537 let agent_stdout = stdout.clone();
538 let agent_pane_id = pane_id.clone();
539 worker_joins.spawn(async move {
540 loop {
541 tokio::time::sleep(Duration::from_secs(2)).await;
542 let state = crate::backend::herdr::cli::agent_get(&agent_pane_id)
543 .await
544 .ok()
545 .and_then(|payload| crate::backend::herdr::cli::parse_agent_state(&payload));
546 let Some(state) = state else {
547 continue;
548 };
549 if send_message(
550 &agent_stdout,
551 &WorkerToDaemon::MuxAgentState {
552 state,
553 agent_name: None,
554 pane_id: agent_pane_id.clone(),
555 message: None,
556 },
557 )
558 .await
559 .is_err()
560 {
561 break;
562 }
563 }
564 });
565 }
566 if !init.prompt.is_empty() && !crate::adapters::passes_initial_prompt_via_args(&init.cli_id) {
567 usage_limit_tracker.lock().await.begin_turn(
568 "",
569 SystemTime::now()
570 .duration_since(UNIX_EPOCH)
571 .unwrap_or_default()
572 .as_millis() as u64,
573 );
574 *current_turn_id.write().await = init
575 .prompt_turn_id
576 .clone()
577 .unwrap_or_else(|| Uuid::new_v4().to_string());
578 *last_uploaded_hash.lock().await = None;
579 let submit = adapter
580 .lock()
581 .await
582 .write_input(backend.as_ref(), &init.prompt)
583 .await?;
584 if let Some(cli_session_id) = submit.cli_session_id {
585 send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
586 }
587 if !submit.submitted {
588 let message = submit
589 .failure_reason
590 .unwrap_or_else(|| "CLI submit could not be confirmed".to_string());
591 send_message(&stdout, &WorkerToDaemon::UserNotify { message }).await?;
592 }
593 }
594
595 {
598 let resolution = {
599 let mut adapter_guard = adapter.lock().await;
600 let resolution = adapter_guard
601 .resolve_transcript_source(backend.as_ref())
602 .await;
603 drop(adapter_guard);
604 resolution
605 };
606 if let Some(resolution) = resolution {
607 let outcome = resolution.unwrap_or_else(|err| {
608 warn!(session = %init.session_id, adapter = %init.cli_id, "resolve_transcript_source error: {:?}", err);
609 ResolveOutcome::NotFound {
610 reason: format!("transcript source resolution failed: {}", err),
611 }
612 });
613 match outcome {
614 ResolveOutcome::Found(source) => {
615 send_message(
616 &stdout,
617 &WorkerToDaemon::CliSessionId {
618 cli_session_id: source.session_id.clone(),
619 },
620 )
621 .await?;
622 adapter
624 .lock()
625 .await
626 .set_transcript_source(&source.session_id);
627 info!(
628 session = %init.session_id, adapter = %init.cli_id,
629 transcript_session = %source.session_id,
630 "transcript source resolved automatically"
631 );
632 }
633 ResolveOutcome::Ambiguous { candidates, .. } => {
634 info!(
635 session = %init.session_id, adapter = %init.cli_id,
636 candidate_count = candidates.len(),
637 "transcript source ambiguous, requesting user choice"
638 );
639 let turn_id = Uuid::new_v4().to_string();
640 let choices: Vec<TranscriptChoice> = candidates
641 .iter()
642 .map(|c| TranscriptChoice {
643 session_id: c.session_id.clone(),
644 label: format!("{} ({})", c.session_id, c.db_path.display()),
645 })
646 .collect();
647 send_message(
648 &stdout,
649 &WorkerToDaemon::TranscriptChoices {
650 candidates: choices,
651 turn_id: turn_id.clone(),
652 },
653 )
654 .await?;
655 }
656 ResolveOutcome::NotFound { reason } => {
657 info!(session = %init.session_id, adapter = %init.cli_id, "transcript source not found: {}", reason);
658 send_message(&stdout, &WorkerToDaemon::UserNotify { message: reason }).await?;
659 }
660 }
661 }
662 }
663
664 let processing_since: Arc<StdMutex<Option<(String, u64)>>> = Arc::new(StdMutex::new(None));
670
671 let (daemon_msg_tx, mut daemon_msg_rx) = mpsc::channel::<DaemonToWorker>(32);
676 std::thread::spawn(move || {
677 use std::io::BufRead as _;
678 let stdin = std::io::stdin();
679 for line in stdin.lock().lines() {
680 let line = match line {
681 Ok(line) => line,
682 Err(_) => break,
683 };
684 if line.trim().is_empty() {
685 continue;
686 }
687 match serde_json::from_str::<DaemonToWorker>(&line) {
688 Ok(msg) => {
689 if daemon_msg_tx.blocking_send(msg).is_err() {
690 break;
691 }
692 }
693 Err(err) => warn!("failed to parse daemon message: {}", err),
694 }
695 }
696 });
697
698 const MESSAGE_PROCESSING_WARN_THRESHOLD: Duration = Duration::from_secs(120);
702 {
703 let watchdog_processing = processing_since.clone();
704 let watchdog_session_id = init.session_id.clone();
705 worker_joins.spawn(async move {
706 loop {
707 tokio::time::sleep(Duration::from_secs(15)).await;
708 let stuck = {
709 let guard = watchdog_processing.lock().unwrap();
710 guard.as_ref().and_then(|(desc, start_ms)| {
711 let elapsed = now_ms().saturating_sub(*start_ms);
712 (elapsed > MESSAGE_PROCESSING_WARN_THRESHOLD.as_millis() as u64)
713 .then(|| (desc.clone(), elapsed))
714 })
715 };
716 if let Some((desc, elapsed_ms)) = stuck {
717 warn!(
718 session = %watchdog_session_id,
719 message = %desc,
720 elapsed_ms,
721 "daemon message processing exceeded {}s threshold",
722 MESSAGE_PROCESSING_WARN_THRESHOLD.as_secs()
723 );
724 }
725 }
726 });
727 }
728
729 {
733 let hb_stdout = stdout.clone();
734 let hb_processing = processing_since.clone();
735 worker_joins.spawn(async move {
736 loop {
737 tokio::time::sleep(Duration::from_secs(10)).await;
738 let since = hb_processing
739 .lock()
740 .unwrap()
741 .as_ref()
742 .map(|(_, start)| *start);
743 if send_message(
744 &hb_stdout,
745 &WorkerToDaemon::Heartbeat {
746 processing_since_ms: since,
747 },
748 )
749 .await
750 .is_err()
751 {
752 break;
753 }
754 }
755 });
756 }
757
758 while let Some(msg) = daemon_msg_rx.recv().await {
759 *processing_since.lock().unwrap() = Some((daemon_message_desc(&msg), now_ms()));
760 match msg {
761 DaemonToWorker::Message { content, turn_id } => {
762 info!(session = %init.session_id, %turn_id, "Message received, sending TurnStarted to coordinator");
763 handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
764 let snapshot = latest_raw_screen.read().await.clone();
765 usage_limit_tracker.lock().await.begin_turn(
766 &snapshot,
767 SystemTime::now()
768 .duration_since(UNIX_EPOCH)
769 .unwrap_or_default()
770 .as_millis() as u64,
771 );
772 *current_turn_id.write().await = turn_id.clone();
773 *last_uploaded_hash.lock().await = None;
774 if trigger_tx
776 .send(Trigger::TurnStarted { turn_id })
777 .await
778 .is_err()
779 {
780 warn!("coordinator channel closed, TurnStarted not sent");
781 }
782 let submit = adapter
783 .lock()
784 .await
785 .write_input(backend.as_ref(), &content)
786 .await?;
787 if let Some(cli_session_id) = submit.cli_session_id {
788 send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
789 }
790 if !submit.submitted {
791 let message = submit
792 .failure_reason
793 .unwrap_or_else(|| "CLI submit could not be confirmed".to_string());
794 send_message(&stdout, &WorkerToDaemon::UserNotify { message }).await?;
795 }
796 }
797 DaemonToWorker::RawInput { content, turn_id } => {
798 info!(session = %init.session_id, %turn_id, "RawInput received, sending TurnStarted to coordinator");
799 handle_tui_prompt_override(&stdout, &analyzer_runtime).await;
800 let snapshot = latest_raw_screen.read().await.clone();
801 usage_limit_tracker.lock().await.begin_turn(
802 &snapshot,
803 SystemTime::now()
804 .duration_since(UNIX_EPOCH)
805 .unwrap_or_default()
806 .as_millis() as u64,
807 );
808 *current_turn_id.write().await = turn_id.clone();
809 *last_uploaded_hash.lock().await = None;
810 if trigger_tx
812 .send(Trigger::TurnStarted { turn_id })
813 .await
814 .is_err()
815 {
816 warn!("coordinator channel closed, TurnStarted not sent");
817 }
818 backend.raw_input(&content).await?;
819 }
820 DaemonToWorker::Close => {
821 backend.destroy_session().await?;
822 break;
823 }
824 DaemonToWorker::Restart => {
825 backend.destroy_session().await?;
826 break;
827 }
828 DaemonToWorker::RefreshScreen => {
829 let screen = backend.capture_viewport().await?;
830 *latest_raw_screen.write().await = screen.clone();
831 let mode = *display_mode.read().await;
832 let rendered = render_screen_for_display_mode(&screen, mode);
833 let now_ms = now_ms();
834 let analyzing = analyzer_runtime.read().await.is_analyzing;
835 let base_status = if analyzing {
836 ScreenStatus::Analyzing
837 } else {
838 ScreenStatus::Working
839 };
840 let (status, usage_limit) =
841 usage_limit_tracker
842 .lock()
843 .await
844 .classify(&screen, base_status, now_ms);
845 *latest_screen.write().await = rendered.clone();
846 let _ = updates.send(rendered.clone());
847 send_message(
848 &stdout,
849 &WorkerToDaemon::ScreenUpdate {
850 content: rendered.clone(),
851 status,
852 usage_limit: usage_limit.clone(),
853 },
854 )
855 .await?;
856 let rendered_hash = lower_hex(&Sha256::digest(rendered.as_bytes()));
857 *last_broadcast_hash.lock().await = Some(rendered_hash);
858 if trigger_tx.send(Trigger::Refresh).await.is_err() {
860 warn!("coordinator channel closed, Refresh not sent");
861 }
862 }
863 DaemonToWorker::SetDisplayMode { mode } => {
864 *display_mode.write().await = mode;
865 let raw = latest_raw_screen.read().await.clone();
866 let rendered = render_screen_for_display_mode(&raw, mode);
867 let now_ms = now_ms();
868 let analyzing = analyzer_runtime.read().await.is_analyzing;
869 let base_status = if analyzing {
870 ScreenStatus::Analyzing
871 } else {
872 ScreenStatus::Working
873 };
874 let (status, usage_limit) =
875 usage_limit_tracker
876 .lock()
877 .await
878 .classify(&raw, base_status, now_ms);
879 *latest_screen.write().await = rendered.clone();
880 let _ = updates.send(rendered.clone());
881 send_message(
882 &stdout,
883 &WorkerToDaemon::ScreenUpdate {
884 content: rendered,
885 status,
886 usage_limit: usage_limit.clone(),
887 },
888 )
889 .await?;
890 if trigger_tx
892 .send(Trigger::SetDisplayMode(mode))
893 .await
894 .is_err()
895 {
896 warn!("coordinator channel closed, SetDisplayMode not sent");
897 }
898 }
899 DaemonToWorker::TermAction { key } => {
900 let keys = term_action_keys(key);
901 backend.send_special_keys(&keys).await?;
902 }
903 DaemonToWorker::SpecialKeys { keys } => {
904 backend.send_special_keys(&keys).await?;
905 }
906 DaemonToWorker::TuiKeys { keys, is_final } => {
907 handle_tui_keys(&backend, &analyzer_runtime, &keys, is_final).await?;
908 }
909 DaemonToWorker::TuiTextInput { keys, text } => {
910 handle_tui_text_input(&backend, &adapter, &analyzer_runtime, &keys, &text).await?;
911 }
912 DaemonToWorker::SetTranscriptSource { cli_session_id } => {
913 let applied = adapter.lock().await.set_transcript_source(&cli_session_id);
914 if applied {
915 info!(session = %init.session_id, adapter = %init.cli_id, transcript_session = %cli_session_id, "transcript source set by user");
916 send_message(&stdout, &WorkerToDaemon::CliSessionId { cli_session_id }).await?;
917 }
918 }
919 DaemonToWorker::Init(_) => {}
920 }
921 *processing_since.lock().unwrap() = None;
922 }
923
924 worker_joins.abort_all();
925 while worker_joins.join_next().await.is_some() {}
926 let _ = backend.kill().await;
927 if let Some(marker) = cli_pid_marker {
928 let _ = tokio::fs::remove_file(marker).await;
929 }
930 info!("worker exiting");
931 Ok(())
932}
933
934fn daemon_message_desc(msg: &DaemonToWorker) -> String {
937 match msg {
938 DaemonToWorker::Message { turn_id, .. } => format!("Message(turn_id={})", turn_id),
939 DaemonToWorker::RawInput { turn_id, .. } => format!("RawInput(turn_id={})", turn_id),
940 DaemonToWorker::Close => "Close".to_string(),
941 DaemonToWorker::Restart => "Restart".to_string(),
942 DaemonToWorker::SetDisplayMode { .. } => "SetDisplayMode".to_string(),
943 DaemonToWorker::TermAction { .. } => "TermAction".to_string(),
944 DaemonToWorker::SpecialKeys { .. } => "SpecialKeys".to_string(),
945 DaemonToWorker::TuiKeys { .. } => "TuiKeys".to_string(),
946 DaemonToWorker::TuiTextInput { .. } => "TuiTextInput".to_string(),
947 DaemonToWorker::RefreshScreen => "RefreshScreen".to_string(),
948 DaemonToWorker::SetTranscriptSource { .. } => "SetTranscriptSource".to_string(),
949 DaemonToWorker::Init(_) => "Init".to_string(),
950 }
951}
952
953#[cfg(test)]
954mod tests {
955 use super::*;
956
957 #[test]
958 fn daemon_message_desc_labels_variants_without_content() {
959 let msg = DaemonToWorker::Message {
960 content: "secret prompt".to_string(),
961 turn_id: "t-1".to_string(),
962 };
963 let desc = daemon_message_desc(&msg);
964 assert!(desc.contains("Message"));
965 assert!(desc.contains("t-1"));
966 assert!(!desc.contains("secret prompt"));
967 assert_eq!(daemon_message_desc(&DaemonToWorker::Close), "Close");
968 assert_eq!(
969 daemon_message_desc(&DaemonToWorker::RefreshScreen),
970 "RefreshScreen"
971 );
972 }
973
974 #[test]
975 fn cli_exit_requires_consecutive_dead_samples() {
976 assert_eq!(CLI_DEAD_CONFIRM_ATTEMPTS, 3);
977 let mut dead = 0;
978 assert!(!should_emit_cli_exit(&mut dead, false));
979 assert_eq!(dead, 1);
980 assert!(!should_emit_cli_exit(&mut dead, false));
981 assert_eq!(dead, 2);
982 assert!(should_emit_cli_exit(&mut dead, false));
983 assert_eq!(dead, 3);
984 assert!(should_emit_cli_exit(&mut dead, false));
985 assert_eq!(dead, 4);
986 }
987
988 #[test]
989 fn cli_exit_counter_resets_on_alive() {
990 let mut dead = 0;
991 assert!(!should_emit_cli_exit(&mut dead, false));
992 assert!(!should_emit_cli_exit(&mut dead, false));
993 assert!(!should_emit_cli_exit(&mut dead, true));
994 assert_eq!(dead, 0);
995 assert!(!should_emit_cli_exit(&mut dead, false));
996 assert_eq!(dead, 1);
997 }
998}