use std::sync::OnceLock;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
use base64::Engine;
use crate::config::Host;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Job {
pub id: JobId,
pub cmd: String,
pub cwd: Option<String>,
pub max_secs: u64,
}
pub fn state_dir(id: &JobId) -> String {
format!("{JOBS_ROOT}/{id}")
}
pub const JOBS_ROOT: &str = "${XDG_STATE_HOME:-$HOME/.local/state}/coop/jobs";
pub fn dispatch_script(host: &Host, job: &Job) -> String {
let dir = state_dir(&job.id);
let command = encode_command(job.cmd.as_bytes());
let cwd = job
.cwd
.as_deref()
.or(host.default_cwd.as_deref())
.unwrap_or("$HOME");
let cd = match home_relative(cwd) {
Some("") => "cd \"$HOME\"".to_string(),
Some(rest) => format!(
"cd \"$HOME/$(printf %s {} | base64 -d)\"",
encode_command(rest.as_bytes())
),
None => format!(
"cd \"$(printf %s {} | base64 -d)\"",
encode_command(cwd.as_bytes())
),
};
let run = if job.max_secs == 0 {
format!("{cd} && printf %s {command} | base64 -d | sh; echo $? > {dir}/rc")
} else {
let watchdog = format!(
"sleep {secs}; if tmux -L {socket} has-session -t coop-{id} 2>/dev/null; then \
if [ ! -f {dir}/rc ]; then echo 124 > {dir}/rc; fi; \
tmux -L {socket} kill-session -t coop-{id}; fi",
socket = host.tmux_socket,
id = job.id,
secs = job.max_secs,
);
format!(
"tmux -L {socket} -f /dev/null new-session -d -s watch-{id} \
\"printf %s {watchdog} | base64 -d | sh\"; \
{cd} && printf %s {command} | base64 -d | sh; rc=$?; \
tmux -L {socket} kill-session -t watch-{id} 2>/dev/null; \
if [ ! -f {dir}/rc ]; then echo $rc > {dir}/rc; fi",
socket = host.tmux_socket,
id = job.id,
watchdog = encode_command(watchdog.as_bytes()),
)
};
format!(
"mkdir -p {dir} && printf %s {command} | base64 -d > {dir}/cmd && \
tmux -L {} -f /dev/null new-session -d -s coop-{} \
'{{ {run}; }} \
| {{ head -c {} > {dir}/log; cat > {dir}/.overflow; \
if [ -s {dir}/.overflow ]; then echo 1 > {dir}/truncated; fi; \
rm -f {dir}/.overflow; }}'",
host.tmux_socket, job.id, host.max_log_bytes
)
}
fn home_relative(path: &str) -> Option<&str> {
for prefix in ["~", "$HOME", "${HOME}"] {
if let Some(rest) = path.strip_prefix(prefix) {
if rest.is_empty() {
return Some("");
}
if let Some(rest) = rest.strip_prefix('/') {
return Some(rest.trim_end_matches('/'));
}
}
}
None
}
pub const ID_HEX_LEN: usize = 6;
pub fn new_id() -> String {
static NEXT: OnceLock<AtomicU64> = OnceLock::new();
let next = NEXT.get_or_init(|| {
let time = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos() as u64;
AtomicU64::new(time ^ u64::from(std::process::id()))
});
let value = next.fetch_add(0x9e37_79b9_7f4a_7c15, Ordering::Relaxed);
format!("{:06x}", value & 0x00ff_ffff)
}
pub fn encode_command(input: &[u8]) -> String {
base64::engine::general_purpose::STANDARD.encode(input)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct JobId(String);
impl JobId {
pub fn as_str(&self) -> &str {
&self.0
}
}
impl std::str::FromStr for JobId {
type Err = anyhow::Error;
fn from_str(raw: &str) -> Result<Self, Self::Err> {
let ok = raw.len() == ID_HEX_LEN
&& raw
.bytes()
.all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b));
if !ok {
anyhow::bail!(
"invalid job id {raw:?}: expected {ID_HEX_LEN} lowercase hex digits, as printed by `coop run`"
);
}
Ok(Self(raw.to_string()))
}
}
impl std::fmt::Display for JobId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}