repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/good-job/__tests__/storage.test.ts
Rawimport {
chmodSync,
existsSync,
mkdtempSync,
readFileSync,
rmSync,
symlinkSync,
writeFileSync,
} from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, describe, expect, it } from "vitest";
import {
feedbackIdFromSeed,
type GoodJobJob,
type GoodJobRecord,
} from "../core.js";
import { openDirectory, openerFor } from "../implementation.js";
import {
atomicWriteJson,
claimJob,
commitFailure,
commitRecord,
configuredDataDirectory,
counts,
enqueueJob,
ensureLayout,
goodJobDataDirectory,
layout,
loadRecords,
migrateStore,
parseJob,
parseRecord,
recoverProcessing,
} from "../storage.js";
const roots: string[] = [];
function temporaryRoot(): string {
const root = mkdtempSync(join(tmpdir(), "good-job-test-"));
roots.push(root);
return root;
}
function job(id = "gj-0000000001"): GoodJobJob {
return {
version: 2,
id,
createdAt: "2026-01-01T00:00:00.000Z",
kind: "praise",
feedback: "gj",
source: {
type: "session",
cwd: "/repo",
sessionPath: "/sessions/one.jsonl",
sessionId: "session-1",
sessionName: "Shared route",
leafId: "leaf-1",
agentModel: { provider: "provider", id: "worker" },
},
analysisModel: { provider: "provider", id: "worker" },
};
}
function record(source = job()): GoodJobRecord {
if (!source.analysisModel) throw new Error("fixture analysis model missing");
return {
version: 2,
id: source.id,
createdAt: source.createdAt,
completedAt: "2026-01-01T00:01:00.000Z",
kind: source.kind,
feedback: source.feedback,
source: source.source,
analysisModel: source.analysisModel,
learning: {
slug: "fixed-shared-route",
summary: "Fixed the shared route",
evidence: ["Tests passed"],
behaviors: ["Inspected callers"],
candidateRule: "Fix shared routes first.",
},
usage: {
input: 10,
output: 5,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 15,
cost: 0.01,
},
};
}
function legacyRecord() {
return {
version: 1,
id: "b0b916cf-20ec-4f50-80f5-0cb83fbbd3c3",
createdAt: "2026-01-01T00:00:00.000Z",
completedAt: "2026-01-01T00:01:00.000Z",
praise: "gj",
source: {
cwd: "/repo",
sessionPath: "/sessions/one.jsonl",
sessionId: "session-1",
sessionName: "Shared route",
leafId: "leaf-1",
agentModel: { provider: "provider", id: "worker" },
},
analysisModel: { provider: "provider", id: "worker" },
learning: {
achievement: "Fixed the shared route",
evidence: ["Tests passed"],
successfulBehaviors: ["Inspected callers"],
candidateRule: "Fix shared routes first.",
},
usage: {
input: 10,
output: 5,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 15,
cost: 0.01,
},
};
}
afterEach(() => {
for (const root of roots.splice(0))
rmSync(root, { recursive: true, force: true });
});
describe("good-job storage", () => {
it("uses platform defaults and resolves configured directories", () => {
expect(
goodJobDataDirectory({ XDG_DATA_HOME: "/data" }, "linux", "/home/o"),
).toBe("/data/pi/good-job");
expect(
goodJobDataDirectory(
{ LOCALAPPDATA: "C:\\Users\\o\\AppData\\Local" },
"win32",
"C:\\Users\\o",
),
).toBe("C:\\Users\\o\\AppData\\Local\\pi\\good-job");
expect(goodJobDataDirectory({}, "darwin", "/Users/o")).toBe(
"/Users/o/Library/Application Support/pi/good-job",
);
expect(configuredDataDirectory("~/shared/gj", "/repo", "/home/o")).toBe(
"/home/o/shared/gj",
);
expect(configuredDataDirectory(".pi/gj", "/repo", "/home/o")).toBe(
"/repo/.pi/gj",
);
});
it("validates current records including nullable import metadata", () => {
const valid = record();
expect(parseRecord(valid)).toEqual(valid);
expect(parseRecord({ ...valid, usage: {} })).toBeUndefined();
expect(
parseRecord({ ...valid, usage: { ...valid.usage, input: "10" } }),
).toBeUndefined();
expect(
parseRecord({ ...valid, id: "too-long-and-not-valid" }),
).toBeUndefined();
expect(parseRecord({ ...valid, id: "wtf-0000000001" })).toBeUndefined();
expect(parseJob({ ...job(), id: "wtf-0000000001" })).toBeUndefined();
const imported: GoodJobRecord = {
...valid,
id: "wtf-0000000002",
kind: "problem",
source: {
type: "import",
importer: "wtf-v1",
fingerprint: "abc123",
importedAt: "2026-01-02T00:00:00.000Z",
},
analysisModel: null,
usage: null,
};
expect(parseRecord(imported)).toEqual(imported);
});
it("claims jobs atomically and commits immutable records", async () => {
const paths = layout(temporaryRoot());
await ensureLayout(paths);
const queued = job();
await enqueueJob(queued, paths);
expect(await counts(paths)).toMatchObject({ pending: 1, processing: 0 });
const claimed = await claimJob(queued.id, paths);
expect(claimed?.job).toEqual(queued);
expect(await claimJob(queued.id, paths)).toBeUndefined();
expect(await counts(paths)).toMatchObject({ pending: 0, processing: 1 });
if (!claimed) throw new Error("missing claimed job");
await commitRecord(record(queued), claimed.path, paths);
expect(await counts(paths)).toEqual({
pending: 0,
processing: 0,
failed: 0,
records: 1,
});
expect((await loadRecords(paths)).records[0]?.learning.slug).toBe(
"fixed-shared-route",
);
});
it("loads valid records concurrently while preserving sort and errors", async () => {
const paths = layout(temporaryRoot());
await ensureLayout(paths);
const older = record(job("gj-0000000001"));
const newer = record(job("gj-0000000002"));
older.completedAt = "2026-01-01T00:01:00.000Z";
newer.completedAt = "2026-01-02T00:01:00.000Z";
await Promise.all([
atomicWriteJson(join(paths.records, `${older.id}.json`), older),
atomicWriteJson(join(paths.records, `${newer.id}.json`), newer),
]);
writeFileSync(join(paths.records, "broken.json"), "not json", "utf8");
const loaded = await loadRecords(paths);
expect(loaded.records.map(({ id }) => id)).toEqual([
"gj-0000000002",
"gj-0000000001",
]);
expect(loaded.errors).toHaveLength(1);
expect(loaded.errors[0]).toContain("broken.json");
});
it("migrates version-1 records into a configured store at runtime", async () => {
const source = temporaryRoot();
const target = temporaryRoot();
const sourcePaths = layout(source);
await ensureLayout(sourcePaths);
const legacy = legacyRecord();
const oldPath = join(sourcePaths.records, `${legacy.id}.json`);
await atomicWriteJson(oldPath, legacy);
const migrated = await migrateStore(source, target);
expect(migrated).toEqual({ jobs: 0, records: 1, failures: 0, active: 0 });
expect(existsSync(oldPath)).toBe(false);
const loaded = await loadRecords(layout(target));
const id = feedbackIdFromSeed("praise", `good-job-v1:${legacy.id}`);
expect(loaded.records[0]).toMatchObject({
id,
kind: "praise",
feedback: "gj",
learning: {
slug: "fixed-the-shared-route",
summary: "Fixed the shared route",
},
});
expect(await migrateStore(source, target)).toEqual({
jobs: 0,
records: 0,
failures: 0,
active: 0,
});
});
it("preserves records when configured and default roots alias", async () => {
const source = temporaryRoot();
const holder = temporaryRoot();
const alias = join(holder, "alias");
const paths = layout(source);
await ensureLayout(paths);
const existing = record();
await atomicWriteJson(join(paths.records, `${existing.id}.json`), existing);
symlinkSync(source, alias, "dir");
expect(await migrateStore(source, alias)).toEqual({
jobs: 0,
records: 0,
failures: 0,
active: 0,
});
expect((await loadRecords(paths)).records).toEqual([existing]);
});
it("refuses migration collisions without changing either record", async () => {
const source = temporaryRoot();
const target = temporaryRoot();
const sourcePaths = layout(source);
const targetPaths = layout(target);
await Promise.all([ensureLayout(sourcePaths), ensureLayout(targetPaths)]);
const original = record();
const conflicting = {
...original,
learning: { ...original.learning, summary: "Conflicting learning" },
};
const sourceFile = join(sourcePaths.records, `${original.id}.json`);
const targetFile = join(targetPaths.records, `${original.id}.json`);
await atomicWriteJson(sourceFile, conflicting);
await atomicWriteJson(targetFile, original);
const originalBytes = readFileSync(targetFile, "utf8");
await expect(migrateStore(source, target)).rejects.toThrow(
"migration collision",
);
expect(readFileSync(targetFile, "utf8")).toBe(originalBytes);
expect(existsSync(sourceFile)).toBe(true);
});
it("relocates invalid-claim diagnostics created during recovery", async () => {
const source = temporaryRoot();
const target = temporaryRoot();
const paths = layout(source);
await ensureLayout(paths);
writeFileSync(join(paths.processing, "invalid.json"), "{}\n", "utf8");
expect(await migrateStore(source, target)).toEqual({
jobs: 0,
records: 0,
failures: 1,
active: 0,
});
const diagnostic = JSON.parse(
readFileSync(join(layout(target).failed, "invalid.json"), "utf8"),
);
expect(diagnostic).toMatchObject({
version: 1,
job: null,
error: "invalid processing claim name",
});
});
it("rejects a queued ID already used by a completed record", async () => {
const paths = layout(temporaryRoot());
await ensureLayout(paths);
const existing = record();
const target = join(paths.records, `${existing.id}.json`);
await atomicWriteJson(target, existing);
const originalBytes = readFileSync(target, "utf8");
await expect(enqueueJob(job(existing.id), paths)).rejects.toThrow(
`duplicate gj id ${existing.id}`,
);
expect(readFileSync(target, "utf8")).toBe(originalBytes);
});
it("persists failed analysis before releasing its processing claim", async () => {
const paths = layout(temporaryRoot());
await ensureLayout(paths);
await enqueueJob(job(), paths);
const claimed = await claimJob("gj-0000000001", paths);
if (!claimed) throw new Error("missing claimed job");
await commitFailure(
claimed.job,
new Error("provider unavailable"),
claimed.path,
paths,
);
expect(await counts(paths)).toMatchObject({
pending: 0,
processing: 0,
failed: 1,
});
});
it("recovers interrupted processing without duplicating completed work", async () => {
const paths = layout(temporaryRoot());
await ensureLayout(paths);
await enqueueJob(job(), paths);
await claimJob("gj-0000000001", paths);
await recoverProcessing(paths);
expect(await counts(paths)).toMatchObject({ pending: 1, processing: 0 });
const claimed = await claimJob("gj-0000000001", paths);
if (!claimed) throw new Error("missing recovered job");
await commitRecord(record(claimed.job), claimed.path, paths);
await recoverProcessing(paths);
expect(await counts(paths)).toMatchObject({ processing: 0, records: 1 });
});
it("resolves Linux openers from PATH and launches the selected opener", async () => {
const root = temporaryRoot();
const executable = join(root, "xdg-open");
writeFileSync(executable, "#!/bin/sh\nexit 0\n", "utf8");
chmodSync(executable, 0o755);
expect(openerFor("linux", { PATH: root })).toEqual({
command: executable,
args: [],
});
const previousPath = process.env.PATH;
try {
process.env.PATH = root;
expect(openerFor("linux", {})).toBeUndefined();
expect(openerFor("linux", { PATH: "" })).toBeUndefined();
} finally {
if (previousPath === undefined) delete process.env.PATH;
else process.env.PATH = previousPath;
}
await expect(
openDirectory(root, {
command: process.execPath,
args: ["-e", "process.exit(0)", "--"],
}),
).resolves.toBeUndefined();
});
});