b1b9321ad2
* WIP job run performance improvements - Added a `perf` tool to better measure job run performance under heavy load - Removed `runFinished` job (not really needed) - startQueuedRuns now uses a jobKey with replace - Fixed an issue with ZodWorker when using jobKey * Publish improvement docker images * fixed the improvement docker publishing * Downgrade back to prisma 4.16.0 because 5.1.x broke docker builds * Changes to how queued runs work - Split the worker into two different workers, one dedicated to performRunExecution - Schedule performRunExecution in a single place, with a queue and using a round robin manually controlled concurrency - Remove startQueuedRuns - All runs are queued before they are started - Setting the worker maxPoolSize to the same as the worker concurrency - Starting to be able to split the docker image * Remove queue name from startRun graphile job * Make the prisma connection pool stuff configurable through env vars * Hardcode (for now) the max concurrent runs limit * Rewrite performRunExecution to be more performant PerformRunExecutionV2: - Does not create and manage jobRunExecution records - Does not reimplement retrying, uses graphile worker retrying instead I’ve kept around PerformRunExecutionV1 so this works when deploying. Definitely needs LOTS of testing * Fix issues with cached tasks - Limit the size of the cached tasks sent when executing a run, using the knapsack problem dynamic programming approach - Actually USE the cached tasks in IO by using the idempotencyKey instead of the task ID - Remove output from all logs - Added a stress test job catalog * Forgot to commit the logger updates * Never log connectionString * Login to docker hub to get around rate limits * Add additional logging to the graphile workers * Fix the *_ENABLED env vars * Allow adding and removing jobs to be done from the webapp * Don’t set the job to failed if it’s being retried * Deprecated queue options in the job and removed startPosition. Now using the job/env combo as the job queue name * Dequeung jobs doesn’t check if the runner is initialized * Fixed issues with retrying a run getting stuck on a cancelled task, and errors from parsing the results of dequeing a job * Remove queued round robin thing that isn’t used anymore * Added slack to job catalog * Better forwards compat * Added long delay * Fixed lock file
118 lines
3.2 KiB
TypeScript
118 lines
3.2 KiB
TypeScript
import { MISSING_CONNECTION_RESOLVED_NOTIFICATION } from "@trigger.dev/core";
|
|
import { PrismaClientOrTransaction, prisma } from "~/db.server";
|
|
import { IngestSendEvent } from "../events/ingestSendEvent.server";
|
|
import { logger } from "../logger.server";
|
|
import { workerQueue } from "../worker.server";
|
|
|
|
export class IntegrationConnectionCreatedService {
|
|
#prismaClient: PrismaClientOrTransaction;
|
|
|
|
constructor(prismaClient: PrismaClientOrTransaction = prisma) {
|
|
this.#prismaClient = prismaClient;
|
|
}
|
|
|
|
public async call(id: string) {
|
|
logger.debug("IntegrationConnectionCreatedService.call", { id });
|
|
|
|
// first, deliver the event through the dispatcher
|
|
const connection = await this.#prismaClient.integrationConnection.findUniqueOrThrow({
|
|
where: {
|
|
id,
|
|
},
|
|
include: {
|
|
externalAccount: true,
|
|
integration: true,
|
|
},
|
|
});
|
|
|
|
const missingConnection = await this.#prismaClient.missingConnection.findUnique({
|
|
where: {
|
|
integrationId_connectionType_accountIdentifier: {
|
|
integrationId: connection.integrationId,
|
|
connectionType: connection.connectionType,
|
|
accountIdentifier: connection.externalAccount
|
|
? connection.externalAccount.id
|
|
: "DEVELOPER",
|
|
},
|
|
},
|
|
include: {
|
|
runs: {
|
|
include: {
|
|
queue: true,
|
|
environment: {
|
|
include: {
|
|
project: true,
|
|
organization: true,
|
|
},
|
|
},
|
|
},
|
|
orderBy: {
|
|
createdAt: "asc",
|
|
},
|
|
},
|
|
integration: true,
|
|
externalAccount: true,
|
|
},
|
|
});
|
|
|
|
if (!missingConnection) {
|
|
return;
|
|
}
|
|
|
|
if (missingConnection.resolved) {
|
|
return;
|
|
}
|
|
|
|
const firstRun = missingConnection.runs[0];
|
|
|
|
if (!firstRun) {
|
|
return;
|
|
}
|
|
|
|
const eventId = `${missingConnection.id}-resolved`;
|
|
|
|
const eventService = new IngestSendEvent();
|
|
|
|
await eventService.call(firstRun.environment, {
|
|
id: eventId,
|
|
name: MISSING_CONNECTION_RESOLVED_NOTIFICATION,
|
|
payload: {
|
|
id: missingConnection.id,
|
|
type: missingConnection.connectionType,
|
|
client: {
|
|
id: missingConnection.integration.slug,
|
|
title: missingConnection.integration.title,
|
|
scopes: missingConnection.integration.scopes,
|
|
createdAt: missingConnection.integration.createdAt,
|
|
updatedAt: missingConnection.integration.updatedAt,
|
|
},
|
|
expiresAt: connection.expiresAt ?? undefined,
|
|
account: missingConnection.externalAccount
|
|
? {
|
|
id: missingConnection.externalAccount.identifier,
|
|
metadata: missingConnection.externalAccount.metadata,
|
|
}
|
|
: undefined,
|
|
},
|
|
context: {},
|
|
});
|
|
|
|
await this.#prismaClient.missingConnection.delete({
|
|
where: {
|
|
id: missingConnection.id,
|
|
},
|
|
});
|
|
|
|
for (const run of missingConnection.runs) {
|
|
logger.debug("[IntegrationConnectionCreatedService] restarting run", {
|
|
run,
|
|
});
|
|
|
|
// We need to start the run again
|
|
await workerQueue.enqueue("startRun", {
|
|
id: run.id,
|
|
});
|
|
}
|
|
}
|
|
}
|