gaoluyang
昨天 b64a0deae5b5d33f9e20671a68936b27f0b9b00b
src/views/im/utils/db.ts
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,541 @@
import type { MessageDO, SettingDO } from '../home/types'
import { toRaw } from 'vue'
import { getCurrentUserId } from '#/views/im/utils/auth'
import { ImConversationType } from './constants'
export const DB_SCHEMA_VERSION = 2
export type DbStoreName =
  | 'channels'
  | 'conversationReads'
  | 'conversations'
  | 'friendRequests'
  | 'friends'
  | 'groupMembers'
  | 'groupRequests'
  | 'groups'
  | 'messages'
  | 'settings'
export type DbTransaction = IDBTransaction
/** IM æœ¬åœ°å­˜å‚¨ key */
export const StorageKeys = {
  localStorage: {
    /** ä¾§è¾¹æ å®½åº¦ï¼Œä¸‰ä¸ª Tab å…±ç”¨ä¸€ä»½è®°å¿† */
    asideWidth: 'im:aside',
    /** ä¼šè¯åˆ—表置顶折叠展开态 */
    conversationPinnedExpanded: 'im:conversation:pinnedExpanded'
  },
  settings: {
    /** ç§èŠæ¶ˆæ¯æ‹‰å–游标 */
    privateMessageMaxId: 'privateMessageMaxId',
    /** ç¾¤èŠæ¶ˆæ¯æ‹‰å–游标 */
    groupMessageMaxId: 'groupMessageMaxId',
    /** é¢‘道消息拉取游标 */
    channelMessageMaxId: 'channelMessageMaxId',
    /** æœ€è¿‘转发会话 key åˆ—表 */
    recentForwardConversationKeys: 'recentForwardConversationKeys',
    // çŠ¶æ€äº‹ä»¶è¡¥å¿å¢žé‡æ‹‰å–æ¸¸æ ‡ï¼šä¸Žä¸Šé¢æ¶ˆæ¯ maxId æ¸¸æ ‡å…±ç”¨åŒä¸€ settings keyspace,统一登记在此避免撞 key;
    // èµ° update_time + id å¤åˆæ¸¸æ ‡ï¼ˆéžå•条 maxId),故用 PullCursor åŽç¼€åŒºåˆ†è¯­ä¹‰
    /** å¥½å‹å…³ç³»å¢žé‡æ‹‰å–游标 */
    friendPullCursor: 'friendPullCursor',
    /** å¥½å‹ç”³è¯·å¢žé‡æ‹‰å–游标 */
    friendRequestPullCursor: 'friendRequestPullCursor',
    /** åŠ ç¾¤ç”³è¯·å¢žé‡æ‹‰å–æ¸¸æ ‡ */
    groupRequestPullCursor: 'groupRequestPullCursor',
    /** ä¼šè¯è¯»ä½ç½®å¢žé‡æ‹‰å–游标 */
    conversationReadPullCursor: 'conversationReadPullCursor'
  }
} as const
let currentDb: IDBDatabase | null = null
let currentUserId: null | number = null
let currentSession = 0
/** æ ¡éªŒå½“前 IM IndexedDB session ä»æœ‰æ•ˆ */
export function isCurrentDbSession(session: number): boolean {
  return session === currentSession
}
/** èŽ·å–å½“å‰ IM IndexedDB session */
export function getDbSession(): number {
  return currentSession
}
/** æ‹¼æŽ¥å½“前用户 IM DB åç§° */
function getDbName(userId: number): string {
  return `im:${userId}`
}
/** åŒ…装 IndexedDB request */
function requestToPromise<T = unknown>(request: IDBRequest<T>): Promise<T> {
  return new Promise((resolve, reject) => {
    request.addEventListener('success', () => resolve(request.result))
    request.addEventListener('error', () => reject(request.error))
  })
}
/** ç­‰å¾…事务完成 */
function transactionDone(transaction: DbTransaction): Promise<void> {
  return new Promise((resolve, reject) => {
    transaction.addEventListener('complete', () => resolve())
    transaction.addEventListener('error', () => reject(transaction.error))
    transaction.addEventListener('abort', () => reject(transaction.error))
  })
}
/** åˆ›å»ºç´¢å¼• */
function createIndex(
  store: IDBObjectStore,
  name: string,
  keyPath: string | string[],
  options?: IDBIndexParameters
) {
  if (!store.indexNames.contains(name)) {
    store.createIndex(name, keyPath, options)
  }
}
/** åˆå§‹åŒ– schema */
function upgradeSchema(db: IDBDatabase) {
  if (!db.objectStoreNames.contains('conversations')) {
    const store = db.createObjectStore('conversations', { keyPath: 'clientConversationId' })
    createIndex(store, 'lastSendTime', 'lastSendTime')
  }
  if (!db.objectStoreNames.contains('conversationReads')) {
    const store = db.createObjectStore('conversationReads', { keyPath: 'clientConversationId' })
    createIndex(store, 'conversationType+targetId', ['conversationType', 'targetId'], {
      unique: true
    })
  }
  if (!db.objectStoreNames.contains('messages')) {
    const store = db.createObjectStore('messages', { keyPath: 'messageKey' })
    createIndex(store, 'clientConversationId', 'clientConversationId')
    createIndex(store, 'clientConversationId+sendTime', ['clientConversationId', 'sendTime'])
    createIndex(store, 'clientMessageId', 'clientMessageId', { unique: true })
  }
  if (!db.objectStoreNames.contains('friends')) {
    const store = db.createObjectStore('friends', { keyPath: 'id' })
    createIndex(store, 'friendUserId', 'friendUserId', { unique: true })
    createIndex(store, 'status', 'status')
  }
  if (!db.objectStoreNames.contains('friendRequests')) {
    const store = db.createObjectStore('friendRequests', { keyPath: 'id' })
    createIndex(store, 'status', 'status')
    createIndex(store, 'createTime', 'createTime')
  }
  if (!db.objectStoreNames.contains('groups')) {
    const store = db.createObjectStore('groups', { keyPath: 'id' })
    createIndex(store, 'name', 'name')
    createIndex(store, 'status', 'status')
  }
  if (!db.objectStoreNames.contains('groupMembers')) {
    const store = db.createObjectStore('groupMembers', { keyPath: 'id' })
    createIndex(store, 'groupId', 'groupId')
    createIndex(store, 'groupId+userId', ['groupId', 'userId'], { unique: true })
  }
  if (!db.objectStoreNames.contains('groupRequests')) {
    const store = db.createObjectStore('groupRequests', { keyPath: 'id' })
    createIndex(store, 'status', 'status')
    createIndex(store, 'createTime', 'createTime')
  }
  if (!db.objectStoreNames.contains('channels')) {
    const store = db.createObjectStore('channels', { keyPath: 'id' })
    createIndex(store, 'status', 'status')
    createIndex(store, 'sort', 'sort')
  }
  if (!db.objectStoreNames.contains('settings')) {
    db.createObjectStore('settings', { keyPath: 'key' })
  }
}
/** æ‰“å¼€ IM IndexedDB */
function openDb(name: string): Promise<IDBDatabase> {
  return new Promise((resolve, reject) => {
    const request = indexedDB.open(name, DB_SCHEMA_VERSION)
    // åˆ›å»ºæˆ–升级对象仓库
    request.addEventListener('upgradeneeded', () => upgradeSchema(request.result))
    // è¿”回可复用连接
    request.addEventListener('success', () => resolve(request.result))
    request.addEventListener('error', () => reject(request.error))
  })
}
/** åˆå§‹åŒ–当前用户 IM DB */
export async function initDb(): Promise<void> {
  const userId = getCurrentUserId()
  if (!Number.isFinite(userId) || userId <= 0) {
    throw new Error('当前用户不存在,无法初始化 IM DB')
  }
  if (currentDb && currentUserId === userId) {
    return
  }
  currentDb?.close()
  currentSession++
  currentUserId = userId
  currentDb = await openDb(getDbName(userId))
}
/** å…³é—­å½“前 IM DB è¿žæŽ¥ */
function closeDbConnection() {
  currentDb?.close()
  currentDb = null
  currentUserId = null
}
/** èŽ·å–å½“å‰ IM DB */
function getRawDb(): IDBDatabase {
  if (!currentDb) {
    throw new Error('IM DB æœªåˆå§‹åŒ–')
  }
  return currentDb
}
/** æ ¡éªŒå•次写入 session */
function guardSession(session: number) {
  if (!isCurrentDbSession(session)) {
    throw new Error('IM DB session å·²å¤±æ•ˆ')
  }
}
/** å…‹éš†å¯å…¥åº“对象 */
function toDbValue<T>(value: T): T {
  return cloneDbValue(value) as T
}
/** è½¬æ¢ä¸º IndexedDB å¯å…‹éš†å¯¹è±¡ */
function cloneDbValue(value: unknown): unknown {
  const raw = toRaw(value)
  if (Array.isArray(raw)) {
    return raw.map((item) => cloneDbValue(item))
  }
  if (!raw || typeof raw !== 'object') {
    return raw
  }
  const prototype = Object.getPrototypeOf(raw)
  if (prototype !== Object.prototype && prototype !== null) {
    return raw
  }
  return Object.fromEntries(
    Object.entries(raw as Record<string, unknown>).map(([key, item]) => [key, cloneDbValue(item)])
  )
}
class DbClient {
  /** æ¸…空 store è®°å½• */
  async clearStore(storeName: DbStoreName, tx?: DbTransaction): Promise<void> {
    if (tx) {
      await requestToPromise(tx.objectStore(storeName).clear())
      return
    }
    await this.transaction([storeName], 'readwrite', (tx) => this.clearStore(storeName, tx))
  }
  /** åˆ é™¤è®°å½• */
  async delete(storeName: DbStoreName, key: IDBValidKey, tx?: DbTransaction): Promise<void> {
    if (tx) {
      await requestToPromise(tx.objectStore(storeName).delete(key))
      return
    }
    await this.transaction([storeName], 'readwrite', (tx) => this.delete(storeName, key, tx))
  }
  /** æŒ‰ç´¢å¼•删除记录 */
  async deleteByIndex(
    storeName: DbStoreName,
    indexName: string,
    query: IDBKeyRange | IDBValidKey,
    tx?: DbTransaction
  ): Promise<void> {
    if (!tx) {
      await this.transaction([storeName], 'readwrite', (tx) =>
        this.deleteByIndex(storeName, indexName, query, tx)
      )
      return
    }
    const index = tx.objectStore(storeName).index(indexName)
    await new Promise<void>((resolve, reject) => {
      const request = index.openCursor(query)
      request.addEventListener('error', () => reject(request.error))
      request.addEventListener('success', () => {
        const cursor = request.result
        if (!cursor) {
          resolve()
          return
        }
        cursor.delete()
        cursor.continue()
      })
    })
  }
  /** èŽ·å–å•æ¡è®°å½• */
  async get<T>(
    storeName: DbStoreName,
    key: IDBValidKey,
    tx?: DbTransaction
  ): Promise<T | undefined> {
    if (tx) {
      return requestToPromise<T | undefined>(tx.objectStore(storeName).get(key))
    }
    return this.transaction<T | undefined>([storeName], 'readonly', (tx) =>
      this.get<T>(storeName, key, tx)
    )
  }
  /** èŽ·å– store å…¨é‡è®°å½• */
  async getAll<T>(storeName: DbStoreName, tx?: DbTransaction): Promise<T[]> {
    if (tx) {
      return requestToPromise<T[]>(tx.objectStore(storeName).getAll())
    }
    return this.transaction<T[]>([storeName], 'readonly', (tx) => this.getAll<T>(storeName, tx))
  }
  /** æŒ‰ç´¢å¼•获取记录列表 */
  async getAllByIndex<T>(
    storeName: DbStoreName,
    indexName: string,
    query?: IDBKeyRange | IDBValidKey,
    tx?: DbTransaction
  ): Promise<T[]> {
    if (tx) {
      return requestToPromise<T[]>(tx.objectStore(storeName).index(indexName).getAll(query))
    }
    return this.transaction<T[]>([storeName], 'readonly', (tx) =>
      this.getAllByIndex<T>(storeName, indexName, query, tx)
    )
  }
  /** æŒ‰å”¯ä¸€ç´¢å¼•获取单条记录 */
  async getByIndex<T>(
    storeName: DbStoreName,
    indexName: string,
    query: IDBKeyRange | IDBValidKey,
    tx?: DbTransaction
  ): Promise<T | undefined> {
    if (tx) {
      return requestToPromise<T | undefined>(tx.objectStore(storeName).index(indexName).get(query))
    }
    return this.transaction<T | undefined>([storeName], 'readonly', (tx) =>
      this.getByIndex<T>(storeName, indexName, query, tx)
    )
  }
  /** æŒ‰ä¼šè¯åˆ†é¡µèŽ·å–æ¶ˆæ¯ */
  async getMessageListByConversation(
    clientConversationId: string,
    options?: { beforeSendTime?: number; limit?: number },
    tx?: DbTransaction
  ): Promise<MessageDO[]> {
    const limit = options?.limit ?? 50
    const upper = options?.beforeSendTime ?? Number.MAX_SAFE_INTEGER
    const range = IDBKeyRange.bound(
      [clientConversationId, 0],
      [clientConversationId, upper],
      false,
      true
    )
    const read = async (tx: DbTransaction): Promise<MessageDO[]> => {
      const index = tx.objectStore('messages').index('clientConversationId+sendTime')
      const out: MessageDO[] = []
      await new Promise<void>((resolve, reject) => {
        // ä»Žæ–°åˆ°æ—§è¯»å–一页
        const request = index.openCursor(range, 'prev')
        request.addEventListener('error', () => reject(request.error))
        request.addEventListener('success', () => {
          const cursor = request.result
          if (!cursor || out.length >= limit) {
            resolve()
            return
          }
          out.push(cursor.value as MessageDO)
          cursor.continue()
        })
      })
      // æ°”泡渲染需要按时间升序
      return out.toReversed()
    }
    if (tx) {
      return read(tx)
    }
    return this.transaction<MessageDO[]>(['messages'], 'readonly', read)
  }
  /** è¯»å–设置 */
  async getSetting<T>(key: string, tx?: DbTransaction): Promise<T | undefined> {
    const item = await this.get<SettingDO<T>>('settings', key, tx)
    return item?.value
  }
  /** å†™å…¥è®°å½• */
  async put<T>(storeName: DbStoreName, value: T, tx?: DbTransaction): Promise<void> {
    if (tx) {
      await requestToPromise(tx.objectStore(storeName).put(toDbValue(value)))
      return
    }
    await this.transaction([storeName], 'readwrite', (tx) => this.put(storeName, value, tx))
  }
  /** å†™å…¥è®¾ç½® */
  async setSetting<T>(key: string, value: T, tx?: DbTransaction): Promise<void> {
    await this.put<SettingDO<T>>('settings', { key, value, updateTime: Date.now() }, tx)
  }
  /** æ‰§è¡Œäº‹åŠ¡ */
  async transaction<T>(
    storeNames: DbStoreName[],
    mode: IDBTransactionMode,
    runner: (tx: DbTransaction) => Promise<T>
  ): Promise<T> {
    // å¼€å¯äº‹åŠ¡å‰æ ¡éªŒ session
    const session = getDbSession()
    guardSession(session)
    const tx = getRawDb().transaction(storeNames, mode)
    const done = transactionDone(tx)
    let result: T
    try {
      // äº‹åŠ¡å†…åªæ‰§è¡Œ IndexedDB request é“¾
      result = await runner(tx)
    } catch (error) {
      try {
        tx.abort()
      } catch {}
      await done.catch(() => undefined)
      throw error
    }
    // commit åŽå†æ¬¡æ ¡éªŒ session
    await done
    guardSession(session)
    return result
  }
}
const dbClient = new DbClient()
/** èŽ·å–å½“å‰ IM DB client */
export function getDb(): DbClient {
  return dbClient
}
/** å½“前用户会话主键 */
export function getClientConversationId(type: number, targetId: number): string {
  return `${type}:${targetId}`
}
/** è§£æžå½“前用户会话主键 */
export function parseClientConversationId(
  clientConversationId: string
): null | { targetId: number; type: number; } {
  const [typeText, targetIdText] = clientConversationId.split(':')
  const type = Number(typeText)
  const targetId = Number(targetIdText)
  if (!Number.isFinite(type) || !Number.isFinite(targetId) || targetId <= 0) {
    return null
  }
  return { type, targetId }
}
/** æœåŠ¡ç«¯æ¶ˆæ¯ä¸»é”® */
export function getServerMessageKey(conversationType: number, id: number): string {
  return `${conversationType}:${id}`
}
/** å®¢æˆ·ç«¯ä¸´æ—¶æ¶ˆæ¯ä¸»é”® */
export function getClientMessageKey(clientMessageId: string): string {
  return `client:${clientMessageId}`
}
/** è§£æžæœ¬åœ°æ¶ˆæ¯ä¸»é”® */
export function parseMessageKey(
  messageKey: string
):
  | null
  | { clientMessageId: string; kind: 'client'; }
  | { conversationType: number; id: number; kind: 'server'; } {
  if (!messageKey) {
    return null
  }
  if (messageKey.startsWith('client:')) {
    const clientMessageId = messageKey.slice('client:'.length)
    return clientMessageId ? { kind: 'client', clientMessageId } : null
  }
  const [conversationTypeText, idText] = messageKey.split(':')
  const conversationType = Number(conversationTypeText)
  const id = Number(idText)
  if (!Number.isFinite(conversationType) || !Number.isFinite(id) || id <= 0) {
    return null
  }
  return { kind: 'server', conversationType, id }
}
/** æ›´æ–°æ¶ˆæ¯æ‹‰å–游标 */
export async function setMessageMaxId(
  conversationType: number,
  maxId: number | undefined,
  tx?: DbTransaction
): Promise<void> {
  if (!maxId) {
    return
  }
  let key: string
  switch (conversationType) {
    case ImConversationType.CHANNEL: {
      key = StorageKeys.settings.channelMessageMaxId
      break
    }
    case ImConversationType.GROUP: {
      key = StorageKeys.settings.groupMessageMaxId
      break
    }
    case ImConversationType.PRIVATE: {
      key = StorageKeys.settings.privateMessageMaxId
      break
    }
    default: {
      throw new Error(`未知 IM ä¼šè¯ç±»åž‹ï¼š${conversationType}`)
    }
  }
  const db = getDb()
  const current = (await db.getSetting<number>(key, tx)) || 0
  if (maxId > current) {
    await db.setSetting(key, maxId, tx)
  }
}
/** åœæ­¢å½“前 IM DB session */
export async function stopRequests(): Promise<void> {
  currentSession++
  const [
    { useMessageStoreWithOut },
    { useConversationStoreWithOut },
    { useFriendStoreWithOut },
    { useGroupStoreWithOut },
    { useChannelStoreWithOut },
    { useGroupRequestStoreWithOut },
    { useFaceStoreWithOut },
    { useRtcStore }
  ] = await Promise.all([
    import('../home/store/messageStore'),
    import('../home/store/conversationStore'),
    import('../home/store/friendStore'),
    import('../home/store/groupStore'),
    import('../home/store/channelStore'),
    import('../home/store/groupRequestStore'),
    import('../home/store/faceStore'),
    import('../home/store/rtcStore')
  ])
  useMessageStoreWithOut().clear()
  useConversationStoreWithOut().clear()
  useFriendStoreWithOut().clear()
  useGroupStoreWithOut().clear()
  useChannelStoreWithOut().clear()
  useGroupRequestStoreWithOut().clear()
  useFaceStoreWithOut().clear()
  useRtcStore().reset()
  useRtcStore().clearGroupCallCache()
  closeDbConnection()
}