use std::fs::{self, File};
use std::io::{self, BufWriter, Write as _};
use std::path::{Path, PathBuf};
use std::sync::{Mutex, PoisonError};
use chrono::Utc;
use serde_json::{Value, json};
use turnframe_core::ids::TurnId;
use turnframe_core::reduce::ReductionPlan;
use turnframe_core::replay::ReplayRecord;
use turnframe_core::response::AssistantTurn;
use turnframe_core::turn::TurnInput;
use turnframe_core::understanding::Understanding;
use turnframe_provider::trace::{CallTrace, TracedCall, TracedOutcome};
use turnframe_understand::Step;
#[derive(Debug)]
#[non_exhaustive]
pub enum TraceEvent<'a> {
Received {
input: &'a TurnInput,
},
Step {
turn: TurnId,
step: &'a Step,
},
Understood {
turn: TurnId,
understanding: &'a Understanding,
},
Reduced {
turn: TurnId,
plan: &'a ReductionPlan,
},
Completed {
reply: &'a AssistantTurn,
record: &'a ReplayRecord,
},
Failed {
turn: TurnId,
code: &'a str,
record: &'a ReplayRecord,
},
}
pub trait TurnTrace: Send + Sync {
fn event(&self, event: &TraceEvent<'_>);
}
#[derive(Debug)]
pub struct JsonlTrace {
path: PathBuf,
file: Mutex<BufWriter<File>>,
}
impl JsonlTrace {
pub fn create(directory: impl AsRef<Path>) -> io::Result<Self> {
let directory = directory.as_ref();
fs::create_dir_all(directory)?;
let name = format!(
"turnframe-{}-{}.jsonl",
Utc::now().format("%Y%m%d-%H%M%S"),
std::process::id()
);
let path = directory.join(name);
let file = File::create(&path)?;
Ok(Self {
path,
file: Mutex::new(BufWriter::new(file)),
})
}
pub fn from_environment(default: impl AsRef<Path>) -> io::Result<Option<Self>> {
match std::env::var("TURNFRAME_TRACE") {
Ok(value) if value == "1" || value.eq_ignore_ascii_case("true") => {
Self::create(default).map(Some)
}
Ok(value) if !value.trim().is_empty() => Self::create(value.trim()).map(Some),
_ => Ok(None),
}
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
fn write(&self, mut line: Value) {
if let Value::Object(map) = &mut line {
map.insert("at".to_owned(), json!(Utc::now()));
}
let mut file = self.file.lock().unwrap_or_else(PoisonError::into_inner);
let written = serde_json::to_writer(&mut *file, &line)
.map_err(io::Error::from)
.and_then(|()| file.write_all(b"\n"))
.and_then(|()| file.flush());
if let Err(error) = written {
tracing::warn!(
target: "turnframe.trace",
error = %error,
"a trace line could not be written"
);
}
}
}
impl TurnTrace for JsonlTrace {
fn event(&self, event: &TraceEvent<'_>) {
self.write(match event {
TraceEvent::Received { input } => json!({
"event": "turn_received", "turn": input.turn_id, "input": input,
}),
TraceEvent::Step { turn, step } => json!({
"event": "step", "turn": turn, "step": step, "says": step.describe(),
}),
TraceEvent::Understood {
turn,
understanding,
} => json!({
"event": "understood", "turn": turn, "understanding": understanding,
}),
TraceEvent::Reduced { turn, plan } => json!({
"event": "reduced", "turn": turn, "plan": plan,
}),
TraceEvent::Completed { reply, record } => json!({
"event": "turn_completed", "turn": reply.turn_id, "reply": reply, "record": record,
}),
TraceEvent::Failed { turn, code, record } => json!({
"event": "turn_failed", "turn": turn, "code": code, "record": record,
}),
});
}
}
impl CallTrace for JsonlTrace {
fn call(&self, call: &TracedCall<'_>) {
let metadata = &call.request.metadata;
let mut line = json!({
"event": "model_call",
"turn": metadata.get(turnframe_tasks::TURN_LABEL),
"task": metadata.get(turnframe_tasks::TASK_LABEL),
"purpose": call.request.purpose.as_str(),
"provider": call.model.provider,
"model": call.model.model,
"latency_ms": u64::try_from(call.latency.as_millis()).unwrap_or(u64::MAX),
"request": call.request,
});
let (key, value) = match &call.outcome {
TracedOutcome::Response(response) => ("response", json!(response)),
TracedOutcome::Streamed {
text,
finish,
usage,
} => (
"streamed",
json!({"text": text, "finish": finish, "usage": usage}),
),
TracedOutcome::Failed(error) => (
"error",
json!({
"kind": error.kind().as_str(),
"code": error.code().map(|code| code.as_str().to_owned()),
"detail": error.detail().map(|detail| detail.as_str().to_owned()),
}),
),
_ => ("outcome", Value::Null),
};
if let Value::Object(map) = &mut line {
map.insert(key.to_owned(), value);
}
self.write(line);
}
}