| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289 |
- import crypto from 'node:crypto';
- import fs from 'node:fs';
- const configPath = process.env.VOC_CONFIG || '/etc/jianhen-voc/config.json';
- const config = JSON.parse(fs.readFileSync(configPath, 'utf8'));
- const parseUrl = process.env.PARSE_URL || 'http://127.0.0.1:3000/parse';
- const relayUrl = process.env.SORFTIME_URL || 'http://127.0.0.1:3000/api/sorftime/forward';
- const appId = config.parse.appId;
- const masterKey = config.parse.masterKey;
- const maxRowsPerResponse = 10;
- const headers = {
- 'Content-Type': 'application/json',
- 'X-Parse-Application-Id': appId,
- 'X-Parse-Master-Key': masterKey,
- };
- const sleep = ms => new Promise(resolve => setTimeout(resolve, ms));
- const shopPointer = objectId => ({ __type: 'Pointer', className: 'Shop', objectId });
- const isoDate = date => date.toISOString().slice(0, 10);
- function log(event, data = {}) {
- console.log(JSON.stringify({ at: new Date().toISOString(), event, ...data }));
- }
- function hash(value) {
- return crypto.createHash('sha256').update(value).digest('hex');
- }
- async function parseRequest(path, method = 'GET', body) {
- const response = await fetch(`${parseUrl}${path}`, {
- method,
- headers,
- 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, 400) || response.status}`);
- }
- return result;
- }
- async function allRows(className, keys = '') {
- const rows = [];
- let skip = 0;
- while (true) {
- const page = await parseRequest(`/classes/${className}?limit=1000&skip=${skip}${keys ? `&keys=${keys}` : ''}`);
- rows.push(...(page.results || []));
- const length = page.results?.length || 0;
- if (length < 1000) return rows;
- skip += length;
- }
- }
- function unwrap(result) {
- let data = result?.data ?? result?.Data ?? result;
- while (data && typeof data === 'object' && !Array.isArray(data)
- && (data.data !== undefined || data.Data !== undefined)) {
- data = data.data ?? data.Data;
- }
- return data;
- }
- async function queryReviews({ asin, domain, start, end, star = '1,2,3,4,5' }) {
- for (let attempt = 1; attempt <= 4; attempt++) {
- try {
- const response = await fetch(relayUrl, {
- method: 'POST',
- headers: { 'Content-Type': 'application/json' },
- body: JSON.stringify({
- path: '/api/ProductReviewsQuery',
- method: 'POST',
- query: { domain },
- body: {
- ASIN: asin,
- PageIndex: 1,
- PageSize: 100,
- OnlyPurchase: 0,
- Star: star,
- QueryStartDt: start,
- QueryEndDt: end,
- },
- }),
- signal: AbortSignal.timeout(90000),
- });
- const result = await response.json();
- if (!response.ok || result?.success === false) {
- throw new Error(result?.message || result?.Message || `Relay HTTP ${response.status}`);
- }
- const data = unwrap(result);
- return data?.Reviews || data?.Items || data?.List || (Array.isArray(data) ? data : []);
- } catch (error) {
- if (attempt === 4) throw error;
- await sleep(attempt * 3000);
- }
- }
- }
- function reviewIdOf(review, target) {
- const id = String(review.ReviewId || review.Id || review.ID || review.id || '').trim();
- if (id) return id;
- const identity = [
- target.domain,
- target.asin,
- review.ReviewsDate || review.ReviewDate || review.Date || '',
- review.ProfileName || review.Author || '',
- review.Content || review.ReviewContent || review.Body || '',
- ].join('|');
- return `derived-${hash(identity)}`;
- }
- function reviewDateOf(review) {
- return String(review.ReviewsDate || review.ReviewDate || review.Date || '').trim();
- }
- function parseReviewDate(value) {
- if (!value) return undefined;
- const normalized = /^\d{8}$/.test(value)
- ? `${value.slice(0, 4)}-${value.slice(4, 6)}-${value.slice(6, 8)}`
- : value;
- const date = new Date(normalized);
- return Number.isNaN(date.getTime()) ? undefined : { __type: 'Date', iso: date.toISOString() };
- }
- async function collectStarBuckets(target, start, end, resultMap, stats) {
- for (const star of ['1', '2', '3', '4', '5']) {
- const starReviews = await queryReviews({ asin: target.asin, domain: target.domain, start, end, star });
- stats.requests++;
- for (const review of starReviews) resultMap.set(reviewIdOf(review, target), review);
- log('review_star_bucket', { asin: target.asin, domain: target.domain, start, end, star, rows: starReviews.length });
- if (starReviews.length >= maxRowsPerResponse) {
- stats.unresolvedSaturatedBuckets.push({ asin: target.asin, domain: target.domain, start, end, star });
- }
- }
- }
- async function saveBatch(requests) {
- for (let start = 0; start < requests.length; start += 25) {
- const result = await parseRequest('/batch', 'POST', { requests: requests.slice(start, start + 25) });
- const failure = result.find?.(item => item.error)?.error;
- if (failure) throw new Error(failure.error || failure.code || 'Parse batch failed');
- }
- }
- const [products, listings, shops, existingReviews] = await Promise.all([
- allRows('Product'),
- allRows('Listing'),
- allRows('Shop'),
- allRows('SorftimeReviews', 'objectId,recordKey,reviewId,asin,domain,shop'),
- ]);
- const shopById = new Map(shops.map(shop => [shop.objectId, shop]));
- const existingByReviewId = new Map();
- for (const review of existingReviews) {
- if (review.reviewId && !existingByReviewId.has(review.reviewId)) {
- existingByReviewId.set(review.reviewId, review);
- }
- }
- const targets = new Map();
- for (const row of [...products, ...listings]) {
- const asin = String(row.asin || row.Asin || '').trim().toUpperCase();
- const shop = shopById.get(row.shop?.objectId);
- if (!asin || !shop) continue;
- const domain = Number(row.domain || shop.domain || 1);
- const marketplaceId = String(row.marketplaceId || shop.marketplaceId || '');
- const key = `${domain}:${asin}`;
- if (!targets.has(key)) targets.set(key, { asin, domain, marketplaceId, shopId: shop.objectId, shopName: shop.name || '' });
- }
- const endDate = isoDate(new Date());
- const start = new Date();
- start.setUTCFullYear(start.getUTCFullYear() - 1);
- const startDate = isoDate(start);
- const runStats = {
- startedAt: new Date().toISOString(), startDate, endDate, targets: targets.size,
- requests: 0, fetchedUnique: 0, associations: 0, withContent: 0, inserted: 0, updated: 0,
- failures: [], unresolvedSaturatedBuckets: [], byTarget: [],
- };
- const globalReviews = new Map();
- const associations = new Map();
- log('review_collection_started', { targets: targets.size, startDate, endDate });
- for (const target of targets.values()) {
- const collected = new Map();
- const targetStats = { requests: 0, unresolvedSaturatedBuckets: [] };
- try {
- await collectStarBuckets(target, startDate, endDate, collected, targetStats);
- let withContent = 0;
- for (const [reviewId, review] of collected) {
- const content = String(review.Content || review.ReviewContent || review.Body || '');
- if (content.trim()) withContent++;
- if (!globalReviews.has(reviewId)) {
- globalReviews.set(reviewId, {
- review,
- asins: new Set(),
- marketplaceIds: new Set(),
- domains: new Set(),
- shopIds: new Set(),
- targets: [],
- });
- }
- const global = globalReviews.get(reviewId);
- global.asins.add(target.asin);
- if (target.marketplaceId) global.marketplaceIds.add(target.marketplaceId);
- global.domains.add(target.domain);
- global.shopIds.add(target.shopId);
- global.targets.push(target);
- const associationKey = `${target.domain}:${target.asin}:${reviewId}`;
- associations.set(associationKey, {
- recordKey: associationKey,
- reviewId,
- asin: target.asin,
- marketplaceId: target.marketplaceId,
- domain: target.domain,
- source: 'sorftime-product-reviews-year',
- capturedAt: { __type: 'Date', iso: new Date().toISOString() },
- shop: shopPointer(target.shopId),
- });
- }
- runStats.requests += targetStats.requests;
- runStats.unresolvedSaturatedBuckets.push(...targetStats.unresolvedSaturatedBuckets);
- runStats.byTarget.push({ ...target, reviews: collected.size, withContent, requests: targetStats.requests });
- log('review_target_completed', { ...target, reviews: collected.size, withContent, requests: targetStats.requests });
- } catch (error) {
- runStats.failures.push({ ...target, error: error.message });
- log('review_target_failed', { ...target, error: error.message });
- }
- await sleep(500);
- }
- const capturedAt = { __type: 'Date', iso: new Date().toISOString() };
- const reviewRequests = [];
- for (const [reviewId, global] of globalReviews) {
- const existing = existingByReviewId.get(reviewId);
- const targetsForReview = global.targets.slice().sort((a, b) => `${a.domain}:${a.asin}`.localeCompare(`${b.domain}:${b.asin}`));
- const primary = targetsForReview.find(target => target.asin === existing?.asin) || targetsForReview[0];
- const review = global.review;
- const content = String(review.Content || review.ReviewContent || review.Body || '');
- reviewRequests.push({
- method: existing ? 'PUT' : 'POST',
- path: existing ? `/parse/classes/SorftimeReviews/${existing.objectId}` : '/parse/classes/SorftimeReviews',
- body: {
- recordKey: reviewId,
- reviewId,
- asin: primary.asin,
- asins: [...global.asins].sort(),
- marketplaceId: primary.marketplaceId,
- marketplaceIds: [...global.marketplaceIds].sort(),
- domain: primary.domain,
- domains: [...global.domains].sort((a, b) => a - b),
- shopIds: [...global.shopIds].sort(),
- source: 'sorftime-product-reviews-year',
- star: Number(review.Star || review.Rating || 0),
- rating: Number(review.Star || review.Rating || 0),
- title: review.Title || review.ReviewTitle || '',
- content,
- reviewDate: reviewDateOf(review),
- reviewPublishedAt: parseReviewDate(reviewDateOf(review)),
- author: review.ProfileName || review.Author || '',
- verified: Boolean(review.OnlyPurchase || review.VerifiedPurchase),
- capturedAt,
- shop: shopPointer(primary.shopId),
- rawData: review,
- },
- });
- if (existing) runStats.updated++; else runStats.inserted++;
- if (content.trim()) runStats.withContent++;
- }
- await saveBatch(reviewRequests);
- const existingAssociations = await allRows('SorftimeReviewProduct', 'objectId,recordKey');
- const associationByKey = new Map(existingAssociations.filter(row => row.recordKey).map(row => [row.recordKey, row.objectId]));
- const associationRequests = [...associations.values()].map(body => {
- const objectId = associationByKey.get(body.recordKey);
- return {
- method: objectId ? 'PUT' : 'POST',
- path: objectId ? `/parse/classes/SorftimeReviewProduct/${objectId}` : '/parse/classes/SorftimeReviewProduct',
- body,
- };
- });
- await saveBatch(associationRequests);
- runStats.fetchedUnique = globalReviews.size;
- runStats.associations = associations.size;
- runStats.finishedAt = new Date().toISOString();
- console.log(JSON.stringify({ event: 'review_collection_finished', ...runStats }));
|