start-relay-client.js 4.4 KB

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