/* eslint-disable no-console, @typescript-eslint/no-explicit-any, @typescript-eslint/no-non-null-assertion */ import { createClient, RedisClientType } from 'redis'; import { AdminConfig } from './admin.types'; import { Favorite, IStorage, PlayRecord, SkipConfig } from './types'; // 搜索历史最大条数 const SEARCH_HISTORY_LIMIT = 20; // 数据类型转换辅助函数 function ensureString(value: any): string { return String(value); } function ensureStringArray(value: any[]): string[] { return value.map((item) => String(item)); } // 连接配置接口 export interface RedisConnectionConfig { url: string; clientName: string; // 用于日志显示,如 "Redis" 或 "Pika" } // 添加Redis操作重试包装器 function createRetryWrapper(clientName: string, getClient: () => RedisClientType) { return async function withRetry( operation: () => Promise, maxRetries = 3 ): Promise { for (let i = 0; i < maxRetries; i++) { try { return await operation(); } catch (err: any) { const isLastAttempt = i === maxRetries - 1; const isConnectionError = err.message?.includes('Connection') || err.message?.includes('ECONNREFUSED') || err.message?.includes('ENOTFOUND') || err.code === 'ECONNRESET' || err.code === 'EPIPE'; if (isConnectionError && !isLastAttempt) { console.log( `${clientName} operation failed, retrying... (${i + 1}/${maxRetries})` ); console.error('Error:', err.message); // 等待一段时间后重试 await new Promise((resolve) => setTimeout(resolve, 1000 * (i + 1))); // 尝试重新连接 try { const client = getClient(); if (!client.isOpen) { await client.connect(); } } catch (reconnectErr) { console.error('Failed to reconnect:', reconnectErr); } continue; } throw err; } } throw new Error('Max retries exceeded'); }; } // 创建客户端的工厂函数 export function createRedisClient(config: RedisConnectionConfig, globalSymbol: symbol): RedisClientType { let client: RedisClientType | undefined = (global as any)[globalSymbol]; if (!client) { if (!config.url) { throw new Error(`${config.clientName}_URL env variable not set`); } // 创建客户端配置 const clientConfig: any = { url: config.url, socket: { // 重连策略:指数退避,最大30秒 reconnectStrategy: (retries: number) => { console.log(`${config.clientName} reconnection attempt ${retries + 1}`); if (retries > 10) { console.error(`${config.clientName} max reconnection attempts exceeded`); return false; // 停止重连 } return Math.min(1000 * Math.pow(2, retries), 30000); // 指数退避,最大30秒 }, connectTimeout: 10000, // 10秒连接超时 // 设置no delay,减少延迟 noDelay: true, }, // 添加其他配置 pingInterval: 30000, // 30秒ping一次,保持连接活跃 }; client = createClient(clientConfig); // 添加错误事件监听 client.on('error', (err) => { console.error(`${config.clientName} client error:`, err); }); client.on('connect', () => { console.log(`${config.clientName} connected`); }); client.on('reconnecting', () => { console.log(`${config.clientName} reconnecting...`); }); client.on('ready', () => { console.log(`${config.clientName} ready`); }); // 初始连接,带重试机制 const connectWithRetry = async () => { try { await client!.connect(); console.log(`${config.clientName} connected successfully`); } catch (err) { console.error(`${config.clientName} initial connection failed:`, err); console.log('Will retry in 5 seconds...'); setTimeout(connectWithRetry, 5000); } }; connectWithRetry(); (global as any)[globalSymbol] = client; } return client; } // 抽象基类,包含所有通用的Redis操作逻辑 export abstract class BaseRedisStorage implements IStorage { protected client: RedisClientType; protected withRetry: (operation: () => Promise, maxRetries?: number) => Promise; constructor(config: RedisConnectionConfig, globalSymbol: symbol) { this.client = createRedisClient(config, globalSymbol); this.withRetry = createRetryWrapper(config.clientName, () => this.client); } // ---------- 播放记录 ---------- private prKey(user: string, key: string) { return `u:${user}:pr:${key}`; // u:username:pr:source+id } async getPlayRecord( userName: string, key: string ): Promise { const val = await this.withRetry(() => this.client.get(this.prKey(userName, key)) ); return val ? (JSON.parse(val) as PlayRecord) : null; } async setPlayRecord( userName: string, key: string, record: PlayRecord ): Promise { await this.withRetry(() => this.client.set(this.prKey(userName, key), JSON.stringify(record)) ); } async getAllPlayRecords( userName: string ): Promise> { const pattern = `u:${userName}:pr:*`; const keys: string[] = await this.withRetry(() => this.client.keys(pattern)); if (keys.length === 0) return {}; const values = await this.withRetry(() => this.client.mGet(keys)); const result: Record = {}; keys.forEach((fullKey: string, idx: number) => { const raw = values[idx]; if (raw) { const rec = JSON.parse(raw) as PlayRecord; // 截取 source+id 部分 const keyPart = ensureString(fullKey.replace(`u:${userName}:pr:`, '')); result[keyPart] = rec; } }); return result; } async deletePlayRecord(userName: string, key: string): Promise { await this.withRetry(() => this.client.del(this.prKey(userName, key))); } // ---------- 收藏 ---------- private favKey(user: string, key: string) { return `u:${user}:fav:${key}`; } async getFavorite(userName: string, key: string): Promise { const val = await this.withRetry(() => this.client.get(this.favKey(userName, key)) ); return val ? (JSON.parse(val) as Favorite) : null; } async setFavorite( userName: string, key: string, favorite: Favorite ): Promise { await this.withRetry(() => this.client.set(this.favKey(userName, key), JSON.stringify(favorite)) ); } async getAllFavorites(userName: string): Promise> { const pattern = `u:${userName}:fav:*`; const keys: string[] = await this.withRetry(() => this.client.keys(pattern)); if (keys.length === 0) return {}; const values = await this.withRetry(() => this.client.mGet(keys)); const result: Record = {}; keys.forEach((fullKey: string, idx: number) => { const raw = values[idx]; if (raw) { const fav = JSON.parse(raw) as Favorite; const keyPart = ensureString(fullKey.replace(`u:${userName}:fav:`, '')); result[keyPart] = fav; } }); return result; } async deleteFavorite(userName: string, key: string): Promise { await this.withRetry(() => this.client.del(this.favKey(userName, key))); } // ---------- 用户注册 / 登录(旧版本,保持兼容) ---------- private userPwdKey(user: string) { return `u:${user}:pwd`; } async registerUser(userName: string, password: string): Promise { // 简单存储明文密码,生产环境应加密 await this.withRetry(() => this.client.set(this.userPwdKey(userName), password)); } async verifyUser(userName: string, password: string): Promise { const stored = await this.withRetry(() => this.client.get(this.userPwdKey(userName)) ); if (stored === null) return false; // 确保比较时都是字符串类型 return ensureString(stored) === password; } // 检查用户是否存在 async checkUserExist(userName: string): Promise { // 使用 EXISTS 判断 key 是否存在 const exists = await this.withRetry(() => this.client.exists(this.userPwdKey(userName)) ); return exists === 1; } // 修改用户密码 async changePassword(userName: string, newPassword: string): Promise { // 简单存储明文密码,生产环境应加密 await this.withRetry(() => this.client.set(this.userPwdKey(userName), newPassword) ); } // 删除用户及其所有数据 async deleteUser(userName: string): Promise { // 删除用户密码 await this.withRetry(() => this.client.del(this.userPwdKey(userName))); // 删除搜索历史 await this.withRetry(() => this.client.del(this.shKey(userName))); // 删除播放记录 const playRecordPattern = `u:${userName}:pr:*`; const playRecordKeys = await this.withRetry(() => this.client.keys(playRecordPattern) ); if (playRecordKeys.length > 0) { await this.withRetry(() => this.client.del(playRecordKeys)); } // 删除收藏夹 const favoritePattern = `u:${userName}:fav:*`; const favoriteKeys = await this.withRetry(() => this.client.keys(favoritePattern) ); if (favoriteKeys.length > 0) { await this.withRetry(() => this.client.del(favoriteKeys)); } // 删除跳过片头片尾配置 const skipConfigPattern = `u:${userName}:skip:*`; const skipConfigKeys = await this.withRetry(() => this.client.keys(skipConfigPattern) ); if (skipConfigKeys.length > 0) { await this.withRetry(() => this.client.del(skipConfigKeys)); } } // ---------- 新版用户存储(使用Hash和Sorted Set) ---------- private userInfoKey(userName: string) { return `user:${userName}:info`; } private userListKey() { return 'user:list'; } private oidcSubKey(oidcSub: string) { return `oidc:sub:${oidcSub}`; } // SHA256加密密码 private async hashPassword(password: string): Promise { const encoder = new TextEncoder(); const data = encoder.encode(password); const hashBuffer = await crypto.subtle.digest('SHA-256', data); const hashArray = Array.from(new Uint8Array(hashBuffer)); return hashArray.map(b => b.toString(16).padStart(2, '0')).join(''); } // 创建新用户(新版本) async createUserV2( userName: string, password: string, role: 'owner' | 'admin' | 'user' = 'user', tags?: string[], oidcSub?: string ): Promise { const hashedPassword = await this.hashPassword(password); const createdAt = Date.now(); // 存储用户信息到Hash const userInfo: Record = { role, banned: 'false', password: hashedPassword, created_at: createdAt.toString(), }; if (tags && tags.length > 0) { userInfo.tags = JSON.stringify(tags); } if (oidcSub) { userInfo.oidcSub = oidcSub; // 创建OIDC映射 await this.withRetry(() => this.client.set(this.oidcSubKey(oidcSub), userName)); } await this.withRetry(() => this.client.hSet(this.userInfoKey(userName), userInfo)); // 添加到用户列表(Sorted Set,按注册时间排序) await this.withRetry(() => this.client.zAdd(this.userListKey(), { score: createdAt, value: userName, })); // 如果创建的是站长用户,清除站长存在状态缓存 if (userName === process.env.USERNAME) { const { ownerExistenceCache } = await import('./user-cache'); ownerExistenceCache.delete(userName); } } // 验证用户密码(新版本) async verifyUserV2(userName: string, password: string): Promise { const userInfo = await this.withRetry(() => this.client.hGetAll(this.userInfoKey(userName)) ); if (!userInfo || !userInfo.password) { return false; } const hashedPassword = await this.hashPassword(password); return userInfo.password === hashedPassword; } // 获取用户信息(新版本) async getUserInfoV2(userName: string): Promise<{ role: 'owner' | 'admin' | 'user'; banned: boolean; tags?: string[]; oidcSub?: string; created_at: number; } | null> { const userInfo = await this.withRetry(() => this.client.hGetAll(this.userInfoKey(userName)) ); if (!userInfo || Object.keys(userInfo).length === 0) { return null; } return { role: (userInfo.role as 'owner' | 'admin' | 'user') || 'user', banned: userInfo.banned === 'true', tags: userInfo.tags ? JSON.parse(userInfo.tags) : undefined, oidcSub: userInfo.oidcSub, created_at: parseInt(userInfo.created_at || '0', 10), }; } // 更新用户信息(新版本) async updateUserInfoV2( userName: string, updates: { role?: 'owner' | 'admin' | 'user'; banned?: boolean; tags?: string[]; oidcSub?: string; } ): Promise { const userInfo: Record = {}; if (updates.role !== undefined) { userInfo.role = updates.role; } if (updates.banned !== undefined) { userInfo.banned = updates.banned ? 'true' : 'false'; } if (updates.tags !== undefined) { if (updates.tags.length > 0) { userInfo.tags = JSON.stringify(updates.tags); } else { // 删除tags字段 await this.withRetry(() => this.client.hDel(this.userInfoKey(userName), 'tags')); } } if (updates.oidcSub !== undefined) { const oldInfo = await this.getUserInfoV2(userName); if (oldInfo?.oidcSub && oldInfo.oidcSub !== updates.oidcSub) { // 删除旧的OIDC映射 await this.withRetry(() => this.client.del(this.oidcSubKey(oldInfo.oidcSub!))); } userInfo.oidcSub = updates.oidcSub; // 创建新的OIDC映射 await this.withRetry(() => this.client.set(this.oidcSubKey(updates.oidcSub!), userName)); } if (Object.keys(userInfo).length > 0) { await this.withRetry(() => this.client.hSet(this.userInfoKey(userName), userInfo)); } } // 修改用户密码(新版本) async changePasswordV2(userName: string, newPassword: string): Promise { const hashedPassword = await this.hashPassword(newPassword); await this.withRetry(() => this.client.hSet(this.userInfoKey(userName), 'password', hashedPassword) ); } // 检查用户是否存在(新版本) async checkUserExistV2(userName: string): Promise { const exists = await this.withRetry(() => this.client.exists(this.userInfoKey(userName)) ); return exists === 1; } // 通过OIDC Sub查找用户名 async getUserByOidcSub(oidcSub: string): Promise { const userName = await this.withRetry(() => this.client.get(this.oidcSubKey(oidcSub)) ); return userName ? ensureString(userName) : null; } // 获取用户列表(分页,新版本) async getUserListV2( offset: number = 0, limit: number = 20, ownerUsername?: string ): Promise<{ users: Array<{ username: string; role: 'owner' | 'admin' | 'user'; banned: boolean; tags?: string[]; created_at: number; }>; total: number; }> { // 获取总数 let total = await this.withRetry(() => this.client.zCard(this.userListKey())); // 检查站长是否在数据库中(使用缓存) let ownerInfo = null; let ownerInDatabase = false; if (ownerUsername) { // 先检查缓存 const { ownerExistenceCache } = await import('./user-cache'); const cachedExists = ownerExistenceCache.get(ownerUsername); if (cachedExists !== null) { // 使用缓存的结果 ownerInDatabase = cachedExists; if (ownerInDatabase) { // 如果站长在数据库中,获取详细信息 ownerInfo = await this.getUserInfoV2(ownerUsername); } } else { // 缓存未命中,查询数据库 ownerInfo = await this.getUserInfoV2(ownerUsername); ownerInDatabase = !!ownerInfo; // 更新缓存 ownerExistenceCache.set(ownerUsername, ownerInDatabase); } // 如果站长不在数据库中,总数+1(无论在哪一页都要加) if (!ownerInDatabase) { total += 1; } } // 如果站长不在数据库中且在第一页,需要调整获取的用户数量和偏移量 let actualOffset = offset; let actualLimit = limit; if (ownerUsername && !ownerInDatabase) { if (offset === 0) { // 第一页:只获取 limit-1 个用户,为站长留出位置 actualLimit = limit - 1; } else { // 其他页:偏移量需要减1,因为站长占据了第一页的一个位置 actualOffset = offset - 1; } } // 获取用户列表(按注册时间升序) const usernames = await this.withRetry(() => this.client.zRange(this.userListKey(), actualOffset, actualOffset + actualLimit - 1) ); const users = []; // 如果有站长且在第一页,确保站长始终在第一位 if (ownerUsername && offset === 0) { // 即使站长不在数据库中,也要添加站长(站长使用环境变量认证) users.push({ username: ownerUsername, role: 'owner' as const, banned: ownerInfo?.banned || false, tags: ownerInfo?.tags, created_at: ownerInfo?.created_at || 0, }); } // 获取其他用户信息 for (const username of usernames) { const usernameStr = ensureString(username); // 跳过站长(已经添加) if (ownerUsername && usernameStr === ownerUsername) { continue; } const userInfo = await this.getUserInfoV2(usernameStr); if (userInfo) { users.push({ username: usernameStr, role: userInfo.role, banned: userInfo.banned, tags: userInfo.tags, created_at: userInfo.created_at, }); } } return { users, total }; } // 删除用户(新版本) async deleteUserV2(userName: string): Promise { // 获取用户信息 const userInfo = await this.getUserInfoV2(userName); // 删除OIDC映射 if (userInfo?.oidcSub) { await this.withRetry(() => this.client.del(this.oidcSubKey(userInfo.oidcSub!))); } // 删除用户信息Hash await this.withRetry(() => this.client.del(this.userInfoKey(userName))); // 从用���列表中移除 await this.withRetry(() => this.client.zRem(this.userListKey(), userName)); // 删除用户的其他数据(播放记录、收藏等) await this.deleteUser(userName); } // ---------- 搜索历史 ---------- private shKey(user: string) { return `u:${user}:sh`; // u:username:sh } async getSearchHistory(userName: string): Promise { const result = await this.withRetry(() => this.client.lRange(this.shKey(userName), 0, -1) ); // 确保返回的都是字符串类型 return ensureStringArray(result as any[]); } async addSearchHistory(userName: string, keyword: string): Promise { const key = this.shKey(userName); // 先去重 await this.withRetry(() => this.client.lRem(key, 0, ensureString(keyword))); // 插入到最前 await this.withRetry(() => this.client.lPush(key, ensureString(keyword))); // 限制最大长度 await this.withRetry(() => this.client.lTrim(key, 0, SEARCH_HISTORY_LIMIT - 1)); } async deleteSearchHistory(userName: string, keyword?: string): Promise { const key = this.shKey(userName); if (keyword) { await this.withRetry(() => this.client.lRem(key, 0, ensureString(keyword))); } else { await this.withRetry(() => this.client.del(key)); } } // ---------- 获取全部用户 ---------- async getAllUsers(): Promise { const keys = await this.withRetry(() => this.client.keys('u:*:pwd')); return keys .map((k) => { const match = k.match(/^u:(.+?):pwd$/); return match ? ensureString(match[1]) : undefined; }) .filter((u): u is string => typeof u === 'string'); } // ---------- 管理员配置 ---------- private adminConfigKey() { return 'admin:config'; } async getAdminConfig(): Promise { const val = await this.withRetry(() => this.client.get(this.adminConfigKey())); return val ? (JSON.parse(val) as AdminConfig) : null; } async setAdminConfig(config: AdminConfig): Promise { await this.withRetry(() => this.client.set(this.adminConfigKey(), JSON.stringify(config)) ); } // ---------- 跳过片头片尾配置 ---------- private skipConfigKey(user: string, source: string, id: string) { return `u:${user}:skip:${source}+${id}`; } private danmakuFilterConfigKey(user: string) { return `u:${user}:danmaku_filter`; } async getSkipConfig( userName: string, source: string, id: string ): Promise { const val = await this.withRetry(() => this.client.get(this.skipConfigKey(userName, source, id)) ); return val ? (JSON.parse(val) as SkipConfig) : null; } async setSkipConfig( userName: string, source: string, id: string, config: SkipConfig ): Promise { await this.withRetry(() => this.client.set( this.skipConfigKey(userName, source, id), JSON.stringify(config) ) ); } async deleteSkipConfig( userName: string, source: string, id: string ): Promise { await this.withRetry(() => this.client.del(this.skipConfigKey(userName, source, id)) ); } async getAllSkipConfigs( userName: string ): Promise<{ [key: string]: SkipConfig }> { const pattern = `u:${userName}:skip:*`; const keys = await this.withRetry(() => this.client.keys(pattern)); if (keys.length === 0) { return {}; } const configs: { [key: string]: SkipConfig } = {}; // 批量获取所有配置 const values = await this.withRetry(() => this.client.mGet(keys)); keys.forEach((key, index) => { const value = values[index]; if (value) { // 从key中提取source+id const match = key.match(/^u:.+?:skip:(.+)$/); if (match) { const sourceAndId = match[1]; configs[sourceAndId] = JSON.parse(value as string) as SkipConfig; } } }); return configs; } // ---------- 弹幕过滤配置 ---------- async getDanmakuFilterConfig( userName: string ): Promise { const val = await this.withRetry(() => this.client.get(this.danmakuFilterConfigKey(userName)) ); return val ? (JSON.parse(val) as import('./types').DanmakuFilterConfig) : null; } async setDanmakuFilterConfig( userName: string, config: import('./types').DanmakuFilterConfig ): Promise { await this.withRetry(() => this.client.set( this.danmakuFilterConfigKey(userName), JSON.stringify(config) ) ); } async deleteDanmakuFilterConfig(userName: string): Promise { await this.withRetry(() => this.client.del(this.danmakuFilterConfigKey(userName)) ); } // 清空所有数据 async clearAllData(): Promise { try { // 获取所有用户 const allUsers = await this.getAllUsers(); // 删除所有用户及其数据 for (const username of allUsers) { await this.deleteUser(username); } // 删除管理员配置 await this.withRetry(() => this.client.del(this.adminConfigKey())); console.log('所有数据已清空'); } catch (error) { console.error('清空数据失败:', error); throw new Error('清空数据失败'); } } // ---------- 通用键值存储 ---------- private globalValueKey(key: string) { return `global:${key}`; } async getGlobalValue(key: string): Promise { const val = await this.withRetry(() => this.client.get(this.globalValueKey(key)) ); return val ? ensureString(val) : null; } async setGlobalValue(key: string, value: string): Promise { await this.withRetry(() => this.client.set(this.globalValueKey(key), ensureString(value)) ); } async deleteGlobalValue(key: string): Promise { await this.withRetry(() => this.client.del(this.globalValueKey(key))); } // ---------- 通知相关 ---------- private notificationsKey(userName: string) { return `u:${userName}:notifications`; } private lastFavoriteCheckKey(userName: string) { return `u:${userName}:last_fav_check`; } async getNotifications(userName: string): Promise { const val = await this.withRetry(() => this.client.get(this.notificationsKey(userName)) ); return val ? (JSON.parse(val) as import('./types').Notification[]) : []; } async addNotification( userName: string, notification: import('./types').Notification ): Promise { const notifications = await this.getNotifications(userName); notifications.unshift(notification); // 新通知放在最前面 // 限制通知数量,最多保留100条 if (notifications.length > 100) { notifications.splice(100); } await this.withRetry(() => this.client.set(this.notificationsKey(userName), JSON.stringify(notifications)) ); } async markNotificationAsRead( userName: string, notificationId: string ): Promise { const notifications = await this.getNotifications(userName); const notification = notifications.find((n) => n.id === notificationId); if (notification) { notification.read = true; await this.withRetry(() => this.client.set(this.notificationsKey(userName), JSON.stringify(notifications)) ); } } async deleteNotification( userName: string, notificationId: string ): Promise { const notifications = await this.getNotifications(userName); const filtered = notifications.filter((n) => n.id !== notificationId); await this.withRetry(() => this.client.set(this.notificationsKey(userName), JSON.stringify(filtered)) ); } async clearAllNotifications(userName: string): Promise { await this.withRetry(() => this.client.del(this.notificationsKey(userName))); } async getUnreadNotificationCount(userName: string): Promise { const notifications = await this.getNotifications(userName); return notifications.filter((n) => !n.read).length; } async getLastFavoriteCheckTime(userName: string): Promise { const val = await this.withRetry(() => this.client.get(this.lastFavoriteCheckKey(userName)) ); return val ? parseInt(val, 10) : 0; } async setLastFavoriteCheckTime( userName: string, timestamp: number ): Promise { await this.withRetry(() => this.client.set(this.lastFavoriteCheckKey(userName), timestamp.toString()) ); } }