collect-sp-api-reports-all.mjs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310
  1. import zlib from 'node:zlib';
  2. import {
  3. activeAmazonShops, allRows, dateValue, forward, hash, log, regionByMarketplace,
  4. shopPointer, sleep, updateObject, upsertMany, writeAudit,
  5. } from './sp-api-collector-lib.mjs';
  6. const allReportDefinitions = [
  7. { type: 'GET_MERCHANT_LISTINGS_ALL_DATA' },
  8. { type: 'GET_FBA_MYI_UNSUPPRESSED_INVENTORY_DATA' },
  9. { type: 'GET_FBA_MYI_ALL_INVENTORY_DATA' },
  10. { type: 'GET_FBA_INVENTORY_AGED_DATA' },
  11. { type: 'GET_FBA_INVENTORY_PLANNING_DATA' },
  12. { type: 'GET_RESERVED_INVENTORY_DATA' },
  13. { type: 'GET_FBA_ESTIMATED_FBA_FEES_TXT_DATA' },
  14. { type: 'GET_FBA_FULFILLMENT_CUSTOMER_RETURNS_DATA', dateRange: true },
  15. { type: 'GET_FBA_REIMBURSEMENTS_DATA', dateRange: true },
  16. { type: 'GET_AMAZON_FULFILLED_SHIPMENTS_DATA_GENERAL', dateRange: true },
  17. { type: 'GET_FBA_FULFILLMENT_INVENTORY_RECEIPTS_DATA', dateRange: true },
  18. { type: 'GET_FBA_FULFILLMENT_INVENTORY_ADJUSTMENTS_DATA', dateRange: true },
  19. { type: 'GET_FBA_FULFILLMENT_REMOVAL_ORDER_DETAIL_DATA', dateRange: true },
  20. { type: 'GET_FBA_FULFILLMENT_REMOVAL_SHIPMENT_DETAIL_DATA', dateRange: true },
  21. {
  22. type: 'GET_SALES_AND_TRAFFIC_REPORT',
  23. dateRange: true,
  24. singleMarketplace: true,
  25. reportOptions: { dateGranularity: 'DAY', asinGranularity: 'SKU' },
  26. },
  27. ];
  28. const requestedReportTypes = new Set(
  29. String(process.env.SP_API_REPORT_TYPES || '')
  30. .split(',')
  31. .map(value => value.trim())
  32. .filter(Boolean),
  33. );
  34. const reportDefinitions = requestedReportTypes.size
  35. ? allReportDefinitions.filter(definition => requestedReportTypes.has(definition.type))
  36. : allReportDefinitions;
  37. const shops = await activeAmazonShops();
  38. const regionShops = new Map();
  39. for (const shop of shops) {
  40. const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
  41. if (!regionShops.has(region)) regionShops.set(region, { shop, marketplaceIds: [] });
  42. regionShops.get(region).marketplaceIds.push(shop.marketplaceId);
  43. }
  44. const dateKey = new Date().toISOString().slice(0, 10);
  45. const dataEndTime = new Date().toISOString();
  46. const dataStartTime = new Date(Date.now() - 89 * 24 * 60 * 60 * 1000).toISOString();
  47. const existingReports = await allRows('Reports');
  48. const existingByKey = new Map(existingReports.filter(report => report.recordKey).map(report => [report.recordKey, report]));
  49. const created = [];
  50. for (const [region, { shop, marketplaceIds }] of regionShops) {
  51. for (const definition of reportDefinitions) {
  52. // Merchant listings do not identify the originating marketplace when a
  53. // regional request contains multiple marketplace IDs. Collect them per
  54. // marketplace below so products can be linked to the correct Shop.
  55. if (region === 'FE' || definition.singleMarketplace || definition.type === 'GET_MERCHANT_LISTINGS_ALL_DATA') continue;
  56. const recordKey = `${region}:${definition.type}:${dateKey}`;
  57. const previous = existingByKey.get(recordKey);
  58. if (previous?.reportId) {
  59. created.push(previous);
  60. log('report_reused', { region, reportType: definition.type, reportId: previous.reportId });
  61. continue;
  62. }
  63. const startedAt = new Date();
  64. try {
  65. const body = {
  66. reportType: definition.type,
  67. marketplaceIds: [...new Set(marketplaceIds)],
  68. ...(definition.dateRange ? { dataStartTime, dataEndTime } : {}),
  69. ...(definition.reportOptions ? { reportOptions: definition.reportOptions } : {}),
  70. };
  71. const result = await forward(shop.objectId, '/reports/2021-06-30/reports', 'POST', body, 5);
  72. if (!result.reportId) throw new Error('Create report returned no reportId');
  73. const row = {
  74. recordKey,
  75. region,
  76. reportId: result.reportId,
  77. reportType: definition.type,
  78. marketplaceIds: body.marketplaceIds,
  79. dataStartTime: definition.dateRange ? dateValue(dataStartTime) : undefined,
  80. dataEndTime: definition.dateRange ? dateValue(dataEndTime) : undefined,
  81. processingStatus: 'IN_QUEUE',
  82. isParsed: false,
  83. rawData: result,
  84. shop: shopPointer(shop.objectId),
  85. };
  86. await upsertMany('Reports', [row]);
  87. created.push(row);
  88. await writeAudit({
  89. dataType: 'report_create', endpoint: '/reports/2021-06-30/reports',
  90. status: 'completed', shop, region, marketplaceId: shop.marketplaceId,
  91. records: 1, pages: 1, startedAt, details: { reportType: definition.type, reportId: result.reportId },
  92. });
  93. log('report_created', { region, reportType: definition.type, reportId: result.reportId });
  94. } catch (error) {
  95. await writeAudit({
  96. dataType: 'report_create', endpoint: '/reports/2021-06-30/reports',
  97. status: 'failed', shop, region, marketplaceId: shop.marketplaceId,
  98. error: error.message, startedAt, details: { reportType: definition.type },
  99. });
  100. log('report_create_failed', { region, reportType: definition.type, error: error.message });
  101. }
  102. await sleep(2500);
  103. }
  104. }
  105. for (const shop of shops) {
  106. const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
  107. const fallbackDefinitions = reportDefinitions.filter(definition =>
  108. definition.singleMarketplace
  109. || region === 'FE'
  110. || definition.type === 'GET_MERCHANT_LISTINGS_ALL_DATA'
  111. || definition.type === 'GET_FBA_FULFILLMENT_INVENTORY_RECEIPTS_DATA'
  112. || definition.type === 'GET_FBA_FULFILLMENT_INVENTORY_ADJUSTMENTS_DATA'
  113. || (region === 'NA' && definition.type === 'GET_FBA_MYI_ALL_INVENTORY_DATA')
  114. );
  115. for (const definition of fallbackDefinitions) {
  116. const recordKey = `MP:${shop.marketplaceId}:${definition.type}:${dateKey}`;
  117. const previous = existingByKey.get(recordKey);
  118. if (previous?.reportId) {
  119. created.push(previous);
  120. log('marketplace_report_reused', { marketplaceId: shop.marketplaceId, reportType: definition.type, reportId: previous.reportId });
  121. continue;
  122. }
  123. const startedAt = new Date();
  124. try {
  125. const body = {
  126. reportType: definition.type,
  127. marketplaceIds: [shop.marketplaceId],
  128. ...(definition.dateRange ? { dataStartTime, dataEndTime } : {}),
  129. ...(definition.reportOptions ? { reportOptions: definition.reportOptions } : {}),
  130. };
  131. const result = await forward(shop.objectId, '/reports/2021-06-30/reports', 'POST', body, 5);
  132. if (!result.reportId) throw new Error('Create report returned no reportId');
  133. const row = {
  134. recordKey,
  135. region,
  136. reportId: result.reportId,
  137. reportType: definition.type,
  138. marketplaceIds: [shop.marketplaceId],
  139. dataStartTime: definition.dateRange ? dateValue(dataStartTime) : undefined,
  140. dataEndTime: definition.dateRange ? dateValue(dataEndTime) : undefined,
  141. processingStatus: 'IN_QUEUE',
  142. isParsed: false,
  143. rawData: result,
  144. shop: shopPointer(shop.objectId),
  145. };
  146. await upsertMany('Reports', [row]);
  147. created.push(row);
  148. await writeAudit({
  149. dataType: 'report_create_marketplace', endpoint: '/reports/2021-06-30/reports',
  150. status: 'completed', shop, region, marketplaceId: shop.marketplaceId,
  151. records: 1, pages: 1, startedAt, details: { reportType: definition.type, reportId: result.reportId },
  152. });
  153. log('marketplace_report_created', { marketplaceId: shop.marketplaceId, reportType: definition.type, reportId: result.reportId });
  154. } catch (error) {
  155. await writeAudit({
  156. dataType: 'report_create_marketplace', endpoint: '/reports/2021-06-30/reports',
  157. status: 'failed', shop, region, marketplaceId: shop.marketplaceId,
  158. error: error.message, startedAt, details: { reportType: definition.type },
  159. });
  160. log('marketplace_report_create_failed', { marketplaceId: shop.marketplaceId, reportType: definition.type, error: error.message });
  161. }
  162. await sleep(2500);
  163. }
  164. }
  165. function parseCsvLine(line, delimiter) {
  166. const values = [];
  167. let value = '';
  168. let quoted = false;
  169. for (let index = 0; index < line.length; index++) {
  170. const char = line[index];
  171. if (char === '"') {
  172. if (quoted && line[index + 1] === '"') {
  173. value += '"';
  174. index++;
  175. } else {
  176. quoted = !quoted;
  177. }
  178. } else if (char === delimiter && !quoted) {
  179. values.push(value);
  180. value = '';
  181. } else {
  182. value += char;
  183. }
  184. }
  185. values.push(value);
  186. return values;
  187. }
  188. function safeKey(value, index) {
  189. const key = String(value || `column_${index + 1}`).trim().replace(/[.$]/g, '_');
  190. return key || `column_${index + 1}`;
  191. }
  192. function parseDocument(content, reportType) {
  193. const trimmed = content.replace(/^\uFEFF/, '').trim();
  194. if (!trimmed) return [];
  195. if (trimmed.startsWith('{') || trimmed.startsWith('[')) {
  196. try {
  197. const parsed = JSON.parse(trimmed);
  198. if (Array.isArray(parsed)) return parsed;
  199. const rows = [];
  200. for (const [section, value] of Object.entries(parsed)) {
  201. if (Array.isArray(value)) rows.push(...value.map(item => ({ _section: section, ...(item && typeof item === 'object' ? item : { value: item }) })));
  202. }
  203. return rows.length ? rows : [parsed];
  204. } catch {}
  205. }
  206. const lines = trimmed.split(/\r?\n/).filter(Boolean);
  207. if (lines.length === 1) return [{ content: lines[0] }];
  208. const delimiter = lines[0].includes('\t') ? '\t' : ',';
  209. const headers = parseCsvLine(lines[0], delimiter).map(safeKey);
  210. return lines.slice(1).map(line => {
  211. const values = parseCsvLine(line, delimiter);
  212. return Object.fromEntries(headers.map((header, index) => [header, values[index] ?? '']));
  213. }).filter(row => Object.values(row).some(value => value !== ''));
  214. }
  215. async function processReport(report, shop, region) {
  216. const document = await forward(shop.objectId,
  217. `/reports/2021-06-30/documents/${encodeURIComponent(report.reportDocumentId)}`,
  218. 'GET', undefined, 5);
  219. if (!document.url) throw new Error('Report document returned no URL');
  220. const response = await fetch(document.url, { signal: AbortSignal.timeout(300000) });
  221. if (!response.ok) throw new Error(`Report download failed: ${response.status}`);
  222. let bytes = Buffer.from(await response.arrayBuffer());
  223. if (document.compressionAlgorithm === 'GZIP') bytes = zlib.gunzipSync(bytes);
  224. const content = new TextDecoder('utf-8').decode(bytes);
  225. const parsedRows = parseDocument(content, report.reportType);
  226. const rows = parsedRows.map(row => {
  227. const rowHash = hash(row);
  228. const marketplaceKey = [...(report.marketplaceIds || [])].sort().join(',') || report.shop?.objectId || region;
  229. return {
  230. recordKey: `${region}:${marketplaceKey}:${report.reportType}:${rowHash}`,
  231. rowHash,
  232. reportType: report.reportType,
  233. reportId: report.reportId,
  234. region,
  235. marketplaceIds: report.marketplaceIds || [],
  236. capturedAt: dateValue(new Date()),
  237. rowData: row,
  238. rawData: row,
  239. shop: shopPointer(shop.objectId),
  240. };
  241. });
  242. const saved = await upsertMany('AmazonReportRow', rows);
  243. await updateObject('Reports', report.objectId, {
  244. isParsed: true,
  245. parsedAt: dateValue(new Date()),
  246. parsedRecords: rows.length,
  247. });
  248. return { rows: rows.length, ...saved, bytes: bytes.length };
  249. }
  250. for (let round = 1; round <= 20; round++) {
  251. const reports = await allRows('Reports');
  252. const targetKeys = new Set(created.map(report => report.recordKey).filter(Boolean));
  253. const pending = reports.filter(report => targetKeys.has(report.recordKey) && report.isParsed !== true);
  254. let waiting = 0;
  255. log('report_poll_round', { round, pending: pending.length });
  256. for (const report of pending) {
  257. const shop = shops.find(item => item.objectId === report.shop?.objectId) || regionShops.get(report.region)?.shop;
  258. if (!shop) continue;
  259. try {
  260. let status = report.processingStatus || 'IN_QUEUE';
  261. let documentId = report.reportDocumentId || '';
  262. if (!['DONE', 'CANCELLED', 'FATAL'].includes(status)) {
  263. const info = await forward(shop.objectId, `/reports/2021-06-30/reports/${report.reportId}`, 'GET', undefined, 5);
  264. status = info.processingStatus || status;
  265. documentId = info.reportDocumentId || documentId;
  266. await updateObject('Reports', report.objectId, {
  267. processingStatus: status,
  268. ...(documentId ? { reportDocumentId: documentId } : {}),
  269. rawData: info,
  270. });
  271. }
  272. if (status === 'DONE' && documentId) {
  273. report.reportDocumentId = documentId;
  274. const result = await processReport(report, shop, report.region);
  275. log('report_processed', { region: report.region, reportType: report.reportType, reportId: report.reportId, ...result });
  276. } else if (!['CANCELLED', 'FATAL'].includes(status)) {
  277. waiting++;
  278. } else {
  279. log('report_terminal_without_data', { region: report.region, reportType: report.reportType, reportId: report.reportId, status });
  280. }
  281. } catch (error) {
  282. waiting++;
  283. log('report_process_failed', { region: report.region, reportType: report.reportType, reportId: report.reportId, error: error.message });
  284. }
  285. await sleep(1000);
  286. }
  287. if (waiting === 0) break;
  288. if (round < 20) await sleep(30000);
  289. }
  290. const finalReports = await allRows('Reports');
  291. const targetKeys = new Set(created.map(report => report.recordKey).filter(Boolean));
  292. const summary = finalReports.filter(report => targetKeys.has(report.recordKey)).reduce((result, report) => {
  293. const key = `${report.processingStatus || 'UNKNOWN'}|parsed=${Boolean(report.isParsed)}`;
  294. result[key] = (result[key] || 0) + 1;
  295. return result;
  296. }, {});
  297. log('report_collection_finished', { created: created.length, summary });