From 9176e755335d39faf0aa5106064be779dc6e386f Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 15 Jul 2026 19:35:16 +0000 Subject: [PATCH 1/2] feat(jobs): optional dedicated worker process Splits pg-boss setup into startClient() (connect + queues + schedules, run everywhere) and registerWorkers() (the job handlers). The web tier registers workers in-process by default; set WORKERS_IN_PROCESS=false to have it only enqueue and run one or more standalone workers instead (src/worker.ts, `npm run worker`) so long extraction/embedding/generation jobs don't compete with request handling. pg-boss distributes jobs across all connected workers. Adds tsx devDependency and a commented worker service in docker-compose. Default behaviour unchanged. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_0191JikNRN8Q2HBtf6fpLXmH --- docker-compose.yml | 24 +++++++++++++++ package-lock.json | 67 ++++++++++++++++++++++++++++++++++++++++ package.json | 2 ++ src/lib/jobs/index.ts | 72 ++++++++++++++++++++++++++++++++++++------- src/worker.ts | 43 ++++++++++++++++++++++++++ 5 files changed, 197 insertions(+), 11 deletions(-) create mode 100644 src/worker.ts diff --git a/docker-compose.yml b/docker-compose.yml index 01bd722..3753bff 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -26,6 +26,30 @@ services: db: condition: service_healthy restart: unless-stopped + # To move background jobs (embedding, summarization, coverage generation) off + # the web tier, set WORKERS_IN_PROCESS: "false" here and run the worker + # service below. By default (unset) the app runs workers in-process — no + # extra service needed. + # environment: + # WORKERS_IN_PROCESS: "false" + + # Optional dedicated job worker. Runs `npm run worker`, which uses tsx — so it + # needs an image that includes the TS sources + devDependencies (a source + # checkout, or a custom image built from the repo). The default prebuilt + # standalone image runs jobs in-process instead, so this service is opt-in. + # worker: + # image: ghcr.io/veniplex/study-helper:${STUDYHELPER_VERSION:-latest} + # command: ["npm", "run", "worker"] + # environment: + # DATABASE_URL: postgres://study:${POSTGRES_PASSWORD:-study}@db:5432/study + # UPLOAD_DIR: /data/uploads + # TZ: ${TZ:-Europe/Berlin} + # volumes: + # - ${DATA_DIR:-./data}/uploads:/data/uploads + # depends_on: + # db: + # condition: service_healthy + # restart: unless-stopped db: image: pgvector/pgvector:pg17 diff --git a/package-lock.json b/package-lock.json index 8db5ee1..04be377 100644 --- a/package-lock.json +++ b/package-lock.json @@ -78,6 +78,7 @@ "jsdom": "^29.1.1", "prettier": "^3.9.4", "tailwindcss": "^4", + "tsx": "^4.23.0", "typescript": "^5", "vitest": "^4.1.10" } @@ -5090,6 +5091,72 @@ "node": ">=14.0.0" } }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/@emnapi/core": { + "version": "1.11.1", + "dev": true, + "inBundle": true, + "license": "MIT", + "optional": true, + "dependencies": { + "@emnapi/wasi-threads": "1.2.2", + "tslib": "^2.4.0" + } + }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/@emnapi/runtime": { + "version": "1.11.1", + "dev": true, + "inBundle": true, + "license": "MIT", + "optional": true, + "dependencies": { + "tslib": "^2.4.0" + } + }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/@emnapi/wasi-threads": { + "version": "1.2.2", + "dev": true, + "inBundle": true, + "license": "MIT", + "optional": true, + "dependencies": { + "tslib": "^2.4.0" + } + }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/@napi-rs/wasm-runtime": { + "version": "1.1.4", + "dev": true, + "inBundle": true, + "license": "MIT", + "optional": true, + "dependencies": { + "@tybys/wasm-util": "^0.10.1" + }, + "funding": { + "type": "github", + "url": "https://github.com/sponsors/Brooooooklyn" + }, + "peerDependencies": { + "@emnapi/core": "^1.7.1", + "@emnapi/runtime": "^1.7.1" + } + }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/@tybys/wasm-util": { + "version": "0.10.2", + "dev": true, + "inBundle": true, + "license": "MIT", + "optional": true, + "dependencies": { + "tslib": "^2.4.0" + } + }, + "node_modules/@tailwindcss/oxide-wasm32-wasi/node_modules/tslib": { + "version": "2.8.1", + "dev": true, + "inBundle": true, + "license": "0BSD", + "optional": true + }, "node_modules/@tailwindcss/oxide-win32-arm64-msvc": { "version": "4.3.2", "resolved": "https://registry.npmjs.org/@tailwindcss/oxide-win32-arm64-msvc/-/oxide-win32-arm64-msvc-4.3.2.tgz", diff --git a/package.json b/package.json index bde34cf..b815bbf 100644 --- a/package.json +++ b/package.json @@ -6,6 +6,7 @@ "dev": "next dev", "build": "next build", "start": "next start", + "worker": "NODE_OPTIONS=--conditions=react-server tsx src/worker.ts", "lint": "eslint", "format": "prettier --write .", "format:check": "prettier --check .", @@ -87,6 +88,7 @@ "jsdom": "^29.1.1", "prettier": "^3.9.4", "tailwindcss": "^4", + "tsx": "^4.23.0", "typescript": "^5", "vitest": "^4.1.10" } diff --git a/src/lib/jobs/index.ts b/src/lib/jobs/index.ts index 53fce92..d190dd6 100644 --- a/src/lib/jobs/index.ts +++ b/src/lib/jobs/index.ts @@ -12,13 +12,43 @@ export const QUEUE_SEND_REMINDERS = "send-reminders" export const QUEUE_DAILY_PLAN = "daily-plan-reminder" export const QUEUE_CHECK_UPDATES = "check-updates" -async function start(): Promise { +const ALL_QUEUES = [ + QUEUE_EMBED_MATERIAL, + QUEUE_SUMMARIZE_MATERIAL, + QUEUE_GENERATE_COVERAGE, + QUEUE_UNPACK_ZIP, + QUEUE_SEND_REMINDERS, + QUEUE_DAILY_PLAN, + QUEUE_CHECK_UPDATES, +] + +/** Whether this process should run job handlers in-process (default true). Set + * WORKERS_IN_PROCESS=false on the web tier when a dedicated worker runs. */ +function workersInProcess(): boolean { + return (process.env.WORKERS_IN_PROCESS ?? "true").toLowerCase() !== "false" +} + +/** + * Connects pg-boss, ensures every queue exists and installs the cron schedules. + * Runs in every process that needs to enqueue jobs (web tier and worker alike). + */ +export async function startClient(): Promise { const boss = new PgBoss({ connectionString: env.DATABASE_URL }) boss.on("error", (error: Error) => console.error("[pg-boss]", error)) await boss.start() - await boss.createQueue(QUEUE_EMBED_MATERIAL) - await boss.createQueue(QUEUE_SEND_REMINDERS) + for (const queue of ALL_QUEUES) await boss.createQueue(queue) + await boss.schedule(QUEUE_SEND_REMINDERS, "*/5 * * * *") + await boss.schedule(QUEUE_DAILY_PLAN, "0 7 * * *") + await boss.schedule(QUEUE_CHECK_UPDATES, "0 6 * * *") + return boss +} +/** + * Registers all job handlers. Runs in whichever process should do the work — + * the web server by default, or a standalone worker (see startDedicatedWorker) + * when the web tier sets WORKERS_IN_PROCESS=false. + */ +export async function registerWorkers(boss: PgBoss): Promise { await boss.work<{ materialId: string }>(QUEUE_EMBED_MATERIAL, async (jobs) => { const { processMaterial } = await import("@/lib/ai/rag") for (const job of jobs) { @@ -28,7 +58,6 @@ async function start(): Promise { } }) - await boss.createQueue(QUEUE_SUMMARIZE_MATERIAL) await boss.work<{ materialId: string }>(QUEUE_SUMMARIZE_MATERIAL, async (jobs) => { const { summarizeMaterial } = await import("@/lib/ai/generation/summarize") for (const job of jobs) { @@ -36,7 +65,6 @@ async function start(): Promise { } }) - await boss.createQueue(QUEUE_GENERATE_COVERAGE) await boss.work<{ jobId: string }>(QUEUE_GENERATE_COVERAGE, async (jobs) => { const { runCoverageGeneration } = await import("@/lib/ai/generation/generate") for (const job of jobs) { @@ -44,7 +72,6 @@ async function start(): Promise { } }) - await boss.createQueue(QUEUE_UNPACK_ZIP) await boss.work(QUEUE_UNPACK_ZIP, async (jobs) => { const { unpackZip } = await import("./unpack-zip") for (const job of jobs) { @@ -56,16 +83,12 @@ async function start(): Promise { const { sendDueReminders } = await import("./reminders") await sendDueReminders() }) - await boss.schedule(QUEUE_SEND_REMINDERS, "*/5 * * * *") - await boss.createQueue(QUEUE_DAILY_PLAN) await boss.work(QUEUE_DAILY_PLAN, async () => { const { sendDailyPlanReminders } = await import("./reminders") await sendDailyPlanReminders() }) - await boss.schedule(QUEUE_DAILY_PLAN, "0 7 * * *") - await boss.createQueue(QUEUE_CHECK_UPDATES) await boss.work(QUEUE_CHECK_UPDATES, async () => { const { checkForUpdate } = await import("@/lib/update-check") try { @@ -74,8 +97,18 @@ async function start(): Promise { console.error("[check-updates]", error) } }) - await boss.schedule(QUEUE_CHECK_UPDATES, "0 6 * * *") +} +async function start(): Promise { + const boss = await startClient() + if (workersInProcess()) { + await registerWorkers(boss) + console.log("[pg-boss] job workers running in-process") + } else { + console.log( + "[pg-boss] in-process workers disabled (WORKERS_IN_PROCESS=false) — expecting a dedicated worker" + ) + } return boss } @@ -84,6 +117,23 @@ export function getBoss(): Promise { return globalForBoss.boss } +/** + * Entry point for a standalone worker process (see src/worker.ts). Registers + * handlers unconditionally and caches the instance so the module's enqueue + * helpers reuse the same connection instead of spinning up a second one. + */ +export function startDedicatedWorker(): Promise { + if (!globalForBoss.boss) { + globalForBoss.boss = (async () => { + const boss = await startClient() + await registerWorkers(boss) + console.log("[worker] pg-boss dedicated worker started") + return boss + })() + } + return globalForBoss.boss +} + export async function enqueueEmbedMaterial(materialId: string): Promise { const boss = await getBoss() // singletonKey coalesces duplicate enqueues for the same material; processing diff --git a/src/worker.ts b/src/worker.ts new file mode 100644 index 0000000..effc0b5 --- /dev/null +++ b/src/worker.ts @@ -0,0 +1,43 @@ +/** + * Standalone job-worker process. + * + * Run this instead of (or alongside) the Next.js server to execute background + * jobs — text extraction, embedding, summarization and coverage generation — + * off the web tier. Start the web server with WORKERS_IN_PROCESS=false so it + * only enqueues, and run one or more of these workers; pg-boss distributes jobs + * across all connected workers. + * + * WORKERS_IN_PROCESS=false (on the web server) + * npm run worker (one or more worker processes) + * + * Requires the same environment as the app (DATABASE_URL, AI keys, UPLOAD_DIR, + * …). Run via tsx with the react-server condition so `server-only` is a no-op + * (see the "worker" npm script). + */ +import { startDedicatedWorker } from "@/lib/jobs" + +async function main(): Promise { + const boss = await startDedicatedWorker() + + let shuttingDown = false + const shutdown = async (signal: string) => { + if (shuttingDown) return + shuttingDown = true + console.log(`[worker] received ${signal}, draining…`) + try { + await boss.stop({ graceful: true }) + } catch (error) { + console.error("[worker] error during shutdown", error) + } + process.exit(0) + } + process.on("SIGTERM", () => void shutdown("SIGTERM")) + process.on("SIGINT", () => void shutdown("SIGINT")) + + console.log("[worker] ready — waiting for jobs") +} + +void main().catch((error) => { + console.error("[worker] fatal", error) + process.exit(1) +}) From d36097b834eb8663ed6d7c002714501c534d777b Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 15 Jul 2026 19:45:39 +0000 Subject: [PATCH 2/2] feat(retrieval): optional pgvector HNSW ANN index MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds a fast approximate-nearest-neighbour vector index for large corpora, opt-in and off by default (the dimensionless sequential scan stays the default, unchanged). - src/lib/ai/ann.ts: builds a typed shadow column embedding_hnsw vector(N) + HNSW index for the active embedding model, dimension inferred from stored data; rebuild handles model/dimension changes cleanly. State tracked in a new non-secret `ai.ann` setting. - Runs as a background job (reindex-vectors queue) triggered by an admin action (startVectorReindex); admin AI page shows status + a build button (AnnIndexCard), i18n de/en. - hybridSearch uses the ANN column when ready for the active model, else the scan; new material/summary chunks keep the ANN column current via populateAnn. All raw SQL is additive and guarded — failure marks the index failed and search falls back to the scan. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_0191JikNRN8Q2HBtf6fpLXmH --- messages/de.json | 11 ++- messages/en.json | 11 ++- src/app/[locale]/(app)/admin/actions.ts | 15 ++++ src/app/[locale]/(app)/admin/ai/page.tsx | 15 ++-- src/components/admin/ann-index-card.tsx | 95 +++++++++++++++++++++ src/lib/ai/ann.ts | 104 +++++++++++++++++++++++ src/lib/ai/generation/summarize.ts | 2 + src/lib/ai/rag.ts | 19 ++++- src/lib/jobs/index.ts | 17 ++++ src/lib/settings.ts | 16 ++++ 10 files changed, 291 insertions(+), 14 deletions(-) create mode 100644 src/components/admin/ann-index-card.tsx create mode 100644 src/lib/ai/ann.ts diff --git a/messages/de.json b/messages/de.json index f9918fc..da24a55 100644 --- a/messages/de.json +++ b/messages/de.json @@ -176,7 +176,16 @@ "usageUser": "Benutzer", "usageInput": "Input-Tokens", "usageOutput": "Output-Tokens", - "usageEmpty": "Noch kein Verbrauch." + "usageEmpty": "Noch kein Verbrauch.", + "annTitle": "Vektor-Index (ANN, schnelle Suche)", + "annDescription": "Optionaler HNSW-Index für schnelle Vektorsuche bei sehr großen Materialmengen. Baut einen typisierten Index für das aktuelle Embedding-Modell auf. Standardmäßig aus; ohne Index läuft die Suche als (langsamerer) sequentieller Scan.", + "annStatusIdle": "Nicht aufgebaut", + "annStatusBuilding": "Wird aufgebaut …", + "annStatusReady": "Aktiv", + "annStatusFailed": "Fehlgeschlagen", + "annRebuild": "Index (neu) aufbauen", + "annRebuilding": "Index-Aufbau gestartet – läuft im Hintergrund.", + "annNeedEmbedding": "Zuerst ein Standard-Embedding-Modell konfigurieren." }, "email": { "title": "E-Mail (SMTP)", diff --git a/messages/en.json b/messages/en.json index 3e5f296..4de09f8 100644 --- a/messages/en.json +++ b/messages/en.json @@ -176,7 +176,16 @@ "usageUser": "User", "usageInput": "Input tokens", "usageOutput": "Output tokens", - "usageEmpty": "No usage yet." + "usageEmpty": "No usage yet.", + "annTitle": "Vector index (ANN, fast search)", + "annDescription": "Optional HNSW index for fast vector search over very large material sets. Builds a typed index for the current embedding model. Off by default; without it, search runs as a (slower) sequential scan.", + "annStatusIdle": "Not built", + "annStatusBuilding": "Building …", + "annStatusReady": "Active", + "annStatusFailed": "Failed", + "annRebuild": "Build / rebuild index", + "annRebuilding": "Index build started — running in the background.", + "annNeedEmbedding": "Configure a default embedding model first." }, "email": { "title": "Email (SMTP)", diff --git a/src/app/[locale]/(app)/admin/actions.ts b/src/app/[locale]/(app)/admin/actions.ts index 3c34188..c82d6ba 100644 --- a/src/app/[locale]/(app)/admin/actions.ts +++ b/src/app/[locale]/(app)/admin/actions.ts @@ -114,3 +114,18 @@ export async function saveAiSettings(value: unknown) { revalidatePath("/", "layout") return { ok: true as const } } + +/** Kicks off a background rebuild of the pgvector HNSW ANN index. */ +export async function startVectorReindex() { + await requireAdmin() + const { enqueueReindexVectors } = await import("@/lib/jobs") + await enqueueReindexVectors() + return { ok: true as const } +} + +/** Current ANN index state (status/model/dimensions), for the admin UI. */ +export async function getAnnStatus() { + await requireAdmin() + const { getSetting } = await import("@/lib/settings") + return (await getSetting("ai.ann")) ?? { status: "idle" as const } +} diff --git a/src/app/[locale]/(app)/admin/ai/page.tsx b/src/app/[locale]/(app)/admin/ai/page.tsx index 0437b59..31c7a0c 100644 --- a/src/app/[locale]/(app)/admin/ai/page.tsx +++ b/src/app/[locale]/(app)/admin/ai/page.tsx @@ -6,17 +6,14 @@ import { requireAdmin } from "@/lib/auth/session" import { getSetting } from "@/lib/settings" import { daysAgo } from "@/lib/utils" import { AiSettingsForm } from "@/components/admin/ai-settings-form" -import { - Card, - CardContent, - CardHeader, - CardTitle, -} from "@/components/ui/card" +import { AnnIndexCard } from "@/components/admin/ann-index-card" +import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/card" export default async function AdminAiPage() { await requireAdmin() const t = await getTranslations("admin.ai") const ai = await getSetting("ai") + const ann = await getSetting("ai.ann") const thirtyDaysAgo = daysAgo(30) const usage = await db @@ -33,8 +30,10 @@ export default async function AdminAiPage() { return (
- + diff --git a/src/components/admin/ann-index-card.tsx b/src/components/admin/ann-index-card.tsx new file mode 100644 index 0000000..09eb8f5 --- /dev/null +++ b/src/components/admin/ann-index-card.tsx @@ -0,0 +1,95 @@ +"use client" + +import * as React from "react" +import { Loader2 } from "lucide-react" +import { useTranslations } from "next-intl" +import { toast } from "sonner" +import { Button } from "@/components/ui/button" +import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/card" +import { getAnnStatus, startVectorReindex } from "@/app/[locale]/(app)/admin/actions" + +type AnnStatus = { + status: "idle" | "building" | "ready" | "failed" + embeddingModel?: string + dimensions?: number + error?: string +} + +/** + * Admin control for the optional pgvector HNSW ANN index: shows its state and + * triggers a background rebuild. While building it polls for status. + */ +export function AnnIndexCard({ + initial, + embeddingConfigured, +}: { + initial: AnnStatus + embeddingConfigured: boolean +}) { + const t = useTranslations("admin.ai") + const [status, setStatus] = React.useState(initial) + const [pending, setPending] = React.useState(false) + + React.useEffect(() => { + if (status.status !== "building") return + let active = true + const timer = setInterval(async () => { + try { + const s = (await getAnnStatus()) as AnnStatus + if (active) setStatus(s) + } catch { + // transient — keep polling + } + }, 3000) + return () => { + active = false + clearInterval(timer) + } + }, [status.status]) + + async function onRebuild() { + setPending(true) + try { + await startVectorReindex() + setStatus((s) => ({ ...s, status: "building" })) + toast.success(t("annRebuilding")) + } catch (error) { + toast.error(error instanceof Error ? error.message : String(error)) + } finally { + setPending(false) + } + } + + const label: Record = { + idle: t("annStatusIdle"), + building: t("annStatusBuilding"), + ready: t("annStatusReady"), + failed: t("annStatusFailed"), + } + const busy = pending || status.status === "building" + + return ( + + + {t("annTitle")} + + +

{t("annDescription")}

+
+ {label[status.status]} + {status.status === "ready" && status.dimensions + ? ` · ${status.embeddingModel} · dim ${status.dimensions}` + : ""} + {status.status === "failed" && status.error ? ` — ${status.error}` : ""} +
+ + {!embeddingConfigured && ( +

{t("annNeedEmbedding")}

+ )} +
+
+ ) +} diff --git a/src/lib/ai/ann.ts b/src/lib/ai/ann.ts new file mode 100644 index 0000000..1253816 --- /dev/null +++ b/src/lib/ai/ann.ts @@ -0,0 +1,104 @@ +import "server-only" +import { sql } from "drizzle-orm" +import { db } from "@/db" +import { getSetting, setSetting } from "@/lib/settings" + +// Optional pgvector HNSW ANN index. The base `material_chunk.embedding` column +// is deliberately dimension-less (so different embedding models can coexist), +// which pgvector cannot index. This module maintains a typed shadow column +// `embedding_hnsw vector(N)` + an HNSW index for the *active* embedding model, +// built on demand by an admin. Everything here is additive and opt-in: unless +// an admin has built the index, `annDimensionFor` returns null and search keeps +// using the sequential cosine scan (unchanged default). + +const ANN_COLUMN = "embedding_hnsw" +const ANN_INDEX = "material_chunk_hnsw_idx" + +/** + * The ANN dimension to use for this embedding model, or null when the ANN index + * isn't ready for it (callers then fall back to the dimensionless scan). + */ +export async function annDimensionFor(embeddingRef: string): Promise { + const ann = await getSetting("ai.ann") + if (!ann || ann.status !== "ready" || ann.embeddingModel !== embeddingRef || !ann.dimensions) { + return null + } + return ann.dimensions +} + +/** Populates the ANN column for a material's freshly-embedded chunks. No-op + * unless the ANN index is ready for the active model. */ +export async function populateAnn(materialId: string, embeddingRef: string): Promise { + const dim = await annDimensionFor(embeddingRef) + if (!dim) return + try { + await db.execute(sql` + UPDATE material_chunk + SET ${sql.raw(ANN_COLUMN)} = embedding::text::vector(${sql.raw(String(dim))}) + WHERE material_id = ${materialId} + AND embedding_model = ${embeddingRef} + AND embedding IS NOT NULL + AND ${sql.raw(ANN_COLUMN)} IS NULL + `) + } catch (error) { + // ANN is best-effort; a failure here must not break ingestion. + console.error("[ann] populate failed", materialId, error) + } +} + +/** + * (Re)builds the HNSW ANN index for the active embedding model. Heavy — runs in + * the background. Infers the embedding dimension from stored data, rebuilds the + * typed column + index from `embedding` (the source of truth), and records the + * result in the `ai.ann` setting. On any failure the state is marked `failed` + * and search continues on the sequential scan. + */ +export async function reindexVectors(): Promise { + const ai = await getSetting("ai") + const embeddingRef = ai?.defaultEmbeddingModel + if (!embeddingRef) { + await setSetting("ai.ann", { status: "failed", error: "No embedding model configured" }) + return + } + await setSetting("ai.ann", { status: "building", embeddingModel: embeddingRef }) + try { + const dimRows = await db.execute<{ dim: number }>(sql` + SELECT vector_dims(embedding) AS dim + FROM material_chunk + WHERE embedding_model = ${embeddingRef} AND embedding IS NOT NULL + LIMIT 1 + `) + const dim = Number(dimRows[0]?.dim) + if (!Number.isInteger(dim) || dim <= 0) { + await setSetting("ai.ann", { + status: "failed", + embeddingModel: embeddingRef, + error: "No embeddings found for the active model yet", + }) + return + } + const dimLit = sql.raw(String(dim)) + // Rebuild from scratch so a model/dimension change is handled cleanly. + await db.execute(sql`DROP INDEX IF EXISTS ${sql.raw(ANN_INDEX)}`) + await db.execute(sql`ALTER TABLE material_chunk DROP COLUMN IF EXISTS ${sql.raw(ANN_COLUMN)}`) + await db.execute( + sql`ALTER TABLE material_chunk ADD COLUMN ${sql.raw(ANN_COLUMN)} vector(${dimLit})` + ) + await db.execute(sql` + UPDATE material_chunk + SET ${sql.raw(ANN_COLUMN)} = embedding::text::vector(${dimLit}) + WHERE embedding_model = ${embeddingRef} AND embedding IS NOT NULL + `) + await db.execute( + sql`CREATE INDEX ${sql.raw(ANN_INDEX)} ON material_chunk USING hnsw (${sql.raw(ANN_COLUMN)} vector_cosine_ops)` + ) + await setSetting("ai.ann", { status: "ready", embeddingModel: embeddingRef, dimensions: dim }) + } catch (error) { + console.error("[ann] reindex failed", error) + await setSetting("ai.ann", { + status: "failed", + embeddingModel: embeddingRef, + error: error instanceof Error ? error.message.slice(0, 500) : "reindex failed", + }) + } +} diff --git a/src/lib/ai/generation/summarize.ts b/src/lib/ai/generation/summarize.ts index 91f9948..037b62f 100644 --- a/src/lib/ai/generation/summarize.ts +++ b/src/lib/ai/generation/summarize.ts @@ -8,6 +8,7 @@ import { readStoredText } from "@/lib/storage" import { getEmbeddingModel, getLanguageModel, resolveModelForUser } from "@/lib/ai/registry" import { recordAiAudit, runAi, type AiUsage } from "@/lib/ai/run" import { chunkText } from "@/lib/ai/rag" +import { populateAnn } from "@/lib/ai/ann" /** Leaf chunks grouped per section-summary call. */ const SECTION_UNITS = 20 @@ -162,6 +163,7 @@ export async function summarizeMaterial(materialId: string): Promise { const ai = await getSetting("ai") const embeddingRef = ai?.defaultEmbeddingModel await storeSummaryChunks(row, sectionSummaries, embeddingRef, acc) + if (embeddingRef) await populateAnn(row.id, embeddingRef) // REDUCE: roll section summaries up into a single document summary. let docSummary = sectionSummaries[0] diff --git a/src/lib/ai/rag.ts b/src/lib/ai/rag.ts index 2532705..9a109b6 100644 --- a/src/lib/ai/rag.ts +++ b/src/lib/ai/rag.ts @@ -8,6 +8,7 @@ import { deleteFile, readStoredText, saveText } from "@/lib/storage" import { getEmbeddingModel } from "./registry" import { recordAiAudit, runAi } from "./run" import { extractText } from "./extract" +import { annDimensionFor, populateAnn } from "./ann" /** Bounded preview of extracted text kept in the DB for quick ILIKE search. */ const PREVIEW_CHARS = 200_000 @@ -239,6 +240,9 @@ async function embedMaterialText( { inputTokens: totalInput, outputTokens: 0, totalTokens: totalInput } ) } + + // Keep the ANN index column current (no-op unless an admin built it). + await populateAnn(materialId, embeddingRef) } export type RagHit = { @@ -329,7 +333,9 @@ async function hybridSearch( .orderBy(sql`ts_rank(${materialChunk.contentTsv}, ${tsq}) DESC`) .limit(pool) - // Vector (cosine) ranking. + // Vector (cosine) ranking — uses the HNSW ANN column when an admin has built + // it for the active model, otherwise the dimensionless sequential scan. + const annDim = await annDimensionFor(embeddingRef) const model = await getEmbeddingModel(embeddingRef, userId) const { embedding } = await runAi( { @@ -342,18 +348,23 @@ async function hybridSearch( () => embed({ model, value: query }) ) const vectorLiteral = `[${embedding.join(",")}]` + const embCol = annDim ? sql.raw("material_chunk.embedding_hnsw") : sql`${materialChunk.embedding}` + const distExpr = sql`${embCol} <=> ${vectorLiteral}::vector` + const vectorWhere = annDim + ? and(base, eq(materialChunk.embeddingModel, embeddingRef), sql`${embCol} IS NOT NULL`) + : and(base, eq(materialChunk.embeddingModel, embeddingRef)) const vector = await db .select({ id: materialChunk.id, content: materialChunk.content, materialName: material.name, materialId: material.id, - similarity: sql`1 - (${materialChunk.embedding} <=> ${vectorLiteral}::vector)`, + similarity: sql`1 - (${distExpr})`, }) .from(materialChunk) .innerJoin(material, eq(materialChunk.materialId, material.id)) - .where(and(base, eq(materialChunk.embeddingModel, embeddingRef))) - .orderBy(sql`${materialChunk.embedding} <=> ${vectorLiteral}::vector`) + .where(vectorWhere) + .orderBy(distExpr) .limit(pool) return rrfFuse(vector, lexical, limit) diff --git a/src/lib/jobs/index.ts b/src/lib/jobs/index.ts index d190dd6..5e62164 100644 --- a/src/lib/jobs/index.ts +++ b/src/lib/jobs/index.ts @@ -7,6 +7,7 @@ const globalForBoss = globalThis as unknown as { boss?: Promise } export const QUEUE_EMBED_MATERIAL = "embed-material" export const QUEUE_SUMMARIZE_MATERIAL = "summarize-material" export const QUEUE_GENERATE_COVERAGE = "generate-coverage" +export const QUEUE_REINDEX_VECTORS = "reindex-vectors" export const QUEUE_UNPACK_ZIP = "unpack-zip" export const QUEUE_SEND_REMINDERS = "send-reminders" export const QUEUE_DAILY_PLAN = "daily-plan-reminder" @@ -16,6 +17,7 @@ const ALL_QUEUES = [ QUEUE_EMBED_MATERIAL, QUEUE_SUMMARIZE_MATERIAL, QUEUE_GENERATE_COVERAGE, + QUEUE_REINDEX_VECTORS, QUEUE_UNPACK_ZIP, QUEUE_SEND_REMINDERS, QUEUE_DAILY_PLAN, @@ -72,6 +74,11 @@ export async function registerWorkers(boss: PgBoss): Promise { } }) + await boss.work(QUEUE_REINDEX_VECTORS, async () => { + const { reindexVectors } = await import("@/lib/ai/ann") + await reindexVectors() + }) + await boss.work(QUEUE_UNPACK_ZIP, async (jobs) => { const { unpackZip } = await import("./unpack-zip") for (const job of jobs) { @@ -165,6 +172,16 @@ export async function enqueueGeneration(jobId: string): Promise { ) } +export async function enqueueReindexVectors(): Promise { + const boss = await getBoss() + // Long-running maintenance; only one at a time. + await boss.send( + QUEUE_REINDEX_VECTORS, + {}, + { retryLimit: 0, singletonKey: "reindex", expireInSeconds: 7200 } + ) +} + export async function enqueueUnpackZip( payload: import("./unpack-zip").UnpackZipPayload ): Promise { diff --git a/src/lib/settings.ts b/src/lib/settings.ts index e44a1cc..2d98d58 100644 --- a/src/lib/settings.ts +++ b/src/lib/settings.ts @@ -100,6 +100,20 @@ export const aiSettingsSchema = z.object({ export type AiSettings = z.infer export type AiProvider = z.infer +/** + * State of the optional pgvector HNSW ANN index (fast vector search). Kept + * separate from the (encrypted) `ai` settings because it's operational state, + * not a credential, and is written from a background reindex job. + */ +export const annSettingsSchema = z.object({ + status: z.enum(["idle", "building", "ready", "failed"]).default("idle"), + /** "providerId:modelId" the index was built for. */ + embeddingModel: z.string().optional(), + /** Embedding dimension the typed ANN column was created with. */ + dimensions: z.number().int().positive().optional(), + error: z.string().optional(), +}) + const settingsSchemas = { "auth.registrationMode": registrationModeSchema, "auth.socialProviders": socialProvidersSchema, @@ -108,6 +122,7 @@ const settingsSchemas = { branding: brandingSchema, uploads: uploadsSchema, ai: aiSettingsSchema, + "ai.ann": annSettingsSchema, "push.vapid": vapidSchema, "system.updateCheck": updateCheckSchema, } as const @@ -132,6 +147,7 @@ const defaults: { [K in SettingKey]: SettingValue } = { branding: { appName: "StudyHelper" }, uploads: { maxUploadMb: 200, storageQuotaMbPerUser: 0 }, ai: { providers: [], monthlyTokenLimitPerUser: 0 }, + "ai.ann": { status: "idle" }, "push.vapid": undefined as never, // generated on first use "system.updateCheck": undefined as never, // set once the first check has run }