use core::fmt;
use std::io::{self, Write};
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::sync::Mutex;
use std::sync::mpsc::{Receiver, Sender, SyncSender, TryRecvError, channel, sync_channel};
use std::thread::JoinHandle;
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use keel_core_api::{ErrorClass, ErrorCode};
use serde::{Deserialize, Serialize};
use tracing::warn;
pub const EVENTS_VERSION: u32 = 1;
pub const EVENTS_SUBDIR: &str = "events";
pub const EVENTS_EXT: &str = "ndjson";
const FLUSH_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Event {
pub v: u32,
pub seq: u64,
pub ms: u64,
#[serde(flatten)]
pub kind: EventKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CacheStore {
Memory,
Persistent,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "event", rename_all = "snake_case")]
pub enum EventKind {
RunStart {
run: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
wall_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pid: Option<u32>,
},
CallStart {
call: String,
target: String,
op: String,
},
CacheHit {
call: String,
target: String,
scope: CacheStore,
},
CacheMiss {
call: String,
target: String,
scope: CacheStore,
},
Throttle {
call: String,
target: String,
wait_ms: u64,
},
BreakerReject { call: String, target: String },
BreakerHalfOpen { call: String, target: String },
BreakerOpen {
call: String,
target: String,
cooldown_ms: u64,
},
BreakerClose { call: String, target: String },
AttemptStart {
call: String,
target: String,
attempt: u32,
},
AttemptError {
call: String,
target: String,
attempt: u32,
class: ErrorClass,
#[serde(default, skip_serializing_if = "Option::is_none")]
http_status: Option<u16>,
},
Backoff {
call: String,
target: String,
attempt: u32,
wait_ms: u64,
},
CallEnd {
call: String,
target: String,
result: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
code: Option<ErrorCode>,
attempts: u32,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TraceRef {
pub run: String,
pub seq: u64,
}
impl TraceRef {
#[must_use]
pub fn file_name(&self) -> String {
format!("{}.{EVENTS_EXT}", self.run)
}
}
impl fmt::Display for TraceRef {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{}#{}", self.run, self.seq)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ParseTraceRefError;
impl fmt::Display for ParseTraceRefError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("trace ref must look like <run>#<seq>, e.g. 019a2b3c4d5-7f2e#12")
}
}
impl std::error::Error for ParseTraceRefError {}
impl FromStr for TraceRef {
type Err = ParseTraceRefError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let (run, seq) = s.rsplit_once('#').ok_or(ParseTraceRefError)?;
if run.is_empty() {
return Err(ParseTraceRefError);
}
let seq = seq.parse().map_err(|_| ParseTraceRefError)?;
Ok(Self {
run: run.to_owned(),
seq,
})
}
}
#[derive(Debug, Clone)]
pub struct EventsEnv {
pub keel_events: Option<String>,
pub base_dir: PathBuf,
}
impl EventsEnv {
#[must_use]
pub fn capture() -> Self {
Self {
keel_events: std::env::var("KEEL_EVENTS").ok(),
base_dir: std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")),
}
}
}
#[must_use]
pub fn resolve_events_dir(env: &EventsEnv) -> Option<PathBuf> {
let keel_dir = env.base_dir.join(".keel");
match env.keel_events.as_deref().map(str::trim) {
Some(v)
if v.is_empty()
|| v.eq_ignore_ascii_case("0")
|| v.eq_ignore_ascii_case("false")
|| v.eq_ignore_ascii_case("off") =>
{
None
}
Some(_) => Some(keel_dir.join(EVENTS_SUBDIR)),
None => keel_dir.is_dir().then(|| keel_dir.join(EVENTS_SUBDIR)),
}
}
enum Msg {
Event(Event),
Flush(SyncSender<()>),
Shutdown,
}
#[derive(Debug)]
struct Emitter {
seq: u64,
tx: Sender<Msg>,
}
#[derive(Debug)]
pub struct EventSink {
emitter: Mutex<Emitter>,
run_id: String,
path: Option<PathBuf>,
writer: Option<JoinHandle<()>>,
}
impl EventSink {
#[must_use]
pub fn from_env() -> Option<Self> {
let dir = resolve_events_dir(&EventsEnv::capture())?;
match Self::open(&dir) {
Ok(sink) => Some(sink),
Err(error) => {
warn!(dir = %dir.display(), error = %error, "event sink unavailable; live events disabled");
None
}
}
}
pub fn open(dir: &Path) -> io::Result<Self> {
std::fs::create_dir_all(dir)?;
let mut collision: io::Error = io::ErrorKind::AlreadyExists.into();
for _ in 0..3 {
let run_id = new_run_id();
let path = dir.join(format!("{run_id}.{EVENTS_EXT}"));
match std::fs::File::create_new(&path) {
Ok(file) => {
return Self::start(
Box::new(io::BufWriter::new(file)),
run_id,
Some(path),
Some(epoch_ms()),
Some(std::process::id()),
);
}
Err(e) if e.kind() == io::ErrorKind::AlreadyExists => collision = e,
Err(e) => return Err(e),
}
}
Err(collision)
}
pub fn to_writer(writer: Box<dyn Write + Send>, run_id: &str) -> io::Result<Self> {
Self::start(writer, run_id.to_owned(), None, None, None)
}
fn start(
writer: Box<dyn Write + Send>,
run_id: String,
path: Option<PathBuf>,
wall_ms: Option<u64>,
pid: Option<u32>,
) -> io::Result<Self> {
let (tx, rx) = channel();
let handle = std::thread::Builder::new()
.name("keel-events".to_owned())
.spawn(move || write_events(&rx, writer))?;
let sink = Self {
emitter: Mutex::new(Emitter { seq: 0, tx }),
run_id: run_id.clone(),
path,
writer: Some(handle),
};
sink.emit(
0,
EventKind::RunStart {
run: run_id,
wall_ms,
pid,
},
);
Ok(sink)
}
#[must_use]
pub fn run_id(&self) -> &str {
&self.run_id
}
#[must_use]
pub fn path(&self) -> Option<&Path> {
self.path.as_deref()
}
pub fn emit(&self, ms: u64, kind: EventKind) -> u64 {
let mut emitter = self.emitter.lock().expect("event sink lock poisoned");
let seq = emitter.seq;
emitter.seq += 1;
let _ = emitter.tx.send(Msg::Event(Event {
v: EVENTS_VERSION,
seq,
ms,
kind,
}));
seq
}
pub fn flush(&self) {
let (ack_tx, ack_rx) = sync_channel(1);
{
let emitter = self.emitter.lock().expect("event sink lock poisoned");
if emitter.tx.send(Msg::Flush(ack_tx)).is_err() {
return; }
}
let _ = ack_rx.recv_timeout(FLUSH_TIMEOUT);
}
}
impl Drop for EventSink {
fn drop(&mut self) {
if let Ok(emitter) = self.emitter.get_mut() {
let _ = emitter.tx.send(Msg::Shutdown);
}
if let Some(handle) = self.writer.take() {
let _ = handle.join();
}
}
}
fn write_events(rx: &Receiver<Msg>, mut out: Box<dyn Write + Send>) {
let mut dirty = false;
loop {
let msg = if dirty {
match rx.try_recv() {
Ok(msg) => msg,
Err(TryRecvError::Empty) => {
let _ = out.flush();
dirty = false;
continue;
}
Err(TryRecvError::Disconnected) => break,
}
} else {
match rx.recv() {
Ok(msg) => msg,
Err(_) => break,
}
};
match msg {
Msg::Event(event) => {
if serde_json::to_writer(&mut out, &event).is_ok() && out.write_all(b"\n").is_ok() {
dirty = true;
}
}
Msg::Flush(ack) => {
let _ = out.flush();
dirty = false;
let _ = ack.send(());
}
Msg::Shutdown => break,
}
}
let _ = out.flush();
}
fn new_run_id() -> String {
format!("{:011x}-{:04x}", epoch_ms(), fastrand::u16(..))
}
fn epoch_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |d| u64::try_from(d.as_millis()).unwrap_or(u64::MAX))
}
#[cfg(test)]
mod tests {
use super::{
EVENTS_SUBDIR, Event, EventKind, EventSink, EventsEnv, ParseTraceRefError, TraceRef,
resolve_events_dir,
};
use std::io::Write;
use std::path::PathBuf;
use std::sync::{Arc, Mutex};
#[derive(Debug, Clone, Default)]
struct SharedBuf(Arc<Mutex<Vec<u8>>>);
impl SharedBuf {
fn contents(&self) -> String {
String::from_utf8(self.0.lock().expect("buf lock").clone()).expect("utf-8 feed")
}
}
impl Write for SharedBuf {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().expect("buf lock").extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
fn env(keel_events: Option<&str>, base_dir: &std::path::Path) -> EventsEnv {
EventsEnv {
keel_events: keel_events.map(str::to_owned),
base_dir: base_dir.to_owned(),
}
}
#[test]
fn activation_decision_table() {
let tmp = tempfile::tempdir().expect("tempdir");
let bare = tmp.path();
assert_eq!(resolve_events_dir(&env(None, bare)), None);
for off in ["0", "false", "off", "FALSE", "Off", "", " "] {
assert_eq!(resolve_events_dir(&env(Some(off), bare)), None, "{off:?}");
}
let expected = bare.join(".keel").join(EVENTS_SUBDIR);
for on in ["1", "true", "on", "yes"] {
assert_eq!(
resolve_events_dir(&env(Some(on), bare)),
Some(expected.clone()),
"{on:?}"
);
}
std::fs::create_dir(bare.join(".keel")).expect("mk .keel");
assert_eq!(resolve_events_dir(&env(None, bare)), Some(expected));
let tmp2 = tempfile::tempdir().expect("tempdir");
std::fs::write(tmp2.path().join(".keel"), b"not a dir").expect("write file");
assert_eq!(resolve_events_dir(&env(None, tmp2.path())), None);
}
#[test]
fn trace_ref_round_trips_and_rejects_junk() {
let r = TraceRef {
run: "019a2b3c4d5-7f2e".to_owned(),
seq: 12,
};
assert_eq!(r.to_string(), "019a2b3c4d5-7f2e#12");
assert_eq!(r.file_name(), "019a2b3c4d5-7f2e.ndjson");
assert_eq!("019a2b3c4d5-7f2e#12".parse::<TraceRef>(), Ok(r));
assert_eq!(
"a#b#3".parse::<TraceRef>(),
Ok(TraceRef {
run: "a#b".to_owned(),
seq: 3
})
);
for bad in ["", "norun", "#7", "run#", "run#x", "run#-1"] {
assert_eq!(bad.parse::<TraceRef>(), Err(ParseTraceRefError), "{bad:?}");
}
}
#[test]
fn sink_writes_header_then_events_in_seq_order_and_drop_flushes() {
let buf = SharedBuf::default();
let sink =
EventSink::to_writer(Box::new(buf.clone()), "run-test").expect("sink must start");
assert_eq!(sink.run_id(), "run-test");
assert_eq!(sink.path(), None);
let seq = sink.emit(
5,
EventKind::CallStart {
call: "t-000001".to_owned(),
target: "api.example.com".to_owned(),
op: "GET api.example.com".to_owned(),
},
);
assert_eq!(seq, 1, "run_start header owns seq 0");
drop(sink);
let lines: Vec<Event> = buf
.contents()
.lines()
.map(|l| serde_json::from_str(l).expect("every line parses"))
.collect();
assert_eq!(lines.len(), 2);
assert_eq!(
lines[0],
Event {
v: 1,
seq: 0,
ms: 0,
kind: EventKind::RunStart {
run: "run-test".to_owned(),
wall_ms: None,
pid: None,
},
}
);
assert_eq!(lines[1].seq, 1);
assert_eq!(lines[1].ms, 5);
assert_eq!(
buf.contents().lines().next().expect("header line"),
r#"{"v":1,"seq":0,"ms":0,"event":"run_start","run":"run-test"}"#
);
}
#[test]
fn flush_makes_the_feed_current_without_dropping_the_sink() {
let buf = SharedBuf::default();
let sink =
EventSink::to_writer(Box::new(buf.clone()), "run-flush").expect("sink must start");
sink.emit(
1,
EventKind::BreakerReject {
call: "t-000001".to_owned(),
target: "api.example.com".to_owned(),
},
);
sink.flush();
let contents = buf.contents();
assert!(
contents.contains(r#""event":"breaker_reject""#),
"flushed feed must contain the event: {contents}"
);
}
#[test]
fn file_backed_sink_writes_run_file_with_wall_header() {
let tmp = tempfile::tempdir().expect("tempdir");
let dir = tmp.path().join(".keel").join(EVENTS_SUBDIR);
let sink = EventSink::open(&dir).expect("open sink");
let run = sink.run_id().to_owned();
let path = sink.path().expect("file-backed").to_owned();
assert_eq!(path, dir.join(format!("{run}.ndjson")));
drop(sink);
let feed = std::fs::read_to_string(&path).expect("feed readable");
let header: Event =
serde_json::from_str(feed.lines().next().expect("header")).expect("header parses");
match header.kind {
EventKind::RunStart {
run: r,
wall_ms,
pid,
} => {
assert_eq!(r, run);
assert!(wall_ms.is_some(), "production header anchors wall time");
assert_eq!(pid, Some(std::process::id()));
}
other => panic!("first line must be run_start, got {other:?}"),
}
}
#[test]
fn events_env_capture_reads_process_state() {
let env = EventsEnv::capture();
assert_ne!(env.base_dir, PathBuf::new(), "cwd captured");
}
}