routes-schedule.ts 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406
  1. import express from "npm:express";
  2. import { spApiScheduler } from './sp-api-schedule.ts';
  3. import { newSpApiScheduler } from './new-sp-api-schedule.ts';
  4. import { sorftimeScheduler } from './sorftime-api-schedule.ts';
  5. /**
  6. * 定时器清洗数据
  7. */
  8. (async () => {
  9. try {
  10. await spApiScheduler.start();
  11. console.log('[API Routes] Amazon SP-API定时任务已启动');
  12. } catch (error) {
  13. console.error('[API Routes] Amazon SP-API定时任务启动失败:', error);
  14. }
  15. if (process.env.ENABLE_SORFTIME_SCHEDULES === 'true') {
  16. try {
  17. await sorftimeScheduler.start();
  18. console.log('[API Routes] Sorftime定时任务已启动');
  19. } catch (error) {
  20. console.error('[API Routes] Sorftime定时任务启动失败:', error);
  21. }
  22. } else {
  23. console.log('[API Routes] Sorftime定时任务已禁用');
  24. }
  25. })();
  26. const router = express.Router();
  27. console.log('加载schedule路由 /api/schedule/');
  28. router.use(express.json({
  29. charset: 'utf-8',
  30. // 额外配置:确保解析URL编码的请求体也支持中文
  31. type: 'application/json'
  32. }));
  33. // 定时任务管理路由
  34. router.get('/sp-api-schedule/status', (req, res) => {
  35. res.json({
  36. success: true,
  37. data: {
  38. spApi: spApiScheduler.getStatus()
  39. }
  40. });
  41. });
  42. router.get('/sp-api-schedule/start', async (req, res) => {
  43. console.log('[API Routes] 启动定时任务');
  44. try {
  45. const result = await spApiScheduler.start()
  46. res.json({
  47. success: true,
  48. message: '定时器执行陈功',
  49. data: result
  50. });
  51. } catch (error) {
  52. res.status(500).json({
  53. success: false,
  54. message: error.message
  55. });
  56. }
  57. });
  58. router.post('/sp-api-schedule/trigger', async (req, res) => {
  59. try {
  60. const shopId = req.body.shopId;
  61. const result = await spApiScheduler.triggerManually(shopId);
  62. res.json({
  63. success: result.success,
  64. message: result.message,
  65. data: result
  66. });
  67. } catch (error) {
  68. res.status(500).json({
  69. success: false,
  70. message: error.message
  71. });
  72. }
  73. });
  74. // ===== 以下为新版独立采集接口(SSE 实时进度推送)=====
  75. // 前端通过 EventSource 或 fetch + ReadableStream 消费 SSE 流
  76. // 事件格式:
  77. // event: progress data: { type, message, processed? }
  78. // event: done data: { success, message, processed, errorCount, errors }
  79. // event: error data: { success: false, message }
  80. // 心跳注释行(': heartbeat')每 15 秒发送一次,防止连接超时
  81. /**
  82. * 创建 SSE 响应并返回辅助函数
  83. */
  84. function createSseResponse(res) {
  85. res.setHeader('Content-Type', 'text/event-stream; charset=utf-8');
  86. res.setHeader('Cache-Control', 'no-cache');
  87. res.setHeader('Connection', 'keep-alive');
  88. res.setHeader('X-Accel-Buffering', 'no');
  89. res.setHeader('Access-Control-Allow-Origin', '*');
  90. res.setHeader('Access-Control-Allow-Methods', 'GET, POST, OPTIONS');
  91. res.setHeader('Access-Control-Allow-Headers', 'Content-Type');
  92. res.flushHeaders();
  93. const send = (event, data) => {
  94. if (!res.writableEnded) {
  95. res.write(`event: ${event}\ndata: ${JSON.stringify(data)}\n\n`);
  96. }
  97. };
  98. const heartbeatTimer = setInterval(() => {
  99. if (!res.writableEnded) res.write(': heartbeat\n\n');
  100. }, 15000);
  101. const close = () => {
  102. clearInterval(heartbeatTimer);
  103. if (!res.writableEnded) res.end();
  104. };
  105. return { send, close };
  106. }
  107. /**
  108. * POST /api/schedule/sp-api-schedule/collect-listing
  109. * 采集 Listing 数据,通过 SSE 流式返回进度
  110. * Body: { shopId: string }
  111. */
  112. router.post('/sp-api-schedule/collect-listing', async (req, res) => {
  113. const { shopId } = req.body || {};
  114. const { send, close } = createSseResponse(res);
  115. req.on('aborted', () => close());
  116. try {
  117. send('progress', { type: 'info', message: `开始采集 Listing 数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  118. const result = await newSpApiScheduler.collectListings(
  119. shopId || undefined,
  120. (event) => send('progress', event)
  121. );
  122. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  123. } catch (error) {
  124. send('error', { success: false, message: error.message || '采集 Listing 数据时发生未知错误' });
  125. } finally {
  126. close();
  127. }
  128. });
  129. /**
  130. * POST /api/schedule/sp-api-schedule/collect-order
  131. * 采集订单数据,通过 SSE 流式返回进度
  132. * Body: { shopId: string }
  133. */
  134. router.post('/sp-api-schedule/collect-order', async (req, res) => {
  135. const { shopId } = req.body || {};
  136. const { send, close } = createSseResponse(res);
  137. req.on('aborted', () => close());
  138. try {
  139. send('progress', { type: 'info', message: `开始采集订单数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  140. const result = await newSpApiScheduler.collectOrders(
  141. shopId || undefined,
  142. (event) => send('progress', event)
  143. );
  144. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  145. } catch (error) {
  146. send('error', { success: false, message: error.message || '采集订单数据时发生未知错误' });
  147. } finally {
  148. close();
  149. }
  150. });
  151. /**
  152. * POST /api/schedule/sp-api-schedule/collect-report
  153. * 采集报表数据(创建报表任务),通过 SSE 流式返回进度
  154. * Body: { shopId: string }
  155. */
  156. router.post('/sp-api-schedule/collect-report', async (req, res) => {
  157. const { shopId } = req.body || {};
  158. const { send, close } = createSseResponse(res);
  159. req.on('aborted', () => close());
  160. try {
  161. send('progress', { type: 'info', message: `开始采集报表数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  162. const result = await newSpApiScheduler.collectReports(
  163. shopId || undefined,
  164. (event) => send('progress', event)
  165. );
  166. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  167. } catch (error) {
  168. send('error', { success: false, message: error.message || '采集报表数据时发生未知错误' });
  169. } finally {
  170. close();
  171. }
  172. });
  173. /**
  174. * POST /api/schedule/sp-api-schedule/parse-report
  175. * 检查未完成报表状态并解析已完成报表,通过 SSE 流式返回进度
  176. * Body: { shopId: string }
  177. */
  178. router.post('/sp-api-schedule/parse-report', async (req, res) => {
  179. const { shopId } = req.body || {};
  180. const { send, close } = createSseResponse(res);
  181. req.on('aborted', () => close());
  182. try {
  183. send('progress', { type: 'info', message: `开始解析报表${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  184. const result = await newSpApiScheduler.parseReports(
  185. shopId || undefined,
  186. (event) => send('progress', event)
  187. );
  188. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  189. } catch (error) {
  190. send('error', { success: false, message: error.message || '解析报表时发生未知错误' });
  191. } finally {
  192. close();
  193. }
  194. });
  195. // Sorftime 定时任务管理路由
  196. router.get('/sorftime-schedule/status', (req, res) => {
  197. res.json({
  198. success: true,
  199. data: {
  200. sorftime: sorftimeScheduler.getStatus()
  201. }
  202. });
  203. });
  204. /**
  205. * POST /api/schedule/sorftime-schedule/trigger-products
  206. * 手动触发全量类目热销产品采集(每类目分4页请求),SSE 实时返回进度
  207. * 对应定时任务:每月1日 00:00(Asia/Shanghai)
  208. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  209. * SSE done: { success, message, successCount, errorCount, errors }
  210. */
  211. router.post('/sorftime-schedule/trigger-products', async (req, res) => {
  212. const { shopId, nodeIds } = req.body || {};
  213. const { send, close } = createSseResponse(res);
  214. req.on('aborted', () => close());
  215. try {
  216. send('progress', { type: 'info', message: '开始采集类目热销产品...' });
  217. const result = await sorftimeScheduler.triggerCategoryProducts(
  218. shopId || undefined,
  219. nodeIds || undefined,
  220. (event) => send('progress', event)
  221. );
  222. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  223. } catch (error) {
  224. send('error', { success: false, message: error.message || '采集类目热销产品时发生未知错误' });
  225. } finally {
  226. close();
  227. }
  228. });
  229. /**
  230. * POST /api/schedule/sorftime-schedule/trigger-keywords
  231. * 手动触发全量类目关键词采集,SSE 实时返回进度
  232. * 对应定时任务:每周一 01:00(Asia/Shanghai)
  233. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  234. * SSE done: { success, message, successCount, errorCount, errors }
  235. */
  236. router.post('/sorftime-schedule/trigger-keywords', async (req, res) => {
  237. const { shopId, nodeIds } = req.body || {};
  238. const { send, close } = createSseResponse(res);
  239. req.on('aborted', () => close());
  240. try {
  241. send('progress', { type: 'info', message: '开始采集类目关键词...' });
  242. const result = await sorftimeScheduler.triggerCategoryKeywords(
  243. shopId || undefined,
  244. nodeIds || undefined,
  245. (event) => send('progress', event)
  246. );
  247. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  248. } catch (error) {
  249. send('error', { success: false, message: error.message || '采集类目关键词时发生未知错误' });
  250. } finally {
  251. close();
  252. }
  253. });
  254. /**
  255. * POST /api/schedule/sorftime-schedule/trigger-category-trends
  256. * 手动触发全量类目市场趋势采集(每类目6种 TrendIndex),SSE 实时返回进度
  257. * 对应定时任务:每月1日 01:00(Asia/Shanghai)
  258. * Body: { shopId?: string, nodeIds?: string[] } 两者均不传则采集全部叶子节点类目
  259. * SSE done: { success, message, successCount, errorCount, errors }
  260. */
  261. router.post('/sorftime-schedule/trigger-category-trends', async (req, res) => {
  262. const { shopId, nodeIds } = req.body || {};
  263. const { send, close } = createSseResponse(res);
  264. req.on('aborted', () => close());
  265. try {
  266. send('progress', { type: 'info', message: '开始采集类目市场趋势...' });
  267. const result = await sorftimeScheduler.triggerCategoryTrends(
  268. shopId || undefined,
  269. nodeIds || undefined,
  270. (event) => send('progress', event)
  271. );
  272. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  273. } catch (error) {
  274. send('error', { success: false, message: error.message || '采集类目市场趋势时发生未知错误' });
  275. } finally {
  276. close();
  277. }
  278. });
  279. /**
  280. * POST /api/schedule/sorftime-schedule/trigger-single-category
  281. * 按需触发单个类目的全量数据采集(热销产品 + 关键词 + 市场趋势),SSE 实时返回进度
  282. * Body: { nodeId: string, domain?: number } domain 默认 1(美国站)
  283. * SSE done: { success, message, tasks: { products, keywords, trends }, duration, errors }
  284. */
  285. router.post('/sorftime-schedule/trigger-single-category', async (req, res) => {
  286. const { nodeId, domain } = req.body || {};
  287. if (!nodeId) {
  288. return res.status(400).json({ success: false, message: '缺少 nodeId 参数' });
  289. }
  290. const { send, close } = createSseResponse(res);
  291. req.on('aborted', () => close());
  292. try {
  293. send('progress', { type: 'info', message: `开始采集类目 ${nodeId} 全量数据(热销产品 + 关键词 + 市场趋势)` });
  294. const result = await sorftimeScheduler.triggerSingleCategory(
  295. nodeId,
  296. domain || 1,
  297. (event) => send('progress', event)
  298. );
  299. send('done', { success: result.success, message: result.message, tasks: result.tasks, duration: result.duration, errors: result.errors });
  300. } catch (error) {
  301. send('error', { success: false, message: error.message || '采集单类目数据时发生未知错误' });
  302. } finally {
  303. close();
  304. }
  305. });
  306. /**
  307. * POST /api/schedule/sorftime-schedule/collect-reviews
  308. * 手动触发产品评论采集(遍历 Product 表 ASIN 请求 ProductReviewsQuery),SSE 实时返回进度
  309. * Body: 无
  310. * SSE done: { success, message, successCount, errorCount, errors }
  311. */
  312. router.post('/sorftime-schedule/collect-reviews', async (req, res) => {
  313. const { send, close } = createSseResponse(res);
  314. req.on('aborted', () => close());
  315. try {
  316. send('progress', { type: 'info', message: '开始采集产品评论数据...' });
  317. const result = await sorftimeScheduler.triggerProductReviews(
  318. (event) => send('progress', event)
  319. );
  320. send('done', { success: result.success, message: result.message, successCount: result.successCount, errorCount: result.errorCount, errors: result.errors });
  321. } catch (error) {
  322. send('error', { success: false, message: error.message || '采集产品评论时发生未知错误' });
  323. } finally {
  324. close();
  325. }
  326. });
  327. /**
  328. * POST /api/schedule/sorftime-schedule/collect-shop-products
  329. * 采集店铺竞品数据(通过 SellerId QueryType=5 分页请求),SSE 实时返回进度
  330. * Body: { shopId?: string }
  331. */
  332. router.post('/sorftime-schedule/collect-shop-products', async (req, res) => {
  333. const { shopId } = req.body || {};
  334. const { send, close } = createSseResponse(res);
  335. req.on('aborted', () => close());
  336. try {
  337. send('progress', { type: 'info', message: `开始采集店铺产品数据${shopId ? ',店铺: ' + shopId : '(全部活跃店铺)'}` });
  338. const result = await sorftimeScheduler.collectShopProducts(
  339. shopId || undefined,
  340. (event) => send('progress', event)
  341. );
  342. send('done', { success: result.success, message: result.message, processed: result.processed, errorCount: result.errorCount, errors: result.errors });
  343. } catch (error) {
  344. send('error', { success: false, message: error.message || '采集店铺产品时发生未知错误' });
  345. } finally {
  346. close();
  347. }
  348. });
  349. router.get('/test', (req, res)=> {
  350. res.json({
  351. message: "测试 sp-api-schedule",
  352. body: req.body
  353. });
  354. })
  355. export default router;