use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Mutex;
use uuid::Uuid;
use super::types::{MessageEnqueueCallback, QueuedUserMessage};
#[derive(Debug, Clone)]
pub struct CmdResult {
pub success: bool,
pub code: i32,
pub output: String,
}
#[derive(Debug, Clone)]
pub struct RunningTask {
pub label: String,
pub started: std::time::Instant,
}
pub struct BackgroundTaskManager {
enqueue: MessageEnqueueCallback,
running: Mutex<HashMap<Uuid, Vec<RunningTask>>>,
}
impl BackgroundTaskManager {
pub fn new(enqueue: MessageEnqueueCallback) -> Self {
Self {
enqueue,
running: Mutex::new(HashMap::new()),
}
}
pub fn running_for(&self, session_id: Uuid) -> usize {
self.running
.lock()
.map(|m| m.get(&session_id).map(Vec::len).unwrap_or(0))
.unwrap_or(0)
}
pub fn running_tasks(&self, session_id: Uuid) -> Vec<RunningTask> {
self.running
.lock()
.map(|m| m.get(&session_id).cloned().unwrap_or_default())
.unwrap_or_default()
}
fn mark_started(&self, session_id: Uuid, label: &str) {
if let Ok(mut m) = self.running.lock() {
m.entry(session_id).or_default().push(RunningTask {
label: label.to_string(),
started: std::time::Instant::now(),
});
}
}
fn mark_finished(&self, session_id: Uuid, label: &str) {
if let Ok(mut m) = self.running.lock()
&& let Some(tasks) = m.get_mut(&session_id)
{
if let Some(pos) = tasks.iter().position(|t| t.label == label) {
tasks.remove(pos);
}
if tasks.is_empty() {
m.remove(&session_id);
}
}
}
pub fn spawn_command(
self: std::sync::Arc<Self>,
session_id: Uuid,
cwd: PathBuf,
label: String,
command: String,
) {
self.mark_started(session_id, &label);
let this = std::sync::Arc::clone(&self);
let task_id = Uuid::new_v4();
tokio::spawn(async move {
tracing::info!(
target: "background_task",
"Background task '{label}' started for session {session_id} \
(id={task_id}, cwd={})",
cwd.display()
);
if let Some(repo) = task_repo() {
let cwd_str = cwd.to_string_lossy().to_string();
if let Err(e) = repo
.record(task_id, session_id, &label, &command, &cwd_str)
.await
{
tracing::error!(
target: "background_task",
"Failed to persist background task '{label}': {e:#}"
);
}
}
let started = std::time::Instant::now();
let result = run_detached(&command, &cwd).await;
tracing::info!(
target: "background_task",
"Background task '{label}' for session {session_id} finished \
(success={}, exit={}, elapsed={:.1}s)",
result.success,
result.code,
started.elapsed().as_secs_f32()
);
let msg = completion_message(&label, &command, &result);
if let Some(repo) = task_repo()
&& let Err(e) = repo.clear(task_id).await
{
tracing::error!(
target: "background_task",
"Failed to clear background task '{label}' after completion: {e:#}"
);
}
this.mark_finished(session_id, &label);
super::session_routes::resolve_route(session_id, &this.enqueue)(session_id, msg);
});
}
}
pub(super) fn task_repo() -> Option<crate::db::BackgroundTaskRepository> {
crate::db::global_pool().map(|p| crate::db::BackgroundTaskRepository::new(p.clone()))
}
async fn run_detached(command: &str, cwd: &std::path::Path) -> CmdResult {
use tokio::process::Command;
let output = Command::new("sh")
.arg("-c")
.arg(command)
.current_dir(cwd)
.output()
.await;
match output {
Ok(out) => {
let mut combined = String::from_utf8_lossy(&out.stdout).into_owned();
let err = String::from_utf8_lossy(&out.stderr);
if !err.trim().is_empty() {
if !combined.is_empty() {
combined.push('\n');
}
combined.push_str(&err);
}
CmdResult {
success: out.status.success(),
code: out.status.code().unwrap_or(-1),
output: combined,
}
}
Err(e) => {
tracing::error!(
target: "background_task",
"Background command could not be launched in {}: {e}",
cwd.display()
);
CmdResult {
success: false,
code: -1,
output: format!("failed to launch: {e}"),
}
}
}
}
pub(crate) fn short_label(command: &str) -> String {
let after_cd = crate::utils::command_label::command_label(command);
let label: String = after_cd.chars().take(60).collect();
if after_cd.chars().count() > 60 {
format!("{label}…")
} else {
label
}
}
pub(crate) fn tail_lines(text: &str, n: usize) -> String {
let lines: Vec<&str> = text.lines().collect();
let start = lines.len().saturating_sub(n);
lines[start..].join("\n")
}
pub(crate) fn completion_message(
label: &str,
command: &str,
result: &CmdResult,
) -> QueuedUserMessage {
let status = if result.success {
"exit 0 (success)".to_string()
} else {
format!("exit {} (failure)", result.code)
};
let tail = tail_lines(&result.output, 50);
let context = format!(
"[System: the background task you started has finished.\n\
Task: {label}\n\
Command: {command}\n\
Status: {status}\n\
Output (last 50 lines):\n{tail}\n\n\
Report the result to the user and continue anything that was waiting on it. \
Do not re-run the command — this IS its result.]"
);
let display = format!(
"🔧 background task {}: {label}",
if result.success { "finished" } else { "failed" }
);
QueuedUserMessage::system(context, display)
}