start-relay-client.js 5.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162
  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 { getProductMode } = require('../mcp/src/core/product-mode');
  16. const {
  17. getRelayBaseUrl,
  18. getTenantApiSecret,
  19. getRelayPrivateKey,
  20. getTenantId,
  21. getRelayDeviceGuid
  22. } = require('../mcp/src/core/relay-config');
  23. const POLL_WAIT_MS = 30000;
  24. const INITIAL_BACKOFF_MS = 1000;
  25. const MAX_BACKOFF_MS = 60000;
  26. function decryptPayload(encryptedPayload, privateKeyPem) {
  27. const key = crypto.createPrivateKey(privateKeyPem);
  28. if (String(encryptedPayload).startsWith('v2:')) {
  29. const envelope = JSON.parse(Buffer.from(String(encryptedPayload).slice(3), 'base64').toString('utf8'));
  30. const aesKey = crypto.privateDecrypt(
  31. { key, oaepHash: 'sha256' },
  32. Buffer.from(envelope.key, 'base64')
  33. );
  34. const decipher = crypto.createDecipheriv(
  35. 'aes-256-gcm',
  36. aesKey,
  37. Buffer.from(envelope.iv, 'base64')
  38. );
  39. decipher.setAuthTag(Buffer.from(envelope.tag, 'base64'));
  40. return Buffer.concat([
  41. decipher.update(Buffer.from(envelope.ciphertext, 'base64')),
  42. decipher.final()
  43. ]).toString('utf8');
  44. }
  45. const buffer = Buffer.from(encryptedPayload, 'base64');
  46. const decrypted = crypto.privateDecrypt({ key, oaepHash: 'sha256' }, buffer);
  47. return decrypted.toString('utf8');
  48. }
  49. async function ackEvents(baseUrl, apiSecret, guid, eventIds) {
  50. if (!eventIds.length) return;
  51. try {
  52. const res = await fetch(`${baseUrl}/api/relay/ack`, {
  53. method: 'POST',
  54. headers: {
  55. 'Content-Type': 'application/json',
  56. Authorization: `Bearer ${apiSecret}`
  57. },
  58. body: JSON.stringify({ guid, eventIds })
  59. });
  60. if (!res.ok) {
  61. console.warn('[RelayClient] ACK 失败:', res.status, await res.text());
  62. return;
  63. }
  64. const data = await res.json();
  65. console.log(`[RelayClient] ACK ${data.ackedCount || eventIds.length} 条事件`);
  66. } catch (err) {
  67. console.warn('[RelayClient] ACK 请求异常:', err.message);
  68. }
  69. }
  70. async function runPollOnce(baseUrl, apiSecret, guid, privateKey) {
  71. const res = await fetch(`${baseUrl}/api/relay/poll`, {
  72. method: 'POST',
  73. headers: {
  74. 'Content-Type': 'application/json',
  75. Authorization: `Bearer ${apiSecret}`
  76. },
  77. body: JSON.stringify({ guid, batchSize: 100, waitMs: POLL_WAIT_MS })
  78. });
  79. if (!res.ok) {
  80. throw new Error(`poll failed: ${res.status} ${await res.text()}`);
  81. }
  82. const data = await res.json();
  83. if (!data.events || !data.events.length) return { acked: 0 };
  84. console.log(`[RelayClient] 取回 ${data.events.length} 条事件`);
  85. const eventIds = [];
  86. for (const event of data.events) {
  87. try {
  88. const decrypted = decryptPayload(event.encryptedPayload, privateKey);
  89. const payload = JSON.parse(decrypted);
  90. const envelope = { code: 0, msg: 'from-relay', data: Array.isArray(payload) ? payload : [payload], __rawBody: decrypted };
  91. await processWebhookEvents(envelope);
  92. eventIds.push(event.eventId);
  93. } catch (err) {
  94. console.error(`[RelayClient] 解密/处理失败 eventId=${event.eventId}:`, err.message);
  95. // 解密失败也要 ACK,避免 Relay 重复投递
  96. eventIds.push(event.eventId);
  97. }
  98. }
  99. await ackEvents(baseUrl, apiSecret, guid, eventIds);
  100. return { acked: eventIds.length };
  101. }
  102. function resolveGuid() {
  103. // 命令行参数 > 环境变量 > relay-config.json > credentials
  104. return process.argv[2] || process.env.RELAY_DEVICE_GUID || getRelayDeviceGuid() || readQiweiGuid() || '';
  105. }
  106. async function main() {
  107. const product = getProductMode();
  108. if (product.mode !== 'enterprise') {
  109. console.error('[RelayClient] 当前为个人版,请使用工作台本地监听;企业 Relay Client 未启动');
  110. process.exit(1);
  111. }
  112. const baseUrl = getRelayBaseUrl();
  113. const apiSecret = getTenantApiSecret();
  114. const privateKey = getRelayPrivateKey();
  115. const tenantId = getTenantId();
  116. const guid = resolveGuid();
  117. if (!apiSecret || !privateKey || !tenantId) {
  118. console.error('[RelayClient] 缺少配置:请检查 .env.local 中的 TENANT_API_SECRET、RELAY_PRIVATE_KEY、TENANT_ID');
  119. process.exit(1);
  120. }
  121. if (!guid) {
  122. console.error('[RelayClient] 缺少 deviceGuid:请通过命令行传入,或配置 RELAY_DEVICE_GUID / relay-config.json / 完成企微登录');
  123. process.exit(1);
  124. }
  125. console.log(`[RelayClient] 启动 Relay 轮询: ${baseUrl}`);
  126. console.log(`[RelayClient] tenantId=${tenantId}, guid=${guid}`);
  127. let backoff = INITIAL_BACKOFF_MS;
  128. while (true) {
  129. try {
  130. await runPollOnce(baseUrl, apiSecret, guid, privateKey);
  131. backoff = INITIAL_BACKOFF_MS;
  132. } catch (err) {
  133. console.error('[RelayClient] 轮询异常:', err.message);
  134. console.log(`[RelayClient] ${backoff}ms 后重试...`);
  135. await new Promise((resolve) => setTimeout(resolve, backoff));
  136. backoff = Math.min(backoff * 2, MAX_BACKOFF_MS);
  137. }
  138. }
  139. }
  140. main().catch((err) => {
  141. console.error('[RelayClient] 致命错误:', err);
  142. process.exit(1);
  143. });