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 }); }