Skip to content

Commit bf33949

Browse files
authored
fix(grok): parse token usage from session updates (#65)
Extract per-prompt token deltas from Grok session update _meta fields, attach usage events to parsed sessions, and backfill via usage_parser_version. Signed-off-by: samzong <samzong.lu@gmail.com>
1 parent 39075ed commit bf33949

1 file changed

Lines changed: 203 additions & 15 deletions

File tree

src/adapters/grok.rs

Lines changed: 203 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,19 @@
1+
use std::collections::HashMap;
12
use std::fs;
23
use std::io::{BufRead, BufReader};
34
use std::path::{Path, PathBuf};
45

56
use serde_json::Value;
67
use tracing::debug;
78

8-
use crate::adapters::file_scan::{self, FileScanEntry};
9+
use crate::adapters::file_scan::{self, FileScanEntry, FileScanOptions};
910
use crate::adapters::{
1011
RawMessage, RawSession, ResumeCommand, SourceAdapter, SyncScanResult, SyncScanStats,
1112
};
1213
use crate::db::store::Store;
13-
use crate::types::Role;
14+
use crate::types::{RawUsageEvent, Role, TokenSource};
15+
16+
const USAGE_PARSER_VERSION: u32 = 1;
1417

1518
pub struct GrokAdapter;
1619

@@ -23,6 +26,10 @@ impl SourceAdapter for GrokAdapter {
2326
"GRK"
2427
}
2528

29+
fn usage_parser_version(&self) -> Option<u32> {
30+
Some(USAGE_PARSER_VERSION)
31+
}
32+
2633
fn resume_command(&self, source_id: &str) -> Option<ResumeCommand> {
2734
Some(ResumeCommand {
2835
program: "grok".to_string(),
@@ -73,6 +80,14 @@ struct GrokSummary {
7380
directory: Option<String>,
7481
started_at: i64,
7582
updated_at: Option<i64>,
83+
current_model_id: Option<String>,
84+
}
85+
86+
struct PromptTokenState {
87+
prompt_id: String,
88+
peak_total_tokens: i64,
89+
timestamp: Option<i64>,
90+
model: Option<String>,
7691
}
7792

7893
#[derive(Clone, Debug, PartialEq, Eq)]
@@ -102,9 +117,17 @@ fn scan_for_sync_impl(
102117
since_ts: Option<i64>,
103118
) -> anyhow::Result<SyncScanResult> {
104119
let (entries, _) = collect_grok_entries(sessions_dir);
105-
file_scan::run_file_scan(store, "grok", since_ts, entries, |entry, mtime_ms| {
106-
parse_grok_session_for_entry(&entry, mtime_ms)
107-
})
120+
file_scan::run_file_scan_with_options(
121+
store,
122+
"grok",
123+
since_ts,
124+
FileScanOptions {
125+
usage_parser_version: Some(USAGE_PARSER_VERSION),
126+
event_parser_version: None,
127+
},
128+
entries,
129+
|entry, mtime_ms| parse_grok_session_for_entry(&entry, mtime_ms),
130+
)
108131
}
109132

110133
fn prune_impl(sessions_dir: &Path, store: &Store) -> anyhow::Result<()> {
@@ -243,7 +266,7 @@ fn parse_grok_session_for_entry(
243266
}
244267
};
245268

246-
let messages = match parse_grok_updates(&entry.stat_target) {
269+
let (messages, mut usage_events) = match parse_grok_updates(&entry.stat_target) {
247270
Ok(parsed) => parsed,
248271
Err(err) => {
249272
debug!("failed to parse Grok updates {}: {err}", entry.stat_target.display());
@@ -255,6 +278,15 @@ fn parse_grok_session_for_entry(
255278
return Ok(None);
256279
}
257280

281+
let fallback_model = summary.current_model_id.clone().unwrap_or_else(|| "grok".to_string());
282+
let source_path = entry.stat_target.to_str().map(str::to_string);
283+
for event in &mut usage_events {
284+
if event.model.is_empty() {
285+
event.model = fallback_model.clone();
286+
}
287+
event.source_path = source_path.clone();
288+
}
289+
258290
let started_at =
259291
summary.started_at.max(messages.first().and_then(|message| message.timestamp).unwrap_or(0));
260292

@@ -265,7 +297,8 @@ fn parse_grok_session_for_entry(
265297
summary.updated_at.or(Some(mtime_ms)),
266298
None,
267299
messages,
268-
);
300+
)
301+
.with_usage(usage_events, USAGE_PARSER_VERSION);
269302

270303
session.updated_at = Some(mtime_ms);
271304
Ok(Some(session))
@@ -292,19 +325,33 @@ fn load_grok_summary(session_dir: &Path, fallback_id: &str) -> anyhow::Result<Gr
292325
.filter(|cwd| !cwd.is_empty())
293326
.map(str::to_string);
294327

328+
let current_model_id = doc
329+
.get("current_model_id")
330+
.and_then(|model| model.as_str())
331+
.map(str::trim)
332+
.filter(|model| !model.is_empty())
333+
.map(str::to_string);
334+
295335
let started_at = parse_rfc3339_ms(doc.get("created_at")).unwrap_or(0);
296336
let updated_at = parse_rfc3339_ms(doc.get("updated_at"));
297-
Ok(GrokSummary { session_id, directory, started_at, updated_at })
337+
Ok(GrokSummary { session_id, directory, started_at, updated_at, current_model_id })
298338
}
299339

300-
fn parse_grok_updates(updates_path: &Path) -> anyhow::Result<Vec<RawMessage>> {
340+
fn parse_grok_updates(
341+
updates_path: &Path,
342+
) -> anyhow::Result<(Vec<RawMessage>, Vec<RawUsageEvent>)> {
301343
let file = fs::File::open(updates_path)?;
302344
parse_grok_updates_reader(BufReader::new(file))
303345
}
304346

305-
fn parse_grok_updates_reader<R: BufRead>(reader: R) -> anyhow::Result<Vec<RawMessage>> {
347+
fn parse_grok_updates_reader<R: BufRead>(
348+
reader: R,
349+
) -> anyhow::Result<(Vec<RawMessage>, Vec<RawUsageEvent>)> {
306350
let mut messages = Vec::new();
307351
let mut pending_agent: Option<PendingAgentMessage> = None;
352+
let mut prompts: Vec<PromptTokenState> = Vec::new();
353+
let mut prompt_index: HashMap<String, usize> = HashMap::new();
354+
let mut last_model: Option<String> = None;
308355
for line in reader.lines() {
309356
let line = line?;
310357
let line = line.trim();
@@ -327,6 +374,7 @@ fn parse_grok_updates_reader<R: BufRead>(reader: R) -> anyhow::Result<Vec<RawMes
327374
let session_update =
328375
update.get("sessionUpdate").and_then(|value| value.as_str()).unwrap_or("");
329376
let timestamp_ms = doc.get("timestamp").and_then(json_i64).map(|ts| ts * 1000);
377+
track_prompt_tokens(params, timestamp_ms, &mut prompts, &mut prompt_index, &mut last_model);
330378

331379
match session_update {
332380
"user_message_chunk" => {
@@ -377,7 +425,83 @@ fn parse_grok_updates_reader<R: BufRead>(reader: R) -> anyhow::Result<Vec<RawMes
377425
}
378426
}
379427

380-
Ok(messages)
428+
Ok((messages, prompt_usage_events(prompts)))
429+
}
430+
431+
fn track_prompt_tokens(
432+
params: &Value,
433+
timestamp_ms: Option<i64>,
434+
prompts: &mut Vec<PromptTokenState>,
435+
prompt_index: &mut HashMap<String, usize>,
436+
last_model: &mut Option<String>,
437+
) {
438+
let Some(meta) = params.get("_meta") else {
439+
return;
440+
};
441+
if let Some(model) = meta.get("modelId").and_then(|value| value.as_str())
442+
&& !model.is_empty()
443+
{
444+
*last_model = Some(model.to_string());
445+
}
446+
let Some(prompt_id) =
447+
meta.get("promptId").and_then(|value| value.as_str()).filter(|id| !id.is_empty())
448+
else {
449+
return;
450+
};
451+
let Some(total_tokens) = meta.get("totalTokens").and_then(json_i64) else {
452+
return;
453+
};
454+
let index = match prompt_index.get(prompt_id) {
455+
Some(index) => *index,
456+
None => {
457+
prompts.push(PromptTokenState {
458+
prompt_id: prompt_id.to_string(),
459+
peak_total_tokens: 0,
460+
timestamp: None,
461+
model: None,
462+
});
463+
prompt_index.insert(prompt_id.to_string(), prompts.len() - 1);
464+
prompts.len() - 1
465+
}
466+
};
467+
let prompt = &mut prompts[index];
468+
prompt.peak_total_tokens = prompt.peak_total_tokens.max(total_tokens);
469+
if timestamp_ms.is_some() {
470+
prompt.timestamp = timestamp_ms;
471+
}
472+
if prompt.model.is_none() {
473+
prompt.model = last_model.clone();
474+
}
475+
}
476+
477+
fn prompt_usage_events(prompts: Vec<PromptTokenState>) -> Vec<RawUsageEvent> {
478+
let mut events = Vec::new();
479+
let mut previous_peak = 0i64;
480+
for prompt in prompts {
481+
let delta = prompt.peak_total_tokens - previous_peak;
482+
previous_peak = prompt.peak_total_tokens;
483+
if delta <= 0 {
484+
continue;
485+
}
486+
events.push(RawUsageEvent {
487+
event_key: format!("prompt:{}", prompt.prompt_id),
488+
event_seq: events.len() as u32,
489+
message_seq: None,
490+
timestamp: prompt.timestamp.unwrap_or_default(),
491+
model: prompt.model.unwrap_or_default(),
492+
provider: "xai".to_string(),
493+
input_tokens: delta,
494+
output_tokens: 0,
495+
cache_read_tokens: 0,
496+
cache_write_tokens: 0,
497+
reasoning_tokens: 0,
498+
token_source: TokenSource::Derived,
499+
parser_version: USAGE_PARSER_VERSION,
500+
source_path: None,
501+
raw_usage_json: Some(format!("{{\"totalTokens\":{}}}", prompt.peak_total_tokens)),
502+
});
503+
}
504+
events
381505
}
382506

383507
fn agent_chunk_key(params: &Value) -> AgentChunkKey {
@@ -584,18 +708,40 @@ mod tests {
584708
#[test]
585709
fn parse_grok_updates_extracts_messages() {
586710
let jsonl = r#"{"timestamp":10,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"hello"}},"_meta":{"totalTokens":100,"turnStartMs":1000}}}
587-
{"timestamp":11,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hi"}},"_meta":{"totalTokens":250,"promptId":"prompt-1","turnStartMs":1000}}}
711+
{"timestamp":11,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"hi"}},"_meta":{"totalTokens":250,"promptId":"prompt-1","modelId":"grok-composer-2.5-fast","turnStartMs":1000}}}
588712
{"timestamp":12,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":" there"}},"_meta":{"totalTokens":255,"promptId":"prompt-1","turnStartMs":1000}}}
589713
{"timestamp":20,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"next"}},"_meta":{"totalTokens":260,"turnStartMs":2000}}}
590-
{"timestamp":21,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"done"}},"_meta":{"totalTokens":400,"turnStartMs":2000}}}
714+
{"timestamp":21,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"done"}},"_meta":{"totalTokens":400,"promptId":"prompt-2","turnStartMs":2000}}}
591715
"#;
592716

593-
let messages = parse_grok_updates_reader(Cursor::new(jsonl)).unwrap();
717+
let (messages, usage_events) = parse_grok_updates_reader(Cursor::new(jsonl)).unwrap();
594718

595719
assert_eq!(messages.len(), 4);
596720
assert_eq!(messages[0].role, Role::User);
597721
assert_eq!(messages[0].content, "hello");
598722
assert_eq!(messages[1].content, "hi there");
723+
assert_eq!(usage_events.len(), 2);
724+
assert_eq!(usage_events[0].event_key, "prompt:prompt-1");
725+
assert_eq!(usage_events[0].input_tokens, 255);
726+
assert_eq!(usage_events[0].model, "grok-composer-2.5-fast");
727+
assert_eq!(usage_events[0].timestamp, 12_000);
728+
assert_eq!(usage_events[1].event_key, "prompt:prompt-2");
729+
assert_eq!(usage_events[1].input_tokens, 145);
730+
}
731+
732+
#[test]
733+
fn parse_grok_updates_resyncs_after_compaction_shrink() {
734+
let jsonl = r#"{"timestamp":10,"method":"session/update","params":{"sessionId":"s","update":{"sessionUpdate":"agent_thought_chunk","content":{"type":"text","text":"a"}},"_meta":{"totalTokens":1000,"promptId":"prompt-a"}}}
735+
{"timestamp":20,"method":"session/update","params":{"sessionId":"s","update":{"sessionUpdate":"agent_thought_chunk","content":{"type":"text","text":"b"}},"_meta":{"totalTokens":800,"promptId":"prompt-b"}}}
736+
{"timestamp":30,"method":"session/update","params":{"sessionId":"s","update":{"sessionUpdate":"agent_thought_chunk","content":{"type":"text","text":"c"}},"_meta":{"totalTokens":900,"promptId":"prompt-c"}}}
737+
"#;
738+
739+
let (_, usage_events) = parse_grok_updates_reader(Cursor::new(jsonl)).unwrap();
740+
741+
assert_eq!(usage_events.len(), 2);
742+
assert_eq!(usage_events[0].event_key, "prompt:prompt-a");
743+
assert_eq!(usage_events[1].event_key, "prompt:prompt-c");
744+
assert_eq!(usage_events[1].input_tokens, 100);
599745
}
600746

601747
#[test]
@@ -607,7 +753,7 @@ mod tests {
607753
{"timestamp":21,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"second"}},"_meta":{"promptId":"prompt-2","turnStartMs":2000}}}
608754
"#;
609755

610-
let messages = parse_grok_updates_reader(Cursor::new(jsonl)).unwrap();
756+
let (messages, _) = parse_grok_updates_reader(Cursor::new(jsonl)).unwrap();
611757

612758
assert_eq!(messages.len(), 4);
613759
assert_eq!(messages[1].role, Role::Assistant);
@@ -721,10 +867,52 @@ mod tests {
721867
let store = setup_store();
722868
store.insert_session(&make_existing_session(session_id, mtime, 1)).unwrap();
723869

870+
let result = scan_for_sync_impl(&root, &store, None).unwrap();
871+
assert_eq!(result.sessions.len(), 1, "missing usage state must trigger a backfill parse");
872+
873+
store
874+
.persist_usage_events_for_existing_session(
875+
"grok",
876+
session_id,
877+
&[],
878+
USAGE_PARSER_VERSION,
879+
Some(mtime),
880+
)
881+
.unwrap();
882+
724883
let result = scan_for_sync_impl(&root, &store, None).unwrap();
725884
assert_eq!(result.sessions.len(), 0);
726885
assert_eq!(result.stats.skipped_sessions, 1);
727886

728887
let _ = fs::remove_dir_all(&root);
729888
}
889+
890+
#[test]
891+
fn parse_grok_session_attaches_usage_with_model_fallback() {
892+
let root = temp_grok_root("usage");
893+
let session_id = "019e9003-1ed9-70e3-803b-1e7f96a072eb";
894+
let updates_path = write_grok_session(
895+
&root,
896+
"%2Ftmp%2Fproject",
897+
session_id,
898+
r#"{"info":{"id":"019e9003-1ed9-70e3-803b-1e7f96a072eb","cwd":"/tmp/project"},"created_at":"2026-06-04T00:00:00Z","current_model_id":"grok-composer-2.5-fast"}"#,
899+
r#"{"timestamp":10,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"user_message_chunk","content":{"type":"text","text":"hey"}},"_meta":{"turnStartMs":1000}}}
900+
{"timestamp":11,"method":"session/update","params":{"sessionId":"019e9003-1ed9-70e3-803b-1e7f96a072eb","update":{"sessionUpdate":"agent_message_chunk","content":{"type":"text","text":"sure"}},"_meta":{"totalTokens":500,"promptId":"prompt-1","turnStartMs":1000}}}"#,
901+
);
902+
let mtime = file_scan::stat_mtime_ms(&updates_path).unwrap();
903+
let entry = FileScanEntry {
904+
session_id: session_id.to_string(),
905+
stat_target: updates_path.clone(),
906+
directory: None,
907+
};
908+
909+
let session = parse_grok_session_for_entry(&entry, mtime).unwrap().unwrap();
910+
911+
assert_eq!(session.usage_parser_version, Some(USAGE_PARSER_VERSION));
912+
assert_eq!(session.usage_events.len(), 1);
913+
assert_eq!(session.usage_events[0].model, "grok-composer-2.5-fast");
914+
assert_eq!(session.usage_events[0].source_path.as_deref(), updates_path.to_str());
915+
916+
let _ = fs::remove_dir_all(&root);
917+
}
730918
}

0 commit comments

Comments
 (0)