collect-sorftime-reviews-year.mjs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  1. import crypto from 'node:crypto';
  2. import fs from 'node:fs';
  3. const configPath = process.env.VOC_CONFIG || '/etc/jianhen-voc/config.json';
  4. const config = JSON.parse(fs.readFileSync(configPath, 'utf8'));
  5. const parseUrl = process.env.PARSE_URL || 'http://127.0.0.1:3000/parse';
  6. const relayUrl = process.env.SORFTIME_URL || 'http://127.0.0.1:3000/api/sorftime/forward';
  7. const appId = config.parse.appId;
  8. const masterKey = config.parse.masterKey;
  9. const maxRowsPerResponse = 10;
  10. const headers = {
  11. 'Content-Type': 'application/json',
  12. 'X-Parse-Application-Id': appId,
  13. 'X-Parse-Master-Key': masterKey,
  14. };
  15. const sleep = ms => new Promise(resolve => setTimeout(resolve, ms));
  16. const shopPointer = objectId => ({ __type: 'Pointer', className: 'Shop', objectId });
  17. const isoDate = date => date.toISOString().slice(0, 10);
  18. function log(event, data = {}) {
  19. console.log(JSON.stringify({ at: new Date().toISOString(), event, ...data }));
  20. }
  21. function hash(value) {
  22. return crypto.createHash('sha256').update(value).digest('hex');
  23. }
  24. async function parseRequest(path, method = 'GET', body) {
  25. const response = await fetch(`${parseUrl}${path}`, {
  26. method,
  27. headers,
  28. body: body === undefined ? undefined : JSON.stringify(body),
  29. });
  30. const text = await response.text();
  31. let result = {};
  32. try { result = JSON.parse(text); } catch {}
  33. if (!response.ok || result.error) {
  34. throw new Error(`${method} ${path}: ${result.error?.message || result.error || text.slice(0, 400) || response.status}`);
  35. }
  36. return result;
  37. }
  38. async function allRows(className, keys = '') {
  39. const rows = [];
  40. let skip = 0;
  41. while (true) {
  42. const page = await parseRequest(`/classes/${className}?limit=1000&skip=${skip}${keys ? `&keys=${keys}` : ''}`);
  43. rows.push(...(page.results || []));
  44. const length = page.results?.length || 0;
  45. if (length < 1000) return rows;
  46. skip += length;
  47. }
  48. }
  49. function unwrap(result) {
  50. let data = result?.data ?? result?.Data ?? result;
  51. while (data && typeof data === 'object' && !Array.isArray(data)
  52. && (data.data !== undefined || data.Data !== undefined)) {
  53. data = data.data ?? data.Data;
  54. }
  55. return data;
  56. }
  57. async function queryReviews({ asin, domain, start, end, star = '1,2,3,4,5' }) {
  58. for (let attempt = 1; attempt <= 4; attempt++) {
  59. try {
  60. const response = await fetch(relayUrl, {
  61. method: 'POST',
  62. headers: { 'Content-Type': 'application/json' },
  63. body: JSON.stringify({
  64. path: '/api/ProductReviewsQuery',
  65. method: 'POST',
  66. query: { domain },
  67. body: {
  68. ASIN: asin,
  69. PageIndex: 1,
  70. PageSize: 100,
  71. OnlyPurchase: 0,
  72. Star: star,
  73. QueryStartDt: start,
  74. QueryEndDt: end,
  75. },
  76. }),
  77. signal: AbortSignal.timeout(90000),
  78. });
  79. const result = await response.json();
  80. if (!response.ok || result?.success === false) {
  81. throw new Error(result?.message || result?.Message || `Relay HTTP ${response.status}`);
  82. }
  83. const data = unwrap(result);
  84. return data?.Reviews || data?.Items || data?.List || (Array.isArray(data) ? data : []);
  85. } catch (error) {
  86. if (attempt === 4) throw error;
  87. await sleep(attempt * 3000);
  88. }
  89. }
  90. }
  91. function reviewIdOf(review, target) {
  92. const id = String(review.ReviewId || review.Id || review.ID || review.id || '').trim();
  93. if (id) return id;
  94. const identity = [
  95. target.domain,
  96. target.asin,
  97. review.ReviewsDate || review.ReviewDate || review.Date || '',
  98. review.ProfileName || review.Author || '',
  99. review.Content || review.ReviewContent || review.Body || '',
  100. ].join('|');
  101. return `derived-${hash(identity)}`;
  102. }
  103. function reviewDateOf(review) {
  104. return String(review.ReviewsDate || review.ReviewDate || review.Date || '').trim();
  105. }
  106. function parseReviewDate(value) {
  107. if (!value) return undefined;
  108. const normalized = /^\d{8}$/.test(value)
  109. ? `${value.slice(0, 4)}-${value.slice(4, 6)}-${value.slice(6, 8)}`
  110. : value;
  111. const date = new Date(normalized);
  112. return Number.isNaN(date.getTime()) ? undefined : { __type: 'Date', iso: date.toISOString() };
  113. }
  114. async function collectStarBuckets(target, start, end, resultMap, stats) {
  115. for (const star of ['1', '2', '3', '4', '5']) {
  116. const starReviews = await queryReviews({ asin: target.asin, domain: target.domain, start, end, star });
  117. stats.requests++;
  118. for (const review of starReviews) resultMap.set(reviewIdOf(review, target), review);
  119. log('review_star_bucket', { asin: target.asin, domain: target.domain, start, end, star, rows: starReviews.length });
  120. if (starReviews.length >= maxRowsPerResponse) {
  121. stats.unresolvedSaturatedBuckets.push({ asin: target.asin, domain: target.domain, start, end, star });
  122. }
  123. }
  124. }
  125. async function saveBatch(requests) {
  126. for (let start = 0; start < requests.length; start += 25) {
  127. const result = await parseRequest('/batch', 'POST', { requests: requests.slice(start, start + 25) });
  128. const failure = result.find?.(item => item.error)?.error;
  129. if (failure) throw new Error(failure.error || failure.code || 'Parse batch failed');
  130. }
  131. }
  132. const [products, listings, shops, existingReviews] = await Promise.all([
  133. allRows('Product'),
  134. allRows('Listing'),
  135. allRows('Shop'),
  136. allRows('SorftimeReviews', 'objectId,recordKey,reviewId,asin,domain,shop'),
  137. ]);
  138. const shopById = new Map(shops.map(shop => [shop.objectId, shop]));
  139. const existingByReviewId = new Map();
  140. for (const review of existingReviews) {
  141. if (review.reviewId && !existingByReviewId.has(review.reviewId)) {
  142. existingByReviewId.set(review.reviewId, review);
  143. }
  144. }
  145. const targets = new Map();
  146. for (const row of [...products, ...listings]) {
  147. const asin = String(row.asin || row.Asin || '').trim().toUpperCase();
  148. const shop = shopById.get(row.shop?.objectId);
  149. if (!asin || !shop) continue;
  150. const domain = Number(row.domain || shop.domain || 1);
  151. const marketplaceId = String(row.marketplaceId || shop.marketplaceId || '');
  152. const key = `${domain}:${asin}`;
  153. if (!targets.has(key)) targets.set(key, { asin, domain, marketplaceId, shopId: shop.objectId, shopName: shop.name || '' });
  154. }
  155. const endDate = isoDate(new Date());
  156. const start = new Date();
  157. start.setUTCFullYear(start.getUTCFullYear() - 1);
  158. const startDate = isoDate(start);
  159. const runStats = {
  160. startedAt: new Date().toISOString(), startDate, endDate, targets: targets.size,
  161. requests: 0, fetchedUnique: 0, associations: 0, withContent: 0, inserted: 0, updated: 0,
  162. failures: [], unresolvedSaturatedBuckets: [], byTarget: [],
  163. };
  164. const globalReviews = new Map();
  165. const associations = new Map();
  166. log('review_collection_started', { targets: targets.size, startDate, endDate });
  167. for (const target of targets.values()) {
  168. const collected = new Map();
  169. const targetStats = { requests: 0, unresolvedSaturatedBuckets: [] };
  170. try {
  171. await collectStarBuckets(target, startDate, endDate, collected, targetStats);
  172. let withContent = 0;
  173. for (const [reviewId, review] of collected) {
  174. const content = String(review.Content || review.ReviewContent || review.Body || '');
  175. if (content.trim()) withContent++;
  176. if (!globalReviews.has(reviewId)) {
  177. globalReviews.set(reviewId, {
  178. review,
  179. asins: new Set(),
  180. marketplaceIds: new Set(),
  181. domains: new Set(),
  182. shopIds: new Set(),
  183. targets: [],
  184. });
  185. }
  186. const global = globalReviews.get(reviewId);
  187. global.asins.add(target.asin);
  188. if (target.marketplaceId) global.marketplaceIds.add(target.marketplaceId);
  189. global.domains.add(target.domain);
  190. global.shopIds.add(target.shopId);
  191. global.targets.push(target);
  192. const associationKey = `${target.domain}:${target.asin}:${reviewId}`;
  193. associations.set(associationKey, {
  194. recordKey: associationKey,
  195. reviewId,
  196. asin: target.asin,
  197. marketplaceId: target.marketplaceId,
  198. domain: target.domain,
  199. source: 'sorftime-product-reviews-year',
  200. capturedAt: { __type: 'Date', iso: new Date().toISOString() },
  201. shop: shopPointer(target.shopId),
  202. });
  203. }
  204. runStats.requests += targetStats.requests;
  205. runStats.unresolvedSaturatedBuckets.push(...targetStats.unresolvedSaturatedBuckets);
  206. runStats.byTarget.push({ ...target, reviews: collected.size, withContent, requests: targetStats.requests });
  207. log('review_target_completed', { ...target, reviews: collected.size, withContent, requests: targetStats.requests });
  208. } catch (error) {
  209. runStats.failures.push({ ...target, error: error.message });
  210. log('review_target_failed', { ...target, error: error.message });
  211. }
  212. await sleep(500);
  213. }
  214. const capturedAt = { __type: 'Date', iso: new Date().toISOString() };
  215. const reviewRequests = [];
  216. for (const [reviewId, global] of globalReviews) {
  217. const existing = existingByReviewId.get(reviewId);
  218. const targetsForReview = global.targets.slice().sort((a, b) => `${a.domain}:${a.asin}`.localeCompare(`${b.domain}:${b.asin}`));
  219. const primary = targetsForReview.find(target => target.asin === existing?.asin) || targetsForReview[0];
  220. const review = global.review;
  221. const content = String(review.Content || review.ReviewContent || review.Body || '');
  222. reviewRequests.push({
  223. method: existing ? 'PUT' : 'POST',
  224. path: existing ? `/parse/classes/SorftimeReviews/${existing.objectId}` : '/parse/classes/SorftimeReviews',
  225. body: {
  226. recordKey: reviewId,
  227. reviewId,
  228. asin: primary.asin,
  229. asins: [...global.asins].sort(),
  230. marketplaceId: primary.marketplaceId,
  231. marketplaceIds: [...global.marketplaceIds].sort(),
  232. domain: primary.domain,
  233. domains: [...global.domains].sort((a, b) => a - b),
  234. shopIds: [...global.shopIds].sort(),
  235. source: 'sorftime-product-reviews-year',
  236. star: Number(review.Star || review.Rating || 0),
  237. rating: Number(review.Star || review.Rating || 0),
  238. title: review.Title || review.ReviewTitle || '',
  239. content,
  240. reviewDate: reviewDateOf(review),
  241. reviewPublishedAt: parseReviewDate(reviewDateOf(review)),
  242. author: review.ProfileName || review.Author || '',
  243. verified: Boolean(review.OnlyPurchase || review.VerifiedPurchase),
  244. capturedAt,
  245. shop: shopPointer(primary.shopId),
  246. rawData: review,
  247. },
  248. });
  249. if (existing) runStats.updated++; else runStats.inserted++;
  250. if (content.trim()) runStats.withContent++;
  251. }
  252. await saveBatch(reviewRequests);
  253. const existingAssociations = await allRows('SorftimeReviewProduct', 'objectId,recordKey');
  254. const associationByKey = new Map(existingAssociations.filter(row => row.recordKey).map(row => [row.recordKey, row.objectId]));
  255. const associationRequests = [...associations.values()].map(body => {
  256. const objectId = associationByKey.get(body.recordKey);
  257. return {
  258. method: objectId ? 'PUT' : 'POST',
  259. path: objectId ? `/parse/classes/SorftimeReviewProduct/${objectId}` : '/parse/classes/SorftimeReviewProduct',
  260. body,
  261. };
  262. });
  263. await saveBatch(associationRequests);
  264. runStats.fetchedUnique = globalReviews.size;
  265. runStats.associations = associations.size;
  266. runStats.finishedAt = new Date().toISOString();
  267. console.log(JSON.stringify({ event: 'review_collection_finished', ...runStats }));