personal-polling.mjs 6.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177
  1. import path from 'node:path';
  2. import { spawn } from 'node:child_process';
  3. import { PACKAGE_ROOT } from './config-loader.mjs';
  4. let processorBridgePromise = null;
  5. function loadProcessorBridge() {
  6. if (!processorBridgePromise) processorBridgePromise = import('./processor-bridge.mjs');
  7. return processorBridgePromise;
  8. }
  9. const defaultListenerApi = {
  10. getStatus: async () => (await loadProcessorBridge()).getPersonalListenerStatus(),
  11. recoverLogin: async () => (await loadProcessorBridge()).recoverPersonalLogin(),
  12. start: async () => (await loadProcessorBridge()).startPersonalListener(),
  13. stop: async options => (await loadProcessorBridge()).stopPersonalListener(options),
  14. };
  15. export class PersonalPollingRuntime {
  16. constructor({ config, workspaceRoot, onState = () => {}, listenerApi = {} }) {
  17. this.config = config;
  18. this.workspaceRoot = workspaceRoot;
  19. this.onState = onState;
  20. this.listenerApi = {
  21. getStatus: listenerApi.getStatus || defaultListenerApi.getStatus,
  22. recoverLogin: listenerApi.recoverLogin || defaultListenerApi.recoverLogin,
  23. start: listenerApi.start || defaultListenerApi.start,
  24. stop: listenerApi.stop || defaultListenerApi.stop,
  25. };
  26. this.running = false;
  27. this.listenerStarted = false;
  28. this.friendWorker = null;
  29. this.loopPromise = null;
  30. this.cancelWait = null;
  31. this.lastRecoveryAttemptAt = 0;
  32. this.lastRecoveryError = '';
  33. }
  34. wait(ms) {
  35. return new Promise(resolve => {
  36. const timer = setTimeout(() => {
  37. this.cancelWait = null;
  38. resolve();
  39. }, ms);
  40. this.cancelWait = () => {
  41. clearTimeout(timer);
  42. this.cancelWait = null;
  43. resolve();
  44. };
  45. });
  46. }
  47. startFriendWorker() {
  48. if (!this.config.friendPollingEnabled || this.friendWorker) return;
  49. const scriptPath = path.join(PACKAGE_ROOT, 'scripts', 'friend-polling-worker.js');
  50. this.friendWorker = spawn(process.execPath, [scriptPath], {
  51. cwd: this.workspaceRoot,
  52. env: { ...process.env, QIWEI_WORKSPACE_ROOT: this.workspaceRoot },
  53. stdio: 'inherit',
  54. windowsHide: true,
  55. });
  56. this.onState({ friendPolling: { status: 'running', pid: this.friendWorker.pid } });
  57. this.friendWorker.once('exit', (code, signal) => {
  58. this.friendWorker = null;
  59. this.onState({ friendPolling: { status: this.running ? 'error' : 'stopped', code, signal } });
  60. });
  61. }
  62. recoveryDue() {
  63. const cooldownMs = Math.max(30000, Number(this.config.retryMs) || 10000);
  64. return Date.now() - this.lastRecoveryAttemptAt >= cooldownMs;
  65. }
  66. async recoverAndRestart(listener = {}) {
  67. if (!this.recoveryDue()) return null;
  68. this.lastRecoveryAttemptAt = Date.now();
  69. this.onState({ personalPolling: { status: 'reconnecting' } });
  70. const recovered = await this.listenerApi.recoverLogin();
  71. const recoveredCode = Number(recovered?.summary?.statusCode);
  72. if (!recovered?.summary?.loggedIn && recoveredCode !== 2) {
  73. this.lastRecoveryError = String(
  74. recovered?.errors?.[0]?.message
  75. || recovered?.assistantMessage
  76. || 'Automatic login recovery is pending.',
  77. );
  78. return null;
  79. }
  80. this.lastRecoveryError = '';
  81. if (listener.running) await this.listenerApi.stop({ preserveAgentState: true });
  82. this.listenerStarted = false;
  83. const started = await this.listenerApi.start();
  84. this.listenerStarted = Boolean(started?.data?.running);
  85. return started;
  86. }
  87. async pollOnce() {
  88. let result = await this.listenerApi.getStatus();
  89. let listener = result?.data?.listener || {};
  90. let account = result?.data?.account || {};
  91. let requiresLogin = account.reason === 'upstream_device_missing';
  92. if (account.online === false && !requiresLogin) {
  93. await this.recoverAndRestart(listener);
  94. result = await this.listenerApi.getStatus();
  95. listener = result?.data?.listener || {};
  96. account = result?.data?.account || {};
  97. requiresLogin = account.reason === 'upstream_device_missing';
  98. }
  99. if (!listener.running && account.online !== false) {
  100. try {
  101. const started = await this.listenerApi.start();
  102. this.listenerStarted = Boolean(started?.data?.running);
  103. } catch (error) {
  104. const recovered = await this.recoverAndRestart(listener);
  105. if (!recovered) throw error;
  106. }
  107. result = await this.listenerApi.getStatus();
  108. listener = result?.data?.listener || {};
  109. account = result?.data?.account || {};
  110. }
  111. this.listenerStarted = Boolean(listener.running);
  112. if (this.listenerStarted) this.startFriendWorker();
  113. const listenerHealthy = listener.running && account.online !== false;
  114. if (listenerHealthy) this.lastRecoveryError = '';
  115. this.onState({
  116. personalPolling: {
  117. status: listenerHealthy ? 'running' : (requiresLogin ? 'needs_login' : 'waiting'),
  118. syncKey: Number(listener.syncKey || 0),
  119. startedAt: listener.startedAt || null,
  120. lastError: listenerHealthy
  121. ? ''
  122. : (requiresLogin ? account.statusText : '') || listener.lastError || this.lastRecoveryError,
  123. },
  124. });
  125. return listener;
  126. }
  127. async loop() {
  128. while (this.running) {
  129. try {
  130. await this.pollOnce();
  131. } catch (error) {
  132. this.listenerStarted = false;
  133. this.onState({ personalPolling: { status: 'waiting', lastError: error.message } });
  134. }
  135. if (this.running) await this.wait(this.config.retryMs);
  136. }
  137. }
  138. start() {
  139. if (this.running) return;
  140. this.running = true;
  141. this.onState({ personalPolling: { status: 'starting', lastError: '' } });
  142. this.loopPromise = this.loop();
  143. }
  144. async stop() {
  145. this.running = false;
  146. if (this.cancelWait) this.cancelWait();
  147. if (this.listenerStarted) {
  148. try { await this.listenerApi.stop({ preserveAgentState: true }); } catch {}
  149. }
  150. this.listenerStarted = false;
  151. if (this.friendWorker) {
  152. try { this.friendWorker.kill(); } catch {}
  153. this.friendWorker = null;
  154. }
  155. this.onState({
  156. personalPolling: { status: 'stopped' },
  157. friendPolling: { status: 'stopped' },
  158. });
  159. await this.loopPromise;
  160. }
  161. }