Files
triggerdotdev--trigger.dev/apps/webapp/app/v3/runEnginePendingVersionLookup.server.ts
Eric Allam 61ca40b4b1 perf(run-engine,webapp): look up PENDING_VERSION runs via ClickHouse (#3707)
## Summary

When a background worker registers, the engine resolves runs that were
queued before the worker was ready (status `PENDING_VERSION`). That
lookup used to scan a Postgres status index on `TaskRun`. Move it to
ClickHouse: query candidate run ids from `task_runs_v2`, then refetch
the actual rows from Postgres by primary key with a `status =
'PENDING_VERSION'` guard for idempotency.

## Design

The lookup is a pluggable interface on the run engine
(`PendingVersionRunIdLookup`). The webapp wires a ClickHouse-backed
implementation through the org-scoped `clickhouseFactory` using a new
`"engine"` client type, configured by `RUN_ENGINE_CLICKHOUSE_*` env
vars. The URL falls back to `CLICKHOUSE_URL` when unset, so self-hosted
deployments don't need new config to keep working.

When the lookup returns no candidates, one bounded retry is scheduled
~5s later to cover ClickHouse replication lag against `task_runs_v2`.
The Postgres status guard on both the candidate refetch and the inner
`updateMany` prevents double-promotion when a retry races with a
concurrent deploy.

Tests cover three existing PENDING_VERSION cases via a small
Postgres-backed test adapter; new ClickHouse-backed integration tests
will follow.
2026-05-22 17:18:41 +01:00

25 lines
1.1 KiB
TypeScript

import { type PendingVersionRunIdLookup } from "@internal/run-engine";
import { clickhouseFactory } from "~/services/clickhouse/clickhouseFactoryInstance.server";
import { logger } from "~/services/logger.server";
import { singleton } from "~/utils/singleton";
import { ClickhousePendingVersionLookup } from "./services/clickhousePendingVersionLookup.server";
/**
* Lookup used by `@internal/run-engine`'s `PendingVersionSystem` to find
* `PENDING_VERSION` TaskRun ids via ClickHouse, removing the need for
* Postgres index #13 (`TaskRun_status_runtimeEnvironmentId_createdAt_id_idx`).
*
* Resolves the ClickHouse client per call via {@link clickhouseFactory}
* using the `"engine"` client type, configured by `RUN_ENGINE_CLICKHOUSE_*`
* env vars and routed per-organization for customers with HIPAA / data
* sovereignty data stores.
*/
export const runEnginePendingVersionLookup = singleton(
"runEnginePendingVersionLookup",
initializeRunEnginePendingVersionLookup
);
function initializeRunEnginePendingVersionLookup(): PendingVersionRunIdLookup {
return new ClickhousePendingVersionLookup({ clickhouseFactory, logger });
}