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

318 lines
12 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, listDirectory, type DirEntry } from "./webdav-client.js";
const CACHE_TTL_MS = 24 * 60 * 60 * 1000;
/** 同名侧车 XML 全文上限,与 importXml 一致 */
const SIDECAR_MAX_BYTES = 20_000_000;
/** 媒体路径 → 所在目录("" 表示挂载根)+ 去扩展名 basename */
function splitDirBase(path: string): { dir: string; base: string } {
const norm = path.replace(/\\/g, "/");
const idx = norm.lastIndexOf("/");
const dir = idx >= 0 ? norm.slice(0, idx) : "";
const file = idx >= 0 ? norm.slice(idx + 1) : norm;
const dot = file.lastIndexOf(".");
return { dir, base: dot > 0 ? file.slice(0, dot) : file };
}
/** 同 basename + .xml,扩展名与 stem 均大小写不敏感 */
function isSidecarName(fileName: string, base: string): boolean {
return fileName.toLowerCase() === `${base}.xml`.toLowerCase();
}
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 = {
/** 同目录同名 .xml 侧车 → 落 local-xml;失败/超限/非弹幕一律静默 false */
async ensureSidecarXml(userId: string, mediaItemId: string): Promise<boolean> {
try {
const item = await mediaItemDao.getByIdForUser(mediaItemId, userId);
if (!item?.mountId) return false;
const mount = await mountDao.getByIdForUser(item.mountId, userId);
if (!mount) return false;
const { dir, base } = splitDirBase(item.path);
if (!base) return false;
const client = createWebdav({
baseUrl: mount.baseUrl,
username: mount.username ?? undefined,
password: mount.secretEnc ? decryptSecret(mount.secretEnc) : undefined,
});
const listAbs = joinWebdavPath(mount.rootPath || "/", dir || "/");
const entries = await listDirectory(client, listAbs);
const hit = entries.find(
(e: DirEntry) => e.type === "file" && isSidecarName(e.basename, base),
);
if (!hit || hit.size > SIDECAR_MAX_BYTES) return false;
const raw = (await client.getFileContents(hit.filename, {
format: "binary",
})) as Uint8Array;
const buf = Buffer.isBuffer(raw) ? raw : Buffer.from(raw);
if (buf.byteLength > SIDECAR_MAX_BYTES) return false;
const xml = buf.toString("utf8");
if (!xml.includes("<d ")) return false;
await danmakuDocsDao.upsert({
userId,
mediaItemId,
source: "local-xml",
xml,
byteSize: Buffer.byteLength(xml, "utf8"),
});
return true;
} catch {
return false;
}
},
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 };
// 同目录同名 .xml 侧车自动关联;命中即按 local-xml 返回
if (await danmakuService.ensureSidecarXml(userId, mediaItemId)) {
const linked = await danmakuDocsDao.get(userId, mediaItemId, "local-xml");
return {
ok: true,
source: "local-xml",
xml: linked?.xml ?? "",
message: "同名弹幕",
};
}
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);
},
};