import { and, asc, desc, eq, inArray, isNotNull, like, or, sql } from "drizzle-orm"; import { mediaItems, mounts, type MediaItem } from "@app/models"; import { db } from "@app/db"; export type MediaItemRow = MediaItem; export const mediaItemDao = { async list( userId: string, params: { limit: number; offset: number; mountId?: string | undefined; scrapeStatus?: string | undefined; q?: string | undefined; }, ): Promise<{ rows: MediaItemRow[]; total: number }> { const conds = [eq(mediaItems.userId, userId)]; if (params.mountId) conds.push(eq(mediaItems.mountId, params.mountId)); if (params.scrapeStatus) conds.push(eq(mediaItems.scrapeStatus, params.scrapeStatus)); if (params.q) { const likeQ = `%${params.q}%`; const qCond = or(like(mediaItems.title, likeQ), like(mediaItems.rawName, likeQ)); if (qCond) conds.push(qCond); } const where = and(...conds); const rows = await db .select() .from(mediaItems) .where(where) .orderBy(desc(mediaItems.updatedAt)) .limit(params.limit) .offset(params.offset); const [countRow] = await db .select({ count: sql`count(*)` }) .from(mediaItems) .where(where); return { rows, total: Number(countRow?.["count"] ?? 0) }; }, /** * 海报墙分组分页:指定挂载下已刮削条目按 bangumiId 分组。 * q 命中任一集即收入该组,但 episodeCount 始终为该组全量集数(与列表筛选「已刮削」对得上)。 */ async listScrapedGroups( userId: string, params: { mountId: string; q?: string | undefined; limit: number; offset: number; }, ): Promise<{ rows: { bangumiId: number; title: string; posterUrl: string | null; episodeCount: number; updatedAt: Date; }[]; total: number; }> { const baseConds = [ eq(mediaItems.userId, userId), eq(mediaItems.mountId, params.mountId), eq(mediaItems.scrapeStatus, "ok"), isNotNull(mediaItems.bangumiId), ]; const fullWhere = and(...baseConds); // 先找出 q 命中的作品 id;命中集合决定「有哪些组」,不缩小组内集数 let matchIds: number[] | null = null; if (params.q) { const likeQ = `%${params.q}%`; const qCond = or(like(mediaItems.title, likeQ), like(mediaItems.rawName, likeQ)); const matched = await db .select({ bangumiId: mediaItems.bangumiId }) .from(mediaItems) .where(and(...baseConds, qCond)) .groupBy(mediaItems.bangumiId); matchIds = matched.map((r) => r.bangumiId).filter((x): x is number => x != null); if (matchIds.length === 0) return { rows: [], total: 0 }; } const scopeWhere = matchIds == null ? fullWhere : and(fullWhere, inArray(mediaItems.bangumiId, matchIds)); // rn=1 即每组 updated_at 最新的一行,其 title/posterUrl 就代表该组 const ranked = db .select({ bangumiId: mediaItems.bangumiId, title: mediaItems.title, posterUrl: mediaItems.posterUrl, updatedAt: mediaItems.updatedAt, rn: sql`row_number() over (partition by ${mediaItems.bangumiId} order by ${mediaItems.updatedAt} desc, ${mediaItems.id} desc)`.as( "rn", ), episodeCount: sql`count(*) over (partition by ${mediaItems.bangumiId})`.as( "episode_count", ), }) .from(mediaItems) .where(scopeWhere) .as("ranked"); const pageRows = await db .select({ bangumiId: ranked.bangumiId, title: ranked.title, posterUrl: ranked.posterUrl, episodeCount: ranked.episodeCount, updatedAt: ranked.updatedAt, }) .from(ranked) .where(eq(ranked.rn, 1)) .orderBy(desc(ranked.updatedAt), desc(ranked.bangumiId)) .limit(params.limit) .offset(params.offset); const total = matchIds != null ? matchIds.length : await this.countScrapedSeries(userId, params.mountId); const rows: { bangumiId: number; title: string; posterUrl: string | null; episodeCount: number; updatedAt: Date; }[] = []; for (const row of pageRows) { const bid = row.bangumiId; if (bid == null) continue; rows.push({ bangumiId: bid, title: row.title, posterUrl: row.posterUrl, episodeCount: Number(row.episodeCount), updatedAt: row.updatedAt, }); } return { rows, total }; }, async countScrapedSeries(userId: string, mountId: string): Promise { const [countRow] = await db .select({ count: sql`count(distinct ${mediaItems.bangumiId})` }) .from(mediaItems) .where( and( eq(mediaItems.userId, userId), eq(mediaItems.mountId, mountId), eq(mediaItems.scrapeStatus, "ok"), isNotNull(mediaItems.bangumiId), ), ); return Number(countRow?.["count"] ?? 0); }, /** 路径前缀下的条目(文件夹级批量刮削/指定用);精确目录边界,避免 /Foo 误吞 /FooBar */ async listByPathPrefix(userId: string, mountId: string, path: string): Promise { const normPrefix = path.replace(/\/+$/, "") || "/"; const base = [eq(mediaItems.userId, userId), eq(mediaItems.mountId, mountId)]; if (normPrefix === "/") { return db .select() .from(mediaItems) .where(and(...base)) .orderBy(asc(mediaItems.path)); } return db .select() .from(mediaItems) .where( and( ...base, or(eq(mediaItems.path, normPrefix), like(mediaItems.path, `${normPrefix}/%`)), ), ) .orderBy(asc(mediaItems.path)); }, async getByIdForUser(id: string, userId: string): Promise { const [row] = await db .select() .from(mediaItems) .where(and(eq(mediaItems.id, id), eq(mediaItems.userId, userId))) .limit(1); return row ?? null; }, /** 同作品(bangumiId)剧集,按集数升序;mountId 可选,库内详情只看该挂载 */ async listByBangumiId( userId: string, bangumiId: number, mountId?: string | undefined, ): Promise { const conds = [eq(mediaItems.userId, userId), eq(mediaItems.bangumiId, bangumiId)]; if (mountId) conds.push(eq(mediaItems.mountId, mountId)); return db .select() .from(mediaItems) .where(and(...conds)) .orderBy(asc(mediaItems.epNumber), asc(mediaItems.rawName)); }, async getByUserPath( userId: string, mountId: string, path: string, ): Promise { const [row] = await db .select() .from(mediaItems) .where( and( eq(mediaItems.userId, userId), eq(mediaItems.mountId, mountId), eq(mediaItems.path, path), ), ) .limit(1); return row ?? null; }, async upsertFromScan(data: { userId: string; mountId: string; path: string; rawName: string; title: string; size: number; mime?: string | undefined; }): Promise { const existing = await this.getByUserPath(data.userId, data.mountId, data.path); if (existing) { const [row] = await db .update(mediaItems) .set({ size: data.size, mime: data.mime, updatedAt: new Date(), scannedAt: new Date(), }) .where(eq(mediaItems.id, existing.id)) .returning(); if (!row) throw new Error("Failed to update media item"); return row; } const rows = await db .insert(mediaItems) .values({ userId: data.userId, mountId: data.mountId, path: data.path, rawName: data.rawName, title: data.title, size: data.size, mime: data.mime, scrapeStatus: "pending", }) .returning(); const row = rows[0]; if (!row) throw new Error("Failed to insert media item"); return row; }, async updateScrape( id: string, userId: string, data: { bangumiId?: number | null; dandanplayEpisodeId?: number | null; matchedHash?: string | null; epNumber?: number | null; scrapeStatus: string; scrapedAt?: Date | null; posterUrl?: string | null; title?: string; }, ): Promise { const [row] = await db .update(mediaItems) .set({ ...data, updatedAt: new Date() }) .where(and(eq(mediaItems.id, id), eq(mediaItems.userId, userId))) .returning(); return row ?? null; }, /** 弹幕识别身份落库(owner 校验):episodeId + 识别用文件指纹,后续播放免重 match */ async setDanmakuMatch( id: string, userId: string, data: { episodeId: number; matchedHash?: string | null }, ): Promise { const [row] = await db .update(mediaItems) .set({ dandanplayEpisodeId: data.episodeId, matchedHash: data.matchedHash ?? null, updatedAt: new Date(), }) .where(and(eq(mediaItems.id, id), eq(mediaItems.userId, userId))) .returning(); return row ?? null; }, /** 删除挂载下的全部条目(挂载删除级联),返回被删条目 id 供清理进度 */ async deleteByMount(mountId: string, userId: string): Promise { const rows = await db .delete(mediaItems) .where(and(eq(mediaItems.mountId, mountId), eq(mediaItems.userId, userId))) .returning({ id: mediaItems.id }); return rows.map((r) => r.id); }, /** 清理挂载已不存在的孤儿条目(历史脏数据兜底),返回删除数量 */ async cleanupOrphans(): Promise { const rows = await db .delete(mediaItems) .where(sql`${mediaItems.mountId} not in (select id from ${mounts})`) .returning({ id: mediaItems.id }); return rows.length; }, };