Skip to content

Commit 9dbcc5f

Browse files
committed
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.
1 parent 3039bc1 commit 9dbcc5f

3 files changed

Lines changed: 86 additions & 18 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: improvement
4+
---
5+
6+
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.

apps/webapp/app/db.server.ts

Lines changed: 68 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,15 @@ async function $transactionInner<R>(
124124

125125
export { Prisma };
126126

127-
function tagDatasource<T extends PrismaClient>(datasource: "writer" | "replica", client: T): T {
127+
type DatasourceLabel =
128+
| "control-plane-writer"
129+
| "control-plane-replica"
130+
| "legacy-run-ops-writer"
131+
| "legacy-run-ops-replica"
132+
| "run-ops-writer"
133+
| "run-ops-replica";
134+
135+
function tagDatasource<T extends PrismaClient>(datasource: DatasourceLabel, client: T): T {
128136
return client.$extends({
129137
name: "datasource-tagger",
130138
query: {
@@ -142,7 +150,7 @@ function tagDatasource<T extends PrismaClient>(datasource: "writer" | "replica",
142150
// Same extension as tagDatasource but typed for RunOpsPrismaClient (different
143151
// generated package — does not extend @trigger.dev/database.PrismaClient).
144152
function tagDatasourceRunOps(
145-
datasource: "writer" | "replica",
153+
datasource: DatasourceLabel,
146154
client: RunOpsPrismaClient
147155
): RunOpsPrismaClient {
148156
return client.$extends({
@@ -168,15 +176,17 @@ function captureInfraErrorsRunOps(client: RunOpsPrismaClient): RunOpsPrismaClien
168176
}
169177

170178
export const prisma = singleton("prisma", () =>
171-
captureInfrastructureErrors(tagDatasource("writer", getClient()))
179+
captureInfrastructureErrors(tagDatasource("control-plane-writer", getClient()))
172180
);
173181

174182
export const $replica: PrismaReplicaClient = singleton("replica", () => {
175183
const replica = getReplicaClient();
176184
// Brand ONLY a real replica so the run-store routing layer keeps replica reads off the primary.
177185
// No replica configured → fall back to the writer `prisma`, which must stay UNBRANDED.
178186
return replica
179-
? markReadReplicaClient(captureInfrastructureErrors(tagDatasource("replica", replica)))
187+
? markReadReplicaClient(
188+
captureInfrastructureErrors(tagDatasource("control-plane-replica", replica))
189+
)
180190
: prisma;
181191
});
182192

@@ -296,27 +306,43 @@ const runOpsTopology: RunOpsTopology = singleton("runOpsTopology", () => {
296306
controlPlane: { writer: prisma, replica: $replica },
297307
buildNewWriter: (url, clientType) =>
298308
captureInfraErrorsRunOps(
299-
tagDatasourceRunOps("writer", buildRunOpsWriterClient({ url, clientType }))
309+
tagDatasourceRunOps("run-ops-writer", buildRunOpsWriterClient({ url, clientType }))
300310
),
301311
// Brand the run-ops replica (only built for a real replica URL) so routed replica reads stay
302312
// off the primary. When no replica URL is set, selectRunOpsTopology reuses the writer here —
303313
// which this callback never touches, so the writer stays unbranded.
304314
buildNewReplica: (url, clientType) =>
305315
markReadReplicaClient(
306316
captureInfraErrorsRunOps(
307-
tagDatasourceRunOps("replica", buildRunOpsReplicaClient({ url, clientType }))
317+
tagDatasourceRunOps("run-ops-replica", buildRunOpsReplicaClient({ url, clientType }))
308318
)
309319
),
310320
// Legacy client shares the exact control-plane wrapper stack (the legacy DB carries the full
311321
// control-plane schema); markReadReplicaClient only on a real replica URL, as with the NEW replica.
312322
buildLegacyWriter: (url, clientType) =>
313323
captureInfrastructureErrors(
314-
tagDatasource("writer", buildWriterClient({ url, clientType }))
324+
tagDatasource(
325+
"legacy-run-ops-writer",
326+
buildWriterClient({
327+
url,
328+
clientType,
329+
poolTimeout: env.RUN_OPS_LEGACY_DATABASE_WRITER_POOL_TIMEOUT,
330+
connectTimeout: env.RUN_OPS_LEGACY_DATABASE_WRITER_CONNECTION_TIMEOUT,
331+
})
332+
)
315333
),
316334
buildLegacyReplica: (url, clientType) =>
317335
markReadReplicaClient(
318336
captureInfrastructureErrors(
319-
tagDatasource("replica", buildReplicaClient({ url, clientType }))
337+
tagDatasource(
338+
"legacy-run-ops-replica",
339+
buildReplicaClient({
340+
url,
341+
clientType,
342+
poolTimeout: env.RUN_OPS_LEGACY_DATABASE_READ_REPLICA_POOL_TIMEOUT,
343+
connectTimeout: env.RUN_OPS_LEGACY_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT,
344+
})
345+
)
320346
)
321347
),
322348
}
@@ -383,7 +409,12 @@ function getClient() {
383409
const url = env.CONTROL_PLANE_DATABASE_URL ?? env.DATABASE_URL;
384410
invariant(typeof url === "string", "neither CONTROL_PLANE_DATABASE_URL nor DATABASE_URL is set");
385411

386-
return buildWriterClient({ url, clientType: "writer" });
412+
return buildWriterClient({
413+
url,
414+
clientType: "writer",
415+
poolTimeout: env.DATABASE_WRITER_POOL_TIMEOUT,
416+
connectTimeout: env.DATABASE_WRITER_CONNECTION_TIMEOUT,
417+
});
387418
}
388419

389420
// Generalized writer builder shared by the control-plane client and the run-ops
@@ -392,14 +423,18 @@ function getClient() {
392423
export function buildWriterClient({
393424
url,
394425
clientType,
426+
poolTimeout,
427+
connectTimeout,
395428
}: {
396429
url: string;
397430
clientType: string;
431+
poolTimeout?: number;
432+
connectTimeout?: number;
398433
}): PrismaClient {
399434
const databaseUrl = buildPrismaConnectionUrl(url, {
400435
connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(),
401-
poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(),
402-
connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
436+
poolTimeout: (poolTimeout ?? env.DATABASE_POOL_TIMEOUT).toString(),
437+
connectTimeout: (connectTimeout ?? env.DATABASE_CONNECTION_TIMEOUT).toString(),
403438
applicationName: env.SERVICE_NAME,
404439
});
405440

@@ -530,7 +565,12 @@ function getReplicaClient() {
530565
return;
531566
}
532567

533-
return buildReplicaClient({ url, clientType: "reader" });
568+
return buildReplicaClient({
569+
url,
570+
clientType: "reader",
571+
poolTimeout: env.DATABASE_READ_REPLICA_POOL_TIMEOUT,
572+
connectTimeout: env.DATABASE_READ_REPLICA_CONNECTION_TIMEOUT,
573+
});
534574
}
535575

536576
// Generalized replica builder shared by the control-plane replica and the run-ops
@@ -539,14 +579,18 @@ function getReplicaClient() {
539579
export function buildReplicaClient({
540580
url,
541581
clientType,
582+
poolTimeout,
583+
connectTimeout,
542584
}: {
543585
url: string;
544586
clientType: string;
587+
poolTimeout?: number;
588+
connectTimeout?: number;
545589
}): PrismaClient {
546590
const replicaUrl = buildPrismaConnectionUrl(url, {
547591
connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(),
548-
poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(),
549-
connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
592+
poolTimeout: (poolTimeout ?? env.DATABASE_POOL_TIMEOUT).toString(),
593+
connectTimeout: (connectTimeout ?? env.DATABASE_CONNECTION_TIMEOUT).toString(),
550594
applicationName: env.SERVICE_NAME,
551595
});
552596

@@ -675,8 +719,10 @@ function buildRunOpsWriterClient({
675719
}): RunOpsPrismaClient {
676720
const databaseUrl = buildPrismaConnectionUrl(url, {
677721
connectionLimit: env.DATABASE_CONNECTION_LIMIT.toString(),
678-
poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(),
679-
connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
722+
poolTimeout: (env.RUN_OPS_DATABASE_WRITER_POOL_TIMEOUT ?? env.DATABASE_POOL_TIMEOUT).toString(),
723+
connectTimeout: (
724+
env.RUN_OPS_DATABASE_WRITER_CONNECTION_TIMEOUT ?? env.DATABASE_CONNECTION_TIMEOUT
725+
).toString(),
680726
applicationName: env.SERVICE_NAME,
681727
});
682728

@@ -728,8 +774,12 @@ function buildRunOpsReplicaClient({
728774
connectionLimit: (
729775
env.RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_LIMIT ?? env.DATABASE_CONNECTION_LIMIT
730776
).toString(),
731-
poolTimeout: env.DATABASE_POOL_TIMEOUT.toString(),
732-
connectTimeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
777+
poolTimeout: (
778+
env.RUN_OPS_DATABASE_READ_REPLICA_POOL_TIMEOUT ?? env.DATABASE_POOL_TIMEOUT
779+
).toString(),
780+
connectTimeout: (
781+
env.RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT ?? env.DATABASE_CONNECTION_TIMEOUT
782+
).toString(),
733783
applicationName: env.SERVICE_NAME,
734784
});
735785

apps/webapp/app/env.server.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,10 @@ const EnvironmentSchema = z
120120
DATABASE_CONNECTION_LIMIT: z.coerce.number().int().default(10),
121121
DATABASE_POOL_TIMEOUT: z.coerce.number().int().default(60),
122122
DATABASE_CONNECTION_TIMEOUT: z.coerce.number().int().default(20),
123+
DATABASE_WRITER_POOL_TIMEOUT: z.coerce.number().int().optional(),
124+
DATABASE_WRITER_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
125+
DATABASE_READ_REPLICA_POOL_TIMEOUT: z.coerce.number().int().optional(),
126+
DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
123127
// Dashboard-agent conversation store. Cloud points this at a dedicated
124128
// database; when unset it falls back to DATABASE_URL (OSS), where
125129
// the tables live in the isolated `trigger_dashboard_agent` schema.
@@ -186,6 +190,14 @@ const EnvironmentSchema = z
186190
.optional(),
187191
// Optional cap for the unpooled new run-ops read replica. Unset falls back to DATABASE_CONNECTION_LIMIT.
188192
RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_LIMIT: z.coerce.number().int().optional(),
193+
RUN_OPS_DATABASE_WRITER_POOL_TIMEOUT: z.coerce.number().int().optional(),
194+
RUN_OPS_DATABASE_WRITER_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
195+
RUN_OPS_DATABASE_READ_REPLICA_POOL_TIMEOUT: z.coerce.number().int().optional(),
196+
RUN_OPS_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
197+
RUN_OPS_LEGACY_DATABASE_WRITER_POOL_TIMEOUT: z.coerce.number().int().optional(),
198+
RUN_OPS_LEGACY_DATABASE_WRITER_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
199+
RUN_OPS_LEGACY_DATABASE_READ_REPLICA_POOL_TIMEOUT: z.coerce.number().int().optional(),
200+
RUN_OPS_LEGACY_DATABASE_READ_REPLICA_CONNECTION_TIMEOUT: z.coerce.number().int().optional(),
189201
// Direct DSN for applying the full @trigger.dev/database migrations to the LEGACY run-ops DB, keeping
190202
// its schema current after the control plane moves off it. Direct, not pooled — migrations never run
191203
// over a pooler. Optional; unset -> the entrypoint's legacy migrate step is skipped.

0 commit comments

Comments
 (0)