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

145 lines
4.1 KiB
TypeScript

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 function extractMediaAssetTextJobId(assetId: string): string {
return `extract: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 };
}
export async function enqueueExtractMediaAssetTextJob(
ctx: MutationCtx,
args: { userId: string; assetId: string; debounceMs?: number },
): Promise<{ id: string }> {
const id = extractMediaAssetTextJobId(args.assetId);
await upsertJob(ctx, {
id,
userId: args.userId,
type: "extract.media_asset_text",
payload: { assetId: args.assetId },
debounceMs: args.debounceMs,
});
return { id };
}