diff --git a/.changeset/poor-starfishes-act.md b/.changeset/poor-starfishes-act.md new file mode 100644 index 000000000..5b6fbc735 --- /dev/null +++ b/.changeset/poor-starfishes-act.md @@ -0,0 +1,5 @@ +--- +"trigger.dev": patch +--- + +Configurable deployed heartbeat interval via HEARTBEAT_INTERVAL_MS env var diff --git a/apps/webapp/app/env.server.ts b/apps/webapp/app/env.server.ts index 0d81c61d8..e4571f736 100644 --- a/apps/webapp/app/env.server.ts +++ b/apps/webapp/app/env.server.ts @@ -177,6 +177,11 @@ const EnvironmentSchema = z.object({ LOOPS_API_KEY: z.string().optional(), MARQS_DISABLE_REBALANCING: z.coerce.boolean().default(false), + MARQS_VISIBILITY_TIMEOUT_MS: z.coerce + .number() + .int() + .default(60 * 1000 * 15), + PROD_TASK_HEARTBEAT_INTERVAL_MS: z.coerce.number().int().optional(), VERBOSE_GRAPHILE_LOGGING: z.string().default("false"), V2_MARQS_ENABLED: z.string().default("0"), diff --git a/apps/webapp/app/v3/environmentVariables/environmentVariablesRepository.server.ts b/apps/webapp/app/v3/environmentVariables/environmentVariablesRepository.server.ts index c0ed29483..cdd845061 100644 --- a/apps/webapp/app/v3/environmentVariables/environmentVariablesRepository.server.ts +++ b/apps/webapp/app/v3/environmentVariables/environmentVariablesRepository.server.ts @@ -798,6 +798,15 @@ async function resolveBuiltInProdVariables(runtimeEnvironment: RuntimeEnvironmen ]); } + if (env.PROD_TASK_HEARTBEAT_INTERVAL_MS) { + result = result.concat([ + { + key: "HEARTBEAT_INTERVAL_MS", + value: String(env.PROD_TASK_HEARTBEAT_INTERVAL_MS), + }, + ]); + } + const commonVariables = await resolveCommonBuiltInVariables(runtimeEnvironment); return [...result, ...commonVariables]; diff --git a/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts index 7ee0953d9..4e80c9804 100644 --- a/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/devQueueConsumer.server.ts @@ -162,8 +162,8 @@ export class DevQueueConsumer { /** * @deprecated Use `taskRunHeartbeat` instead */ - public async taskHeartbeat(workerId: string, id: string, seconds: number = 60) { - logger.debug("[DevQueueConsumer] taskHeartbeat()", { id, seconds }); + public async taskHeartbeat(workerId: string, id: string) { + logger.debug("[DevQueueConsumer] taskHeartbeat()", { id }); const taskRunAttempt = await prisma.taskRunAttempt.findUnique({ where: { friendlyId: id }, @@ -173,13 +173,13 @@ export class DevQueueConsumer { return; } - await marqs?.heartbeatMessage(taskRunAttempt.taskRunId, seconds); + await marqs?.heartbeatMessage(taskRunAttempt.taskRunId); } - public async taskRunHeartbeat(workerId: string, id: string, seconds: number = 60) { - logger.debug("[DevQueueConsumer] taskRunHeartbeat()", { id, seconds }); + public async taskRunHeartbeat(workerId: string, id: string) { + logger.debug("[DevQueueConsumer] taskRunHeartbeat()", { id }); - await marqs?.heartbeatMessage(id, seconds); + await marqs?.heartbeatMessage(id); } public async stop(reason: string = "CLI disconnected") { diff --git a/apps/webapp/app/v3/marqs/index.server.ts b/apps/webapp/app/v3/marqs/index.server.ts index 5ebe087e8..81fec8e19 100644 --- a/apps/webapp/app/v3/marqs/index.server.ts +++ b/apps/webapp/app/v3/marqs/index.server.ts @@ -698,8 +698,8 @@ export class MarQS { } // This should increment by the number of seconds, but with a max value of Date.now() + visibilityTimeoutInMs - public async heartbeatMessage(messageId: string, seconds: number = 30) { - await this.options.visibilityTimeoutStrategy.heartbeat(messageId, seconds * 1000); + public async heartbeatMessage(messageId: string) { + await this.options.visibilityTimeoutStrategy.heartbeat(messageId, this.visibilityTimeoutInMs); } get visibilityTimeoutInMs() { @@ -1871,7 +1871,7 @@ function getMarQSClient() { redis: redisOptions, defaultEnvConcurrency: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_LIMIT, defaultOrgConcurrency: env.DEFAULT_ORG_EXECUTION_CONCURRENCY_LIMIT, - visibilityTimeoutInMs: 120 * 1000, // 2 minutes, + visibilityTimeoutInMs: env.MARQS_VISIBILITY_TIMEOUT_MS, enableRebalancing: !env.MARQS_DISABLE_REBALANCING, subscriber: concurrencyTracker, }); diff --git a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts index 28ce217b7..606d9ca92 100644 --- a/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts +++ b/apps/webapp/app/v3/marqs/sharedQueueConsumer.server.ts @@ -1169,8 +1169,8 @@ class SharedQueueTasks { } satisfies TaskRunExecutionLazyAttemptPayload; } - async taskHeartbeat(attemptFriendlyId: string, seconds: number = 60) { - logger.debug("[SharedQueueConsumer] taskHeartbeat()", { id: attemptFriendlyId, seconds }); + async taskHeartbeat(attemptFriendlyId: string) { + logger.debug("[SharedQueueConsumer] taskHeartbeat()", { id: attemptFriendlyId }); const taskRunAttempt = await prisma.taskRunAttempt.findUnique({ where: { friendlyId: attemptFriendlyId }, @@ -1180,13 +1180,13 @@ class SharedQueueTasks { return; } - await marqs?.heartbeatMessage(taskRunAttempt.taskRunId, seconds); + await marqs?.heartbeatMessage(taskRunAttempt.taskRunId); } - async taskRunHeartbeat(runId: string, seconds: number = 60) { - logger.debug("[SharedQueueConsumer] taskRunHeartbeat()", { runId, seconds }); + async taskRunHeartbeat(runId: string) { + logger.debug("[SharedQueueConsumer] taskRunHeartbeat()", { runId }); - await marqs?.heartbeatMessage(runId, seconds); + await marqs?.heartbeatMessage(runId); } public async taskRunFailed(completion: TaskRunFailedExecutionResult) { diff --git a/packages/cli-v3/src/entryPoints/deploy-run-worker.ts b/packages/cli-v3/src/entryPoints/deploy-run-worker.ts index 177694c73..e29a3d321 100644 --- a/packages/cli-v3/src/entryPoints/deploy-run-worker.ts +++ b/packages/cli-v3/src/entryPoints/deploy-run-worker.ts @@ -77,12 +77,13 @@ process.on("uncaughtException", function (error, origin) { } }); -const heartbeatIntervalMs = getEnvVar("USAGE_HEARTBEAT_INTERVAL_MS"); +const usageIntervalMs = getEnvVar("USAGE_HEARTBEAT_INTERVAL_MS"); const usageEventUrl = getEnvVar("USAGE_EVENT_URL"); const triggerJWT = getEnvVar("TRIGGER_JWT"); +const heartbeatIntervalMs = getEnvVar("HEARTBEAT_INTERVAL_MS"); const prodUsageManager = new ProdUsageManager(new DevUsageManager(), { - heartbeatIntervalMs: heartbeatIntervalMs ? parseInt(heartbeatIntervalMs, 10) : undefined, + heartbeatIntervalMs: usageIntervalMs ? parseInt(usageIntervalMs, 10) : undefined, url: usageEventUrl, jwt: triggerJWT, }); @@ -383,7 +384,9 @@ runtime.setGlobalRuntimeManager(prodRuntimeManager); process.title = "trigger-dev-worker"; -for await (const _ of setInterval(15_000)) { +const heartbeatInterval = parseInt(heartbeatIntervalMs ?? "30000", 10); + +for await (const _ of setInterval(heartbeatInterval)) { if (_isRunning && _execution) { try { await zodIpc.send("TASK_HEARTBEAT", { id: _execution.attempt.id });