use std::time::{SystemTime, UNIX_EPOCH};
pub fn now_unix_nanos() -> u128 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0)
}
pub struct RunSpan {
#[cfg(feature = "otel")]
inner: Option<imp::RunSpan>,
}
pub fn run_begin(trace_id: Option<&str>, start_unix_nanos: u128) -> RunSpan {
#[cfg(feature = "otel")]
{
RunSpan {
inner: imp::RunSpan::begin(trace_id, start_unix_nanos),
}
}
#[cfg(not(feature = "otel"))]
{
let _ = (trace_id, start_unix_nanos);
RunSpan {}
}
}
impl RunSpan {
pub fn record_chat(
&mut self,
model: &str,
input_tokens: u64,
output_tokens: u64,
ok: bool,
start_unix_nanos: u128,
) {
#[cfg(feature = "otel")]
if let Some(i) = self.inner.as_mut() {
i.record_chat(model, input_tokens, output_tokens, ok, start_unix_nanos);
}
#[cfg(not(feature = "otel"))]
let _ = (model, input_tokens, output_tokens, ok, start_unix_nanos);
}
pub fn record_tool(&mut self, tool_name: &str, ok: bool, start_unix_nanos: u128) {
#[cfg(feature = "otel")]
if let Some(i) = self.inner.as_mut() {
i.record_tool(tool_name, ok, start_unix_nanos);
}
#[cfg(not(feature = "otel"))]
let _ = (tool_name, ok, start_unix_nanos);
}
pub fn finish(self, model: &str, input_tokens: u64, output_tokens: u64, ok: bool) {
#[cfg(feature = "otel")]
if let Some(i) = self.inner {
i.finish(model, input_tokens, output_tokens, ok);
}
#[cfg(not(feature = "otel"))]
let _ = (model, input_tokens, output_tokens, ok);
}
}
pub fn capture_log(unix_nanos: u128, level: &str, event: &str, fields: &serde_json::Value) {
#[cfg(feature = "otel")]
imp::capture_log(unix_nanos, level, event, fields);
#[cfg(not(feature = "otel"))]
let _ = (unix_nanos, level, event, fields);
}
pub fn arm_logs(endpoint: &str, service: &str, version: &str) {
#[cfg(feature = "otel")]
imp::arm_logs(endpoint, service, version);
#[cfg(not(feature = "otel"))]
let _ = (endpoint, service, version);
}
#[cfg(feature = "otel")]
mod imp {
use crate::net::http::{self, Url};
use serde_json::{Value, json};
use std::time::Duration;
pub(super) struct Span {
pub trace_id: String,
pub span_id: String,
pub parent_span_id: Option<String>,
pub name: String,
pub start_unix_nanos: u128,
pub end_unix_nanos: u128,
pub ok: bool,
pub attrs: Vec<(&'static str, Value)>,
}
pub(super) struct RunSpan {
trace_id: String,
span_id: String,
start_unix_nanos: u128,
endpoint: String,
children: Vec<Span>,
}
pub(super) fn str_val(s: impl Into<String>) -> Value {
json!({ "stringValue": s.into() })
}
pub(super) fn int_val(n: u64) -> Value {
json!({ "intValue": n.to_string() })
}
impl RunSpan {
pub(super) fn begin(trace_id: Option<&str>, start_unix_nanos: u128) -> Option<RunSpan> {
let endpoint = std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT")
.ok()
.filter(|s| !s.is_empty())?;
let trace_id = trace_id?.to_string();
Some(RunSpan {
span_id: crate::obs::trace::new_span_id(),
trace_id,
start_unix_nanos,
endpoint,
children: Vec::new(),
})
}
pub(super) fn record_chat(
&mut self,
model: &str,
input_tokens: u64,
output_tokens: u64,
ok: bool,
start_unix_nanos: u128,
) {
self.children.push(Span {
trace_id: self.trace_id.clone(),
span_id: crate::obs::trace::new_span_id(),
parent_span_id: Some(self.span_id.clone()),
name: "chat".into(),
start_unix_nanos,
end_unix_nanos: super::now_unix_nanos(),
ok,
attrs: vec![
("gen_ai.operation.name", str_val("chat")),
("gen_ai.request.model", str_val(model)),
("gen_ai.usage.input_tokens", int_val(input_tokens)),
("gen_ai.usage.output_tokens", int_val(output_tokens)),
],
});
}
pub(super) fn record_tool(&mut self, tool_name: &str, ok: bool, start_unix_nanos: u128) {
self.children.push(Span {
trace_id: self.trace_id.clone(),
span_id: crate::obs::trace::new_span_id(),
parent_span_id: Some(self.span_id.clone()),
name: "execute_tool".into(),
start_unix_nanos,
end_unix_nanos: super::now_unix_nanos(),
ok,
attrs: vec![
("gen_ai.operation.name", str_val("execute_tool")),
("gen_ai.tool.name", str_val(tool_name)),
],
});
}
pub(super) fn finish(
mut self,
model: &str,
input_tokens: u64,
output_tokens: u64,
ok: bool,
) {
let run = Span {
trace_id: self.trace_id.clone(),
span_id: self.span_id.clone(),
parent_span_id: None,
name: "invoke_agent".into(),
start_unix_nanos: self.start_unix_nanos,
end_unix_nanos: super::now_unix_nanos(),
ok,
attrs: vec![
("gen_ai.operation.name", str_val("invoke_agent")),
("gen_ai.request.model", str_val(model)),
("gen_ai.usage.input_tokens", int_val(input_tokens)),
("gen_ai.usage.output_tokens", int_val(output_tokens)),
],
};
self.children.push(run);
let _ = export(
&self.endpoint,
&to_otlp_json(&self.children, "agentd", crate::VERSION),
);
}
}
pub(super) fn to_otlp_json(spans: &[Span], service: &str, version: &str) -> Value {
let encoded: Vec<Value> = spans.iter().map(encode_span).collect();
json!({
"resourceSpans": [{
"resource": { "attributes": [
{ "key": "service.name", "value": str_val(service) },
{ "key": "service.version", "value": str_val(version) },
]},
"scopeSpans": [{ "scope": { "name": "agentd" }, "spans": encoded }]
}]
})
}
fn encode_span(s: &Span) -> Value {
let attrs: Vec<Value> = s
.attrs
.iter()
.map(|(k, v)| json!({ "key": k, "value": v }))
.collect();
let mut span = json!({
"traceId": s.trace_id,
"spanId": s.span_id,
"name": s.name,
"kind": 1, "startTimeUnixNano": s.start_unix_nanos.to_string(),
"endTimeUnixNano": s.end_unix_nanos.to_string(),
"status": { "code": if s.ok { 1 } else { 2 } }, "attributes": attrs,
});
if let Some(p) = &s.parent_span_id {
span["parentSpanId"] = json!(p);
}
span
}
fn export(endpoint: &str, body: &Value) -> Result<(), String> {
let base = endpoint.trim_end_matches('/');
let target = if base.ends_with("/v1/traces") {
base.to_string()
} else {
format!("{base}/v1/traces")
};
let url = Url::parse(&target).map_err(|e| format!("otel: bad endpoint '{target}': {e}"))?;
if url.is_tls() {
return Err("otel: https OTLP endpoints need --features tls".into());
}
let bytes = serde_json::to_vec(body).map_err(|e| e.to_string())?;
let mut stream = http::connect_tcp(&url.host, url.port, Duration::from_secs(5))
.map_err(|e| e.to_string())?;
let headers = [("content-type", "application/json")];
let resp = http::send(
&mut stream,
&url.host_header(),
"POST",
&url.path,
&headers,
&bytes,
)
.map_err(|e| e.to_string())?;
if resp.is_success() {
Ok(())
} else {
Err(format!("otel: collector returned HTTP {}", resp.status))
}
}
use std::sync::{Mutex, OnceLock};
static LOGS: OnceLock<Mutex<Vec<Value>>> = OnceLock::new();
pub(super) fn capture_log(unix_nanos: u128, level: &str, event: &str, fields: &Value) {
let Some(buf) = LOGS.get() else { return };
let mut b = buf.lock().unwrap_or_else(|e| e.into_inner());
if b.len() >= 8192 {
return;
}
b.push(json!({
"timeUnixNano": unix_nanos.to_string(),
"severityText": level.to_ascii_uppercase(),
"body": { "stringValue": event },
"attributes": [ { "key": "log.fields", "value": { "stringValue": fields.to_string() } } ],
}));
}
pub(super) fn arm_logs(endpoint: &str, service: &str, version: &str) {
if LOGS.set(Mutex::new(Vec::new())).is_err() {
return; }
let (endpoint, service, version) = (
endpoint.to_string(),
service.to_string(),
version.to_string(),
);
std::thread::Builder::new()
.name("otel-logs".into())
.spawn(move || {
loop {
std::thread::sleep(Duration::from_secs(5));
let batch = {
let mut b = LOGS
.get()
.expect("armed")
.lock()
.unwrap_or_else(|e| e.into_inner());
std::mem::take(&mut *b)
};
if batch.is_empty() {
continue;
}
let body = to_otlp_logs_json(&batch, &service, &version);
let _ = export_signal(&endpoint, "logs", &body);
}
})
.ok();
}
pub(super) fn to_otlp_logs_json(records: &[Value], service: &str, version: &str) -> Value {
json!({
"resourceLogs": [{
"resource": { "attributes": [
{ "key": "service.name", "value": { "stringValue": service } },
{ "key": "service.version", "value": { "stringValue": version } },
]},
"scopeLogs": [{ "scope": { "name": "agentd" }, "logRecords": records }]
}]
})
}
fn export_signal(endpoint: &str, signal: &str, body: &Value) -> Result<(), String> {
let base = endpoint.trim_end_matches('/');
let suffix = format!("/v1/{signal}");
let target = if base.ends_with(&suffix) {
base.to_string()
} else {
format!("{base}{suffix}")
};
let url = Url::parse(&target).map_err(|e| format!("otel: bad endpoint '{target}': {e}"))?;
if url.is_tls() {
return Err("otel: https OTLP endpoints need --features tls".into());
}
let bytes = serde_json::to_vec(body).map_err(|e| e.to_string())?;
let mut stream = http::connect_tcp(&url.host, url.port, Duration::from_secs(5))
.map_err(|e| e.to_string())?;
let headers = [("content-type", "application/json")];
let resp = http::send(
&mut stream,
&url.host_header(),
"POST",
&url.path,
&headers,
&bytes,
)
.map_err(|e| e.to_string())?;
if resp.is_success() {
Ok(())
} else {
Err(format!("otel: collector returned HTTP {}", resp.status))
}
}
#[cfg(test)]
mod tests {
use super::*;
fn span() -> Span {
Span {
trace_id: "4bf92f3577b34da6a3ce929d0e0e4736".into(),
span_id: "00f067aa0ba902b7".into(),
parent_span_id: None,
name: "invoke_agent".into(),
start_unix_nanos: 1_700_000_000_000_000_000,
end_unix_nanos: 1_700_000_001_000_000_000,
ok: true,
attrs: vec![
("gen_ai.operation.name", str_val("invoke_agent")),
("gen_ai.usage.input_tokens", int_val(1234)),
],
}
}
#[test]
fn otlp_json_has_the_expected_shape() {
let v = to_otlp_json(&[span()], "agentd", "0.1.0");
let rs = &v["resourceSpans"][0];
assert_eq!(
rs["resource"]["attributes"][0]["value"]["stringValue"],
"agentd"
);
let sp = &rs["scopeSpans"][0]["spans"][0];
assert_eq!(sp["traceId"], "4bf92f3577b34da6a3ce929d0e0e4736");
assert_eq!(sp["name"], "invoke_agent");
assert_eq!(sp["status"]["code"], 1); assert_eq!(sp["startTimeUnixNano"], "1700000000000000000"); assert_eq!(sp["attributes"][1]["value"]["intValue"], "1234"); assert!(sp.get("parentSpanId").is_none());
}
#[test]
fn error_span_sets_status_error_and_parent() {
let mut s = span();
s.ok = false;
s.parent_span_id = Some("aaaaaaaaaaaaaaaa".into());
let sp = to_otlp_json(&[s], "agentd", "0.1.0");
let sp = &sp["resourceSpans"][0]["scopeSpans"][0]["spans"][0];
assert_eq!(sp["status"]["code"], 2); assert_eq!(sp["parentSpanId"], "aaaaaaaaaaaaaaaa");
}
#[test]
fn a_run_records_chat_and_tool_children_under_the_run_span() {
let mut run = RunSpan {
trace_id: "4bf92f3577b34da6a3ce929d0e0e4736".into(),
span_id: "00f067aa0ba902b7".into(),
start_unix_nanos: 1_700_000_000_000_000_000,
endpoint: "http://127.0.0.1:4318".into(),
children: Vec::new(),
};
run.record_chat("m", 10, 20, true, 1_700_000_000_000_000_000);
run.record_tool("resource.read", true, 1_700_000_000_500_000_000);
assert_eq!(run.children.len(), 2);
let mut batch = std::mem::take(&mut run.children);
batch.push(Span {
trace_id: run.trace_id.clone(),
span_id: run.span_id.clone(),
parent_span_id: None,
name: "invoke_agent".into(),
start_unix_nanos: run.start_unix_nanos,
end_unix_nanos: 1_700_000_001_000_000_000,
ok: true,
attrs: vec![],
});
let v = to_otlp_json(&batch, "agentd", "0.1.0");
let spans = &v["resourceSpans"][0]["scopeSpans"][0]["spans"];
assert_eq!(spans[0]["name"], "chat");
assert_eq!(spans[0]["parentSpanId"], "00f067aa0ba902b7");
assert_eq!(spans[0]["attributes"][1]["value"]["stringValue"], "m");
assert_eq!(spans[1]["name"], "execute_tool");
assert_eq!(
spans[1]["attributes"][1]["value"]["stringValue"],
"resource.read"
);
assert_eq!(spans[1]["parentSpanId"], "00f067aa0ba902b7");
assert_eq!(spans[2]["name"], "invoke_agent");
assert!(spans[2].get("parentSpanId").is_none()); assert_eq!(spans[0]["traceId"], spans[2]["traceId"]);
}
#[test]
fn otlp_logs_json_has_the_expected_shape() {
let rec = json!({
"timeUnixNano": "1700000000000000000",
"severityText": "INFO",
"body": { "stringValue": "turn.start" },
"attributes": [ { "key": "log.fields", "value": { "stringValue": "{}" } } ],
});
let v = to_otlp_logs_json(&[rec], "agentd", "2.0.0");
let rl = &v["resourceLogs"][0];
assert_eq!(
rl["resource"]["attributes"][0]["value"]["stringValue"],
"agentd"
);
let lr = &rl["scopeLogs"][0]["logRecords"][0];
assert_eq!(lr["body"]["stringValue"], "turn.start");
assert_eq!(lr["severityText"], "INFO");
assert_eq!(lr["timeUnixNano"], "1700000000000000000");
}
}
}