use std::io::Write;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::Mutex;
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use crate::error::{Error, Result};
use crate::observe::{Flow, Observer, RunEvent, EVENT_NAMES};
const DEFAULT_TIMEOUT_MS: u64 = 5_000;
const POLL: Duration = Duration::from_millis(5);
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum OnFailure {
#[default]
Continue,
Cancel,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct Hook {
#[serde(default)]
on: Vec<String>,
#[serde(default)]
append: Option<PathBuf>,
#[serde(default)]
run: Option<Vec<String>>,
#[serde(default)]
on_failure: OnFailure,
#[serde(default)]
timeout_ms: Option<u64>,
}
impl Hook {
fn check(&self, index: usize, path: &Path) -> Result<()> {
let at = format!("{}: key `hook[{index}]`", path.display());
match (&self.append, &self.run) {
(None, None) => {
return Err(Error::Config(format!(
"{at}: a hook needs an action — set `append` to a path or `run` to an argv"
)))
}
(Some(_), Some(_)) => {
return Err(Error::Config(format!(
"{at}: a hook has one action — set `append` or `run`, not both"
)))
}
_ => {}
}
if self.run.as_ref().is_some_and(Vec::is_empty) {
return Err(Error::Config(format!("{at}: `run` names no program")));
}
for name in &self.on {
if !EVENT_NAMES.contains(&name.as_str()) {
return Err(Error::Config(format!(
"{at}: `{name}` is not an event this crate emits. It emits: {}",
EVENT_NAMES.join(", ")
)));
}
}
Ok(())
}
fn wants(&self, tag: &str) -> bool {
self.on.is_empty() || self.on.iter().any(|n| n == tag)
}
fn fire(&self, dir: &Path, line: &str, lock: &Mutex<()>) -> Result<()> {
if let Some(rel) = &self.append {
let at = dir.join(rel);
let _guard = lock.lock().unwrap_or_else(|e| e.into_inner());
let mut f = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&at)
.map_err(|e| {
Error::Config(format!("hook cannot append to {}: {e}", at.display()))
})?;
writeln!(f, "{line}").map_err(|e| {
Error::Config(format!("hook cannot append to {}: {e}", at.display()))
})?;
return Ok(());
}
let argv = self.run.as_ref().expect("check() proved one action exists");
let limit = Duration::from_millis(self.timeout_ms.unwrap_or(DEFAULT_TIMEOUT_MS));
let mut child = Command::new(&argv[0])
.args(&argv[1..])
.current_dir(dir)
.stdin(Stdio::piped())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.map_err(|e| Error::Config(format!("hook could not run `{}`: {e}", argv[0])))?;
if let Some(mut stdin) = child.stdin.take() {
let _ = writeln!(stdin, "{line}");
}
match wait_bounded(&mut child, limit) {
Some(status) if status.success() => Ok(()),
Some(status) => Err(Error::Config(format!(
"hook `{}` exited with {status}",
argv[0]
))),
None => Err(Error::Config(format!(
"hook `{}` did not finish within {}ms and was killed",
argv[0],
limit.as_millis()
))),
}
}
}
fn wait_bounded(child: &mut Child, limit: Duration) -> Option<std::process::ExitStatus> {
let deadline = Instant::now() + limit;
loop {
match child.try_wait() {
Ok(Some(status)) => return Some(status),
Err(_) => return None,
Ok(None) => {}
}
if Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
return None;
}
std::thread::sleep(POLL);
}
}
#[derive(Debug)]
pub struct Hooks {
hooks: Vec<Hook>,
dir: PathBuf,
lock: Mutex<()>,
}
impl Hooks {
pub(crate) fn new(hooks: Vec<Hook>, dir: impl Into<PathBuf>) -> Self {
let dir = dir.into();
for hook in &hooks {
let Some(rel) = &hook.append else { continue };
let at = dir.join(rel);
if at.exists() {
continue;
}
if let Err(e) = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&at)
{
tracing::warn!("hook cannot create {}: {e}", at.display());
}
}
Self {
hooks,
dir,
lock: Mutex::new(()),
}
}
pub fn is_empty(&self) -> bool {
self.hooks.is_empty()
}
pub(crate) fn check(hooks: &[Hook], path: &Path) -> Result<()> {
for (i, hook) in hooks.iter().enumerate() {
hook.check(i, path)?;
}
Ok(())
}
}
impl Observer for Hooks {
fn event(&self, event: &RunEvent) -> Flow {
let Ok(value) = serde_json::to_value(event) else {
return Flow::Continue;
};
let Some(tag) = value.get("event").and_then(serde_json::Value::as_str) else {
return Flow::Continue;
};
let line = value.to_string();
let mut flow = Flow::Continue;
for (i, hook) in self.hooks.iter().enumerate() {
if !hook.wants(tag) {
continue;
}
if let Err(why) = hook.fire(&self.dir, &line, &self.lock) {
tracing::warn!("hook[{i}] on `{tag}` failed: {why}");
if hook.on_failure == OnFailure::Cancel {
flow = Flow::Cancel;
}
}
}
flow
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::observe::EventKind;
fn hook(toml: &str) -> Hook {
toml::from_str(toml).unwrap()
}
#[test]
fn every_event_the_crate_emits_is_a_name_a_hook_may_use() {
let names = EVENT_NAMES
.iter()
.map(|n| format!("\"{n}\""))
.collect::<Vec<_>>()
.join(", ");
let h = hook(&format!("on = [{names}]\nappend = \"a.jsonl\"\n"));
h.check(0, Path::new("io.local.toml")).unwrap();
}
#[test]
fn an_event_the_crate_does_not_emit_is_refused_naming_it() {
let h = hook("on = [\"finshed\"]\nappend = \"a.jsonl\"\n");
let err = h.check(3, Path::new("io.local.toml")).unwrap_err();
assert!(err.to_string().contains("finshed"), "{err}");
assert!(err.to_string().contains("hook[3]"), "{err}");
}
#[test]
fn a_hook_needs_exactly_one_action() {
let none = hook("on = [\"stalled\"]\n");
assert!(none
.check(0, Path::new("io.local.toml"))
.unwrap_err()
.to_string()
.contains("needs an action"));
let both = hook("append = \"a.jsonl\"\nrun = [\"true\"]\n");
assert!(both
.check(1, Path::new("io.local.toml"))
.unwrap_err()
.to_string()
.contains("not both"));
let empty = hook("run = []\n");
assert!(empty
.check(2, Path::new("io.local.toml"))
.unwrap_err()
.to_string()
.contains("names no program"));
}
#[test]
fn a_hook_with_no_filter_wants_everything() {
let all = hook("append = \"a.jsonl\"\n");
for name in EVENT_NAMES {
assert!(all.wants(name), "{name}");
}
let one = hook("on = [\"stalled\"]\nappend = \"a.jsonl\"\n");
assert!(one.wants("stalled"));
assert!(!one.wants("finished"));
}
#[test]
fn an_append_hook_writes_one_json_line_per_matching_event() {
let dir = tempfile::tempdir().unwrap();
let hooks = Hooks::new(vec![hook("append = \"audit.jsonl\"\n")], dir.path());
hooks.event(&RunEvent::new(1, 1, EventKind::Stalled));
hooks.event(&RunEvent::new(1, 2, EventKind::Replan { window: 3 }));
let log = std::fs::read_to_string(dir.path().join("audit.jsonl")).unwrap();
assert_eq!(log.lines().count(), 2, "{log}");
assert!(log
.lines()
.next()
.unwrap()
.contains("\"event\":\"stalled\""));
assert!(log.lines().nth(1).unwrap().contains("\"event\":\"replan\""));
}
#[test]
fn a_hook_that_matches_nothing_leaves_an_empty_file_rather_than_none() {
let dir = tempfile::tempdir().unwrap();
let hooks = Hooks::new(
vec![hook(
"on = [\"question_asked\"]\nappend = \"audit.jsonl\"\n",
)],
dir.path(),
);
hooks.event(&RunEvent::new(1, 1, EventKind::Stalled));
let at = dir.path().join("audit.jsonl");
assert!(at.exists(), "an installed hook creates its log");
assert_eq!(std::fs::read_to_string(&at).unwrap(), "");
}
}