run-legacy-sync-worker.mjs 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  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 DEPENDENCIES = { 'memory-records': 'learning-records', lessons: 'appointments' };
  11. const ONCE = process.argv.includes('--once');
  12. if (!SESSION_TOKEN) throw new Error('缺少 XIAOSHU_SYNC_OPERATOR_TOKEN;同步运行器必须使用超级管理员会话');
  13. if (!DATASETS.length) throw new Error('XIAOSHU_SYNC_DATASETS 不能为空');
  14. for (const [dataset,dependency] of Object.entries(DEPENDENCIES)) if (DATASETS.includes(dataset) && !DATASETS.includes(dependency)) throw new Error(`${dataset} 同步必须同时包含前置数据集 ${dependency}`);
  15. const status = {
  16. startedAt: new Date().toISOString(),
  17. running: false,
  18. lastStartedAt: '',
  19. lastFinishedAt: '',
  20. lastSuccessAt: '',
  21. consecutiveFailures: 0,
  22. catchingUp: false,
  23. slowTicks: 0,
  24. datasets: {},
  25. };
  26. async function call(params) {
  27. const response = await fetch(`${FUNCTION_URL}/xiaoshu/ops/gateway-v3`, {
  28. method: 'POST',
  29. headers: { 'Content-Type': 'application/json', 'X-Parse-Application-Id': APP_ID },
  30. body: JSON.stringify({ token: SESSION_TOKEN, params }),
  31. signal: AbortSignal.timeout(90_000),
  32. });
  33. const payload = await response.json().catch(() => ({}));
  34. if (!response.ok || payload.success !== true) throw new Error(`${response.status}: ${payload.message || payload.error || '同步云函数执行失败'}`);
  35. return payload.data;
  36. }
  37. async function syncDataset(dataset) {
  38. const startedAt = Date.now();
  39. try {
  40. const result = await call({ operation: 'ops/sync/run', dataset, maxPages: 2, pageSize: 500 });
  41. status.datasets[dataset] = { ok: !result.failures && !result.hasMore, state: result.failures ? 'error' : result.hasMore ? 'catching_up' : 'healthy', elapsedMs: Date.now() - startedAt, ...result };
  42. return result;
  43. } catch (error) {
  44. status.datasets[dataset] = { ok: false, elapsedMs: Date.now() - startedAt, error: String(error?.message || error).slice(0, 500) };
  45. throw error;
  46. }
  47. }
  48. async function tick() {
  49. if (status.running) return;
  50. const tickStartedAt = Date.now();
  51. status.running = true;
  52. status.lastStartedAt = new Date().toISOString();
  53. const independent = DATASETS.filter((dataset) => !DEPENDENCIES[dataset]);
  54. const results = await Promise.allSettled(independent.map(syncDataset));
  55. const byDataset = new Map(independent.map((dataset,index) => [dataset,results[index]]));
  56. const dependents = DATASETS.filter((dataset) => DEPENDENCIES[dataset]);
  57. const dependentResults = await Promise.allSettled(dependents.map(async (dataset) => {
  58. const dependency = DEPENDENCIES[dataset],prerequisite = byDataset.get(dependency);
  59. if (prerequisite?.status !== 'fulfilled' || prerequisite.value.hasMore || prerequisite.value.failures) {
  60. status.datasets[dataset] = { ok:false,state:'waiting_for_dependency',dependency };
  61. return { hasMore:true,skipped:true,dependency };
  62. }
  63. return syncDataset(dataset);
  64. }));
  65. results.push(...dependentResults);
  66. const failures = results.filter((result) => result.status === 'rejected');
  67. const catchingUp = results.some((result) => result.status === 'fulfilled' && (result.value.hasMore || result.value.failures));
  68. status.lastFinishedAt = new Date().toISOString();
  69. status.running = false;
  70. status.catchingUp = catchingUp;
  71. if (failures.length) {
  72. status.consecutiveFailures += 1;
  73. console.error(`[legacy-sync] ${status.lastFinishedAt} ${failures.length}/${DATASETS.length} 个数据集失败`);
  74. process.exitCode = ONCE ? 1 : 0;
  75. } else {
  76. status.consecutiveFailures = 0;
  77. if (!catchingUp) status.lastSuccessAt = status.lastFinishedAt;
  78. const elapsedMs = Date.now() - tickStartedAt;
  79. if (elapsedMs > 10_000) {
  80. status.slowTicks += 1;
  81. console.warn(`[legacy-sync] ${status.lastFinishedAt} 同步耗时 ${elapsedMs}ms,已超过 10 秒目标`);
  82. } else {
  83. status.slowTicks = 0;
  84. }
  85. console.log(`[legacy-sync] ${status.lastFinishedAt} ${DATASETS.length} 个数据集已处理${catchingUp ? '(仍在追赶,尚未追平)' : '并追平'},耗时 ${elapsedMs}ms`);
  86. }
  87. }
  88. if (!ONCE && HEALTH_PORT) {
  89. createServer((request, response) => {
  90. if (request.url !== '/health') {
  91. response.writeHead(404).end();
  92. return;
  93. }
  94. const lastSuccessAge = status.lastSuccessAt ? Math.floor((Date.now() - new Date(status.lastSuccessAt).getTime()) / 1000) : null;
  95. const healthy = lastSuccessAge !== null && lastSuccessAge <= 15 && status.consecutiveFailures === 0 && !status.catchingUp;
  96. response.writeHead(healthy ? 200 : 503, { 'Content-Type': 'application/json; charset=utf-8', 'Cache-Control': 'no-store' });
  97. response.end(JSON.stringify({ healthy, lastSuccessAge, intervalMs: INTERVAL_MS, ...status }));
  98. }).listen(HEALTH_PORT, '127.0.0.1', () => console.log(`[legacy-sync] 健康检查监听 http://127.0.0.1:${HEALTH_PORT}/health`));
  99. }
  100. await tick();
  101. if (!ONCE) setInterval(() => void tick(), INTERVAL_MS);