7c95ee498e
## Summary
Stamp every Prisma span with `db.datasource: "writer" | "replica"` so
traces can distinguish which client the query went through.
Both `PrismaClient` instances share the same global
`@prisma/instrumentation`, so their spans come out with identical names
and attributes today. This makes them trivially filterable.
## How
Two pieces in `apps/webapp/app/`:
1. **`v3/tracer.server.ts`** — a `DatasourceAttributeSpanProcessor`
reads an OTel context key in `onStart` and calls
`span.setAttribute("db.datasource", value)`. Registered as the first
span processor.
2. **`db.server.ts`** — `tagDatasource(datasource, client)` wraps each
`PrismaClient` with `$extends({ query: { $allOperations } })`. The
middleware sets the context key around the query and directly tags the
active span (to catch `prisma:client:operation`, which Prisma creates
before the middleware fires).
### Context-propagation gotcha
`PrismaPromise` is lazy — `query(args)` returns a thenable that only
starts when someone `.then()`s it. The naive `context.with(ctx, () =>
query(args))` restores ALS synchronously, so when Prisma's internal code
awaits the thenable later, the engine spans fire with the original ALS.
Wrapping as `async () => await query(args)` forces the `.then()` inside
the `context.with` callback, so ALS stays on our context for the engine
spans.
### Coverage
- **Tagged**: all `prisma:engine:*` (`connection`, `db_query`,
`serialize`, `query`, etc.), `prisma:client:operation`,
`prisma:client:serialize`, `prisma:client:connect`
- **Not tagged**: `prisma:client:load_engine` — one-time startup, fires
before any query
Concurrent `Promise.all([writer.x, replica.y])` correctly tags each pool
separately (ALS isolates per-Promise chain).
### Performance
One `context.with` (~200ns) and one `setAttribute` per span (effectively
free per OTel JS benchmarks) per Prisma op. Negligible against a query
path measured in milliseconds.
## Test plan
- [ ] Verify `db.datasource` appears on `prisma:engine:connection` spans
after the webapp is restarted
- [ ] Spot-check a handful of real traces carry the attribute
415 lines
11 KiB
TypeScript
415 lines
11 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 { 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,
|
|
};
|
|
|
|
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> {
|
|
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),
|
|
(error) => {
|
|
logger.error("prisma.$transaction error", {
|
|
code: error.code,
|
|
meta: error.meta,
|
|
stack: error.stack,
|
|
message: error.message,
|
|
name: error.name,
|
|
});
|
|
},
|
|
options
|
|
);
|
|
});
|
|
} else {
|
|
return transac(
|
|
prisma,
|
|
fnOrName,
|
|
(error) => {
|
|
logger.error("prisma.$transaction error", {
|
|
code: error.code,
|
|
meta: error.meta,
|
|
stack: error.stack,
|
|
message: error.message,
|
|
name: error.name,
|
|
});
|
|
},
|
|
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", () => tagDatasource("writer", getClient()));
|
|
|
|
export const $replica: PrismaReplicaClient = singleton("replica", () => {
|
|
const replica = getReplicaClient();
|
|
return replica ? 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
|
|
client.$connect();
|
|
|
|
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
|
|
replicaClient.$connect();
|
|
|
|
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}`]);
|