callback-runtime-smoke-test.mjs 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198
  1. import assert from 'node:assert/strict';
  2. import crypto from 'node:crypto';
  3. import fs from 'node:fs';
  4. import os from 'node:os';
  5. import path from 'node:path';
  6. import { pathToFileURL } from 'node:url';
  7. import { decryptPayload, EnterpriseRelayRuntime } from '../runtime/callback-service/src/enterprise-relay-client.mjs';
  8. import { PersonalPollingRuntime } from '../runtime/callback-service/src/personal-polling.mjs';
  9. import { loadRuntimeConfig, resolveRuntimeMode } from '../runtime/callback-service/src/config-loader.mjs';
  10. import { PortraitQueueWorker } from '../runtime/callback-service/src/portrait-queue-worker.mjs';
  11. import {
  12. clearRuntimeStopRequest,
  13. readRuntimeState,
  14. readRuntimeStopRequest,
  15. requestRuntimeStop,
  16. writeRuntimeState,
  17. } from '../runtime/callback-service/src/runtime-state.mjs';
  18. function encryptV2(payload, publicKey) {
  19. const aesKey = crypto.randomBytes(32);
  20. const iv = crypto.randomBytes(12);
  21. const cipher = crypto.createCipheriv('aes-256-gcm', aesKey, iv);
  22. const ciphertext = Buffer.concat([cipher.update(payload, 'utf8'), cipher.final()]);
  23. const envelope = {
  24. key: crypto.publicEncrypt({ key: publicKey, oaepHash: 'sha256' }, aesKey).toString('base64'),
  25. iv: iv.toString('base64'),
  26. tag: cipher.getAuthTag().toString('base64'),
  27. ciphertext: ciphertext.toString('base64'),
  28. };
  29. return `v2:${Buffer.from(JSON.stringify(envelope)).toString('base64')}`;
  30. }
  31. const tempRoot = fs.mkdtempSync(path.join(os.tmpdir(), 'qiwei-runtime-smoke-'));
  32. try {
  33. const configPath = path.join(tempRoot, 'runtime-config.mjs');
  34. fs.writeFileSync(configPath, "export default { edition: 'personal', dashboard: { port: 4399 } };\n", 'utf8');
  35. const loaded = await loadRuntimeConfig({ workspaceRoot: tempRoot, configPath });
  36. assert.equal(loaded.mode, 'personal');
  37. assert.equal(loaded.config.dashboard.port, 4399);
  38. assert.equal(loaded.config.enterprise.relay.batchSize, 100);
  39. assert.equal(resolveRuntimeMode({ edition: 'enterprise' }), 'enterprise');
  40. const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', { modulusLength: 2048 });
  41. const payload = JSON.stringify({ cmd: 15000, content: 'runtime-smoke' });
  42. const encrypted = encryptV2(payload, publicKey);
  43. assert.equal(decryptPayload(encrypted, privateKey), payload);
  44. const relayConfig = { batchSize: 10, waitMs: 1000, retryMinMs: 250, retryMaxMs: 1000 };
  45. const relayContext = {
  46. baseUrl: 'http://relay.test',
  47. apiSecret: 'test-secret',
  48. privateKey,
  49. tenantId: 'tenant-test',
  50. guid: 'guid-test',
  51. };
  52. let ackBody = null;
  53. let processedEnvelope = null;
  54. const successFetch = async (url, options) => {
  55. if (url.endsWith('/api/relay/poll')) {
  56. return new Response(JSON.stringify({ events: [{ eventId: 'event-1', encryptedPayload: encrypted }] }), {
  57. status: 200,
  58. headers: { 'Content-Type': 'application/json' },
  59. });
  60. }
  61. ackBody = JSON.parse(options.body);
  62. return new Response(JSON.stringify({ ackedCount: 1 }), {
  63. status: 200,
  64. headers: { 'Content-Type': 'application/json' },
  65. });
  66. };
  67. const relay = new EnterpriseRelayRuntime({
  68. config: relayConfig,
  69. fetchImpl: successFetch,
  70. processor: async envelope => { processedEnvelope = envelope; },
  71. });
  72. const relayResult = await relay.pollOnce(relayContext);
  73. assert.equal(relayResult.acked, 1);
  74. assert.deepEqual(ackBody.eventIds, ['event-1']);
  75. assert.equal(processedEnvelope.data[0].content, 'runtime-smoke');
  76. let failureAckCalled = false;
  77. const failureRelay = new EnterpriseRelayRuntime({
  78. config: relayConfig,
  79. fetchImpl: async url => {
  80. if (url.endsWith('/api/relay/ack')) failureAckCalled = true;
  81. return new Response(JSON.stringify({ events: [{ eventId: 'event-2', encryptedPayload: encrypted }] }), {
  82. status: 200,
  83. headers: { 'Content-Type': 'application/json' },
  84. });
  85. },
  86. processor: async () => { throw new Error('processor-test-failure'); },
  87. });
  88. await assert.rejects(() => failureRelay.pollOnce(relayContext), /retained 1 failed event/);
  89. assert.equal(failureAckCalled, false);
  90. let releasePortraitRun;
  91. let portraitRuns = 0;
  92. const portraitStates = [];
  93. const portraitWorker = new PortraitQueueWorker({
  94. intervalMs: 60_000,
  95. processor: async limit => {
  96. portraitRuns += 1;
  97. assert.equal(limit, 5);
  98. await new Promise(resolve => { releasePortraitRun = resolve; });
  99. return { processed: 2, remaining: 3 };
  100. },
  101. onState: patch => portraitStates.push(patch.portraitQueue),
  102. });
  103. portraitWorker.start();
  104. assert.equal(portraitRuns, 1);
  105. assert.deepEqual(await portraitWorker.runOnce(), { skipped: true, reason: 'in-flight' });
  106. releasePortraitRun();
  107. await portraitWorker.stop();
  108. assert.equal(portraitStates.some(item => item.processed === 2 && item.remaining === 3), true);
  109. assert.equal(portraitStates.at(-1).status, 'stopped');
  110. const statePath = path.join(tempRoot, 'runtime-state.json');
  111. writeRuntimeState({ pid: 123, status: 'running', apiSecret: 'hidden', components: { relay: { token: 'hidden' } } }, statePath);
  112. const state = readRuntimeState(statePath);
  113. assert.equal(state.status, 'running');
  114. assert.equal('apiSecret' in state, false);
  115. assert.equal('token' in state.components.relay, false);
  116. requestRuntimeStop(statePath);
  117. assert.equal(readRuntimeStopRequest(statePath).targetPid, 123);
  118. clearRuntimeStopRequest(statePath);
  119. assert.deepEqual(readRuntimeStopRequest(statePath), {});
  120. const runtimeIndexSource = fs.readFileSync(new URL('../runtime/callback-service/src/index.mjs', import.meta.url), 'utf8');
  121. assert.match(runtimeIndexSource, /alreadyRunning: true/);
  122. let listener = { running: true, syncKey: 7, lastError: 'offline', startedAt: 1 };
  123. let stopOptions = null;
  124. let recoveryCount = 0;
  125. const personalStates = [];
  126. const personal = new PersonalPollingRuntime({
  127. config: { retryMs: 3000, friendPollingEnabled: false },
  128. workspaceRoot: tempRoot,
  129. onState: statePatch => personalStates.push(statePatch),
  130. listenerApi: {
  131. getStatus: async () => ({
  132. data: {
  133. account: { online: recoveryCount > 0 },
  134. listener,
  135. },
  136. }),
  137. recoverLogin: async () => {
  138. recoveryCount += 1;
  139. return { summary: { loggedIn: true, statusCode: 2 } };
  140. },
  141. stop: async options => { stopOptions = options; listener = { ...listener, running: false }; },
  142. start: async () => {
  143. listener = { running: true, syncKey: 7, lastError: '', startedAt: 2 };
  144. return { data: listener };
  145. },
  146. },
  147. });
  148. const recoveredListener = await personal.pollOnce();
  149. assert.equal(recoveryCount, 1);
  150. assert.deepEqual(stopOptions, { preserveAgentState: true });
  151. assert.equal(recoveredListener.running, true);
  152. assert.equal(recoveredListener.lastError, '');
  153. assert.equal(personalStates.at(-1).personalPolling.status, 'running');
  154. assert.equal(personalStates.at(-1).personalPolling.lastError, '');
  155. let missingDeviceRecoveryCount = 0;
  156. const missingDeviceStates = [];
  157. const missingDeviceRuntime = new PersonalPollingRuntime({
  158. config: { retryMs: 3000, friendPollingEnabled: false },
  159. workspaceRoot: tempRoot,
  160. onState: statePatch => missingDeviceStates.push(statePatch),
  161. listenerApi: {
  162. getStatus: async () => ({
  163. data: {
  164. account: {
  165. online: false,
  166. reason: 'upstream_device_missing',
  167. statusText: '设备实例已失效,请重新扫码登录',
  168. },
  169. listener: { running: false, syncKey: 8, lastError: '旧错误' },
  170. },
  171. }),
  172. recoverLogin: async () => {
  173. missingDeviceRecoveryCount += 1;
  174. return { summary: { loggedIn: false, statusCode: -1 } };
  175. },
  176. start: async () => ({ data: { running: false } }),
  177. stop: async () => {},
  178. },
  179. });
  180. await missingDeviceRuntime.pollOnce();
  181. assert.equal(missingDeviceRecoveryCount, 0);
  182. assert.equal(missingDeviceStates.at(-1).personalPolling.status, 'needs_login');
  183. assert.equal(missingDeviceStates.at(-1).personalPolling.lastError, '设备实例已失效,请重新扫码登录');
  184. process.stdout.write(`${JSON.stringify({ status: 'ok', checks: 26 }, null, 2)}\n`);
  185. } finally {
  186. fs.rmSync(tempRoot, { recursive: true, force: true });
  187. }