| ¶Ô±ÈÐÂÎļþ |
| | |
| | | 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 è¾¾å°å轮翻页ä¸éï¼æåç»ææ¬è½®') |
| | | } |