Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/good-job/__tests__/storage.test.ts

Raw
import {
	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();
	});
});