#!/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);