zhangwencui
2 天以前 a2002ba0e8aa2a0e4eee61b5205748d8f6e454fc
src/views/im/utils/pull.ts
¶Ô±ÈÐÂÎļþ
@@ -0,0 +1,155 @@
import { getDb } from './db'
/**
 * IM çŠ¶æ€äº‹ä»¶è¡¥å¿ï¼ˆå¢žé‡æ‹‰å–ï¼‰é€šç”¨ç¼–æŽ’
 *
 * å„业务模块(好友、群、申请、读位置)共用一套 update_time + id æ­£å‘游标:
 * ä»ŽæŒä¹…化游标出发循环拉取变更并合并进本地 store,直到某页不满(没有更多),
 * ç”¨äºŽ WebSocket æ¼æŽ¨åŽçš„兜底补偿(进入 IM、断线重连时触发)
 */
/** å¢žé‡æ‹‰å–游标:上次拉到的位置 */
export interface PullCursor {
  lastUpdateTime?: number
  lastId?: number
}
/** å¯ä½œä¸ºæ¸¸æ ‡çš„æ‹‰å–记录:服务端按 update_time + id è¿”回,客户端取最后一条推进游标 */
interface PullRecord {
  id: number
  updateTime?: number
}
/** å•次拉取条数(与后端 limit ä¸Šé™å¯¹é½ï¼‰ */
const PULL_PAGE_SIZE = 100
/** å•轮最多翻页数,兜底防御异常游标导致的死循环 */
const PULL_MAX_PAGES = 100
/** çŠ¶æ€äº‹ä»¶æ‹‰å–å›žæ‰«çª—å£ï¼Œè¦†ç›–åŒç§’å†…æ—§è¡Œæ›´æ–°å’Œå®¢æˆ·ç«¯ / æœåŠ¡ç«¯æ—¶é’Ÿç²¾åº¦å·® */
const PULL_OVERLAP_MS = 5000
/** æ¶ˆæ¯ç±» minId æ‹‰å–的单轮翻页上限,兜底防御异常游标死翻;消息量可能远大于状态事件,放宽到 1000 */
const MIN_ID_PULL_MAX_PAGES = 1000
/** è¯»å–某模块的拉取游标;无则返回空游标(首次拉全量) */
export async function getPullCursor(key: string): Promise<PullCursor> {
  return (await getDb().getSetting<PullCursor>(key)) ?? {}
}
/**
 * é€šç”¨å¢žé‡æ‹‰å–:从持久化游标出发,循环拉取并应用变更,直到某页不满
 *
 * @param cursorKey æ¸¸æ ‡åœ¨ settings é‡Œçš„ key(每个模块一份)
 * @param fetchPage æŒ‰æ¸¸æ ‡æ‹‰ä¸€é¡µï¼ˆè°ƒå¯¹åº” pull æŽ¥å£ï¼Œè¿”回 VO åˆ—表)
 * @param apply æŠŠä¸€é¡µå˜æ›´åˆå¹¶è¿›æœ¬åœ° store;返回 false è¡¨ç¤ºæœ¬é¡µæœªå®Œå…¨è½åœ°ï¼ˆå¦‚账号已切、依赖资源补拉失败),
 *              æ­¤æ—¶ä¸æŽ¨è¿›æ¸¸æ ‡å¹¶ç»ˆæ­¢æœ¬è½®ï¼Œé¿å…æ¸¸æ ‡è¶Šè¿‡æœªè½åœ°çš„记录、导致后续增量永久漏拉。可返回 Promise
 */
export async function runIncrementalPull<T extends PullRecord>(
  cursorKey: string,
  fetchPage: (params: { lastId?: number; lastUpdateTime?: number; limit: number }) => Promise<T[]>,
  apply: (records: T[]) => boolean | Promise<boolean>,
  isActive?: () => boolean
): Promise<void> {
  const storedCursor = await getPullCursor(cursorKey)
  const highWater = { ...storedCursor }
  let cursor =
    storedCursor.lastUpdateTime == null
      ? {}
      : { lastUpdateTime: Math.max(0, storedCursor.lastUpdateTime - PULL_OVERLAP_MS), lastId: 0 }
  for (let page = 0; page < PULL_MAX_PAGES; page++) {
    if (isActive && !isActive()) {
      return
    }
    const list = await fetchPage({
      lastUpdateTime: cursor.lastUpdateTime,
      lastId: cursor.lastId,
      limit: PULL_PAGE_SIZE
    })
    if (isActive && !isActive()) {
      return
    }
    if (list.length > 0) {
      // apply æœªå®Œå…¨è½åœ°ï¼ˆè¿”回 false)时直接终止:游标只能跟着已落地的数据走,否则会跳过本页记录
      if ((await apply(list)) === false) {
        return
      }
      if (isActive && !isActive()) {
        return
      }
      // æŽ¨è¿›æ¸¸æ ‡åˆ°æœ¬é¡µæœ€åŽä¸€æ¡å¹¶æŒä¹…化:下次从这里接着拉
      const last = list[list.length - 1]
      if (!last) {
        return
      }
      if (last.updateTime == null) {
        return
      }
      cursor = { lastUpdateTime: last.updateTime, lastId: last.id }
      if (
        highWater.lastUpdateTime == null ||
        cursor.lastUpdateTime > highWater.lastUpdateTime ||
        (cursor.lastUpdateTime === highWater.lastUpdateTime &&
          cursor.lastId > (highWater.lastId ?? 0))
      ) {
        highWater.lastUpdateTime = cursor.lastUpdateTime
        highWater.lastId = cursor.lastId
        await getDb().setSetting(cursorKey, highWater)
      }
    }
    // ä¸æ»¡ä¸€é¡µ = æ²¡æœ‰æ›´å¤šå˜æ›´
    if (list.length < PULL_PAGE_SIZE) {
      return
    }
  }
  console.warn(`[IM pull] ${cursorKey} è¾¾åˆ°å•轮翻页上限,提前结束本轮补偿`)
}
/**
 * æ¶ˆæ¯ç±»å¢žé‡æ‹‰å–通用编排:按单调 minId æ¸¸æ ‡å¾ªçŽ¯ç¿»é¡µï¼Œç›´åˆ°ç©ºé¡µæˆ–æ¸¸æ ‡ä¸å†å‰è¿›
 *
 * ä¸Ž runIncrementalPull çš„区别:消息游标是单调消息 id(非 update_time + id å¤åˆæ¸¸æ ‡ï¼‰ï¼Œä¸” messageMaxId çš„æŒä¹…化
 * è·Ÿéšæ¶ˆæ¯å…¥åº“在 applyPage å†…完成(messageStore),故本函数只负责「循环翻页机制」,不碰 settings,也不掺消息业务。
 *
 * @param initialMinId èµ·å§‹æ¸¸æ ‡ï¼šä¸Šæ¬¡å…¥åº“的最大消息 id
 * @param pageSize     å•页条数
 * @param fetchPage    æŒ‰ minId æ‹‰ä¸€é¡µ
 * @param applyPage    å¤„理本页(建会话 / åˆ†å‘ / å…¥åº“ + æŽ¨è¿› messageMaxId);返回 false è¡¨ç¤ºæœ¬é¡µæœªè½åœ°ï¼Œåœæ­¢ä¸”不推进游标
 * @param isActive     æ¯æ¬¡ await å‰åŽè‡ªæ£€ï¼šå–消 / åˆ‡è´¦å·æ—¶è¿”回 false,丢弃本批不入库、不再翻页
 * @param maxPages     å•轮翻页上限
 */
export async function runMinIdPull<T extends { id?: number }>(options: {
  applyPage: (records: T[], nextMinId?: number) => Promise<boolean> | Promise<void>
  fetchPage: (params: { minId: number; size: number }) => Promise<T[]>
  initialMinId: number
  isActive?: () => boolean
  maxPages?: number
  pageSize: number
}): Promise<void> {
  const { initialMinId, pageSize, fetchPage, applyPage, isActive } = options
  const maxPages = options.maxPages ?? MIN_ID_PULL_MAX_PAGES
  let minId = initialMinId || 0
  for (let page = 0; page < maxPages; page++) {
    if (isActive && !isActive()) {
      return
    }
    const list = await fetchPage({ minId, size: pageSize })
    // æ‹‰å–期间取消 / åˆ‡è´¦å·ï¼šä¸¢å¼ƒæœ¬æ‰¹ä¸å…¥åº“,也不再翻页
    if (isActive && !isActive()) {
      return
    }
    if (!list || list.length === 0) {
      return
    }
    // æœ¬æ‰¹æœ€å¤§æ¶ˆæ¯ id ä½œä¸ºä¸‹æ¬¡æ¸¸æ ‡ï¼›æ— æœ‰æ•ˆ id åˆ™æœ¬æ‰¹ apply åŽåœï¼ˆæ— æ³•继续翻页)
    const validIds = list.map((record) => record.id).filter((id): id is number => id != null)
    const nextMinId = validIds.length > 0 ? Math.max(...validIds) : undefined
    // applyPage è¿”回 false:本页未落地(如入库失败),不推进游标并终止,避免漏消息
    if ((await applyPage(list, nextMinId)) === false) {
      return
    }
    // æ— æœ‰æ•ˆ id,或游标没前进(后端契约是 id > minId,理论不会出现):停,防御死翻
    if (nextMinId == null || nextMinId <= minId) {
      return
    }
    minId = nextMinId
  }
  console.warn('[IM pull] runMinIdPull è¾¾åˆ°å•轮翻页上限,提前结束本轮')
}