From 61d286eadfe080c63fabf4a076a6c328ad1ece39 Mon Sep 17 00:00:00 2001 From: Andrew Kent Date: Wed, 9 Sep 2026 12:42:10 -0600 Subject: [PATCH] add data enrichments for grok --- bt-daemon/src/translate/grok.rs | 31 ++++++++++++-- bt-daemon/src/translate/mod.rs | 2 +- bt-daemon/tests/grok_translator.rs | 67 ++++++++++++++++++++++++++++++ 3 files changed, 95 insertions(+), 5 deletions(-) diff --git a/bt-daemon/src/translate/grok.rs b/bt-daemon/src/translate/grok.rs index 899a23e..ba699f8 100644 --- a/bt-daemon/src/translate/grok.rs +++ b/bt-daemon/src/translate/grok.rs @@ -4,13 +4,17 @@ //! trace data source; `events.jsonl` is independently tailed, best-effort tool //! enrichment. Both are mirrored into daemon-owned storage at hook boundaries. +use super::git::GitMetadataCache; use super::recent::{RecentMap, RecentSet}; -use super::{AgentTranslator, SessionCtx, SpanOp, SpanRow, SpanType, TranslatorFactory}; +use super::{ + local_username, AgentTranslator, SessionCtx, SpanOp, SpanRow, SpanType, TranslatorFactory, +}; use crate::ids; use crate::wire::Envelope; use serde_json::{json, Map, Value}; use std::collections::BTreeMap; use std::io::{Read, Seek, SeekFrom}; +use std::sync::Arc; const CATCH_UP_BYTE_BUDGET: u64 = 64 * 1024; const CATCH_UP_RECORD_BUDGET: usize = 256; @@ -21,7 +25,15 @@ const MAX_OUTPUT_CHUNKS: usize = 2_048; const MAX_OUTPUT_BYTES: usize = 2 * 1024 * 1024; const MAX_SYSTEM_PROMPT_BYTES: u64 = 2 * 1024 * 1024; -pub struct GrokTranslatorFactory; +pub struct GrokTranslatorFactory { + git: Arc, +} + +impl GrokTranslatorFactory { + pub(super) fn new(git: Arc) -> Self { + Self { git } + } +} impl TranslatorFactory for GrokTranslatorFactory { fn source(&self) -> &str { @@ -29,7 +41,7 @@ impl TranslatorFactory for GrokTranslatorFactory { } fn create(&self, session_id: &str) -> Box { - Box::new(GrokTranslator::new(session_id)) + Box::new(GrokTranslator::new(session_id, self.git.clone())) } } @@ -214,11 +226,13 @@ struct GrokTranslator { first_llm_span_id: Option, first_llm_user_input: Option, pending: Option, + cwd: Option, + git: Arc, last_ts_ms: i64, } impl GrokTranslator { - fn new(session_id: &str) -> Self { + fn new(session_id: &str, git: Arc) -> Self { let root = ids::span_id(session_id, "session"); Self { session_id: session_id.to_string(), @@ -241,6 +255,8 @@ impl GrokTranslator { first_llm_span_id: None, first_llm_user_input: None, pending: None, + cwd: None, + git, last_ts_ms: 0, } } @@ -305,6 +321,7 @@ impl GrokTranslator { metadata.insert("source".into(), json!("grok")); metadata.insert("session_id".into(), json!(self.session_id)); metadata.insert("trace_source".into(), json!("session_transcript")); + metadata.insert("username".into(), json!(local_username())); for field in ["cwd", "workspaceRoot", "permissionMode", "transcriptPath"] { if let Some(value) = event.payload.get(field) { metadata.insert(field.into(), value.clone()); @@ -1180,6 +1197,9 @@ impl AgentTranslator for GrokTranslator { self.pending.is_none(), "Grok translator has pending catch-up work; drain it before handling another event" ); + if let Some(cwd) = event.payload.get("cwd").and_then(Value::as_str) { + self.cwd = Some(cwd.to_string()); + } let (mut ops, complete) = self.process_transcripts(event, ctx)?; if complete { self.finish_hook(event, &mut ops); @@ -1188,6 +1208,7 @@ impl AgentTranslator for GrokTranslator { event: event.clone(), }); } + self.git.enrich_rows(self.cwd.as_deref(), &mut ops); Ok(ops) } @@ -1200,6 +1221,7 @@ impl AgentTranslator for GrokTranslator { self.pending = None; self.finish_hook(&pending.event, &mut ops); } + self.git.enrich_rows(self.cwd.as_deref(), &mut ops); Ok(Some(ops)) } @@ -1221,6 +1243,7 @@ impl AgentTranslator for GrokTranslator { ..Default::default() })); } + self.git.enrich_rows(self.cwd.as_deref(), &mut ops); Ok(ops) } } diff --git a/bt-daemon/src/translate/mod.rs b/bt-daemon/src/translate/mod.rs index 473e996..038350d 100644 --- a/bt-daemon/src/translate/mod.rs +++ b/bt-daemon/src/translate/mod.rs @@ -162,7 +162,7 @@ impl Registry { r.register(Box::new(AntigravityTranslatorFactory::new(git.clone()))); r.register(Box::new(ClaudeTranslatorFactory::new(git.clone()))); r.register(Box::new(CodexTranslatorFactory::new(git.clone()))); - r.register(Box::new(GrokTranslatorFactory)); + r.register(Box::new(GrokTranslatorFactory::new(git.clone()))); r.register(Box::new(OpenCodeTranslatorFactory::new(git.clone()))); r.register(Box::new(PiTranslatorFactory::new(git))); r diff --git a/bt-daemon/tests/grok_translator.rs b/bt-daemon/tests/grok_translator.rs index 761ac3f..d2c01c8 100644 --- a/bt-daemon/tests/grok_translator.rs +++ b/bt-daemon/tests/grok_translator.rs @@ -49,6 +49,11 @@ fn ctx() -> SessionCtx { config: None, } } +fn expected_username() -> String { + std::env::var("USER") + .or_else(|_| std::env::var("USERNAME")) + .unwrap_or_default() +} fn point_at(event: &mut Envelope, updates: &Path, events: &Path) { event.payload["_bt_grok_transcript_mirrors"]["updates"]["mirror"] = json!(updates); @@ -226,6 +231,68 @@ fn grok_transcript_builds_turn_llm_and_tool_spans_with_aggregate_usage() { }) ); } +#[test] +fn grok_enriches_emitted_spans_from_the_captured_working_directory() { + let temp = tempfile::tempdir().unwrap(); + let repo = temp.path().join("repo"); + std::fs::create_dir(&repo).unwrap(); + let git = |args: &[&str]| { + assert!(std::process::Command::new("git") + .arg("-C") + .arg(&repo) + .args(args) + .status() + .unwrap() + .success()); + }; + git(&["init", "-b", "main"]); + git(&["config", "user.email", "test@example.com"]); + git(&["config", "user.name", "Test"]); + std::fs::write(repo.join("README.md"), "test").unwrap(); + git(&["add", "README.md"]); + git(&["commit", "-m", "initial"]); + git(&[ + "remote", + "add", + "origin", + "https://secret@example.com/acme/app.git", + ]); + let commit = std::process::Command::new("git") + .arg("-C") + .arg(&repo) + .args(["rev-parse", "HEAD"]) + .output() + .unwrap(); + let commit = String::from_utf8(commit.stdout).unwrap().trim().to_string(); + + let updates = std::fs::metadata(fixture("updates.jsonl")).unwrap().len(); + let events = std::fs::metadata(fixture("events.jsonl")).unwrap().len(); + let mut event = envelope(updates, events); + event.payload["cwd"] = json!(repo); + let registry = Registry::default_agents(); + let mut translator = registry.create("grok", "grok-session"); + let ops = translator.handle(&event, &ctx()).unwrap(); + + let inserts: Vec<_> = ops + .iter() + .filter_map(|op| match op { + SpanOp::Insert(row) => Some(row), + SpanOp::Merge(_) => None, + }) + .collect(); + let root = inserts.iter().find(|row| row.name == "Grok").unwrap(); + assert_eq!( + root.metadata.as_ref().unwrap()["username"], + json!(expected_username()) + ); + assert!(inserts.iter().all(|row| { + row.metadata.as_ref().is_some_and(|metadata| { + metadata["git_origin_url"] == "https://example.com/acme/app.git" + && metadata["git_branch"] == "main" + && metadata["git_commit_sha"] == commit + }) + })); +} #[test] fn grok_transcript_reads_incrementally_and_replay_is_deterministic() {