Skip to main content

turnframe_runtime/
trace.rs

1//! A trace of whole turns, for a person debugging one: the message that arrived, each
2//! step of its understanding, what was understood and decided, every model call with
3//! its prompts and answer, and the reply with its replay record.
4//!
5//! Opt-in and local. A trace holds the users' words and every prompt, so it belongs on
6//! a developer's disk or in a tool the deployment chose, never on by default. Set a
7//! [`TurnTrace`] with [`OrchestratorBuilder::trace`](crate::orchestrator::OrchestratorBuilder::trace)
8//! and wrap each provider in a [`TracedProvider`](turnframe_provider::trace::TracedProvider);
9//! [`JsonlTrace`] is both, one JSON object per line, ready for `jq` or an importer.
10
11use 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/// One thing a turn did, as a trace receives it.
28#[derive(Debug)]
29#[non_exhaustive]
30pub enum TraceEvent<'a> {
31    /// A turn arrived.
32    Received {
33        /// What arrived.
34        input: &'a TurnInput,
35    },
36    /// Understanding decided something.
37    Step {
38        /// The turn.
39        turn: TurnId,
40        /// What it decided.
41        step: &'a Step,
42    },
43    /// What the message was understood to say, card acts included.
44    Understood {
45        /// The turn.
46        turn: TurnId,
47        /// The understanding reduced.
48        understanding: &'a Understanding,
49    },
50    /// What the reducer decided.
51    Reduced {
52        /// The turn.
53        turn: TurnId,
54        /// The plan.
55        plan: &'a ReductionPlan,
56    },
57    /// The turn finished and this is the reply the user got.
58    Completed {
59        /// The reply.
60        reply: &'a AssistantTurn,
61        /// Everything the turn recorded.
62        record: &'a ReplayRecord,
63    },
64    /// The turn failed.
65    Failed {
66        /// The turn.
67        turn: TurnId,
68        /// The stable error code.
69        code: &'a str,
70        /// What it recorded before it failed.
71        record: &'a ReplayRecord,
72    },
73}
74
75/// Where a runtime reports its turns. It must not block for long: it runs on the turn.
76pub trait TurnTrace: Send + Sync {
77    /// Records one event.
78    fn event(&self, event: &TraceEvent<'_>);
79}
80
81/// Appends every event and model call to one JSON Lines file, flushed line by line
82/// so a crash keeps what came before it.
83#[derive(Debug)]
84pub struct JsonlTrace {
85    path: PathBuf,
86    file: Mutex<BufWriter<File>>,
87}
88
89impl JsonlTrace {
90    /// A new file in `directory`, created if missing, named after the moment it opens.
91    ///
92    /// # Errors
93    ///
94    /// The [`io::Error`] creating the directory or the file.
95    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    /// The trace the `TURNFRAME_TRACE` variable asks for: `1` or `true` writes into
112    /// `default`, any other value names the directory, unset or empty means none.
113    ///
114    /// # Errors
115    ///
116    /// The [`io::Error`] creating the directory or the file.
117    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    /// The file being written.
128    #[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}