use serde::Serialize;
use serde_json::{Map, Value};
use std::fs::OpenOptions;
use std::io::Write;
use std::path::{Path, PathBuf};
pub const MAX_EVENT_PAYLOAD_BYTES: usize = 500;
pub const ROTATE_AT_BYTES: u64 = 8 * 1024 * 1024;
#[derive(Debug, thiserror::Error)]
pub enum EmitError {
#[error("event io error: {0}")]
Io(#[from] std::io::Error),
#[error("event payload was not a JSON object")]
NotAnObject,
}
#[derive(Debug, Clone)]
pub struct EventEmitter {
path: PathBuf,
source: String,
}
impl EventEmitter {
pub fn new(path: impl Into<PathBuf>, source: impl Into<String>) -> Self {
EventEmitter {
path: path.into(),
source: source.into(),
}
}
pub fn emit<P: Serialize>(&self, kind: &str, payload: &P) -> Result<(), EmitError> {
let value = serde_json::to_value(payload).map_err(|_| EmitError::NotAnObject)?;
let obj = match value {
Value::Object(m) => m,
Value::Null => Map::new(),
_ => return Err(EmitError::NotAnObject),
};
let payload_len = serde_json::to_string(&obj).map(|s| s.len()).unwrap_or(0);
if payload_len > MAX_EVENT_PAYLOAD_BYTES {
let mut meta = Map::new();
meta.insert("intended_kind".into(), Value::String(kind.to_string()));
meta.insert("size".into(), Value::Number(payload_len.into()));
return self.write_line("event_payload_too_large", meta);
}
self.write_line(kind, obj)
}
pub fn emit_fields(&self, kind: &str, fields: Map<String, Value>) -> Result<(), EmitError> {
let payload_len = serde_json::to_string(&fields).map(|s| s.len()).unwrap_or(0);
if payload_len > MAX_EVENT_PAYLOAD_BYTES {
let mut meta = Map::new();
meta.insert("intended_kind".into(), Value::String(kind.to_string()));
meta.insert("size".into(), Value::Number(payload_len.into()));
return self.write_line("event_payload_too_large", meta);
}
self.write_line(kind, fields)
}
fn write_line(&self, event_type: &str, payload: Map<String, Value>) -> Result<(), EmitError> {
let mut obj = Map::new();
obj.insert("ts".into(), Value::String(now_rfc3339()));
obj.insert("type".into(), Value::String(event_type.to_string()));
obj.insert("source".into(), Value::String(self.source.clone()));
obj.insert("data".into(), Value::Object(payload));
let mut line = serde_json::to_string(&Value::Object(obj))
.map_err(|e| EmitError::Io(std::io::Error::new(std::io::ErrorKind::InvalidData, e)))?;
line.push('\n');
self.maybe_rotate()?;
if let Some(parent) = self.path.parent() {
std::fs::create_dir_all(parent)?;
}
let mut f = OpenOptions::new()
.create(true)
.append(true)
.open(&self.path)?;
f.write_all(line.as_bytes())?;
Ok(())
}
fn maybe_rotate(&self) -> Result<(), EmitError> {
let size = match std::fs::metadata(&self.path) {
Ok(m) => m.len(),
Err(_) => return Ok(()), };
if size <= ROTATE_AT_BYTES {
return Ok(());
}
let rotated = rotated_path(&self.path);
let _ = std::fs::rename(&self.path, rotated);
Ok(())
}
pub fn path(&self) -> &Path {
&self.path
}
}
fn rotated_path(path: &Path) -> PathBuf {
let mut s = path.as_os_str().to_os_string();
s.push(".1");
PathBuf::from(s)
}
pub(crate) fn now_rfc3339() -> String {
use std::time::{SystemTime, UNIX_EPOCH};
let dur = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default();
let secs = dur.as_secs();
let millis = dur.subsec_millis();
let (year, month, day, hour, min, sec) = civil_from_unix(secs);
format!("{year:04}-{month:02}-{day:02}T{hour:02}:{min:02}:{sec:02}.{millis:03}Z")
}
fn civil_from_unix(secs: u64) -> (i64, u32, u32, u32, u32, u32) {
let days = (secs / 86_400) as i64;
let rem = secs % 86_400;
let hour = (rem / 3600) as u32;
let min = ((rem % 3600) / 60) as u32;
let sec = (rem % 60) as u32;
let z = days + 719_468;
let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
let doe = z - era * 146_097; let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); let mp = (5 * doy + 2) / 153; let d = (doy - (153 * mp + 2) / 5 + 1) as u32; let m = if mp < 10 { mp + 3 } else { mp - 9 } as u32; let year = if m <= 2 { y + 1 } else { y };
(year, m, d, hour, min, sec)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn temp_events_path(tag: &str) -> PathBuf {
let mut p = std::env::temp_dir();
p.push(format!(
"fno-agents-events-test-{}-{}-{}.jsonl",
tag,
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
p
}
fn read_lines(path: &Path) -> Vec<Value> {
std::fs::read_to_string(path)
.unwrap_or_default()
.lines()
.map(|l| serde_json::from_str::<Value>(l).expect("each line is valid json"))
.collect()
}
#[test]
fn emits_line_with_ts_type_source_data() {
let path = temp_events_path("basic");
let em = EventEmitter::new(&path, "daemon");
em.emit("daemon_started", &json!({"pid": 4242, "version": "0.1.0"}))
.unwrap();
let lines = read_lines(&path);
assert_eq!(lines.len(), 1);
let l = &lines[0];
assert_eq!(l["type"], "daemon_started");
assert_eq!(l["source"], "daemon");
assert_eq!(l["data"]["pid"], 4242);
assert!(l.get("kind").is_none(), "no legacy kind field");
assert!(l["ts"].as_str().unwrap().ends_with('Z'));
assert!(l["ts"].as_str().unwrap().starts_with("20"));
std::fs::remove_file(&path).ok();
}
#[test]
fn oversized_payload_becomes_meta_event_not_silence() {
let path = temp_events_path("oversize");
let em = EventEmitter::new(&path, "daemon");
let huge = "x".repeat(2000);
em.emit("agent_spawned", &json!({"blob": huge})).unwrap();
let lines = read_lines(&path);
assert_eq!(lines.len(), 1, "exactly one line: the meta-event");
let l = &lines[0];
assert_eq!(l["type"], "event_payload_too_large");
assert_eq!(l["data"]["intended_kind"], "agent_spawned");
assert!(l["data"]["size"].as_u64().unwrap() > MAX_EVENT_PAYLOAD_BYTES as u64);
std::fs::remove_file(&path).ok();
}
#[test]
fn appends_preserve_fifo_order() {
let path = temp_events_path("fifo");
let em = EventEmitter::new(&path, "daemon");
for i in 0..10 {
em.emit("tick", &json!({"seq": i})).unwrap();
}
let lines = read_lines(&path);
let seqs: Vec<u64> = lines
.iter()
.map(|l| l["data"]["seq"].as_u64().unwrap())
.collect();
assert_eq!(seqs, (0..10).collect::<Vec<_>>());
std::fs::remove_file(&path).ok();
}
#[test]
fn null_payload_is_allowed_as_empty_object() {
let path = temp_events_path("null");
let em = EventEmitter::new(&path, "worker:wkA");
em.emit("heartbeat", &Value::Null).unwrap();
let lines = read_lines(&path);
assert_eq!(lines[0]["type"], "heartbeat");
assert_eq!(lines[0]["source"], "worker:wkA");
assert_eq!(lines[0]["data"], json!({}));
std::fs::remove_file(&path).ok();
}
#[test]
fn civil_date_matches_known_epoch_points() {
assert_eq!(civil_from_unix(0), (1970, 1, 1, 0, 0, 0));
assert_eq!(civil_from_unix(1_700_000_000), (2023, 11, 14, 22, 13, 20));
}
}