refresh-competitor-listings.ts 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166
  1. import 'dotenv/config';
  2. import { resolve } from 'node:path';
  3. import { pathToFileURL } from 'node:url';
  4. import { loadConfig } from '../src/config/env.js';
  5. import { ParseRestClient } from '../src/db/parse-rest.client.js';
  6. import { ApiError } from '../src/http/api-error.js';
  7. import { CompetitorListingMonitorService } from '../src/modules/competitor-listing-monitor/competitor-listing-monitor.service.js';
  8. import type { CompetitorListingRefreshRun } from '../src/modules/competitor-listing-monitor/domain.js';
  9. import { ParseRestCompetitorListingMonitorRepository } from '../src/modules/competitor-listing-monitor/repositories/parse-rest-competitor-listing-monitor.repository.js';
  10. import { FmodeVocEcommerceClient } from '../src/modules/domestic-voc/upstream/fmode-client.js';
  11. export const COMPETITOR_LISTING_REFRESH_EXIT = {
  12. completed: 0,
  13. failed: 1,
  14. partial: 2,
  15. alreadyRunning: 3,
  16. } as const;
  17. type RefreshCliService = Pick<CompetitorListingMonitorService, 'startRefresh' | 'getRun'>;
  18. export interface CompetitorListingRefreshCliOptions {
  19. service: RefreshCliService;
  20. workspaceId: string;
  21. pollMs?: number;
  22. timeoutMs?: number;
  23. sleep?: (milliseconds: number) => Promise<void>;
  24. write?: (line: string) => void;
  25. writeError?: (line: string) => void;
  26. }
  27. export async function runCompetitorListingRefreshCli(
  28. options: CompetitorListingRefreshCliOptions,
  29. ): Promise<number> {
  30. const pollMs = positiveInteger(options.pollMs ?? 2_000, 'pollMs');
  31. const timeoutMs = positiveInteger(options.timeoutMs ?? 6 * 60 * 60 * 1_000, 'timeoutMs');
  32. const sleep = options.sleep ?? ((milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds)));
  33. const write = options.write ?? console.log;
  34. const writeError = options.writeError ?? console.error;
  35. let queued: CompetitorListingRefreshRun;
  36. try {
  37. queued = await options.service.startRefresh(options.workspaceId, 'jd', 'scheduled');
  38. } catch (error) {
  39. if (error instanceof ApiError && error.code === 'competitor_listing_refresh_running') {
  40. writeError(JSON.stringify({
  41. status: 'already_running',
  42. workspaceId: options.workspaceId,
  43. error: error.code,
  44. }));
  45. return COMPETITOR_LISTING_REFRESH_EXIT.alreadyRunning;
  46. }
  47. writeError(JSON.stringify({
  48. status: 'failed',
  49. workspaceId: options.workspaceId,
  50. error: 'competitor_listing_refresh_start_failed',
  51. }));
  52. return COMPETITOR_LISTING_REFRESH_EXIT.failed;
  53. }
  54. write(JSON.stringify({
  55. status: queued.status,
  56. workspaceId: queued.workspaceId,
  57. runId: queued.id,
  58. total: queued.total,
  59. requestedAt: queued.requestedAt,
  60. }));
  61. const deadline = Date.now() + timeoutMs;
  62. for (;;) {
  63. let run: CompetitorListingRefreshRun | null;
  64. try {
  65. run = await options.service.getRun(options.workspaceId, queued.id);
  66. } catch {
  67. writeError(JSON.stringify({
  68. status: 'failed',
  69. workspaceId: options.workspaceId,
  70. runId: queued.id,
  71. error: 'competitor_listing_refresh_status_failed',
  72. }));
  73. return COMPETITOR_LISTING_REFRESH_EXIT.failed;
  74. }
  75. if (!run) {
  76. writeError(JSON.stringify({
  77. status: 'failed',
  78. workspaceId: options.workspaceId,
  79. runId: queued.id,
  80. error: 'competitor_listing_refresh_run_not_found',
  81. }));
  82. return COMPETITOR_LISTING_REFRESH_EXIT.failed;
  83. }
  84. if (run.status === 'completed' || run.status === 'partial' || run.status === 'failed') {
  85. const output = {
  86. status: run.status,
  87. workspaceId: run.workspaceId,
  88. runId: run.id,
  89. total: run.total,
  90. completed: run.completed,
  91. baseline: run.baseline,
  92. unchanged: run.unchanged,
  93. changed: run.changed,
  94. failed: run.failed,
  95. requestedAt: run.requestedAt,
  96. startedAt: run.startedAt,
  97. completedAt: run.completedAt,
  98. };
  99. const line = JSON.stringify(output);
  100. if (run.status === 'failed') writeError(line);
  101. else write(line);
  102. return run.status === 'completed'
  103. ? COMPETITOR_LISTING_REFRESH_EXIT.completed
  104. : run.status === 'partial'
  105. ? COMPETITOR_LISTING_REFRESH_EXIT.partial
  106. : COMPETITOR_LISTING_REFRESH_EXIT.failed;
  107. }
  108. if (Date.now() >= deadline) {
  109. writeError(JSON.stringify({
  110. status: 'failed',
  111. workspaceId: options.workspaceId,
  112. runId: queued.id,
  113. error: 'competitor_listing_refresh_timeout',
  114. }));
  115. return COMPETITOR_LISTING_REFRESH_EXIT.failed;
  116. }
  117. await sleep(pollMs);
  118. }
  119. }
  120. export async function main(args = process.argv.slice(2)): Promise<number> {
  121. const config = loadConfig();
  122. if (config.storageDriver !== 'parse_rest') {
  123. console.error(JSON.stringify({ status: 'failed', error: 'parse_rest_storage_required' }));
  124. return COMPETITOR_LISTING_REFRESH_EXIT.failed;
  125. }
  126. const workspaceId = args.find((argument) => !argument.startsWith('--')) || config.auth.defaultWorkspaceId;
  127. const pollMs = readIntegerArgument(args, '--poll-ms=', 2_000);
  128. const timeoutMs = readIntegerArgument(args, '--timeout-ms=', 6 * 60 * 60 * 1_000);
  129. const client = new ParseRestClient({
  130. serverUrl: config.parse.serverUrl,
  131. appId: config.parse.appId,
  132. masterKey: config.parse.masterKey,
  133. timeoutMs: config.parse.timeoutMs,
  134. });
  135. const service = new CompetitorListingMonitorService(
  136. new ParseRestCompetitorListingMonitorRepository(client),
  137. new FmodeVocEcommerceClient(config.fmode),
  138. );
  139. return runCompetitorListingRefreshCli({ service, workspaceId, pollMs, timeoutMs });
  140. }
  141. function readIntegerArgument(args: string[], prefix: string, fallback: number): number {
  142. const raw = args.find((argument) => argument.startsWith(prefix))?.slice(prefix.length);
  143. return raw === undefined ? fallback : positiveInteger(Number(raw), prefix.slice(2, -1));
  144. }
  145. function positiveInteger(value: number, name: string): number {
  146. if (!Number.isSafeInteger(value) || value <= 0) throw new Error(`${name} must be a positive integer`);
  147. return value;
  148. }
  149. const entry = process.argv[1];
  150. if (entry && pathToFileURL(resolve(entry)).href === import.meta.url) {
  151. main().then((exitCode) => {
  152. process.exitCode = exitCode;
  153. }).catch(() => {
  154. console.error(JSON.stringify({ status: 'failed', error: 'competitor_listing_refresh_cli_failed' }));
  155. process.exitCode = COMPETITOR_LISTING_REFRESH_EXIT.failed;
  156. });
  157. }