use serde_json::{Map, Value};
use std::collections::VecDeque;
use std::io::Write;
use std::sync::Mutex;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Level {
Trace,
Debug,
Info,
Warn,
Error,
}
impl Level {
pub fn as_str(self) -> &'static str {
match self {
Level::Trace => "trace",
Level::Debug => "debug",
Level::Info => "info",
Level::Warn => "warn",
Level::Error => "error",
}
}
pub fn parse(s: &str) -> Option<Level> {
match s.to_ascii_lowercase().as_str() {
"trace" => Some(Level::Trace),
"debug" => Some(Level::Debug),
"info" => Some(Level::Info),
"warn" | "warning" => Some(Level::Warn),
"error" => Some(Level::Error),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy)]
pub enum Comp {
Supervisor,
Agent,
Mcp,
Intel,
}
impl Comp {
fn as_str(self) -> &'static str {
match self {
Comp::Supervisor => "supervisor",
Comp::Agent => "agent",
Comp::Mcp => "mcp",
Comp::Intel => "intel",
}
}
}
#[derive(Debug, Clone)]
pub struct LogCtx {
pub run_id: String,
pub agent_id: String,
pub agent_path: String,
pub comp: Comp,
pub pid: u32,
pub trace_id: Option<String>,
}
#[derive(Clone)]
pub struct Logger {
ctx: LogCtx,
min: Level,
log_content: bool,
}
static STDERR_LOCK: Mutex<()> = Mutex::new(());
pub const EVENTS_SCHEMA: &str = "1.0";
pub const EVENTS_RING_DEFAULT: usize = 1024;
struct RingEntry {
seq: u64,
level: &'static str,
event: String,
line: Value,
}
struct EventRing {
buf: VecDeque<RingEntry>,
cap: usize,
dropped: u64,
}
impl EventRing {
fn new(cap: usize) -> EventRing {
let cap = cap.max(1);
EventRing {
buf: VecDeque::with_capacity(cap),
cap,
dropped: 0,
}
}
fn push(&mut self, entry: RingEntry) {
if self.buf.len() == self.cap {
self.buf.pop_front();
self.dropped = self.dropped.saturating_add(1);
}
self.buf.push_back(entry);
}
}
static EVENT_RING: Mutex<Option<EventRing>> = Mutex::new(None);
static RING_INSTALLED: AtomicU64 = AtomicU64::new(0);
static RING_SEQ: AtomicU64 = AtomicU64::new(0);
static EVENTS_DIRTY: AtomicU64 = AtomicU64::new(0);
pub fn take_events_dirty() -> bool {
EVENTS_DIRTY.swap(0, Ordering::Relaxed) != 0
}
pub fn install_event_ring(cap: usize) {
let mut g = EVENT_RING.lock().unwrap_or_else(|e| e.into_inner());
*g = Some(EventRing::new(cap));
RING_INSTALLED.store(1, Ordering::Relaxed);
}
pub const EVENT_FAMILIES: &[&str] = &[
"a2a",
"agent",
"audit",
"breaker",
"budget",
"cgroup",
"child",
"config",
"context",
"drain",
"goal",
"human",
"inbox",
"instance",
"instruction",
"intel",
"interface",
"knowledge",
"lifecycle",
"limit",
"loop",
"mcp",
"message",
"otel",
"plan",
"preflight",
"pressure",
"proc",
"prompt",
"registry",
"restore",
"run",
"skill",
"skills",
"start",
"step",
"store",
"stream",
"subagent",
"test",
"timer",
"tool",
"turn",
"wait",
"webhook",
"workflow",
];
pub struct TappedEvent {
pub stream: String,
pub subject: String,
pub data: Value,
}
const SAMPLE_EVERY: u64 = 16;
struct RuntimeTap {
stream: String,
all: Vec<String>,
sampled: Vec<String>,
queue: VecDeque<TappedEvent>,
cap: usize,
dropped: u64,
seen: u64,
}
static RUNTIME_TAP: Mutex<Option<RuntimeTap>> = Mutex::new(None);
static TAP_ARMED: AtomicU64 = AtomicU64::new(0);
pub fn install_runtime_tap(stream: &str, all: Vec<String>, sampled: Vec<String>, cap: usize) {
let mut g = RUNTIME_TAP.lock().unwrap_or_else(|e| e.into_inner());
*g = Some(RuntimeTap {
stream: stream.to_string(),
all,
sampled,
queue: VecDeque::new(),
cap: cap.max(1),
dropped: 0,
seen: 0,
});
TAP_ARMED.store(1, Ordering::Relaxed);
}
pub fn runtime_tap_armed() -> bool {
TAP_ARMED.load(Ordering::Relaxed) != 0
}
pub fn tap_direct(stream: &str, subject: &str, data: Value) {
if !runtime_tap_armed() {
return;
}
let mut g = RUNTIME_TAP.lock().unwrap_or_else(|e| e.into_inner());
if let Some(t) = g.as_mut() {
push_bounded(t, stream.to_string(), subject.to_string(), data);
}
}
fn push_bounded(t: &mut RuntimeTap, stream: String, subject: String, data: Value) {
if t.queue.len() >= t.cap {
t.dropped = t.dropped.saturating_add(1);
return;
}
t.queue.push_back(TappedEvent {
stream,
subject,
data,
});
}
static TAP_DRAINING: AtomicU64 = AtomicU64::new(0);
pub fn tap_drain_guard() -> TapDrainGuard {
TAP_DRAINING.store(1, Ordering::Relaxed);
TapDrainGuard
}
pub struct TapDrainGuard;
impl Drop for TapDrainGuard {
fn drop(&mut self) {
TAP_DRAINING.store(0, Ordering::Relaxed);
}
}
fn capture_to_tap(event: &str, line: &Value) {
if !runtime_tap_armed() || TAP_DRAINING.load(Ordering::Relaxed) != 0 {
return;
}
let family = event.split('.').next().unwrap_or(event);
let mut g = RUNTIME_TAP.lock().unwrap_or_else(|e| e.into_inner());
let Some(t) = g.as_mut() else { return };
let full = t.all.iter().any(|f| f == family);
let sampled = !full && t.sampled.iter().any(|f| f == family);
if !full && !sampled {
return;
}
if sampled {
t.seen = t.seen.wrapping_add(1);
if t.seen % SAMPLE_EVERY != 0 {
return;
}
}
let stream = t.stream.clone();
push_bounded(t, stream, event.to_string(), line.clone());
}
pub fn drain_runtime_tap() -> (Vec<TappedEvent>, u64) {
if !runtime_tap_armed() {
return (Vec::new(), 0);
}
let mut g = RUNTIME_TAP.lock().unwrap_or_else(|e| e.into_inner());
let Some(t) = g.as_mut() else {
return (Vec::new(), 0);
};
let dropped = std::mem::take(&mut t.dropped);
(t.queue.drain(..).collect(), dropped)
}
pub struct EventWindow {
pub events: Vec<Value>,
pub oldest_seq: u64,
pub newest_seq: u64,
pub dropped: u64,
}
pub fn read_event_window(
after: u64,
limit: usize,
level: Option<&str>,
event_prefixes: &[&str],
) -> Option<EventWindow> {
if RING_INSTALLED.load(Ordering::Relaxed) == 0 {
return None;
}
let g = EVENT_RING.lock().unwrap_or_else(|e| e.into_inner());
let ring = g.as_ref()?;
let oldest_seq = ring.buf.front().map(|e| e.seq).unwrap_or(0);
let newest_seq = ring.buf.back().map(|e| e.seq).unwrap_or(0);
let mut events = Vec::new();
for entry in ring.buf.iter() {
if entry.seq <= after {
continue;
}
if let Some(want) = level
&& entry.level != want
{
continue;
}
if !event_prefixes.is_empty() && !event_prefixes.iter().any(|p| entry.event.starts_with(p))
{
continue;
}
let mut line = match &entry.line {
Value::Object(m) => m.clone(),
_ => Map::new(),
};
line.insert("seq".into(), Value::Number(entry.seq.into()));
events.push(Value::Object(line));
if events.len() >= limit {
break;
}
}
Some(EventWindow {
events,
oldest_seq,
newest_seq,
dropped: ring.dropped,
})
}
fn capture_to_ring(level: &'static str, event: &str, line: &Value) {
if RING_INSTALLED.load(Ordering::Relaxed) == 0 {
return;
}
let seq = RING_SEQ.fetch_add(1, Ordering::Relaxed) + 1;
let mut g = EVENT_RING.lock().unwrap_or_else(|e| e.into_inner());
if let Some(ring) = g.as_mut() {
ring.push(RingEntry {
seq,
level,
event: event.to_string(),
line: line.clone(),
});
EVENTS_DIRTY.store(1, Ordering::Relaxed);
}
}
impl Logger {
pub fn new(ctx: LogCtx, min: Level) -> Self {
Logger {
ctx,
min,
log_content: false,
}
}
pub fn with_content(mut self, on: bool) -> Self {
self.log_content = on;
self
}
pub fn content_capture(&self) -> bool {
self.log_content
}
pub fn ctx(&self) -> &LogCtx {
&self.ctx
}
pub fn event(&self, level: Level, event: &str, fields: Value) {
if level < self.min {
return;
}
let mut m = Map::new();
m.insert(
"ts".into(),
Value::String(rfc3339_millis(SystemTime::now())),
);
m.insert("level".into(), Value::String(level.as_str().into()));
m.insert("event".into(), Value::String(event.into()));
m.insert("run_id".into(), Value::String(self.ctx.run_id.clone()));
m.insert("agent_id".into(), Value::String(self.ctx.agent_id.clone()));
m.insert(
"agent_path".into(),
Value::String(self.ctx.agent_path.clone()),
);
m.insert("comp".into(), Value::String(self.ctx.comp.as_str().into()));
m.insert("pid".into(), Value::Number(self.ctx.pid.into()));
if let Some(tid) = &self.ctx.trace_id {
m.insert("trace_id".into(), Value::String(tid.clone()));
}
if let Value::Object(extra) = fields {
for (k, v) in extra {
m.insert(k, v);
}
}
let value = Value::Object(m);
capture_to_ring(level.as_str(), event, &value);
capture_to_tap(event, &value);
crate::obs::otel::capture_log(
crate::obs::otel::now_unix_nanos(),
level.as_str(),
event,
&value,
);
let mut line = serde_json::to_vec(&value).unwrap_or_else(|_| b"{}".to_vec());
line.push(b'\n');
let _guard = STDERR_LOCK.lock().unwrap_or_else(|e| e.into_inner());
let _ = std::io::stderr().write_all(&line);
}
pub fn info(&self, event: &str, fields: Value) {
self.event(Level::Info, event, fields);
}
pub fn warn(&self, event: &str, fields: Value) {
self.event(Level::Warn, event, fields);
}
pub fn error(&self, event: &str, fields: Value) {
self.event(Level::Error, event, fields);
}
pub fn debug(&self, event: &str, fields: Value) {
self.event(Level::Debug, event, fields);
}
}
pub fn rfc3339_millis(t: SystemTime) -> String {
let dur = t.duration_since(UNIX_EPOCH).unwrap_or_default();
let secs = dur.as_secs() as i64;
let millis = dur.subsec_millis();
let days = secs.div_euclid(86_400);
let secs_of_day = secs.rem_euclid(86_400);
let (y, m, d) = civil_from_days(days);
let hh = secs_of_day / 3600;
let mm = (secs_of_day % 3600) / 60;
let ss = secs_of_day % 60;
format!("{y:04}-{m:02}-{d:02}T{hh:02}:{mm:02}:{ss:02}.{millis:03}Z")
}
pub(crate) fn civil_from_days(z: i64) -> (i64, i64, i64) {
let z = z + 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; let m = if mp < 10 { mp + 3 } else { mp - 9 }; (if m <= 2 { y + 1 } else { y }, m, d)
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
#[test]
fn rfc3339_known_timestamps() {
assert_eq!(rfc3339_millis(UNIX_EPOCH), "1970-01-01T00:00:00.000Z");
let t = UNIX_EPOCH + Duration::from_secs(1_700_000_000);
assert_eq!(rfc3339_millis(t), "2023-11-14T22:13:20.000Z");
let t = UNIX_EPOCH + Duration::from_millis(1_700_000_000_123);
assert_eq!(rfc3339_millis(t), "2023-11-14T22:13:20.123Z");
}
#[test]
fn leap_year_day() {
let t = UNIX_EPOCH + Duration::from_secs(19_782 * 86_400);
assert_eq!(&rfc3339_millis(t)[..10], "2024-02-29");
}
#[test]
fn level_ordering_filters() {
assert!(Level::Debug < Level::Info);
assert!(Level::Error > Level::Warn);
}
fn ring_entry(seq: u64, level: &'static str, event: &str) -> RingEntry {
RingEntry {
seq,
level,
event: event.to_string(),
line: serde_json::json!({"event": event, "level": level}),
}
}
#[test]
fn ring_evicts_oldest_and_counts_dropped() {
let mut r = EventRing::new(2);
r.push(ring_entry(1, "info", "loop.step"));
r.push(ring_entry(2, "info", "loop.step"));
assert_eq!(r.dropped, 0);
r.push(ring_entry(3, "warn", "limit.exceeded"));
assert_eq!(r.dropped, 1);
let seqs: Vec<u64> = r.buf.iter().map(|e| e.seq).collect();
assert_eq!(seqs, vec![2, 3]);
}
#[test]
fn ring_zero_cap_clamps_to_one() {
let mut r = EventRing::new(0);
r.push(ring_entry(1, "info", "a"));
r.push(ring_entry(2, "info", "b"));
assert_eq!(r.buf.len(), 1);
assert_eq!(r.buf.back().unwrap().seq, 2);
assert_eq!(r.dropped, 1);
}
#[test]
fn install_then_read_window_with_cursor_and_filters() {
install_event_ring(64);
let base = RING_SEQ.load(Ordering::Relaxed); let log = Logger::new(
LogCtx {
run_id: "r".into(),
agent_id: "0".into(),
agent_path: "0".into(),
comp: Comp::Supervisor,
pid: 1,
trace_id: None,
},
Level::Trace,
);
log.info("loop.step", serde_json::json!({"step": 1}));
log.warn("limit.exceeded", serde_json::json!({"limit": "steps"}));
log.info("subagent.spawn", serde_json::json!({"node": 1}));
let w = read_event_window(base, 100, None, &[]).expect("ring installed");
assert!(w.events.len() >= 3);
assert!(w.events.iter().all(|e| e.get("seq").is_some()));
assert!(w.newest_seq >= w.oldest_seq);
let w = read_event_window(base, 100, Some("warn"), &[]).expect("ring");
assert!(w.events.iter().all(|e| e["level"] == "warn"));
assert!(w.events.iter().any(|e| e["event"] == "limit.exceeded"));
let w = read_event_window(base, 100, None, &["subagent."]).expect("ring");
assert!(
w.events
.iter()
.all(|e| e["event"].as_str().unwrap().starts_with("subagent."))
);
let w = read_event_window(base, 1, None, &[]).expect("ring");
assert_eq!(w.events.len(), 1);
}
#[test]
fn families_cover_the_emitted_vocabulary() {
fn walk(dir: &std::path::Path, out: &mut Vec<std::path::PathBuf>) {
for e in std::fs::read_dir(dir).into_iter().flatten().flatten() {
let p = e.path();
if p.is_dir() {
walk(&p, out);
} else if p.extension().is_some_and(|x| x == "rs") {
out.push(p);
}
}
}
let mut files = Vec::new();
walk(
&std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("src"),
&mut files,
);
let mut missing: Vec<String> = Vec::new();
for f in files {
let src = std::fs::read_to_string(&f).unwrap_or_default();
for m in ["\n.info(\"", ".warn(\"", ".error(\"", ".debug(\""] {
let needle = m.trim_start_matches('\n');
let mut rest = src.as_str();
while let Some(i) = rest.find(needle) {
rest = &rest[i + needle.len()..];
let Some(end) = rest.find('"') else { break };
let name = &rest[..end];
if name.is_empty()
|| !name.chars().all(|c| {
c.is_ascii_lowercase() || c.is_ascii_digit() || c == '.' || c == '_'
})
{
continue;
}
let family = name.split('.').next().unwrap_or(name);
if !EVENT_FAMILIES.contains(&family) && !missing.contains(&family.to_string()) {
missing.push(family.to_string());
}
}
}
}
assert!(
missing.is_empty(),
"these event families are emitted but missing from EVENT_FAMILIES: {missing:?}"
);
}
#[test]
fn the_family_list_is_sorted_and_unique() {
let mut sorted = EVENT_FAMILIES.to_vec();
sorted.sort_unstable();
sorted.dedup();
assert_eq!(sorted.as_slice(), EVENT_FAMILIES);
}
}