import { internal } from "../_generated/api"; import { nowIso } from "./time"; import type { MutationCtx } from "../_generated/server"; type JobStatus = "queued" | "running" | "succeeded" | "failed"; async function upsertJob( ctx: MutationCtx, args: { id: string; userId: string; type: string; payload: unknown; debounceMs?: number; }, ): Promise<{ id: string; status: JobStatus }> { const debounceMs = Math.max(0, Math.floor(args.debounceMs ?? 1200)); const ts = nowIso(); const existing = await ctx.db .query("jobs") .withIndex("by_job_id", (q) => q.eq("id", args.id)) .first(); if (existing) { if (existing.user_id !== args.userId) { // 说明:避免不同用户复用同一 job id 导致“互相覆盖”。 return { id: args.id, status: existing.status as JobStatus }; } const status = String(existing.status) as JobStatus; if (status === "running" || status === "queued") { await ctx.db.patch(existing._id, { updated_at: ts }); // 说明:start 内部会检查状态,重复 schedule 不会造成重复执行。 await ctx.scheduler.runAfter(debounceMs, internal.jobs.start, { id: args.id }); return { id: args.id, status }; } await ctx.db.patch(existing._id, { type: args.type, status: "queued", payload: args.payload, result: null, error: null, updated_at: ts, started_at: null, finished_at: null, }); await ctx.scheduler.runAfter(debounceMs, internal.jobs.start, { id: args.id }); return { id: args.id, status: "queued" }; } await ctx.db.insert("jobs", { id: args.id, user_id: args.userId, type: args.type, status: "queued", payload: args.payload, result: null, error: null, created_at: ts, updated_at: ts, started_at: null, finished_at: null, }); await ctx.scheduler.runAfter(debounceMs, internal.jobs.start, { id: args.id }); return { id: args.id, status: "queued" }; } export function ingestDocumentJobId(documentId: string): string { return `ingest:document:${documentId}`; } export function ingestMindmapJobId(docId: string, mindmapId: string): string { return `ingest:mindmap:${docId}:${mindmapId}`; } export function ingestMediaAssetJobId(assetId: string): string { return `ingest:media:${assetId}`; } export async function enqueueIngestDocumentJob( ctx: MutationCtx, args: { userId: string; documentId: string; debounceMs?: number }, ): Promise<{ id: string }> { const id = ingestDocumentJobId(args.documentId); await upsertJob(ctx, { id, userId: args.userId, type: "ingest.rag_index_document", payload: { documentId: args.documentId }, debounceMs: args.debounceMs, }); return { id }; } export async function enqueueIngestMindmapJob( ctx: MutationCtx, args: { userId: string; docId: string; mindmapId: string; debounceMs?: number }, ): Promise<{ id: string }> { const id = ingestMindmapJobId(args.docId, args.mindmapId); await upsertJob(ctx, { id, userId: args.userId, type: "ingest.rag_index_mindmap", payload: { docId: args.docId, mindmapId: args.mindmapId }, debounceMs: args.debounceMs, }); return { id }; } export async function enqueueIngestMediaAssetJob( ctx: MutationCtx, args: { userId: string; assetId: string; debounceMs?: number }, ): Promise<{ id: string }> { const id = ingestMediaAssetJobId(args.assetId); await upsertJob(ctx, { id, userId: args.userId, type: "ingest.rag_index_media_asset", payload: { assetId: args.assetId }, debounceMs: args.debounceMs, }); return { id }; }