import { activeAmazonShops, allRows, cleanObject, dateValue, forward, hash, log, orderPointer, regionByMarketplace, runId, shopPointer, sleep, updateObject, upsertMany, writeAudit, } from './sp-api-collector-lib.mjs'; const now = new Date(); const capturedAt = dateValue(now); const shops = await activeAmazonShops(); const orders = await allRows('Order', 'objectId,platformOrderId,orderId,orderDate,itemsSyncedAt,shop'); const listings = await allRows('Listing', 'objectId,sku,asin,marketplaceId,shop'); const products = await allRows('Product', 'objectId,asin,marketplaceId,shop'); const shopsByRegion = new Map(); for (const shop of shops) { const region = regionByMarketplace[shop.marketplaceId] || shop.config?.SpApiConfig?.region || 'UNKNOWN'; if (!shopsByRegion.has(region)) shopsByRegion.set(region, shop); } async function runTask({ dataType, endpoint, shop, execute }) { const startedAt = new Date(); const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN'; try { const result = await execute(region); await writeAudit({ dataType, endpoint, status: 'completed', shop, region, marketplaceId: shop.marketplaceId, records: result.records || 0, pages: result.pages || 1, startedAt, details: result, }); log('task_completed', { dataType, endpoint, shop: shop.name, region, ...result }); return result; } catch (error) { await writeAudit({ dataType, endpoint, status: 'failed', shop, region, marketplaceId: shop.marketplaceId, error: error.message, startedAt, }); log('task_failed', { dataType, endpoint, shop: shop.name, region, error: error.message }); return { records: 0, pages: 0, error: error.message }; } } for (const shop of shops) { const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN'; const accountResult = await runTask({ dataType: 'seller_account', endpoint: '/sellers/v1/account', shop, execute: async () => { const result = await forward(shop.objectId, '/sellers/v1/account'); const payload = result.payload || result; if (payload.errors?.length) throw new Error(payload.errors[0].message || payload.errors[0].code || 'Seller account query failed'); const sellerId = payload.sellerId || shop.config?.SpApiConfig?.sellerID || ''; await upsertMany('SpApiSellerAccount', [{ recordKey: `${sellerId || shop.objectId}:${region}`, sellerId, marketplaceId: shop.marketplaceId, businessType: payload.businessType || '', region, capturedAt, rawData: payload, shop: shopPointer(shop.objectId), }]); return { records: 1, pages: 1, sellerId }; }, }); await runTask({ dataType: 'marketplace_participations', endpoint: '/sellers/v1/marketplaceParticipations', shop, execute: async () => { const result = await forward(shop.objectId, '/sellers/v1/marketplaceParticipations'); const payload = result.payload || result; const participations = Array.isArray(payload) ? payload : (payload.marketplaceParticipations || []); const sellerId = accountResult.sellerId || shop.config?.SpApiConfig?.sellerID || `${region}:${shop.objectId}`; const rows = participations.map(item => ({ recordKey: `${sellerId}:${item.marketplace?.id || ''}`, marketplaceId: item.marketplace?.id || '', name: item.marketplace?.name || '', countryCode: item.marketplace?.countryCode || '', region, participation: item.participation || {}, capturedAt, rawData: item, shop: shopPointer(shop.objectId), })); await upsertMany('SpApiMarketplaceParticipation', rows); return { records: rows.length, pages: 1 }; }, }); } for (const shop of shops) { await runTask({ dataType: 'fba_inventory', endpoint: '/fba/inventory/v1/summaries', shop, execute: async () => { let token = ''; let pages = 0; let records = 0; const seen = new Set(); do { pages++; const path = token ? `/fba/inventory/v1/summaries?nextToken=${encodeURIComponent(token)}` : `/fba/inventory/v1/summaries?granularityType=Marketplace&granularityId=${encodeURIComponent(shop.marketplaceId)}&marketplaceIds=${encodeURIComponent(shop.marketplaceId)}&details=true`; const result = await forward(shop.objectId, path, 'GET', undefined, 2); const payload = result.payload || result; const rows = (payload.inventorySummaries || []).map(item => ({ recordKey: `${shop.objectId}:${item.sellerSku || item.fnSku || item.asin || hash(item)}`, marketplaceId: shop.marketplaceId, asin: item.asin || '', sellerSku: item.sellerSku || '', fnSku: item.fnSku || '', condition: item.condition || '', productName: item.productName || '', totalQuantity: Number(item.totalQuantity || 0), inventoryDetails: item.inventoryDetails || {}, lastUpdatedTime: dateValue(item.lastUpdatedTime), capturedAt, rawData: item, shop: shopPointer(shop.objectId), })); await upsertMany('AmazonInventorySummary', rows); records += rows.length; token = String(result.pagination?.nextToken || payload.nextToken || '').trim(); if (token) { if (seen.has(token)) throw new Error('Repeated inventory nextToken'); seen.add(token); await sleep(1200); } } while (token); return { records, pages }; }, }); await runTask({ dataType: 'sales_metrics', endpoint: '/sales/v1/orderMetrics', shop, execute: async () => { const start = new Date(now.getTime() - 729 * 24 * 60 * 60 * 1000).toISOString(); const path = `/sales/v1/orderMetrics?marketplaceIds=${encodeURIComponent(shop.marketplaceId)}` + `&interval=${encodeURIComponent(`${start}--${now.toISOString()}`)}&granularity=Day`; const result = await forward(shop.objectId, path, 'GET', undefined, 2); const metrics = result.payload || result; const rows = (Array.isArray(metrics) ? metrics : []).map(metric => ({ recordKey: `${shop.objectId}:${metric.interval || hash(metric)}`, marketplaceId: shop.marketplaceId, interval: metric.interval || '', granularity: 'Day', unitCount: Number(metric.unitCount || 0), orderItemCount: Number(metric.orderItemCount || 0), orderCount: Number(metric.orderCount || 0), averageUnitPrice: metric.averageUnitPrice || {}, totalSales: metric.totalSales || {}, capturedAt, rawData: metric, shop: shopPointer(shop.objectId), })); await upsertMany('AmazonSalesMetric', rows); return { records: rows.length, pages: 1 }; }, }); } const catalogTargets = new Map(); for (const row of [...listings, ...products]) { const asin = String(row.asin || '').trim().toUpperCase(); const shop = shops.find(item => item.objectId === row.shop?.objectId); if (!asin || !shop) continue; catalogTargets.set(`${shop.objectId}:${asin}`, { asin, shop }); } for (const { asin, shop } of catalogTargets.values()) { await runTask({ dataType: 'catalog_item', endpoint: '/catalog/2022-04-01/items/{asin}', shop, execute: async () => { const included = 'attributes,dimensions,identifiers,images,productTypes,relationships,salesRanks,summaries'; const result = await forward(shop.objectId, `/catalog/2022-04-01/items/${encodeURIComponent(asin)}?marketplaceIds=${encodeURIComponent(shop.marketplaceId)}&includedData=${included}`, 'GET', undefined, 2); const summary = result.summaries?.[0] || {}; const row = { recordKey: `${shop.objectId}:${asin}`, marketplaceId: shop.marketplaceId, asin, title: summary.itemName || '', brand: summary.brand || '', productType: result.productTypes?.[0]?.productType || '', attributes: result.attributes || {}, dimensions: result.dimensions || [], identifiers: result.identifiers || [], images: result.images || [], relationships: result.relationships || [], salesRanks: result.salesRanks || [], summaries: result.summaries || [], capturedAt, rawData: result, shop: shopPointer(shop.objectId), }; await upsertMany('AmazonCatalogItem', [row]); return { records: 1, pages: 1 }; }, }); await sleep(500); } for (const shop of shops) { const asins = [...new Set(listings .filter(item => item.shop?.objectId === shop.objectId && item.asin) .map(item => String(item.asin).trim().toUpperCase()))]; for (let start = 0; start < asins.length; start += 20) { const batch = asins.slice(start, start + 20); await runTask({ dataType: 'product_pricing', endpoint: '/products/pricing/v0/price', shop, execute: async () => { const result = await forward(shop.objectId, `/products/pricing/v0/price?MarketplaceId=${encodeURIComponent(shop.marketplaceId)}&ItemType=Asin&Asins=${batch.join(',')}`, 'GET', undefined, 2); const payload = result.payload || result; const rows = (Array.isArray(payload) ? payload : []).map(item => ({ recordKey: `${shop.objectId}:${item.ASIN || item.Sku || hash(item)}`, marketplaceId: shop.marketplaceId, asin: item.ASIN || '', sku: item.Sku || '', status: item.status || '', product: item.Product || {}, capturedAt, rawData: item, shop: shopPointer(shop.objectId), })); await upsertMany('AmazonProductPricing', rows); return { records: rows.length, pages: 1, batch: start / 20 + 1 }; }, }); await sleep(1200); } } const ordersByShop = new Map(); for (const order of orders) { if (!order.shop?.objectId || !order.platformOrderId) continue; if (!ordersByShop.has(order.shop.objectId)) ordersByShop.set(order.shop.objectId, []); ordersByShop.get(order.shop.objectId).push(order); } for (const shop of shops) { const shopOrders = (ordersByShop.get(shop.objectId) || []).filter(order => !order.itemsSyncedAt); let processed = 0; let failures = 0; const startedAt = new Date(); for (const order of shopOrders) { try { let token = ''; const allItems = []; do { const path = `/orders/v0/orders/${encodeURIComponent(order.platformOrderId)}/orderItems` + (token ? `?NextToken=${encodeURIComponent(token)}` : ''); const result = await forward(shop.objectId, path, 'GET', undefined, 3); const payload = result.payload || result; allItems.push(...(payload.OrderItems || [])); token = String(payload.NextToken || '').trim(); if (token) await sleep(500); } while (token); const rows = allItems.map(item => ({ recordKey: `${shop.objectId}:${item.OrderItemId}`, marketplaceId: shop.marketplaceId, amazonOrderId: order.platformOrderId, orderItemId: item.OrderItemId || '', asin: item.ASIN || '', sellerSku: item.SellerSKU || '', title: item.Title || '', quantityOrdered: Number(item.QuantityOrdered || 0), quantityShipped: Number(item.QuantityShipped || 0), itemPrice: item.ItemPrice || {}, itemTax: item.ItemTax || {}, promotionDiscount: item.PromotionDiscount || {}, promotionIds: item.PromotionIds || [], capturedAt, rawData: item, shop: shopPointer(shop.objectId), order: orderPointer(order.objectId), })); await upsertMany('AmazonOrderItem', rows); await updateObject('Order', order.objectId, { itemCount: rows.length, itemsSyncedAt: capturedAt }); processed += rows.length; } catch (error) { failures++; log('order_items_failed', { shop: shop.name, orderId: order.platformOrderId, error: error.message }); } if ((processed + failures) % 50 === 0) log('order_items_progress', { shop: shop.name, processed, failures }); await sleep(650); } await writeAudit({ dataType: 'order_items', endpoint: '/orders/v0/orders/{orderId}/orderItems', status: failures ? 'partial' : 'completed', shop, region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN', marketplaceId: shop.marketplaceId, records: processed, pages: shopOrders.length, error: failures ? `${failures} orders failed` : '', startedAt, details: { orders: shopOrders.length, failures }, }); log('order_items_completed', { shop: shop.name, orders: shopOrders.length, records: processed, failures }); } log('core_collection_finished', { runId, shops: shops.length, regions: [...shopsByRegion.keys()], catalogTargets: catalogTargets.size, orders: orders.length, });