diff --git a/app/web/src/pages/api/stream.ts b/app/web/src/pages/api/stream.ts new file mode 100644 index 0000000..7885e74 --- /dev/null +++ b/app/web/src/pages/api/stream.ts @@ -0,0 +1,31 @@ +/** 认证后的 Range 流代理:鉴权在本路由,webdav 细节封装在 @app/trpc 的 openRemoteFileStream。 */ +import type { APIRoute } from "astro"; +import { COOKIE_SESSION } from "@app/types"; +import { authService, openRemoteFileStream } from "@app/trpc"; + +export const prerender = false; + +export const GET: APIRoute = async ({ request, cookies, url }) => { + const token = cookies.get(COOKIE_SESSION)?.value; + const user = await authService.getSessionUser(token); + if (!user) { + return new Response("Unauthorized", { status: 401 }); + } + const mediaItemId = url.searchParams.get("id"); + if (!mediaItemId) return new Response("id required", { status: 400 }); + + try { + const result = await openRemoteFileStream( + user.id, + mediaItemId, + request.headers.get("range"), + ); + return new Response(result.body, { + status: result.status, + headers: result.headers, + }); + } catch (e) { + const msg = e instanceof Error ? e.message : "stream failed"; + return new Response(msg, { status: 502 }); + } +}; diff --git a/packages/trpc/src/index.ts b/packages/trpc/src/index.ts index 0eb77f4..00bfbe3 100644 --- a/packages/trpc/src/index.ts +++ b/packages/trpc/src/index.ts @@ -43,3 +43,7 @@ export type { BangumiSearchHit } from "./services/scrape.service.js"; // Playback / danmaku export { playbackService } from "./services/playback.service.js"; export { danmakuService } from "./services/danmaku.service.js"; + +// Stream range proxy +export { openRemoteFileStream } from "./services/stream.service.js"; +export type { RemoteFileStream } from "./services/stream.service.js"; diff --git a/packages/trpc/src/services/stream.service.ts b/packages/trpc/src/services/stream.service.ts new file mode 100644 index 0000000..a9c81b2 --- /dev/null +++ b/packages/trpc/src/services/stream.service.ts @@ -0,0 +1,71 @@ +import { mediaItemDao, mountDao } from "@app/dao"; +import { decryptSecret } from "./secret.js"; +import { createWebdav } from "./webdav-client.js"; + +export type RemoteFileStream = { + status: number; + headers: Headers; + body: ReadableStream | null; +}; + +/** 打开远端媒体的 Range 流;webdav 细节只在本层出现,web 只拿 Headers + ReadableStream。 */ +export async function openRemoteFileStream( + userId: string, + mediaItemId: string, + rangeHeader: string | null, +): Promise { + const item = await mediaItemDao.getByIdForUser(mediaItemId, userId); + if (!item?.mountId) throw new Error("媒体或挂载不存在"); + const mount = await mountDao.getByIdForUser(item.mountId, userId); + if (!mount) throw new Error("挂载不存在"); + + const client = createWebdav({ + baseUrl: mount.baseUrl, + username: mount.username ?? undefined, + password: mount.secretEnc ? decryptSecret(mount.secretEnc) : undefined, + }); + const abs = mount.rootPath.replace(/\/+$/, "") + item.path; + + const headers = new Headers(); + headers.set("Content-Type", item.mime || "application/octet-stream"); + headers.set("Accept-Ranges", "bytes"); + headers.set("Cache-Control", "private, max-age=0"); + + let status = 200; + let start = 0; + let end = item.size > 0 ? item.size - 1 : 0; + + if (rangeHeader && item.size > 0) { + const m = /bytes=(\d*)-(\d*)/.exec(rangeHeader); + if (m) { + if (m[1]) start = Number(m[1]); + if (m[2]) end = Math.min(Number(m[2]), item.size - 1); + status = 206; + headers.set("Content-Range", `bytes ${start}-${end}/${item.size}`); + headers.set("Content-Length", String(end - start + 1)); + } + } else if (item.size > 0) { + headers.set("Content-Length", String(item.size)); + } + + // createReadStream 同步返回 ReadableLike,range 选项嵌套在 { range } 内 + const nodeStream = + item.size > 0 + ? client.createReadStream(abs, { range: { start, end } }) + : client.createReadStream(abs); + + const body = new ReadableStream({ + start(controller) { + nodeStream.on("data", (chunk: Buffer | string) => { + controller.enqueue(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + }); + nodeStream.on("end", () => controller.close()); + nodeStream.on("error", (err: unknown) => controller.error(err)); + }, + cancel() { + nodeStream.destroy(); + }, + }); + + return { status, headers, body }; +}