repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/good-job/storage.ts
Rawimport { 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,
};
}