message-sync.service.ts 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119
  1. import Parse from '../../../../shared/db/parse-client.js';
  2. import { syncMsgPage } from './qiwe-api.service.js';
  3. import { persistSyncMsgItem } from './message-persist.service.js';
  4. import { backfillGroupChatsFromMessages } from './groups.service.js';
  5. const DEFAULT_PAGE_LIMIT = 50;
  6. const DEFAULT_MAX_PAGES = 30;
  7. function cursorKey(guid: string): string {
  8. return `msg_sync_seq:${guid}`;
  9. }
  10. export async function getMessageSyncCursor(guid: string): Promise<number> {
  11. const query = new Parse.Query('SystemMigration');
  12. query.equalTo('key', cursorKey(guid));
  13. const row = await query.first({ useMasterKey: true });
  14. return row?.get('valueNumber') ?? 0;
  15. }
  16. export async function setMessageSyncCursor(guid: string, msgSeq: number): Promise<void> {
  17. const query = new Parse.Query('SystemMigration');
  18. query.equalTo('key', cursorKey(guid));
  19. let row = await query.first({ useMasterKey: true });
  20. if (!row) {
  21. row = new Parse.Object('SystemMigration');
  22. row.set('key', cursorKey(guid));
  23. }
  24. row.set('valueNumber', msgSeq);
  25. row.set('completedAt', new Date());
  26. await row.save(null, { useMasterKey: true });
  27. }
  28. export interface MessageSyncResult {
  29. guid: string;
  30. startSeq: number;
  31. endSeq: number;
  32. pages: number;
  33. created: number;
  34. skipped: number;
  35. hasMore: boolean;
  36. groupsBackfilled?: number;
  37. errors: string[];
  38. }
  39. export async function syncMessagesFromQiWe(options: {
  40. guid: string;
  41. msgSeq?: number;
  42. limit?: number;
  43. maxPages?: number;
  44. resetCursor?: boolean;
  45. }): Promise<MessageSyncResult> {
  46. const guid = options.guid;
  47. const limit = options.limit ?? DEFAULT_PAGE_LIMIT;
  48. const maxPages = options.maxPages ?? DEFAULT_MAX_PAGES;
  49. let msgSeq = options.resetCursor ? 0 : (options.msgSeq ?? await getMessageSyncCursor(guid));
  50. const result: MessageSyncResult = {
  51. guid,
  52. startSeq: msgSeq,
  53. endSeq: msgSeq,
  54. pages: 0,
  55. created: 0,
  56. skipped: 0,
  57. hasMore: false,
  58. errors: [],
  59. };
  60. console.log(`[MsgSync] 开始同步 guid=${guid} msgSeq=${msgSeq} limit=${limit}`);
  61. while (result.pages < maxPages) {
  62. let pageData;
  63. try {
  64. pageData = await syncMsgPage(guid, msgSeq, limit);
  65. } catch (err: unknown) {
  66. const message = err instanceof Error ? err.message : String(err);
  67. result.errors.push(message);
  68. break;
  69. }
  70. result.pages++;
  71. const list = pageData.syncMsgList || [];
  72. for (const raw of list) {
  73. try {
  74. const status = await persistSyncMsgItem(guid, raw);
  75. if (status === 'created') result.created++;
  76. else result.skipped++;
  77. } catch (err: unknown) {
  78. const message = err instanceof Error ? err.message : String(err);
  79. result.errors.push(`persist: ${message}`);
  80. }
  81. }
  82. const nextSeq = pageData.travelSyncKey ?? msgSeq;
  83. result.endSeq = nextSeq;
  84. await setMessageSyncCursor(guid, nextSeq);
  85. const hasMore = pageData.hasMore === 1 || pageData.hasMore === true;
  86. result.hasMore = hasMore;
  87. console.log(
  88. `[MsgSync] 第 ${result.pages} 页: batch=${list.length} created=${result.created} skipped=${result.skipped} nextSeq=${nextSeq} hasMore=${hasMore}`,
  89. );
  90. if (!hasMore) break;
  91. msgSeq = nextSeq;
  92. }
  93. const groupsBackfilled = await backfillGroupChatsFromMessages(guid);
  94. result.groupsBackfilled = groupsBackfilled;
  95. console.log(
  96. `[MsgSync] 完成 pages=${result.pages} created=${result.created} skipped=${result.skipped} endSeq=${result.endSeq} groupsBackfilled=${groupsBackfilled}`,
  97. );
  98. return result;
  99. }