reparse-sales-traffic-reports.mjs 3.0 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980
  1. import zlib from 'node:zlib';
  2. import {
  3. activeAmazonShops, allRows, dateValue, forward, hash, log,
  4. shopPointer, sleep, updateObject, upsertMany,
  5. } from './sp-api-collector-lib.mjs';
  6. function parseDocument(content) {
  7. const trimmed = content.replace(/^\uFEFF/, '').trim();
  8. if (!trimmed) return [];
  9. const parsed = JSON.parse(trimmed);
  10. const rows = [];
  11. for (const [section, value] of Object.entries(parsed)) {
  12. if (!Array.isArray(value)) continue;
  13. for (const item of value) {
  14. rows.push({ _section: section, ...(item && typeof item === 'object' ? item : { value: item }) });
  15. }
  16. }
  17. return rows.length ? rows : [parsed];
  18. }
  19. const shops = await activeAmazonShops();
  20. const reports = (await allRows('Reports')).filter(report =>
  21. report.recordKey?.startsWith('MP:')
  22. && report.reportType === 'GET_SALES_AND_TRAFFIC_REPORT'
  23. && report.isParsed !== true
  24. && report.processingStatus === 'DONE'
  25. && report.reportDocumentId
  26. );
  27. for (const report of reports) {
  28. const shop = shops.find(item => item.objectId === report.shop?.objectId);
  29. if (!shop) {
  30. log('sales_report_skipped', { reportId: report.reportId, reason: 'shop not found' });
  31. continue;
  32. }
  33. try {
  34. const document = await forward(shop.objectId,
  35. `/reports/2021-06-30/documents/${encodeURIComponent(report.reportDocumentId)}`,
  36. 'GET', undefined, 5);
  37. if (!document.url) throw new Error('Report document returned no URL');
  38. const response = await fetch(document.url, { signal: AbortSignal.timeout(300000) });
  39. if (!response.ok) throw new Error(`Report download failed: ${response.status}`);
  40. let bytes = Buffer.from(await response.arrayBuffer());
  41. if (document.compressionAlgorithm === 'GZIP') bytes = zlib.gunzipSync(bytes);
  42. const rows = parseDocument(new TextDecoder('utf-8').decode(bytes));
  43. const marketplaceKey = [...(report.marketplaceIds || [])].sort().join(',') || shop.marketplaceId;
  44. const objects = rows.map(row => {
  45. const rowHash = hash(row);
  46. return {
  47. recordKey: `${report.region || 'UNKNOWN'}:${marketplaceKey}:${report.reportType}:${rowHash}`,
  48. rowHash,
  49. reportType: report.reportType,
  50. reportId: report.reportId,
  51. region: report.region || '',
  52. marketplaceIds: report.marketplaceIds || [],
  53. capturedAt: dateValue(new Date()),
  54. rowData: row,
  55. rawData: row,
  56. shop: shopPointer(shop.objectId),
  57. };
  58. });
  59. const saved = await upsertMany('AmazonReportRow', objects);
  60. await updateObject('Reports', report.objectId, {
  61. isParsed: true,
  62. parsedAt: dateValue(new Date()),
  63. parsedRecords: objects.length,
  64. });
  65. log('sales_report_reparsed', {
  66. marketplaceIds: report.marketplaceIds,
  67. reportId: report.reportId,
  68. rows: objects.length,
  69. ...saved,
  70. });
  71. } catch (error) {
  72. log('sales_report_reparse_failed', { reportId: report.reportId, error: error.message });
  73. }
  74. await sleep(1000);
  75. }
  76. log('sales_report_reparse_finished', { reports: reports.length });