Add attempt metrics in dev (prod WIP). Added max concurrent runs setting to dev using p-limit (#1766)

This commit is contained in:
Eric Allam
2025-03-06 11:57:01 +00:00
committed by GitHub
parent 722fae566c
commit cf4c4d7b65
13 changed files with 129 additions and 19 deletions
+3
View File
@@ -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),
@@ -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", {
@@ -925,6 +925,7 @@ export class RunEngine {
return {
version: "1" as const,
dequeuedAt: new Date(),
snapshot: {
id: newSnapshot.id,
friendlyId: newSnapshot.friendlyId,
+1
View File
@@ -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",
+5
View File
@@ -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<typeof DevCommandOptions>;
@@ -37,6 +38,10 @@ export function configureDevCommand(program: Command) {
"--env-file <env file>",
"Path to the .env file to use for the dev session. Defaults to .env in the project directory."
)
.option(
"--max-concurrent-runs <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(
+30 -4
View File
@@ -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<string>();
private runLimiter?: ReturnType<typeof pLimit>;
constructor(public readonly options: WorkerRuntimeOptions) {}
async init(): Promise<void> {
@@ -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);
@@ -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 () => {});
}
@@ -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);
@@ -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);
+1
View File
@@ -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<typeof DevConfigResponseBody>;
@@ -139,6 +139,7 @@ export type ExecutionResult = z.infer<typeof ExecutionResult>;
export const DequeuedMessage = z.object({
version: z.literal("1"),
snapshot: ExecutionSnapshot,
dequeuedAt: z.coerce.date(),
image: z.string().optional(),
checkpoint: z
.object({
+3
View File
@@ -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
+2 -2
View File
@@ -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"
}
}
}