| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310 |
- import zlib from 'node:zlib';
- import {
- activeAmazonShops, allRows, dateValue, forward, hash, log, regionByMarketplace,
- shopPointer, sleep, updateObject, upsertMany, writeAudit,
- } from './sp-api-collector-lib.mjs';
- const allReportDefinitions = [
- { type: 'GET_MERCHANT_LISTINGS_ALL_DATA' },
- { type: 'GET_FBA_MYI_UNSUPPRESSED_INVENTORY_DATA' },
- { type: 'GET_FBA_MYI_ALL_INVENTORY_DATA' },
- { type: 'GET_FBA_INVENTORY_AGED_DATA' },
- { type: 'GET_FBA_INVENTORY_PLANNING_DATA' },
- { type: 'GET_RESERVED_INVENTORY_DATA' },
- { type: 'GET_FBA_ESTIMATED_FBA_FEES_TXT_DATA' },
- { type: 'GET_FBA_FULFILLMENT_CUSTOMER_RETURNS_DATA', dateRange: true },
- { type: 'GET_FBA_REIMBURSEMENTS_DATA', dateRange: true },
- { type: 'GET_AMAZON_FULFILLED_SHIPMENTS_DATA_GENERAL', dateRange: true },
- { type: 'GET_FBA_FULFILLMENT_INVENTORY_RECEIPTS_DATA', dateRange: true },
- { type: 'GET_FBA_FULFILLMENT_INVENTORY_ADJUSTMENTS_DATA', dateRange: true },
- { type: 'GET_FBA_FULFILLMENT_REMOVAL_ORDER_DETAIL_DATA', dateRange: true },
- { type: 'GET_FBA_FULFILLMENT_REMOVAL_SHIPMENT_DETAIL_DATA', dateRange: true },
- {
- type: 'GET_SALES_AND_TRAFFIC_REPORT',
- dateRange: true,
- singleMarketplace: true,
- reportOptions: { dateGranularity: 'DAY', asinGranularity: 'SKU' },
- },
- ];
- const requestedReportTypes = new Set(
- String(process.env.SP_API_REPORT_TYPES || '')
- .split(',')
- .map(value => value.trim())
- .filter(Boolean),
- );
- const reportDefinitions = requestedReportTypes.size
- ? allReportDefinitions.filter(definition => requestedReportTypes.has(definition.type))
- : allReportDefinitions;
- const shops = await activeAmazonShops();
- const regionShops = new Map();
- for (const shop of shops) {
- const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
- if (!regionShops.has(region)) regionShops.set(region, { shop, marketplaceIds: [] });
- regionShops.get(region).marketplaceIds.push(shop.marketplaceId);
- }
- const dateKey = new Date().toISOString().slice(0, 10);
- const dataEndTime = new Date().toISOString();
- const dataStartTime = new Date(Date.now() - 89 * 24 * 60 * 60 * 1000).toISOString();
- const existingReports = await allRows('Reports');
- const existingByKey = new Map(existingReports.filter(report => report.recordKey).map(report => [report.recordKey, report]));
- const created = [];
- for (const [region, { shop, marketplaceIds }] of regionShops) {
- for (const definition of reportDefinitions) {
- // Merchant listings do not identify the originating marketplace when a
- // regional request contains multiple marketplace IDs. Collect them per
- // marketplace below so products can be linked to the correct Shop.
- if (region === 'FE' || definition.singleMarketplace || definition.type === 'GET_MERCHANT_LISTINGS_ALL_DATA') continue;
- const recordKey = `${region}:${definition.type}:${dateKey}`;
- const previous = existingByKey.get(recordKey);
- if (previous?.reportId) {
- created.push(previous);
- log('report_reused', { region, reportType: definition.type, reportId: previous.reportId });
- continue;
- }
- const startedAt = new Date();
- try {
- const body = {
- reportType: definition.type,
- marketplaceIds: [...new Set(marketplaceIds)],
- ...(definition.dateRange ? { dataStartTime, dataEndTime } : {}),
- ...(definition.reportOptions ? { reportOptions: definition.reportOptions } : {}),
- };
- const result = await forward(shop.objectId, '/reports/2021-06-30/reports', 'POST', body, 5);
- if (!result.reportId) throw new Error('Create report returned no reportId');
- const row = {
- recordKey,
- region,
- reportId: result.reportId,
- reportType: definition.type,
- marketplaceIds: body.marketplaceIds,
- dataStartTime: definition.dateRange ? dateValue(dataStartTime) : undefined,
- dataEndTime: definition.dateRange ? dateValue(dataEndTime) : undefined,
- processingStatus: 'IN_QUEUE',
- isParsed: false,
- rawData: result,
- shop: shopPointer(shop.objectId),
- };
- await upsertMany('Reports', [row]);
- created.push(row);
- await writeAudit({
- dataType: 'report_create', endpoint: '/reports/2021-06-30/reports',
- status: 'completed', shop, region, marketplaceId: shop.marketplaceId,
- records: 1, pages: 1, startedAt, details: { reportType: definition.type, reportId: result.reportId },
- });
- log('report_created', { region, reportType: definition.type, reportId: result.reportId });
- } catch (error) {
- await writeAudit({
- dataType: 'report_create', endpoint: '/reports/2021-06-30/reports',
- status: 'failed', shop, region, marketplaceId: shop.marketplaceId,
- error: error.message, startedAt, details: { reportType: definition.type },
- });
- log('report_create_failed', { region, reportType: definition.type, error: error.message });
- }
- await sleep(2500);
- }
- }
- for (const shop of shops) {
- const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
- const fallbackDefinitions = reportDefinitions.filter(definition =>
- definition.singleMarketplace
- || region === 'FE'
- || definition.type === 'GET_MERCHANT_LISTINGS_ALL_DATA'
- || definition.type === 'GET_FBA_FULFILLMENT_INVENTORY_RECEIPTS_DATA'
- || definition.type === 'GET_FBA_FULFILLMENT_INVENTORY_ADJUSTMENTS_DATA'
- || (region === 'NA' && definition.type === 'GET_FBA_MYI_ALL_INVENTORY_DATA')
- );
- for (const definition of fallbackDefinitions) {
- const recordKey = `MP:${shop.marketplaceId}:${definition.type}:${dateKey}`;
- const previous = existingByKey.get(recordKey);
- if (previous?.reportId) {
- created.push(previous);
- log('marketplace_report_reused', { marketplaceId: shop.marketplaceId, reportType: definition.type, reportId: previous.reportId });
- continue;
- }
- const startedAt = new Date();
- try {
- const body = {
- reportType: definition.type,
- marketplaceIds: [shop.marketplaceId],
- ...(definition.dateRange ? { dataStartTime, dataEndTime } : {}),
- ...(definition.reportOptions ? { reportOptions: definition.reportOptions } : {}),
- };
- const result = await forward(shop.objectId, '/reports/2021-06-30/reports', 'POST', body, 5);
- if (!result.reportId) throw new Error('Create report returned no reportId');
- const row = {
- recordKey,
- region,
- reportId: result.reportId,
- reportType: definition.type,
- marketplaceIds: [shop.marketplaceId],
- dataStartTime: definition.dateRange ? dateValue(dataStartTime) : undefined,
- dataEndTime: definition.dateRange ? dateValue(dataEndTime) : undefined,
- processingStatus: 'IN_QUEUE',
- isParsed: false,
- rawData: result,
- shop: shopPointer(shop.objectId),
- };
- await upsertMany('Reports', [row]);
- created.push(row);
- await writeAudit({
- dataType: 'report_create_marketplace', endpoint: '/reports/2021-06-30/reports',
- status: 'completed', shop, region, marketplaceId: shop.marketplaceId,
- records: 1, pages: 1, startedAt, details: { reportType: definition.type, reportId: result.reportId },
- });
- log('marketplace_report_created', { marketplaceId: shop.marketplaceId, reportType: definition.type, reportId: result.reportId });
- } catch (error) {
- await writeAudit({
- dataType: 'report_create_marketplace', endpoint: '/reports/2021-06-30/reports',
- status: 'failed', shop, region, marketplaceId: shop.marketplaceId,
- error: error.message, startedAt, details: { reportType: definition.type },
- });
- log('marketplace_report_create_failed', { marketplaceId: shop.marketplaceId, reportType: definition.type, error: error.message });
- }
- await sleep(2500);
- }
- }
- function parseCsvLine(line, delimiter) {
- const values = [];
- let value = '';
- let quoted = false;
- for (let index = 0; index < line.length; index++) {
- const char = line[index];
- if (char === '"') {
- if (quoted && line[index + 1] === '"') {
- value += '"';
- index++;
- } else {
- quoted = !quoted;
- }
- } else if (char === delimiter && !quoted) {
- values.push(value);
- value = '';
- } else {
- value += char;
- }
- }
- values.push(value);
- return values;
- }
- function safeKey(value, index) {
- const key = String(value || `column_${index + 1}`).trim().replace(/[.$]/g, '_');
- return key || `column_${index + 1}`;
- }
- function parseDocument(content, reportType) {
- const trimmed = content.replace(/^\uFEFF/, '').trim();
- if (!trimmed) return [];
- if (trimmed.startsWith('{') || trimmed.startsWith('[')) {
- try {
- const parsed = JSON.parse(trimmed);
- if (Array.isArray(parsed)) return parsed;
- const rows = [];
- for (const [section, value] of Object.entries(parsed)) {
- if (Array.isArray(value)) rows.push(...value.map(item => ({ _section: section, ...(item && typeof item === 'object' ? item : { value: item }) })));
- }
- return rows.length ? rows : [parsed];
- } catch {}
- }
- const lines = trimmed.split(/\r?\n/).filter(Boolean);
- if (lines.length === 1) return [{ content: lines[0] }];
- const delimiter = lines[0].includes('\t') ? '\t' : ',';
- const headers = parseCsvLine(lines[0], delimiter).map(safeKey);
- return lines.slice(1).map(line => {
- const values = parseCsvLine(line, delimiter);
- return Object.fromEntries(headers.map((header, index) => [header, values[index] ?? '']));
- }).filter(row => Object.values(row).some(value => value !== ''));
- }
- async function processReport(report, shop, region) {
- const document = await forward(shop.objectId,
- `/reports/2021-06-30/documents/${encodeURIComponent(report.reportDocumentId)}`,
- 'GET', undefined, 5);
- if (!document.url) throw new Error('Report document returned no URL');
- const response = await fetch(document.url, { signal: AbortSignal.timeout(300000) });
- if (!response.ok) throw new Error(`Report download failed: ${response.status}`);
- let bytes = Buffer.from(await response.arrayBuffer());
- if (document.compressionAlgorithm === 'GZIP') bytes = zlib.gunzipSync(bytes);
- const content = new TextDecoder('utf-8').decode(bytes);
- const parsedRows = parseDocument(content, report.reportType);
- const rows = parsedRows.map(row => {
- const rowHash = hash(row);
- const marketplaceKey = [...(report.marketplaceIds || [])].sort().join(',') || report.shop?.objectId || region;
- return {
- recordKey: `${region}:${marketplaceKey}:${report.reportType}:${rowHash}`,
- rowHash,
- reportType: report.reportType,
- reportId: report.reportId,
- region,
- marketplaceIds: report.marketplaceIds || [],
- capturedAt: dateValue(new Date()),
- rowData: row,
- rawData: row,
- shop: shopPointer(shop.objectId),
- };
- });
- const saved = await upsertMany('AmazonReportRow', rows);
- await updateObject('Reports', report.objectId, {
- isParsed: true,
- parsedAt: dateValue(new Date()),
- parsedRecords: rows.length,
- });
- return { rows: rows.length, ...saved, bytes: bytes.length };
- }
- for (let round = 1; round <= 20; round++) {
- const reports = await allRows('Reports');
- const targetKeys = new Set(created.map(report => report.recordKey).filter(Boolean));
- const pending = reports.filter(report => targetKeys.has(report.recordKey) && report.isParsed !== true);
- let waiting = 0;
- log('report_poll_round', { round, pending: pending.length });
- for (const report of pending) {
- const shop = shops.find(item => item.objectId === report.shop?.objectId) || regionShops.get(report.region)?.shop;
- if (!shop) continue;
- try {
- let status = report.processingStatus || 'IN_QUEUE';
- let documentId = report.reportDocumentId || '';
- if (!['DONE', 'CANCELLED', 'FATAL'].includes(status)) {
- const info = await forward(shop.objectId, `/reports/2021-06-30/reports/${report.reportId}`, 'GET', undefined, 5);
- status = info.processingStatus || status;
- documentId = info.reportDocumentId || documentId;
- await updateObject('Reports', report.objectId, {
- processingStatus: status,
- ...(documentId ? { reportDocumentId: documentId } : {}),
- rawData: info,
- });
- }
- if (status === 'DONE' && documentId) {
- report.reportDocumentId = documentId;
- const result = await processReport(report, shop, report.region);
- log('report_processed', { region: report.region, reportType: report.reportType, reportId: report.reportId, ...result });
- } else if (!['CANCELLED', 'FATAL'].includes(status)) {
- waiting++;
- } else {
- log('report_terminal_without_data', { region: report.region, reportType: report.reportType, reportId: report.reportId, status });
- }
- } catch (error) {
- waiting++;
- log('report_process_failed', { region: report.region, reportType: report.reportType, reportId: report.reportId, error: error.message });
- }
- await sleep(1000);
- }
- if (waiting === 0) break;
- if (round < 20) await sleep(30000);
- }
- const finalReports = await allRows('Reports');
- const targetKeys = new Set(created.map(report => report.recordKey).filter(Boolean));
- const summary = finalReports.filter(report => targetKeys.has(report.recordKey)).reduce((result, report) => {
- const key = `${report.processingStatus || 'UNKNOWN'}|parsed=${Boolean(report.isParsed)}`;
- result[key] = (result[key] || 0) + 1;
- return result;
- }, {});
- log('report_collection_finished', { created: created.length, summary });
|