From bc5906aeb487a76dc49b7ef5c3709c6dfa1ac3ef Mon Sep 17 00:00:00 2001 From: roboomp Date: Thu, 23 Jul 2026 18:47:35 +0000 Subject: [PATCH] fix(stats): prevented malformed entries from stalling sync - Guarded malformed persisted content blocks so later project files continue ingesting. - Marked full-session migrations complete only after a successful sync pass. Fixes #6373 --- packages/stats/CHANGELOG.md | 4 ++ packages/stats/src/aggregator.ts | 12 ++++-- packages/stats/src/db.ts | 37 ++++++++----------- packages/stats/src/parser.ts | 10 ++++- packages/stats/test/behavior-backfill.test.ts | 21 +++++++++++ .../test/parser-malformed-entries.test.ts | 15 ++++++++ 6 files changed, 72 insertions(+), 27 deletions(-) diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index ec8e4f048..b5f132ea6 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Fixed malformed persisted content blocks aborting stats ingestion before later projects and settled pending full-session migrations after successful backfills ([#6373](https://github.com/can1357/oh-my-pi/issues/6373)). + ## [17.0.6] - 2026-07-20 ### Changed diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index fb37e9305..50fc25561 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -24,6 +24,7 @@ import { insertMessageStats, insertToolCalls, insertUserMessageStats, + markSessionBackfillsComplete, setFileOffset, updateToolResults, updateUserMessageLinks, @@ -209,12 +210,15 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: await initDb(); const files = await listAllSessionFiles(); - if (files.length === 0) return { processed: 0, files: 0 }; - let totalProcessed = 0; let filesProcessed = 0; let completed = 0; let cursor = 0; + const finish = () => { + markSessionBackfillsComplete(); + return { processed: totalProcessed, files: filesProcessed }; + }; + if (files.length === 0) return finish(); const report = (sessionFile: string) => { completed++; @@ -259,7 +263,7 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: for (const sessionFile of files) { await processFile(sessionFile, parseSessionFile); } - return { processed: totalProcessed, files: filesProcessed }; + return finish(); } const poolSize = Math.min(files.length, requestedWorkers); @@ -282,7 +286,7 @@ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: for (const handle of handles) handle.worker.terminate(); } - return { processed: totalProcessed, files: filesProcessed }; + return finish(); } const HOUR_MS = 60 * 60 * 1000; diff --git a/packages/stats/src/db.ts b/packages/stats/src/db.ts index 7f24d1e72..085313545 100644 --- a/packages/stats/src/db.ts +++ b/packages/stats/src/db.ts @@ -1079,28 +1079,23 @@ function backfillPriorityPremiumRequests(database: Database): void { .run(PRIORITY_PREMIUM_REQUESTS_BACKFILL_KEY, BACKFILL_PENDING); } -export function markPriorityPremiumRequestsBackfillComplete(): void { +/** + * Settle every full-session backfill after a successful sync pass. + */ +export function markSessionBackfillsComplete(): void { if (!db) return; - db.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)").run( - PRIORITY_PREMIUM_REQUESTS_BACKFILL_KEY, - BACKFILL_COMPLETE, - ); -} - -export function markUserMessagesBackfillComplete(): void { - if (!db) return; - db.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)").run( - USER_MESSAGES_BACKFILL_KEY, - BACKFILL_COMPLETE, - ); -} - -export function markUserMessageLinksRepairComplete(): void { - if (!db) return; - db.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)").run( - USER_MESSAGE_LINKS_REPAIR_KEY, - BACKFILL_COMPLETE, - ); + const markComplete = db.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)"); + const apply = db.transaction(() => { + for (const key of [ + USER_MESSAGES_BACKFILL_KEY, + TOOL_CALLS_BACKFILL_KEY, + USER_MESSAGE_LINKS_REPAIR_KEY, + PRIORITY_PREMIUM_REQUESTS_BACKFILL_KEY, + ]) { + markComplete.run(key, BACKFILL_COMPLETE); + } + }); + apply(); } /** diff --git a/packages/stats/src/parser.ts b/packages/stats/src/parser.ts index 897de549f..b12725596 100644 --- a/packages/stats/src/parser.ts +++ b/packages/stats/src/parser.ts @@ -240,7 +240,11 @@ function extractToolCalls( const blocks = msg.content.filter( (block): block is ToolCall => - block.type === "toolCall" && typeof block.id === "string" && typeof block.name === "string", + block !== null && + typeof block === "object" && + block.type === "toolCall" && + typeof block.id === "string" && + typeof block.name === "string", ); if (blocks.length === 0) return []; @@ -277,7 +281,9 @@ function extractToolResultLink(sessionFile: string, entry: SessionMessageEntry): let resultChars = 0; if (Array.isArray(msg.content)) { for (const block of msg.content) { - if (block.type === "text" && typeof block.text === "string") resultChars += block.text.length; + if (block && typeof block === "object" && block.type === "text" && typeof block.text === "string") { + resultChars += block.text.length; + } } } return { diff --git a/packages/stats/test/behavior-backfill.test.ts b/packages/stats/test/behavior-backfill.test.ts index 135c3a214..055e7a201 100644 --- a/packages/stats/test/behavior-backfill.test.ts +++ b/packages/stats/test/behavior-backfill.test.ts @@ -95,4 +95,25 @@ describe("behavior backfill", () => { expect(getBehaviorOverall(null).totalMessages).toBe(1); expect(getFileOffset(sessionFile)).not.toBeNull(); }); + + it("marks full-session backfills complete after a successful sync", async () => { + await writeSessionFile(); + await syncAllSessions({ workers: 1 }); + closeDb(); + + const database = new Database(getStatsDbPath(), { readonly: true }); + const rows = database + .query( + "SELECT key, value FROM meta WHERE key IN ('user_messages_v8', 'tool_calls_v1', 'user_message_links_v1', 'premium_requests_priority_v1') ORDER BY key", + ) + .all() as { key: string; value: string }[]; + database.close(); + + expect(rows).toEqual([ + { key: "premium_requests_priority_v1", value: "complete" }, + { key: "tool_calls_v1", value: "complete" }, + { key: "user_message_links_v1", value: "complete" }, + { key: "user_messages_v8", value: "complete" }, + ]); + }); }); diff --git a/packages/stats/test/parser-malformed-entries.test.ts b/packages/stats/test/parser-malformed-entries.test.ts index ee18a19ef..033ac9558 100644 --- a/packages/stats/test/parser-malformed-entries.test.ts +++ b/packages/stats/test/parser-malformed-entries.test.ts @@ -98,6 +98,21 @@ describe("malformed session entries", () => { expect(result.stats.map(s => s.entryId)).toEqual(["ok"]); }); + it("ignores malformed content blocks without aborting later entries", async () => { + const file = await writeSession([ + assistantEntry("a1", { + content: [null, { type: "toolCall", id: "call-1", name: "bash", arguments: {} }], + usage: USAGE, + timestamp: 1752000000000, + }), + assistantEntry("a2", { content: [], usage: USAGE, timestamp: 1752000001000 }), + ]); + + const result = await parseSessionFile(file); + expect(result.stats.map(s => s.entryId)).toEqual(["a1", "a2"]); + expect(result.toolCalls.map(c => c.toolCallId)).toEqual(["call-1"]); + }); + it("keeps tool_calls insertable when the turn lacks a message timestamp", async () => { const file = await writeSession([ assistantEntry("a1", {