import { internalAction, internalMutation, internalQuery, mutation, query } from "./_generated/server"; import { v } from "convex/values"; import { nowIso } from "./_utils/time"; import { api, internal } from "./_generated/api"; import { lightragIngestText } from "./_utils/lightrag"; import { extractTextFromDocumentContent, extractTextFromMindmapData } from "./_utils/text"; import { enqueueIngestDocumentJob, enqueueIngestMediaAssetJob, enqueueIngestMindmapJob } from "./_utils/ingestJobs"; export const get = query({ args: { userId: v.string(), id: v.string() }, handler: async (ctx, args) => { const job = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (!job) return null; if (job.user_id !== args.userId) return null; return { id: job.id, type: job.type, status: job.status, payload: job.payload, result: job.result, error: job.error, created_at: job.created_at, updated_at: job.updated_at, started_at: job.started_at, finished_at: job.finished_at, }; }, }); export const enqueueDemo = mutation({ args: { userId: v.string(), id: v.string(), ms: v.optional(v.number()) }, handler: async (ctx, args) => { const ts = nowIso(); const payload = { ms: args.ms ?? 800 }; await ctx.db.insert("jobs", { id: args.id, user_id: args.userId, type: "demo.sleep", status: "queued", payload, result: null, error: null, created_at: ts, updated_at: ts, started_at: null, finished_at: null, }); // 说明:阶段 5 骨架——用 scheduler 触发内部 mutation,再由内部 action 执行耗时逻辑。 await ctx.scheduler.runAfter(0, internal.jobs.start, { id: args.id }); return { ok: true, id: args.id }; }, }); export const enqueueRagIndexDocument = mutation({ args: { userId: v.string(), documentId: v.string() }, handler: async (ctx, args) => { const { id } = await enqueueIngestDocumentJob(ctx, { userId: args.userId, documentId: args.documentId, debounceMs: 0 }); return { ok: true, id }; }, }); export const enqueueRagIndexMindmap = mutation({ args: { userId: v.string(), docId: v.string(), mindmapId: v.string() }, handler: async (ctx, args) => { const { id } = await enqueueIngestMindmapJob(ctx, { userId: args.userId, docId: args.docId, mindmapId: args.mindmapId, debounceMs: 0, }); return { ok: true, id }; }, }); export const enqueueRagIndexMediaAsset = mutation({ args: { userId: v.string(), assetId: v.string() }, handler: async (ctx, args) => { const { id } = await enqueueIngestMediaAssetJob(ctx, { userId: args.userId, assetId: args.assetId, debounceMs: 0 }); return { ok: true, id }; }, }); export const start = internalMutation({ args: { id: v.string() }, handler: async (ctx, args) => { const job = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (!job) return; if (job.status !== "queued") return; const ts = nowIso(); await ctx.db.patch(job._id, { status: "running", started_at: ts, updated_at: ts }); await ctx.scheduler.runAfter(0, internal.jobs.run, { id: args.id }); }, }); export const run = internalAction({ args: { id: v.string() }, handler: async (ctx, args) => { const job = await ctx.runQuery(internal.jobs._getInternal, { id: args.id }); if (!job) return; if (job.status !== "running") return; try { if (job.type === "demo.sleep") { const ms = typeof job.payload?.ms === "number" ? job.payload.ms : 800; await new Promise((r) => setTimeout(r, ms)); await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, sleptMs: ms }, }); return; } if (job.type === "ingest.rag_index_document") { const documentId = String(job.payload?.documentId ?? "").trim(); if (!documentId) throw new Error("缺少 documentId"); const [meta, contentRes] = await Promise.all([ ctx.runQuery(internal.documents.getMetaForIngest, { userId: job.user_id, id: documentId }), ctx.runQuery(internal.documents.getContentForIngest, { userId: job.user_id, id: documentId }), ]); if (!meta) throw new Error("页面不存在或无权限"); const title = meta.title ?? "无标题"; const text = extractTextFromDocumentContent(contentRes?.content ?? null); const fileSource = `document:${documentId}`; const { trackId } = await lightragIngestText({ fileSource, text: `# ${title}\n\n${text}` }); await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, kind: "document", documentId, trackId }, }); return; } if (job.type === "ingest.rag_index_mindmap") { const docId = String(job.payload?.docId ?? "").trim(); const mindmapId = String(job.payload?.mindmapId ?? "").trim(); if (!docId) throw new Error("缺少 docId"); if (!mindmapId) throw new Error("缺少 mindmapId"); const [docMeta, mindmapRes] = await Promise.all([ ctx.runQuery(internal.documents.getMetaForIngest, { userId: job.user_id, id: docId }), ctx.runQuery(internal.mindmaps.getForIngest, { userId: job.user_id, docId, mindmapId }), ]); if (!docMeta) throw new Error("页面不存在或无权限"); if (!mindmapRes?.ok) throw new Error("导图读取失败"); if (!mindmapRes.meta?.exists) { await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, kind: "mindmap", docId, mindmapId, skipped: true, reason: "mindmap_not_exists" }, }); return; } const title = docMeta.title ?? "无标题"; const text = extractTextFromMindmapData(mindmapRes.data ?? null); const fileSource = `mindmap:${docId}:${mindmapId}`; const { trackId } = await lightragIngestText({ fileSource, text: `# ${title}\n\n${text}` }); await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, kind: "mindmap", docId, mindmapId, trackId }, }); return; } if (job.type === "ingest.rag_index_media_asset") { const assetId = String(job.payload?.assetId ?? "").trim(); if (!assetId) throw new Error("缺少 assetId"); const asset = await ctx.runQuery(api.mediaAssets.getById, { userId: job.user_id, id: assetId }); if (!asset) throw new Error("资源不存在或无权限"); const title = asset.file_name ?? asset.id; const text = String(asset.ocr_text ?? "").trim(); if (!text) { await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, kind: "media_asset", assetId, skipped: true, reason: "empty_ocr_text" }, }); return; } const fileSource = `media:${assetId}`; const { trackId } = await lightragIngestText({ fileSource, text: `# ${title}\n\n${text}` }); await ctx.runMutation(internal.jobs.finishSuccess, { id: args.id, result: { ok: true, kind: "media_asset", assetId, trackId }, }); return; } throw new Error(`未知任务类型:${job.type}`); } catch (err) { const message = err instanceof Error ? err.message : String(err); await ctx.runMutation(internal.jobs.finishFailure, { id: args.id, error: message }); } }, }); export const _getInternal = internalQuery({ args: { id: v.string() }, handler: async (ctx, args) => { const job = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (!job) return null; return { id: job.id, user_id: job.user_id, type: job.type, status: job.status, payload: job.payload, }; }, }); export const finishSuccess = internalMutation({ args: { id: v.string(), result: v.any() }, handler: async (ctx, args) => { const job = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (!job) return; const ts = nowIso(); await ctx.db.patch(job._id, { status: "succeeded", result: args.result, error: null, finished_at: ts, updated_at: ts, }); }, }); export const finishFailure = internalMutation({ args: { id: v.string(), error: v.string() }, handler: async (ctx, args) => { const job = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (!job) return; const ts = nowIso(); await ctx.db.patch(job._id, { status: "failed", result: null, error: args.error, finished_at: ts, updated_at: ts, }); }, });