use std::io::Write;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::error::{Error, Result};
pub const LOG_HEAD: usize = 192 << 10;
pub const LOG_TAIL: usize = 64 << 10;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RunStatus {
Running,
Succeeded,
Failed,
Skipped,
}
impl RunStatus {
pub fn finished(self) -> bool {
self != RunStatus::Running
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RunTrigger {
Schedule,
Missed,
Manual,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Run {
pub id: u64,
pub kind: String,
pub trigger: RunTrigger,
pub by: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub scheduled_for: Option<i64>,
pub status: RunStatus,
pub started_at: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub finished_at: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub duration_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub exit_code: Option<i32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default)]
pub output_bytes: u64,
#[serde(default, skip_serializing_if = "Value::is_null")]
pub detail: Value,
}
impl Run {
pub fn finish(&mut self, status: RunStatus) {
let now = crate::stack::controller::now_ms();
self.status = status;
self.finished_at = Some(now);
self.duration_ms = Some(now.saturating_sub(self.started_at));
}
}
pub struct RunLog {
file: Option<std::fs::File>,
written: usize,
tail: std::collections::VecDeque<u8>,
pub seen: u64,
}
impl RunLog {
fn open(path: &Path) -> Result<RunLog> {
use std::os::unix::fs::OpenOptionsExt;
let file = std::fs::OpenOptions::new()
.create(true)
.append(true)
.mode(0o600)
.open(path)?;
Ok(RunLog {
file: Some(file),
written: 0,
tail: Default::default(),
seen: 0,
})
}
pub fn sink() -> RunLog {
RunLog {
file: None,
written: 0,
tail: Default::default(),
seen: 0,
}
}
pub fn write(&mut self, b: &[u8]) {
self.seen += b.len() as u64;
let head = (LOG_HEAD - self.written.min(LOG_HEAD)).min(b.len());
if head > 0 {
if let Some(f) = &mut self.file {
let _ = f.write_all(&b[..head]);
}
self.written += head;
}
for &c in &b[head..] {
if self.tail.len() == LOG_TAIL {
self.tail.pop_front();
}
self.tail.push_back(c);
}
}
pub fn line(&mut self, l: &str) {
self.write(format!("{l}\n").as_bytes());
}
pub fn close(&mut self) {
let dropped = self.seen - self.written as u64 - self.tail.len() as u64;
if let Some(f) = &mut self.file {
if dropped > 0 {
let _ = writeln!(f, "\n[... {dropped} bytes of output not kept ...]");
}
let (a, b) = self.tail.as_slices();
let _ = f.write_all(a);
let _ = f.write_all(b);
let _ = f.sync_all();
}
self.tail.clear();
self.file = None;
}
}
impl Drop for RunLog {
fn drop(&mut self) {
if self.file.is_some() {
self.close();
}
}
}
pub struct RunStore {
dir: PathBuf,
}
impl RunStore {
pub fn new(dir: PathBuf) -> RunStore {
RunStore { dir }
}
fn json(&self, id: u64) -> PathBuf {
self.dir.join(format!("{id}.json"))
}
fn log_path(&self, id: u64) -> PathBuf {
self.dir.join(format!("{id}.log"))
}
fn ids(&self) -> Vec<u64> {
let mut ids: Vec<u64> = match std::fs::read_dir(&self.dir) {
Ok(rd) => rd
.flatten()
.filter_map(|e| e.file_name().to_str()?.strip_suffix(".json")?.parse().ok())
.collect(),
Err(_) => vec![],
};
ids.sort_unstable();
ids
}
pub fn save(&self, r: &Run) -> Result<()> {
crate::app::write_atomic(&self.json(r.id), &serde_json::to_vec_pretty(r)?)
}
pub fn start(
&self,
kind: &str,
trigger: RunTrigger,
by: &str,
slot: Option<i64>,
keep: usize,
) -> Result<(Run, RunLog)> {
std::fs::create_dir_all(&self.dir)?;
let id = self.ids().last().copied().unwrap_or(0) + 1;
let r = Run {
id,
kind: kind.into(),
trigger,
by: by.into(),
scheduled_for: slot,
status: RunStatus::Running,
started_at: crate::stack::controller::now_ms(),
finished_at: None,
duration_ms: None,
exit_code: None,
error: None,
output_bytes: 0,
detail: Value::Null,
};
self.save(&r)?;
let log = RunLog::open(&self.log_path(id))?;
self.prune(keep);
Ok((r, log))
}
pub fn finish(&self, r: &mut Run, log: &mut RunLog) -> Result<()> {
log.close();
r.output_bytes = log.seen;
self.save(r)
}
pub fn get(&self, id: u64) -> Result<Run> {
match std::fs::read(self.json(id)) {
Ok(b) => Ok(serde_json::from_slice(&b)?),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
Err(Error::NotFound(format!("run {id}")))
}
Err(e) => Err(e.into()),
}
}
pub fn list(&self, limit: usize) -> Vec<Run> {
self.ids()
.into_iter()
.rev()
.take(limit)
.filter_map(|id| self.get(id).ok())
.collect()
}
pub fn last(&self) -> Option<Run> {
self.list(1).into_iter().next()
}
pub fn log(&self, id: u64, offset: u64) -> Result<(String, u64, bool)> {
let r = self.get(id)?;
let b = std::fs::read(self.log_path(id)).unwrap_or_default();
let start = (offset as usize).min(b.len());
Ok((
String::from_utf8_lossy(&b[start..]).into_owned(),
b.len() as u64,
r.status.finished(),
))
}
pub fn prune(&self, keep: usize) {
let ids = self.ids();
let excess = ids.len().saturating_sub(keep.max(1));
for id in ids.into_iter().take(excess) {
let _ = std::fs::remove_file(self.json(id));
let _ = std::fs::remove_file(self.log_path(id));
}
}
pub fn recover(&self) {
for id in self.ids() {
if let Ok(mut r) = self.get(id) {
if r.status == RunStatus::Running {
r.error = Some("interrupted: the daemon stopped during it".into());
r.finish(RunStatus::Failed);
let _ = self.save(&r);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn records_logs_and_retention() {
let d = tempfile::tempdir().unwrap();
let s = RunStore::new(d.path().join("runs"));
for i in 0..5 {
let (mut r, mut log) = s.start("job", RunTrigger::Manual, "t", Some(i), 3).unwrap();
log.line(&format!("run {i}"));
r.exit_code = Some(0);
r.finish(RunStatus::Succeeded);
s.finish(&mut r, &mut log).unwrap();
}
let l = s.list(10);
assert_eq!(l.iter().map(|r| r.id).collect::<Vec<_>>(), [5, 4, 3]);
assert_eq!(l[0].output_bytes, 6);
assert!(l[0].duration_ms.is_some());
let (text, off, done) = s.log(5, 0).unwrap();
assert_eq!((text.as_str(), off, done), ("run 4\n", 6, true));
assert_eq!(s.log(5, 4).unwrap().0, "4\n");
assert!(s.get(1).is_err(), "pruned");
let (r, _log) = s.start("job", RunTrigger::Schedule, "t", None, 3).unwrap();
s.recover();
let r = s.get(r.id).unwrap();
assert_eq!(r.status, RunStatus::Failed);
assert!(r.error.unwrap().contains("interrupted"));
}
#[test]
fn log_keeps_head_and_tail() {
let d = tempfile::tempdir().unwrap();
let s = RunStore::new(d.path().to_path_buf());
let (mut r, mut log) = s.start("job", RunTrigger::Manual, "t", None, 5).unwrap();
let chunk = vec![b'a'; 64 << 10];
for _ in 0..10 {
log.write(&chunk);
}
log.write(b"THE END");
r.finish(RunStatus::Succeeded);
s.finish(&mut r, &mut log).unwrap();
let (text, len, _) = s.log(r.id, 0).unwrap();
assert!(len < (LOG_HEAD + LOG_TAIL + 100) as u64, "{len}");
assert!(text.ends_with("THE END"));
assert!(text.contains("bytes of output not kept"));
assert_eq!(s.get(r.id).unwrap().output_bytes, (640 << 10) + 7);
}
}