fix(stats): dedupe forked-session entries to stop double-counting

SessionManager.fork() and createBranchedSession() deep-copy a parent's
entries into the child JSONL — same entry_id, timestamp, model,
responseId, token counts, cost. Stats keyed uniqueness on
(session_file, entry_id), so both files contributed to request, token,
and cost totals.

insertMessageStats and insertUserMessageStats now skip rows whose
(entry_id, timestamp) already exists under a different session_file
(first-write-wins). A one-shot migration on initDb, gated by the
fork_dedupe_v1 meta sentinel, collapses pre-existing duplicate rows
in messages and user_messages by keeping the lowest-id row per
(entry_id, timestamp) group.

Fixes #3370
This commit is contained in:
roboomp
2026-06-24 14:14:25 +00:00
parent c053afa087
commit b5769a4e10
3 changed files with 378 additions and 7 deletions
+4
View File
@@ -2,6 +2,10 @@
## [Unreleased]
### Fixed
- Stats sync counted the same provider request multiple times when a forked or branched session file copied the parent's entries verbatim. Inserts now skip rows whose `(entry_id, timestamp)` already exists under a different `session_file`, and a one-shot migration on the next `omp stats` run collapses any pre-existing duplicates ([#3370](https://github.com/can1357/oh-my-pi/issues/3370)).
## [16.1.15] - 2026-06-22
### Added
+79 -7
View File
@@ -45,6 +45,7 @@ const USER_MESSAGES_BACKFILL_KEY = "user_messages_v6";
const USER_MESSAGE_LINKS_REPAIR_KEY = "user_message_links_v1";
const PRIORITY_PREMIUM_REQUESTS_BACKFILL_KEY = "premium_requests_priority_v1";
const AGENT_TYPE_BACKFILL_KEY = "agent_type_v1";
const FORK_DEDUPE_KEY = "fork_dedupe_v1";
function shouldResetBackfill(value: string | undefined): boolean {
return value !== BACKFILL_COMPLETE && value !== BACKFILL_PENDING;
}
@@ -218,6 +219,7 @@ export async function initDb(): Promise<Database> {
backfillPriorityPremiumRequests(db);
backfillAgentType(db);
backfillMissingCatalogCosts(db);
backfillForkDuplicates(db);
return db;
}
@@ -334,22 +336,34 @@ export function setFileOffset(sessionFile: string, offset: number, lastModified:
/**
* Insert message stats into the database.
*
* Forked / branched sessions (see `SessionManager.fork()` and
* `createBranchedSession()` in `@oh-my-pi/pi-coding-agent`) deep-copy a parent
* session's entries into a new JSONL — same `entry_id`, `timestamp`, `model`,
* `provider`, token counts, and `responseId`. The `UNIQUE(session_file,
* entry_id)` constraint alone keys each row by file, so without the guard
* below the same provider request would land twice and inflate every
* aggregate. The `WHERE NOT EXISTS` clause skips inserts whose
* `(entry_id, timestamp)` already exists under a different `session_file` —
* first-write-wins across the lineage. Same-file re-syncs still hit the
* `ON CONFLICT(session_file, entry_id)` upsert below so historical
* `premium_requests` fix-ups continue to work.
*/
export function insertMessageStats(stats: MessageStats[]): number {
if (!db || stats.length === 0) return 0;
// Use UPSERT so a re-sync can fix up `premium_requests` for rows persisted
// before priority service-tier traffic was counted as premium. The guard
// `WHERE messages.premium_requests < excluded.premium_requests` keeps every
// other column immutable and never demotes an existing count (e.g. when a
// later parse drops back to 0 for the same row).
const stmt = db.prepare(`
INSERT INTO messages (
session_file, entry_id, folder, model, provider, api, timestamp,
duration, ttft, stop_reason, error_message,
input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, total_tokens, premium_requests,
cost_input, cost_output, cost_cache_read, cost_cache_write, cost_total, agent_type
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
)
SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
WHERE NOT EXISTS (
SELECT 1 FROM messages
WHERE entry_id = ? AND timestamp = ? AND session_file <> ?
)
ON CONFLICT(session_file, entry_id) DO UPDATE SET
premium_requests = excluded.premium_requests
WHERE messages.premium_requests < excluded.premium_requests
@@ -383,6 +397,11 @@ export function insertMessageStats(stats: MessageStats[]): number {
cost.cacheWrite,
cost.total,
s.agentType,
// `WHERE NOT EXISTS` binds: skip when a different session_file
// already holds this (entry_id, timestamp).
s.entryId,
s.timestamp,
s.sessionFile,
);
if (result.changes > 0) inserted++;
}
@@ -893,6 +912,46 @@ function backfillAgentType(database: Database): void {
apply();
}
/**
* One-shot collapse of forked-session duplicates that landed under the old
* `UNIQUE(session_file, entry_id)`-only invariant. `SessionManager.fork()`
* and `createBranchedSession()` deep-copy a parent's entries into the new
* JSONL — same `entry_id`, `timestamp`, `model`, `responseId`, token counts,
* cost — and the previous insert path counted both files toward request /
* token / cost totals. The migration keeps the lowest-`id` row per
* `(entry_id, timestamp)` group (almost always the parent — sessions are
* filename-timestamped and sync processes them in name order, so the
* originating file lands first) and drops every other copy. Same fix on
* `user_messages` since forks copy user entries too. Idempotent and
* crash-safe: enrolled at module-load via the `meta` sentinel, marked
* COMPLETE inside the same transaction so an aborted run rolls back and
* retries on the next init.
*/
function backfillForkDuplicates(database: Database): void {
const row = database.prepare("SELECT value FROM meta WHERE key = ?").get(FORK_DEDUPE_KEY) as
| { value: string }
| undefined;
if (row?.value === BACKFILL_COMPLETE) return;
const markComplete = database.prepare("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)");
const apply = database.transaction(() => {
database.run(`
DELETE FROM messages
WHERE id NOT IN (
SELECT MIN(id) FROM messages GROUP BY entry_id, timestamp
)
`);
database.run(`
DELETE FROM user_messages
WHERE id NOT IN (
SELECT MIN(id) FROM user_messages GROUP BY entry_id, timestamp
)
`);
markComplete.run(FORK_DEDUPE_KEY, BACKFILL_COMPLETE);
});
apply();
}
/**
* One-shot wipe of `file_offsets` to force `parseSessionFile` to re-parse
* every session from byte zero. We don't touch `user_messages`; the parser
@@ -961,6 +1020,9 @@ export function markUserMessageLinksRepairComplete(): void {
/**
* Insert user-message stats. Idempotent via UNIQUE(session_file, entry_id).
* The `WHERE NOT EXISTS` clause matches {@link insertMessageStats}: forks
* copy user entries verbatim into the child JSONL, so the same
* `(entry_id, timestamp)` must not land twice across different session files.
*/
export function insertUserMessageStats(stats: UserMessageStats[]): number {
if (!db || stats.length === 0) return 0;
@@ -970,7 +1032,12 @@ export function insertUserMessageStats(stats: UserMessageStats[]): number {
session_file, entry_id, folder, timestamp, model, provider,
chars, words, yelling, profanity, anguish,
negation, repetition, blame
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
)
SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
WHERE NOT EXISTS (
SELECT 1 FROM user_messages
WHERE entry_id = ? AND timestamp = ? AND session_file <> ?
)
`);
let inserted = 0;
@@ -991,6 +1058,11 @@ export function insertUserMessageStats(stats: UserMessageStats[]): number {
s.negation,
s.repetition,
s.blame,
// `WHERE NOT EXISTS` binds: skip when a different session_file
// already holds this (entry_id, timestamp).
s.entryId,
s.timestamp,
s.sessionFile,
);
if (result.changes > 0) inserted++;
}
+295
View File
@@ -0,0 +1,295 @@
import { Database } from "bun:sqlite";
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs/promises";
import * as os from "node:os";
import * as path from "node:path";
import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator";
import { closeDb, getOverallStats, getRecentRequests, initDb, insertMessageStats } from "@oh-my-pi/omp-stats/db";
import type { MessageStats } from "@oh-my-pi/omp-stats/types";
import { getAgentDir, getSessionsDir, getStatsDbPath, setAgentDir, TempDir } from "@oh-my-pi/pi-utils";
const XDG_KEYS = ["XDG_DATA_HOME", "XDG_STATE_HOME", "XDG_CACHE_HOME"] as const;
const originalConfigDir = process.env.PI_CONFIG_DIR;
const originalAgentDir = getAgentDir();
const originalXdg: Record<string, string | undefined> = {};
let tempDir: TempDir | null = null;
beforeEach(() => {
tempDir = TempDir.createSync("@pi-stats-fork-dedup-");
for (const key of XDG_KEYS) {
originalXdg[key] = process.env[key];
delete process.env[key];
}
const configDir = path.relative(os.homedir(), tempDir.join("config"));
process.env.PI_CONFIG_DIR = configDir;
setAgentDir(path.join(os.homedir(), configDir, "agent"));
});
afterEach(() => {
closeDb();
if (originalConfigDir === undefined) {
delete process.env.PI_CONFIG_DIR;
} else {
process.env.PI_CONFIG_DIR = originalConfigDir;
}
for (const key of XDG_KEYS) {
const prior = originalXdg[key];
if (prior === undefined) delete process.env[key];
else process.env[key] = prior;
}
setAgentDir(originalAgentDir);
tempDir?.removeSync();
tempDir = null;
});
interface AssistantOptions {
entryId: string;
parentId?: string | null;
timestamp: string;
}
function buildUserEntry(entryId: string, timestamp: string, content: string) {
return {
type: "message",
id: entryId,
parentId: null,
timestamp,
message: { role: "user", content },
};
}
function buildAssistantEntry(opts: AssistantOptions) {
return {
type: "message",
id: opts.entryId,
parentId: opts.parentId ?? null,
timestamp: opts.timestamp,
message: {
role: "assistant",
content: [{ type: "text", text: "ok" }],
api: "openai-responses",
provider: "openai",
model: "gpt-5.4",
responseId: `resp-${opts.entryId}`,
usage: {
input: 100,
output: 50,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 150,
cost: { input: 0.001, output: 0.002, cacheRead: 0, cacheWrite: 0, total: 0.003 },
},
stopReason: "stop",
timestamp: Date.parse(opts.timestamp),
duration: 10,
ttft: 5,
},
};
}
async function writeSessionFile(
folderSlug: string,
fileName: string,
header: { id: string; cwd: string; parentSession?: string },
entries: unknown[],
): Promise<string> {
const sessionDir = path.join(getSessionsDir(), folderSlug);
await fs.mkdir(sessionDir, { recursive: true });
const sessionFile = path.join(sessionDir, fileName);
const headerEntry = {
type: "session",
version: 3,
id: header.id,
timestamp: new Date().toISOString(),
cwd: header.cwd,
...(header.parentSession ? { parentSession: header.parentSession } : {}),
};
const lines = [headerEntry, ...entries].map(entry => JSON.stringify(entry)).join("\n");
await Bun.write(sessionFile, `${lines}\n`);
return sessionFile;
}
function makeStat(sessionFile: string, entryId: string, timestamp: number, premium = 0): MessageStats {
return {
sessionFile,
entryId,
folder: "/tmp/fork-dedup",
model: "gpt-5.4",
provider: "openai",
api: "openai-responses",
timestamp,
duration: 10,
ttft: 5,
stopReason: "stop",
errorMessage: null,
usage: {
input: 100,
output: 50,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 150,
premiumRequests: premium,
cost: { input: 0.001, output: 0.002, cacheRead: 0, cacheWrite: 0, total: 0.003 },
},
agentType: "main",
};
}
describe("stats sync deduplicates forked-session entries", () => {
it("counts each provider request once even when a fork copied the entries", async () => {
const ts = new Date("2026-06-24T10:00:00.000Z").toISOString();
const userEntry = buildUserEntry("user01ab", ts, "hello");
const assistantEntry = buildAssistantEntry({ entryId: "asst01ab", parentId: "user01ab", timestamp: ts });
const parentFile = await writeSessionFile(
"--tmp--fork-dedup",
"01_parent.jsonl",
{ id: "parent00", cwd: "/tmp/project" },
[userEntry, assistantEntry],
);
// `SessionManager.createBranchedSession` / `forkFrom` deep-copy the
// parent's entries into the child file with `parentSession` in the
// header. Earlier stats sync keyed uniqueness on (session_file,
// entry_id), so both files contributed the same provider request to
// every aggregate.
await writeSessionFile(
"--tmp--fork-dedup",
"02_fork.jsonl",
{ id: "fork0000", cwd: "/tmp/project", parentSession: parentFile },
[userEntry, assistantEntry],
);
await syncAllSessions({ workers: 1 });
const assistantRequests = getRecentRequests(10).filter(r => r.entryId === "asst01ab");
expect(assistantRequests).toHaveLength(1);
expect(assistantRequests[0].sessionFile).toBe(parentFile);
const overall = getOverallStats();
expect(overall.totalRequests).toBe(1);
expect(overall.totalInputTokens).toBe(100);
expect(overall.totalOutputTokens).toBe(50);
expect(overall.totalCost).toBeCloseTo(0.003, 8);
});
it("admits new entries appended in the forked session", async () => {
const ts = new Date("2026-06-24T10:00:00.000Z").toISOString();
const userEntry = buildUserEntry("user01ab", ts, "hello");
const assistantEntry = buildAssistantEntry({ entryId: "asst01ab", parentId: "user01ab", timestamp: ts });
const parentFile = await writeSessionFile(
"--tmp--fork-dedup",
"01_parent.jsonl",
{ id: "parent00", cwd: "/tmp/project" },
[userEntry, assistantEntry],
);
// After fork, the child file appends a fresh user+assistant pair.
// Inherited entries must dedupe; new ones must still count.
const newTs = new Date("2026-06-24T10:05:00.000Z").toISOString();
const newUserEntry = buildUserEntry("user02cd", newTs, "follow-up");
const newAssistantEntry = buildAssistantEntry({
entryId: "asst02cd",
parentId: "user02cd",
timestamp: newTs,
});
await writeSessionFile(
"--tmp--fork-dedup",
"02_fork.jsonl",
{ id: "fork0000", cwd: "/tmp/project", parentSession: parentFile },
[userEntry, assistantEntry, newUserEntry, newAssistantEntry],
);
await syncAllSessions({ workers: 1 });
const overall = getOverallStats();
expect(overall.totalRequests).toBe(2);
expect(overall.totalInputTokens).toBe(200);
expect(overall.totalOutputTokens).toBe(100);
expect(overall.totalCost).toBeCloseTo(0.006, 8);
});
it("collapses pre-existing duplicate rows on init (one-shot migration)", async () => {
await initDb();
closeDb();
// Clear the fork-dedupe sentinel so the migration runs against the
// hand-inserted dupes, then plant two messages and two user_messages
// rows that share (entry_id, timestamp) under different session_files
// — exactly the shape `SessionManager.fork()` leaves behind.
const database = new Database(getStatsDbPath());
database.prepare("DELETE FROM meta WHERE key = ?").run("fork_dedupe_v1");
const ts = Date.now();
const insertMessage = database.prepare(`
INSERT INTO messages (
session_file, entry_id, folder, model, provider, api, timestamp,
duration, ttft, stop_reason, error_message,
input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, total_tokens, premium_requests,
cost_input, cost_output, cost_cache_read, cost_cache_write, cost_total, agent_type
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
const sharedArgs = [
"/tmp/fork-dedup/project",
"gpt-5.4",
"openai",
"openai-responses",
ts,
10,
5,
"stop",
null,
100,
50,
0,
0,
150,
0,
0.001,
0.002,
0,
0,
0.003,
"main",
];
insertMessage.run("/tmp/parent.jsonl", "asst01ab", ...sharedArgs);
insertMessage.run("/tmp/fork.jsonl", "asst01ab", ...sharedArgs);
const insertUser = database.prepare(`
INSERT INTO user_messages (
session_file, entry_id, folder, timestamp, model, provider,
chars, words, yelling, profanity, anguish,
negation, repetition, blame
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`);
const sharedUser = ["/tmp/fork-dedup/project", ts, "gpt-5.4", "openai", 5, 1, 0, 0, 0, 0, 0, 0];
insertUser.run("/tmp/parent.jsonl", "user01ab", ...sharedUser);
insertUser.run("/tmp/fork.jsonl", "user01ab", ...sharedUser);
database.close();
await initDb();
const messageRows = getRecentRequests(10).filter(r => r.entryId === "asst01ab");
expect(messageRows).toHaveLength(1);
expect(messageRows[0].sessionFile).toBe("/tmp/parent.jsonl");
const overall = getOverallStats();
expect(overall.totalRequests).toBe(1);
expect(overall.totalCost).toBeCloseTo(0.003, 8);
// The migration is idempotent. Re-running init must not delete the
// surviving row.
closeDb();
await initDb();
expect(getOverallStats().totalRequests).toBe(1);
});
it("still upserts premium_requests for re-syncs of the same session file", async () => {
await initDb();
const stat = makeStat("/tmp/session-a.jsonl", "asst01ab", Date.now(), 0);
insertMessageStats([stat]);
const upgraded = { ...stat, usage: { ...stat.usage, premiumRequests: 1 } };
insertMessageStats([upgraded]);
const requests = getRecentRequests(10).filter(r => r.entryId === "asst01ab");
expect(requests).toHaveLength(1);
expect(requests[0].usage.premiumRequests).toBe(1);
});
});