From 7d439fc83be5c09a1d155667a6a6aeccda7d256b Mon Sep 17 00:00:00 2001 From: Semih Yesilyurt Date: Sun, 27 Sep 2026 19:06:31 +0300 Subject: [PATCH] =?UTF-8?q?feat(prefetch):=20PL24=20b=C3=BCt=C3=A7esini=20?= =?UTF-8?q?par=C3=A7aya=20kayd=C4=B1r=20=E2=80=94=20tavan-=C3=BCst=C3=BC?= =?UTF-8?q?=20no-op=20k=C4=B1rpma,=20%50=20a=C4=9Fa=C3=A7=20pay=C4=B1,=20a?= =?UTF-8?q?=C3=A7=C4=B1lan-par=C3=A7as=C4=B1z=20Faz-0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Prod ölçümü (2026-09-27): PL24 günlük bütçesinin %93'ü ağaç genişletmeye, %7'si parçaya gidiyordu; park etmiş 3,5k pl24 fast-lane işinin 426'sı tavan-üstü children no-op'uydu (ertesi günün 600'lük bütçesinin %71'i); kullanıcıların açtığı 239 PL24 aracı hâlâ parçasızdı (fast lane 1. seviyede durur, P5 yaprakları 2-3. seviyede; bulk backfill kapalı). - queueCategoryJob artık ŞERİDİN tavanına (maxDepthFor) bakar, MAX_DEPTH'e değil: fast lane'de 1. seviyeden sonra children işi hiç üretilmez. - process(): tavan-üstü children işi hiçbir gate/sayaçtan önce settle edilir (incrementCompleted ile zincir muhasebesi kapanır) — park etmiş eski no-op'lar bütçe yakmaz. - PL24 ana şerit derinliği 3 (PREFETCH_PL24_MAX_DEPTH), fast 1 kalır. - Günlük bütçe bölüşümü: ağaç genişletme (init+children) günlük bütçenin en fazla %50'sini kullanır (PREFETCH_PL24_CHILDREN_SHARE), kalan parçaya; ayrı sayaç prefetch:daily:pl24::children. - Scan Faz-0: kullanıcı açmış (user_vehicles, son 90 gün) + parçasız pl24 araçları dalga başına 3 adet, fast kuyrukta (lifo → önce çalışır, taze decode yine öne geçer), depthCap=3, backfill:true (pace+jitter, ANA bütçe eşiği 480 → 120 taze decode'a kalır). Bulk PL24 anahtarından bağımsız, PREFETCH_PL24_OPENED_ENABLED=false ile kapanır. Co-Authored-By: Claude Fable 5.1 --- .../src/jobs/prefetch-budget-split.spec.ts | 397 ++++++++++++++++++ .../src/jobs/prefetch-worker.service.spec.ts | 2 +- apps/api/src/jobs/prefetch-worker.service.ts | 301 ++++++++++--- apps/api/src/jobs/prefetch.types.ts | 15 + docker-compose.coolify.yml | 5 + 5 files changed, 652 insertions(+), 68 deletions(-) create mode 100644 apps/api/src/jobs/prefetch-budget-split.spec.ts diff --git a/apps/api/src/jobs/prefetch-budget-split.spec.ts b/apps/api/src/jobs/prefetch-budget-split.spec.ts new file mode 100644 index 0000000..f92b117 --- /dev/null +++ b/apps/api/src/jobs/prefetch-budget-split.spec.ts @@ -0,0 +1,397 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +/** + * PL24 bütçe bölüşümü + açılan-parçasız öncelik (2026-09-27). + * + * Prod ölçümü: PL24 günlük bütçesinin %93'ü ağaç genişletmeye (children), %7'si + * parçaya gidiyordu; kuyrukta park etmiş 3,5k pl24 işinin 426'sı tavan-üstü + * children no-op'uydu (ertesi günün 600'lük bütçesinin %71'i boşa yanacaktı); + * kullanıcıların açtığı 239 PL24 aracı parçasızdı (fast lane 1. seviyede durur, + * P5 yaprakları 2-3. seviyede). + * + * Kilitlenen davranışlar: + * 1. Tavan-üstü children işi hiç kuyruğa girmez; girmiş olan hiçbir sayaca + * dokunmadan (pencere/rate/bütçe) muhasebesi kapatılıp düşer. + * 2. Ağaç genişletme günlük bütçenin en fazla payını (varsayılan %50) yer; + * parça işleri kalanı kullanır. + * 3. Scan Faz-0: kullanıcı açmış + parçasız pl24 araçları, bulk anahtarı kapalı + * olsa da fast kuyrukta (lifo), depthCap=3 ve `backfill: true` ile en SONA + * eklenir (lifo → ilk çalışır). + */ + +const REAL_ENV = { ...process.env }; + +async function load(env: Record) { + vi.resetModules(); + for (const [k, v] of Object.entries(env)) { + if (v === undefined) delete process.env[k]; + else process.env[k] = v; + } + return await import("./prefetch-worker.service"); +} + +const dedupeStub = { + tryCopyChildren: async () => 0, + tryCopyParts: async () => 0, + dedupeVehicle: async () => ({ + copiedCategories: 0, + copiedParts: 0, + copiedPics: 0, + leavesFilled: 0, + donorVehicleId: null, + }), +}; + +function makeQueue(name: string) { + return { + name, + add: vi.fn(async (..._a: unknown[]) => undefined), + getJob: vi.fn(async (..._a: unknown[]): Promise => null), + getJobCounts: vi.fn(async () => ({ waiting: 0, delayed: 0, active: 0, prioritized: 0 })), + toKey: (t: string) => `bull:${name}:${t}`, + client: Promise.resolve({ zcount: vi.fn(async () => 0) }), + }; +} + +function makeRedis(overrides: Record = {}) { + const incrCalls: string[] = []; + const redis = { + exists: vi.fn(async (..._a: unknown[]) => false), + get: vi.fn(async (..._a: unknown[]): Promise => null), + set: vi.fn(async (..._a: unknown[]) => undefined), + del: vi.fn(async (..._a: unknown[]) => undefined), + incr: vi.fn(async (k: string) => { + incrCalls.push(k); + return 1; + }), + expire: vi.fn(async () => undefined), + ttl: vi.fn(async () => -2), + setNx: vi.fn(async () => true), + getJson: vi.fn(async (..._a: unknown[]): Promise => null), + setJson: vi.fn(async () => undefined), + ...overrides, + }; + return { redis, incrCalls }; +} + +/** Fake drizzle chain: every terminal `.limit()` yields the next queued rows. */ +function makeDb(limitResults: unknown[][]) { + const queued = [...limitResults]; + const chain: Record = {}; + for (const m of ["select", "from", "where", "orderBy", "innerJoin", "groupBy"]) { + chain[m] = vi.fn(() => chain); + } + chain.limit = vi.fn(() => queued.shift() ?? []); + return chain; +} + +function makeJob(name: string, data: Record) { + return { + name, + data, + moveToDelayed: vi.fn(async (..._a: unknown[]) => undefined), + attemptsMade: 0, + }; +} + +type Svc = { + process: (j: unknown, t?: string) => Promise; + queueCategoryJob: (c: unknown, v: string, s: string, d: number, f: boolean) => Promise; + checkSourceDailyBudget: (s: string, lane: string, kind?: string) => Promise; + processBackfillScan: () => Promise; +}; + +function build( + mod: Awaited>, + deps: { redis: unknown; db?: unknown; queue?: unknown; fastQueue?: unknown }, +) { + const queue = deps.queue ?? makeQueue("catalog-prefetch"); + const fastQueue = deps.fastQueue ?? makeQueue("catalog-prefetch-fast"); + const svc = new mod.PrefetchWorkerService( + queue as never, + fastQueue as never, + { getChildren: vi.fn(async () => []), getCategoryWithParts: vi.fn(async () => ({})) } as never, + deps.redis as never, + { payload: vi.fn(async () => ({})) } as never, + (deps.db ?? makeDb([])) as never, + dedupeStub as never, + ); + return svc as never as Svc; +} + +afterEach(() => { + process.env = { ...REAL_ENV }; + vi.useRealTimers(); +}); + +describe("PL24 derinlik tavanı — ana şerit 3, fast 1, chain depthCap kazanır", () => { + it("varsayılanlar", async () => { + const mod = await load({ PREFETCH_PL24_MAX_DEPTH: "", PREFETCH_PL24_FAST_DEPTH: "" }); + const fn = mod.__testables.maxDepthFor; + expect(fn("pl24", true)).toBe(1); + expect(fn("pl24", false)).toBe(3); + expect(fn("emex", false)).toBe(12); + // Faz-0 zinciri: fast kuyrukta ama depthCap ile 3'e iner/çıkar + expect(fn("pl24", true, 3)).toBe(3); + // depthCap asla MAX_DEPTH'i aşmaz + expect(fn("pl24", false, 99)).toBe(12); + }); + + it("envShare: 0/1/boş → kapalı ya da varsayılan", async () => { + const mod = await load({}); + const { envShare } = mod.__testables; + process.env.X_SHARE = "0.3"; + expect(envShare("X_SHARE", 0.5)).toBe(0.3); + process.env.X_SHARE = "0"; + expect(envShare("X_SHARE", 0.5)).toBe(0); + process.env.X_SHARE = "1"; + expect(envShare("X_SHARE", 0.5)).toBe(0); + process.env.X_SHARE = ""; + expect(envShare("X_SHARE", 0.5)).toBe(0.5); + }); +}); + +describe("tavan-üstü children işi", () => { + const nonLeaf = { + id: "c1", + linkPath: "/p5vwag/extern/groups/x", + source: "pl24", + unavailable: false, + hasSubgroups: false, + }; + + it("fast şeritte 1. seviyeden sonra kuyruğa GİRMEZ, 0. seviyede girer", async () => { + const mod = await load({ PREFETCH_PL24_FAST_DEPTH: "" }); + const fastQueue = makeQueue("catalog-prefetch-fast"); + const svc = build(mod, { redis: makeRedis().redis, db: makeDb([[], []]), fastQueue }); + // depth 1 == fast tavanı → no-op olurdu, iş üretilmemeli + expect(await svc.queueCategoryJob(nonLeaf, "v1", "pl24", 1, true)).toBe(0); + expect(fastQueue.add).not.toHaveBeenCalled(); + // depth 0 → çocuklarını çekmeye değer + expect(await svc.queueCategoryJob(nonLeaf, "v1", "pl24", 0, true)).toBe(1); + expect(fastQueue.add).toHaveBeenCalledTimes(1); + expect(fastQueue.add.mock.calls[0][0]).toBe("prefetch-children"); + }); + + it("ana şeritte 3. seviyede durur (eski MAX_DEPTH=12 değil)", async () => { + const mod = await load({ PREFETCH_PL24_MAX_DEPTH: "" }); + const queue = makeQueue("catalog-prefetch"); + const svc = build(mod, { redis: makeRedis().redis, db: makeDb([[], []]), queue }); + expect(await svc.queueCategoryJob(nonLeaf, "v1", "pl24", 3, false)).toBe(0); + expect(await svc.queueCategoryJob(nonLeaf, "v1", "pl24", 2, false)).toBe(1); + expect(queue.add).toHaveBeenCalledTimes(1); + }); + + it("kuyruğa girmiş eski no-op: pencere/rate/bütçe sayaçlarına dokunmadan muhasebeyi kapatır", async () => { + // Pencere KAPALI (09-18 dışı) ve bütçe DOLU olsa bile iş hiçbir gate'e uğramaz. + vi.useFakeTimers(); + vi.setSystemTime(Date.UTC(2026, 8, 27, 0, 0, 0)); // 03:00 İstanbul — pencere dışı + const mod = await load({ + PREFETCH_PL24_START: "9", + PREFETCH_PL24_END: "18", + PREFETCH_DAILY_PL24: "600", + }); + const { redis, incrCalls } = makeRedis({ + get: vi.fn(async (k: string) => (String(k).startsWith("prefetch:daily") ? "600" : null)), + getJson: vi.fn(async () => ({ completed: 4, total: 5 })), + }); + const svc = build(mod, { redis }); + const job = makeJob("prefetch-children", { + vehicleId: "v1", + categoryId: "c1", + source: "pl24", + depth: 1, + fast: true, + }); + await expect(svc.process(job, "tok")).resolves.toBeUndefined(); + expect(job.moveToDelayed).not.toHaveBeenCalled(); + // ne rate ne daily sayaç (finalize'nin noresult sayacı bu işin değil, zincirin kapanışıdır) + expect( + incrCalls.filter((k) => k.startsWith("prefetch:rate") || k.startsWith("prefetch:daily")), + ).toHaveLength(0); + // zincir muhasebesi: completed 4→5 = total → araç finalize edilir + const progressWrites = (redis.setJson as ReturnType).mock.calls.filter((c) => + String(c[0]).startsWith("prefetch:progress:"), + ); + expect(progressWrites.length).toBeGreaterThan(0); + expect((progressWrites[0][1] as { completed: number }).completed).toBe(5); + }); +}); + +describe("günlük bütçe bölüşümü — ağaç genişletme payı", () => { + it("children payı dolunca children ertelenir, parça işi geçer; sayaçlar ayrı", async () => { + vi.useFakeTimers(); + vi.setSystemTime(Date.UTC(2026, 8, 27, 9, 0, 0)); // 12:00 İstanbul + const mod = await load({ PREFETCH_DAILY_PL24: "600", PREFETCH_PL24_CHILDREN_SHARE: "" }); + const state: Record = {}; + const { redis, incrCalls } = makeRedis({ + get: vi.fn(async (k: string) => state[k] ?? null), + }); + const svc = build(mod, { redis }); + + // Toplam 100 harcanmış, children 300 (= %50 × 600) → children durur + for (const k of Object.keys(state)) delete state[k]; + const day = Math.floor(Date.now() / 86_400_000); + state[`prefetch:daily:pl24:${day}`] = "100"; + state[`prefetch:daily:pl24:${day}:children`] = "300"; + await expect(svc.checkSourceDailyBudget("pl24", "fast", "children")).rejects.toThrow( + /Rate limited/, + ); + expect(incrCalls).toHaveLength(0); + + // Aynı durumda parça işi geçer ve YALNIZ toplam sayacı artırır + await svc.checkSourceDailyBudget("pl24", "fast", "parts"); + expect(incrCalls).toEqual([`prefetch:daily:pl24:${day}`]); + + // children payı altında: iki sayaç da artar + state[`prefetch:daily:pl24:${day}:children`] = "10"; + incrCalls.length = 0; + await svc.checkSourceDailyBudget("pl24", "fast", "children"); + expect(incrCalls).toEqual([ + `prefetch:daily:pl24:${day}`, + `prefetch:daily:pl24:${day}:children`, + ]); + }); + + it("pay kapalıysa (0) children sınırsız ve ek sayaç yok", async () => { + const mod = await load({ PREFETCH_DAILY_PL24: "600", PREFETCH_PL24_CHILDREN_SHARE: "0" }); + const day = Math.floor(Date.now() / 86_400_000); + const { redis, incrCalls } = makeRedis({ + get: vi.fn(async (k: string) => (String(k).endsWith(":children") ? "9999" : "0")), + }); + const svc = build(mod, { redis }); + await svc.checkSourceDailyBudget("pl24", "main", "children"); + expect(incrCalls).toEqual([`prefetch:daily:pl24:${day}`]); + }); + + it("process(): children işi 'children' türüyle, parça işi 'parts' türüyle bütçelenir", async () => { + vi.useFakeTimers(); + vi.setSystemTime(Date.UTC(2026, 8, 27, 9, 0, 0)); + const mod = await load({ PREFETCH_DAILY_PL24: "600" }); + const day = Math.floor(Date.now() / 86_400_000); + const { redis, incrCalls } = makeRedis({ + get: vi.fn(async (k: string) => (String(k).endsWith(":children") ? "300" : "100")), + }); + const svc = build(mod, { redis }); + const children = makeJob("prefetch-children", { + vehicleId: "v1", + categoryId: "c1", + source: "pl24", + depth: 0, + fast: true, + }); + await expect(svc.process(children, "tok")).rejects.toThrow(); + expect(children.moveToDelayed).toHaveBeenCalledTimes(1); + // Ertelenen iş günlük sayaca yazmaz (rate sayacı yazar) + expect(incrCalls.filter((k) => k.startsWith("prefetch:daily"))).toHaveLength(0); + + incrCalls.length = 0; + const parts = makeJob("prefetch-parts", { + vehicleId: "v1", + categoryId: "c2", + source: "pl24", + depth: 1, + fast: true, + }); + await svc.process(parts, "tok"); + expect(incrCalls.filter((k) => k.startsWith("prefetch:daily"))).toEqual([ + `prefetch:daily:pl24:${day}`, + ]); + }); + + it("Faz-0 zinciri (fast + backfill) ANA eşiğe tabi: 480'de durur, saf fast 600'e kadar gider", async () => { + vi.useFakeTimers(); + vi.setSystemTime(Date.UTC(2026, 8, 27, 9, 0, 0)); + const mod = await load({ PREFETCH_DAILY_PL24: "600" }); + const { redis } = makeRedis({ + get: vi.fn(async (k: string) => (String(k).endsWith(":children") ? "0" : "480")), + }); + const svc = build(mod, { redis }); + const phase0 = makeJob("prefetch-parts", { + vehicleId: "v1", + categoryId: "c1", + source: "pl24", + depth: 2, + fast: true, + backfill: true, + depthCap: 3, + }); + await expect(svc.process(phase0, "tok")).rejects.toThrow(); + expect(phase0.moveToDelayed).toHaveBeenCalledTimes(1); + + const fresh = makeJob("prefetch-parts", { + vehicleId: "v2", + categoryId: "c2", + source: "pl24", + depth: 1, + fast: true, + }); + await expect(svc.process(fresh, "tok")).resolves.toBeUndefined(); + }); +}); + +describe("scan Faz-0 — kullanıcı açmış, parçasız pl24 araçları", () => { + beforeEach(() => { + process.env.CATALOG_BACKFILL_ENABLED = "true"; + }); + + it("bulk pl24 anahtarı KAPALIYKEN bile fast kuyruğa depthCap=3 + backfill ile ve EN SONA eklenir", async () => { + const mod = await load({ + PL24_BACKFILL_ENABLED: undefined, + PL24_TR_DISABLED: "true", + PREFETCH_PL24_OPENED_ENABLED: "", + PREFETCH_PL24_OPENED_PER_WAVE: "2", + }); + // Faz-1 emex parçasız aracı, Faz-2 yok (baskı yok ama cursor sorgusu boş), Faz-0 iki pl24 aracı + const db = makeDb([ + [{ id: "e1", source: "emex" }], // Faz-1 (emex eligible) + [], // Faz-2 rolling + [ + { id: "p1", source: "pl24", lastOpen: new Date() }, + { id: "p2", source: "pl24", lastOpen: new Date() }, + { id: "p3", source: "pl24", lastOpen: new Date() }, // per-wave 2 → dışarıda kalır + ], + ]); + const fastQueue = makeQueue("catalog-prefetch-fast"); + const svc = build(mod, { redis: makeRedis().redis, db, fastQueue }); + await svc.processBackfillScan(); + + const adds = fastQueue.add.mock.calls.map((c) => c[1] as Record); + expect(adds.map((a) => a.vehicleId)).toEqual(["e1", "p1", "p2"]); + expect(adds[0]).not.toHaveProperty("backfill"); + expect(adds[1]).toMatchObject({ fast: true, backfill: true, depthCap: 3, source: "pl24" }); + expect(adds[2]).toMatchObject({ fast: true, backfill: true, depthCap: 3 }); + // Faz-0 sorgusu innerJoin + groupBy kullandı + expect(db.innerJoin).toHaveBeenCalled(); + expect(db.groupBy).toHaveBeenCalled(); + }); + + it("anahtar 'false' ise Faz-0 sorgusu hiç çalışmaz", async () => { + const mod = await load({ + PL24_BACKFILL_ENABLED: undefined, + PREFETCH_PL24_OPENED_ENABLED: "false", + }); + const db = makeDb([[], [], [{ id: "p1", source: "pl24", lastOpen: new Date() }]]); + const fastQueue = makeQueue("catalog-prefetch-fast"); + const svc = build(mod, { redis: makeRedis().redis, db, fastQueue }); + await svc.processBackfillScan(); + expect(fastQueue.add).not.toHaveBeenCalled(); + expect(db.innerJoin).not.toHaveBeenCalled(); + }); + + it("pl24 ana bütçe eşiği dolduysa Faz-0 o dalgada beslenmez", async () => { + const mod = await load({ + PL24_BACKFILL_ENABLED: undefined, + PREFETCH_PL24_OPENED_ENABLED: "", + PREFETCH_DAILY_PL24: "600", + }); + const { redis } = makeRedis({ + get: vi.fn(async (k: string) => (String(k).startsWith("prefetch:daily:pl24") ? "480" : null)), + }); + const db = makeDb([[], [], [{ id: "p1", source: "pl24", lastOpen: new Date() }]]); + const fastQueue = makeQueue("catalog-prefetch-fast"); + const svc = build(mod, { redis, db, fastQueue }); + await svc.processBackfillScan(); + expect(fastQueue.add).not.toHaveBeenCalled(); + }); +}); diff --git a/apps/api/src/jobs/prefetch-worker.service.spec.ts b/apps/api/src/jobs/prefetch-worker.service.spec.ts index 79d7385..d626022 100644 --- a/apps/api/src/jobs/prefetch-worker.service.spec.ts +++ b/apps/api/src/jobs/prefetch-worker.service.spec.ts @@ -9,7 +9,7 @@ import { PrefetchWorkerService } from "./prefetch-worker.service"; function makeDb(limitResults: unknown[][]) { const queued = [...limitResults]; const chain: Record = {}; - for (const m of ["select", "from", "where", "orderBy"]) { + for (const m of ["select", "from", "where", "orderBy", "innerJoin", "groupBy"]) { chain[m] = vi.fn(() => chain); } chain.limit = vi.fn(() => queued.shift() ?? []); diff --git a/apps/api/src/jobs/prefetch-worker.service.ts b/apps/api/src/jobs/prefetch-worker.service.ts index 2a07fdd..0b826ad 100644 --- a/apps/api/src/jobs/prefetch-worker.service.ts +++ b/apps/api/src/jobs/prefetch-worker.service.ts @@ -6,11 +6,11 @@ import { type OnModuleInit, } from "@nestjs/common"; import { DelayedError, type Job, type Queue, Worker } from "bullmq"; -import { and, asc, eq, gt, inArray, isNull, notExists, sql } from "drizzle-orm"; +import { and, asc, desc, eq, gt, gte, inArray, isNull, notExists, sql } from "drizzle-orm"; import { CatalogDedupe } from "../catalog/catalog-dedupe.service"; import { CategoriesService } from "../categories/categories.service"; import { DATABASE, type Database } from "../database/database.provider"; -import { categories, parts, vehicles } from "../database/schema/core"; +import { categories, parts, userVehicles, vehicles } from "../database/schema/core"; import { isPl24LeafNode, isPl24PartDetailNode } from "../integrations/pl24/pl24-tree"; import { PostHogService } from "../posthog/posthog.service"; import { RedisService } from "../redis/redis.service"; @@ -212,9 +212,53 @@ const PL24_PACE_MS = Number(process.env.PREFETCH_PL24_DELAY_MS) || 8_000; */ const PL24_FAST_MAX_DEPTH = Number(process.env.PREFETCH_PL24_FAST_DEPTH) || 1; -/** Depth ceiling for this source+lane. */ -function maxDepthFor(source: string, fast: boolean): number { - if (source === "pl24" && fast) return PL24_FAST_MAX_DEPTH; +/** + * Depth ceiling for PL24 jobs that are NOT the fresh-decode fast lane (budgeted + * backfill + the opened-partless Phase-0 below). 3 levels reach the illustration + * leaves of every P5 tree we measured (VW-group main → sub → illustration, BMW + * hg → btnr, Renault group → sub → illustration); anything deeper is fetched + * when a user opens the node. Was MAX_DEPTH (12): measured on prod 2026-09-27, + * tree expansion took 93% of the PL24 budget and parts only 7%, and even the + * most recent vehicles ended half-drilled (T-Roc 110/643 leaves, Passat 50/670). + */ +const PL24_MAIN_MAX_DEPTH = Number(process.env.PREFETCH_PL24_MAX_DEPTH) || 3; + +/** + * Share of a source's DAILY budget that tree expansion (init + children jobs) + * may consume; the rest is reserved for parts jobs. Parts are what users search + * for, and one leaf's parts list costs one request, so once the share is spent + * the queue drains the leaves it already knows instead of discovering more. + * "0" or "1" disables the split for that source. + */ +function envShare(name: string, dflt: number): number { + const raw = process.env[name]; + if (raw === undefined || raw === "") return dflt; + const v = Number(raw); + return v > 0 && v < 1 ? v : 0; +} +const SOURCE_CHILDREN_SHARE: Record = { + pl24: envShare("PREFETCH_PL24_CHILDREN_SHARE", 0.5), +}; + +/** + * Phase-0 of the scan: PL24 vehicles a user OPENED (user_vehicles) that still + * have zero parts. The fresh-decode fast lane stops at depth 1, so P5 leaves + * (depth 2-3) were never reached and the user saw groups but no parts. Prod + * 2026-09-27: 239 such vehicles (15 opened in the last week). Capped per hourly + * wave, paced, and charged against the MAIN budget threshold — user demand, not + * bulk load, so it runs even while the bulk backfill switch is off. + */ +const PL24_OPENED_ENABLED = process.env.PREFETCH_PL24_OPENED_ENABLED !== "false"; +const PL24_OPENED_PER_WAVE = Number(process.env.PREFETCH_PL24_OPENED_PER_WAVE) || 3; +const PL24_OPENED_DAYS = Number(process.env.PREFETCH_PL24_OPENED_DAYS) || 90; + +/** Extra chain data inherited by every sub-job of an init (see prefetch.types). */ +type ChainOpts = { depthCap?: number; backfill?: boolean }; + +/** Depth ceiling for this source+lane (an explicit chain cap wins, never above MAX_DEPTH). */ +function maxDepthFor(source: string, fast: boolean, depthCap?: number): number { + if (depthCap && depthCap > 0) return Math.min(depthCap, MAX_DEPTH); + if (source === "pl24") return fast ? PL24_FAST_MAX_DEPTH : PL24_MAIN_MAX_DEPTH; return MAX_DEPTH; } @@ -224,7 +268,7 @@ function jitter(ms: number): number { } /** Test-only surface for the pure helpers above. */ -export const __testables = { maxDepthFor, jitter, isPl24BackfillEnabled }; +export const __testables = { maxDepthFor, jitter, isPl24BackfillEnabled, envShare }; // ── Phase-1 residue exclusion ── /** @@ -388,6 +432,19 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { } } } + // Depth-ceiling no-op: a children job at/over its lane's ceiling can never + // do upstream work, so settle it BEFORE any gate or counter (pure + // arithmetic, no I/O). Prod 2026-09-27: 426 of the 3.5k parked pl24 + // fast-lane jobs were such no-ops — 71% of the next day's 600-call budget + // would have been charged for nothing. + if (job.name === "prefetch-children") { + const d = job.data as PrefetchCategoryJobData; + if (d.source && d.depth >= maxDepthFor(d.source, !!d.fast, d.depthCap)) { + return await this.processChildren(job as Job); + } + } + const fast = !!(job.data as { fast?: boolean }).fast; + const backfill = !!(job.data as { backfill?: boolean }).backfill; // Per-source rate gate FIRST (before the pcat pace) so we don't burn the 15s // sleep on a job we're about to defer. Scan jobs are exempt. if ( @@ -396,7 +453,12 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { job.name === "prefetch-children" || job.name === "prefetch-parts") ) { - const lane = (job.data as { fast?: boolean }).fast ? "fast" : "main"; + const lane = fast ? "fast" : "main"; + // Budget threshold: Phase-0 rides the fast QUEUE (ordering) but is + // background work, so it stops at the main-lane threshold and leaves the + // fast-lane reserve to live decodes. + const budgetLane = fast && !backfill ? "fast" : "main"; + const kind = job.name === "prefetch-parts" ? "parts" : "children"; // GATE ORDER IS LOAD-BEARING — cheapest first, and every gate that can // reject the job must run BEFORE any counter is debited. // @@ -421,16 +483,17 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { // here, so only jobs about to do real work are counted. The lane // decides which threshold applies (backfill stops at the main limit, // the user's fast lane may use the full budget). - await this.checkSourceDailyBudget(data.source, lane); + await this.checkSourceDailyBudget(data.source, budgetLane, kind); } if (data.source === "parts-catalogs" && PCAT_PACE_MS > 0) { await new Promise((r) => setTimeout(r, PCAT_PACE_MS)); } - // PL24 backfill only: pace + jitter. The fast (user) lane is never delayed. + // PL24 background work only (backfill lane + Phase-0): pace + jitter. A + // user's fresh-decode chain is never delayed. if ( data.source === "pl24" && PL24_PACE_MS > 0 && - !(job.data as { fast?: boolean }).fast && + (!fast || backfill) && (job.name === "prefetch-children" || job.name === "prefetch-parts") ) { await new Promise((r) => setTimeout(r, jitter(PL24_PACE_MS))); @@ -480,7 +543,13 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { */ private async processInit(job: Job): Promise { const { vehicleId, source, fast = false } = job.data; - this.logger.log(`[prefetch] Init for vehicle=${vehicleId}, source=${source}`); + const chain: ChainOpts = { depthCap: job.data.depthCap, backfill: job.data.backfill }; + this.logger.log( + `[prefetch] Init for vehicle=${vehicleId}, source=${source}` + + (chain.depthCap || chain.backfill + ? ` (depthCap=${chain.depthCap ?? "-"}, backfill=${chain.backfill ? 1 : 0})` + : ""), + ); // Cooldown + time-window are enforced in process() before the daily budget // is debited — see the comment there; re-checking here would be a no-op. @@ -586,7 +655,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { if (child.unavailable) continue; // Count jobs ACTUALLY queued, not nodes walked — see addJob's doc: an // inflated progress.total makes the chain never reach "finished". - queued += await this.queueCategoryJob(child, vehicleId, source, 1, fast); + queued += await this.queueCategoryJob(child, vehicleId, source, 1, fast, chain); } } else if (this.isLeafLinkPath(cat.linkPath, cat.source, cat.hasSubgroups, cat.linkWid)) { // Leaf — check if parts already fetched @@ -606,6 +675,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { action: "parts" as const, depth: 0, fast, + ...chain, })) ) { queued++; @@ -620,6 +690,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { action: "children" as const, depth: 0, fast, + ...chain, })) ) { queued++; @@ -644,19 +715,24 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { */ private async processChildren(job: Job): Promise { const { vehicleId, categoryId, source, depth, fast = false } = job.data; + const chain: ChainOpts = { depthCap: job.data.depthCap, backfill: job.data.backfill }; this.logger.log(`[prefetch] Children for category=${categoryId}, depth=${depth}`); // Cooldown + time-window are enforced in process() before the daily budget // is debited — see the comment there; re-checking here would be a no-op. - const depthCeiling = maxDepthFor(source, fast); + const depthCeiling = maxDepthFor(source, fast, chain.depthCap); if (depth >= depthCeiling) { - // For the PL24 fast lane this is the normal stopping point, not a problem: - // deeper nodes are drilled lazily on user click or by the backfill lane. - const level = source === "pl24" && fast ? "log" : "warn"; + // For PL24 this is the normal stopping point, not a problem: deeper nodes + // are drilled lazily on user click or by the budgeted lane. The job still + // counts towards the chain's denominator (queueCategoryJob counted it), so + // settle it — otherwise the vehicle never reaches "finished" and the scan + // re-picks it every wave. + const level = source === "pl24" ? "log" : "warn"; this.logger[level]( `[prefetch] Depth ceiling ${depthCeiling} reached for category=${categoryId} (source=${source}, lane=${fast ? "fast" : "main"})`, ); + await this.incrementCompleted(vehicleId); return; } @@ -686,7 +762,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { if (child.unavailable) continue; // Real queued-job count (see addJob) — walking a node that dedupes or // already has parts must not inflate progress.total. - queued += await this.queueCategoryJob(child, vehicleId, source, depth + 1, fast); + queued += await this.queueCategoryJob(child, vehicleId, source, depth + 1, fast, chain); } if (queued > 0) { @@ -807,26 +883,19 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { // Background backfill is opt-in per source; pl24 defaults to OFF so bulk // load can never burn the one surviving account by accident. if (s === "pl24" && !isPl24BackfillEnabled()) continue; - if (await this.redis.exists(`prefetch:activity:${s}`)) continue; - if (cfg.businessHoursOnly !== false && !isWithinTimeWindow(s)) continue; - const mainLimit = this.dailyMainLimit(s); - if (mainLimit > 0) { - const spent = Number((await this.redis.get(this.dailyKey(s))) ?? 0); - if (spent >= mainLimit) { - this.logger.log( - `[backfill] ${s} daily budget spent (${spent}/${mainLimit}) — source skipped this wave`, - ); - continue; - } - } - eligible.push(s); + if (await this.sourceGatesOpen(s, cfg.businessHoursOnly)) eligible.push(s); } - if (eligible.length === 0) { + // Phase-0 (opened-partless pl24) has its own switch and does NOT need the + // bulk-backfill switch: it is user demand, capped and budgeted (see below). + const phase0Open = + PL24_OPENED_ENABLED && + (eligible.includes("pl24") || (await this.sourceGatesOpen("pl24", cfg.businessHoursOnly))); + if (eligible.length === 0 && !phase0Open) { this.logger.log("[backfill] Skip — no eligible sources (cooldown / off-hours)"); return; } - const picked: Array<{ id: string; source: string; fast: boolean }> = []; + const picked: Array<{ id: string; source: string; fast: boolean; chain?: ChainOpts }> = []; const seen = new Set(); const overfetch = batchSize * 4; // headroom for in-flight skips @@ -834,6 +903,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { v: { id: string; source: string | null }, fast: boolean, limit: number = batchSize, + chain?: ChainOpts, ): Promise => { if (picked.length >= limit || seen.has(v.id) || !v.source) return; // TWO different in-flight guards, both must be clear: @@ -856,25 +926,27 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { // Skip poison vehicles (generic-model / over-cap catalog explosions). if (await this.redis.exists(this.poisonKey(v.id))) return; seen.add(v.id); - picked.push({ id: v.id, source: v.source, fast }); + picked.push({ id: v.id, source: v.source, fast, ...(chain ? { chain } : {}) }); }; // Phase 1 — clear the obvious backlog first: decoded vehicles with zero parts. - const noParts = await this.db - .select({ id: vehicles.id, source: vehicles.source }) - .from(vehicles) - .where( - and( - inArray(vehicles.source, eligible), - notExists( - this.db.select({ one: sql`1` }).from(parts).where(eq(parts.vehicleId, vehicles.id)), + if (eligible.length > 0) { + const noParts = await this.db + .select({ id: vehicles.id, source: vehicles.source }) + .from(vehicles) + .where( + and( + inArray(vehicles.source, eligible), + notExists( + this.db.select({ one: sql`1` }).from(parts).where(eq(parts.vehicleId, vehicles.id)), + ), ), - ), - ) - .orderBy(asc(vehicles.createdAt)) - .limit(overfetch); + ) + .orderBy(asc(vehicles.createdAt)) + .limit(overfetch); - for (const v of noParts) await tryPick(v, true); + for (const v of noParts) await tryPick(v, true); + } // Phase 2 — rolling rescan of vehicles that are NOT fully fetched, to gap-fill // partials. Targets the DURABLE `fullyFetched` flag (indexed) instead of @@ -883,7 +955,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { // vehicle at once. The fullyFetchedAt clause keeps the periodic re-validation // the TTL used to provide. Gated by the ceilings above + a JOB-unit budget. const phase2Limit = Math.min(batchSize, picked.length + phase2Budget); - if (phase2Allowed && picked.length < phase2Limit) { + if (eligible.length > 0 && phase2Allowed && picked.length < phase2Limit) { const cursorObj = await this.redis.getJson<{ ts: string }>(BACKFILL_CURSOR_KEY); const cursor = cursorObj?.ts ? new Date(cursorObj.ts) : new Date(0); @@ -919,26 +991,91 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { } } + // Phase 0 — PL24 vehicles a user OPENED that still have zero parts. They got + // the fresh-decode fast lane (depth 1: groups, no leaves) and nothing since, + // because the bulk pl24 backfill is switched off. Drill them to + // PL24_MAIN_MAX_DEPTH: fast QUEUE so lifo puts them ahead of the deep + // backlog (a later fresh decode still preempts them), but `backfill: true` + // so they pace, jitter and stop at the MAIN budget threshold, leaving the + // fast-lane reserve to live decodes. Most recently opened first. Picked and + // enqueued LAST on purpose: lifo means last-in runs first. + let phase0Count = 0; + if (phase0Open) { + const since = new Date(Date.now() - PL24_OPENED_DAYS * 86_400_000); + const lastOpen = sql`max(${userVehicles.lastAccessedAt})`; + const opened = await this.db + .select({ id: vehicles.id, source: vehicles.source, lastOpen }) + .from(vehicles) + .innerJoin(userVehicles, eq(userVehicles.vehicleId, vehicles.id)) + .where( + and( + eq(vehicles.source, "pl24"), + gte(userVehicles.lastAccessedAt, since), + notExists( + this.db.select({ one: sql`1` }).from(parts).where(eq(parts.vehicleId, vehicles.id)), + ), + ), + ) + .groupBy(vehicles.id, vehicles.source) + .orderBy(desc(lastOpen)) + .limit(PL24_OPENED_PER_WAVE * 4); + const before = picked.length; + for (const v of opened) { + await tryPick(v, true, before + PL24_OPENED_PER_WAVE, { + depthCap: PL24_MAIN_MAX_DEPTH, + backfill: true, + }); + } + phase0Count = picked.length - before; + } + if (picked.length === 0) { this.logger.log("[backfill] No candidates this wave"); return; } - for (const v of picked) await this.enqueueInit(v.id, v.source, v.fast); + for (const v of picked) await this.enqueueInit(v.id, v.source, v.fast, v.chain); const fastCount = picked.filter((v) => v.fast).length; this.logger.log( `[backfill] Queued ${picked.length} vehicle(s) (${fastCount} fast-lane, ` + + `${phase0Count} opened-partless pl24, ` + `sources=${eligible.join(",")}, pressure=${depth.pressure}, total=${depth.total}, ` + `phase2Budget=${phase2Budget})`, ); } + /** + * Cooldown + scrape window + main-lane daily budget for one source — the gates + * a wave must clear before it may feed that source. Shared by the bulk + * eligibility loop and Phase-0 (which deliberately skips the bulk switch). + */ + private async sourceGatesOpen(s: string, businessHoursOnly?: boolean): Promise { + if (await this.redis.exists(`prefetch:activity:${s}`)) return false; + if (businessHoursOnly !== false && !isWithinTimeWindow(s)) return false; + const mainLimit = this.dailyMainLimit(s); + if (mainLimit > 0) { + const spent = Number((await this.redis.get(this.dailyKey(s))) ?? 0); + if (spent >= mainLimit) { + this.logger.log( + `[backfill] ${s} daily budget spent (${spent}/${mainLimit}) — source skipped this wave`, + ); + return false; + } + } + return true; + } + /** Queue a prefetch-init for a vehicle and set the in-flight guard. */ - private async enqueueInit(vehicleId: string, source: string, fast = false): Promise { + private async enqueueInit( + vehicleId: string, + source: string, + fast = false, + chain?: ChainOpts, + ): Promise { const q = fast ? this.fastQueue : this.queue; await q.add( "prefetch-init", - { vehicleId, source: source as PrefetchInitJobData["source"], fast }, + { vehicleId, source: source as PrefetchInitJobData["source"], fast, ...(chain ?? {}) }, { removeOnComplete: { count: 1000 }, removeOnFail: { count: 5000 }, @@ -966,6 +1103,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { source: string, depth: number, fast = false, + chain: ChainOpts = {}, ): Promise { if (cat.unavailable) return 0; @@ -996,17 +1134,20 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { action: "parts" as const, depth, fast, + ...chain, })) ? 1 : 0; } return 0; } - if (cat.linkPath && depth < MAX_DEPTH) { - // Non-leaf within the depth cap — explore children. The `depth < MAX_DEPTH` - // gate mirrors processChildren's early-return: without it we'd enqueue a - // prefetch-children job that processChildren just drops, burning a rate-limit - // slot on a no-op (this was ~85% of the queue at MAX_DEPTH=2). + // Non-leaf within THIS LANE's depth ceiling — explore children. The gate + // mirrors processChildren's early-return: without it we'd enqueue a + // prefetch-children job that processChildren just drops, burning a rate slot + // AND a daily-budget unit on a no-op (~85% of the queue at MAX_DEPTH=2; and + // again 426 of 3.5k parked pl24 jobs on 2026-09-27 when this compared + // against MAX_DEPTH instead of the fast lane's ceiling of 1). + if (cat.linkPath && depth < maxDepthFor(source, fast, chain.depthCap)) { const [childCheck] = await this.db .select({ id: categories.id }) .from(categories) @@ -1023,7 +1164,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { let n = 0; for (const child of children) { if (child.unavailable) continue; - n += await this.queueCategoryJob(child, vehicleId, source, depth + 1, fast); + n += await this.queueCategoryJob(child, vehicleId, source, depth + 1, fast, chain); } return n; } @@ -1034,6 +1175,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { action: "children" as const, depth, fast, + ...chain, })) ? 1 : 0; @@ -1294,6 +1436,7 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { private async checkSourceDailyBudget( source: string, lane: "main" | "fast" = "main", + kind: "parts" | "children" = "parts", ): Promise { const max = SOURCE_DAILY_MAX[source] ?? 0; if (max <= 0) return; @@ -1301,15 +1444,10 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { const dayMs = 86_400_000; const now = Date.now(); const key = this.dailyKey(source, now); + const childKey = `${key}:children`; - // READ-then-INCR (was INCR-then-check). A REJECTED attempt must not count: - // the old order inflated the counter with every defer (observed 48531 against - // a 30000 budget), which (a) made the number useless for capacity decisions - // and (b) — now that the scan reads the same counter to stop feeding a spent - // source — would let pure defer churn lock the source out. Worst-case - // overshoot under the read/incr race is WORKER_CONCURRENCY jobs: acceptable. - const n = Number((await this.redis.get(key)) ?? 0); - if (n >= limit) { + /** Defer to the next in-window slot after the UTC rollover (see below). */ + const deferToNextDay = (why: string): never => { // Defer to the next UTC day, plus up to 45min of JITTER. Without jitter every // deferred job wakes in the SAME millisecond (observed: 11495 jobs all at // 00:00:01 UTC) — the promotion lands as one burst and the pressure signal @@ -1321,16 +1459,45 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy { // other half of the deadlock fixed in process(). alignToWindow is a no-op // when no window is configured (the default). const msLeft = Math.max(1000, alignToWindow(source, rollover) - now); - if (n === limit) { + if (why) { this.logger.warn( - `[prefetch] ${source} daily budget hit (lane=${lane}, ${n}/${limit} of ${max}) — ` + - `deferring ~${Math.round(msLeft / 3_600_000)}h to the next in-window slot`, + `[prefetch] ${source} ${why} — deferring ~${Math.round(msLeft / 3_600_000)}h to the next in-window slot`, ); } throw new RateLimitError(msLeft, "source-rate"); + }; + + // Tree-expansion share: children/init jobs stop at their slice of the daily + // budget so the remainder is guaranteed to parts jobs (what users search for). + const share = SOURCE_CHILDREN_SHARE[source] ?? 0; + if (kind === "children" && share > 0) { + const childLimit = Math.floor(max * share); + const c = Number((await this.redis.get(childKey)) ?? 0); + if (c >= childLimit) { + deferToNextDay( + c === childLimit + ? `tree-expansion share spent (lane=${lane}, ${c}/${childLimit} of ${max}; rest reserved for parts)` + : "", + ); + } + } + + // READ-then-INCR (was INCR-then-check). A REJECTED attempt must not count: + // the old order inflated the counter with every defer (observed 48531 against + // a 30000 budget), which (a) made the number useless for capacity decisions + // and (b) — now that the scan reads the same counter to stop feeding a spent + // source — would let pure defer churn lock the source out. Worst-case + // overshoot under the read/incr race is WORKER_CONCURRENCY jobs: acceptable. + const n = Number((await this.redis.get(key)) ?? 0); + if (n >= limit) { + deferToNextDay(n === limit ? `daily budget hit (lane=${lane}, ${n}/${limit} of ${max})` : ""); } const after = await this.redis.incr(key); if (after === 1) await this.redis.expire(key, 90_000); // ~25h, outlives the window + if (kind === "children" && share > 0) { + const ca = await this.redis.incr(childKey); + if (ca === 1) await this.redis.expire(childKey, 90_000); + } } /** Redis key for a source's UTC-day budget counter (shared by both lanes). */ diff --git a/apps/api/src/jobs/prefetch.types.ts b/apps/api/src/jobs/prefetch.types.ts index ffede71..86fe946 100644 --- a/apps/api/src/jobs/prefetch.types.ts +++ b/apps/api/src/jobs/prefetch.types.ts @@ -11,6 +11,17 @@ export interface PrefetchInitJobData { * vehicles aren't starved behind the rolling rescan. See PrefetchWorkerService. */ fast?: boolean; + /** + * Explicit depth ceiling for this chain (Phase-0 opened-partless drills use + * PL24_MAIN_MAX_DEPTH while riding the fast queue). Unset = the lane default. + */ + depthCap?: number; + /** + * Background work that happens to ride the fast queue for ORDERING only: it + * paces/jitters like backfill and stops at the MAIN daily-budget threshold, so + * the fast-lane reserve stays with the user's fresh-decode chains. + */ + backfill?: boolean; } /** Per-category job: fetches children OR parts */ @@ -22,4 +33,8 @@ export interface PrefetchCategoryJobData { depth: number; /** Inherited from the init job — keeps the whole chain in the fast lane. */ fast?: boolean; + /** Inherited from the init job (see PrefetchInitJobData). */ + depthCap?: number; + /** Inherited from the init job (see PrefetchInitJobData). */ + backfill?: boolean; } diff --git a/docker-compose.coolify.yml b/docker-compose.coolify.yml index 0ae546e..acdda47 100644 --- a/docker-compose.coolify.yml +++ b/docker-compose.coolify.yml @@ -166,6 +166,11 @@ services: - PREFETCH_RATE_PL24=${PREFETCH_RATE_PL24:-} - PREFETCH_DAILY_PL24=${PREFETCH_DAILY_PL24:-} - PREFETCH_PL24_FAST_DEPTH=${PREFETCH_PL24_FAST_DEPTH:-} + - PREFETCH_PL24_MAX_DEPTH=${PREFETCH_PL24_MAX_DEPTH:-} + - PREFETCH_PL24_CHILDREN_SHARE=${PREFETCH_PL24_CHILDREN_SHARE:-} + - PREFETCH_PL24_OPENED_ENABLED=${PREFETCH_PL24_OPENED_ENABLED:-} + - PREFETCH_PL24_OPENED_PER_WAVE=${PREFETCH_PL24_OPENED_PER_WAVE:-} + - PREFETCH_PL24_OPENED_DAYS=${PREFETCH_PL24_OPENED_DAYS:-} - PREFETCH_PL24_DELAY_MS=${PREFETCH_PL24_DELAY_MS:-} - PL24_HTTP_DAILY_MAX=${PL24_HTTP_DAILY_MAX:-} - PL24_HTTP_USER_RESERVE=${PL24_HTTP_USER_RESERVE:-}