diff --git a/.hermes/plans/phase3-plan.md b/.hermes/plans/phase3-plan.md new file mode 100644 index 0000000..2cc5666 --- /dev/null +++ b/.hermes/plans/phase3-plan.md @@ -0,0 +1,114 @@ +# MCPedia — Phase 3 "Async + Scale" Implementation Plan + +Status: Phase 1 (MVP) + Phase 2 (Semantic+API) DONE. Phase 3 adds async +background work, git-driven reindex, document revision history, and MCP +Resources. All logic stays in `@mcpedia/core`; new `packages/queue` wires +BullMQ; `apps/worker` runs the worker process; the existing API gets a git-sync +webhook + job-status procedures; the MCP server gains Resources. + +## Scope (4 features from PHASES.md) + +1. **Redis + BullMQ background indexing/embedding workers** +2. **Git synchronization hook** (auto-reindex on push via webhook) +3. **Document revision system** (`document_revisions`) +4. **MCP Resources** (`mcpedia://docs/...`) alongside existing tools + +## Architecture decisions (locked) + +- **Redis**: shared imrnes Redis `100.121.180.82:6379`, no auth (verified + `+PONG`). `REDIS_URL` env (default `redis://100.121.180.82:6379`), optional + `REDIS_PASSWORD`. BullMQ key prefix `mcpedia:` to avoid collisions on the + shared instance. +- **Queue lib**: `bullmq@6.1.2` + `ioredis@6.0.0` (BullMQ peer dep). Pass an + ioredis instance; BullMQ duplicates it for blocking commands. +- **Single source of truth preserved**: per-doc indexing logic moves into + `@mcpedia/core` as `indexContentFile(relPath, reason?)`. The script, the + worker, and the git hook ALL call this. Revisions are snapshotted inside it. +- **Revisions**: created only when body actually changes vs the latest revision + (avoids bloat on every sync). Stored in `document_revisions`. + +## Files touched + +### packages/config +- `src/index.ts`: add `REDIS_URL`, `REDIS_PASSWORD`, `QUEUE_PREFIX`. + +### packages/db +- `src/schema.ts`: add `documentRevisions` table + (id, documentId→documents.id cascade, slug, revisionNo int, title, body, + meta jsonb, reason text, createdAt). Index (document_id, revision_no DESC), + (slug). +- `drizzle/0002_document_revisions.sql`: migration (applied via psql). +- `drizzle/meta/0002_snapshot.json` + `_journal.json` entry (keeps drizzle-kit + consistent even though we apply manually). + +### packages/core (new) +- `src/index.service.ts`: + - `indexContentFile(relPath: string, reason = "index")` — parse → upsert + `documents` → `indexChunks` → snapshot revision (if changed). + - `runFullIndex(reason?)` — walk content, index each, return counts. +- `src/revision.service.ts`: + - `createRevision(...)`, `listRevisions(slug, limit)`, + `getRevision(id)`, `latestRevisionBody(slug)`, `restoreRevision(id)`. +- `src/index.ts`: export both. + +### packages/queue (NEW) +- `package.json` (@mcpedia/queue): deps bullmq, ioredis, @mcpedia/core, + @mcpedia/db, @mcpedia/config. +- `src/client.ts`: ioredis instance factory from config. +- `src/queue.ts`: + - `INDEX_QUEUE = "mcpedia-index"`. + - `enqueueIndexDoc(slug, absPath, reason)`, `enqueueFullIndex(reason)`. + - `getQueue()` lazy singleton. +- `src/worker.ts`: `startWorker()` — BullMQ Worker with 3 job types: + `index-doc` (single), `index-all` (full), `reindex` (full, reason=git-push). + Graceful shutdown on SIGINT/SIGTERM. Job progress + error handling. + +### apps/worker (NEW) +- `package.json` (@mcpedia/worker): script `start: bun src/index.ts`. +- `src/index.ts`: `startWorker()` + heartbeat log. + +### apps/api +- `src/index.ts`: add `POST /hooks/reindex` (full) and + `POST /hooks/index?slug=` (single) webhook routes → enqueue jobs. Mount + AFTER /trpc. +- `src/router.ts`: add `jobStatus` (id→state/prev/failedData), + `queueStatus` (waiting/active/completed/failed counts), + `revisions` (slug→list), `restoreRevision` (id→new slug/doc). +- `package.json`: add `@mcpedia/queue` dep, `hooks` reused. + +### apps/mcp +- `src/index.ts`: register Resources: + - `mcpedia://docs` (list all metas) + - `mcpedia://docs/{slug}` (full body from disk) + - `mcpedia://docs/{slug}/chunks` (chunk previews) + - `mcpedia://docs/{slug}/revisions` (revision list) +- `src/smoke.test.ts`: add `listResources` + read `mcpedia://docs` assertion. + +### scripts +- `scripts/indexer.ts`: refactor `main()` to call `runFullIndex()`. + +### Root +- `package.json`: add `"worker": "bun --cwd apps/worker run start"`, + `"reindex": "bun run scripts/worker.ts"`? No — `worker` runs the listener; + triggering reindex = `bun run api` webhook or `enqueueFullIndex` helper. + Add `"enqueue-index": "bun run scripts/enqueue.ts"` (one-shot enqueue). +- `.env.example`: add `REDIS_URL`, `REDIS_PASSWORD`, `QUEUE_PREFIX`. + +### Docs +- `PHASES.md`: mark Phase 3 items DONE with notes. +- `README.md`: document worker, webhook, revisions, MCP resources. + +## Verification (real, not claimed) + +1. `bun install` picks up new deps. +2. `bunx turbo run build` + `typecheck` green across workspace. +3. **Real BullMQ e2e against imrnes Redis**: script that enqueues an + `index-doc` job, starts a Worker, asserts the job completes and the doc row + + chunks + a revision row appear in Postgres. Verifies Redis+ioredis+bullmq + + db + core all wired correctly. +4. `bun --cwd apps/mcp run smoke` passes (incl. new resources). +5. `bun run index` (runFullIndex) green; verify `documents`, + `document_chunks`, `document_revisions` row counts via psql. +6. API webhook: `curl -XPOST localhost:4020/hooks/reindex` enqueues; worker + processes; `curl localhost:4020/trpc/queueStatus` reflects counts. +7. MCP resource read returns real content. diff --git a/PHASES.md b/PHASES.md index 24a0406..d1905e1 100644 --- a/PHASES.md +++ b/PHASES.md @@ -27,12 +27,50 @@ Legend: ✅ built · 🟡 partial · ⬜ deferred - [x] MCP server — added `semantic_search` + `hybrid_search` tools (6 total). - [x] Web search — keyword/hybrid toggle (`?mode=hybrid`), hybrid reaches semantically-related docs keyword misses. -## Phase 3 — Async + Scale +## Phase 3 — Async + Scale ✅ DONE -- [ ] Redis + BullMQ background indexing / embedding workers -- [ ] Git synchronization hook (auto-reindex on push) -- [ ] Document revision system (`document_revisions`) -- [ ] MCP Resources (`mcpedia://docs/...`) in addition to tools +- [x] **Redis + BullMQ background indexing / embedding workers** — + `packages/queue` (ioredis singleton + BullMQ `Queue`/`Worker`, prefix + `mcpedia:` on shared imrnes Redis `:6379`); `apps/worker` runs + `startWorker()`. Three job types: `index-doc`, `index-all`, `reindex`. + Single indexing entry point `indexContentFile`/`runFullIndex` in + `@mcpedia/core` shared by the script, worker, and git hook. Verified + end-to-end against live Redis (job enqueue → worker → Postgres write). +- [x] **Git synchronization hook (auto-reindex on push)** — API webhook + `POST /hooks/reindex` (full) and `POST /hooks/index?slug=` (single) enqueue + BullMQ jobs. Wire a Git provider (GitHub/Gitea) post-receive / webhook to + `POST /hooks/reindex` to auto-reindex on push. `scripts/enqueue.ts` is a + one-shot enqueue helper (`bun run enqueue --all` / ``). +- [x] **Document revision system (`document_revisions`)** — `packages/db` + migration `0002_document_revisions.sql`. Indexer snapshots a revision only + when the body actually changes vs the latest revision (pure metadata edits + don't bloat history). `listRevisions` / `getRevision` / `restoreRevision` + in `@mcpedia/core`; exposed as tRPC `revisions` / `getRevision` / + `restoreRevision` and the `mcpedia://docs/{+slug}/revisions` MCP Resource. +- [x] **MCP Resources (`mcpedia://docs/...`)** — alongside the 6 tools: + `mcpedia://docs` (list), `mcpedia://docs/{+slug}` (body from disk), + `mcpedia://docs/{+slug}/chunks` (chunk preview), + `mcpedia://docs/{+slug}/revisions` (history). `{+slug}` uses RFC 6570 + reserved expansion so slugs containing `/` match. + +### New/changed commands +``` +bun run index # full reindex (runFullIndex, writes revisions) +bun run enqueue --all # enqueue a full reindex job (no worker needed) +bun run enqueue # enqueue a single-doc reindex job +bun run worker # start the BullMQ indexing worker (long-running) +bun run api # Hono+tRPC API on :4020 (added /hooks/* webhooks) +``` + +### Verification done (real, against imrnes Redis + Postgres) +- `turbo run typecheck` green across all 13 packages. +- BullMQ e2e: enqueue `index-doc` → worker completes → `documents` + + `document_chunks` + `document_revisions` rows present. +- Revision dedup proven: editing a body creates a new revision; metadata-only + reindex does not; `restoreRevision` writes history back into the live row. +- MCP smoke test passes (tools + all 4 resources). +- API webhook `POST /hooks/reindex` enqueues → worker drains queue → + `queueStatus` reflects counts. ## Phase 4 — Scale-out (only if needed) diff --git a/README.md b/README.md index 3735ae8..e728cf6 100644 --- a/README.md +++ b/README.md @@ -14,16 +14,19 @@ column) and served through a single **Core** layer that every interface mcpedia/ ├── apps/ │ ├── web/ # Next.js 16 (Turbopack) — human-facing docs UI + search -│ └── mcp/ # MCP server (stdio) — AI-agent interface +│ ├── mcp/ # MCP server (stdio) — AI-agent interface (tools + resources) +│ └── api/ # Hono + tRPC v11 API on :4020 (+ /hooks/* git-sync webhooks) ├── packages/ │ ├── types/ # shared domain types (DocSection, Document, SearchHit, ...) │ ├── config/ # loads .env (repo root) as authoritative dev config │ ├── db/ # Drizzle ORM schema + client + drizzle-kit config │ ├── parser/ # frontmatter (gray-matter) parsing │ ├── search/ # Postgres FTS query (ts_rank + ts_headline) -│ └── core/ # Document/Content/Search services — the only business logic +│ ├── embeddings/ # embedding provider + chunker +│ ├── queue/ # Redis (ioredis) + BullMQ worker/queue (Phase 3) +│ └── core/ # Document/Content/Search/Index/Revision — the only business logic ├── content/ # docs/ writeups/ research/ notes/ (the knowledge base) -└── scripts/ # indexer.ts (walks content/ -> upserts into Postgres) +└── scripts/ # indexer.ts (full reindex), enqueue.ts (one-shot job enqueue) ``` ## Architecture principle @@ -101,23 +104,45 @@ the DB stores metadata + the search vector. | `list_documents` | List, optionally filtered by section | | `get_related_documents` | Docs sharing tags with a given slug | +### MCP Resources + +| URI | Purpose | +| -------------------------------- | ---------------------------------------- | +| `mcpedia://docs` | List all published documents | +| `mcpedia://docs/{+slug}` | Full markdown body (read from disk) | +| `mcpedia://docs/{+slug}/chunks` | Preview of embedded semantic chunks | +| `mcpedia://docs/{+slug}/revisions` | Revision history summary | + +(`{+slug}` uses RFC 6570 reserved expansion so a slug like +`docs/websocket/contract` matches the template.) + Smoke test (in-memory transport, real JSON-RPC): ```bash bun --cwd apps/mcp run smoke ``` -## API (Phase 2) +## API (Phase 2 + Phase 3) -A tRPC v11 API is also exposed via Hono on **:4020** (all procedures mirror the -MCP tools): +A tRPC v11 API is exposed via Hono on **:4020** (all procedures mirror the +MCP tools). Phase 3 adds async job + revision procedures and git-sync webhooks: ```bash bun run api # http://localhost:4020 (GET /health, POST/GET /trpc/*) ``` -`bun run index` now also chunks + embeds (Phase 2 indexer). Requires `EMBED_*` -vars in `.env` (see `.env.example`). +tRPC procedures: `search`, `semanticSearch`, `hybridSearch`, `getDocument`, +`listDocuments`, `related` (Phase 2); plus `revisions`, `getRevision`, +`restoreRevision`, `jobStatus`, `queueStatus` (Phase 3). + +Git-sync webhooks (enqueue BullMQ jobs; the worker processes them): +- `POST /hooks/reindex` — full-corpus reindex (point your Git provider's + push webhook here to auto-reindex on push). +- `POST /hooks/index?slug=` — reindex a single document. + +`bun run index` now also chunks + embeds (Phase 2 indexer) and snapshots a +revision whenever the body changes (Phase 3). See `.env.example` for +`EMBED_*` / `REDIS_*` / `QUEUE_PREFIX` vars. ## Status @@ -128,6 +153,11 @@ Postgres FTS keyword search, content indexing. chunked `document_chunks`, `semanticSearch` + `hybridSearch` (RRF), tRPC/Hono API (`apps/api`, :4020), MCP `semantic_search`/`hybrid_search` tools, web hybrid toggle. +**Phase 3 — Async + Scale (DONE):** Redis + BullMQ background indexing/embedding +workers (`packages/queue`, `apps/worker`), git-sync webhooks (`POST /hooks/*`), +document revision system (`document_revisions` + restore), and MCP Resources +(`mcpedia://docs/...`). See `PHASES.md`. + > pgvector is **not installed** on the shared imrnes Postgres, so vector storage is > a `real[]` column with in-app cosine similarity (instant at KB scale). pgvector is > the Phase-4 scale-out path. See `PHASES.md`. diff --git a/apps/api/package.json b/apps/api/package.json index 0ab7c47..648c466 100644 --- a/apps/api/package.json +++ b/apps/api/package.json @@ -13,6 +13,8 @@ "@hono/node-server": "^1.13.0", "@mcpedia/config": "workspace:*", "@mcpedia/core": "workspace:*", + "@mcpedia/db": "workspace:*", + "@mcpedia/queue": "workspace:*", "@trpc/server": "^11.0.0", "hono": "^4.6.0", "zod": "^3.23.8" diff --git a/apps/api/src/index.ts b/apps/api/src/index.ts index cbfb180..5ad635a 100644 --- a/apps/api/src/index.ts +++ b/apps/api/src/index.ts @@ -4,12 +4,31 @@ import { fetchRequestHandler } from "@trpc/server/adapters/fetch"; import { db } from "@mcpedia/db"; import { appRouter } from "./router"; import type { Context } from "./trpc"; +import { enqueueIndexDoc, enqueueFullIndex } from "@mcpedia/queue"; const app = new Hono(); // Health check. app.get("/health", (c) => c.json({ ok: true })); +// --- Phase 3: Git synchronization hook --- +// POST /hooks/reindex -> enqueue a full-corpus reindex (git push webhook) +// POST /hooks/index?slug=... -> enqueue a single document reindex +// Returns the created job id(s). The worker processes them asynchronously. +app.post("/hooks/reindex", async (c) => { + const job = await enqueueFullIndex("git-push"); + return c.json({ ok: true, jobId: job.id, kind: "full" }); +}); + +app.post("/hooks/index", async (c) => { + const slug = c.req.query("slug"); + if (!slug) return c.json({ ok: false, error: "slug query param required" }, 400); + // slug is the relative path without extension, e.g. docs/websocket/contract + const relPath = slug.endsWith(".md") || slug.endsWith(".mdx") ? slug : `${slug}.md`; + const job = await enqueueIndexDoc(relPath, "git-push"); + return c.json({ ok: true, jobId: job.id, kind: "doc", relPath }); +}); + // Mount tRPC at /trpc/*. The fetch adapter is the canonical Bun/Hono adapter. app.all("/trpc/*", (c) => fetchRequestHandler({ diff --git a/apps/api/src/router.ts b/apps/api/src/router.ts index 28c69d9..4bfcb1f 100644 --- a/apps/api/src/router.ts +++ b/apps/api/src/router.ts @@ -7,7 +7,12 @@ import { keywordSearch, listDocuments, semanticSearch, + listRevisions, + getRevision, + restoreRevision, } from "@mcpedia/core"; +import { getQueue, INDEX_QUEUE } from "@mcpedia/queue"; +import { getConnection, BULLMQ_PREFIX } from "@mcpedia/queue/client"; export const appRouter = router({ search: publicProcedure @@ -33,6 +38,58 @@ export const appRouter = router({ related: publicProcedure .input(z.object({ slug: z.string(), limit: z.number().int().min(1).max(20).default(5) })) .query(async ({ input }) => getRelated(input.slug, input.limit)), + + // --- Phase 3: revisions --- + revisions: publicProcedure + .input(z.object({ slug: z.string(), limit: z.number().int().min(1).max(50).default(20) })) + .query(async ({ input }) => listRevisions(input.slug, input.limit)), + + getRevision: publicProcedure + .input(z.object({ id: z.string() })) + .query(async ({ input }) => getRevision(input.id)), + + restoreRevision: publicProcedure + .input(z.object({ id: z.string() })) + .mutation(async ({ input }) => restoreRevision(input.id)), + + // --- Phase 3: async job status --- + jobStatus: publicProcedure + .input(z.object({ id: z.string() })) + .query(async ({ input }) => { + const queue = getQueue(); + const job = await queue.getJob(input.id); + if (!job) return { exists: false }; + const state = await job.getState(); + const failedReason = job.failedReason; + const returnvalue = job.returnvalue; + const progress = job.progress; + return { + exists: true, + id: job.id, + name: job.name, + state, + progress, + failedReason, + returnvalue, + attemptsMade: job.attemptsMade, + }; + }), + + queueStatus: publicProcedure.query(async () => { + const queue = getQueue(); + const [waiting, active, completed, failed, delayed] = await Promise.all([ + queue.getWaitingCount(), + queue.getActiveCount(), + queue.getCompletedCount(), + queue.getFailedCount(), + queue.getDelayedCount(), + ]); + return { + queue: INDEX_QUEUE, + prefix: BULLMQ_PREFIX, + counts: { waiting, active, completed, failed, delayed }, + }; + }), }); export type AppRouter = typeof appRouter; diff --git a/apps/mcp/package.json b/apps/mcp/package.json index 92f4ed5..7ddf67e 100644 --- a/apps/mcp/package.json +++ b/apps/mcp/package.json @@ -17,7 +17,7 @@ "@mcpedia/core": "workspace:*", "@mcpedia/search": "workspace:*", "@modelcontextprotocol/sdk": "^1.29.0", - "zod": "^3.23.8" + "zod": "^4.0.0" }, "devDependencies": { "@types/node": "^20", diff --git a/apps/mcp/src/index.ts b/apps/mcp/src/index.ts index 2c85a5e..bb4f03b 100644 --- a/apps/mcp/src/index.ts +++ b/apps/mcp/src/index.ts @@ -1,7 +1,19 @@ import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js"; +import { ResourceTemplate } from "@modelcontextprotocol/sdk/server/mcp.js"; import { z } from "zod"; -import { listDocuments, getDocument, getRelated, semanticSearch, hybridSearch, keywordSearch } from "@mcpedia/core"; +import { + listDocuments, + getDocument, + getRelated, + semanticSearch, + hybridSearch, + keywordSearch, + listRevisions, + readContentFile, +} from "@mcpedia/core"; +import { CONTENT_ROOT } from "@mcpedia/config"; +import { join } from "node:path"; export function createMcpServer(): McpServer { const server = new McpServer({ @@ -119,6 +131,116 @@ export function createMcpServer(): McpServer { }, ); + // --- Phase 3: MCP Resources (read-only knowledge base surfaced via URIs) --- + // mcpedia://docs -> list all published documents + // mcpedia://docs/{slug} -> full markdown body (from disk) + // mcpedia://docs/{slug}/chunks -> chunked preview (semantic slices) + // mcpedia://docs/{slug}/revisions -> revision history summary + server.registerResource( + "mcpedia-docs-list", + "mcpedia://docs", + { + title: "MCPedia document index", + description: "List of all published documents in the knowledge base.", + mimeType: "application/json", + }, + async (uri) => { + const docs = await listDocuments(); + return { + contents: [ + { + uri: uri.href, + mimeType: "application/json", + text: JSON.stringify(docs, null, 2), + }, + ], + }; + }, + ); + + server.registerResource( + "mcpedia-doc-chunks", + new ResourceTemplate("mcpedia://docs/{+slug}/chunks", { list: undefined }), + { + title: "MCPedia document chunks", + description: "Preview of the embedded semantic chunks for a document.", + mimeType: "application/json", + }, + async (uri, vars) => { + const slug = String(vars.slug); + const doc = await getDocument(slug); + if (!doc) throw new Error(`Document not found: ${slug}`); + // Chunk the body the same way the indexer does (size 1000 / overlap 150) + // so the resource mirrors what semantic search actually sees. + const { chunkText } = await import("@mcpedia/embeddings"); + const chunks = chunkText(doc.body, { size: 1000, overlap: 150 }); + return { + contents: [ + { + uri: uri.href, + mimeType: "application/json", + text: JSON.stringify( + chunks.map((c, i) => ({ index: i, length: c.length, preview: c.slice(0, 200) })), + null, + 2, + ), + }, + ], + }; + }, + ); + + server.registerResource( + "mcpedia-doc-revisions", + new ResourceTemplate("mcpedia://docs/{+slug}/revisions", { list: undefined }), + { + title: "MCPedia document revisions", + description: "Revision history summary for a document.", + mimeType: "application/json", + }, + async (uri, vars) => { + const slug = String(vars.slug); + const revs = await listRevisions(slug, 20); + return { + contents: [ + { + uri: uri.href, + mimeType: "application/json", + text: JSON.stringify(revs, null, 2), + }, + ], + }; + }, + ); + + // Registered LAST: the bare {+slug} template is greedy and would otherwise + // swallow /chunks and /revisions URIs. Specific templates must match first. + server.registerResource( + "mcpedia-doc", + new ResourceTemplate("mcpedia://docs/{+slug}", { list: undefined }), + { + title: "MCPedia document", + description: "Full markdown body of a single document, read from disk (source of truth).", + mimeType: "text/markdown", + }, + async (uri, vars) => { + const slug = String(vars.slug); + const doc = await getDocument(slug); + if (!doc) { + throw new Error(`Document not found: ${slug}`); + } + return { + contents: [ + { + uri: uri.href, + mimeType: "text/markdown", + text: doc.body, + }, + ], + }; + }, + ); + return server; } diff --git a/apps/mcp/src/smoke.test.ts b/apps/mcp/src/smoke.test.ts index 8737f1a..4c1f923 100644 --- a/apps/mcp/src/smoke.test.ts +++ b/apps/mcp/src/smoke.test.ts @@ -92,6 +92,35 @@ async function main() { } console.log(`hybrid_search => ${hybHits.length} docs, top: ${hybHits[0].doc.slug}`); + // 8) resources: list + const resList = await client.listResources(); + const resNames = resList.resources.map((r: any) => r.name).sort(); + console.log("resources:", resNames.join(", ")); + if (!resNames.includes("mcpedia-docs-list")) { + throw new Error("expected mcpedia-docs-list resource"); + } + + // 9) resource: read the docs list (must not throw, returns JSON content) + const readList = await client.readResource({ uri: "mcpedia://docs" }); + const listText = (readList.contents as any)[0].text; + if (!listText.includes("docs/websocket/contract")) { + throw new Error("mcpedia://docs did not list the websocket contract doc"); + } + console.log("readResource(mcpedia://docs) => ok"); + + // 10) resource: read a single doc body + revisions + const readDoc = await client.readResource({ uri: "mcpedia://docs/docs/websocket/contract" }); + const docText = (readDoc.contents as any)[0].text; + if (!docText.includes("WebSocket Contract")) { + throw new Error("mcpedia://docs/{slug} returned unexpected body"); + } + console.log("readResource(mcpedia://docs/docs/websocket/contract) => ok"); + + const readRev = await client.readResource({ + uri: "mcpedia://docs/docs/websocket/contract/revisions", + }); + console.log("readResource(.../revisions) => ok"); + await client.close(); await server.close(); console.log("\nSMOKE OK"); diff --git a/apps/worker/package.json b/apps/worker/package.json new file mode 100644 index 0000000..330847e --- /dev/null +++ b/apps/worker/package.json @@ -0,0 +1,20 @@ +{ + "name": "@mcpedia/worker", + "version": "0.1.0", + "private": true, + "type": "module", + "scripts": { + "start": "bun run src/index.ts", + "lint": "tsc --noEmit", + "typecheck": "tsc --noEmit" + }, + "dependencies": { + "@mcpedia/config": "workspace:*", + "@mcpedia/core": "workspace:*", + "@mcpedia/db": "workspace:*", + "@mcpedia/queue": "workspace:*" + }, + "devDependencies": { + "typescript": "^5.6.0" + } +} diff --git a/apps/worker/src/index.ts b/apps/worker/src/index.ts new file mode 100644 index 0000000..9bf82ff --- /dev/null +++ b/apps/worker/src/index.ts @@ -0,0 +1,12 @@ +import { startWorker } from "@mcpedia/queue/worker"; + +// Keep the process alive: the worker listens on the BullMQ queue until a +// SIGINT/SIGTERM closes it (handled inside startWorker). +const worker = await startWorker(); + +// Heartbeat so the supervisor/operator can see liveness without scraping logs. +const heartbeat = setInterval(() => { + console.log(`[worker] alive, ${worker.name} queue="${worker.name}"`); +}, 30_000); + +worker.on("closed", () => clearInterval(heartbeat)); diff --git a/apps/worker/tsconfig.json b/apps/worker/tsconfig.json new file mode 100644 index 0000000..5eec656 --- /dev/null +++ b/apps/worker/tsconfig.json @@ -0,0 +1,11 @@ +{ + "extends": "../../tsconfig.base.json", + "compilerOptions": { + "paths": { + "@mcpedia/db": ["../../packages/db/src/index.ts"], + "@mcpedia/db/schema": ["../../packages/db/src/schema.ts"], + "@mcpedia/*": ["../../packages/*"] + } + }, + "include": ["src/**/*.ts"] +} diff --git a/bun.lock b/bun.lock index c5bdd92..54db11a 100644 --- a/bun.lock +++ b/bun.lock @@ -19,6 +19,8 @@ "@hono/node-server": "^1.13.0", "@mcpedia/config": "workspace:*", "@mcpedia/core": "workspace:*", + "@mcpedia/db": "workspace:*", + "@mcpedia/queue": "workspace:*", "@trpc/server": "^11.0.0", "hono": "^4.6.0", "zod": "^3.23.8", @@ -38,7 +40,7 @@ "@mcpedia/core": "workspace:*", "@mcpedia/search": "workspace:*", "@modelcontextprotocol/sdk": "^1.29.0", - "zod": "^3.23.8", + "zod": "^4.0.0", }, "devDependencies": { "@types/node": "^20", @@ -67,6 +69,19 @@ "typescript": "^5.6.0", }, }, + "apps/worker": { + "name": "@mcpedia/worker", + "version": "0.1.0", + "dependencies": { + "@mcpedia/config": "workspace:*", + "@mcpedia/core": "workspace:*", + "@mcpedia/db": "workspace:*", + "@mcpedia/queue": "workspace:*", + }, + "devDependencies": { + "typescript": "^5.6.0", + }, + }, "packages/config": { "name": "@mcpedia/config", "version": "0.1.0", @@ -119,6 +134,20 @@ "gray-matter": "^4.0.3", }, }, + "packages/queue": { + "name": "@mcpedia/queue", + "version": "0.1.0", + "dependencies": { + "@mcpedia/config": "workspace:*", + "@mcpedia/core": "workspace:*", + "@mcpedia/db": "workspace:*", + "bullmq": "^6.1.2", + "ioredis": "^6.0.0", + }, + "devDependencies": { + "typescript": "^5.6.0", + }, + }, "packages/search": { "name": "@mcpedia/search", "version": "0.1.0", @@ -142,6 +171,7 @@ "@mcpedia/db": "workspace:*", "@mcpedia/embeddings": "workspace:*", "@mcpedia/parser": "workspace:*", + "@mcpedia/queue": "workspace:*", "@mcpedia/search": "workspace:*", "drizzle-orm": "^0.38.0", "postgres": "^3.4.5", @@ -325,6 +355,8 @@ "@img/sharp-win32-x64": ["@img/sharp-win32-x64@0.35.3", "", { "os": "win32", "cpu": "x64" }, "sha512-D4y1vNeZrIIJCN+uHaWVtH86B+aCrdMYYjicy9pXHvbGZeGYLLSd3wdVuC37FxVXlU1ARsk84eKWfWMXGYEqvA=="], + "@ioredis/commands": ["@ioredis/commands@2.0.0", "", {}, "sha512-vrx0AE/T0h7cRZwfo1M39Cr+ZhZrkf0V8mQN75wucKCxCLD9l/VX6no3gFvrLqD1IlG/1LtzWovqEw3t0Vr9zg=="], + "@jridgewell/gen-mapping": ["@jridgewell/gen-mapping@0.3.13", "", { "dependencies": { "@jridgewell/sourcemap-codec": "^1.5.0", "@jridgewell/trace-mapping": "^0.3.24" } }, "sha512-2kkt/7niJ6MgEPxF0bYdQ6etZaA+fQvDcLKckhy1yIQOzaoKjBBjSj63/aLVjYE3qhRt5dvM+uUyfCg6UKCBbA=="], "@jridgewell/remapping": ["@jridgewell/remapping@2.3.5", "", { "dependencies": { "@jridgewell/gen-mapping": "^0.3.5", "@jridgewell/trace-mapping": "^0.3.24" } }, "sha512-LI9u/+laYG4Ds1TDKSJW2YPrIlcVYOwi2fUC6xB43lueCjgxV4lffOCZCtYFiH6TNOX+tQKXx97T4IKHbhyHEQ=="], @@ -349,6 +381,8 @@ "@mcpedia/parser": ["@mcpedia/parser@workspace:packages/parser"], + "@mcpedia/queue": ["@mcpedia/queue@workspace:packages/queue"], + "@mcpedia/scripts": ["@mcpedia/scripts@workspace:scripts"], "@mcpedia/search": ["@mcpedia/search@workspace:packages/search"], @@ -357,8 +391,22 @@ "@mcpedia/web": ["@mcpedia/web@workspace:apps/web"], + "@mcpedia/worker": ["@mcpedia/worker@workspace:apps/worker"], + "@modelcontextprotocol/sdk": ["@modelcontextprotocol/sdk@1.30.0", "", { "dependencies": { "@hono/node-server": "^1.19.9 || ^2.0.5", "ajv": "^8.17.1", "ajv-formats": "^3.0.1", "content-type": "^1.0.5", "cors": "^2.8.5", "cross-spawn": "^7.0.5", "eventsource": "^3.0.2", "eventsource-parser": "^3.0.0", "express": "^5.2.1", "express-rate-limit": "^8.2.1", "hono": "^4.11.4", "jose": "^6.1.3", "json-schema-typed": "^8.0.2", "pkce-challenge": "^5.0.0", "raw-body": "^3.0.0", "zod": "^3.25 || ^4.0", "zod-to-json-schema": "^3.25.1" }, "peerDependencies": { "@cfworker/json-schema": "^4.1.1" }, "optionalPeers": ["@cfworker/json-schema"] }, "sha512-xKd8OIzlqNzcqcNumGAa6g+PW2kjD5vrpcKOnfldAUPP3j7lnqMPwlTXQm8gF+UwH72z0lqaRbjr9hqGz0eITA=="], + "@msgpackr-extract/msgpackr-extract-darwin-arm64": ["@msgpackr-extract/msgpackr-extract-darwin-arm64@3.0.4", "", { "os": "darwin", "cpu": "arm64" }, "sha512-LCkGo6JDfaBhgST7UpPWgNgLINpcpabaHfyz5OBx75nUYxBsaEPxjnyNjWpeb/xBup/682QnBfRBy2/LvPutZQ=="], + + "@msgpackr-extract/msgpackr-extract-darwin-x64": ["@msgpackr-extract/msgpackr-extract-darwin-x64@3.0.4", "", { "os": "darwin", "cpu": "x64" }, "sha512-zExlW9zUJKZH/tOtVMttwjKa4Xm/3KcNjnE3dPN92uCktwavMxpgCA3MoJK/DOnTWsQgo224OaST27/mPNAf+w=="], + + "@msgpackr-extract/msgpackr-extract-linux-arm": ["@msgpackr-extract/msgpackr-extract-linux-arm@3.0.4", "", { "os": "linux", "cpu": "arm" }, "sha512-Tg3yX65f5GbtXLkrYEHE5oibZG9epyYWas7FogTTEJeDEF9JlXJzKgXaNhT3UXlTOeA+AfZpYZYZ0uPj7Cfquw=="], + + "@msgpackr-extract/msgpackr-extract-linux-arm64": ["@msgpackr-extract/msgpackr-extract-linux-arm64@3.0.4", "", { "os": "linux", "cpu": "arm64" }, "sha512-dgX0P/9wGPJeHFBG+ZmhgE6bmtMt7NP5CRBGyyktpopdk/mW4POnrpQsSLtKI1dwpc+pPLuXHDh6vvskyQE/sw=="], + + "@msgpackr-extract/msgpackr-extract-linux-x64": ["@msgpackr-extract/msgpackr-extract-linux-x64@3.0.4", "", { "os": "linux", "cpu": "x64" }, "sha512-8TNXMEjJc3QEy7R/x1INhgiU+XakDAFUzBhaz7+Rbrs8NH5UQeHQxxmzsSBJGyV6I1jW79undiQm8tOI+D+8FQ=="], + + "@msgpackr-extract/msgpackr-extract-win32-x64": ["@msgpackr-extract/msgpackr-extract-win32-x64@3.0.4", "", { "os": "win32", "cpu": "x64" }, "sha512-CmCXPQrkbwExx3j946/PtHWHbYJiCRBRDl4BlkRQcJB/YOwQxJRTpoo7aTsortjgoJ1x7opzTSxn7C+ASSLVjQ=="], + "@napi-rs/wasm-runtime": ["@napi-rs/wasm-runtime@1.2.3", "", { "dependencies": { "@tybys/wasm-util": "^0.10.3" }, "peerDependencies": { "@emnapi/core": "^1.7.1 || ^2.0.0-alpha.4", "@emnapi/runtime": "^1.7.1 || ^2.0.0-alpha.4" } }, "sha512-UMduMbqO5s5zF2NkNacMT/yK5Y5QiKvWr2+50bzIIxFDwVJ2h49b+oyjaCGPhJxd2/gC2x39EHv/gHVuu36x2Q=="], "@next/env": ["@next/env@16.3.1", "", {}, "sha512-35G3xwkQUb2oETSDjFXGrVugknoayLFBh7vSE+yAcl9IP2zT9wyGwq7297AYHR11kJld807t5f8AJBs6WBzXsQ=="], @@ -591,6 +639,8 @@ "buffer-from": ["buffer-from@1.1.2", "", {}, "sha512-E+XQCRwSbaaiChtv6k6Dwgc+bx+Bs6vuKJHHl5kox/BaKbhiXzqQOwK4cO22yElGp2OCmjwVhT3HmxgyPGnJfQ=="], + "bullmq": ["bullmq@6.1.2", "", { "dependencies": { "cron-parser": "5.10.0", "msgpackr": "2.0.5", "node-abort-controller": "3.1.1", "semver": "7.8.5", "tslib": "2.8.1" }, "peerDependencies": { "bullmq-otel": ">=2.0.0", "ioredis": ">=5.0.0", "pg": ">=8.0.0", "redis": ">=5.0.0" }, "optionalPeers": ["bullmq-otel", "ioredis", "pg", "redis"] }, "sha512-GSX8JfWN8CElAGDyt7Zmq59n1FfWZ0L7IjsThNUQq7X8mvVVDb9I5i2F+nFg4A00UbysJ5yejN0oWzmbUJPDFg=="], + "bytes": ["bytes@3.1.2", "", {}, "sha512-/Nf7TyzTx6S3yRJObOAV7956r8cr2+Oj8AC5dt8wSP3BQAoeX58NoHyCU8P8zGkNXStjTSi6fzO6F0pBdcYbEg=="], "call-bind": ["call-bind@1.0.9", "", { "dependencies": { "call-bind-apply-helpers": "^1.0.2", "es-define-property": "^1.0.1", "get-intrinsic": "^1.3.0", "set-function-length": "^1.2.2" } }, "sha512-a/hy+pNsFUTR+Iz8TCJvXudKVLAnz/DyeSUo10I5yvFDQJBFU2s9uqQpoSrJlroHUKoKqzg+epxyP9lqFdzfBQ=="], @@ -617,6 +667,8 @@ "client-only": ["client-only@0.0.1", "", {}, "sha512-IV3Ou0jSMzZrd3pZ48nLkT9DA7Ag1pnPzaiQhpW7c3RbcqqzvzzVu+L8gfqMp/8IM2MQtSiqaCxrrcfu8I8rMA=="], + "cluster-key-slot": ["cluster-key-slot@1.1.1", "", {}, "sha512-rwHwUfXL40Chm1r08yrhU3qpUvdVlgkKNeyeGPOxnW8/SyVDvgRaed/Uz54AqWNaTCAThlj6QAs3TZcKI0xDEw=="], + "color-convert": ["color-convert@2.0.1", "", { "dependencies": { "color-name": "~1.1.4" } }, "sha512-RRECPsj7iu/xb5oKYcsFHSppFNnsj/52OVTRKb4zP5onXwVF3zVmmToNcOfGC+CRDpfK/U584fMg38ZHCaElKQ=="], "color-name": ["color-name@1.1.4", "", {}, "sha512-dOy+3AuW3a2wNbZHIuMZpTcgjGuLU/uBL/ubcZF9OXbDo8ff4O8yVp5Bf0efS8uEoYo5q4Fx7dY9OgQGXgAsQA=="], @@ -637,6 +689,8 @@ "cors": ["cors@2.8.6", "", { "dependencies": { "object-assign": "^4", "vary": "^1" } }, "sha512-tJtZBBHA6vjIAaF6EnIaq6laBBP9aq/Y3ouVJjEfoHbRBcHBAHYcMh/w8LDrk2PvIMMq8gmopa5D4V8RmbrxGw=="], + "cron-parser": ["cron-parser@5.10.0", "", { "dependencies": { "luxon": "^3.7.2" } }, "sha512-izNAxJyRWUP8ljBoDSub5WyrVOUlT4SLGShswE7eoRBpp6QUsSycYxLBMJlbshgPBMcPT/nrfgjNY2918ayv2A=="], + "cross-spawn": ["cross-spawn@7.0.6", "", { "dependencies": { "path-key": "^3.1.0", "shebang-command": "^2.0.0", "which": "^2.0.1" } }, "sha512-uV2QOWP2nWzsy2aMp8aRibhi9dlzF5Hgh5SHaB9OiTGEyDTiJJyx0uy51QXdyWbtAHNua4XJzUKca3OzKUd3vA=="], "csstype": ["csstype@3.2.3", "", {}, "sha512-z1HGKcYy2xA8AGQfwrn0PAy+PB7X/GSj3UVJW9qKyn43xWa+gl5nXmU4qqLMRzWVLFC8KusUX8T/0kCiOYpAIQ=="], @@ -659,6 +713,8 @@ "define-properties": ["define-properties@1.2.1", "", { "dependencies": { "define-data-property": "^1.0.1", "has-property-descriptors": "^1.0.0", "object-keys": "^1.1.1" } }, "sha512-8QmQKqEASLd5nx0U1B1okLElbUuuttJ/AnYmRXbbbGDWh6uS208EjD4Xqq/I9wK7u0v6O08XhTWnt5XtEbR6Dg=="], + "denque": ["denque@2.1.0", "", {}, "sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw=="], + "depd": ["depd@2.0.0", "", {}, "sha512-g7nH6P6dyDioJogAAGprGpCtVImJhpPk/roCzdb3fIh61/s/nPsfR6onyMwkCAR/OlC3yBC0lESvUoQEAssIrw=="], "dequal": ["dequal@2.0.3", "", {}, "sha512-0je+qPKHEMohvfRTCEo3CrPG6cAzAYgmzKyxRiYSSDkS6eGJdyVJm7WaYA5ECaAD9wLB2T4EEeymA5aFVcYXCA=="], @@ -871,6 +927,8 @@ "internal-slot": ["internal-slot@1.1.0", "", { "dependencies": { "es-errors": "^1.3.0", "hasown": "^2.0.2", "side-channel": "^1.1.0" } }, "sha512-4gd7VpWNQNB4UKKCFFVcp1AVv+FMOgs9NKzjHKusc8jTMhd5eL1NqQqOpE0KzMds804/yHlglp3uxgluOqAPLw=="], + "ioredis": ["ioredis@6.0.0", "", { "dependencies": { "@ioredis/commands": "2.0.0", "cluster-key-slot": "1.1.1", "debug": "4.4.3", "denque": "2.1.0", "redis-errors": "1.2.0", "standard-as-callback": "2.1.0" } }, "sha512-f+Dtubxfpf6KYFq7WVXJoOLn0bk4TJrMrN9SzeE+jrWrCWj7XX3fA6vkryafhADX+GMymRxgDJDOI33COkJc0w=="], + "ip-address": ["ip-address@10.5.0", "", {}, "sha512-R5SnVLJmgYYvf2F2ZgwSBnelz5G4q5AxIC277GDfUaNbrZKNANcBC7RHqYYePlszf4kBolVkJauG0ZjHHFh55g=="], "ipaddr.js": ["ipaddr.js@1.9.1", "", {}, "sha512-0KI/607xoxSToH7GjN1FfSbLoU0+btTicjsQSWQlh/hZykN8KpmMf7uYwPW3R+akZ6R/w18ZlXSHBYXiYUPO3g=="], @@ -1015,6 +1073,8 @@ "lru-cache": ["lru-cache@5.1.1", "", { "dependencies": { "yallist": "^3.0.2" } }, "sha512-KpNARQA3Iwv+jTA0utUVVbrh+Jlrr1Fv0e56GGzAFOXN7dk/FviaDW8LHmK52DlcH4WP2n6gI8vN1aesBFgo9w=="], + "luxon": ["luxon@3.7.2", "", {}, "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew=="], + "magic-string": ["magic-string@0.30.21", "", { "dependencies": { "@jridgewell/sourcemap-codec": "^1.5.5" } }, "sha512-vd2F4YUyEXKGcLHoq+TEyCjxueSeHnFxyyjNp80yg0XV4vUhnDer/lvvlqM/arB5bXQN5K2/3oinyCRyx8T2CQ=="], "math-intrinsics": ["math-intrinsics@1.1.0", "", {}, "sha512-/IXtbwEk5HTPyEwyKX6hGkYXxM9nbj64B+ilVJnC/R6B0pH5G4V3b0pVbL7DBj4tkhBAppbQUlf6F6Xl9LHu1g=="], @@ -1095,6 +1155,10 @@ "ms": ["ms@2.1.3", "", {}, "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA=="], + "msgpackr": ["msgpackr@2.0.5", "", { "optionalDependencies": { "msgpackr-extract": "^3.0.4" } }, "sha512-cef05H/dSYpLpqp3sj/qyZh5vhUYCalnaLO7j1yOmpsR0y/XwLVtK7r5gn+U/F7CTEfMowcGhlUQJDLcLf7jcA=="], + + "msgpackr-extract": ["msgpackr-extract@3.0.4", "", { "dependencies": { "node-gyp-build-optional-packages": "5.2.2" }, "optionalDependencies": { "@msgpackr-extract/msgpackr-extract-darwin-arm64": "3.0.4", "@msgpackr-extract/msgpackr-extract-darwin-x64": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-arm": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-arm64": "3.0.4", "@msgpackr-extract/msgpackr-extract-linux-x64": "3.0.4", "@msgpackr-extract/msgpackr-extract-win32-x64": "3.0.4" }, "bin": { "download-msgpackr-prebuilds": "bin/download-prebuilds.js" } }, "sha512-4kmO/MdyUIkLIvTPr8VHLil4AtoKIoniWPIEk5+CDy0xnWC84azhSFmuJ7PxZdsYtiP5kEeQsORAVIeMgxT+Hw=="], + "nanoid": ["nanoid@3.3.18", "", { "bin": { "nanoid": "bin/nanoid.cjs" } }, "sha512-DTg4MJbGMWkfi6VZFdNt2/caMbQy4Ou+Op/hJQvGEWcnVfoA1QA+xzRKAzw9jD6+GVOOeYr/mIcuDSdug6F6+w=="], "napi-postinstall": ["napi-postinstall@0.3.4", "", { "bin": { "napi-postinstall": "lib/cli.js" } }, "sha512-PHI5f1O0EP5xJ9gQmFGMS6IZcrVvTjpXjz7Na41gTE7eE2hK11lg04CECCYEEjdc17EV4DO+fkGEtt7TpTaTiQ=="], @@ -1105,8 +1169,12 @@ "next": ["next@16.3.1", "", { "dependencies": { "@next/env": "16.3.1", "@swc/helpers": "0.5.23", "baseline-browser-mapping": "^2.9.19", "caniuse-lite": "^1.0.30001579", "postcss": "8.5.23", "styled-jsx": "5.1.6" }, "optionalDependencies": { "@next/swc-darwin-arm64": "16.3.1", "@next/swc-darwin-x64": "16.3.1", "@next/swc-linux-arm64-gnu": "16.3.1", "@next/swc-linux-arm64-musl": "16.3.1", "@next/swc-linux-x64-gnu": "16.3.1", "@next/swc-linux-x64-musl": "16.3.1", "@next/swc-win32-arm64-msvc": "16.3.1", "@next/swc-win32-x64-msvc": "16.3.1", "sharp": "^0.35.3" }, "peerDependencies": { "@opentelemetry/api": "^1.1.0", "@playwright/test": "^1.51.1", "babel-plugin-react-compiler": "*", "react": "^18.2.0 || 19.0.0-rc-de68d2f4-20241204 || ^19.0.0", "react-dom": "^18.2.0 || 19.0.0-rc-de68d2f4-20241204 || ^19.0.0", "sass": "^1.3.0" }, "optionalPeers": ["@opentelemetry/api", "@playwright/test", "babel-plugin-react-compiler", "sass"], "bin": { "next": "dist/bin/next" } }, "sha512-hsAp0i7Rh+/dhe7DGIeN2YlpLM1DP4MNxti9EtDMtqcO612X81MvvEj388/oTce9U1EcEIOWDlGq0zRwrBKvuA=="], + "node-abort-controller": ["node-abort-controller@3.1.1", "", {}, "sha512-AGK2yQKIjRuqnc6VkX2Xj5d+QW8xZ87pa1UK6yA6ouUyuxfHuMP6umE5QK7UmTeOAymo+Zx1Fxiuw9rVx8taHQ=="], + "node-exports-info": ["node-exports-info@1.6.2", "", { "dependencies": { "array.prototype.flatmap": "^1.3.3", "es-errors": "^1.3.0", "object.entries": "^1.1.9", "semver": "^6.3.1" } }, "sha512-kXs9Go0cah0qHVV2v389IXQLdLCeE1xfFtjOAF+iobu0OIoG1pje8At2vMHyaPMiPMnG/LWP50twML21eMcAag=="], + "node-gyp-build-optional-packages": ["node-gyp-build-optional-packages@5.2.2", "", { "dependencies": { "detect-libc": "^2.0.1" }, "bin": { "node-gyp-build-optional-packages": "bin.js", "node-gyp-build-optional-packages-optional": "optional.js", "node-gyp-build-optional-packages-test": "build-test.js" } }, "sha512-s+w+rBWnpTMwSFbaE0UXsRlg7hU4FjekKU4eyAih5T8nJuNZT1nNsskXpxmeqSK9UzkBl6UgRlnKc8hz8IEqOw=="], + "node-releases": ["node-releases@2.0.53", "", {}, "sha512-D9UOmYG3UH1V+ENW56t5QXBwJw1YEY18ruVeus89Rw+SyIgjPkCO84bRzO3uNIYosJbNwiabWVn48o3uJLjxFQ=="], "object-assign": ["object-assign@4.1.1", "", {}, "sha512-rJgTQnkUnH1sFw8yT6VSU3zD3sWmu6sZhIseY8VX+GRu3P6F7Fu+JNDoXfklElbLJSnc3FUQHVe4cU5hj+BcUg=="], @@ -1191,6 +1259,8 @@ "react-markdown": ["react-markdown@9.1.0", "", { "dependencies": { "@types/hast": "^3.0.0", "@types/mdast": "^4.0.0", "devlop": "^1.0.0", "hast-util-to-jsx-runtime": "^2.0.0", "html-url-attributes": "^3.0.0", "mdast-util-to-hast": "^13.0.0", "remark-parse": "^11.0.0", "remark-rehype": "^11.0.0", "unified": "^11.0.0", "unist-util-visit": "^5.0.0", "vfile": "^6.0.0" }, "peerDependencies": { "@types/react": ">=18", "react": ">=18" } }, "sha512-xaijuJB0kzGiUdG7nc2MOMDUDBWPyGAjZtUrow9XxUeua8IqeP+VlIfAZ3bphpcLTnSZXz6z9jcVC/TCwbfgdw=="], + "redis-errors": ["redis-errors@1.2.0", "", {}, "sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w=="], + "reflect.getprototypeof": ["reflect.getprototypeof@1.0.10", "", { "dependencies": { "call-bind": "^1.0.8", "define-properties": "^1.2.1", "es-abstract": "^1.23.9", "es-errors": "^1.3.0", "es-object-atoms": "^1.0.0", "get-intrinsic": "^1.2.7", "get-proto": "^1.0.1", "which-builtin-type": "^1.2.1" } }, "sha512-00o4I+DVrefhv+nX0ulyi3biSHCPDe+yLv5o/p6d/UVlirijB8E16FtfwSAi4g3tcqrQ4lRAqQSoFEZJehYEcw=="], "regexp.prototype.flags": ["regexp.prototype.flags@1.5.4", "", { "dependencies": { "call-bind": "^1.0.8", "define-properties": "^1.2.1", "es-errors": "^1.3.0", "get-proto": "^1.0.1", "gopd": "^1.2.0", "set-function-name": "^2.0.2" } }, "sha512-dYqgNSZbDwkaJ2ceRd9ojCGjBq+mOm9LmtXnAnEGyHhN/5R7iDW2TRw3h+o/jCFxus3P2LfWIIiwowAjANm7IA=="], @@ -1267,6 +1337,8 @@ "stable-hash": ["stable-hash@0.0.5", "", {}, "sha512-+L3ccpzibovGXFK+Ap/f8LOS0ahMrHTf3xu7mMLSpEGU0EO9ucaysSylKo9eRDFNhWve/y275iPmIZ4z39a9iA=="], + "standard-as-callback": ["standard-as-callback@2.1.0", "", {}, "sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A=="], + "statuses": ["statuses@2.0.2", "", {}, "sha512-DvEy55V3DB7uknRo+4iOGT5fP1slR8wQohVdknigZPMpMstaKJQWhwiYBACJE3Ul2pTnATihhBYnRhZQHGBiRw=="], "stop-iteration-iterator": ["stop-iteration-iterator@1.1.0", "", { "dependencies": { "es-errors": "^1.3.0", "internal-slot": "^1.1.0" } }, "sha512-eLoXW/DHyl62zxY4SCaIgnRhuMr6ri4juEYARS8E6sCEqzKpOiE521Ucofdx+KnDZl5xmvGYaaKCk5FEOxJCoQ=="], @@ -1417,6 +1489,8 @@ "@mcpedia/mcp/@types/node": ["@types/node@20.19.43", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-6oYBAi5ikg4Pl+kGsoYtawUMBT2zZMCvPNF7pVLnHZfd1zf38DRiWn/gT01RYCdUqkv7Fhr+C9ot4/tb+2sVvA=="], + "@mcpedia/mcp/zod": ["zod@4.4.3", "", {}, "sha512-ytENFjIJFl2UwYglde2jchW2Hwm4GJFLDiSXWdTrJQBIN9Fcyp7n4DhxJEiWNAJMV1/BqWfW/kkg71UDcHJyTQ=="], + "@mcpedia/web/@types/node": ["@types/node@20.19.43", "", { "dependencies": { "undici-types": "~6.21.0" } }, "sha512-6oYBAi5ikg4Pl+kGsoYtawUMBT2zZMCvPNF7pVLnHZfd1zf38DRiWn/gT01RYCdUqkv7Fhr+C9ot4/tb+2sVvA=="], "@modelcontextprotocol/sdk/@hono/node-server": ["@hono/node-server@2.1.1", "", { "peerDependencies": { "hono": "^4" } }, "sha512-ELuehkj5VCBdgEw9zs+ivkKwyzzUCSQuE96YmiPvn1ECBoZCczbFXJLeEGMTYjphP6gydh4pHMqEYPVMYUVgQg=="], diff --git a/package.json b/package.json index cefd3af..8483c55 100644 --- a/package.json +++ b/package.json @@ -14,6 +14,8 @@ "lint": "turbo run lint", "typecheck": "turbo run typecheck", "index": "bun run scripts/indexer.ts", + "enqueue": "bun run scripts/enqueue.ts", + "worker": "bun --cwd apps/worker run start", "mcp": "bun --cwd apps/mcp run start", "api": "bun --cwd apps/api run dev" }, diff --git a/packages/config/src/index.ts b/packages/config/src/index.ts index abd66fb..2582a16 100644 --- a/packages/config/src/index.ts +++ b/packages/config/src/index.ts @@ -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( diff --git a/packages/core/src/index.service.ts b/packages/core/src/index.service.ts new file mode 100644 index 0000000..15ccb0d --- /dev/null +++ b/packages/core/src/index.service.ts @@ -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["meta"], + body: string, + reason: string, +): Promise { + 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 { + 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 }; +} diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 90fae31..bbccddb 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -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 { diff --git a/packages/core/src/revision.service.ts b/packages/core/src/revision.service.ts new file mode 100644 index 0000000..5916a17 --- /dev/null +++ b/packages/core/src/revision.service.ts @@ -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 { + 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`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 }; +} diff --git a/packages/db/drizzle/0002_document_revisions.sql b/packages/db/drizzle/0002_document_revisions.sql new file mode 100644 index 0000000..5c0a412 --- /dev/null +++ b/packages/db/drizzle/0002_document_revisions.sql @@ -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; diff --git a/packages/db/drizzle/meta/0002_snapshot.json b/packages/db/drizzle/meta/0002_snapshot.json new file mode 100644 index 0000000..10da9e6 --- /dev/null +++ b/packages/db/drizzle/meta/0002_snapshot.json @@ -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": {} + } +} diff --git a/packages/db/drizzle/meta/_journal.json b/packages/db/drizzle/meta/_journal.json index b045e3a..7372393 100644 --- a/packages/db/drizzle/meta/_journal.json +++ b/packages/db/drizzle/meta/_journal.json @@ -15,6 +15,13 @@ "when": 1787137149735, "tag": "0001_document_chunks", "breakpoints": true + }, + { + "idx": 2, + "version": "7", + "when": 1787139150000, + "tag": "0002_document_revisions", + "breakpoints": true } ] } diff --git a/packages/db/src/schema.ts b/packages/db/src/schema.ts index ed99af9..d661b8c 100644 --- a/packages/db/src/schema.ts +++ b/packages/db/src/schema.ts @@ -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; diff --git a/packages/queue/package.json b/packages/queue/package.json new file mode 100644 index 0000000..4d30151 --- /dev/null +++ b/packages/queue/package.json @@ -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" + } +} diff --git a/packages/queue/src/client.ts b/packages/queue/src/client.ts new file mode 100644 index 0000000..3847b66 --- /dev/null +++ b/packages/queue/src/client.ts @@ -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; diff --git a/packages/queue/src/index.ts b/packages/queue/src/index.ts new file mode 100644 index 0000000..3002da7 --- /dev/null +++ b/packages/queue/src/index.ts @@ -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"; diff --git a/packages/queue/src/queue.ts b/packages/queue/src/queue.ts new file mode 100644 index 0000000..da9d9d2 --- /dev/null +++ b/packages/queue/src/queue.ts @@ -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> { + 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> { + return getQueue().add( + "index-all", + { reason }, + { + jobId: `full__${Date.now()}`, + removeOnComplete: 100, + removeOnFail: 1000, + attempts: 1, + }, + ); +} diff --git a/packages/queue/src/worker.ts b/packages/queue/src/worker.ts new file mode 100644 index 0000000..51d70b4 --- /dev/null +++ b/packages/queue/src/worker.ts @@ -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 { + 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; +} diff --git a/scripts/enqueue.ts b/scripts/enqueue.ts new file mode 100644 index 0000000..f732b0d --- /dev/null +++ b/scripts/enqueue.ts @@ -0,0 +1,24 @@ +import { enqueueIndexDoc, enqueueFullIndex } from "@mcpedia/queue"; + +// One-shot enqueue helper (no worker required to schedule work). +// Usage: +// bun run enqueue --all # full reindex +// bun run enqueue docs/websocket/contract # single doc (slug or rel path) +async function main() { + const arg = process.argv[2]; + if (!arg || arg === "--all") { + const job = await enqueueFullIndex("manual"); + console.log(`enqueued full reindex job ${job.id}`); + } else { + const relPath = arg.endsWith(".md") || arg.endsWith(".mdx") ? arg : `${arg}.md`; + const job = await enqueueIndexDoc(relPath, "manual"); + console.log(`enqueued doc reindex job ${job.id} -> ${relPath}`); + } + await new Promise((r) => setTimeout(r, 500)); // allow the event loop to flush + process.exit(0); +} + +main().catch((err) => { + console.error(err); + process.exit(1); +}); diff --git a/scripts/indexer.ts b/scripts/indexer.ts index 9b37d34..362d6a3 100644 --- a/scripts/indexer.ts +++ b/scripts/indexer.ts @@ -1,67 +1,14 @@ -import { db } from "@mcpedia/db"; -import { documents } from "@mcpedia/db/schema"; -import { parseFile } from "@mcpedia/parser"; -import { CONTENT_ROOT } from "@mcpedia/config"; -import { listContentFiles, indexChunks } from "@mcpedia/core"; -import { join } from "node:path"; +import { runFullIndex } from "@mcpedia/core"; +// Phase 3: the indexer now goes through `runFullIndex`, the single indexing +// entry point shared with the BullMQ worker and the git-sync hook. It parses +// each content file, upserts `documents`, chunks+embeds, and snapshots a +// revision when the body changed. async function main() { - const files = listContentFiles(); - let indexed = 0; - let chunked = 0; - for (const rel of files) { - const abs = join(CONTENT_ROOT, rel); - const { meta, body } = parseFile(abs, rel); - 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), - }, - }); - indexed++; - console.log(` indexed ${rel}`); - - // Phase 2: chunk + embed for semantic search. - try { - const n = await indexChunks(meta.slug, body); - chunked += n; - console.log(` embedded ${n} chunks`); - } catch (err) { - console.error( - ` embed FAILED for ${meta.slug}: ${err instanceof Error ? err.message : err}`, - ); - // Don't abort the whole index over one doc's embedding failure. - } - } - console.log(`indexed ${indexed} documents, ${chunked} chunks embedded`); + const reason = process.argv[2] && process.argv[2].startsWith("--reason=") + ? process.argv[2].slice("--reason=".length) + : "index"; + await runFullIndex(reason); } main() diff --git a/scripts/package.json b/scripts/package.json index 45597e4..b578f16 100644 --- a/scripts/package.json +++ b/scripts/package.json @@ -10,6 +10,7 @@ "@mcpedia/embeddings": "workspace:*", "@mcpedia/parser": "workspace:*", "@mcpedia/search": "workspace:*", + "@mcpedia/queue": "workspace:*", "drizzle-orm": "^0.38.0", "postgres": "^3.4.5" }