start-relay-client.js 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139
  1. #!/usr/bin/env node
  2. /**
  3. * Relay 长轮询客户端
  4. *
  5. * 独立进程运行,从中央 Relay 拉取属于本租户的加密事件,
  6. * 用本地 RSA 私钥解密后落盘到 outputs/webhook/。
  7. *
  8. * 启动方式:
  9. * node scripts/start-relay-client.js [device-guid]
  10. * npm run relay
  11. */
  12. const crypto = require('crypto');
  13. const { saveWebhookEvent } = require('../mcp/src/core/webhook-server');
  14. const {
  15. getRelayBaseUrl,
  16. getTenantApiSecret,
  17. getRelayPrivateKey,
  18. getTenantId,
  19. getRelayDeviceGuid
  20. } = require('../mcp/src/core/relay-config');
  21. const POLL_WAIT_MS = 30000;
  22. const INITIAL_BACKOFF_MS = 1000;
  23. const MAX_BACKOFF_MS = 60000;
  24. function saveEvent(eventId, payload) {
  25. const filePath = saveWebhookEvent({ eventId, ...payload }, 'relay');
  26. return filePath;
  27. }
  28. function decryptPayload(encryptedPayload, privateKeyPem) {
  29. const key = crypto.createPrivateKey(privateKeyPem);
  30. const buffer = Buffer.from(encryptedPayload, 'base64');
  31. const decrypted = crypto.privateDecrypt({ key, oaepHash: 'sha256' }, buffer);
  32. return decrypted.toString('utf8');
  33. }
  34. async function ackEvents(baseUrl, apiSecret, guid, eventIds) {
  35. if (!eventIds.length) return;
  36. try {
  37. const res = await fetch(`${baseUrl}/api/relay/ack`, {
  38. method: 'POST',
  39. headers: {
  40. 'Content-Type': 'application/json',
  41. Authorization: `Bearer ${apiSecret}`
  42. },
  43. body: JSON.stringify({ guid, eventIds })
  44. });
  45. if (!res.ok) {
  46. console.warn('[RelayClient] ACK 失败:', res.status, await res.text());
  47. return;
  48. }
  49. const data = await res.json();
  50. console.log(`[RelayClient] ACK ${data.ackedCount || eventIds.length} 条事件`);
  51. } catch (err) {
  52. console.warn('[RelayClient] ACK 请求异常:', err.message);
  53. }
  54. }
  55. async function runPollOnce(baseUrl, apiSecret, guid, privateKey) {
  56. const res = await fetch(`${baseUrl}/api/relay/poll`, {
  57. method: 'POST',
  58. headers: {
  59. 'Content-Type': 'application/json',
  60. Authorization: `Bearer ${apiSecret}`
  61. },
  62. body: JSON.stringify({ guid, batchSize: 100, waitMs: POLL_WAIT_MS })
  63. });
  64. if (!res.ok) {
  65. throw new Error(`poll failed: ${res.status} ${await res.text()}`);
  66. }
  67. const data = await res.json();
  68. if (!data.events || !data.events.length) return;
  69. console.log(`[RelayClient] 取回 ${data.events.length} 条事件`);
  70. const eventIds = [];
  71. for (const event of data.events) {
  72. try {
  73. const decrypted = decryptPayload(event.encryptedPayload, privateKey);
  74. const payload = JSON.parse(decrypted);
  75. const filePath = saveEvent(event.eventId, payload);
  76. console.log(`[RelayClient] 已解密落盘: ${filePath}`);
  77. eventIds.push(event.eventId);
  78. } catch (err) {
  79. console.error(`[RelayClient] 解密/落盘失败 eventId=${event.eventId}:`, err.message);
  80. // 解密失败也要 ACK,避免 Relay 重复投递
  81. eventIds.push(event.eventId);
  82. }
  83. }
  84. await ackEvents(baseUrl, apiSecret, guid, eventIds);
  85. }
  86. async function main() {
  87. const baseUrl = getRelayBaseUrl();
  88. const apiSecret = getTenantApiSecret();
  89. const privateKey = getRelayPrivateKey();
  90. const tenantId = getTenantId();
  91. // guid 优先级:命令行参数 > 环境变量 > relay-config.json
  92. const guid = process.argv[2] || process.env.RELAY_DEVICE_GUID || getRelayDeviceGuid();
  93. if (!apiSecret || !privateKey || !tenantId) {
  94. console.error('[RelayClient] 缺少配置:请检查 .env.local 中的 TENANT_API_SECRET、RELAY_PRIVATE_KEY、TENANT_ID');
  95. process.exit(1);
  96. }
  97. if (!guid) {
  98. console.error('[RelayClient] 缺少 deviceGuid:请通过命令行传入,或配置 RELAY_DEVICE_GUID / relay-config.json');
  99. process.exit(1);
  100. }
  101. console.log(`[RelayClient] 启动 Relay 轮询: ${baseUrl}`);
  102. console.log(`[RelayClient] tenantId=${tenantId}, guid=${guid}`);
  103. let backoff = INITIAL_BACKOFF_MS;
  104. while (true) {
  105. try {
  106. await runPollOnce(baseUrl, apiSecret, guid, privateKey);
  107. backoff = INITIAL_BACKOFF_MS;
  108. } catch (err) {
  109. console.error('[RelayClient] 轮询异常:', err.message);
  110. console.log(`[RelayClient] ${backoff}ms 后重试...`);
  111. await new Promise((resolve) => setTimeout(resolve, backoff));
  112. backoff = Math.min(backoff * 2, MAX_BACKOFF_MS);
  113. }
  114. }
  115. }
  116. main().catch((err) => {
  117. console.error('[RelayClient] 致命错误:', err);
  118. process.exit(1);
  119. });