fix(trpc): bound scrape concurrency and surface DB errors

This commit is contained in:
noelorin 2026-09-25 02:56:23 +08:00
parent 94d389b855
commit c37f9dcf62
2 changed files with 42 additions and 33 deletions

View File

@ -5,6 +5,10 @@ import { createWebdav, isVideoFilename, joinWebdavPath, listDirectory } from "./
import { scrapeService } from "./scrape.service.js"; import { scrapeService } from "./scrape.service.js";
const MAX_SCAN_FILES = 2000; const MAX_SCAN_FILES = 2000;
/** 同时执行的 upsert+刮削 job 上限,避免最坏 2000 并发打 Bangumi 被限流 */
const SCRAPE_CONCURRENCY = 4;
type ScanJob = () => Promise<void>;
async function walkVideos( async function walkVideos(
userId: string, userId: string,
@ -12,7 +16,7 @@ async function walkVideos(
root: string, root: string,
rel: string, rel: string,
budget: { left: number }, budget: { left: number },
out: Promise<void>[], out: ScanJob[],
): Promise<void> { ): Promise<void> {
if (budget.left <= 0) return; if (budget.left <= 0) return;
const mount = await mountDao.getByIdForUser(mountId, userId); const mount = await mountDao.getByIdForUser(mountId, userId);
@ -38,8 +42,8 @@ async function walkVideos(
if (!isVideoFilename(entry.basename)) continue; if (!isVideoFilename(entry.basename)) continue;
budget.left -= 1; budget.left -= 1;
const title = entry.basename.replace(/\.[a-z0-9]+$/i, ""); const title = entry.basename.replace(/\.[a-z0-9]+$/i, "");
out.push( // 延迟执行:扫描阶段只收集 job,随后分批限流
(async () => { out.push(async () => {
const item = await mediaItemDao.upsertFromScan({ const item = await mediaItemDao.upsertFromScan({
userId, userId,
mountId, mountId,
@ -49,8 +53,7 @@ async function walkVideos(
size: entry.size, size: entry.size,
}); });
await scrapeService.autoScrape(userId, item.id); await scrapeService.autoScrape(userId, item.id);
})(), });
);
} }
} }
@ -78,9 +81,12 @@ export const libraryService = {
if (!mount) throw new TRPCError({ code: "NOT_FOUND", message: "挂载不存在" }); if (!mount) throw new TRPCError({ code: "NOT_FOUND", message: "挂载不存在" });
if (!mount.enabled) throw new TRPCError({ code: "BAD_REQUEST", message: "挂载已禁用" }); if (!mount.enabled) throw new TRPCError({ code: "BAD_REQUEST", message: "挂载已禁用" });
const budget = { left: MAX_SCAN_FILES }; const budget = { left: MAX_SCAN_FILES };
const jobs: Promise<void>[] = []; const jobs: ScanJob[] = [];
await walkVideos(userId, mountId, mount.rootPath || "/", "/", budget, jobs); 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: "扫描完成" }; return { scanned: MAX_SCAN_FILES - budget.left, message: "扫描完成" };
}, },
}; };

View File

@ -104,8 +104,17 @@ export const scrapeService = {
}); });
return; return;
} }
// 只包网络搜索:DB 写失败必须向外抛,避免被误标 failed 后吞掉
let hits: BangumiSearchHit[];
try { try {
const hits = await bangumiSearch(q); hits = await bangumiSearch(q);
} catch {
await mediaItemDao.updateScrape(mediaItemId, userId, {
scrapeStatus: "failed",
scrapedAt: new Date(),
});
return;
}
const hit = hits[0]; const hit = hits[0];
if (!hit) { if (!hit) {
await mediaItemDao.updateScrape(mediaItemId, userId, { await mediaItemDao.updateScrape(mediaItemId, userId, {
@ -122,11 +131,5 @@ export const scrapeService = {
title: hit.nameCn || hit.name, title: hit.nameCn || hit.name,
epNumber: parseEpisodeFromFilename(row.rawName), epNumber: parseEpisodeFromFilename(row.rawName),
}); });
} catch {
await mediaItemDao.updateScrape(mediaItemId, userId, {
scrapeStatus: "failed",
scrapedAt: new Date(),
});
}
}, },
}; };