diff --git a/apps/api/src/modules/admin/__tests__/queue-monitor.service.spec.ts b/apps/api/src/modules/admin/__tests__/queue-monitor.service.spec.ts new file mode 100644 index 0000000..3bce1ba --- /dev/null +++ b/apps/api/src/modules/admin/__tests__/queue-monitor.service.spec.ts @@ -0,0 +1,217 @@ +import type { ClockService } from "@cv/core"; +import type { Mock } from "vitest"; +import { beforeEach, describe, expect, it, vi } from "vitest"; +import type { PrismaService } from "@/modules/database/prisma.service"; +import { QueueMonitorService } from "../queue-monitor.service"; + +type PrismaStub = { + message: { count: Mock; aggregate: Mock; findMany: Mock }; + workerHeartbeat: { findMany: Mock }; +}; + +const buildPrisma = (): PrismaStub => ({ + message: { + count: vi.fn(), + aggregate: vi.fn(), + findMany: vi.fn(), + }, + workerHeartbeat: { findMany: vi.fn() }, +}); + +const now = new Date("2026-05-17T12:00:00.000Z"); +const clock: ClockService = { now: () => now }; + +const buildService = (prisma: PrismaStub): QueueMonitorService => + new QueueMonitorService(prisma as unknown as PrismaService, clock); + +describe("QueueMonitorService", () => { + describe("getStats", () => { + let prisma: PrismaStub; + let service: QueueMonitorService; + + beforeEach(() => { + prisma = buildPrisma(); + service = buildService(prisma); + }); + + it("returns zeros and null oldest when the queue is empty", async () => { + prisma.message.count.mockResolvedValue(0); + prisma.message.aggregate.mockResolvedValue({ + _min: { availableAt: null }, + }); + + const result = await service.getStats(); + + expect(result).toEqual({ + pending: 0, + scheduled: 0, + processing: 0, + oldestPendingSeconds: null, + }); + }); + + it("splits PENDING rows into pending vs scheduled by availableAt", async () => { + prisma.message.count + .mockResolvedValueOnce(3) + .mockResolvedValueOnce(2) + .mockResolvedValueOnce(1); + prisma.message.aggregate.mockResolvedValue({ + _min: { availableAt: new Date(now.getTime() - 12_000) }, + }); + + const result = await service.getStats(); + + expect(result).toEqual({ + pending: 3, + scheduled: 2, + processing: 1, + oldestPendingSeconds: 12, + }); + + const [pendingCall, scheduledCall, processingCall] = + prisma.message.count.mock.calls; + expect(pendingCall?.[0]?.where).toMatchObject({ + status: { name: "PENDING" }, + availableAt: { lte: now }, + deliveredAt: null, + }); + expect(scheduledCall?.[0]?.where).toMatchObject({ + status: { name: "PENDING" }, + availableAt: { gt: now }, + deliveredAt: null, + }); + expect(processingCall?.[0]?.where).toMatchObject({ + status: { name: "PENDING" }, + deliveredAt: { not: null }, + }); + }); + }); + + describe("getMessages", () => { + let prisma: PrismaStub; + let service: QueueMonitorService; + + beforeEach(() => { + prisma = buildPrisma(); + service = buildService(prisma); + }); + + const baseRow = { + id: "msg-1", + queue: "default", + name: "parse-cv", + availableAt: new Date(now.getTime() - 5_000), + deliveredAt: null, + createdAt: new Date(now.getTime() - 10_000), + status: { name: "PENDING" }, + }; + + it("derives status from name + availableAt + deliveredAt", async () => { + prisma.message.findMany.mockResolvedValue([ + baseRow, + { + ...baseRow, + id: "msg-2", + availableAt: new Date(now.getTime() + 5_000), + }, + { ...baseRow, id: "msg-3", deliveredAt: now }, + { ...baseRow, id: "msg-4", status: { name: "SUCCESS" } }, + { ...baseRow, id: "msg-5", status: { name: "FAILURE" } }, + ]); + + const result = await service.getMessages(); + + expect(result.map((r) => [r.id, r.status])).toEqual([ + ["msg-1", "pending"], + ["msg-2", "scheduled"], + ["msg-3", "processing"], + ["msg-4", "success"], + ["msg-5", "failure"], + ]); + }); + + it("orders by createdAt desc and clamps the limit", async () => { + prisma.message.findMany.mockResolvedValue([]); + + await service.getMessages(10_000); + + expect(prisma.message.findMany).toHaveBeenCalledWith({ + include: { status: true }, + orderBy: { createdAt: "desc" }, + take: 500, + }); + + await service.getMessages(0); + + expect(prisma.message.findMany).toHaveBeenLastCalledWith({ + include: { status: true }, + orderBy: { createdAt: "desc" }, + take: 1, + }); + }); + + it("maps queue/name onto the domain shape", async () => { + prisma.message.findMany.mockResolvedValue([baseRow]); + + const [result] = await service.getMessages(); + + expect(result).toEqual({ + id: "msg-1", + queueName: "default", + messageName: "parse-cv", + status: "pending", + createdAt: baseRow.createdAt, + }); + }); + }); + + describe("getWorkers", () => { + let prisma: PrismaStub; + let service: QueueMonitorService; + + beforeEach(() => { + prisma = buildPrisma(); + service = buildService(prisma); + }); + + it("derives healthy / stale / dead from lastSeenAt age", async () => { + prisma.workerHeartbeat.findMany.mockResolvedValue([ + { + workerId: "fresh", + startedAt: new Date(now.getTime() - 60_000), + lastSeenAt: new Date(now.getTime() - 30_000), + }, + { + workerId: "stalish", + startedAt: new Date(now.getTime() - 600_000), + lastSeenAt: new Date(now.getTime() - 120_000), + }, + { + workerId: "gone", + startedAt: new Date(now.getTime() - 10_000_000), + lastSeenAt: new Date(now.getTime() - 1_000_000), + }, + ]); + + const result = await service.getWorkers(); + + expect(result.map((r) => [r.workerId, r.status])).toEqual([ + ["fresh", "healthy"], + ["stalish", "stale"], + ["gone", "dead"], + ]); + + expect(prisma.workerHeartbeat.findMany).toHaveBeenCalledWith({ + orderBy: { startedAt: "desc" }, + }); + }); + + it("returns an empty list when nothing has registered", async () => { + prisma.workerHeartbeat.findMany.mockResolvedValue([]); + + const result = await service.getWorkers(); + + expect(result).toEqual([]); + }); + }); +}); diff --git a/apps/api/src/modules/admin/queue-monitor.service.ts b/apps/api/src/modules/admin/queue-monitor.service.ts index a453e5d..1018cd1 100644 --- a/apps/api/src/modules/admin/queue-monitor.service.ts +++ b/apps/api/src/modules/admin/queue-monitor.service.ts @@ -2,28 +2,6 @@ import { ClockService } from "@cv/core"; import { Injectable } from "@nestjs/common"; import { PrismaService } from "@/modules/database/prisma.service"; -interface QueueStatsRow { - pending: bigint; - scheduled: bigint; - processing: bigint; - oldest_pending_seconds: number | null; -} - -interface QueueMessageRow { - id: string; - queue_name: string; - message_name: string | null; - created_at: Date; - available_at: Date; - delivered_at: Date | null; -} - -interface WorkerHeartbeatRow { - worker_id: string; - started_at: Date; - last_seen_at: Date; -} - export interface QueueStatsResult { pending: number; scheduled: number; @@ -35,7 +13,7 @@ export interface QueueMessageResult { id: string; queueName: string; messageName: string | null; - status: "pending" | "scheduled" | "processing"; + status: "pending" | "scheduled" | "processing" | "success" | "failure"; createdAt: Date; } @@ -48,29 +26,41 @@ export interface WorkerHealthResult { const HEALTHY_THRESHOLD_SECONDS = 60; const STALE_THRESHOLD_SECONDS = 300; +const DEFAULT_MESSAGE_LIMIT = 50; +const MAX_MESSAGE_LIMIT = 500; const deriveWorkerStatus = ( lastSeenAt: Date, now: Date, -): "healthy" | "stale" | "dead" => { +): WorkerHealthResult["status"] => { const ageSeconds = (now.getTime() - lastSeenAt.getTime()) / 1000; - if (ageSeconds <= HEALTHY_THRESHOLD_SECONDS) return "healthy"; - if (ageSeconds <= STALE_THRESHOLD_SECONDS) return "stale"; + if (ageSeconds <= HEALTHY_THRESHOLD_SECONDS) { + return "healthy"; + } + if (ageSeconds <= STALE_THRESHOLD_SECONDS) { + return "stale"; + } return "dead"; }; const deriveMessageStatus = ( - row: QueueMessageRow, + statusName: string, + availableAt: Date, + deliveredAt: Date | null, now: Date, -): "pending" | "scheduled" | "processing" => { - if (row.delivered_at) return "processing"; - return row.available_at > now ? "scheduled" : "pending"; +): QueueMessageResult["status"] => { + if (statusName === "SUCCESS") { + return "success"; + } + if (statusName === "FAILURE") { + return "failure"; + } + if (deliveredAt) { + return "processing"; + } + return availableAt > now ? "scheduled" : "pending"; }; -/** - * Read-only monitoring service for the project-q queue tables. - * Uses raw SQL since queue tables live in the `queue` schema (not managed by Prisma). - */ @Injectable() export class QueueMonitorService { constructor( @@ -79,54 +69,83 @@ export class QueueMonitorService { ) {} async getStats(): Promise { - const rows = await this.prisma.$queryRaw` - SELECT - COUNT(*) FILTER (WHERE delivered_at IS NULL AND available_at <= now()) AS pending, - COUNT(*) FILTER (WHERE delivered_at IS NULL AND available_at > now()) AS scheduled, - COUNT(*) FILTER (WHERE delivered_at IS NOT NULL) AS processing, - EXTRACT(EPOCH FROM now() - MIN(available_at) FILTER (WHERE delivered_at IS NULL AND available_at <= now())) AS oldest_pending_seconds - FROM queue.messages - `; + const now = this.clock.now(); - const row = rows[0]; - return { - pending: Number(row?.pending ?? 0), - scheduled: Number(row?.scheduled ?? 0), - processing: Number(row?.processing ?? 0), - oldestPendingSeconds: row?.oldest_pending_seconds ?? null, - }; + const [pending, scheduled, processing, oldest] = await Promise.all([ + this.prisma.message.count({ + where: { + status: { name: "PENDING" }, + availableAt: { lte: now }, + deliveredAt: null, + }, + }), + this.prisma.message.count({ + where: { + status: { name: "PENDING" }, + availableAt: { gt: now }, + deliveredAt: null, + }, + }), + this.prisma.message.count({ + where: { + status: { name: "PENDING" }, + deliveredAt: { not: null }, + }, + }), + this.prisma.message.aggregate({ + _min: { availableAt: true }, + where: { + status: { name: "PENDING" }, + availableAt: { lte: now }, + deliveredAt: null, + }, + }), + ]); + + const oldestAvailable = oldest._min.availableAt; + const oldestPendingSeconds = oldestAvailable + ? (now.getTime() - oldestAvailable.getTime()) / 1000 + : null; + + return { pending, scheduled, processing, oldestPendingSeconds }; } - async getMessages(limit = 50): Promise { - const clampedLimit = Math.max(1, Math.min(limit, 500)); - const rows = await this.prisma.$queryRaw` - SELECT id, queue_name, body->'message'->>'name' AS message_name, - created_at, available_at, delivered_at - FROM queue.messages ORDER BY created_at DESC LIMIT ${clampedLimit} - `; + async getMessages( + limit = DEFAULT_MESSAGE_LIMIT, + ): Promise { + const clampedLimit = Math.max(1, Math.min(limit, MAX_MESSAGE_LIMIT)); + const rows = await this.prisma.message.findMany({ + include: { status: true }, + orderBy: { createdAt: "desc" }, + take: clampedLimit, + }); const now = this.clock.now(); return rows.map((row) => ({ - id: String(row.id), - queueName: row.queue_name, - messageName: row.message_name, - status: deriveMessageStatus(row, now), - createdAt: row.created_at, + id: row.id, + queueName: row.queue, + messageName: row.name, + status: deriveMessageStatus( + row.status.name, + row.availableAt, + row.deliveredAt, + now, + ), + createdAt: row.createdAt, })); } async getWorkers(): Promise { - const rows = await this.prisma.$queryRaw` - SELECT worker_id, started_at, last_seen_at - FROM queue.worker_heartbeats ORDER BY started_at DESC - `; + const rows = await this.prisma.workerHeartbeat.findMany({ + orderBy: { startedAt: "desc" }, + }); const now = this.clock.now(); return rows.map((row) => ({ - workerId: row.worker_id, - status: deriveWorkerStatus(row.last_seen_at, now), - startedAt: row.started_at, - lastSeenAt: row.last_seen_at, + workerId: row.workerId, + status: deriveWorkerStatus(row.lastSeenAt, now), + startedAt: row.startedAt, + lastSeenAt: row.lastSeenAt, })); } } diff --git a/apps/worker/package.json b/apps/worker/package.json index 6abe541..dd18f21 100644 --- a/apps/worker/package.json +++ b/apps/worker/package.json @@ -8,6 +8,8 @@ "start": "node dist/main.js project-q:work async", "dev": "nodemon --watch src -e ts --exec \"ts-node -r tsconfig-paths/register src/main.ts project-q:work async\"", "typecheck": "tsc -p tsconfig.build.json --noEmit", + "test": "vitest run", + "test:watch": "vitest", "lint": "biome check .", "lint:fix": "biome check --write ." }, @@ -27,7 +29,6 @@ "@riotbyte-com/project-q-prisma": "1.0.19-rc.5", "eventemitter2": "^6.4.9", "nest-commander": "^3.16.0", - "pg": "^8.16.3", "playwright": "^1.52.0", "playwright-core": "^1.52.0", "reflect-metadata": "^0.2.2", @@ -39,10 +40,11 @@ "@cv/tsconfig": "*", "@swc/core": "^1.15.13", "@types/node": "^22.7.5", - "@types/pg": "^8.15.0", "nodemon": "^3.1.7", "ts-node": "^10.9.2", "tsconfig-paths": "^4.2.0", - "typescript": "^5.6.3" + "typescript": "^5.6.3", + "unplugin-swc": "^1.5.9", + "vitest": "^4.0.16" } } diff --git a/apps/worker/src/config.ts b/apps/worker/src/config.ts index b681337..c244b85 100644 --- a/apps/worker/src/config.ts +++ b/apps/worker/src/config.ts @@ -34,7 +34,9 @@ export const config = { 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"), + heartbeatDbIntervalMs: Number( + process.env["HEARTBEAT_DB_INTERVAL_MS"] ?? "10000", + ), // `--no-sandbox` (CVG-64): Chromium's setuid sandbox conflicts with the // container's user namespaces, so we disable it to keep the worker image // simple and free of CAP_SYS_ADMIN. The defense-in-depth this normally diff --git a/apps/worker/src/heartbeat/__tests__/db-heartbeat.strategy.spec.ts b/apps/worker/src/heartbeat/__tests__/db-heartbeat.strategy.spec.ts new file mode 100644 index 0000000..1a35e05 --- /dev/null +++ b/apps/worker/src/heartbeat/__tests__/db-heartbeat.strategy.spec.ts @@ -0,0 +1,88 @@ +import type { PrismaService } from "@cv/core"; +import type { Mock } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { DbHeartbeatStrategy } from "../db-heartbeat.strategy"; + +type PrismaStub = { + workerHeartbeat: { upsert: Mock }; +}; + +const buildPrisma = (): PrismaStub => ({ + workerHeartbeat: { upsert: vi.fn().mockResolvedValue(undefined) }, +}); + +const buildStrategy = ( + prisma: PrismaStub, + intervalMs = 10_000, +): DbHeartbeatStrategy => + new DbHeartbeatStrategy(prisma as unknown as PrismaService, intervalMs); + +describe("DbHeartbeatStrategy", () => { + let prisma: PrismaStub; + + beforeEach(() => { + prisma = buildPrisma(); + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-05-17T12:00:00.000Z")); + }); + + afterEach(() => { + vi.useRealTimers(); + }); + + it("upserts an initial heartbeat keyed on workerId during onStart", async () => { + const strategy = buildStrategy(prisma); + + await strategy.onStart("worker-1"); + + expect(prisma.workerHeartbeat.upsert).toHaveBeenCalledTimes(1); + const call = prisma.workerHeartbeat.upsert.mock.calls[0]?.[0]; + expect(call?.where).toEqual({ workerId: "worker-1" }); + expect(call?.create).toMatchObject({ workerId: "worker-1" }); + }); + + it("does not write again until the interval has elapsed", async () => { + const strategy = buildStrategy(prisma, 10_000); + await strategy.onStart("worker-1"); + prisma.workerHeartbeat.upsert.mockClear(); + + vi.advanceTimersByTime(5_000); + await strategy.onRunning(); + expect(prisma.workerHeartbeat.upsert).not.toHaveBeenCalled(); + + vi.advanceTimersByTime(5_000); + await strategy.onRunning(); + expect(prisma.workerHeartbeat.upsert).toHaveBeenCalledTimes(1); + }); + + it("advances lastSeenAt on each interval-window write", async () => { + const strategy = buildStrategy(prisma, 10_000); + await strategy.onStart("worker-1"); + prisma.workerHeartbeat.upsert.mockClear(); + + vi.advanceTimersByTime(10_000); + await strategy.onRunning(); + + vi.advanceTimersByTime(10_000); + await strategy.onRunning(); + + const lastSeenValues = prisma.workerHeartbeat.upsert.mock.calls.map( + ([args]) => (args.update as { lastSeenAt: Date }).lastSeenAt.getTime(), + ); + expect(lastSeenValues).toEqual([ + new Date("2026-05-17T12:00:10.000Z").getTime(), + new Date("2026-05-17T12:00:20.000Z").getTime(), + ]); + for (const call of prisma.workerHeartbeat.upsert.mock.calls) { + expect(call[0].where).toEqual({ workerId: "worker-1" }); + } + }); + + it("is a no-op before onStart has bound a workerId", async () => { + const strategy = buildStrategy(prisma); + + await strategy.onRunning(); + + expect(prisma.workerHeartbeat.upsert).not.toHaveBeenCalled(); + }); +}); diff --git a/apps/worker/src/heartbeat/db-heartbeat.strategy.ts b/apps/worker/src/heartbeat/db-heartbeat.strategy.ts index ed0123b..37754aa 100644 --- a/apps/worker/src/heartbeat/db-heartbeat.strategy.ts +++ b/apps/worker/src/heartbeat/db-heartbeat.strategy.ts @@ -1,24 +1,37 @@ -import pg from "pg"; +import { PrismaService } from "@cv/core"; +import { Inject, Injectable } from "@nestjs/common"; import { HeartbeatStrategy } from "./heartbeat.strategy"; +import { HEARTBEAT_DB_INTERVAL_MS } from "./heartbeat.tokens"; -/** Writes periodic heartbeats to a Postgres table for distributed monitoring. */ +/** Upserts the worker's `WorkerHeartbeat` row on a fixed interval. */ +@Injectable() export class DbHeartbeatStrategy implements HeartbeatStrategy { - private pool: pg.Pool | null = null; private workerId: string | null = null; + private startedAt: Date | null = null; private lastWriteAt = 0; constructor( - private readonly databaseUrl: string, - private readonly intervalMs: number, + private readonly prisma: PrismaService, + @Inject(HEARTBEAT_DB_INTERVAL_MS) private readonly intervalMs: number, ) {} async onStart(workerId: string): Promise { - this.pool = new pg.Pool({ connectionString: this.databaseUrl }); this.workerId = workerId; + this.startedAt = new Date(); + await this.prisma.workerHeartbeat.upsert({ + where: { workerId }, + update: { startedAt: this.startedAt, lastSeenAt: this.startedAt }, + create: { + workerId, + startedAt: this.startedAt, + lastSeenAt: this.startedAt, + }, + }); + this.lastWriteAt = this.startedAt.getTime(); } async onRunning(): Promise { - if (!(this.pool && this.workerId)) { + if (!this.workerId) { return; } @@ -26,23 +39,19 @@ export class DbHeartbeatStrategy implements HeartbeatStrategy { if (now - this.lastWriteAt < this.intervalMs) { return; } - this.lastWriteAt = now; - await this.pool.query( - `INSERT INTO worker_heartbeats (worker_id, last_seen_at) - VALUES ($1, now()) - ON CONFLICT (worker_id) - DO UPDATE SET last_seen_at = now()`, - [this.workerId], - ); + const lastSeenAt = new Date(now); + await this.prisma.workerHeartbeat.upsert({ + where: { workerId: this.workerId }, + update: { lastSeenAt }, + create: { + workerId: this.workerId, + startedAt: this.startedAt ?? lastSeenAt, + lastSeenAt, + }, + }); } - async onStop(): Promise { - if (!this.pool) { - return; - } - await this.pool.end(); - this.pool = null; - } + async onStop(): Promise {} } diff --git a/apps/worker/src/heartbeat/heartbeat.module.ts b/apps/worker/src/heartbeat/heartbeat.module.ts index 513b2cb..5e16e67 100644 --- a/apps/worker/src/heartbeat/heartbeat.module.ts +++ b/apps/worker/src/heartbeat/heartbeat.module.ts @@ -1,34 +1,41 @@ -import { DynamicModule, Module } from "@nestjs/common"; +import { DatabaseModule } from "@cv/core"; +import { DynamicModule, Module, Provider } from "@nestjs/common"; import { WorkerConfig } from "../config"; import { DbHeartbeatStrategy } from "./db-heartbeat.strategy"; import { FileHeartbeatStrategy } from "./file-heartbeat.strategy"; import { HeartbeatListener } from "./heartbeat.listener"; import { HEARTBEAT_STRATEGIES, HeartbeatStrategy } from "./heartbeat.strategy"; +import { HEARTBEAT_DB_INTERVAL_MS } from "./heartbeat.tokens"; @Module({}) export class HeartbeatModule { static register(config: WorkerConfig): DynamicModule { - const strategies: HeartbeatStrategy[] = []; + const dbEnabled = config.heartbeatDbIntervalMs > 0; + const fileStrategy = config.heartbeatFilePath + ? new FileHeartbeatStrategy(config.heartbeatFilePath) + : null; - if (config.heartbeatFilePath) { - strategies.push(new FileHeartbeatStrategy(config.heartbeatFilePath)); - } - - if (config.heartbeatDbIntervalMs > 0) { - strategies.push( - new DbHeartbeatStrategy( - config.databaseUrl, - config.heartbeatDbIntervalMs, - ), - ); - } + const providers: Provider[] = [ + { + provide: HEARTBEAT_DB_INTERVAL_MS, + useValue: config.heartbeatDbIntervalMs, + }, + ...(dbEnabled ? [DbHeartbeatStrategy] : []), + { + provide: HEARTBEAT_STRATEGIES, + inject: dbEnabled ? [DbHeartbeatStrategy] : [], + useFactory: (db?: DbHeartbeatStrategy): HeartbeatStrategy[] => [ + ...(fileStrategy ? [fileStrategy] : []), + ...(db ? [db] : []), + ], + }, + HeartbeatListener, + ]; return { module: HeartbeatModule, - providers: [ - { provide: HEARTBEAT_STRATEGIES, useValue: strategies }, - HeartbeatListener, - ], + imports: [DatabaseModule], + providers, exports: [HEARTBEAT_STRATEGIES], }; } diff --git a/apps/worker/src/heartbeat/heartbeat.tokens.ts b/apps/worker/src/heartbeat/heartbeat.tokens.ts new file mode 100644 index 0000000..242fa01 --- /dev/null +++ b/apps/worker/src/heartbeat/heartbeat.tokens.ts @@ -0,0 +1 @@ +export const HEARTBEAT_DB_INTERVAL_MS = Symbol("HEARTBEAT_DB_INTERVAL_MS"); diff --git a/apps/worker/tsconfig.build.json b/apps/worker/tsconfig.build.json index b671829..79521d1 100644 --- a/apps/worker/tsconfig.build.json +++ b/apps/worker/tsconfig.build.json @@ -5,7 +5,7 @@ "sourceMap": true, "rootDir": "src" }, - "exclude": ["node_modules", "dist"], + "exclude": ["node_modules", "dist", "**/*.spec.ts", "**/__tests__/**"], "references": [ { "path": "../../packages/core" }, { "path": "../../packages/file-storage" }, diff --git a/apps/worker/vitest.config.ts b/apps/worker/vitest.config.ts new file mode 100644 index 0000000..46cfcf9 --- /dev/null +++ b/apps/worker/vitest.config.ts @@ -0,0 +1,22 @@ +import swc from "unplugin-swc"; +import { defineConfig } from "vitest/config"; + +export default defineConfig({ + plugins: [ + swc.vite({ + jsc: { + parser: { syntax: "typescript", decorators: true }, + transform: { + legacyDecorator: true, + decoratorMetadata: true, + }, + }, + }), + ], + test: { + globals: true, + environment: "node", + include: ["src/**/*.{test,spec}.{ts,tsx}"], + exclude: ["**/node_modules/**", "**/dist/**", "**/build/**"], + }, +}); diff --git a/packages/core/prisma/migrations/20260517084853_add_worker_heartbeats_prisma/migration.sql b/packages/core/prisma/migrations/20260517084853_add_worker_heartbeats_prisma/migration.sql new file mode 100644 index 0000000..07d783e --- /dev/null +++ b/packages/core/prisma/migrations/20260517084853_add_worker_heartbeats_prisma/migration.sql @@ -0,0 +1,11 @@ +-- CreateTable +CREATE TABLE "worker_heartbeats" ( + "workerId" TEXT NOT NULL, + "startedAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + "lastSeenAt" TIMESTAMP(3) NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT "worker_heartbeats_pkey" PRIMARY KEY ("workerId") +); + +-- CreateIndex +CREATE INDEX "worker_heartbeats_lastSeenAt_idx" ON "worker_heartbeats"("lastSeenAt"); diff --git a/packages/core/prisma/models/worker-heartbeat.prisma b/packages/core/prisma/models/worker-heartbeat.prisma new file mode 100644 index 0000000..dd24f3c --- /dev/null +++ b/packages/core/prisma/models/worker-heartbeat.prisma @@ -0,0 +1,12 @@ +/// Per-worker liveness signal. Workers upsert their own row on a fixed +/// interval (see `apps/worker/.../db-heartbeat.strategy.ts`); monitoring +/// derives `healthy` / `stale` / `dead` from `lastSeenAt`. `workerId` is the +/// primary key so the upsert is a single round-trip with no auxiliary index. +model WorkerHeartbeat { + workerId String @id + startedAt DateTime @default(now()) + lastSeenAt DateTime @default(now()) + + @@index([lastSeenAt]) + @@map("worker_heartbeats") +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 6d40219..2d079d4 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -558,9 +558,6 @@ importers: nest-commander: specifier: ^3.16.0 version: 3.16.0(@nestjs/common@11.1.18(class-transformer@0.5.1)(class-validator@0.14.3)(reflect-metadata@0.2.2)(rxjs@7.8.2))(@nestjs/core@11.1.18)(@types/inquirer@8.2.12)(typescript@5.9.3) - pg: - specifier: ^8.16.3 - version: 8.16.3 playwright: specifier: ^1.52.0 version: 1.58.2 @@ -589,9 +586,6 @@ importers: '@types/node': specifier: ^22.7.5 version: 22.19.3 - '@types/pg': - specifier: ^8.15.0 - version: 8.16.0 nodemon: specifier: ^3.1.7 version: 3.1.11 @@ -604,6 +598,12 @@ 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.9.0) packages/ai-parser: dependencies: