From 005b45f7d46ef27298a1ee8948898b4f65ade00a Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Wed, 25 Feb 2026 19:50:57 +0100 Subject: [PATCH] feat(server): add data import module with file upload and import jobs --- .../migration.sql | 39 +++ .../migration.sql | 2 + apps/server/prisma/models/data-import.prisma | 37 +++ .../cv-parser/graphql/upload-file.input.ts | 16 ++ .../cv-parser/graphql/upload-file.resolver.ts | 79 ++++++ .../data-import-source.interface.ts | 11 + .../modules/data-import/data-import.module.ts | 29 ++ .../events/user-file-created.event.ts | 13 + .../data-import/graphql/user-file.resolver.ts | 54 ++++ .../data-import/graphql/user-file.type.ts | 66 +++++ .../modules/data-import/import-job.entity.ts | 18 ++ .../modules/data-import/import-job.mapper.ts | 36 +++ .../modules/data-import/import-job.policy.ts | 28 ++ .../src/modules/data-import/import.service.ts | 86 ++++++ .../listeners/import-job.listener.ts | 268 ++++++++++++++++++ .../data-import/sources/file-import-source.ts | 152 ++++++++++ .../modules/data-import/user-file.entity.ts | 21 ++ .../modules/data-import/user-file.mapper.ts | 39 +++ .../modules/data-import/user-file.policy.ts | 12 + 19 files changed, 1006 insertions(+) create mode 100644 apps/server/prisma/migrations/20260209211505_add_data_import_models/migration.sql create mode 100644 apps/server/prisma/migrations/20260217120000_add_file_hash/migration.sql create mode 100644 apps/server/prisma/models/data-import.prisma create mode 100644 apps/server/src/modules/cv-parser/graphql/upload-file.input.ts create mode 100644 apps/server/src/modules/cv-parser/graphql/upload-file.resolver.ts create mode 100644 apps/server/src/modules/data-import/data-import-source.interface.ts create mode 100644 apps/server/src/modules/data-import/data-import.module.ts create mode 100644 apps/server/src/modules/data-import/events/user-file-created.event.ts create mode 100644 apps/server/src/modules/data-import/graphql/user-file.resolver.ts create mode 100644 apps/server/src/modules/data-import/graphql/user-file.type.ts create mode 100644 apps/server/src/modules/data-import/import-job.entity.ts create mode 100644 apps/server/src/modules/data-import/import-job.mapper.ts create mode 100644 apps/server/src/modules/data-import/import-job.policy.ts create mode 100644 apps/server/src/modules/data-import/import.service.ts create mode 100644 apps/server/src/modules/data-import/listeners/import-job.listener.ts create mode 100644 apps/server/src/modules/data-import/sources/file-import-source.ts create mode 100644 apps/server/src/modules/data-import/user-file.entity.ts create mode 100644 apps/server/src/modules/data-import/user-file.mapper.ts create mode 100644 apps/server/src/modules/data-import/user-file.policy.ts diff --git a/apps/server/prisma/migrations/20260209211505_add_data_import_models/migration.sql b/apps/server/prisma/migrations/20260209211505_add_data_import_models/migration.sql new file mode 100644 index 0000000..14ec9ee --- /dev/null +++ b/apps/server/prisma/migrations/20260209211505_add_data_import_models/migration.sql @@ -0,0 +1,39 @@ +-- CreateTable +CREATE TABLE "user_files" ( + "id" TEXT NOT NULL, + "userId" TEXT NOT NULL, + "fileName" TEXT NOT NULL, + "mimeType" TEXT NOT NULL, + "sizeBytes" INTEGER NOT NULL, + "source" TEXT NOT NULL, + "status" TEXT NOT NULL DEFAULT 'pending', + "statusMessage" TEXT NOT NULL DEFAULT '', + "resultJson" TEXT, + "error" TEXT, + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "user_files_pkey" PRIMARY KEY ("id") +); + +-- CreateTable +CREATE TABLE "import_jobs" ( + "id" TEXT NOT NULL, + "userFileId" TEXT NOT NULL, + "source" TEXT NOT NULL, + "status" TEXT NOT NULL DEFAULT 'pending', + "statusMessage" TEXT NOT NULL DEFAULT '', + "error" TEXT, + "startedAt" TIMESTAMP(3), + "completedAt" TIMESTAMP(3), + "createdAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "updatedAt" TIMESTAMP(3) NOT NULL, + + CONSTRAINT "import_jobs_pkey" PRIMARY KEY ("id") +); + +-- AddForeignKey +ALTER TABLE "user_files" ADD CONSTRAINT "user_files_userId_fkey" FOREIGN KEY ("userId") REFERENCES "users"("id") ON DELETE CASCADE ON UPDATE CASCADE; + +-- AddForeignKey +ALTER TABLE "import_jobs" ADD CONSTRAINT "import_jobs_userFileId_fkey" FOREIGN KEY ("userFileId") REFERENCES "user_files"("id") ON DELETE CASCADE ON UPDATE CASCADE; diff --git a/apps/server/prisma/migrations/20260217120000_add_file_hash/migration.sql b/apps/server/prisma/migrations/20260217120000_add_file_hash/migration.sql new file mode 100644 index 0000000..1dfd058 --- /dev/null +++ b/apps/server/prisma/migrations/20260217120000_add_file_hash/migration.sql @@ -0,0 +1,2 @@ +-- Add fingerprint for file deduplication +ALTER TABLE "user_files" ADD COLUMN "fingerprint" TEXT; diff --git a/apps/server/prisma/models/data-import.prisma b/apps/server/prisma/models/data-import.prisma new file mode 100644 index 0000000..135dd87 --- /dev/null +++ b/apps/server/prisma/models/data-import.prisma @@ -0,0 +1,37 @@ +model UserFile { + id String @id @default(cuid()) + profileId String + fingerprint String? + fileName String + mimeType String + sizeBytes Int + source String + status String @default("pending") + statusMessage String @default("") + resultJson String? + error String? + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + profile Profile @relation(fields: [profileId], references: [id], onDelete: Cascade) + importJobs ImportJob[] + + @@map("user_files") +} + +model ImportJob { + id String @id @default(cuid()) + userFileId String + source String + status String @default("pending") + statusMessage String @default("") + error String? + startedAt DateTime? + completedAt DateTime? + createdAt DateTime @default(now()) + updatedAt DateTime @updatedAt + + userFile UserFile @relation(fields: [userFileId], references: [id], onDelete: Cascade) + + @@map("import_jobs") +} diff --git a/apps/server/src/modules/cv-parser/graphql/upload-file.input.ts b/apps/server/src/modules/cv-parser/graphql/upload-file.input.ts new file mode 100644 index 0000000..218b030 --- /dev/null +++ b/apps/server/src/modules/cv-parser/graphql/upload-file.input.ts @@ -0,0 +1,16 @@ +import { Field, ID, InputType } from "@nestjs/graphql"; + +@InputType() +export class UploadFileInput { + @Field(() => ID, { nullable: true, description: "Target profile. If omitted, auto-creates a Default profile." }) + profileId?: string; + + @Field() + fileName!: string; + + @Field() + mimeType!: string; + + @Field({ description: "Base64-encoded file content" }) + content!: string; +} diff --git a/apps/server/src/modules/cv-parser/graphql/upload-file.resolver.ts b/apps/server/src/modules/cv-parser/graphql/upload-file.resolver.ts new file mode 100644 index 0000000..9c0b95a --- /dev/null +++ b/apps/server/src/modules/cv-parser/graphql/upload-file.resolver.ts @@ -0,0 +1,79 @@ +import { createHash } from "node:crypto"; +import type { User as DomainUser } from "@cv/auth"; +import { JwtAuthGuard, VerifiedScopeGuard } from "@cv/auth"; +import { validateFile } from "@cv/file-upload"; +import { BadRequestException, UseGuards } from "@nestjs/common"; +import { Args, Mutation, Resolver } from "@nestjs/graphql"; +import { CurrentUser } from "@/modules/current-user/current-user.decorator"; +import { UserFileType } from "@/modules/data-import/graphql/user-file.type"; +import { ImportService } from "@/modules/data-import/import.service"; +import { FileImportSource } from "@/modules/data-import/sources/file-import-source"; +import { ProfileService } from "@/modules/profile/profile.service"; +import { UploadFileInput } from "./upload-file.input"; + +const MAX_DECODED_SIZE = 10 * 1024 * 1024; // 10MB + +@Resolver() +@UseGuards(JwtAuthGuard, VerifiedScopeGuard) +export class UploadFileResolver { + constructor( + private readonly importService: ImportService, + private readonly fileImportSource: FileImportSource, + private readonly profileService: ProfileService, + ) {} + + @Mutation(() => UserFileType) + async uploadFile( + @CurrentUser() user: DomainUser, + @Args("input") input: UploadFileInput, + ): Promise { + const buffer = Buffer.from(input.content, "base64"); + + if (buffer.byteLength > MAX_DECODED_SIZE) { + throw new BadRequestException( + `File exceeds maximum size of ${MAX_DECODED_SIZE / (1024 * 1024)}MB`, + ); + } + + const validation = validateFile({ + buffer, + mimeType: input.mimeType, + originalName: input.fileName, + sizeBytes: buffer.byteLength, + }); + + if (!validation.valid) { + throw new BadRequestException(validation.error); + } + + const fingerprint = createHash("sha256").update(buffer).digest("hex"); + + const duplicate = await this.importService.findDuplicateForUser( + user.id, + fingerprint, + ); + + if (duplicate) { + return UserFileType.fromDomain(duplicate, { isDuplicate: true }); + } + + const profile = input.profileId + ? await this.profileService.findByIdAndUserOrFail(input.profileId, user.id) + : await this.profileService.getOrCreateDefaultProfile(user.id); + + const userFile = await this.importService.createImport( + user, + profile.id, + this.fileImportSource, + { + fileName: input.fileName, + mimeType: input.mimeType, + sizeBytes: buffer.byteLength, + fingerprint, + }, + { buffer, mimeType: input.mimeType }, + ); + + return UserFileType.fromDomain(userFile); + } +} diff --git a/apps/server/src/modules/data-import/data-import-source.interface.ts b/apps/server/src/modules/data-import/data-import-source.interface.ts new file mode 100644 index 0000000..4fe87e3 --- /dev/null +++ b/apps/server/src/modules/data-import/data-import-source.interface.ts @@ -0,0 +1,11 @@ +import type { ParsedCVData } from "@cv/ai-parser"; +import type { User } from "@cv/auth"; + +export interface DataImportSource { + readonly name: string; + execute( + user: User, + params: Record, + onStatus: (message: string) => Promise, + ): Promise; +} diff --git a/apps/server/src/modules/data-import/data-import.module.ts b/apps/server/src/modules/data-import/data-import.module.ts new file mode 100644 index 0000000..a2210e8 --- /dev/null +++ b/apps/server/src/modules/data-import/data-import.module.ts @@ -0,0 +1,29 @@ +import { AuthorizationModule } from "@cv/auth"; +import { DatabaseModule } from "@cv/system"; +import { Module } from "@nestjs/common"; +import { EventEmitterModule } from "@nestjs/event-emitter"; +import { EntityResolverService } from "@/modules/cv-parser/entity-resolver.service"; +import { ProfileModule } from "@/modules/profile/profile.module"; +import { UserFileResolver } from "./graphql/user-file.resolver"; +import { ImportService } from "./import.service"; +import { ImportJobMapper } from "./import-job.mapper"; +import { ImportJobPolicy } from "./import-job.policy"; +import { ImportJobListener } from "./listeners/import-job.listener"; +import { UserFileMapper } from "./user-file.mapper"; +import { UserFilePolicy } from "./user-file.policy"; + +@Module({ + imports: [DatabaseModule, AuthorizationModule, EventEmitterModule.forRoot(), ProfileModule], + providers: [ + ImportService, + EntityResolverService, + UserFileMapper, + UserFilePolicy, + ImportJobMapper, + ImportJobPolicy, + UserFileResolver, + ImportJobListener, + ], + exports: [ImportService, UserFileMapper, ImportJobMapper], +}) +export class DataImportModule {} diff --git a/apps/server/src/modules/data-import/events/user-file-created.event.ts b/apps/server/src/modules/data-import/events/user-file-created.event.ts new file mode 100644 index 0000000..57fdbb6 --- /dev/null +++ b/apps/server/src/modules/data-import/events/user-file-created.event.ts @@ -0,0 +1,13 @@ +import type { User } from "@cv/auth"; +import type { DataImportSource } from "../data-import-source.interface"; + +export class UserFileCreatedEvent { + static readonly event = "user-file.created"; + + constructor( + public readonly userFileId: string, + public readonly user: User, + public readonly source: DataImportSource, + public readonly params: Record, + ) {} +} diff --git a/apps/server/src/modules/data-import/graphql/user-file.resolver.ts b/apps/server/src/modules/data-import/graphql/user-file.resolver.ts new file mode 100644 index 0000000..476fc3e --- /dev/null +++ b/apps/server/src/modules/data-import/graphql/user-file.resolver.ts @@ -0,0 +1,54 @@ +import type { User as DomainUser } from "@cv/auth"; +import { + AuthorizationService, + JwtAuthGuard, + VerifiedScopeGuard, +} from "@cv/auth"; +import { UseGuards } from "@nestjs/common"; +import { Args, Mutation, Query, Resolver } from "@nestjs/graphql"; +import { CurrentUser } from "@/modules/current-user/current-user.decorator"; +import { ImportService } from "../import.service"; +import { UserFile as UserFileEntity } from "../user-file.entity"; +import { UserFileType } from "./user-file.type"; + +@Resolver(() => UserFileType) +@UseGuards(JwtAuthGuard, VerifiedScopeGuard) +export class UserFileResolver { + constructor( + private readonly importService: ImportService, + private readonly authorizationService: AuthorizationService, + ) {} + + @Query(() => UserFileType, { nullable: true }) + async userFile( + @CurrentUser() user: DomainUser, + @Args("id") id: string, + ): Promise { + const file = await this.importService.findUserFileById(id); + if (!file) return null; + + await this.authorizationService.canView(user, file, UserFileEntity); + + return UserFileType.fromDomain(file); + } + + @Query(() => [UserFileType]) + async myUserFiles(@CurrentUser() user: DomainUser): Promise { + const files = await this.importService.findUserFilesForUser(user.id); + return files.map((f) => UserFileType.fromDomain(f)); + } + + @Mutation(() => Boolean) + async deleteUserFile( + @CurrentUser() user: DomainUser, + @Args("id") id: string, + ): Promise { + const file = await this.importService.findUserFileById(id); + if (!file) return false; + + await this.authorizationService.canDelete(user, file, UserFileEntity); + await this.importService.deleteUserFile(id); + + return true; + } +} diff --git a/apps/server/src/modules/data-import/graphql/user-file.type.ts b/apps/server/src/modules/data-import/graphql/user-file.type.ts new file mode 100644 index 0000000..dedaf26 --- /dev/null +++ b/apps/server/src/modules/data-import/graphql/user-file.type.ts @@ -0,0 +1,66 @@ +import { Field, Int, ObjectType } from "@nestjs/graphql"; +import GraphQLJSON from "graphql-type-json"; +import type { UserFile as UserFileDomain } from "../user-file.entity"; + +@ObjectType() +export class UserFileType { + @Field() + id!: string; + + @Field() + profileId!: string; + + @Field() + fileName!: string; + + @Field() + mimeType!: string; + + @Field(() => Int) + sizeBytes!: number; + + @Field() + source!: string; + + @Field() + status!: string; + + @Field() + statusMessage!: string; + + @Field(() => GraphQLJSON, { nullable: true }) + result?: unknown; + + @Field(() => String, { nullable: true }) + error?: string | null; + + @Field(() => Boolean) + isDuplicate!: boolean; + + @Field() + createdAt!: Date; + + @Field() + updatedAt!: Date; + + static fromDomain( + entity: UserFileDomain, + opts?: { isDuplicate?: boolean }, + ): UserFileType { + const type = new UserFileType(); + type.id = entity.id; + type.profileId = entity.profileId; + type.fileName = entity.fileName; + type.mimeType = entity.mimeType; + type.sizeBytes = entity.sizeBytes; + type.source = entity.source; + type.status = entity.status; + type.statusMessage = entity.statusMessage; + type.result = entity.resultJson ? JSON.parse(entity.resultJson) : undefined; + type.error = entity.error; + type.isDuplicate = opts?.isDuplicate ?? false; + type.createdAt = entity.createdAt; + type.updatedAt = entity.updatedAt; + return type; + } +} diff --git a/apps/server/src/modules/data-import/import-job.entity.ts b/apps/server/src/modules/data-import/import-job.entity.ts new file mode 100644 index 0000000..48d1bed --- /dev/null +++ b/apps/server/src/modules/data-import/import-job.entity.ts @@ -0,0 +1,18 @@ +import { BaseEntity } from "@cv/system"; + +export class ImportJob extends BaseEntity { + constructor( + id: string, + public readonly userFileId: string, + public readonly source: string, + public readonly status: string, + public readonly statusMessage: string, + public readonly error: string | null, + public readonly startedAt: Date | null, + public readonly completedAt: Date | null, + createdAt: Date, + updatedAt: Date, + ) { + super(id, createdAt, updatedAt); + } +} diff --git a/apps/server/src/modules/data-import/import-job.mapper.ts b/apps/server/src/modules/data-import/import-job.mapper.ts new file mode 100644 index 0000000..7b34657 --- /dev/null +++ b/apps/server/src/modules/data-import/import-job.mapper.ts @@ -0,0 +1,36 @@ +import type { BaseMapper } from "@cv/system"; +import { Injectable } from "@nestjs/common"; +import type { Prisma } from "@prisma/client"; + +type PrismaImportJob = Prisma.ImportJobGetPayload; + +import { ImportJob } from "./import-job.entity"; + +@Injectable() +export class ImportJobMapper implements BaseMapper { + toDomain(prisma: null): null; + toDomain(prisma: PrismaImportJob): ImportJob; + toDomain(prisma: PrismaImportJob | null): ImportJob | null; + toDomain(prisma: PrismaImportJob | null): ImportJob | null { + if (!prisma) return null; + + return new ImportJob( + prisma.id, + prisma.userFileId, + prisma.source, + prisma.status, + prisma.statusMessage, + prisma.error, + prisma.startedAt, + prisma.completedAt, + prisma.createdAt, + prisma.updatedAt, + ); + } + + mapToDomain(items: PrismaImportJob[]): ImportJob[] { + return items + .map((item) => this.toDomain(item)) + .filter((item): item is ImportJob => item !== null); + } +} diff --git a/apps/server/src/modules/data-import/import-job.policy.ts b/apps/server/src/modules/data-import/import-job.policy.ts new file mode 100644 index 0000000..9ff7041 --- /dev/null +++ b/apps/server/src/modules/data-import/import-job.policy.ts @@ -0,0 +1,28 @@ +import type { IPolicy, User } from "@cv/auth"; +import { Policy } from "@cv/auth"; +import { Injectable } from "@nestjs/common"; +import { ImportJob } from "./import-job.entity"; + +/** + * ImportJob is an internal/admin-only entity. + * All user-facing operations go through UserFile instead. + */ +@Injectable() +@Policy(ImportJob) +export class ImportJobPolicy implements IPolicy { + view(_user: User, _resource: ImportJob): boolean { + return false; + } + + create(_user: User, _resource?: Partial): boolean { + return false; + } + + update(_user: User, _resource: ImportJob): boolean { + return false; + } + + delete(_user: User, _resource: ImportJob): boolean { + return false; + } +} diff --git a/apps/server/src/modules/data-import/import.service.ts b/apps/server/src/modules/data-import/import.service.ts new file mode 100644 index 0000000..4a3673f --- /dev/null +++ b/apps/server/src/modules/data-import/import.service.ts @@ -0,0 +1,86 @@ +import type { User } from "@cv/auth"; +import { PrismaService } from "@cv/system"; +import { Injectable } from "@nestjs/common"; +import { EventEmitter2 } from "@nestjs/event-emitter"; +import type { DataImportSource } from "./data-import-source.interface"; +import { UserFileCreatedEvent } from "./events/user-file-created.event"; +import { UserFile } from "./user-file.entity"; +import { UserFileMapper } from "./user-file.mapper"; + +@Injectable() +export class ImportService { + constructor( + private readonly prisma: PrismaService, + private readonly userFileMapper: UserFileMapper, + private readonly eventEmitter: EventEmitter2, + ) {} + + /** + * Find an existing completed UserFile with the same fingerprint for this user. + */ + async findDuplicateForUser( + userId: string, + fingerprint: string, + ): Promise { + const record = await this.prisma.userFile.findFirst({ + where: { + fingerprint, + profile: { userId }, + status: "completed", + }, + orderBy: { createdAt: "desc" }, + }); + return this.userFileMapper.toDomain(record); + } + + /** + * Create a UserFile and emit an event for async processing. + */ + async createImport( + user: User, + profileId: string, + source: DataImportSource, + file: { fileName: string; mimeType: string; sizeBytes: number; fingerprint?: string }, + params: Record, + ): Promise { + const record = await this.prisma.userFile.create({ + data: { + profile: { connect: { id: profileId } }, + fingerprint: file.fingerprint ?? null, + fileName: file.fileName, + mimeType: file.mimeType, + sizeBytes: file.sizeBytes, + source: source.name, + status: "pending", + statusMessage: "Queued for processing", + }, + }); + + this.eventEmitter.emit( + UserFileCreatedEvent.event, + new UserFileCreatedEvent(record.id, user, source, params), + ); + + return this.userFileMapper.toDomain(record) as UserFile; + } + + /** + * Delete a UserFile and its cascading ImportJobs. + */ + async deleteUserFile(id: string): Promise { + await this.prisma.userFile.delete({ where: { id } }); + } + + async findUserFileById(id: string): Promise { + const record = await this.prisma.userFile.findUnique({ where: { id } }); + return this.userFileMapper.toDomain(record); + } + + async findUserFilesForUser(userId: string): Promise { + const records = await this.prisma.userFile.findMany({ + where: { profile: { userId } }, + orderBy: { createdAt: "desc" }, + }); + return this.userFileMapper.mapToDomain(records); + } +} diff --git a/apps/server/src/modules/data-import/listeners/import-job.listener.ts b/apps/server/src/modules/data-import/listeners/import-job.listener.ts new file mode 100644 index 0000000..49ff425 --- /dev/null +++ b/apps/server/src/modules/data-import/listeners/import-job.listener.ts @@ -0,0 +1,268 @@ +import type { ParsedCVData } from "@cv/ai-parser"; +import { PrismaService } from "@cv/system"; +import { Injectable, Logger } from "@nestjs/common"; +import { OnEvent } from "@nestjs/event-emitter"; +import { + EntityResolverService, + type ResolvedEducation, + type ResolvedJobExperience, +} from "@/modules/cv-parser/entity-resolver.service"; +import { ProfileService } from "@/modules/profile/profile.service"; +import { UserFileCreatedEvent } from "../events/user-file-created.event"; + +interface ResolvedResult { + personalInfo?: { + name?: string | undefined; + headline?: string | undefined; + introduction?: string | undefined; + city?: string | undefined; + country?: string | undefined; + phone?: string | undefined; + website?: string | undefined; + linkedInUrl?: string | undefined; + }; + jobExperiences: ResolvedJobExperience[]; + education: ResolvedEducation[]; +} + +@Injectable() +export class ImportJobListener { + private readonly logger = new Logger(ImportJobListener.name); + + constructor( + private readonly prisma: PrismaService, + private readonly entityResolver: EntityResolverService, + private readonly profileService: ProfileService, + ) {} + + @OnEvent(UserFileCreatedEvent.event, { async: true }) + async handleUserFileCreated(event: UserFileCreatedEvent): Promise { + const job = await this.prisma.importJob.create({ + data: { + userFileId: event.userFileId, + source: event.source.name, + status: "pending", + statusMessage: "Queued", + }, + }); + + await this.processJob(job.id, event); + } + + private async processJob( + jobId: string, + event: UserFileCreatedEvent, + ): Promise { + const { userFileId, user, source, params } = event; + const tag = `job=${jobId} source=${source.name}`; + + try { + await this.startJob(jobId, userFileId); + + const userFile = await this.prisma.userFile.findUniqueOrThrow({ + where: { id: userFileId }, + select: { profileId: true }, + }); + + const onStatus = async (message: string) => { + this.logger.log(`${tag} ${message}`); + await this.updateStatusMessage(jobId, userFileId, message); + }; + + const parsed: ParsedCVData = await source.execute(user, params, onStatus); + + this.logger.log( + `${tag} Parsed: ${parsed.jobExperiences.length} job(s), ` + + `${parsed.education.length} education(s), ${parsed.skills.length} skill(s)`, + ); + this.logger.debug(`${tag} Parsed result: ${JSON.stringify(parsed, null, 2)}`); + + await this.updateStatusMessage(jobId, userFileId, "Resolving entities"); + const resolved = await this.resolveEntities(parsed); + + if (parsed.personalInfo) { + resolved.personalInfo = parsed.personalInfo; + } + + this.logger.log(`${tag} Entity resolution complete`); + + await this.populateProfile(userFile.profileId, parsed, tag); + await this.completeJob(jobId, userFileId, resolved); + } catch (err) { + const message = err instanceof Error ? err.message : "Processing failed"; + this.logger.error(`${tag} Import failed: ${message}`, err instanceof Error ? err.stack : undefined); + await this.failJob(jobId, userFileId, message); + } + } + + private async resolveEntities(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.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, + }; + }), + ), + Promise.all( + parsed.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 { jobExperiences, education }; + } + + /** + * Auto-populate empty profile fields from imported personalInfo. + * Only fills fields that are currently empty -- never overwrites existing data. + */ + private async populateProfile( + profileId: string, + parsed: ParsedCVData, + tag: string, + ): Promise { + const info = parsed.personalInfo; + if (!info) return; + + try { + const existing = await this.profileService.findByIdOrFail(profileId); + const updates: Record = {}; + + const fieldMap: Record = { + fullName: info.name, + headline: info.headline, + summary: info.introduction, + city: info.city, + country: info.country, + phone: info.phone, + website: info.website, + linkedInUrl: info.linkedInUrl, + }; + + for (const [field, value] of Object.entries(fieldMap)) { + if (value && !existing[field as keyof typeof existing]) { + updates[field] = value; + } + } + + if (Object.keys(updates).length === 0) { + this.logger.debug(`${tag} Profile already has data, skipping auto-populate`); + return; + } + + await this.profileService.updateProfile(profileId, updates); + this.logger.log(`${tag} Auto-populated profile: ${Object.keys(updates).join(", ")}`); + } catch (err) { + this.logger.warn(`${tag} Failed to auto-populate profile (non-fatal): ${err instanceof Error ? err.message : err}`); + } + } + + private async startJob(jobId: string, userFileId: string): Promise { + await Promise.all([ + this.prisma.importJob.update({ + where: { id: jobId }, + data: { + status: "processing", + statusMessage: "Starting import", + startedAt: new Date(), + }, + }), + this.prisma.userFile.update({ + where: { id: userFileId }, + data: { status: "processing", statusMessage: "Starting import" }, + }), + ]); + } + + private async updateStatusMessage( + jobId: string, + userFileId: string, + message: string, + ): Promise { + await Promise.all([ + this.prisma.importJob.update({ + where: { id: jobId }, + data: { statusMessage: message }, + }), + this.prisma.userFile.update({ + where: { id: userFileId }, + data: { statusMessage: message }, + }), + ]); + } + + private async completeJob( + jobId: string, + userFileId: string, + result: ResolvedResult, + ): Promise { + await Promise.all([ + this.prisma.importJob.update({ + where: { id: jobId }, + data: { + status: "completed", + statusMessage: "Import completed", + completedAt: new Date(), + }, + }), + this.prisma.userFile.update({ + where: { id: userFileId }, + data: { + status: "completed", + statusMessage: "Import completed", + resultJson: JSON.stringify(result), + }, + }), + ]); + } + + private async failJob( + jobId: string, + userFileId: string, + error: string, + ): Promise { + await Promise.all([ + this.prisma.importJob.update({ + where: { id: jobId }, + data: { + status: "failed", + statusMessage: "Import failed", + error, + completedAt: new Date(), + }, + }), + this.prisma.userFile.update({ + where: { id: userFileId }, + data: { status: "failed", statusMessage: "Import failed", error }, + }), + ]); + } +} diff --git a/apps/server/src/modules/data-import/sources/file-import-source.ts b/apps/server/src/modules/data-import/sources/file-import-source.ts new file mode 100644 index 0000000..d009395 --- /dev/null +++ b/apps/server/src/modules/data-import/sources/file-import-source.ts @@ -0,0 +1,152 @@ +import { + CVParserService as CVParser, + type ExistingUserContext, + type ParsedCVData, +} from "@cv/ai-parser"; +import type { User } from "@cv/auth"; +import { PrismaService } from "@cv/system"; +import { TextExtractorRegistry, TEXT_EXTRACTOR_REGISTRY } from "@cv/file-upload"; +import { Inject, Injectable, Logger } from "@nestjs/common"; +import { EducationService } from "@/modules/education/education.service"; +import { AIProviderResolverService } from "@/modules/cv-parser/ai-provider-resolver.service"; +import { UserJobExperienceService } from "@/modules/job-experience/employment/user-job-experience.service"; +import { ProfileService } from "@/modules/profile/profile.service"; +import type { DataImportSource } from "../data-import-source.interface"; + +/** + * Build ExistingUserContext from domain models for AI prompt enrichment. + */ +const buildExistingUserContext = ( + profile: { fullName?: string | null; headline?: string | null; city?: string | null; country?: string | null } | null, + jobs: Array<{ company: { name: string }; role: { name: string }; startDate: Date; endDate?: Date }>, + educations: Array<{ institution: { name: string }; degree: string; startDate: Date; endDate: Date | null }>, +): ExistingUserContext | undefined => { + const context: ExistingUserContext = {}; + + if (profile?.fullName) context.name = profile.fullName; + if (profile?.headline) context.headline = profile.headline; + if (profile?.city) context.city = profile.city; + if (profile?.country) context.country = profile.country; + + if (jobs.length > 0) { + context.jobs = jobs.map((j) => { + const entry: { company: string; role: string; startDate: string; endDate?: string } = { + company: j.company.name, + role: j.role.name, + startDate: j.startDate.toISOString().slice(0, 10), + }; + if (j.endDate) entry.endDate = j.endDate.toISOString().slice(0, 10); + return entry; + }); + } + + if (educations.length > 0) { + context.education = educations.map((e) => { + const entry: { institution: string; degree: string; startDate: string; endDate?: string } = { + institution: e.institution.name, + degree: e.degree, + startDate: e.startDate.toISOString().slice(0, 10), + }; + if (e.endDate) entry.endDate = e.endDate.toISOString().slice(0, 10); + return entry; + }); + } + + const hasData = Object.keys(context).length > 0; + return hasData ? context : undefined; +}; + +/** + * Wraps the existing file extraction + AI parsing pipeline + * as a DataImportSource. + */ +@Injectable() +export class FileImportSource implements DataImportSource { + readonly name = "file-upload"; + private readonly logger = new Logger(FileImportSource.name); + + constructor( + @Inject(TEXT_EXTRACTOR_REGISTRY) + private readonly textExtractorRegistry: TextExtractorRegistry, + private readonly providerResolver: AIProviderResolverService, + private readonly profileService: ProfileService, + private readonly jobExperienceService: UserJobExperienceService, + private readonly educationService: EducationService, + private readonly prisma: PrismaService, + ) {} + + async execute( + user: User, + params: Record, + onStatus: (message: string) => Promise, + ): Promise { + const buffer = params["buffer"] as Buffer; + const mimeType = params["mimeType"] as string; + + await onStatus("Extracting text from file"); + + const extraction = await this.textExtractorRegistry.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"); + } + + this.logger.log(`Extracted ${extraction.text.length} chars from ${mimeType}`); + this.logger.debug(`Extracted text preview: ${extraction.text.slice(0, 500)}`); + + await onStatus("Gathering existing profile data"); + + const existingContext = await this.gatherUserContext(user); + + await onStatus("Analyzing with AI"); + + const provider = await this.providerResolver.resolveForUser(user); + this.logger.log(`AI provider: ${provider.constructor.name}`); + + const parser = new CVParser(provider); + const result = await parser.parseCVText(extraction.text, existingContext); + + return result; + } + + private async gatherUserContext(user: User): Promise { + try { + const firstProfile = await this.prisma.profile.findFirst({ + where: { userId: user.id }, + select: { id: true }, + }); + + if (!firstProfile) return undefined; + + const [profile, jobs, educationResult] = await Promise.all([ + this.profileService.findByIdOrFail(firstProfile.id), + this.jobExperienceService.findForProfile(firstProfile.id), + this.educationService.findManyForProfile(firstProfile.id), + ]); + + const context = buildExistingUserContext( + profile, + jobs, + educationResult.edges.map((e) => e.node), + ); + + if (context) { + this.logger.debug(`User context: ${JSON.stringify(context)}`); + } + + return context; + } catch (err) { + this.logger.warn( + `Failed to gather user context (non-fatal): ${err instanceof Error ? err.message : err}`, + ); + return undefined; + } + } +} diff --git a/apps/server/src/modules/data-import/user-file.entity.ts b/apps/server/src/modules/data-import/user-file.entity.ts new file mode 100644 index 0000000..bdeea5a --- /dev/null +++ b/apps/server/src/modules/data-import/user-file.entity.ts @@ -0,0 +1,21 @@ +import { BaseEntity } from "@cv/system"; + +export class UserFile extends BaseEntity { + constructor( + id: string, + public readonly profileId: string, + public readonly fingerprint: string | null, + public readonly fileName: string, + public readonly mimeType: string, + public readonly sizeBytes: number, + public readonly source: string, + public readonly status: string, + public readonly statusMessage: string, + public readonly resultJson: string | null, + public readonly error: string | null, + createdAt: Date, + updatedAt: Date, + ) { + super(id, createdAt, updatedAt); + } +} diff --git a/apps/server/src/modules/data-import/user-file.mapper.ts b/apps/server/src/modules/data-import/user-file.mapper.ts new file mode 100644 index 0000000..3e494e6 --- /dev/null +++ b/apps/server/src/modules/data-import/user-file.mapper.ts @@ -0,0 +1,39 @@ +import type { BaseMapper } from "@cv/system"; +import { Injectable } from "@nestjs/common"; +import type { Prisma } from "@prisma/client"; + +type PrismaUserFile = Prisma.UserFileGetPayload; + +import { UserFile } from "./user-file.entity"; + +@Injectable() +export class UserFileMapper implements BaseMapper { + toDomain(prisma: null): null; + toDomain(prisma: PrismaUserFile): UserFile; + toDomain(prisma: PrismaUserFile | null): UserFile | null; + toDomain(prisma: PrismaUserFile | null): UserFile | null { + if (!prisma) return null; + + return new UserFile( + prisma.id, + prisma.profileId, + prisma.fingerprint, + prisma.fileName, + prisma.mimeType, + prisma.sizeBytes, + prisma.source, + prisma.status, + prisma.statusMessage, + prisma.resultJson, + prisma.error, + prisma.createdAt, + prisma.updatedAt, + ); + } + + mapToDomain(items: PrismaUserFile[]): UserFile[] { + return items + .map((item) => this.toDomain(item)) + .filter((item): item is UserFile => item !== null); + } +} diff --git a/apps/server/src/modules/data-import/user-file.policy.ts b/apps/server/src/modules/data-import/user-file.policy.ts new file mode 100644 index 0000000..bb0a727 --- /dev/null +++ b/apps/server/src/modules/data-import/user-file.policy.ts @@ -0,0 +1,12 @@ +import { Policy, ProfileOwnedResourcePolicy } from "@cv/auth"; +import { PrismaService } from "@cv/system"; +import { Injectable } from "@nestjs/common"; +import { UserFile } from "./user-file.entity"; + +@Injectable() +@Policy(UserFile) +export class UserFilePolicy extends ProfileOwnedResourcePolicy { + constructor(prisma: PrismaService) { + super(prisma); + } +} -- 2.51.2