修复upstash方式数据导入

This commit is contained in:
mtvpls
2026-01-12 20:10:37 +08:00
parent 83b3315fba
commit f1ced03517
3 changed files with 118 additions and 102 deletions
@@ -155,8 +155,8 @@ async function getUserPasswordV2(username: string): Promise<string | null> {
// 直接调用hGetAll获取完整用户信息(包括密码) // 直接调用hGetAll获取完整用户信息(包括密码)
const userInfoKey = `user:${username}:info`; const userInfoKey = `user:${username}:info`;
if (typeof storage.withRetry === 'function' && storage.client?.hGetAll) { if (typeof storage.withRetry === 'function' && storage.client?.hgetall) {
const userInfo = await storage.withRetry(() => storage.client.hGetAll(userInfoKey)); const userInfo = await storage.withRetry(() => storage.client.hgetall(userInfoKey));
if (userInfo && userInfo.password) { if (userInfo && userInfo.password) {
return userInfo.password; return userInfo.password;
} }
@@ -149,10 +149,10 @@ export async function POST(req: NextRequest) {
} }
// 使用storage.withRetry直接设置用户信息 // 使用storage.withRetry直接设置用户信息
await storage.withRetry(() => storage.client.hSet(userInfoKey, userInfo)); await storage.withRetry(() => storage.client.hset(userInfoKey, userInfo));
// 添加到用户列表 // 添加到用户列表
await storage.withRetry(() => storage.client.zAdd('user:list', { await storage.withRetry(() => storage.client.zadd('user:list', {
score: createdAt, score: createdAt,
value: username, value: username,
})); }));
+114 -98
View File
@@ -58,10 +58,26 @@ async function withRetry<T>(
} }
export class UpstashRedisStorage implements IStorage { export class UpstashRedisStorage implements IStorage {
private client: Redis; private _client: Redis;
client: any;
constructor() { constructor() {
this.client = getUpstashRedisClient(); this._client = getUpstashRedisClient();
// 创建兼容Redis API的client包装器(支持camelCase和lowercase
this.client = {
hSet: (key: string, data: Record<string, any>) => this._client.hset(key, data),
hset: (key: string, data: Record<string, any>) => this._client.hset(key, data),
zAdd: (key: string, member: { score: number; value: string }) => this._client.zadd(key, { score: member.score, member: member.value }),
zadd: (key: string, member: { score: number; value: string }) => this._client.zadd(key, { score: member.score, member: member.value }),
set: (key: string, value: string) => this._client.set(key, value),
hGetAll: (key: string) => this._client.hgetall(key),
hgetall: (key: string) => this._client.hgetall(key),
};
}
// 公开withRetry方法供外部使用
withRetry<T>(operation: () => Promise<T>, maxRetries = 3): Promise<T> {
return withRetry(operation, maxRetries);
} }
// ---------- 播放记录 ---------- // ---------- 播放记录 ----------
@@ -79,7 +95,7 @@ export class UpstashRedisStorage implements IStorage {
key: string key: string
): Promise<PlayRecord | null> { ): Promise<PlayRecord | null> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.hget(this.prHashKey(userName), key) this._client.hget(this.prHashKey(userName), key)
); );
return val ? (val as PlayRecord) : null; return val ? (val as PlayRecord) : null;
} }
@@ -89,14 +105,14 @@ export class UpstashRedisStorage implements IStorage {
key: string, key: string,
record: PlayRecord record: PlayRecord
): Promise<void> { ): Promise<void> {
await withRetry(() => this.client.hset(this.prHashKey(userName), { [key]: record })); await withRetry(() => this._client.hset(this.prHashKey(userName), { [key]: record }));
} }
async getAllPlayRecords( async getAllPlayRecords(
userName: string userName: string
): Promise<Record<string, PlayRecord>> { ): Promise<Record<string, PlayRecord>> {
const hashData = await withRetry(() => const hashData = await withRetry(() =>
this.client.hgetall(this.prHashKey(userName)) this._client.hgetall(this.prHashKey(userName))
); );
if (!hashData || Object.keys(hashData).length === 0) return {}; if (!hashData || Object.keys(hashData).length === 0) return {};
@@ -111,7 +127,7 @@ export class UpstashRedisStorage implements IStorage {
} }
async deletePlayRecord(userName: string, key: string): Promise<void> { async deletePlayRecord(userName: string, key: string): Promise<void> {
await withRetry(() => this.client.hdel(this.prHashKey(userName), key)); await withRetry(() => this._client.hdel(this.prHashKey(userName), key));
} }
// 清理超出限制的旧播放记录 // 清理超出限制的旧播放记录
@@ -210,13 +226,13 @@ export class UpstashRedisStorage implements IStorage {
// 2. 获取旧结构的所有播放记录key // 2. 获取旧结构的所有播放记录key
const pattern = `u:${userName}:pr:*`; const pattern = `u:${userName}:pr:*`;
const oldKeys: string[] = await withRetry(() => this.client.keys(pattern)); const oldKeys: string[] = await withRetry(() => this._client.keys(pattern));
if (oldKeys.length === 0) { if (oldKeys.length === 0) {
console.log(`用户 ${userName} 没有旧的播放记录,标记为已迁移`); console.log(`用户 ${userName} 没有旧的播放记录,标记为已迁移`);
// 即使没有数据也标记为已迁移 // 即使没有数据也标记为已迁移
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { playrecord_migrated: true }) this._client.hset(this.userInfoKey(userName), { playrecord_migrated: true })
); );
// 清除用户信息缓存 // 清除用户信息缓存
userInfoCache?.delete(userName); userInfoCache?.delete(userName);
@@ -228,7 +244,7 @@ export class UpstashRedisStorage implements IStorage {
// 3. 批量获取旧数据并转换为hash格式 // 3. 批量获取旧数据并转换为hash格式
const hashData: Record<string, any> = {}; const hashData: Record<string, any> = {};
for (const fullKey of oldKeys) { for (const fullKey of oldKeys) {
const value = await withRetry(() => this.client.get(fullKey)); const value = await withRetry(() => this._client.get(fullKey));
if (value) { if (value) {
// 提取 source+id 部分作为hash的field // 提取 source+id 部分作为hash的field
const keyPart = ensureString(fullKey.replace(`u:${userName}:pr:`, '')); const keyPart = ensureString(fullKey.replace(`u:${userName}:pr:`, ''));
@@ -239,18 +255,18 @@ export class UpstashRedisStorage implements IStorage {
// 4. 写入新的hash结构 // 4. 写入新的hash结构
if (Object.keys(hashData).length > 0) { if (Object.keys(hashData).length > 0) {
await withRetry(() => await withRetry(() =>
this.client.hset(this.prHashKey(userName), hashData) this._client.hset(this.prHashKey(userName), hashData)
); );
console.log(`成功迁移 ${Object.keys(hashData).length} 条播放记录到hash结构`); console.log(`成功迁移 ${Object.keys(hashData).length} 条播放记录到hash结构`);
} }
// 5. 删除旧的key // 5. 删除旧的key
await withRetry(() => this.client.del(...oldKeys)); await withRetry(() => this._client.del(...oldKeys));
console.log(`删除了 ${oldKeys.length} 个旧的播放记录key`); console.log(`删除了 ${oldKeys.length} 个旧的播放记录key`);
// 6. 标记迁移完成 // 6. 标记迁移完成
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { playrecord_migrated: true }) this._client.hset(this.userInfoKey(userName), { playrecord_migrated: true })
); );
// 7. 清除用户信息缓存,确保下次获取时能读取到最新的迁移标识 // 7. 清除用户信息缓存,确保下次获取时能读取到最新的迁移标识
@@ -271,7 +287,7 @@ export class UpstashRedisStorage implements IStorage {
async getFavorite(userName: string, key: string): Promise<Favorite | null> { async getFavorite(userName: string, key: string): Promise<Favorite | null> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.hget(this.favHashKey(userName), key) this._client.hget(this.favHashKey(userName), key)
); );
return val ? (val as Favorite) : null; return val ? (val as Favorite) : null;
} }
@@ -281,12 +297,12 @@ export class UpstashRedisStorage implements IStorage {
key: string, key: string,
favorite: Favorite favorite: Favorite
): Promise<void> { ): Promise<void> {
await withRetry(() => this.client.hset(this.favHashKey(userName), { [key]: favorite })); await withRetry(() => this._client.hset(this.favHashKey(userName), { [key]: favorite }));
} }
async getAllFavorites(userName: string): Promise<Record<string, Favorite>> { async getAllFavorites(userName: string): Promise<Record<string, Favorite>> {
const hashData = await withRetry(() => const hashData = await withRetry(() =>
this.client.hgetall(this.favHashKey(userName)) this._client.hgetall(this.favHashKey(userName))
); );
if (!hashData || Object.keys(hashData).length === 0) return {}; if (!hashData || Object.keys(hashData).length === 0) return {};
@@ -301,7 +317,7 @@ export class UpstashRedisStorage implements IStorage {
} }
async deleteFavorite(userName: string, key: string): Promise<void> { async deleteFavorite(userName: string, key: string): Promise<void> {
await withRetry(() => this.client.hdel(this.favHashKey(userName), key)); await withRetry(() => this._client.hdel(this.favHashKey(userName), key));
} }
// 迁移收藏:从旧的多key结构迁移到新的hash结构 // 迁移收藏:从旧的多key结构迁移到新的hash结构
@@ -339,13 +355,13 @@ export class UpstashRedisStorage implements IStorage {
// 2. 获取旧结构的所有收藏key // 2. 获取旧结构的所有收藏key
const pattern = `u:${userName}:fav:*`; const pattern = `u:${userName}:fav:*`;
const oldKeys: string[] = await withRetry(() => this.client.keys(pattern)); const oldKeys: string[] = await withRetry(() => this._client.keys(pattern));
if (oldKeys.length === 0) { if (oldKeys.length === 0) {
console.log(`用户 ${userName} 没有旧的收藏,标记为已迁移`); console.log(`用户 ${userName} 没有旧的收藏,标记为已迁移`);
// 即使没有数据也标记为已迁移 // 即使没有数据也标记为已迁移
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { favorite_migrated: true }) this._client.hset(this.userInfoKey(userName), { favorite_migrated: true })
); );
// 清除用户信息缓存 // 清除用户信息缓存
userInfoCache?.delete(userName); userInfoCache?.delete(userName);
@@ -357,7 +373,7 @@ export class UpstashRedisStorage implements IStorage {
// 3. 批量获取旧数据并转换为hash格式 // 3. 批量获取旧数据并转换为hash格式
const hashData: Record<string, any> = {}; const hashData: Record<string, any> = {};
for (const fullKey of oldKeys) { for (const fullKey of oldKeys) {
const value = await withRetry(() => this.client.get(fullKey)); const value = await withRetry(() => this._client.get(fullKey));
if (value) { if (value) {
// 提取 source+id 部分作为hash的field // 提取 source+id 部分作为hash的field
const keyPart = ensureString(fullKey.replace(`u:${userName}:fav:`, '')); const keyPart = ensureString(fullKey.replace(`u:${userName}:fav:`, ''));
@@ -368,18 +384,18 @@ export class UpstashRedisStorage implements IStorage {
// 4. 写入新的hash结构 // 4. 写入新的hash结构
if (Object.keys(hashData).length > 0) { if (Object.keys(hashData).length > 0) {
await withRetry(() => await withRetry(() =>
this.client.hset(this.favHashKey(userName), hashData) this._client.hset(this.favHashKey(userName), hashData)
); );
console.log(`成功迁移 ${Object.keys(hashData).length} 条收藏到hash结构`); console.log(`成功迁移 ${Object.keys(hashData).length} 条收藏到hash结构`);
} }
// 5. 删除旧的key // 5. 删除旧的key
await withRetry(() => this.client.del(...oldKeys)); await withRetry(() => this._client.del(...oldKeys));
console.log(`删除了 ${oldKeys.length} 个旧的收藏key`); console.log(`删除了 ${oldKeys.length} 个旧的收藏key`);
// 6. 标记迁移完成 // 6. 标记迁移完成
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { favorite_migrated: true }) this._client.hset(this.userInfoKey(userName), { favorite_migrated: true })
); );
// 7. 清除用户信息缓存,确保下次获取时能读取到最新的迁移标识 // 7. 清除用户信息缓存,确保下次获取时能读取到最新的迁移标识
@@ -395,7 +411,7 @@ export class UpstashRedisStorage implements IStorage {
async verifyUser(userName: string, password: string): Promise<boolean> { async verifyUser(userName: string, password: string): Promise<boolean> {
const stored = await withRetry(() => const stored = await withRetry(() =>
this.client.get(this.userPwdKey(userName)) this._client.get(this.userPwdKey(userName))
); );
if (stored === null) return false; if (stored === null) return false;
// 确保比较时都是字符串类型 // 确保比较时都是字符串类型
@@ -406,7 +422,7 @@ export class UpstashRedisStorage implements IStorage {
async checkUserExist(userName: string): Promise<boolean> { async checkUserExist(userName: string): Promise<boolean> {
// 使用 EXISTS 判断 key 是否存在 // 使用 EXISTS 判断 key 是否存在
const exists = await withRetry(() => const exists = await withRetry(() =>
this.client.exists(this.userPwdKey(userName)) this._client.exists(this.userPwdKey(userName))
); );
return exists === 1; return exists === 1;
} }
@@ -415,52 +431,52 @@ export class UpstashRedisStorage implements IStorage {
async changePassword(userName: string, newPassword: string): Promise<void> { async changePassword(userName: string, newPassword: string): Promise<void> {
// 简单存储明文密码,生产环境应加密 // 简单存储明文密码,生产环境应加密
await withRetry(() => await withRetry(() =>
this.client.set(this.userPwdKey(userName), newPassword) this._client.set(this.userPwdKey(userName), newPassword)
); );
} }
// 删除用户及其所有数据 // 删除用户及其所有数据
async deleteUser(userName: string): Promise<void> { async deleteUser(userName: string): Promise<void> {
// 删除用户密码 // 删除用户密码
await withRetry(() => this.client.del(this.userPwdKey(userName))); await withRetry(() => this._client.del(this.userPwdKey(userName)));
// 删除搜索历史 // 删除搜索历史
await withRetry(() => this.client.del(this.shKey(userName))); await withRetry(() => this._client.del(this.shKey(userName)));
// 删除播放记录(新hash结构) // 删除播放记录(新hash结构)
await withRetry(() => this.client.del(this.prHashKey(userName))); await withRetry(() => this._client.del(this.prHashKey(userName)));
// 删除旧的播放记录key(如果有) // 删除旧的播放记录key(如果有)
const playRecordPattern = `u:${userName}:pr:*`; const playRecordPattern = `u:${userName}:pr:*`;
const playRecordKeys = await withRetry(() => const playRecordKeys = await withRetry(() =>
this.client.keys(playRecordPattern) this._client.keys(playRecordPattern)
); );
if (playRecordKeys.length > 0) { if (playRecordKeys.length > 0) {
await withRetry(() => this.client.del(...playRecordKeys)); await withRetry(() => this._client.del(...playRecordKeys));
} }
// 删除收藏夹(新hash结构) // 删除收藏夹(新hash结构)
await withRetry(() => this.client.del(this.favHashKey(userName))); await withRetry(() => this._client.del(this.favHashKey(userName)));
// 删除旧的收藏key(如果有) // 删除旧的收藏key(如果有)
const favoritePattern = `u:${userName}:fav:*`; const favoritePattern = `u:${userName}:fav:*`;
const favoriteKeys = await withRetry(() => const favoriteKeys = await withRetry(() =>
this.client.keys(favoritePattern) this._client.keys(favoritePattern)
); );
if (favoriteKeys.length > 0) { if (favoriteKeys.length > 0) {
await withRetry(() => this.client.del(...favoriteKeys)); await withRetry(() => this._client.del(...favoriteKeys));
} }
// 删除跳过片头片尾配置(新hash结构) // 删除跳过片头片尾配置(新hash结构)
await withRetry(() => this.client.del(this.skipHashKey(userName))); await withRetry(() => this._client.del(this.skipHashKey(userName)));
// 删除旧的跳过配置key(如果有) // 删除旧的跳过配置key(如果有)
const skipConfigPattern = `u:${userName}:skip:*`; const skipConfigPattern = `u:${userName}:skip:*`;
const skipConfigKeys = await withRetry(() => const skipConfigKeys = await withRetry(() =>
this.client.keys(skipConfigPattern) this._client.keys(skipConfigPattern)
); );
if (skipConfigKeys.length > 0) { if (skipConfigKeys.length > 0) {
await withRetry(() => this.client.del(...skipConfigKeys)); await withRetry(() => this._client.del(...skipConfigKeys));
} }
} }
@@ -497,7 +513,7 @@ export class UpstashRedisStorage implements IStorage {
): Promise<void> { ): Promise<void> {
// 先检查用户是否已存在(原子性检查) // 先检查用户是否已存在(原子性检查)
const exists = await withRetry(() => const exists = await withRetry(() =>
this.client.exists(this.userInfoKey(userName)) this._client.exists(this.userInfoKey(userName))
); );
if (exists === 1) { if (exists === 1) {
throw new Error('用户已存在'); throw new Error('用户已存在');
@@ -521,17 +537,17 @@ export class UpstashRedisStorage implements IStorage {
if (oidcSub) { if (oidcSub) {
userInfo.oidcSub = oidcSub; userInfo.oidcSub = oidcSub;
// 创建OIDC映射 // 创建OIDC映射
await withRetry(() => this.client.set(this.oidcSubKey(oidcSub), userName)); await withRetry(() => this._client.set(this.oidcSubKey(oidcSub), userName));
} }
if (enabledApis && enabledApis.length > 0) { if (enabledApis && enabledApis.length > 0) {
userInfo.enabledApis = JSON.stringify(enabledApis); userInfo.enabledApis = JSON.stringify(enabledApis);
} }
await withRetry(() => this.client.hset(this.userInfoKey(userName), userInfo)); await withRetry(() => this._client.hset(this.userInfoKey(userName), userInfo));
// 添加到用户列表(Sorted Set,按注册时间排序) // 添加到用户列表(Sorted Set,按注册时间排序)
await withRetry(() => this.client.zadd(this.userListKey(), { await withRetry(() => this._client.zadd(this.userListKey(), {
score: createdAt, score: createdAt,
member: userName, member: userName,
})); }));
@@ -546,7 +562,7 @@ export class UpstashRedisStorage implements IStorage {
// 验证用户密码(新版本) // 验证用户密码(新版本)
async verifyUserV2(userName: string, password: string): Promise<boolean> { async verifyUserV2(userName: string, password: string): Promise<boolean> {
const userInfo = await withRetry(() => const userInfo = await withRetry(() =>
this.client.hgetall(this.userInfoKey(userName)) this._client.hgetall(this.userInfoKey(userName))
); );
if (!userInfo || !userInfo.password) { if (!userInfo || !userInfo.password) {
@@ -576,7 +592,7 @@ export class UpstashRedisStorage implements IStorage {
} }
const userInfo = await withRetry(() => const userInfo = await withRetry(() =>
this.client.hgetall(this.userInfoKey(userName)) this._client.hgetall(this.userInfoKey(userName))
); );
if (!userInfo || Object.keys(userInfo).length === 0) { if (!userInfo || Object.keys(userInfo).length === 0) {
@@ -696,7 +712,7 @@ export class UpstashRedisStorage implements IStorage {
userInfo.tags = JSON.stringify(updates.tags); userInfo.tags = JSON.stringify(updates.tags);
} else { } else {
// 删除tags字段 // 删除tags字段
await withRetry(() => this.client.hdel(this.userInfoKey(userName), 'tags')); await withRetry(() => this._client.hdel(this.userInfoKey(userName), 'tags'));
} }
} }
@@ -705,7 +721,7 @@ export class UpstashRedisStorage implements IStorage {
userInfo.enabledApis = JSON.stringify(updates.enabledApis); userInfo.enabledApis = JSON.stringify(updates.enabledApis);
} else { } else {
// 删除enabledApis字段 // 删除enabledApis字段
await withRetry(() => this.client.hdel(this.userInfoKey(userName), 'enabledApis')); await withRetry(() => this._client.hdel(this.userInfoKey(userName), 'enabledApis'));
} }
} }
@@ -713,15 +729,15 @@ export class UpstashRedisStorage implements IStorage {
const oldInfo = await this.getUserInfoV2(userName); const oldInfo = await this.getUserInfoV2(userName);
if (oldInfo?.oidcSub && oldInfo.oidcSub !== updates.oidcSub) { if (oldInfo?.oidcSub && oldInfo.oidcSub !== updates.oidcSub) {
// 删除旧的OIDC映射 // 删除旧的OIDC映射
await withRetry(() => this.client.del(this.oidcSubKey(oldInfo.oidcSub!))); await withRetry(() => this._client.del(this.oidcSubKey(oldInfo.oidcSub!)));
} }
userInfo.oidcSub = updates.oidcSub; userInfo.oidcSub = updates.oidcSub;
// 创建新的OIDC映射 // 创建新的OIDC映射
await withRetry(() => this.client.set(this.oidcSubKey(updates.oidcSub!), userName)); await withRetry(() => this._client.set(this.oidcSubKey(updates.oidcSub!), userName));
} }
if (Object.keys(userInfo).length > 0) { if (Object.keys(userInfo).length > 0) {
await withRetry(() => this.client.hset(this.userInfoKey(userName), userInfo)); await withRetry(() => this._client.hset(this.userInfoKey(userName), userInfo));
} }
// 清除缓存 // 清除缓存
@@ -732,7 +748,7 @@ export class UpstashRedisStorage implements IStorage {
async changePasswordV2(userName: string, newPassword: string): Promise<void> { async changePasswordV2(userName: string, newPassword: string): Promise<void> {
const hashedPassword = await this.hashPassword(newPassword); const hashedPassword = await this.hashPassword(newPassword);
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { password: hashedPassword }) this._client.hset(this.userInfoKey(userName), { password: hashedPassword })
); );
// 清除缓存 // 清除缓存
@@ -742,7 +758,7 @@ export class UpstashRedisStorage implements IStorage {
// 检查用户是否存在(新版本) // 检查用户是否存在(新版本)
async checkUserExistV2(userName: string): Promise<boolean> { async checkUserExistV2(userName: string): Promise<boolean> {
const exists = await withRetry(() => const exists = await withRetry(() =>
this.client.exists(this.userInfoKey(userName)) this._client.exists(this.userInfoKey(userName))
); );
return exists === 1; return exists === 1;
} }
@@ -750,7 +766,7 @@ export class UpstashRedisStorage implements IStorage {
// 通过OIDC Sub查找用户名 // 通过OIDC Sub查找用户名
async getUserByOidcSub(oidcSub: string): Promise<string | null> { async getUserByOidcSub(oidcSub: string): Promise<string | null> {
const userName = await withRetry(() => const userName = await withRetry(() =>
this.client.get(this.oidcSubKey(oidcSub)) this._client.get(this.oidcSubKey(oidcSub))
); );
return userName ? ensureString(userName) : null; return userName ? ensureString(userName) : null;
} }
@@ -763,7 +779,7 @@ export class UpstashRedisStorage implements IStorage {
let cursor: number | string = 0; let cursor: number | string = 0;
do { do {
const result = await withRetry(() => const result = await withRetry(() =>
this.client.scan(cursor as number, { match: 'user:*:info', count: 100 }) this._client.scan(cursor as number, { match: 'user:*:info', count: 100 })
); );
cursor = result[0]; cursor = result[0];
@@ -771,7 +787,7 @@ export class UpstashRedisStorage implements IStorage {
// 检查每个用户的 tags // 检查每个用户的 tags
for (const key of keys) { for (const key of keys) {
const userInfo = await withRetry(() => this.client.hgetall(key)); const userInfo = await withRetry(() => this._client.hgetall(key));
if (userInfo && userInfo.tags) { if (userInfo && userInfo.tags) {
const tags = JSON.parse(userInfo.tags as string); const tags = JSON.parse(userInfo.tags as string);
if (tags.includes(tagName)) { if (tags.includes(tagName)) {
@@ -804,7 +820,7 @@ export class UpstashRedisStorage implements IStorage {
total: number; total: number;
}> { }> {
// 获取总数 // 获取总数
let total = await withRetry(() => this.client.zcard(this.userListKey())); let total = await withRetry(() => this._client.zcard(this.userListKey()));
// 检查站长是否在数据库中(使用缓存) // 检查站长是否在数据库中(使用缓存)
let ownerInfo = null; let ownerInfo = null;
@@ -851,7 +867,7 @@ export class UpstashRedisStorage implements IStorage {
// 获取用户列表(按注册时间升序) // 获取用户列表(按注册时间升序)
const usernames = await withRetry(() => const usernames = await withRetry(() =>
this.client.zrange(this.userListKey(), actualOffset, actualOffset + actualLimit - 1) this._client.zrange(this.userListKey(), actualOffset, actualOffset + actualLimit - 1)
); );
const users = []; const users = [];
@@ -902,14 +918,14 @@ export class UpstashRedisStorage implements IStorage {
// 删除OIDC映射 // 删除OIDC映射
if (userInfo?.oidcSub) { if (userInfo?.oidcSub) {
await withRetry(() => this.client.del(this.oidcSubKey(userInfo.oidcSub!))); await withRetry(() => this._client.del(this.oidcSubKey(userInfo.oidcSub!)));
} }
// 删除用户信息Hash // 删除用户信息Hash
await withRetry(() => this.client.del(this.userInfoKey(userName))); await withRetry(() => this._client.del(this.userInfoKey(userName)));
// 从用户列表中移除 // 从用户列表中移除
await withRetry(() => this.client.zrem(this.userListKey(), userName)); await withRetry(() => this._client.zrem(this.userListKey(), userName));
// 删除用户的其他数据(播放记录、收藏等) // 删除用户的其他数据(播放记录、收藏等)
await this.deleteUser(userName); await this.deleteUser(userName);
@@ -925,7 +941,7 @@ export class UpstashRedisStorage implements IStorage {
async getSearchHistory(userName: string): Promise<string[]> { async getSearchHistory(userName: string): Promise<string[]> {
const result = await withRetry(() => const result = await withRetry(() =>
this.client.lrange(this.shKey(userName), 0, -1) this._client.lrange(this.shKey(userName), 0, -1)
); );
// 确保返回的都是字符串类型 // 确保返回的都是字符串类型
return ensureStringArray(result as any[]); return ensureStringArray(result as any[]);
@@ -934,19 +950,19 @@ export class UpstashRedisStorage implements IStorage {
async addSearchHistory(userName: string, keyword: string): Promise<void> { async addSearchHistory(userName: string, keyword: string): Promise<void> {
const key = this.shKey(userName); const key = this.shKey(userName);
// 先去重 // 先去重
await withRetry(() => this.client.lrem(key, 0, ensureString(keyword))); await withRetry(() => this._client.lrem(key, 0, ensureString(keyword)));
// 插入到最前 // 插入到最前
await withRetry(() => this.client.lpush(key, ensureString(keyword))); await withRetry(() => this._client.lpush(key, ensureString(keyword)));
// 限制最大长度 // 限制最大长度
await withRetry(() => this.client.ltrim(key, 0, SEARCH_HISTORY_LIMIT - 1)); await withRetry(() => this._client.ltrim(key, 0, SEARCH_HISTORY_LIMIT - 1));
} }
async deleteSearchHistory(userName: string, keyword?: string): Promise<void> { async deleteSearchHistory(userName: string, keyword?: string): Promise<void> {
const key = this.shKey(userName); const key = this.shKey(userName);
if (keyword) { if (keyword) {
await withRetry(() => this.client.lrem(key, 0, ensureString(keyword))); await withRetry(() => this._client.lrem(key, 0, ensureString(keyword)));
} else { } else {
await withRetry(() => this.client.del(key)); await withRetry(() => this._client.del(key));
} }
} }
@@ -955,7 +971,7 @@ export class UpstashRedisStorage implements IStorage {
// 从新版用户列表获取 // 从新版用户列表获取
const userListKey = this.userListKey(); const userListKey = this.userListKey();
const users = await withRetry(() => const users = await withRetry(() =>
this.client.zrange(userListKey, 0, -1) this._client.zrange(userListKey, 0, -1)
); );
const userList = users.map(u => ensureString(u)); const userList = users.map(u => ensureString(u));
@@ -974,12 +990,12 @@ export class UpstashRedisStorage implements IStorage {
} }
async getAdminConfig(): Promise<AdminConfig | null> { async getAdminConfig(): Promise<AdminConfig | null> {
const val = await withRetry(() => this.client.get(this.adminConfigKey())); const val = await withRetry(() => this._client.get(this.adminConfigKey()));
return val ? (val as AdminConfig) : null; return val ? (val as AdminConfig) : null;
} }
async setAdminConfig(config: AdminConfig): Promise<void> { async setAdminConfig(config: AdminConfig): Promise<void> {
await withRetry(() => this.client.set(this.adminConfigKey(), config)); await withRetry(() => this._client.set(this.adminConfigKey(), config));
} }
// ---------- 跳过片头片尾配置 ---------- // ---------- 跳过片头片尾配置 ----------
@@ -998,7 +1014,7 @@ export class UpstashRedisStorage implements IStorage {
): Promise<SkipConfig | null> { ): Promise<SkipConfig | null> {
const key = `${source}+${id}`; const key = `${source}+${id}`;
const val = await withRetry(() => const val = await withRetry(() =>
this.client.hget(this.skipHashKey(userName), key) this._client.hget(this.skipHashKey(userName), key)
); );
return val ? (val as SkipConfig) : null; return val ? (val as SkipConfig) : null;
} }
@@ -1011,7 +1027,7 @@ export class UpstashRedisStorage implements IStorage {
): Promise<void> { ): Promise<void> {
const key = `${source}+${id}`; const key = `${source}+${id}`;
await withRetry(() => await withRetry(() =>
this.client.hset(this.skipHashKey(userName), { [key]: config }) this._client.hset(this.skipHashKey(userName), { [key]: config })
); );
} }
@@ -1022,7 +1038,7 @@ export class UpstashRedisStorage implements IStorage {
): Promise<void> { ): Promise<void> {
const key = `${source}+${id}`; const key = `${source}+${id}`;
await withRetry(() => await withRetry(() =>
this.client.hdel(this.skipHashKey(userName), key) this._client.hdel(this.skipHashKey(userName), key)
); );
} }
@@ -1030,7 +1046,7 @@ export class UpstashRedisStorage implements IStorage {
userName: string userName: string
): Promise<{ [key: string]: SkipConfig }> { ): Promise<{ [key: string]: SkipConfig }> {
const hashData = await withRetry(() => const hashData = await withRetry(() =>
this.client.hgetall<Record<string, SkipConfig>>(this.skipHashKey(userName)) this._client.hgetall<Record<string, SkipConfig>>(this.skipHashKey(userName))
); );
return hashData || {}; return hashData || {};
@@ -1065,18 +1081,18 @@ export class UpstashRedisStorage implements IStorage {
} }
const pattern = `u:${userName}:skip:*`; const pattern = `u:${userName}:skip:*`;
const oldKeys: string[] = await withRetry(() => this.client.keys(pattern)); const oldKeys: string[] = await withRetry(() => this._client.keys(pattern));
if (oldKeys.length === 0) { if (oldKeys.length === 0) {
console.log(`用户 ${userName} 没有旧的跳过配置,标记为已迁移`); console.log(`用户 ${userName} 没有旧的跳过配置,标记为已迁移`);
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { skip_migrated: 'true' }) this._client.hset(this.userInfoKey(userName), { skip_migrated: 'true' })
); );
userInfoCache?.delete(userName); userInfoCache?.delete(userName);
return; return;
} }
const values = await withRetry(() => this.client.mget(oldKeys)); const values = await withRetry(() => this._client.mget(oldKeys));
const hashData: Record<string, SkipConfig> = {}; const hashData: Record<string, SkipConfig> = {};
oldKeys.forEach((key, index) => { oldKeys.forEach((key, index) => {
@@ -1092,16 +1108,16 @@ export class UpstashRedisStorage implements IStorage {
if (Object.keys(hashData).length > 0) { if (Object.keys(hashData).length > 0) {
await withRetry(() => await withRetry(() =>
this.client.hset(this.skipHashKey(userName), hashData) this._client.hset(this.skipHashKey(userName), hashData)
); );
console.log(`成功迁移 ${Object.keys(hashData).length} 条跳过配置到hash结构`); console.log(`成功迁移 ${Object.keys(hashData).length} 条跳过配置到hash结构`);
} }
await withRetry(() => this.client.del(...oldKeys)); await withRetry(() => this._client.del(...oldKeys));
console.log(`删除了 ${oldKeys.length} 个旧的跳过配置key`); console.log(`删除了 ${oldKeys.length} 个旧的跳过配置key`);
await withRetry(() => await withRetry(() =>
this.client.hset(this.userInfoKey(userName), { skip_migrated: 'true' }) this._client.hset(this.userInfoKey(userName), { skip_migrated: 'true' })
); );
userInfoCache?.delete(userName); userInfoCache?.delete(userName);
@@ -1113,7 +1129,7 @@ export class UpstashRedisStorage implements IStorage {
userName: string userName: string
): Promise<import('./types').DanmakuFilterConfig | null> { ): Promise<import('./types').DanmakuFilterConfig | null> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.get(this.danmakuFilterConfigKey(userName)) this._client.get(this.danmakuFilterConfigKey(userName))
); );
return val ? (val as import('./types').DanmakuFilterConfig) : null; return val ? (val as import('./types').DanmakuFilterConfig) : null;
} }
@@ -1123,13 +1139,13 @@ export class UpstashRedisStorage implements IStorage {
config: import('./types').DanmakuFilterConfig config: import('./types').DanmakuFilterConfig
): Promise<void> { ): Promise<void> {
await withRetry(() => await withRetry(() =>
this.client.set(this.danmakuFilterConfigKey(userName), config) this._client.set(this.danmakuFilterConfigKey(userName), config)
); );
} }
async deleteDanmakuFilterConfig(userName: string): Promise<void> { async deleteDanmakuFilterConfig(userName: string): Promise<void> {
await withRetry(() => await withRetry(() =>
this.client.del(this.danmakuFilterConfigKey(userName)) this._client.del(this.danmakuFilterConfigKey(userName))
); );
} }
@@ -1145,7 +1161,7 @@ export class UpstashRedisStorage implements IStorage {
} }
// 删除管理员配置 // 删除管理员配置
await withRetry(() => this.client.del(this.adminConfigKey())); await withRetry(() => this._client.del(this.adminConfigKey()));
console.log('所有数据已清空'); console.log('所有数据已清空');
} catch (error) { } catch (error) {
@@ -1161,7 +1177,7 @@ export class UpstashRedisStorage implements IStorage {
async getGlobalValue(key: string): Promise<string | null> { async getGlobalValue(key: string): Promise<string | null> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.get(this.globalValueKey(key)) this._client.get(this.globalValueKey(key))
); );
// Upstash 会自动反序列化 JSON,如果值是对象,需要重新序列化为字符串 // Upstash 会自动反序列化 JSON,如果值是对象,需要重新序列化为字符串
if (val === null) return null; if (val === null) return null;
@@ -1172,12 +1188,12 @@ export class UpstashRedisStorage implements IStorage {
async setGlobalValue(key: string, value: string): Promise<void> { async setGlobalValue(key: string, value: string): Promise<void> {
await withRetry(() => await withRetry(() =>
this.client.set(this.globalValueKey(key), value) this._client.set(this.globalValueKey(key), value)
); );
} }
async deleteGlobalValue(key: string): Promise<void> { async deleteGlobalValue(key: string): Promise<void> {
await withRetry(() => this.client.del(this.globalValueKey(key))); await withRetry(() => this._client.del(this.globalValueKey(key)));
} }
// ---------- 通知相关 ---------- // ---------- 通知相关 ----------
@@ -1191,7 +1207,7 @@ export class UpstashRedisStorage implements IStorage {
async getNotifications(userName: string): Promise<import('./types').Notification[]> { async getNotifications(userName: string): Promise<import('./types').Notification[]> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.get(this.notificationsKey(userName)) this._client.get(this.notificationsKey(userName))
); );
return val ? (val as import('./types').Notification[]) : []; return val ? (val as import('./types').Notification[]) : [];
} }
@@ -1207,7 +1223,7 @@ export class UpstashRedisStorage implements IStorage {
notifications.splice(100); notifications.splice(100);
} }
await withRetry(() => await withRetry(() =>
this.client.set(this.notificationsKey(userName), notifications) this._client.set(this.notificationsKey(userName), notifications)
); );
} }
@@ -1220,7 +1236,7 @@ export class UpstashRedisStorage implements IStorage {
if (notification) { if (notification) {
notification.read = true; notification.read = true;
await withRetry(() => await withRetry(() =>
this.client.set(this.notificationsKey(userName), notifications) this._client.set(this.notificationsKey(userName), notifications)
); );
} }
} }
@@ -1232,12 +1248,12 @@ export class UpstashRedisStorage implements IStorage {
const notifications = await this.getNotifications(userName); const notifications = await this.getNotifications(userName);
const filtered = notifications.filter((n) => n.id !== notificationId); const filtered = notifications.filter((n) => n.id !== notificationId);
await withRetry(() => await withRetry(() =>
this.client.set(this.notificationsKey(userName), filtered) this._client.set(this.notificationsKey(userName), filtered)
); );
} }
async clearAllNotifications(userName: string): Promise<void> { async clearAllNotifications(userName: string): Promise<void> {
await withRetry(() => this.client.del(this.notificationsKey(userName))); await withRetry(() => this._client.del(this.notificationsKey(userName)));
} }
async getUnreadNotificationCount(userName: string): Promise<number> { async getUnreadNotificationCount(userName: string): Promise<number> {
@@ -1247,7 +1263,7 @@ export class UpstashRedisStorage implements IStorage {
async getLastFavoriteCheckTime(userName: string): Promise<number> { async getLastFavoriteCheckTime(userName: string): Promise<number> {
const val = await withRetry(() => const val = await withRetry(() =>
this.client.get(this.lastFavoriteCheckKey(userName)) this._client.get(this.lastFavoriteCheckKey(userName))
); );
return val ? (val as number) : 0; return val ? (val as number) : 0;
} }
@@ -1257,7 +1273,7 @@ export class UpstashRedisStorage implements IStorage {
timestamp: number timestamp: number
): Promise<void> { ): Promise<void> {
await withRetry(() => await withRetry(() =>
this.client.set(this.lastFavoriteCheckKey(userName), timestamp) this._client.set(this.lastFavoriteCheckKey(userName), timestamp)
); );
} }
@@ -1271,42 +1287,42 @@ export class UpstashRedisStorage implements IStorage {
} }
async getAllMovieRequests(): Promise<import('./types').MovieRequest[]> { async getAllMovieRequests(): Promise<import('./types').MovieRequest[]> {
const data = await withRetry(() => this.client.hgetall(this.movieRequestsKey())); const data = await withRetry(() => this._client.hgetall(this.movieRequestsKey()));
if (!data) return []; if (!data) return [];
return Object.values(data) as import('./types').MovieRequest[]; return Object.values(data) as import('./types').MovieRequest[];
} }
async getMovieRequest(requestId: string): Promise<import('./types').MovieRequest | null> { async getMovieRequest(requestId: string): Promise<import('./types').MovieRequest | null> {
const val = await withRetry(() => this.client.hget(this.movieRequestsKey(), requestId)); const val = await withRetry(() => this._client.hget(this.movieRequestsKey(), requestId));
return val ? (val as import('./types').MovieRequest) : null; return val ? (val as import('./types').MovieRequest) : null;
} }
async createMovieRequest(request: import('./types').MovieRequest): Promise<void> { async createMovieRequest(request: import('./types').MovieRequest): Promise<void> {
await withRetry(() => this.client.hset(this.movieRequestsKey(), { [request.id]: request })); await withRetry(() => this._client.hset(this.movieRequestsKey(), { [request.id]: request }));
} }
async updateMovieRequest(requestId: string, updates: Partial<import('./types').MovieRequest>): Promise<void> { async updateMovieRequest(requestId: string, updates: Partial<import('./types').MovieRequest>): Promise<void> {
const existing = await this.getMovieRequest(requestId); const existing = await this.getMovieRequest(requestId);
if (!existing) throw new Error('Movie request not found'); if (!existing) throw new Error('Movie request not found');
const updated = { ...existing, ...updates }; const updated = { ...existing, ...updates };
await withRetry(() => this.client.hset(this.movieRequestsKey(), { [requestId]: updated })); await withRetry(() => this._client.hset(this.movieRequestsKey(), { [requestId]: updated }));
} }
async deleteMovieRequest(requestId: string): Promise<void> { async deleteMovieRequest(requestId: string): Promise<void> {
await withRetry(() => this.client.hdel(this.movieRequestsKey(), requestId)); await withRetry(() => this._client.hdel(this.movieRequestsKey(), requestId));
} }
async getUserMovieRequests(userName: string): Promise<string[]> { async getUserMovieRequests(userName: string): Promise<string[]> {
const val = await withRetry(() => this.client.smembers(this.userMovieRequestsKey(userName))); const val = await withRetry(() => this._client.smembers(this.userMovieRequestsKey(userName)));
return val ? ensureStringArray(val) : []; return val ? ensureStringArray(val) : [];
} }
async addUserMovieRequest(userName: string, requestId: string): Promise<void> { async addUserMovieRequest(userName: string, requestId: string): Promise<void> {
await withRetry(() => this.client.sadd(this.userMovieRequestsKey(userName), requestId)); await withRetry(() => this._client.sadd(this.userMovieRequestsKey(userName), requestId));
} }
async removeUserMovieRequest(userName: string, requestId: string): Promise<void> { async removeUserMovieRequest(userName: string, requestId: string): Promise<void> {
await withRetry(() => this.client.srem(this.userMovieRequestsKey(userName), requestId)); await withRetry(() => this._client.srem(this.userMovieRequestsKey(userName), requestId));
} }
} }