sync-sp-api-pages.mjs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281
  1. const parseUrl = process.env.PARSE_URL || 'http://127.0.0.1:3000/parse';
  2. const apiUrl = process.env.API_URL || 'http://127.0.0.1:3000/api/amazon';
  3. const appId = process.env.PARSE_APP_ID;
  4. const masterKey = process.env.PARSE_MASTER_KEY;
  5. const targetSellerId = process.env.SELLER_ID;
  6. if (!appId || !masterKey || !targetSellerId) {
  7. throw new Error('Missing PARSE_APP_ID, PARSE_MASTER_KEY, or SELLER_ID');
  8. }
  9. const parseHeaders = {
  10. 'Content-Type': 'application/json',
  11. 'X-Parse-Application-Id': appId,
  12. 'X-Parse-Master-Key': masterKey,
  13. };
  14. const sleep = ms => new Promise(resolve => setTimeout(resolve, ms));
  15. const pointer = objectId => ({ __type: 'Pointer', className: 'Shop', objectId });
  16. const parseDate = value => value ? { __type: 'Date', iso: new Date(value).toISOString() } : undefined;
  17. function cleanObject(value) {
  18. return Object.fromEntries(Object.entries(value).filter(([, item]) => item !== undefined));
  19. }
  20. function log(event, data = {}) {
  21. console.log(JSON.stringify({ at: new Date().toISOString(), event, ...data }));
  22. }
  23. async function parseRequest(path, method = 'GET', body) {
  24. const response = await fetch(`${parseUrl}${path}`, {
  25. method,
  26. headers: parseHeaders,
  27. body: body === undefined ? undefined : JSON.stringify(body),
  28. });
  29. const text = await response.text();
  30. let result = {};
  31. try { result = JSON.parse(text); } catch {}
  32. if (!response.ok || result.error) {
  33. throw new Error(`${method} ${path}: ${result.error?.message || result.error || text.slice(0, 300) || response.status}`);
  34. }
  35. return result;
  36. }
  37. async function forward(shopId, path) {
  38. const response = await fetch(`${apiUrl}/forward`, {
  39. method: 'POST',
  40. headers: { 'Content-Type': 'application/json', 'shop-objectid': shopId },
  41. body: JSON.stringify({ path, method: 'GET' }),
  42. signal: AbortSignal.timeout(180000),
  43. });
  44. const text = await response.text();
  45. let result = {};
  46. try { result = JSON.parse(text); } catch {}
  47. if (!response.ok || result.success === false) {
  48. throw new Error(result.message || result.errors?.[0]?.message || text.slice(0, 300) || `HTTP ${response.status}`);
  49. }
  50. return result.data || result;
  51. }
  52. async function loadExisting(className, shopId, identityKey) {
  53. const where = encodeURIComponent(JSON.stringify({ shop: pointer(shopId) }));
  54. const rows = new Map();
  55. let skip = 0;
  56. while (true) {
  57. const page = await parseRequest(`/classes/${className}?where=${where}&keys=objectId,${identityKey}&limit=1000&skip=${skip}`);
  58. for (const row of page.results || []) {
  59. const identity = String(row[identityKey] || '').trim();
  60. if (identity && !rows.has(identity)) rows.set(identity, row.objectId);
  61. }
  62. const length = page.results?.length || 0;
  63. if (length < 1000) return rows;
  64. skip += length;
  65. }
  66. }
  67. async function saveRows(className, identityKey, rows, existing) {
  68. let inserted = 0;
  69. let updated = 0;
  70. for (let start = 0; start < rows.length; start += 25) {
  71. const batch = rows.slice(start, start + 25);
  72. const requests = batch.map(row => {
  73. const identity = String(row[identityKey] || '').trim();
  74. const objectId = existing.get(identity);
  75. return {
  76. method: objectId ? 'PUT' : 'POST',
  77. path: objectId ? `/parse/classes/${className}/${objectId}` : `/parse/classes/${className}`,
  78. body: row,
  79. };
  80. });
  81. const result = await parseRequest('/batch', 'POST', { requests });
  82. for (let index = 0; index < batch.length; index++) {
  83. const response = result[index];
  84. if (response?.error) throw new Error(`${className} batch save: ${response.error.error || response.error.code}`);
  85. const identity = String(batch[index][identityKey] || '').trim();
  86. if (existing.has(identity)) {
  87. updated++;
  88. } else {
  89. existing.set(identity, response?.success?.objectId);
  90. inserted++;
  91. }
  92. }
  93. }
  94. return { inserted, updated };
  95. }
  96. function mapOrder(raw, shopId, isNew) {
  97. const shippingAddress = raw.ShippingAddress || undefined;
  98. const body = cleanObject({
  99. platformOrderId: raw.AmazonOrderId,
  100. orderId: raw.AmazonOrderId,
  101. orderDate: parseDate(raw.PurchaseDate),
  102. lastUpdateDate: parseDate(raw.LastUpdateDate),
  103. status: raw.OrderStatus,
  104. marketplaceId: raw.MarketplaceId,
  105. totalAmount: raw.OrderTotal?.Amount === undefined ? undefined : Number(raw.OrderTotal.Amount),
  106. currency: raw.OrderTotal?.CurrencyCode,
  107. buyerInfo: raw.BuyerInfo,
  108. shippingAddress,
  109. customerRegion: shippingAddress?.CountryCode,
  110. shippingCountry: shippingAddress?.CountryCode,
  111. shippingState: shippingAddress?.StateOrRegion,
  112. shippingCity: shippingAddress?.City,
  113. fulfillmentChannel: raw.FulfillmentChannel,
  114. orderType: raw.OrderType,
  115. isPrime: Boolean(raw.IsPrime),
  116. isBusinessOrder: Boolean(raw.IsBusinessOrder),
  117. earliestShipDate: parseDate(raw.EarliestShipDate),
  118. latestShipDate: parseDate(raw.LatestShipDate),
  119. shop: pointer(shopId),
  120. });
  121. if (isNew) {
  122. if (body.totalAmount === undefined) body.totalAmount = 0;
  123. body.items = [];
  124. }
  125. return body;
  126. }
  127. function mapListing(raw, shopId, marketplaceId) {
  128. const summary = (raw.summaries || []).find(item => item.marketplaceId === marketplaceId) || raw.summaries?.[0] || {};
  129. const price = raw.attributes?.list_price?.[0];
  130. const availability = raw.fulfillmentAvailability?.[0];
  131. const parentRelation = raw.attributes?.child_parent_sku_relationship?.[0];
  132. return cleanObject({
  133. sku: raw.sku,
  134. sellerSku: raw.sku,
  135. asin: summary.asin,
  136. parentSku: parentRelation?.parent_sku,
  137. title: summary.itemName,
  138. mainImage: summary.mainImage?.link,
  139. productType: summary.productType,
  140. status: summary.status?.[0],
  141. createdDate: parseDate(summary.createdDate),
  142. lastUpdatedDate: parseDate(summary.lastUpdatedDate),
  143. marketplaceId: summary.marketplaceId || marketplaceId,
  144. attributes: raw.attributes,
  145. issues: raw.issues,
  146. fulfillmentAvailability: raw.fulfillmentAvailability,
  147. quantity: availability?.quantity,
  148. price: price?.value === undefined ? undefined : Number(price.value),
  149. currency: price?.currency,
  150. shop: pointer(shopId),
  151. });
  152. }
  153. async function syncOrders(shop, ordersStart) {
  154. const existing = await loadExisting('Order', shop.objectId, 'platformOrderId');
  155. const seenTokens = new Set();
  156. let token = '';
  157. let page = 0;
  158. let fetched = 0;
  159. let inserted = 0;
  160. let updated = 0;
  161. do {
  162. page++;
  163. if (page > 500) throw new Error('Orders pagination exceeded 500 pages');
  164. const path = token
  165. ? `/orders/v0/orders?NextToken=${encodeURIComponent(token)}`
  166. : `/orders/v0/orders?MarketplaceIds=${encodeURIComponent(shop.marketplaceId)}&CreatedAfter=${encodeURIComponent(ordersStart)}`;
  167. const result = await forward(shop.objectId, path);
  168. const rows = result.payload?.Orders || [];
  169. const bodies = rows
  170. .filter(row => row.AmazonOrderId)
  171. .map(row => mapOrder(row, shop.objectId, !existing.has(row.AmazonOrderId)));
  172. const saved = await saveRows('Order', 'platformOrderId', bodies, existing);
  173. fetched += rows.length;
  174. inserted += saved.inserted;
  175. updated += saved.updated;
  176. log('orders_page', { shopId: shop.objectId, shopName: shop.name, page, fetched: rows.length, inserted: saved.inserted, updated: saved.updated });
  177. token = String(result.payload?.NextToken || result.NextToken || '').trim();
  178. if (token) {
  179. if (seenTokens.has(token)) throw new Error('Orders API returned a repeated NextToken');
  180. seenTokens.add(token);
  181. await sleep(1200);
  182. }
  183. } while (token);
  184. return { pages: page, fetched, inserted, updated, total: existing.size };
  185. }
  186. async function syncListings(shop, listingsStart) {
  187. const existing = await loadExisting('Listing', shop.objectId, 'sku');
  188. const basePath = `/listings/2021-08-01/items/${encodeURIComponent(targetSellerId)}`;
  189. const baseQuery = `marketplaceIds=${encodeURIComponent(shop.marketplaceId)}`
  190. + '&includedData=summaries,attributes,issues,fulfillmentAvailability'
  191. + '&withStatus=BUYABLE,DISCOVERABLE'
  192. + '&sortBy=lastUpdatedDate&sortOrder=ASC&pageSize=10'
  193. + `&lastUpdatedAfter=${encodeURIComponent(listingsStart)}`;
  194. const seenTokens = new Set();
  195. let token = '';
  196. let page = 0;
  197. let fetched = 0;
  198. let inserted = 0;
  199. let updated = 0;
  200. do {
  201. page++;
  202. if (page > 500) throw new Error('Listings pagination exceeded 500 pages');
  203. const path = `${basePath}?${baseQuery}${token ? `&pageToken=${encodeURIComponent(token)}` : ''}`;
  204. const result = await forward(shop.objectId, path);
  205. const rows = result.items || [];
  206. const bodies = rows.filter(row => row.sku).map(row => mapListing(row, shop.objectId, shop.marketplaceId));
  207. const saved = await saveRows('Listing', 'sku', bodies, existing);
  208. fetched += rows.length;
  209. inserted += saved.inserted;
  210. updated += saved.updated;
  211. log('listings_page', { shopId: shop.objectId, shopName: shop.name, page, fetched: rows.length, inserted: saved.inserted, updated: saved.updated });
  212. token = String(result.pagination?.nextToken || '').trim();
  213. if (token) {
  214. if (seenTokens.has(token)) throw new Error('Listings API returned a repeated pageToken');
  215. seenTokens.add(token);
  216. await sleep(1200);
  217. }
  218. } while (token);
  219. return { pages: page, fetched, inserted, updated, total: existing.size };
  220. }
  221. const startedAt = new Date();
  222. const ordersStart = new Date(startedAt.getTime() - 729 * 24 * 60 * 60 * 1000).toISOString();
  223. const listingsStart = new Date(startedAt.getTime() - 365 * 24 * 60 * 60 * 1000).toISOString();
  224. const shopsResult = await parseRequest('/classes/Shop?limit=1000');
  225. const shops = (shopsResult.results || []).filter(shop => {
  226. const sp = shop.config?.SpApiConfig;
  227. return shop.platform === 'amazon'
  228. && shop.status === 'active'
  229. && shop.disable !== true
  230. && shop.deleted !== true
  231. && sp?.sellerID === targetSellerId;
  232. });
  233. const details = [];
  234. log('paged_sync_started', { sellerId: targetSellerId, shops: shops.length, ordersStart, listingsStart });
  235. for (const shop of shops) {
  236. const detail = { shopId: shop.objectId, shopName: shop.name, marketplaceId: shop.marketplaceId, orders: null, listings: null, errors: [] };
  237. details.push(detail);
  238. try {
  239. detail.orders = await syncOrders(shop, ordersStart);
  240. log('orders_completed', { shopId: shop.objectId, shopName: shop.name, ...detail.orders });
  241. } catch (error) {
  242. detail.errors.push(`orders: ${error.message}`);
  243. log('orders_failed', { shopId: shop.objectId, shopName: shop.name, error: error.message });
  244. }
  245. if (shop.config.SpApiConfig.listingEnabled === false) {
  246. detail.listings = { skipped: true, reason: 'regional seller ID not authorized' };
  247. } else {
  248. try {
  249. detail.listings = await syncListings(shop, listingsStart);
  250. log('listings_completed', { shopId: shop.objectId, shopName: shop.name, ...detail.listings });
  251. } catch (error) {
  252. detail.errors.push(`listings: ${error.message}`);
  253. log('listings_failed', { shopId: shop.objectId, shopName: shop.name, error: error.message });
  254. }
  255. }
  256. }
  257. log('paged_sync_finished', {
  258. sellerId: targetSellerId,
  259. startedAt: startedAt.toISOString(),
  260. finishedAt: new Date().toISOString(),
  261. shops: shops.length,
  262. errors: details.flatMap(item => item.errors.map(error => `${item.shopName}: ${error}`)),
  263. details,
  264. });