145 lines
4.1 KiB
TypeScript
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 };
|
|
}
|