| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119 |
- 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<number> {
- 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<void> {
- 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<MessageSyncResult> {
- 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;
- }
|