import { initTRPC } from "@trpc/server"; import { EventEmitter, on } from "node:events"; import { getConfig } from "./config"; import { errorMessage } from "@/lib/utils/error"; import { PublishInputSchema, type ProgressEvent, type PlatformResult, type PublishResult, platforms, type PlatformContext, type PlatformAdaptor, type PublishInput, PublishResultSchema, ProgressEventSchema, } from "./publish.schemas"; import { publishModrinthVersion } from "./platforms/modrinth"; import { publishCurseForgeVersion } from "./platforms/curseforge"; import { computeVersionRanges, type McVersionEntry } from "@/lib/utils/mcVersion"; import { getCurseForgeMcVersions, getModrinthMcVersions } from "./meta"; import { readAsFile } from "@/lib/utils/file"; import { Template } from "@/lib/utils/template"; import { buildJarPath } from "@/lib/utils/format"; export const platformAdaptors: { [key in typeof platforms[number]]: PlatformAdaptor } = { modrinth: { publish: publishModrinthVersion, getGameVersions: getModrinthMcVersions, }, curseforge: { publish: publishCurseForgeVersion, getGameVersions: getCurseForgeMcVersions, }, }; // ── tRPC init ────────────────────────────────────────── const t = initTRPC.create(); // ── Shared EventEmitter for progress ─────────────────── const ee = new EventEmitter(); // ── Platform orchestrator ────────────────────────────── async function runPlatform( platform: typeof platforms[number], mcVersions: PublishInput["mcVersions"][keyof PublishInput["mcVersions"]], ctx: PlatformContext, ): Promise { const result: PlatformResult = { success: 0, fail: 0, errors: [] }; const total = mcVersions.length; if (total === 0) { ee.emit("progress", { platform, current: 0, total: 0, status: "completed", errors: [], } satisfies ProgressEvent); return result; } // Emit initial running state ee.emit("progress", { platform, current: 0, total, status: "running", errors: [], } satisfies ProgressEvent); const publishFn = platformAdaptors[platform].publish; const getGameVersionsFn = platformAdaptors[platform].getGameVersions; // 平台准备阶段(元数据拉取 + 版本范围计算)。此处异常不在逐版本的 // try/catch 内,必须单独兜底:记为该平台全部失败并发出 failed 事件, // 否则前端进度条会停在 running 假象里,错误也被吞掉。 let mcVersionEntries: McVersionEntry[]; try { mcVersionEntries = computeVersionRanges( mcVersions.map(v => v.mc_version), await getGameVersionsFn(), ctx.input.cutoffMcVersion ); } catch (err) { const msg = errorMessage(err); console.error(`[publish] ${platform} 准备阶段失败:`, err); result.fail = total; result.errors.push(msg); ee.emit("progress", { platform, current: 0, total, status: "failed", errors: result.errors, } satisfies ProgressEvent); return result; } const primaryFileTemplate = new Template(ctx.config.filename_format); const sourceFileTemplate = new Template(ctx.config.source_filename_format); for (let i = 0; i < mcVersionEntries.length; i++) { const mcVersionEntry = mcVersionEntries[i]; try { const primaryFile = await readAsFile(buildJarPath(ctx.config.project_dir, primaryFileTemplate, ctx.input.version, mcVersionEntry.version)); const sourceFile = await readAsFile(buildJarPath(ctx.config.project_dir, sourceFileTemplate, ctx.input.version, mcVersionEntry.version)); await publishFn(ctx, mcVersionEntry, { primary: primaryFile, source: sourceFile, }); result.success++; const isLast = i + 1 === total; ee.emit("progress", { platform, current: i + 1, total, status: isLast ? "completed" : "running", errors: result.errors, } satisfies ProgressEvent); } catch (err) { result.fail++; const msg = `${mcVersionEntry.version}: ${errorMessage(err)}`; result.errors.push(msg); ee.emit("progress", { platform, current: i, total, status: "failed", errors: result.errors, } satisfies ProgressEvent); // Cancel remaining for this platform only break; } } return result; } // ── Router ───────────────────────────────────────────── export const appRouter = t.router({ publish: t.procedure .input(PublishInputSchema) .mutation(async ({ input }): Promise => { const config = await getConfig(input.configFile); const ctx: PlatformContext = { input, config }; const results = await Promise.all(platforms.map(async (platform) => { try { const result = await runPlatform(platform, input.mcVersions[platform], ctx); return {platform, result}; } catch (err) { console.error(`[publish] ${platform} 执行异常:`, err); return { platform, result: { success: 0, fail: 0, errors: [errorMessage(err)] } } } })); return PublishResultSchema.parse(Object.fromEntries(results.map(result => [result.platform, result.result]))); }), progress: t.procedure.subscription(async function* (opts) { for await (const [data] of on(ee, "progress", { signal: opts.signal, })) { const parsed = ProgressEventSchema.safeParse(data); if (!parsed.success) { // 坏事件不应中断订阅(客户端没有重连机制),跳过并记录 console.error("[publish] invalid progress event:", parsed.error); continue; } yield parsed.data; } }), }); export type AppRouter = typeof appRouter;