diff --git a/apps/worker/package.json b/apps/worker/package.json new file mode 100644 index 0000000..e756df6 --- /dev/null +++ b/apps/worker/package.json @@ -0,0 +1,39 @@ +{ + "name": "@cv/worker", + "version": "0.0.0", + "private": true, + "type": "commonjs", + "scripts": { + "build": "tsc -p tsconfig.build.json", + "start": "node dist/main.js project-q:work async", + "dev": "nodemon --watch src -e ts --exec \"ts-node src/main.ts project-q:work async\"", + "typecheck": "tsc -p tsconfig.build.json --noEmit", + "lint": "biome check .", + "lint:fix": "biome check --write ." + }, + "dependencies": { + "@cv/system": "*", + "@nestjs/common": "^10.4.7", + "@nestjs/config": "^3.2.0", + "@nestjs/core": "^10.4.7", + "@riotbyte/project-q-core": "link:/Users/niels/Developer/riotbyte/project-q/packages/core", + "@riotbyte/project-q-nestjs": "link:/Users/niels/Developer/riotbyte/project-q/packages/framework/nest", + "@riotbyte/project-q-prisma": "link:/Users/niels/Developer/riotbyte/project-q/packages/transport/prisma", + "@riotbyte/nest-service-locator": "link:/Users/niels/Developer/riotbyte/nest-service-locator", + "@nestjs/event-emitter": "^3.0.1", + "nest-commander": "^3.16.0", + "eventemitter2": "^6.4.9", + "zod": "^4.3.6", + "playwright": "^1.52.0", + "playwright-core": "^1.52.0", + "reflect-metadata": "^0.2.2", + "rxjs": "^7.8.1" + }, + "devDependencies": { + "@cv/tsconfig": "*", + "@types/node": "^22.7.5", + "nodemon": "^3.1.7", + "ts-node": "^10.9.2", + "typescript": "^5.6.3" + } +} diff --git a/apps/worker/src/config.ts b/apps/worker/src/config.ts new file mode 100644 index 0000000..3730286 --- /dev/null +++ b/apps/worker/src/config.ts @@ -0,0 +1,23 @@ +const requireEnv = (key: string): string => { + const value = process.env[key]; + if (!value) { + throw new Error(`Missing required env var: ${key}`); + } + return value; +}; + +export type WorkerConfig = typeof config; + +export const config = { + get databaseUrl() { + return requireEnv("DATABASE_URL"); + }, + queueSchema: process.env["QUEUE_SCHEMA"] ?? "queue", + queueName: process.env["QUEUE_NAME"] ?? "default", + pollIntervalMs: Number(process.env["POLL_INTERVAL_MS"] ?? "1000"), + pdfOutputDir: process.env["PDF_OUTPUT_DIR"] ?? "./pdf-output", + pdfTimeoutMs: Number(process.env["PDF_TIMEOUT_MS"] ?? "30000"), + heartbeatFilePath: + process.env["HEARTBEAT_FILE_PATH"] ?? "/tmp/worker-heartbeat", + heartbeatDbIntervalMs: Number(process.env["HEARTBEAT_DB_INTERVAL_MS"] ?? "0"), +} as const; diff --git a/apps/worker/src/handlers/render-pdf.handler.ts b/apps/worker/src/handlers/render-pdf.handler.ts new file mode 100644 index 0000000..4257acc --- /dev/null +++ b/apps/worker/src/handlers/render-pdf.handler.ts @@ -0,0 +1,41 @@ +import * as fs from "node:fs/promises"; +import * as path from "node:path"; +import { Inject, Injectable, Logger } from "@nestjs/common"; +import type { Envelope, Handler } from "@riotbyte/project-q-core"; +import { HandlerTag } from "@riotbyte/project-q-nestjs"; +import type { WorkerConfig } from "../config"; +import { HtmlToPdfService } from "../pdf/html-to-pdf.service"; +import { WORKER_CONFIG } from "../worker.module"; + +type RenderPdfData = { + cvId: string; + html: string; + requestedBy: string; +}; + +@Injectable() +@HandlerTag.decorator({ handles: "render-pdf" }) +export class RenderPdfHandler implements Handler { + private readonly logger = new Logger(RenderPdfHandler.name); + + constructor( + private readonly pdfService: HtmlToPdfService, + @Inject(WORKER_CONFIG) private readonly config: WorkerConfig, + ) {} + + async handle(envelope: Envelope): Promise { + const { cvId, html, requestedBy } = envelope.message.data as RenderPdfData; + + this.logger.log( + `Rendering PDF for CV ${cvId} (requested by ${requestedBy})`, + ); + + const pdf = await this.pdfService.convert(html, this.config.pdfTimeoutMs); + + await fs.mkdir(this.config.pdfOutputDir, { recursive: true }); + const outputPath = path.join(this.config.pdfOutputDir, `${cvId}.pdf`); + await fs.writeFile(outputPath, pdf); + + this.logger.log(`PDF written to ${outputPath} (${pdf.length} bytes)`); + } +} diff --git a/apps/worker/src/heartbeat/db-heartbeat.strategy.ts b/apps/worker/src/heartbeat/db-heartbeat.strategy.ts new file mode 100644 index 0000000..52bced7 --- /dev/null +++ b/apps/worker/src/heartbeat/db-heartbeat.strategy.ts @@ -0,0 +1,48 @@ +import pg from "pg"; +import type { HeartbeatStrategy } from "./heartbeat.strategy"; + +/** Writes periodic heartbeats to a Postgres table for distributed monitoring. */ +export class DbHeartbeatStrategy implements HeartbeatStrategy { + private pool: pg.Pool | null = null; + private workerId: string | null = null; + private lastWriteAt = 0; + + constructor( + private readonly databaseUrl: string, + private readonly intervalMs: number, + ) {} + + async onStart(workerId: string): Promise { + this.pool = new pg.Pool({ connectionString: this.databaseUrl }); + this.workerId = workerId; + } + + async onRunning(): Promise { + if (!(this.pool && this.workerId)) { + return; + } + + const now = Date.now(); + if (now - this.lastWriteAt < this.intervalMs) { + return; + } + + this.lastWriteAt = now; + + await this.pool.query( + `INSERT INTO queue.worker_heartbeats (worker_id, last_seen_at) + VALUES ($1, now()) + ON CONFLICT (worker_id) + DO UPDATE SET last_seen_at = now()`, + [this.workerId], + ); + } + + async onStop(): Promise { + if (!this.pool) { + return; + } + await this.pool.end(); + this.pool = null; + } +} diff --git a/apps/worker/src/heartbeat/file-heartbeat.strategy.ts b/apps/worker/src/heartbeat/file-heartbeat.strategy.ts new file mode 100644 index 0000000..a762534 --- /dev/null +++ b/apps/worker/src/heartbeat/file-heartbeat.strategy.ts @@ -0,0 +1,18 @@ +import * as fs from "node:fs/promises"; +import type { HeartbeatStrategy } from "./heartbeat.strategy"; + +/** Touches a file on each heartbeat tick for Docker/orchestrator health checks. */ +export class FileHeartbeatStrategy implements HeartbeatStrategy { + constructor(private readonly filePath: string) {} + + async onStart(): Promise { + await fs.writeFile(this.filePath, ""); + } + + async onRunning(): Promise { + const now = new Date(); + await fs.utimes(this.filePath, now, now); + } + + async onStop(): Promise {} +} diff --git a/apps/worker/src/heartbeat/heartbeat.listener.ts b/apps/worker/src/heartbeat/heartbeat.listener.ts new file mode 100644 index 0000000..9764105 --- /dev/null +++ b/apps/worker/src/heartbeat/heartbeat.listener.ts @@ -0,0 +1,32 @@ +import { Inject, Injectable, type OnModuleInit } from "@nestjs/common"; +import { + WorkerRunningEvent, + WorkerStartedEvent, + WorkerStoppedEvent, +} from "@riotbyte/project-q-core"; +import { EventEmitter2 } from "eventemitter2"; +import type { HeartbeatStrategy } from "./heartbeat.strategy"; +import { HEARTBEAT_STRATEGIES } from "./heartbeat.strategy"; + +@Injectable() +export class HeartbeatListener implements OnModuleInit { + constructor( + private readonly emitter: EventEmitter2, + @Inject(HEARTBEAT_STRATEGIES) + private readonly strategies: HeartbeatStrategy[], + ) {} + + onModuleInit(): void { + this.emitter.on(WorkerStartedEvent.name, async (event: WorkerStartedEvent) => { + await Promise.all(this.strategies.map((s) => s.onStart(event.worker.id))); + }); + + this.emitter.on(WorkerRunningEvent.name, async () => { + await Promise.all(this.strategies.map((s) => s.onRunning())); + }); + + this.emitter.on(WorkerStoppedEvent.name, async () => { + await Promise.all(this.strategies.map((s) => s.onStop())); + }); + } +} diff --git a/apps/worker/src/heartbeat/heartbeat.module.ts b/apps/worker/src/heartbeat/heartbeat.module.ts new file mode 100644 index 0000000..8049258 --- /dev/null +++ b/apps/worker/src/heartbeat/heartbeat.module.ts @@ -0,0 +1,37 @@ +import type { DynamicModule } from "@nestjs/common"; +import { Module } from "@nestjs/common"; +import type { WorkerConfig } from "../config"; +import { DbHeartbeatStrategy } from "./db-heartbeat.strategy"; +import { FileHeartbeatStrategy } from "./file-heartbeat.strategy"; +import { HeartbeatListener } from "./heartbeat.listener"; +import type { HeartbeatStrategy } from "./heartbeat.strategy"; +import { HEARTBEAT_STRATEGIES } from "./heartbeat.strategy"; + +@Module({}) +export class HeartbeatModule { + static register(config: WorkerConfig): DynamicModule { + const strategies: HeartbeatStrategy[] = []; + + if (config.heartbeatFilePath) { + strategies.push(new FileHeartbeatStrategy(config.heartbeatFilePath)); + } + + if (config.heartbeatDbIntervalMs > 0) { + strategies.push( + new DbHeartbeatStrategy( + config.databaseUrl, + config.heartbeatDbIntervalMs, + ), + ); + } + + return { + module: HeartbeatModule, + providers: [ + { provide: HEARTBEAT_STRATEGIES, useValue: strategies }, + HeartbeatListener, + ], + exports: [HEARTBEAT_STRATEGIES], + }; + } +} diff --git a/apps/worker/src/heartbeat/heartbeat.strategy.ts b/apps/worker/src/heartbeat/heartbeat.strategy.ts new file mode 100644 index 0000000..afb5416 --- /dev/null +++ b/apps/worker/src/heartbeat/heartbeat.strategy.ts @@ -0,0 +1,8 @@ +export const HEARTBEAT_STRATEGIES = Symbol("HEARTBEAT_STRATEGIES"); + +/** Lifecycle-aware heartbeat strategy for worker health monitoring. */ +export interface HeartbeatStrategy { + onStart(workerId: string): Promise; + onRunning(): Promise; + onStop(): Promise; +} diff --git a/apps/worker/src/logger.provider.ts b/apps/worker/src/logger.provider.ts new file mode 100644 index 0000000..58a7e1d --- /dev/null +++ b/apps/worker/src/logger.provider.ts @@ -0,0 +1,18 @@ +import { Logger as NestLogger } from "@nestjs/common"; +import type { Logger } from "@riotbyte/project-q-core"; + +export class NestProjectQLogger implements Logger { + private readonly logger = new NestLogger("ProjectQ"); + + info(message: string): void { + this.logger.log(message); + } + + error(message: string): void { + this.logger.error(message); + } + + debug(message: string): void { + this.logger.debug(message); + } +} diff --git a/apps/worker/src/main.ts b/apps/worker/src/main.ts new file mode 100644 index 0000000..3dc26f1 --- /dev/null +++ b/apps/worker/src/main.ts @@ -0,0 +1,18 @@ +import "reflect-metadata"; +import { Logger } from "@nestjs/common"; +import { CommandFactory } from "nest-commander"; +import { WorkerModule } from "./worker.module"; + +const logger = new Logger("Bootstrap"); + +async function bootstrap(): Promise { + await CommandFactory.run(WorkerModule, { + logger: ["log", "error", "warn", "debug"], + cliName: "worker", + }); +} + +bootstrap().catch((error) => { + logger.error("Worker crashed", error); + process.exit(1); +}); diff --git a/apps/worker/src/pdf/html-to-pdf.service.ts b/apps/worker/src/pdf/html-to-pdf.service.ts new file mode 100644 index 0000000..2ace953 --- /dev/null +++ b/apps/worker/src/pdf/html-to-pdf.service.ts @@ -0,0 +1,50 @@ +import { Injectable, Logger, OnModuleDestroy } from "@nestjs/common"; +import { chromium } from "playwright"; +import type { Browser } from "playwright-core"; + +/** Manages a Playwright browser singleton for HTML-to-PDF conversion. */ +@Injectable() +export class HtmlToPdfService implements OnModuleDestroy { + private readonly logger = new Logger(HtmlToPdfService.name); + private browser: Browser | null = null; + + private async ensureBrowser(): Promise { + if (this.browser) { + return this.browser; + } + + this.logger.log("Launching Chromium"); + this.browser = await chromium.launch({ + args: ["--no-sandbox", "--disable-dev-shm-usage"], + }); + + return this.browser; + } + + /** Renders HTML to an A4 PDF buffer. */ + async convert(html: string, timeoutMs = 30_000): Promise { + const browser = await this.ensureBrowser(); + const page = await browser.newPage(); + page.setDefaultTimeout(timeoutMs); + + try { + await page.setContent(html, { + waitUntil: "networkidle", + timeout: timeoutMs, + }); + + return await page.pdf({ format: "A4", printBackground: true }); + } finally { + await page.close(); + } + } + + async onModuleDestroy(): Promise { + if (!this.browser) { + return; + } + this.logger.log("Closing Chromium"); + await this.browser.close(); + this.browser = null; + } +} diff --git a/apps/worker/src/worker.module.ts b/apps/worker/src/worker.module.ts new file mode 100644 index 0000000..a306b93 --- /dev/null +++ b/apps/worker/src/worker.module.ts @@ -0,0 +1,41 @@ +import { Module } from "@nestjs/common"; +import { ConfigModule } from "@nestjs/config"; +import { EventEmitterModule } from "@nestjs/event-emitter"; +import { DatabaseModule, PrismaService } from "@cv/system"; +import { LoggerProvider, MessengerModule } from "@riotbyte/project-q-nestjs"; +import { PrismaClientToken, PrismaTransportFactory } from "@riotbyte/project-q-prisma"; +import type { WorkerConfig } from "./config"; +import { config } from "./config"; +import { RenderPdfHandler } from "./handlers/render-pdf.handler"; +import { HeartbeatModule } from "./heartbeat/heartbeat.module"; +import { NestProjectQLogger } from "./logger.provider"; +import { HtmlToPdfService } from "./pdf/html-to-pdf.service"; + +export const WORKER_CONFIG = Symbol("WORKER_CONFIG"); + +@Module({ + imports: [ + ConfigModule.forRoot({ + isGlobal: true, + load: [() => ({ DATABASE_URL: config.databaseUrl })], + }), + DatabaseModule, + EventEmitterModule.forRoot(), + HeartbeatModule.register(config), + MessengerModule.forRoot({ + transports: { + async: { dsn: `prisma://?queue=${config.queueName}`, retry: true }, + }, + routing: new Map([["render-pdf", "async"]]), + }), + ], + providers: [ + { provide: WORKER_CONFIG, useValue: config satisfies WorkerConfig }, + { provide: PrismaClientToken, useExisting: PrismaService }, + { provide: LoggerProvider, useClass: NestProjectQLogger }, + PrismaTransportFactory, + RenderPdfHandler, + HtmlToPdfService, + ], +}) +export class WorkerModule {} diff --git a/apps/worker/tsconfig.build.json b/apps/worker/tsconfig.build.json new file mode 100644 index 0000000..af81f49 --- /dev/null +++ b/apps/worker/tsconfig.build.json @@ -0,0 +1,8 @@ +{ + "extends": "./tsconfig.json", + "compilerOptions": { + "declaration": false, + "sourceMap": true + }, + "exclude": ["node_modules", "dist"] +} diff --git a/apps/worker/tsconfig.json b/apps/worker/tsconfig.json new file mode 100644 index 0000000..33d1d3b --- /dev/null +++ b/apps/worker/tsconfig.json @@ -0,0 +1,10 @@ +{ + "extends": "../../packages/tsconfig/tsconfig.node.json", + "compilerOptions": { + "outDir": "dist", + "preserveSymlinks": true, + "types": ["node"] + }, + "include": ["src/**/*.ts"], + "exclude": ["node_modules", "dist"] +} diff --git a/biome.json b/biome.json index adcff73..b315c11 100644 --- a/biome.json +++ b/biome.json @@ -76,6 +76,7 @@ { "includes": [ "apps/server/**/*", + "apps/worker/**/*", "packages/auth/**/*", "packages/system/**/*", "packages/ai-parser/**/*",