e97704d904
* bump worker version * Suggested glossary for the RunEngine, TBC * Removed BatchTaskRun changes from this branch, they were done in main * Set the BatchTaskRun status to completed when all runs are completed * When dequeuing respect passed in maxResources * Ported over the new run props: idempotencyKeyExpiresAt, versions, oneTimeUseToken, maxDurationInSeconds * Didn’t hit save… the new props when triggering tasks passed through * Idempotency expiration + waitpoint edge case * WIP on creating checkpoint, parking for now * fix worker routes * upgrade webapp node types to support generic event emitter * separate event bus handler singleton and run failure alerts * duration waits * fix execution snapshot debug spans * task waits * fix event bus types * temporary fix for react hook run handle type * disable run notifications for now * convert any typecasts to expect errors to more easily fix later * fix webapp types after node types upgrade * updateEnvConcurrencyLimits across marqs and the runqueue * Pass proper values into the run engine * RunQueue settings and removed unused rebalancing workers * Remove rebalancing prop * Tidied more things up * Update/remove queue limits for MARQS and RunQueue * taskQueue/concurrencyLimit changes ported back into the RunEngine * Reworked completing waitpoints to improve performance and reduce race conditions * Improved test robustness * Down to a single run lock only when a run is totally unblocked and ready to continue * warm starts, worker notifications, wait fixes * Fix for Run Engine poll interval env var * Expect the waitpoint to be completed quickly * If a run is locked then it’s too late to expire it * Added VALKEY_ env vars and plugged them into the run engine * Extracted and updated the guard queue function so it can be used when batching * Added logging and universal concurrency changes to trigger task v1 * Added notes back in * Bump @trigger.dev/worker to 3.3.7 * reportInvocationUsage for the runAttemptStarted event * improve execution snapshot span debug span start times * Unfriendly IDs * update lockfile * Created a shared determineEngineVersion function * disable unfinished commands * save new cli config to different location, misc fixes * add basic engine version check via current deploy * new run engine will default to node 22 runtime * block some actions for projects on previous run engine * fix worker group tests * fix triggerAndWait test * one typescript version to rule them all * redlock type patch * fix type issues caused by ts-reset * improve cleanup scripts * add missing socket.io dep * fix run notification handler type * fix worker group test again * generate prisma client for e2e tests * remove worker group tests for now * prevent image pull rate limits during unit tests * increase timeout for queue concurrency limit test * generate prisma client for preview release * same node types everywhere * Updated engine readme, removed legacy system notes * use default machine preset from platform package * worker instances plural in schema * disable pnpm update notifications * return worker group details from connect call * add workers admin route * fix heartbeat route return type * move deployment labels to core apps * refactor run controller env schema * Add firstAttemptStartedAt to TaskRun * RunEngine 2.0 batch trigger support (#1581) * Make it clear when BatchTriggerV2Service is used * Copy of BatchTriggerV2Service * WIP batch triggering * Allow blocking a run with multiple waitpoints at once. Made it atomic * Removed unused param * New batch service * Pass through the parentRunId and resumeParentOnCompletion * Use the new batch service, and correct trigger task version * Force V1 engine if using BatchTriggerV2Service, we’ve already done the check at this point * Removed the $transaction and early exit if nothing changed * Adedd a simple batch task to the hello world reference catalog * Fix for batch waits not working * Added parentRunId in a couple more places * Removed waitForBatch log * Added another parentRunId * Expanded the example to include all the different triggers * More changes to blocking to support continuing after idempotent completed runs * Fix for the wrong type when blocking a run * remove @map * optimise worker auth query * add engine version header to core api client requests * remove unique constraint for default group id * consolidate migrations * the first managed worker becomes the global default * Debug events off by default, added an admin toggle to show them * worker group name can't be an empty string * add exec helper to core * move machine resources to core * add pre-dequeue callback to determine max resources * optionally skip dequeue * bump worker package * move worker to core * fix ReadableStream type error * fix another type issue * update a few more tsconfigs * add metadata changes introduced in #1563 * Run Engine 2.0 trigger idempotency (#1613) * Return isCached from the trigger API endpoint * Fix for the wrong type when blocking a run * Render the idempotent run in the inspector * Event repository for idempotency * Debug events off by default, added an admin toggle to show them * triggerAndWait idempotency span * Some improvements to the reference idempotency task * Removed the cached tracing from the SDK * Server-side creating cached span * Improved idempotency test task * Create cached task spans in a better way * Idempotency span support inc batch trigger * Simplified how the spans are done, using more of the existing code * Improved the idempotency test task * Added Waitpoint Batch type, add to TaskRunWaitpoint with order * Pass batch ids through to the run engine when triggering * Added batchIndex * Better batch support in the run engine * Added settings to batch trigger service, before major overhaul * Allow the longer run/batch ids in the filters * Changed how batching works, includes breaking changes in CLI * Removed batch idempotency because it gets put on the runs instead * Added `runs` to the batch.retrieve call/API * Set firstAttemptStartedAt when creating the first attempt * Do nothing when receiving a BATCH waitpoint * Some fixes in the new batch trigger service… mostly just passing missing optional params through * Tweaked the idempotency test task for more situations * Only block with a batch if it’s a batchTriggerAndWait… 🤦♂️ * Added another case to the idempotency test task: multiple of the same idempotencyKey in a single batch * Support for the same run multiple times in the same batch * Small tweaks * Make sure to complete batches, even if they’re not andWait ones * Export RunDuplicateIdempotencyKeyError from the run engine * Latest lockfile * Trigger with a machine (old run engine) * RE2, allow setting machine when triggering * Fix for new glob patterns * add max run count to dequeue from version route * add worker instance name env var and header * queue consumer pre skip callback * poll for more runs after final execution errors * fix dequeue search param schema * add shortcut to debug switch * expose run engine timeouts as env vars * make warm start durations configurable * add optional status to json reply helper * fix preSkip hook, add debug logs * BLOCKED_BY_WAITPOINTS -> SUSPENDED * exit controller when run suspended * check if already replied before http reply * run controller will wait for next run after the current one is suspended * cancel run button shortcut * minimal event repository environment type * fix update metadata call * run suspension and misc fixes wip * change debug shortcut to shift + D * Started work on the Dev supervisor * Formatting * Fix for bad imports * Before rebuilding SSE * Presence updating from the CLI working via SSE * add worker notification debug logs * send run:stop when exiting run phase * skip current snapshot poll on worker notification * add more logs and route to submit run debug logs * add worker and runner ids to snapshots * improve run notification debug logs * add workload debug log route * misc run controller fixes and refactor * prevent parallel execution of critical functions * update bun to 1.2.1 * WIP with dev dequeuing * Method to convert friendlyIds to non-friendly, do nothing with actual ids * Set the engine on BackgroundWorker, lazily upgrade projects to engine V2 * Runs with ttls were getting immediately expired… oops. * Pass the Waiting for deploy reason through, so we have it on the execution snapshots * Fixed the logic for getting the right background worker for a run * Use the correct ID when dequeuing… * determineEngineVersion is now fully functional * Rate limiter ignores the dev endpoints * Retrieving a batch gives you the runIds * Set a unique version for the RE2 BatchTaskRun * add provisional changeset * The start of dev run execution is working * First dev run working * Moved the dev run controller closer to what Nick did with the managed one * export exec output type * Heartbeat fix: don’t heartbeat if _isHeartbeating == false * Dev runs get notifications, some dev bug fixes * Improved logging or dequeuing * We need to dequeue runs from the latest version too, for triggerAndWait * Ported Eric’s validateWorkerManifest with nicer errors * When flattening an idempotency key if part is undefined, return undefined * Dev logging fixes * Remove sigterm listener * Deprecating workers. Don’t specify a BackgroundWorker when dequeuing an environment * Deleted some old files. Renamed “managed” to “deploy” * When a build finishes, always copy the build dir (otherwise the first one gets trampled on by the 2nd) * Dev master queues should work differently * Deleting old workers * Added debounce function to core * Improvement to canceling * WIP on debounce canceling on socket disconnection * Added environment data to execution snapshots * Dev runs that have stalled get “Canceled” with a reason explaining why * Show CLI messaged when a connection to the platform is lost/restored * Fix TriggerTask after merge * Add trigger task v2 max attempts, replace some findUniques * Port the new queue logic to the run engine * More fixes post-merge * We weren’t setting a `retryConfig` up for the tests… it’s now required * Start the Redis worker inside the Run Engine… 🤦♂️ * Trying to make the testcontainers more reliable * Added keyPrefix: "engine:” * Badly placed bracket in trigger task * Better Redis namespacing * Fix for expired run not getting removed from the queue * Don’t create a redis client in the testcontainers, return the redisOptions instead * Cleanup redis client in the run lock tests * Fix for the RunQueue not supporting keyPrefix * Updated more of the RunQueue scripts rebalancing * Trying to make Redis more robust in the tests… * Improved test resiliciency more * Fix for delays (checkpoint check) * Increase the timeout slightly to fix ttl test * Added priority support when triggering * More wip trying to make test containers more reliable * batchTriggerAndWait test is still failing… some wip to try fix it * Fixed redis tests now we’re not providing a client * Separate Redis clients for the run engine worker/queue/runlock * Made the wait for duration test more resilient * Added idempotencyKeyExpiresAt to Waitpoints * Waitpoint timeouts and idempotency expiry * Use finishWaitpoint, removed extra worker job * Added waitpoint idempotency tests * Creating resume tokens is working * Some improvements to the resume tokens * Moved resumeTokens to just be wait functions 🥳 * Delete old RuntimeManagers * Wait for token is working * Better test for the wait tokens * Improved the test task some more * Hide the accessories in the span inspector * WIP on waitpoint inspector * WIP on complete waitpoint form * Span overview panel can be changed based on the entity type * Improved the waitpoint display * WIP on completing waitpoint form * Use the existing CodeBlock for the tip * Style improvements * Complete waitpoint * All waitpoint sidebar variants * Waits now use a pause icon * Durations waits use the API to create/block with a waitpoint, not the runtime * Fix for engine.blockRunWithWaitpoint required org id * Removed old wait code from the run controllers/task run process * Form action for skipping a datetime waitpoint * Move testDockerCheckpoint to a separate core package export (it can’t be bundled on the client) * Fix for glitchy hourglass animation * Completed waitpoints display better * Increase Redis maxRetriesPerRequest to 20 (default) * Completing and skipping waitpoints is working * Remove the database prisma dev command, since we need to use create only now. Updated docs * Added skip timeout, reworked the UI * Tweaked spacing * Added payload limit to waitpoint token completion from dashboard * Test idempotency works on wait.for and wait.until * Moved the worker-actions to /engine/ from /api/ * Moved dev engine endpoints to /engine/ from /api/ * Separate /engine/ rate limiter * Added parallel wait prevention, it’s working for duration waits but not well for triggerAndWait yet * WIP post-merge conflicts * Set taskEventStore column in the new engine * Remove duplicate keys * Post-merge fixes * Fix for span merge layout * Use executedAt instead of firstAttemptStartedAt --------- Co-authored-by: Matt Aitken <matt@mattaitken.com>
334 lines
7.9 KiB
TypeScript
334 lines
7.9 KiB
TypeScript
// Bulk adds data to the database for testing
|
|
// Call it like this
|
|
// 1. pnpm run build:db:populate
|
|
// 2. pnpm run db:populate -- --projectRef=proj_liazlkfgmfcusswwgohl --taskIdentifier=child-task --runCount=100000
|
|
import { generateFriendlyId } from "~/v3/friendlyIdentifiers";
|
|
import { prisma } from "../app/db.server";
|
|
import { createHash } from "crypto";
|
|
import {
|
|
BackgroundWorker,
|
|
BackgroundWorkerTask,
|
|
RuntimeEnvironmentType,
|
|
WorkerInstanceGroupType,
|
|
} from "@trigger.dev/database";
|
|
import { nanoid } from "nanoid";
|
|
|
|
async function populate() {
|
|
if (process.env.NODE_ENV !== "development") {
|
|
return;
|
|
}
|
|
|
|
const project = await getProject();
|
|
|
|
await generateRuns(project);
|
|
await createWorkerGroup(project);
|
|
const { worker, tasks } = await createBackgroundWorker(project, getEnvTypeFromArg());
|
|
await createWorkerDeployment(project, worker, getEnvTypeFromArg());
|
|
}
|
|
|
|
function getEnvironment(
|
|
project: ProjectWithEnvironment,
|
|
envType: RuntimeEnvironmentType = "PRODUCTION"
|
|
) {
|
|
const env = project.environments.find((e) => e.type === envType);
|
|
|
|
if (!env) {
|
|
throw new Error(`No environment of type "${envType}" found for project ${project.id}`);
|
|
}
|
|
|
|
return env;
|
|
}
|
|
|
|
async function createWorkerDeployment(
|
|
project: ProjectWithEnvironment,
|
|
worker: BackgroundWorker,
|
|
envType: RuntimeEnvironmentType = "PRODUCTION"
|
|
) {
|
|
const env = getEnvironment(project, envType);
|
|
const deploymentId = `cm3c821sk00032v6is7ufqy3d-${env.slug}`;
|
|
|
|
if (env.type === "DEVELOPMENT") {
|
|
console.warn("Skipping deployment creation for development environment");
|
|
return;
|
|
}
|
|
|
|
let deployment = await prisma.workerDeployment.findUnique({
|
|
where: {
|
|
id: deploymentId,
|
|
},
|
|
});
|
|
|
|
if (deployment) {
|
|
console.log(`Deployment "${deploymentId}" already exists`);
|
|
return deployment;
|
|
}
|
|
|
|
const firstOrgMember = project.organization.members[0];
|
|
|
|
deployment = await prisma.workerDeployment.create({
|
|
data: {
|
|
id: deploymentId,
|
|
friendlyId: generateFriendlyId("deployment"),
|
|
contentHash: worker.contentHash,
|
|
version: worker.version,
|
|
shortCode: nanoid(8),
|
|
imageReference: `trigger/${project.externalRef}:${worker.version}.${env.slug}`,
|
|
status: "DEPLOYING",
|
|
projectId: project.id,
|
|
environmentId: env.id,
|
|
workerId: worker.id,
|
|
triggeredById: firstOrgMember.userId,
|
|
},
|
|
});
|
|
|
|
console.log(`Created deployment "${deploymentId}"`);
|
|
|
|
return deployment;
|
|
}
|
|
|
|
async function createBackgroundWorker(
|
|
project: ProjectWithEnvironment,
|
|
envType: RuntimeEnvironmentType = "PRODUCTION"
|
|
) {
|
|
const env = getEnvironment(project, envType);
|
|
const taskIdentifier = "seed-task";
|
|
const backgroundWorkerId = `cm3c8fmiv00042v6imoqwxst1-${env.slug}`;
|
|
|
|
let worker = await prisma.backgroundWorker.findUnique({
|
|
where: {
|
|
id: backgroundWorkerId,
|
|
},
|
|
include: {
|
|
tasks: true,
|
|
},
|
|
});
|
|
|
|
if (worker) {
|
|
console.log(`Worker "${backgroundWorkerId}" already exists`);
|
|
|
|
return {
|
|
worker,
|
|
tasks: worker.tasks,
|
|
};
|
|
}
|
|
|
|
worker = await prisma.backgroundWorker.create({
|
|
data: {
|
|
id: backgroundWorkerId,
|
|
friendlyId: generateFriendlyId("worker"),
|
|
contentHash: "hash",
|
|
projectId: project.id,
|
|
runtimeEnvironmentId: env.id,
|
|
version: "20241111.1",
|
|
metadata: {},
|
|
},
|
|
include: {
|
|
tasks: true,
|
|
},
|
|
});
|
|
|
|
console.log(`Created worker "${backgroundWorkerId}"`);
|
|
|
|
const taskIdentifiers = Array.isArray(taskIdentifier) ? taskIdentifier : [taskIdentifier];
|
|
|
|
const tasks: BackgroundWorkerTask[] = [];
|
|
|
|
for (const identifier of taskIdentifiers) {
|
|
const task = await prisma.backgroundWorkerTask.create({
|
|
data: {
|
|
friendlyId: generateFriendlyId("task"),
|
|
slug: identifier,
|
|
filePath: `/trigger/${identifier}.ts`,
|
|
exportName: identifier,
|
|
workerId: worker.id,
|
|
runtimeEnvironmentId: env.id,
|
|
projectId: project.id,
|
|
},
|
|
});
|
|
|
|
tasks.push(task);
|
|
}
|
|
|
|
return {
|
|
worker,
|
|
tasks,
|
|
};
|
|
}
|
|
|
|
async function createWorkerGroup(project: ProjectWithEnvironment) {
|
|
const workerGroupName = "seed-unmanaged";
|
|
const rawToken = "tr_wgt_15480aa1712cae4b8db8c7a49707d69d";
|
|
|
|
const existingWorkerGroup = await prisma.workerInstanceGroup.findFirst({
|
|
where: {
|
|
projectId: project.id,
|
|
name: workerGroupName,
|
|
},
|
|
});
|
|
|
|
if (existingWorkerGroup) {
|
|
console.log(`Worker group "${workerGroupName}" already exists`);
|
|
|
|
await setAsDefaultWorkerGroup(project, existingWorkerGroup.id);
|
|
|
|
return existingWorkerGroup;
|
|
}
|
|
|
|
const token = await prisma.workerGroupToken.create({
|
|
data: {
|
|
tokenHash: createHash("sha256").update(rawToken).digest("hex"),
|
|
},
|
|
});
|
|
|
|
const workerGroup = await prisma.workerInstanceGroup.create({
|
|
data: {
|
|
projectId: project.id,
|
|
organizationId: project.organizationId,
|
|
type: WorkerInstanceGroupType.UNMANAGED,
|
|
masterQueue: `${project.id}-${workerGroupName}`,
|
|
tokenId: token.id,
|
|
description: "Seeded worker group",
|
|
name: workerGroupName,
|
|
},
|
|
});
|
|
|
|
await setAsDefaultWorkerGroup(project, workerGroup.id);
|
|
|
|
return workerGroup;
|
|
}
|
|
|
|
async function setAsDefaultWorkerGroup(project: ProjectWithEnvironment, workerGroupId: string) {
|
|
// Set as default worker group
|
|
await prisma.project.update({
|
|
where: {
|
|
id: project.id,
|
|
},
|
|
data: {
|
|
defaultWorkerGroupId: workerGroupId,
|
|
},
|
|
});
|
|
}
|
|
|
|
async function getProject() {
|
|
const projectRef = getArg("projectRef");
|
|
if (!projectRef) {
|
|
throw new Error("projectRef is required");
|
|
}
|
|
|
|
const project = await prisma.project.findUnique({
|
|
include: {
|
|
environments: true,
|
|
organization: {
|
|
include: {
|
|
members: true,
|
|
},
|
|
},
|
|
},
|
|
where: {
|
|
externalRef: projectRef,
|
|
},
|
|
});
|
|
|
|
if (!project) {
|
|
throw new Error("Project not found");
|
|
}
|
|
|
|
return project;
|
|
}
|
|
|
|
type ProjectWithEnvironment = Awaited<ReturnType<typeof getProject>>;
|
|
|
|
async function generateRuns(project: ProjectWithEnvironment) {
|
|
const taskIdentifier = getArg("taskIdentifier");
|
|
if (!taskIdentifier) {
|
|
throw new Error("taskIdentifier is required");
|
|
}
|
|
|
|
const runCount = parseInt(getArg("runCount") || "100");
|
|
|
|
const task = await prisma.backgroundWorkerTask.findFirst({
|
|
where: {
|
|
projectId: project.id,
|
|
slug: taskIdentifier,
|
|
},
|
|
orderBy: {
|
|
createdAt: "desc",
|
|
},
|
|
});
|
|
|
|
if (!task) {
|
|
throw new Error("Task not found");
|
|
}
|
|
|
|
const runs = await prisma.taskRun.createMany({
|
|
data: Array(runCount)
|
|
.fill(0)
|
|
.map((_, index) => {
|
|
const friendlyId = generateFriendlyId("run");
|
|
|
|
return {
|
|
status: "CANCELED",
|
|
completedAt: new Date(),
|
|
number: index + 1,
|
|
friendlyId,
|
|
runtimeEnvironmentId: project.environments[randomIndex(project.environments)].id,
|
|
projectId: project.id,
|
|
taskIdentifier,
|
|
payload: JSON.stringify({ foo: "bar" }),
|
|
traceId: "traceId",
|
|
spanId: "spanId",
|
|
queue: "task/${taskIdentifier}",
|
|
};
|
|
}),
|
|
skipDuplicates: true,
|
|
});
|
|
|
|
console.log(`Added ${runs.count} runs`);
|
|
}
|
|
|
|
function getEnvTypeFromArg(): RuntimeEnvironmentType {
|
|
const env = getArg("env");
|
|
|
|
if (!env) {
|
|
return RuntimeEnvironmentType.PRODUCTION;
|
|
}
|
|
|
|
switch (env) {
|
|
case "dev":
|
|
return RuntimeEnvironmentType.DEVELOPMENT;
|
|
case "prod":
|
|
return RuntimeEnvironmentType.PRODUCTION;
|
|
case "stg":
|
|
return RuntimeEnvironmentType.STAGING;
|
|
default:
|
|
throw new Error(`Invalid environment: ${env}`);
|
|
}
|
|
}
|
|
|
|
function getArg(name: string) {
|
|
const args = process.argv.slice(2);
|
|
|
|
let value = "";
|
|
|
|
args.forEach((val) => {
|
|
if (val.startsWith(`--${name}=`)) {
|
|
value = val.split("=")[1];
|
|
}
|
|
});
|
|
|
|
return !value ? undefined : value;
|
|
}
|
|
|
|
function randomIndex<T>(array: T[]) {
|
|
return Math.floor(Math.random() * array.length);
|
|
}
|
|
|
|
populate()
|
|
.catch((e) => {
|
|
console.error(e);
|
|
process.exit(1);
|
|
})
|
|
.finally(async () => {
|
|
await prisma.$disconnect();
|
|
});
|