import { activeAmazonShops, dateValue, forward, hash, log, regionByMarketplace, shopPointer, sleep, upsertMany, writeAudit, } from './sp-api-collector-lib.mjs'; 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); } const postedAfter = new Date(Date.now() - 729 * 24 * 60 * 60 * 1000).toISOString(); const capturedAt = dateValue(new Date()); for (const [region, shop] of regionShops) { const startedAt = new Date(); let token = ''; let pages = 0; let records = 0; const seenTokens = new Set(); const typeCounts = {}; try { do { pages++; if (pages > 1000) throw new Error('Finances pagination exceeded 1000 pages'); const path = token ? `/finances/v0/financialEvents?NextToken=${encodeURIComponent(token)}` : `/finances/v0/financialEvents?PostedAfter=${encodeURIComponent(postedAfter)}&MaxResultsPerPage=100`; const result = await forward(shop.objectId, path); const payload = result.payload || result; const financialEvents = payload.FinancialEvents || {}; const rows = []; for (const [eventListName, events] of Object.entries(financialEvents)) { if (!Array.isArray(events)) continue; const eventType = eventListName.replace(/List$/, ''); typeCounts[eventType] = (typeCounts[eventType] || 0) + events.length; for (const event of events) { const eventHash = hash(event); rows.push({ recordKey: `${region}:${eventType}:${eventHash}`, region, eventType, amazonOrderId: event.AmazonOrderId || event.amazonOrderId || event.OrderId || '', sellerOrderId: event.SellerOrderId || event.MerchantOrderId || '', marketplaceName: event.MarketplaceName || '', postedDate: dateValue(event.PostedDate || event.postedDate), capturedAt, rawData: event, shop: shopPointer(shop.objectId), }); } } await upsertMany('AmazonFinancialEvent', rows); records += rows.length; log('finances_page', { region, page: pages, records: rows.length, total: records }); token = String(payload.NextToken || result.NextToken || '').trim(); if (token) { if (seenTokens.has(token)) throw new Error('Repeated finances NextToken'); seenTokens.add(token); await sleep(1800); } } while (token); await writeAudit({ dataType: 'financial_events', endpoint: '/finances/v0/financialEvents', status: 'completed', shop, region, marketplaceId: shop.marketplaceId, records, pages, startedAt, details: { postedAfter, typeCounts }, }); log('finances_completed', { region, records, pages, typeCounts }); } catch (error) { await writeAudit({ dataType: 'financial_events', endpoint: '/finances/v0/financialEvents', status: records ? 'partial' : 'failed', shop, region, marketplaceId: shop.marketplaceId, records, pages, error: error.message, startedAt, details: { postedAfter, typeCounts }, }); log('finances_failed', { region, records, pages, typeCounts, error: error.message }); } }