app-template/packages/trpc/src/services/danmaku.service.ts

250 lines
9.3 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import { createHash } from "node:crypto";
import { TRPCError } from "@trpc/server";
import { danmakuCacheDao, danmakuDocsDao, danmakuPrefsDao, mediaItemDao, mountDao } from "@app/dao";
import {
HASH_HEAD_BYTES,
hasOpenCredentials,
headMd5,
matchFileName,
openCommentXml,
openMatch,
rankCandidates,
} from "@app/danmaku";
import type { DanmakuFetchOutput, DanmakuMatchOutput, DanmakuSettingsOutput } from "@app/types";
import { decryptSecret } from "./secret.js";
import { createWebdav, joinWebdavPath } from "./webdav-client.js";
const CACHE_TTL_MS = 24 * 60 * 60 * 1000;
function matchKeyFor(path: string, size: number): string {
return createHash("sha256").update(`${path}:${size}`).digest("hex");
}
/** 取 WebDAV 文件前 16MB 的 MD5(开放网络文件识别用) */
export async function remoteFileHash(userId: string, mediaId: string): Promise<string | null> {
const item = await mediaItemDao.getByIdForUser(mediaId, userId);
if (!item?.mountId) return null;
const mount = await mountDao.getByIdForUser(item.mountId, userId);
if (!mount) return null;
const client = createWebdav({
baseUrl: mount.baseUrl,
username: mount.username ?? undefined,
password: mount.secretEnc ? decryptSecret(mount.secretEnc) : undefined,
});
try {
const abs = joinWebdavPath(mount.rootPath || "/", item.path);
const stream = client.createReadStream(abs, {
range: { start: 0, end: HASH_HEAD_BYTES - 1 },
});
const sample = await new Promise<Buffer>((resolve, reject) => {
const chunks: Buffer[] = [];
stream.on("data", (chunk: Buffer) => {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk)));
});
stream.on("end", () => resolve(Buffer.concat(chunks)));
stream.on("error", reject);
});
return headMd5(sample.subarray(0, HASH_HEAD_BYTES));
} catch {
return null;
}
}
async function persistOpenXml(userId: string, mediaItemId: string, xml: string): Promise<void> {
await danmakuDocsDao.upsert({
userId,
mediaItemId,
source: "open-network",
xml,
byteSize: Buffer.byteLength(xml, "utf8"),
});
}
async function persistOpenCache(episodeId: number, xml: string): Promise<void> {
await danmakuCacheDao.upsert({
matchKey: String(episodeId),
payload: xml,
source: "open-network",
expiresAt: new Date(Date.now() + CACHE_TTL_MS),
});
}
export const danmakuService = {
async fetch(userId: string, mediaItemId: string): Promise<DanmakuFetchOutput> {
const item = await mediaItemDao.getByIdForUser(mediaItemId, userId);
if (!item) return { ok: false, source: "none", xml: "", message: "媒体不存在" };
const local = await danmakuDocsDao.get(userId, mediaItemId, "local-xml");
if (local) return { ok: true, source: "local-xml", xml: local.xml };
if (!hasOpenCredentials()) {
return {
ok: true,
source: "none",
xml: "",
message:
"未配置开放弹幕网络:服务端需设置 OPEN_DANMAKU_APP_ID / OPEN_DANMAKU_APP_SECRET",
};
}
const cacheKey =
item.dandanplayEpisodeId != null
? String(item.dandanplayEpisodeId)
: matchKeyFor(item.path, item.size);
const cached = await danmakuCacheDao.getValid(cacheKey);
if (cached) return { ok: true, source: "cache", xml: cached.payload };
let episodeId = item.dandanplayEpisodeId ?? null;
if (episodeId == null) {
const hash = await remoteFileHash(userId, mediaItemId);
const outcome = await openMatch(matchFileName(item.path), item.size, hash);
const ranked = rankCandidates(outcome.candidates);
if (!outcome.ok || !outcome.isMatched || ranked.candidates.length === 0) {
return {
ok: true,
source: "none",
xml: "",
message: outcome.errorMessage ?? "开放网络未匹配到弹幕",
};
}
// 仅高置信才自动绑;否则等手选(fetch 不落库)
if (!ranked.highConfidence) {
return {
ok: true,
source: "none",
xml: "",
message: "匹配不明确,请在播放页手选弹幕集",
};
}
episodeId = ranked.candidates[0]?.episodeId ?? null;
if (episodeId == null) {
return { ok: true, source: "none", xml: "", message: "开放网络未匹配到弹幕" };
}
await mediaItemDao.setDanmakuMatch(mediaItemId, userId, {
episodeId,
matchedHash: hash,
source: "auto",
});
}
const xml = await openCommentXml(episodeId);
if (xml === null) {
return { ok: true, source: "none", xml: "", message: "弹幕库拉取失败" };
}
await persistOpenXml(userId, mediaItemId, xml);
await persistOpenCache(episodeId, xml);
return { ok: true, source: "open-network", xml };
},
async matchCandidates(userId: string, mediaItemId: string): Promise<DanmakuMatchOutput> {
const item = await mediaItemDao.getByIdForUser(mediaItemId, userId);
if (!item) {
return {
ok: false,
matched: false,
message: "媒体不存在",
candidates: [],
highConfidence: false,
};
}
if (!hasOpenCredentials()) {
return {
ok: false,
matched: false,
message: "未配置开放弹幕网络",
candidates: [],
highConfidence: false,
};
}
const hash = await remoteFileHash(userId, mediaItemId);
const outcome = await openMatch(matchFileName(item.path), item.size, hash);
const ranked = rankCandidates(outcome.candidates);
return {
ok: outcome.ok,
matched: outcome.isMatched && ranked.candidates.length > 0,
message: outcome.errorMessage ?? undefined,
candidates: ranked.candidates,
highConfidence: ranked.highConfidence,
};
},
async selectMatch(
userId: string,
mediaItemId: string,
episodeId: number,
): Promise<DanmakuFetchOutput> {
const item = await mediaItemDao.getByIdForUser(mediaItemId, userId);
if (!item) return { ok: false, source: "none", xml: "", message: "媒体不存在" };
await mediaItemDao.setDanmakuMatch(mediaItemId, userId, {
episodeId,
matchedHash: item.matchedHash,
source: "manual",
});
const xml = await openCommentXml(episodeId);
if (xml === null) {
return { ok: true, source: "none", xml: "", message: "弹幕库拉取失败" };
}
await persistOpenXml(userId, mediaItemId, xml);
await persistOpenCache(episodeId, xml);
return { ok: true, source: "open-network", xml };
},
async importXml(
userId: string,
mediaItemId: string,
xml: string,
): Promise<{ ok: true; count: number }> {
const item = await mediaItemDao.getByIdForUser(mediaItemId, userId);
if (!item) throw new TRPCError({ code: "NOT_FOUND", message: "媒体不存在" });
if (xml.length > 20_000_000) {
throw new TRPCError({ code: "BAD_REQUEST", message: "XML 过大" });
}
if (!xml.includes("<d ")) {
throw new TRPCError({ code: "BAD_REQUEST", message: "不是有效的弹幕 XML" });
}
await danmakuDocsDao.upsert({
userId,
mediaItemId,
source: "local-xml",
xml,
byteSize: Buffer.byteLength(xml, "utf8"),
});
return { ok: true, count: (xml.match(/<d /g) ?? []).length };
},
async getPrefs(userId: string): Promise<DanmakuSettingsOutput> {
const row = await danmakuPrefsDao.getByUser(userId);
return {
enabled: row?.enabled ?? true,
opacity: (row?.opacity ?? 80) / 100,
density: (row?.density ?? 100) / 100,
blockKeywords: row ? (JSON.parse(row.blockKeywords) as string[]) : [],
blockTypes: row
? (JSON.parse(row.blockTypes) as Array<"scroll" | "top" | "bottom">)
: [],
openNetworkConfigured: hasOpenCredentials(),
};
},
async savePrefs(
userId: string,
input: {
enabled: boolean;
opacity: number;
density: number;
blockKeywords: string[];
blockTypes: Array<"scroll" | "top" | "bottom">;
},
): Promise<DanmakuSettingsOutput> {
await danmakuPrefsDao.upsert({
userId,
enabled: input.enabled,
opacity: Math.round(input.opacity * 100),
density: Math.round(input.density * 100),
blockKeywords: JSON.stringify(input.blockKeywords),
blockTypes: JSON.stringify(input.blockTypes),
});
return danmakuService.getPrefs(userId);
},
};