collect-sp-api-core.mjs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304
  1. import {
  2. activeAmazonShops, allRows, cleanObject, dateValue, forward, hash, log,
  3. orderPointer, regionByMarketplace, runId, shopPointer, sleep, updateObject,
  4. upsertMany, writeAudit,
  5. } from './sp-api-collector-lib.mjs';
  6. const now = new Date();
  7. const capturedAt = dateValue(now);
  8. const shops = await activeAmazonShops();
  9. const orders = await allRows('Order', 'objectId,platformOrderId,orderId,orderDate,itemsSyncedAt,shop');
  10. const listings = await allRows('Listing', 'objectId,sku,asin,marketplaceId,shop');
  11. const products = await allRows('Product', 'objectId,asin,marketplaceId,shop');
  12. const shopsByRegion = new Map();
  13. for (const shop of shops) {
  14. const region = regionByMarketplace[shop.marketplaceId] || shop.config?.SpApiConfig?.region || 'UNKNOWN';
  15. if (!shopsByRegion.has(region)) shopsByRegion.set(region, shop);
  16. }
  17. async function runTask({ dataType, endpoint, shop, execute }) {
  18. const startedAt = new Date();
  19. const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
  20. try {
  21. const result = await execute(region);
  22. await writeAudit({
  23. dataType, endpoint, status: 'completed', shop, region,
  24. marketplaceId: shop.marketplaceId, records: result.records || 0,
  25. pages: result.pages || 1, startedAt, details: result,
  26. });
  27. log('task_completed', { dataType, endpoint, shop: shop.name, region, ...result });
  28. return result;
  29. } catch (error) {
  30. await writeAudit({
  31. dataType, endpoint, status: 'failed', shop, region,
  32. marketplaceId: shop.marketplaceId, error: error.message, startedAt,
  33. });
  34. log('task_failed', { dataType, endpoint, shop: shop.name, region, error: error.message });
  35. return { records: 0, pages: 0, error: error.message };
  36. }
  37. }
  38. for (const shop of shops) {
  39. const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
  40. const accountResult = await runTask({
  41. dataType: 'seller_account', endpoint: '/sellers/v1/account', shop,
  42. execute: async () => {
  43. const result = await forward(shop.objectId, '/sellers/v1/account');
  44. const payload = result.payload || result;
  45. if (payload.errors?.length) throw new Error(payload.errors[0].message || payload.errors[0].code || 'Seller account query failed');
  46. const sellerId = payload.sellerId || shop.config?.SpApiConfig?.sellerID || '';
  47. await upsertMany('SpApiSellerAccount', [{
  48. recordKey: `${sellerId || shop.objectId}:${region}`,
  49. sellerId,
  50. marketplaceId: shop.marketplaceId,
  51. businessType: payload.businessType || '',
  52. region,
  53. capturedAt,
  54. rawData: payload,
  55. shop: shopPointer(shop.objectId),
  56. }]);
  57. return { records: 1, pages: 1, sellerId };
  58. },
  59. });
  60. await runTask({
  61. dataType: 'marketplace_participations', endpoint: '/sellers/v1/marketplaceParticipations', shop,
  62. execute: async () => {
  63. const result = await forward(shop.objectId, '/sellers/v1/marketplaceParticipations');
  64. const payload = result.payload || result;
  65. const participations = Array.isArray(payload) ? payload : (payload.marketplaceParticipations || []);
  66. const sellerId = accountResult.sellerId || shop.config?.SpApiConfig?.sellerID || `${region}:${shop.objectId}`;
  67. const rows = participations.map(item => ({
  68. recordKey: `${sellerId}:${item.marketplace?.id || ''}`,
  69. marketplaceId: item.marketplace?.id || '',
  70. name: item.marketplace?.name || '',
  71. countryCode: item.marketplace?.countryCode || '',
  72. region,
  73. participation: item.participation || {},
  74. capturedAt,
  75. rawData: item,
  76. shop: shopPointer(shop.objectId),
  77. }));
  78. await upsertMany('SpApiMarketplaceParticipation', rows);
  79. return { records: rows.length, pages: 1 };
  80. },
  81. });
  82. }
  83. for (const shop of shops) {
  84. await runTask({
  85. dataType: 'fba_inventory', endpoint: '/fba/inventory/v1/summaries', shop,
  86. execute: async () => {
  87. let token = '';
  88. let pages = 0;
  89. let records = 0;
  90. const seen = new Set();
  91. do {
  92. pages++;
  93. const path = token
  94. ? `/fba/inventory/v1/summaries?nextToken=${encodeURIComponent(token)}`
  95. : `/fba/inventory/v1/summaries?granularityType=Marketplace&granularityId=${encodeURIComponent(shop.marketplaceId)}&marketplaceIds=${encodeURIComponent(shop.marketplaceId)}&details=true`;
  96. const result = await forward(shop.objectId, path, 'GET', undefined, 2);
  97. const payload = result.payload || result;
  98. const rows = (payload.inventorySummaries || []).map(item => ({
  99. recordKey: `${shop.objectId}:${item.sellerSku || item.fnSku || item.asin || hash(item)}`,
  100. marketplaceId: shop.marketplaceId,
  101. asin: item.asin || '',
  102. sellerSku: item.sellerSku || '',
  103. fnSku: item.fnSku || '',
  104. condition: item.condition || '',
  105. productName: item.productName || '',
  106. totalQuantity: Number(item.totalQuantity || 0),
  107. inventoryDetails: item.inventoryDetails || {},
  108. lastUpdatedTime: dateValue(item.lastUpdatedTime),
  109. capturedAt,
  110. rawData: item,
  111. shop: shopPointer(shop.objectId),
  112. }));
  113. await upsertMany('AmazonInventorySummary', rows);
  114. records += rows.length;
  115. token = String(result.pagination?.nextToken || payload.nextToken || '').trim();
  116. if (token) {
  117. if (seen.has(token)) throw new Error('Repeated inventory nextToken');
  118. seen.add(token);
  119. await sleep(1200);
  120. }
  121. } while (token);
  122. return { records, pages };
  123. },
  124. });
  125. await runTask({
  126. dataType: 'sales_metrics', endpoint: '/sales/v1/orderMetrics', shop,
  127. execute: async () => {
  128. const start = new Date(now.getTime() - 729 * 24 * 60 * 60 * 1000).toISOString();
  129. const path = `/sales/v1/orderMetrics?marketplaceIds=${encodeURIComponent(shop.marketplaceId)}`
  130. + `&interval=${encodeURIComponent(`${start}--${now.toISOString()}`)}&granularity=Day`;
  131. const result = await forward(shop.objectId, path, 'GET', undefined, 2);
  132. const metrics = result.payload || result;
  133. const rows = (Array.isArray(metrics) ? metrics : []).map(metric => ({
  134. recordKey: `${shop.objectId}:${metric.interval || hash(metric)}`,
  135. marketplaceId: shop.marketplaceId,
  136. interval: metric.interval || '',
  137. granularity: 'Day',
  138. unitCount: Number(metric.unitCount || 0),
  139. orderItemCount: Number(metric.orderItemCount || 0),
  140. orderCount: Number(metric.orderCount || 0),
  141. averageUnitPrice: metric.averageUnitPrice || {},
  142. totalSales: metric.totalSales || {},
  143. capturedAt,
  144. rawData: metric,
  145. shop: shopPointer(shop.objectId),
  146. }));
  147. await upsertMany('AmazonSalesMetric', rows);
  148. return { records: rows.length, pages: 1 };
  149. },
  150. });
  151. }
  152. const catalogTargets = new Map();
  153. for (const row of [...listings, ...products]) {
  154. const asin = String(row.asin || '').trim().toUpperCase();
  155. const shop = shops.find(item => item.objectId === row.shop?.objectId);
  156. if (!asin || !shop) continue;
  157. catalogTargets.set(`${shop.objectId}:${asin}`, { asin, shop });
  158. }
  159. for (const { asin, shop } of catalogTargets.values()) {
  160. await runTask({
  161. dataType: 'catalog_item', endpoint: '/catalog/2022-04-01/items/{asin}', shop,
  162. execute: async () => {
  163. const included = 'attributes,dimensions,identifiers,images,productTypes,relationships,salesRanks,summaries';
  164. const result = await forward(shop.objectId,
  165. `/catalog/2022-04-01/items/${encodeURIComponent(asin)}?marketplaceIds=${encodeURIComponent(shop.marketplaceId)}&includedData=${included}`,
  166. 'GET', undefined, 2);
  167. const summary = result.summaries?.[0] || {};
  168. const row = {
  169. recordKey: `${shop.objectId}:${asin}`,
  170. marketplaceId: shop.marketplaceId,
  171. asin,
  172. title: summary.itemName || '',
  173. brand: summary.brand || '',
  174. productType: result.productTypes?.[0]?.productType || '',
  175. attributes: result.attributes || {},
  176. dimensions: result.dimensions || [],
  177. identifiers: result.identifiers || [],
  178. images: result.images || [],
  179. relationships: result.relationships || [],
  180. salesRanks: result.salesRanks || [],
  181. summaries: result.summaries || [],
  182. capturedAt,
  183. rawData: result,
  184. shop: shopPointer(shop.objectId),
  185. };
  186. await upsertMany('AmazonCatalogItem', [row]);
  187. return { records: 1, pages: 1 };
  188. },
  189. });
  190. await sleep(500);
  191. }
  192. for (const shop of shops) {
  193. const asins = [...new Set(listings
  194. .filter(item => item.shop?.objectId === shop.objectId && item.asin)
  195. .map(item => String(item.asin).trim().toUpperCase()))];
  196. for (let start = 0; start < asins.length; start += 20) {
  197. const batch = asins.slice(start, start + 20);
  198. await runTask({
  199. dataType: 'product_pricing', endpoint: '/products/pricing/v0/price', shop,
  200. execute: async () => {
  201. const result = await forward(shop.objectId,
  202. `/products/pricing/v0/price?MarketplaceId=${encodeURIComponent(shop.marketplaceId)}&ItemType=Asin&Asins=${batch.join(',')}`,
  203. 'GET', undefined, 2);
  204. const payload = result.payload || result;
  205. const rows = (Array.isArray(payload) ? payload : []).map(item => ({
  206. recordKey: `${shop.objectId}:${item.ASIN || item.Sku || hash(item)}`,
  207. marketplaceId: shop.marketplaceId,
  208. asin: item.ASIN || '',
  209. sku: item.Sku || '',
  210. status: item.status || '',
  211. product: item.Product || {},
  212. capturedAt,
  213. rawData: item,
  214. shop: shopPointer(shop.objectId),
  215. }));
  216. await upsertMany('AmazonProductPricing', rows);
  217. return { records: rows.length, pages: 1, batch: start / 20 + 1 };
  218. },
  219. });
  220. await sleep(1200);
  221. }
  222. }
  223. const ordersByShop = new Map();
  224. for (const order of orders) {
  225. if (!order.shop?.objectId || !order.platformOrderId) continue;
  226. if (!ordersByShop.has(order.shop.objectId)) ordersByShop.set(order.shop.objectId, []);
  227. ordersByShop.get(order.shop.objectId).push(order);
  228. }
  229. for (const shop of shops) {
  230. const shopOrders = (ordersByShop.get(shop.objectId) || []).filter(order => !order.itemsSyncedAt);
  231. let processed = 0;
  232. let failures = 0;
  233. const startedAt = new Date();
  234. for (const order of shopOrders) {
  235. try {
  236. let token = '';
  237. const allItems = [];
  238. do {
  239. const path = `/orders/v0/orders/${encodeURIComponent(order.platformOrderId)}/orderItems`
  240. + (token ? `?NextToken=${encodeURIComponent(token)}` : '');
  241. const result = await forward(shop.objectId, path, 'GET', undefined, 3);
  242. const payload = result.payload || result;
  243. allItems.push(...(payload.OrderItems || []));
  244. token = String(payload.NextToken || '').trim();
  245. if (token) await sleep(500);
  246. } while (token);
  247. const rows = allItems.map(item => ({
  248. recordKey: `${shop.objectId}:${item.OrderItemId}`,
  249. marketplaceId: shop.marketplaceId,
  250. amazonOrderId: order.platformOrderId,
  251. orderItemId: item.OrderItemId || '',
  252. asin: item.ASIN || '',
  253. sellerSku: item.SellerSKU || '',
  254. title: item.Title || '',
  255. quantityOrdered: Number(item.QuantityOrdered || 0),
  256. quantityShipped: Number(item.QuantityShipped || 0),
  257. itemPrice: item.ItemPrice || {},
  258. itemTax: item.ItemTax || {},
  259. promotionDiscount: item.PromotionDiscount || {},
  260. promotionIds: item.PromotionIds || [],
  261. capturedAt,
  262. rawData: item,
  263. shop: shopPointer(shop.objectId),
  264. order: orderPointer(order.objectId),
  265. }));
  266. await upsertMany('AmazonOrderItem', rows);
  267. await updateObject('Order', order.objectId, { itemCount: rows.length, itemsSyncedAt: capturedAt });
  268. processed += rows.length;
  269. } catch (error) {
  270. failures++;
  271. log('order_items_failed', { shop: shop.name, orderId: order.platformOrderId, error: error.message });
  272. }
  273. if ((processed + failures) % 50 === 0) log('order_items_progress', { shop: shop.name, processed, failures });
  274. await sleep(650);
  275. }
  276. await writeAudit({
  277. dataType: 'order_items', endpoint: '/orders/v0/orders/{orderId}/orderItems',
  278. status: failures ? 'partial' : 'completed', shop,
  279. region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN', marketplaceId: shop.marketplaceId,
  280. records: processed, pages: shopOrders.length, error: failures ? `${failures} orders failed` : '',
  281. startedAt, details: { orders: shopOrders.length, failures },
  282. });
  283. log('order_items_completed', { shop: shop.name, orders: shopOrders.length, records: processed, failures });
  284. }
  285. log('core_collection_finished', {
  286. runId,
  287. shops: shops.length,
  288. regions: [...shopsByRegion.keys()],
  289. catalogTargets: catalogTargets.size,
  290. orders: orders.length,
  291. });