mirror of
https://github.com/badlogic/pi-mono.git
synced 2026-08-19 05:33:45 +00:00
305 lines
10 KiB
TypeScript
305 lines
10 KiB
TypeScript
import { access, mkdtemp, rm, writeFile } from "node:fs/promises";
|
|
import { tmpdir } from "node:os";
|
|
import { join } from "node:path";
|
|
import * as storedValues from "@earendil-works/pi-agent-core";
|
|
import * as sessionWrites from "@earendil-works/pi-agent-core";
|
|
import { BACKGROUND_CONTEXT } from "@earendil-works/pi-agent-core";
|
|
import { describe, expect, it } from "vitest";
|
|
import type { SqliteDatabase } from "../src/index.ts";
|
|
import { createNodeSqliteFactory, SqliteSessionRepo, sql } from "../src/index.ts";
|
|
|
|
async function withTempDir<T>(run: (directory: string) => Promise<T>): Promise<T> {
|
|
const directory = await mkdtemp(join(tmpdir(), "pi-sqlite-session-"));
|
|
try {
|
|
return await run(directory);
|
|
} finally {
|
|
await rm(directory, { recursive: true, force: true });
|
|
}
|
|
}
|
|
|
|
async function withDb<T>(path: string, run: (db: SqliteDatabase) => Promise<T> | T): Promise<T> {
|
|
const db = await createNodeSqliteFactory().open(path);
|
|
try {
|
|
return await run(db);
|
|
} finally {
|
|
db.close();
|
|
}
|
|
}
|
|
|
|
describe("SqliteSessionRepo", () => {
|
|
it("creates one initialized session file with main lane values", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1_700_000_000_000,
|
|
});
|
|
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
const metadata = session.metadata;
|
|
expect(metadata).toMatchObject({
|
|
id: "session",
|
|
createdAt: 1_700_000_000_000,
|
|
storageVersion: 1,
|
|
});
|
|
expect(metadata.path).toBe(join(directory, "session.sqlite"));
|
|
|
|
await withDb(metadata.path, (db) => {
|
|
expect(sql`SELECT COUNT(*) AS count FROM session`.get<{ count: number }>(db)).toEqual({ count: 1 });
|
|
expect(sql`SELECT message_count, usage_payload, next_seq FROM session`.get(db)).toEqual({
|
|
message_count: 0,
|
|
usage_payload: JSON.stringify({
|
|
input: 0,
|
|
output: 0,
|
|
cacheRead: 0,
|
|
cacheWrite: 0,
|
|
totalTokens: 0,
|
|
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
}),
|
|
next_seq: 3,
|
|
});
|
|
expect(sql`SELECT namespace, key, seq, value FROM scalar_values ORDER BY seq`.all(db)).toEqual([
|
|
{
|
|
namespace: storedValues.laneLeaf("main").namespace,
|
|
key: "main",
|
|
seq: 1,
|
|
value: "null",
|
|
},
|
|
{
|
|
namespace: storedValues.laneState("main").namespace,
|
|
key: "main",
|
|
seq: 2,
|
|
value: JSON.stringify({ currentOperationId: null, pendingNextRun: [] }),
|
|
},
|
|
]);
|
|
expect(sql`SELECT COUNT(*) AS count FROM list_values`.get(db)).toEqual({ count: 0 });
|
|
});
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
});
|
|
});
|
|
|
|
it("exposes explicit branch scans through the open-session facade", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1_700_000_000_000,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
await session.mutate(
|
|
"main",
|
|
(mutator) =>
|
|
mutator.commit(
|
|
[
|
|
sessionWrites.insertEntry({ id: "root", parentId: null, type: "custom", customType: "root" }),
|
|
sessionWrites.insertEntry({ id: "child", parentId: "root", type: "custom", customType: "child" }),
|
|
],
|
|
BACKGROUND_CONTEXT,
|
|
),
|
|
BACKGROUND_CONTEXT,
|
|
);
|
|
|
|
expect(await session.scanBranch({ start: "child", order: "oldestFirst" }, BACKGROUND_CONTEXT)).toMatchObject([
|
|
{ id: "root" },
|
|
{ id: "child" },
|
|
]);
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
await expect(session.scanBranch({ start: "child" }, BACKGROUND_CONTEXT)).rejects.toThrow("Session is closed");
|
|
});
|
|
});
|
|
|
|
it("rejects duplicate create without deleting the existing database", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
const { metadata } = session;
|
|
|
|
await expect(repo.create({ id: "session" }, BACKGROUND_CONTEXT)).rejects.toThrow();
|
|
await withDb(metadata.path, (db) => {
|
|
expect(sql`SELECT COUNT(*) AS count FROM session`.get<{ count: number }>(db)).toEqual({ count: 1 });
|
|
});
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
});
|
|
});
|
|
|
|
it("lists sessions without changing the writer lease", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
const { metadata } = session;
|
|
let leaseBeforeList: unknown[] = [];
|
|
await withDb(metadata.path, (db) => {
|
|
leaseBeforeList = sql`SELECT owner_id, fence FROM writer_lease`.all(db);
|
|
});
|
|
|
|
await expect(repo.list(undefined, BACKGROUND_CONTEXT)).resolves.toMatchObject([
|
|
{ id: "session", path: metadata.path },
|
|
]);
|
|
await withDb(metadata.path, (db) => {
|
|
expect(sql`SELECT owner_id, fence FROM writer_lease`.all(db)).toEqual(leaseBeforeList);
|
|
});
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
});
|
|
});
|
|
|
|
it("skips corrupt and incompatible files during list discovery", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
await writeFile(join(directory, "corrupt.sqlite"), "not a sqlite database");
|
|
await withDb(session.metadata.path, (db) => {
|
|
sql`UPDATE session SET storage_version = ${999}`.run(db);
|
|
});
|
|
|
|
expect(await repo.list(undefined, BACKGROUND_CONTEXT)).toEqual([]);
|
|
});
|
|
});
|
|
|
|
it("does not remove a pre-existing non-database file when create fails", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const path = join(directory, "session.sqlite");
|
|
await writeFile(path, "not a sqlite database");
|
|
|
|
await expect(repo.create({ id: "session" }, BACKGROUND_CONTEXT)).rejects.toThrow();
|
|
await expect(access(path)).resolves.toBeUndefined();
|
|
});
|
|
});
|
|
|
|
it("rejects delete for missing files and live external writer leases", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
const { metadata } = session;
|
|
await session.close(BACKGROUND_CONTEXT);
|
|
|
|
await withDb(metadata.path, (db) => {
|
|
sql`INSERT INTO writer_lease (owner_id, fence, expires_at_ms) VALUES (${"external"}, ${1}, ${1_000})`.run(
|
|
db,
|
|
);
|
|
});
|
|
await expect(repo.delete(metadata, BACKGROUND_CONTEXT)).rejects.toThrow("already claimed");
|
|
await withDb(metadata.path, (db) => {
|
|
sql`DELETE FROM writer_lease`.run(db);
|
|
});
|
|
|
|
await expect(
|
|
repo.delete({ ...metadata, path: join(directory, "missing.sqlite") }, BACKGROUND_CONTEXT),
|
|
).rejects.toThrow();
|
|
await repo.delete(metadata, BACKGROUND_CONTEXT);
|
|
await expect(access(metadata.path)).rejects.toThrow();
|
|
});
|
|
});
|
|
|
|
it("closes open sessions through repo close and rejects later operations", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const session = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
|
|
await repo.close(BACKGROUND_CONTEXT);
|
|
|
|
await expect(session.getStats(BACKGROUND_CONTEXT)).rejects.toThrow("Session is closed");
|
|
await expect(repo.list(undefined, BACKGROUND_CONTEXT)).rejects.toThrow("SqliteSessionRepo is closed");
|
|
await expect(repo.create({ id: "other" }, BACKGROUND_CONTEXT)).rejects.toThrow("SqliteSessionRepo is closed");
|
|
await expect(repo.close(BACKGROUND_CONTEXT)).resolves.toBeUndefined();
|
|
});
|
|
});
|
|
|
|
it("opens a session through the version gate and rejects a live external writer lease", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1,
|
|
});
|
|
const created = await repo.create({ id: "session" }, BACKGROUND_CONTEXT);
|
|
const { metadata } = created;
|
|
await created.close(BACKGROUND_CONTEXT);
|
|
|
|
const opened = await repo.open(metadata, BACKGROUND_CONTEXT);
|
|
expect(opened.metadata).toMatchObject({ id: "session", path: metadata.path });
|
|
await opened.close(BACKGROUND_CONTEXT);
|
|
|
|
await withDb(metadata.path, (db) => {
|
|
sql`INSERT INTO writer_lease (owner_id, fence, expires_at_ms) VALUES (${"external"}, ${1}, ${1_000})`.run(
|
|
db,
|
|
);
|
|
});
|
|
await expect(repo.open(metadata, BACKGROUND_CONTEXT)).rejects.toThrow("already claimed");
|
|
});
|
|
});
|
|
|
|
it("does not copy usage ledger rows when forking", async () => {
|
|
await withTempDir(async (directory) => {
|
|
const repo = new SqliteSessionRepo({
|
|
directory,
|
|
databaseFactory: createNodeSqliteFactory(),
|
|
now: () => 1_700_000_000_000,
|
|
});
|
|
const source = await repo.create({ id: "source" }, BACKGROUND_CONTEXT);
|
|
await source.mutate(
|
|
"main",
|
|
(mutator) =>
|
|
mutator.commit(
|
|
[
|
|
sessionWrites.insertEntry({ id: "root", parentId: null, type: "custom", customType: "root" }),
|
|
sessionWrites.insertEntry({
|
|
id: "child",
|
|
parentId: "root",
|
|
type: "message",
|
|
message: { role: "user", content: "child", timestamp: 1 },
|
|
}),
|
|
storedValues.setValue(storedValues.laneLeaf("main"), "child"),
|
|
sessionWrites.insertUsage({
|
|
id: "usage",
|
|
adjustment: false,
|
|
usage: {
|
|
input: 1,
|
|
output: 1,
|
|
cacheRead: 0,
|
|
cacheWrite: 0,
|
|
totalTokens: 2,
|
|
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
},
|
|
}),
|
|
],
|
|
BACKGROUND_CONTEXT,
|
|
),
|
|
BACKGROUND_CONTEXT,
|
|
);
|
|
|
|
await source.close(BACKGROUND_CONTEXT);
|
|
const fork = await repo.fork(source.metadata, { id: "fork", entryId: "child" }, BACKGROUND_CONTEXT);
|
|
|
|
await withDb(fork.metadata.path, (db) => {
|
|
expect(sql`SELECT COUNT(*) AS count FROM usage_ledger`.get<{ count: number }>(db)).toEqual({ count: 0 });
|
|
});
|
|
await Promise.all([source.close(BACKGROUND_CONTEXT), fork.close(BACKGROUND_CONTEXT)]);
|
|
});
|
|
});
|
|
});
|