#!/usr/bin/env node /** * Relay 长轮询客户端 * * 独立进程运行,从中央 Relay 拉取属于本租户的加密事件, * 用本地 RSA 私钥解密后落盘到 outputs/webhook/。 * * 启动方式: * node scripts/start-relay-client.js [device-guid] * npm run relay */ const crypto = require('crypto'); const { saveWebhookEvent } = require('../mcp/src/core/webhook-server'); 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 saveEvent(eventId, payload) { const filePath = saveWebhookEvent({ eventId, ...payload }, 'relay'); return filePath; } function decryptPayload(encryptedPayload, privateKeyPem) { const key = crypto.createPrivateKey(privateKeyPem); 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; 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 filePath = saveEvent(event.eventId, payload); console.log(`[RelayClient] 已解密落盘: ${filePath}`); 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); } async function main() { const baseUrl = getRelayBaseUrl(); const apiSecret = getTenantApiSecret(); const privateKey = getRelayPrivateKey(); const tenantId = getTenantId(); // guid 优先级:命令行参数 > 环境变量 > relay-config.json const guid = process.argv[2] || process.env.RELAY_DEVICE_GUID || getRelayDeviceGuid(); 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); });