import express from "npm:express"; import { spApiScheduler } from './sp-api-schedule.ts'; import { newSpApiScheduler } from './new-sp-api-schedule.ts'; import { sorftimeScheduler } from './sorftime-api-schedule.ts'; /** * 定时器清洗数据 */ (async () => { try { await spApiScheduler.start(); console.log('[API Routes] Amazon SP-API定时任务已启动'); } catch (error) { console.error('[API Routes] Amazon SP-API定时任务启动失败:', error); } try { await sorftimeScheduler.start(); console.log('[API Routes] Sorftime定时任务已启动'); } catch (error) { console.error('[API Routes] Sorftime定时任务启动失败:', error); } })(); const router = express.Router(); console.log('加载schedule路由 /api/schedule/'); router.use(express.json({ charset: 'utf-8', // 额外配置:确保解析URL编码的请求体也支持中文 type: 'application/json' })); // 定时任务管理路由 router.get('/sp-api-schedule/status', (req, res) => { res.json({ success: true, data: { spApi: spApiScheduler.getStatus() } }); }); router.get('/sp-api-schedule/start', async (req, res) => { console.log('[API Routes] 启动定时任务'); try { const result = await spApiScheduler.start() res.json({ success: true, message: '定时器执行陈功', data: result }); } catch (error) { res.status(500).json({ success: false, message: error.message }); } }); router.post('/sp-api-schedule/trigger', async (req, res) => { try { const shopId = req.body.shopId; const result = await spApiScheduler.triggerManually(shopId); res.json({ success: result.success, message: result.message, data: result }); } catch (error) { res.status(500).json({ success: false, message: error.message }); } }); // ===== 以下为新版独立采集接口(SSE 实时进度推送)===== // 前端通过 EventSource 或 fetch + ReadableStream 消费 SSE 流 // 事件格式: // event: progress data: { type, message, processed? } // event: done data: { success, message, processed, errorCount, errors } // event: error data: { success: false, message } // 心跳注释行(': heartbeat')每 15 秒发送一次,防止连接超时 /** * 创建 SSE 响应并返回辅助函数 */ function createSseResponse(res) { res.setHeader('Content-Type', 'text/event-stream; charset=utf-8'); res.setHeader('Cache-Control', 'no-cache'); res.setHeader('Connection', 'keep-alive'); res.setHeader('X-Accel-Buffering', 'no'); res.setHeader('Access-Control-Allow-Origin', '*'); res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS'); res.setHeader('Access-Control-Allow-Headers', 'Content-Type'); res.flushHeaders(); const send = (event, data) => { if (!res.writableEnded) { res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`); } }; const heartbeatTimer = setInterval(() => { if (!res.writableEnded) res.write(': heartbeat\n\n'); }, 15000); const close = () => { clearInterval(heartbeatTimer); if (!res.writableEnded) res.end(); }; return { send, close }; } /** * POST /api/schedule/sp-api-schedule/collect-listing * 采集 Listing 数据,通过 SSE 流式返回进度 * Body: { shopId: string } */ router.post('/sp-api-schedule/collect-listing', async (req, res) => { const { shopId } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始采集 Listing 数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` }); const result = await newSpApiScheduler.collectListings( shopId || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集 Listing 数据时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sp-api-schedule/collect-order * 采集订单数据,通过 SSE 流式返回进度 * Body: { shopId: string } */ router.post('/sp-api-schedule/collect-order', async (req, res) => { const { shopId } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始采集订单数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` }); const result = await newSpApiScheduler.collectOrders( shopId || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集订单数据时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sp-api-schedule/collect-report * 采集报表数据(创建报表任务),通过 SSE 流式返回进度 * Body: { shopId: string } */ router.post('/sp-api-schedule/collect-report', async (req, res) => { const { shopId } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始采集报表数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` }); const result = await newSpApiScheduler.collectReports( shopId || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集报表数据时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sp-api-schedule/parse-report * 检查未完成报表状态并解析已完成报表,通过 SSE 流式返回进度 * Body: { shopId: string } */ router.post('/sp-api-schedule/parse-report', async (req, res) => { const { shopId } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始解析报表${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` }); const result = await newSpApiScheduler.parseReports( shopId || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '解析报表时发生未知错误' }); } finally { close(); } }); // Sorftime 定时任务管理路由 router.get('/sorftime-schedule/status', (req, res) => { res.json({ success: true, data: { sorftime: sorftimeScheduler.getStatus() } }); }); /** * POST /api/schedule/sorftime-schedule/trigger-products * 手动触发全量类目热销产品采集(每类目分4页请求),SSE 实时返回进度 * 对应定时任务:每月1日 00:00(Asia/Shanghai) * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目 * SSE done: { success, message, successCount, errorCount, errors } */ router.post('/sorftime-schedule/trigger-products', async (req, res) => { const { shopId, nodeIds } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: '开始采集类目热销产品...' }); const result = await sorftimeScheduler.triggerCategoryProducts( shopId || undefined, nodeIds || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集类目热销产品时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sorftime-schedule/trigger-keywords * 手动触发全量类目关键词采集,SSE 实时返回进度 * 对应定时任务:每周一 01:00(Asia/Shanghai) * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目 * SSE done: { success, message, successCount, errorCount, errors } */ router.post('/sorftime-schedule/trigger-keywords', async (req, res) => { const { shopId, nodeIds } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: '开始采集类目关键词...' }); const result = await sorftimeScheduler.triggerCategoryKeywords( shopId || undefined, nodeIds || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集类目关键词时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sorftime-schedule/trigger-category-trends * 手动触发全量类目市场趋势采集(每类目6种 TrendIndex),SSE 实时返回进度 * 对应定时任务:每月1日 01:00(Asia/Shanghai) * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目 * SSE done: { success, message, successCount, errorCount, errors } */ router.post('/sorftime-schedule/trigger-category-trends', async (req, res) => { const { shopId, nodeIds } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: '开始采集类目市场趋势...' }); const result = await sorftimeScheduler.triggerCategoryTrends( shopId || undefined, nodeIds || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集类目市场趋势时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sorftime-schedule/trigger-single-category * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势),SSE 实时返回进度 * Body: { nodeId: string, domain?: number } domain 默认 1(美国站) * SSE done: { success, message, tasks: { products, keywords, trends }, duration, errors } */ router.post('/sorftime-schedule/trigger-single-category', async (req, res) => { const { nodeId, domain } = req.body || {}; if (!nodeId) { return res.status(400).json({ success: false, message: '缺少 nodeId 参数' }); } const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始采集类目 ${nodeId} 全量数据(热销产品 + 关键词 + 市场趋势)` }); const result = await sorftimeScheduler.triggerSingleCategory( nodeId, domain || 1, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, tasks: result.tasks, duration: result.duration, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集单类目数据时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sorftime-schedule/collect-reviews * 手动触发产品评论采集(遍历 Product 表 ASIN 请求 ProductReviewsQuery),SSE 实时返回进度 * Body: 无 * SSE done: { success, message, successCount, errorCount, errors } */ router.post('/sorftime-schedule/collect-reviews', async (req, res) => { const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: '开始采集产品评论数据...' }); const result = await sorftimeScheduler.triggerProductReviews( (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集产品评论时发生未知错误' }); } finally { close(); } }); /** * POST /api/schedule/sorftime-schedule/collect-shop-products * 采集店铺竞品数据(通过 SellerId QueryType=5 分页请求),SSE 实时返回进度 * Body: { shopId?: string } */ router.post('/sorftime-schedule/collect-shop-products', async (req, res) => { const { shopId } = req.body || {}; const { send, close } = createSseResponse(res); req.on('close', () => close()); try { send('progress', { type: 'info', message: `开始采集店铺产品数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` }); const result = await sorftimeScheduler.collectShopProducts( shopId || undefined, (event) => send('progress', event) ); send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors }); } catch (error) { send('error', { success: false, message: error.message || '采集店铺产品时发生未知错误' }); } finally { close(); } }); router.get('/test', (req, res)=> { res.json({ message: "测试 sp-api-schedule", body: req.body }); }) export default router;