use std::io::Write;
use std::time::{Duration, Instant};
use anyhow::{Result, bail};
use crate::config::Host;
use crate::errors::CoopError;
use crate::probe::{From as ProbeFrom, State, next_interval, probe};
use crate::transport::Transport;
use crate::wrapper::{JobId, state_dir};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Selection {
LastBytes,
All,
Lines(u64),
}
pub fn once(
transport: &dyn Transport,
host: &Host,
id: &JobId,
selection: Selection,
out: &mut dyn Write,
) -> Result<()> {
crate::errors::require_master(transport, host)?;
let dir = state_dir(id);
let read = match selection {
Selection::LastBytes => format!("tail -c 65536 {dir}/log"),
Selection::All => format!("cat {dir}/log"),
Selection::Lines(lines) => format!("tail -n {lines} {dir}/log"),
};
let script = format!("{read}; printf '\\037%s' \"$(cat {dir}/truncated 2>/dev/null)\"");
let output = transport.run(host, &script)?;
if output.code != 0 {
bail!("tail failed: {}", output.stderr.trim());
}
let (body, truncated) = match output.stdout.iter().rposition(|&b| b == 0x1f) {
Some(i) => (&output.stdout[..i], output.stdout[i + 1..] == *b"1"),
None => (&output.stdout[..], false),
};
out.write_all(body)?;
if truncated {
eprintln!(
"coop: log was capped at {} bytes; the job ran to completion but \
later output was discarded\n \
the log holds stdout and stderr merged, in the order the job \
wrote them; redirect inside your command to separate them",
host.max_log_bytes
);
}
Ok(())
}
pub fn follow(
transport: &dyn Transport,
host: &Host,
id: &JobId,
from: u64,
out: &mut dyn Write,
) -> Result<i32> {
wait_loop(
transport,
host,
id,
ProbeFrom::Offset(from),
Some(out),
None,
)
}
pub fn follow_deferred(
transport: &dyn Transport,
host: &Host,
id: &JobId,
out: &mut dyn Write,
) -> Result<i32> {
let code = wait_only(transport, host, id, None)?;
once(transport, host, id, Selection::All, out)?;
Ok(code)
}
pub fn wait_only(
transport: &dyn Transport,
host: &Host,
id: &JobId,
timeout: Option<u64>,
) -> Result<i32> {
wait_loop(
transport,
host,
id,
ProbeFrom::StateOnly,
None,
timeout.map(Duration::from_secs),
)
}
fn wait_loop(
transport: &dyn Transport,
host: &Host,
id: &JobId,
mut from: ProbeFrom,
mut out: Option<&mut dyn Write>,
timeout: Option<Duration>,
) -> Result<i32> {
let started = Instant::now();
let mut interval = Duration::from_secs(1);
loop {
let result = probe(transport, host, id, from)
.map_err(|error| error.context(CoopError::Dropped { id: id.to_string() }))?;
let new_bytes = !result.bytes.is_empty();
if let Some(writer) = out.as_deref_mut() {
writer.write_all(&result.bytes)?;
}
if let ProbeFrom::Offset(offset) = from {
from = ProbeFrom::Offset(offset.saturating_add(result.bytes.len() as u64));
}
match result.state {
State::Done(code) => {
if let Some(writer) = out.as_deref_mut()
&& let ProbeFrom::Offset(offset) = from
{
let tail = probe(transport, host, id, ProbeFrom::Offset(offset))?;
writer.write_all(&tail.bytes)?;
}
return Ok(code);
}
State::Orphan => {
return Err(CoopError::Orphan { id: id.to_string() }.into());
}
State::Running => {}
}
if timeout.is_some_and(|limit| started.elapsed() >= limit) {
return Err(CoopError::Timeout { id: id.to_string() }.into());
}
let sleep = timeout
.map(|limit| interval.min(limit.saturating_sub(started.elapsed())))
.unwrap_or(interval);
std::thread::sleep(sleep);
interval = next_interval(interval, new_bytes);
}
}