From c37f9dcf625b487276c6d29c51a1376ec109fb1b Mon Sep 17 00:00:00 2001 From: noelorin Date: Fri, 25 Sep 2026 02:56:23 +0800 Subject: [PATCH] fix(trpc): bound scrape concurrency and surface DB errors --- packages/trpc/src/services/library.service.ts | 38 +++++++++++-------- packages/trpc/src/services/scrape.service.ts | 37 +++++++++--------- 2 files changed, 42 insertions(+), 33 deletions(-) diff --git a/packages/trpc/src/services/library.service.ts b/packages/trpc/src/services/library.service.ts index 09a481a..da4c1c3 100644 --- a/packages/trpc/src/services/library.service.ts +++ b/packages/trpc/src/services/library.service.ts @@ -5,6 +5,10 @@ import { createWebdav, isVideoFilename, joinWebdavPath, listDirectory } from "./ import { scrapeService } from "./scrape.service.js"; const MAX_SCAN_FILES = 2000; +/** 同时执行的 upsert+刮削 job 上限,避免最坏 2000 并发打 Bangumi 被限流 */ +const SCRAPE_CONCURRENCY = 4; + +type ScanJob = () => Promise; async function walkVideos( userId: string, @@ -12,7 +16,7 @@ async function walkVideos( root: string, rel: string, budget: { left: number }, - out: Promise[], + out: ScanJob[], ): Promise { if (budget.left <= 0) return; const mount = await mountDao.getByIdForUser(mountId, userId); @@ -38,19 +42,18 @@ async function walkVideos( if (!isVideoFilename(entry.basename)) continue; budget.left -= 1; const title = entry.basename.replace(/\.[a-z0-9]+$/i, ""); - out.push( - (async () => { - const item = await mediaItemDao.upsertFromScan({ - userId, - mountId, - path: childRel, - rawName: entry.basename, - title, - size: entry.size, - }); - await scrapeService.autoScrape(userId, item.id); - })(), - ); + // 延迟执行:扫描阶段只收集 job,随后分批限流 + out.push(async () => { + const item = await mediaItemDao.upsertFromScan({ + userId, + mountId, + path: childRel, + rawName: entry.basename, + title, + size: entry.size, + }); + await scrapeService.autoScrape(userId, item.id); + }); } } @@ -78,9 +81,12 @@ export const libraryService = { if (!mount) throw new TRPCError({ code: "NOT_FOUND", message: "挂载不存在" }); if (!mount.enabled) throw new TRPCError({ code: "BAD_REQUEST", message: "挂载已禁用" }); const budget = { left: MAX_SCAN_FILES }; - const jobs: Promise[] = []; + const jobs: ScanJob[] = []; await walkVideos(userId, mountId, mount.rootPath || "/", "/", budget, jobs); - await Promise.allSettled(jobs); + for (let i = 0; i < jobs.length; i += SCRAPE_CONCURRENCY) { + const batch = jobs.slice(i, i + SCRAPE_CONCURRENCY).map((job) => job()); + await Promise.allSettled(batch); + } return { scanned: MAX_SCAN_FILES - budget.left, message: "扫描完成" }; }, }; diff --git a/packages/trpc/src/services/scrape.service.ts b/packages/trpc/src/services/scrape.service.ts index fbc2676..eb8e6bf 100644 --- a/packages/trpc/src/services/scrape.service.ts +++ b/packages/trpc/src/services/scrape.service.ts @@ -104,29 +104,32 @@ export const scrapeService = { }); return; } + // 只包网络搜索:DB 写失败必须向外抛,避免被误标 failed 后吞掉 + let hits: BangumiSearchHit[]; try { - const hits = await bangumiSearch(q); - const hit = hits[0]; - if (!hit) { - await mediaItemDao.updateScrape(mediaItemId, userId, { - scrapeStatus: "unmatched", - scrapedAt: new Date(), - }); - return; - } - await mediaItemDao.updateScrape(mediaItemId, userId, { - bangumiId: hit.id, - scrapeStatus: "ok", - scrapedAt: new Date(), - posterUrl: hit.image || null, - title: hit.nameCn || hit.name, - epNumber: parseEpisodeFromFilename(row.rawName), - }); + hits = await bangumiSearch(q); } catch { await mediaItemDao.updateScrape(mediaItemId, userId, { scrapeStatus: "failed", scrapedAt: new Date(), }); + return; } + const hit = hits[0]; + if (!hit) { + await mediaItemDao.updateScrape(mediaItemId, userId, { + scrapeStatus: "unmatched", + scrapedAt: new Date(), + }); + return; + } + await mediaItemDao.updateScrape(mediaItemId, userId, { + bangumiId: hit.id, + scrapeStatus: "ok", + scrapedAt: new Date(), + posterUrl: hit.image || null, + title: hit.nameCn || hit.name, + epNumber: parseEpisodeFromFilename(row.rawName), + }); }, };