import Parse from '../../../../shared/db/parse-client.js'; import { syncMsgPage } from './qiwe-api.service.js'; import { persistSyncMsgItem } from './message-persist.service.js'; import { backfillGroupChatsFromMessages } from './groups.service.js'; const DEFAULT_PAGE_LIMIT = 50; const DEFAULT_MAX_PAGES = 30; function cursorKey(guid: string): string { return `msg_sync_seq:${guid}`; } export async function getMessageSyncCursor(guid: string): Promise { const query = new Parse.Query('SystemMigration'); query.equalTo('key', cursorKey(guid)); const row = await query.first({ useMasterKey: true }); return row?.get('valueNumber') ?? 0; } export async function setMessageSyncCursor(guid: string, msgSeq: number): Promise { const query = new Parse.Query('SystemMigration'); query.equalTo('key', cursorKey(guid)); let row = await query.first({ useMasterKey: true }); if (!row) { row = new Parse.Object('SystemMigration'); row.set('key', cursorKey(guid)); } row.set('valueNumber', msgSeq); row.set('completedAt', new Date()); await row.save(null, { useMasterKey: true }); } export interface MessageSyncResult { guid: string; startSeq: number; endSeq: number; pages: number; created: number; skipped: number; hasMore: boolean; groupsBackfilled?: number; errors: string[]; } export async function syncMessagesFromQiWe(options: { guid: string; msgSeq?: number; limit?: number; maxPages?: number; resetCursor?: boolean; }): Promise { const guid = options.guid; const limit = options.limit ?? DEFAULT_PAGE_LIMIT; const maxPages = options.maxPages ?? DEFAULT_MAX_PAGES; let msgSeq = options.resetCursor ? 0 : (options.msgSeq ?? await getMessageSyncCursor(guid)); const result: MessageSyncResult = { guid, startSeq: msgSeq, endSeq: msgSeq, pages: 0, created: 0, skipped: 0, hasMore: false, errors: [], }; console.log(`[MsgSync] 开始同步 guid=${guid} msgSeq=${msgSeq} limit=${limit}`); while (result.pages < maxPages) { let pageData; try { pageData = await syncMsgPage(guid, msgSeq, limit); } catch (err: unknown) { const message = err instanceof Error ? err.message : String(err); result.errors.push(message); break; } result.pages++; const list = pageData.syncMsgList || []; for (const raw of list) { try { const status = await persistSyncMsgItem(guid, raw); if (status === 'created') result.created++; else result.skipped++; } catch (err: unknown) { const message = err instanceof Error ? err.message : String(err); result.errors.push(`persist: ${message}`); } } const nextSeq = pageData.travelSyncKey ?? msgSeq; result.endSeq = nextSeq; await setMessageSyncCursor(guid, nextSeq); const hasMore = pageData.hasMore === 1 || pageData.hasMore === true; result.hasMore = hasMore; console.log( `[MsgSync] 第 ${result.pages} 页: batch=${list.length} created=${result.created} skipped=${result.skipped} nextSeq=${nextSeq} hasMore=${hasMore}`, ); if (!hasMore) break; msgSeq = nextSeq; } const groupsBackfilled = await backfillGroupChatsFromMessages(guid); result.groupsBackfilled = groupsBackfilled; console.log( `[MsgSync] 完成 pages=${result.pages} created=${result.created} skipped=${result.skipped} endSeq=${result.endSeq} groupsBackfilled=${groupsBackfilled}`, ); return result; }