diff --git a/packages/dao/src/danmaku-cache.ts b/packages/dao/src/danmaku-cache.ts new file mode 100644 index 0000000..ace6ef9 --- /dev/null +++ b/packages/dao/src/danmaku-cache.ts @@ -0,0 +1,45 @@ +import { and, eq, gt, isNull, or } from "drizzle-orm"; +import { danmakuCache, type DanmakuCacheRow } from "@app/models"; +import { db } from "@app/db"; + +export const danmakuCacheDao = { + async getValid(matchKey: string): Promise { + const now = new Date(); + const [row] = await db + .select() + .from(danmakuCache) + .where( + and( + eq(danmakuCache.matchKey, matchKey), + or(isNull(danmakuCache.expiresAt), gt(danmakuCache.expiresAt, now)), + ), + ) + .limit(1); + return row ?? null; + }, + async upsert(data: { + matchKey: string; + payload: string; + source: string; + expiresAt: Date | null; + }): Promise { + const existing = await db + .select() + .from(danmakuCache) + .where(eq(danmakuCache.matchKey, data.matchKey)) + .limit(1); + if (existing[0]) { + const [row] = await db + .update(danmakuCache) + .set({ payload: data.payload, source: data.source, expiresAt: data.expiresAt }) + .where(eq(danmakuCache.id, existing[0].id)) + .returning(); + if (!row) throw new Error("Failed to update danmaku cache"); + return row; + } + const rows = await db.insert(danmakuCache).values(data).returning(); + const row = rows[0]; + if (!row) throw new Error("Failed to insert danmaku cache"); + return row; + }, +}; diff --git a/packages/dao/src/index.ts b/packages/dao/src/index.ts index 338e53b..8edc9fb 100644 --- a/packages/dao/src/index.ts +++ b/packages/dao/src/index.ts @@ -5,3 +5,9 @@ export { sessionDao } from "./sessions.js"; export type { SessionRow } from "./sessions.js"; export { passwordResetDao } from "./password-resets.js"; export type { PasswordResetRow } from "./password-resets.js"; +export { mountDao } from "./mounts.js"; +export type { MountRow } from "./mounts.js"; +export { mediaItemDao } from "./media-items.js"; +export type { MediaItemRow } from "./media-items.js"; +export { playbackProgressDao } from "./playback-progress.js"; +export { danmakuCacheDao } from "./danmaku-cache.js"; diff --git a/packages/dao/src/media-items.ts b/packages/dao/src/media-items.ts new file mode 100644 index 0000000..034ec39 --- /dev/null +++ b/packages/dao/src/media-items.ts @@ -0,0 +1,114 @@ +import { and, desc, eq, like, or, sql } from "drizzle-orm"; +import { mediaItems, 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; + scrapeStatus?: string | undefined; + q?: string | undefined; + }, + ): Promise<{ rows: MediaItemRow[]; total: number }> { + const conds = [eq(mediaItems.userId, userId)]; + 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) }; + }, + 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; + }, + async getByUserPath(userId: string, path: string): Promise { + const [row] = await db + .select() + .from(mediaItems) + .where(and(eq(mediaItems.userId, userId), 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.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; + 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; + }, +}; diff --git a/packages/dao/src/mounts.ts b/packages/dao/src/mounts.ts new file mode 100644 index 0000000..a4b0fae --- /dev/null +++ b/packages/dao/src/mounts.ts @@ -0,0 +1,44 @@ +import { and, desc, eq } from "drizzle-orm"; +import { mounts, type Mount, type NewMount } from "@app/models"; +import { db } from "@app/db"; + +export type MountRow = Mount; + +export const mountDao = { + async listByUser(userId: string): Promise { + return db + .select() + .from(mounts) + .where(eq(mounts.userId, userId)) + .orderBy(desc(mounts.createdAt)); + }, + async getByIdForUser(id: string, userId: string): Promise { + const [row] = await db + .select() + .from(mounts) + .where(and(eq(mounts.id, id), eq(mounts.userId, userId))) + .limit(1); + return row ?? null; + }, + async create(data: NewMount): Promise { + const rows = await db.insert(mounts).values(data).returning(); + const row = rows[0]; + if (!row) throw new Error("Failed to create mount"); + return row; + }, + async update( + id: string, + userId: string, + data: Partial>, + ): Promise { + const [row] = await db + .update(mounts) + .set({ ...data, updatedAt: new Date() }) + .where(and(eq(mounts.id, id), eq(mounts.userId, userId))) + .returning(); + return row ?? null; + }, + async delete(id: string, userId: string): Promise { + await db.delete(mounts).where(and(eq(mounts.id, id), eq(mounts.userId, userId))); + }, +}; diff --git a/packages/dao/src/playback-progress.ts b/packages/dao/src/playback-progress.ts new file mode 100644 index 0000000..32897fe --- /dev/null +++ b/packages/dao/src/playback-progress.ts @@ -0,0 +1,43 @@ +import { and, eq } from "drizzle-orm"; +import { playbackProgress, type PlaybackProgress } from "@app/models"; +import { db } from "@app/db"; + +export const playbackProgressDao = { + async get(userId: string, mediaItemId: string): Promise { + const [row] = await db + .select() + .from(playbackProgress) + .where( + and( + eq(playbackProgress.userId, userId), + eq(playbackProgress.mediaItemId, mediaItemId), + ), + ) + .limit(1); + return row ?? null; + }, + async upsert( + userId: string, + mediaItemId: string, + positionMs: number, + durationMs: number, + ): Promise { + const existing = await this.get(userId, mediaItemId); + if (existing) { + const [row] = await db + .update(playbackProgress) + .set({ positionMs, durationMs, updatedAt: new Date() }) + .where(eq(playbackProgress.id, existing.id)) + .returning(); + if (!row) throw new Error("Failed to update progress"); + return row; + } + const rows = await db + .insert(playbackProgress) + .values({ userId, mediaItemId, positionMs, durationMs }) + .returning(); + const row = rows[0]; + if (!row) throw new Error("Failed to insert progress"); + return row; + }, +};