| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281 |
- const parseUrl = process.env.PARSE_URL || 'http://127.0.0.1:3000/parse';
- const apiUrl = process.env.API_URL || 'http://127.0.0.1:3000/api/amazon';
- const appId = process.env.PARSE_APP_ID;
- const masterKey = process.env.PARSE_MASTER_KEY;
- const targetSellerId = process.env.SELLER_ID;
- if (!appId || !masterKey || !targetSellerId) {
- throw new Error('Missing PARSE_APP_ID, PARSE_MASTER_KEY, or SELLER_ID');
- }
- const parseHeaders = {
- 'Content-Type': 'application/json',
- 'X-Parse-Application-Id': appId,
- 'X-Parse-Master-Key': masterKey,
- };
- const sleep = ms => new Promise(resolve => setTimeout(resolve, ms));
- const pointer = objectId => ({ __type: 'Pointer', className: 'Shop', objectId });
- const parseDate = value => value ? { __type: 'Date', iso: new Date(value).toISOString() } : undefined;
- function cleanObject(value) {
- return Object.fromEntries(Object.entries(value).filter(([, item]) => item !== undefined));
- }
- function log(event, data = {}) {
- console.log(JSON.stringify({ at: new Date().toISOString(), event, ...data }));
- }
- async function parseRequest(path, method = 'GET', body) {
- const response = await fetch(`${parseUrl}${path}`, {
- method,
- headers: parseHeaders,
- body: body === undefined ? undefined : JSON.stringify(body),
- });
- const text = await response.text();
- let result = {};
- try { result = JSON.parse(text); } catch {}
- if (!response.ok || result.error) {
- throw new Error(`${method} ${path}: ${result.error?.message || result.error || text.slice(0, 300) || response.status}`);
- }
- return result;
- }
- async function forward(shopId, path) {
- const response = await fetch(`${apiUrl}/forward`, {
- method: 'POST',
- headers: { 'Content-Type': 'application/json', 'shop-objectid': shopId },
- body: JSON.stringify({ path, method: 'GET' }),
- signal: AbortSignal.timeout(180000),
- });
- const text = await response.text();
- let result = {};
- try { result = JSON.parse(text); } catch {}
- if (!response.ok || result.success === false) {
- throw new Error(result.message || result.errors?.[0]?.message || text.slice(0, 300) || `HTTP ${response.status}`);
- }
- return result.data || result;
- }
- async function loadExisting(className, shopId, identityKey) {
- const where = encodeURIComponent(JSON.stringify({ shop: pointer(shopId) }));
- const rows = new Map();
- let skip = 0;
- while (true) {
- const page = await parseRequest(`/classes/${className}?where=${where}&keys=objectId,${identityKey}&limit=1000&skip=${skip}`);
- for (const row of page.results || []) {
- const identity = String(row[identityKey] || '').trim();
- if (identity && !rows.has(identity)) rows.set(identity, row.objectId);
- }
- const length = page.results?.length || 0;
- if (length < 1000) return rows;
- skip += length;
- }
- }
- async function saveRows(className, identityKey, rows, existing) {
- let inserted = 0;
- let updated = 0;
- for (let start = 0; start < rows.length; start += 25) {
- const batch = rows.slice(start, start + 25);
- const requests = batch.map(row => {
- const identity = String(row[identityKey] || '').trim();
- const objectId = existing.get(identity);
- return {
- method: objectId ? 'PUT' : 'POST',
- path: objectId ? `/parse/classes/${className}/${objectId}` : `/parse/classes/${className}`,
- body: row,
- };
- });
- const result = await parseRequest('/batch', 'POST', { requests });
- for (let index = 0; index < batch.length; index++) {
- const response = result[index];
- if (response?.error) throw new Error(`${className} batch save: ${response.error.error || response.error.code}`);
- const identity = String(batch[index][identityKey] || '').trim();
- if (existing.has(identity)) {
- updated++;
- } else {
- existing.set(identity, response?.success?.objectId);
- inserted++;
- }
- }
- }
- return { inserted, updated };
- }
- function mapOrder(raw, shopId, isNew) {
- const shippingAddress = raw.ShippingAddress || undefined;
- const body = cleanObject({
- platformOrderId: raw.AmazonOrderId,
- orderId: raw.AmazonOrderId,
- orderDate: parseDate(raw.PurchaseDate),
- lastUpdateDate: parseDate(raw.LastUpdateDate),
- status: raw.OrderStatus,
- marketplaceId: raw.MarketplaceId,
- totalAmount: raw.OrderTotal?.Amount === undefined ? undefined : Number(raw.OrderTotal.Amount),
- currency: raw.OrderTotal?.CurrencyCode,
- buyerInfo: raw.BuyerInfo,
- shippingAddress,
- customerRegion: shippingAddress?.CountryCode,
- shippingCountry: shippingAddress?.CountryCode,
- shippingState: shippingAddress?.StateOrRegion,
- shippingCity: shippingAddress?.City,
- fulfillmentChannel: raw.FulfillmentChannel,
- orderType: raw.OrderType,
- isPrime: Boolean(raw.IsPrime),
- isBusinessOrder: Boolean(raw.IsBusinessOrder),
- earliestShipDate: parseDate(raw.EarliestShipDate),
- latestShipDate: parseDate(raw.LatestShipDate),
- shop: pointer(shopId),
- });
- if (isNew) {
- if (body.totalAmount === undefined) body.totalAmount = 0;
- body.items = [];
- }
- return body;
- }
- function mapListing(raw, shopId, marketplaceId) {
- const summary = (raw.summaries || []).find(item => item.marketplaceId === marketplaceId) || raw.summaries?.[0] || {};
- const price = raw.attributes?.list_price?.[0];
- const availability = raw.fulfillmentAvailability?.[0];
- const parentRelation = raw.attributes?.child_parent_sku_relationship?.[0];
- return cleanObject({
- sku: raw.sku,
- sellerSku: raw.sku,
- asin: summary.asin,
- parentSku: parentRelation?.parent_sku,
- title: summary.itemName,
- mainImage: summary.mainImage?.link,
- productType: summary.productType,
- status: summary.status?.[0],
- createdDate: parseDate(summary.createdDate),
- lastUpdatedDate: parseDate(summary.lastUpdatedDate),
- marketplaceId: summary.marketplaceId || marketplaceId,
- attributes: raw.attributes,
- issues: raw.issues,
- fulfillmentAvailability: raw.fulfillmentAvailability,
- quantity: availability?.quantity,
- price: price?.value === undefined ? undefined : Number(price.value),
- currency: price?.currency,
- shop: pointer(shopId),
- });
- }
- async function syncOrders(shop, ordersStart) {
- const existing = await loadExisting('Order', shop.objectId, 'platformOrderId');
- const seenTokens = new Set();
- let token = '';
- let page = 0;
- let fetched = 0;
- let inserted = 0;
- let updated = 0;
- do {
- page++;
- if (page > 500) throw new Error('Orders pagination exceeded 500 pages');
- const path = token
- ? `/orders/v0/orders?NextToken=${encodeURIComponent(token)}`
- : `/orders/v0/orders?MarketplaceIds=${encodeURIComponent(shop.marketplaceId)}&CreatedAfter=${encodeURIComponent(ordersStart)}`;
- const result = await forward(shop.objectId, path);
- const rows = result.payload?.Orders || [];
- const bodies = rows
- .filter(row => row.AmazonOrderId)
- .map(row => mapOrder(row, shop.objectId, !existing.has(row.AmazonOrderId)));
- const saved = await saveRows('Order', 'platformOrderId', bodies, existing);
- fetched += rows.length;
- inserted += saved.inserted;
- updated += saved.updated;
- log('orders_page', { shopId: shop.objectId, shopName: shop.name, page, fetched: rows.length, inserted: saved.inserted, updated: saved.updated });
- token = String(result.payload?.NextToken || result.NextToken || '').trim();
- if (token) {
- if (seenTokens.has(token)) throw new Error('Orders API returned a repeated NextToken');
- seenTokens.add(token);
- await sleep(1200);
- }
- } while (token);
- return { pages: page, fetched, inserted, updated, total: existing.size };
- }
- async function syncListings(shop, listingsStart) {
- const existing = await loadExisting('Listing', shop.objectId, 'sku');
- const basePath = `/listings/2021-08-01/items/${encodeURIComponent(targetSellerId)}`;
- const baseQuery = `marketplaceIds=${encodeURIComponent(shop.marketplaceId)}`
- + '&includedData=summaries,attributes,issues,fulfillmentAvailability'
- + '&withStatus=BUYABLE,DISCOVERABLE'
- + '&sortBy=lastUpdatedDate&sortOrder=ASC&pageSize=10'
- + `&lastUpdatedAfter=${encodeURIComponent(listingsStart)}`;
- const seenTokens = new Set();
- let token = '';
- let page = 0;
- let fetched = 0;
- let inserted = 0;
- let updated = 0;
- do {
- page++;
- if (page > 500) throw new Error('Listings pagination exceeded 500 pages');
- const path = `${basePath}?${baseQuery}${token ? `&pageToken=${encodeURIComponent(token)}` : ''}`;
- const result = await forward(shop.objectId, path);
- const rows = result.items || [];
- const bodies = rows.filter(row => row.sku).map(row => mapListing(row, shop.objectId, shop.marketplaceId));
- const saved = await saveRows('Listing', 'sku', bodies, existing);
- fetched += rows.length;
- inserted += saved.inserted;
- updated += saved.updated;
- log('listings_page', { shopId: shop.objectId, shopName: shop.name, page, fetched: rows.length, inserted: saved.inserted, updated: saved.updated });
- token = String(result.pagination?.nextToken || '').trim();
- if (token) {
- if (seenTokens.has(token)) throw new Error('Listings API returned a repeated pageToken');
- seenTokens.add(token);
- await sleep(1200);
- }
- } while (token);
- return { pages: page, fetched, inserted, updated, total: existing.size };
- }
- const startedAt = new Date();
- const ordersStart = new Date(startedAt.getTime() - 729 * 24 * 60 * 60 * 1000).toISOString();
- const listingsStart = new Date(startedAt.getTime() - 365 * 24 * 60 * 60 * 1000).toISOString();
- const shopsResult = await parseRequest('/classes/Shop?limit=1000');
- const shops = (shopsResult.results || []).filter(shop => {
- const sp = shop.config?.SpApiConfig;
- return shop.platform === 'amazon'
- && shop.status === 'active'
- && shop.disable !== true
- && shop.deleted !== true
- && sp?.sellerID === targetSellerId;
- });
- const details = [];
- log('paged_sync_started', { sellerId: targetSellerId, shops: shops.length, ordersStart, listingsStart });
- for (const shop of shops) {
- const detail = { shopId: shop.objectId, shopName: shop.name, marketplaceId: shop.marketplaceId, orders: null, listings: null, errors: [] };
- details.push(detail);
- try {
- detail.orders = await syncOrders(shop, ordersStart);
- log('orders_completed', { shopId: shop.objectId, shopName: shop.name, ...detail.orders });
- } catch (error) {
- detail.errors.push(`orders: ${error.message}`);
- log('orders_failed', { shopId: shop.objectId, shopName: shop.name, error: error.message });
- }
- if (shop.config.SpApiConfig.listingEnabled === false) {
- detail.listings = { skipped: true, reason: 'regional seller ID not authorized' };
- } else {
- try {
- detail.listings = await syncListings(shop, listingsStart);
- log('listings_completed', { shopId: shop.objectId, shopName: shop.name, ...detail.listings });
- } catch (error) {
- detail.errors.push(`listings: ${error.message}`);
- log('listings_failed', { shopId: shop.objectId, shopName: shop.name, error: error.message });
- }
- }
- }
- log('paged_sync_finished', {
- sellerId: targetSellerId,
- startedAt: startedAt.toISOString(),
- finishedAt: new Date().toISOString(),
- shops: shops.length,
- errors: details.flatMap(item => item.errors.map(error => `${item.shopName}: ${error}`)),
- details,
- });
|