use std::path::PathBuf;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::Duration;
use vessel::protocol::{AgentState, AttachEndReason};
use vessel::runtime;
use vessel::{Client, Request, Response, Server};
static TEST_COUNTER: AtomicU32 = AtomicU32::new(0);
fn unique_socket_path() -> PathBuf {
let id = TEST_COUNTER.fetch_add(1, Ordering::SeqCst);
let pid = std::process::id();
PathBuf::from(format!("/tmp/vessel-test-{pid}-{id}.sock"))
}
struct SocketCleanup(PathBuf);
impl Drop for SocketCleanup {
fn drop(&mut self) {
std::fs::remove_file(&self.0).ok();
}
}
vessel::async_test! {
async fn test_server_ping_pong() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client.request(Request::Ping)
.await
.expect("request failed");
assert!(matches!(response, Response::Pong));
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_spawn_and_list() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Spawn {
cmd: vec!["sleep".into(), "10".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, pid } => {
assert!(pid > 0);
id
}
other => panic!("expected Spawned, got {:?}", other),
};
let response = client.request(Request::List { labels: vec![] }).await.expect("list failed");
match response {
Response::Agents { agents } => {
assert_eq!(agents.len(), 1);
assert_eq!(agents[0].id, agent_id);
assert_eq!(agents[0].command, vec!["sleep", "10"]);
assert_eq!(agents[0].cwd, None);
}
other => panic!("expected Agents, got {:?}", other),
}
let response = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 15,
proc_filter: None,
})
.await
.expect("kill failed");
assert!(matches!(response, Response::Ok));
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_spawn_send_snapshot() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Spawn {
cmd: vec!["bash".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(200)).await;
let response = client
.request(Request::Send {
id: Some(agent_id.clone()),
labels: Vec::new(),
all: false,
proc_filter: None,
data: "echo VESSEL_TEST_OUTPUT".into(),
newline: true,
enter: false,
submit_delay_ms: None,
paste: false,
})
.await
.expect("send failed");
assert!(matches!(response, Response::Ok));
runtime::time::sleep(Duration::from_millis(300)).await;
let response = client
.request(Request::Snapshot {
id: agent_id.clone(),
strip_colors: true,
})
.await
.expect("snapshot failed");
match response {
Response::Snapshot { content, .. } => {
assert!(
content.contains("VESSEL_TEST_OUTPUT"),
"snapshot should contain our output: {}",
content
);
}
other => panic!("expected Snapshot, got {:?}", other),
}
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_agent_not_found() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Snapshot {
id: "nonexistent-agent".into(),
strip_colors: true,
})
.await
.expect("request failed");
match response {
Response::Error { message } => {
assert!(message.contains("not found"));
}
other => panic!("expected Error, got {:?}", other),
}
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_screen_cursor_movement() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Spawn {
cmd: vec![
"sh".into(),
"-c".into(),
r#"printf "ABC\rX"; sleep 10"#.into(),
],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(200)).await;
let response = client
.request(Request::Snapshot {
id: agent_id.clone(),
strip_colors: true,
})
.await
.expect("snapshot failed");
match response {
Response::Snapshot { content, .. } => {
assert!(
content.contains("XBC"),
"cursor movement should produce XBC: {}",
content
);
}
other => panic!("expected Snapshot, got {:?}", other),
}
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_transcript_tail() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Spawn {
cmd: vec![
"sh".into(),
"-c".into(),
"echo LINE_ONE; echo LINE_TWO; sleep 10".into(),
],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(200)).await;
let response = client
.request(Request::Tail {
id: agent_id.clone(),
lines: 10,
follow: false,
})
.await
.expect("tail failed");
match response {
Response::Output { data, .. } => {
let text = String::from_utf8_lossy(&data);
assert!(text.contains("LINE_ONE"), "should contain LINE_ONE: {}", text);
assert!(text.contains("LINE_TWO"), "should contain LINE_TWO: {}", text);
}
other => panic!("expected Output, got {:?}", other),
}
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_attach_and_detach() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path.clone());
let response = client
.request(Request::Spawn {
cmd: vec!["bash".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(100)).await;
let mut stream = runtime::net::UnixStream::connect(&socket_path)
.await
.expect("connect failed");
let attach_req = Request::Attach {
id: agent_id.clone(),
readonly: false,
};
let mut json = serde_json::to_string(&attach_req).unwrap();
json.push('\n');
use runtime::io::AsyncWriteExt;
stream.write_all(json.as_bytes()).await.expect("write failed");
use runtime::io::AsyncBufReadExt;
let mut reader = runtime::io::BufReader::new(&mut stream);
let mut line = String::new();
reader.read_line(&mut line).await.expect("read failed");
let response: Response = serde_json::from_str(&line).expect("parse failed");
match response {
Response::AttachStarted { id, size } => {
assert_eq!(id, agent_id);
assert_eq!(size, (24, 80));
}
other => panic!("expected AttachStarted, got {:?}", other),
}
drop(reader);
drop(stream);
runtime::time::sleep(Duration::from_millis(100)).await;
let response = client.request(Request::List { labels: vec![] }).await.expect("list failed");
match response {
Response::Agents { agents } => {
assert_eq!(agents.len(), 1);
assert_eq!(agents[0].id, agent_id);
}
other => panic!("expected Agents, got {:?}", other),
}
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_attach_readonly_mode() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path.clone());
let response = client
.request(Request::Spawn {
cmd: vec!["sh".into(), "-c".into(), "echo HELLO; sleep 10".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(200)).await;
let mut stream = runtime::net::UnixStream::connect(&socket_path)
.await
.expect("connect failed");
let attach_req = Request::Attach {
id: agent_id.clone(),
readonly: true,
};
let mut json = serde_json::to_string(&attach_req).unwrap();
json.push('\n');
use runtime::io::AsyncWriteExt;
use runtime::io::AsyncReadExt;
stream.write_all(json.as_bytes()).await.expect("write failed");
let mut buf = vec![0u8; 65536];
let n = stream.read(&mut buf).await.expect("read failed");
let newline_pos = buf[..n].iter().position(|&b| b == b'\n').expect("no newline");
let response: Response = serde_json::from_slice(&buf[..newline_pos]).expect("parse failed");
assert!(matches!(response, Response::AttachStarted { .. }));
drop(stream);
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_attach_nonexistent_agent() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut stream = runtime::net::UnixStream::connect(&socket_path)
.await
.expect("connect failed");
let attach_req = Request::Attach {
id: "nonexistent-agent".into(),
readonly: false,
};
let mut json = serde_json::to_string(&attach_req).unwrap();
json.push('\n');
use runtime::io::AsyncWriteExt;
stream.write_all(json.as_bytes()).await.expect("write failed");
use runtime::io::AsyncBufReadExt;
let mut reader = runtime::io::BufReader::new(&mut stream);
let mut line = String::new();
reader.read_line(&mut line).await.expect("read failed");
let response: Response = serde_json::from_str(&line).expect("parse failed");
match response {
Response::Error { message } => {
assert!(message.contains("not found"));
}
other => panic!("expected Error, got {:?}", other),
}
let mut client = Client::new(socket_path);
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_attach_receives_output() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path.clone());
let response = client
.request(Request::Spawn {
cmd: vec!["bash".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
runtime::time::sleep(Duration::from_millis(100)).await;
let mut stream = runtime::net::UnixStream::connect(&socket_path)
.await
.expect("connect failed");
let attach_req = Request::Attach {
id: agent_id.clone(),
readonly: false,
};
let mut json = serde_json::to_string(&attach_req).unwrap();
json.push('\n');
use runtime::io::AsyncWriteExt;
use runtime::io::AsyncReadExt;
stream.write_all(json.as_bytes()).await.expect("write failed");
let mut buf = vec![0u8; 4096];
let n = stream.read(&mut buf).await.expect("read failed");
let line_end = buf[..n].iter().position(|&b| b == b'\n').unwrap_or(n);
let response: Response = serde_json::from_slice(&buf[..line_end]).expect("parse failed");
assert!(matches!(response, Response::AttachStarted { .. }));
let cmd = b"echo ATTACH_TEST_OUTPUT\n";
stream.write_all(cmd).await.expect("write failed");
let mut output = Vec::new();
let deadline = runtime::time::Instant::now() + Duration::from_secs(2);
while runtime::time::Instant::now() < deadline {
runtime::time::sleep(Duration::from_millis(50)).await;
let mut chunk = vec![0u8; 4096];
match runtime::time::timeout(Duration::from_millis(50), stream.read(&mut chunk)).await {
Ok(Ok(n)) if n > 0 => {
output.extend_from_slice(&chunk[..n]);
let text = String::from_utf8_lossy(&output);
if text.contains("ATTACH_TEST_OUTPUT") {
break;
}
}
_ => {}
}
}
let text = String::from_utf8_lossy(&output);
assert!(
text.contains("ATTACH_TEST_OUTPUT"),
"should receive command output through attach: {}",
text
);
drop(stream);
let _ = client
.request(Request::Kill {
id: Some(agent_id),
labels: vec![],
all: false,
signal: 9,
proc_filter: None,
})
.await;
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_attach_agent_exit() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path.clone());
let response = client
.request(Request::Spawn {
cmd: vec!["sh".into(), "-c".into(), "sleep 0.5; exit 42".into()],
rows: 24,
cols: 80,
name: None,
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
let agent_id = match response {
Response::Spawned { id, .. } => id,
other => panic!("expected Spawned, got {:?}", other),
};
let mut stream = runtime::net::UnixStream::connect(&socket_path)
.await
.expect("connect failed");
let attach_req = Request::Attach {
id: agent_id.clone(),
readonly: false,
};
let mut json = serde_json::to_string(&attach_req).unwrap();
json.push('\n');
use runtime::io::AsyncWriteExt;
use runtime::io::AsyncReadExt;
stream.write_all(json.as_bytes()).await.expect("write failed");
let mut buf = vec![0u8; 4096];
let n = stream.read(&mut buf).await.expect("read failed");
let line_end = buf[..n].iter().position(|&b| b == b'\n').unwrap_or(n);
let response: Response = serde_json::from_slice(&buf[..line_end]).expect("parse failed");
assert!(matches!(response, Response::AttachStarted { .. }));
let mut received_end = false;
let deadline = runtime::time::Instant::now() + Duration::from_secs(3);
while runtime::time::Instant::now() < deadline && !received_end {
match runtime::time::timeout(Duration::from_millis(100), stream.read(&mut buf)).await {
Ok(Ok(n)) if n > 0 => {
if buf[0] == b'{' {
if let Ok(response) = serde_json::from_slice::<Response>(&buf[..n]) {
if let Response::AttachEnded { reason } = response {
match reason {
AttachEndReason::AgentExited { exit_code } => {
assert_eq!(exit_code, Some(42));
received_end = true;
}
other => panic!("expected AgentExited, got {:?}", other),
}
}
}
}
}
Ok(Ok(0)) => break, _ => {}
}
}
assert!(received_end, "should receive AttachEnded when agent exits");
drop(stream);
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_kill_all() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
for i in 0..3 {
let response = client
.request(Request::Spawn {
cmd: vec!["sleep".into(), "10".into()],
rows: 24,
cols: 80,
name: Some(format!("agent-{i}")),
labels: vec![],
timeout: None,
max_output: None,
env: vec![],
cwd: None,
no_resize: false,
record: false,
memory_limit: None,
})
.await
.expect("spawn failed");
assert!(matches!(response, Response::Spawned { .. }));
}
let response = client.request(Request::List { labels: vec![] }).await.expect("list failed");
match &response {
Response::Agents { agents } => {
assert_eq!(agents.len(), 3, "should have 3 agents");
}
other => panic!("expected Agents, got {:?}", other),
}
let response = client
.request(Request::Kill {
id: None,
labels: vec![],
all: true,
signal: 9,
proc_filter: None,
})
.await
.expect("kill --all failed");
assert!(matches!(response, Response::Ok));
runtime::time::sleep(Duration::from_millis(200)).await;
let response = client.request(Request::List { labels: vec![] }).await.expect("list failed");
match response {
Response::Agents { agents } => {
let running: Vec<_> = agents.iter().filter(|a| a.state == AgentState::Running).collect();
assert!(running.is_empty(), "no agents should be running after kill --all, got: {:?}", running);
}
other => panic!("expected Agents, got {:?}", other),
}
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_kill_all_no_agents() {
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut client = Client::new(socket_path);
let response = client
.request(Request::Kill {
id: None,
labels: vec![],
all: true,
signal: 9,
proc_filter: None,
})
.await
.expect("request failed");
match response {
Response::Error { message } => {
assert!(message.contains("no running agents"), "should say no running agents: {}", message);
}
other => panic!("expected Error, got {:?}", other),
}
let _ = client.request(Request::Shutdown).await;
}
}
vessel::async_test! {
async fn test_oversized_frame_rejected_without_dispatch() {
use vessel::runtime::io::{AsyncReadExt, AsyncWriteExt};
use vessel::runtime::net::UnixStream;
let socket_path = unique_socket_path();
let _cleanup = SocketCleanup(socket_path.clone());
let server_socket = socket_path.clone();
runtime::task::spawn(async move {
let mut server = Server::new(server_socket);
let _ = server.run().await;
});
runtime::time::sleep(Duration::from_millis(100)).await;
let mut stream = UnixStream::connect(&socket_path).await.expect("connect");
let blob = vec![b'x'; 2 * 1024 * 1024];
let _ = stream.write_all(&blob).await;
let mut resp = Vec::new();
let _ = stream.read_to_end(&mut resp).await;
if !resp.is_empty() {
let text = String::from_utf8_lossy(&resp);
assert!(
text.contains("exceeds maximum size"),
"unexpected response to oversized frame: {text}"
);
}
drop(stream);
let mut client = Client::new(socket_path);
let response = client
.request(Request::Ping)
.await
.expect("server should still respond after rejecting oversized frame");
assert!(matches!(response, Response::Pong));
let _ = client.request(Request::Shutdown).await;
}
}