Files
plp/apps/worker/src/main.ts
T
root fa0fa78312 fix(治理): 完成任务 8 安全与通知闭环修复
- 串行化拉黑、处罚、投瓶和匹配策略检查\n- 完成异步站内通知、未读统计和偏好并发语义\n- 补齐后台查询审计、处罚恢复和隐私测试\n- 稳定 Redis 恢复、匹配锁序及超时测试
2026-09-17 13:08:55 +08:00

114 lines
3.0 KiB
TypeScript

import { pathToFileURL } from "node:url";
import { PrismaClient } from "@prisma/client";
import { ModerationWorker } from "./moderation-worker.js";
import { LeaseReaper } from "./lease-reaper.processor.js";
import { NotificationProcessor } from "./notification.processor.js";
type SignalSource = {
once(signal: "SIGTERM" | "SIGINT", listener: () => void): unknown;
removeListener(signal: "SIGTERM" | "SIGINT", listener: () => void): unknown;
};
type RunWorkerOptions = {
worker: { runOnce(): Promise<boolean> };
connect: () => Promise<unknown>;
disconnect: () => Promise<unknown>;
once: boolean;
sleep?: () => Promise<unknown>;
signals?: SignalSource;
};
type Worker = { runOnce(): Promise<boolean> };
export function createCombinedWorker(
reaper: Worker,
moderation: Worker,
notification?: Worker,
): Worker {
return {
async runOnce() {
const reaped = await reaper.runOnce();
const moderated = await moderation.runOnce();
const notified = notification ? await notification.runOnce() : false;
return reaped || moderated || notified;
},
};
}
export async function runWorker(options: RunWorkerOptions) {
const signals = options.signals ?? process;
const sleep =
options.sleep ??
(() => new Promise((resolve) => setTimeout(resolve, 1000)));
let stopping = false;
const stop = () => {
stopping = true;
};
if (!options.once) {
signals.once("SIGTERM", stop);
signals.once("SIGINT", stop);
}
try {
await options.connect();
do {
const handled = await options.worker.runOnce();
if (options.once || stopping) break;
if (!handled) await sleep();
} while (!stopping);
} finally {
if (!options.once) {
signals.removeListener("SIGTERM", stop);
signals.removeListener("SIGINT", stop);
}
await options.disconnect();
}
}
export async function main() {
const prisma = new PrismaClient();
const moderation = new ModerationWorker(prisma);
const reaper = new LeaseReaper(prisma);
const notification = new NotificationProcessor(prisma);
const worker = createCombinedWorker(reaper, moderation, notification);
await runWorker({
worker,
connect: () => prisma.$connect(),
disconnect: () => prisma.$disconnect(),
once: process.argv.includes("--once"),
});
}
export function logTopLevelError(
error: unknown,
logger: Pick<Console, "error"> = console,
) {
logger.error(
JSON.stringify({
scope: "moderation-worker",
errorClass:
error instanceof Error && error.name ? error.name : "UnknownError",
}),
);
}
type TopLevelErrorOptions = {
logger?: Pick<Console, "error">;
processState?: { exitCode?: number };
};
export function handleTopLevelError(
error: unknown,
options: TopLevelErrorOptions = {},
) {
logTopLevelError(error, options.logger);
(options.processState ?? process).exitCode = 1;
}
const isEntrypoint =
process.argv[1] !== undefined &&
import.meta.url === pathToFileURL(process.argv[1]).href;
if (isEntrypoint) {
main().catch(handleTopLevelError);
}