diff --git a/apps/api/src/integrations/carcatonline/carcatonline.client.ts b/apps/api/src/integrations/carcatonline/carcatonline.client.ts index b60d3ee..003ca22 100644 --- a/apps/api/src/integrations/carcatonline/carcatonline.client.ts +++ b/apps/api/src/integrations/carcatonline/carcatonline.client.ts @@ -225,7 +225,12 @@ export class CarcatonlineClient { for (let attempt = 0; attempt < 2; attempt += 1) { this.opts.onRequest?.(path); const res = await this.doFetch(`${this.apiUrl}${path}`, { - headers: { Authorization: `Bearer ${token.token}`, "User-Agent": USER_AGENT }, + // Turkish where the API has it (model lists); group/part names stay English. + headers: { + Authorization: `Bearer ${token.token}`, + "User-Agent": USER_AGENT, + "Accept-Language": "tr", + }, }); if (res.status === 429) throw new CarcatonlineRateLimitError(path); if (res.status === 401 && attempt === 0) { diff --git a/apps/api/src/integrations/carcatonline/carcatonline.matcher.spec.ts b/apps/api/src/integrations/carcatonline/carcatonline.matcher.spec.ts index 17b6b2a..cb44eac 100644 --- a/apps/api/src/integrations/carcatonline/carcatonline.matcher.spec.ts +++ b/apps/api/src/integrations/carcatonline/carcatonline.matcher.spec.ts @@ -52,6 +52,21 @@ describe("carcatonline model matcher", () => { expect(matchCarcatModel("W124", list)?.tier).toBe("code"); }); + it("bridges Turkish vs English/French region words and roman vs digit generations", () => { + const list = models([ + "ARKANA RUSSIA", + "DUSTER 3 / BIGSTER", + "KADJAR CHINE", + "CAPTUR II CHINA", + "KOLEOS 2 - CHINA", + ]); + expect(matchCarcatModel("ARKANA RUSYA", list)?.model.name).toBe("ARKANA RUSSIA"); + expect(matchCarcatModel("DUSTER III / BIGSTER", list)?.model.name).toBe("DUSTER 3 / BIGSTER"); + expect(matchCarcatModel("KADJAR ÇİN", list)?.model.name).toBe("KADJAR CHINE"); + expect(matchCarcatModel("CAPTUR II ÇİN", list)?.model.name).toBe("CAPTUR II CHINA"); + expect(matchCarcatModel("KOLEOS 2 - ÇİN", list)?.model.name).toBe("KOLEOS 2 - CHINA"); + }); + it("returns null when ambiguous or absent", () => { const list = models(["Mercedes-Benz W205 C-Class", "Mercedes-Benz W205 AMG"]); expect(matchCarcatModel("W205", list)).toBeNull(); diff --git a/apps/api/src/integrations/carcatonline/carcatonline.matcher.ts b/apps/api/src/integrations/carcatonline/carcatonline.matcher.ts index 2409f8e..d10c3f0 100644 --- a/apps/api/src/integrations/carcatonline/carcatonline.matcher.ts +++ b/apps/api/src/integrations/carcatonline/carcatonline.matcher.ts @@ -79,6 +79,35 @@ export function isNonVehicleModel(label: string): boolean { } /** Uppercase, diacritic-folded, punctuation collapsed, PL24 backslash escapes removed. */ +const ROMAN: Record = { + I: "1", + II: "2", + III: "3", + IV: "4", + V: "5", + VI: "6", + VII: "7", + VIII: "8", +}; + +/** Region words differ between our Turkish PL24 labels and carcatonline's English/French ones. */ +const TOKEN_SYNONYMS: Record = { + RUSSIA: "RUSYA", + RUSSIE: "RUSYA", + CHINA: "CIN", + CHINE: "CIN", + EUROPA: "EUROPE", + AVRUPA: "EUROPE", + TURKEY: "TURKIYE", + INDIA: "HINDISTAN", + BRAZIL: "BREZILYA", + KOREA: "KORE", +}; + +function canonicalToken(t: string): string { + return TOKEN_SYNONYMS[t] ?? ROMAN[t] ?? t; +} + export function normalizeLabel(s: string | null | undefined): string { return (s ?? "") .toUpperCase() @@ -86,7 +115,11 @@ export function normalizeLabel(s: string | null | undefined): string { .replace(/\p{M}/gu, "") .replace(/\\/g, "") .replace(/[^A-Z0-9]+/g, " ") - .trim(); + .trim() + .split(" ") + .filter(Boolean) + .map(canonicalToken) + .join(" "); } /** Drop region suffixes "(EL)/(ER)", year ranges "(2012- 2017)", "2004 - 2014", and 4-digit years. */ diff --git a/apps/api/src/integrations/carcatonline/carcatonline.pacing.spec.ts b/apps/api/src/integrations/carcatonline/carcatonline.pacing.spec.ts index 49ecab5..7bc1fb6 100644 --- a/apps/api/src/integrations/carcatonline/carcatonline.pacing.spec.ts +++ b/apps/api/src/integrations/carcatonline/carcatonline.pacing.spec.ts @@ -18,6 +18,7 @@ import { function fakeRedis(): RedisLike & { store: Map; ttl: Map } { const store = new Map(); const ttl = new Map(); + const zs = new Map(); const set = async (key: string, value: string, mode: "EX" | "PX", n: number, cond?: "NX") => { if (cond === "NX" && store.has(key)) return null; store.set(key, value); @@ -42,6 +43,22 @@ function fakeRedis(): RedisLike & { store: Map; ttl: Map s > Number(max)), + ); + }, + async zcard(k) { + return (zs.get(k) ?? []).length; + }, + async zrange(k) { + const sorted = [...(zs.get(k) ?? [])].sort((a, b) => a[1] - b[1]); + return sorted.length ? [sorted[0][0], String(sorted[0][1])] : []; + }, }; } @@ -141,9 +158,9 @@ describe("hourly request cap", () => { expect(msUntilNextHour(at("2026-09-26T19:40:00Z"))).toBe(20 * 60_000 + 30_000); }); - it("stops before the upstream hourly limit and reports the wait", async () => { + it("stops before the upstream limit counted over a SLIDING 60 minutes, and reports when the oldest call ages out", async () => { const redis = fakeRedis(); - const t = at("2026-09-26T19:40:00Z"); + let t = at("2026-09-26T19:40:00Z"); const throttle = new CarcatonlineThrottle( redis, base, @@ -152,11 +169,22 @@ describe("hourly request cap", () => { redis.store.delete(CARCAT_KEYS.lastCall); }, ); - await throttle.beforeCall(); + await throttle.beforeCall(); // 19:40 + redis.store.delete(CARCAT_KEYS.lastCall); + t = at("2026-09-26T19:50:00Z"); + await throttle.beforeCall(); // 19:50 + expect(await throttle.callsInWindow()).toBe(2); + // a clock-hour bucket would reset at 20:00; the sliding window must not + t = at("2026-09-26T20:05:00Z"); + await expect(throttle.beforeCall()).rejects.toMatchObject({ + name: "CarcatonlineHourlyCapError", + retryAfterMs: 35 * 60_000 + 15_000, // oldest (19:40) + 60 min − 20:05 + margin + }); + // once the 19:40 call ages out, a slot frees up + t = at("2026-09-26T20:41:00Z"); redis.store.delete(CARCAT_KEYS.lastCall); await throttle.beforeCall(); - expect(await throttle.callsThisHour()).toBe(2); - await expect(throttle.beforeCall()).rejects.toBeInstanceOf(CarcatonlineHourlyCapError); + expect(await throttle.callsInWindow()).toBe(2); }); }); diff --git a/apps/api/src/integrations/carcatonline/carcatonline.pacing.ts b/apps/api/src/integrations/carcatonline/carcatonline.pacing.ts index d14a5ed..a867dbc 100644 --- a/apps/api/src/integrations/carcatonline/carcatonline.pacing.ts +++ b/apps/api/src/integrations/carcatonline/carcatonline.pacing.ts @@ -18,6 +18,11 @@ export interface RedisLike { expire(key: string, seconds: number): Promise; /** SET key value PX ms NX — returns "OK" or null. */ set(key: string, value: string, mode: "PX", ttlMs: number, cond: "NX"): Promise; + zadd(key: string, score: number, member: string): Promise; + zremrangebyscore(key: string, min: number | string, max: number | string): Promise; + zcard(key: string): Promise; + /** ZRANGE key 0 0 WITHSCORES → [member, score] */ + zrange(key: string, start: number, stop: number, withScores: "WITHSCORES"): Promise; } export const CARCAT_KEYS = { @@ -26,6 +31,7 @@ export const CARCAT_KEYS = { lastCall: "carcatonline:last-call", callsPrefix: "carcatonline:calls:", hourlyPrefix: "carcatonline:calls-hour:", + window: "carcatonline:calls-window", groupsPrefix: "carcatonline:groups:", modelsPrefix: "carcatonline:models:", } as const; @@ -56,6 +62,7 @@ export function carcatConfigFromEnv(env: NodeJS.ProcessEnv = process.env): Carca } const ISTANBUL_OFFSET_MS = 3 * 60 * 60 * 1000; // fixed UTC+3 (no DST since 2016) +const WINDOW_MS = 60 * 60 * 1000; /** Hour (0-23) and "YYYY-MM-DD" in Europe/Istanbul. */ export function istanbulParts(date: Date): { hour: number; minute: number; day: string } { @@ -125,7 +132,7 @@ export class CarcatonlineHourlyCapError extends Error { readonly retryAfterMs: number, ) { super( - `carcatonline hourly call cap reached (${used}/${cap}); next hour in ${Math.round(retryAfterMs / 60000)} min`, + `carcatonline hourly call cap reached (${used}/${cap} in the last 60 min); window frees in ${Math.round(retryAfterMs / 60000)} min`, ); this.name = "CarcatonlineHourlyCapError"; } @@ -193,19 +200,29 @@ export class CarcatonlineThrottle { return Number((await this.redis.get(hourlyCallsKey(this.now()))) ?? 0); } + /** Calls made in the last 60 minutes (the upstream limit is a sliding window). */ + async callsInWindow(): Promise { + const nowMs = this.now().getTime(); + await this.redis.zremrangebyscore(CARCAT_KEYS.window, "-inf", nowMs - WINDOW_MS); + return this.redis.zcard(CARCAT_KEYS.window); + } + async beforeCall(): Promise { const locked = await this.lockoutRemainingMs(); if (locked > 0) throw new CarcatonlineLockedError(locked); const used = await this.callsToday(); if (used >= this.cfg.dailyCallCap) throw new CarcatonlineBudgetError(used, this.cfg.dailyCallCap); - const usedHour = await this.callsThisHour(); - if (usedHour >= this.cfg.hourlyCallCap) { - throw new CarcatonlineHourlyCapError( - usedHour, - this.cfg.hourlyCallCap, - msUntilNextHour(this.now()), - ); + // Upstream "Customer Hourly Request Limit" is a SLIDING 60-minute window, + // so a clock-hour bucket lets two bursts straddle a boundary and trip it. + // Count our own calls in the last hour and wait until the oldest one ages out. + const inWindow = await this.callsInWindow(); + if (inWindow >= this.cfg.hourlyCallCap) { + const nowMs = this.now().getTime(); + const oldest = await this.redis.zrange(CARCAT_KEYS.window, 0, 0, "WITHSCORES"); + const oldestMs = Number(oldest[1] ?? nowMs); + const retryAfterMs = Math.max(oldestMs + WINDOW_MS - nowMs, 0) + 15_000; + throw new CarcatonlineHourlyCapError(inWindow, this.cfg.hourlyCallCap, retryAfterMs); } // Shared min-interval slot: SET NX PX; spin (bounded) until acquired. for (let i = 0; i < 50; i += 1) { @@ -225,5 +242,12 @@ export class CarcatonlineThrottle { const hourKey = hourlyCallsKey(this.now()); const h = await this.redis.incr(hourKey); if (h === 1) await this.redis.expire(hourKey, 2 * 3600); + const stamp = this.now().getTime(); + await this.redis.zadd( + CARCAT_KEYS.window, + stamp, + `${stamp}-${Math.random().toString(36).slice(2, 8)}`, + ); + await this.redis.expire(CARCAT_KEYS.window, 2 * 3600); } } diff --git a/apps/api/src/jobs/processors/carcatonline-backfill.processor.spec.ts b/apps/api/src/jobs/processors/carcatonline-backfill.processor.spec.ts index 2d8ac73..857f118 100644 --- a/apps/api/src/jobs/processors/carcatonline-backfill.processor.spec.ts +++ b/apps/api/src/jobs/processors/carcatonline-backfill.processor.spec.ts @@ -19,6 +19,7 @@ const CFG = { function fakeRedis() { const store = new Map(); + const zs = new Map(); return { store, async get(k: string) { @@ -40,6 +41,22 @@ function fakeRedis() { return n; }, async expire() {}, + async zadd(k: string, score: number, member: string) { + zs.set(k, [...(zs.get(k) ?? []), [member, score]]); + }, + async zremrangebyscore(k: string, _min: unknown, max: number) { + zs.set( + k, + (zs.get(k) ?? []).filter(([, s]) => s > Number(max)), + ); + }, + async zcard(k: string) { + return (zs.get(k) ?? []).length; + }, + async zrange(k: string) { + const sorted = [...(zs.get(k) ?? [])].sort((a, b) => a[1] - b[1]); + return sorted.length ? [sorted[0][0], String(sorted[0][1])] : []; + }, async keys(pattern: string) { const prefix = pattern.replace(/\*$/, ""); return [...store.keys()].filter((k) => k.startsWith(prefix)); @@ -163,7 +180,7 @@ describe("processCarcatonlineBackfill", () => { expect(night.queue.add).toHaveBeenCalledWith( CARCAT_JOB.vehicle, { catalogVehicleId: "cv-1" }, - { jobId: "carcat-cv-1-2026-09-26" }, + { jobId: "carcat-cv-1-2026-09-26", priority: 3 }, ); }); @@ -256,7 +273,7 @@ describe("processCarcatonlineBackfill", () => { expect(j.moveToDelayed.mock.calls[0][0]).toBeGreaterThan(NIGHT.getTime() + 2400 * 1000); }); - it("vehicle: stops at the hourly cap and re-schedules to the next clock hour", async () => { + it("vehicle: stops at the sliding hourly cap and re-schedules until the window frees", async () => { const { db } = fakeDb([[cv], []]); const redis = fakeRedis(); const client = fakeClient({ @@ -265,9 +282,9 @@ describe("processCarcatonlineBackfill", () => { const j = job(CARCAT_JOB.vehicle, { catalogVehicleId: "cv-1" }); const deps = { ...baseDeps(db, redis, client, NIGHT), config: { ...CFG, hourlyCallCap: 1 } }; await expect(processCarcatonlineBackfill(j, "tok", deps)).rejects.toBeInstanceOf(DelayedError); - // models (1 call) fits the cap; the cascade's first call trips it → delayed to the next hour + margin + // models (1 call) fits the cap; the cascade's first call trips it → delayed until that call ages out of the 60-min window (+15 s margin) const [ts] = j.moveToDelayed.mock.calls[0]; - expect(ts).toBe(Date.parse("2026-09-26T20:00:00Z") + 30_000); + expect(ts).toBe(NIGHT.getTime() + 3600_000 + 15_000); }); it("vehicle: a re-run resumes from the matched car and cached group responses without new calls", async () => { diff --git a/apps/api/src/jobs/processors/carcatonline-backfill.processor.ts b/apps/api/src/jobs/processors/carcatonline-backfill.processor.ts index 7b3087c..dc07fb7 100644 --- a/apps/api/src/jobs/processors/carcatonline-backfill.processor.ts +++ b/apps/api/src/jobs/processors/carcatonline-backfill.processor.ts @@ -38,6 +38,42 @@ import { CARCAT_JOB, type CarcatonlineVehicleJobData } from "../queues/carcatonl type Database = PostgresJsDatabase>; +/** + * Expected tree cost by brand: 1 = 2-level PartsLink24 trees (VW group; measured + * Skoda 12–38 calls), 3 = 3-level trees with hundreds of subgroup calls + * (Renault/Dacia/PSA — Austral needed 264+), 2 = unknown/mid. Cheap trees are + * scanned first so a night's ~1.000-call budget completes many catalogs instead + * of stalling on a few big ones; in-flight (matched) vehicles always go first. + */ +export const BRAND_TREE_RANK: Record = { + volkswagen: 1, + skoda: 1, + seat: 1, + cupra: 1, + audi: 1, + porsche: 1, + bentley: 1, + renault: 3, + dacia: 3, + alpine: 3, + citroen: 3, + peugeot: 3, + ds: 3, + opel: 3, +}; + +export function brandTreeRank(brandName: string | null | undefined): 1 | 2 | 3 { + return BRAND_TREE_RANK[(brandName ?? "").trim().toLowerCase()] ?? 2; +} + +/** BullMQ priority (lower runs first): resume in-flight vehicles, then cheap trees. */ +export function jobPriority( + status: string | null | undefined, + brandName: string | null | undefined, +): number { + return status === "matched" ? 1 : brandTreeRank(brandName) + 1; +} + /** Only architectures whose catalog page is served from the DB tree (see CatalogService.getCategoryTree). */ export const CARCAT_BACKFILL_ARCHITECTURES = ["P5_MODERN"]; @@ -120,8 +156,17 @@ async function scan( if (pending >= batch) return { skipped: "queue-busy", pending }; const cutoff = new Date(now().getTime() - RETRY_AFTER_DAYS * 24 * 3600 * 1000).toISOString(); + const rankCase = sql.raw( + `CASE lower(brand_name) ${Object.entries(BRAND_TREE_RANK) + .map(([b, r]) => `WHEN '${b}' THEN ${r}`) + .join(" ")} ELSE 2 END`, + ); const rows = await deps.db - .select({ id: catalogVehicles.id }) + .select({ + id: catalogVehicles.id, + brandName: catalogVehicles.brandName, + status: sql`${catalogVehicles.metadata}->'carcatonline'->>'status'`, + }) .from(catalogVehicles) .where( and( @@ -138,7 +183,11 @@ async function scan( ), ) .orderBy( - sql`${catalogVehicles.metadata}->'carcatonline'->>'at' NULLS FIRST, ${catalogVehicles.createdAt}`, + // in-flight first, then cheap trees, then never-tried before retries + sql`CASE WHEN (${catalogVehicles.metadata}->'carcatonline'->>'status') = 'matched' THEN 0 ELSE 1 END`, + rankCase, + sql`${catalogVehicles.metadata}->'carcatonline'->>'at' NULLS FIRST`, + catalogVehicles.createdAt, ) .limit(batch - pending); @@ -146,7 +195,10 @@ async function scan( await deps.queue.add( CARCAT_JOB.vehicle, { catalogVehicleId: r.id } satisfies CarcatonlineVehicleJobData, - { jobId: `carcat-${r.id}-${now().toISOString().slice(0, 10)}` }, + { + jobId: `carcat-${r.id}-${now().toISOString().slice(0, 10)}`, + priority: jobPriority(r.status, r.brandName), + }, ); } log.log(`[carcatonline] scan enqueued ${rows.length} vehicle(s) (pending before: ${pending})`); @@ -212,7 +264,7 @@ async function fillVehicle( `${CARCAT_KEYS.groupsPrefix}${key}`, JSON.stringify(groups), "EX", - 24 * 3600, + 7 * 24 * 3600, ); }, },