#!/usr/bin/env node /** * Relay 长轮询客户端 * * 独立进程运行,从中央 Relay 拉取属于本租户的加密事件, * 用本地 RSA 私钥解密后构造 v1 envelope 并交给 processWebhookEvents 处理。 * * 启动方式: * node scripts/start-relay-client.js [device-guid] * npm run relay */ const crypto = require('crypto'); const { processWebhookEvents } = require('../mcp/src/core/webhook-server'); const { readQiweiGuid } = require('../mcp/src/core/credentials'); const { getProductMode } = require('../mcp/src/core/product-mode'); const { getRelayBaseUrl, getTenantApiSecret, getRelayPrivateKey, getTenantId, getRelayDeviceGuid } = require('../mcp/src/core/relay-config'); const POLL_WAIT_MS = 30000; const INITIAL_BACKOFF_MS = 1000; const MAX_BACKOFF_MS = 60000; function decryptPayload(encryptedPayload, privateKeyPem) { const key = crypto.createPrivateKey(privateKeyPem); if (String(encryptedPayload).startsWith('v2:')) { const envelope = JSON.parse(Buffer.from(String(encryptedPayload).slice(3), 'base64').toString('utf8')); const aesKey = crypto.privateDecrypt( { key, oaepHash: 'sha256' }, Buffer.from(envelope.key, 'base64') ); const decipher = crypto.createDecipheriv( 'aes-256-gcm', aesKey, Buffer.from(envelope.iv, 'base64') ); decipher.setAuthTag(Buffer.from(envelope.tag, 'base64')); return Buffer.concat([ decipher.update(Buffer.from(envelope.ciphertext, 'base64')), decipher.final() ]).toString('utf8'); } const buffer = Buffer.from(encryptedPayload, 'base64'); const decrypted = crypto.privateDecrypt({ key, oaepHash: 'sha256' }, buffer); return decrypted.toString('utf8'); } async function ackEvents(baseUrl, apiSecret, guid, eventIds) { if (!eventIds.length) return; try { const res = await fetch(`${baseUrl}/api/relay/ack`, { method: 'POST', headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${apiSecret}` }, body: JSON.stringify({ guid, eventIds }) }); if (!res.ok) { console.warn('[RelayClient] ACK 失败:', res.status, await res.text()); return; } const data = await res.json(); console.log(`[RelayClient] ACK ${data.ackedCount || eventIds.length} 条事件`); } catch (err) { console.warn('[RelayClient] ACK 请求异常:', err.message); } } async function runPollOnce(baseUrl, apiSecret, guid, privateKey) { const res = await fetch(`${baseUrl}/api/relay/poll`, { method: 'POST', headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${apiSecret}` }, body: JSON.stringify({ guid, batchSize: 100, waitMs: POLL_WAIT_MS }) }); if (!res.ok) { throw new Error(`poll failed: ${res.status} ${await res.text()}`); } const data = await res.json(); if (!data.events || !data.events.length) return { acked: 0 }; console.log(`[RelayClient] 取回 ${data.events.length} 条事件`); const eventIds = []; for (const event of data.events) { try { const decrypted = decryptPayload(event.encryptedPayload, privateKey); const payload = JSON.parse(decrypted); const envelope = { code: 0, msg: 'from-relay', data: Array.isArray(payload) ? payload : [payload], __rawBody: decrypted }; await processWebhookEvents(envelope); eventIds.push(event.eventId); } catch (err) { console.error(`[RelayClient] 解密/处理失败 eventId=${event.eventId}:`, err.message); // 解密失败也要 ACK,避免 Relay 重复投递 eventIds.push(event.eventId); } } await ackEvents(baseUrl, apiSecret, guid, eventIds); return { acked: eventIds.length }; } function resolveGuid() { // 命令行参数 > 环境变量 > relay-config.json > credentials return process.argv[2] || process.env.RELAY_DEVICE_GUID || getRelayDeviceGuid() || readQiweiGuid() || ''; } async function main() { const product = getProductMode(); if (product.mode !== 'enterprise') { console.error('[RelayClient] 当前为个人版,请使用工作台本地监听;企业 Relay Client 未启动'); process.exit(1); } const baseUrl = getRelayBaseUrl(); const apiSecret = getTenantApiSecret(); const privateKey = getRelayPrivateKey(); const tenantId = getTenantId(); const guid = resolveGuid(); if (!apiSecret || !privateKey || !tenantId) { console.error('[RelayClient] 缺少配置:请检查 .env.local 中的 TENANT_API_SECRET、RELAY_PRIVATE_KEY、TENANT_ID'); process.exit(1); } if (!guid) { console.error('[RelayClient] 缺少 deviceGuid:请通过命令行传入,或配置 RELAY_DEVICE_GUID / relay-config.json / 完成企微登录'); process.exit(1); } console.log(`[RelayClient] 启动 Relay 轮询: ${baseUrl}`); console.log(`[RelayClient] tenantId=${tenantId}, guid=${guid}`); let backoff = INITIAL_BACKOFF_MS; while (true) { try { await runPollOnce(baseUrl, apiSecret, guid, privateKey); backoff = INITIAL_BACKOFF_MS; } catch (err) { console.error('[RelayClient] 轮询异常:', err.message); console.log(`[RelayClient] ${backoff}ms 后重试...`); await new Promise((resolve) => setTimeout(resolve, backoff)); backoff = Math.min(backoff * 2, MAX_BACKOFF_MS); } } } main().catch((err) => { console.error('[RelayClient] 致命错误:', err); process.exit(1); });