import 'dotenv/config'; import { resolve } from 'node:path'; import { pathToFileURL } from 'node:url'; import { loadConfig } from '../src/config/env.js'; import { ParseRestClient } from '../src/db/parse-rest.client.js'; import { ApiError } from '../src/http/api-error.js'; import { CompetitorListingMonitorService } from '../src/modules/competitor-listing-monitor/competitor-listing-monitor.service.js'; import type { CompetitorListingRefreshRun } from '../src/modules/competitor-listing-monitor/domain.js'; import { ParseRestCompetitorListingMonitorRepository } from '../src/modules/competitor-listing-monitor/repositories/parse-rest-competitor-listing-monitor.repository.js'; import { FmodeVocEcommerceClient } from '../src/modules/domestic-voc/upstream/fmode-client.js'; export const COMPETITOR_LISTING_REFRESH_EXIT = { completed: 0, failed: 1, partial: 2, alreadyRunning: 3, } as const; type RefreshCliService = Pick; export interface CompetitorListingRefreshCliOptions { service: RefreshCliService; workspaceId: string; pollMs?: number; timeoutMs?: number; sleep?: (milliseconds: number) => Promise; write?: (line: string) => void; writeError?: (line: string) => void; } export async function runCompetitorListingRefreshCli( options: CompetitorListingRefreshCliOptions, ): Promise { const pollMs = positiveInteger(options.pollMs ?? 2_000, 'pollMs'); const timeoutMs = positiveInteger(options.timeoutMs ?? 6 * 60 * 60 * 1_000, 'timeoutMs'); const sleep = options.sleep ?? ((milliseconds) => new Promise((resolve) => setTimeout(resolve, milliseconds))); const write = options.write ?? console.log; const writeError = options.writeError ?? console.error; let queued: CompetitorListingRefreshRun; try { queued = await options.service.startRefresh(options.workspaceId, 'jd', 'scheduled'); } catch (error) { if (error instanceof ApiError && error.code === 'competitor_listing_refresh_running') { writeError(JSON.stringify({ status: 'already_running', workspaceId: options.workspaceId, error: error.code, })); return COMPETITOR_LISTING_REFRESH_EXIT.alreadyRunning; } writeError(JSON.stringify({ status: 'failed', workspaceId: options.workspaceId, error: 'competitor_listing_refresh_start_failed', })); return COMPETITOR_LISTING_REFRESH_EXIT.failed; } write(JSON.stringify({ status: queued.status, workspaceId: queued.workspaceId, runId: queued.id, total: queued.total, requestedAt: queued.requestedAt, })); const deadline = Date.now() + timeoutMs; for (;;) { let run: CompetitorListingRefreshRun | null; try { run = await options.service.getRun(options.workspaceId, queued.id); } catch { writeError(JSON.stringify({ status: 'failed', workspaceId: options.workspaceId, runId: queued.id, error: 'competitor_listing_refresh_status_failed', })); return COMPETITOR_LISTING_REFRESH_EXIT.failed; } if (!run) { writeError(JSON.stringify({ status: 'failed', workspaceId: options.workspaceId, runId: queued.id, error: 'competitor_listing_refresh_run_not_found', })); return COMPETITOR_LISTING_REFRESH_EXIT.failed; } if (run.status === 'completed' || run.status === 'partial' || run.status === 'failed') { const output = { status: run.status, workspaceId: run.workspaceId, runId: run.id, total: run.total, completed: run.completed, baseline: run.baseline, unchanged: run.unchanged, changed: run.changed, failed: run.failed, requestedAt: run.requestedAt, startedAt: run.startedAt, completedAt: run.completedAt, }; const line = JSON.stringify(output); if (run.status === 'failed') writeError(line); else write(line); return run.status === 'completed' ? COMPETITOR_LISTING_REFRESH_EXIT.completed : run.status === 'partial' ? COMPETITOR_LISTING_REFRESH_EXIT.partial : COMPETITOR_LISTING_REFRESH_EXIT.failed; } if (Date.now() >= deadline) { writeError(JSON.stringify({ status: 'failed', workspaceId: options.workspaceId, runId: queued.id, error: 'competitor_listing_refresh_timeout', })); return COMPETITOR_LISTING_REFRESH_EXIT.failed; } await sleep(pollMs); } } export async function main(args = process.argv.slice(2)): Promise { const config = loadConfig(); if (config.storageDriver !== 'parse_rest') { console.error(JSON.stringify({ status: 'failed', error: 'parse_rest_storage_required' })); return COMPETITOR_LISTING_REFRESH_EXIT.failed; } const workspaceId = args.find((argument) => !argument.startsWith('--')) || config.auth.defaultWorkspaceId; const pollMs = readIntegerArgument(args, '--poll-ms=', 2_000); const timeoutMs = readIntegerArgument(args, '--timeout-ms=', 6 * 60 * 60 * 1_000); const client = new ParseRestClient({ serverUrl: config.parse.serverUrl, appId: config.parse.appId, masterKey: config.parse.masterKey, timeoutMs: config.parse.timeoutMs, }); const service = new CompetitorListingMonitorService( new ParseRestCompetitorListingMonitorRepository(client), new FmodeVocEcommerceClient(config.fmode), ); return runCompetitorListingRefreshCli({ service, workspaceId, pollMs, timeoutMs }); } function readIntegerArgument(args: string[], prefix: string, fallback: number): number { const raw = args.find((argument) => argument.startsWith(prefix))?.slice(prefix.length); return raw === undefined ? fallback : positiveInteger(Number(raw), prefix.slice(2, -1)); } function positiveInteger(value: number, name: string): number { if (!Number.isSafeInteger(value) || value <= 0) throw new Error(`${name} must be a positive integer`); return value; } const entry = process.argv[1]; if (entry && pathToFileURL(resolve(entry)).href === import.meta.url) { main().then((exitCode) => { process.exitCode = exitCode; }).catch(() => { console.error(JSON.stringify({ status: 'failed', error: 'competitor_listing_refresh_cli_failed' })); process.exitCode = COMPETITOR_LISTING_REFRESH_EXIT.failed; }); }