| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291 |
- import assert from 'node:assert/strict';
- import crypto from 'node:crypto';
- import fs from 'node:fs';
- import os from 'node:os';
- import path from 'node:path';
- import { pathToFileURL } from 'node:url';
- import { decryptPayload, EnterpriseRelayRuntime } from '../runtime/callback-service/src/enterprise-relay-client.mjs';
- import { PersonalPollingRuntime } from '../runtime/callback-service/src/personal-polling.mjs';
- import { loadRuntimeConfig, resolveRuntimeMode } from '../runtime/callback-service/src/config-loader.mjs';
- import { PortraitQueueWorker } from '../runtime/callback-service/src/portrait-queue-worker.mjs';
- import {
- acquireRuntimeLock,
- clearRuntimeStopRequest,
- readRuntimeState,
- readRuntimeStopRequest,
- releaseRuntimeLock,
- requestRuntimeStop,
- writeRuntimeState,
- } from '../runtime/callback-service/src/runtime-state.mjs';
- function encryptV2(payload, publicKey) {
- const aesKey = crypto.randomBytes(32);
- const iv = crypto.randomBytes(12);
- const cipher = crypto.createCipheriv('aes-256-gcm', aesKey, iv);
- const ciphertext = Buffer.concat([cipher.update(payload, 'utf8'), cipher.final()]);
- const envelope = {
- key: crypto.publicEncrypt({ key: publicKey, oaepHash: 'sha256' }, aesKey).toString('base64'),
- iv: iv.toString('base64'),
- tag: cipher.getAuthTag().toString('base64'),
- ciphertext: ciphertext.toString('base64'),
- };
- return `v2:${Buffer.from(JSON.stringify(envelope)).toString('base64')}`;
- }
- const tempRoot = fs.mkdtempSync(path.join(os.tmpdir(), 'qiwei-runtime-smoke-'));
- try {
- const configPath = path.join(tempRoot, 'runtime-config.mjs');
- fs.writeFileSync(configPath, "export default { edition: 'personal', dashboard: { port: 4399 } };\n", 'utf8');
- const loaded = await loadRuntimeConfig({ workspaceRoot: tempRoot, configPath });
- assert.equal(loaded.mode, 'personal');
- assert.equal(loaded.config.dashboard.port, 4399);
- assert.equal(loaded.config.enterprise.relay.batchSize, 100);
- assert.equal(resolveRuntimeMode({ edition: 'enterprise' }), 'enterprise');
- const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', { modulusLength: 2048 });
- const payload = JSON.stringify({ cmd: 15000, content: 'runtime-smoke' });
- const encrypted = encryptV2(payload, publicKey);
- assert.equal(decryptPayload(encrypted, privateKey), payload);
- const relayConfig = { batchSize: 10, waitMs: 1000, retryMinMs: 250, retryMaxMs: 1000 };
- const relayContext = {
- baseUrl: 'http://relay.test',
- apiSecret: 'test-secret',
- privateKey,
- tenantId: 'tenant-test',
- guid: 'guid-test',
- };
- let ackBody = null;
- let processedEnvelope = null;
- const successFetch = async (url, options) => {
- if (url.endsWith('/api/relay/poll')) {
- return new Response(JSON.stringify({ events: [{ eventId: 'event-1', encryptedPayload: encrypted }] }), {
- status: 200,
- headers: { 'Content-Type': 'application/json' },
- });
- }
- ackBody = JSON.parse(options.body);
- return new Response(JSON.stringify({ ackedCount: 1 }), {
- status: 200,
- headers: { 'Content-Type': 'application/json' },
- });
- };
- const relay = new EnterpriseRelayRuntime({
- config: relayConfig,
- fetchImpl: successFetch,
- processor: async envelope => { processedEnvelope = envelope; },
- });
- const relayResult = await relay.pollOnce(relayContext);
- assert.equal(relayResult.acked, 1);
- assert.deepEqual(ackBody.eventIds, ['event-1']);
- assert.equal(processedEnvelope.data[0].content, 'runtime-smoke');
- let recoveryCalls = 0;
- const recoveryStates = [];
- const recoveryRelay = new EnterpriseRelayRuntime({
- config: relayConfig,
- recovery: async () => {
- recoveryCalls += 1;
- return { status: 'completed', recovered: 2 };
- },
- onState: patch => recoveryStates.push(patch),
- });
- await recoveryRelay.triggerRecovery();
- assert.equal(recoveryCalls, 1);
- assert.equal(recoveryRelay.recoveryMessages, 2);
- assert.equal(recoveryStates.at(-1).enterpriseRelay.generationRecoveryMessages, 2);
- // An acknowledged callback may still be generating when the process is
- // restarted. Recovery must run after an otherwise empty relay long-poll;
- // do not rely on another inbound event arriving to unblock the customer.
- let emptyPollRecoveryCalls = 0;
- let emptyPolls = 0;
- let signalEmptyPollRecovery;
- const emptyPollRecovery = new Promise(resolve => { signalEmptyPollRecovery = resolve; });
- const emptyPollStates = [];
- const emptyPollRelay = new EnterpriseRelayRuntime({
- config: { ...relayConfig, retryMinMs: 25, retryMaxMs: 100 },
- guid: 'guid-empty-poll',
- recovery: async () => {
- emptyPollRecoveryCalls += 1;
- signalEmptyPollRecovery();
- return { status: 'completed', recovered: 1 };
- },
- onState: patch => emptyPollStates.push(patch),
- fetchImpl: async (url, options = {}) => {
- if (url.endsWith('/api/tenant/status')) {
- return new Response(JSON.stringify({ devices: [] }), { status: 200, headers: { 'Content-Type': 'application/json' } });
- }
- if (url.endsWith('/api/relay/poll')) {
- emptyPolls += 1;
- if (emptyPolls === 1) {
- return new Response(JSON.stringify({ events: [] }), { status: 200, headers: { 'Content-Type': 'application/json' } });
- }
- return new Promise((resolve, reject) => {
- if (options.signal?.aborted) {
- reject(new DOMException('Aborted', 'AbortError'));
- return;
- }
- options.signal?.addEventListener('abort', () => reject(new DOMException('Aborted', 'AbortError')), { once: true });
- });
- }
- throw new Error(`unexpected empty-poll request: ${url}`);
- },
- });
- // Bypass the separate upstream reconnect fixture. This test owns the relay
- // poll loop and verifies the post-poll recovery hook itself.
- emptyPollRelay.contextIdentity = 'http://relay.fixture|tenant-fixture|guid-empty-poll';
- emptyPollRelay.nextUpstreamReconnectAt = Date.now() + 60_000;
- const originalRelayEnv = {
- RELAY_BASE_URL: process.env.RELAY_BASE_URL,
- TENANT_API_SECRET: process.env.TENANT_API_SECRET,
- TENANT_ID: process.env.TENANT_ID,
- RELAY_PRIVATE_KEY: process.env.RELAY_PRIVATE_KEY,
- };
- process.env.RELAY_BASE_URL = 'http://relay.fixture';
- process.env.TENANT_API_SECRET = 'fixture-secret';
- process.env.TENANT_ID = 'tenant-fixture';
- process.env.RELAY_PRIVATE_KEY = 'fixture-private-key';
- try {
- emptyPollRelay.start();
- await Promise.race([
- emptyPollRecovery,
- new Promise((_, reject) => setTimeout(() => reject(new Error('empty relay poll did not trigger generation recovery')), 1000)),
- ]);
- await emptyPollRelay.stop();
- assert.equal(emptyPolls >= 1, true);
- assert.equal(emptyPollRecoveryCalls, 1);
- assert.equal(emptyPollStates.some(patch => patch.enterpriseRelay?.lastFetched === 0), true);
- } finally {
- for (const [key, value] of Object.entries(originalRelayEnv)) {
- if (value === undefined) delete process.env[key];
- else process.env[key] = value;
- }
- }
- let failureAckCalled = false;
- const failureRelay = new EnterpriseRelayRuntime({
- config: relayConfig,
- fetchImpl: async url => {
- if (url.endsWith('/api/relay/ack')) failureAckCalled = true;
- return new Response(JSON.stringify({ events: [{ eventId: 'event-2', encryptedPayload: encrypted }] }), {
- status: 200,
- headers: { 'Content-Type': 'application/json' },
- });
- },
- processor: async () => { throw new Error('processor-test-failure'); },
- });
- await assert.rejects(() => failureRelay.pollOnce(relayContext), /retained 1 failed event/);
- assert.equal(failureAckCalled, false);
- let releasePortraitRun;
- let portraitRuns = 0;
- const portraitStates = [];
- const portraitWorker = new PortraitQueueWorker({
- intervalMs: 60_000,
- processor: async limit => {
- portraitRuns += 1;
- assert.equal(limit, 5);
- await new Promise(resolve => { releasePortraitRun = resolve; });
- return { processed: 2, remaining: 3 };
- },
- onState: patch => portraitStates.push(patch.portraitQueue),
- });
- portraitWorker.start();
- assert.equal(portraitRuns, 1);
- assert.deepEqual(await portraitWorker.runOnce(), { skipped: true, reason: 'in-flight' });
- releasePortraitRun();
- await portraitWorker.stop();
- assert.equal(portraitStates.some(item => item.processed === 2 && item.remaining === 3), true);
- assert.equal(portraitStates.at(-1).status, 'stopped');
- const statePath = path.join(tempRoot, 'runtime-state.json');
- writeRuntimeState({ pid: 123, status: 'running', apiSecret: 'hidden', components: { relay: { token: 'hidden' } } }, statePath);
- const state = readRuntimeState(statePath);
- assert.equal(state.status, 'running');
- assert.equal('apiSecret' in state, false);
- assert.equal('token' in state.components.relay, false);
- const runtimeLock = acquireRuntimeLock(statePath);
- assert.throws(() => acquireRuntimeLock(statePath), /already running/);
- assert.equal(releaseRuntimeLock(runtimeLock), true);
- const staleLockPath = `${statePath}.lock`;
- fs.writeFileSync(staleLockPath, JSON.stringify({ pid: 999999, acquiredAt: 'fixture' }), 'utf8');
- const reclaimedLock = acquireRuntimeLock(statePath);
- assert.equal(releaseRuntimeLock(reclaimedLock), true);
- assert.equal(fs.existsSync(staleLockPath), false);
- requestRuntimeStop(statePath);
- assert.equal(readRuntimeStopRequest(statePath).targetPid, 123);
- clearRuntimeStopRequest(statePath);
- assert.deepEqual(readRuntimeStopRequest(statePath), {});
- const runtimeIndexSource = fs.readFileSync(new URL('../runtime/callback-service/src/index.mjs', import.meta.url), 'utf8');
- assert.match(runtimeIndexSource, /alreadyRunning: true/);
- let listener = { running: true, syncKey: 7, lastError: 'offline', startedAt: 1 };
- let stopOptions = null;
- let recoveryCount = 0;
- const personalStates = [];
- const personal = new PersonalPollingRuntime({
- config: { retryMs: 3000, friendPollingEnabled: false },
- workspaceRoot: tempRoot,
- onState: statePatch => personalStates.push(statePatch),
- listenerApi: {
- getStatus: async () => ({
- data: {
- account: { online: recoveryCount > 0 },
- listener,
- },
- }),
- recoverLogin: async () => {
- recoveryCount += 1;
- return { summary: { loggedIn: true, statusCode: 2 } };
- },
- stop: async options => { stopOptions = options; listener = { ...listener, running: false }; },
- start: async () => {
- listener = { running: true, syncKey: 7, lastError: '', startedAt: 2 };
- return { data: listener };
- },
- },
- });
- const recoveredListener = await personal.pollOnce();
- assert.equal(recoveryCount, 1);
- assert.deepEqual(stopOptions, { preserveAgentState: true });
- assert.equal(recoveredListener.running, true);
- assert.equal(recoveredListener.lastError, '');
- assert.equal(personalStates.at(-1).personalPolling.status, 'running');
- assert.equal(personalStates.at(-1).personalPolling.lastError, '');
- let missingDeviceRecoveryCount = 0;
- const missingDeviceStates = [];
- const missingDeviceRuntime = new PersonalPollingRuntime({
- config: { retryMs: 3000, friendPollingEnabled: false },
- workspaceRoot: tempRoot,
- onState: statePatch => missingDeviceStates.push(statePatch),
- listenerApi: {
- getStatus: async () => ({
- data: {
- account: {
- online: false,
- reason: 'upstream_device_missing',
- statusText: '设备实例已失效,请重新扫码登录',
- },
- listener: { running: false, syncKey: 8, lastError: '旧错误' },
- },
- }),
- recoverLogin: async () => {
- missingDeviceRecoveryCount += 1;
- return { summary: { loggedIn: false, statusCode: -1 } };
- },
- start: async () => ({ data: { running: false } }),
- stop: async () => {},
- },
- });
- await missingDeviceRuntime.pollOnce();
- assert.equal(missingDeviceRecoveryCount, 0);
- assert.equal(missingDeviceStates.at(-1).personalPolling.status, 'needs_login');
- assert.equal(missingDeviceStates.at(-1).personalPolling.lastError, '设备实例已失效,请重新扫码登录');
- process.stdout.write(`${JSON.stringify({ status: 'ok', checks: 34 }, null, 2)}\n`);
- } finally {
- fs.rmSync(tempRoot, { recursive: true, force: true });
- }
|