diff --git a/packages/core/prisma/migrations/20260513130000_add_async_jobs/migration.sql b/packages/core/prisma/migrations/20260513130000_add_async_jobs/migration.sql new file mode 100644 index 0000000..80b42fb --- /dev/null +++ b/packages/core/prisma/migrations/20260513130000_add_async_jobs/migration.sql @@ -0,0 +1,20 @@ +-- CreateTable +CREATE TABLE "async_jobs" ( + "id" TEXT NOT NULL, + "userId" TEXT NOT NULL, + "kind" TEXT NOT NULL, + "result" JSONB, + "error" TEXT, + "completedAt" TIMESTAMP(3), + + CONSTRAINT "async_jobs_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE INDEX "async_jobs_userId_idx" ON "async_jobs"("userId"); + +-- CreateIndex +CREATE INDEX "async_jobs_kind_idx" ON "async_jobs"("kind"); + +-- AddForeignKey +ALTER TABLE "async_jobs" ADD CONSTRAINT "async_jobs_userId_fkey" FOREIGN KEY ("userId") REFERENCES "users"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/packages/core/prisma/models/async-job.prisma b/packages/core/prisma/models/async-job.prisma new file mode 100644 index 0000000..243bbd6 --- /dev/null +++ b/packages/core/prisma/models/async-job.prisma @@ -0,0 +1,25 @@ +/// Sidecar table that adds what project-q's `Message` does not own: +/// user association, result payload, and a free-text error message. +/// +/// `id` matches the project-q `Message.id` of the dispatched envelope - the +/// API creates an `AsyncJob` row and dispatches the envelope (carrying +/// `{ jobId }`) in the same transaction. Status, input (envelope), timestamps, +/// and retry state are all read from the related `Message` row; we don't +/// duplicate them. +/// +/// `kind` is a string discriminator (e.g. "parse-cv"). Handlers validate the +/// envelope payload against a zod schema keyed by `kind`. +model AsyncJob { + id String @id + userId String + kind String + result Json? + error String? + completedAt DateTime? + + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + + @@index([userId]) + @@index([kind]) + @@map("async_jobs") +} diff --git a/packages/core/prisma/models/user.prisma b/packages/core/prisma/models/user.prisma index 5a3716f..3090ad5 100644 --- a/packages/core/prisma/models/user.prisma +++ b/packages/core/prisma/models/user.prisma @@ -35,6 +35,9 @@ model User { // AI call logs aiCallLogs AiCallLog[] + // Async jobs (parse-cv, future kinds) + asyncJobs AsyncJob[] + @@map("users") } diff --git a/packages/core/src/modules/messenger/index.ts b/packages/core/src/modules/messenger/index.ts index 93d96b9..edd479b 100644 --- a/packages/core/src/modules/messenger/index.ts +++ b/packages/core/src/modules/messenger/index.ts @@ -5,6 +5,7 @@ export { InjectMessageBus, } from "@riotbyte-com/project-q-nestjs"; export { NestProjectQLogger } from "./logger.provider"; +export { ParseCVMessage } from "./messages/parse-cv.message"; export { RenderPdfMessage } from "./messages/render-pdf.message"; export { ProjectQMessagingModule, diff --git a/packages/core/src/modules/messenger/messages/parse-cv.message.ts b/packages/core/src/modules/messenger/messages/parse-cv.message.ts new file mode 100644 index 0000000..d48c0af --- /dev/null +++ b/packages/core/src/modules/messenger/messages/parse-cv.message.ts @@ -0,0 +1,24 @@ +import { defineZodMessage } from "@riotbyte-com/project-q-core"; +import { z } from "zod/v4"; + +/// `jobId` is the matching `AsyncJob.id`. `input` carries the source data the +/// handler needs - the file's storage key (for uploads pre-staged to +/// file-storage before dispatch) or the inline story text. +export const ParseCVMessage = defineZodMessage( + "parse-cv", + z.object({ + jobId: z.string(), + input: z.discriminatedUnion("source", [ + z.object({ + source: z.literal("file"), + fileKey: z.string(), + mimeType: z.string(), + originalName: z.string(), + }), + z.object({ + source: z.literal("story"), + text: z.string(), + }), + ]), + }), +); diff --git a/packages/core/src/modules/messenger/messenger.module.ts b/packages/core/src/modules/messenger/messenger.module.ts index 1f8753f..1d227b9 100644 --- a/packages/core/src/modules/messenger/messenger.module.ts +++ b/packages/core/src/modules/messenger/messenger.module.ts @@ -30,7 +30,10 @@ export class ProjectQMessagingModule { transports: { async: { dsn: `prisma://?queue=${queueName}`, retry: true }, }, - routing: new Map([["render-pdf", "async"]]), + routing: new Map([ + ["render-pdf", "async"], + ["parse-cv", "async"], + ]), }), ], providers: [{ provide: LoggerProvider, useClass: NestProjectQLogger }], diff --git a/packages/handlers/package.json b/packages/handlers/package.json index 91868e4..082008c 100644 --- a/packages/handlers/package.json +++ b/packages/handlers/package.json @@ -9,8 +9,11 @@ "build": "tsc -b" }, "dependencies": { + "@cv/ai-parser": "workspace:*", + "@cv/ai-provider": "workspace:*", "@cv/core": "workspace:*", "@cv/file-storage": "workspace:*", + "@cv/file-upload": "workspace:*", "@nestjs/common": "^11.1.18", "@riotbyte-com/project-q-core": "1.0.19-rc.5", "@riotbyte-com/project-q-nestjs": "1.0.19-rc.5", diff --git a/packages/handlers/src/handlers.module.ts b/packages/handlers/src/handlers.module.ts index a89dda2..7efa223 100644 --- a/packages/handlers/src/handlers.module.ts +++ b/packages/handlers/src/handlers.module.ts @@ -1,4 +1,7 @@ +import { BaseModule, DatabaseModule, UserModule } from "@cv/core"; +import { FileExtractionModule } from "@cv/file-upload"; import { DynamicModule, Module } from "@nestjs/common"; +import { ParseCvHandler } from "./parse-cv.handler"; import { HtmlToPdfService } from "./pdf/html-to-pdf.service"; import { PDF_SERVICE_CONFIG, PdfServiceConfig } from "./pdf/pdf-service.config"; import { RenderPdfHandler } from "./render-pdf.handler"; @@ -10,12 +13,19 @@ export class HandlersModule { static forRoot(config: PdfServiceConfig): DynamicModule { return { module: HandlersModule, + imports: [ + BaseModule, + DatabaseModule, + UserModule, + FileExtractionModule.forRoot(), + ], providers: [ { provide: PDF_SERVICE_CONFIG, useValue: config }, HtmlToPdfService, RenderPdfHandler, + ParseCvHandler, ], - exports: [HtmlToPdfService, RenderPdfHandler], + exports: [HtmlToPdfService, RenderPdfHandler, ParseCvHandler], }; } } diff --git a/packages/handlers/src/index.ts b/packages/handlers/src/index.ts index a000af7..0da7f7f 100644 --- a/packages/handlers/src/index.ts +++ b/packages/handlers/src/index.ts @@ -4,3 +4,4 @@ export { type PdfServiceConfig, } from "./pdf/html-to-pdf.service"; export { RenderPdfHandler } from "./render-pdf.handler"; +export { ParseCvHandler } from "./parse-cv.handler"; diff --git a/packages/handlers/src/parse-cv.handler.ts b/packages/handlers/src/parse-cv.handler.ts new file mode 100644 index 0000000..9e29d85 --- /dev/null +++ b/packages/handlers/src/parse-cv.handler.ts @@ -0,0 +1,147 @@ +import { CVParserService } from "@cv/ai-parser"; +import { + AIProvider, + AnthropicProvider, + LlamaCppProvider, + OpenAIProvider, +} from "@cv/ai-provider"; +import { + ClockService, + ParseCVMessage, + PrismaService, + TokenEncryptionService, +} from "@cv/core"; +import { FILE_STORAGE, type FileStorage } from "@cv/file-storage"; +import { TEXT_EXTRACTOR_REGISTRY, TextExtractorRegistry } from "@cv/file-upload"; +import { Inject, Injectable, Logger } from "@nestjs/common"; +import { type Envelope, type Handler } from "@riotbyte-com/project-q-core"; +import { HandlerTag } from "@riotbyte-com/project-q-nestjs"; + +const DEFAULT_BASE_URLS: Record = { + anthropic: "https://api.anthropic.com", + openai: "https://api.openai.com", +}; + +@Injectable() +@HandlerTag.decorator({ handles: "parse-cv" }) +export class ParseCvHandler implements Handler { + private readonly logger = new Logger(ParseCvHandler.name); + + constructor( + private readonly prisma: PrismaService, + private readonly encryption: TokenEncryptionService, + private readonly clock: ClockService, + @Inject(FILE_STORAGE) private readonly storage: FileStorage, + @Inject(TEXT_EXTRACTOR_REGISTRY) + private readonly extractors: TextExtractorRegistry, + ) {} + + async handle(envelope: Envelope): Promise { + const { jobId, input } = ParseCVMessage.parse(envelope.message).data; + + const job = await this.prisma.asyncJob.findUniqueOrThrow({ + where: { id: jobId }, + include: { user: true }, + }); + + this.logger.log(`Processing parse-cv job ${jobId} (kind=${job.kind})`); + + try { + const text = + input.source === "file" + ? await this.extractFromFile(input.fileKey, input.mimeType) + : input.text; + + const provider = await this.resolveProvider(job.userId); + const parser = new CVParserService(provider); + const result = await parser.parseCVText(text); + + await this.prisma.asyncJob.update({ + where: { id: jobId }, + data: { result, error: null, completedAt: this.clock.now() }, + }); + + this.logger.log(`Completed parse-cv job ${jobId}`); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + await this.prisma.asyncJob.update({ + where: { id: jobId }, + data: { error: message, completedAt: this.clock.now() }, + }); + this.logger.error(`Failed parse-cv job ${jobId}: ${message}`); + throw err; + } + } + + private async extractFromFile( + fileKey: string, + mimeType: string, + ): Promise { + const buffer = await this.storage.read(fileKey); + const extraction = await this.extractors.extract(buffer, mimeType); + + if (!extraction.success) { + throw new Error(`Text extraction failed: ${extraction.error}`); + } + + if (!extraction.text || extraction.text.trim().length === 0) { + throw new Error("Could not extract any text from the file"); + } + + return extraction.text; + } + + private async resolveProvider(userId: string): Promise { + const settings = await this.prisma.userAiSettings.upsert({ + where: { userId }, + create: { userId, aiPreference: "NO_AI" }, + update: {}, + }); + + if (settings.aiPreference === "NO_AI") { + throw new Error( + "AI is disabled for this user. Enable it in profile settings.", + ); + } + + if (settings.aiPreference === "PLATFORM") { + const type = process.env["AI_PROVIDER"] ?? "llama-cpp"; + if (type !== "llama-cpp") { + throw new Error( + `Platform provider type "${type}" is not configured in the worker.`, + ); + } + return new LlamaCppProvider({ + baseUrl: process.env["LLAMA_CPP_BASE_URL"] ?? "http://localhost:8080", + model: process.env["LLAMA_CPP_MODEL"] ?? "", + }); + } + + if (!settings.activeProviderId) { + throw new Error( + "No active AI provider configured. Set one in profile settings.", + ); + } + + const provider = await this.prisma.userAiProvider.findFirstOrThrow({ + where: { id: settings.activeProviderId, userId }, + }); + + const baseUrl = + provider.baseUrl ?? DEFAULT_BASE_URLS[provider.providerType] ?? ""; + const apiKey = this.encryption.decrypt(provider.encryptedApiKey); + const config = { + baseUrl, + apiKey, + ...(provider.model != null && { model: provider.model }), + }; + + if (provider.providerType === "anthropic") { + return new AnthropicProvider(config); + } + if (provider.providerType === "openai") { + return new OpenAIProvider(config); + } + throw new Error(`Unsupported provider type: ${provider.providerType}`); + } +} diff --git a/packages/handlers/tsconfig.json b/packages/handlers/tsconfig.json index 660dedd..7575f3e 100644 --- a/packages/handlers/tsconfig.json +++ b/packages/handlers/tsconfig.json @@ -2,7 +2,12 @@ "extends": "@cv/tsconfig/tsconfig.library.json", "compilerOptions": { "outDir": "./dist", - "rootDir": "./src" + "rootDir": "./src", + "baseUrl": ".", + "paths": { + "#core/*": ["../core/src/*"], + "#file-upload/*": ["../file-upload/src/*"] + } }, "include": [ "src/**/*" diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 99116dd..73816de 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -999,12 +999,21 @@ importers: packages/handlers: dependencies: + '@cv/ai-parser': + specifier: workspace:* + version: link:../ai-parser + '@cv/ai-provider': + specifier: workspace:* + version: link:../ai-provider '@cv/core': specifier: workspace:* version: link:../core '@cv/file-storage': specifier: workspace:* version: link:../file-storage + '@cv/file-upload': + specifier: workspace:* + version: link:../file-upload '@nestjs/common': specifier: ^11.1.18 version: 11.1.18(class-transformer@0.5.1)(class-validator@0.14.3)(reflect-metadata@0.2.2)(rxjs@7.8.2)