collect-sp-api-feedback-inbound.mjs 5.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118
  1. import {
  2. activeAmazonShops, allRows, dateValue, forward, hash, log, regionByMarketplace,
  3. shopPointer, sleep, upsertMany, writeAudit,
  4. } from './sp-api-collector-lib.mjs';
  5. const shops = await activeAmazonShops();
  6. const listings = await allRows('Listing', 'asin,shop');
  7. const capturedAt = dateValue(new Date());
  8. for (const shop of shops) {
  9. const asins = [...new Set(listings
  10. .filter(item => item.shop?.objectId === shop.objectId && item.asin)
  11. .map(item => String(item.asin).trim().toUpperCase()))];
  12. let records = 0;
  13. let emptyResponses = 0;
  14. let failures = 0;
  15. const startedAt = new Date();
  16. for (const asin of asins) {
  17. const endpoints = [
  18. ['item_browse_node', `/customerFeedback/2024-06-01/items/${asin}/browseNode?marketplaceId=${shop.marketplaceId}`],
  19. ['item_review_topics', `/customerFeedback/2024-06-01/items/${asin}/reviews/topics?marketplaceId=${shop.marketplaceId}&sortBy=MENTIONS&sortOrder=DESC`],
  20. ['item_review_trends', `/customerFeedback/2024-06-01/items/${asin}/reviews/trends?marketplaceId=${shop.marketplaceId}`],
  21. ];
  22. for (const [feedbackType, endpoint] of endpoints) {
  23. try {
  24. const result = await forward(shop.objectId, endpoint, 'GET', undefined, 2);
  25. if (!result || Object.keys(result).length === 0) {
  26. emptyResponses++;
  27. continue;
  28. }
  29. const browseNodeId = String(result.browseNodeId || result.browseNode?.id || '');
  30. await upsertMany('AmazonCustomerFeedback', [{
  31. recordKey: `${shop.objectId}:${asin}:${feedbackType}`,
  32. marketplaceId: shop.marketplaceId,
  33. feedbackType,
  34. asin,
  35. browseNodeId,
  36. capturedAt,
  37. rawData: result,
  38. shop: shopPointer(shop.objectId),
  39. }]);
  40. records++;
  41. } catch (error) {
  42. failures++;
  43. log('feedback_failed', { shop: shop.name, asin, feedbackType, error: error.message });
  44. }
  45. await sleep(500);
  46. }
  47. }
  48. await writeAudit({
  49. dataType: 'customer_feedback', endpoint: '/customerFeedback/2024-06-01',
  50. status: failures ? 'partial' : 'completed', shop,
  51. region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN', marketplaceId: shop.marketplaceId,
  52. records, pages: asins.length * 3, error: failures ? `${failures} calls failed` : '',
  53. startedAt, details: { asins: asins.length, emptyResponses, failures },
  54. });
  55. log('feedback_completed', { shop: shop.name, asins: asins.length, records, emptyResponses, failures });
  56. }
  57. for (const shop of shops) {
  58. const startedAt = new Date();
  59. const statuses = ['WORKING', 'SHIPPED', 'RECEIVING', 'CLOSED', 'CANCELLED', 'DELETED', 'ERROR', 'IN_TRANSIT', 'DELIVERED', 'CHECKED_IN'];
  60. const after = new Date(Date.now() - 729 * 24 * 60 * 60 * 1000).toISOString();
  61. const before = new Date().toISOString();
  62. let token = '';
  63. let pages = 0;
  64. let records = 0;
  65. const seenTokens = new Set();
  66. try {
  67. do {
  68. pages++;
  69. const path = token
  70. ? `/fba/inbound/v0/shipments?QueryType=NEXT_TOKEN&MarketplaceId=${shop.marketplaceId}&NextToken=${encodeURIComponent(token)}`
  71. : `/fba/inbound/v0/shipments?QueryType=DATE_RANGE&MarketplaceId=${shop.marketplaceId}`
  72. + `&ShipmentStatusList=${statuses.join(',')}`
  73. + `&LastUpdatedAfter=${encodeURIComponent(after)}&LastUpdatedBefore=${encodeURIComponent(before)}`;
  74. const result = await forward(shop.objectId, path, 'GET', undefined, 2);
  75. const payload = result.payload || result;
  76. const shipments = payload.ShipmentData || payload.shipments || [];
  77. const rows = shipments.map(item => ({
  78. recordKey: `${regionByMarketplace[shop.marketplaceId] || 'UNKNOWN'}:${item.ShipmentId || hash(item)}`,
  79. region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN',
  80. marketplaceId: shop.marketplaceId,
  81. shipmentId: item.ShipmentId || '',
  82. shipmentName: item.ShipmentName || '',
  83. shipmentStatus: item.ShipmentStatus || '',
  84. destinationFulfillmentCenterId: item.DestinationFulfillmentCenterId || '',
  85. labelPrepType: item.LabelPrepType || '',
  86. capturedAt,
  87. rawData: item,
  88. shop: shopPointer(shop.objectId),
  89. }));
  90. await upsertMany('AmazonInboundShipment', rows);
  91. records += rows.length;
  92. token = String(payload.NextToken || '').trim();
  93. if (token) {
  94. if (seenTokens.has(token)) throw new Error('Repeated inbound NextToken');
  95. seenTokens.add(token);
  96. await sleep(1200);
  97. }
  98. } while (token);
  99. await writeAudit({
  100. dataType: 'inbound_shipments', endpoint: '/fba/inbound/v0/shipments',
  101. status: 'completed', shop, region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN',
  102. marketplaceId: shop.marketplaceId, records, pages, startedAt, details: { after, before },
  103. });
  104. log('inbound_completed', { shop: shop.name, records, pages });
  105. } catch (error) {
  106. await writeAudit({
  107. dataType: 'inbound_shipments', endpoint: '/fba/inbound/v0/shipments',
  108. status: records ? 'partial' : 'failed', shop, region: regionByMarketplace[shop.marketplaceId] || 'UNKNOWN',
  109. marketplaceId: shop.marketplaceId, records, pages, error: error.message, startedAt,
  110. });
  111. log('inbound_failed', { shop: shop.name, records, pages, error: error.message });
  112. }
  113. }