| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198 |
- 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 {
- clearRuntimeStopRequest,
- readRuntimeState,
- readRuntimeStopRequest,
- 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 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);
- 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: 26 }, null, 2)}\n`);
- } finally {
- fs.rmSync(tempRoot, { recursive: true, force: true });
- }
|