1use serde::Serialize;
26use serde_json::{Map, Value};
27use std::fs::OpenOptions;
28use std::io::Write;
29use std::path::{Path, PathBuf};
30
31pub const MAX_EVENT_PAYLOAD_BYTES: usize = 500;
35
36pub const ROTATE_AT_BYTES: u64 = 8 * 1024 * 1024;
40
41#[derive(Debug, thiserror::Error)]
45pub enum EmitError {
46 #[error("event io error: {0}")]
47 Io(#[from] std::io::Error),
48 #[error("event payload was not a JSON object")]
49 NotAnObject,
50}
51
52#[derive(Debug, Clone)]
56pub struct EventEmitter {
57 path: PathBuf,
58 source: String,
59}
60
61impl EventEmitter {
62 pub fn new(path: impl Into<PathBuf>, source: impl Into<String>) -> Self {
66 EventEmitter {
67 path: path.into(),
68 source: source.into(),
69 }
70 }
71
72 pub fn emit<P: Serialize>(&self, kind: &str, payload: &P) -> Result<(), EmitError> {
81 let value = serde_json::to_value(payload).map_err(|_| EmitError::NotAnObject)?;
82 let obj = match value {
83 Value::Object(m) => m,
84 Value::Null => Map::new(),
85 _ => return Err(EmitError::NotAnObject),
86 };
87
88 let payload_len = serde_json::to_string(&obj).map(|s| s.len()).unwrap_or(0);
92 if payload_len > MAX_EVENT_PAYLOAD_BYTES {
93 let mut meta = Map::new();
94 meta.insert("intended_kind".into(), Value::String(kind.to_string()));
95 meta.insert("size".into(), Value::Number(payload_len.into()));
96 return self.write_line("event_payload_too_large", meta);
97 }
98
99 self.write_line(kind, obj)
100 }
101
102 pub fn emit_fields(&self, kind: &str, fields: Map<String, Value>) -> Result<(), EmitError> {
105 let payload_len = serde_json::to_string(&fields).map(|s| s.len()).unwrap_or(0);
106 if payload_len > MAX_EVENT_PAYLOAD_BYTES {
107 let mut meta = Map::new();
108 meta.insert("intended_kind".into(), Value::String(kind.to_string()));
109 meta.insert("size".into(), Value::Number(payload_len.into()));
110 return self.write_line("event_payload_too_large", meta);
111 }
112 self.write_line(kind, fields)
113 }
114
115 fn write_line(&self, event_type: &str, payload: Map<String, Value>) -> Result<(), EmitError> {
116 let mut obj = Map::new();
121 obj.insert("ts".into(), Value::String(now_rfc3339()));
122 obj.insert("type".into(), Value::String(event_type.to_string()));
123 obj.insert("source".into(), Value::String(self.source.clone()));
124 obj.insert("data".into(), Value::Object(payload));
125 let mut line = serde_json::to_string(&Value::Object(obj))
126 .map_err(|e| EmitError::Io(std::io::Error::new(std::io::ErrorKind::InvalidData, e)))?;
127 line.push('\n');
128
129 self.maybe_rotate()?;
130 if let Some(parent) = self.path.parent() {
131 std::fs::create_dir_all(parent)?;
132 }
133 let mut f = OpenOptions::new()
136 .create(true)
137 .append(true)
138 .open(&self.path)?;
139 f.write_all(line.as_bytes())?;
140 Ok(())
141 }
142
143 fn maybe_rotate(&self) -> Result<(), EmitError> {
148 let size = match std::fs::metadata(&self.path) {
149 Ok(m) => m.len(),
150 Err(_) => return Ok(()), };
152 if size <= ROTATE_AT_BYTES {
153 return Ok(());
154 }
155 let rotated = rotated_path(&self.path);
156 let _ = std::fs::rename(&self.path, rotated);
159 Ok(())
160 }
161
162 pub fn path(&self) -> &Path {
164 &self.path
165 }
166}
167
168fn rotated_path(path: &Path) -> PathBuf {
169 let mut s = path.as_os_str().to_os_string();
170 s.push(".1");
171 PathBuf::from(s)
172}
173
174pub(crate) fn now_rfc3339() -> String {
179 use std::time::{SystemTime, UNIX_EPOCH};
180 let dur = SystemTime::now()
181 .duration_since(UNIX_EPOCH)
182 .unwrap_or_default();
183 let secs = dur.as_secs();
184 let millis = dur.subsec_millis();
185 let (year, month, day, hour, min, sec) = civil_from_unix(secs);
186 format!("{year:04}-{month:02}-{day:02}T{hour:02}:{min:02}:{sec:02}.{millis:03}Z")
187}
188
189fn civil_from_unix(secs: u64) -> (i64, u32, u32, u32, u32, u32) {
192 let days = (secs / 86_400) as i64;
193 let rem = secs % 86_400;
194 let hour = (rem / 3600) as u32;
195 let min = ((rem % 3600) / 60) as u32;
196 let sec = (rem % 60) as u32;
197
198 let z = days + 719_468;
200 let era = if z >= 0 { z } else { z - 146_096 } / 146_097;
201 let doe = z - era * 146_097; let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365; let y = yoe + era * 400;
204 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 };
209 (year, m, d, hour, min, sec)
210}
211
212#[cfg(test)]
213mod tests {
214 use super::*;
215 use serde_json::json;
216
217 fn temp_events_path(tag: &str) -> PathBuf {
218 let mut p = std::env::temp_dir();
219 p.push(format!(
220 "fno-agents-events-test-{}-{}-{}.jsonl",
221 tag,
222 std::process::id(),
223 std::time::SystemTime::now()
225 .duration_since(std::time::UNIX_EPOCH)
226 .unwrap()
227 .as_nanos()
228 ));
229 p
230 }
231
232 fn read_lines(path: &Path) -> Vec<Value> {
233 std::fs::read_to_string(path)
234 .unwrap_or_default()
235 .lines()
236 .map(|l| serde_json::from_str::<Value>(l).expect("each line is valid json"))
237 .collect()
238 }
239
240 #[test]
241 fn emits_line_with_ts_type_source_data() {
242 let path = temp_events_path("basic");
243 let em = EventEmitter::new(&path, "daemon");
244 em.emit("daemon_started", &json!({"pid": 4242, "version": "0.1.0"}))
245 .unwrap();
246
247 let lines = read_lines(&path);
248 assert_eq!(lines.len(), 1);
249 let l = &lines[0];
250 assert_eq!(l["type"], "daemon_started");
251 assert_eq!(l["source"], "daemon");
252 assert_eq!(l["data"]["pid"], 4242);
253 assert!(l.get("kind").is_none(), "no legacy kind field");
254 assert!(l["ts"].as_str().unwrap().ends_with('Z'));
255 assert!(l["ts"].as_str().unwrap().starts_with("20"));
256 std::fs::remove_file(&path).ok();
257 }
258
259 #[test]
260 fn oversized_payload_becomes_meta_event_not_silence() {
261 let path = temp_events_path("oversize");
262 let em = EventEmitter::new(&path, "daemon");
263 let huge = "x".repeat(2000);
264 em.emit("agent_spawned", &json!({"blob": huge})).unwrap();
265
266 let lines = read_lines(&path);
267 assert_eq!(lines.len(), 1, "exactly one line: the meta-event");
268 let l = &lines[0];
269 assert_eq!(l["type"], "event_payload_too_large");
270 assert_eq!(l["data"]["intended_kind"], "agent_spawned");
271 assert!(l["data"]["size"].as_u64().unwrap() > MAX_EVENT_PAYLOAD_BYTES as u64);
272 std::fs::remove_file(&path).ok();
273 }
274
275 #[test]
276 fn appends_preserve_fifo_order() {
277 let path = temp_events_path("fifo");
278 let em = EventEmitter::new(&path, "daemon");
279 for i in 0..10 {
280 em.emit("tick", &json!({"seq": i})).unwrap();
281 }
282 let lines = read_lines(&path);
283 let seqs: Vec<u64> = lines
284 .iter()
285 .map(|l| l["data"]["seq"].as_u64().unwrap())
286 .collect();
287 assert_eq!(seqs, (0..10).collect::<Vec<_>>());
288 std::fs::remove_file(&path).ok();
289 }
290
291 #[test]
292 fn null_payload_is_allowed_as_empty_object() {
293 let path = temp_events_path("null");
294 let em = EventEmitter::new(&path, "worker:wkA");
295 em.emit("heartbeat", &Value::Null).unwrap();
296 let lines = read_lines(&path);
297 assert_eq!(lines[0]["type"], "heartbeat");
298 assert_eq!(lines[0]["source"], "worker:wkA");
299 assert_eq!(lines[0]["data"], json!({}));
300 std::fs::remove_file(&path).ok();
301 }
302
303 #[test]
304 fn civil_date_matches_known_epoch_points() {
305 assert_eq!(civil_from_unix(0), (1970, 1, 1, 0, 0, 0));
307 assert_eq!(civil_from_unix(1_700_000_000), (2023, 11, 14, 22, 13, 20));
309 }
310}