Files
mcpedia/packages/queue/src/worker.ts
T
asepharyana 8f2229d447 feat(mcpedia): Phase 3 — async indexing (BullMQ), git-sync webhook, revisions, MCP Resources
- packages/queue: ioredis singleton + BullMQ Queue/Worker (prefix mcpedia:
  on shared imrnes Redis :6379); apps/worker runs startWorker()
- @mcpedia/core: indexContentFile/runFullIndex (single indexing entry point
  shared by script/worker/hook) + revision.service (list/get/restore)
- document_revisions table (migration 0002) — snapshots only on body change
- apps/api: POST /hooks/reindex + /hooks/index webhooks; tRPC revisions,
  getRevision, restoreRevision, jobStatus, queueStatus
- apps/mcp: register MCP Resources mcpedia://docs{/,+slug/chunks/revisions}
  ({+slug} RFC6570 reserved expansion for slugs containing /)
- apps/mcp zod pinned to ^4 to match MCP SDK 1.30 compiled types
  (resolves registerTool TS2589/ShapeOutput skew)
- scripts/enqueue.ts one-shot job enqueue helper; indexer refactored to runFullIndex
- PHASES.md/README/.env.example/docs updated
2026-08-19 20:18:28 +07:00

64 lines
1.8 KiB
TypeScript

import { Worker, type Job } from "bullmq";
import { getConnection, BULLMQ_PREFIX } from "./client";
import { INDEX_QUEUE } from "./queue";
import { indexContentFile, runFullIndex } from "@mcpedia/core";
export function createWorker(): Worker {
const worker = new Worker(
INDEX_QUEUE,
async (job: Job) => {
switch (job.name) {
case "index-doc": {
const { relPath, reason } = job.data as {
relPath: string;
reason: string;
};
await job.log(`indexing ${relPath}`);
const r = await indexContentFile(relPath, reason);
return r;
}
case "index-all": {
const { reason } = job.data as { reason: string };
await job.log(`full index (${reason})`);
return await runFullIndex(reason);
}
default:
throw new Error(`unknown job type: ${job.name}`);
}
},
{
connection: getConnection(),
prefix: BULLMQ_PREFIX,
concurrency: 4,
},
);
worker.on("completed", (job) => {
console.log(`[worker] completed ${job.name} (${job.id})`);
});
worker.on("failed", (job, err) => {
console.error(`[worker] failed ${job?.name} (${job?.id}): ${err.message}`);
});
worker.on("error", (err) => {
console.error(`[worker] error:`, err.message);
});
return worker;
}
/** Start the worker and wire graceful shutdown. */
export async function startWorker(): Promise<Worker> {
const worker = createWorker();
console.log("[worker] indexing worker started");
const shutdown = async (sig: string) => {
console.log(`[worker] ${sig} received, closing...`);
await worker.close();
process.exit(0);
};
process.on("SIGINT", () => void shutdown("SIGINT"));
process.on("SIGTERM", () => void shutdown("SIGTERM"));
return worker;
}