#![allow(clippy::cast_sign_loss)]
#![allow(clippy::cast_possible_truncation)]
#![allow(clippy::option_if_let_else)]
#![allow(clippy::too_many_lines)]
#![allow(clippy::significant_drop_tightening)]
mod agent;
mod manager;
mod screen;
mod transcript;
pub use agent::{Agent, AgentState as InternalAgentState};
pub use manager::AgentManager;
pub use screen::Screen;
pub use transcript::Transcript;
use crate::protocol::{
AgentInfo, AgentState, AttachEndReason, DumpFormat, Event, ExitReason, Request, Response,
TranscriptEntry,
};
use crate::pty;
use crate::runtime::io::{AsyncReadExt, AsyncWriteExt};
use crate::runtime::net::{OwnedReadHalf, OwnedWriteHalf, UnixStream};
use crate::runtime::sync::{Mutex, broadcast};
use nix::sys::signal::Signal;
#[cfg(unix)]
use std::os::unix::fs::FileTypeExt;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use thiserror::Error;
use tracing::{debug, error, info, instrument, warn};
#[derive(Debug, Error)]
pub enum ServerError {
#[error("failed to bind socket: {0}")]
Bind(#[source] std::io::Error),
#[error("failed to accept connection: {0}")]
Accept(#[source] std::io::Error),
#[error("agent not found: {0}")]
AgentNotFound(String),
#[error("failed to spawn agent: {0}")]
Spawn(#[source] crate::pty::PtyError),
#[error("I/O error: {0}")]
Io(#[source] std::io::Error),
#[error("another server is already running on this socket")]
AlreadyRunning,
}
pub struct Server {
socket_path: PathBuf,
manager: Arc<Mutex<AgentManager>>,
shutdown_tx: broadcast::Sender<()>,
event_tx: broadcast::Sender<Event>,
}
impl Server {
#[must_use]
pub fn new(socket_path: PathBuf) -> Self {
let (shutdown_tx, _) = broadcast::channel(1);
let (event_tx, _) = broadcast::channel(1024);
Self {
socket_path,
manager: Arc::new(Mutex::new(AgentManager::new())),
shutdown_tx,
event_tx,
}
}
#[instrument(skip(self), fields(socket = %self.socket_path.display()))]
pub async fn run(&mut self) -> Result<(), ServerError> {
if self.socket_path.exists() {
let metadata = std::fs::symlink_metadata(&self.socket_path).map_err(ServerError::Io)?;
if metadata.file_type().is_symlink() {
return Err(ServerError::Bind(std::io::Error::other(
"socket path is a symlink - possible security attack",
)));
}
if metadata.file_type().is_socket() {
if UnixStream::connect(&self.socket_path).await.is_ok() {
return Err(ServerError::AlreadyRunning);
}
std::fs::remove_file(&self.socket_path).ok();
} else if metadata.file_type().is_file() {
std::fs::remove_file(&self.socket_path).ok();
}
}
if let Some(parent) = self.socket_path.parent() {
std::fs::create_dir_all(parent).map_err(ServerError::Io)?;
}
let listener = {
#[cfg(unix)]
let _umask_guard = UmaskGuard::new(0o177);
crate::runtime::net::bind_unix_listener(&self.socket_path)
.await
.map_err(ServerError::Bind)?
};
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let perms = std::fs::Permissions::from_mode(0o600);
std::fs::set_permissions(&self.socket_path, perms).map_err(ServerError::Io)?;
}
info!("Server listening on {:?}", self.socket_path);
let manager = Arc::clone(&self.manager);
let event_tx = self.event_tx.clone();
let mut pty_shutdown = self.shutdown_tx.subscribe();
crate::runtime::task::spawn(async move {
crate::runtime::select! {
() = pty_reader_task(manager, event_tx) => {}
_ = pty_shutdown.recv() => {}
}
});
let mut shutdown_rx = self.shutdown_tx.subscribe();
let mut sigterm =
crate::runtime::signal::signal(crate::runtime::signal::SignalKind::terminate())
.map_err(ServerError::Io)?;
let mut sigint =
crate::runtime::signal::signal(crate::runtime::signal::SignalKind::interrupt())
.map_err(ServerError::Io)?;
let mut sighup =
crate::runtime::signal::signal(crate::runtime::signal::SignalKind::hangup())
.map_err(ServerError::Io)?;
loop {
crate::runtime::select! {
result = listener.accept() => {
match result {
Ok((stream, _addr)) => {
debug!("Accepted connection");
let manager = Arc::clone(&self.manager);
let shutdown_tx = self.shutdown_tx.clone();
let event_tx = self.event_tx.clone();
crate::runtime::task::spawn(async move {
if let Err(e) = handle_connection(stream, manager, shutdown_tx, event_tx).await {
error!("Connection error: {}", e);
}
});
}
Err(e) => {
error!("Accept error: {}", e);
}
}
}
_ = shutdown_rx.recv() => {
info!("Shutdown signal received (internal)");
break;
}
_ = sigterm.recv() => {
let mgr = self.manager.lock().await;
let running = mgr.list().filter(|a| a.is_running()).count();
drop(mgr);
if running > 0 {
warn!("SIGTERM received but {} agents still running — ignoring \
(use `vessel shutdown` to force)", running);
} else {
info!("SIGTERM received with no running agents, shutting down");
break;
}
}
_ = sigint.recv() => {
let mgr = self.manager.lock().await;
let running = mgr.list().filter(|a| a.is_running()).count();
drop(mgr);
if running > 0 {
warn!("SIGINT received but {} agents still running — ignoring \
(use `vessel shutdown` to force)", running);
} else {
info!("SIGINT received with no running agents, shutting down");
break;
}
}
_ = sighup.recv() => {
let mgr = self.manager.lock().await;
let running = mgr.list().filter(|a| a.is_running()).count();
drop(mgr);
if running > 0 {
info!("SIGHUP received but {} agents still running, ignoring", running);
} else {
info!("SIGHUP received with no running agents, shutting down");
break;
}
}
}
}
{
let mgr = self.manager.lock().await;
let running: Vec<String> = mgr
.list()
.filter(|a| a.is_running())
.map(|a| a.id.clone())
.collect();
if !running.is_empty() {
info!("Sending SIGTERM to {} running agent(s)", running.len());
for id in &running {
if let Some(agent) = mgr.get(id) {
let _ = agent.pty.signal(Signal::SIGTERM);
}
}
drop(mgr);
let deadline = crate::runtime::time::Instant::now() + Duration::from_secs(5);
loop {
crate::runtime::time::sleep(Duration::from_millis(100)).await;
let mgr = self.manager.lock().await;
let still_running = running
.iter()
.filter(|id| mgr.get(id).is_some_and(agent::Agent::is_running))
.count();
drop(mgr);
if still_running == 0 {
info!("All agents exited gracefully");
break;
}
if crate::runtime::time::Instant::now() >= deadline {
warn!(
"{} agent(s) did not exit in time, sending SIGKILL",
still_running
);
let mgr = self.manager.lock().await;
for id in &running {
if let Some(agent) = mgr.get(id)
&& agent.is_running()
{
let _ = agent.pty.signal(Signal::SIGKILL);
}
}
break;
}
}
}
}
std::fs::remove_file(&self.socket_path).ok();
info!("Server shut down");
Ok(())
}
pub fn shutdown(&self) {
let _ = self.shutdown_tx.send(());
}
}
const MAX_FRAME_BYTES: usize = 1024 * 1024;
const PTY_WRITE_TIMEOUT: Duration = Duration::from_secs(5);
const PTY_WRITE_RETRY: Duration = Duration::from_millis(1);
fn wrap_bracketed_paste(text: &str) -> Vec<u8> {
use crate::protocol::{PASTE_END, PASTE_START};
let mut out = Vec::with_capacity(text.len() + PASTE_START.len() + PASTE_END.len());
out.extend_from_slice(PASTE_START);
let bytes = text.as_bytes();
let mut i = 0;
while i < bytes.len() {
if bytes[i..].starts_with(PASTE_END) {
i += PASTE_END.len();
} else {
out.push(bytes[i]);
i += 1;
}
}
out.extend_from_slice(PASTE_END);
out
}
fn resolve_targets(
mgr: &AgentManager,
id: Option<&str>,
all: bool,
labels: &[String],
proc_filter: Option<&str>,
empty_all_msg: &str,
) -> Result<Vec<String>, String> {
if let Some(agent_id) = id {
return Ok(vec![agent_id.to_string()]);
}
if !all && labels.is_empty() && proc_filter.is_none() {
return Err("must specify agent ID, --label, --proc, or --all".to_string());
}
let matched: Vec<String> = if all {
mgr.list()
.filter(|a| a.is_running())
.map(|a| a.id.clone())
.collect()
} else {
mgr.list()
.filter(|a| {
if !a.is_running() {
return false;
}
if !labels.is_empty() && !a.has_labels(labels) {
return false;
}
if let Some(pf) = proc_filter
&& !a.command.join(" ").contains(pf)
{
return false;
}
true
})
.map(|a| a.id.clone())
.collect()
};
if matched.is_empty() {
if all {
return Err(empty_all_msg.to_string());
}
if proc_filter.is_some() && !labels.is_empty() {
return Err("no agents match the specified process filter and labels".to_string());
}
if proc_filter.is_some() {
return Err("no agents match the specified process filter".to_string());
}
return Err("no agents match the specified labels".to_string());
}
Ok(matched)
}
struct SendTarget {
id: String,
fd: std::os::fd::OwnedFd,
write_lock: Arc<crate::runtime::sync::Mutex<()>>,
}
#[allow(clippy::too_many_arguments)]
async fn collect_send_targets(
manager: &Arc<Mutex<AgentManager>>,
id: Option<&str>,
all: bool,
labels: &[String],
proc_filter: Option<&str>,
selector_used: bool,
command: &str,
recorded_payload: &str,
) -> Result<(Vec<SendTarget>, Vec<crate::protocol::SendOutcome>), String> {
use crate::protocol::SendOutcome;
let mut mgr = manager.lock().await;
let ids = resolve_targets(
&mgr,
id,
all,
labels,
proc_filter,
"no running agents to send to",
)?;
let mut targets = Vec::with_capacity(ids.len());
let mut settled = Vec::new();
for target_id in ids {
let Some(agent) = mgr.get_mut(&target_id) else {
if selector_used {
settled.push(SendOutcome::failed(
target_id,
"agent disappeared".to_string(),
));
continue;
}
return Err(format!("agent not found: {target_id}"));
};
agent.record_command(command, recorded_payload);
match dup_pty_fd(agent) {
Ok(fd) => targets.push(SendTarget {
id: target_id,
fd,
write_lock: Arc::clone(&agent.write_lock),
}),
Err(e) => {
if !selector_used {
return Err(e);
}
settled.push(SendOutcome::failed(target_id, e));
}
}
}
Ok((targets, settled))
}
async fn fan_out_writes(
targets: Vec<SendTarget>,
body: Arc<Vec<u8>>,
submit_key: Arc<Vec<u8>>,
delay_ms: u64,
) -> Vec<crate::protocol::SendOutcome> {
use crate::protocol::SendOutcome;
let handles: Vec<_> = targets
.into_iter()
.map(|target| {
let body = Arc::clone(&body);
let submit_key = Arc::clone(&submit_key);
let id = target.id.clone();
let handle = crate::runtime::task::spawn(async move {
let _write_guard = target.write_lock.lock().await;
write_all_pty(&target.fd, &body).await?;
if !submit_key.is_empty() {
if delay_ms > 0 {
crate::runtime::time::sleep(Duration::from_millis(delay_ms)).await;
}
write_all_pty(&target.fd, &submit_key).await?;
}
Ok::<(), String>(())
});
(id, handle)
})
.collect();
let mut results = Vec::with_capacity(handles.len());
for (id, handle) in handles {
match handle.await {
Ok(Ok(())) => results.push(SendOutcome::delivered(id)),
Ok(Err(e)) => results.push(SendOutcome::failed(id, e)),
Err(e) => results.push(SendOutcome::failed(id, format!("write task failed: {e}"))),
}
}
results
}
fn send_response(
results: Vec<crate::protocol::SendOutcome>,
selector_used: bool,
) -> crate::protocol::Response {
if selector_used {
return Response::SendResults { results };
}
match results.into_iter().next() {
Some(outcome) => match outcome.error {
Some(e) => Response::error(e),
None => Response::Ok,
},
None => Response::error("no agents matched"),
}
}
fn dup_pty_fd(agent: &Agent) -> Result<std::os::fd::OwnedFd, String> {
nix::unistd::dup(crate::sys::borrow_fd(agent.pty.master_fd()))
.map_err(|e| format!("failed to duplicate PTY descriptor: {e}"))
}
async fn write_all_pty(fd: &std::os::fd::OwnedFd, buf: &[u8]) -> Result<(), String> {
let deadline = Instant::now() + PTY_WRITE_TIMEOUT;
let mut written = 0;
while written < buf.len() {
match nix::unistd::write(fd, &buf[written..]) {
Ok(0) => {
return Err(format!(
"write failed: PTY accepted 0 of {} remaining bytes",
buf.len() - written
));
}
Ok(n) => written += n,
Err(nix::errno::Errno::EAGAIN | nix::errno::Errno::EINTR) => {
if Instant::now() >= deadline {
return Err(format!(
"write timed out after {}s: PTY accepted {written} of {} bytes",
PTY_WRITE_TIMEOUT.as_secs(),
buf.len()
));
}
crate::runtime::time::sleep(PTY_WRITE_RETRY).await;
}
Err(e) => return Err(format!("write failed: {e}")),
}
}
Ok(())
}
enum FrameError {
Io(std::io::Error),
TooLarge,
}
struct FrameReader {
reader: OwnedReadHalf,
buf: Vec<u8>,
}
impl FrameReader {
const fn new(reader: OwnedReadHalf) -> Self {
Self {
reader,
buf: Vec::new(),
}
}
async fn next_frame(&mut self) -> Result<Option<Vec<u8>>, FrameError> {
const CHUNK: usize = 16 * 1024;
loop {
if let Some(pos) = self.buf.iter().position(|&b| b == b'\n') {
let mut frame: Vec<u8> = self.buf.drain(..=pos).collect();
frame.pop(); return Ok(Some(frame));
}
if self.buf.len() > MAX_FRAME_BYTES {
return Err(FrameError::TooLarge);
}
let start = self.buf.len();
self.buf.resize(start + CHUNK, 0);
let n = self
.reader
.read(&mut self.buf[start..])
.await
.map_err(FrameError::Io)?;
self.buf.truncate(start + n);
if n == 0 {
return if self.buf.is_empty() {
Ok(None)
} else {
Ok(Some(std::mem::take(&mut self.buf)))
};
}
}
}
fn into_inner(self) -> OwnedReadHalf {
self.reader
}
}
#[instrument(skip_all)]
async fn handle_connection(
stream: UnixStream,
manager: Arc<Mutex<AgentManager>>,
shutdown_tx: broadcast::Sender<()>,
event_tx: broadcast::Sender<Event>,
) -> Result<(), ServerError> {
let (reader, writer) = stream.into_split();
let mut reader = FrameReader::new(reader);
let mut writer = writer;
loop {
let frame = match reader.next_frame().await {
Ok(Some(frame)) => frame,
Ok(None) => {
debug!("Client disconnected");
break;
}
Err(FrameError::Io(e)) => return Err(ServerError::Io(e)),
Err(FrameError::TooLarge) => {
let response = Response::error(format!(
"request frame exceeds maximum size of {MAX_FRAME_BYTES} bytes"
));
let mut json = serde_json::to_string(&response)
.expect("Response serialization should never fail");
json.push('\n');
writer.write_all(json.as_bytes()).await.ok();
debug!("Closing connection: request frame exceeded size limit");
break;
}
};
let request: Request = match serde_json::from_slice(&frame) {
Ok(req) => req,
Err(e) => {
let response = Response::error(format!("invalid request: {e}"));
let mut json = serde_json::to_string(&response)
.expect("Response serialization should never fail");
json.push('\n');
writer.write_all(json.as_bytes()).await.ok();
continue;
}
};
debug!(?request, "Received request");
if let Request::Attach { id, readonly } = &request {
let attach_result = handle_attach(
id.clone(),
*readonly,
reader.into_inner(),
writer,
&manager,
&event_tx,
)
.await;
match attach_result {
Ok(()) => {
debug!("Attach session ended normally");
}
Err(e) => {
if let ServerError::Io(ref io_err) = e {
if io_err.kind() == std::io::ErrorKind::BrokenPipe {
debug!(
"Attach session ended: broken pipe (expected when tmux kills pane)"
);
} else {
warn!("Attach session error: {}", e);
}
} else {
warn!("Attach session error: {}", e);
}
}
}
return Ok(());
}
if let Request::Events {
filter,
include_output,
} = &request
{
let events_result =
handle_events(filter.clone(), *include_output, writer, &event_tx).await;
match events_result {
Ok(()) => {
debug!("Events stream ended normally");
}
Err(e) => {
warn!("Events stream error: {}", e);
}
}
return Ok(());
}
let is_shutdown = matches!(request, Request::Shutdown);
let response = handle_request(request, &manager, &event_tx).await;
let mut json =
serde_json::to_string(&response).expect("Response serialization should never fail");
json.push('\n');
writer
.write_all(json.as_bytes())
.await
.map_err(ServerError::Io)?;
if is_shutdown {
let _ = shutdown_tx.send(());
break;
}
}
Ok(())
}
#[instrument(skip_all)]
async fn handle_request(
request: Request,
manager: &Arc<Mutex<AgentManager>>,
event_tx: &broadcast::Sender<Event>,
) -> Response {
match request {
Request::Ping => Response::Pong,
Request::Spawn {
cmd,
rows,
cols,
name,
labels,
timeout,
max_output,
env,
cwd,
no_resize,
record,
memory_limit,
} => {
if cmd.is_empty() {
return Response::error("command is empty");
}
let resolved_cwd = match cwd {
Some(requested) => match std::fs::canonicalize(&requested) {
Ok(path) => match path.into_os_string().into_string() {
Ok(path) => Some(path),
Err(_) => {
return Response::error(
"spawn failed: resolved working directory is not valid UTF-8",
);
}
},
Err(error) => {
return Response::error(format!(
"spawn failed: cannot resolve working directory {requested:?}: {error}"
));
}
},
None => None,
};
let env_vars: Vec<(String, String)> = env
.iter()
.filter_map(|s| {
let mut parts = s.splitn(2, '=');
match (parts.next(), parts.next()) {
(Some(key), Some(value)) if !key.is_empty() => {
Some((key.to_string(), value.to_string()))
}
_ => None, }
})
.collect();
let limits = if timeout.is_some() || max_output.is_some() {
Some(crate::protocol::ResourceLimits {
timeout,
max_output,
})
} else {
None
};
let wrap_memory_limit = memory_limit.clone();
let mut mgr = manager.lock().await;
let id = if let Some(custom_name) = name {
if custom_name.is_empty() {
return Response::error("agent name cannot be empty");
}
if !custom_name
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '/')
{
return Response::error(
"agent name must contain only alphanumeric characters, hyphens, underscores, and slashes",
);
}
if custom_name.starts_with('/')
|| custom_name.ends_with('/')
|| custom_name.contains("//")
{
return Response::error(
"agent name must not start/end with '/' or contain '//'",
);
}
if custom_name.len() > 64 {
return Response::error("agent name must be 64 characters or fewer");
}
if let Some(existing) = mgr.get(&custom_name) {
if existing.is_running() {
return Response::error(format!(
"agent name already in use: {custom_name}"
));
}
mgr.remove(&custom_name);
}
custom_name
} else {
mgr.generate_id()
};
let effective_cmd = if crate::has_systemd_run() {
let unit_id = id.replace('/', "-");
let mut wrapped = vec![
"systemd-run".to_string(),
"--user".to_string(),
"--scope".to_string(),
"--collect".to_string(),
format!("--unit=vessel-agent-{unit_id}"),
];
if let Some(ref limit) = wrap_memory_limit {
wrapped.extend([
"-p".to_string(),
format!("MemoryMax={limit}"),
"-p".to_string(),
"MemorySwapMax=0".to_string(),
]);
info!(%limit, %id, "Wrapping spawn with systemd-run scope + memory limit");
} else {
info!(%id, "Wrapping spawn with systemd-run scope");
}
wrapped.push("--".to_string());
wrapped.extend(cmd.iter().cloned());
wrapped
} else {
if wrap_memory_limit.is_some() {
warn!(
"--memory-limit requested but systemd-run not available; spawning without cgroup limits"
);
}
cmd.clone()
};
let spawn_env = pty::SpawnEnv { vars: env_vars };
match pty::spawn_with_env(
&effective_cmd,
rows,
cols,
&spawn_env,
resolved_cwd.as_deref(),
) {
Ok(pty_process) => {
let pid = pty_process.pid.as_raw() as u32;
let agent = Agent::new(
id.clone(),
cmd.clone(),
resolved_cwd,
labels.clone(),
limits,
pty_process,
rows,
cols,
no_resize,
record,
);
mgr.add(agent);
info!(%id, %pid, ?labels, ?limits, "Spawned agent");
let _ = event_tx.send(Event::AgentSpawned {
id: id.clone(),
pid,
command: cmd,
labels,
});
Response::Spawned { id, pid }
}
Err(e) => Response::error(format!("spawn failed: {e}")),
}
}
Request::List { labels } => {
let mgr = manager.lock().await;
let agents: Vec<AgentInfo> = mgr
.list()
.filter(|agent| labels.is_empty() || agent.has_labels(&labels))
.map(|agent| {
let elapsed = agent.started_at.elapsed();
let now_millis = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as u64;
let started_at = now_millis.saturating_sub(elapsed.as_millis() as u64);
let rss_bytes = if agent.is_running() {
get_process_tree_rss(agent.pid())
} else {
None
};
AgentInfo {
id: agent.id.clone(),
pid: agent.pid(),
state: match agent.state {
InternalAgentState::Running => AgentState::Running,
InternalAgentState::Exited { .. } => AgentState::Exited,
},
command: agent.command.clone(),
cwd: agent.cwd.clone(),
labels: agent.labels.clone(),
size: agent.screen.size(),
started_at,
exit_code: agent.exit_code(),
exit_reason: agent.exit_reason,
limits: agent.limits,
no_resize: agent.no_resize,
rss_bytes,
}
})
.collect();
Response::Agents { agents }
}
Request::Kill {
id,
labels,
all,
signal,
proc_filter,
} => {
if !(1..=31).contains(&signal) {
return Response::error(format!("invalid signal number: {signal} (must be 1-31)"));
}
let mgr = manager.lock().await;
let targets = match resolve_targets(
&mgr,
id.as_deref(),
all,
&labels,
proc_filter.as_deref(),
"no running agents to kill",
) {
Ok(targets) => targets,
Err(e) => return Response::error(e),
};
let sig = Signal::try_from(signal).unwrap_or(Signal::SIGTERM);
let mut errors = Vec::new();
let mut killed = 0;
for target_id in targets {
if let Some(agent) = mgr.get(&target_id) {
if !agent.is_running() {
info!(%target_id, "Agent already exited, nothing to kill");
continue;
}
match agent.pty.signal(sig) {
Ok(()) => {
info!(%target_id, ?sig, "Sent signal to agent");
killed += 1;
}
Err(e) => {
errors.push(format!("{target_id}: {e}"));
}
}
}
}
if !errors.is_empty() {
Response::error(format!("failed to kill some agents: {}", errors.join(", ")))
} else if let (0, Some(id)) = (killed, &id) {
Response::error(format!("agent not found: {id}"))
} else {
Response::Ok
}
}
Request::Send {
id,
labels,
all,
proc_filter,
data,
newline,
enter,
submit_delay_ms,
paste,
} => {
let selector_used = all || !labels.is_empty() || proc_filter.is_some();
let mut submit_key = Vec::new();
if newline {
submit_key.push(b'\n');
}
if enter {
submit_key.push(b'\r');
}
let body = if paste {
wrap_bracketed_paste(&data)
} else {
data.clone().into_bytes()
};
let recorded = if submit_key.is_empty() {
data.clone()
} else {
format!("{data}\n")
};
let (targets, mut results) = match collect_send_targets(
manager,
id.as_deref(),
all,
&labels,
proc_filter.as_deref(),
selector_used,
"send",
&recorded,
)
.await
{
Ok(v) => v,
Err(e) => return Response::error(e),
};
let delay = submit_delay_ms.unwrap_or(crate::protocol::DEFAULT_SUBMIT_DELAY_MS);
results
.extend(fan_out_writes(targets, Arc::new(body), Arc::new(submit_key), delay).await);
send_response(results, selector_used)
}
Request::SendBytes {
id,
labels,
all,
proc_filter,
data,
} => {
let selector_used = all || !labels.is_empty() || proc_filter.is_some();
let recorded = hex::encode(&data);
let (targets, mut results) = match collect_send_targets(
manager,
id.as_deref(),
all,
&labels,
proc_filter.as_deref(),
selector_used,
"send_bytes",
&recorded,
)
.await
{
Ok(v) => v,
Err(e) => return Response::error(e),
};
results.extend(fan_out_writes(targets, Arc::new(data), Arc::new(Vec::new()), 0).await);
send_response(results, selector_used)
}
Request::Tail {
id,
lines,
follow: _,
} => {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&id) {
let data = agent.transcript.tail_lines(lines);
let exited = !agent.is_running();
Response::Output { data, exited }
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::Dump { id, since, format } => {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&id) {
let entries: Vec<TranscriptEntry> = if let Some(ts) = since {
agent
.transcript
.since(ts)
.into_iter()
.map(|e| TranscriptEntry {
timestamp: e.timestamp,
data: e.data.clone(),
})
.collect()
} else {
agent
.transcript
.all()
.map(|e| TranscriptEntry {
timestamp: e.timestamp,
data: e.data.clone(),
})
.collect()
};
match format {
DumpFormat::Jsonl => Response::Transcript { entries },
DumpFormat::Text => {
let data: Vec<u8> = entries.iter().flat_map(|e| e.data.clone()).collect();
let exited = !agent.is_running();
Response::Output { data, exited }
}
}
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::Snapshot { id, strip_colors } => {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&id) {
let content = if strip_colors {
agent.screen.snapshot()
} else {
agent.screen.contents_formatted()
};
let cursor = agent.screen.cursor_position();
let size = agent.screen.size();
Response::Snapshot {
content,
cursor,
size,
}
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::Attach { id, readonly: _ } => {
let mgr = manager.lock().await;
if mgr.get(&id).is_some() {
Response::error("attach request should not reach handle_request")
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::Events { .. } => {
Response::error("events request should not reach handle_request")
}
Request::Resize {
id,
rows,
cols,
clear_transcript,
} => {
const MIN_SIZE: u16 = 1;
const MAX_SIZE: u16 = 500;
if !(MIN_SIZE..=MAX_SIZE).contains(&rows) || !(MIN_SIZE..=MAX_SIZE).contains(&cols) {
return Response::error(format!(
"invalid dimensions: {cols}x{rows} (must be {MIN_SIZE}-{MAX_SIZE})"
));
}
let mut mgr = manager.lock().await;
if let Some(agent) = mgr.get_mut(&id) {
if let Err(e) = agent.pty.resize(rows, cols) {
return Response::error(format!("resize failed: {e}"));
}
agent.screen.resize(rows, cols);
if clear_transcript {
agent.transcript.clear();
agent.screen_cleared_at = Some(std::time::Instant::now());
if let Err(e) = agent.pty.signal(nix::sys::signal::Signal::SIGWINCH) {
warn!(%id, "Failed to send SIGWINCH after transcript clear: {e}");
}
info!(%id, %rows, %cols, "Resized agent and cleared transcript");
} else {
info!(%id, %rows, %cols, "Resized agent");
}
Response::Ok
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::GetRecording { id } => {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&id) {
if agent.recording {
Response::Recording {
agent_id: id,
commands: agent.recorded_commands.clone(),
}
} else {
Response::error(format!("recording not enabled for agent: {id}"))
}
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::GetEnv { id } => {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&id) {
if agent.is_running() {
let pid = agent.pid();
drop(mgr); match read_proc_environ(pid) {
Ok(env) => Response::AgentEnv { id, env },
Err(e) => {
Response::error(format!("failed to read environment for {id}: {e}"))
}
}
} else {
Response::error(format!(
"agent {id} has exited — environment no longer available"
))
}
} else {
Response::error(format!("agent not found: {id}"))
}
}
Request::Shutdown => {
info!("Shutdown requested");
Response::Ok
}
}
}
#[instrument(skip(reader, writer, manager, event_tx))]
async fn handle_attach(
agent_id: String,
readonly: bool,
mut reader: OwnedReadHalf,
mut writer: OwnedWriteHalf,
manager: &Arc<Mutex<AgentManager>>,
event_tx: &broadcast::Sender<Event>,
) -> Result<(), ServerError> {
let size = {
let mut mgr = manager.lock().await;
if let Some(agent) = mgr.get_mut(&agent_id) {
if !agent.is_running() {
let response = Response::error(format!("agent {agent_id} has exited"));
let mut json = serde_json::to_string(&response)
.expect("Response serialization should never fail");
json.push('\n');
writer.write_all(json.as_bytes()).await.ok();
return Ok(());
}
agent.attached = true;
agent.screen.size()
} else {
let response = Response::error(format!("agent not found: {agent_id}"));
let mut json =
serde_json::to_string(&response).expect("Response serialization should never fail");
json.push('\n');
writer.write_all(json.as_bytes()).await.ok();
return Ok(());
}
};
let response = Response::AttachStarted {
id: agent_id.clone(),
size,
};
let mut json =
serde_json::to_string(&response).expect("Response serialization should never fail");
json.push('\n');
writer
.write_all(json.as_bytes())
.await
.map_err(ServerError::Io)?;
info!("Attach started for agent {agent_id}");
{
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(&agent_id) {
let recently_cleared = agent
.screen_cleared_at
.is_some_and(|t| t.elapsed() < std::time::Duration::from_secs(1));
if recently_cleared {
info!("Screen recently cleared, sending clear screen instead of stale render");
drop(mgr);
writer
.write_all(b"\x1b[2J\x1b[H") .await
.map_err(ServerError::Io)?;
writer.flush().await.map_err(ServerError::Io)?;
} else {
let initial_screen = agent.screen.render_full_screen();
info!(
"Sending initial screen render: {} bytes",
initial_screen.len()
);
drop(mgr); writer
.write_all(&initial_screen)
.await
.map_err(ServerError::Io)?;
writer.flush().await.map_err(ServerError::Io)?;
info!("Initial screen render sent");
}
}
}
let result = run_attach_bridge(&agent_id, readonly, &mut reader, &mut writer, manager).await;
let end_reason = {
let mut mgr = manager.lock().await;
if let Some(agent) = mgr.get_mut(&agent_id) {
agent.attached = false;
}
match &result {
Ok(reason) => reason.clone(),
Err(e) => AttachEndReason::Error {
message: e.to_string(),
},
}
};
if let AttachEndReason::AgentExited { exit_code } = &end_reason {
let _ = event_tx.send(Event::AgentExited {
id: agent_id.clone(),
exit_code: *exit_code,
});
}
let response = Response::AttachEnded { reason: end_reason };
let mut json =
serde_json::to_string(&response).expect("Response serialization should never fail");
json.push('\n');
writer.write_all(json.as_bytes()).await.ok();
info!("Attach ended for agent {}", agent_id);
result.map(|_| ())
}
#[instrument(skip(writer, event_tx))]
async fn handle_events(
filter: Vec<String>,
include_output: bool,
mut writer: OwnedWriteHalf,
event_tx: &broadcast::Sender<Event>,
) -> Result<(), ServerError> {
let mut event_rx = event_tx.subscribe();
info!(?filter, %include_output, "Events subscription started");
loop {
match event_rx.recv().await {
Ok(event) => {
let agent_id = match &event {
Event::AgentSpawned { id, .. }
| Event::AgentOutput { id, .. }
| Event::AgentExited { id, .. } => id,
};
if !filter.is_empty() && !filter.contains(agent_id) {
continue;
}
if !include_output && matches!(event, Event::AgentOutput { .. }) {
continue;
}
let response = Response::Event(event);
let mut json = serde_json::to_string(&response)
.expect("Response serialization should never fail");
json.push('\n');
if writer.write_all(json.as_bytes()).await.is_err() {
debug!("Events client disconnected");
break;
}
}
Err(broadcast::error::RecvError::Closed) => {
debug!("Events channel closed");
break;
}
Err(broadcast::error::RecvError::Lagged(n)) => {
warn!("Events subscriber lagged, missed {n} events");
}
Err(broadcast::error::RecvError::Cancelled) => {
debug!("Events recv cancelled");
break;
}
Err(broadcast::error::RecvError::PolledAfterCompletion) => {
debug!("Events recv polled after completion");
break;
}
}
}
info!("Events subscription ended");
Ok(())
}
async fn run_attach_bridge(
agent_id: &str,
readonly: bool,
reader: &mut OwnedReadHalf,
writer: &mut OwnedWriteHalf,
manager: &Arc<Mutex<AgentManager>>,
) -> Result<AttachEndReason, ServerError> {
let mut input_buf = [0u8; 4096];
let mut output_buf = [0u8; 4096];
let mut poll_interval = crate::runtime::time::interval(Duration::from_millis(10));
loop {
crate::runtime::select! {
result = reader.read(&mut input_buf) => {
match result {
Ok(0) => {
debug!("Client disconnected during attach");
return Ok(AttachEndReason::Detached);
}
Ok(n) => {
if !readonly {
let mgr = manager.lock().await;
if let Some(agent) = mgr.get(agent_id) {
let pty_fd = agent.pty.master_fd();
let borrowed_fd = crate::sys::borrow_fd(pty_fd);
if let Err(e) = nix::unistd::write(borrowed_fd, &input_buf[..n]) {
warn!("Failed to write to PTY: {e}");
return Ok(AttachEndReason::Error {
message: format!("PTY write error: {e}"),
});
}
} else {
return Ok(AttachEndReason::Error {
message: "agent no longer exists".to_string(),
});
}
}
}
Err(e) => {
return Err(ServerError::Io(e));
}
}
}
() = poll_interval.tick() => {
let mut mgr = manager.lock().await;
if let Some(agent) = mgr.get_mut(agent_id) {
if let Ok(Some(code)) = agent.pty.try_wait() {
agent.state = InternalAgentState::Exited { code };
return Ok(AttachEndReason::AgentExited { exit_code: Some(code) });
}
if !agent.is_running() {
return Ok(AttachEndReason::AgentExited {
exit_code: agent.exit_code(),
});
}
let pty_fd = agent.pty.master_fd();
let borrowed_fd = crate::sys::borrow_fd(pty_fd);
match nix::unistd::read(borrowed_fd, &mut output_buf) {
Ok(n) if n > 0 => {
let data = &output_buf[..n];
agent.transcript.append(data);
agent.screen.process(data);
drop(mgr); writer.write_all(data).await.map_err(ServerError::Io)?;
}
Ok(_) | Err(nix::Error::EAGAIN) => {}
Err(nix::Error::EIO) => {
if let Ok(Some(code)) = agent.pty.try_wait() {
agent.state = InternalAgentState::Exited { code };
return Ok(AttachEndReason::AgentExited { exit_code: Some(code) });
}
}
Err(e) => {
warn!("PTY read error: {e}");
}
}
} else {
return Ok(AttachEndReason::Error {
message: "agent no longer exists".to_string(),
});
}
}
}
}
}
async fn pty_reader_task(manager: Arc<Mutex<AgentManager>>, event_tx: broadcast::Sender<Event>) {
use crate::runtime::time::{Duration, interval};
let mut poll_interval = interval(Duration::from_millis(10));
loop {
poll_interval.tick().await;
let mut mgr = manager.lock().await;
let ids: Vec<String> = mgr.list().map(|a| a.id.clone()).collect();
for id in ids {
if let Some(agent) = mgr.get_mut(&id) {
if !agent.is_running() || agent.attached {
continue;
}
if agent.is_timed_out() {
if !agent.sigterm_sent {
info!(%id, "Agent timeout - sending SIGTERM");
let _ = agent.pty.signal(Signal::SIGTERM);
agent.sigterm_sent = true;
agent.sigterm_sent_at = Some(std::time::Instant::now());
} else if agent.should_sigkill() {
info!(%id, "Agent timeout grace period expired - sending SIGKILL");
let _ = agent.pty.signal(Signal::SIGKILL);
}
}
let fd = agent.pty.master_fd();
let mut buf = [0u8; 4096];
let borrowed_fd = crate::sys::borrow_fd(fd);
match nix::unistd::read(borrowed_fd, &mut buf) {
Ok(n) if n > 0 => {
let data = &buf[..n];
agent.transcript.append(data);
agent.screen.process(data);
let _ = event_tx.send(Event::AgentOutput {
id: id.clone(),
data: data.to_vec(),
});
}
Ok(_) | Err(nix::Error::EAGAIN) => {}
Err(nix::Error::EIO) => {
if let Ok(Some(code)) = agent.pty.try_wait() {
agent.state = InternalAgentState::Exited { code };
agent.exit_reason =
Some(if agent.sigterm_sent && (code == 143 || code == 137) {
ExitReason::Timeout
} else {
ExitReason::Normal
});
info!(%id, %code, exit_reason = ?agent.exit_reason, "Agent exited");
let _ = event_tx.send(Event::AgentExited {
id: id.clone(),
exit_code: Some(code),
});
}
}
Err(e) => {
warn!(%id, %e, "PTY read error");
}
}
if agent.is_running()
&& let Ok(Some(code)) = agent.pty.try_wait()
{
agent.state = InternalAgentState::Exited { code };
agent.exit_reason =
Some(if agent.sigterm_sent && (code == 143 || code == 137) {
ExitReason::Timeout
} else {
ExitReason::Normal
});
info!(%id, %code, exit_reason = ?agent.exit_reason, "Agent exited");
let _ = event_tx.send(Event::AgentExited {
id: id.clone(),
exit_code: Some(code),
});
}
}
}
}
}
fn read_proc_environ(pid: u32) -> Result<Vec<(String, String)>, std::io::Error> {
let path = format!("/proc/{pid}/environ");
let data = std::fs::read(&path)?;
let mut env = Vec::new();
for entry in data.split(|&b| b == 0) {
if entry.is_empty() {
continue;
}
let s = String::from_utf8_lossy(entry);
if let Some((key, value)) = s.split_once('=') {
env.push((key.to_string(), value.to_string()));
}
}
env.sort_by(|a, b| a.0.cmp(&b.0));
Ok(env)
}
fn get_process_tree_rss(pid: u32) -> Option<u64> {
let mut total_rss: u64 = 0;
let mut stack = vec![pid];
let page_size = crate::sys::page_size();
while let Some(p) = stack.pop() {
if let Ok(stat) = std::fs::read_to_string(format!("/proc/{p}/stat")) {
if let Some(after_comm) = stat.rfind(')') {
let fields: Vec<&str> = stat[after_comm + 2..].split_whitespace().collect();
if let Some(rss_pages) = fields.get(21).and_then(|s| s.parse::<u64>().ok()) {
total_rss += rss_pages * page_size;
}
}
}
let task_path = format!("/proc/{p}/task");
if let Ok(tasks) = std::fs::read_dir(&task_path) {
for task in tasks.flatten() {
let children_path = task.path().join("children");
if let Ok(children) = std::fs::read_to_string(&children_path) {
for child_pid in children.split_whitespace() {
if let Ok(cpid) = child_pid.parse::<u32>() {
stack.push(cpid);
}
}
}
}
}
}
if total_rss > 0 { Some(total_rss) } else { None }
}
pub async fn is_server_running(socket_path: &Path) -> bool {
UnixStream::connect(socket_path).await.is_ok()
}
#[cfg(unix)]
struct UmaskGuard(nix::sys::stat::Mode);
#[cfg(unix)]
impl UmaskGuard {
fn new(mask: libc::mode_t) -> Self {
let mode = nix::sys::stat::Mode::from_bits_truncate(mask);
Self(nix::sys::stat::umask(mode))
}
}
#[cfg(unix)]
impl Drop for UmaskGuard {
fn drop(&mut self) {
nix::sys::stat::umask(self.0);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::{PASTE_END, PASTE_START};
#[test]
fn wraps_text_in_paste_markers() {
let out = wrap_bracketed_paste("line one\nline two");
assert!(out.starts_with(PASTE_START));
assert!(out.ends_with(PASTE_END));
let inner = &out[PASTE_START.len()..out.len() - PASTE_END.len()];
assert_eq!(inner, b"line one\nline two");
}
#[test]
fn empty_text_still_produces_a_well_formed_envelope() {
let out = wrap_bracketed_paste("");
assert_eq!(out, [PASTE_START, PASTE_END].concat());
}
#[test]
fn strips_embedded_paste_terminator() {
let text = "before\x1b[201~after";
let out = wrap_bracketed_paste(text);
let inner = &out[PASTE_START.len()..out.len() - PASTE_END.len()];
assert_eq!(inner, b"beforeafter");
let count = out
.windows(PASTE_END.len())
.filter(|w| *w == PASTE_END)
.count();
assert_eq!(count, 1, "payload must contain exactly one terminator");
}
#[test]
fn strips_repeated_and_adjacent_terminators() {
let out = wrap_bracketed_paste("a\x1b[201~\x1b[201~b");
let inner = &out[PASTE_START.len()..out.len() - PASTE_END.len()];
assert_eq!(inner, b"ab");
}
#[test]
fn leaves_the_paste_introducer_alone() {
let out = wrap_bracketed_paste("a\x1b[200~b");
let inner = &out[PASTE_START.len()..out.len() - PASTE_END.len()];
assert_eq!(inner, b"a\x1b[200~b");
}
}