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