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
This commit is contained in:
@@ -40,6 +40,12 @@ export const EMBED_BASE_URL = process.env.EMBED_BASE_URL ?? "";
|
||||
export const EMBED_API_KEY = process.env.EMBED_API_KEY ?? "";
|
||||
export const EMBED_MODEL = process.env.EMBED_MODEL ?? "";
|
||||
|
||||
// Phase 3: Redis + BullMQ (shared imrnes Redis, no auth by default).
|
||||
export const REDIS_URL = process.env.REDIS_URL ?? "redis://100.121.180.82:6379";
|
||||
export const REDIS_PASSWORD = process.env.REDIS_PASSWORD ?? "";
|
||||
// BullMQ key prefix to namespace jobs on the shared Redis instance.
|
||||
export const QUEUE_PREFIX = process.env.QUEUE_PREFIX ?? "mcpedia";
|
||||
|
||||
if (!DATABASE_URL) {
|
||||
// Fail fast with an explicit message instead of a cryptic driver error.
|
||||
throw new Error(
|
||||
|
||||
@@ -0,0 +1,155 @@
|
||||
import { db } from "@mcpedia/db";
|
||||
import { documents, documentRevisions, documentChunks } from "@mcpedia/db/schema";
|
||||
import { parseFile } from "@mcpedia/parser";
|
||||
import { CONTENT_ROOT } from "@mcpedia/config";
|
||||
import { listContentFiles } from "./content.service";
|
||||
import { indexChunks } from "./document.service";
|
||||
import { toMeta } from "./row-map";
|
||||
import { eq, desc, and, sql } from "drizzle-orm";
|
||||
import { join } from "node:path";
|
||||
|
||||
export interface IndexResult {
|
||||
indexed: number;
|
||||
chunks: number;
|
||||
revisions: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Index a single content file: parse → upsert `documents` → chunk+embed →
|
||||
* snapshot a revision if the body changed since the last indexed revision.
|
||||
*
|
||||
* This is THE single indexing entry point shared by the CLI script, the
|
||||
* BullMQ worker, and the git-sync hook — no business logic is duplicated.
|
||||
*
|
||||
* @param relPath path relative to CONTENT_ROOT (e.g. "docs/websocket/contract")
|
||||
* @param reason provenance tag for the revision ("index" | "git-push" | "reindex")
|
||||
*/
|
||||
export async function indexContentFile(
|
||||
relPath: string,
|
||||
reason = "index",
|
||||
): Promise<{ indexed: boolean; chunks: number; revision: boolean }> {
|
||||
const abs = join(CONTENT_ROOT, relPath);
|
||||
const { meta, body } = parseFile(abs, relPath);
|
||||
const nowIso =
|
||||
meta.updatedAt && meta.updatedAt !== ""
|
||||
? meta.updatedAt
|
||||
: new Date().toISOString();
|
||||
|
||||
await db
|
||||
.insert(documents)
|
||||
.values({
|
||||
id: meta.id,
|
||||
slug: meta.slug,
|
||||
title: meta.title,
|
||||
type: meta.type,
|
||||
section: meta.section,
|
||||
status: meta.status,
|
||||
author: meta.author,
|
||||
tags: meta.tags,
|
||||
path: meta.path,
|
||||
body,
|
||||
createdAt: new Date(meta.createdAt || nowIso),
|
||||
updatedAt: new Date(nowIso),
|
||||
})
|
||||
.onConflictDoUpdate({
|
||||
target: documents.slug,
|
||||
set: {
|
||||
title: meta.title,
|
||||
type: meta.type,
|
||||
section: meta.section,
|
||||
status: meta.status,
|
||||
author: meta.author,
|
||||
tags: meta.tags,
|
||||
path: meta.path,
|
||||
body,
|
||||
updatedAt: new Date(nowIso),
|
||||
},
|
||||
});
|
||||
|
||||
// Semantic chunks (embedding). A failure here must not abort the whole
|
||||
// index — log and continue; FTS still works without embeddings.
|
||||
let chunks = 0;
|
||||
try {
|
||||
chunks = await indexChunks(meta.slug, body);
|
||||
} catch (err) {
|
||||
console.error(
|
||||
` embed FAILED for ${meta.slug}: ${err instanceof Error ? err.message : err}`,
|
||||
);
|
||||
}
|
||||
|
||||
// Snapshot a revision only when the body actually changed vs the latest
|
||||
// revision. Pure metadata/index changes (tags/title) won't create noise.
|
||||
const revision = await snapshotRevision(meta.slug, meta, body, reason);
|
||||
|
||||
return { indexed: true, chunks, revision };
|
||||
}
|
||||
|
||||
/**
|
||||
* Compare the incoming body against the latest revision's body; if different
|
||||
* (or no prior revision exists), create a new revision with an incremented
|
||||
* per-document revisionNo.
|
||||
*/
|
||||
async function snapshotRevision(
|
||||
slug: string,
|
||||
meta: ReturnType<typeof parseFile>["meta"],
|
||||
body: string,
|
||||
reason: string,
|
||||
): Promise<boolean> {
|
||||
const [doc] = await db
|
||||
.select({ id: documents.id })
|
||||
.from(documents)
|
||||
.where(eq(documents.slug, slug));
|
||||
if (!doc) return false;
|
||||
|
||||
const [latest] = await db
|
||||
.select({ body: documentRevisions.body, revisionNo: documentRevisions.revisionNo })
|
||||
.from(documentRevisions)
|
||||
.where(eq(documentRevisions.documentId, doc.id))
|
||||
.orderBy(desc(documentRevisions.revisionNo))
|
||||
.limit(1);
|
||||
|
||||
if (latest && latest.body === body) {
|
||||
return false; // unchanged → no new revision
|
||||
}
|
||||
|
||||
const nextNo = (latest?.revisionNo ?? 0) + 1;
|
||||
await db.insert(documentRevisions).values({
|
||||
documentId: doc.id,
|
||||
slug,
|
||||
revisionNo: nextNo,
|
||||
title: meta.title,
|
||||
body,
|
||||
meta: {
|
||||
type: meta.type,
|
||||
section: meta.section,
|
||||
status: meta.status,
|
||||
author: meta.author,
|
||||
tags: meta.tags,
|
||||
},
|
||||
reason,
|
||||
});
|
||||
return true;
|
||||
}
|
||||
|
||||
/**
|
||||
* Walk the entire content tree and index every file. Returns aggregate counts.
|
||||
*/
|
||||
export async function runFullIndex(reason = "index"): Promise<IndexResult> {
|
||||
const files = listContentFiles();
|
||||
let indexed = 0;
|
||||
let chunks = 0;
|
||||
let revisions = 0;
|
||||
for (const rel of files) {
|
||||
const r = await indexContentFile(rel, reason);
|
||||
indexed++;
|
||||
chunks += r.chunks;
|
||||
if (r.revision) revisions++;
|
||||
console.log(
|
||||
` indexed ${rel}${r.chunks ? ` (${r.chunks} chunks)` : ""}${r.revision ? " [revision]" : ""}`,
|
||||
);
|
||||
}
|
||||
console.log(
|
||||
`indexed ${indexed} documents, ${chunks} chunks, ${revisions} new revisions`,
|
||||
);
|
||||
return { indexed, chunks, revisions };
|
||||
}
|
||||
@@ -1,6 +1,8 @@
|
||||
export * from "./content.service";
|
||||
export * from "./document.service";
|
||||
export * from "./search.service";
|
||||
export * from "./index.service";
|
||||
export * from "./revision.service";
|
||||
export { toMeta } from "./row-map";
|
||||
|
||||
export type {
|
||||
|
||||
@@ -0,0 +1,116 @@
|
||||
import { db } from "@mcpedia/db";
|
||||
import { documents, documentRevisions, documentChunks } from "@mcpedia/db/schema";
|
||||
import { eq, desc, and, sql } from "drizzle-orm";
|
||||
import { toMeta } from "./row-map";
|
||||
import type { DocumentMeta } from "@mcpedia/types";
|
||||
|
||||
export interface RevisionSummary {
|
||||
id: string;
|
||||
slug: string;
|
||||
revisionNo: number;
|
||||
title: string;
|
||||
reason: string;
|
||||
createdAt: string;
|
||||
bodyLength: number;
|
||||
}
|
||||
|
||||
/** List revisions for a slug, newest first. */
|
||||
export async function listRevisions(
|
||||
slug: string,
|
||||
limit = 20,
|
||||
): Promise<RevisionSummary[]> {
|
||||
const [doc] = await db
|
||||
.select({ id: documents.id })
|
||||
.from(documents)
|
||||
.where(eq(documents.slug, slug));
|
||||
if (!doc) return [];
|
||||
|
||||
const rows = await db
|
||||
.select({
|
||||
id: documentRevisions.id,
|
||||
slug: documentRevisions.slug,
|
||||
revisionNo: documentRevisions.revisionNo,
|
||||
title: documentRevisions.title,
|
||||
reason: documentRevisions.reason,
|
||||
createdAt: documentRevisions.createdAt,
|
||||
bodyLength: sql<number>`length(${documentRevisions.body})`,
|
||||
})
|
||||
.from(documentRevisions)
|
||||
.where(eq(documentRevisions.documentId, doc.id))
|
||||
.orderBy(desc(documentRevisions.revisionNo))
|
||||
.limit(limit);
|
||||
|
||||
return rows.map((r) => ({
|
||||
id: r.id,
|
||||
slug: r.slug,
|
||||
revisionNo: r.revisionNo,
|
||||
title: r.title,
|
||||
reason: r.reason,
|
||||
createdAt: r.createdAt.toISOString(),
|
||||
bodyLength: r.bodyLength,
|
||||
}));
|
||||
}
|
||||
|
||||
/** Fetch a single revision's full body. */
|
||||
export async function getRevision(
|
||||
id: string,
|
||||
): Promise<{ id: string; revisionNo: number; body: string; meta: unknown } | null> {
|
||||
const [row] = await db
|
||||
.select({
|
||||
id: documentRevisions.id,
|
||||
revisionNo: documentRevisions.revisionNo,
|
||||
body: documentRevisions.body,
|
||||
meta: documentRevisions.meta,
|
||||
})
|
||||
.from(documentRevisions)
|
||||
.where(eq(documentRevisions.id, id));
|
||||
if (!row) return null;
|
||||
return {
|
||||
id: row.id,
|
||||
revisionNo: row.revisionNo,
|
||||
body: row.body,
|
||||
meta: row.meta,
|
||||
};
|
||||
}
|
||||
|
||||
/** Restore a revision: write its body+metadata back into the live `documents` row. */
|
||||
export async function restoreRevision(
|
||||
id: string,
|
||||
): Promise<{ slug: string; documentId: string } | null> {
|
||||
const [rev] = await db
|
||||
.select({
|
||||
id: documentRevisions.id,
|
||||
documentId: documentRevisions.documentId,
|
||||
slug: documentRevisions.slug,
|
||||
title: documentRevisions.title,
|
||||
body: documentRevisions.body,
|
||||
meta: documentRevisions.meta,
|
||||
})
|
||||
.from(documentRevisions)
|
||||
.where(eq(documentRevisions.id, id));
|
||||
if (!rev) return null;
|
||||
|
||||
const m = rev.meta as {
|
||||
type?: string;
|
||||
section?: string;
|
||||
status?: string;
|
||||
author?: string;
|
||||
tags?: string[];
|
||||
};
|
||||
|
||||
await db
|
||||
.update(documents)
|
||||
.set({
|
||||
title: rev.title,
|
||||
type: (m.type as any) ?? "documentation",
|
||||
section: (m.section as any) ?? "docs",
|
||||
status: (m.status as any) ?? "published",
|
||||
author: m.author ?? "",
|
||||
tags: m.tags ?? [],
|
||||
body: rev.body,
|
||||
updatedAt: new Date(),
|
||||
})
|
||||
.where(eq(documents.id, rev.documentId));
|
||||
|
||||
return { slug: rev.slug, documentId: rev.documentId };
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
CREATE TABLE "document_revisions" (
|
||||
"id" uuid PRIMARY KEY DEFAULT gen_random_uuid() NOT NULL,
|
||||
"document_id" text NOT NULL,
|
||||
"slug" text NOT NULL,
|
||||
"revision_no" integer NOT NULL,
|
||||
"title" text NOT NULL,
|
||||
"body" text NOT NULL,
|
||||
"meta" jsonb NOT NULL,
|
||||
"reason" text DEFAULT 'index' NOT NULL,
|
||||
"created_at" timestamp with time zone DEFAULT now() NOT NULL
|
||||
);
|
||||
--> statement-breakpoint
|
||||
CREATE INDEX "document_revisions_document_id_idx" ON "document_revisions" USING btree ("document_id");
|
||||
--> statement-breakpoint
|
||||
CREATE INDEX "document_revisions_slug_idx" ON "document_revisions" USING btree ("slug");
|
||||
--> statement-breakpoint
|
||||
CREATE INDEX "document_revisions_doc_rev_idx" ON "document_revisions" USING btree ("document_id", "revision_no" DESC);
|
||||
--> statement-breakpoint
|
||||
ALTER TABLE "document_revisions" ADD CONSTRAINT "document_revisions_document_id_documents_id_fk" FOREIGN KEY ("document_id") REFERENCES "public"."documents"("id") ON DELETE cascade;
|
||||
@@ -0,0 +1,110 @@
|
||||
{
|
||||
"id": "0002_document_revisions",
|
||||
"prevId": "0001_document_chunks",
|
||||
"version": "7",
|
||||
"dialect": "postgresql",
|
||||
"tables": {
|
||||
"document_revisions": {
|
||||
"name": "document_revisions",
|
||||
"columns": {
|
||||
"id": {
|
||||
"name": "id",
|
||||
"type": "uuid",
|
||||
"primaryKey": true,
|
||||
"notNull": true,
|
||||
"default": "gen_random_uuid()"
|
||||
},
|
||||
"document_id": {
|
||||
"name": "document_id",
|
||||
"type": "text",
|
||||
"notNull": true
|
||||
},
|
||||
"slug": {
|
||||
"name": "slug",
|
||||
"type": "text",
|
||||
"notNull": true
|
||||
},
|
||||
"revision_no": {
|
||||
"name": "revision_no",
|
||||
"type": "integer",
|
||||
"notNull": true
|
||||
},
|
||||
"title": {
|
||||
"name": "title",
|
||||
"type": "text",
|
||||
"notNull": true
|
||||
},
|
||||
"body": {
|
||||
"name": "body",
|
||||
"type": "text",
|
||||
"notNull": true
|
||||
},
|
||||
"meta": {
|
||||
"name": "meta",
|
||||
"type": "jsonb",
|
||||
"notNull": true
|
||||
},
|
||||
"reason": {
|
||||
"name": "reason",
|
||||
"type": "text",
|
||||
"notNull": true,
|
||||
"default": "'index'"
|
||||
},
|
||||
"created_at": {
|
||||
"name": "created_at",
|
||||
"type": "timestamp",
|
||||
"notNull": true,
|
||||
"default": "now()"
|
||||
}
|
||||
},
|
||||
"indexes": {
|
||||
"document_revisions_document_id_idx": {
|
||||
"name": "document_revisions_document_id_idx",
|
||||
"columns": [
|
||||
{ "name": "document_id", "asc": true }
|
||||
],
|
||||
"isUnique": false
|
||||
},
|
||||
"document_revisions_slug_idx": {
|
||||
"name": "document_revisions_slug_idx",
|
||||
"columns": [
|
||||
{ "name": "slug", "asc": true }
|
||||
],
|
||||
"isUnique": false
|
||||
},
|
||||
"document_revisions_doc_rev_idx": {
|
||||
"name": "document_revisions_doc_rev_idx",
|
||||
"columns": [
|
||||
{ "name": "document_id", "asc": true },
|
||||
{ "name": "revision_no", "asc": false }
|
||||
],
|
||||
"isUnique": false
|
||||
}
|
||||
},
|
||||
"foreignKeys": {
|
||||
"document_revisions_document_id_documents_id_fk": {
|
||||
"name": "document_revisions_document_id_documents_id_fk",
|
||||
"columns": ["document_id"],
|
||||
"referenceTable": "documents",
|
||||
"referenceColumns": ["id"],
|
||||
"onDelete": "cascade"
|
||||
}
|
||||
},
|
||||
"compositePrimaryKeys": {},
|
||||
"uniqueConstraints": {},
|
||||
"policies": {}
|
||||
}
|
||||
},
|
||||
"enums": {},
|
||||
"schemas": {},
|
||||
"sequences": {},
|
||||
"roles": {},
|
||||
"policies": {},
|
||||
"views": {},
|
||||
"extensions": {},
|
||||
"_meta": {
|
||||
"columns": {},
|
||||
"schemas": {},
|
||||
"tables": {}
|
||||
}
|
||||
}
|
||||
@@ -15,6 +15,13 @@
|
||||
"when": 1787137149735,
|
||||
"tag": "0001_document_chunks",
|
||||
"breakpoints": true
|
||||
},
|
||||
{
|
||||
"idx": 2,
|
||||
"version": "7",
|
||||
"when": 1787139150000,
|
||||
"tag": "0002_document_revisions",
|
||||
"breakpoints": true
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ import {
|
||||
customType,
|
||||
index,
|
||||
integer,
|
||||
jsonb,
|
||||
pgTable,
|
||||
real,
|
||||
text,
|
||||
@@ -77,5 +78,48 @@ export const documentChunks = pgTable(
|
||||
export type DocumentChunkRow = typeof documentChunks.$inferSelect;
|
||||
export type NewDocumentChunkRow = typeof documentChunks.$inferInsert;
|
||||
|
||||
// Phase 3: document revision system. Each row is an immutable snapshot of a
|
||||
// document's body + metadata at a point in time (taken by the indexer whenever
|
||||
// the body actually changes). revisionNo is per-document and monotonically
|
||||
// increasing so the latest revision is always max(revision_no).
|
||||
export const documentRevisions = pgTable(
|
||||
"document_revisions",
|
||||
{
|
||||
id: uuid("id").primaryKey().defaultRandom(),
|
||||
documentId: text("document_id")
|
||||
.notNull()
|
||||
.references(() => documents.id, { onDelete: "cascade" }),
|
||||
slug: text("slug").notNull(),
|
||||
revisionNo: integer("revision_no").notNull(),
|
||||
title: text("title").notNull(),
|
||||
body: text("body").notNull(),
|
||||
// Metadata snapshot (type/section/status/author/tags) as JSON so a revision
|
||||
// is self-describing even if the live document is later restructured.
|
||||
meta: jsonb("meta").notNull().$type<{
|
||||
type: string;
|
||||
section: string;
|
||||
status: string;
|
||||
author: string;
|
||||
tags: string[];
|
||||
}>(),
|
||||
// Why this revision was created (e.g. "index", "git-push", "restore").
|
||||
reason: text("reason").notNull().default("index"),
|
||||
createdAt: timestamp("created_at", { withTimezone: true })
|
||||
.notNull()
|
||||
.defaultNow(),
|
||||
},
|
||||
(t) => ({
|
||||
docIdx: index("document_revisions_document_id_idx").on(t.documentId),
|
||||
slugIdx: index("document_revisions_slug_idx").on(t.slug),
|
||||
docRevIdx: index("document_revisions_doc_rev_idx").on(
|
||||
t.documentId,
|
||||
sql`${t.revisionNo} desc`,
|
||||
),
|
||||
}),
|
||||
);
|
||||
|
||||
export type DocumentRevisionRow = typeof documentRevisions.$inferSelect;
|
||||
export type NewDocumentRevisionRow = typeof documentRevisions.$inferInsert;
|
||||
|
||||
export type DocumentRow = typeof documents.$inferSelect;
|
||||
export type NewDocumentRow = typeof documents.$inferInsert;
|
||||
|
||||
@@ -0,0 +1,22 @@
|
||||
{
|
||||
"name": "@mcpedia/queue",
|
||||
"version": "0.1.0",
|
||||
"private": true,
|
||||
"type": "module",
|
||||
"exports": {
|
||||
".": "./src/index.ts",
|
||||
"./client": "./src/client.ts",
|
||||
"./queue": "./src/queue.ts",
|
||||
"./worker": "./src/worker.ts"
|
||||
},
|
||||
"dependencies": {
|
||||
"@mcpedia/config": "workspace:*",
|
||||
"@mcpedia/core": "workspace:*",
|
||||
"@mcpedia/db": "workspace:*",
|
||||
"bullmq": "^6.1.2",
|
||||
"ioredis": "^6.0.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
import { REDIS_URL, REDIS_PASSWORD, QUEUE_PREFIX } from "@mcpedia/config";
|
||||
import IORedis, { type RedisOptions } from "ioredis";
|
||||
|
||||
/**
|
||||
* Shared ioredis connection for BullMQ. BullMQ requires an ioredis instance and
|
||||
* internally duplicates it for blocking commands, so we keep the option objects
|
||||
* explicit (maxRetriesPerRequest: null is REQUIRED for the blocking
|
||||
* connection — a finite retry count causes "Connection in key mode" errors).
|
||||
*/
|
||||
function buildOptions(): RedisOptions {
|
||||
const opts: RedisOptions = {
|
||||
maxRetriesPerRequest: null,
|
||||
lazyConnect: true,
|
||||
enableOfflineQueue: true,
|
||||
};
|
||||
if (REDIS_PASSWORD) opts.password = REDIS_PASSWORD;
|
||||
return opts;
|
||||
}
|
||||
|
||||
let _connection: IORedis | null = null;
|
||||
|
||||
/** Lazily-created singleton ioredis connection. */
|
||||
export function getConnection(): IORedis {
|
||||
if (!_connection) {
|
||||
_connection = new IORedis(REDIS_URL, buildOptions());
|
||||
_connection.on("error", (err) => {
|
||||
// Log but don't crash the process on transient Redis errors.
|
||||
console.error("[queue] redis error:", err.message);
|
||||
});
|
||||
}
|
||||
return _connection;
|
||||
}
|
||||
|
||||
export const BULLMQ_PREFIX = QUEUE_PREFIX;
|
||||
@@ -0,0 +1,9 @@
|
||||
export { getConnection, BULLMQ_PREFIX } from "./client";
|
||||
export {
|
||||
getQueue,
|
||||
enqueueIndexDoc,
|
||||
enqueueFullIndex,
|
||||
INDEX_QUEUE,
|
||||
} from "./queue";
|
||||
export type { IndexDocJobData, IndexAllJobData, JobType } from "./queue";
|
||||
export { createWorker, startWorker } from "./worker";
|
||||
@@ -0,0 +1,66 @@
|
||||
import { Queue, type Job } from "bullmq";
|
||||
import { getConnection, BULLMQ_PREFIX } from "./client";
|
||||
|
||||
export const INDEX_QUEUE = "mcpedia-index";
|
||||
|
||||
/** Lazily-created singleton BullMQ queue. */
|
||||
let _queue: Queue | null = null;
|
||||
|
||||
export function getQueue(): Queue {
|
||||
if (!_queue) {
|
||||
_queue = new Queue(INDEX_QUEUE, {
|
||||
connection: getConnection(),
|
||||
prefix: BULLMQ_PREFIX,
|
||||
});
|
||||
}
|
||||
return _queue;
|
||||
}
|
||||
|
||||
export interface IndexDocJobData {
|
||||
relPath: string;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
export interface IndexAllJobData {
|
||||
reason: string;
|
||||
}
|
||||
|
||||
export type JobType = "index-doc" | "index-all";
|
||||
|
||||
/**
|
||||
* Enqueue a single-document reindex job. Keyed by slug so repeated edits
|
||||
* collapse into one pending job (BullMQ dedup by jobId within the window).
|
||||
*/
|
||||
export async function enqueueIndexDoc(
|
||||
relPath: string,
|
||||
reason = "index",
|
||||
): Promise<Job<IndexDocJobData>> {
|
||||
const slug = relPath.replace(/\.mdx?$/, "");
|
||||
return getQueue().add(
|
||||
"index-doc",
|
||||
{ relPath, reason },
|
||||
{
|
||||
jobId: `doc__${slug}`,
|
||||
removeOnComplete: 1000,
|
||||
removeOnFail: 5000,
|
||||
attempts: 3,
|
||||
backoff: { type: "exponential", delay: 2000 },
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
/** Enqueue a full-corpus reindex (used by the git-sync hook). */
|
||||
export async function enqueueFullIndex(
|
||||
reason = "reindex",
|
||||
): Promise<Job<IndexAllJobData>> {
|
||||
return getQueue().add(
|
||||
"index-all",
|
||||
{ reason },
|
||||
{
|
||||
jobId: `full__${Date.now()}`,
|
||||
removeOnComplete: 100,
|
||||
removeOnFail: 1000,
|
||||
attempts: 1,
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,63 @@
|
||||
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;
|
||||
}
|
||||
Reference in New Issue
Block a user