use std::io::Write;
use anyhow::{Context, Result, bail};
use crate::config::Host;
use crate::transport::Transport;
use crate::wrapper::{Job, JobId, dispatch_script, new_id};
pub fn dispatch(
transport: &dyn Transport,
host: &Host,
cmd: &str,
cwd: Option<&str>,
max_secs: Option<u64>,
) -> Result<JobId> {
dispatch_with_warnings(transport, host, cmd, cwd, max_secs, &mut std::io::stderr())
}
#[doc(hidden)]
pub fn dispatch_with_warnings(
transport: &dyn Transport,
host: &Host,
cmd: &str,
cwd: Option<&str>,
max_secs: Option<u64>,
warnings: &mut dyn Write,
) -> Result<JobId> {
if !transport.master_alive(host) {
return Err(crate::errors::no_master(host));
}
let job = Job {
id: new_id().parse::<JobId>().expect("generated ids are valid"),
cmd: cmd.to_owned(),
cwd: cwd.map(str::to_owned),
max_secs: max_secs.unwrap_or(host.max_job_secs),
};
let script = format!(
"{}; {} && {{ tmux -L {} list-sessions -F '#{{session_name}}' 2>/dev/null | grep -c '^coop-' || true; }}",
crate::jobs::prune(host),
dispatch_script(host, &job),
host.tmux_socket
);
let output = transport.run(host, &script)?;
if output.code != 0 {
bail!(
"dispatch failed for {}: {}",
host.name,
output.stderr.trim()
);
}
let running: u32 = output
.text()
.trim()
.parse()
.context("invalid running-session count in dispatch reply")?;
if running > host.max_running {
writeln!(
warnings,
"coop: dispatched {}; {running} now running on {}, cap {}",
job.id, host.name, host.max_running
)?;
}
Ok(job.id)
}