| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110 |
- #!/usr/bin/env node
- import { createServer } from 'node:http';
- const APP_ID = process.env.XIAOSHU_PARSE_APP_ID || '7pIbDBJmKx_main';
- const FUNCTION_URL = (process.env.XIAOSHU_FUNCTION_URL || 'https://server.xiaoshu.pro/api/functions').replace(/\/$/, '');
- const SESSION_TOKEN = String(process.env.XIAOSHU_SYNC_OPERATOR_TOKEN || '').trim();
- const INTERVAL_MS = Math.max(5_000, Number(process.env.XIAOSHU_SYNC_INTERVAL_MS || 5_000));
- const HEALTH_PORT = Math.max(0, Number(process.env.XIAOSHU_SYNC_HEALTH_PORT || 9087));
- const DATASETS = String(process.env.XIAOSHU_SYNC_DATASETS || 'learning-records,practice-records,memory-records,assessments,course-bindings,appointments,lessons')
- .split(',').map((value) => value.trim()).filter(Boolean);
- const DEPENDENCIES = { 'memory-records': 'learning-records', lessons: 'appointments' };
- const ONCE = process.argv.includes('--once');
- if (!SESSION_TOKEN) throw new Error('缺少 XIAOSHU_SYNC_OPERATOR_TOKEN;同步运行器必须使用超级管理员会话');
- if (!DATASETS.length) throw new Error('XIAOSHU_SYNC_DATASETS 不能为空');
- for (const [dataset,dependency] of Object.entries(DEPENDENCIES)) if (DATASETS.includes(dataset) && !DATASETS.includes(dependency)) throw new Error(`${dataset} 同步必须同时包含前置数据集 ${dependency}`);
- const status = {
- startedAt: new Date().toISOString(),
- running: false,
- lastStartedAt: '',
- lastFinishedAt: '',
- lastSuccessAt: '',
- consecutiveFailures: 0,
- catchingUp: false,
- slowTicks: 0,
- datasets: {},
- };
- async function call(params) {
- const response = await fetch(`${FUNCTION_URL}/xiaoshu/ops/gateway-v3`, {
- method: 'POST',
- headers: { 'Content-Type': 'application/json', 'X-Parse-Application-Id': APP_ID },
- body: JSON.stringify({ token: SESSION_TOKEN, params }),
- signal: AbortSignal.timeout(90_000),
- });
- const payload = await response.json().catch(() => ({}));
- if (!response.ok || payload.success !== true) throw new Error(`${response.status}: ${payload.message || payload.error || '同步云函数执行失败'}`);
- return payload.data;
- }
- async function syncDataset(dataset) {
- const startedAt = Date.now();
- try {
- const result = await call({ operation: 'ops/sync/run', dataset, maxPages: 2, pageSize: 500 });
- status.datasets[dataset] = { ok: !result.failures && !result.hasMore, state: result.failures ? 'error' : result.hasMore ? 'catching_up' : 'healthy', elapsedMs: Date.now() - startedAt, ...result };
- return result;
- } catch (error) {
- status.datasets[dataset] = { ok: false, elapsedMs: Date.now() - startedAt, error: String(error?.message || error).slice(0, 500) };
- throw error;
- }
- }
- async function tick() {
- if (status.running) return;
- const tickStartedAt = Date.now();
- status.running = true;
- status.lastStartedAt = new Date().toISOString();
- const independent = DATASETS.filter((dataset) => !DEPENDENCIES[dataset]);
- const results = await Promise.allSettled(independent.map(syncDataset));
- const byDataset = new Map(independent.map((dataset,index) => [dataset,results[index]]));
- const dependents = DATASETS.filter((dataset) => DEPENDENCIES[dataset]);
- const dependentResults = await Promise.allSettled(dependents.map(async (dataset) => {
- const dependency = DEPENDENCIES[dataset],prerequisite = byDataset.get(dependency);
- if (prerequisite?.status !== 'fulfilled' || prerequisite.value.hasMore || prerequisite.value.failures) {
- status.datasets[dataset] = { ok:false,state:'waiting_for_dependency',dependency };
- return { hasMore:true,skipped:true,dependency };
- }
- return syncDataset(dataset);
- }));
- results.push(...dependentResults);
- const failures = results.filter((result) => result.status === 'rejected');
- const catchingUp = results.some((result) => result.status === 'fulfilled' && (result.value.hasMore || result.value.failures));
- status.lastFinishedAt = new Date().toISOString();
- status.running = false;
- status.catchingUp = catchingUp;
- if (failures.length) {
- status.consecutiveFailures += 1;
- console.error(`[legacy-sync] ${status.lastFinishedAt} ${failures.length}/${DATASETS.length} 个数据集失败`);
- process.exitCode = ONCE ? 1 : 0;
- } else {
- status.consecutiveFailures = 0;
- if (!catchingUp) status.lastSuccessAt = status.lastFinishedAt;
- const elapsedMs = Date.now() - tickStartedAt;
- if (elapsedMs > 10_000) {
- status.slowTicks += 1;
- console.warn(`[legacy-sync] ${status.lastFinishedAt} 同步耗时 ${elapsedMs}ms,已超过 10 秒目标`);
- } else {
- status.slowTicks = 0;
- }
- console.log(`[legacy-sync] ${status.lastFinishedAt} ${DATASETS.length} 个数据集已处理${catchingUp ? '(仍在追赶,尚未追平)' : '并追平'},耗时 ${elapsedMs}ms`);
- }
- }
- if (!ONCE && HEALTH_PORT) {
- createServer((request, response) => {
- if (request.url !== '/health') {
- response.writeHead(404).end();
- return;
- }
- const lastSuccessAge = status.lastSuccessAt ? Math.floor((Date.now() - new Date(status.lastSuccessAt).getTime()) / 1000) : null;
- const healthy = lastSuccessAge !== null && lastSuccessAge <= 15 && status.consecutiveFailures === 0 && !status.catchingUp;
- response.writeHead(healthy ? 200 : 503, { 'Content-Type': 'application/json; charset=utf-8', 'Cache-Control': 'no-store' });
- response.end(JSON.stringify({ healthy, lastSuccessAge, intervalMs: INTERVAL_MS, ...status }));
- }).listen(HEALTH_PORT, '127.0.0.1', () => console.log(`[legacy-sync] 健康检查监听 http://127.0.0.1:${HEALTH_PORT}/health`));
- }
- await tick();
- if (!ONCE) setInterval(() => void tick(), INTERVAL_MS);
|