Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/good-job/storage.ts

Raw
import { randomUUID } from "node:crypto";
import {
	link,
	mkdir,
	open,
	readdir,
	readFile,
	realpath,
	rename,
	rm,
} from "node:fs/promises";
import { homedir } from "node:os";
import { basename, dirname, join, posix, resolve, win32 } from "node:path";
import type {
	GoodJobFailure,
	GoodJobJob,
	GoodJobRecord,
	ImportedLearningSource,
	LearningSource,
	ModelRef,
	SessionLearningSource,
	TokenUsage,
} from "./core.js";
import {
	boundedError,
	FAILURE_VERSION,
	feedbackIdFromSeed,
	JOB_VERSION,
	learningSlug,
	RECORD_VERSION,
} from "./core.js";

export interface GoodJobLayout {
	root: string;
	pending: string;
	processing: string;
	failed: string;
	records: string;
}

export interface GoodJobCounts {
	pending: number;
	processing: number;
	failed: number;
	records: number;
}

export interface ClaimedJob {
	job: GoodJobJob;
	path: string;
}

export interface LoadedRecords {
	records: GoodJobRecord[];
	errors: string[];
}

export interface MigrationResult {
	jobs: number;
	records: number;
	failures: number;
	active: number;
}

export function goodJobDataDirectory(
	env: NodeJS.ProcessEnv = process.env,
	platform = process.platform,
	home = homedir(),
): string {
	if (platform === "win32") {
		const base = env.LOCALAPPDATA?.trim() || env.APPDATA?.trim();
		return win32.join(
			base || win32.join(home, "AppData", "Local"),
			"pi",
			"good-job",
		);
	}
	if (platform === "darwin")
		return posix.join(home, "Library", "Application Support", "pi", "good-job");
	return posix.join(
		env.XDG_DATA_HOME?.trim() || posix.join(home, ".local", "share"),
		"pi",
		"good-job",
	);
}

export function configuredDataDirectory(
	value: string | undefined,
	cwd: string,
	home = homedir(),
): string {
	if (!value) return goodJobDataDirectory();
	const trimmed = value.trim();
	if (trimmed === "~") return home;
	if (/^~[\\/]/.test(trimmed)) return resolve(home, trimmed.slice(2));
	return resolve(cwd, trimmed);
}

export function layout(root = goodJobDataDirectory()): GoodJobLayout {
	return {
		root,
		pending: join(root, "queue", "pending"),
		processing: join(root, "queue", "processing"),
		failed: join(root, "queue", "failed"),
		records: join(root, "records"),
	};
}

export async function ensureLayout(paths: GoodJobLayout): Promise<void> {
	await Promise.all(
		[paths.pending, paths.processing, paths.failed, paths.records].map(
			(directory) => mkdir(directory, { recursive: true, mode: 0o700 }),
		),
	);
}

function isNodeError(error: unknown): error is NodeJS.ErrnoException {
	return error instanceof Error && "code" in error;
}

async function jsonFiles(directory: string): Promise<string[]> {
	try {
		return (await readdir(directory))
			.filter((file) => file.endsWith(".json"))
			.sort();
	} catch (error) {
		if (isNodeError(error) && error.code === "ENOENT") return [];
		throw error;
	}
}

async function syncDirectory(directory: string): Promise<void> {
	if (process.platform === "win32") return;
	try {
		const handle = await open(directory, "r");
		try {
			await handle.sync();
		} finally {
			await handle.close();
		}
	} catch {
		// File fsync plus rename is the portable durability floor.
	}
}

async function temporaryJson(target: string, value: unknown): Promise<string> {
	const directory = dirname(target);
	await mkdir(directory, { recursive: true, mode: 0o700 });
	const temporary = join(
		directory,
		`.${basename(target)}.${process.pid}.${randomUUID()}.tmp`,
	);
	const file = await open(temporary, "wx", 0o600);
	try {
		await file.writeFile(`${JSON.stringify(value)}\n`, "utf8");
		await file.sync();
	} catch (error) {
		await file.close().catch(() => {});
		await rm(temporary, { force: true });
		throw error;
	}
	await file.close();
	return temporary;
}

export async function atomicWriteJson(
	target: string,
	value: unknown,
): Promise<void> {
	const temporary = await temporaryJson(target, value);
	try {
		await rename(temporary, target);
		await syncDirectory(dirname(target));
	} catch (error) {
		await rm(temporary, { force: true });
		throw error;
	}
}

function isLinkUnsupported(error: unknown): boolean {
	if (!isNodeError(error)) return false;
	return (
		error.code === "EACCES" ||
		error.code === "EPERM" ||
		error.code === "ENOSYS" ||
		error.code === "EOPNOTSUPP" ||
		error.code === "EXDEV" ||
		error.code === "EMLINK"
	);
}

async function exclusiveCreateJson(
	target: string,
	value: unknown,
): Promise<void> {
	const file = await open(target, "wx", 0o600);
	try {
		await file.writeFile(`${JSON.stringify(value)}\n`, "utf8");
		await file.sync();
	} finally {
		await file.close();
	}
}

export async function atomicCreateJson(
	target: string,
	value: unknown,
): Promise<void> {
	const temporary = await temporaryJson(target, value);
	try {
		try {
			await link(temporary, target);
		} catch (error) {
			// Android SELinux and some filesystems forbid hard links.
			if (!isLinkUnsupported(error)) throw error;
			await exclusiveCreateJson(target, value);
		}
		await rm(temporary, { force: true });
		await syncDirectory(dirname(target));
	} catch (error) {
		await rm(temporary, { force: true });
		throw error;
	}
}

function isObject(value: unknown): value is Record<string, unknown> {
	return value !== null && typeof value === "object" && !Array.isArray(value);
}

function modelRef(value: unknown): value is ModelRef {
	return (
		isObject(value) &&
		typeof value.provider === "string" &&
		value.provider.length > 0 &&
		typeof value.id === "string" &&
		value.id.length > 0
	);
}

function sessionSourceFields(
	value: unknown,
): value is Omit<SessionLearningSource, "type"> {
	return (
		isObject(value) &&
		typeof value.cwd === "string" &&
		(value.sessionPath === null || typeof value.sessionPath === "string") &&
		typeof value.sessionId === "string" &&
		(value.sessionName === null || typeof value.sessionName === "string") &&
		(value.leafId === null || typeof value.leafId === "string") &&
		(value.agentModel === null || modelRef(value.agentModel))
	);
}

function sessionSource(value: unknown): value is SessionLearningSource {
	return (
		isObject(value) && value.type === "session" && sessionSourceFields(value)
	);
}

function importedSource(value: unknown): value is ImportedLearningSource {
	return (
		isObject(value) &&
		value.type === "import" &&
		typeof value.importer === "string" &&
		value.importer.length > 0 &&
		typeof value.fingerprint === "string" &&
		value.fingerprint.length > 0 &&
		typeof value.importedAt === "string"
	);
}

function source(value: unknown): value is LearningSource {
	return sessionSource(value) || importedSource(value);
}

function legacySource(
	value: unknown,
): Omit<SessionLearningSource, "type"> | undefined {
	return sessionSourceFields(value)
		? (value as Omit<SessionLearningSource, "type">)
		: undefined;
}

function stringArray(value: unknown): value is string[] {
	return (
		Array.isArray(value) &&
		value.length > 0 &&
		value.every((item) => typeof item === "string" && item.length > 0)
	);
}

function nonnegativeInteger(value: unknown): value is number {
	return typeof value === "number" && Number.isSafeInteger(value) && value >= 0;
}

function recordUsage(value: unknown): value is TokenUsage {
	return (
		isObject(value) &&
		nonnegativeInteger(value.input) &&
		nonnegativeInteger(value.output) &&
		nonnegativeInteger(value.cacheRead) &&
		nonnegativeInteger(value.cacheWrite) &&
		nonnegativeInteger(value.totalTokens) &&
		typeof value.cost === "number" &&
		Number.isFinite(value.cost) &&
		value.cost >= 0
	);
}

function feedbackKind(value: unknown): value is GoodJobJob["kind"] {
	return value === "praise" || value === "problem";
}

const CURRENT_ID = /^(gj|wtf)-[0-9a-hjkmnp-tv-z]{10}$/;

function currentId(value: unknown, kind: GoodJobJob["kind"]): value is string {
	const prefix = kind === "praise" ? "gj" : "wtf";
	return typeof value === "string" && CURRENT_ID.exec(value)?.[1] === prefix;
}

export function parseJob(value: unknown): GoodJobJob | undefined {
	if (
		!isObject(value) ||
		value.version !== JOB_VERSION ||
		!feedbackKind(value.kind) ||
		!currentId(value.id, value.kind) ||
		typeof value.createdAt !== "string" ||
		typeof value.feedback !== "string" ||
		!sessionSource(value.source) ||
		!(value.analysisModel === null || modelRef(value.analysisModel))
	)
		return undefined;
	return value as unknown as GoodJobJob;
}

export function parseRecord(value: unknown): GoodJobRecord | undefined {
	if (
		!isObject(value) ||
		value.version !== RECORD_VERSION ||
		!feedbackKind(value.kind) ||
		!currentId(value.id, value.kind) ||
		typeof value.createdAt !== "string" ||
		typeof value.completedAt !== "string" ||
		typeof value.feedback !== "string" ||
		!source(value.source) ||
		!(value.analysisModel === null || modelRef(value.analysisModel)) ||
		!isObject(value.learning) ||
		typeof value.learning.slug !== "string" ||
		learningSlug(value.learning.slug) !== value.learning.slug ||
		typeof value.learning.summary !== "string" ||
		value.learning.summary.length === 0 ||
		!stringArray(value.learning.evidence) ||
		!stringArray(value.learning.behaviors) ||
		!(
			value.learning.candidateRule === null ||
			typeof value.learning.candidateRule === "string"
		) ||
		!(value.usage === null || recordUsage(value.usage))
	)
		return undefined;
	return value as unknown as GoodJobRecord;
}

function migrateLegacyJob(value: unknown): GoodJobJob | undefined {
	if (
		!isObject(value) ||
		value.version !== 1 ||
		typeof value.id !== "string" ||
		typeof value.createdAt !== "string" ||
		typeof value.praise !== "string" ||
		!(value.analysisModel === null || modelRef(value.analysisModel))
	)
		return undefined;
	const oldSource = legacySource(value.source);
	if (!oldSource) return undefined;
	return {
		version: JOB_VERSION,
		id: feedbackIdFromSeed("praise", `good-job-v1:${value.id}`),
		createdAt: value.createdAt,
		kind: "praise",
		feedback: value.praise,
		source: { type: "session", ...oldSource },
		analysisModel: value.analysisModel,
	};
}

function migrateLegacyRecord(value: unknown): GoodJobRecord | undefined {
	if (
		!isObject(value) ||
		value.version !== 1 ||
		typeof value.id !== "string" ||
		typeof value.createdAt !== "string" ||
		typeof value.completedAt !== "string" ||
		typeof value.praise !== "string" ||
		!modelRef(value.analysisModel) ||
		!isObject(value.learning) ||
		typeof value.learning.achievement !== "string" ||
		value.learning.achievement.length === 0 ||
		!stringArray(value.learning.evidence) ||
		!stringArray(value.learning.successfulBehaviors) ||
		!(
			value.learning.candidateRule === null ||
			typeof value.learning.candidateRule === "string"
		) ||
		!recordUsage(value.usage)
	)
		return undefined;
	const oldSource = legacySource(value.source);
	if (!oldSource) return undefined;
	return {
		version: RECORD_VERSION,
		id: feedbackIdFromSeed("praise", `good-job-v1:${value.id}`),
		createdAt: value.createdAt,
		completedAt: value.completedAt,
		kind: "praise",
		feedback: value.praise,
		source: { type: "session", ...oldSource },
		analysisModel: value.analysisModel,
		learning: {
			slug: learningSlug(value.learning.achievement),
			summary: value.learning.achievement,
			evidence: value.learning.evidence,
			behaviors: value.learning.successfulBehaviors,
			candidateRule: value.learning.candidateRule,
		},
		usage: value.usage,
	};
}

async function readJson(path: string): Promise<unknown> {
	return JSON.parse(await readFile(path, "utf8"));
}

async function exists(path: string): Promise<boolean> {
	try {
		const file = await open(path, "r");
		await file.close();
		return true;
	} catch (error) {
		if (isNodeError(error) && error.code === "ENOENT") return false;
		throw error;
	}
}

async function moveValue(
	sourcePath: string,
	targetPath: string,
	value: unknown,
): Promise<void> {
	if (sourcePath === targetPath) {
		await atomicWriteJson(targetPath, value);
		return;
	}
	try {
		await atomicCreateJson(targetPath, value);
	} catch (error) {
		if (!(isNodeError(error) && error.code === "EEXIST")) throw error;
		const existing = await readJson(targetPath);
		if (JSON.stringify(existing) !== JSON.stringify(value))
			throw new Error(`migration collision at ${targetPath}`);
	}
	await rm(sourcePath, { force: true });
}

async function migrateJobs(
	from: GoodJobLayout,
	to: GoodJobLayout,
): Promise<number> {
	let migrated = 0;
	for (const file of await jsonFiles(from.pending)) {
		const sourcePath = join(from.pending, file);
		const value = await readJson(sourcePath);
		const job = parseJob(value) ?? migrateLegacyJob(value);
		if (!job) throw new Error(`cannot migrate queued job ${sourcePath}`);
		const targetPath = join(to.pending, `${job.id}.json`);
		if (sourcePath !== targetPath || !parseJob(value)) {
			await moveValue(sourcePath, targetPath, job);
			migrated++;
		}
	}
	return migrated;
}

async function migrateRecords(
	from: GoodJobLayout,
	to: GoodJobLayout,
): Promise<number> {
	let migrated = 0;
	for (const file of await jsonFiles(from.records)) {
		const sourcePath = join(from.records, file);
		const value = await readJson(sourcePath);
		const record = parseRecord(value) ?? migrateLegacyRecord(value);
		if (!record)
			throw new Error(`cannot migrate learning record ${sourcePath}`);
		const targetPath = join(to.records, `${record.id}.json`);
		if (sourcePath !== targetPath || !parseRecord(value)) {
			await moveValue(sourcePath, targetPath, record);
			migrated++;
		}
	}
	return migrated;
}

async function migrateFailures(
	from: GoodJobLayout,
	to: GoodJobLayout,
): Promise<number> {
	let migrated = 0;
	for (const file of await jsonFiles(from.failed)) {
		const sourcePath = join(from.failed, file);
		const value = await readJson(sourcePath);
		if (
			!isObject(value) ||
			value.version !== FAILURE_VERSION ||
			typeof value.failedAt !== "string" ||
			typeof value.error !== "string"
		)
			throw new Error(`cannot migrate failure ${sourcePath}`);
		if (value.job === null) {
			const targetPath = join(to.failed, file);
			if (sourcePath !== targetPath) {
				await moveValue(sourcePath, targetPath, value);
				migrated++;
			}
			continue;
		}
		const job = parseJob(value.job) ?? migrateLegacyJob(value.job);
		if (!job) throw new Error(`cannot migrate failure ${sourcePath}`);
		const failure: GoodJobFailure = {
			version: FAILURE_VERSION,
			job,
			failedAt: value.failedAt,
			error: value.error,
		};
		const targetPath = join(to.failed, `${job.id}.json`);
		if (sourcePath !== targetPath || !parseJob(value.job)) {
			await moveValue(sourcePath, targetPath, failure);
			migrated++;
		}
	}
	return migrated;
}

async function sameDirectory(left: string, right: string): Promise<boolean> {
	if (resolve(left) === resolve(right)) return true;
	try {
		return (await realpath(left)) === (await realpath(right));
	} catch (error) {
		if (isNodeError(error) && error.code === "ENOENT") return false;
		throw error;
	}
}

export async function migrateStore(
	sourceRoot: string,
	targetRoot: string,
): Promise<MigrationResult> {
	const from = layout(sourceRoot);
	const destination = layout(targetRoot);
	const before = await counts(from);
	if (before.pending + before.processing + before.failed + before.records === 0)
		return { jobs: 0, records: 0, failures: 0, active: 0 };
	await ensureLayout(destination);
	const to = (await sameDirectory(sourceRoot, targetRoot)) ? from : destination;
	await recoverProcessing(from);
	const [jobs, records, failures] = await Promise.all([
		migrateJobs(from, to),
		migrateRecords(from, to),
		migrateFailures(from, to),
	]);
	return {
		jobs,
		records,
		failures,
		active: (await jsonFiles(from.processing)).length,
	};
}

export async function enqueueJob(
	job: GoodJobJob,
	paths: GoodJobLayout,
): Promise<void> {
	if (!parseJob(job)) throw new Error("refusing to persist invalid gj job");
	const file = `${job.id}.json`;
	if (
		(await exists(join(paths.records, file))) ||
		(await exists(join(paths.failed, file))) ||
		(await jsonFiles(paths.processing)).some((claim) =>
			claim.startsWith(`${job.id}--`),
		)
	)
		throw new Error(`duplicate gj id ${job.id}`);
	try {
		await atomicCreateJson(join(paths.pending, file), job);
	} catch (error) {
		if (isNodeError(error) && error.code === "EEXIST")
			throw new Error(`duplicate gj id ${job.id}`);
		throw error;
	}
}

export async function claimJob(
	id: string,
	paths: GoodJobLayout,
): Promise<ClaimedJob | undefined> {
	const pending = join(paths.pending, `${id}.json`);
	const processing = join(
		paths.processing,
		`${id}--${process.pid}--${randomUUID()}.json`,
	);
	await mkdir(paths.processing, { recursive: true, mode: 0o700 });
	try {
		await rename(pending, processing);
	} catch (error) {
		if (isNodeError(error) && error.code === "ENOENT") return undefined;
		throw error;
	}
	const job = parseJob(await readJson(processing));
	if (job) return { job, path: processing };
	await moveInvalidJob(processing, id, paths, "invalid queued job schema");
	return undefined;
}

async function moveInvalidJob(
	path: string,
	id: string,
	paths: GoodJobLayout,
	error: string,
): Promise<void> {
	const failed = {
		version: FAILURE_VERSION,
		job: null,
		failedAt: new Date().toISOString(),
		error,
	};
	await atomicCreateJson(join(paths.failed, `${id}.json`), failed);
	await rm(path, { force: true });
}

export async function pendingJobIds(paths: GoodJobLayout): Promise<string[]> {
	return (await jsonFiles(paths.pending)).map((file) => file.slice(0, -5));
}

export async function recoverProcessing(paths: GoodJobLayout): Promise<void> {
	for (const file of await jsonFiles(paths.processing)) {
		const claim = parseClaimFile(file);
		const processing = join(paths.processing, file);
		if (!claim) {
			await moveInvalidJob(
				processing,
				file.slice(0, -5),
				paths,
				"invalid processing claim name",
			);
			continue;
		}
		const recordFile = `${claim.id}.json`;
		if (
			(await exists(join(paths.records, recordFile))) ||
			(await exists(join(paths.failed, recordFile)))
		) {
			await rm(processing, { force: true });
			continue;
		}
		if (claim.pid !== process.pid && processExists(claim.pid)) continue;
		const pending = join(paths.pending, recordFile);
		if (await exists(pending)) {
			await rm(processing, { force: true });
			continue;
		}
		await rename(processing, pending);
	}
}

function parseClaimFile(file: string): { id: string; pid: number } | undefined {
	const match = /^(.*)--([1-9][0-9]*)--[^/\\]+\.json$/.exec(file);
	if (!match?.[1] || !match[2]) return undefined;
	const pid = Number(match[2]);
	return Number.isSafeInteger(pid) ? { id: match[1], pid } : undefined;
}

function processExists(pid: number): boolean {
	try {
		process.kill(pid, 0);
		return true;
	} catch (error) {
		return !(isNodeError(error) && error.code === "ESRCH");
	}
}

export async function commitRecord(
	record: GoodJobRecord,
	claimPath: string,
	paths: GoodJobLayout,
): Promise<void> {
	if (!parseRecord(record))
		throw new Error("refusing to persist invalid gj record");
	try {
		await atomicCreateJson(join(paths.records, `${record.id}.json`), record);
	} catch (error) {
		if (isNodeError(error) && error.code === "EEXIST")
			throw new Error(`duplicate gj id ${record.id}`);
		throw error;
	}
	await rm(claimPath, { force: true });
}

export async function commitFailure(
	job: GoodJobJob,
	error: unknown,
	claimPath: string,
	paths: GoodJobLayout,
): Promise<void> {
	const failure: GoodJobFailure = {
		version: FAILURE_VERSION,
		job,
		failedAt: new Date().toISOString(),
		error: boundedError(error),
	};
	try {
		await atomicCreateJson(join(paths.failed, `${job.id}.json`), failure);
	} catch (error) {
		if (isNodeError(error) && error.code === "EEXIST")
			throw new Error(`duplicate gj id ${job.id}`);
		throw error;
	}
	await rm(claimPath, { force: true });
}

export async function loadRecords(
	paths: GoodJobLayout,
): Promise<LoadedRecords> {
	const files = await jsonFiles(paths.records);
	const loaded = await Promise.all(
		files.map(async (file) => {
			try {
				const record = parseRecord(await readJson(join(paths.records, file)));
				if (!record) throw new Error("invalid record schema");
				return { record };
			} catch (error) {
				return { error: `${file}: ${boundedError(error)}` };
			}
		}),
	);
	const records: GoodJobRecord[] = [];
	const errors: string[] = [];
	for (const result of loaded) {
		if (result.record) records.push(result.record);
		else if (result.error) errors.push(result.error);
	}
	records.sort((left, right) =>
		right.completedAt.localeCompare(left.completedAt),
	);
	return { records, errors };
}

export async function counts(paths: GoodJobLayout): Promise<GoodJobCounts> {
	const [pending, processing, failed, records] = await Promise.all([
		jsonFiles(paths.pending),
		jsonFiles(paths.processing),
		jsonFiles(paths.failed),
		jsonFiles(paths.records),
	]);
	return {
		pending: pending.length,
		processing: processing.length,
		failed: failed.length,
		records: records.length,
	};
}