Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .server-changes/per-client-connection-timeouts.md
Original file line number Diff line number Diff line change
@@ -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.
86 changes: 68 additions & 18 deletions apps/webapp/app/db.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,15 @@ async function $transactionInner<R>(

export { Prisma };

function tagDatasource<T extends PrismaClient>(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";
Comment thread
ericallam marked this conversation as resolved.

function tagDatasource<T extends PrismaClient>(datasource: DatasourceLabel, client: T): T {
return client.$extends({
name: "datasource-tagger",
query: {
Expand All @@ -142,7 +150,7 @@ function tagDatasource<T extends PrismaClient>(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({
Expand All @@ -168,15 +176,17 @@ 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", () => {
const replica = getReplicaClient();
// 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;
});

Expand Down Expand Up @@ -296,27 +306,43 @@ 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 —
// which this callback never touches, so the writer stays unbranded.
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,
})
)
)
),
}
Expand Down Expand Up @@ -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
Expand All @@ -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,
});

Expand Down Expand Up @@ -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
Expand All @@ -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,
});

Expand Down Expand Up @@ -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,
});

Expand Down Expand Up @@ -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,
});

Expand Down
20 changes: 19 additions & 1 deletion apps/webapp/app/env.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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")]),
Expand All @@ -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,
Comment thread
ericallam marked this conversation as resolved.
// 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.
Expand Down Expand Up @@ -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,
Comment thread
ericallam marked this conversation as resolved.
// 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.
Expand Down
4 changes: 4 additions & 0 deletions docs/self-hosting/env/webapp.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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. |
Comment thread
coderabbitai[bot] marked this conversation as resolved.
| **Redis** | | | |
| `REDIS_HOST` | Yes | — | Redis host. |
| `REDIS_PORT` | Yes | — | Redis port. |
Expand Down
Loading