diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index c6d246518..9a108a940 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -571,6 +571,9 @@ const EnvironmentSchema = z.object({ /** The max number of runs per API call that we'll dequeue in DEV */ DEV_DEQUEUE_MAX_RUNS_PER_PULL: z.coerce.number().int().default(10), + /** The maximum concurrent local run processes executing at once in dev */ + DEV_MAX_CONCURRENT_RUNS: z.coerce.number().int().default(25), + LEGACY_RUN_ENGINE_WORKER_ENABLED: z.string().default(process.env.WORKER_ENABLED ?? "true"), LEGACY_RUN_ENGINE_WORKER_CONCURRENCY_WORKERS: z.coerce.number().int().default(2), LEGACY_RUN_ENGINE_WORKER_CONCURRENCY_TASKS_PER_WORKER: z.coerce.number().int().default(1), diff --git a/apps/webapp/app/routes/engine.v1.dev.config.ts b/apps/webapp/app/routes/engine.v1.dev.config.ts index 501b6e80d..0a4c8e4ba 100644 --- a/apps/webapp/app/routes/engine.v1.dev.config.ts +++ b/apps/webapp/app/routes/engine.v1.dev.config.ts @@ -20,6 +20,7 @@ export const loader = createLoaderApiRoute( environmentId: authentication.environment.id, dequeueIntervalWithRun: env.DEV_DEQUEUE_INTERVAL_WITH_RUN, dequeueIntervalWithoutRun: env.DEV_DEQUEUE_INTERVAL_WITHOUT_RUN, + maxConcurrentRuns: env.DEV_MAX_CONCURRENT_RUNS, }); } catch (error) { logger.error("Failed to get dev settings", { diff --git a/internal-packages/run-engine/src/engine/index.ts b/internal-packages/run-engine/src/engine/index.ts index 0433c6d18..cddf3994d 100644 --- a/internal-packages/run-engine/src/engine/index.ts +++ b/internal-packages/run-engine/src/engine/index.ts @@ -925,6 +925,7 @@ export class RunEngine { return { version: "1" as const, + dequeuedAt: new Date(), snapshot: { id: newSnapshot.id, friendlyId: newSnapshot.friendlyId, diff --git a/packages/cli-v3/package.json b/packages/cli-v3/package.json index 7e7c42231..1992498e5 100644 --- a/packages/cli-v3/package.json +++ b/packages/cli-v3/package.json @@ -110,6 +110,7 @@ "nypm": "^0.3.9", "object-hash": "^3.0.0", "open": "^10.0.3", + "p-limit": "^6.2.0", "p-retry": "^6.1.0", "partysocket": "^1.0.2", "pkg-types": "^1.1.3", diff --git a/packages/cli-v3/src/commands/dev.ts b/packages/cli-v3/src/commands/dev.ts index 7a9726dc5..479a47217 100644 --- a/packages/cli-v3/src/commands/dev.ts +++ b/packages/cli-v3/src/commands/dev.ts @@ -19,6 +19,7 @@ const DevCommandOptions = CommonCommandOptions.extend({ skipUpdateCheck: z.boolean().default(false), envFile: z.string().optional(), keepTmpFiles: z.boolean().default(false), + maxConcurrentRuns: z.coerce.number().optional(), }); export type DevCommandOptions = z.infer; @@ -37,6 +38,10 @@ export function configureDevCommand(program: Command) { "--env-file ", "Path to the .env file to use for the dev session. Defaults to .env in the project directory." ) + .option( + "--max-concurrent-runs ", + "The maximum number of concurrent runs to allow in the dev session" + ) .option("--debug-otel", "Enable OpenTelemetry debugging") .option("--skip-update-check", "Skip checking for @trigger.dev package updates") .option( diff --git a/packages/cli-v3/src/dev/devSupervisor.ts b/packages/cli-v3/src/dev/devSupervisor.ts index 0c74d2ab5..75865cc1f 100644 --- a/packages/cli-v3/src/dev/devSupervisor.ts +++ b/packages/cli-v3/src/dev/devSupervisor.ts @@ -24,6 +24,7 @@ import { WorkerClientToServerEvents, WorkerServerToClientEvents, } from "@trigger.dev/core/v3/workers"; +import pLimit from "p-limit"; export type WorkerRuntimeOptions = { name: string | undefined; @@ -65,6 +66,8 @@ class DevSupervisor implements WorkerRuntime { private socketConnections = new Set(); + private runLimiter?: ReturnType; + constructor(public readonly options: WorkerRuntimeOptions) {} async init(): Promise { @@ -81,6 +84,15 @@ class DevSupervisor implements WorkerRuntime { logger.debug("[DevSupervisor] Got dev settings", { settings: settings.data }); this.config = settings.data; + const maxConcurrentRuns = Math.min( + this.config.maxConcurrentRuns, + this.options.args.maxConcurrentRuns ?? this.config.maxConcurrentRuns + ); + + logger.debug("[DevSupervisor] Using maxConcurrentRuns", { maxConcurrentRuns }); + + this.runLimiter = pLimit(maxConcurrentRuns); + this.#createSocket(); //start an SSE connection for presence @@ -178,6 +190,14 @@ class DevSupervisor implements WorkerRuntime { return; } + if ( + this.runLimiter && + this.runLimiter.activeCount + this.runLimiter.pendingCount > this.runLimiter.concurrency + ) { + logger.debug(`[DevSupervisor] dequeueRuns. Run limit reached, trying again later`); + setTimeout(() => this.#dequeueRuns(), this.config.dequeueIntervalWithoutRun); + } + //get relevant versions //ignore deprecated and the latest worker const oldWorkerIds = this.#getActiveOldWorkers(); @@ -287,10 +307,16 @@ class DevSupervisor implements WorkerRuntime { this.runControllers.set(message.run.friendlyId, runController); - //don't await for run completion, we want to dequeue more runs - runController.start(message).then(() => { - logger.debug("[DevSupervisor] Run started", { runId: message.run.friendlyId }); - }); + if (this.runLimiter) { + this.runLimiter(() => runController.start(message)).then(() => { + logger.debug("[DevSupervisor] Run started", { runId: message.run.friendlyId }); + }); + } else { + //don't await for run completion, we want to dequeue more runs + runController.start(message).then(() => { + logger.debug("[DevSupervisor] Run started", { runId: message.run.friendlyId }); + }); + } } setTimeout(() => this.#dequeueRuns(), this.config.dequeueIntervalWithRun); diff --git a/packages/cli-v3/src/entryPoints/dev-run-controller.ts b/packages/cli-v3/src/entryPoints/dev-run-controller.ts index 46cfdfc0a..68c051d8f 100644 --- a/packages/cli-v3/src/entryPoints/dev-run-controller.ts +++ b/packages/cli-v3/src/entryPoints/dev-run-controller.ts @@ -5,6 +5,7 @@ import { LogLevel, RunExecutionData, TaskRunExecution, + TaskRunExecutionMetrics, TaskRunExecutionResult, TaskRunFailedExecutionResult, } from "@trigger.dev/core/v3"; @@ -475,10 +476,12 @@ export class DevRunController { private async startAndExecuteRunAttempt({ runFriendlyId, snapshotFriendlyId, + dequeuedAt, isWarmStart = false, }: { runFriendlyId: string; snapshotFriendlyId: string; + dequeuedAt?: Date; isWarmStart?: boolean; }) { this.subscribeToRunNotifications({ @@ -486,6 +489,8 @@ export class DevRunController { snapshot: { friendlyId: snapshotFriendlyId }, }); + const attemptStartedAt = Date.now(); + const start = await this.httpClient.dev.startRunAttempt(runFriendlyId, snapshotFriendlyId); if (!start.success) { @@ -495,6 +500,8 @@ export class DevRunController { return; } + const attemptDuration = Date.now() - attemptStartedAt; + const { run, snapshot, execution, envVars } = start.data; eventBus.emit("runStarted", this.opts.worker, execution); @@ -508,8 +515,28 @@ export class DevRunController { // This is the only case where incrementing the attempt number is allowed this.enterRunPhase(run, snapshot); + const metrics = [ + { + name: "start", + event: "create_attempt", + timestamp: attemptStartedAt, + duration: attemptDuration, + }, + ].concat( + dequeuedAt + ? [ + { + name: "start", + event: "dequeue", + timestamp: dequeuedAt.getTime(), + duration: 0, + }, + ] + : [] + ); + try { - return await this.executeRun({ run, snapshot, execution, envVars }); + return await this.executeRun({ run, snapshot, execution, envVars, metrics }); } catch (error) { // TODO: Handle the case where we're in the warm start phase or executing a new run // This can happen if we kill the run while it's still executing, e.g. after receiving an attempt number mismatch @@ -566,7 +593,10 @@ export class DevRunController { snapshot, execution, envVars, - }: WorkloadRunAttemptStartResponseBody) { + metrics, + }: WorkloadRunAttemptStartResponseBody & { + metrics?: TaskRunExecutionMetrics; + }) { if (!this.opts.worker.serverWorker) { throw new Error(`No server worker for Dev ${run.friendlyId}`); } @@ -594,6 +624,7 @@ export class DevRunController { payload: { execution, traceContext: execution.run.traceContext ?? {}, + metrics, }, messageId: run.friendlyId, }); @@ -753,6 +784,7 @@ export class DevRunController { await this.startAndExecuteRunAttempt({ runFriendlyId: dequeueMessage.run.friendlyId, snapshotFriendlyId: dequeueMessage.snapshot.friendlyId, + dequeuedAt: dequeueMessage.dequeuedAt, }).finally(async () => {}); } diff --git a/packages/cli-v3/src/entryPoints/dev-run-worker.ts b/packages/cli-v3/src/entryPoints/dev-run-worker.ts index d196707ec..30ff92ddc 100644 --- a/packages/cli-v3/src/entryPoints/dev-run-worker.ts +++ b/packages/cli-v3/src/entryPoints/dev-run-worker.ts @@ -17,6 +17,7 @@ import { waitUntil, WorkerManifest, WorkerToExecutorMessageCatalog, + runTimelineMetrics, } from "@trigger.dev/core/v3"; import { TriggerTracer } from "@trigger.dev/core/v3/tracer"; import { @@ -36,6 +37,7 @@ import { TracingSDK, usage, UsageTimeoutManager, + StandardRunTimelineMetricsManager, } from "@trigger.dev/core/v3/workers"; import { ZodIpcConnection } from "@trigger.dev/core/v3/zodIpc"; import { readFile } from "node:fs/promises"; @@ -87,6 +89,10 @@ process.on("uncaughtException", function (error, origin) { const heartbeatIntervalMs = getEnvVar("HEARTBEAT_INTERVAL_MS"); +const standardRunTimelineMetricsManager = new StandardRunTimelineMetricsManager(); +runTimelineMetrics.setGlobalManager(standardRunTimelineMetricsManager); +standardRunTimelineMetricsManager.seedMetricsFromEnvironment(); + const devUsageManager = new DevUsageManager(); usage.setGlobalUsageManager(devUsageManager); timeout.setGlobalManager(new UsageTimeoutManager(devUsageManager)); @@ -189,9 +195,11 @@ const zodIpc = new ZodIpcConnection({ emitSchema: ExecutorToWorkerMessageCatalog, process, handlers: { - EXECUTE_TASK_RUN: async ({ execution, traceContext, metadata }, sender) => { + EXECUTE_TASK_RUN: async ({ execution, traceContext, metadata, metrics }, sender) => { log(`[${new Date().toISOString()}] Received EXECUTE_TASK_RUN`, execution); + standardRunTimelineMetricsManager.registerMetricsFromExecution(metrics); + if (_isRunning) { logError("Worker is already running a task"); @@ -246,11 +254,22 @@ const zodIpc = new ZodIpcConnection({ } try { - const beforeImport = performance.now(); - await import(normalizeImportPath(taskManifest.entryPoint)); - const durationMs = performance.now() - beforeImport; + await runTimelineMetrics.measureMetric( + "trigger.dev/start", + "import", + { + entryPoint: taskManifest.entryPoint, + }, + async () => { + const beforeImport = performance.now(); + await import(normalizeImportPath(taskManifest.entryPoint)); + const durationMs = performance.now() - beforeImport; - log(`Imported task ${execution.task.id} [${taskManifest.entryPoint}] in ${durationMs}ms`); + log( + `Imported task ${execution.task.id} [${taskManifest.entryPoint}] in ${durationMs}ms` + ); + } + ); } catch (err) { logError(`Failed to import task ${execution.task.id}`, err); diff --git a/packages/cli-v3/src/entryPoints/managed-run-worker.ts b/packages/cli-v3/src/entryPoints/managed-run-worker.ts index c196f3ca0..8749801b4 100644 --- a/packages/cli-v3/src/entryPoints/managed-run-worker.ts +++ b/packages/cli-v3/src/entryPoints/managed-run-worker.ts @@ -17,6 +17,7 @@ import { runMetadata, waitUntil, apiClientManager, + runTimelineMetrics, } from "@trigger.dev/core/v3"; import { TriggerTracer } from "@trigger.dev/core/v3/tracer"; import { @@ -37,6 +38,7 @@ import { StandardMetadataManager, StandardWaitUntilManager, ManagedRuntimeManager, + StandardRunTimelineMetricsManager, } from "@trigger.dev/core/v3/workers"; import { ZodIpcConnection } from "@trigger.dev/core/v3/zodIpc"; import { readFile } from "node:fs/promises"; @@ -91,6 +93,10 @@ const usageEventUrl = getEnvVar("USAGE_EVENT_URL"); const triggerJWT = getEnvVar("TRIGGER_JWT"); const heartbeatIntervalMs = getEnvVar("HEARTBEAT_INTERVAL_MS"); +const standardRunTimelineMetricsManager = new StandardRunTimelineMetricsManager(); +runTimelineMetrics.setGlobalManager(standardRunTimelineMetricsManager); +standardRunTimelineMetricsManager.seedMetricsFromEnvironment(); + const devUsageManager = new DevUsageManager(); const prodUsageManager = new ProdUsageManager(devUsageManager, { heartbeatIntervalMs: usageIntervalMs ? parseInt(usageIntervalMs, 10) : undefined, @@ -199,7 +205,9 @@ const zodIpc = new ZodIpcConnection({ emitSchema: ExecutorToWorkerMessageCatalog, process, handlers: { - EXECUTE_TASK_RUN: async ({ execution, traceContext, metadata }, sender) => { + EXECUTE_TASK_RUN: async ({ execution, traceContext, metadata, metrics }, sender) => { + standardRunTimelineMetricsManager.registerMetricsFromExecution(metrics); + console.log(`[${new Date().toISOString()}] Received EXECUTE_TASK_RUN`, execution); if (_isRunning) { @@ -256,12 +264,21 @@ const zodIpc = new ZodIpcConnection({ } try { - const beforeImport = performance.now(); - await import(normalizeImportPath(taskManifest.entryPoint)); - const durationMs = performance.now() - beforeImport; + await runTimelineMetrics.measureMetric( + "trigger.dev/start", + "import", + { + entryPoint: taskManifest.entryPoint, + }, + async () => { + const beforeImport = performance.now(); + await import(normalizeImportPath(taskManifest.entryPoint)); + const durationMs = performance.now() - beforeImport; - console.log( - `Imported task ${execution.task.id} [${taskManifest.entryPoint}] in ${durationMs}ms` + console.log( + `Imported task ${execution.task.id} [${taskManifest.entryPoint}] in ${durationMs}ms` + ); + } ); } catch (err) { console.error(`Failed to import task ${execution.task.id}`, err); diff --git a/packages/core/src/v3/schemas/api.ts b/packages/core/src/v3/schemas/api.ts index 205d9b6a7..ff2028f7c 100644 --- a/packages/core/src/v3/schemas/api.ts +++ b/packages/core/src/v3/schemas/api.ts @@ -427,6 +427,7 @@ export const DevConfigResponseBody = z.object({ environmentId: z.string(), dequeueIntervalWithRun: z.number(), dequeueIntervalWithoutRun: z.number(), + maxConcurrentRuns: z.number(), }); export type DevConfigResponseBody = z.infer; diff --git a/packages/core/src/v3/schemas/runEngine.ts b/packages/core/src/v3/schemas/runEngine.ts index a93ca96a1..397630cb6 100644 --- a/packages/core/src/v3/schemas/runEngine.ts +++ b/packages/core/src/v3/schemas/runEngine.ts @@ -139,6 +139,7 @@ export type ExecutionResult = z.infer; export const DequeuedMessage = z.object({ version: z.literal("1"), snapshot: ExecutionSnapshot, + dequeuedAt: z.coerce.date(), image: z.string().optional(), checkpoint: z .object({ diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index d3e027be5..dc1941178 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -1206,6 +1206,9 @@ importers: open: specifier: ^10.0.3 version: 10.0.3 + p-limit: + specifier: ^6.2.0 + version: 6.2.0 p-retry: specifier: ^6.1.0 version: 6.1.0 diff --git a/references/v3-catalog/package.json b/references/v3-catalog/package.json index 52732b334..ee15d2e5e 100644 --- a/references/v3-catalog/package.json +++ b/references/v3-catalog/package.json @@ -6,7 +6,7 @@ "schema": "./prisma/schema.zmodel" }, "scripts": { - "dev:trigger": "trigger dev", + "dev": "trigger dev", "deploy": "trigger deploy --self-hosted --load-image", "management": "tsx -r dotenv/config ./src/management.ts", "queues": "ts-node -r dotenv/config -r tsconfig-paths/register ./src/queues.ts", @@ -84,4 +84,4 @@ "ts-node": "^10.9.2", "tsconfig-paths": "^4.2.0" } -} +} \ No newline at end of file