Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 51 additions & 4 deletions src/lib/opencode-storage.ts
Original file line number Diff line number Diff line change
Expand Up @@ -310,6 +310,50 @@ function buildCompletedAssistantQuery(params: {
};
}

function forkCopyFingerprint(message: OpenCodeMessage, completed: number): string {
const tokens = message.tokens;
return JSON.stringify([
message.time?.created,
completed,
message.providerID,
message.modelID,
tokens?.input,
tokens?.output,
tokens?.reasoning,
tokens?.cache?.read,
tokens?.cache?.write,
message.cost,
]);
}

/**
* OpenCode copies a parent's finished messages into a fork with new ids but
* identical times, model, tokens, and cost. Count each finished fingerprint as
* many times as the session that has it most often, so copies are not added on
* top of the original. Unfinished rows are always kept.
*/
function dropForkCopies(messages: OpenCodeMessage[]): OpenCodeMessage[] {
const countBySessionAndFingerprint = new Map<string, number>();
const keptCountByFingerprint = new Map<string, number>();
const out: OpenCodeMessage[] = [];
for (const message of messages) {
const completed = completedAt(message);
if (completed === null) {
out.push(message);
continue;
}
const fingerprint = forkCopyFingerprint(message, completed);
const sessionKey = JSON.stringify([message.sessionID, fingerprint]);
const sessionCount = (countBySessionAndFingerprint.get(sessionKey) ?? 0) + 1;
countBySessionAndFingerprint.set(sessionKey, sessionCount);
const keptCount = keptCountByFingerprint.get(fingerprint) ?? 0;
if (sessionCount <= keptCount) continue;
keptCountByFingerprint.set(fingerprint, sessionCount);
out.push(message);
}
return out;
}

function compareMessageOrder(a: OpenCodeMessage, b: OpenCodeMessage): number {
const aCreated = typeof a.time?.created === "number" ? a.time.created : Number.MAX_SAFE_INTEGER;
const bCreated = typeof b.time?.created === "number" ? b.time.created : Number.MAX_SAFE_INTEGER;
Expand Down Expand Up @@ -369,7 +413,7 @@ export async function iterAssistantMessages(params: {
try {
const q = buildMessageQuery({ sinceMs: params.sinceMs, untilMs: params.untilMs });
const rows = conn.all<MessageRow>(q.sql, q.args);
return mapAssistantMessages(rows);
return dropForkCopies(mapAssistantMessages(rows));
} finally {
conn.close();
}
Expand All @@ -392,13 +436,15 @@ export async function iterCompletedAssistantMessages(params: {
try {
if (await hasJsonExtract(conn)) {
const query = buildCompletedAssistantQuery(params);
return mapCompletedAssistantMessages(conn.all<MessageRow>(query.sql, query.args));
return dropForkCopies(
mapCompletedAssistantMessages(conn.all<MessageRow>(query.sql, query.args)),
);
}

const rows = conn.all<MessageRow>(
`SELECT id, session_id, time_created, time_updated, data FROM "message"`,
);
return mapCompletedAssistantMessages(rows)
const messages = mapCompletedAssistantMessages(rows)
.filter((message) => {
const atMs = completedAt(message);
if (atMs === null) return false;
Expand All @@ -411,6 +457,7 @@ export async function iterCompletedAssistantMessages(params: {
return true;
})
.sort(compareCompletedMessageOrder);
return dropForkCopies(messages);
} finally {
conn.close();
}
Expand Down Expand Up @@ -483,7 +530,7 @@ export async function iterAssistantMessagesForSessions(params: {
}

messages.sort(compareMessageOrder);
return messages;
return dropForkCopies(messages);
} finally {
conn.close();
}
Expand Down
227 changes: 227 additions & 0 deletions tests/lib.opencode-storage.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,3 +142,230 @@ describe("opencode storage multi-session reads", () => {
expect(conn.close).toHaveBeenCalledTimes(1);
});
});

describe("opencode storage forked session copies", () => {
function messageRow(params: {
id: string;
sessionID: string;
created?: number;
completed?: number | null;
input?: number;
output?: number;
cost?: number;
}) {
const created = params.created ?? 100;
const completed = params.completed === undefined ? 150 : params.completed;
return {
id: params.id,
session_id: params.sessionID,
time_created: created,
data: JSON.stringify({
role: "assistant",
providerID: "openai",
modelID: "gpt-5",
tokens: {
input: params.input ?? 1000,
output: params.output ?? 500,
reasoning: 10,
cache: { read: 20, write: 30 },
},
cost: params.cost ?? 10,
time: completed === null ? { created } : { created, completed },
agent: "build",
}),
};
}

const originalRow = messageRow({ id: "msg_01original", sessionID: "ses_parent" });
// Message ids ascend with time. OpenCode copies finished parent rows into the
// fork with a newer id and identical data.
const forkCopyRow = messageRow({ id: "msg_05forkevent_1", sessionID: "ses_fork" });
const forkNewRow = messageRow({
id: "msg_06forknew",
sessionID: "ses_fork",
created: 200,
completed: 250,
input: 40,
cost: 1,
});

function mockConnection(
rowsForSql: (sql: string, params?: unknown[]) => unknown[],
options: { jsonExtract?: boolean } = {},
) {
const jsonExtractResult = options.jsonExtract === false ? null : "assistant";
sqliteMocks.openOpenCodeSqliteReadOnly.mockResolvedValue({
get: vi.fn((sql: string) =>
sql.includes('FROM "session"') ? { ok: 1 } : { r: jsonExtractResult },
),
all: vi.fn(rowsForSql),
close: vi.fn(),
});
}

beforeEach(() => {
vi.clearAllMocks();
vi.resetModules();
fsMocks.existsSync.mockReturnValue(true);
});

it("counts a forked session's copied message once in all-history reads", async () => {
mockConnection(() => [originalRow, forkCopyRow, forkNewRow]);

const { iterAssistantMessages, iterCompletedAssistantMessages } = await import(
"../src/lib/opencode-storage.js"
);

const messages = await iterAssistantMessages({});
expect(messages.map((message) => message.id)).toEqual(["msg_01original", "msg_06forknew"]);

const completed = await iterCompletedAssistantMessages({});
expect(completed.map((message) => message.id)).toEqual(["msg_01original", "msg_06forknew"]);
});

it("counts a forked session's copied message once in completed reads without json_extract", async () => {
mockConnection(() => [forkNewRow, forkCopyRow, originalRow], { jsonExtract: false });

const { iterCompletedAssistantMessages } = await import("../src/lib/opencode-storage.js");
const completed = await iterCompletedAssistantMessages({});

expect(completed.map((message) => message.id)).toEqual(["msg_01original", "msg_06forknew"]);
});

it("counts a forked session's copied message once across a session set", async () => {
mockConnection((_sql, params) => {
const sessionIDs = params ?? [];
return [originalRow, forkCopyRow, forkNewRow].filter((row) =>
sessionIDs.includes(row.session_id),
);
});

const { iterAssistantMessagesForSessions } = await import("../src/lib/opencode-storage.js");
const messages = await iterAssistantMessagesForSessions({
sessionIDs: ["ses_parent", "ses_fork"],
});

expect(messages.map((message) => message.id)).toEqual(["msg_01original", "msg_06forknew"]);
});

it("keeps a fork's copied history in its own single-session read", async () => {
mockConnection(() => [forkCopyRow, forkNewRow]);

const { iterAssistantMessagesForSession } = await import("../src/lib/opencode-storage.js");
const messages = await iterAssistantMessagesForSession({ sessionID: "ses_fork" });

expect(messages.map((message) => message.id)).toEqual(["msg_05forkevent_1", "msg_06forknew"]);
});

it("keeps finished messages that differ only in tokens and identical unfinished messages across sessions", async () => {
const otherTokensRow = messageRow({
id: "msg_other_tokens",
sessionID: "ses_fork",
output: 501,
});
const unfinishedRow = messageRow({
id: "msg_unfinished_a",
sessionID: "ses_parent",
created: 300,
completed: null,
});
const unfinishedTwinRow = { ...unfinishedRow, id: "msg_unfinished_b", session_id: "ses_fork" };
mockConnection(() => [originalRow, otherTokensRow, unfinishedRow, unfinishedTwinRow]);

const { iterAssistantMessages } = await import("../src/lib/opencode-storage.js");
const messages = await iterAssistantMessages({});

expect(messages.map((message) => message.id)).toEqual([
"msg_01original",
"msg_other_tokens",
"msg_unfinished_a",
"msg_unfinished_b",
]);
});

it("keeps identical finished messages that are in the same session", async () => {
const sameSessionTwinRow = { ...originalRow, id: "msg_02sametwin" };
mockConnection(() => [originalRow, sameSessionTwinRow]);

const { iterAssistantMessages, iterCompletedAssistantMessages } = await import(
"../src/lib/opencode-storage.js"
);

const messages = await iterAssistantMessages({});
expect(messages.map((message) => message.id)).toEqual(["msg_01original", "msg_02sametwin"]);

const completed = await iterCompletedAssistantMessages({});
expect(completed.map((message) => message.id)).toEqual(["msg_01original", "msg_02sametwin"]);
});

it("drops both fork copies when the parent has two identical messages", async () => {
const parentTwinRow = { ...originalRow, id: "msg_02parenttwin" };
const forkTwinCopyRow = { ...originalRow, id: "msg_05forkevent_2", session_id: "ses_fork" };
mockConnection(() => [originalRow, parentTwinRow, forkCopyRow, forkTwinCopyRow]);

const { iterAssistantMessages } = await import("../src/lib/opencode-storage.js");
const messages = await iterAssistantMessages({});

expect(messages.map((message) => message.id)).toEqual(["msg_01original", "msg_02parenttwin"]);
});

it("keeps two rows when an earlier fork copied one twin and a later fork copied both", async () => {
// The parent had two identical finished messages and is not part of the read.
const earlyForkCopyRow = {
...originalRow,
id: "msg_03earlyfork_1",
session_id: "ses_fork_early",
};
const lateForkCopyRow = { ...originalRow, id: "msg_07latefork_1", session_id: "ses_fork_late" };
const lateForkTwinCopyRow = {
...originalRow,
id: "msg_07latefork_2",
session_id: "ses_fork_late",
};
const rows = [earlyForkCopyRow, lateForkCopyRow, lateForkTwinCopyRow];
mockConnection((_sql, params) => {
const sessionIDs = params ?? [];
return rows.filter((row) => sessionIDs.includes(row.session_id));
});

const { iterAssistantMessagesForSessions } = await import("../src/lib/opencode-storage.js");
const sessionMessages = await iterAssistantMessagesForSessions({
sessionIDs: ["ses_fork_early", "ses_fork_late"],
});
expect(sessionMessages.map((message) => message.id)).toEqual([
"msg_03earlyfork_1",
"msg_07latefork_2",
]);

mockConnection(() => rows);
const { iterAssistantMessages } = await import("../src/lib/opencode-storage.js");
const allMessages = await iterAssistantMessages({});
expect(allMessages.map((message) => message.id)).toEqual([
"msg_03earlyfork_1",
"msg_07latefork_2",
]);
});

it("counts a fork copy once when the original and copy come from different query chunks", async () => {
const fillerSessionIDs = Array.from({ length: 899 }, (_, index) => `ses_filler_${index}`);
const queriedChunks: unknown[][] = [];
mockConnection((_sql, params) => {
const sessionIDs = params ?? [];
queriedChunks.push(sessionIDs);
return [originalRow, forkCopyRow, forkNewRow].filter((row) =>
sessionIDs.includes(row.session_id),
);
});

const { iterAssistantMessagesForSessions } = await import("../src/lib/opencode-storage.js");
const messages = await iterAssistantMessagesForSessions({
sessionIDs: ["ses_parent", ...fillerSessionIDs, "ses_fork"],
});

expect(queriedChunks).toHaveLength(2);
expect(queriedChunks[0]).toContain("ses_parent");
expect(queriedChunks[0]).not.toContain("ses_fork");
expect(queriedChunks[1]).toEqual(["ses_fork"]);
expect(messages.map((message) => message.id)).toEqual(["msg_01original", "msg_06forknew"]);
});
});
Loading