use std::collections::BTreeMap;
use serde_json::{Value, json};
use crate::error::{Result, TinyAgentsError};
use crate::harness::events::AgentEvent;
use crate::harness::ids::now_ms;
use crate::harness::observability::AgentObservation;
use crate::harness::usage::Usage;
mod types;
pub use types::{
LangfuseAuth, LangfuseClient, LangfuseScore, LangfuseScoreValue, LangfuseTraceConfig,
};
impl LangfuseClient {
pub fn new(base_url: impl Into<String>, auth: LangfuseAuth) -> Result<Self> {
let endpoint = match &auth {
LangfuseAuth::Basic { .. } => normalize_langfuse_endpoint(base_url.into())?,
LangfuseAuth::Bearer { .. } => normalize_proxy_endpoint(base_url.into())?,
};
Ok(Self {
endpoint,
auth,
client: reqwest::Client::new(),
})
}
pub fn direct(
base_url: impl Into<String>,
public_key: impl Into<String>,
secret_key: impl Into<String>,
) -> Result<Self> {
Self::new(
base_url,
LangfuseAuth::Basic {
public_key: public_key.into(),
secret_key: secret_key.into(),
},
)
}
pub fn proxy(base_url: impl Into<String>, token: impl Into<String>) -> Result<Self> {
Self::new(
base_url,
LangfuseAuth::Bearer {
token: token.into(),
},
)
}
pub fn from_env() -> Result<Self> {
if let Ok(url) = std::env::var("TINYHUMANS_LANGFUSE_PROXY_URL")
&& !url.trim().is_empty()
{
let token = std::env::var("TINYHUMANS_AUTH_TOKEN").map_err(|_| {
TinyAgentsError::Validation(
"TINYHUMANS_AUTH_TOKEN is required for Langfuse proxy mode".to_string(),
)
})?;
return Self::proxy(url, token);
}
let base_url = std::env::var("LANGFUSE_BASE_URL").map_err(|_| {
TinyAgentsError::Validation("LANGFUSE_BASE_URL is required".to_string())
})?;
let public_key = std::env::var("LANGFUSE_PUBLIC_KEY").map_err(|_| {
TinyAgentsError::Validation("LANGFUSE_PUBLIC_KEY is required".to_string())
})?;
let secret_key = std::env::var("LANGFUSE_SECRET_KEY").map_err(|_| {
TinyAgentsError::Validation("LANGFUSE_SECRET_KEY is required".to_string())
})?;
Self::direct(base_url, public_key, secret_key)
}
pub fn endpoint(&self) -> &str {
&self.endpoint
}
pub fn build_ingestion_batch(
&self,
trace: LangfuseTraceConfig,
observations: &[AgentObservation],
) -> Result<Value> {
let trace_id = resolve_trace_id(&trace, observations)?;
let timestamp = observations
.first()
.map(|obs| iso_ms(obs.ts_ms))
.unwrap_or_else(|| iso_ms(now_ms()));
let metadata = trace_metadata(&trace, observations);
let call_models = collect_call_models(observations);
let mut batch = Vec::with_capacity(observations.len() + 2);
batch.push(json!({
"id": format!("{}:trace", trace_id),
"timestamp": timestamp,
"type": "trace-create",
"body": clean_nulls(json!({
"id": trace_id,
"timestamp": timestamp,
"name": trace.name,
"userId": trace.user_id,
"sessionId": trace.session_id,
"environment": trace.environment,
"release": trace.release,
"version": trace.version,
"tags": if trace.tags.is_empty() { Value::Null } else { json!(trace.tags) },
"metadata": metadata,
})),
}));
for run in run_spans(&trace_id, observations) {
batch.push(run);
}
for obs in observations {
if is_run_lifecycle(&obs.event) {
continue;
}
batch.push(observation_event(&trace_id, obs, &call_models));
}
Ok(json!({ "batch": batch }))
}
pub fn build_score_batch(&self, score: LangfuseScore) -> Value {
let timestamp = iso_ms(now_ms());
let score_id = score.id.clone().unwrap_or_else(|| default_score_id(&score));
let event_id = format!("{score_id}:score");
json!({
"batch": [json!({
"id": event_id,
"timestamp": timestamp,
"type": "score-create",
"body": clean_nulls(json!({
"id": score_id,
"traceId": score.trace_id,
"observationId": score.observation_id,
"name": score.name,
"value": score.value.to_value(),
"dataType": score.value.data_type(),
"comment": score.comment,
})),
})]
})
}
pub async fn create_score(&self, score: LangfuseScore) -> Result<Value> {
let payload = self.build_score_batch(score);
self.send_batch(payload).await
}
pub async fn send_observations(
&self,
trace: LangfuseTraceConfig,
observations: &[AgentObservation],
) -> Result<Value> {
let payload = self.build_ingestion_batch(trace, observations)?;
self.send_batch(payload).await
}
pub async fn send_batch(&self, payload: Value) -> Result<Value> {
let mut req = self.client.post(&self.endpoint).json(&payload);
req = match &self.auth {
LangfuseAuth::Basic {
public_key,
secret_key,
} => req.basic_auth(public_key, Some(secret_key)),
LangfuseAuth::Bearer { token } => req.bearer_auth(token),
};
let response = req
.send()
.await
.map_err(|e| TinyAgentsError::Model(format!("Langfuse request failed: {e}")))?;
let status = response.status();
let body = response
.text()
.await
.map_err(|e| TinyAgentsError::Model(format!("Langfuse response read failed: {e}")))?;
let parsed = serde_json::from_str(&body).unwrap_or_else(|_| json!({ "message": body }));
if !status.is_success() && status.as_u16() != 207 {
return Err(TinyAgentsError::Model(format!(
"Langfuse ingestion returned {status}: {parsed}"
)));
}
if status.as_u16() == 207
&& let Some(errors) = parsed.get("errors").and_then(Value::as_array)
&& !errors.is_empty()
{
return Err(TinyAgentsError::Model(format!(
"Langfuse ingestion partially failed ({} rejected): {}",
errors.len(),
json!(errors)
)));
}
Ok(parsed)
}
}
fn normalize_langfuse_endpoint(raw: String) -> Result<String> {
let trimmed = raw.trim().trim_end_matches('/');
if trimmed.is_empty() {
return Err(TinyAgentsError::Validation(
"Langfuse URL must not be empty".to_string(),
));
}
if trimmed.ends_with("/api/public/ingestion")
|| trimmed.ends_with("/telemetry/langfuse/ingestion")
{
return Ok(trimmed.to_string());
}
Ok(format!("{trimmed}/api/public/ingestion"))
}
fn normalize_proxy_endpoint(raw: String) -> Result<String> {
let trimmed = raw.trim().trim_end_matches('/');
if trimmed.is_empty() {
return Err(TinyAgentsError::Validation(
"Langfuse proxy URL must not be empty".to_string(),
));
}
if trimmed.ends_with("/api/public/ingestion")
|| trimmed.ends_with("/telemetry/langfuse/ingestion")
{
return Ok(trimmed.to_string());
}
Ok(format!("{trimmed}/telemetry/langfuse/ingestion"))
}
fn resolve_trace_id(
trace: &LangfuseTraceConfig,
observations: &[AgentObservation],
) -> Result<String> {
if let Some(id) = &trace.trace_id
&& !id.trim().is_empty()
{
return Ok(id.clone());
}
observations
.first()
.map(|obs| obs.root_run_id.as_str().to_string())
.ok_or_else(|| {
TinyAgentsError::Validation("at least one observation is required".to_string())
})
}
fn trace_metadata(trace: &LangfuseTraceConfig, observations: &[AgentObservation]) -> Value {
let mut metadata = serde_json::Map::new();
if let Some(first) = observations.first() {
metadata.insert("root_run_id".to_string(), json!(first.root_run_id.as_str()));
metadata.insert("run_id".to_string(), json!(first.run_id.as_str()));
if let Some(parent) = &first.parent_run_id {
metadata.insert("parent_run_id".to_string(), json!(parent.as_str()));
}
}
if let Value::Object(extra) = &trace.metadata {
for (k, v) in extra {
metadata.insert(k.clone(), v.clone());
}
}
if metadata.is_empty() {
Value::Null
} else {
Value::Object(metadata)
}
}
fn is_run_lifecycle(event: &AgentEvent) -> bool {
matches!(
event,
AgentEvent::RunStarted { .. }
| AgentEvent::RunCompleted { .. }
| AgentEvent::RunFailed { .. }
)
}
fn run_span_id(trace_id: &str, run_id: &str) -> String {
format!("{trace_id}:run:{run_id}")
}
struct RunAcc<'a> {
parent_run_id: Option<&'a str>,
root_run_id: &'a str,
first_ts: u64,
last_ts: u64,
end_ts: Option<u64>,
error: Option<&'a str>,
}
impl<'a> RunAcc<'a> {
fn new(obs: &'a AgentObservation) -> Self {
Self {
parent_run_id: obs.parent_run_id.as_ref().map(|id| id.as_str()),
root_run_id: obs.root_run_id.as_str(),
first_ts: obs.ts_ms,
last_ts: obs.ts_ms,
end_ts: None,
error: None,
}
}
fn observe(&mut self, obs: &'a AgentObservation) {
self.first_ts = self.first_ts.min(obs.ts_ms);
self.last_ts = self.last_ts.max(obs.ts_ms);
match &obs.event {
AgentEvent::RunCompleted { .. } => self.end_ts = Some(obs.ts_ms),
AgentEvent::RunFailed { error, .. } => {
self.end_ts = Some(obs.ts_ms);
self.error = Some(error.as_str());
}
_ => {}
}
}
}
fn run_spans(trace_id: &str, observations: &[AgentObservation]) -> Vec<Value> {
let mut order: Vec<&str> = Vec::new();
let mut runs: BTreeMap<&str, RunAcc> = BTreeMap::new();
for obs in observations {
let rid = obs.run_id.as_str();
runs.entry(rid)
.and_modify(|acc| acc.observe(obs))
.or_insert_with(|| {
order.push(rid);
let mut acc = RunAcc::new(obs);
acc.observe(obs);
acc
});
}
order
.into_iter()
.map(|rid| {
let acc = &runs[rid];
let span_id = run_span_id(trace_id, rid);
let is_root = rid == acc.root_run_id;
let name = if is_root { "agent" } else { "sub-agent" };
let end_iso = iso_ms(acc.end_ts.unwrap_or(acc.last_ts));
json!({
"id": span_id,
"timestamp": end_iso,
"type": "span-create",
"body": clean_nulls(json!({
"id": span_id,
"traceId": trace_id,
"parentObservationId": acc.parent_run_id.map(|p| run_span_id(trace_id, p)),
"name": name,
"startTime": iso_ms(acc.first_ts),
"endTime": end_iso,
"level": acc.error.map(|_| "ERROR"),
"statusMessage": acc.error,
"metadata": json!({
"run_id": rid,
"root_run_id": acc.root_run_id,
"parent_run_id": acc.parent_run_id,
}),
})),
})
})
.collect()
}
fn collect_call_models(observations: &[AgentObservation]) -> BTreeMap<&str, &str> {
let mut models = BTreeMap::new();
for obs in observations {
if let AgentEvent::ModelStarted { call_id, model } = &obs.event {
models.insert(call_id.as_str(), model.as_str());
}
}
models
}
fn observation_event(
trace_id: &str,
obs: &AgentObservation,
call_models: &BTreeMap<&str, &str>,
) -> Value {
let timestamp = iso_ms(obs.ts_ms);
let parent = run_span_id(trace_id, obs.run_id.as_str());
let metadata = json!({
"run_id": obs.run_id.as_str(),
"root_run_id": obs.root_run_id.as_str(),
"parent_run_id": obs.parent_run_id.as_ref().map(|id| id.as_str()),
"offset": obs.offset,
"event_kind": obs.event.kind(),
});
match &obs.event {
AgentEvent::ModelCompleted {
call_id,
started_at_ms,
usage,
input,
output,
} => {
let metadata = with_call_id(metadata, call_id.as_str());
json!({
"id": obs.event_id.as_str(),
"timestamp": timestamp,
"type": "generation-create",
"body": clean_nulls(json!({
"id": scoped_observation_id(trace_id, call_id.as_str()),
"traceId": trace_id,
"parentObservationId": parent,
"name": "model",
"model": call_models.get(call_id.as_str()).copied(),
"startTime": started_at_ms.map(iso_ms).unwrap_or_else(|| timestamp.clone()),
"endTime": timestamp,
"usage": usage.map(langfuse_usage),
"input": input,
"output": output,
"metadata": metadata,
})),
})
}
AgentEvent::ToolCompleted {
call_id,
tool_name,
started_at_ms,
input,
output,
duration_ms,
output_bytes,
error,
} => {
let end_time = match (started_at_ms, duration_ms) {
(Some(start), Some(dur)) => iso_ms(start.saturating_add(*dur)),
_ => timestamp.clone(),
};
let mut tool_metadata = with_call_id(metadata.clone(), call_id.as_str());
if let (Some(map), Some(bytes)) = (tool_metadata.as_object_mut(), output_bytes) {
map.insert("output_bytes".into(), json!(bytes));
}
json!({
"id": obs.event_id.as_str(),
"timestamp": timestamp,
"type": "span-create",
"body": clean_nulls(json!({
"id": scoped_observation_id(trace_id, call_id.as_str()),
"traceId": trace_id,
"parentObservationId": parent,
"name": tool_name,
"startTime": started_at_ms.map(iso_ms).unwrap_or_else(|| timestamp.clone()),
"endTime": end_time,
"input": input,
"output": output,
"level": error.as_ref().map(|_| "ERROR"),
"statusMessage": error,
"metadata": tool_metadata,
})),
})
}
_ => json!({
"id": obs.event_id.as_str(),
"timestamp": timestamp,
"type": "event-create",
"body": clean_nulls(json!({
"id": obs.event_id.as_str(),
"traceId": trace_id,
"parentObservationId": parent,
"name": obs.event.kind(),
"startTime": timestamp,
"metadata": metadata,
})),
}),
}
}
fn default_score_id(score: &LangfuseScore) -> String {
match &score.observation_id {
Some(obs) => format!("{}:{}:score:{}", score.trace_id, obs, score.name),
None => format!("{}:score:{}", score.trace_id, score.name),
}
}
fn scoped_observation_id(trace_id: &str, call_id: &str) -> String {
format!("{trace_id}:{call_id}")
}
fn with_call_id(mut metadata: Value, call_id: &str) -> Value {
if let Some(map) = metadata.as_object_mut() {
map.insert("call_id".into(), json!(call_id));
}
metadata
}
fn langfuse_usage(usage: Usage) -> Value {
json!({
"input": usage.input_tokens,
"output": usage.output_tokens,
"total": usage.total_tokens,
"unit": "TOKENS",
})
}
pub(crate) fn clean_nulls(mut value: Value) -> Value {
if let Value::Object(map) = &mut value {
map.retain(|_, v| !v.is_null());
}
value
}
pub(crate) fn iso_ms(ms: u64) -> String {
use std::time::{Duration, UNIX_EPOCH};
let system_time = UNIX_EPOCH + Duration::from_millis(ms);
let duration = system_time
.duration_since(UNIX_EPOCH)
.unwrap_or(Duration::from_secs(0));
let secs = duration.as_secs();
let millis = duration.subsec_millis();
format_unix_iso(secs, millis)
}
fn format_unix_iso(secs: u64, millis: u32) -> String {
let days = (secs / 86_400) as i64;
let day_secs = secs % 86_400;
let (year, month, day) = civil_from_days(days);
let hour = day_secs / 3_600;
let minute = (day_secs % 3_600) / 60;
let second = day_secs % 60;
format!("{year:04}-{month:02}-{day:02}T{hour:02}:{minute:02}:{second:02}.{millis:03}Z")
}
fn civil_from_days(days: i64) -> (i32, u32, 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 / 1_460 + 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 = mp + if mp < 10 { 3 } else { -9 };
let year = y + if m <= 2 { 1 } else { 0 };
(year as i32, m as u32, d as u32)
}
#[cfg(test)]
mod test;