run-legacy-sync-worker.mjs 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293
  1. #!/usr/bin/env node
  2. import { createServer } from 'node:http';
  3. const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
  4. const FUNCTION_URL = (process.env.XIAOSHU_FUNCTION_URL || 'https://server.xiaoshu.pro/api/functions').replace(/\/$/, '');
  5. const SESSION_TOKEN = String(process.env.XIAOSHU_SYNC_OPERATOR_TOKEN || '').trim();
  6. const INTERVAL_MS = Math.max(5_000, Number(process.env.XIAOSHU_SYNC_INTERVAL_MS || 5_000));
  7. const HEALTH_PORT = Math.max(0, Number(process.env.XIAOSHU_SYNC_HEALTH_PORT || 9087));
  8. const DATASETS = String(process.env.XIAOSHU_SYNC_DATASETS || 'learning-records,practice-records,memory-records,assessments,course-bindings,appointments,lessons')
  9. .split(',').map((value) => value.trim()).filter(Boolean);
  10. const ONCE = process.argv.includes('--once');
  11. if (!SESSION_TOKEN) throw new Error('缺少 XIAOSHU_SYNC_OPERATOR_TOKEN;同步运行器必须使用超级管理员会话');
  12. if (!DATASETS.length) throw new Error('XIAOSHU_SYNC_DATASETS 不能为空');
  13. const status = {
  14. startedAt: new Date().toISOString(),
  15. running: false,
  16. lastStartedAt: '',
  17. lastFinishedAt: '',
  18. lastSuccessAt: '',
  19. consecutiveFailures: 0,
  20. slowTicks: 0,
  21. datasets: {},
  22. };
  23. async function call(params) {
  24. const response = await fetch(`${FUNCTION_URL}/xiaoshu/ops/gateway-v3`, {
  25. method: 'POST',
  26. headers: { 'Content-Type': 'application/json', 'X-Parse-Application-Id': APP_ID },
  27. body: JSON.stringify({ token: SESSION_TOKEN, params }),
  28. signal: AbortSignal.timeout(90_000),
  29. });
  30. const payload = await response.json().catch(() => ({}));
  31. if (!response.ok || payload.success !== true) throw new Error(`${response.status}: ${payload.message || payload.error || '同步云函数执行失败'}`);
  32. return payload.data;
  33. }
  34. async function syncDataset(dataset) {
  35. const startedAt = Date.now();
  36. try {
  37. const result = await call({ operation: 'ops/sync/run', dataset, maxPages: 2, pageSize: 500 });
  38. status.datasets[dataset] = { ok: true, elapsedMs: Date.now() - startedAt, ...result };
  39. return result;
  40. } catch (error) {
  41. status.datasets[dataset] = { ok: false, elapsedMs: Date.now() - startedAt, error: String(error?.message || error).slice(0, 500) };
  42. throw error;
  43. }
  44. }
  45. async function tick() {
  46. if (status.running) return;
  47. const tickStartedAt = Date.now();
  48. status.running = true;
  49. status.lastStartedAt = new Date().toISOString();
  50. const results = await Promise.allSettled(DATASETS.map(syncDataset));
  51. const failures = results.filter((result) => result.status === 'rejected');
  52. status.lastFinishedAt = new Date().toISOString();
  53. status.running = false;
  54. if (failures.length) {
  55. status.consecutiveFailures += 1;
  56. console.error(`[legacy-sync] ${status.lastFinishedAt} ${failures.length}/${DATASETS.length} 个数据集失败`);
  57. process.exitCode = ONCE ? 1 : 0;
  58. } else {
  59. status.consecutiveFailures = 0;
  60. status.lastSuccessAt = status.lastFinishedAt;
  61. const elapsedMs = Date.now() - tickStartedAt;
  62. if (elapsedMs > 10_000) {
  63. status.slowTicks += 1;
  64. console.warn(`[legacy-sync] ${status.lastFinishedAt} 同步耗时 ${elapsedMs}ms,已超过 10 秒目标`);
  65. } else {
  66. status.slowTicks = 0;
  67. }
  68. console.log(`[legacy-sync] ${status.lastFinishedAt} ${DATASETS.length} 个数据集同步完成,耗时 ${elapsedMs}ms`);
  69. }
  70. }
  71. if (!ONCE && HEALTH_PORT) {
  72. createServer((request, response) => {
  73. if (request.url !== '/health') {
  74. response.writeHead(404).end();
  75. return;
  76. }
  77. const lastSuccessAge = status.lastSuccessAt ? Math.floor((Date.now() - new Date(status.lastSuccessAt).getTime()) / 1000) : null;
  78. const healthy = lastSuccessAge !== null && lastSuccessAge <= 15 && status.consecutiveFailures === 0;
  79. response.writeHead(healthy ? 200 : 503, { 'Content-Type': 'application/json; charset=utf-8', 'Cache-Control': 'no-store' });
  80. response.end(JSON.stringify({ healthy, lastSuccessAge, intervalMs: INTERVAL_MS, ...status }));
  81. }).listen(HEALTH_PORT, '127.0.0.1', () => console.log(`[legacy-sync] 健康检查监听 http://127.0.0.1:${HEALTH_PORT}/health`));
  82. }
  83. await tick();
  84. if (!ONCE) setInterval(() => void tick(), INTERVAL_MS);