feat: migrate Prisma from 6.14.0 to 7.7.0 with driver adapters
- Bump prisma, @prisma/client to 7.7.0, add @prisma/adapter-pg - Switch to engine-less client (engineType = 'client') with PrismaPg adapter - Remove binaryTargets and metrics preview feature from schema.prisma - Remove url/directUrl from datasource block (Prisma 7 requirement) - Create prisma.config.ts for CLI tools (migrations) - Rewrite db.server.ts to use PrismaPg adapter for writer + replica clients - Drop $metrics: remove from metrics.ts, delete configurePrismaMetrics from tracer.server.ts - Update PrismaClientKnownRequestError import path (runtime/library -> runtime/client) - Update all PrismaClient instantiation sites to use adapter pattern: testcontainers, tests/utils.ts, scripts, benchmark producer - Exclude prisma.config.ts from TypeScript build Co-Authored-By: Eric Allam <eallam@icloud.com>
This commit is contained in:
@@ -7,6 +7,7 @@ import {
|
||||
type PrismaTransactionClient,
|
||||
type PrismaTransactionOptions,
|
||||
} from "@trigger.dev/database";
|
||||
import { PrismaPg } from "@prisma/adapter-pg";
|
||||
import invariant from "tiny-invariant";
|
||||
import { z } from "zod";
|
||||
import { env } from "./env.server";
|
||||
@@ -109,21 +110,22 @@ 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,
|
||||
});
|
||||
const databaseUrl = new URL(DATABASE_URL);
|
||||
|
||||
// Set application_name as a query param on the connection string (pg understands this)
|
||||
databaseUrl.searchParams.set("application_name", env.SERVICE_NAME);
|
||||
|
||||
console.log(`🔌 setting up prisma client to ${redactUrlSecrets(databaseUrl)}`);
|
||||
|
||||
const adapter = new PrismaPg({
|
||||
connectionString: databaseUrl.href,
|
||||
max: env.DATABASE_CONNECTION_LIMIT,
|
||||
idleTimeoutMillis: env.DATABASE_POOL_TIMEOUT * 1000,
|
||||
connectionTimeoutMillis: env.DATABASE_CONNECTION_TIMEOUT * 1000,
|
||||
});
|
||||
|
||||
const client = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url: databaseUrl.href,
|
||||
},
|
||||
},
|
||||
adapter,
|
||||
log: [
|
||||
// events
|
||||
{
|
||||
@@ -233,21 +235,20 @@ function getReplicaClient() {
|
||||
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,
|
||||
});
|
||||
const replicaUrl = new URL(env.DATABASE_READ_REPLICA_URL);
|
||||
replicaUrl.searchParams.set("application_name", env.SERVICE_NAME);
|
||||
|
||||
console.log(`🔌 setting up read replica connection to ${redactUrlSecrets(replicaUrl)}`);
|
||||
|
||||
const adapter = new PrismaPg({
|
||||
connectionString: replicaUrl.href,
|
||||
max: env.DATABASE_CONNECTION_LIMIT,
|
||||
idleTimeoutMillis: env.DATABASE_POOL_TIMEOUT * 1000,
|
||||
connectionTimeoutMillis: env.DATABASE_CONNECTION_TIMEOUT * 1000,
|
||||
});
|
||||
|
||||
const replicaClient = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url: replicaUrl.href,
|
||||
},
|
||||
},
|
||||
adapter,
|
||||
log: [
|
||||
// events
|
||||
{
|
||||
@@ -350,19 +351,6 @@ function getReplicaClient() {
|
||||
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 = "";
|
||||
|
||||
@@ -1,5 +1,4 @@
|
||||
import { LoaderFunctionArgs } from "@remix-run/server-runtime";
|
||||
import { prisma } from "~/db.server";
|
||||
import { metricsRegister } from "~/metrics.server";
|
||||
|
||||
export async function loader({ request }: LoaderFunctionArgs) {
|
||||
@@ -13,14 +12,9 @@ export async function loader({ request }: LoaderFunctionArgs) {
|
||||
}
|
||||
}
|
||||
|
||||
// We need to remove empty lines from the prisma metrics, grafana doesn't like them
|
||||
const prismaMetrics = (await prisma.$metrics.prometheus()).replace(/^\s*[\r\n]/gm, "");
|
||||
const coreMetrics = await metricsRegister.metrics();
|
||||
|
||||
// Order matters, core metrics end with `# EOF`, prisma metrics don't
|
||||
const metrics = prismaMetrics + coreMetrics;
|
||||
|
||||
return new Response(metrics, {
|
||||
return new Response(coreMetrics, {
|
||||
headers: {
|
||||
"Content-Type": metricsRegister.contentType,
|
||||
},
|
||||
|
||||
@@ -54,9 +54,7 @@ import { LoggerSpanExporter } from "./telemetry/loggerExporter.server";
|
||||
import { CompactMetricExporter } from "./telemetry/compactMetricExporter.server";
|
||||
import { logger } from "~/services/logger.server";
|
||||
import { flattenAttributes } from "@trigger.dev/core/v3";
|
||||
import { prisma } from "~/db.server";
|
||||
import { metricsRegister } from "~/metrics.server";
|
||||
import type { Prisma } from "@trigger.dev/database";
|
||||
import { performance } from "node:perf_hooks";
|
||||
|
||||
export const SEMINTATTRS_FORCE_RECORDING = "forceRecording";
|
||||
@@ -330,221 +328,12 @@ function setupMetrics() {
|
||||
|
||||
const meter = meterProvider.getMeter("trigger.dev", "3.3.12");
|
||||
|
||||
configurePrismaMetrics({ meter });
|
||||
configureNodejsMetrics({ meter });
|
||||
configureHostMetrics({ meterProvider });
|
||||
|
||||
return meter;
|
||||
}
|
||||
|
||||
function configurePrismaMetrics({ meter }: { meter: Meter }) {
|
||||
// Counters
|
||||
const queriesTotal = meter.createObservableCounter("db.client.queries.total", {
|
||||
description: "Total number of Prisma Client queries executed",
|
||||
unit: "queries",
|
||||
});
|
||||
const datasourceQueriesTotal = meter.createObservableCounter("db.datasource.queries.total", {
|
||||
description: "Total number of datasource queries executed",
|
||||
unit: "queries",
|
||||
});
|
||||
const connectionsOpenedTotal = meter.createObservableCounter("db.pool.connections.opened.total", {
|
||||
description: "Total number of pool connections opened",
|
||||
unit: "connections",
|
||||
});
|
||||
const connectionsClosedTotal = meter.createObservableCounter("db.pool.connections.closed.total", {
|
||||
description: "Total number of pool connections closed",
|
||||
unit: "connections",
|
||||
});
|
||||
|
||||
// Gauges
|
||||
const queriesActive = meter.createObservableGauge("db.client.queries.active", {
|
||||
description: "Number of currently active Prisma Client queries",
|
||||
unit: "queries",
|
||||
});
|
||||
const queriesWait = meter.createObservableGauge("db.client.queries.wait", {
|
||||
description: "Number of queries currently waiting for a connection",
|
||||
unit: "queries",
|
||||
});
|
||||
const totalGauge = meter.createObservableGauge("db.pool.connections.total", {
|
||||
description: "Open Prisma-pool connections",
|
||||
unit: "connections",
|
||||
});
|
||||
const busyGauge = meter.createObservableGauge("db.pool.connections.busy", {
|
||||
description: "Connections currently executing queries",
|
||||
unit: "connections",
|
||||
});
|
||||
const freeGauge = meter.createObservableGauge("db.pool.connections.free", {
|
||||
description: "Idle (free) connections in the pool",
|
||||
unit: "connections",
|
||||
});
|
||||
|
||||
// Histogram statistics as gauges
|
||||
const queriesWaitTimeCount = meter.createObservableGauge("db.client.queries.wait_time.count", {
|
||||
description: "Number of wait time observations",
|
||||
unit: "observations",
|
||||
});
|
||||
const queriesWaitTimeSum = meter.createObservableGauge("db.client.queries.wait_time.sum", {
|
||||
description: "Total wait time across all observations",
|
||||
unit: "ms",
|
||||
});
|
||||
const queriesWaitTimeMean = meter.createObservableGauge("db.client.queries.wait_time.mean", {
|
||||
description: "Average wait time for a connection",
|
||||
unit: "ms",
|
||||
});
|
||||
|
||||
const queriesDurationCount = meter.createObservableGauge("db.client.queries.duration.count", {
|
||||
description: "Number of query duration observations",
|
||||
unit: "observations",
|
||||
});
|
||||
const queriesDurationSum = meter.createObservableGauge("db.client.queries.duration.sum", {
|
||||
description: "Total query duration across all observations",
|
||||
unit: "ms",
|
||||
});
|
||||
const queriesDurationMean = meter.createObservableGauge("db.client.queries.duration.mean", {
|
||||
description: "Average duration of Prisma Client queries",
|
||||
unit: "ms",
|
||||
});
|
||||
|
||||
const datasourceQueriesDurationCount = meter.createObservableGauge(
|
||||
"db.datasource.queries.duration.count",
|
||||
{
|
||||
description: "Number of datasource query duration observations",
|
||||
unit: "observations",
|
||||
}
|
||||
);
|
||||
const datasourceQueriesDurationSum = meter.createObservableGauge(
|
||||
"db.datasource.queries.duration.sum",
|
||||
{
|
||||
description: "Total datasource query duration across all observations",
|
||||
unit: "ms",
|
||||
}
|
||||
);
|
||||
const datasourceQueriesDurationMean = meter.createObservableGauge(
|
||||
"db.datasource.queries.duration.mean",
|
||||
{
|
||||
description: "Average duration of datasource queries",
|
||||
unit: "ms",
|
||||
}
|
||||
);
|
||||
|
||||
// Single helper so we hit Prisma only once per scrape ---------------------
|
||||
async function readPrismaMetrics() {
|
||||
const metrics = await prisma.$metrics.json();
|
||||
|
||||
// Extract counter values
|
||||
const counters: Record<string, number> = {};
|
||||
for (const counter of metrics.counters) {
|
||||
counters[counter.key] = counter.value;
|
||||
}
|
||||
|
||||
// Extract gauge values
|
||||
const gauges: Record<string, number> = {};
|
||||
for (const gauge of metrics.gauges) {
|
||||
gauges[gauge.key] = gauge.value;
|
||||
}
|
||||
|
||||
// Extract histogram values
|
||||
const histograms: Record<string, Prisma.MetricHistogram> = {};
|
||||
for (const histogram of metrics.histograms) {
|
||||
histograms[histogram.key] = histogram.value;
|
||||
}
|
||||
|
||||
return {
|
||||
counters: {
|
||||
queriesTotal: counters["prisma_client_queries_total"] ?? 0,
|
||||
datasourceQueriesTotal: counters["prisma_datasource_queries_total"] ?? 0,
|
||||
connectionsOpenedTotal: counters["prisma_pool_connections_opened_total"] ?? 0,
|
||||
connectionsClosedTotal: counters["prisma_pool_connections_closed_total"] ?? 0,
|
||||
},
|
||||
gauges: {
|
||||
queriesActive: gauges["prisma_client_queries_active"] ?? 0,
|
||||
queriesWait: gauges["prisma_client_queries_wait"] ?? 0,
|
||||
connectionsOpen: gauges["prisma_pool_connections_open"] ?? 0,
|
||||
connectionsBusy: gauges["prisma_pool_connections_busy"] ?? 0,
|
||||
connectionsIdle: gauges["prisma_pool_connections_idle"] ?? 0,
|
||||
},
|
||||
histograms: {
|
||||
queriesWait: histograms["prisma_client_queries_wait_histogram_ms"],
|
||||
queriesDuration: histograms["prisma_client_queries_duration_histogram_ms"],
|
||||
datasourceQueriesDuration: histograms["prisma_datasource_queries_duration_histogram_ms"],
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
meter.addBatchObservableCallback(
|
||||
async (res) => {
|
||||
const { counters, gauges, histograms } = await readPrismaMetrics();
|
||||
|
||||
// Observe counters
|
||||
res.observe(queriesTotal, counters.queriesTotal);
|
||||
res.observe(datasourceQueriesTotal, counters.datasourceQueriesTotal);
|
||||
res.observe(connectionsOpenedTotal, counters.connectionsOpenedTotal);
|
||||
res.observe(connectionsClosedTotal, counters.connectionsClosedTotal);
|
||||
|
||||
// Observe gauges
|
||||
res.observe(queriesActive, gauges.queriesActive);
|
||||
res.observe(queriesWait, gauges.queriesWait);
|
||||
res.observe(totalGauge, gauges.connectionsOpen);
|
||||
res.observe(busyGauge, gauges.connectionsBusy);
|
||||
res.observe(freeGauge, gauges.connectionsIdle);
|
||||
|
||||
// Observe histogram statistics as gauges
|
||||
if (histograms.queriesWait) {
|
||||
res.observe(queriesWaitTimeCount, histograms.queriesWait.count);
|
||||
res.observe(queriesWaitTimeSum, histograms.queriesWait.sum);
|
||||
res.observe(
|
||||
queriesWaitTimeMean,
|
||||
histograms.queriesWait.count > 0
|
||||
? histograms.queriesWait.sum / histograms.queriesWait.count
|
||||
: 0
|
||||
);
|
||||
}
|
||||
|
||||
if (histograms.queriesDuration) {
|
||||
res.observe(queriesDurationCount, histograms.queriesDuration.count);
|
||||
res.observe(queriesDurationSum, histograms.queriesDuration.sum);
|
||||
res.observe(
|
||||
queriesDurationMean,
|
||||
histograms.queriesDuration.count > 0
|
||||
? histograms.queriesDuration.sum / histograms.queriesDuration.count
|
||||
: 0
|
||||
);
|
||||
}
|
||||
|
||||
if (histograms.datasourceQueriesDuration) {
|
||||
res.observe(datasourceQueriesDurationCount, histograms.datasourceQueriesDuration.count);
|
||||
res.observe(datasourceQueriesDurationSum, histograms.datasourceQueriesDuration.sum);
|
||||
res.observe(
|
||||
datasourceQueriesDurationMean,
|
||||
histograms.datasourceQueriesDuration.count > 0
|
||||
? histograms.datasourceQueriesDuration.sum / histograms.datasourceQueriesDuration.count
|
||||
: 0
|
||||
);
|
||||
}
|
||||
},
|
||||
[
|
||||
queriesTotal,
|
||||
datasourceQueriesTotal,
|
||||
connectionsOpenedTotal,
|
||||
connectionsClosedTotal,
|
||||
queriesActive,
|
||||
queriesWait,
|
||||
totalGauge,
|
||||
busyGauge,
|
||||
freeGauge,
|
||||
queriesWaitTimeCount,
|
||||
queriesWaitTimeSum,
|
||||
queriesWaitTimeMean,
|
||||
queriesDurationCount,
|
||||
queriesDurationSum,
|
||||
queriesDurationMean,
|
||||
datasourceQueriesDurationCount,
|
||||
datasourceQueriesDurationSum,
|
||||
datasourceQueriesDurationMean,
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
function configureNodejsMetrics({ meter }: { meter: Meter }) {
|
||||
if (!env.INTERNAL_OTEL_NODEJS_METRICS_ENABLED) {
|
||||
return;
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
*/
|
||||
|
||||
import { PrismaClient } from "@trigger.dev/database";
|
||||
import { PrismaPg } from "@prisma/adapter-pg";
|
||||
import { performance } from "node:perf_hooks";
|
||||
|
||||
interface ProducerConfig {
|
||||
@@ -91,13 +92,8 @@ function generateError() {
|
||||
}
|
||||
|
||||
async function runProducer(config: ProducerConfig) {
|
||||
const prisma = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url: config.postgresUrl,
|
||||
},
|
||||
},
|
||||
});
|
||||
const adapter = new PrismaPg(config.postgresUrl);
|
||||
const prisma = new PrismaClient({ adapter });
|
||||
|
||||
try {
|
||||
console.log(
|
||||
|
||||
@@ -5,12 +5,13 @@
|
||||
"main": "./dist/index.js",
|
||||
"types": "./dist/index.d.ts",
|
||||
"dependencies": {
|
||||
"@prisma/client": "6.14.0",
|
||||
"@prisma/adapter-pg": "7.7.0",
|
||||
"@prisma/client": "7.7.0",
|
||||
"decimal.js": "^10.6.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/decimal.js": "^7.4.3",
|
||||
"prisma": "6.14.0",
|
||||
"prisma": "7.7.0",
|
||||
"rimraf": "6.0.1"
|
||||
},
|
||||
"scripts": {
|
||||
@@ -25,4 +26,4 @@
|
||||
"build": "pnpm run clean && tsc --noEmit false --outDir dist --declaration",
|
||||
"dev": "tsc --noEmit false --outDir dist --declaration --watch"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
import path from "node:path";
|
||||
import { defineConfig } from "prisma/config";
|
||||
|
||||
export default defineConfig({
|
||||
schema: path.join(__dirname, "prisma", "schema.prisma"),
|
||||
engine: "classic",
|
||||
datasource: {
|
||||
url: process.env.DATABASE_URL ?? "postgresql://localhost:5432/trigger",
|
||||
directUrl: process.env.DIRECT_URL,
|
||||
},
|
||||
});
|
||||
@@ -1,14 +1,11 @@
|
||||
datasource db {
|
||||
provider = "postgresql"
|
||||
url = env("DATABASE_URL")
|
||||
directUrl = env("DIRECT_URL")
|
||||
provider = "postgresql"
|
||||
}
|
||||
|
||||
generator client {
|
||||
provider = "prisma-client-js"
|
||||
output = "../generated/prisma"
|
||||
binaryTargets = ["native", "debian-openssl-1.1.x"]
|
||||
previewFeatures = ["metrics"]
|
||||
provider = "prisma-client-js"
|
||||
output = "../generated/prisma"
|
||||
engineType = "client"
|
||||
}
|
||||
|
||||
model User {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { PrismaClient } from "../generated/prisma";
|
||||
import { Decimal } from "decimal.js";
|
||||
import { PrismaClientKnownRequestError } from "@prisma/client/runtime/library";
|
||||
import { PrismaClientKnownRequestError } from "@prisma/client/runtime/client";
|
||||
|
||||
// Define the isolation levels manually
|
||||
type TransactionIsolationLevel =
|
||||
|
||||
@@ -9,5 +9,5 @@
|
||||
"noEmit": true,
|
||||
"strict": true
|
||||
},
|
||||
"exclude": ["node_modules"]
|
||||
"exclude": ["node_modules", "prisma.config.ts"]
|
||||
}
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
"dependencies": {
|
||||
"@clickhouse/client": "^1.11.1",
|
||||
"@opentelemetry/api": "^1.9.0",
|
||||
"@prisma/adapter-pg": "7.7.0",
|
||||
"@trigger.dev/database": "workspace:*",
|
||||
"ioredis": "^5.3.2"
|
||||
},
|
||||
@@ -21,4 +22,4 @@
|
||||
"scripts": {
|
||||
"typecheck": "tsc --noEmit"
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import { StartedPostgreSqlContainer } from "@testcontainers/postgresql";
|
||||
import { StartedRedisContainer } from "@testcontainers/redis";
|
||||
import { PrismaClient } from "@trigger.dev/database";
|
||||
import { PrismaPg } from "@prisma/adapter-pg";
|
||||
import { RedisOptions } from "ioredis";
|
||||
import { Network, type StartedNetwork } from "testcontainers";
|
||||
import { TaskContext, test } from "vitest";
|
||||
@@ -106,13 +107,8 @@ export const prisma = async (
|
||||
|
||||
console.log("Initializing Prisma with URL:", url);
|
||||
|
||||
const prisma = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url,
|
||||
},
|
||||
},
|
||||
});
|
||||
const adapter = new PrismaPg(url);
|
||||
const prisma = new PrismaClient({ adapter });
|
||||
try {
|
||||
await use(prisma);
|
||||
} finally {
|
||||
|
||||
Generated
+494
-15
File diff suppressed because it is too large
Load Diff
@@ -46,6 +46,7 @@
|
||||
*/
|
||||
|
||||
import { PrismaClient, TaskRunExecutionStatus } from "@trigger.dev/database";
|
||||
import { PrismaPg } from "@prisma/adapter-pg";
|
||||
import { createRedisClient } from "@internal/redis";
|
||||
|
||||
interface StuckRun {
|
||||
@@ -71,18 +72,20 @@ async function main() {
|
||||
const [environmentId, postgresUrl, redisReadUrl, redisWriteUrl] = process.argv.slice(2);
|
||||
|
||||
if (!environmentId || !postgresUrl || !redisReadUrl) {
|
||||
console.error("Usage: tsx scripts/recover-stuck-runs.ts <environmentId> <postgresUrl> <redisReadUrl> [redisWriteUrl]");
|
||||
console.error(
|
||||
"Usage: tsx scripts/recover-stuck-runs.ts <environmentId> <postgresUrl> <redisReadUrl> [redisWriteUrl]"
|
||||
);
|
||||
console.error("");
|
||||
console.error("Dry-run mode when no redisWriteUrl is provided (read-only).");
|
||||
console.error("Execute mode when redisWriteUrl is provided (makes actual changes).");
|
||||
console.error("");
|
||||
console.error("Example (dry-run):");
|
||||
console.error(' tsx scripts/recover-stuck-runs.ts env_1234567890 \\');
|
||||
console.error(" tsx scripts/recover-stuck-runs.ts env_1234567890 \\");
|
||||
console.error(' "postgresql://user:pass@localhost:5432/triggerdev" \\');
|
||||
console.error(' "redis://readonly.example.com:6379"');
|
||||
console.error("");
|
||||
console.error("Example (execute):");
|
||||
console.error(' tsx scripts/recover-stuck-runs.ts env_1234567890 \\');
|
||||
console.error(" tsx scripts/recover-stuck-runs.ts env_1234567890 \\");
|
||||
console.error(' "postgresql://user:pass@localhost:5432/triggerdev" \\');
|
||||
console.error(' "redis://readonly.example.com:6379" \\');
|
||||
console.error(' "redis://writeonly.example.com:6379"');
|
||||
@@ -100,13 +103,8 @@ async function main() {
|
||||
console.log(`🔍 Scanning for stuck runs in environment: ${environmentId}`);
|
||||
|
||||
// Create Prisma client with the provided connection URL
|
||||
const prisma = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url: postgresUrl,
|
||||
},
|
||||
},
|
||||
});
|
||||
const adapter = new PrismaPg(postgresUrl);
|
||||
const prisma = new PrismaClient({ adapter });
|
||||
|
||||
try {
|
||||
// Get environment details
|
||||
@@ -259,7 +257,9 @@ async function main() {
|
||||
}
|
||||
|
||||
// Prepare recovery operations
|
||||
console.log(`\n⚡ ${executeMode ? "Executing" : "Planning"} recovery for ${stuckRuns.length} stuck runs`);
|
||||
console.log(
|
||||
`\n⚡ ${executeMode ? "Executing" : "Planning"} recovery for ${stuckRuns.length} stuck runs`
|
||||
);
|
||||
console.log(`This will:`);
|
||||
console.log(` 1. Add each run back to its specific queue sorted set`);
|
||||
console.log(` 2. Remove each run from the queue-specific currentConcurrency set`);
|
||||
|
||||
+3
-8
@@ -1,18 +1,13 @@
|
||||
import { PrismaClient } from "@trigger.dev/database";
|
||||
import { PrismaPg } from "@prisma/adapter-pg";
|
||||
|
||||
type SetDBCallback = (prisma: PrismaClient) => Promise<void>;
|
||||
|
||||
export const setDB = async (cb: SetDBCallback) => {
|
||||
const { DATABASE_URL } = process.env;
|
||||
|
||||
const prisma = new PrismaClient({
|
||||
datasources: {
|
||||
db: {
|
||||
url: DATABASE_URL,
|
||||
// We can't set directUrl here, and we don't have to
|
||||
},
|
||||
},
|
||||
});
|
||||
const adapter = new PrismaPg(DATABASE_URL ?? "postgresql://localhost:5432/trigger");
|
||||
const prisma = new PrismaClient({ adapter });
|
||||
|
||||
await prisma.$connect();
|
||||
await cb(prisma);
|
||||
|
||||
Reference in New Issue
Block a user