Fix for TASK_RUN_HEARTBEAT errors in deployed and dev works

This commit is contained in:
Eric Allam
2024-09-16 19:10:19 +01:00
parent 00668ff39d
commit de135e4885
7 changed files with 40 additions and 18 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"trigger.dev": patch
---
Configurable deployed heartbeat interval via HEARTBEAT_INTERVAL_MS env var
+5
View File
@@ -177,6 +177,11 @@ const EnvironmentSchema = z.object({
LOOPS_API_KEY: z.string().optional(), LOOPS_API_KEY: z.string().optional(),
MARQS_DISABLE_REBALANCING: z.coerce.boolean().default(false), 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"), VERBOSE_GRAPHILE_LOGGING: z.string().default("false"),
V2_MARQS_ENABLED: z.string().default("0"), V2_MARQS_ENABLED: z.string().default("0"),
@@ -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); const commonVariables = await resolveCommonBuiltInVariables(runtimeEnvironment);
return [...result, ...commonVariables]; return [...result, ...commonVariables];
@@ -162,8 +162,8 @@ export class DevQueueConsumer {
/** /**
* @deprecated Use `taskRunHeartbeat` instead * @deprecated Use `taskRunHeartbeat` instead
*/ */
public async taskHeartbeat(workerId: string, id: string, seconds: number = 60) { public async taskHeartbeat(workerId: string, id: string) {
logger.debug("[DevQueueConsumer] taskHeartbeat()", { id, seconds }); logger.debug("[DevQueueConsumer] taskHeartbeat()", { id });
const taskRunAttempt = await prisma.taskRunAttempt.findUnique({ const taskRunAttempt = await prisma.taskRunAttempt.findUnique({
where: { friendlyId: id }, where: { friendlyId: id },
@@ -173,13 +173,13 @@ export class DevQueueConsumer {
return; return;
} }
await marqs?.heartbeatMessage(taskRunAttempt.taskRunId, seconds); await marqs?.heartbeatMessage(taskRunAttempt.taskRunId);
} }
public async taskRunHeartbeat(workerId: string, id: string, seconds: number = 60) { public async taskRunHeartbeat(workerId: string, id: string) {
logger.debug("[DevQueueConsumer] taskRunHeartbeat()", { id, seconds }); logger.debug("[DevQueueConsumer] taskRunHeartbeat()", { id });
await marqs?.heartbeatMessage(id, seconds); await marqs?.heartbeatMessage(id);
} }
public async stop(reason: string = "CLI disconnected") { public async stop(reason: string = "CLI disconnected") {
+3 -3
View File
@@ -698,8 +698,8 @@ export class MarQS {
} }
// This should increment by the number of seconds, but with a max value of Date.now() + visibilityTimeoutInMs // 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) { public async heartbeatMessage(messageId: string) {
await this.options.visibilityTimeoutStrategy.heartbeat(messageId, seconds * 1000); await this.options.visibilityTimeoutStrategy.heartbeat(messageId, this.visibilityTimeoutInMs);
} }
get visibilityTimeoutInMs() { get visibilityTimeoutInMs() {
@@ -1871,7 +1871,7 @@ function getMarQSClient() {
redis: redisOptions, redis: redisOptions,
defaultEnvConcurrency: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_LIMIT, defaultEnvConcurrency: env.DEFAULT_ENV_EXECUTION_CONCURRENCY_LIMIT,
defaultOrgConcurrency: env.DEFAULT_ORG_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, enableRebalancing: !env.MARQS_DISABLE_REBALANCING,
subscriber: concurrencyTracker, subscriber: concurrencyTracker,
}); });
@@ -1169,8 +1169,8 @@ class SharedQueueTasks {
} satisfies TaskRunExecutionLazyAttemptPayload; } satisfies TaskRunExecutionLazyAttemptPayload;
} }
async taskHeartbeat(attemptFriendlyId: string, seconds: number = 60) { async taskHeartbeat(attemptFriendlyId: string) {
logger.debug("[SharedQueueConsumer] taskHeartbeat()", { id: attemptFriendlyId, seconds }); logger.debug("[SharedQueueConsumer] taskHeartbeat()", { id: attemptFriendlyId });
const taskRunAttempt = await prisma.taskRunAttempt.findUnique({ const taskRunAttempt = await prisma.taskRunAttempt.findUnique({
where: { friendlyId: attemptFriendlyId }, where: { friendlyId: attemptFriendlyId },
@@ -1180,13 +1180,13 @@ class SharedQueueTasks {
return; return;
} }
await marqs?.heartbeatMessage(taskRunAttempt.taskRunId, seconds); await marqs?.heartbeatMessage(taskRunAttempt.taskRunId);
} }
async taskRunHeartbeat(runId: string, seconds: number = 60) { async taskRunHeartbeat(runId: string) {
logger.debug("[SharedQueueConsumer] taskRunHeartbeat()", { runId, seconds }); logger.debug("[SharedQueueConsumer] taskRunHeartbeat()", { runId });
await marqs?.heartbeatMessage(runId, seconds); await marqs?.heartbeatMessage(runId);
} }
public async taskRunFailed(completion: TaskRunFailedExecutionResult) { public async taskRunFailed(completion: TaskRunFailedExecutionResult) {
@@ -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 usageEventUrl = getEnvVar("USAGE_EVENT_URL");
const triggerJWT = getEnvVar("TRIGGER_JWT"); const triggerJWT = getEnvVar("TRIGGER_JWT");
const heartbeatIntervalMs = getEnvVar("HEARTBEAT_INTERVAL_MS");
const prodUsageManager = new ProdUsageManager(new DevUsageManager(), { const prodUsageManager = new ProdUsageManager(new DevUsageManager(), {
heartbeatIntervalMs: heartbeatIntervalMs ? parseInt(heartbeatIntervalMs, 10) : undefined, heartbeatIntervalMs: usageIntervalMs ? parseInt(usageIntervalMs, 10) : undefined,
url: usageEventUrl, url: usageEventUrl,
jwt: triggerJWT, jwt: triggerJWT,
}); });
@@ -383,7 +384,9 @@ runtime.setGlobalRuntimeManager(prodRuntimeManager);
process.title = "trigger-dev-worker"; 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) { if (_isRunning && _execution) {
try { try {
await zodIpc.send("TASK_HEARTBEAT", { id: _execution.attempt.id }); await zodIpc.send("TASK_HEARTBEAT", { id: _execution.attempt.id });