diff --git a/packages/trpc/src/index.ts b/packages/trpc/src/index.ts index 50acc58..0ede111 100644 --- a/packages/trpc/src/index.ts +++ b/packages/trpc/src/index.ts @@ -38,3 +38,7 @@ export { buildSearchQuery, } from "./services/scrape.service.js"; export type { BangumiSearchHit } from "./services/scrape.service.js"; + +// Playback / danmaku +export { playbackService } from "./services/playback.service.js"; +export { danmakuService } from "./services/danmaku.service.js"; diff --git a/packages/trpc/src/services/danmaku.service.ts b/packages/trpc/src/services/danmaku.service.ts new file mode 100644 index 0000000..e8ba06f --- /dev/null +++ b/packages/trpc/src/services/danmaku.service.ts @@ -0,0 +1,114 @@ +import { createHash } from "node:crypto"; +import { danmakuCacheDao, mediaItemDao, mountDao } from "@app/dao"; +import type { DanmakuFetchOutput } from "@app/types"; +import { decryptSecret } from "./secret.js"; +import { createWebdav } from "./webdav-client.js"; + +const CACHE_TTL_MS = 24 * 60 * 60 * 1000; + +function openApiBase(): string { + return process.env["OPEN_DANMAKU_API_BASE"] ?? "https://api.dandanplay.net"; +} + +function matchKeyFor(path: string, size: number): string { + return createHash("sha256").update(`${path}:${size}`).digest("hex"); +} + +/** 取 WebDAV 文件 hash(首 64KiB + 尺寸),用于开放网络识别 */ +async function remoteSampleHash(userId: string, mediaId: string): Promise { + 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 = mount.rootPath.replace(/\/+$/, "") + item.path; + const stream = client.createReadStream(abs, { range: { start: 0, end: 65535 } }); + const sample = await new Promise((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 createHash("sha256") + .update(sample) + .update(String(item.size)) + .digest("hex") + .toUpperCase(); + } catch { + return null; + } +} + +async function fetchOpenNetworkXml( + path: string, + size: number, + hash: string | null, +): Promise { + const base = openApiBase(); + try { + const identifyRes = await fetch(`${base}/api/v2/identify`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + fileHash: hash ?? "", + fileSize: size, + fileName: path.split("/").pop() ?? path, + }), + }); + if (!identifyRes.ok) return null; + const identified = (await identifyRes.json()) as { + isMatched?: boolean; + matched?: boolean; + episodeId?: number; + }; + const episodeId = identified.episodeId; + if (!episodeId || !(identified.isMatched || identified.matched)) return null; + const commentRes = await fetch(`${base}/api/v2/comment/${episodeId}?withResponses=false`); + if (!commentRes.ok) return null; + const body = (await commentRes.json()) as { danmaku?: string; xml?: string }; + return body.danmaku ?? body.xml ?? null; + } catch { + return null; + } +} + +export const danmakuService = { + async fetch(userId: string, mediaItemId: string): Promise { + const item = await mediaItemDao.getByIdForUser(mediaItemId, userId); + if (!item) return { ok: false, source: "none", xml: "", message: "媒体不存在" }; + + const matchKey = matchKeyFor(item.path, item.size); + const cached = await danmakuCacheDao.getValid(matchKey); + if (cached) { + return { ok: true, source: "cache", xml: cached.payload }; + } + + const hash = await remoteSampleHash(userId, mediaItemId); + const xml = await fetchOpenNetworkXml(item.path, item.size, hash); + if (!xml) { + return { ok: true, source: "none", xml: "", message: "开放网络未匹配到弹幕" }; + } + await danmakuCacheDao.upsert({ + matchKey, + payload: xml, + source: "open-network", + expiresAt: new Date(Date.now() + CACHE_TTL_MS), + }); + return { ok: true, source: "open-network", xml }; + }, + + /** 本地 XML 由前端解析后直接喂播放器;此处仅确认媒体存在并记 size 日志约束 */ + async importMeta(userId: string, mediaItemId: string, byteSize: number): Promise<{ ok: true }> { + const item = await mediaItemDao.getByIdForUser(mediaItemId, userId); + if (!item) return Promise.reject(new Error("媒体不存在")); + if (byteSize > 20_000_000) return Promise.reject(new Error("XML 过大")); + return { ok: true }; + }, +}; diff --git a/packages/trpc/src/services/playback.service.ts b/packages/trpc/src/services/playback.service.ts new file mode 100644 index 0000000..6d3291d --- /dev/null +++ b/packages/trpc/src/services/playback.service.ts @@ -0,0 +1,29 @@ +import { playbackProgressDao } from "@app/dao"; +import type { PlaybackReportInput } from "@app/types"; + +export const playbackService = { + async report(userId: string, input: PlaybackReportInput) { + const row = await playbackProgressDao.upsert( + userId, + input.mediaItemId, + input.positionMs, + input.durationMs, + ); + return { + mediaItemId: row.mediaItemId, + positionMs: row.positionMs, + durationMs: row.durationMs, + updatedAt: row.updatedAt, + }; + }, + async get(userId: string, mediaItemId: string) { + const row = await playbackProgressDao.get(userId, mediaItemId); + if (!row) return { mediaItemId, positionMs: 0, durationMs: 0, updatedAt: null }; + return { + mediaItemId: row.mediaItemId, + positionMs: row.positionMs, + durationMs: row.durationMs, + updatedAt: row.updatedAt, + }; + }, +};