Skip to main content

ryu_tracing/
lib.rs

1//! Per-run observability trace store (spec unit #178 / M4).
2//!
3//! Persists ordered spans keyed by `conversation_id` (the run id) in a local
4//! SQLite database (`~/.ryu/traces.db`).  Each span records:
5//!   - `kind`        — `"tool-call"` or `"model-call"`
6//!   - `name`        — tool name or model id
7//!   - `args_hash`   — SHA-256 hex of the tool input (tool-call spans only;
8//!                     NOT the raw payload — privacy by default)
9//!   - `started_at`  — Unix milliseconds (write time)
10//!   - `ended_at`    — Unix milliseconds (finish time; `None` while running)
11//!   - `error`       — non-`None` when the span ended with an error
12//!   - `session_id`  — nullable link to the gateway audit row (populated when
13//!                     #176 threads the id; left `None` until then)
14//!
15//! Placement rationale (Core vs Gateway, see CLAUDE.md §1): span ordering and
16//! tool-call sequencing are *what ran* (orchestration) — Core.  Token counts,
17//! cost, and provider-latency are *what is measured/paid* — Gateway audit only.
18//! Core's trace intentionally stores NO tokens/cost fields.
19
20use std::path::PathBuf;
21use std::sync::Arc;
22
23use anyhow::{Context, Result};
24use rusqlite::{params, Connection};
25use serde::{Deserialize, Serialize};
26use tokio::sync::Mutex;
27
28/// A single ordered span within a run.
29#[derive(Debug, Clone, Serialize, Deserialize)]
30pub struct Span {
31    /// Stable row id (UUID).
32    pub id: String,
33    /// The run this span belongs to (`conversation_id` from the chat request).
34    pub conversation_id: String,
35    /// `"tool-call"` or `"model-call"`.
36    pub kind: String,
37    /// Tool name (tool-call) or model id (model-call).
38    pub name: String,
39    /// SHA-256 hex of the raw tool-input JSON (tool-call only).  Never the raw
40    /// payload — protects sensitive args from being stored in Core.
41    pub args_hash: Option<String>,
42    /// Unix milliseconds — when the span was opened.
43    pub started_at: i64,
44    /// Unix milliseconds — when the span was closed.  `None` while in-flight.
45    pub ended_at: Option<i64>,
46    /// Error message if the span ended with a failure.
47    pub error: Option<String>,
48    /// Nullable link to the gateway audit row via `x-ryu-session` (populated
49    /// by #176; `None` until that thread lands).
50    pub session_id: Option<String>,
51    /// Autoincrement ordering key (monotonically increasing within the DB).
52    pub seq: i64,
53}
54
55/// SQLite-backed trace store.  Cheap to clone — wraps an `Arc<Mutex<Connection>>`.
56#[derive(Clone)]
57pub struct TraceStore {
58    conn: Arc<Mutex<Connection>>,
59}
60
61fn now_millis() -> i64 {
62    chrono::Utc::now().timestamp_millis()
63}
64
65/// SHA-256 hex of a JSON value's canonical representation.
66pub fn hash_args(value: &serde_json::Value) -> String {
67    use std::fmt::Write;
68    // Canonical form: serialize the value (serde_json is deterministic for
69    // the same in-memory Value; sufficient for a privacy-safe fingerprint).
70    let bytes = serde_json::to_vec(value).unwrap_or_default();
71    let digest = sha2_digest(&bytes);
72    let mut out = String::with_capacity(64);
73    for b in digest {
74        let _ = write!(out, "{b:02x}");
75    }
76    out
77}
78
79/// Minimal SHA-256 implementation using the `sha2` crate (already in the
80/// dependency tree via rustls/ring).  Isolated here so the rest of the module
81/// has no sha2 import noise.
82fn sha2_digest(data: &[u8]) -> [u8; 32] {
83    use sha2::Digest;
84    let mut h = sha2::Sha256::new();
85    h.update(data);
86    h.finalize().into()
87}
88
89impl TraceStore {
90    /// Open (or create) the trace store at a specific path.
91    pub fn open(path: PathBuf) -> Result<Self> {
92        if let Some(parent) = path.parent() {
93            std::fs::create_dir_all(parent)
94                .with_context(|| format!("creating trace db dir {}", parent.display()))?;
95        }
96        let conn = Connection::open(&path)
97            .with_context(|| format!("opening trace db {}", path.display()))?;
98        Self::init_schema(&conn)?;
99        Ok(Self {
100            conn: Arc::new(Mutex::new(conn)),
101        })
102    }
103
104    /// Open an in-memory store (tests / ephemeral use).
105    pub fn open_in_memory() -> Result<Self> {
106        let conn = Connection::open_in_memory().context("opening in-memory trace db")?;
107        Self::init_schema(&conn)?;
108        Ok(Self {
109            conn: Arc::new(Mutex::new(conn)),
110        })
111    }
112
113    fn init_schema(conn: &Connection) -> Result<()> {
114        conn.execute_batch(
115            "PRAGMA journal_mode = WAL;
116             CREATE TABLE IF NOT EXISTS spans (
117                 seq             INTEGER PRIMARY KEY AUTOINCREMENT,
118                 id              TEXT NOT NULL UNIQUE,
119                 conversation_id TEXT NOT NULL,
120                 kind            TEXT NOT NULL,
121                 name            TEXT NOT NULL,
122                 args_hash       TEXT,
123                 started_at      INTEGER NOT NULL,
124                 ended_at        INTEGER,
125                 error           TEXT,
126                 session_id      TEXT
127             );
128             CREATE INDEX IF NOT EXISTS idx_spans_conversation
129                 ON spans(conversation_id, seq);",
130        )
131        .context("initializing trace schema")?;
132
133        // Additive migration guard — safe to call on every startup.
134        let existing: std::collections::HashSet<String> = {
135            let mut stmt = conn.prepare("PRAGMA table_info(spans)")?;
136            let names = stmt.query_map([], |row| row.get::<_, String>(1))?;
137            names.filter_map(|r| r.ok()).collect()
138        };
139        if !existing.contains("session_id") {
140            conn.execute_batch("ALTER TABLE spans ADD COLUMN session_id TEXT")
141                .context("adding session_id column")?;
142        }
143
144        Ok(())
145    }
146
147    /// Open a span (tool-call or model-call).  Returns the new span id.
148    ///
149    /// `args_hash` should be `Some(hash_args(&input))` for tool-call spans and
150    /// `None` for model-call spans.
151    pub async fn open_span(
152        &self,
153        conversation_id: &str,
154        kind: &str,
155        name: &str,
156        args_hash: Option<&str>,
157        session_id: Option<&str>,
158    ) -> Result<String> {
159        let span_id = uuid::Uuid::new_v4().to_string();
160        let now = now_millis();
161        let conn = self.conn.lock().await;
162        conn.execute(
163            "INSERT INTO spans (id, conversation_id, kind, name, args_hash, started_at, session_id)
164             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
165            params![
166                span_id,
167                conversation_id,
168                kind,
169                name,
170                args_hash,
171                now,
172                session_id
173            ],
174        )
175        .context("inserting span")?;
176        Ok(span_id)
177    }
178
179    /// Close a span — set `ended_at` and optionally record an error.
180    pub async fn close_span(&self, span_id: &str, error: Option<&str>) -> Result<()> {
181        let now = now_millis();
182        let conn = self.conn.lock().await;
183        conn.execute(
184            "UPDATE spans SET ended_at = ?1, error = ?2 WHERE id = ?3",
185            params![now, error, span_id],
186        )
187        .context("closing span")?;
188        Ok(())
189    }
190
191    /// Return all spans for a run in ascending `seq` order.
192    pub async fn get_spans(&self, conversation_id: &str) -> Result<Vec<Span>> {
193        let conn = self.conn.lock().await;
194        let mut stmt = conn.prepare(
195            "SELECT seq, id, conversation_id, kind, name, args_hash,
196                    started_at, ended_at, error, session_id
197             FROM spans
198             WHERE conversation_id = ?1
199             ORDER BY seq ASC",
200        )?;
201        let rows = stmt.query_map(params![conversation_id], |row| {
202            Ok(Span {
203                seq: row.get(0)?,
204                id: row.get(1)?,
205                conversation_id: row.get(2)?,
206                kind: row.get(3)?,
207                name: row.get(4)?,
208                args_hash: row.get(5)?,
209                started_at: row.get(6)?,
210                ended_at: row.get(7)?,
211                error: row.get(8)?,
212                session_id: row.get(9)?,
213            })
214        })?;
215        rows.collect::<std::result::Result<Vec<_>, _>>()
216            .context("reading spans")
217    }
218}
219
220#[cfg(test)]
221mod tests {
222    use super::*;
223
224    #[tokio::test]
225    async fn write_and_read_back_tool_call_span() {
226        let store = TraceStore::open_in_memory().unwrap();
227        let conv_id = "test-conv-1";
228
229        // Open a tool-call span.
230        let input = serde_json::json!({ "path": "/tmp/foo.txt" });
231        let ah = hash_args(&input);
232        let span_id = store
233            .open_span(conv_id, "tool-call", "read_file", Some(&ah), None)
234            .await
235            .unwrap();
236
237        // Close it successfully.
238        store.close_span(&span_id, None).await.unwrap();
239
240        // Read back.
241        let spans = store.get_spans(conv_id).await.unwrap();
242        assert_eq!(spans.len(), 1);
243        let s = &spans[0];
244        assert_eq!(s.conversation_id, conv_id);
245        assert_eq!(s.kind, "tool-call");
246        assert_eq!(s.name, "read_file");
247        assert!(s.args_hash.is_some());
248        assert!(s.ended_at.is_some());
249        assert!(s.error.is_none());
250    }
251
252    #[tokio::test]
253    async fn error_span_records_message() {
254        let store = TraceStore::open_in_memory().unwrap();
255        let span_id = store
256            .open_span("conv-2", "tool-call", "bash", None, None)
257            .await
258            .unwrap();
259        store
260            .close_span(&span_id, Some("permission denied"))
261            .await
262            .unwrap();
263        let spans = store.get_spans("conv-2").await.unwrap();
264        assert_eq!(spans[0].error.as_deref(), Some("permission denied"));
265    }
266
267    #[tokio::test]
268    async fn multiple_spans_ordered_by_seq() {
269        let store = TraceStore::open_in_memory().unwrap();
270        let conv = "conv-order";
271        for name in ["alpha", "beta", "gamma"] {
272            let id = store
273                .open_span(conv, "tool-call", name, None, None)
274                .await
275                .unwrap();
276            store.close_span(&id, None).await.unwrap();
277        }
278        let spans = store.get_spans(conv).await.unwrap();
279        assert_eq!(spans.len(), 3);
280        assert!(spans[0].seq < spans[1].seq);
281        assert!(spans[1].seq < spans[2].seq);
282        assert_eq!(spans[0].name, "alpha");
283        assert_eq!(spans[2].name, "gamma");
284    }
285}