1use std::io::Read as _;
18use std::sync::atomic::{AtomicBool, Ordering};
19use std::sync::Arc;
20use std::sync::Mutex;
21use std::time::{Instant, SystemTime, UNIX_EPOCH};
22
23#[cfg(unix)]
24use libc;
25
26use base64::Engine as _;
27use portable_pty::{native_pty_system, CommandBuilder, PtySize};
28use tokio::sync::mpsc;
29
30use super::{
31 BlockControllerRuntimeStatus, BlockInputUnion, Controller, META_KEY_CMD, META_KEY_CMD_ARGS,
32 META_KEY_CMD_CLEAR_ON_START, META_KEY_CMD_CLOSE_ON_EXIT, META_KEY_CMD_CLOSE_ON_EXIT_DELAY,
33 META_KEY_CMD_CLOSE_ON_EXIT_FORCE, META_KEY_CMD_ENV, META_KEY_CMD_RUN_ONCE,
34 META_KEY_CMD_RUN_ON_START, META_KEY_CONNECTION, STATUS_DONE, STATUS_INIT, STATUS_RUNNING,
35};
36use crate::backend::eventbus::EventBus;
37use crate::backend::shellexec::{ConnInterface, ShellProc};
38use crate::backend::storage::filestore::FileStore;
39use crate::backend::storage::store::Store;
40use crate::backend::obj::{self, MetaMapType, RuntimeOpts};
41use crate::backend::wps;
42
43const SHELL_INPUT_CH_SIZE: usize = 256;
55
56#[cfg(windows)]
63fn detect_local_shell_path_windows() -> String {
64 use std::os::windows::process::CommandExt;
65 use std::process::Command;
66 const CREATE_NO_WINDOW: u32 = 0x08000000;
67 if Command::new("where")
69 .arg("pwsh")
70 .creation_flags(CREATE_NO_WINDOW)
71 .output()
72 .map(|o| o.status.success())
73 .unwrap_or(false)
74 {
75 return "pwsh".to_string();
76 }
77 if Command::new("where")
79 .arg("powershell")
80 .creation_flags(CREATE_NO_WINDOW)
81 .output()
82 .map(|o| o.status.success())
83 .unwrap_or(false)
84 {
85 return "powershell".to_string();
86 }
87 "cmd.exe".to_string()
88}
89
90#[cfg(not(windows))]
92fn detect_local_shell_path_windows() -> String {
93 "cmd.exe".to_string()
94}
95
96const PTY_READ_BUF_SIZE: usize = 4096;
98
99#[allow(dead_code)]
102const KILL_GRACE_SECS: u64 = 5;
103
104struct ShellControllerInner {
105 proc_status: String,
107 proc_exit_code: i32,
109 status_version: i32,
111 conn_name: String,
113 input_tx: Option<mpsc::UnboundedSender<BlockInputUnion>>,
116 #[allow(dead_code)]
118 input_rx: Option<mpsc::UnboundedReceiver<BlockInputUnion>>,
119 child_pid: Option<u32>,
121 spawn_ts_ms: Option<i64>,
123 last_pty_output: Option<Instant>,
125 is_agent_pane: bool,
127 input_seq_next: u64,
129 input_seq_buf: std::collections::BTreeMap<u64, BlockInputUnion>,
131}
132
133pub type ConnFactory =
136 Box<dyn Fn(&str, &MetaMapType) -> Result<Box<dyn ConnInterface>, String> + Send + Sync>;
137
138pub struct ShellController {
140 controller_type: String,
142 tab_id: String,
143 block_id: String,
144 run_lock: Arc<AtomicBool>,
146 inner: Arc<Mutex<ShellControllerInner>>,
148 conn_factory: Mutex<Option<ConnFactory>>,
150 broker: Option<Arc<wps::Broker>>,
152 #[allow(dead_code)]
154 event_bus: Option<Arc<EventBus>>,
155 wstore: Option<Arc<Store>>,
157}
158
159impl ShellController {
160 pub fn new(
162 controller_type: String,
163 tab_id: String,
164 block_id: String,
165 broker: Option<Arc<wps::Broker>>,
166 event_bus: Option<Arc<EventBus>>,
167 wstore: Option<Arc<Store>>,
168 ) -> Self {
169 Self {
170 controller_type,
171 tab_id,
172 block_id,
173 run_lock: Arc::new(AtomicBool::new(false)),
174 inner: Arc::new(Mutex::new(ShellControllerInner {
175 proc_status: STATUS_INIT.to_string(),
176 proc_exit_code: 0,
177 status_version: 0,
178 conn_name: String::new(),
179 input_tx: None,
180 input_rx: None,
181 child_pid: None,
182 spawn_ts_ms: None,
183 last_pty_output: None,
184 is_agent_pane: false,
185 input_seq_next: 0,
186 input_seq_buf: std::collections::BTreeMap::new(),
187 })),
188 conn_factory: Mutex::new(None),
189 broker,
190 event_bus,
191 wstore,
192 }
193 }
194
195 #[allow(dead_code)]
197 pub fn set_conn_factory(&self, factory: ConnFactory) {
198 *self.conn_factory.lock().unwrap() = Some(factory);
199 }
200
201 fn try_lock_run(&self) -> bool {
203 self.run_lock
204 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
205 .is_ok()
206 }
207
208 fn unlock_run(&self) {
210 self.run_lock.store(false, Ordering::SeqCst);
211 }
212
213 fn set_status(inner: &mut ShellControllerInner, status: &str) {
215 inner.proc_status = status.to_string();
216 inner.status_version += 1;
217 }
218
219 fn get_status_snapshot(&self) -> BlockControllerRuntimeStatus {
221 let inner = self.inner.lock().unwrap();
222 BlockControllerRuntimeStatus {
223 blockid: self.block_id.clone(),
224 version: inner.status_version,
225 shellprocstatus: inner.proc_status.clone(),
226 shellprocconnname: inner.conn_name.clone(),
227 shellprocexitcode: inner.proc_exit_code,
228 spawn_ts_ms: inner.spawn_ts_ms,
229 is_agent_pane: inner.is_agent_pane,
230 }
231 }
232
233 pub fn last_output_secs_ago(&self) -> Option<u64> {
235 self.inner.lock().unwrap().last_pty_output.map(|t| t.elapsed().as_secs())
236 }
237
238 #[allow(dead_code)]
240 pub fn is_agent_pane(&self) -> bool {
241 self.inner.lock().unwrap().is_agent_pane
242 }
243
244 fn should_run_on_start(meta: &MetaMapType) -> bool {
246 obj::meta_get_bool(meta, META_KEY_CMD_RUN_ON_START, true)
247 }
248
249 #[allow(dead_code)]
251 fn should_run_once(meta: &MetaMapType) -> bool {
252 obj::meta_get_bool(meta, META_KEY_CMD_RUN_ONCE, false)
253 }
254
255 #[allow(dead_code)]
257 fn should_clear_on_start(meta: &MetaMapType) -> bool {
258 obj::meta_get_bool(meta, META_KEY_CMD_CLEAR_ON_START, false)
259 }
260
261 #[allow(dead_code)]
263 fn should_close_on_exit(meta: &MetaMapType) -> bool {
264 obj::meta_get_bool(meta, META_KEY_CMD_CLOSE_ON_EXIT, false)
265 }
266
267 #[allow(dead_code)]
269 fn should_close_on_exit_force(meta: &MetaMapType) -> bool {
270 obj::meta_get_bool(meta, META_KEY_CMD_CLOSE_ON_EXIT_FORCE, false)
271 }
272
273 #[allow(dead_code)]
275 fn close_on_exit_delay_ms(meta: &MetaMapType) -> u64 {
276 match meta.get(META_KEY_CMD_CLOSE_ON_EXIT_DELAY) {
277 Some(serde_json::Value::Number(n)) => n.as_u64().unwrap_or(2000),
278 _ => 2000,
279 }
280 }
281
282 fn get_conn_name(meta: &MetaMapType) -> String {
284 obj::meta_get_string(meta, META_KEY_CONNECTION, "local")
285 }
286
287 fn get_cmd_str(meta: &MetaMapType) -> String {
289 obj::meta_get_string(meta, META_KEY_CMD, "")
290 }
291
292 fn get_cmd_args(meta: &MetaMapType) -> Vec<String> {
294 match meta.get(META_KEY_CMD_ARGS) {
295 Some(serde_json::Value::Array(arr)) => arr
296 .iter()
297 .filter_map(|v| v.as_str().map(|s| s.to_string()))
298 .collect(),
299 _ => vec![],
300 }
301 }
302
303 fn is_interactive(meta: &MetaMapType) -> bool {
305 obj::meta_get_bool(meta, "cmd:interactive", false)
306 }
307
308 fn publish_status(&self) {
310 if let Some(ref broker) = self.broker {
311 let status = self.get_status_snapshot();
312 super::publish_controller_status(broker, &status);
313 }
314 }
315
316 fn pty_size_from_rt_opts(rt_opts: &Option<serde_json::Value>) -> PtySize {
334 const DEFAULT_PTY_ROWS: u16 = 25;
337 const DEFAULT_PTY_COLS: u16 = 200;
338 let (mut rows, mut cols) = (DEFAULT_PTY_ROWS, DEFAULT_PTY_COLS);
339 if let Some(v) = rt_opts {
340 if let Ok(rt) = serde_json::from_value::<RuntimeOpts>(v.clone()) {
341 let ts = &rt.termsize;
342 if !(ts.rows == 0 && ts.cols == 0) {
344 if ts.cols > 0 {
345 cols = ts.cols.clamp(1, 1000) as u16;
346 }
347 if ts.rows > 0 {
348 rows = ts.rows.clamp(1, 1000) as u16;
349 }
350 }
351 }
352 }
353 PtySize {
354 rows,
355 cols,
356 pixel_width: 0,
357 pixel_height: 0,
358 }
359 }
360}
361
362impl Controller for ShellController {
363 fn start(
364 &self,
365 block_meta: MetaMapType,
366 rt_opts: Option<serde_json::Value>,
367 force: bool,
368 ) -> Result<(), String> {
369 let cmd_str_preview = Self::get_cmd_str(&block_meta);
370 let interactive_preview = Self::is_interactive(&block_meta);
371 tracing::info!(
372 block_id = %self.block_id,
373 controller = %self.controller_type,
374 cmd = %cmd_str_preview,
375 interactive = interactive_preview,
376 force = force,
377 "block start requested"
378 );
379
380 if !force && !Self::should_run_on_start(&block_meta) {
382 tracing::info!(block_id = %self.block_id, "skipping start: run_on_start is false");
383 return Ok(());
384 }
385
386 if !self.try_lock_run() {
388 return Err("controller is already running".to_string());
389 }
390
391 let conn_name = Self::get_conn_name(&block_meta);
393
394 {
396 let mut inner = self.inner.lock().unwrap();
397 Self::set_status(&mut inner, STATUS_RUNNING);
398 inner.conn_name = conn_name.clone();
399 }
400
401 let (input_tx, input_rx) = mpsc::unbounded_channel();
404 {
405 let mut inner = self.inner.lock().unwrap();
406 inner.input_tx = Some(input_tx);
407 inner.input_seq_next = 0;
408 inner.input_seq_buf.clear();
409 }
410
411 self.publish_status();
419
420 let has_factory = self.conn_factory.lock().unwrap().is_some();
422
423 if has_factory {
424 let conn_result = {
426 let factory = self.conn_factory.lock().unwrap();
427 factory.as_ref().unwrap()(&conn_name, &block_meta)
428 };
429
430 let mut conn = match conn_result {
431 Ok(c) => c,
432 Err(e) => {
433 let mut inner = self.inner.lock().unwrap();
434 Self::set_status(&mut inner, STATUS_DONE);
435 inner.proc_exit_code = -1;
436 inner.input_tx = None;
437 self.unlock_run();
438 return Err(format!("failed to create connection: {e}"));
439 }
440 };
441
442 if let Err(e) = conn.start() {
443 let mut inner = self.inner.lock().unwrap();
444 Self::set_status(&mut inner, STATUS_DONE);
445 inner.proc_exit_code = -1;
446 inner.input_tx = None;
447 self.unlock_run();
448 return Err(format!("failed to start process: {e}"));
449 }
450
451 let mut shell_proc = ShellProc::new(conn_name, conn);
452 let _done_rx = shell_proc.take_done_rx();
453 let exit_code = shell_proc.wait_and_signal();
454
455 {
456 let mut inner = self.inner.lock().unwrap();
457 inner.proc_exit_code = exit_code;
458 Self::set_status(&mut inner, STATUS_DONE);
459 inner.input_tx = None;
460 }
461 self.publish_status();
462 self.unlock_run();
463 return Ok(());
464 }
465
466 let pty_system = native_pty_system();
478 let pty_size = Self::pty_size_from_rt_opts(&rt_opts);
479
480 let pair = pty_system.openpty(pty_size).map_err(|e| {
481 tracing::error!(block_id = %self.block_id, error = %e, "failed to open PTY");
482 let mut inner = self.inner.lock().unwrap();
483 Self::set_status(&mut inner, STATUS_DONE);
484 inner.proc_exit_code = -1;
485 inner.input_tx = None;
486 self.unlock_run();
487 format!("failed to open PTY: {e}")
488 })?;
489 tracing::info!(block_id = %self.block_id, rows = pty_size.rows, cols = pty_size.cols, "PTY opened");
490
491 let cmd_str = Self::get_cmd_str(&block_meta);
493 let cmd_args = Self::get_cmd_args(&block_meta);
494 let interactive = Self::is_interactive(&block_meta);
495
496 let agent_id_for_jekt: Option<String> = block_meta
499 .get(META_KEY_CMD_ENV)
500 .and_then(|m| m.as_object())
501 .and_then(|obj| obj.get("AGENTMUX_AGENT_ID"))
502 .and_then(|v| v.as_str())
503 .map(|s| s.to_string())
504 .or_else(|| {
505 let cfg = crate::backend::wconfig::ConfigWatcher::with_config(
506 crate::backend::wconfig::build_default_config(),
507 );
508 cfg.get_settings().cmd_env.get("AGENTMUX_AGENT_ID").cloned()
509 })
510 .or_else(|| std::env::var("WAVEMUX_AGENT_ID").ok());
511
512 let mut cmd = if !cmd_str.is_empty() && (!cmd_args.is_empty() || interactive) {
513 tracing::info!(block_id = %self.block_id, cmd = %cmd_str, args = ?cmd_args, "direct spawn path");
516 let mut c = CommandBuilder::new(&cmd_str);
517 if !cmd_args.is_empty() {
518 let arg_refs: Vec<&str> = cmd_args.iter().map(|s| s.as_str()).collect();
519 c.args(arg_refs);
520 }
521 c
522 } else if !cmd_str.is_empty() {
523 tracing::info!(block_id = %self.block_id, cmd = %cmd_str, "shell-wrapped spawn path");
525 if cfg!(windows) {
526 let mut c = CommandBuilder::new("cmd.exe");
527 c.args(["/C", &cmd_str]);
528 c
529 } else {
530 let mut c = CommandBuilder::new("/bin/sh");
531 c.args(["-c", &cmd_str]);
532 c
533 }
534 } else {
535 let shell_path = if cfg!(windows) {
538 detect_local_shell_path_windows()
539 } else {
540 std::env::var("SHELL").unwrap_or_else(|_| "/bin/bash".to_string())
541 };
542
543 let shell_type = crate::backend::shellintegration::detect_shell_type(&shell_path);
544
545 let shell_home = crate::backend::base::get_home_dir().join(".agentmux");
551 crate::backend::shellintegration::deploy_scripts(&shell_home);
552
553 tracing::info!(block_id = %self.block_id, shell = %shell_path, shell_type = ?shell_type, "interactive shell path");
554
555 let mut c = CommandBuilder::new(&shell_path);
556
557 if let Some(startup) = crate::backend::shellintegration::get_shell_startup(shell_type, &shell_home) {
559 for arg in &startup.extra_args {
560 c.arg(arg);
561 }
562 for (k, v) in &startup.env_vars {
563 c.env(k, v);
564 }
565 }
566
567 c.env("TERM", "xterm-256color");
572 c.env("COLORTERM", "truecolor");
573 c.env("TERM_PROGRAM", "agentmux");
574 c.env("AGENTMUX_BLOCKID", &self.block_id);
575 c.env("AGENTMUX_TABID", &self.tab_id);
576 c.env("AGENTMUX_VERSION", env!("CARGO_PKG_VERSION"));
577
578 let log_dir = dirs::home_dir()
581 .unwrap_or_default()
582 .join(".agentmux")
583 .join("logs");
584 c.env("AGENTMUX_LOG_DIR", log_dir.to_string_lossy().as_ref());
585
586 if let Ok(local_url) = std::env::var("AGENTMUX_LOCAL_URL") {
589 c.env("AGENTMUX_LOCAL_URL", &local_url);
590 }
591
592 c.env("AGENTMUX", "1");
597
598 {
617 let sep = if cfg!(windows) { ";" } else { ":" };
618 let current_path = std::env::var("PATH").unwrap_or_default();
619 let mut prepend: Vec<String> = Vec::new();
620 let mut append: Vec<String> = Vec::new();
621
622 if let Some(bundled_bin) = crate::backend::tool_store::bundled_tools_dir() {
624 if bundled_bin.exists() {
625 let bashwrap_exe = if cfg!(windows) {
633 "agentmux-bashwrap.exe"
634 } else {
635 "agentmux-bashwrap"
636 };
637 let bw = bundled_bin.join(bashwrap_exe);
638 if bw.exists() {
639 tracing::info!(
640 target: "agent-tools",
641 path = %bw.display(),
642 "agent bashwrap: bundled (version-locked, prepended to PATH)"
643 );
644 } else {
645 tracing::warn!(
646 target: "agent-tools",
647 dir = %bundled_bin.display(),
648 "agent bashwrap: bundled store present but agentmux-bashwrap MISSING — agent will resolve via system PATH (risk of a stale copy; see RETRO_BASHWRAP_STALE_BUNDLE_2026_06_13.md)"
649 );
650 }
651 prepend.push(bundled_bin.to_string_lossy().into_owned());
652 } else {
653 tracing::warn!(
654 target: "agent-tools",
655 "agent bashwrap: no bundled tools dir — agent will resolve agentmux-bashwrap via system PATH (risk of a stale copy; see RETRO_BASHWRAP_STALE_BUNDLE_2026_06_13.md)"
656 );
657 }
658 }
659
660 if let Some(user_bin) = crate::backend::tool_store::user_tools_dir() {
662 if user_bin.exists() {
663 append.push(user_bin.to_string_lossy().into_owned());
664 }
665 }
666
667 if !prepend.is_empty() || !append.is_empty() {
668 let mut parts = prepend;
669 if !current_path.is_empty() {
670 parts.push(current_path);
671 }
672 parts.extend(append);
673 c.env("PATH", parts.join(sep));
674 }
675 }
676
677 let mut has_agent_id = false;
681
682 let config = crate::backend::wconfig::ConfigWatcher::with_config(
684 crate::backend::wconfig::build_default_config(),
685 );
686 let settings = config.get_settings();
687 for (k, v) in &settings.cmd_env {
688 if k == "AGENTMUX_AGENT_ID" {
689 has_agent_id = true;
690 }
691 let expanded = crate::backend::base::expand_home_dir_safe(v);
692 c.env(k, expanded.to_string_lossy().as_ref());
693 }
694
695 if let Some(env_map) = block_meta.get(META_KEY_CMD_ENV) {
697 if let Some(obj) = env_map.as_object() {
698 for (k, v) in obj {
699 if let Some(val) = v.as_str() {
700 if k == "AGENTMUX_AGENT_ID" {
701 has_agent_id = true;
702 }
703 let expanded = crate::backend::base::expand_home_dir_safe(val);
704 c.env(k, expanded.to_string_lossy().as_ref());
705 }
706 }
707 }
708 }
709
710 if !has_agent_id {
716 c.env_remove("AGENTMUX_AGENT_ID");
717 c.env_remove("AGENTMUX_AGENT_COLOR");
718 c.env_remove("AGENTMUX_AGENT_TEXT_COLOR");
719 c.env_remove("WAVEMUX_AGENT_ID");
720 c.env_remove("WAVEMUX_AGENT_COLOR");
721 }
722
723 c
724 };
725
726 let cwd = obj::meta_get_string(&block_meta, super::META_KEY_CMD_CWD, "");
728 if !cwd.is_empty() {
729 cmd.cwd(&cwd);
730 }
731
732 let mut child = pair.slave.spawn_command(cmd).map_err(|e| {
733 tracing::error!(block_id = %self.block_id, error = %e, cmd = %cmd_str, "spawn failed");
734 let mut inner = self.inner.lock().unwrap();
735 Self::set_status(&mut inner, STATUS_DONE);
736 inner.proc_exit_code = -1;
737 inner.input_tx = None;
738 self.unlock_run();
739 format!("failed to spawn command: {e}")
740 })?;
741 tracing::info!(block_id = %self.block_id, "process spawned successfully");
742
743 let is_agent = agent_id_for_jekt.is_some()
745 || cmd_str.to_lowercase().contains("claude")
746 || cmd_str.to_lowercase().contains("codex")
747 || cmd_str.to_lowercase().contains("gemini")
748 || cmd_str.to_lowercase().contains("qwen");
749
750 let spawn_ts_ms = SystemTime::now()
752 .duration_since(UNIX_EPOCH)
753 .map(|d| d.as_millis() as i64)
754 .unwrap_or(0);
755 {
756 let mut inner = self.inner.lock().unwrap();
757 if let Some(pid) = child.process_id() {
758 super::pidregistry::register(&self.block_id, pid);
759 inner.child_pid = Some(pid);
760 }
761 inner.spawn_ts_ms = Some(spawn_ts_ms);
762 inner.is_agent_pane = is_agent;
763 }
764
765 if let Some(ref agent_id) = agent_id_for_jekt {
769 match crate::backend::reactive::get_global_handler()
770 .register_agent(agent_id, &self.block_id, Some(&self.tab_id))
771 {
772 Ok(()) => {
773 tracing::info!(
774 block_id = %self.block_id,
775 agent_id = %agent_id,
776 "jekt: auto-registered"
777 );
778 if let Ok(local_url) = std::env::var("AGENTMUX_LOCAL_URL") {
780 let data_dir = crate::backend::base::get_wave_data_dir();
781 crate::backend::reactive::registry::write(
782 &data_dir,
783 agent_id,
784 &local_url,
785 &self.block_id,
786 );
787 }
788 }
789 Err(e) => tracing::warn!(
790 block_id = %self.block_id,
791 agent_id = %agent_id,
792 error = %e,
793 "jekt: auto-register failed"
794 ),
795 }
796 }
797 tracing::info!(
798 block_id = %self.block_id,
799 wstore_present = self.wstore.is_some(),
800 event_bus_present = self.event_bus.is_some(),
801 "[dnd-debug] pre-seed state after spawn"
802 );
803
804 if let Some(ref store) = self.wstore {
807 let effective_cwd = if !cwd.is_empty() {
808 cwd.clone()
809 } else {
810 std::env::current_dir()
811 .map(|p| p.to_string_lossy().to_string())
812 .unwrap_or_default()
813 };
814 tracing::debug!(block_id = %self.block_id, cwd = %effective_cwd, "seeding cmd:cwd");
815 if !effective_cwd.is_empty() {
816 let oref_str = format!("block:{}", self.block_id);
817 let mut meta_update = MetaMapType::new();
818 meta_update.insert(
819 super::META_KEY_CMD_CWD.to_string(),
820 serde_json::Value::String(effective_cwd),
821 );
822 match store.must_get::<crate::backend::obj::Block>(&self.block_id) {
824 Ok(block) if obj::meta_get_string(&block.meta, super::META_KEY_CMD_CWD, "").is_empty() => {
825 match crate::server::service::update_object_meta(store, &oref_str, &meta_update) {
826 Ok(()) => {
827 if let Ok(updated_block) = store.must_get::<crate::backend::obj::Block>(&self.block_id) {
831 if let Some(ref event_bus) = self.event_bus {
832 let update_data = serde_json::to_value(&obj::WaveObjUpdate {
833 updatetype: "update".into(),
834 otype: "block".into(),
835 oid: self.block_id.clone(),
836 obj: Some(obj::wave_obj_to_value(&updated_block)),
837 }).ok();
838 event_bus.broadcast_event(&crate::backend::eventbus::WSEventType {
839 eventtype: "waveobj:update".to_string(),
840 oref: oref_str.clone(),
841 data: update_data,
842 });
843 tracing::info!(block_id = %self.block_id, "cmd:cwd seeded and broadcast to frontend");
844 } else {
845 tracing::warn!(block_id = %self.block_id, "cmd:cwd written to store but no event_bus to broadcast — frontend won't update");
846 }
847 }
848 }
849 Err(e) => {
850 tracing::warn!(block_id = %self.block_id, error = %e, "failed to seed cmd:cwd in store");
851 }
852 }
853 }
854 Ok(_) => {
855 tracing::debug!(block_id = %self.block_id, "cmd:cwd already set, skipping seed");
856 }
857 Err(e) => {
858 tracing::warn!(block_id = %self.block_id, error = %e, "failed to read block for cmd:cwd seed");
859 }
860 }
861 }
862 }
863
864 let reader = pair.master.try_clone_reader().map_err(|e| {
866 let _ = child.kill();
867 let mut inner = self.inner.lock().unwrap();
868 Self::set_status(&mut inner, STATUS_DONE);
869 inner.proc_exit_code = -1;
870 inner.input_tx = None;
871 self.unlock_run();
872 format!("failed to clone PTY reader: {e}")
873 })?;
874
875 let writer = pair.master.take_writer().map_err(|e| {
876 let _ = child.kill();
877 let mut inner = self.inner.lock().unwrap();
878 Self::set_status(&mut inner, STATUS_DONE);
879 inner.proc_exit_code = -1;
880 inner.input_tx = None;
881 self.unlock_run();
882 format!("failed to take PTY writer: {e}")
883 })?;
884
885 let block_id_read = self.block_id.clone();
887 let broker_read = self.broker.clone();
888 let inner_read = self.inner.clone();
889 let is_agent_read = is_agent;
890 tokio::task::spawn_blocking(move || {
891 let mut reader = reader;
892 let mut buf = [0u8; PTY_READ_BUF_SIZE];
893
894 let mut translator: Option<crate::agents::translator::claude::ClaudeTranslator> =
906 if is_agent_read {
907 Some(crate::agents::translator::claude::ClaudeTranslator::new())
908 } else {
909 None
910 };
911 let mut line_buf: Vec<u8> = Vec::new();
916
917 let mut osc_extractor: Option<crate::backend::osc_extractor::OscExtractor> =
923 if is_agent_read {
924 Some(crate::backend::osc_extractor::OscExtractor::new())
925 } else {
926 None
927 };
928
929 loop {
930 match reader.read(&mut buf) {
931 Ok(0) => break, Ok(n) => {
933 inner_read.lock().unwrap().last_pty_output = Some(Instant::now());
934 if let Some(ref broker) = broker_read {
935 let raw = &buf[..n];
936
937 let mut cleaned_storage: Vec<u8> = Vec::new();
941 let mut osc_events: Vec<crate::backend::osc_extractor::OscEvent> = Vec::new();
942 if let Some(ref mut ext) = osc_extractor {
943 let (cleaned, evs) = ext.feed(raw);
944 cleaned_storage = cleaned;
945 osc_events = evs;
946 }
947 let chunk: &[u8] = if osc_extractor.is_some() {
948 &cleaned_storage
949 } else {
950 raw
951 };
952
953 handle_append_block_file(
954 broker,
955 &block_id_read,
956 "term",
957 chunk,
958 None, None, );
961
962 for ev in &osc_events {
963 wps::publish_block_activity(broker, &block_id_read, &ev.payload);
964 }
965
966 if let Some(ref mut t) = translator {
967 accumulate_and_translate(
968 broker,
969 &block_id_read,
970 &mut line_buf,
971 chunk,
972 t,
973 );
974 }
975 }
976 }
977 Err(e) => {
978 tracing::debug!("PTY read error for {}: {}", block_id_read, e);
979 break;
980 }
981 }
982 }
983 });
984
985 let master = pair.master;
988 tokio::spawn(async move {
989 let mut writer = writer;
990 let mut input_rx = input_rx;
991 while let Some(input) = input_rx.recv().await {
992 if let Some(data) = input.input_data {
993 use std::io::Write;
994 if let Err(e) = writer.write_all(&data) {
995 tracing::debug!("PTY write error: {}", e);
996 break;
997 }
998 }
999 if let Some(ref size) = input.term_size {
1000 let pty_size = PtySize {
1001 rows: size.rows as u16,
1002 cols: size.cols as u16,
1003 pixel_width: 0,
1004 pixel_height: 0,
1005 };
1006 if let Err(e) = master.resize(pty_size) {
1007 tracing::debug!("PTY resize error: {}", e);
1008 }
1009 }
1010 if input.sig_name.is_some() {
1011 break;
1013 }
1014 }
1015 });
1017
1018 let inner_wait = Arc::clone(&self.inner);
1020 let block_id_wait = self.block_id.clone();
1021 let agent_id_wait = agent_id_for_jekt.clone();
1022 let broker_wait = self.broker.clone();
1023 let run_lock = Arc::clone(&self.run_lock);
1024 tokio::task::spawn_blocking(move || {
1025 let mut child = child;
1026
1027 let exit_status = child.wait();
1029 let exit_code = match exit_status {
1030 Ok(status) => {
1031 if status.success() {
1032 0
1033 } else {
1034 1
1036 }
1037 }
1038 Err(e) => {
1039 tracing::warn!("wait error for block {}: {}", block_id_wait, e);
1040 -1
1041 }
1042 };
1043
1044 tracing::info!(block_id = %block_id_wait, exit_code = exit_code, "process exited");
1045
1046 super::pidregistry::unregister(&block_id_wait);
1048
1049 crate::backend::reactive::get_global_handler().unregister_block(&block_id_wait);
1052
1053 if let Some(ref agent_id) = agent_id_wait {
1055 let data_dir = crate::backend::base::get_wave_data_dir();
1056 crate::backend::reactive::registry::remove(&data_dir, agent_id);
1057 if let Some(sub) = crate::muxbus::cloud_subscriber::get_global_subscriber() {
1058 sub.remove_agent(agent_id);
1059 }
1060 }
1061
1062 {
1064 let mut inner = inner_wait.lock().unwrap();
1065 inner.proc_exit_code = exit_code;
1066 ShellController::set_status(&mut inner, STATUS_DONE);
1067 inner.input_tx = None;
1068 }
1069
1070 if let Some(ref broker) = broker_wait {
1072 let status = {
1073 let inner = inner_wait.lock().unwrap();
1074 BlockControllerRuntimeStatus {
1075 blockid: block_id_wait.clone(),
1076 version: inner.status_version,
1077 shellprocstatus: inner.proc_status.clone(),
1078 shellprocconnname: inner.conn_name.clone(),
1079 shellprocexitcode: inner.proc_exit_code,
1080 spawn_ts_ms: inner.spawn_ts_ms,
1081 is_agent_pane: inner.is_agent_pane,
1082 }
1083 };
1084 super::publish_controller_status(broker, &status);
1085 }
1086
1087 run_lock.store(false, Ordering::SeqCst);
1089 });
1090
1091 Ok(())
1093 }
1094
1095 fn stop(&self, _graceful: bool, new_status: &str) -> Result<(), String> {
1096 #[allow(unused_variables)] let pid_to_kill = {
1099 let mut inner = self.inner.lock().unwrap();
1100 if inner.proc_status == new_status {
1101 return Ok(());
1102 }
1103 let pid = inner.child_pid;
1104 inner.input_tx = None;
1107 Self::set_status(&mut inner, new_status);
1108 pid
1109 };
1110
1111 #[cfg(unix)]
1117 if let Some(pid) = pid_to_kill {
1118 unsafe { libc::kill(-(pid as libc::pid_t), libc::SIGTERM) };
1120 tokio::spawn(async move {
1121 tokio::time::sleep(tokio::time::Duration::from_secs(KILL_GRACE_SECS)).await;
1122 unsafe { libc::kill(-(pid as libc::pid_t), libc::SIGKILL) };
1123 });
1124 }
1125
1126 Ok(())
1127 }
1128
1129 fn get_runtime_status(&self) -> BlockControllerRuntimeStatus {
1130 self.get_status_snapshot()
1131 }
1132
1133 fn send_input(&self, input: BlockInputUnion, seq: Option<u64>) -> Result<(), String> {
1134 let mut inner = self.inner.lock().unwrap();
1135 let tx = match &inner.input_tx {
1136 Some(tx) => tx.clone(),
1137 None => return Err("controller is not running".to_string()),
1138 };
1139 match seq {
1140 None => tx.send(input).map_err(|e| format!("send_input: {e}")),
1141 Some(s) => {
1142 if s == 0 && inner.input_seq_next > 0 {
1148 tracing::info!(
1149 block_id = %self.block_id,
1150 prev_next = inner.input_seq_next,
1151 new_seq = s,
1152 "input seq reset (session reset detected)"
1153 );
1154 inner.input_seq_next = s;
1155 inner.input_seq_buf.clear();
1156 }
1157
1158 let next = inner.input_seq_next;
1159 if s == next {
1160 inner.input_seq_next += 1;
1163 if let Err(e) = tx.send(input) {
1164 tracing::warn!(
1167 block_id = %self.block_id,
1168 seq = s,
1169 "send_input: input channel closed, discarding packet: {e}"
1170 );
1171 return Ok(());
1172 }
1173 loop {
1175 let expected = inner.input_seq_next;
1176 match inner.input_seq_buf.remove(&expected) {
1177 Some(buffered) => {
1178 inner.input_seq_next += 1;
1179 if let Err(e) = tx.send(buffered) {
1180 tracing::warn!(
1181 block_id = %self.block_id,
1182 seq = expected,
1183 "send_input drain: input channel closed, discarding buffered packet: {e}"
1184 );
1185 }
1186 }
1187 None => break,
1188 }
1189 }
1190 Ok(())
1191 } else if s > next {
1192 if inner.input_seq_buf.len() < SHELL_INPUT_CH_SIZE {
1193 inner.input_seq_buf.insert(s, input);
1194 } else {
1195 tracing::warn!(block_id = %self.block_id, seq = s, "input reorder buffer full, dropping");
1196 }
1197 Ok(())
1198 } else {
1199 tracing::warn!(block_id = %self.block_id, seq = s, next, "duplicate input seq, discarding");
1200 Ok(())
1201 }
1202 }
1203 }
1204 }
1205
1206 fn controller_type(&self) -> &str {
1207 &self.controller_type
1208 }
1209
1210 fn block_id(&self) -> &str {
1211 &self.block_id
1212 }
1213
1214 fn as_any(&self) -> &dyn std::any::Any {
1215 self
1216 }
1217}
1218
1219const AGENT_LINE_BUFFER_CAP: usize = 1024 * 1024;
1227
1228fn extract_agent_events(
1246 line_buf: &mut Vec<u8>,
1247 chunk: &[u8],
1248 translator: &mut crate::agents::translator::claude::ClaudeTranslator,
1249) -> Vec<crate::agents::types::AgentEvent> {
1250 use crate::agents::translator::Translator as _;
1251 let mut out: Vec<crate::agents::types::AgentEvent> = Vec::new();
1252 line_buf.extend_from_slice(chunk);
1253 if line_buf.len() > AGENT_LINE_BUFFER_CAP {
1254 line_buf.clear();
1257 return out;
1258 }
1259 while let Some(nl) = line_buf.iter().position(|&b| b == b'\n') {
1260 let line_bytes: Vec<u8> = line_buf.drain(..=nl).collect();
1261 let line = String::from_utf8_lossy(&line_bytes);
1265 let trimmed = line.trim_end_matches(['\n', '\r']);
1266 if !trimmed.starts_with('{') {
1267 continue;
1270 }
1271 let Ok(frame) = serde_json::from_str::<serde_json::Value>(trimmed) else {
1272 continue;
1273 };
1274 out.extend(translator.translate(frame));
1275 }
1276 out
1277}
1278
1279fn accumulate_and_translate(
1284 broker: &wps::Broker,
1285 block_id: &str,
1286 line_buf: &mut Vec<u8>,
1287 chunk: &[u8],
1288 translator: &mut crate::agents::translator::claude::ClaudeTranslator,
1289) {
1290 for event in extract_agent_events(line_buf, chunk, translator) {
1291 broker.publish(wps::WaveEvent {
1292 event: format!("agent_event:{}", block_id),
1293 scopes: vec![],
1294 sender: String::new(),
1295 persist: 0,
1296 data: Some(serde_json::to_value(&event).unwrap_or_default()),
1297 });
1298 }
1299}
1300
1301pub fn persist_to_blockfile_silent(
1310 block_id: &str,
1311 filename: &str,
1312 data: &[u8],
1313 filestore: Option<&Arc<FileStore>>,
1314 global_output_zone: Option<&str>,
1315) {
1316 if let Some(fs) = filestore {
1317 let needs_create = match fs.stat(block_id, filename) {
1318 Ok(None) => true,
1319 Ok(Some(_)) => false,
1320 Err(e) => {
1321 tracing::warn!(
1322 block_id = %block_id, filename = %filename, error = %e,
1323 "persist_silent: stat failed; skipping"
1324 );
1325 return;
1326 }
1327 };
1328 if needs_create {
1329 if let Err(e) = fs.make_file(
1330 block_id,
1331 filename,
1332 std::collections::HashMap::new(),
1333 crate::backend::storage::filestore::FileOpts::default(),
1334 ) {
1335 use crate::backend::storage::error::StoreError;
1336 if !matches!(e, StoreError::AlreadyExists) {
1337 tracing::warn!(
1338 block_id = %block_id, filename = %filename, error = %e,
1339 "persist_silent: make_file failed; skipping"
1340 );
1341 return;
1342 }
1343 }
1344 }
1345 if let Err(e) = fs.append_data(block_id, filename, data) {
1346 tracing::warn!(
1347 block_id = %block_id, filename = %filename, error = %e,
1348 "persist_silent: append_data failed"
1349 );
1350 }
1351 }
1352 if let Some(zone) = global_output_zone {
1353 if let Some(gfs) = crate::backend::agent_session::global_transcript_store() {
1354 mirror_append_to_global(gfs, zone, data);
1355 }
1356 }
1357}
1358
1359pub fn handle_append_block_file(
1365 broker: &wps::Broker,
1366 block_id: &str,
1367 filename: &str,
1368 data: &[u8],
1369 filestore: Option<&Arc<FileStore>>,
1370 global_output_zone: Option<&str>,
1371) {
1372 let data64 = base64::engine::general_purpose::STANDARD.encode(data);
1373
1374 let event_data = wps::WSFileEventData {
1375 zoneid: block_id.to_string(),
1376 filename: filename.to_string(),
1377 fileop: wps::FILE_OP_APPEND.to_string(),
1378 data64,
1379 };
1380
1381 let event = wps::WaveEvent {
1382 event: wps::EVENT_BLOCK_FILE.to_string(),
1383 scopes: vec![format!("block:{block_id}")],
1384 sender: String::new(),
1385 persist: 0,
1386 data: serde_json::to_value(&event_data).ok(),
1387 };
1388
1389 broker.publish(event);
1390
1391 if let Some(fs) = filestore {
1395 let needs_create = match fs.stat(block_id, filename) {
1396 Ok(None) => true,
1397 Ok(Some(_)) => false,
1398 Err(e) => {
1399 tracing::warn!(
1400 block_id = %block_id,
1401 filename = %filename,
1402 error = %e,
1403 "filestore stat failed; skipping write-through"
1404 );
1405 return;
1406 }
1407 };
1408
1409 if needs_create {
1410 if let Err(e) = fs.make_file(
1411 block_id,
1412 filename,
1413 std::collections::HashMap::new(),
1414 crate::backend::storage::filestore::FileOpts::default(),
1415 ) {
1416 use crate::backend::storage::error::StoreError;
1419 if !matches!(e, StoreError::AlreadyExists) {
1420 tracing::warn!(
1421 block_id = %block_id,
1422 filename = %filename,
1423 error = %e,
1424 "filestore make_file failed; skipping write-through"
1425 );
1426 return;
1427 }
1428 }
1429 }
1430
1431 if let Err(e) = fs.append_data(block_id, filename, data) {
1432 tracing::warn!(
1433 block_id = %block_id,
1434 filename = %filename,
1435 error = %e,
1436 "filestore append_data failed"
1437 );
1438 }
1439 }
1446
1447 if let Some(zone) = global_output_zone {
1454 if let Some(gfs) = crate::backend::agent_session::global_transcript_store() {
1455 mirror_append_to_global(gfs, zone, data);
1456 }
1457 }
1458}
1459
1460fn mirror_append_to_global(gfs: &Arc<FileStore>, zone: &str, data: &[u8]) {
1465 use crate::backend::agent_session::OUTPUT_FILE;
1466 use crate::backend::storage::error::StoreError;
1467
1468 match gfs.stat(zone, OUTPUT_FILE) {
1469 Ok(None) => {
1470 if let Err(e) = gfs.make_file(
1471 zone,
1472 OUTPUT_FILE,
1473 std::collections::HashMap::new(),
1474 crate::backend::storage::filestore::FileOpts::default(),
1475 ) {
1476 if !matches!(e, StoreError::AlreadyExists) {
1477 tracing::warn!(zone = %zone, error = %e, "global transcripts: make_file failed; skipping mirror");
1478 return;
1479 }
1480 }
1481 }
1482 Ok(Some(_)) => {}
1483 Err(e) => {
1484 tracing::warn!(zone = %zone, error = %e, "global transcripts: stat failed; skipping mirror");
1485 return;
1486 }
1487 }
1488 if let Err(e) = gfs.append_data(zone, OUTPUT_FILE, data) {
1489 tracing::warn!(zone = %zone, error = %e, "global transcripts: append_data failed");
1490 }
1491}
1492
1493pub(crate) fn resolve_global_output_zone(
1499 wstore: &Option<Arc<crate::backend::storage::store::Store>>,
1500 block_id: &str,
1501) -> Option<String> {
1502 let store = wstore.as_ref()?;
1503 let block = store
1504 .must_get::<crate::backend::obj::Block>(block_id)
1505 .ok()?;
1506 crate::backend::agent_session::agent_zone_for_block_meta(&block.meta)
1507}
1508
1509pub(crate) const OUTPUT_IDX_HEADER_LEN: i64 = 8;
1513
1514pub(crate) fn rebuild_output_idx(
1530 fs: &FileStore,
1531 block_id: &str,
1532 output_size: u64,
1533) -> Option<u64> {
1534 const IDX: &str = "output.idx";
1535 const WIN: i64 = 1 << 20; let mut buf: Vec<u8> = Vec::new();
1539 buf.extend_from_slice(&output_size.to_le_bytes());
1540
1541 let mut line_count: u64 = 0;
1542 let mut cursor: u64 = 0; let mut line_buf: Vec<u8> = Vec::new(); let mut read_pos: i64 = 0;
1545
1546 let flush_line = |line_buf: &mut Vec<u8>,
1547 cursor: &mut u64,
1548 buf: &mut Vec<u8>,
1549 line_count: &mut u64,
1550 had_newline: bool| {
1551 let is_blank = String::from_utf8_lossy(line_buf).trim().is_empty();
1553 if !is_blank {
1554 buf.extend_from_slice(&cursor.to_le_bytes());
1555 *line_count += 1;
1556 }
1557 *cursor += line_buf.len() as u64 + if had_newline { 1 } else { 0 };
1559 line_buf.clear();
1560 };
1561
1562 while read_pos < output_size as i64 {
1563 let (_, chunk) = fs.read_at(block_id, "output", read_pos, WIN).ok()?;
1564 if chunk.is_empty() {
1565 break;
1566 }
1567 for &b in &chunk {
1568 if b == b'\n' {
1569 flush_line(&mut line_buf, &mut cursor, &mut buf, &mut line_count, true);
1570 } else {
1571 line_buf.push(b);
1572 }
1573 }
1574 read_pos += chunk.len() as i64;
1575 }
1576 if !line_buf.is_empty() {
1578 flush_line(&mut line_buf, &mut cursor, &mut buf, &mut line_count, false);
1579 }
1580
1581 if let Ok(None) = fs.stat(block_id, IDX) {
1582 let _ = fs.make_file(
1583 block_id,
1584 IDX,
1585 std::collections::HashMap::new(),
1586 crate::backend::storage::filestore::FileOpts::default(),
1587 );
1588 }
1589 match fs.write_file(block_id, IDX, &buf) {
1590 Ok(()) => {
1591 tracing::info!(block_id = %block_id, lines = line_count, covered = output_size, "output.idx rebuilt");
1592 Some(line_count)
1593 }
1594 Err(e) => {
1595 tracing::warn!(block_id = %block_id, error = %e, "output.idx rebuild write failed");
1596 None
1597 }
1598 }
1599}
1600
1601#[allow(dead_code)]
1604pub fn handle_truncate_block_file(broker: &wps::Broker, block_id: &str, filename: &str) {
1605 let event_data = wps::WSFileEventData {
1606 zoneid: block_id.to_string(),
1607 filename: filename.to_string(),
1608 fileop: wps::FILE_OP_TRUNCATE.to_string(),
1609 data64: String::new(),
1610 };
1611
1612 let event = wps::WaveEvent {
1613 event: wps::EVENT_BLOCK_FILE.to_string(),
1614 scopes: vec![format!("block:{block_id}")],
1615 sender: String::new(),
1616 persist: 0,
1617 data: serde_json::to_value(&event_data).ok(),
1618 };
1619
1620 broker.publish(event);
1621}
1622
1623#[cfg(test)]
1624mod tests {
1625 use super::*;
1626 use crate::backend::shellexec::MockConn;
1627 use std::sync::Arc;
1628
1629 fn make_shell_meta() -> MetaMapType {
1630 let mut meta = MetaMapType::new();
1631 meta.insert(
1632 "controller".to_string(),
1633 serde_json::Value::String("shell".to_string()),
1634 );
1635 meta
1636 }
1637
1638 fn make_cmd_meta(cmd: &str) -> MetaMapType {
1639 let mut meta = MetaMapType::new();
1640 meta.insert(
1641 "controller".to_string(),
1642 serde_json::Value::String("cmd".to_string()),
1643 );
1644 meta.insert(
1645 "cmd".to_string(),
1646 serde_json::Value::String(cmd.to_string()),
1647 );
1648 meta
1649 }
1650
1651 #[test]
1657 fn pty_size_defaults_when_rt_opts_absent() {
1658 let sz = ShellController::pty_size_from_rt_opts(&None);
1659 assert_eq!((sz.rows, sz.cols), (25, 200));
1660 }
1661
1662 #[test]
1663 fn pty_size_defaults_when_termsize_is_serde_default() {
1664 let v = serde_json::json!({ "termsize": { "rows": 0, "cols": 0 } });
1666 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1667 assert_eq!((sz.rows, sz.cols), (25, 200));
1668 }
1669
1670 #[test]
1671 fn pty_size_honors_supplied_termsize() {
1672 let v = serde_json::json!({ "termsize": { "rows": 50, "cols": 130 } });
1673 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1674 assert_eq!((sz.rows, sz.cols), (50, 130));
1675 }
1676
1677 #[test]
1678 fn pty_size_keeps_default_rows_for_cols_only_payload() {
1679 let v = serde_json::json!({ "termsize": { "rows": 0, "cols": 130 } });
1680 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1681 assert_eq!((sz.rows, sz.cols), (25, 130));
1682 }
1683
1684 #[test]
1685 fn pty_size_clamps_oversized_values() {
1686 let v = serde_json::json!({ "termsize": { "rows": 99999, "cols": 99999 } });
1687 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1688 assert_eq!((sz.rows, sz.cols), (1000, 1000));
1689 }
1690
1691 #[test]
1692 fn pty_size_defaults_on_unparseable_rt_opts() {
1693 let v = serde_json::json!({ "totally": "unrelated" });
1695 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1696 assert_eq!((sz.rows, sz.cols), (25, 200));
1697 }
1698
1699 #[test]
1700 fn pty_size_ignores_non_positive_axes() {
1701 let v = serde_json::json!({ "termsize": { "rows": -5, "cols": -1 } });
1703 let sz = ShellController::pty_size_from_rt_opts(&Some(v));
1704 assert_eq!((sz.rows, sz.cols), (25, 200));
1705 }
1706
1707 #[test]
1708 fn test_shell_controller_new() {
1709 let ctrl = ShellController::new(
1710 "shell".to_string(),
1711 "tab-1".to_string(),
1712 "block-1".to_string(),
1713 None,
1714 None,
1715 None,
1716 );
1717 assert_eq!(ctrl.controller_type(), "shell");
1718 assert_eq!(ctrl.block_id(), "block-1");
1719
1720 let status = ctrl.get_runtime_status();
1721 assert_eq!(status.shellprocstatus, STATUS_INIT);
1722 assert_eq!(status.blockid, "block-1");
1723 assert_eq!(status.version, 0);
1724 }
1725
1726 #[test]
1727 fn test_shell_controller_start_stop() {
1728 let ctrl = ShellController::new(
1729 "shell".to_string(),
1730 "tab-1".to_string(),
1731 "block-1".to_string(),
1732 None,
1733 None,
1734 None,
1735 );
1736
1737 ctrl.set_conn_factory(Box::new(|_conn_name, _meta| {
1739 Ok(Box::new(MockConn::new(0)) as Box<dyn ConnInterface>)
1740 }));
1741
1742 let meta = make_shell_meta();
1743 let result = ctrl.start(meta, None, false);
1744 assert!(result.is_ok());
1745
1746 let status = ctrl.get_runtime_status();
1748 assert_eq!(status.shellprocstatus, STATUS_DONE);
1749
1750 let result = ctrl.stop(true, STATUS_DONE);
1752 assert!(result.is_ok());
1753 }
1754
1755 #[test]
1756 fn test_shell_controller_run_on_start_false() {
1757 let ctrl = ShellController::new(
1758 "shell".to_string(),
1759 "tab-1".to_string(),
1760 "block-1".to_string(),
1761 None,
1762 None,
1763 None,
1764 );
1765
1766 let mut meta = make_shell_meta();
1767 meta.insert(
1768 META_KEY_CMD_RUN_ON_START.to_string(),
1769 serde_json::Value::Bool(false),
1770 );
1771
1772 let result = ctrl.start(meta, None, false);
1773 assert!(result.is_ok());
1774
1775 let status = ctrl.get_runtime_status();
1777 assert_eq!(status.shellprocstatus, STATUS_INIT);
1778 }
1779
1780 #[test]
1781 fn test_shell_controller_force_start() {
1782 let ctrl = ShellController::new(
1783 "shell".to_string(),
1784 "tab-1".to_string(),
1785 "block-1".to_string(),
1786 None,
1787 None,
1788 None,
1789 );
1790
1791 ctrl.set_conn_factory(Box::new(|_conn_name, _meta| {
1792 Ok(Box::new(MockConn::new(0)) as Box<dyn ConnInterface>)
1793 }));
1794
1795 let mut meta = make_shell_meta();
1796 meta.insert(
1797 META_KEY_CMD_RUN_ON_START.to_string(),
1798 serde_json::Value::Bool(false),
1799 );
1800
1801 let result = ctrl.start(meta, None, true);
1803 assert!(result.is_ok());
1804
1805 let status = ctrl.get_runtime_status();
1806 assert_eq!(status.shellprocstatus, STATUS_DONE);
1808 }
1809
1810 #[test]
1811 fn test_shell_controller_with_conn_factory() {
1812 let ctrl = ShellController::new(
1813 "cmd".to_string(),
1814 "tab-1".to_string(),
1815 "block-1".to_string(),
1816 None,
1817 None,
1818 None,
1819 );
1820
1821 ctrl.set_conn_factory(Box::new(|_conn_name, _meta| {
1823 Ok(Box::new(MockConn::new(42)) as Box<dyn ConnInterface>)
1824 }));
1825
1826 let meta = make_cmd_meta("echo hello");
1827 let result = ctrl.start(meta, None, true);
1828 assert!(result.is_ok());
1829
1830 let status = ctrl.get_runtime_status();
1831 assert_eq!(status.shellprocstatus, STATUS_DONE);
1832 assert_eq!(status.shellprocexitcode, 42);
1833 }
1834
1835 #[test]
1836 fn test_shell_controller_conn_factory_error() {
1837 let ctrl = ShellController::new(
1838 "shell".to_string(),
1839 "tab-1".to_string(),
1840 "block-1".to_string(),
1841 None,
1842 None,
1843 None,
1844 );
1845
1846 ctrl.set_conn_factory(Box::new(|_conn_name, _meta| {
1847 Err("connection refused".to_string())
1848 }));
1849
1850 let meta = make_shell_meta();
1851 let result = ctrl.start(meta, None, true);
1852 assert!(result.is_err());
1853 assert!(result.unwrap_err().contains("connection refused"));
1854
1855 let status = ctrl.get_runtime_status();
1856 assert_eq!(status.shellprocstatus, STATUS_DONE);
1857 assert_eq!(status.shellprocexitcode, -1);
1858 }
1859
1860 #[test]
1861 fn test_shell_controller_send_input_not_running() {
1862 let ctrl = ShellController::new(
1863 "shell".to_string(),
1864 "tab-1".to_string(),
1865 "block-1".to_string(),
1866 None,
1867 None,
1868 None,
1869 );
1870
1871 let result = ctrl.send_input(BlockInputUnion::data(b"hello".to_vec()), None);
1872 assert!(result.is_err());
1873 assert!(result.unwrap_err().contains("not running"));
1874 }
1875
1876 #[test]
1877 fn test_shell_controller_status_version_increments() {
1878 let ctrl = ShellController::new(
1879 "shell".to_string(),
1880 "tab-1".to_string(),
1881 "block-1".to_string(),
1882 None,
1883 None,
1884 None,
1885 );
1886
1887 ctrl.set_conn_factory(Box::new(|_conn_name, _meta| {
1888 Ok(Box::new(MockConn::new(0)) as Box<dyn ConnInterface>)
1889 }));
1890
1891 let v0 = ctrl.get_runtime_status().version;
1892
1893 let meta = make_shell_meta();
1894 ctrl.start(meta, None, true).unwrap();
1895
1896 let v_after = ctrl.get_runtime_status().version;
1897 assert!(v_after > v0);
1899 }
1900
1901 #[test]
1902 fn test_controller_trait_as_arc() {
1903 let ctrl: Arc<dyn Controller> = Arc::new(ShellController::new(
1904 "shell".to_string(),
1905 "tab-1".to_string(),
1906 "block-1".to_string(),
1907 None,
1908 None,
1909 None,
1910 ));
1911
1912 assert_eq!(ctrl.controller_type(), "shell");
1913 assert_eq!(ctrl.block_id(), "block-1");
1914 let status = ctrl.get_runtime_status();
1915 assert_eq!(status.shellprocstatus, STATUS_INIT);
1916 }
1917
1918 #[test]
1919 fn test_meta_helpers() {
1920 let mut meta = MetaMapType::new();
1921 assert!(ShellController::should_run_on_start(&meta)); assert!(!ShellController::should_run_once(&meta)); assert!(!ShellController::should_clear_on_start(&meta)); assert!(!ShellController::should_close_on_exit(&meta)); meta.insert(
1927 META_KEY_CMD_RUN_ON_START.to_string(),
1928 serde_json::Value::Bool(false),
1929 );
1930 assert!(!ShellController::should_run_on_start(&meta));
1931
1932 meta.insert(
1933 META_KEY_CMD_RUN_ONCE.to_string(),
1934 serde_json::Value::Bool(true),
1935 );
1936 assert!(ShellController::should_run_once(&meta));
1937
1938 meta.insert(
1939 META_KEY_CMD_CLEAR_ON_START.to_string(),
1940 serde_json::Value::Bool(true),
1941 );
1942 assert!(ShellController::should_clear_on_start(&meta));
1943 }
1944
1945 #[test]
1946 fn test_close_on_exit_delay() {
1947 let mut meta = MetaMapType::new();
1948 assert_eq!(ShellController::close_on_exit_delay_ms(&meta), 2000); meta.insert(
1951 META_KEY_CMD_CLOSE_ON_EXIT_DELAY.to_string(),
1952 serde_json::json!(5000),
1953 );
1954 assert_eq!(ShellController::close_on_exit_delay_ms(&meta), 5000);
1955 }
1956
1957 #[test]
1958 fn test_conn_name_from_meta() {
1959 let mut meta = MetaMapType::new();
1960 assert_eq!(ShellController::get_conn_name(&meta), "local"); meta.insert(
1963 META_KEY_CONNECTION.to_string(),
1964 serde_json::Value::String("user@host".to_string()),
1965 );
1966 assert_eq!(ShellController::get_conn_name(&meta), "user@host");
1967 }
1968
1969 #[test]
1970 fn test_handle_append_block_file() {
1971 let broker = wps::Broker::new();
1972
1973 broker.subscribe(
1975 "test-route",
1976 wps::SubscriptionRequest {
1977 event: wps::EVENT_BLOCK_FILE.to_string(),
1978 scopes: vec!["block:block-1".to_string()],
1979 allscopes: false,
1980 },
1981 );
1982
1983 handle_append_block_file(&broker, "block-1", "term", b"hello world", None, None);
1984
1985 let _history = broker.read_event_history(wps::EVENT_BLOCK_FILE, "block:block-1", 10);
1987 }
1990
1991 #[cfg(test)]
1994 fn read_idx(fs: &FileStore, block_id: &str) -> (u64, Vec<u64>) {
1995 let raw = fs.read_file(block_id, "output.idx").unwrap().unwrap();
1996 let covered = u64::from_le_bytes(raw[0..8].try_into().unwrap());
1997 let offsets = raw[8..]
1998 .chunks_exact(8)
1999 .map(|c| u64::from_le_bytes(c.try_into().unwrap()))
2000 .collect();
2001 (covered, offsets)
2002 }
2003
2004 #[test]
2005 fn test_rebuild_output_idx_basic() {
2006 use crate::backend::storage::filestore::FileStore;
2007 let fs = FileStore::open_in_memory().expect("filestore");
2008 let bid = "idx-block";
2009 let data = b"line0\nline1\nline2\n";
2010 fs.make_file(bid, "output", Default::default(), Default::default()).unwrap();
2011 fs.append_data(bid, "output", data).unwrap();
2012
2013 let n = rebuild_output_idx(&fs, bid, data.len() as u64).unwrap();
2014 assert_eq!(n, 3);
2015 let (covered, offsets) = read_idx(&fs, bid);
2016 assert_eq!(covered, data.len() as u64);
2017 assert_eq!(offsets, vec![0, 6, 12]);
2019 }
2020
2021 #[test]
2022 fn test_rebuild_output_idx_blank_and_crlf_and_no_trailing_nl() {
2023 use crate::backend::storage::filestore::FileStore;
2024 let fs = FileStore::open_in_memory().expect("filestore");
2025 let bid = "idx-block2";
2026 let data = b"a\n \nb\r\n\ntail";
2035 fs.make_file(bid, "output", Default::default(), Default::default()).unwrap();
2036 fs.append_data(bid, "output", data).unwrap();
2037
2038 let n = rebuild_output_idx(&fs, bid, data.len() as u64).unwrap();
2039 assert_eq!(n, 3, "a, b(crlf), tail are the 3 non-blank lines");
2040 let (_covered, offsets) = read_idx(&fs, bid);
2041 assert_eq!(offsets, vec![0, 6, 10]);
2042
2043 let full = fs.read_file(bid, "output").unwrap().unwrap();
2045 assert_eq!(&full[0..1], b"a");
2046 assert_eq!(&full[6..7], b"b");
2047 assert_eq!(&full[10..14], b"tail");
2048 }
2049
2050 #[test]
2051 fn test_rebuild_output_idx_empty() {
2052 use crate::backend::storage::filestore::FileStore;
2053 let fs = FileStore::open_in_memory().expect("filestore");
2054 let bid = "idx-empty";
2055 fs.make_file(bid, "output", Default::default(), Default::default()).unwrap();
2056 let n = rebuild_output_idx(&fs, bid, 0).unwrap();
2057 assert_eq!(n, 0);
2058 let (covered, offsets) = read_idx(&fs, bid);
2059 assert_eq!(covered, 0);
2060 assert!(offsets.is_empty());
2061 }
2062
2063 #[test]
2064 fn test_handle_truncate_block_file() {
2065 let broker = wps::Broker::new();
2066 handle_truncate_block_file(&broker, "block-1", "term");
2068 }
2069
2070 #[test]
2071 fn test_register_and_get_controller() {
2072 let ctrl: Arc<dyn Controller> = Arc::new(ShellController::new(
2073 "shell".to_string(),
2074 "tab-1".to_string(),
2075 "test-register-block".to_string(),
2076 None,
2077 None,
2078 None,
2079 ));
2080
2081 super::super::register_controller("test-register-block", ctrl.clone());
2082
2083 let retrieved = super::super::get_controller("test-register-block");
2084 assert!(retrieved.is_some());
2085 assert_eq!(retrieved.unwrap().block_id(), "test-register-block");
2086
2087 super::super::delete_controller("test-register-block");
2089 assert!(super::super::get_controller("test-register-block").is_none());
2090 }
2091
2092 #[test]
2093 fn test_resync_creates_shell_controller() {
2094 use crate::backend::obj::Block;
2095
2096 let mut meta = MetaMapType::new();
2097 meta.insert(
2098 "controller".to_string(),
2099 serde_json::Value::String("shell".to_string()),
2100 );
2101 meta.insert(
2103 META_KEY_CMD_RUN_ON_START.to_string(),
2104 serde_json::Value::Bool(false),
2105 );
2106
2107 let block = Block {
2108 oid: "resync-test-block".to_string(),
2109 version: 1,
2110 meta,
2111 ..Default::default()
2112 };
2113
2114 let result = super::super::resync_controller(&block, "tab-1", None, false, None, None, None, None);
2115 assert!(result.is_ok());
2116
2117 let ctrl = super::super::get_controller("resync-test-block");
2118 assert!(ctrl.is_some());
2119 assert_eq!(ctrl.unwrap().controller_type(), "shell");
2120
2121 super::super::delete_controller("resync-test-block");
2123 }
2124
2125 #[test]
2128 fn test_handle_append_block_file_writes_to_filestore() {
2129 use crate::backend::storage::filestore::FileStore;
2130 use std::sync::Arc;
2131
2132 let broker = wps::Broker::new();
2133 let fs = Arc::new(FileStore::open_in_memory().expect("open in-memory filestore"));
2134
2135 let block_id = "test-block-fs";
2136 let filename = "output";
2137
2138 let line1 = b"line one\n";
2140 handle_append_block_file(&broker, block_id, filename, line1, Some(&fs), None);
2141
2142 let line2 = b"line two\n";
2144 handle_append_block_file(&broker, block_id, filename, line2, Some(&fs), None);
2145
2146 let data = fs.read_file(block_id, filename)
2148 .expect("read_file ok")
2149 .expect("data present");
2150
2151 let text = String::from_utf8(data).expect("valid utf8");
2152 assert!(text.contains("line one"), "expected 'line one' in {:?}", text);
2153 assert!(text.contains("line two"), "expected 'line two' in {:?}", text);
2154
2155 let stat = fs.stat(block_id, filename).unwrap().unwrap();
2157 assert_eq!(stat.size, (line1.len() + line2.len()) as i64);
2158
2159 broker.subscribe(
2161 "test-route-fs",
2162 wps::SubscriptionRequest {
2163 event: wps::EVENT_BLOCK_FILE.to_string(),
2164 scopes: vec![format!("block:{}", block_id)],
2165 allscopes: false,
2166 },
2167 );
2168 handle_append_block_file(&broker, block_id, filename, b"line three\n", Some(&fs), None);
2170 let stat_after = fs.stat(block_id, filename).unwrap().unwrap();
2171 assert_eq!(stat_after.size, (line1.len() + line2.len() + b"line three\n".len()) as i64);
2172 }
2173
2174 #[test]
2179 fn mirror_append_to_global_creates_and_appends() {
2180 use crate::backend::agent_session::OUTPUT_FILE;
2181 let gfs = Arc::new(FileStore::open_in_memory().expect("global filestore"));
2182 let zone = "agent:def-mirror-1:current";
2183
2184 mirror_append_to_global(&gfs, zone, b"{\"type\":\"user\"}\n");
2186 mirror_append_to_global(&gfs, zone, b"{\"type\":\"assistant\"}\n");
2188
2189 let data = gfs
2190 .read_file(zone, OUTPUT_FILE)
2191 .expect("read ok")
2192 .expect("present");
2193 let text = String::from_utf8(data).unwrap();
2194 assert!(text.contains("\"user\""), "got {text:?}");
2195 assert!(text.contains("\"assistant\""), "got {text:?}");
2196 assert_eq!(text.lines().filter(|l| !l.trim().is_empty()).count(), 2);
2198 }
2199
2200 #[test]
2201 fn resolve_global_output_zone_maps_agent_block() {
2202 let wstore = Arc::new(Store::open_in_memory().expect("wstore"));
2203
2204 let oid = uuid::Uuid::new_v4().to_string();
2206 let mut meta = MetaMapType::new();
2207 meta.insert("view".to_string(), serde_json::json!("agent"));
2208 meta.insert("agentId".to_string(), serde_json::json!("def-zone-1"));
2209 let mut block = obj::Block {
2210 oid: oid.clone(),
2211 parentoref: String::new(),
2212 version: 1,
2213 runtimeopts: None,
2214 stickers: None,
2215 meta,
2216 subblockids: None,
2217 };
2218 wstore.insert(&mut block).expect("insert block");
2219
2220 let some = Some(wstore.clone());
2221 assert_eq!(
2222 resolve_global_output_zone(&some, &oid).as_deref(),
2223 Some("agent:def-zone-1:current"),
2224 );
2225
2226 assert_eq!(resolve_global_output_zone(&some, "no-such-block"), None);
2228 assert_eq!(resolve_global_output_zone(&None, &oid), None);
2230 }
2231
2232 use crate::agents::translator::claude::ClaudeTranslator;
2237 use crate::agents::types::AgentEvent;
2238
2239 #[test]
2240 fn extract_agent_events_full_line_translates() {
2241 let mut t = ClaudeTranslator::new();
2242 let mut buf: Vec<u8> = Vec::new();
2243 let line =
2244 br#"{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"hello"}}}
2245"#;
2246 let events = extract_agent_events(&mut buf, line, &mut t);
2247 assert_eq!(events.len(), 1);
2248 match &events[0] {
2249 AgentEvent::AssistantText { delta } => assert_eq!(delta, "hello"),
2250 other => panic!("expected AssistantText, got {other:?}"),
2251 }
2252 assert!(buf.is_empty());
2254 }
2255
2256 #[test]
2257 fn extract_agent_events_chunked_line_accumulates() {
2258 let mut t = ClaudeTranslator::new();
2262 let mut buf: Vec<u8> = Vec::new();
2263 let events = extract_agent_events(
2265 &mut buf,
2266 br#"{"type":"stream_event","event":{"type":"content_"#,
2267 &mut t,
2268 );
2269 assert!(events.is_empty());
2270 let events = extract_agent_events(
2272 &mut buf,
2273 br#"block_delta","delta":{"type":"text_delta","text":"hi"}}}
2274"#,
2275 &mut t,
2276 );
2277 assert_eq!(events.len(), 1);
2278 match &events[0] {
2279 AgentEvent::AssistantText { delta } => assert_eq!(delta, "hi"),
2280 other => panic!("expected AssistantText, got {other:?}"),
2281 }
2282 }
2283
2284 #[test]
2285 fn extract_agent_events_drops_non_json_lines() {
2286 let mut t = ClaudeTranslator::new();
2289 let mut buf: Vec<u8> = Vec::new();
2290 let pty_text = b"\x1b[2K\x1b[0;0H> some prompt\n[m\nplain text\n";
2291 let events = extract_agent_events(&mut buf, pty_text, &mut t);
2292 assert!(events.is_empty(), "got unexpected events: {events:?}");
2293 assert!(buf.is_empty());
2296 }
2297
2298 #[test]
2299 fn extract_agent_events_drops_carriage_returns() {
2300 let mut t = ClaudeTranslator::new();
2303 let mut buf: Vec<u8> = Vec::new();
2304 let line = br#"{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"crlf"}}}
2305"#;
2306 let mut bytes: Vec<u8> = line.to_vec();
2308 let last = bytes.len() - 1;
2309 bytes.insert(last, b'\r');
2310 let events = extract_agent_events(&mut buf, &bytes, &mut t);
2311 assert_eq!(events.len(), 1);
2312 }
2313
2314 #[test]
2315 fn extract_agent_events_drops_malformed_json() {
2316 let mut t = ClaudeTranslator::new();
2319 let mut buf: Vec<u8> = Vec::new();
2320 let events = extract_agent_events(&mut buf, b"{not_valid_json\n", &mut t);
2321 assert!(events.is_empty());
2322 }
2323
2324 #[test]
2325 fn extract_agent_events_resets_oversized_buffer() {
2326 let mut t = ClaudeTranslator::new();
2329 let mut buf: Vec<u8> = Vec::new();
2330 let chunk = vec![b'x'; AGENT_LINE_BUFFER_CAP + 1];
2331 let events = extract_agent_events(&mut buf, &chunk, &mut t);
2332 assert!(events.is_empty());
2333 assert!(
2334 buf.is_empty(),
2335 "buffer should reset past cap, was {} bytes",
2336 buf.len()
2337 );
2338 }
2339
2340 #[test]
2341 fn extract_agent_events_preserves_utf8_across_read_boundary() {
2342 let mut t = ClaudeTranslator::new();
2351 let mut buf: Vec<u8> = Vec::new();
2352 let frame = r#"{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"こんにちは"}}}
2353"#;
2354 let frame_bytes = frame.as_bytes();
2355 let split = frame_bytes
2360 .iter()
2361 .position(|&b| b == 0xe3)
2362 .expect("expected to find a multi-byte codepoint")
2363 + 1; let (a, b) = frame_bytes.split_at(split);
2365 let events = extract_agent_events(&mut buf, a, &mut t);
2366 assert!(events.is_empty());
2368 let events = extract_agent_events(&mut buf, b, &mut t);
2369 assert_eq!(events.len(), 1);
2370 match &events[0] {
2371 AgentEvent::AssistantText { delta } => {
2372 assert_eq!(
2373 delta, "こんにちは",
2374 "UTF-8 must round-trip cleanly across read boundary; got {delta:?}"
2375 );
2376 assert!(!delta.contains('\u{FFFD}'), "no replacement chars");
2377 }
2378 other => panic!("expected AssistantText, got {other:?}"),
2379 }
2380 }
2381
2382 #[test]
2383 fn extract_agent_events_two_lines_one_chunk() {
2384 let mut t = ClaudeTranslator::new();
2387 let mut buf: Vec<u8> = Vec::new();
2388 let two_lines = br#"{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"a"}}}
2389{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"b"}}}
2390"#;
2391 let events = extract_agent_events(&mut buf, two_lines, &mut t);
2392 assert_eq!(events.len(), 2);
2393 }
2394}