feat(phase3): worker container + event bus + scheduled jobs
apps/worker: - BullMQ nightly scheduler (cron 0 3 * * *) - Redis Streams consumer-group per wired/active project - Persists events to Event model schema: - Event model (streamId unique, project + type indexed)
This commit is contained in:
26
apps/worker/Dockerfile
Normal file
26
apps/worker/Dockerfile
Normal file
@@ -0,0 +1,26 @@
|
||||
FROM node:22-alpine AS base
|
||||
RUN apk add --no-cache libc6-compat openssl
|
||||
WORKDIR /app
|
||||
RUN corepack enable && corepack prepare pnpm@9.12.0 --activate
|
||||
|
||||
# ---- builder: install everything (need apps/web prisma client) ----
|
||||
FROM base AS builder
|
||||
ENV NODE_ENV=development
|
||||
COPY package.json pnpm-lock.yaml* pnpm-workspace.yaml turbo.json ./
|
||||
COPY apps/web/package.json apps/web/package.json
|
||||
COPY apps/worker/package.json apps/worker/package.json
|
||||
RUN pnpm install --frozen-lockfile --prod=false || pnpm install --prod=false
|
||||
COPY . .
|
||||
WORKDIR /app/apps/web
|
||||
RUN npx --no-install prisma generate
|
||||
|
||||
# ---- runner ----
|
||||
FROM base AS runner
|
||||
ENV NODE_ENV=production
|
||||
WORKDIR /app
|
||||
COPY --from=builder /app/package.json /app/pnpm-workspace.yaml /app/turbo.json ./
|
||||
COPY --from=builder /app/node_modules ./node_modules
|
||||
COPY --from=builder /app/apps/web ./apps/web
|
||||
COPY --from=builder /app/apps/worker ./apps/worker
|
||||
|
||||
CMD ["pnpm","--filter","@panel/worker","start"]
|
||||
21
apps/worker/package.json
Normal file
21
apps/worker/package.json
Normal file
@@ -0,0 +1,21 @@
|
||||
{
|
||||
"name": "@panel/worker",
|
||||
"version": "0.1.0",
|
||||
"private": true,
|
||||
"scripts": {
|
||||
"dev": "tsx watch src/index.ts",
|
||||
"start": "tsx src/index.ts",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@panel/web": "workspace:*",
|
||||
"@prisma/client": "^5.22.0",
|
||||
"bullmq": "^5.34.0",
|
||||
"ioredis": "^5.4.1",
|
||||
"tsx": "^4.19.2"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
94
apps/worker/src/consumers/event-bus.ts
Normal file
94
apps/worker/src/consumers/event-bus.ts
Normal file
@@ -0,0 +1,94 @@
|
||||
import { redis } from "../redis";
|
||||
import { prisma } from "../db";
|
||||
|
||||
const GROUP = "panel";
|
||||
const CONSUMER = `panel-worker-${process.env.HOSTNAME ?? "1"}`;
|
||||
|
||||
type RawEvent = {
|
||||
event_type: string;
|
||||
version?: string | number;
|
||||
occurred_at?: string;
|
||||
payload?: string;
|
||||
};
|
||||
|
||||
async function ensureGroup(stream: string) {
|
||||
try {
|
||||
await redis.xgroup("CREATE", stream, GROUP, "0", "MKSTREAM");
|
||||
console.log(`[event-bus] group created for ${stream}`);
|
||||
} catch (e) {
|
||||
const msg = e instanceof Error ? e.message : String(e);
|
||||
if (!msg.includes("BUSYGROUP")) throw e;
|
||||
}
|
||||
}
|
||||
|
||||
async function consume(stream: string, projectKey: string) {
|
||||
await ensureGroup(stream);
|
||||
while (true) {
|
||||
try {
|
||||
const res = await redis.xreadgroup(
|
||||
"GROUP", GROUP, CONSUMER,
|
||||
"COUNT", 32,
|
||||
"BLOCK", 15000,
|
||||
"STREAMS", stream, ">",
|
||||
) as Array<[string, Array<[string, string[]]>]> | null;
|
||||
|
||||
if (!res) continue;
|
||||
|
||||
for (const [, entries] of res) {
|
||||
for (const [streamId, fields] of entries) {
|
||||
const obj = fieldsToObj(fields);
|
||||
try {
|
||||
await persist(streamId, projectKey, obj);
|
||||
await redis.xack(stream, GROUP, streamId);
|
||||
} catch (err) {
|
||||
console.error(`[event-bus] persist failed ${stream} ${streamId}:`, err);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
console.error(`[event-bus] xreadgroup error for ${stream}:`, err);
|
||||
await new Promise((r) => setTimeout(r, 2000));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function fieldsToObj(fields: string[]): RawEvent {
|
||||
const o: Record<string, string> = {};
|
||||
for (let i = 0; i + 1 < fields.length; i += 2) o[fields[i]] = fields[i + 1];
|
||||
return o as RawEvent;
|
||||
}
|
||||
|
||||
async function persist(streamId: string, projectKey: string, raw: RawEvent) {
|
||||
let payload: unknown = null;
|
||||
try {
|
||||
payload = raw.payload ? JSON.parse(raw.payload) : null;
|
||||
} catch {
|
||||
payload = raw.payload ?? null;
|
||||
}
|
||||
await prisma.event.create({
|
||||
data: {
|
||||
streamId,
|
||||
projectKey,
|
||||
eventType: raw.event_type ?? "unknown",
|
||||
version: raw.version ? Number(raw.version) : 1,
|
||||
occurredAt: raw.occurred_at ? new Date(raw.occurred_at) : new Date(),
|
||||
payload: payload as object,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
export async function startEventBus() {
|
||||
const projects = await prisma.project.findMany({
|
||||
where: { status: { in: ["wired", "active"] } },
|
||||
select: { key: true },
|
||||
});
|
||||
if (projects.length === 0) {
|
||||
console.log("[event-bus] no active projects; idle");
|
||||
return;
|
||||
}
|
||||
for (const p of projects) {
|
||||
const stream = `${p.key}:events`;
|
||||
console.log(`[event-bus] consuming ${stream}`);
|
||||
void consume(stream, p.key);
|
||||
}
|
||||
}
|
||||
3
apps/worker/src/db.ts
Normal file
3
apps/worker/src/db.ts
Normal file
@@ -0,0 +1,3 @@
|
||||
import { PrismaClient } from "@prisma/client";
|
||||
|
||||
export const prisma = new PrismaClient({ log: ["error"] });
|
||||
32
apps/worker/src/index.ts
Normal file
32
apps/worker/src/index.ts
Normal file
@@ -0,0 +1,32 @@
|
||||
import { startEventBus } from "./consumers/event-bus";
|
||||
import { startScheduledJobs } from "./schedulers/nightly";
|
||||
import { redis } from "./redis";
|
||||
import { prisma } from "./db";
|
||||
|
||||
async function main() {
|
||||
console.log("[worker] starting…");
|
||||
await redis.ping();
|
||||
console.log("[worker] redis ok");
|
||||
await prisma.$queryRaw`SELECT 1`;
|
||||
console.log("[worker] panel-db ok");
|
||||
|
||||
await startScheduledJobs();
|
||||
await startEventBus();
|
||||
|
||||
console.log("[worker] up.");
|
||||
}
|
||||
|
||||
const shutdown = async (sig: string) => {
|
||||
console.log(`[worker] ${sig} — shutting down`);
|
||||
await prisma.$disconnect().catch(() => {});
|
||||
await redis.quit().catch(() => {});
|
||||
process.exit(0);
|
||||
};
|
||||
|
||||
process.on("SIGTERM", () => void shutdown("SIGTERM"));
|
||||
process.on("SIGINT", () => void shutdown("SIGINT"));
|
||||
|
||||
main().catch((e) => {
|
||||
console.error("[worker] fatal:", e);
|
||||
process.exit(1);
|
||||
});
|
||||
13
apps/worker/src/redis.ts
Normal file
13
apps/worker/src/redis.ts
Normal file
@@ -0,0 +1,13 @@
|
||||
import IORedis from "ioredis";
|
||||
|
||||
const url = process.env.REDIS_URL;
|
||||
if (!url) throw new Error("REDIS_URL not set");
|
||||
|
||||
export const redis = new IORedis(url, {
|
||||
maxRetriesPerRequest: null,
|
||||
enableReadyCheck: true,
|
||||
});
|
||||
|
||||
redis.on("error", (err) => {
|
||||
console.error("[redis]", err.message);
|
||||
});
|
||||
30
apps/worker/src/schedulers/nightly.ts
Normal file
30
apps/worker/src/schedulers/nightly.ts
Normal file
@@ -0,0 +1,30 @@
|
||||
import { Queue, Worker, type Job } from "bullmq";
|
||||
import { redis } from "../redis";
|
||||
import { prisma } from "../db";
|
||||
|
||||
const QUEUE = "nightly";
|
||||
|
||||
const queue = new Queue(QUEUE, { connection: redis });
|
||||
|
||||
async function nightlyRefresh(_job: Job) {
|
||||
// Placeholder: when materialized views are added (Phase 4+), refresh them here.
|
||||
const projectCount = await prisma.project.count();
|
||||
console.log(`[nightly] refresh run at ${new Date().toISOString()} — ${projectCount} projects`);
|
||||
return { ok: true, projects: projectCount };
|
||||
}
|
||||
|
||||
export async function startScheduledJobs() {
|
||||
// BullMQ repeat: every day at 03:00 server time
|
||||
await queue.upsertJobScheduler(
|
||||
"nightly-refresh",
|
||||
{ pattern: "0 3 * * *" },
|
||||
{
|
||||
name: "nightly-refresh",
|
||||
data: {},
|
||||
opts: { removeOnComplete: 50, removeOnFail: 50 },
|
||||
},
|
||||
);
|
||||
|
||||
new Worker(QUEUE, nightlyRefresh, { connection: redis, concurrency: 1 });
|
||||
console.log("[scheduler] nightly job armed (0 3 * * *)");
|
||||
}
|
||||
16
apps/worker/tsconfig.json
Normal file
16
apps/worker/tsconfig.json
Normal file
@@ -0,0 +1,16 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "esnext",
|
||||
"moduleResolution": "bundler",
|
||||
"esModuleInterop": true,
|
||||
"resolveJsonModule": true,
|
||||
"strict": true,
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true,
|
||||
"isolatedModules": true,
|
||||
"lib": ["ES2022"],
|
||||
"types": ["node"]
|
||||
},
|
||||
"include": ["src/**/*.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user