collect-sp-api-finances.mjs 3.2 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879
  1. import {
  2. activeAmazonShops, dateValue, forward, hash, log, regionByMarketplace,
  3. shopPointer, sleep, upsertMany, writeAudit,
  4. } from './sp-api-collector-lib.mjs';
  5. const shops = await activeAmazonShops();
  6. const regionShops = new Map();
  7. for (const shop of shops) {
  8. const region = regionByMarketplace[shop.marketplaceId] || 'UNKNOWN';
  9. if (!regionShops.has(region)) regionShops.set(region, shop);
  10. }
  11. const postedAfter = new Date(Date.now() - 729 * 24 * 60 * 60 * 1000).toISOString();
  12. const capturedAt = dateValue(new Date());
  13. for (const [region, shop] of regionShops) {
  14. const startedAt = new Date();
  15. let token = '';
  16. let pages = 0;
  17. let records = 0;
  18. const seenTokens = new Set();
  19. const typeCounts = {};
  20. try {
  21. do {
  22. pages++;
  23. if (pages > 1000) throw new Error('Finances pagination exceeded 1000 pages');
  24. const path = token
  25. ? `/finances/v0/financialEvents?NextToken=${encodeURIComponent(token)}`
  26. : `/finances/v0/financialEvents?PostedAfter=${encodeURIComponent(postedAfter)}&MaxResultsPerPage=100`;
  27. const result = await forward(shop.objectId, path);
  28. const payload = result.payload || result;
  29. const financialEvents = payload.FinancialEvents || {};
  30. const rows = [];
  31. for (const [eventListName, events] of Object.entries(financialEvents)) {
  32. if (!Array.isArray(events)) continue;
  33. const eventType = eventListName.replace(/List$/, '');
  34. typeCounts[eventType] = (typeCounts[eventType] || 0) + events.length;
  35. for (const event of events) {
  36. const eventHash = hash(event);
  37. rows.push({
  38. recordKey: `${region}:${eventType}:${eventHash}`,
  39. region,
  40. eventType,
  41. amazonOrderId: event.AmazonOrderId || event.amazonOrderId || event.OrderId || '',
  42. sellerOrderId: event.SellerOrderId || event.MerchantOrderId || '',
  43. marketplaceName: event.MarketplaceName || '',
  44. postedDate: dateValue(event.PostedDate || event.postedDate),
  45. capturedAt,
  46. rawData: event,
  47. shop: shopPointer(shop.objectId),
  48. });
  49. }
  50. }
  51. await upsertMany('AmazonFinancialEvent', rows);
  52. records += rows.length;
  53. log('finances_page', { region, page: pages, records: rows.length, total: records });
  54. token = String(payload.NextToken || result.NextToken || '').trim();
  55. if (token) {
  56. if (seenTokens.has(token)) throw new Error('Repeated finances NextToken');
  57. seenTokens.add(token);
  58. await sleep(1800);
  59. }
  60. } while (token);
  61. await writeAudit({
  62. dataType: 'financial_events', endpoint: '/finances/v0/financialEvents',
  63. status: 'completed', shop, region, marketplaceId: shop.marketplaceId,
  64. records, pages, startedAt, details: { postedAfter, typeCounts },
  65. });
  66. log('finances_completed', { region, records, pages, typeCounts });
  67. } catch (error) {
  68. await writeAudit({
  69. dataType: 'financial_events', endpoint: '/finances/v0/financialEvents',
  70. status: records ? 'partial' : 'failed', shop, region, marketplaceId: shop.marketplaceId,
  71. records, pages, error: error.message, startedAt, details: { postedAfter, typeCounts },
  72. });
  73. log('finances_failed', { region, records, pages, typeCounts, error: error.message });
  74. }
  75. }