From 1b0b18c7a6846e618a4b7dcdbaa2ec34ce06a335 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 12 Jul 2026 12:42:32 +0200 Subject: [PATCH] fix(stats): handled malformed session entries to prevent sync crashes - Coerce missing `stopReason`, token counts, and timestamps at the parser boundary to satisfy database NOT NULL constraints. - Skip assistant entries missing essential model attribution or usage data instead of aborting the sync process. - Filter invalid tool call blocks that lack required identifiers to ensure successful database insertion. --- packages/stats/CHANGELOG.md | 4 + packages/stats/src/parser.ts | 58 ++++++++- .../test/parser-malformed-entries.test.ts | 119 ++++++++++++++++++ 3 files changed, 175 insertions(+), 6 deletions(-) create mode 100644 packages/stats/test/parser-malformed-entries.test.ts diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index 5897b909c..76f3de847 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed a `SQLITE_CONSTRAINT_NOTNULL` crash (`messages.stop_reason`) aborting the entire session sync when a persisted assistant message lacks a `stopReason`. Malformed entries — missing stop reason, token counts, or message timestamp — are now coerced at the parser boundary, and entries with no usage or model attribution are skipped instead of failing the batch insert. + ## [16.4.2] - 2026-07-10 ### Fixed diff --git a/packages/stats/src/parser.ts b/packages/stats/src/parser.ts index ce590119f..897de549f 100644 --- a/packages/stats/src/parser.ts +++ b/packages/stats/src/parser.ts @@ -6,7 +6,9 @@ import { getPriorityPremiumRequests, resolveModelServiceTier, type ServiceTierByFamily, + type ToolCall, type ToolResultMessage, + type Usage, } from "@oh-my-pi/pi-ai"; import { getSessionsDir, isEnoent, readLines } from "@oh-my-pi/pi-utils"; import type { @@ -142,6 +144,14 @@ function extractUserStats(sessionFile: string, folder: string, entry: SessionMes /** * Extract stats from an assistant message entry. + * + * Session JSONL on disk is not guaranteed to match the current + * `AssistantMessage` shape: crash-truncated turns, sessions written by older + * versions, and foreign producers all flow through this parser. Every field + * returned here feeds a NOT NULL column in stats.db, so malformed entries are + * coerced (missing `stopReason`, token counts, `timestamp`) or skipped + * (missing `model`/`provider`/`api`/`usage`) instead of crashing the whole + * sync with a constraint violation. */ function extractStats( sessionFile: string, @@ -152,6 +162,9 @@ function extractStats( ): MessageStats | null { const msg = entry.message as AssistantMessage; if (msg?.role !== "assistant") return null; + if (typeof msg.model !== "string" || typeof msg.provider !== "string" || typeof msg.api !== "string") return null; + const rawUsage = msg.usage as Partial | undefined; + if (!rawUsage || typeof rawUsage !== "object") return null; // Backfill: when the session recorded `priority` as the active service tier // at this point but the AI usage payload was captured before priority @@ -159,11 +172,29 @@ function extractStats( // "Premium Reqs" stat aggregates priority traffic on re-sync. Trust any // non-zero value already in `usage.premiumRequests` (Copilot multipliers or // the new AI code path) and only synthesise when the field is missing/zero. - const recorded = msg.usage.premiumRequests ?? 0; + const recorded = rawUsage.premiumRequests ?? 0; const model = { provider: msg.provider, api: msg.api, id: msg.model }; const tier = resolveModelServiceTier(currentServiceTier, model); const derived = recorded > 0 ? recorded : getPriorityPremiumRequests(tier, model); - const usage = derived === recorded ? msg.usage : { ...msg.usage, premiumRequests: derived }; + const wellFormed = + typeof rawUsage.input === "number" && + typeof rawUsage.output === "number" && + typeof rawUsage.cacheRead === "number" && + typeof rawUsage.cacheWrite === "number" && + typeof rawUsage.totalTokens === "number"; + const usage: Usage = + wellFormed && derived === recorded + ? (rawUsage as Usage) + : { + ...rawUsage, + input: rawUsage.input ?? 0, + output: rawUsage.output ?? 0, + cacheRead: rawUsage.cacheRead ?? 0, + cacheWrite: rawUsage.cacheWrite ?? 0, + totalTokens: rawUsage.totalTokens ?? 0, + cost: rawUsage.cost ?? { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + premiumRequests: derived, + }; return { sessionFile, @@ -172,16 +203,25 @@ function extractStats( model: msg.model, provider: msg.provider, api: msg.api, - timestamp: msg.timestamp, + timestamp: coerceEntryTimestamp(msg.timestamp, entry), duration: msg.duration ?? null, ttft: msg.ttft ?? null, - stopReason: msg.stopReason, + // A message persisted without a terminal stop reason never completed + // normally: classify by whether it carried an error. + stopReason: msg.stopReason ?? (msg.errorMessage ? "error" : "aborted"), errorMessage: msg.errorMessage ?? null, usage, agentType, }; } +/** Message timestamp, falling back to the entry's ISO timestamp, then 0. */ +function coerceEntryTimestamp(timestamp: number | undefined, entry: SessionMessageEntry): number { + if (typeof timestamp === "number" && Number.isFinite(timestamp)) return timestamp; + const ts = Date.parse(entry.timestamp); + return Number.isFinite(ts) ? ts : 0; +} + /** * Extract one {@link ToolCallStats} per `toolCall` content block of an * assistant message. Returns an empty array for turns without tool calls. @@ -194,8 +234,14 @@ function extractToolCalls( ): ToolCallStats[] { const msg = entry.message as AssistantMessage; if (msg?.role !== "assistant" || !Array.isArray(msg.content)) return []; + // `tool_calls` columns are NOT NULL: skip turns that can't be attributed + // (malformed persisted entries — see extractStats) and blocks missing ids. + if (typeof msg.model !== "string" || typeof msg.provider !== "string") return []; - const blocks = msg.content.filter(block => block.type === "toolCall"); + const blocks = msg.content.filter( + (block): block is ToolCall => + block.type === "toolCall" && typeof block.id === "string" && typeof block.name === "string", + ); if (blocks.length === 0) return []; return blocks.map(block => { @@ -213,7 +259,7 @@ function extractToolCalls( toolName: block.name, model: msg.model, provider: msg.provider, - timestamp: msg.timestamp, + timestamp: coerceEntryTimestamp(msg.timestamp, entry), agentType, callsInTurn: blocks.length, argsChars, diff --git a/packages/stats/test/parser-malformed-entries.test.ts b/packages/stats/test/parser-malformed-entries.test.ts new file mode 100644 index 000000000..ee18a19ef --- /dev/null +++ b/packages/stats/test/parser-malformed-entries.test.ts @@ -0,0 +1,119 @@ +import { describe, expect, it } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as path from "node:path"; +import { initDb, insertMessageStats, insertToolCalls } from "@oh-my-pi/omp-stats/db"; +import { parseSessionFile } from "@oh-my-pi/omp-stats/parser"; +import { getSessionsDir } from "@oh-my-pi/pi-utils"; +import { installStatsTestIsolation } from "./helpers/temp-agent"; + +installStatsTestIsolation("@pi-stats-malformed-"); + +const USAGE = { + input: 10, + output: 20, + cacheRead: 0, + cacheWrite: 0, + totalTokens: 30, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, +}; + +function assistantEntry(id: string, message: Record): string { + return JSON.stringify({ + type: "message", + id, + timestamp: "2026-07-12T00:00:00.000Z", + message: { + role: "assistant", + api: "anthropic-messages", + provider: "anthropic", + model: "claude-fable-5", + ...message, + }, + }); +} + +async function writeSession(lines: string[]): Promise { + const dir = path.join(getSessionsDir(), "--tmp--malformed"); + await fs.mkdir(dir, { recursive: true }); + const file = path.join(dir, "session.jsonl"); + await Bun.write(file, `${lines.join("\n")}\n`); + return file; +} + +// Regression: a single persisted assistant message missing `stopReason` (or +// usage/token fields) used to bind NULL into stats.db's NOT NULL columns and +// crash the entire sync with SQLITE_CONSTRAINT_NOTNULL. The parser must +// coerce or skip malformed entries so the batch always inserts. +describe("malformed session entries", () => { + it("coerces a missing stopReason instead of failing the NOT NULL insert", async () => { + const file = await writeSession([ + assistantEntry("a1", { content: [{ type: "text", text: "hi" }], usage: USAGE, timestamp: 1752000000000 }), + assistantEntry("a2", { + content: [], + usage: USAGE, + timestamp: 1752000001000, + errorMessage: "boom", + }), + ]); + + const result = await parseSessionFile(file); + expect(result.stats.map(s => s.stopReason)).toEqual(["aborted", "error"]); + + await initDb(); + expect(insertMessageStats(result.stats)).toBe(2); + }); + + it("zero-fills missing token counts and falls back to the entry timestamp", async () => { + const file = await writeSession([ + assistantEntry("a1", { + content: [], + stopReason: "stop", + usage: { cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 } }, + }), + ]); + + const result = await parseSessionFile(file); + expect(result.stats).toHaveLength(1); + const stats = result.stats[0]; + expect(stats.usage.totalTokens).toBe(0); + expect(stats.timestamp).toBe(Date.parse("2026-07-12T00:00:00.000Z")); + + await initDb(); + expect(insertMessageStats(result.stats)).toBe(1); + }); + + it("skips assistant entries with no usage or model attribution", async () => { + const file = await writeSession([ + assistantEntry("a1", { content: [], stopReason: "stop" }), + JSON.stringify({ + type: "message", + id: "a2", + timestamp: "2026-07-12T00:00:00.000Z", + message: { role: "assistant", content: [], stopReason: "stop", usage: USAGE }, + }), + assistantEntry("ok", { content: [], stopReason: "stop", usage: USAGE, timestamp: 1752000002000 }), + ]); + + const result = await parseSessionFile(file); + expect(result.stats.map(s => s.entryId)).toEqual(["ok"]); + }); + + it("keeps tool_calls insertable when the turn lacks a message timestamp", async () => { + const file = await writeSession([ + assistantEntry("a1", { + content: [ + { type: "toolCall", id: "call-1", name: "bash", arguments: { command: "ls" } }, + { type: "toolCall", name: "broken" }, // no id: unattributable, must be skipped + ], + usage: USAGE, + }), + ]); + + const result = await parseSessionFile(file); + expect(result.toolCalls.map(c => c.toolCallId)).toEqual(["call-1"]); + expect(result.toolCalls[0].timestamp).toBe(Date.parse("2026-07-12T00:00:00.000Z")); + + await initDb(); + expect(insertToolCalls(result.toolCalls)).toBe(1); + }); +});