Files
triggerdotdev--trigger.dev/apps/webapp/app/db.server.ts
Chris Arderne 4fde283e76 chore: format and lint webapp also (#4056)
#3977 added formatting and linting everywhere else.

This extends it to the webapp.
2026-06-26 13:02:53 +01:00

442 lines
12 KiB
TypeScript

import {
Prisma,
PrismaClient,
$transaction as transac,
type PrismaClientOrTransaction,
type PrismaReplicaClient,
type PrismaTransactionClient,
type PrismaTransactionOptions,
} from "@trigger.dev/database";
import invariant from "tiny-invariant";
import { z } from "zod";
import { env } from "./env.server";
import { logger } from "./services/logger.server";
import { isValidDatabaseUrl } from "./utils/db";
import {
captureInfrastructureErrors,
infraErrorAlreadyLogged,
logTransactionInfrastructureError,
} from "./utils/prismaErrors";
import { singleton } from "./utils/singleton";
import { DATASOURCE_CONTEXT_KEY, startActiveSpan } from "./v3/tracer.server";
import { context, Span, trace } from "@opentelemetry/api";
import { queryPerformanceMonitor } from "./utils/queryPerformanceMonitor.server";
export type {
PrismaTransactionClient,
PrismaClientOrTransaction,
PrismaTransactionOptions,
PrismaReplicaClient,
};
// Boundary logger for transac(): skips an error the client extension already
// logged (and tagged) at the statement level, so a single failure is logged
// once. Shared by both $transaction overloads so the guard can't drift.
function logTransactionPrismaError(error: Prisma.PrismaClientKnownRequestError) {
if (infraErrorAlreadyLogged(error)) {
return;
}
logger.error("prisma.$transaction error", {
code: error.code,
meta: error.meta,
stack: error.stack,
message: error.message,
name: error.name,
});
}
export async function $transaction<R>(
prisma: PrismaClientOrTransaction,
name: string,
fn: (prisma: PrismaTransactionClient, span?: Span) => Promise<R>,
options?: PrismaTransactionOptions
): Promise<R | undefined>;
export async function $transaction<R>(
prisma: PrismaClientOrTransaction,
fn: (prisma: PrismaTransactionClient) => Promise<R>,
options?: PrismaTransactionOptions
): Promise<R | undefined>;
export async function $transaction<R>(
prisma: PrismaClientOrTransaction,
fnOrName: ((prisma: PrismaTransactionClient) => Promise<R>) | string,
fnOrOptions?: ((prisma: PrismaTransactionClient) => Promise<R>) | PrismaTransactionOptions,
options?: PrismaTransactionOptions
): Promise<R | undefined> {
try {
return await $transactionInner(prisma, fnOrName, fnOrOptions, options);
} catch (error) {
// transac()'s callback only logs coded Prisma errors; infra errors such as
// PrismaClientInitializationError reach the boundary without a `.code`.
logTransactionInfrastructureError(error);
throw error;
}
}
async function $transactionInner<R>(
prisma: PrismaClientOrTransaction,
fnOrName: ((prisma: PrismaTransactionClient) => Promise<R>) | string,
fnOrOptions?: ((prisma: PrismaTransactionClient) => Promise<R>) | PrismaTransactionOptions,
options?: PrismaTransactionOptions
): Promise<R | undefined> {
if (typeof fnOrName === "string") {
return await startActiveSpan(fnOrName, async (span) => {
span.setAttribute("$transaction", true);
if (options?.isolationLevel) {
span.setAttribute("isolation_level", options.isolationLevel);
}
if (options?.timeout) {
span.setAttribute("timeout", options.timeout);
}
if (options?.maxWait) {
span.setAttribute("max_wait", options.maxWait);
}
if (options?.swallowPrismaErrors) {
span.setAttribute("swallow_prisma_errors", options.swallowPrismaErrors);
}
const fn = fnOrOptions as (prisma: PrismaTransactionClient, span: Span) => Promise<R>;
return transac(prisma, (client) => fn(client, span), logTransactionPrismaError, options);
});
} else {
return transac(
prisma,
fnOrName,
logTransactionPrismaError,
typeof fnOrOptions === "function" ? undefined : fnOrOptions
);
}
}
export { Prisma };
function tagDatasource<T extends PrismaClient>(datasource: "writer" | "replica", client: T): T {
return client.$extends({
name: "datasource-tagger",
query: {
$allOperations: ({ query, args }) => {
trace.getActiveSpan()?.setAttribute("db.datasource", datasource);
return context.with(
context.active().setValue(DATASOURCE_CONTEXT_KEY, datasource),
async () => await query(args)
);
},
},
}) as unknown as T;
}
export const prisma = singleton("prisma", () =>
captureInfrastructureErrors(tagDatasource("writer", getClient()))
);
export const $replica: PrismaReplicaClient = singleton("replica", () => {
const replica = getReplicaClient();
return replica ? captureInfrastructureErrors(tagDatasource("replica", replica)) : prisma;
});
function getClient() {
const { DATABASE_URL } = process.env;
invariant(typeof DATABASE_URL === "string", "DATABASE_URL env var not set");
const databaseUrl = extendQueryParams(DATABASE_URL, {
connection_limit: env.DATABASE_CONNECTION_LIMIT.toString(),
pool_timeout: env.DATABASE_POOL_TIMEOUT.toString(),
connection_timeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
application_name: env.SERVICE_NAME,
});
console.log(`🔌 setting up prisma client to ${redactUrlSecrets(databaseUrl)}`);
const client = new PrismaClient({
datasources: {
db: {
url: databaseUrl.href,
},
},
log: [
// events
{
emit: "event",
level: "error",
},
{
emit: "event",
level: "info",
},
{
emit: "event",
level: "warn",
},
// stdout
...((process.env.PRISMA_LOG_TO_STDOUT === "1"
? [
{
emit: "stdout",
level: "error",
},
{
emit: "stdout",
level: "info",
},
{
emit: "stdout",
level: "warn",
},
]
: []) satisfies Prisma.LogDefinition[]),
// Query performance monitoring
...((process.env.VERBOSE_PRISMA_LOGS === "1" ||
process.env.VERY_SLOW_QUERY_THRESHOLD_MS !== undefined
? [
{
emit: "event",
level: "query",
},
]
: []) satisfies Prisma.LogDefinition[]),
// verbose
...((process.env.VERBOSE_PRISMA_LOGS === "1"
? [
{
emit: "stdout",
level: "query",
},
]
: []) satisfies Prisma.LogDefinition[]),
],
});
// Only use structured logging if we're not already logging to stdout
if (process.env.PRISMA_LOG_TO_STDOUT !== "1") {
client.$on("info", (log) => {
logger.info("PrismaClient info", {
clientType: "writer",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
});
});
client.$on("warn", (log) => {
logger.warn("PrismaClient warn", {
clientType: "writer",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
});
});
client.$on("error", (log) => {
logger.error("PrismaClient error", {
clientType: "writer",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
ignoreError: true,
});
});
}
// Add query performance monitoring
client.$on("query", (log) => {
queryPerformanceMonitor.onQuery("writer", log);
});
// Connect eagerly; Prisma will connect on use anyway.
// Swallow the error when testing (DB likely unavailable)
const connectPromise = client.$connect();
if (env.NODE_ENV === "test") {
connectPromise.catch((error) => {
logger.warn("Failed to eagerly connect prisma client (writer)", { error });
});
}
console.log(`🔌 prisma client connected`);
return client;
}
function getReplicaClient() {
if (!env.DATABASE_READ_REPLICA_URL) {
console.log(`🔌 No database replica, using the regular client`);
return;
}
const replicaUrl = extendQueryParams(env.DATABASE_READ_REPLICA_URL, {
connection_limit: env.DATABASE_CONNECTION_LIMIT.toString(),
pool_timeout: env.DATABASE_POOL_TIMEOUT.toString(),
connection_timeout: env.DATABASE_CONNECTION_TIMEOUT.toString(),
application_name: env.SERVICE_NAME,
});
console.log(`🔌 setting up read replica connection to ${redactUrlSecrets(replicaUrl)}`);
const replicaClient = new PrismaClient({
datasources: {
db: {
url: replicaUrl.href,
},
},
log: [
// events
{
emit: "event",
level: "error",
},
{
emit: "event",
level: "info",
},
{
emit: "event",
level: "warn",
},
// stdout
...((process.env.PRISMA_LOG_TO_STDOUT === "1"
? [
{
emit: "stdout",
level: "error",
},
{
emit: "stdout",
level: "info",
},
{
emit: "stdout",
level: "warn",
},
]
: []) satisfies Prisma.LogDefinition[]),
// Query performance monitoring
...((process.env.VERBOSE_PRISMA_LOGS === "1" ||
process.env.VERY_SLOW_QUERY_THRESHOLD_MS !== undefined
? [
{
emit: "event",
level: "query",
},
]
: []) satisfies Prisma.LogDefinition[]),
// verbose
...((process.env.VERBOSE_PRISMA_LOGS === "1"
? [
{
emit: "stdout",
level: "query",
},
]
: []) satisfies Prisma.LogDefinition[]),
],
});
// Only use structured logging if we're not already logging to stdout
if (process.env.PRISMA_LOG_TO_STDOUT !== "1") {
replicaClient.$on("info", (log) => {
logger.info("PrismaClient info", {
clientType: "reader",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
});
});
replicaClient.$on("warn", (log) => {
logger.warn("PrismaClient warn", {
clientType: "reader",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
});
});
replicaClient.$on("error", (log) => {
logger.error("PrismaClient error", {
clientType: "reader",
event: {
timestamp: log.timestamp,
message: log.message,
target: log.target,
},
});
});
}
// Add query performance monitoring for replica client
replicaClient.$on("query", (log) => {
queryPerformanceMonitor.onQuery("replica", log);
});
// Connect eagerly; Prisma will connect on use anyway.
// Swallow the error when testing (DB likely unavailable)
const connectPromise = replicaClient.$connect();
if (env.NODE_ENV === "test") {
connectPromise.catch((error) => {
logger.warn("Failed to eagerly connect prisma client (replica)", { error });
});
}
console.log(`🔌 read replica connected`);
return replicaClient;
}
function extendQueryParams(hrefOrUrl: string | URL, queryParams: Record<string, string>) {
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({
code: z.string(),
});
function getDatabaseSchema() {
if (!isValidDatabaseUrl(env.DATABASE_URL)) {
throw new Error("Invalid Database URL");
}
const databaseUrl = new URL(env.DATABASE_URL);
const schemaFromSearchParam = databaseUrl.searchParams.get("schema");
if (!schemaFromSearchParam) {
console.debug("❗ database schema unspecified, will default to `public` schema");
return "public";
}
return schemaFromSearchParam;
}
export const DATABASE_SCHEMA = singleton("DATABASE_SCHEMA", getDatabaseSchema);
export const sqlDatabaseSchema = Prisma.sql([`${DATABASE_SCHEMA}`]);