1use std::fs::{self, File};
12use std::io::{self, BufWriter, Write as _};
13use std::path::{Path, PathBuf};
14use std::sync::{Mutex, PoisonError};
15
16use chrono::Utc;
17use serde_json::{Value, json};
18use turnframe_core::ids::TurnId;
19use turnframe_core::reduce::ReductionPlan;
20use turnframe_core::replay::ReplayRecord;
21use turnframe_core::response::AssistantTurn;
22use turnframe_core::turn::TurnInput;
23use turnframe_core::understanding::Understanding;
24use turnframe_provider::trace::{CallTrace, TracedCall, TracedOutcome};
25use turnframe_understand::Step;
26
27#[derive(Debug)]
29#[non_exhaustive]
30pub enum TraceEvent<'a> {
31 Received {
33 input: &'a TurnInput,
35 },
36 Step {
38 turn: TurnId,
40 step: &'a Step,
42 },
43 Understood {
45 turn: TurnId,
47 understanding: &'a Understanding,
49 },
50 Reduced {
52 turn: TurnId,
54 plan: &'a ReductionPlan,
56 },
57 Completed {
59 reply: &'a AssistantTurn,
61 record: &'a ReplayRecord,
63 },
64 Failed {
66 turn: TurnId,
68 code: &'a str,
70 record: &'a ReplayRecord,
72 },
73}
74
75pub trait TurnTrace: Send + Sync {
77 fn event(&self, event: &TraceEvent<'_>);
79}
80
81#[derive(Debug)]
84pub struct JsonlTrace {
85 path: PathBuf,
86 file: Mutex<BufWriter<File>>,
87}
88
89impl JsonlTrace {
90 pub fn create(directory: impl AsRef<Path>) -> io::Result<Self> {
96 let directory = directory.as_ref();
97 fs::create_dir_all(directory)?;
98 let name = format!(
99 "turnframe-{}-{}.jsonl",
100 Utc::now().format("%Y%m%d-%H%M%S"),
101 std::process::id()
102 );
103 let path = directory.join(name);
104 let file = File::create(&path)?;
105 Ok(Self {
106 path,
107 file: Mutex::new(BufWriter::new(file)),
108 })
109 }
110
111 pub fn from_environment(default: impl AsRef<Path>) -> io::Result<Option<Self>> {
118 match std::env::var("TURNFRAME_TRACE") {
119 Ok(value) if value == "1" || value.eq_ignore_ascii_case("true") => {
120 Self::create(default).map(Some)
121 }
122 Ok(value) if !value.trim().is_empty() => Self::create(value.trim()).map(Some),
123 _ => Ok(None),
124 }
125 }
126
127 #[must_use]
129 pub fn path(&self) -> &Path {
130 &self.path
131 }
132
133 fn write(&self, mut line: Value) {
134 if let Value::Object(map) = &mut line {
135 map.insert("at".to_owned(), json!(Utc::now()));
136 }
137 let mut file = self.file.lock().unwrap_or_else(PoisonError::into_inner);
138 let written = serde_json::to_writer(&mut *file, &line)
139 .map_err(io::Error::from)
140 .and_then(|()| file.write_all(b"\n"))
141 .and_then(|()| file.flush());
142 if let Err(error) = written {
143 tracing::warn!(
144 target: "turnframe.trace",
145 error = %error,
146 "a trace line could not be written"
147 );
148 }
149 }
150}
151
152impl TurnTrace for JsonlTrace {
153 fn event(&self, event: &TraceEvent<'_>) {
154 self.write(match event {
155 TraceEvent::Received { input } => json!({
156 "event": "turn_received", "turn": input.turn_id, "input": input,
157 }),
158 TraceEvent::Step { turn, step } => json!({
159 "event": "step", "turn": turn, "step": step, "says": step.describe(),
160 }),
161 TraceEvent::Understood {
162 turn,
163 understanding,
164 } => json!({
165 "event": "understood", "turn": turn, "understanding": understanding,
166 }),
167 TraceEvent::Reduced { turn, plan } => json!({
168 "event": "reduced", "turn": turn, "plan": plan,
169 }),
170 TraceEvent::Completed { reply, record } => json!({
171 "event": "turn_completed", "turn": reply.turn_id, "reply": reply, "record": record,
172 }),
173 TraceEvent::Failed { turn, code, record } => json!({
174 "event": "turn_failed", "turn": turn, "code": code, "record": record,
175 }),
176 });
177 }
178}
179
180impl CallTrace for JsonlTrace {
181 fn call(&self, call: &TracedCall<'_>) {
182 let metadata = &call.request.metadata;
183 let mut line = json!({
184 "event": "model_call",
185 "turn": metadata.get(turnframe_tasks::TURN_LABEL),
186 "task": metadata.get(turnframe_tasks::TASK_LABEL),
187 "purpose": call.request.purpose.as_str(),
188 "provider": call.model.provider,
189 "model": call.model.model,
190 "latency_ms": u64::try_from(call.latency.as_millis()).unwrap_or(u64::MAX),
191 "request": call.request,
192 });
193 let (key, value) = match &call.outcome {
194 TracedOutcome::Response(response) => ("response", json!(response)),
195 TracedOutcome::Streamed {
196 text,
197 finish,
198 usage,
199 } => (
200 "streamed",
201 json!({"text": text, "finish": finish, "usage": usage}),
202 ),
203 TracedOutcome::Failed(error) => (
204 "error",
205 json!({
206 "kind": error.kind().as_str(),
207 "code": error.code().map(|code| code.as_str().to_owned()),
208 "detail": error.detail().map(|detail| detail.as_str().to_owned()),
209 }),
210 ),
211 _ => ("outcome", Value::Null),
212 };
213 if let Value::Object(map) = &mut line {
214 map.insert(key.to_owned(), value);
215 }
216 self.write(line);
217 }
218}