perf(backfill): stop the prefetch worker from throttling itself (Tier 1)
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
Some checks failed
QA Gate (P0/P1) / Test affected app (pull_request) Has been cancelled
The worker's own upstream fetches called touchActivity(), setting the prefetch:activity:<source> cooldown key (TTL 300s) that checkCooldown then honoured — so after each fetch the worker paused itself for up to ~5 minutes (re-checking every 60s, ~5 empty cycles per key). At ~1 fetch / 5 min the 3387-job backlog needed ~6 days to drain. - Wrap each worker job in an AsyncLocalStorage backfill context; touchActivity skips the cooldown key when invoked from the worker, so the cooldown reflects only real user requests (worker yields to users, never to itself). - Cooldown TTL 300s -> 90s (a request 5 min ago isn't "active"). - checkCooldown pauses for the key's actual remaining TTL (one wait) instead of a fixed 60s re-check loop. No extra upstream load — only removes the worker's self-imposed idle time. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -20,6 +20,7 @@ import {
|
||||
import { ConfigService } from "@nestjs/config";
|
||||
|
||||
import { ProxyAgent } from "undici";
|
||||
import { isBackfillContext } from "../../jobs/prefetch-context";
|
||||
import { RedisService } from "../../redis/redis.service";
|
||||
import { EmexBrowserService } from "./emex.browser";
|
||||
import { createEmptyDecodedVehicle, mapEmexResponse } from "./emex.mapper";
|
||||
@@ -143,8 +144,11 @@ export class EmexService {
|
||||
|
||||
/** Mark EMEX as actively used (5min TTL) to defer prefetch worker */
|
||||
private async touchActivity(): Promise<void> {
|
||||
// Worker-originated fetches must NOT register as user activity, or the
|
||||
// backfill worker throttles itself via checkCooldown.
|
||||
if (isBackfillContext()) return;
|
||||
try {
|
||||
await this.redis.set("prefetch:activity:emex", String(Date.now()), 300);
|
||||
await this.redis.set("prefetch:activity:emex", String(Date.now()), 90);
|
||||
} catch {
|
||||
// Non-critical — don't break the request
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
*/
|
||||
|
||||
import { Injectable, Logger } from "@nestjs/common";
|
||||
import { isBackfillContext } from "../../jobs/prefetch-context";
|
||||
import { RedisService } from "../../redis/redis.service";
|
||||
import { PartsCatalogsAuthService } from "./parts-catalogs-auth.service";
|
||||
import {
|
||||
@@ -43,8 +44,11 @@ export class PartsCatalogsService {
|
||||
|
||||
/** Mark parts-catalogs as actively used (5min TTL) to defer prefetch worker */
|
||||
private async touchActivity(): Promise<void> {
|
||||
// Worker-originated fetches must NOT register as user activity, or the
|
||||
// backfill worker throttles itself via checkCooldown.
|
||||
if (isBackfillContext()) return;
|
||||
try {
|
||||
await this.redis.set("prefetch:activity:parts-catalogs", String(Date.now()), 300);
|
||||
await this.redis.set("prefetch:activity:parts-catalogs", String(Date.now()), 90);
|
||||
} catch {
|
||||
// Non-critical
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ import {
|
||||
ServiceUnavailableException,
|
||||
} from "@nestjs/common";
|
||||
import { ConfigService } from "@nestjs/config";
|
||||
import { isBackfillContext } from "../../jobs/prefetch-context";
|
||||
import { RedisService } from "../../redis/redis.service";
|
||||
import { StorageService } from "../../storage/storage.service";
|
||||
import { PL24AuthService } from "./pl24-auth.service";
|
||||
@@ -827,8 +828,11 @@ export class PL24Service {
|
||||
|
||||
/** Mark PL24 as actively used (5min TTL) to defer prefetch worker */
|
||||
private async touchActivity(): Promise<void> {
|
||||
// Worker-originated fetches must NOT register as user activity, or the
|
||||
// backfill worker throttles itself via checkCooldown.
|
||||
if (isBackfillContext()) return;
|
||||
try {
|
||||
await this.redis.set("prefetch:activity:pl24", String(Date.now()), 300);
|
||||
await this.redis.set("prefetch:activity:pl24", String(Date.now()), 90);
|
||||
} catch {
|
||||
// Non-critical — don't break the request
|
||||
}
|
||||
|
||||
21
apps/api/src/jobs/prefetch-context.ts
Normal file
21
apps/api/src/jobs/prefetch-context.ts
Normal file
@@ -0,0 +1,21 @@
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
|
||||
/**
|
||||
* Marks the current async call-chain as originating from the catalog backfill /
|
||||
* prefetch worker.
|
||||
*
|
||||
* Source services (pl24 / emex / parts-catalogs) call `touchActivity()` on every
|
||||
* fetch to set a `prefetch:activity:<source>` cooldown key, which the worker's
|
||||
* `checkCooldown()` then honours by pausing. That key is meant to signal *user*
|
||||
* activity so backfill yields to live users — but the worker's OWN fetches set it
|
||||
* too, so the worker throttles itself (fetch → set 300s key → next job sees it →
|
||||
* pause → repeat). Wrapping each worker job in this context lets `touchActivity()`
|
||||
* skip the key when the call originates from the worker, so the cooldown reflects
|
||||
* only real user requests.
|
||||
*/
|
||||
export const backfillContext = new AsyncLocalStorage<true>();
|
||||
|
||||
/** True when executing inside a backfill / prefetch worker job. */
|
||||
export function isBackfillContext(): boolean {
|
||||
return backfillContext.getStore() === true;
|
||||
}
|
||||
@@ -41,15 +41,19 @@ export class RateLimitError extends Error {
|
||||
}
|
||||
|
||||
/**
|
||||
* Check if a user is actively using the source.
|
||||
* Throws RateLimitError (1min retry) if cooldown key exists.
|
||||
* Check if a *user* is actively using the source (the cooldown key is set only
|
||||
* by user-facing requests; worker fetches skip it via the backfill context).
|
||||
* Throws RateLimitError (retry after the key's remaining lifetime) if set.
|
||||
*/
|
||||
export async function checkCooldown(redis: RedisService, source: string): Promise<void> {
|
||||
const key = `prefetch:activity:${source}`;
|
||||
const exists = await redis.exists(key);
|
||||
if (exists) {
|
||||
const retryMs = source === "parts-catalogs" ? 120_000 : 60_000;
|
||||
throw new RateLimitError(retryMs, "cooldown");
|
||||
// ttl() returns the key's remaining lifetime in seconds (-2 = no key,
|
||||
// -1 = no expiry). Pause for exactly that long so the worker waits once —
|
||||
// until the real user-activity window ends — instead of re-checking every
|
||||
// fixed 60s and burning ~5 empty cycles per key lifetime.
|
||||
const ttl = await redis.ttl(key);
|
||||
if (ttl > 0) {
|
||||
throw new RateLimitError(ttl * 1000, "cooldown");
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ import { DATABASE, type Database } from "../database/database.provider";
|
||||
import { categories, parts, vehicles } from "../database/schema/core";
|
||||
import { RedisService } from "../redis/redis.service";
|
||||
import { QUEUE_NAMES, getBullConnection } from "./bull.config";
|
||||
import { backfillContext } from "./prefetch-context";
|
||||
import {
|
||||
RateLimitError,
|
||||
checkCooldown,
|
||||
@@ -59,7 +60,9 @@ export class PrefetchWorkerService implements OnModuleInit, OnModuleDestroy {
|
||||
onModuleInit() {
|
||||
this.worker = new Worker(
|
||||
QUEUE_NAMES.CATALOG_PREFETCH,
|
||||
(job, token) => this.process(job, token),
|
||||
// Run every job inside the backfill context so source services' touchActivity()
|
||||
// skips the user-activity cooldown key — the worker must not throttle itself.
|
||||
(job, token) => backfillContext.run(true, () => this.process(job, token)),
|
||||
{
|
||||
connection: getBullConnection(),
|
||||
concurrency: 1,
|
||||
|
||||
Reference in New Issue
Block a user