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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 | undefined { return sessionSourceFields(value) ? (value as Omit) : 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 { return JSON.parse(await readFile(path, "utf8")); } async function exists(path: string): Promise { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { 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 { return (await jsonFiles(paths.pending)).map((file) => file.slice(0, -5)); } export async function recoverProcessing(paths: GoodJobLayout): Promise { 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 { 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 { 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 { 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 { 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, }; }