callback-runtime-smoke-test.mjs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291
  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. acquireRuntimeLock,
  13. clearRuntimeStopRequest,
  14. readRuntimeState,
  15. readRuntimeStopRequest,
  16. releaseRuntimeLock,
  17. requestRuntimeStop,
  18. writeRuntimeState,
  19. } from '../runtime/callback-service/src/runtime-state.mjs';
  20. function encryptV2(payload, publicKey) {
  21. const aesKey = crypto.randomBytes(32);
  22. const iv = crypto.randomBytes(12);
  23. const cipher = crypto.createCipheriv('aes-256-gcm', aesKey, iv);
  24. const ciphertext = Buffer.concat([cipher.update(payload, 'utf8'), cipher.final()]);
  25. const envelope = {
  26. key: crypto.publicEncrypt({ key: publicKey, oaepHash: 'sha256' }, aesKey).toString('base64'),
  27. iv: iv.toString('base64'),
  28. tag: cipher.getAuthTag().toString('base64'),
  29. ciphertext: ciphertext.toString('base64'),
  30. };
  31. return `v2:${Buffer.from(JSON.stringify(envelope)).toString('base64')}`;
  32. }
  33. const tempRoot = fs.mkdtempSync(path.join(os.tmpdir(), 'qiwei-runtime-smoke-'));
  34. try {
  35. const configPath = path.join(tempRoot, 'runtime-config.mjs');
  36. fs.writeFileSync(configPath, "export default { edition: 'personal', dashboard: { port: 4399 } };\n", 'utf8');
  37. const loaded = await loadRuntimeConfig({ workspaceRoot: tempRoot, configPath });
  38. assert.equal(loaded.mode, 'personal');
  39. assert.equal(loaded.config.dashboard.port, 4399);
  40. assert.equal(loaded.config.enterprise.relay.batchSize, 100);
  41. assert.equal(resolveRuntimeMode({ edition: 'enterprise' }), 'enterprise');
  42. const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', { modulusLength: 2048 });
  43. const payload = JSON.stringify({ cmd: 15000, content: 'runtime-smoke' });
  44. const encrypted = encryptV2(payload, publicKey);
  45. assert.equal(decryptPayload(encrypted, privateKey), payload);
  46. const relayConfig = { batchSize: 10, waitMs: 1000, retryMinMs: 250, retryMaxMs: 1000 };
  47. const relayContext = {
  48. baseUrl: 'http://relay.test',
  49. apiSecret: 'test-secret',
  50. privateKey,
  51. tenantId: 'tenant-test',
  52. guid: 'guid-test',
  53. };
  54. let ackBody = null;
  55. let processedEnvelope = null;
  56. const successFetch = async (url, options) => {
  57. if (url.endsWith('/api/relay/poll')) {
  58. return new Response(JSON.stringify({ events: [{ eventId: 'event-1', encryptedPayload: encrypted }] }), {
  59. status: 200,
  60. headers: { 'Content-Type': 'application/json' },
  61. });
  62. }
  63. ackBody = JSON.parse(options.body);
  64. return new Response(JSON.stringify({ ackedCount: 1 }), {
  65. status: 200,
  66. headers: { 'Content-Type': 'application/json' },
  67. });
  68. };
  69. const relay = new EnterpriseRelayRuntime({
  70. config: relayConfig,
  71. fetchImpl: successFetch,
  72. processor: async envelope => { processedEnvelope = envelope; },
  73. });
  74. const relayResult = await relay.pollOnce(relayContext);
  75. assert.equal(relayResult.acked, 1);
  76. assert.deepEqual(ackBody.eventIds, ['event-1']);
  77. assert.equal(processedEnvelope.data[0].content, 'runtime-smoke');
  78. let recoveryCalls = 0;
  79. const recoveryStates = [];
  80. const recoveryRelay = new EnterpriseRelayRuntime({
  81. config: relayConfig,
  82. recovery: async () => {
  83. recoveryCalls += 1;
  84. return { status: 'completed', recovered: 2 };
  85. },
  86. onState: patch => recoveryStates.push(patch),
  87. });
  88. await recoveryRelay.triggerRecovery();
  89. assert.equal(recoveryCalls, 1);
  90. assert.equal(recoveryRelay.recoveryMessages, 2);
  91. assert.equal(recoveryStates.at(-1).enterpriseRelay.generationRecoveryMessages, 2);
  92. // An acknowledged callback may still be generating when the process is
  93. // restarted. Recovery must run after an otherwise empty relay long-poll;
  94. // do not rely on another inbound event arriving to unblock the customer.
  95. let emptyPollRecoveryCalls = 0;
  96. let emptyPolls = 0;
  97. let signalEmptyPollRecovery;
  98. const emptyPollRecovery = new Promise(resolve => { signalEmptyPollRecovery = resolve; });
  99. const emptyPollStates = [];
  100. const emptyPollRelay = new EnterpriseRelayRuntime({
  101. config: { ...relayConfig, retryMinMs: 25, retryMaxMs: 100 },
  102. guid: 'guid-empty-poll',
  103. recovery: async () => {
  104. emptyPollRecoveryCalls += 1;
  105. signalEmptyPollRecovery();
  106. return { status: 'completed', recovered: 1 };
  107. },
  108. onState: patch => emptyPollStates.push(patch),
  109. fetchImpl: async (url, options = {}) => {
  110. if (url.endsWith('/api/tenant/status')) {
  111. return new Response(JSON.stringify({ devices: [] }), { status: 200, headers: { 'Content-Type': 'application/json' } });
  112. }
  113. if (url.endsWith('/api/relay/poll')) {
  114. emptyPolls += 1;
  115. if (emptyPolls === 1) {
  116. return new Response(JSON.stringify({ events: [] }), { status: 200, headers: { 'Content-Type': 'application/json' } });
  117. }
  118. return new Promise((resolve, reject) => {
  119. if (options.signal?.aborted) {
  120. reject(new DOMException('Aborted', 'AbortError'));
  121. return;
  122. }
  123. options.signal?.addEventListener('abort', () => reject(new DOMException('Aborted', 'AbortError')), { once: true });
  124. });
  125. }
  126. throw new Error(`unexpected empty-poll request: ${url}`);
  127. },
  128. });
  129. // Bypass the separate upstream reconnect fixture. This test owns the relay
  130. // poll loop and verifies the post-poll recovery hook itself.
  131. emptyPollRelay.contextIdentity = 'http://relay.fixture|tenant-fixture|guid-empty-poll';
  132. emptyPollRelay.nextUpstreamReconnectAt = Date.now() + 60_000;
  133. const originalRelayEnv = {
  134. RELAY_BASE_URL: process.env.RELAY_BASE_URL,
  135. TENANT_API_SECRET: process.env.TENANT_API_SECRET,
  136. TENANT_ID: process.env.TENANT_ID,
  137. RELAY_PRIVATE_KEY: process.env.RELAY_PRIVATE_KEY,
  138. };
  139. process.env.RELAY_BASE_URL = 'http://relay.fixture';
  140. process.env.TENANT_API_SECRET = 'fixture-secret';
  141. process.env.TENANT_ID = 'tenant-fixture';
  142. process.env.RELAY_PRIVATE_KEY = 'fixture-private-key';
  143. try {
  144. emptyPollRelay.start();
  145. await Promise.race([
  146. emptyPollRecovery,
  147. new Promise((_, reject) => setTimeout(() => reject(new Error('empty relay poll did not trigger generation recovery')), 1000)),
  148. ]);
  149. await emptyPollRelay.stop();
  150. assert.equal(emptyPolls >= 1, true);
  151. assert.equal(emptyPollRecoveryCalls, 1);
  152. assert.equal(emptyPollStates.some(patch => patch.enterpriseRelay?.lastFetched === 0), true);
  153. } finally {
  154. for (const [key, value] of Object.entries(originalRelayEnv)) {
  155. if (value === undefined) delete process.env[key];
  156. else process.env[key] = value;
  157. }
  158. }
  159. let failureAckCalled = false;
  160. const failureRelay = new EnterpriseRelayRuntime({
  161. config: relayConfig,
  162. fetchImpl: async url => {
  163. if (url.endsWith('/api/relay/ack')) failureAckCalled = true;
  164. return new Response(JSON.stringify({ events: [{ eventId: 'event-2', encryptedPayload: encrypted }] }), {
  165. status: 200,
  166. headers: { 'Content-Type': 'application/json' },
  167. });
  168. },
  169. processor: async () => { throw new Error('processor-test-failure'); },
  170. });
  171. await assert.rejects(() => failureRelay.pollOnce(relayContext), /retained 1 failed event/);
  172. assert.equal(failureAckCalled, false);
  173. let releasePortraitRun;
  174. let portraitRuns = 0;
  175. const portraitStates = [];
  176. const portraitWorker = new PortraitQueueWorker({
  177. intervalMs: 60_000,
  178. processor: async limit => {
  179. portraitRuns += 1;
  180. assert.equal(limit, 5);
  181. await new Promise(resolve => { releasePortraitRun = resolve; });
  182. return { processed: 2, remaining: 3 };
  183. },
  184. onState: patch => portraitStates.push(patch.portraitQueue),
  185. });
  186. portraitWorker.start();
  187. assert.equal(portraitRuns, 1);
  188. assert.deepEqual(await portraitWorker.runOnce(), { skipped: true, reason: 'in-flight' });
  189. releasePortraitRun();
  190. await portraitWorker.stop();
  191. assert.equal(portraitStates.some(item => item.processed === 2 && item.remaining === 3), true);
  192. assert.equal(portraitStates.at(-1).status, 'stopped');
  193. const statePath = path.join(tempRoot, 'runtime-state.json');
  194. writeRuntimeState({ pid: 123, status: 'running', apiSecret: 'hidden', components: { relay: { token: 'hidden' } } }, statePath);
  195. const state = readRuntimeState(statePath);
  196. assert.equal(state.status, 'running');
  197. assert.equal('apiSecret' in state, false);
  198. assert.equal('token' in state.components.relay, false);
  199. const runtimeLock = acquireRuntimeLock(statePath);
  200. assert.throws(() => acquireRuntimeLock(statePath), /already running/);
  201. assert.equal(releaseRuntimeLock(runtimeLock), true);
  202. const staleLockPath = `${statePath}.lock`;
  203. fs.writeFileSync(staleLockPath, JSON.stringify({ pid: 999999, acquiredAt: 'fixture' }), 'utf8');
  204. const reclaimedLock = acquireRuntimeLock(statePath);
  205. assert.equal(releaseRuntimeLock(reclaimedLock), true);
  206. assert.equal(fs.existsSync(staleLockPath), false);
  207. requestRuntimeStop(statePath);
  208. assert.equal(readRuntimeStopRequest(statePath).targetPid, 123);
  209. clearRuntimeStopRequest(statePath);
  210. assert.deepEqual(readRuntimeStopRequest(statePath), {});
  211. const runtimeIndexSource = fs.readFileSync(new URL('../runtime/callback-service/src/index.mjs', import.meta.url), 'utf8');
  212. assert.match(runtimeIndexSource, /alreadyRunning: true/);
  213. let listener = { running: true, syncKey: 7, lastError: 'offline', startedAt: 1 };
  214. let stopOptions = null;
  215. let recoveryCount = 0;
  216. const personalStates = [];
  217. const personal = new PersonalPollingRuntime({
  218. config: { retryMs: 3000, friendPollingEnabled: false },
  219. workspaceRoot: tempRoot,
  220. onState: statePatch => personalStates.push(statePatch),
  221. listenerApi: {
  222. getStatus: async () => ({
  223. data: {
  224. account: { online: recoveryCount > 0 },
  225. listener,
  226. },
  227. }),
  228. recoverLogin: async () => {
  229. recoveryCount += 1;
  230. return { summary: { loggedIn: true, statusCode: 2 } };
  231. },
  232. stop: async options => { stopOptions = options; listener = { ...listener, running: false }; },
  233. start: async () => {
  234. listener = { running: true, syncKey: 7, lastError: '', startedAt: 2 };
  235. return { data: listener };
  236. },
  237. },
  238. });
  239. const recoveredListener = await personal.pollOnce();
  240. assert.equal(recoveryCount, 1);
  241. assert.deepEqual(stopOptions, { preserveAgentState: true });
  242. assert.equal(recoveredListener.running, true);
  243. assert.equal(recoveredListener.lastError, '');
  244. assert.equal(personalStates.at(-1).personalPolling.status, 'running');
  245. assert.equal(personalStates.at(-1).personalPolling.lastError, '');
  246. let missingDeviceRecoveryCount = 0;
  247. const missingDeviceStates = [];
  248. const missingDeviceRuntime = new PersonalPollingRuntime({
  249. config: { retryMs: 3000, friendPollingEnabled: false },
  250. workspaceRoot: tempRoot,
  251. onState: statePatch => missingDeviceStates.push(statePatch),
  252. listenerApi: {
  253. getStatus: async () => ({
  254. data: {
  255. account: {
  256. online: false,
  257. reason: 'upstream_device_missing',
  258. statusText: '设备实例已失效,请重新扫码登录',
  259. },
  260. listener: { running: false, syncKey: 8, lastError: '旧错误' },
  261. },
  262. }),
  263. recoverLogin: async () => {
  264. missingDeviceRecoveryCount += 1;
  265. return { summary: { loggedIn: false, statusCode: -1 } };
  266. },
  267. start: async () => ({ data: { running: false } }),
  268. stop: async () => {},
  269. },
  270. });
  271. await missingDeviceRuntime.pollOnce();
  272. assert.equal(missingDeviceRecoveryCount, 0);
  273. assert.equal(missingDeviceStates.at(-1).personalPolling.status, 'needs_login');
  274. assert.equal(missingDeviceStates.at(-1).personalPolling.lastError, '设备实例已失效,请重新扫码登录');
  275. process.stdout.write(`${JSON.stringify({ status: 'ok', checks: 34 }, null, 2)}\n`);
  276. } finally {
  277. fs.rmSync(tempRoot, { recursive: true, force: true });
  278. }