From ad1b512b064723abdf9532be4c40bcded27521d1 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Tue, 5 Dec 2023 11:16:10 +0000 Subject: [PATCH] Add read replica support --- apps/webapp/app/db.server.ts | 82 ++++++++++++++++++++++++++++++----- apps/webapp/app/env.server.ts | 1 + 2 files changed, 72 insertions(+), 11 deletions(-) diff --git a/apps/webapp/app/db.server.ts b/apps/webapp/app/db.server.ts index 8c4959390..aabb00901 100644 --- a/apps/webapp/app/db.server.ts +++ b/apps/webapp/app/db.server.ts @@ -69,23 +69,22 @@ export { Prisma }; export const prisma = singleton("prisma", getClient); +export const $replica: Omit = singleton( + "replica", + () => getReplicaClient() ?? prisma +); + function getClient() { const { DATABASE_URL } = process.env; invariant(typeof DATABASE_URL === "string", "DATABASE_URL env var not set"); - const databaseUrl = new URL(DATABASE_URL); - // We need to add the connection_limit and pool_timeout query params to the url, in a way that works if the DATABASE_URL already has query params - const query = databaseUrl.searchParams; - query.set("connection_limit", env.DATABASE_CONNECTION_LIMIT.toString()); - query.set("pool_timeout", env.DATABASE_POOL_TIMEOUT.toString()); - databaseUrl.search = query.toString(); + const databaseUrl = extendQueryParams(DATABASE_URL, { + connection_limit: env.DATABASE_CONNECTION_LIMIT.toString(), + pool_timeout: env.DATABASE_POOL_TIMEOUT.toString(), + }); - // Remove the username:password in the url and print that to the console - const urlWithoutCredentials = new URL(databaseUrl.href); - urlWithoutCredentials.password = ""; - - console.log(`🔌 setting up prisma client to ${urlWithoutCredentials.toString()}`); + console.log(`🔌 setting up prisma client to ${redactUrlSecrets(databaseUrl.href)}`); const client = new PrismaClient({ datasources: { @@ -127,6 +126,67 @@ function getClient() { return client; } +function getReplicaClient() { + if (!env.DATABASE_READ_REPLICA_URL) { + return; + } + + const replicaUrl = extendQueryParams(env.DATABASE_READ_REPLICA_URL, { + connection_limit: env.DATABASE_CONNECTION_LIMIT.toString(), + pool_timeout: env.DATABASE_POOL_TIMEOUT.toString(), + }); + + console.log(`🔌 setting up read replica connection to ${redactUrlSecrets(replicaUrl)}`); + + const replicaClient = new PrismaClient({ + datasources: { + db: { + url: replicaUrl.href, + }, + }, + log: [ + { + emit: "stdout", + level: "error", + }, + { + emit: "stdout", + level: "info", + }, + { + emit: "stdout", + level: "warn", + }, + ], + }); + + // connect eagerly + replicaClient.$connect(); + + console.log(`🔌 read replica connected`); + + return replicaClient; +} + +function extendQueryParams(hrefOrUrl: string | URL, queryParams: Record) { + const url = new URL(hrefOrUrl); + const query = url.searchParams; + + for (const [key, val] of Object.entries(queryParams)) { + query.set(key, val); + } + + url.search = query.toString(); + + return url; +} + +function redactUrlSecrets(hrefOrUrl: string | URL) { + const url = new URL(hrefOrUrl); + url.password = ""; + return url.href; +} + export type { PrismaClient } from "@trigger.dev/database"; export const PrismaErrorSchema = z.object({ diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 2027f880d..a6d0555ee 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -5,6 +5,7 @@ import { isValidRegex } from "./utils/regex"; const EnvironmentSchema = z.object({ NODE_ENV: z.union([z.literal("development"), z.literal("production"), z.literal("test")]), DATABASE_URL: z.string(), + DATABASE_READ_REPLICA_URL: z.string().optional(), DATABASE_CONNECTION_LIMIT: z.coerce.number().int().default(10), DATABASE_POOL_TIMEOUT: z.coerce.number().int().default(60), DIRECT_URL: z.string(),