sync-jd-listing-reviews.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351
  1. import 'dotenv/config';
  2. import { createHash } from 'node:crypto';
  3. import { mkdir, readFile, writeFile } from 'node:fs/promises';
  4. import { dirname, resolve } from 'node:path';
  5. import { z } from 'zod';
  6. import { loadConfig } from '../src/config/env.js';
  7. import { ParseRestClient, parseDate, type ParseObject } from '../src/db/parse-rest.client.js';
  8. import { VOC_PARSE_CLASSES } from '../src/db/parse-rest.schema.js';
  9. import { makeReviewKey } from '../src/modules/domestic-voc/domain/identity.js';
  10. import type { ListingSourceSnapshot } from '../src/modules/listing-ai/domain.js';
  11. type JsonRecord = Record<string, unknown>;
  12. interface StoredSource extends ParseObject { productId?: string; payload: ListingSourceSnapshot }
  13. interface StoredReview extends ParseObject { naturalKey: string; productId: string; reviewDate?: unknown }
  14. interface Checkpoint {
  15. workspaceId: string;
  16. cohort: string;
  17. nextPage: number;
  18. totalItems: number;
  19. totalPages: number;
  20. pagesProcessed: number;
  21. commentsScanned: number;
  22. commentsMatched: number;
  23. commentsCreated: number;
  24. commentsUpdated: number;
  25. commentsDuplicate: number;
  26. commentsUnmatched: number;
  27. matchedProducts: string[];
  28. startedAt: string;
  29. updatedAt: string;
  30. status: 'running' | 'completed';
  31. phase?: 'windowed';
  32. historyStart?: string;
  33. historyEnd?: string;
  34. windowStart?: string;
  35. windowEnd?: string;
  36. windowsProcessed?: number;
  37. }
  38. const ACTIVE_COHORT = 'listing-jd-v3-formal-625';
  39. const METHOD = 'jingdong.pop.PopCommentJsfService.getVenderCommentsForJos';
  40. const PAGE_SIZE = 50;
  41. const DEFAULT_HISTORY_START = '2026-06-01 00:00:00';
  42. const DEFAULT_HISTORY_END = '2026-08-29 00:00:00';
  43. const args = new Map(process.argv.slice(2).map((arg) => {
  44. const [key, ...rest] = arg.split('=');
  45. return [key!, rest.join('=') || 'true'];
  46. }));
  47. const workspaceId = args.get('--workspace') ?? process.env.SAAS_DEFAULT_WORKSPACE_ID ?? 'demashi';
  48. const checkpointPath = resolve(args.get('--checkpoint') ?? 'logs/jd-listing-review-sync-2026-06-01-to-2026-08-28.json');
  49. const resume = args.get('--resume') === 'true';
  50. const delayMs = Math.max(100, Number(args.get('--delay-ms') ?? 200));
  51. const maxPages = Math.max(1, Number(args.get('--max-pages') ?? Number.POSITIVE_INFINITY));
  52. const historyStartArg = args.get('--history-start') ?? DEFAULT_HISTORY_START;
  53. const historyEndArg = args.get('--history-end') ?? DEFAULT_HISTORY_END;
  54. const windowDays = Math.max(1, Number(args.get('--window-days') ?? 7));
  55. const pruneOutsideRange = args.get('--prune-outside-range') !== 'false';
  56. const sourceSchema = z.object({
  57. JD_SOURCE_PARSE_URL: z.url(),
  58. JD_SOURCE_PARSE_APP_ID: z.string().min(1),
  59. JD_SOURCE_PARSE_MASTER_KEY: z.string().min(1),
  60. JD_APP_KEY: z.string().min(1),
  61. JD_APP_SECRET: z.string().min(1),
  62. });
  63. async function main(): Promise<void> {
  64. const sourceConfig = sourceSchema.parse(process.env);
  65. const config = loadConfig();
  66. const target = new ParseRestClient({
  67. serverUrl: config.parse.serverUrl,
  68. appId: config.parse.appId,
  69. masterKey: config.parse.masterKey,
  70. timeoutMs: config.parse.timeoutMs,
  71. });
  72. const authSource = new ParseRestClient({
  73. serverUrl: sourceConfig.JD_SOURCE_PARSE_URL,
  74. appId: sourceConfig.JD_SOURCE_PARSE_APP_ID,
  75. masterKey: sourceConfig.JD_SOURCE_PARSE_MASTER_KEY,
  76. timeoutMs: config.jdListing.timeoutMs,
  77. });
  78. const authRow = (await authSource.find<JsonRecord>('EcomAuth', {
  79. where: { platform: 'jd', type: 'access_token' }, order: '-createdAt', limit: 1,
  80. })).results[0];
  81. const authData = record(authRow?.['data']);
  82. const accessToken = text(authData['access_token']);
  83. if (!accessToken) throw new Error('jd_authorization_missing');
  84. const currentSources = await target.findAll<StoredSource>(VOC_PARSE_CLASSES.listingSourceSnapshot, {
  85. workspaceId, platform: 'jd', isCurrent: true, catalogIncluded: true, catalogCohort: ACTIVE_COHORT,
  86. });
  87. const cohortProductIds = new Set(currentSources.map((row) => row.payload.productId));
  88. if (cohortProductIds.size !== 625) throw new Error(`listing_review_cohort_count_mismatch:${cohortProductIds.size}:625`);
  89. const allSources = await target.findAll<StoredSource>(VOC_PARSE_CLASSES.listingSourceSnapshot, { workspaceId, platform: 'jd' });
  90. const skuToProduct = buildSkuMap(allSources, cohortProductIds);
  91. if (!skuToProduct.size) throw new Error('listing_review_sku_map_empty');
  92. const conflicts = findSkuConflicts(allSources, cohortProductIds);
  93. if (conflicts.length) throw new Error(`listing_review_sku_conflicts:${conflicts.slice(0, 5).join(',')}`);
  94. let existingReviews = await target.findAll<StoredReview>(VOC_PARSE_CLASSES.review, { workspaceId, platform: 'jd' });
  95. if (pruneOutsideRange) {
  96. const outsideRange = existingReviews.filter((review) => cohortProductIds.has(review.productId)
  97. && !isWithinJdRange(review.reviewDate, historyStartArg, historyEndArg));
  98. const deleteRequests = outsideRange.map((review) => ({
  99. method: 'DELETE' as const,
  100. path: `/classes/${VOC_PARSE_CLASSES.review}/${review.objectId}`,
  101. }));
  102. for (let index = 0; index < deleteRequests.length; index += 50) {
  103. await writeFmodeBatch(target, deleteRequests.slice(index, index + 50));
  104. }
  105. existingReviews = existingReviews.filter((review) => !outsideRange.includes(review));
  106. console.log(JSON.stringify({ event: 'listing_review_prune', deleted: outsideRange.length, historyStart: historyStartArg, historyEndExclusive: historyEndArg }));
  107. }
  108. const existingByNaturalKey = new Map(existingReviews.map((row) => [row.naturalKey, row]));
  109. const seen = new Set(existingByNaturalKey.keys());
  110. let checkpoint = resume ? await readCheckpoint(checkpointPath) : null;
  111. if (checkpoint?.status === 'completed') {
  112. console.log(JSON.stringify({ mode: 'already_completed', checkpointPath, ...checkpoint }, null, 2));
  113. return;
  114. }
  115. const now = new Date().toISOString();
  116. checkpoint ??= {
  117. workspaceId, cohort: ACTIVE_COHORT, nextPage: 1, totalItems: 0, totalPages: 0,
  118. pagesProcessed: 0, commentsScanned: 0, commentsMatched: 0, commentsCreated: 0,
  119. commentsUpdated: 0, commentsDuplicate: 0, commentsUnmatched: 0, matchedProducts: [],
  120. startedAt: now, updatedAt: now, status: 'running',
  121. };
  122. if (checkpoint.workspaceId !== workspaceId || checkpoint.cohort !== ACTIVE_COHORT) {
  123. throw new Error('listing_review_checkpoint_scope_mismatch');
  124. }
  125. if (checkpoint.phase === 'windowed'
  126. && (checkpoint.historyStart !== normalizeJdDateTime(historyStartArg)
  127. || checkpoint.historyEnd !== normalizeJdDateTime(historyEndArg))) {
  128. throw new Error(`listing_review_checkpoint_range_mismatch:${checkpoint.historyStart}:${checkpoint.historyEnd}`);
  129. }
  130. const matchedProducts = new Set(checkpoint.matchedProducts);
  131. const client = new JdJosReviewClient(sourceConfig.JD_APP_KEY, sourceConfig.JD_APP_SECRET, accessToken);
  132. if (!checkpoint.totalItems) {
  133. const metadata = await client.page(1, PAGE_SIZE);
  134. checkpoint.totalItems = metadata.totalItem;
  135. checkpoint.totalPages = Math.ceil(metadata.totalItem / PAGE_SIZE);
  136. }
  137. if (checkpoint.phase !== 'windowed') {
  138. const oldestStoredReview = existingReviews
  139. .filter((review) => cohortProductIds.has(review.productId))
  140. .map((review) => storedDateIso(review.reviewDate))
  141. .filter((value): value is string => Boolean(value))
  142. .sort()[0];
  143. const historyEnd = normalizeJdDateTime(historyEndArg || oldestStoredReview || jdTomorrow());
  144. const historyStart = normalizeJdDateTime(historyStartArg);
  145. checkpoint.phase = 'windowed';
  146. checkpoint.historyStart = historyStart;
  147. checkpoint.historyEnd = historyEnd;
  148. checkpoint.windowStart = historyStart;
  149. checkpoint.windowEnd = minJdDateTime(addJdDays(historyStart, windowDays), historyEnd);
  150. checkpoint.windowsProcessed = 0;
  151. checkpoint.nextPage = 1;
  152. checkpoint.updatedAt = new Date().toISOString();
  153. await persistCheckpoint(checkpointPath, checkpoint);
  154. console.log(JSON.stringify({ event: 'listing_review_window_migration', historyStart, historyEnd, retainedRecentPages: checkpoint.pagesProcessed }));
  155. }
  156. let pagesThisRun = 0;
  157. while (checkpoint.windowStart && checkpoint.windowEnd && checkpoint.historyEnd
  158. && jdEpoch(checkpoint.windowStart) < jdEpoch(checkpoint.historyEnd)
  159. && pagesThisRun < maxPages) {
  160. const page = checkpoint.nextPage;
  161. const response = await client.page(page, PAGE_SIZE, {
  162. beginTime: checkpoint.windowStart,
  163. endTime: checkpoint.windowEnd,
  164. });
  165. const windowPages = Math.ceil(response.totalItem / PAGE_SIZE);
  166. if (windowPages > 200) throw new Error(`listing_review_window_too_large:${checkpoint.windowStart}:${checkpoint.windowEnd}:${response.totalItem}`);
  167. if (!response.comments.length && page <= windowPages) throw new Error(`listing_review_empty_page:${page}/${windowPages}`);
  168. const requests: Array<{ method: 'POST' | 'PUT'; path: string; body: unknown }> = [];
  169. for (const comment of response.comments) {
  170. checkpoint.commentsScanned += 1;
  171. const skuId = text(comment['skuid'] ?? comment['skuId']);
  172. const productId = skuToProduct.get(skuId);
  173. if (!productId) { checkpoint.commentsUnmatched += 1; continue; }
  174. const content = text(comment['content']);
  175. if (!content) { checkpoint.commentsUnmatched += 1; continue; }
  176. const reviewId = text(comment['commentId']);
  177. const reviewDate = dateIso(comment['creationTime']);
  178. const reviewKey = makeReviewKey({ platform: 'jd', productId, reviewId, content, reviewDate });
  179. const naturalKey = [workspaceId, 'jd', reviewKey].map(encodeURIComponent).join('|');
  180. const existing = existingByNaturalKey.get(naturalKey);
  181. if (!existing && seen.has(naturalKey)) { checkpoint.commentsDuplicate += 1; continue; }
  182. const body = {
  183. naturalKey, workspaceId, platform: 'jd', productId, sourceReviewId: reviewId || null,
  184. reviewKey, rating: rating(comment['score']), content,
  185. reviewDate: reviewDate ? parseDate(reviewDate) : null,
  186. rawPayload: sanitizeComment(comment),
  187. };
  188. if (existing) {
  189. requests.push({ method: 'PUT', path: `/classes/${VOC_PARSE_CLASSES.review}/${existing.objectId}`, body });
  190. checkpoint.commentsUpdated += 1;
  191. } else {
  192. requests.push({ method: 'POST', path: `/classes/${VOC_PARSE_CLASSES.review}`, body });
  193. checkpoint.commentsCreated += 1;
  194. }
  195. seen.add(naturalKey);
  196. matchedProducts.add(productId);
  197. checkpoint.commentsMatched += 1;
  198. }
  199. for (let index = 0; index < requests.length; index += 50) {
  200. await writeFmodeBatch(target, requests.slice(index, index + 50));
  201. }
  202. checkpoint.pagesProcessed += 1;
  203. pagesThisRun += 1;
  204. const windowFinished = page >= Math.max(1, windowPages);
  205. if (windowFinished) {
  206. checkpoint.windowsProcessed = (checkpoint.windowsProcessed ?? 0) + 1;
  207. checkpoint.windowStart = checkpoint.windowEnd;
  208. checkpoint.windowEnd = minJdDateTime(addJdDays(checkpoint.windowStart, windowDays), checkpoint.historyEnd);
  209. checkpoint.nextPage = 1;
  210. } else {
  211. checkpoint.nextPage = page + 1;
  212. }
  213. checkpoint.matchedProducts = [...matchedProducts].sort();
  214. checkpoint.updatedAt = new Date().toISOString();
  215. await persistCheckpoint(checkpointPath, checkpoint);
  216. if (checkpoint.pagesProcessed % 25 === 0 || windowFinished && (checkpoint.windowsProcessed ?? 0) % 25 === 0) {
  217. console.log(JSON.stringify({ event: 'listing_review_progress', window: checkpoint.windowsProcessed, windowStart: checkpoint.windowStart, page, windowPages, scanned: checkpoint.commentsScanned, matched: checkpoint.commentsMatched, created: checkpoint.commentsCreated, updated: checkpoint.commentsUpdated, matchedProducts: matchedProducts.size }));
  218. }
  219. await wait(delayMs);
  220. }
  221. if (checkpoint.windowStart && checkpoint.historyEnd && jdEpoch(checkpoint.windowStart) >= jdEpoch(checkpoint.historyEnd)) checkpoint.status = 'completed';
  222. checkpoint.updatedAt = new Date().toISOString();
  223. checkpoint.matchedProducts = [...matchedProducts].sort();
  224. await persistCheckpoint(checkpointPath, checkpoint);
  225. console.log(JSON.stringify({
  226. mode: checkpoint.status, workspaceId, cohortProducts: cohortProductIds.size, mappedSkus: skuToProduct.size,
  227. totalItems: checkpoint.totalItems, totalPages: checkpoint.totalPages, pagesProcessed: checkpoint.pagesProcessed,
  228. commentsScanned: checkpoint.commentsScanned, commentsMatched: checkpoint.commentsMatched,
  229. commentsCreated: checkpoint.commentsCreated, commentsUpdated: checkpoint.commentsUpdated,
  230. commentsDuplicate: checkpoint.commentsDuplicate, commentsUnmatched: checkpoint.commentsUnmatched,
  231. matchedProducts: matchedProducts.size, checkpointPath,
  232. }, null, 2));
  233. }
  234. class JdJosReviewClient {
  235. constructor(private readonly appKey: string, private readonly appSecret: string, private readonly accessToken: string) {}
  236. async page(page: number, pageSize: number, filters: JsonRecord = {}): Promise<{ totalItem: number; comments: JsonRecord[] }> {
  237. for (let attempt = 0; attempt < 6; attempt += 1) {
  238. try {
  239. const params: Record<string, string> = {
  240. method: METHOD, access_token: this.accessToken, app_key: this.appKey,
  241. timestamp: jdTime(), v: '2.0', sign_method: 'md5',
  242. '360buy_param_json': JSON.stringify({ page, pageSize, ...filters }),
  243. };
  244. const plain = Object.keys(params).sort().map((key) => `${key}${params[key] ?? ''}`).join('');
  245. params['sign'] = createHash('md5').update(`${this.appSecret}${plain}${this.appSecret}`).digest('hex').toUpperCase();
  246. const response = await fetch(`https://api.jd.com/routerjson?${new URLSearchParams(params)}`, { signal: AbortSignal.timeout(30_000) });
  247. const body = await response.json() as JsonRecord;
  248. if (!response.ok) throw new Error(`jd_review_http_${response.status}`);
  249. const root = record(body[Object.keys(body)[0] ?? '']);
  250. if (text(root['code']) !== '0' || text(root['resultCode']) !== '200') {
  251. const detail = text(root['resultMessage'] ?? root['resultMsg'] ?? root['message'] ?? root['msg'] ?? root['errorMessage']);
  252. const diagnostic = JSON.stringify(Object.fromEntries(Object.entries(root).filter(([key]) => key !== 'comments'))).slice(0, 1_000);
  253. throw new Error(`jd_review_api_${text(root['code'])}_${text(root['resultCode'])}${detail ? `:${detail}` : ''}:${diagnostic}`);
  254. }
  255. return { totalItem: Math.max(0, Number(root['totalItem']) || 0), comments: list(root['comments']).map(record) };
  256. } catch (error) {
  257. if (attempt >= 5) throw error;
  258. await wait(Math.min(8_000, 500 * 2 ** attempt));
  259. }
  260. }
  261. throw new Error('jd_review_retry_exhausted');
  262. }
  263. }
  264. function buildSkuMap(rows: StoredSource[], cohortProductIds: Set<string>): Map<string, string> {
  265. const output = new Map<string, string>();
  266. for (const row of rows) {
  267. const source = row.payload;
  268. if (!cohortProductIds.has(source.productId)) continue;
  269. output.set(source.productId, source.productId);
  270. for (const sku of source.skus ?? []) if (sku.skuId) output.set(String(sku.skuId), source.productId);
  271. }
  272. return output;
  273. }
  274. function findSkuConflicts(rows: StoredSource[], cohortProductIds: Set<string>): string[] {
  275. const owners = new Map<string, string>(); const conflicts = new Set<string>();
  276. for (const row of rows) {
  277. const source = row.payload; if (!cohortProductIds.has(source.productId)) continue;
  278. for (const sku of source.skus ?? []) {
  279. const skuId = String(sku.skuId || ''); if (!skuId) continue;
  280. const previous = owners.get(skuId); if (previous && previous !== source.productId) conflicts.add(skuId); else owners.set(skuId, source.productId);
  281. }
  282. }
  283. return [...conflicts];
  284. }
  285. function sanitizeComment(comment: JsonRecord): JsonRecord {
  286. const keys = ['commentId', 'creationTime', 'content', 'skuName', 'score', 'skuid', 'images', 'isVenderReply', 'replyCount', 'usefulCount', 'skuImage', 'status', 'videos'];
  287. return Object.fromEntries(keys.filter((key) => comment[key] !== undefined).map((key) => [key, comment[key]]));
  288. }
  289. async function writeFmodeBatch(
  290. client: ParseRestClient,
  291. requests: Array<{ method: 'POST' | 'PUT' | 'DELETE'; path: string; body?: unknown }>,
  292. ): Promise<void> {
  293. // Fmode exposes Parse at /backend/{appId}/data, but its batch router expects
  294. // nested request paths to start at /data rather than repeat the full mount.
  295. const results = await client.request<Array<{ error?: { code?: number; error?: string } }>>('/batch', {
  296. method: 'POST',
  297. body: {
  298. requests: requests.map((request) => ({ ...request, path: `/data${request.path}` })),
  299. },
  300. });
  301. const failure = results.find((result) => result.error)?.error;
  302. if (failure) throw new Error(`listing_review_batch_${failure.code ?? 'unknown'}:${failure.error ?? 'failed'}`);
  303. }
  304. function storedDateIso(value: unknown): string | null {
  305. if (typeof value === 'string') return dateIso(value);
  306. const iso = text(record(value)['iso']);
  307. return iso ? dateIso(iso) : null;
  308. }
  309. function isWithinJdRange(value: unknown, start: string, endExclusive: string): boolean {
  310. const iso = storedDateIso(value);
  311. if (!iso) return false;
  312. const epoch = new Date(iso).valueOf();
  313. return epoch >= jdEpoch(normalizeJdDateTime(start)) && epoch < jdEpoch(normalizeJdDateTime(endExclusive));
  314. }
  315. function normalizeJdDateTime(value: string): string {
  316. if (/^\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}$/.test(value)) return value;
  317. const date = new Date(value);
  318. if (Number.isNaN(date.valueOf())) throw new Error(`listing_review_invalid_date:${value}`);
  319. return formatJdDateTime(date);
  320. }
  321. function jdEpoch(value: string): number {
  322. const epoch = new Date(`${value.replace(' ', 'T')}+08:00`).valueOf();
  323. if (Number.isNaN(epoch)) throw new Error(`listing_review_invalid_jd_date:${value}`);
  324. return epoch;
  325. }
  326. function addJdDays(value: string, days: number): string { return formatJdDateTime(new Date(jdEpoch(value) + days * 86_400_000)); }
  327. function minJdDateTime(left: string, right: string): string { return jdEpoch(left) <= jdEpoch(right) ? left : right; }
  328. function jdTomorrow(): string { return `${formatJdDateTime(new Date(Date.now() + 86_400_000)).slice(0, 10)} 00:00:00`; }
  329. function rating(value: unknown): number { const parsed = Number(value); return Number.isFinite(parsed) && parsed >= 0 && parsed <= 5 ? parsed : 0; }
  330. function dateIso(value: unknown): string | null { const parsed = Number(value); const date = Number.isFinite(parsed) ? new Date(parsed < 10_000_000_000 ? parsed * 1_000 : parsed) : new Date(String(value ?? '')); return Number.isNaN(date.valueOf()) ? null : date.toISOString(); }
  331. function record(value: unknown): JsonRecord { return value && typeof value === 'object' && !Array.isArray(value) ? value as JsonRecord : {}; }
  332. function list(value: unknown): unknown[] { return Array.isArray(value) ? value : []; }
  333. function text(value: unknown): string { return value === null || value === undefined ? '' : String(value).trim(); }
  334. function wait(milliseconds: number): Promise<void> { return new Promise((resolve) => setTimeout(resolve, milliseconds)); }
  335. function formatJdDateTime(date: Date): string { const parts = new Intl.DateTimeFormat('sv-SE', { timeZone: 'Asia/Shanghai', year: 'numeric', month: '2-digit', day: '2-digit', hour: '2-digit', minute: '2-digit', second: '2-digit', hourCycle: 'h23' }).formatToParts(date); const get = (type: Intl.DateTimeFormatPartTypes) => parts.find((part) => part.type === type)?.value ?? ''; return `${get('year')}-${get('month')}-${get('day')} ${get('hour')}:${get('minute')}:${get('second')}`; }
  336. function jdTime(): string { return formatJdDateTime(new Date()); }
  337. async function readCheckpoint(path: string): Promise<Checkpoint | null> { try { return JSON.parse(await readFile(path, 'utf8')) as Checkpoint; } catch { return null; } }
  338. async function persistCheckpoint(path: string, checkpoint: Checkpoint): Promise<void> { await mkdir(dirname(path), { recursive: true }); await writeFile(path, JSON.stringify(checkpoint, null, 2)); }
  339. main().catch((error) => { console.error(`[sync-jd-listing-reviews] ${error instanceof Error ? error.message : error}`); process.exitCode = 1; });