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 });