| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139 |
- #!/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);
- });
|