Merge PR #6400: fix(stats): prevent malformed entries from stalling sync (@roboomp)

This commit is contained in:
can1357
2026-07-23 22:15:24 +02:00
6 changed files with 72 additions and 27 deletions
+4
View File
@@ -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
+8 -4
View File
@@ -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;
+16 -21
View File
@@ -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();
}
/**
+8 -2
View File
@@ -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 {
@@ -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" },
]);
});
});
@@ -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", {