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.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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<Usage> | 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,
|
||||
|
||||
@@ -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, unknown>): 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<string> {
|
||||
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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user