Files
ditto.site/packages/worker/src/env.ts
T
Samraaj BathandClaude Fable 5 4c1dea77e1 Parallelize the clone queue: in-flight dedup + WORKER_CONCURRENCY
Two changes that multiply effective queue throughput:

- In-flight dedup: POST /v1/clones now attaches identical submits (same
  cacheKey) to the already queued/running job instead of enqueueing a
  duplicate capture. Retry-spam previously multiplied queue load (observed
  4x same-URL submissions in prod); now it's free. Best-effort — a
  same-instant race can still double-submit. Gated on !noCache.

- WORKER_CONCURRENCY (default 1): the worker registers N independent
  pg-boss pollers, running up to N captures concurrently per process with
  no head-of-line blocking. Each slot gets its own verify harness dir
  (harnessDir/slot-i) so concurrent framework builds never collide;
  concurrency 1 keeps the flat harnessDir layout.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-06 22:44:42 -07:00

33 lines
1.4 KiB
TypeScript

import { join } from "node:path";
import { parseDuration } from "./duration.js";
export type WorkerEnv = {
databaseUrl: string;
artifactsDir: string;
/** cache time-to-stale (CACHE_STALE_AFTER, default 24h; "0" disables caching). */
cacheStaleAfterMs: number;
/** this worker's isolated build harness dir for verify jobs (M5). */
harnessDir: string;
/** persistent entry-capture cache dir (the single→multi speed path). "" disables it. */
captureCacheDir: string;
tier: string;
/** concurrent clone jobs in this process (WORKER_CONCURRENCY, default 1). Each
* slot runs a full Chromium capture — size the container accordingly. */
concurrency: number;
};
export function loadWorkerEnv(): WorkerEnv {
const databaseUrl = process.env.DATABASE_URL;
if (!databaseUrl) throw new Error("DATABASE_URL is required for the worker");
const concurrency = Math.max(1, parseInt(process.env.WORKER_CONCURRENCY ?? "1", 10) || 1);
return {
databaseUrl,
artifactsDir: process.env.ARTIFACTS_DIR ?? join(process.cwd(), "local-data", "artifacts"),
cacheStaleAfterMs: parseDuration(process.env.CACHE_STALE_AFTER, 24 * 60 * 60 * 1000),
harnessDir: process.env.HARNESS_DIR ?? join(process.cwd(), "local-data", "harness"),
captureCacheDir: process.env.CAPTURE_CACHE_DIR ?? join(process.cwd(), "local-data", "capture-cache"),
tier: process.env.VERIFY_TIER ?? "stage2",
concurrency,
};
}