From de14062dc76ed6ce85b7626962872bd6ee60fc63 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Wed, 5 Aug 2026 16:20:24 +0100 Subject: [PATCH] feat(webapp): per-client database pool and connect timeout overrides Adds optional per-client env overrides for the Prisma pool_timeout and connect_timeout, one pair each for the writer and read replica of the control-plane, legacy run-ops, and run-ops databases, falling back to the shared DATABASE_POOL_TIMEOUT / DATABASE_CONNECTION_TIMEOUT when unset. This lets each database be tuned independently, e.g. a fail-fast connect timeout on one without changing the others. It also tags each client's queries with its specific datasource (control-plane, legacy-run-ops, or run-ops; writer or replica) so telemetry can attribute connection behavior per database. No behavior change until an override is set. --- .../per-client-connection-timeouts.md | 6 ++ apps/webapp/app/db.server.ts | 86 +++++++++++++++---- apps/webapp/app/env.server.ts | 20 ++++- docs/self-hosting/env/webapp.mdx | 4 + 4 files changed, 97 insertions(+), 19 deletions(-) create mode 100644 .server-changes/per-client-connection-timeouts.md diff --git a/.server-changes/per-client-connection-timeouts.md b/.server-changes/per-client-connection-timeouts.md new file mode 100644 index 00000000000..60da30a6554 --- /dev/null +++ b/.server-changes/per-client-connection-timeouts.md @@ -0,0 +1,6 @@ +--- +area: webapp +type: improvement +--- + +Allow the database connection pool and connect timeouts to be tuned separately for each database's writer and read replica, falling back to the shared defaults when unset. diff --git a/apps/webapp/app/db.server.ts b/apps/webapp/app/db.server.ts index 002a84d6bce..5525697f673 100644 --- a/apps/webapp/app/db.server.ts +++ b/apps/webapp/app/db.server.ts @@ -124,7 +124,15 @@ async function $transactionInner( export { Prisma }; -function tagDatasource(datasource: "writer" | "replica", client: T): T { +type DatasourceLabel = + | "control-plane-writer" + | "control-plane-replica" + | "legacy-run-ops-writer" + | "legacy-run-ops-replica" + | "run-ops-writer" + | "run-ops-replica"; + +function tagDatasource(datasource: DatasourceLabel, client: T): T { return client.$extends({ name: "datasource-tagger", query: { @@ -142,7 +150,7 @@ function tagDatasource(datasource: "writer" | "replica", // Same extension as tagDatasource but typed for RunOpsPrismaClient (different // generated package — does not extend @trigger.dev/database.PrismaClient). function tagDatasourceRunOps( - datasource: "writer" | "replica", + datasource: DatasourceLabel, client: RunOpsPrismaClient ): RunOpsPrismaClient { return client.$extends({ @@ -168,7 +176,7 @@ function captureInfraErrorsRunOps(client: RunOpsPrismaClient): RunOpsPrismaClien } export const prisma = singleton("prisma", () => - captureInfrastructureErrors(tagDatasource("writer", getClient())) + captureInfrastructureErrors(tagDatasource("control-plane-writer", getClient())) ); export const $replica: PrismaReplicaClient = singleton("replica", () => { @@ -176,7 +184,9 @@ export const $replica: PrismaReplicaClient = singleton("replica", () => { // Brand ONLY a real replica so the run-store routing layer keeps replica reads off the primary. // No replica configured → fall back to the writer `prisma`, which must stay UNBRANDED. return replica - ? markReadReplicaClient(captureInfrastructureErrors(tagDatasource("replica", replica))) + ? markReadReplicaClient( + captureInfrastructureErrors(tagDatasource("control-plane-replica", replica)) + ) : prisma; }); @@ -296,7 +306,7 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => { controlPlane: { writer: prisma, replica: $replica }, buildNewWriter: (url, clientType) => captureInfraErrorsRunOps( - tagDatasourceRunOps("writer", buildRunOpsWriterClient({ url, clientType })) + tagDatasourceRunOps("run-ops-writer", buildRunOpsWriterClient({ url, clientType })) ), // Brand the run-ops replica (only built for a real replica URL) so routed replica reads stay // off the primary. When no replica URL is set, selectRunOpsTopology reuses the writer here — @@ -304,19 +314,35 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => { buildNewReplica: (url, clientType) => markReadReplicaClient( captureInfraErrorsRunOps( - tagDatasourceRunOps("replica", buildRunOpsReplicaClient({ url, clientType })) + tagDatasourceRunOps("run-ops-replica", buildRunOpsReplicaClient({ url, clientType })) ) ), // Legacy client shares the exact control-plane wrapper stack (the legacy DB carries the full // control-plane schema); markReadReplicaClient only on a real replica URL, as with the NEW replica. buildLegacyWriter: (url, clientType) => captureInfrastructureErrors( - tagDatasource("writer", buildWriterClient({ url, clientType })) + tagDatasource( + "legacy-run-ops-writer", + buildWriterClient({ + url, + clientType, + poolTimeout: env.RUN_OPS_LEGACY_DATABASE_WRITER_POOL_TIMEOUT, + connectTimeout: env.RUN_OPS_LEGACY_DATABASE_WRITER_CONNECTION_TIMEOUT, + }) + ) ), buildLegacyReplica: (url, clientType) => markReadReplicaClient( captureInfrastructureErrors( - tagDatasource("replica", buildReplicaClient({ url, clientType })) + tagDatasource( + "legacy-run-ops-replica", + buildReplicaClient({ + url, + clientType, + poolTimeout: env.RUN_OPS_LEGACY_DATABASE_READ_REPLICA_POOL_TIMEOUT, + connectTimeout: env.RUN_OPS_LEGACY_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT, + }) + ) ) ), } @@ -383,7 +409,12 @@ function getClient() { const url = env.CONTROL_PLANE_DATABASE_URL ?? env.DATABASE_URL; invariant(typeof url === "string", "neither CONTROL_PLANE_DATABASE_URL nor DATABASE_URL is set"); - return buildWriterClient({ url, clientType: "writer" }); + return buildWriterClient({ + url, + clientType: "writer", + poolTimeout: env.DATABASE_WRITER_POOL_TIMEOUT, + connectTimeout: env.DATABASE_WRITER_CONNECTION_TIMEOUT, + }); } // Generalized writer builder shared by the control-plane client and the run-ops @@ -392,14 +423,18 @@ function getClient() { export function buildWriterClient({ url, clientType, + poolTimeout, + connectTimeout, }: { url: string; clientType: string; + poolTimeout?: number; + connectTimeout?: number; }): PrismaClient { const databaseUrl = buildPrismaConnectionUrl(url, { connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(), - poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(), - connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(), + poolTimeout: (poolTimeout ?? env.DATABASE_POOL_TIMEOUT).toString(), + connectTimeout: (connectTimeout ?? env.DATABASE_CONNECTION_TIMEOUT).toString(), applicationName: env.SERVICE_NAME, }); @@ -530,7 +565,12 @@ function getReplicaClient() { return; } - return buildReplicaClient({ url, clientType: "reader" }); + return buildReplicaClient({ + url, + clientType: "reader", + poolTimeout: env.DATABASE_READ_REPLICA_POOL_TIMEOUT, + connectTimeout: env.DATABASE_READ_REPLICA_CONNECTION_TIMEOUT, + }); } // Generalized replica builder shared by the control-plane replica and the run-ops @@ -539,14 +579,18 @@ function getReplicaClient() { export function buildReplicaClient({ url, clientType, + poolTimeout, + connectTimeout, }: { url: string; clientType: string; + poolTimeout?: number; + connectTimeout?: number; }): PrismaClient { const replicaUrl = buildPrismaConnectionUrl(url, { connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(), - poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(), - connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(), + poolTimeout: (poolTimeout ?? env.DATABASE_POOL_TIMEOUT).toString(), + connectTimeout: (connectTimeout ?? env.DATABASE_CONNECTION_TIMEOUT).toString(), applicationName: env.SERVICE_NAME, }); @@ -675,8 +719,10 @@ function buildRunOpsWriterClient({ }): RunOpsPrismaClient { const databaseUrl = buildPrismaConnectionUrl(url, { connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(), - poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(), - connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(), + poolTimeout: (env.RUN_OPS_DATABASE_WRITER_POOL_TIMEOUT ?? env.DATABASE_POOL_TIMEOUT).toString(), + connectTimeout: ( + env.RUN_OPS_DATABASE_WRITER_CONNECTION_TIMEOUT ?? env.DATABASE_CONNECTION_TIMEOUT + ).toString(), applicationName: env.SERVICE_NAME, }); @@ -728,8 +774,12 @@ function buildRunOpsReplicaClient({ connectionLimit: ( env.RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_LIMIT ?? env.DATABASE_CONNECTION_LIMIT ).toString(), - poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(), - connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(), + poolTimeout: ( + env.RUN_OPS_DATABASE_READ_REPLICA_POOL_TIMEOUT ?? env.DATABASE_POOL_TIMEOUT + ).toString(), + connectTimeout: ( + env.RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT ?? env.DATABASE_CONNECTION_TIMEOUT + ).toString(), applicationName: env.SERVICE_NAME, }); diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 541ff36fa6e..295b3a75daf 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -108,6 +108,12 @@ const isNotInsecureSecret = (value: string) => const INSECURE_SECRET_MESSAGE = "must not be a known-insecure published default; set a strong, unique value. If you cannot rotate it yet (e.g. it protects existing encrypted data or active sessions), set ALLOW_INSECURE_DEFAULT_SECRETS=1 to boot while you migrate."; +/** Optional int env var; blank/whitespace normalises to undefined (z.coerce turns "" into 0). */ +const OptionalIntEnv = z.preprocess( + (v) => (typeof v === "string" && v.trim() === "" ? undefined : v), + z.coerce.number().int().optional() +); + const EnvironmentSchema = z .object({ NODE_ENV: z.union([z.literal("development"), z.literal("production"), z.literal("test")]), @@ -120,6 +126,10 @@ const EnvironmentSchema = z DATABASE_CONNECTION_LIMIT: z.coerce.number().int().default(10), DATABASE_POOL_TIMEOUT: z.coerce.number().int().default(60), DATABASE_CONNECTION_TIMEOUT: z.coerce.number().int().default(20), + DATABASE_WRITER_POOL_TIMEOUT: OptionalIntEnv, + DATABASE_WRITER_CONNECTION_TIMEOUT: OptionalIntEnv, + DATABASE_READ_REPLICA_POOL_TIMEOUT: OptionalIntEnv, + DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: OptionalIntEnv, // Dashboard-agent conversation store. Cloud points this at a dedicated // database; when unset it falls back to DATABASE_URL (OSS), where // the tables live in the isolated `trigger_dashboard_agent` schema. @@ -185,7 +195,15 @@ const EnvironmentSchema = z .refine(isValidDatabaseUrl, "RUN_OPS_LEGACY_DATABASE_READ_REPLICA_URL is invalid") .optional(), // Optional cap for the unpooled new run-ops read replica. Unset falls back to DATABASE_CONNECTION_LIMIT. - RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_LIMIT: z.coerce.number().int().optional(), + RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_LIMIT: OptionalIntEnv, + RUN_OPS_DATABASE_WRITER_POOL_TIMEOUT: OptionalIntEnv, + RUN_OPS_DATABASE_WRITER_CONNECTION_TIMEOUT: OptionalIntEnv, + RUN_OPS_DATABASE_READ_REPLICA_POOL_TIMEOUT: OptionalIntEnv, + RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: OptionalIntEnv, + RUN_OPS_LEGACY_DATABASE_WRITER_POOL_TIMEOUT: OptionalIntEnv, + RUN_OPS_LEGACY_DATABASE_WRITER_CONNECTION_TIMEOUT: OptionalIntEnv, + RUN_OPS_LEGACY_DATABASE_READ_REPLICA_POOL_TIMEOUT: OptionalIntEnv, + RUN_OPS_LEGACY_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: OptionalIntEnv, // Direct DSN for applying the full @trigger.dev/database migrations to the LEGACY run-ops DB, keeping // its schema current after the control plane moves off it. Direct, not pooled — migrations never run // over a pooler. Optional; unset -> the entrypoint's legacy migrate step is skipped. diff --git a/docs/self-hosting/env/webapp.mdx b/docs/self-hosting/env/webapp.mdx index 82dd6958a60..f8820886e49 100644 --- a/docs/self-hosting/env/webapp.mdx +++ b/docs/self-hosting/env/webapp.mdx @@ -26,7 +26,11 @@ mode: "wide" | `DATABASE_CONNECTION_LIMIT` | No | 10 | Max DB connections. | | `DATABASE_POOL_TIMEOUT` | No | 60 | DB pool timeout (s). | | `DATABASE_CONNECTION_TIMEOUT` | No | 20 | DB connect timeout (s). | +| `DATABASE_WRITER_POOL_TIMEOUT` | No | `DATABASE_POOL_TIMEOUT` | Writer pool timeout (s); overrides the shared default for the writer only. | +| `DATABASE_WRITER_CONNECTION_TIMEOUT` | No | `DATABASE_CONNECTION_TIMEOUT` | Writer connect timeout (s); overrides the shared default for the writer only. | | `DATABASE_READ_REPLICA_URL` | No | `DATABASE_URL` | Read-replica DB string. | +| `DATABASE_READ_REPLICA_POOL_TIMEOUT` | No | `DATABASE_POOL_TIMEOUT` | Read-replica pool timeout (s); overrides the shared default for the replica only. | +| `DATABASE_READ_REPLICA_CONNECTION_TIMEOUT` | No | `DATABASE_CONNECTION_TIMEOUT` | Read-replica connect timeout (s); overrides the shared default for the replica only. | | **Redis** | | | | | `REDIS_HOST` | Yes | — | Redis host. | | `REDIS_PORT` | Yes | — | Redis port. |