-
Notifications
You must be signed in to change notification settings - Fork 15
Expand file tree
/
Copy pathworker.ts
More file actions
102 lines (92 loc) · 5.21 KB
/
Copy pathworker.ts
File metadata and controls
102 lines (92 loc) · 5.21 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
// 生产任务 worker 进程入口:注册三轨 cron + 起 BullMQ 消费者。
// 运行:BEACON_QUEUE=bullmq REDIS_URL=redis://... npx tsx worker.ts
// (dev 不需要此进程——进程内队列在按钮点击时同步执行)
import { getQueue, BullMqQueue, MAIN_QUEUE_NAME, AGENT_QUEUE_NAME } from './lib/jobs/queue';
import { SCHEDULES, SCHEDULE_TZ } from './lib/jobs/schedule-config';
import { initObservability, log, flushReports } from './lib/logger';
async function main() {
initObservability('worker');
// 没配 Sentry 时把错误上报接到运维告警出口(显式 webhook,或用户已配的机器人集成)。
// worker 尤其需要:它没有页面,任务失败以前只在容器日志里,除非有人主动去看,否则永远不知道。
await (await import('./lib/ops/install')).installOpsAlerting();
if (process.env.BEACON_QUEUE !== 'bullmq') {
log.error(
'worker 需要 BEACON_QUEUE=bullmq(SaaS / 私有化)。'
+ '整机版用 BEACON_QUEUE=local——定时跑在 web 进程里(单文件 SQLite 不能有第二个写进程),不要起本进程。'
+ '本机开发不设这个变量,不跑定时。',
);
await flushReports();
process.exit(1);
}
const q = getQueue() as BullMqQueue;
// 先清后注册:旧的 repeat key(改了 cron 或时区之前留下的那条)不会被覆盖,
// 不清就会新老两条并存,各按各的时间跑。见 resetSchedules 的注释。
const removed = await q.resetSchedules();
if (removed) log.info('已清理旧定时任务', { removed });
// 注册三轨定时任务(时间口径为北京时间,见 SCHEDULE_TZ)
for (const s of SCHEDULES) {
await q.schedule(s.name, s.cron);
log.info('已注册定时任务', { jobName: s.name, cron: s.cron, tz: SCHEDULE_TZ, note: s.note });
}
// 上面「先清后注册」会把停机期间到点却没跑的那次一起清掉。补回来(只补宽限窗内、周期 ≥6h 的,各一次)。
// 放在注册之后:补跑本身也会经 withRun 记 JobRun,下一次启动就不会再补同一次。
const { runCatchUps } = await import('./lib/jobs/catch-up');
const caught = await runCatchUps(q);
if (caught.length) log.info('已补跑停机期间错过的定时任务', { jobs: caught });
// 起消费者。**两条队列各一个**:
// AI 执行的 loop 一跑就是几十分钟,与它共用槽位的话,热榜采集、定时智能体、
// 清理、发布回流会全部排在几个用户的 AI 任务后面——一个人用助手,全平台定时停摆。
const mainConcurrency = Number(process.env.BEACON_WORKER_CONCURRENCY || 4);
const agentConcurrency = Number(process.env.BEACON_AGENT_CONCURRENCY || 4);
const workers = [
{ name: MAIN_QUEUE_NAME, w: await q.startWorker(MAIN_QUEUE_NAME, mainConcurrency), concurrency: mainConcurrency },
{ name: AGENT_QUEUE_NAME, w: await q.startWorker(AGENT_QUEUE_NAME, agentConcurrency), concurrency: agentConcurrency },
];
for (const { name, w } of workers) {
w.on('completed', (job) => {
log.info('任务完成', { queue: name, jobName: job.name, jobId: job.id, durationMs: job.finishedOn ? job.finishedOn - job.processedOn! : undefined });
});
w.on('failed', (job, err) => {
// 任务失败是要被看见的:error 级 + 自动上报
log.error('任务失败', { queue: name, jobName: job?.name, jobId: job?.id, attempts: job?.attemptsMade, err });
});
w.on('error', (err) => {
log.error('worker 内部错误', { queue: name, err });
});
}
log.info('beacon worker 已启动', {
queues: workers.map((x) => `${x.name}(并发 ${x.concurrency})`).join(' + '),
});
// 微信 iLink 收信:长轮询常驻循环,只在这一个进程里跑(游标是消费性的,两处同时拉会互吞消息;
// 整机版没有本进程,由 instrumentation.node.ts 在 web 进程里起——两处互斥,判据都是 BEACON_QUEUE)
const { startIlinkSupervisor, stopIlinkSupervisor } = await import('./lib/bot/wechat-ilink-poller');
startIlinkSupervisor();
// 企微智能机器人长连接:同理只在这一个进程里连(协议每个机器人只许一条活连接,两处连会互踢)
const { startAibotSupervisor, stopAibotSupervisor } = await import('./lib/bot/wecom-aibot-poller');
startAibotSupervisor();
const shutdown = async (signal: string) => {
log.info('收到退出信号,关闭 worker…', { signal });
stopIlinkSupervisor();
stopAibotSupervisor();
await Promise.all(workers.map(({ w }) => w.close()));
await q.close();
await flushReports();
process.exit(0);
};
process.on('SIGINT', () => void shutdown('SIGINT'));
process.on('SIGTERM', () => void shutdown('SIGTERM'));
// 兜底:未捕获异常必须留痕再退出,让编排器重启(不要吞掉继续裸跑)
process.on('uncaughtException', async (err) => {
log.error('未捕获异常,进程退出', { err });
await flushReports();
process.exit(1);
});
process.on('unhandledRejection', (reason) => {
log.error('未处理的 Promise 拒绝', { err: reason });
});
}
main().catch(async (e) => {
log.error('worker 启动失败', { err: e, phase: 'startup' });
await flushReports();
process.exit(1);
});