From 40d794c87f33c12d25e9931d06197d24da5854da Mon Sep 17 00:00:00 2001 From: robertn702 <8119609+robertn702@users.noreply.github.com> Date: Thu, 13 Aug 2026 16:54:33 +0000 Subject: [PATCH] feat(evaluations): add worker lifecycle --- .../migration.sql | 38 +++++ prisma/schema.prisma | 39 +++++ src/actions/mcpToken.actions.ts | 16 +- .../job-evaluations/[id]/complete/route.ts | 16 ++ .../api/job-evaluations/[id]/fail/route.ts | 17 ++ src/app/api/job-evaluations/route.ts | 45 ++++++ src/components/settings/McpAccessSettings.tsx | 14 +- src/lib/ai/tools/preprocessing.ts | 8 +- src/lib/jobEvaluations/http.ts | 38 +++++ src/lib/jobEvaluations/service.test.ts | 29 ++++ src/lib/jobEvaluations/service.ts | 148 ++++++++++++++++++ src/lib/jobs/resumeDetailInclude.ts | 1 + src/lib/mcp/auth.ts | 13 +- src/middleware.ts | 2 +- src/models/profile.model.ts | 7 + 15 files changed, 422 insertions(+), 9 deletions(-) create mode 100644 prisma/migrations/20260813180000_job_evaluations/migration.sql create mode 100644 src/app/api/job-evaluations/[id]/complete/route.ts create mode 100644 src/app/api/job-evaluations/[id]/fail/route.ts create mode 100644 src/app/api/job-evaluations/route.ts create mode 100644 src/lib/jobEvaluations/http.ts create mode 100644 src/lib/jobEvaluations/service.test.ts create mode 100644 src/lib/jobEvaluations/service.ts diff --git a/prisma/migrations/20260813180000_job_evaluations/migration.sql b/prisma/migrations/20260813180000_job_evaluations/migration.sql new file mode 100644 index 00000000..12b10e31 --- /dev/null +++ b/prisma/migrations/20260813180000_job_evaluations/migration.sql @@ -0,0 +1,38 @@ +-- Durable, source-independent asynchronous evaluation lifecycle. +CREATE TABLE "JobEvaluation" ( + "id" TEXT NOT NULL PRIMARY KEY, + "userId" TEXT NOT NULL, + "jobId" TEXT NOT NULL, + "resumeId" TEXT NOT NULL, + "evaluatorKey" TEXT NOT NULL, + "evaluatorVersion" TEXT NOT NULL, + "definitionHash" TEXT NOT NULL, + "inputHash" TEXT NOT NULL, + "jobHash" TEXT NOT NULL, + "resumeHash" TEXT NOT NULL, + "inputSnapshot" TEXT NOT NULL, + "isCurrent" BOOLEAN NOT NULL DEFAULT true, + "status" TEXT NOT NULL DEFAULT 'pending', + "attemptCount" INTEGER NOT NULL DEFAULT 0, + "nextAttemptAt" DATETIME, + "leaseToken" TEXT, + "leaseExpiresAt" DATETIME, + "resultJson" TEXT, + "resultHash" TEXT, + "lastError" TEXT, + "evaluatedAt" DATETIME, + "createdAt" DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" DATETIME NOT NULL, + CONSTRAINT "JobEvaluation_userId_fkey" FOREIGN KEY ("userId") REFERENCES "User" ("id") ON DELETE CASCADE ON UPDATE CASCADE, + CONSTRAINT "JobEvaluation_jobId_fkey" FOREIGN KEY ("jobId") REFERENCES "Job" ("id") ON DELETE CASCADE ON UPDATE CASCADE, + CONSTRAINT "JobEvaluation_resumeId_fkey" FOREIGN KEY ("resumeId") REFERENCES "Resume" ("id") ON DELETE CASCADE ON UPDATE CASCADE +); + +CREATE UNIQUE INDEX "JobEvaluation_jobId_resumeId_evaluatorKey_evaluatorVersion_inputHash_key" + ON "JobEvaluation"("jobId", "resumeId", "evaluatorKey", "evaluatorVersion", "inputHash"); +CREATE INDEX "JobEvaluation_userId_isCurrent_evaluatorKey_evaluatorVersion_idx" + ON "JobEvaluation"("userId", "isCurrent", "evaluatorKey", "evaluatorVersion"); +CREATE INDEX "JobEvaluation_userId_status_nextAttemptAt_idx" + ON "JobEvaluation"("userId", "status", "nextAttemptAt"); +CREATE INDEX "JobEvaluation_status_leaseExpiresAt_idx" ON "JobEvaluation"("status", "leaseExpiresAt"); +CREATE INDEX "JobEvaluation_jobId_isCurrent_idx" ON "JobEvaluation"("jobId", "isCurrent"); diff --git a/prisma/schema.prisma b/prisma/schema.prisma index 9183f3ea..9a4eeaff 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -35,6 +35,7 @@ model User { Tag Tag[] Question Question[] McpAccessToken McpAccessToken[] + JobEvaluation JobEvaluation[] ChatConversation ChatConversation? defaultResumeId String? defaultResume Resume? @relation("UserDefaultResume", fields: [defaultResumeId], references: [id], onDelete: SetNull) @@ -86,6 +87,7 @@ model Resume { File File? @relation(fields: [FileId], references: [id]) FileId String? @unique Job Job[] + JobEvaluation JobEvaluation[] Automation Automation[] defaultForUser User[] @relation("UserDefaultResume") reviewData String? @@ -312,6 +314,7 @@ model Job { CoverLetter CoverLetter? @relation(fields: [coverLetterId], references: [id]) coverLetterId String? Notes Note[] + evaluations JobEvaluation[] tags Tag[] // Automation discovery fields @@ -530,6 +533,42 @@ model McpAccessToken { @@index([userId]) } +model JobEvaluation { + id String @id @default(uuid()) + userId String + jobId String + resumeId String + evaluatorKey String + evaluatorVersion String + definitionHash String + inputHash String + jobHash String + resumeHash String + inputSnapshot String + isCurrent Boolean @default(true) + status String @default("pending") + attemptCount Int @default(0) + nextAttemptAt DateTime? + leaseToken String? + leaseExpiresAt DateTime? + resultJson String? + resultHash String? + lastError String? + evaluatedAt DateTime? + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + job Job @relation(fields: [jobId], references: [id], onDelete: Cascade) + resume Resume @relation(fields: [resumeId], references: [id], onDelete: Cascade) + + @@unique([jobId, resumeId, evaluatorKey, evaluatorVersion, inputHash]) + @@index([userId, isCurrent, evaluatorKey, evaluatorVersion]) + @@index([userId, status, nextAttemptAt]) + @@index([status, leaseExpiresAt]) + @@index([jobId, isCurrent]) +} + model ChatConversation { id String @id @default(uuid()) userId String @unique diff --git a/src/actions/mcpToken.actions.ts b/src/actions/mcpToken.actions.ts index c3af178f..8fbe720c 100644 --- a/src/actions/mcpToken.actions.ts +++ b/src/actions/mcpToken.actions.ts @@ -19,6 +19,7 @@ export interface PublicTokenMeta { export async function createMcpToken(input: { name: string; expiryDays: 30 | 90 | 365; + type?: "agent" | "evaluation-worker"; }): Promise< | { success: true; token: string; record: PublicTokenMeta } | { success: false; message: string } @@ -44,7 +45,9 @@ export async function createMcpToken(input: { name: input.name.trim(), tokenHash: hash, tokenPrefix: prefix, - scopes: JSON.stringify(["jobs:write", "questions:write", "resume:write"]), + scopes: JSON.stringify(input.type === "evaluation-worker" + ? ["evaluations:worker"] + : ["jobs:write", "questions:write", "resume:write"]), expiresAt, }, }); @@ -105,9 +108,18 @@ function toPublicMeta(record: { id: record.id, name: record.name, tokenPrefix: record.tokenPrefix, - scopes: JSON.parse(record.scopes) as string[], + scopes: parseScopes(record.scopes), expiresAt: record.expiresAt, lastUsedAt: record.lastUsedAt, createdAt: record.createdAt, }; } + +function parseScopes(scopes: string): string[] { + try { + const parsed: unknown = JSON.parse(scopes); + return Array.isArray(parsed) && parsed.every((scope) => typeof scope === "string") ? parsed : []; + } catch { + return []; + } +} diff --git a/src/app/api/job-evaluations/[id]/complete/route.ts b/src/app/api/job-evaluations/[id]/complete/route.ts new file mode 100644 index 00000000..655461fe --- /dev/null +++ b/src/app/api/job-evaluations/[id]/complete/route.ts @@ -0,0 +1,16 @@ +import { NextResponse } from "next/server"; +import { completeEvaluation } from "@/lib/jobEvaluations/service"; +import { boundedJson, stringField, workerAuth } from "@/lib/jobEvaluations/http"; + +export async function POST(request: Request, { params }: { params: Promise<{ id: string }> }) { + const auth = await workerAuth(request); + if ("error" in auth) return auth.error; + const parsed = await boundedJson(request); + if ("error" in parsed) return parsed.error; + const body = parsed.value as Record; + const leaseToken = stringField(body.leaseToken, 200); + const inputHash = stringField(body.inputHash, 128); + if (!leaseToken || !inputHash || !body.result || typeof body.result !== "object" || Array.isArray(body.result) || JSON.stringify(body.result).length > 96_000) return NextResponse.json({ error: "Invalid completion" }, { status: 400 }); + const outcome = await completeEvaluation(auth.userId, (await params).id, leaseToken, inputHash, body.result as Record); + return outcome === "completed" || outcome === "idempotent" ? NextResponse.json({ status: outcome }) : NextResponse.json({ error: "Stale or conflicting completion" }, { status: 409 }); +} diff --git a/src/app/api/job-evaluations/[id]/fail/route.ts b/src/app/api/job-evaluations/[id]/fail/route.ts new file mode 100644 index 00000000..8ec4612b --- /dev/null +++ b/src/app/api/job-evaluations/[id]/fail/route.ts @@ -0,0 +1,17 @@ +import { NextResponse } from "next/server"; +import { failEvaluation } from "@/lib/jobEvaluations/service"; +import { boundedJson, stringField, workerAuth } from "@/lib/jobEvaluations/http"; + +export async function POST(request: Request, { params }: { params: Promise<{ id: string }> }) { + const auth = await workerAuth(request); + if ("error" in auth) return auth.error; + const parsed = await boundedJson(request, 8_000); + if ("error" in parsed) return parsed.error; + const body = parsed.value as Record; + const leaseToken = stringField(body.leaseToken, 200); + const inputHash = stringField(body.inputHash, 128); + const error = stringField(body.error, 2_000); + if (!leaseToken || !inputHash || !error) return NextResponse.json({ error: "Invalid failure" }, { status: 400 }); + const outcome = await failEvaluation(auth.userId, (await params).id, leaseToken, inputHash, error); + return outcome === "retry" || outcome === "failed" || outcome === "idempotent" ? NextResponse.json({ status: outcome }) : NextResponse.json({ error: "Stale or conflicting failure" }, { status: 409 }); +} diff --git a/src/app/api/job-evaluations/route.ts b/src/app/api/job-evaluations/route.ts new file mode 100644 index 00000000..f75754ff --- /dev/null +++ b/src/app/api/job-evaluations/route.ts @@ -0,0 +1,45 @@ +import { NextResponse } from "next/server"; +import prisma from "@/lib/db"; +import { materializeCurrentEvaluations, claimEvaluations, type Evaluator } from "@/lib/jobEvaluations/service"; +import { boundedJson, stringField, workerAuth } from "@/lib/jobEvaluations/http"; + +function evaluatorFrom(value: unknown): Evaluator | null { + if (!value || typeof value !== "object" || Array.isArray(value)) return null; + const body = value as Record; + const evaluatorKey = stringField(body.evaluatorKey, 100); + const evaluatorVersion = stringField(body.evaluatorVersion, 100); + if (!evaluatorKey || !evaluatorVersion || body.evaluatorDefinition === undefined || JSON.stringify(body.evaluatorDefinition).length > 32_000) return null; + return { evaluatorKey, evaluatorVersion, evaluatorDefinition: body.evaluatorDefinition }; +} + +export async function POST(request: Request) { + const auth = await workerAuth(request); + if ("error" in auth) return auth.error; + const parsed = await boundedJson(request); + if ("error" in parsed) return parsed.error; + const body = parsed.value as Record; + const evaluator = evaluatorFrom(body); + const batch = typeof body.batch === "number" ? body.batch : NaN; + const leaseSeconds = typeof body.leaseSeconds === "number" ? body.leaseSeconds : NaN; + if (!evaluator || !Number.isInteger(batch) || batch < 1 || batch > 10 || !Number.isInteger(leaseSeconds) || leaseSeconds < 15 || leaseSeconds > 3600) { + return NextResponse.json({ error: "Invalid claim request" }, { status: 400 }); + } + const rows = await claimEvaluations(auth.userId, evaluator, batch, leaseSeconds * 1000); + return NextResponse.json({ evaluations: rows.map((row) => ({ id: row.id, inputHash: row.inputHash, leaseToken: row.leaseToken, leaseExpiresAt: row.leaseExpiresAt, input: JSON.parse(row.inputSnapshot) })) }); +} + +export async function GET(request: Request) { + const auth = await workerAuth(request); + if ("error" in auth) return auth.error; + const searchParams = new URL(request.url).searchParams; + const evaluatorKey = searchParams.get("evaluatorKey"); + const evaluatorVersion = searchParams.get("evaluatorVersion"); + const definition = searchParams.get("definition"); + if (!evaluatorKey || !evaluatorVersion || !definition) return NextResponse.json({ error: "evaluatorKey, evaluatorVersion, and definition are required" }, { status: 400 }); + let evaluatorDefinition: unknown; + try { evaluatorDefinition = JSON.parse(definition); } catch { return NextResponse.json({ error: "Invalid definition" }, { status: 400 }); } + await materializeCurrentEvaluations(auth.userId, { evaluatorKey, evaluatorVersion, evaluatorDefinition }); + const jobId = new URL(request.url).searchParams.get("jobId"); + const rows = await prisma.jobEvaluation.findMany({ where: { userId: auth.userId, evaluatorKey, evaluatorVersion, ...(jobId ? { jobId } : {}) }, orderBy: { createdAt: "desc" } }); + return NextResponse.json({ evaluations: rows.filter((row) => row.isCurrent), history: rows.filter((row) => !row.isCurrent).map((row) => ({ id: row.id, jobId: row.jobId, status: row.status, evaluatedAt: row.evaluatedAt })) }); +} diff --git a/src/components/settings/McpAccessSettings.tsx b/src/components/settings/McpAccessSettings.tsx index e621c0b8..0504a9b7 100644 --- a/src/components/settings/McpAccessSettings.tsx +++ b/src/components/settings/McpAccessSettings.tsx @@ -153,6 +153,7 @@ export default function McpAccessSettings() { const [expiryDays, setExpiryDays] = useState<30 | 90 | 365>( APP_CONSTANTS.MCP_TOKEN_EXPIRY_DEFAULT_DAYS as 30 | 90 | 365, ); + const [tokenType, setTokenType] = useState<"agent" | "evaluation-worker">("agent"); const [revealedToken, setRevealedToken] = useState<{ token: string; name: string } | null>(null); @@ -173,7 +174,7 @@ export default function McpAccessSettings() { const handleGenerate = async () => { if (!tokenName.trim()) return; setGenerating(true); - const result = await createMcpToken({ name: tokenName.trim(), expiryDays }); + const result = await createMcpToken({ name: tokenName.trim(), expiryDays, type: tokenType }); setGenerating(false); if (!result.success) { toastError(result.message); @@ -183,6 +184,7 @@ export default function McpAccessSettings() { setShowGenerateDialog(false); setTokenName(""); setExpiryDays(APP_CONSTANTS.MCP_TOKEN_EXPIRY_DEFAULT_DAYS as 30 | 90 | 365); + setTokenType("agent"); await fetchTokens(); }; @@ -305,6 +307,16 @@ export default function McpAccessSettings() { onKeyDown={(e) => e.key === "Enter" && handleGenerate()} /> +
+ + +