Files
mnote/wolai-frontend/convex/jobs.ts
T
2026-01-24 12:32:51 +08:00

365 lines
12 KiB
TypeScript

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";
import { extractTextFromAttachment } from "./_utils/attachmentExtract";
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;
}
if (job.type === "extract.media_asset_text") {
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("资源不存在或无权限");
// 说明:只处理 file 类型的常见附件(pdf/docx/pptx/xlsx)。
if (String(asset.asset_type ?? "") !== "file") {
await ctx.runMutation(api.mediaAssets.patchById, {
userId: job.user_id,
id: assetId,
patch: {
ocr_status: "skipped",
ocr_payload: { reason: "non_file_asset" },
ocr_strategy: "attachment_extract",
},
});
await ctx.runMutation(internal.jobs.finishSuccess, {
id: args.id,
result: { ok: true, kind: "extract_media_asset_text", assetId, skipped: true, reason: "non_file_asset" },
});
return;
}
const fileSize = typeof asset.file_size === "number" ? asset.file_size : null;
const maxBytes = 25 * 1024 * 1024;
if (typeof fileSize === "number" && fileSize > maxBytes) {
await ctx.runMutation(api.mediaAssets.patchById, {
userId: job.user_id,
id: assetId,
patch: {
ocr_status: "failed",
ocr_payload: { error: `文件过大(${fileSize} bytes),暂不解析`, maxBytes },
ocr_strategy: "attachment_extract",
},
});
await ctx.runMutation(internal.jobs.finishFailure, {
id: args.id,
error: "文件过大,暂不解析",
});
return;
}
// 说明:Convex Files 的 getUrl 可能过期,先刷新并获取当前可用链接。
const refreshed = await ctx.runMutation(api.mediaAssets.refreshUrl, { userId: job.user_id, id: assetId });
const url = String((refreshed as any)?.signedUrl ?? asset.file_url ?? "").trim();
if (!url) throw new Error("缺少可用文件链接");
const res = await fetch(url);
if (!res.ok) {
throw new Error(`下载附件失败:${res.status}`);
}
const bytes = await res.arrayBuffer();
const extracted = await extractTextFromAttachment({
mimeType: (asset as any).mime_type ?? null,
fileName: (asset as any).file_name ?? null,
bytes,
});
if (!extracted.ok) {
await ctx.runMutation(api.mediaAssets.patchById, {
userId: job.user_id,
id: assetId,
patch: {
ocr_status: "failed",
ocr_payload: { error: extracted.reason, meta: extracted.meta ?? null },
ocr_strategy: extracted.strategy,
},
});
await ctx.runMutation(internal.jobs.finishFailure, { id: args.id, error: extracted.reason });
return;
}
await ctx.runMutation(api.mediaAssets.patchById, {
userId: job.user_id,
id: assetId,
patch: {
ocr_text: extracted.text,
ocr_status: "completed",
ocr_payload: extracted.meta ?? null,
ocr_strategy: extracted.strategy,
},
});
await ctx.runMutation(internal.jobs.finishSuccess, {
id: args.id,
result: {
ok: true,
kind: "extract_media_asset_text",
assetId,
strategy: extracted.strategy,
chars: extracted.text.length,
},
});
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,
});
},
});