From 2cf773fefc748315458c9b3818a62ccc287808ff Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Thu, 14 May 2026 22:16:32 +0200 Subject: [PATCH] feat(CVG-84): move CV parsing onto worker queue (parse-cv end-to-end) (#17) --- apps/api/package.json | 1 + apps/api/schema.gql | 32 +++ apps/api/src/modules/admin/admin.module.ts | 3 +- apps/api/src/modules/admin/admin.resolver.ts | 3 +- .../src/modules/admin/ai-call-log.module.ts | 12 -- apps/api/src/modules/app.module.ts | 6 +- .../__tests__/typed-async-job.bundle.spec.ts | 193 +++++++++++++++++ .../async-job/async-job-graphql.module.ts | 20 ++ .../modules/async-job/async-job.resolver.ts | 21 ++ .../src/modules/async-job/async-job.type.ts | 44 ++++ .../async-job/typed-async-job.bundle.ts | 131 ++++++++++++ .../async-job/typed-async-job.interceptor.ts | 113 ++++++++++ .../async-job/typed-async-job.shared.ts | 53 +++++ .../cv-parser-dispatch.service.spec.ts | 201 ++++++++++++++++++ .../__tests__/cv-parser.service.spec.ts | 65 ++++-- .../cv-parser/ai-provider-resolver.service.ts | 122 ----------- .../cv-parser/cv-parser-dispatch.service.ts | 120 +++++++++++ .../src/modules/cv-parser/cv-parser.module.ts | 13 +- .../modules/cv-parser/cv-parser.service.ts | 95 ++------- .../cv-parser/enqueue-parse-cv.resolver.ts | 40 ++++ .../cv-parser/entity-resolver.service.ts | 75 ++++++- .../cv-parser/parse-cv-job.resolver.ts | 46 ++++ .../modules/cv-parser/parse-cv-job.type.ts | 62 ++++++ .../data-import/sources/file-import-source.ts | 4 +- .../user-settings/user-ai-settings.service.ts | 31 --- apps/api/vitest.config.ts | 18 ++ .../mutations/enqueue-parse-file.graphql | 13 ++ .../mutations/useEnqueueParseFileMutation.ts | 33 +++ .../onboarding/queries/parse-cv-job.graphql | 46 ++++ .../onboarding/queries/useParseCvJobQuery.ts | 26 +++ .../migration.sql | 32 +++ packages/core/prisma/models/async-job.prisma | 37 ++++ packages/core/prisma/models/user.prisma | 3 + packages/core/src/index.ts | 2 + .../ai-call-log-persistence.service.ts | 6 +- .../ai-resolution}/ai-call-log.service.ts | 0 .../ai-provider-resolver.service.ts | 136 ++++++++++++ .../ai-resolution/ai-resolution.module.ts | 32 +++ .../core/src/modules/ai-resolution/index.ts | 17 ++ .../ai-resolution}/logging-ai-provider.ts | 4 +- .../ai-resolution/user-ai-settings.reader.ts | 60 ++++++ .../__tests__/async-job.service.spec.ts | 79 +++++++ .../__tests__/async-job.store.spec.ts | 85 ++++++++ .../async-job/async-job-store.module.ts | 16 ++ .../src/modules/async-job/async-job.entity.ts | 56 +++++ .../src/modules/async-job/async-job.mapper.ts | 34 +++ .../src/modules/async-job/async-job.module.ts | 21 ++ .../src/modules/async-job/async-job.policy.ts | 8 + .../modules/async-job/async-job.service.ts | 36 ++++ .../src/modules/async-job/async-job.store.ts | 42 ++++ packages/core/src/modules/async-job/index.ts | 7 + .../authorization/authorization.service.ts | 33 ++- .../cv-template/cv-renderer.service.ts | 2 +- packages/core/src/modules/messenger/index.ts | 7 +- .../messenger/message-bus.interface.ts | 11 + .../messenger/messages/parse-cv.message.ts | 28 +++ .../src/modules/messenger/messenger.module.ts | 5 +- .../file-upload/src/extractor-registry.ts | 11 +- packages/file-upload/src/index.ts | 3 +- packages/file-upload/src/validators.ts | 21 ++ packages/handlers/package.json | 10 +- .../src/__tests__/parse-cv.handler.spec.ts | 168 +++++++++++++++ packages/handlers/src/handlers.module.ts | 22 +- packages/handlers/src/index.ts | 1 + packages/handlers/src/parse-cv.handler.ts | 64 ++++++ packages/handlers/src/text-source.resolver.ts | 51 +++++ .../handlers/src/user-cv-parser.service.ts | 30 +++ packages/handlers/tsconfig.json | 7 +- packages/handlers/vitest.config.ts | 11 + pnpm-lock.yaml | 51 +++++ 70 files changed, 2610 insertions(+), 281 deletions(-) delete mode 100644 apps/api/src/modules/admin/ai-call-log.module.ts create mode 100644 apps/api/src/modules/async-job/__tests__/typed-async-job.bundle.spec.ts create mode 100644 apps/api/src/modules/async-job/async-job-graphql.module.ts create mode 100644 apps/api/src/modules/async-job/async-job.resolver.ts create mode 100644 apps/api/src/modules/async-job/async-job.type.ts create mode 100644 apps/api/src/modules/async-job/typed-async-job.bundle.ts create mode 100644 apps/api/src/modules/async-job/typed-async-job.interceptor.ts create mode 100644 apps/api/src/modules/async-job/typed-async-job.shared.ts create mode 100644 apps/api/src/modules/cv-parser/__tests__/cv-parser-dispatch.service.spec.ts delete mode 100644 apps/api/src/modules/cv-parser/ai-provider-resolver.service.ts create mode 100644 apps/api/src/modules/cv-parser/cv-parser-dispatch.service.ts create mode 100644 apps/api/src/modules/cv-parser/enqueue-parse-cv.resolver.ts create mode 100644 apps/api/src/modules/cv-parser/parse-cv-job.resolver.ts create mode 100644 apps/api/src/modules/cv-parser/parse-cv-job.type.ts create mode 100644 apps/client/src/features/onboarding/mutations/enqueue-parse-file.graphql create mode 100644 apps/client/src/features/onboarding/mutations/useEnqueueParseFileMutation.ts create mode 100644 apps/client/src/features/onboarding/queries/parse-cv-job.graphql create mode 100644 apps/client/src/features/onboarding/queries/useParseCvJobQuery.ts create mode 100644 packages/core/prisma/migrations/20260513130000_add_async_jobs/migration.sql create mode 100644 packages/core/prisma/models/async-job.prisma rename {apps/api/src/modules/admin => packages/core/src/modules/ai-resolution}/ai-call-log-persistence.service.ts (90%) rename {apps/api/src/modules/admin => packages/core/src/modules/ai-resolution}/ai-call-log.service.ts (100%) create mode 100644 packages/core/src/modules/ai-resolution/ai-provider-resolver.service.ts create mode 100644 packages/core/src/modules/ai-resolution/ai-resolution.module.ts create mode 100644 packages/core/src/modules/ai-resolution/index.ts rename {apps/api/src/modules/admin => packages/core/src/modules/ai-resolution}/logging-ai-provider.ts (95%) create mode 100644 packages/core/src/modules/ai-resolution/user-ai-settings.reader.ts create mode 100644 packages/core/src/modules/async-job/__tests__/async-job.service.spec.ts create mode 100644 packages/core/src/modules/async-job/__tests__/async-job.store.spec.ts create mode 100644 packages/core/src/modules/async-job/async-job-store.module.ts create mode 100644 packages/core/src/modules/async-job/async-job.entity.ts create mode 100644 packages/core/src/modules/async-job/async-job.mapper.ts create mode 100644 packages/core/src/modules/async-job/async-job.module.ts create mode 100644 packages/core/src/modules/async-job/async-job.policy.ts create mode 100644 packages/core/src/modules/async-job/async-job.service.ts create mode 100644 packages/core/src/modules/async-job/async-job.store.ts create mode 100644 packages/core/src/modules/async-job/index.ts create mode 100644 packages/core/src/modules/messenger/message-bus.interface.ts create mode 100644 packages/core/src/modules/messenger/messages/parse-cv.message.ts create mode 100644 packages/handlers/src/__tests__/parse-cv.handler.spec.ts create mode 100644 packages/handlers/src/parse-cv.handler.ts create mode 100644 packages/handlers/src/text-source.resolver.ts create mode 100644 packages/handlers/src/user-cv-parser.service.ts create mode 100644 packages/handlers/vitest.config.ts diff --git a/apps/api/package.json b/apps/api/package.json index 739f162..db15fa8 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -115,6 +115,7 @@ "tsconfig-paths": "^4.2.0", "tslib": "^2.8.0", "typescript": "^5.6.3", + "unplugin-swc": "^1.5.9", "vitest": "^4.0.16" } } diff --git a/apps/api/schema.gql b/apps/api/schema.gql index 434786e..63a0ea1 100644 --- a/apps/api/schema.gql +++ b/apps/api/schema.gql @@ -103,6 +103,21 @@ type ApplicationStatus { updatedAt: Date! } +type AsyncJob { + completedAt: DateTime + error: String + id: ID! + kind: String! + result: JSON + status: AsyncJobStatus! +} + +enum AsyncJobStatus { + COMPLETED + FAILED + PENDING +} + type AuthenticationResponse { accessTokenExpiration: TokenExpiration! user: User! @@ -270,6 +285,10 @@ type EducationEdge { node: Education! } +type EnqueueAsyncJobResult { + jobId: ID! +} + type HealthResponse { status: String! timestamp: String! @@ -383,6 +402,8 @@ type Mutation { deleteSkill(id: String!): Boolean! deleteUserFile(id: String!): Boolean! deleteVacancy(id: String!): Boolean! + enqueueParseFile(content: String!, fileName: String!, mimeType: String!): EnqueueAsyncJobResult! + enqueueParseStory(storyText: String!): EnqueueAsyncJobResult! generatePdf(cvId: String!): Boolean! login(email: String!, password: String!): AuthenticationResponse! logout: Boolean! @@ -409,6 +430,7 @@ type Mutation { updateProfile(id: ID!, input: UpdateProfileInput!): ProfileType! updateRole(description: String, id: String!, name: String): Role! updateSkill(description: String, id: String!, name: String): Skill! + updateVacancy(applicationUrl: String, companyId: String, deadline: DateTime, description: String, id: String!, isActive: Boolean, isPublic: Boolean, jobTypeId: String, levelId: String, locationId: String, maxSalary: Float, minSalary: Float, requirements: String, roleId: String, title: String): Vacancy! uploadFile(input: UploadFileInput!): UserFileType! verifyEmail(email: String!, token: String!): AuthenticationResponse! } @@ -453,6 +475,14 @@ type PageInfo { startCursor: String } +type ParseCvJob { + completedAt: DateTime + error: String + id: ID! + result: ParsedCVDataWithResolutionType + status: AsyncJobStatus! +} + type ParsedCVDataWithResolutionType { education: [DraftEducationType!]! jobExperiences: [DraftJobExperienceType!]! @@ -492,6 +522,7 @@ type Query { aiCallLog(limit: Int, providerName: String, status: String): [AiCallLogEntryType!]! aiCallLogHistory(limit: Int, providerName: String, status: String): [AiCallLogEntryType!]! application(id: String!): Application + asyncJob(id: ID!): AsyncJob! companies(after: String, before: String, first: Int, last: Int, searchTerm: String): CompanyConnection! company(id: String!): Company! cv(id: String!): CV @@ -508,6 +539,7 @@ type Query { myUserFiles: [UserFileType!]! onboardingStatus: [OnboardingStepResult!]! organization(id: String!): Organization + parseCvJob(id: ID!): ParseCvJob! platformAiAvailable: PlatformAiStatus! platformCapabilities: PlatformCapabilities! profile(id: ID!): ProfileType diff --git a/apps/api/src/modules/admin/admin.module.ts b/apps/api/src/modules/admin/admin.module.ts index 495577c..f225f31 100644 --- a/apps/api/src/modules/admin/admin.module.ts +++ b/apps/api/src/modules/admin/admin.module.ts @@ -1,5 +1,5 @@ import { AIModule } from "@cv/ai-provider"; -import { BaseModule } from "@cv/core"; +import { AIResolutionModule, BaseModule } from "@cv/core"; import { Module } from "@nestjs/common"; import { ApplicationStatusModule } from "@/modules/application/application-status/application-status.module"; import { DatabaseModule } from "@/modules/database/database.module"; @@ -17,6 +17,7 @@ import { QueueMonitorService } from "./queue-monitor.service"; @Module({ imports: [ AIModule.forConfig(), + AIResolutionModule, BaseModule, DatabaseModule, SkillModule, diff --git a/apps/api/src/modules/admin/admin.resolver.ts b/apps/api/src/modules/admin/admin.resolver.ts index 120edf6..62c372f 100644 --- a/apps/api/src/modules/admin/admin.resolver.ts +++ b/apps/api/src/modules/admin/admin.resolver.ts @@ -4,7 +4,7 @@ import { registeredProviderTypes, } from "@cv/ai-provider"; import { AdminGuard, JwtAuthGuard, VerifiedScopeGuard } from "@cv/auth"; -import { PageInfo } from "@cv/core"; +import { AiCallLogService, PageInfo } from "@cv/core"; import { Inject, Optional, UseGuards } from "@nestjs/common"; import { Args, Int, Query, Resolver } from "@nestjs/graphql"; import { @@ -15,7 +15,6 @@ import { SystemStatus, WorkerHealth, } from "./admin.type"; -import { AiCallLogService } from "./ai-call-log.service"; import { QueueMonitorService } from "./queue-monitor.service"; @Resolver() diff --git a/apps/api/src/modules/admin/ai-call-log.module.ts b/apps/api/src/modules/admin/ai-call-log.module.ts deleted file mode 100644 index 04e377a..0000000 --- a/apps/api/src/modules/admin/ai-call-log.module.ts +++ /dev/null @@ -1,12 +0,0 @@ -import { Global, Module } from "@nestjs/common"; -import { DatabaseModule } from "@/modules/database/database.module"; -import { AiCallLogService } from "./ai-call-log.service"; -import { AiCallLogPersistenceService } from "./ai-call-log-persistence.service"; - -@Global() -@Module({ - imports: [DatabaseModule], - providers: [AiCallLogPersistenceService, AiCallLogService], - exports: [AiCallLogService], -}) -export class AiCallLogModule {} diff --git a/apps/api/src/modules/app.module.ts b/apps/api/src/modules/app.module.ts index 3a7a0f8..3ea3b36 100644 --- a/apps/api/src/modules/app.module.ts +++ b/apps/api/src/modules/app.module.ts @@ -1,5 +1,6 @@ import { AuthModule } from "@cv/auth"; import { + AIResolutionModule, AuthorizationModule, BaseModule, DatabaseModule, @@ -34,9 +35,9 @@ import { } from "@/config/throttler.guard"; import { SeedModule } from "@/seed/seed.module"; import { AdminModule } from "./admin/admin.module"; -import { AiCallLogModule } from "./admin/ai-call-log.module"; import { AppModule as AppModuleComponent } from "./app/app.module"; import { ApplicationModule } from "./application/application.module"; +import { AsyncJobGraphQLModule } from "./async-job/async-job-graphql.module"; import { ApplicationStatusModule } from "./application/application-status/application-status.module"; import { AuthenticationModule } from "./authentication/authentication.module"; import { CurrentUserModule } from "./current-user/current-user.module"; @@ -107,7 +108,7 @@ import { VacancyModule } from "./vacancies/vacancy.module"; useClass: FileStorageConfig, }), EventsModule, - AiCallLogModule, + AIResolutionModule, AppConfigModule, BaseModule, DatabaseModule, @@ -128,6 +129,7 @@ import { VacancyModule } from "./vacancies/vacancy.module"; ApplicationStatusModule, ApplicationModule, CVParserModule, + AsyncJobGraphQLModule, CVTemplateModule, DataImportModule, ProjectQMessagingModule.forRoot({ queueName: "default" }), diff --git a/apps/api/src/modules/async-job/__tests__/typed-async-job.bundle.spec.ts b/apps/api/src/modules/async-job/__tests__/typed-async-job.bundle.spec.ts new file mode 100644 index 0000000..3b48ac8 --- /dev/null +++ b/apps/api/src/modules/async-job/__tests__/typed-async-job.bundle.spec.ts @@ -0,0 +1,193 @@ +import { ApolloDriver, type ApolloDriverConfig } from "@nestjs/apollo"; +import type { INestApplication } from "@nestjs/common"; +import { + Field, + GraphQLModule, + ID, + ObjectType, + Resolver, +} from "@nestjs/graphql"; +import { Test } from "@nestjs/testing"; +import { type AsyncJob, AsyncJobKind } from "@prisma/client"; +import request from "supertest"; +import { z } from "zod/v4"; +import { AsyncJobService } from "@cv/core"; +import { + type CompletedPayloadOf, + createTypedAsyncJobBundle, +} from "../typed-async-job.bundle"; +import { TypedAsyncJobInterceptor } from "../typed-async-job.interceptor"; + +// Test-only GraphQL type. Mirrors the AsyncJobOutputType contract +// (`pending` / `failed` statics + a class-as-value reference) so the +// bundle factory accepts it. `completed` is local convenience used by +// the resolver body to surface the payload into the response shape. +@ObjectType() +class TestJobPayloadType { + @Field(() => ID) id!: string; + @Field({ nullable: true }) value?: string; + @Field({ nullable: true }) error?: string; + @Field({ nullable: true }) completedAt?: Date; + + static pending(id: string): TestJobPayloadType { + return Object.assign(new TestJobPayloadType(), { id }); + } + static failed( + id: string, + error: string, + completedAt: Date, + ): TestJobPayloadType { + return Object.assign(new TestJobPayloadType(), { id, error, completedAt }); + } + static completed( + id: string, + value: string, + completedAt: Date, + ): TestJobPayloadType { + return Object.assign(new TestJobPayloadType(), { id, value, completedAt }); + } +} + +const TestJob = createTypedAsyncJobBundle({ + kind: AsyncJobKind.PARSE_CV, + resultSchema: z.object({ value: z.string() }), + type: TestJobPayloadType, +}); +type TestPayload = CompletedPayloadOf; + +@Resolver(() => TestJobPayloadType) +class TestJobResolver { + // If NestJS ever silently reverts to binding `args.id` (string) to this + // param, the `.row.id` / `.data.value` reads below blow up at runtime. + // That's the load-bearing assertion of this spec. + @TestJob.Query() + async testJob( + @TestJob.Completed() job: TestPayload, + ): Promise { + return TestJobPayloadType.completed( + job.row.id, + job.data.value, + job.completedAt, + ); + } +} + +describe("createTypedAsyncJobBundle (integration)", () => { + let app: INestApplication; + const completedAt = new Date("2026-05-13T16:00:00.000Z"); + + const completedRow: AsyncJob = { + id: "job-123", + kind: AsyncJobKind.PARSE_CV, + userId: "user-1", + result: { value: "hello" }, + error: null, + completedAt, + createdAt: new Date("2026-05-13T15:59:00.000Z"), + updatedAt: completedAt, + }; + + beforeAll(async () => { + const module = await Test.createTestingModule({ + imports: [ + GraphQLModule.forRoot({ + driver: ApolloDriver, + autoSchemaFile: true, + // Inject a fake authed user on every request - the interceptor + // needs `req.user.id` to call AsyncJobService.findByIdForUser. + context: ({ req }: { req: Record }) => ({ + req: { ...req, user: { id: "user-1" } }, + }), + }), + ], + providers: [ + TestJobResolver, + TypedAsyncJobInterceptor, + { + provide: AsyncJobService, + useValue: { + findByIdForUser: vi.fn().mockResolvedValue(completedRow), + }, + }, + ], + }).compile(); + + app = module.createNestApplication(); + await app.init(); + }); + + afterAll(async () => { + await app.close(); + }); + + it("binds the validated payload to @Completed(), not the raw id arg", async () => { + const response = await request(app.getHttpServer()) + .post("/graphql") + .send({ + query: + "query TestJob($id: ID!) { testJob(id: $id) { id value completedAt } }", + variables: { id: "job-123" }, + }); + + expect(response.body.errors).toBeUndefined(); + expect(response.body.data.testJob).toEqual({ + id: "job-123", + value: "hello", + completedAt: completedAt.toISOString(), + }); + }); + + it("declares `id: ID!` on the GraphQL schema for queries using the bundle", async () => { + const introspection = await request(app.getHttpServer()) + .post("/graphql") + .send({ + query: `{ + __type(name: "Query") { + fields { name args { name type { kind name ofType { name } } } } + } + }`, + }); + + interface IntrospectionTypeArg { + name: string; + type: { kind: string; name: string | null; ofType: { name: string } | null }; + } + interface IntrospectionTypeField { + name: string; + args: IntrospectionTypeArg[]; + } + const fields: IntrospectionTypeField[] = + introspection.body.data.__type.fields; + const testJobField = fields.find((f) => f.name === "testJob"); + expect(testJobField).toBeDefined(); + expect(testJobField?.args).toHaveLength(1); + expect(testJobField?.args[0]).toEqual({ + name: "id", + type: { + kind: "NON_NULL", + name: null, + ofType: { name: "ID" }, + }, + }); + }); + + // Load-bearing invariant from createTypedAsyncJobBundle: `Args` must be + // applied before `bindPayload` inside the composite param decorator so that + // the schema arg is registered AND the runtime value gets overwritten by + // the payload (last-write-wins on params[index]). If a future refactor + // flips the order, the integration test above catches it eventually - this + // test catches it immediately by reading the bundle's source file. + it("applies Args before bindPayload in the composite Completed() decorator", async () => { + const { readFile } = await import("node:fs/promises"); + const { join } = await import("node:path"); + const source = await readFile( + join(__dirname, "..", "typed-async-job.bundle.ts"), + "utf-8", + ); + const argsIdx = source.indexOf('Args("id"'); + const bindIdx = source.indexOf("bindPayload()(target"); + expect(argsIdx).toBeGreaterThanOrEqual(0); + expect(bindIdx).toBeGreaterThanOrEqual(0); + expect(argsIdx).toBeLessThan(bindIdx); + }); +}); diff --git a/apps/api/src/modules/async-job/async-job-graphql.module.ts b/apps/api/src/modules/async-job/async-job-graphql.module.ts new file mode 100644 index 0000000..3a9637a --- /dev/null +++ b/apps/api/src/modules/async-job/async-job-graphql.module.ts @@ -0,0 +1,20 @@ +import { AsyncJobModule } from "@cv/core"; +import { Module } from "@nestjs/common"; +import { AsyncJobResolver } from "./async-job.resolver"; +import { TypedAsyncJobInterceptor } from "./typed-async-job.interceptor"; + +/** + * GraphQL surface for AsyncJob. The data layer (entity, mapper, policy, + * service) lives in `@cv/core` as `AsyncJobModule`; this module is the + * api-only wiring for the `asyncJob` query and the bundle interceptor. + * + * Consumers needing the data layer (e.g. resolvers using the bundle in + * other feature modules) should import both this module (for the + * interceptor) AND `AsyncJobModule` from `@cv/core` directly. + */ +@Module({ + imports: [AsyncJobModule], + providers: [AsyncJobResolver, TypedAsyncJobInterceptor], + exports: [TypedAsyncJobInterceptor], +}) +export class AsyncJobGraphQLModule {} diff --git a/apps/api/src/modules/async-job/async-job.resolver.ts b/apps/api/src/modules/async-job/async-job.resolver.ts new file mode 100644 index 0000000..369f1b1 --- /dev/null +++ b/apps/api/src/modules/async-job/async-job.resolver.ts @@ -0,0 +1,21 @@ +import { JwtAuthGuard } from "@cv/auth"; +import { AsyncJobService, User as DomainUser } from "@cv/core"; +import { UseGuards } from "@nestjs/common"; +import { Args, ID, Query, Resolver } from "@nestjs/graphql"; +import { CurrentUser } from "@/modules/current-user/current-user.decorator"; +import { AsyncJobType } from "./async-job.type"; + +@Resolver(() => AsyncJobType) +@UseGuards(JwtAuthGuard) +export class AsyncJobResolver { + constructor(private readonly service: AsyncJobService) {} + + @Query(() => AsyncJobType) + async asyncJob( + @CurrentUser() user: DomainUser, + @Args("id", { type: () => ID }) id: string, + ): Promise { + const entity = await this.service.findByIdForUser(id, user); + return AsyncJobType.from(entity); + } +} diff --git a/apps/api/src/modules/async-job/async-job.type.ts b/apps/api/src/modules/async-job/async-job.type.ts new file mode 100644 index 0000000..d031b0b --- /dev/null +++ b/apps/api/src/modules/async-job/async-job.type.ts @@ -0,0 +1,44 @@ +import { AsyncJobEntity, AsyncJobStatus } from "@cv/core"; +import { Field, ID, ObjectType, registerEnumType } from "@nestjs/graphql"; +import GraphQLJSON from "graphql-type-json"; + +export { AsyncJobStatus }; +registerEnumType(AsyncJobStatus, { name: "AsyncJobStatus" }); + +@ObjectType("AsyncJob") +export class AsyncJobType { + @Field(() => ID) + id!: string; + + @Field() + kind!: string; + + @Field(() => AsyncJobStatus) + status!: AsyncJobStatus; + + @Field(() => GraphQLJSON, { nullable: true }) + result?: unknown; + + @Field({ nullable: true }) + error?: string; + + @Field({ nullable: true }) + completedAt?: Date; + + static from(entity: AsyncJobEntity): AsyncJobType { + const t = new AsyncJobType(); + t.id = entity.id; + t.kind = entity.kind; + t.status = entity.status; + if (entity.result !== null) { + t.result = entity.result; + } + if (entity.error !== null) { + t.error = entity.error; + } + if (entity.completedAt !== null) { + t.completedAt = entity.completedAt; + } + return t; + } +} diff --git a/apps/api/src/modules/async-job/typed-async-job.bundle.ts b/apps/api/src/modules/async-job/typed-async-job.bundle.ts new file mode 100644 index 0000000..86339b9 --- /dev/null +++ b/apps/api/src/modules/async-job/typed-async-job.bundle.ts @@ -0,0 +1,131 @@ +import { raise } from "@cv/utils"; +import { + applyDecorators, + createParamDecorator, + type ExecutionContext, + SetMetadata, + UseInterceptors, +} from "@nestjs/common"; +import { Args, GqlExecutionContext, ID, Query } from "@nestjs/graphql"; +import type { z } from "zod/v4"; +import { TypedAsyncJobInterceptor } from "./typed-async-job.interceptor"; +import { + type CompletedPayload, + type CompletedPayloadOf, + TYPED_ASYNC_JOB_META, + TYPED_ASYNC_JOB_PAYLOAD, + type TypedAsyncJobBundleOpts, +} from "./typed-async-job.shared"; + +export type { + AsyncJobOutputType, + AsyncJobTypeFactory, + CompletedPayload, + CompletedPayloadOf, + TypedAsyncJobBundleOpts, +} from "./typed-async-job.shared"; + +/** + * Builds a paired `{ Query, Completed }` decorator set bound to one + * async-job kind. The method decorator owns the GraphQL surface + lifecycle + * (pending / failed / schema-mismatch return early). The param decorator + * surfaces the validated payload to the resolver body, which only runs on + * completed. + * + * @example + * const ParseCvJob = createTypedAsyncJobBundle({ + * kind: AsyncJobKind.PARSE_CV, + * resultSchema: ParsedCVDataSchema, + * type: ParseCvJobType, + * }); + * type ParseCvJobPayload = CompletedPayloadOf; + * + * ${'@'}Resolver() + * class Resolver { + * ${'@'}ParseCvJob.Query() + * async parseCvJob(${'@'}ParseCvJob.Completed() job: ParseCvJobPayload) { + * ... + * } + * } + * + * `Completed()` is a composite param decorator: it declares the `id: ID!` + * arg on the GraphQL field AND binds the validated payload at runtime, so + * resolvers do not need a ghost `_id` param just to satisfy the schema. + */ +export const createTypedAsyncJobBundle = ( + opts: TypedAsyncJobBundleOpts, +) => { + type Payload = CompletedPayload>; + + const QueryDecorator = () => + applyDecorators( + Query(() => opts.type), + SetMetadata(TYPED_ASYNC_JOB_META, opts), + UseInterceptors(TypedAsyncJobInterceptor), + ); + + /** + * Composite param decorator: declares `id: ID!` on the GraphQL field AND + * binds the validated completed-payload to the resolver param at runtime. + * Folds away the ghost `_id` param that resolvers would otherwise need + * just for schema declaration. + * + * Implementation: NestJS resolves each method param by iterating the + * `__routeArguments__`-style metadata array in insertion order, with the + * last matching entry winning at `params[index]`. `@Args(...)` is applied + * first - it emits the schema arg AND a runtime binding to `gqlArgs.id`. + * The payload-binding param decorator runs second and overwrites the + * runtime value with the validated payload. The schema builder still + * sees the `@Args` entry, so `id: ID!` shows up on the field. + * + * The decorator-application order inside this factory is load-bearing. + * Don't reorder without re-running typed-async-job.bundle.spec.ts - + * the spec asserts the resolver param is the payload object, not the + * raw `args.id` string, so it fails loudly if NestJS's param-resolution + * semantics change on upgrade. + */ + const bindPayload = createParamDecorator( + (_data, ctx: ExecutionContext) => { + const gqlCtx = GqlExecutionContext.create(ctx).getContext< + Record + >(); + const payload = gqlCtx[TYPED_ASYNC_JOB_PAYLOAD]; + if (!payload) { + raise( + "TypedAsyncJobInterceptor did not populate the completed payload", + ); + } + return payload; + }, + ); + + const CompletedDecorator = (): ParameterDecorator => { + return (target, key, index) => { + Args("id", { type: () => ID })(target, key, index); + bindPayload()(target, key, index); + }; + }; + + return { + Query: QueryDecorator, + Completed: CompletedDecorator, + /** + * Phantom property: declared but never assigned. Carries the payload + * type so `CompletedPayloadOf` can extract it without + * relying on a runtime value. `readonly` + optional means there is no + * `undefined` hiding under the type either. + */ + __payload: undefined as unknown as Payload, + } as const satisfies BundleShape; +}; + +/** + * Internal: shape of the bundle's return value. Exposed so the satisfies + * check above pins the return without leaking the optional marker into the + * inferred type at the call site. + */ +interface BundleShape { + Query: () => MethodDecorator; + Completed: (...args: unknown[]) => ParameterDecorator; + readonly __payload?: TPayload; +} diff --git a/apps/api/src/modules/async-job/typed-async-job.interceptor.ts b/apps/api/src/modules/async-job/typed-async-job.interceptor.ts new file mode 100644 index 0000000..dc5b0c6 --- /dev/null +++ b/apps/api/src/modules/async-job/typed-async-job.interceptor.ts @@ -0,0 +1,113 @@ +import { AsyncJobService, AsyncJobStatus } from "@cv/core"; +import { raise } from "@cv/utils"; +import { + BadRequestException, + type CallHandler, + type ExecutionContext, + Injectable, + type NestInterceptor, +} from "@nestjs/common"; +import { Reflector } from "@nestjs/core"; +import { GqlExecutionContext } from "@nestjs/graphql"; +import { type Observable, of } from "rxjs"; +import type { z } from "zod/v4"; +import { + type CompletedPayload, + TYPED_ASYNC_JOB_META, + TYPED_ASYNC_JOB_PAYLOAD, + type TypedAsyncJobBundleOpts, +} from "./typed-async-job.shared"; + +/** + * Reads the bundle metadata attached by `@bundle.Query()`, loads the AsyncJob + * entity, and routes pending / failed / schema-mismatch branches back to the + * caller via `of(...)` so the resolver method body is bypassed entirely. + * Only the completed branch lets `next.handle()` proceed; the validated + * payload is stashed on the GraphQL context for `@bundle.Completed()` to + * read. + */ +@Injectable() +export class TypedAsyncJobInterceptor implements NestInterceptor { + constructor( + private readonly asyncJobs: AsyncJobService, + private readonly reflector: Reflector, + ) {} + + async intercept( + ctx: ExecutionContext, + next: CallHandler, + ): Promise> { + const meta = this.reflector.get< + TypedAsyncJobBundleOpts + >(TYPED_ASYNC_JOB_META, ctx.getHandler()); + + if (!meta) { + return next.handle(); + } + + const gql = GqlExecutionContext.create(ctx); + // The resolver hosting `@bundle.Query()` is expected to apply JwtAuthGuard + // and the schema declares `id: ID!`. Failures here are server bugs, not + // client errors - the loud invariant violation is the point. + const user = + gql.getContext().req?.user ?? + raise( + "TypedAsyncJobInterceptor: req.user missing; resolver must apply JwtAuthGuard", + ); + const id = + gql.getArgs<{ id?: string }>().id ?? + raise( + "TypedAsyncJobInterceptor: id arg missing; bundle should have declared id: ID!", + ); + + const entity = await this.asyncJobs.findByIdForUser(id, user); + + if (entity.kind !== meta.kind) { + throw new BadRequestException( + `Job ${id} is not a ${meta.kind} job (kind=${entity.kind})`, + ); + } + + if (entity.status === AsyncJobStatus.PENDING) { + return of(meta.type.pending(entity.id)); + } + + // AsyncJobEntity guarantees completedAt is non-null for FAILED + COMPLETED. + const completedAt = + entity.completedAt ?? + raise( + "AsyncJobEntity status invariant broken: completedAt null on non-PENDING", + ); + + if (entity.status === AsyncJobStatus.FAILED) { + return of( + meta.type.failed( + entity.id, + entity.error ?? + raise("AsyncJobEntity status invariant broken: FAILED with null error"), + completedAt, + ), + ); + } + + const parsed = meta.resultSchema.safeParse(entity.result); + if (!parsed.success) { + return of( + meta.type.failed( + entity.id, + `Result failed schema validation: ${parsed.error.message}`, + completedAt, + ), + ); + } + + const gqlCtx = gql.getContext>(); + gqlCtx[TYPED_ASYNC_JOB_PAYLOAD] = { + data: parsed.data, + row: entity, + completedAt, + } satisfies CompletedPayload; + + return next.handle(); + } +} diff --git a/apps/api/src/modules/async-job/typed-async-job.shared.ts b/apps/api/src/modules/async-job/typed-async-job.shared.ts new file mode 100644 index 0000000..dad35b3 --- /dev/null +++ b/apps/api/src/modules/async-job/typed-async-job.shared.ts @@ -0,0 +1,53 @@ +import type { AsyncJobEntity } from "@cv/core"; +import type { Type } from "@nestjs/common"; +import type { AsyncJobKind } from "@prisma/client"; +import type { z } from "zod/v4"; + +/** Reflector metadata key carrying the bundle opts onto the @Query method. */ +export const TYPED_ASYNC_JOB_META = Symbol("TYPED_ASYNC_JOB_META"); + +/** + * GraphQL context key where the interceptor stashes the validated completed + * payload for the param decorator to read. Symbol-keyed so it can't collide + * with user-supplied context fields. + */ +export const TYPED_ASYNC_JOB_PAYLOAD = Symbol("TYPED_ASYNC_JOB_PAYLOAD"); + +export interface AsyncJobTypeFactory { + pending(id: string): TJob; + failed(id: string, error: string, completedAt: Date): TJob; +} + +export interface CompletedPayload { + data: TData; + row: AsyncJobEntity; + completedAt: Date; +} + +/** + * The GraphQL output type doubles as the factory for non-completed states. + * Constraining to both `Type` (Nest's class-as-value) and the factory + * interface lets the bundle accept a single `type` reference instead of + * passing the same class twice. + */ +export type AsyncJobOutputType = Type & AsyncJobTypeFactory; + +export interface TypedAsyncJobBundleOpts { + kind: AsyncJobKind; + resultSchema: TSchema; + type: AsyncJobOutputType; +} + +/** + * Extract the validated completed-payload type from a bundle so call sites + * can annotate the `@bundle.Completed()` param without re-deriving the + * schema's `z.output`. + * + * const ParseCvJob = createTypedAsyncJobBundle({ ... }); + * type ParseCvJobPayload = CompletedPayloadOf; + */ +export type CompletedPayloadOf = Bundle extends { + readonly __payload?: infer P; +} + ? P + : never; diff --git a/apps/api/src/modules/cv-parser/__tests__/cv-parser-dispatch.service.spec.ts b/apps/api/src/modules/cv-parser/__tests__/cv-parser-dispatch.service.spec.ts new file mode 100644 index 0000000..c6fec11 --- /dev/null +++ b/apps/api/src/modules/cv-parser/__tests__/cv-parser-dispatch.service.spec.ts @@ -0,0 +1,201 @@ +import { + type MessageBus, + ParseCVMessage, + PrismaService, + User, +} from "@cv/core"; +import type { FileStorage } from "@cv/file-storage"; +import { BadRequestException, Logger } from "@nestjs/common"; +import { AsyncJobKind } from "@prisma/client"; +import type { Mocked } from "vitest"; +import { CVParserDispatchService } from "../cv-parser-dispatch.service"; + +describe("CVParserDispatchService", () => { + let service: CVParserDispatchService; + let prisma: { asyncJob: { create: Mock; delete: Mock } }; + let storage: Mocked; + let messageBus: Mocked; + + const user = new User( + "user-1", + "Test User", + new Date(), + new Date(), + null, + ); + + const validPdf = (): Buffer => { + const header = Buffer.from("%PDF-1.4\n"); + return Buffer.concat([header, Buffer.alloc(1024)]); + }; + + beforeEach(() => { + prisma = { + asyncJob: { + create: vi.fn().mockResolvedValue({ id: "job-1" }), + delete: vi.fn().mockResolvedValue({ id: "job-1" }), + }, + }; + storage = { + write: vi.fn().mockResolvedValue(undefined), + read: vi.fn(), + delete: vi.fn(), + }; + messageBus = { + dispatch: vi.fn().mockResolvedValue(undefined), + }; + + service = new CVParserDispatchService( + prisma as unknown as PrismaService, + storage, + messageBus, + ); + }); + + describe("enqueueFile", () => { + it("creates a job, writes the file, and dispatches a parse-cv message", async () => { + const buffer = validPdf(); + + const result = await service.enqueueFile(user, { + buffer, + mimeType: "application/pdf", + originalName: "cv.pdf", + }); + + expect(result).toEqual({ jobId: "job-1" }); + expect(prisma.asyncJob.create).toHaveBeenCalledWith({ + data: { userId: "user-1", kind: AsyncJobKind.PARSE_CV }, + }); + expect(storage.write).toHaveBeenCalledWith( + "cv-uploads/job-1/cv.pdf", + buffer, + ); + expect(messageBus.dispatch).toHaveBeenCalledTimes(1); + const envelope = messageBus.dispatch.mock.calls[0][0]; + const parsed = ParseCVMessage.parse(envelope).data; + expect(parsed).toEqual({ + jobId: "job-1", + input: { + source: "file", + fileKey: "cv-uploads/job-1/cv.pdf", + mimeType: "application/pdf", + originalName: "cv.pdf", + }, + }); + expect(prisma.asyncJob.delete).not.toHaveBeenCalled(); + }); + + it("rejects invalid files before creating a row", async () => { + await expect( + service.enqueueFile(user, { + buffer: Buffer.from("hi"), + mimeType: "application/x-msdownload", + originalName: "virus.exe", + }), + ).rejects.toThrow(BadRequestException); + + expect(prisma.asyncJob.create).not.toHaveBeenCalled(); + expect(storage.write).not.toHaveBeenCalled(); + expect(messageBus.dispatch).not.toHaveBeenCalled(); + }); + + it("rolls the row back and rethrows when storage.write fails", async () => { + const writeErr = new Error("S3 unavailable"); + storage.write.mockRejectedValueOnce(writeErr); + + await expect( + service.enqueueFile(user, { + buffer: validPdf(), + mimeType: "application/pdf", + originalName: "cv.pdf", + }), + ).rejects.toBe(writeErr); + + expect(prisma.asyncJob.delete).toHaveBeenCalledWith({ + where: { id: "job-1" }, + }); + expect(messageBus.dispatch).not.toHaveBeenCalled(); + }); + + it("rolls the row back and rethrows when dispatch fails", async () => { + const dispatchErr = new Error("bus down"); + messageBus.dispatch.mockRejectedValueOnce(dispatchErr); + + await expect( + service.enqueueFile(user, { + buffer: validPdf(), + mimeType: "application/pdf", + originalName: "cv.pdf", + }), + ).rejects.toBe(dispatchErr); + + expect(prisma.asyncJob.delete).toHaveBeenCalledWith({ + where: { id: "job-1" }, + }); + }); + }); + + describe("enqueueStory", () => { + it("creates a job and dispatches a story-shaped message", async () => { + const result = await service.enqueueStory(user, "my career story"); + + expect(result).toEqual({ jobId: "job-1" }); + expect(storage.write).not.toHaveBeenCalled(); + const envelope = messageBus.dispatch.mock.calls[0][0]; + const parsed = ParseCVMessage.parse(envelope).data; + expect(parsed).toEqual({ + jobId: "job-1", + input: { source: "story", text: "my career story" }, + }); + }); + + it("rejects empty story text before creating a row", async () => { + await expect(service.enqueueStory(user, " ")).rejects.toThrow( + BadRequestException, + ); + expect(prisma.asyncJob.create).not.toHaveBeenCalled(); + }); + + it("rejects story text over 50k characters", async () => { + await expect( + service.enqueueStory(user, "a".repeat(50_001)), + ).rejects.toThrow(/maximum length/); + expect(prisma.asyncJob.create).not.toHaveBeenCalled(); + }); + + it("rolls the row back and rethrows when dispatch fails", async () => { + const dispatchErr = new Error("bus down"); + messageBus.dispatch.mockRejectedValueOnce(dispatchErr); + + await expect( + service.enqueueStory(user, "my career story"), + ).rejects.toBe(dispatchErr); + + expect(prisma.asyncJob.delete).toHaveBeenCalledWith({ + where: { id: "job-1" }, + }); + }); + }); + + describe("rollback failure", () => { + it("logs a warning but rethrows the original error when the delete itself fails", async () => { + const dispatchErr = new Error("bus down"); + const rollbackErr = new Error("db gone too"); + messageBus.dispatch.mockRejectedValueOnce(dispatchErr); + prisma.asyncJob.delete.mockRejectedValueOnce(rollbackErr); + + const warnSpy = vi.spyOn(Logger.prototype, "warn").mockImplementation( + () => undefined, + ); + + await expect( + service.enqueueStory(user, "my career story"), + ).rejects.toBe(dispatchErr); + + expect(warnSpy).toHaveBeenCalledWith( + expect.stringContaining("Failed to roll back orphan AsyncJob job-1"), + ); + warnSpy.mockRestore(); + }); + }); +}); diff --git a/apps/api/src/modules/cv-parser/__tests__/cv-parser.service.spec.ts b/apps/api/src/modules/cv-parser/__tests__/cv-parser.service.spec.ts index 28cc841..7964649 100644 --- a/apps/api/src/modules/cv-parser/__tests__/cv-parser.service.spec.ts +++ b/apps/api/src/modules/cv-parser/__tests__/cv-parser.service.spec.ts @@ -1,11 +1,10 @@ import { readFileSync } from "node:fs"; import { join } from "node:path"; -import { User } from "@cv/core"; -import { TextExtractorRegistry } from "@cv/file-upload"; +import { type AIProviderResolver, User } from "@cv/core"; +import type { ExtractsText } from "@cv/file-upload"; import { Mocked } from "vitest"; -import { AIProviderResolverService } from "../ai-provider-resolver.service"; import { CVParserService } from "../cv-parser.service"; -import { EntityResolverService } from "../entity-resolver.service"; +import type { EntityResolver } from "../entity-resolver.service"; import { mockJaneDoeParsedCV, mockJohnSmithParsedCV, @@ -25,10 +24,8 @@ vi.mock("@cv/ai-parser", async (importOriginal) => { describe("CVParserService", () => { let service: CVParserService; - let mockTextExtractorRegistry: { - extract: Mock; - }; - let mockEntityResolver: Mocked; + let mockTextExtractorRegistry: Mocked; + let mockEntityResolver: Mocked; const mockUser = new User( "user-1", @@ -59,16 +56,60 @@ describe("CVParserService", () => { resolveSkill: vi.fn(), resolveSkills: vi.fn(), resolveInstitution: vi.fn(), - } as unknown as Mocked; + resolveAll: vi.fn(), + }; - const mockProviderResolver = { + // Service now delegates to resolveAll, which in production calls the + // per-entity resolvers. Mirror that fan-out here so the existing + // per-entity mocks still drive the test. + mockEntityResolver.resolveAll.mockImplementation(async (parsed) => { + const jobExperiences = await Promise.all( + parsed.jobExperiences.map(async (job) => { + const [company, role, level, skills] = await Promise.all([ + mockEntityResolver.resolveCompany(job.companyName), + mockEntityResolver.resolveRole(job.roleName), + mockEntityResolver.resolveLevel(job.levelName), + mockEntityResolver.resolveSkills(job.skills ?? []), + ]); + return { + company, + role, + level, + skills, + startDate: new Date(job.startDate), + endDate: job.endDate ? new Date(job.endDate) : null, + description: job.description ?? null, + }; + }), + ); + const education = await Promise.all( + parsed.education.map(async (edu) => { + const [institution, skills] = await Promise.all([ + mockEntityResolver.resolveInstitution(edu.institutionName), + mockEntityResolver.resolveSkills(edu.skills ?? []), + ]); + return { + institution, + degree: edu.degree, + fieldOfStudy: edu.fieldOfStudy ?? null, + skills, + startDate: new Date(edu.startDate), + endDate: edu.endDate ? new Date(edu.endDate) : null, + description: edu.description ?? null, + }; + }), + ); + return { jobExperiences, education }; + }); + + const mockProviderResolver: AIProviderResolver = { resolveForUser: vi.fn().mockResolvedValue({ generateText: vi.fn(), }), - } as unknown as AIProviderResolverService; + }; service = new CVParserService( - mockTextExtractorRegistry as unknown as TextExtractorRegistry, + mockTextExtractorRegistry, mockEntityResolver, mockProviderResolver, ); diff --git a/apps/api/src/modules/cv-parser/ai-provider-resolver.service.ts b/apps/api/src/modules/cv-parser/ai-provider-resolver.service.ts deleted file mode 100644 index f2a6d12..0000000 --- a/apps/api/src/modules/cv-parser/ai-provider-resolver.service.ts +++ /dev/null @@ -1,122 +0,0 @@ -import { randomUUID } from "node:crypto"; -import { - AI_PROVIDER, - AIProvider, - AnthropicProvider, - OpenAIProvider, -} from "@cv/ai-provider"; -import { ClockService, User } from "@cv/core"; -import { - BadRequestException, - ForbiddenException, - Inject, - Injectable, - Optional, -} from "@nestjs/common"; -import { AiCallLogService } from "@/modules/admin/ai-call-log.service"; -import { LoggingAIProvider } from "@/modules/admin/logging-ai-provider"; -import { UserAiSettingsService } from "@/modules/user-settings/user-ai-settings.service"; - -const DEFAULT_BASE_URLS: Record = { - anthropic: "https://api.anthropic.com", - openai: "https://api.openai.com", -}; - -/** - * Resolves the correct AIProvider for a given user based on their AI preference. - */ -@Injectable() -export class AIProviderResolverService { - constructor( - @Inject(AI_PROVIDER) - @Optional() - private readonly globalProvider: AIProvider | undefined, - private readonly userAiSettingsService: UserAiSettingsService, - private readonly aiCallLogService: AiCallLogService, - private readonly clock: ClockService, - ) {} - - async resolveForUser(user: User): Promise { - try { - return await this.resolveProvider(user); - } catch (err) { - this.aiCallLogService.record({ - id: randomUUID(), - timestamp: this.clock.now().toISOString(), - providerName: "resolution", - durationMs: 0, - status: "error", - error: err instanceof Error ? err.message : String(err), - userId: user.id, - source: "resolution", - }); - throw err; - } - } - - private async resolveProvider(user: User): Promise { - const settings = await this.userAiSettingsService.getOrCreateSettings(user); - - if (settings.aiPreference === "NO_AI") { - throw new ForbiddenException( - "AI is disabled. Enable it in profile settings.", - ); - } - - if (settings.aiPreference === "PLATFORM") { - if (!this.globalProvider) { - throw new BadRequestException( - "Platform AI is not available. Configure your own provider in profile settings.", - ); - } - return new LoggingAIProvider( - this.globalProvider, - this.aiCallLogService, - this.clock, - { - userId: user.id, - source: "platform", - }, - ); - } - - const resolved = - await this.userAiSettingsService.resolveActiveProvider(user); - - return this.instantiateProvider(resolved, user.id); - } - - private instantiateProvider( - config: { - providerType: string; - decryptedApiKey: string; - model: string | null; - baseUrl: string | null; - }, - userId: string, - ): AIProvider { - const baseUrl = - config.baseUrl || DEFAULT_BASE_URLS[config.providerType] || ""; - - const providerConfig = { - baseUrl, - apiKey: config.decryptedApiKey, - ...(config.model != null && { model: config.model }), - }; - - const provider = (() => { - if (config.providerType === "anthropic") - return new AnthropicProvider(providerConfig); - if (config.providerType === "openai") - return new OpenAIProvider(providerConfig); - throw new BadRequestException( - `Unsupported provider type: ${config.providerType}`, - ); - })(); - - return new LoggingAIProvider(provider, this.aiCallLogService, this.clock, { - userId, - source: "byok", - }); - } -} diff --git a/apps/api/src/modules/cv-parser/cv-parser-dispatch.service.ts b/apps/api/src/modules/cv-parser/cv-parser-dispatch.service.ts new file mode 100644 index 0000000..70bcacf --- /dev/null +++ b/apps/api/src/modules/cv-parser/cv-parser-dispatch.service.ts @@ -0,0 +1,120 @@ +import { + InjectMessageBus, + type MessageBus, + type ParseCVInput, + ParseCVMessage, + PrismaService, + User, +} from "@cv/core"; +import { FILE_STORAGE, type FileStorage } from "@cv/file-storage"; +import { validateFile, validateStoryText } from "@cv/file-upload"; +import { BadRequestException, Inject, Injectable, Logger } from "@nestjs/common"; +import { AsyncJobKind } from "@prisma/client"; + +const fileStorageKey = (jobId: string, originalName: string): string => + `cv-uploads/${jobId}/${originalName}`; + +/** + * Dispatcher for the parse-cv background job. Returns the new `AsyncJob` id + * to the caller; the worker eventually writes `result`/`error`/`completedAt` + * back to the same row. + */ +@Injectable() +export class CVParserDispatchService { + private readonly logger = new Logger(CVParserDispatchService.name); + + constructor( + private readonly prisma: PrismaService, + @Inject(FILE_STORAGE) private readonly storage: FileStorage, + @InjectMessageBus() private readonly messageBus: MessageBus, + ) {} + + async enqueueFile( + user: User, + file: { + buffer: Buffer; + mimeType: string; + originalName: string; + }, + ): Promise<{ jobId: string }> { + const validation = validateFile({ + buffer: file.buffer, + mimeType: file.mimeType, + originalName: file.originalName, + sizeBytes: file.buffer.length, + }); + + if (!validation.valid) { + throw new BadRequestException(validation.error); + } + + return this.dispatch(user, async (jobId) => { + const fileKey = fileStorageKey(jobId, file.originalName); + await this.storage.write(fileKey, file.buffer); + return { + source: "file" as const, + fileKey, + mimeType: file.mimeType, + originalName: file.originalName, + }; + }); + } + + async enqueueStory(user: User, storyText: string): Promise<{ jobId: string }> { + const validation = validateStoryText(storyText); + if (!validation.valid) { + throw new BadRequestException(validation.error); + } + + return this.dispatch(user, async () => ({ + source: "story" as const, + text: storyText, + })); + } + + /** + * Create the `AsyncJob` row, run the caller-supplied prep (storage write + * for file source, nothing for story source), dispatch the message. Any + * failure between row creation and successful dispatch rolls the row back + * via `rollbackJob`. + */ + private async dispatch( + user: User, + prepareInput: (jobId: string) => Promise, + ): Promise<{ jobId: string }> { + const job = await this.prisma.asyncJob.create({ + data: { userId: user.id, kind: AsyncJobKind.PARSE_CV }, + }); + + try { + const input = await prepareInput(job.id); + await this.messageBus.dispatch( + ParseCVMessage.create({ jobId: job.id, input }), + ); + return { jobId: job.id }; + } catch (err) { + await this.rollbackJob(job.id); + throw err; + } + } + + /** + * If `messageBus.dispatch` (or the file-storage write that precedes it) + * fails after the row insert, the AsyncJob would sit in PENDING forever + * with no message backing it. Best-effort delete to avoid orphan rows. + * + * If the delete itself fails (DB went away at the same time as the + * dispatch failure), log it but swallow - we want the original dispatch + * error to reach the caller, not the secondary cleanup error. Operators + * pick up the orphan via the logged warning. + */ + private async rollbackJob(id: string): Promise { + await this.prisma.asyncJob.delete({ where: { id } }).catch((err) => { + this.logger.warn( + `Failed to roll back orphan AsyncJob ${id}: ${ + err instanceof Error ? err.message : String(err) + }`, + ); + }); + } +} diff --git a/apps/api/src/modules/cv-parser/cv-parser.module.ts b/apps/api/src/modules/cv-parser/cv-parser.module.ts index 36b8e5a..5cde03f 100644 --- a/apps/api/src/modules/cv-parser/cv-parser.module.ts +++ b/apps/api/src/modules/cv-parser/cv-parser.module.ts @@ -3,7 +3,7 @@ import { mkdirSync } from "node:fs"; import { join } from "node:path"; import { CVParserModule as CVParserCoreModule } from "@cv/ai-parser"; import { AIModule } from "@cv/ai-provider"; -import { BaseModule, DatabaseModule } from "@cv/core"; +import { AIResolutionModule, BaseModule, DatabaseModule } from "@cv/core"; import { FileExtractionModule } from "@cv/file-upload"; import { Module } from "@nestjs/common"; import { MulterModule } from "@nestjs/platform-express"; @@ -14,13 +14,16 @@ import { EducationModule } from "@/modules/education/education.module"; import { EmploymentModule } from "@/modules/job-experience/employment/employment.module"; import { ProfileModule } from "@/modules/profile/profile.module"; import { UserAiSettingsModule } from "@/modules/user-settings/user-ai-settings.module"; -import { AIProviderResolverService } from "./ai-provider-resolver.service"; +import { AsyncJobGraphQLModule } from "@/modules/async-job/async-job-graphql.module"; +import { CVParserDispatchService } from "./cv-parser-dispatch.service"; import { CVParserResolver } from "./cv-parser.resolver"; import { CVParserService } from "./cv-parser.service"; +import { EnqueueParseCvResolver } from "./enqueue-parse-cv.resolver"; import { EntityResolverService } from "./entity-resolver.service"; import { FileUploadController } from "./file-upload.controller"; import { UploadFileResolver } from "./graphql/upload-file.resolver"; import { ImportOnboardingStep } from "./onboarding/import.step"; +import { ParseCvJobResolver } from "./parse-cv-job.resolver"; @Module({ imports: [ @@ -29,11 +32,13 @@ import { ImportOnboardingStep } from "./onboarding/import.step"; AIModule.forConfig(), FileExtractionModule.forRoot(), DatabaseModule, + AIResolutionModule, UserAiSettingsModule, ProfileModule, EmploymentModule, EducationModule, DataImportModule, + AsyncJobGraphQLModule, MulterModule.register({ storage: diskStorage({ destination: (_req, _file, cb) => { @@ -73,9 +78,11 @@ import { ImportOnboardingStep } from "./onboarding/import.step"; ], providers: [ EntityResolverService, - AIProviderResolverService, CVParserService, + CVParserDispatchService, CVParserResolver, + EnqueueParseCvResolver, + ParseCvJobResolver, UploadFileResolver, ImportOnboardingStep, FileImportSource, diff --git a/apps/api/src/modules/cv-parser/cv-parser.service.ts b/apps/api/src/modules/cv-parser/cv-parser.service.ts index fa226dc..855d553 100644 --- a/apps/api/src/modules/cv-parser/cv-parser.service.ts +++ b/apps/api/src/modules/cv-parser/cv-parser.service.ts @@ -1,25 +1,22 @@ import { CVParserService as CVParser, ParsedCVData } from "@cv/ai-parser"; -import { User } from "@cv/core"; import { + type AIProviderResolver, + AIProviderResolverService, + User, +} from "@cv/core"; +import { + type ExtractsText, TEXT_EXTRACTOR_REGISTRY, - TextExtractorRegistry, validateFile, } from "@cv/file-upload"; import { Inject, Injectable } from "@nestjs/common"; -import { AIProviderResolverService } from "./ai-provider-resolver.service"; import { + type EntityResolver, EntityResolverService, - ResolvedEducation, - ResolvedJobExperience, + ParsedCVDataWithResolution, } from "./entity-resolver.service"; -/** - * Parsed CV data with resolved entities - */ -export interface ParsedCVDataWithResolution { - jobExperiences: ResolvedJobExperience[]; - education: ResolvedEducation[]; -} +export type { ParsedCVDataWithResolution }; const MAX_STORY_TEXT_LENGTH = 50_000; @@ -27,14 +24,16 @@ const MAX_STORY_TEXT_LENGTH = 50_000; export class CVParserService { constructor( @Inject(TEXT_EXTRACTOR_REGISTRY) - private readonly textExtractorRegistry: TextExtractorRegistry, - private readonly entityResolver: EntityResolverService, - private readonly providerResolver: AIProviderResolverService, + private readonly textExtractorRegistry: ExtractsText, + @Inject(EntityResolverService) + private readonly entityResolver: EntityResolver, + @Inject(AIProviderResolverService) + private readonly providerResolver: AIProviderResolver, ) {} /** Create a per-request CVParser using the user's resolved provider */ private async createParser(user: User): Promise { - const provider = await this.providerResolver.resolveForUser(user); + const provider = await this.providerResolver.resolveForUser(user.id); return new CVParser(provider); } @@ -112,71 +111,9 @@ export class CVParserService { return this.resolveEntities(parsed); } - /** - * Resolve entity names to existing database records - */ private async resolveEntities( parsed: ParsedCVData, ): Promise { - const [jobExperiences, education] = await Promise.all([ - this.resolveJobExperiences(parsed.jobExperiences), - this.resolveEducation(parsed.education), - ]); - - return { jobExperiences, education }; - } - - /** - * Resolve job experience entities - */ - private async resolveJobExperiences( - jobs: ParsedCVData["jobExperiences"], - ): Promise { - return Promise.all( - jobs.map(async (job) => { - const [company, role, level, skills] = await Promise.all([ - this.entityResolver.resolveCompany(job.companyName), - this.entityResolver.resolveRole(job.roleName), - this.entityResolver.resolveLevel(job.levelName), - this.entityResolver.resolveSkills(job.skills ?? []), - ]); - - return { - company, - role, - level, - skills, - startDate: new Date(job.startDate), - endDate: job.endDate ? new Date(job.endDate) : null, - description: job.description ?? null, - }; - }), - ); - } - - /** - * Resolve education entities - */ - private async resolveEducation( - education: ParsedCVData["education"], - ): Promise { - return Promise.all( - education.map(async (edu) => { - const [institution, skills] = await Promise.all([ - this.entityResolver.resolveInstitution(edu.institutionName), - this.entityResolver.resolveSkills(edu.skills ?? []), - ]); - - return { - institution, - degree: edu.degree, - fieldOfStudy: edu.fieldOfStudy ?? null, - skills, - startDate: new Date(edu.startDate), - endDate: edu.endDate ? new Date(edu.endDate) : null, - description: edu.description ?? null, - }; - }), - ); + return this.entityResolver.resolveAll(parsed); } } diff --git a/apps/api/src/modules/cv-parser/enqueue-parse-cv.resolver.ts b/apps/api/src/modules/cv-parser/enqueue-parse-cv.resolver.ts new file mode 100644 index 0000000..abf720d --- /dev/null +++ b/apps/api/src/modules/cv-parser/enqueue-parse-cv.resolver.ts @@ -0,0 +1,40 @@ +import { JwtAuthGuard } from "@cv/auth"; +import { User as DomainUser } from "@cv/core"; +import { UseGuards } from "@nestjs/common"; +import { Args, Field, ID, Mutation, ObjectType, Resolver } from "@nestjs/graphql"; +import { CurrentUser } from "@/modules/current-user/current-user.decorator"; +import { CVParserDispatchService } from "./cv-parser-dispatch.service"; + +@ObjectType("EnqueueAsyncJobResult") +class EnqueueAsyncJobResult { + @Field(() => ID) + jobId!: string; +} + +@Resolver() +@UseGuards(JwtAuthGuard) +export class EnqueueParseCvResolver { + constructor(private readonly dispatch: CVParserDispatchService) {} + + @Mutation(() => EnqueueAsyncJobResult) + async enqueueParseStory( + @CurrentUser() user: DomainUser, + @Args("storyText") storyText: string, + ): Promise { + return this.dispatch.enqueueStory(user, storyText); + } + + @Mutation(() => EnqueueAsyncJobResult) + async enqueueParseFile( + @CurrentUser() user: DomainUser, + @Args("fileName") fileName: string, + @Args("mimeType") mimeType: string, + @Args("content") base64Content: string, + ): Promise { + return this.dispatch.enqueueFile(user, { + buffer: Buffer.from(base64Content, "base64"), + mimeType, + originalName: fileName, + }); + } +} diff --git a/apps/api/src/modules/cv-parser/entity-resolver.service.ts b/apps/api/src/modules/cv-parser/entity-resolver.service.ts index 88a1a94..8f1d465 100644 --- a/apps/api/src/modules/cv-parser/entity-resolver.service.ts +++ b/apps/api/src/modules/cv-parser/entity-resolver.service.ts @@ -1,3 +1,4 @@ +import { ParsedCVData } from "@cv/ai-parser"; import { PrismaService } from "@cv/core"; import { Injectable } from "@nestjs/common"; @@ -9,6 +10,14 @@ export interface DraftEntity { name: string; } +/** + * Parsed CV data with all entities resolved against the database. + */ +export interface ParsedCVDataWithResolution { + jobExperiences: ResolvedJobExperience[]; + education: ResolvedEducation[]; +} + /** * Resolved job experience with draft entities */ @@ -35,12 +44,27 @@ export interface ResolvedEducation { description: string | null; } +/** + * Public surface of `EntityResolverService`. Consumers depend on this so tests + * can supply structural mocks without the class's private fields forcing + * nominal typing. + */ +export interface EntityResolver { + resolveCompany(name: string): Promise; + resolveRole(name: string): Promise; + resolveLevel(name?: string): Promise; + resolveSkill(name: string): Promise; + resolveSkills(names: string[]): Promise; + resolveInstitution(name: string): Promise; + resolveAll(parsed: ParsedCVData): Promise; +} + /** * Service for resolving entity names to existing database entities * Returns stubs with null IDs for unmatched names */ @Injectable() -export class EntityResolverService { +export class EntityResolverService implements EntityResolver { constructor(private readonly prisma: PrismaService) {} /** @@ -117,4 +141,53 @@ export class EntityResolverService { return Promise.all(uniqueNames.map((name) => this.resolveSkill(name))); } + + /** + * Resolve every entity referenced in a parsed CV. Used both inline by the + * sync parse path and at read-time by the async parseCvJob query. + */ + async resolveAll( + parsed: ParsedCVData, + ): Promise { + const [jobExperiences, education] = await Promise.all([ + Promise.all( + parsed.jobExperiences.map(async (job) => { + const [company, role, level, skills] = await Promise.all([ + this.resolveCompany(job.companyName), + this.resolveRole(job.roleName), + this.resolveLevel(job.levelName), + this.resolveSkills(job.skills ?? []), + ]); + return { + company, + role, + level, + skills, + startDate: new Date(job.startDate), + endDate: job.endDate ? new Date(job.endDate) : null, + description: job.description ?? null, + }; + }), + ), + Promise.all( + parsed.education.map(async (edu) => { + const [institution, skills] = await Promise.all([ + this.resolveInstitution(edu.institutionName), + this.resolveSkills(edu.skills ?? []), + ]); + return { + institution, + degree: edu.degree, + fieldOfStudy: edu.fieldOfStudy ?? null, + skills, + startDate: new Date(edu.startDate), + endDate: edu.endDate ? new Date(edu.endDate) : null, + description: edu.description ?? null, + }; + }), + ), + ]); + + return { jobExperiences, education }; + } } diff --git a/apps/api/src/modules/cv-parser/parse-cv-job.resolver.ts b/apps/api/src/modules/cv-parser/parse-cv-job.resolver.ts new file mode 100644 index 0000000..e1d239c --- /dev/null +++ b/apps/api/src/modules/cv-parser/parse-cv-job.resolver.ts @@ -0,0 +1,46 @@ +import { ParsedCVDataSchema } from "@cv/ai-parser"; +import { JwtAuthGuard } from "@cv/auth"; +import { UseGuards } from "@nestjs/common"; +import { Resolver } from "@nestjs/graphql"; +import { AsyncJobKind } from "@prisma/client"; +import { + type CompletedPayloadOf, + createTypedAsyncJobBundle, +} from "@/modules/async-job/typed-async-job.bundle"; +import { EntityResolverService } from "./entity-resolver.service"; +import { ParseCvJobType } from "./parse-cv-job.type"; +import { + DraftEducationType, + DraftJobExperienceType, +} from "./types/draft.types"; + +const ParseCvJob = createTypedAsyncJobBundle({ + kind: AsyncJobKind.PARSE_CV, + resultSchema: ParsedCVDataSchema, + type: ParseCvJobType, +}); + +type ParseCvJobPayload = CompletedPayloadOf; + +@Resolver(() => ParseCvJobType) +@UseGuards(JwtAuthGuard) +export class ParseCvJobResolver { + constructor(private readonly entityResolver: EntityResolverService) {} + + @ParseCvJob.Query() + async parseCvJob( + @ParseCvJob.Completed() job: ParseCvJobPayload, + ): Promise { + const resolved = await this.entityResolver.resolveAll(job.data); + return ParseCvJobType.completed( + job.row.id, + { + jobExperiences: resolved.jobExperiences.map( + DraftJobExperienceType.fromResolved, + ), + education: resolved.education.map(DraftEducationType.fromResolved), + }, + job.completedAt, + ); + } +} diff --git a/apps/api/src/modules/cv-parser/parse-cv-job.type.ts b/apps/api/src/modules/cv-parser/parse-cv-job.type.ts new file mode 100644 index 0000000..9c5b48a --- /dev/null +++ b/apps/api/src/modules/cv-parser/parse-cv-job.type.ts @@ -0,0 +1,62 @@ +import { Field, ID, ObjectType } from "@nestjs/graphql"; +import { AsyncJobStatus } from "@/modules/async-job/async-job.type"; +import { + DraftEducationType, + DraftJobExperienceType, + ParsedCVDataWithResolutionType, +} from "./types/draft.types"; + +export { AsyncJobStatus }; + +@ObjectType("ParseCvJob") +export class ParseCvJobType { + @Field(() => ID) + id!: string; + + @Field(() => AsyncJobStatus) + status!: AsyncJobStatus; + + @Field(() => ParsedCVDataWithResolutionType, { nullable: true }) + result?: ParsedCVDataWithResolutionType; + + @Field({ nullable: true }) + error?: string; + + @Field({ nullable: true }) + completedAt?: Date; + + static pending(id: string): ParseCvJobType { + const t = new ParseCvJobType(); + t.id = id; + t.status = AsyncJobStatus.PENDING; + return t; + } + + static failed(id: string, error: string, completedAt: Date): ParseCvJobType { + const t = new ParseCvJobType(); + t.id = id; + t.status = AsyncJobStatus.FAILED; + t.error = error; + t.completedAt = completedAt; + return t; + } + + static completed( + id: string, + result: { + jobExperiences: ReturnType[]; + education: ReturnType[]; + }, + completedAt: Date, + ): ParseCvJobType { + const t = new ParseCvJobType(); + t.id = id; + t.status = AsyncJobStatus.COMPLETED; + t.result = new ParsedCVDataWithResolutionType( + result.jobExperiences, + result.education, + ); + t.completedAt = completedAt; + return t; + } +} diff --git a/apps/api/src/modules/data-import/sources/file-import-source.ts b/apps/api/src/modules/data-import/sources/file-import-source.ts index 23d9aee..bfc6aa5 100644 --- a/apps/api/src/modules/data-import/sources/file-import-source.ts +++ b/apps/api/src/modules/data-import/sources/file-import-source.ts @@ -4,6 +4,7 @@ import { ParsedCVData, } from "@cv/ai-parser"; import { + AIProviderResolverService, EducationService, PrismaService, ProfileService, @@ -15,7 +16,6 @@ import { TextExtractorRegistry, } from "@cv/file-upload"; import { Inject, Injectable, Logger } from "@nestjs/common"; -import { AIProviderResolverService } from "@/modules/cv-parser/ai-provider-resolver.service"; import { DataImportSource } from "../data-import-source.interface"; /** @@ -134,7 +134,7 @@ export class FileImportSource implements DataImportSource { await onStatus("Analyzing with AI"); - const provider = await this.providerResolver.resolveForUser(user); + const provider = await this.providerResolver.resolveForUser(user.id); this.logger.log(`AI provider: ${provider.constructor.name}`); const parser = new CVParser(provider); diff --git a/apps/api/src/modules/user-settings/user-ai-settings.service.ts b/apps/api/src/modules/user-settings/user-ai-settings.service.ts index 7b8844e..34bd946 100644 --- a/apps/api/src/modules/user-settings/user-ai-settings.service.ts +++ b/apps/api/src/modules/user-settings/user-ai-settings.service.ts @@ -195,35 +195,4 @@ export class UserAiSettingsService { }, }); } - - /** - * Resolve the active provider's decrypted config for internal use. - * Not exposed via GraphQL. - */ - async resolveActiveProvider(user: User) { - const settings = await this.getOrCreateSettings(user); - - if (!settings.activeProviderId) { - throw new BadRequestException( - "No active AI provider. Configure one in profile settings.", - ); - } - - const provider = await this.prisma.userAiProvider.findFirst({ - where: { id: settings.activeProviderId, userId: user.id }, - }); - - if (!provider) { - throw new BadRequestException( - "Active AI provider no longer exists. Configure a new one in profile settings.", - ); - } - - return { - providerType: provider.providerType, - decryptedApiKey: this.encryption.decrypt(provider.encryptedApiKey), - model: provider.model, - baseUrl: provider.baseUrl, - }; - } } diff --git a/apps/api/vitest.config.ts b/apps/api/vitest.config.ts index 7d8be84..1d22183 100644 --- a/apps/api/vitest.config.ts +++ b/apps/api/vitest.config.ts @@ -1,7 +1,25 @@ import { resolve } from "node:path"; +import swc from "unplugin-swc"; import { defineConfig } from "vitest/config"; export default defineConfig({ + // NestJS's decorator-driven DI (`@Injectable`, `@Args`, `@Field`, ...) + // depends on `design:paramtypes` metadata emitted by the TS compiler. Vite's + // default esbuild transformer ignores tsconfig's `emitDecoratorMetadata`, + // which crashes the schema builder inside `reflectTypeFromMetadata` for + // anything that boots a `GraphQLModule`. SWC honors it, so we route TS + // through it for tests. + plugins: [ + swc.vite({ + jsc: { + parser: { syntax: "typescript", decorators: true }, + transform: { + legacyDecorator: true, + decoratorMetadata: true, + }, + }, + }), + ], resolve: { alias: { "@/": `${resolve(__dirname, "src")}/`, diff --git a/apps/client/src/features/onboarding/mutations/enqueue-parse-file.graphql b/apps/client/src/features/onboarding/mutations/enqueue-parse-file.graphql new file mode 100644 index 0000000..c7268fd --- /dev/null +++ b/apps/client/src/features/onboarding/mutations/enqueue-parse-file.graphql @@ -0,0 +1,13 @@ +mutation EnqueueParseFile( + $fileName: String! + $mimeType: String! + $content: String! +) { + enqueueParseFile( + fileName: $fileName + mimeType: $mimeType + content: $content + ) { + jobId + } +} diff --git a/apps/client/src/features/onboarding/mutations/useEnqueueParseFileMutation.ts b/apps/client/src/features/onboarding/mutations/useEnqueueParseFileMutation.ts new file mode 100644 index 0000000..116e62e --- /dev/null +++ b/apps/client/src/features/onboarding/mutations/useEnqueueParseFileMutation.ts @@ -0,0 +1,33 @@ +import { useMutation } from "@tanstack/react-query"; +import { + type EnqueueParseFileMutation, + useEnqueueParseFileMutation as useCodegenEnqueueParseFile, +} from "@/generated/graphql"; + +const readFileAsBase64 = (file: File): Promise => + new Promise((resolve, reject) => { + const reader = new FileReader(); + reader.onload = () => { + const result = reader.result as string; + resolve(result.split(",")[1] ?? ""); + }; + reader.onerror = () => reject(new Error("Failed to read file")); + reader.readAsDataURL(file); + }); + +/** + * Replacement for `useUploadFileMutation`. Enqueues a `parse-cv` async job + * and resolves with the job id; callers poll `useParseCvJobQuery(jobId)`. + */ +export const useEnqueueParseFileMutation = () => + useMutation({ + mutationKey: useCodegenEnqueueParseFile.getKey(), + mutationFn: async (file: File) => { + const content = await readFileAsBase64(file); + return useCodegenEnqueueParseFile.fetcher({ + fileName: file.name, + mimeType: file.type, + content, + })(); + }, + }); diff --git a/apps/client/src/features/onboarding/queries/parse-cv-job.graphql b/apps/client/src/features/onboarding/queries/parse-cv-job.graphql new file mode 100644 index 0000000..d7653ac --- /dev/null +++ b/apps/client/src/features/onboarding/queries/parse-cv-job.graphql @@ -0,0 +1,46 @@ +query ParseCvJob($id: ID!) { + parseCvJob(id: $id) { + id + status + error + completedAt + result { + jobExperiences { + company { + id + name + } + role { + id + name + } + level { + id + name + } + skills { + id + name + } + startDate + endDate + description + } + education { + institution { + id + name + } + degree + fieldOfStudy + skills { + id + name + } + startDate + endDate + description + } + } + } +} diff --git a/apps/client/src/features/onboarding/queries/useParseCvJobQuery.ts b/apps/client/src/features/onboarding/queries/useParseCvJobQuery.ts new file mode 100644 index 0000000..dd554c6 --- /dev/null +++ b/apps/client/src/features/onboarding/queries/useParseCvJobQuery.ts @@ -0,0 +1,26 @@ +import { + AsyncJobStatus, + type ParseCvJobQuery, + useParseCvJobQuery as useCodegenParseCvJob, +} from "@/generated/graphql"; + +export type { ParseCvJobQuery }; +export { AsyncJobStatus }; + +/** + * Polls the `parseCvJob` query while `jobId` is set and the job has not + * reached a terminal status. Defaults to a 2s interval; pass + * `refetchInterval: false` to disable polling once the caller has handled + * the terminal state. + */ +export const useParseCvJobQuery = ( + jobId: string | null, + options?: { enabled?: boolean; refetchInterval?: number | false }, +) => + useCodegenParseCvJob( + { id: jobId ?? "" }, + { + enabled: options?.enabled ?? Boolean(jobId), + refetchInterval: options?.refetchInterval, + }, + ); 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..b5f4440 --- /dev/null +++ b/packages/core/prisma/migrations/20260513130000_add_async_jobs/migration.sql @@ -0,0 +1,32 @@ +-- CreateEnum +CREATE TYPE "AsyncJobKind" AS ENUM ('PARSE_CV'); + +-- CreateTable +CREATE TABLE "async_jobs" ( + "id" TEXT NOT NULL, + "userId" TEXT NOT NULL, + "kind" "AsyncJobKind" NOT NULL, + "result" JSONB, + "error" TEXT, + "completedAt" TIMESTAMP(3), + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "async_jobs_pkey" PRIMARY KEY ("id") +); + +-- CreateIndex +CREATE INDEX "async_jobs_userId_kind_idx" ON "async_jobs"("userId", "kind"); + +-- AddForeignKey +ALTER TABLE "async_jobs" ADD CONSTRAINT "async_jobs_userId_fkey" FOREIGN KEY ("userId") REFERENCES "users"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddCheckConstraint: FAILED (error set) and COMPLETED (result set) both +-- require completedAt to be set. PENDING has all three null. The entity +-- constructor enforces these too, but the SQL check is defense-in-depth +-- against direct writes that bypass the ORM. +ALTER TABLE "async_jobs" ADD CONSTRAINT "async_jobs_status_invariants_check" + CHECK ( + ("error" IS NULL OR "completedAt" IS NOT NULL) + AND ("result" IS NULL OR "completedAt" IS NOT NULL) + ); diff --git a/packages/core/prisma/models/async-job.prisma b/packages/core/prisma/models/async-job.prisma new file mode 100644 index 0000000..5d6621b --- /dev/null +++ b/packages/core/prisma/models/async-job.prisma @@ -0,0 +1,37 @@ +/// Discriminator for the kind of work an `AsyncJob` represents. Adding a +/// new kind requires a Prisma migration; the trade-off vs a free-form +/// string is that the database refuses invalid values and downstream +/// callers (bundle decorator, handler routing) get a typed surface. +enum AsyncJobKind { + PARSE_CV +} + +/// Sidecar table that adds what project-q's `Message` does not own: user +/// association, result payload, and a free-text error message. The queue +/// itself (delivery, retry, the input envelope, queue-internal timestamps) +/// lives in `Message`. +/// +/// Lifecycle: the API generates `id` (cuid), inserts this row, then dispatches +/// the parse-cv envelope referencing the same id in its `jobId` payload field. +/// The handler looks the row up by id, runs the work, and updates +/// `result`/`error`/`completedAt`. project-q acks the related `Message` +/// independently. +/// +/// User-facing status is derived from the fields on this row alone: +/// `result` set => COMPLETED, `error` set => FAILED, both null => PENDING. +/// `Message.statusId` is queue-internal and not consulted by the API. +model AsyncJob { + id String @id @default(cuid()) + userId String + kind AsyncJobKind + result Json? + error String? + completedAt DateTime? + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + user User @relation(fields: [userId], references: [id], onDelete: Cascade) + + @@index([userId, 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/index.ts b/packages/core/src/index.ts index bd45763..0171fc7 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -1,4 +1,6 @@ +export * from "./modules/ai-resolution"; export * from "./modules/application"; +export * from "./modules/async-job"; export * from "./modules/auth"; export * from "./modules/authentication"; export * from "./modules/cv-template"; diff --git a/apps/api/src/modules/admin/ai-call-log-persistence.service.ts b/packages/core/src/modules/ai-resolution/ai-call-log-persistence.service.ts similarity index 90% rename from apps/api/src/modules/admin/ai-call-log-persistence.service.ts rename to packages/core/src/modules/ai-resolution/ai-call-log-persistence.service.ts index 770f5dc..645d126 100644 --- a/apps/api/src/modules/admin/ai-call-log-persistence.service.ts +++ b/packages/core/src/modules/ai-resolution/ai-call-log-persistence.service.ts @@ -1,5 +1,5 @@ import { Injectable, Logger } from "@nestjs/common"; -import { PrismaService } from "@/modules/database/prisma.service"; +import { PrismaService } from "../database/prisma.service"; import { AiCallLogEntry } from "./ai-call-log.service"; interface QueryOptions { @@ -45,7 +45,9 @@ export class AiCallLogPersistenceService { return this.prisma.aiCallLog.findMany({ where: { ...(options.status ? { status: options.status } : {}), - ...(options.providerName ? { providerName: options.providerName } : {}), + ...(options.providerName + ? { providerName: options.providerName } + : {}), }, orderBy: { createdAt: "desc" }, take: options.limit ?? 100, diff --git a/apps/api/src/modules/admin/ai-call-log.service.ts b/packages/core/src/modules/ai-resolution/ai-call-log.service.ts similarity index 100% rename from apps/api/src/modules/admin/ai-call-log.service.ts rename to packages/core/src/modules/ai-resolution/ai-call-log.service.ts diff --git a/packages/core/src/modules/ai-resolution/ai-provider-resolver.service.ts b/packages/core/src/modules/ai-resolution/ai-provider-resolver.service.ts new file mode 100644 index 0000000..e78c851 --- /dev/null +++ b/packages/core/src/modules/ai-resolution/ai-provider-resolver.service.ts @@ -0,0 +1,136 @@ +import { randomUUID } from "node:crypto"; +import { + AI_PROVIDER, + AIProvider, + AnthropicProvider, + OpenAIProvider, +} from "@cv/ai-provider"; +import { + BadRequestException, + ForbiddenException, + Inject, + Injectable, + Optional, +} from "@nestjs/common"; +import { AiPreference } from "@prisma/client"; +import { ClockService } from "../../shared/clock.service"; +import { AiCallLogService } from "./ai-call-log.service"; +import { LoggingAIProvider } from "./logging-ai-provider"; +import { UserAiSettingsReader } from "./user-ai-settings.reader"; + +interface ProviderConfig { + baseUrl: string; + apiKey: string; + model?: string; +} + +const DEFAULT_BASE_URLS: Record = { + anthropic: "https://api.anthropic.com", + openai: "https://api.openai.com", +}; + +const PROVIDER_FACTORIES: Record AIProvider> = { + anthropic: (cfg) => new AnthropicProvider(cfg), + openai: (cfg) => new OpenAIProvider(cfg), +}; + +/** + * Public surface of `AIProviderResolverService`. Consumers depend on this so + * tests can supply structural mocks without the class's private fields + * forcing nominal typing. + */ +export interface AIProviderResolver { + resolveForUser(userId: string): Promise; +} + +/** + * Resolves the correct AIProvider for a given user based on their AI + * preference. The returned provider is always wrapped in `LoggingAIProvider` + * so every completion is recorded to `AiCallLog`. + */ +@Injectable() +export class AIProviderResolverService implements AIProviderResolver { + constructor( + @Inject(AI_PROVIDER) + @Optional() + private readonly globalProvider: AIProvider | undefined, + private readonly settingsReader: UserAiSettingsReader, + private readonly aiCallLogService: AiCallLogService, + private readonly clock: ClockService, + ) {} + + async resolveForUser(userId: string): Promise { + try { + return await this.resolveProvider(userId); + } catch (err) { + this.aiCallLogService.record({ + id: randomUUID(), + timestamp: this.clock.now().toISOString(), + providerName: "resolution", + durationMs: 0, + status: "error", + error: err instanceof Error ? err.message : String(err), + userId, + source: "resolution", + }); + throw err; + } + } + + private async resolveProvider(userId: string): Promise { + const settings = await this.settingsReader.getOrCreateSettings(userId); + + if (settings.aiPreference === AiPreference.NO_AI) { + throw new ForbiddenException( + "AI is disabled. Enable it in profile settings.", + ); + } + + if (settings.aiPreference === AiPreference.PLATFORM) { + if (!this.globalProvider) { + throw new BadRequestException( + "Platform AI is not available. Configure your own provider in profile settings.", + ); + } + return new LoggingAIProvider( + this.globalProvider, + this.aiCallLogService, + this.clock, + { userId, source: "platform" }, + ); + } + + const resolved = await this.settingsReader.resolveActiveProvider(userId); + return this.instantiateProvider(resolved, userId); + } + + private instantiateProvider( + config: { + providerType: string; + decryptedApiKey: string; + model: string | null; + baseUrl: string | null; + }, + userId: string, + ): AIProvider { + const factory = PROVIDER_FACTORIES[config.providerType]; + if (!factory) { + throw new BadRequestException( + `Unsupported provider type: ${config.providerType}`, + ); + } + + const providerConfig: ProviderConfig = { + baseUrl: config.baseUrl || DEFAULT_BASE_URLS[config.providerType] || "", + apiKey: config.decryptedApiKey, + ...(config.model != null && { model: config.model }), + }; + + return new LoggingAIProvider( + factory(providerConfig), + this.aiCallLogService, + this.clock, + { userId, source: "byok" }, + ); + } +} diff --git a/packages/core/src/modules/ai-resolution/ai-resolution.module.ts b/packages/core/src/modules/ai-resolution/ai-resolution.module.ts new file mode 100644 index 0000000..d3c7484 --- /dev/null +++ b/packages/core/src/modules/ai-resolution/ai-resolution.module.ts @@ -0,0 +1,32 @@ +import { AIModule } from "@cv/ai-provider"; +import { Module } from "@nestjs/common"; +import { BaseModule } from "../../shared/base.module"; +import { UserModule } from "../auth/user/user.module"; +import { DatabaseModule } from "../database/database.module"; +import { AiCallLogPersistenceService } from "./ai-call-log-persistence.service"; +import { AiCallLogService } from "./ai-call-log.service"; +import { AIProviderResolverService } from "./ai-provider-resolver.service"; +import { UserAiSettingsReader } from "./user-ai-settings.reader"; + +/** + * Provides the shared AI provider resolution chain: reads user AI settings, + * instantiates the configured provider (Anthropic / OpenAI / platform + * LlamaCpp), wraps in `LoggingAIProvider`, and records each call to + * `AiCallLog`. Imported by `apps/api` (cv-parser, admin) and the worker + * `HandlersModule` so both audit-log the same way. + * + * The optional global `AI_PROVIDER` (for PLATFORM users) comes from + * `AIModule.forConfig()`; if `AI_PROVIDER` env var is unset, PLATFORM + * resolution fails with a BadRequestException. + */ +@Module({ + imports: [BaseModule, DatabaseModule, UserModule, AIModule.forConfig()], + providers: [ + UserAiSettingsReader, + AiCallLogPersistenceService, + AiCallLogService, + AIProviderResolverService, + ], + exports: [UserAiSettingsReader, AiCallLogService, AIProviderResolverService], +}) +export class AIResolutionModule {} diff --git a/packages/core/src/modules/ai-resolution/index.ts b/packages/core/src/modules/ai-resolution/index.ts new file mode 100644 index 0000000..f262fb6 --- /dev/null +++ b/packages/core/src/modules/ai-resolution/index.ts @@ -0,0 +1,17 @@ +export { + AiCallLogPersistenceService, +} from "./ai-call-log-persistence.service"; +export { + type AiCallLogEntry, + AiCallLogService, +} from "./ai-call-log.service"; +export { + type AIProviderResolver, + AIProviderResolverService, +} from "./ai-provider-resolver.service"; +export { AIResolutionModule } from "./ai-resolution.module"; +export { LoggingAIProvider } from "./logging-ai-provider"; +export { + type ResolvedProviderConfig, + UserAiSettingsReader, +} from "./user-ai-settings.reader"; diff --git a/apps/api/src/modules/admin/logging-ai-provider.ts b/packages/core/src/modules/ai-resolution/logging-ai-provider.ts similarity index 95% rename from apps/api/src/modules/admin/logging-ai-provider.ts rename to packages/core/src/modules/ai-resolution/logging-ai-provider.ts index 8ef70ea..9bb1cbb 100644 --- a/apps/api/src/modules/admin/logging-ai-provider.ts +++ b/packages/core/src/modules/ai-resolution/logging-ai-provider.ts @@ -5,7 +5,7 @@ import { AIProvider, AIProviderStatus, } from "@cv/ai-provider"; -import { ClockService } from "@cv/core"; +import { ClockService } from "../../shared/clock.service"; import { AiCallLogService } from "./ai-call-log.service"; interface LoggingOptions { @@ -15,7 +15,7 @@ interface LoggingOptions { /** * Thin wrapper around an AIProvider that records call timing and - * response metadata to the in-memory AiCallLogService. + * response metadata to AiCallLogService. */ export class LoggingAIProvider implements AIProvider { readonly name: string; diff --git a/packages/core/src/modules/ai-resolution/user-ai-settings.reader.ts b/packages/core/src/modules/ai-resolution/user-ai-settings.reader.ts new file mode 100644 index 0000000..53a357b --- /dev/null +++ b/packages/core/src/modules/ai-resolution/user-ai-settings.reader.ts @@ -0,0 +1,60 @@ +import { BadRequestException, Injectable } from "@nestjs/common"; +import { AiPreference } from "@prisma/client"; +import { TokenEncryptionService } from "../auth/user/token-encryption.service"; +import { PrismaService } from "../database/prisma.service"; + +export interface ResolvedProviderConfig { + providerType: string; + decryptedApiKey: string; + model: string | null; + baseUrl: string | null; +} + +/** + * Read-only view of `UserAiSettings` and `UserAiProvider` used by the AI + * provider resolver. CRUD (add/remove/update providers, set preference) stays + * in `apps/api`'s `UserAiSettingsService`; this reader is the slice the worker + * also needs. + */ +@Injectable() +export class UserAiSettingsReader { + constructor( + private readonly prisma: PrismaService, + private readonly encryption: TokenEncryptionService, + ) {} + + async getOrCreateSettings(userId: string) { + return this.prisma.userAiSettings.upsert({ + where: { userId }, + create: { userId, aiPreference: AiPreference.NO_AI }, + update: {}, + }); + } + + async resolveActiveProvider(userId: string): Promise { + const settings = await this.getOrCreateSettings(userId); + + if (!settings.activeProviderId) { + throw new BadRequestException( + "No active AI provider configured. Set one in profile settings.", + ); + } + + const provider = await this.prisma.userAiProvider.findFirst({ + where: { id: settings.activeProviderId, userId }, + }); + + if (!provider) { + throw new BadRequestException( + "Active AI provider no longer exists. Configure a new one in profile settings.", + ); + } + + return { + providerType: provider.providerType, + decryptedApiKey: this.encryption.decrypt(provider.encryptedApiKey), + model: provider.model, + baseUrl: provider.baseUrl, + }; + } +} diff --git a/packages/core/src/modules/async-job/__tests__/async-job.service.spec.ts b/packages/core/src/modules/async-job/__tests__/async-job.service.spec.ts new file mode 100644 index 0000000..137d3d8 --- /dev/null +++ b/packages/core/src/modules/async-job/__tests__/async-job.service.spec.ts @@ -0,0 +1,79 @@ +import { AsyncJobKind } from "@prisma/client"; +import { describe, expect, it, vi } from "vitest"; +import type { Mocked } from "vitest"; +import type { Authorizer } from "../../auth/authorization/authorization.service"; +import { CannotViewError } from "../../auth/errors/authorization.error"; +import { EntityNotFoundError } from "../../auth/errors/not-found.util"; +import { User } from "../../auth/user/user.entity"; +import { AsyncJobEntity } from "../async-job.entity"; +import { AsyncJobService } from "../async-job.service"; +import type { AsyncJobStore } from "../async-job.store"; + +describe("AsyncJobService.findByIdForUser", () => { + let service: AsyncJobService; + let store: { findById: Mocked }; + let authorization: Mocked; + + const owner = new User("user-1", "Owner", new Date(), new Date(), null); + const otherUser = new User("user-2", "Other", new Date(), new Date(), null); + + const completedAt = new Date("2026-05-13T16:00:00.000Z"); + const entity = new AsyncJobEntity( + "job-123", + "user-1", + AsyncJobKind.PARSE_CV, + { value: "hello" }, + null, + completedAt, + new Date("2026-05-13T15:59:00.000Z"), + completedAt, + ); + + beforeEach(() => { + store = { findById: vi.fn() }; + authorization = { + canView: vi.fn().mockResolvedValue(undefined), + canCreate: vi.fn(), + canUpdate: vi.fn(), + canDelete: vi.fn(), + isSameUser: vi.fn(), + }; + service = new AsyncJobService( + store as unknown as AsyncJobStore, + authorization, + ); + }); + + it("returns the entity when the user owns the row", async () => { + store.findById.mockResolvedValueOnce(entity); + + const result = await service.findByIdForUser("job-123", owner); + + expect(result).toBe(entity); + expect(authorization.canView).toHaveBeenCalledWith( + owner, + entity, + AsyncJobEntity, + ); + }); + + it("throws EntityNotFoundError without consulting authz when missing", async () => { + store.findById.mockResolvedValueOnce(null); + + await expect(service.findByIdForUser("missing", owner)).rejects.toThrow( + EntityNotFoundError, + ); + expect(authorization.canView).not.toHaveBeenCalled(); + }); + + it("lets the policy denial bubble up (CannotViewError) for a non-owner", async () => { + store.findById.mockResolvedValueOnce(entity); + authorization.canView.mockRejectedValueOnce( + new CannotViewError("AsyncJobEntity"), + ); + + await expect( + service.findByIdForUser("job-123", otherUser), + ).rejects.toThrow(CannotViewError); + }); +}); diff --git a/packages/core/src/modules/async-job/__tests__/async-job.store.spec.ts b/packages/core/src/modules/async-job/__tests__/async-job.store.spec.ts new file mode 100644 index 0000000..32a5b63 --- /dev/null +++ b/packages/core/src/modules/async-job/__tests__/async-job.store.spec.ts @@ -0,0 +1,85 @@ +import { AsyncJobKind, Prisma } from "@prisma/client"; +import { describe, expect, it, vi } from "vitest"; +import type { Mock } from "vitest"; +import type { ClockService } from "../../../shared/clock.service"; +import type { PrismaService } from "../../database/prisma.service"; +import { AsyncJobEntity } from "../async-job.entity"; +import { AsyncJobMapper } from "../async-job.mapper"; +import { AsyncJobStore } from "../async-job.store"; + +describe("AsyncJobStore", () => { + let store: AsyncJobStore; + let prisma: { asyncJob: { findUnique: Mock; update: Mock } }; + + const now = new Date("2026-05-14T13:00:00.000Z"); + const clock: ClockService = { now: () => now }; + + const completedAt = new Date("2026-05-13T16:00:00.000Z"); + const row = { + id: "job-123", + userId: "user-1", + kind: AsyncJobKind.PARSE_CV, + result: { value: "hello" }, + error: null, + completedAt, + createdAt: new Date("2026-05-13T15:59:00.000Z"), + updatedAt: completedAt, + }; + + beforeEach(() => { + prisma = { + asyncJob: { + findUnique: vi.fn(), + update: vi.fn().mockResolvedValue(row), + }, + }; + store = new AsyncJobStore( + prisma as unknown as PrismaService, + new AsyncJobMapper(), + clock, + ); + }); + + describe("findById", () => { + it("returns the mapped entity when present", async () => { + prisma.asyncJob.findUnique.mockResolvedValueOnce(row); + + const result = await store.findById("job-123"); + + expect(result).toBeInstanceOf(AsyncJobEntity); + expect(result?.id).toBe("job-123"); + }); + + it("returns null when missing (no throw)", async () => { + prisma.asyncJob.findUnique.mockResolvedValueOnce(null); + + const result = await store.findById("missing"); + + expect(result).toBeNull(); + }); + }); + + describe("markCompleted", () => { + it("sets result + completedAt and clears error", async () => { + const result = { jobExperiences: [], education: [] }; + + await store.markCompleted("job-123", result); + + expect(prisma.asyncJob.update).toHaveBeenCalledWith({ + where: { id: "job-123" }, + data: { result, error: null, completedAt: now }, + }); + }); + }); + + describe("markFailed", () => { + it("sets error + completedAt and clears any stale result", async () => { + await store.markFailed("job-123", "boom"); + + expect(prisma.asyncJob.update).toHaveBeenCalledWith({ + where: { id: "job-123" }, + data: { error: "boom", result: Prisma.DbNull, completedAt: now }, + }); + }); + }); +}); diff --git a/packages/core/src/modules/async-job/async-job-store.module.ts b/packages/core/src/modules/async-job/async-job-store.module.ts new file mode 100644 index 0000000..55ec73b --- /dev/null +++ b/packages/core/src/modules/async-job/async-job-store.module.ts @@ -0,0 +1,16 @@ +import { Module } from "@nestjs/common"; +import { DatabaseModule } from "../database/database.module"; +import { AsyncJobMapper } from "./async-job.mapper"; +import { AsyncJobStore } from "./async-job.store"; + +/** + * Authorization-free data layer for `AsyncJob`. Workers depend on this + * directly so they don't pull in `AuthorizationModule` + policy registry + * for findById / markCompleted / markFailed. + */ +@Module({ + imports: [DatabaseModule], + providers: [AsyncJobStore, AsyncJobMapper], + exports: [AsyncJobStore, AsyncJobMapper], +}) +export class AsyncJobStoreModule {} diff --git a/packages/core/src/modules/async-job/async-job.entity.ts b/packages/core/src/modules/async-job/async-job.entity.ts new file mode 100644 index 0000000..e7e9410 --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.entity.ts @@ -0,0 +1,56 @@ +import type { AsyncJobKind, Prisma } from "@prisma/client"; +import { BaseEntity } from "../../shared/base.entity"; + +export enum AsyncJobStatus { + PENDING = "PENDING", + COMPLETED = "COMPLETED", + FAILED = "FAILED", +} + +export class AsyncJobEntity extends BaseEntity { + constructor( + id: string, + public readonly userId: string, + public readonly kind: AsyncJobKind, + public readonly result: Prisma.JsonValue | null, + public readonly error: string | null, + public readonly completedAt: Date | null, + createdAt: Date, + updatedAt: Date, + ) { + super(id, createdAt, updatedAt); + // Enforce the FAILED/COMPLETED invariants at construction so downstream + // consumers (interceptor, type mapper) can trust `status` without + // null-fallback fabrication. A row that violates these is a data anomaly + // we want surfaced loudly at read time. + if (error !== null && completedAt === null) { + throw new Error( + `AsyncJobEntity invariant: error set but completedAt is null (id=${id})`, + ); + } + if (result !== null && completedAt === null) { + throw new Error( + `AsyncJobEntity invariant: result set but completedAt is null (id=${id})`, + ); + } + } + + /** + * Single source of truth for status derivation. The constructor enforces + * the underlying invariants so this getter is a pure projection over the + * fields (no fabrication needed at the call site). + * + * error set => FAILED (completedAt is guaranteed set) + * result set => COMPLETED (completedAt is guaranteed set) + * neither => PENDING + */ + get status(): AsyncJobStatus { + if (this.error !== null) { + return AsyncJobStatus.FAILED; + } + if (this.result !== null) { + return AsyncJobStatus.COMPLETED; + } + return AsyncJobStatus.PENDING; + } +} diff --git a/packages/core/src/modules/async-job/async-job.mapper.ts b/packages/core/src/modules/async-job/async-job.mapper.ts new file mode 100644 index 0000000..9dc1522 --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.mapper.ts @@ -0,0 +1,34 @@ +import { Injectable } from "@nestjs/common"; +import { Prisma } from "@prisma/client"; +import { BaseMapper } from "../../shared/mapper.interface"; +import { AsyncJobEntity } from "./async-job.entity"; + +type PrismaAsyncJob = Prisma.AsyncJobGetPayload; + +@Injectable() +export class AsyncJobMapper implements BaseMapper { + toDomain(prisma: null): null; + toDomain(prisma: PrismaAsyncJob): AsyncJobEntity; + toDomain(prisma: PrismaAsyncJob | null): AsyncJobEntity | null; + toDomain(prisma: PrismaAsyncJob | null): AsyncJobEntity | null { + if (!prisma) { + return null; + } + return new AsyncJobEntity( + prisma.id, + prisma.userId, + prisma.kind, + prisma.result, + prisma.error, + prisma.completedAt, + prisma.createdAt, + prisma.updatedAt, + ); + } + + mapToDomain(items: PrismaAsyncJob[]): AsyncJobEntity[] { + return items + .map((item) => this.toDomain(item)) + .filter((item): item is AsyncJobEntity => item !== null); + } +} diff --git a/packages/core/src/modules/async-job/async-job.module.ts b/packages/core/src/modules/async-job/async-job.module.ts new file mode 100644 index 0000000..61e7a03 --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.module.ts @@ -0,0 +1,21 @@ +import { Module } from "@nestjs/common"; +import { AuthorizationModule } from "../auth/authorization/authorization.module"; +import { AsyncJobStoreModule } from "./async-job-store.module"; +import { AsyncJobPolicy } from "./async-job.policy"; +import { AsyncJobService } from "./async-job.service"; + +/** + * Full AsyncJob domain layer for consumers that need authz-wrapped access. + * Imports `AsyncJobStoreModule` for the underlying CRUD and adds the + * ownership policy + `AsyncJobService`. Re-exports the store so consumers + * can inject `AsyncJobStore` directly when they don't need authz. + * + * Workers should import `AsyncJobStoreModule` directly to avoid pulling in + * `AuthorizationModule`. + */ +@Module({ + imports: [AsyncJobStoreModule, AuthorizationModule], + providers: [AsyncJobService, AsyncJobPolicy], + exports: [AsyncJobService, AsyncJobStoreModule], +}) +export class AsyncJobModule {} diff --git a/packages/core/src/modules/async-job/async-job.policy.ts b/packages/core/src/modules/async-job/async-job.policy.ts new file mode 100644 index 0000000..f9285bc --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.policy.ts @@ -0,0 +1,8 @@ +import { Injectable } from "@nestjs/common"; +import { Policy } from "../auth/authorization/policy.decorator"; +import { UserOwnedResourcePolicy } from "../auth/authorization/user-owned-resource.policy"; +import { AsyncJobEntity } from "./async-job.entity"; + +@Injectable() +@Policy(AsyncJobEntity) +export class AsyncJobPolicy extends UserOwnedResourcePolicy {} diff --git a/packages/core/src/modules/async-job/async-job.service.ts b/packages/core/src/modules/async-job/async-job.service.ts new file mode 100644 index 0000000..214b90a --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.service.ts @@ -0,0 +1,36 @@ +import { Inject, Injectable } from "@nestjs/common"; +import { + type Authorizer, + AuthorizationService, +} from "../auth/authorization/authorization.service"; +import { EntityNotFoundError } from "../auth/errors/not-found.util"; +import { User } from "../auth/user/user.entity"; +import { AsyncJobEntity } from "./async-job.entity"; +import { AsyncJobStore } from "./async-job.store"; + +/** + * User-context authz wrapper around `AsyncJobStore`. The api uses this for + * the `asyncJob` query + the bundle interceptor. Workers depend on + * `AsyncJobStore` directly since they don't have a `User` to check against. + */ +@Injectable() +export class AsyncJobService { + constructor( + private readonly store: AsyncJobStore, + @Inject(AuthorizationService) private readonly authorization: Authorizer, + ) {} + + /** + * Load + authorize. Throws `EntityNotFoundError` if missing, + * `CannotViewError` if the user doesn't own the row. The global exception + * filter maps both to the right HTTP status. + */ + async findByIdForUser(id: string, user: User): Promise { + const entity = await this.store.findById(id); + if (!entity) { + throw new EntityNotFoundError("AsyncJob", "id", id); + } + await this.authorization.canView(user, entity, AsyncJobEntity); + return entity; + } +} diff --git a/packages/core/src/modules/async-job/async-job.store.ts b/packages/core/src/modules/async-job/async-job.store.ts new file mode 100644 index 0000000..99cef31 --- /dev/null +++ b/packages/core/src/modules/async-job/async-job.store.ts @@ -0,0 +1,42 @@ +import { Injectable } from "@nestjs/common"; +import { Prisma } from "@prisma/client"; +import { ClockService } from "../../shared/clock.service"; +import { PrismaService } from "../database/prisma.service"; +import { AsyncJobEntity } from "./async-job.entity"; +import { AsyncJobMapper } from "./async-job.mapper"; + +/** + * Authorization-free CRUD on `AsyncJob`. Used by workers (no `User` context) + * and composed by `AsyncJobService` for authz-wrapped api access. The method + * shapes guarantee the entity's status invariants by construction; SQL + * CHECK + the entity constructor are the other two defense layers. + */ +@Injectable() +export class AsyncJobStore { + constructor( + private readonly prisma: PrismaService, + private readonly mapper: AsyncJobMapper, + private readonly clock: ClockService, + ) {} + + async findById(id: string): Promise { + const row = await this.prisma.asyncJob.findUnique({ where: { id } }); + return row ? this.mapper.toDomain(row) : null; + } + + /** Set result + completedAt, clear error. Idempotent on retries. */ + async markCompleted(id: string, result: Prisma.InputJsonValue): Promise { + await this.prisma.asyncJob.update({ + where: { id }, + data: { result, error: null, completedAt: this.clock.now() }, + }); + } + + /** Set error + completedAt, clear any stale result. */ + async markFailed(id: string, error: string): Promise { + await this.prisma.asyncJob.update({ + where: { id }, + data: { error, result: Prisma.DbNull, completedAt: this.clock.now() }, + }); + } +} diff --git a/packages/core/src/modules/async-job/index.ts b/packages/core/src/modules/async-job/index.ts new file mode 100644 index 0000000..1b4b8d0 --- /dev/null +++ b/packages/core/src/modules/async-job/index.ts @@ -0,0 +1,7 @@ +export { AsyncJobEntity, AsyncJobStatus } from "./async-job.entity"; +export { AsyncJobMapper } from "./async-job.mapper"; +export { AsyncJobModule } from "./async-job.module"; +export { AsyncJobPolicy } from "./async-job.policy"; +export { AsyncJobService } from "./async-job.service"; +export { AsyncJobStore } from "./async-job.store"; +export { AsyncJobStoreModule } from "./async-job-store.module"; diff --git a/packages/core/src/modules/auth/authorization/authorization.service.ts b/packages/core/src/modules/auth/authorization/authorization.service.ts index 5ece219..8d7080d 100644 --- a/packages/core/src/modules/auth/authorization/authorization.service.ts +++ b/packages/core/src/modules/auth/authorization/authorization.service.ts @@ -9,8 +9,39 @@ import { import { User } from "../user/user.entity"; import { PolicyRegistry } from "./policy-registry.service"; +/** + * Public surface of `AuthorizationService`, redeclared as an interface so + * consumers (and tests) can depend on the structural type instead of the + * concrete class with private members. The class implements this; consumers + * type their constructor params against `Authorizer` and Nest still injects + * the concrete class. + */ +export interface Authorizer { + canView( + userContext: User, + resource: TResource, + resourceType?: Type, + ): Promise; + canCreate( + userContext: User, + resourceType: Type, + resource?: Partial, + ): Promise; + canUpdate( + userContext: User, + resource: TResource, + resourceType?: Type, + ): Promise; + canDelete( + userContext: User, + resource: TResource, + resourceType?: Type, + ): Promise; + isSameUser(userContext: User, userId: string): boolean; +} + @Injectable() -export class AuthorizationService { +export class AuthorizationService implements Authorizer { constructor(private readonly policyRegistry: PolicyRegistry) {} async canView( diff --git a/packages/core/src/modules/cv-template/cv-renderer.service.ts b/packages/core/src/modules/cv-template/cv-renderer.service.ts index 9d3d34f..9a66ecf 100644 --- a/packages/core/src/modules/cv-template/cv-renderer.service.ts +++ b/packages/core/src/modules/cv-template/cv-renderer.service.ts @@ -4,8 +4,8 @@ import { wrapInDocument, } from "@cv/cv-renderer"; import { Injectable } from "@nestjs/common"; -import type { MessageBus } from "@riotbyte-com/project-q-core"; import { InjectMessageBus } from "@riotbyte-com/project-q-nestjs"; +import type { MessageBus } from "../messenger/message-bus.interface"; import { RenderPdfMessage } from "../messenger/messages/render-pdf.message"; import { CVService } from "./cv.service"; import { CVDataAssemblerService } from "./cv-data-assembler.service"; diff --git a/packages/core/src/modules/messenger/index.ts b/packages/core/src/modules/messenger/index.ts index 93d96b9..5ef8e8f 100644 --- a/packages/core/src/modules/messenger/index.ts +++ b/packages/core/src/modules/messenger/index.ts @@ -1,10 +1,15 @@ // Re-exports so apps don't need a direct dep on @riotbyte-com/project-q-*. -export type { Envelope, MessageBus } from "@riotbyte-com/project-q-core"; +export type { Envelope } from "@riotbyte-com/project-q-core"; +export type { MessageBus } from "./message-bus.interface"; export { HandlerTag, InjectMessageBus, } from "@riotbyte-com/project-q-nestjs"; export { NestProjectQLogger } from "./logger.provider"; +export { + type ParseCVInput, + ParseCVMessage, +} from "./messages/parse-cv.message"; export { RenderPdfMessage } from "./messages/render-pdf.message"; export { ProjectQMessagingModule, diff --git a/packages/core/src/modules/messenger/message-bus.interface.ts b/packages/core/src/modules/messenger/message-bus.interface.ts new file mode 100644 index 0000000..3704a45 --- /dev/null +++ b/packages/core/src/modules/messenger/message-bus.interface.ts @@ -0,0 +1,11 @@ +import type { Envelope, Message } from "@riotbyte-com/project-q-core"; + +/** + * Public surface of project-q's `MessageBus`, redefined locally so consumers + * (and tests) can depend on the interface without project-q's class private + * members forcing nominal typing. The concrete `MessageBus` exported by + * `@riotbyte-com/project-q-core` structurally satisfies this interface. + */ +export interface MessageBus { + dispatch(messageOrEnvelope: Envelope | Message): Promise; +} 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..552c8b0 --- /dev/null +++ b/packages/core/src/modules/messenger/messages/parse-cv.message.ts @@ -0,0 +1,28 @@ +import { defineZodMessage } from "@riotbyte-com/project-q-core"; +import { z } from "zod/v4"; + +const ParseCVInputSchema = 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(), + }), +]); + +export type ParseCVInput = z.infer; + +/// `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: ParseCVInputSchema, + }), +); 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/file-upload/src/extractor-registry.ts b/packages/file-upload/src/extractor-registry.ts index 37ad8d9..22090a5 100644 --- a/packages/file-upload/src/extractor-registry.ts +++ b/packages/file-upload/src/extractor-registry.ts @@ -3,6 +3,15 @@ import { ServiceLocator } from "@riotbyte-com/nest-service-locator"; import { TextExtractorTag } from "./text-extractor.tag"; import type { TextExtractionResult, TextExtractor } from "./types"; +/** + * Narrow public surface used by consumers that only need to route a buffer + * through the registered extractors. Tests stub this directly without the + * class's private `locator` field forcing nominal typing. + */ +export interface ExtractsText { + extract(buffer: Buffer, mimeType: string): Promise; +} + /** * Routes file-extraction requests to the registered extractor that supports * the request's mime type. Discovery is via `ServiceLocator.tagged(...)`: @@ -10,7 +19,7 @@ import type { TextExtractionResult, TextExtractor } from "./types"; * is picked up automatically. */ @Injectable() -export class TextExtractorRegistry { +export class TextExtractorRegistry implements ExtractsText { constructor(private readonly locator: ServiceLocator) {} getAll(): TextExtractor[] { diff --git a/packages/file-upload/src/index.ts b/packages/file-upload/src/index.ts index 90e80fb..76daa50 100644 --- a/packages/file-upload/src/index.ts +++ b/packages/file-upload/src/index.ts @@ -1,7 +1,7 @@ // Types and schemas // Extractor registry -export { TextExtractorRegistry } from "./extractor-registry"; +export { type ExtractsText, TextExtractorRegistry } from "./extractor-registry"; export { BaseTextExtractor } from "./extractors/base-extractor"; export { DOCXExtractor } from "./extractors/docx.extractor"; @@ -50,4 +50,5 @@ export { validateFileName, validateFileSize, validateMimeType, + validateStoryText, } from "./validators"; diff --git a/packages/file-upload/src/validators.ts b/packages/file-upload/src/validators.ts index 2853a3f..dc5ea9a 100644 --- a/packages/file-upload/src/validators.ts +++ b/packages/file-upload/src/validators.ts @@ -7,6 +7,7 @@ import { } from "./types"; const MAX_FILE_SIZE = 10 * 1024 * 1024; // 10MB +const MAX_STORY_TEXT_LENGTH = 50_000; /** * Magic-byte signatures for the supported MIME types. The client-supplied @@ -135,6 +136,26 @@ export const validateMimeType = (mimeType: string): FileValidationResult => { return { valid: true }; }; +/** + * Validate inline story text (used by the parse-cv "story" source). Mirrors + * `validateFile` shape so callers can route both file and story inputs + * through the same validation surface. + */ +export const validateStoryText = (text: string): FileValidationResult => { + if (!text || text.trim().length === 0) { + return { valid: false, error: "Story text cannot be empty" }; + } + + if (text.length > MAX_STORY_TEXT_LENGTH) { + return { + valid: false, + error: `Story text exceeds maximum length of ${MAX_STORY_TEXT_LENGTH} characters`, + }; + } + + return { valid: true }; +}; + /** * Validate file name */ diff --git a/packages/handlers/package.json b/packages/handlers/package.json index 91868e4..4c6c561 100644 --- a/packages/handlers/package.json +++ b/packages/handlers/package.json @@ -6,11 +6,16 @@ "types": "./dist/index.d.ts", "scripts": { "typecheck": "tsc --noEmit", - "build": "tsc -b" + "build": "tsc -b", + "test": "vitest run", + "test:watch": "vitest" }, "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", @@ -20,7 +25,8 @@ "devDependencies": { "@cv/tsconfig": "*", "@types/node": "^22.7.5", - "typescript": "^5.3.3" + "typescript": "^5.3.3", + "vitest": "^4.0.0" }, "files": [ "dist/" diff --git a/packages/handlers/src/__tests__/parse-cv.handler.spec.ts b/packages/handlers/src/__tests__/parse-cv.handler.spec.ts new file mode 100644 index 0000000..d60b8cc --- /dev/null +++ b/packages/handlers/src/__tests__/parse-cv.handler.spec.ts @@ -0,0 +1,168 @@ +import type { AsyncJobEntity, AsyncJobStore } from "@cv/core"; +import type { Envelope } from "@riotbyte-com/project-q-core"; +import { describe, expect, it, vi } from "vitest"; +import { ParseCvHandler } from "../parse-cv.handler"; +import type { TextSource } from "../text-source.resolver"; +import type { UserCvParser } from "../user-cv-parser.service"; + +const fakeEnvelope = (data: unknown): Envelope => + ({ message: { name: "parse-cv", data } }) as Envelope; + +interface HandlerSetup { + asyncJobs: { + findById: ReturnType; + markCompleted: ReturnType; + markFailed: ReturnType; + }; + textSource: { resolve: ReturnType }; + parser: { parseForUser: ReturnType }; + handler: ParseCvHandler; +} + +const setup = (overrides: { + job?: AsyncJobEntity | null; + parseResult?: unknown; + parseError?: Error; + markCompletedError?: Error; + markFailedError?: Error; +} = {}): HandlerSetup => { + const job = + overrides.job === undefined + ? ({ + id: "job-1", + userId: "user-1", + kind: "PARSE_CV", + result: null, + error: null, + completedAt: null, + createdAt: new Date(), + updatedAt: new Date(), + } as unknown as AsyncJobEntity) + : overrides.job; + + const asyncJobs = { + findById: vi.fn().mockResolvedValue(job), + markCompleted: overrides.markCompletedError + ? vi.fn().mockRejectedValue(overrides.markCompletedError) + : vi.fn().mockResolvedValue(undefined), + markFailed: overrides.markFailedError + ? vi.fn().mockRejectedValue(overrides.markFailedError) + : vi.fn().mockResolvedValue(undefined), + }; + + const textSource = { + resolve: vi.fn().mockResolvedValue("extracted text"), + }; + + const parser = { + parseForUser: overrides.parseError + ? vi.fn().mockRejectedValue(overrides.parseError) + : vi.fn().mockResolvedValue( + overrides.parseResult ?? { jobExperiences: [], education: [] }, + ), + }; + + const handler = new ParseCvHandler( + asyncJobs as unknown as AsyncJobStore, + textSource satisfies TextSource, + parser satisfies UserCvParser, + ); + + return { asyncJobs, textSource, parser, handler }; +}; + +describe("ParseCvHandler", () => { + it("marks the job completed with the parsed result on success", async () => { + const { handler, asyncJobs, textSource, parser } = setup(); + + await handler.handle( + fakeEnvelope({ + jobId: "job-1", + input: { source: "story", text: "I worked at ACME as an engineer." }, + }), + ); + + expect(textSource.resolve).toHaveBeenCalledWith({ + source: "story", + text: "I worked at ACME as an engineer.", + }); + expect(parser.parseForUser).toHaveBeenCalledWith("user-1", "extracted text"); + expect(asyncJobs.markCompleted).toHaveBeenCalledWith("job-1", { + jobExperiences: [], + education: [], + }); + expect(asyncJobs.markFailed).not.toHaveBeenCalled(); + }); + + it("marks the job failed and rethrows when the parse fails", async () => { + const { handler, asyncJobs } = setup({ + parseError: new Error("provider broken"), + }); + + await expect( + handler.handle( + fakeEnvelope({ + jobId: "job-1", + input: { source: "story", text: "irrelevant" }, + }), + ), + ).rejects.toThrow("provider broken"); + + expect(asyncJobs.markFailed).toHaveBeenCalledWith("job-1", "provider broken"); + expect(asyncJobs.markCompleted).not.toHaveBeenCalled(); + }); + + it("does not swallow the original parse error if the markFailed write also fails", async () => { + const { handler } = setup({ + parseError: new Error("provider broken"), + markFailedError: new Error("db unreachable"), + }); + + // Rethrows the ORIGINAL parse error, not the db-unreachable secondary failure. + await expect( + handler.handle( + fakeEnvelope({ + jobId: "job-1", + input: { source: "story", text: "irrelevant" }, + }), + ), + ).rejects.toThrow("provider broken"); + }); + + it("surfaces a raw error from markCompleted without calling markFailed", async () => { + // A parse-succeeded + write-failed path should NOT call markFailed + // (that would conflate a DB issue with a parse failure). project-q + // retries the message; the row stays untouched. + const { handler, asyncJobs } = setup({ + markCompletedError: new Error("db unreachable"), + }); + + await expect( + handler.handle( + fakeEnvelope({ + jobId: "job-1", + input: { source: "story", text: "irrelevant" }, + }), + ), + ).rejects.toThrow("db unreachable"); + + expect(asyncJobs.markFailed).not.toHaveBeenCalled(); + }); + + it("returns without throwing when the AsyncJob row is gone (no retry-forever loop)", async () => { + const { handler, asyncJobs, parser } = setup({ job: null }); + + await expect( + handler.handle( + fakeEnvelope({ + jobId: "job-gone", + input: { source: "story", text: "irrelevant" }, + }), + ), + ).resolves.toBeUndefined(); + + expect(parser.parseForUser).not.toHaveBeenCalled(); + expect(asyncJobs.markCompleted).not.toHaveBeenCalled(); + expect(asyncJobs.markFailed).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/handlers/src/handlers.module.ts b/packages/handlers/src/handlers.module.ts index a89dda2..399903d 100644 --- a/packages/handlers/src/handlers.module.ts +++ b/packages/handlers/src/handlers.module.ts @@ -1,7 +1,17 @@ +import { + AIResolutionModule, + AsyncJobStoreModule, + BaseModule, + DatabaseModule, +} 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"; +import { TextSourceResolver } from "./text-source.resolver"; +import { UserCvParserService } from "./user-cv-parser.service"; export { PDF_SERVICE_CONFIG }; @@ -10,12 +20,22 @@ export class HandlersModule { static forRoot(config: PdfServiceConfig): DynamicModule { return { module: HandlersModule, + imports: [ + BaseModule, + DatabaseModule, + AIResolutionModule, + AsyncJobStoreModule, + FileExtractionModule.forRoot(), + ], providers: [ { provide: PDF_SERVICE_CONFIG, useValue: config }, HtmlToPdfService, RenderPdfHandler, + TextSourceResolver, + UserCvParserService, + 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..3d9531b --- /dev/null +++ b/packages/handlers/src/parse-cv.handler.ts @@ -0,0 +1,64 @@ +import { AsyncJobStore, ParseCVMessage } from "@cv/core"; +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"; +import { type TextSource, TextSourceResolver } from "./text-source.resolver"; +import { type UserCvParser, UserCvParserService } from "./user-cv-parser.service"; + +@Injectable() +@HandlerTag.decorator({ handles: "parse-cv" }) +export class ParseCvHandler implements Handler { + private readonly logger = new Logger(ParseCvHandler.name); + + constructor( + private readonly asyncJobs: AsyncJobStore, + @Inject(TextSourceResolver) private readonly textSource: TextSource, + @Inject(UserCvParserService) private readonly parser: UserCvParser, + ) {} + + async handle(envelope: Envelope): Promise { + const { jobId, input } = ParseCVMessage.parse(envelope.message).data; + + const job = await this.asyncJobs.findById(jobId); + + // Row gone (cancelled, cascade-deleted, dispatch-rollback won a race). + // The message has nothing left to update so retrying is pointless - ack + // the envelope by returning normally and let project-q drop the message. + if (!job) { + this.logger.warn( + `Skipping parse-cv job ${jobId}: AsyncJob row no longer exists`, + ); + return; + } + + this.logger.log(`Processing parse-cv job ${jobId}`); + + // Two-stage error handling. Parse failures go through `markFailed` so + // the row reflects the user-facing cause. The terminal `markCompleted` + // write is OUTSIDE the parse try/catch so a DB error there doesn't + // get misreported as a parse failure - it surfaces raw and project-q + // retries. + let result: Awaited>; + try { + const text = await this.textSource.resolve(input); + result = await this.parser.parseForUser(job.userId, text); + } catch (err) { + const message = err instanceof Error ? err.message : String(err); + this.logger.error(`Failed parse-cv job ${jobId}: ${message}`); + // Best-effort error write-back. If this fails too (DB gone, row + // deleted out from under us), log it but rethrow the ORIGINAL error + // so project-q sees the real cause - not "row missing". + await this.asyncJobs.markFailed(jobId, message).catch((writeErr) => { + this.logger.error( + `Failed to persist error onto AsyncJob ${jobId}: ${ + writeErr instanceof Error ? writeErr.message : String(writeErr) + }`, + ); + }); + throw err; + } + + await this.asyncJobs.markCompleted(jobId, result); + this.logger.log(`Completed parse-cv job ${jobId}`); + } +} diff --git a/packages/handlers/src/text-source.resolver.ts b/packages/handlers/src/text-source.resolver.ts new file mode 100644 index 0000000..b1ad474 --- /dev/null +++ b/packages/handlers/src/text-source.resolver.ts @@ -0,0 +1,51 @@ +import type { ParseCVInput } from "@cv/core"; +import { FILE_STORAGE, type FileStorage } from "@cv/file-storage"; +import { TEXT_EXTRACTOR_REGISTRY, TextExtractorRegistry } from "@cv/file-upload"; +import { Inject, Injectable } from "@nestjs/common"; + +/** + * Public surface used by `ParseCvHandler`. Tests stub this directly without + * the class's private storage/extractor fields forcing nominal typing. + */ +export interface TextSource { + resolve(input: ParseCVInput): Promise; +} + +/** + * Converts a `ParseCVInput` (discriminated union of file / story sources) + * into the raw CV text the parser consumes. Owns the file-storage + extractor + * collaboration so the handler doesn't have to. + */ +@Injectable() +export class TextSourceResolver implements TextSource { + constructor( + @Inject(FILE_STORAGE) private readonly storage: FileStorage, + @Inject(TEXT_EXTRACTOR_REGISTRY) + private readonly extractors: TextExtractorRegistry, + ) {} + + async resolve(input: ParseCVInput): Promise { + if (input.source === "story") { + return input.text; + } + return this.extractFromFile(input.fileKey, input.mimeType); + } + + 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; + } +} diff --git a/packages/handlers/src/user-cv-parser.service.ts b/packages/handlers/src/user-cv-parser.service.ts new file mode 100644 index 0000000..620bfb8 --- /dev/null +++ b/packages/handlers/src/user-cv-parser.service.ts @@ -0,0 +1,30 @@ +import { CVParserService, type ParsedCVData } from "@cv/ai-parser"; +import { AIProviderResolverService } from "@cv/core"; +import { Injectable } from "@nestjs/common"; + +/** + * Public surface used by `ParseCvHandler`. Tests stub this directly without + * the class's private providerResolver field forcing nominal typing. + */ +export interface UserCvParser { + parseForUser(userId: string, text: string): Promise; +} + +/** + * Resolves the AI provider for a user (with audit logging) and parses CV + * text against it. Wraps the inline `new CVParserService(provider)` that + * would otherwise live in every caller, and keeps the `AIProviderResolverService` + * dependency off the leaf handler. + */ +@Injectable() +export class UserCvParserService implements UserCvParser { + constructor( + private readonly providerResolver: AIProviderResolverService, + ) {} + + async parseForUser(userId: string, text: string): Promise { + const provider = await this.providerResolver.resolveForUser(userId); + const parser = new CVParserService(provider); + return parser.parseCVText(text); + } +} 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/packages/handlers/vitest.config.ts b/packages/handlers/vitest.config.ts new file mode 100644 index 0000000..76e960e --- /dev/null +++ b/packages/handlers/vitest.config.ts @@ -0,0 +1,11 @@ +import { defineConfig } from "vitest/config"; + +export default defineConfig({ + test: { + root: "./src", + globals: true, + environment: "node", + include: ["**/*.{test,spec}.{ts,tsx}"], + exclude: ["**/node_modules/**", "**/dist/**"], + }, +}); diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 1e0cbed..ae1b0d5 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -293,6 +293,9 @@ importers: typescript: specifier: ^5.6.3 version: 5.9.3 + unplugin-swc: + specifier: ^1.5.9 + version: 1.5.9(@swc/core@1.15.13)(rollup@4.60.3) vitest: specifier: ^4.0.16 version: 4.0.16(@opentelemetry/api@1.9.1)(@types/node@22.19.3)(@vitest/ui@4.0.16)(jiti@2.7.0)(jsdom@28.1.0(@noble/hashes@1.8.0)(canvas@3.2.1))(lightningcss@1.30.2)(yaml@2.8.2) @@ -999,12 +1002,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) @@ -1030,6 +1042,9 @@ importers: typescript: specifier: ^5.3.3 version: 5.9.3 + vitest: + specifier: ^4.0.0 + version: 4.0.16(@opentelemetry/api@1.9.1)(@types/node@22.19.3)(@vitest/ui@4.0.16)(jiti@2.7.0)(jsdom@28.1.0(@noble/hashes@1.8.0)(canvas@3.2.1))(lightningcss@1.30.2)(yaml@2.8.2) packages/mail: dependencies: @@ -6488,6 +6503,10 @@ packages: resolution: {integrity: sha512-gUD/epcRms75Cw8RT1pUdHugZYM5ce64ucs2GEISABwkRsOQr0q2wm/MV2TKThycIe5e0ytRweW2RZxclogCdQ==} engines: {node: '>=8'} + load-tsconfig@0.2.5: + resolution: {integrity: sha512-IXO6OCs9yg8tMKzfPZ1YmheJbZCiEsnBdcB03l0OcfK9prKnJb96siuHCr5Fl37/yo9DnKU+TLpxzTUspw9shg==} + engines: {node: ^12.20.0 || ^14.13.1 || >=16.0.0} + locate-path@2.0.0: resolution: {integrity: sha512-NCI2kiDkyR7VeEKm27Kda/iQHyKJe1Bu0FlTbYp3CqJu+9IFe9bLyAjMxf5ZDDbEg+iMPzB5zYyUTSm8wVTKmA==} engines: {node: '>=4'} @@ -8554,6 +8573,15 @@ packages: resolution: {integrity: sha512-pjy2bYhSsufwWlKwPc+l3cN7+wuJlK6uz0YdJEOlQDbl6jo/YlPi4mb8agUkVC8BF7V8NuzeyPNqRksA3hztKQ==} engines: {node: '>= 0.8'} + unplugin-swc@1.5.9: + resolution: {integrity: sha512-RKwK3yf0M+MN17xZfF14bdKqfx0zMXYdtOdxLiE6jHAoidupKq3jGdJYANyIM1X/VmABhh1WpdO+/f4+Ol89+g==} + peerDependencies: + '@swc/core': ^1.2.108 + + unplugin@2.3.11: + resolution: {integrity: sha512-5uKD0nqiYVzlmCRs01Fhs2BdkEgBS3SAVP6ndrBsuK42iC2+JHyxM05Rm9G8+5mkmRtzMZGY8Ct5+mliZxU/Ww==} + engines: {node: '>=18.12.0'} + upath@2.0.1: resolution: {integrity: sha512-1uEe95xksV1O0CYKXo8vQvN1JEbtJp7lb7C5U9HMsIp6IVwntkH/oNUzyVNQSd4S1sYk2FpSSW44FqMc8qee5w==} engines: {node: '>=4'} @@ -8757,6 +8785,9 @@ packages: resolution: {integrity: sha512-BMhLD/Sw+GbJC21C/UgyaZX41nPt8bUTg+jWyDeg7e7YN4xOM05YPSIXceACnXVtqyEw/LMClUQMtMZ+PGGpqQ==} engines: {node: '>=20'} + webpack-virtual-modules@0.6.2: + resolution: {integrity: sha512-66/V2i5hQanC51vBQKPH4aI8NMAcBW59FVBs+rC7eGHupMyfn34q7rZIE+ETlJ+XTevqfUhVVBgSUNSW2flEUQ==} + whatwg-mimetype@4.0.0: resolution: {integrity: sha512-QaKxh0eNIi2mE9p2vEdzfagOKHCcj1pJ56EEHGQOVxp8r9/iszLUUV7v89x9O1p/T+NlTM5W7jW6+cz4Fq1YVg==} engines: {node: '>=18'} @@ -15320,6 +15351,8 @@ snapshots: strip-bom: 4.0.0 type-fest: 0.6.0 + load-tsconfig@0.2.5: {} + locate-path@2.0.0: dependencies: p-locate: 2.0.0 @@ -17820,6 +17853,22 @@ snapshots: unpipe@1.0.0: {} + unplugin-swc@1.5.9(@swc/core@1.15.13)(rollup@4.60.3): + dependencies: + '@rollup/pluginutils': 5.3.0(rollup@4.60.3) + '@swc/core': 1.15.13 + load-tsconfig: 0.2.5 + unplugin: 2.3.11 + transitivePeerDependencies: + - rollup + + unplugin@2.3.11: + dependencies: + '@jridgewell/remapping': 2.3.5 + acorn: 8.15.0 + picomatch: 4.0.4 + webpack-virtual-modules: 0.6.2 + upath@2.0.1: {} update-browserslist-db@1.2.3(browserslist@4.28.1): @@ -17964,6 +18013,8 @@ snapshots: webidl-conversions@8.0.1: {} + webpack-virtual-modules@0.6.2: {} + whatwg-mimetype@4.0.0: {} whatwg-mimetype@5.0.0: {} -- 2.51.2