use anyhow::{anyhow, Context, Result};
use serde::Serialize;
use std::env;
use std::fs;
use std::path::{Path, PathBuf};
use std::process::Stdio;
use tokio::io::AsyncWriteExt;
use tokio::process::Command;
use crate::shell_integration::{detect_session_shell, rcfile_snippet, Shell};
fn shell_quote(s: &str) -> String {
format!("'{}'", s.replace('\'', r"'\''"))
}
fn ssh_exec_command(ssh_target: &str, remote_cmd: &str) -> Command {
let mut cmd = Command::new("ssh");
cmd.args([
"-o",
"BatchMode=yes",
"-o",
"ConnectTimeout=3",
ssh_target,
"--",
remote_cmd,
]);
cmd
}
pub fn tmux_command(target: Option<&str>, tmux_args: &[&str]) -> Command {
match target {
None => {
let mut cmd = Command::new("tmux");
if let Ok(socket) = std::env::var("MOBUX_TMUX_SOCKET") {
if !socket.is_empty() {
cmd.arg("-L").arg(socket);
}
}
cmd.args(tmux_args);
cmd
}
Some(ssh_target) => {
let quoted: Vec<String> = tmux_args.iter().map(|a| shell_quote(a)).collect();
let remote_cmd = format!("tmux {}", quoted.join(" "));
ssh_exec_command(ssh_target, &remote_cmd)
}
}
}
fn pipe_pane_fifo_path(tmux_bin: &str, session: &str) -> String {
let socket_tag: String = tmux_bin
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '_' })
.collect();
format!("/tmp/mobux-histpipe-{socket_tag}-{session}.fifo")
}
fn tmux_program_and_args(tmux_bin: &str) -> Option<(&str, Vec<&str>)> {
let mut parts = tmux_bin.split_whitespace();
let program = parts.next()?;
Some((program, parts.collect()))
}
fn parse_pane_id(stdout: &[u8]) -> Option<String> {
let id = String::from_utf8_lossy(stdout).trim().to_string();
(!id.is_empty()).then_some(id)
}
pub async fn active_pane_id(tmux_bin: &str, session: &str) -> Option<String> {
let (program, args) = tmux_program_and_args(tmux_bin)?;
let output = Command::new(program)
.args(&args)
.args(["display-message", "-p", "-t", session, "#{pane_id}"])
.output()
.await
.ok()?;
if !output.status.success() {
return None;
}
parse_pane_id(&output.stdout)
}
fn pane_still_current(tapped_pane: &str, current: Option<&str>) -> bool {
match current {
Some(id) => id == tapped_pane,
None => true,
}
}
const PANE_CHECK_INTERVAL: std::time::Duration = std::time::Duration::from_secs(2);
async fn drain_pipe_pane(
mut child: tokio::process::Child,
tx: tokio::sync::mpsc::UnboundedSender<Vec<u8>>,
tmux_bin: String,
session: String,
tapped_pane: String,
check_interval: std::time::Duration,
) {
let Some(mut stdout) = child.stdout.take() else {
eprintln!("history: pipe-pane tap for '{session}' has no stdout handle");
return;
};
let mut stderr = child.stderr.take();
let mut buf = [0u8; 8192];
let mut interval = tokio::time::interval(check_interval);
interval.tick().await;
loop {
tokio::select! {
result = tokio::io::AsyncReadExt::read(&mut stdout, &mut buf) => {
match result {
Ok(0) => {
eprintln!("history: pipe-pane tap for '{session}' ended (reader saw EOF)");
break;
}
Ok(n) => {
if tx.send(buf[..n].to_vec()).is_err() {
break;
}
}
Err(e) => {
eprintln!("history: pipe-pane tap for '{session}' read error: {e}");
break;
}
}
}
_ = interval.tick() => {
let current = active_pane_id(&tmux_bin, &session).await;
if !pane_still_current(&tapped_pane, current.as_deref()) {
eprintln!(
"history: pipe-pane tap for '{session}' stopped — active pane changed from {tapped_pane} to {current:?}"
);
break;
}
}
}
}
let _ = child.kill().await;
let _ = child.wait().await;
if let Some(stderr) = stderr.as_mut() {
let mut err_buf = Vec::new();
let _ = tokio::io::AsyncReadExt::read_to_end(stderr, &mut err_buf).await;
if !err_buf.is_empty() {
eprintln!(
"history: pipe-pane tap for '{session}' stderr: {}",
String::from_utf8_lossy(&err_buf).trim()
);
}
}
}
pub struct PanePipeTap {
tmux_bin: String,
pane_id: String,
fifo_path: String,
}
impl PanePipeTap {
pub async fn start(
tmux_bin: &str,
session: &str,
) -> Option<(Self, tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>)> {
let pane_id = match active_pane_id(tmux_bin, session).await {
Some(id) => id,
None => {
eprintln!(
"history: pipe-pane tap for '{session}' not started — could not resolve an active pane"
);
return None;
}
};
let fifo_path = pipe_pane_fifo_path(tmux_bin, session);
let _ = tokio::fs::remove_file(&fifo_path).await;
match tokio::process::Command::new("mkfifo")
.arg(&fifo_path)
.status()
.await
{
Ok(status) if status.success() => {}
other => {
eprintln!(
"history: pipe-pane tap for '{session}' not started — mkfifo {fifo_path} failed: {other:?}"
);
return None;
}
}
let script = format!(
"{tmux_bin} pipe-pane -t {pane} {cmd} && exec cat {fifo}",
pane = shell_quote(&pane_id),
cmd = shell_quote(&format!("cat >> {fifo_path}")),
fifo = shell_quote(&fifo_path),
);
let child = match tokio::process::Command::new("sh")
.arg("-c")
.arg(&script)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true)
.spawn()
{
Ok(child) => child,
Err(e) => {
eprintln!("history: pipe-pane tap for '{session}' failed to spawn: {e}");
let _ = tokio::fs::remove_file(&fifo_path).await;
return None;
}
};
let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
tokio::spawn(drain_pipe_pane(
child,
tx,
tmux_bin.to_string(),
session.to_string(),
pane_id.clone(),
PANE_CHECK_INTERVAL,
));
Some((
Self {
tmux_bin: tmux_bin.to_string(),
pane_id,
fifo_path,
},
rx,
))
}
}
impl Drop for PanePipeTap {
fn drop(&mut self) {
if let Some((program, args)) = tmux_program_and_args(&self.tmux_bin) {
let _ = std::process::Command::new(program)
.args(&args)
.args(["pipe-pane", "-t", &self.pane_id])
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status();
}
let _ = std::fs::remove_file(&self.fifo_path);
}
}
fn build_write_remote_command(
ssh_target: &str,
dest_dir: &str,
filename: &str,
) -> (Command, String) {
let dest = format!("{}/{}", dest_dir.trim_end_matches('/'), filename);
let remote_cmd = format!(
"mkdir -p {} && cat > {}",
shell_quote(dest_dir),
shell_quote(&dest)
);
(ssh_exec_command(ssh_target, &remote_cmd), dest)
}
async fn stream_stdin_then_wait(mut cmd: Command, data: &[u8], target_desc: &str) -> Result<()> {
cmd.stdin(Stdio::piped());
cmd.stdout(Stdio::piped());
cmd.stderr(Stdio::piped());
let mut child = cmd
.spawn()
.with_context(|| format!("failed to spawn upload to {target_desc}"))?;
let mut stdin = child.stdin.take().context("stdin unavailable for upload")?;
let write_result = stdin.write_all(data).await;
drop(stdin);
let output = child
.wait_with_output()
.await
.with_context(|| format!("upload to {target_desc} failed"))?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("remote upload to {target_desc} failed: {msg}"));
}
write_result.with_context(|| format!("failed to stream upload bytes to {target_desc}"))?;
Ok(())
}
pub async fn write_remote_file(
ssh_target: &str,
dest_dir: &str,
filename: &str,
data: &[u8],
) -> Result<String> {
let (cmd, dest) = build_write_remote_command(ssh_target, dest_dir, filename);
stream_stdin_then_wait(cmd, data, ssh_target).await?;
Ok(dest)
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Session {
pub name: String,
pub windows: i32,
pub attached: i32,
pub created_unix: i64,
}
fn is_no_server_error(msg: &str) -> bool {
msg.contains("failed to connect to server")
|| msg.contains("no server running")
|| msg.contains("error connecting to")
}
pub async fn list_sessions(target: Option<&str>) -> Result<Vec<Session>> {
let output = tmux_command(
target,
&[
"list-sessions",
"-F",
"#{session_windows}:#{session_attached}:#{session_created}:#{session_name}",
],
)
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
if is_no_server_error(&msg) {
return Ok(vec![]);
}
return Err(anyhow!("tmux list-sessions failed: {}", msg));
}
let stdout = String::from_utf8_lossy(&output.stdout);
let mut out = vec![];
for line in stdout.lines() {
let parts: Vec<&str> = line.splitn(4, ':').collect();
if parts.len() != 4 {
continue;
}
out.push(Session {
name: parts[3].to_string(),
windows: parts[0].parse().unwrap_or(0),
attached: parts[1].parse().unwrap_or(0),
created_unix: parts[2].parse().unwrap_or(0),
});
}
out.sort_by(|a, b| a.name.cmp(&b.name));
Ok(out)
}
fn new_session_args(name: &str, shell_cmd: Option<&str>, home: Option<&Path>) -> Vec<String> {
let mut args: Vec<String> = vec!["new-session".into(), "-d".into()];
if let Some(home) = home {
args.push("-c".into());
args.push(home.display().to_string());
}
args.push("-s".into());
args.push(name.into());
if let Some(cmd) = shell_cmd {
args.push(cmd.into());
}
args
}
pub async fn new_session(name: &str, target: Option<&str>) -> Result<()> {
let Some(ssh_target) = target else {
let (shell_type, shell_path) = detect_session_shell();
let shell_cmd = prepare_shell_with_osc133(shell_type, &shell_path)?;
let home = home_dir().ok();
let args = new_session_args(name, Some(&shell_cmd), home.as_deref());
let args_ref: Vec<&str> = args.iter().map(String::as_str).collect();
let output = tmux_command(None, &args_ref)
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux new-session failed: {}", msg));
}
let _ = tmux_command(
None,
&["set-option", "-t", name, "default-command", &shell_cmd],
)
.output()
.await;
return Ok(());
};
let args = new_session_args(name, None, None);
let args_ref: Vec<&str> = args.iter().map(String::as_str).collect();
let output = tmux_command(Some(ssh_target), &args_ref)
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux new-session failed: {}", msg));
}
Ok(())
}
fn prepare_shell_with_osc133(shell: Shell, shell_path: &str) -> Result<String> {
let data_dir = resolve_shell_init_dir()?;
fs::create_dir_all(&data_dir)
.with_context(|| format!("creating shell-init dir: {}", data_dir.display()))?;
match shell {
Shell::Bash => prepare_bash_rcfile(&data_dir, shell_path),
Shell::Zsh => prepare_zsh_zdotdir(&data_dir, shell_path),
Shell::Fish => prepare_fish_command(shell_path),
}
}
fn resolve_shell_init_dir() -> Result<PathBuf> {
let data_dir = if let Ok(override_dir) = env::var("MOBUX_DATA_DIR") {
PathBuf::from(override_dir)
} else {
let dirs = directories::ProjectDirs::from("", "", "mobux")
.ok_or_else(|| anyhow!("could not resolve user home for shell-init dir"))?;
dirs.data_dir().to_path_buf()
};
Ok(data_dir.join("shell-init"))
}
fn prepare_bash_rcfile(shell_init_dir: &Path, shell_path: &str) -> Result<String> {
let rcfile_path = shell_init_dir.join("mobux-bashrc");
let user_bashrc = home_dir()?.join(".bashrc");
let mut content = String::new();
if user_bashrc.exists() {
content.push_str(&format!("source {:?}\n", user_bashrc.display().to_string()));
}
content.push_str(
"
# mobux OSC 133 injection (session-scoped, lazy activation)
_mobux_osc133_ready=0
_mobux_activate_osc133() {
if [[ $_mobux_osc133_ready -eq 0 ]]; then
_mobux_osc133_ready=1
return
fi
unset PROMPT_COMMAND
",
);
content.push_str(rcfile_snippet(Shell::Bash));
content.push_str(
"
}
PROMPT_COMMAND=_mobux_activate_osc133
",
);
fs::write(&rcfile_path, content)
.with_context(|| format!("writing {}", rcfile_path.display()))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mut perms = fs::metadata(&rcfile_path)?.permissions();
perms.set_mode(0o600);
fs::set_permissions(&rcfile_path, perms)?;
}
Ok(format!(
"{} --rcfile {:?}",
shell_path,
rcfile_path.display().to_string()
))
}
fn prepare_zsh_zdotdir(shell_init_dir: &Path, shell_path: &str) -> Result<String> {
let zdotdir = shell_init_dir.join("mobux-zsh");
fs::create_dir_all(&zdotdir)
.with_context(|| format!("creating ZDOTDIR: {}", zdotdir.display()))?;
let zshrc_path = zdotdir.join(".zshrc");
let user_zshrc = home_dir()?.join(".zshrc");
let mut content = String::new();
if user_zshrc.exists() {
content.push_str(&format!("source {:?}\n", user_zshrc.display().to_string()));
}
content.push_str(
"
# mobux OSC 133 injection (session-scoped, lazy activation)
_mobux_osc133_ready=0
_mobux_activate_osc133() {
if [[ $_mobux_osc133_ready -eq 0 ]]; then
_mobux_osc133_ready=1
return
fi
unset -f precmd
",
);
content.push_str(rcfile_snippet(Shell::Zsh));
content.push_str(
"
}
precmd() { _mobux_activate_osc133 }
",
);
fs::write(&zshrc_path, content).with_context(|| format!("writing {}", zshrc_path.display()))?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mut perms = fs::metadata(&zshrc_path)?.permissions();
perms.set_mode(0o600);
fs::set_permissions(&zshrc_path, perms)?;
}
Ok(format!(
"ZDOTDIR={:?} {}",
zdotdir.display().to_string(),
shell_path
))
}
fn prepare_fish_command(shell_path: &str) -> Result<String> {
Ok(shell_path.to_string())
}
fn home_dir() -> Result<PathBuf> {
env::var("HOME")
.map(PathBuf::from)
.map_err(|_| anyhow!("HOME not set"))
}
pub async fn kill_session(name: &str, target: Option<&str>) -> Result<()> {
let output = tmux_command(target, &["kill-session", "-t", name])
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux kill-session failed: {}", msg));
}
Ok(())
}
pub async fn rename_session(old_name: &str, new_name: &str, target: Option<&str>) -> Result<()> {
let output = tmux_command(target, &["rename-session", "-t", old_name, new_name])
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux rename-session failed: {}", msg));
}
Ok(())
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Pane {
pub id: String,
pub index: String,
pub title: String,
pub active: bool,
}
pub async fn list_panes(session: &str, target: Option<&str>) -> Result<Vec<Pane>> {
let output = tmux_command(
target,
&[
"list-windows",
"-t",
session,
"-F",
"#{window_id}:#{window_index}:#{window_active}:#{window_name}",
],
)
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux list-windows failed: {}", msg));
}
let stdout = String::from_utf8_lossy(&output.stdout);
let mut out = vec![];
for line in stdout.lines() {
let parts: Vec<&str> = line.splitn(4, ':').collect();
if parts.len() != 4 {
continue;
}
out.push(Pane {
id: parts[0].to_string(),
index: parts[1].to_string(),
title: parts[3].to_string(),
active: parts[2] == "1",
});
}
Ok(out)
}
pub async fn select_pane(session: &str, window_index: &str, target: Option<&str>) -> Result<()> {
let window_target = format!("{}:{}", session, window_index);
let output = tmux_command(target, &["select-window", "-t", &window_target])
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux select-window failed: {}", msg));
}
Ok(())
}
pub async fn run_command(session: &str, command: &str, target: Option<&str>) -> Result<String> {
let win_target = format!("{}:", session);
let args: Vec<String> = match command {
"new-window" => vec!["new-window".into(), "-t".into(), win_target],
"kill-window" => vec!["kill-window".into(), "-t".into(), win_target],
"split-h" => vec!["split-window".into(), "-h".into(), "-t".into(), win_target],
"split-v" => vec!["split-window".into(), "-v".into(), "-t".into(), win_target],
"next-window" => vec!["next-window".into(), "-t".into(), win_target],
"prev-window" => vec!["previous-window".into(), "-t".into(), win_target],
"next-pane" => vec!["select-pane".into(), "-t".into(), format!("{}:+", session)],
"prev-pane" => vec!["select-pane".into(), "-t".into(), format!("{}:-", session)],
"kill-pane" => vec!["kill-pane".into(), "-t".into(), win_target],
"zoom-pane" => vec!["resize-pane".into(), "-Z".into(), "-t".into(), win_target],
_ => return Err(anyhow!("unknown command: {}", command)),
};
let args_ref: Vec<&str> = args.iter().map(|s| s.as_str()).collect();
let output = tmux_command(target, &args_ref)
.output()
.await
.context("failed to execute tmux")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
if msg.contains("no remaining")
|| msg.contains("session not found")
|| msg.contains("can't find")
|| msg.contains("no current")
{
return Ok(msg);
}
return Err(anyhow!("tmux {} failed: {}", command, msg));
}
Ok(String::from_utf8_lossy(&output.stdout).to_string())
}
pub async fn install_bell_hook(port: u16, token: &str) -> Result<()> {
let hook_cmd = format!(
"run-shell -b 'curl -fsS --max-time 2 \
-H \"X-Mobux-Token: {token}\" \
-X POST \
\"http://127.0.0.1:{port}/internal/trigger?kind=bell&session=#{{hook_session_name}}&window=#{{window_index}}\" \
>/dev/null 2>&1 || true'"
);
let output = Command::new("tmux")
.args(["set-hook", "-g", "alert-bell", &hook_cmd])
.output()
.await
.context("failed to execute tmux set-hook")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux set-hook alert-bell failed: {}", msg));
}
Ok(())
}
pub async fn capture_history(session: &str, lines: i32, target: Option<&str>) -> Result<String> {
let start = format!("-{}", lines);
let output = tmux_command(
target,
&[
"capture-pane",
"-p", "-e", "-S",
&start, "-t",
session,
],
)
.output()
.await
.context("failed to execute tmux capture-pane")?;
if !output.status.success() {
let msg = String::from_utf8_lossy(&output.stderr).trim().to_string();
return Err(anyhow!("tmux capture-pane failed: {}", msg));
}
Ok(String::from_utf8_lossy(&output.stdout).to_string())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn no_server_error_matches_known_tmux_phrasings() {
assert!(is_no_server_error("failed to connect to server"));
assert!(is_no_server_error(
"no server running on /tmp/tmux-1000/default"
));
assert!(is_no_server_error(
"error connecting to /tmp/tmux-1000/default (No such file or directory)"
));
assert!(!is_no_server_error("unknown command: list-sessionz"));
}
#[test]
fn new_session_starts_in_home_directory() {
let args = new_session_args("work", Some("bash"), Some(Path::new("/home/alice")));
assert_eq!(
args,
vec![
"new-session",
"-d",
"-c",
"/home/alice",
"-s",
"work",
"bash"
]
);
}
#[test]
fn new_session_without_home_omits_start_directory() {
let args = new_session_args("work", Some("bash"), None);
assert_eq!(args, vec!["new-session", "-d", "-s", "work", "bash"]);
}
#[test]
fn new_session_remote_omits_shell_cmd_and_home() {
let args = new_session_args("work", None, None);
assert_eq!(args, vec!["new-session", "-d", "-s", "work"]);
}
fn argv(cmd: &Command) -> Vec<String> {
cmd.as_std()
.get_args()
.map(|a| a.to_string_lossy().into_owned())
.collect()
}
#[test]
fn tmux_command_local_is_plain_tmux() {
std::env::remove_var("MOBUX_TMUX_SOCKET");
let cmd = tmux_command(None, &["list-sessions"]);
assert_eq!(
cmd.as_std().get_program().to_string_lossy(),
"tmux",
"no target => plain local tmux"
);
assert_eq!(argv(&cmd), vec!["list-sessions"]);
}
#[test]
fn tmux_command_remote_wraps_in_batch_mode_ssh() {
let cmd = tmux_command(Some("mvhenten@devbox"), &["kill-session", "-t", "work"]);
assert_eq!(cmd.as_std().get_program().to_string_lossy(), "ssh");
assert_eq!(
argv(&cmd),
vec![
"-o",
"BatchMode=yes",
"-o",
"ConnectTimeout=3",
"mvhenten@devbox",
"--",
"tmux 'kill-session' '-t' 'work'",
]
);
}
#[test]
fn tmux_command_remote_quotes_tab_containing_format_strings() {
let cmd = tmux_command(
Some("devbox"),
&["list-sessions", "-F", "#{session_name}\t#{session_windows}"],
);
let remote_cmd = argv(&cmd).pop().expect("remote command string");
assert_eq!(
remote_cmd,
"tmux 'list-sessions' '-F' '#{session_name}\t#{session_windows}'",
);
assert_eq!(remote_cmd.matches('\'').count(), 6, "3 quoted words");
}
#[test]
fn shell_quote_escapes_embedded_single_quotes() {
assert_eq!(shell_quote("plain"), "'plain'");
assert_eq!(shell_quote("it's"), r"'it'\''s'");
}
#[test]
fn pipe_pane_fifo_path_scoped_by_socket_and_session() {
let default_a = pipe_pane_fifo_path("tmux", "work");
let default_b = pipe_pane_fifo_path("tmux", "other");
let socketed = pipe_pane_fifo_path("tmux -L histflake-test", "work");
assert_ne!(
default_a, default_b,
"different sessions must not share a fifo"
);
assert_ne!(
default_a, socketed,
"different tmux sockets must not share a fifo, same session name"
);
assert_eq!(
pipe_pane_fifo_path("tmux -L histflake-test", "work"),
socketed,
"same socket + session is deterministic"
);
}
#[test]
fn write_remote_command_targets_the_node_not_the_hub() {
let (cmd, dest) = build_write_remote_command(
"mvhenten@devbox",
"/tmp/mobux-uploads",
"1700000000000-photo.jpg",
);
assert_eq!(dest, "/tmp/mobux-uploads/1700000000000-photo.jpg");
assert_eq!(cmd.as_std().get_program().to_string_lossy(), "ssh");
assert_eq!(
argv(&cmd),
vec![
"-o",
"BatchMode=yes",
"-o",
"ConnectTimeout=3",
"mvhenten@devbox",
"--",
"mkdir -p '/tmp/mobux-uploads' && cat > '/tmp/mobux-uploads/1700000000000-photo.jpg'",
]
);
}
#[test]
fn tmux_program_and_args_splits_the_socket_flag() {
assert_eq!(
tmux_program_and_args("tmux -L histflake-test"),
Some(("tmux", vec!["-L", "histflake-test"]))
);
assert_eq!(tmux_program_and_args("tmux"), Some(("tmux", vec![])));
assert_eq!(tmux_program_and_args(""), None);
}
#[test]
fn parse_pane_id_trims_and_rejects_empty() {
assert_eq!(parse_pane_id(b"%3\n"), Some("%3".to_string()));
assert_eq!(parse_pane_id(b" \n"), None);
assert_eq!(parse_pane_id(b""), None);
}
#[test]
fn pane_still_current_matches_only_the_tapped_pane() {
assert!(
pane_still_current("%0", Some("%0")),
"same pane => still current"
);
assert!(
!pane_still_current("%0", Some("%1")),
"a confirmed different pane => not current"
);
assert!(
pane_still_current("%0", None),
"a failed query is inconclusive, not a mismatch"
);
}
#[tokio::test]
async fn pipe_pane_tap_does_not_start_when_tmux_is_unreachable() {
let started = PanePipeTap::start("definitely-not-a-real-tmux-xyz", "irrelevant").await;
assert!(
started.is_none(),
"a tmux_bin that can't even resolve the active pane must not start a tap"
);
}
#[tokio::test]
async fn drain_pipe_pane_forwards_bytes_then_closes_on_child_exit() {
let child = tokio::process::Command::new("sh")
.arg("-c")
.arg("printf hello")
.stdout(std::process::Stdio::piped())
.spawn()
.expect("spawn stand-in child");
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
tokio::time::timeout(
std::time::Duration::from_secs(5),
drain_pipe_pane(
child,
tx,
"definitely-not-a-real-tmux-xyz".to_string(),
"irrelevant".to_string(),
"%0".to_string(),
std::time::Duration::from_secs(60),
),
)
.await
.expect("drain must return once the child's stdout hits EOF, not hang");
let mut received = Vec::new();
while let Ok(chunk) = rx.try_recv() {
received.extend(chunk);
}
assert_eq!(received, b"hello");
assert!(
rx.recv().await.is_none(),
"the channel closes once the child exits — a dead feed must be visible \
to the caller, not silently starved"
);
}
#[tokio::test]
async fn drain_pipe_pane_stops_when_the_active_pane_changes() {
let dir = tempfile::tempdir().expect("tempdir");
let stub = dir.path().join("fake-tmux");
std::fs::write(&stub, "#!/bin/sh\necho %99\n").expect("write stub");
let mut perms = std::fs::metadata(&stub).unwrap().permissions();
std::os::unix::fs::PermissionsExt::set_mode(&mut perms, 0o755);
std::fs::set_permissions(&stub, perms).expect("chmod stub");
let child = tokio::process::Command::new("sleep")
.arg("30")
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.spawn()
.expect("spawn stand-in child");
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
tokio::time::timeout(
std::time::Duration::from_secs(5),
drain_pipe_pane(
child,
tx,
stub.to_string_lossy().into_owned(),
"irrelevant".to_string(),
"%0".to_string(),
std::time::Duration::from_millis(20),
),
)
.await
.expect("drain must return once the pane mismatch is detected, not hang forever");
assert!(rx.recv().await.is_none());
}
#[test]
fn write_remote_command_quotes_a_directory_with_a_trailing_slash() {
let (_, dest) = build_write_remote_command("devbox", "/tmp/mobux-uploads/", "f.txt");
assert_eq!(dest, "/tmp/mobux-uploads/f.txt");
}
#[test]
fn write_remote_command_survives_a_single_quote_in_the_filename() {
let (cmd, dest) = build_write_remote_command("devbox", "/tmp/mobux-uploads", "it's.txt");
assert_eq!(dest, "/tmp/mobux-uploads/it's.txt");
let remote_cmd = argv(&cmd).pop().expect("remote command string");
assert_eq!(
remote_cmd,
r"mkdir -p '/tmp/mobux-uploads' && cat > '/tmp/mobux-uploads/it'\''s.txt'",
);
}
#[tokio::test]
async fn write_remote_command_is_injection_safe_against_a_real_shell() {
let dir = tempfile::tempdir().expect("tempdir");
let dest_dir = dir.path().to_str().expect("utf8 tempdir");
let canary = dir.path().join("PWNED");
let adversarial_filename = "$(touch PWNED)`touch PWNED`;touch PWNED\ntouch-PWNED.txt";
let (cmd, dest) = build_write_remote_command("devbox", dest_dir, adversarial_filename);
let remote_cmd = argv(&cmd).pop().expect("remote command string");
let mut sh = tokio::process::Command::new("sh")
.arg("-c")
.arg(&remote_cmd)
.current_dir(dir.path())
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.expect("spawn sh");
sh.stdin
.take()
.expect("stdin")
.write_all(b"payload")
.await
.expect("write payload");
let output = sh.wait_with_output().await.expect("sh -c exited");
assert!(
output.status.success(),
"sh -c failed: {}",
String::from_utf8_lossy(&output.stderr)
);
assert!(
!canary.exists(),
"injected command executed: PWNED canary file exists"
);
assert_eq!(dest, format!("{dest_dir}/{adversarial_filename}"));
let written = fs::read_to_string(&dest)
.expect("uploaded file exists under its literal, unexecuted name");
assert_eq!(written, "payload");
}
#[tokio::test]
async fn stream_stdin_then_wait_prefers_remote_stderr_over_a_broken_pipe_write() {
let mut cmd = Command::new("sh");
cmd.arg("-c").arg("echo 'disk full' >&2; exit 1");
let data = vec![b'x'; 8 * 1024 * 1024];
let err = stream_stdin_then_wait(cmd, &data, "devbox")
.await
.expect_err("a failing remote command must surface as an error");
let msg = err.to_string();
assert!(
msg.contains("disk full"),
"must surface the remote's real stderr, not a generic write error: {msg}"
);
assert!(
!msg.contains("failed to stream"),
"must not report the broken-pipe write error when the remote's own \
stderr explains the failure: {msg}"
);
}
#[tokio::test]
async fn stream_stdin_then_wait_surfaces_the_write_error_when_the_process_still_succeeds() {
let mut cmd = Command::new("sh");
cmd.arg("-c").arg("exec 0<&-; exit 0");
let data = vec![b'x'; 8 * 1024 * 1024];
let err = stream_stdin_then_wait(cmd, &data, "devbox")
.await
.expect_err("a write failure must not be silently swallowed");
assert!(
err.to_string().contains("failed to stream"),
"falls back to the write error when the process reports success: {err}"
);
}
}