feat(prefetch): PL24 bütçesini parçaya kaydır — tavan-üstü no-op kırpma, %50 ağaç payı, açılan-parçasız Faz-0
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled

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:<gün>: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 <noreply@anthropic.com>
This commit is contained in:
2026-09-27 19:06:31 +03:00
parent 19d7f995b9
commit 7d439fc83b
5 changed files with 652 additions and 68 deletions

View File

@@ -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<string, string | undefined>) {
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<unknown> => 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<string, unknown> = {}) {
const incrCalls: string[] = [];
const redis = {
exists: vi.fn(async (..._a: unknown[]) => false),
get: vi.fn(async (..._a: unknown[]): Promise<string | null> => 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<unknown> => 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<string, unknown> = {};
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<string, unknown>) {
return {
name,
data,
moveToDelayed: vi.fn(async (..._a: unknown[]) => undefined),
attemptsMade: 0,
};
}
type Svc = {
process: (j: unknown, t?: string) => Promise<void>;
queueCategoryJob: (c: unknown, v: string, s: string, d: number, f: boolean) => Promise<number>;
checkSourceDailyBudget: (s: string, lane: string, kind?: string) => Promise<void>;
processBackfillScan: () => Promise<void>;
};
function build(
mod: Awaited<ReturnType<typeof load>>,
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<typeof vi.fn>).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<string, string> = {};
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<string, unknown>);
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();
});
});

View File

@@ -9,7 +9,7 @@ import { PrefetchWorkerService } from "./prefetch-worker.service";
function makeDb(limitResults: unknown[][]) {
const queued = [...limitResults];
const chain: Record<string, unknown> = {};
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() ?? []);

View File

@@ -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<string, number> = {
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<PrefetchCategoryJobData>);
}
}
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<PrefetchInitJobData>): Promise<void> {
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<PrefetchCategoryJobData>): Promise<void> {
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<string>();
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<void> => {
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<Date>`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<boolean> {
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<void> {
private async enqueueInit(
vehicleId: string,
source: string,
fast = false,
chain?: ChainOpts,
): Promise<void> {
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<number> {
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<void> {
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). */

View File

@@ -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;
}

View File

@@ -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:-}