From 584c7da5dfb1b54a2df93f5ada74b525dcf6b23b Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 18 Apr 2024 15:46:14 +0100 Subject: [PATCH] v3: prod worker graceful shutdown (#1034) * graceful exit with timeout * handle and display graceful timeout errors * fix for very long waits * changeset * increase termination grace period to an hour --- .changeset/tiny-elephants-scream.md | 7 +++ apps/kubernetes-provider/src/index.ts | 1 + .../app/v3/services/completeAttempt.server.ts | 54 ++++++++++++++++--- .../cli-v3/src/workers/prod/entry-point.ts | 37 +++++++++++-- .../cli-v3/src/workers/prod/worker-facade.ts | 20 +++++++ .../core/src/v3/runtime/devRuntimeManager.ts | 9 ++-- .../core/src/v3/runtime/prodRuntimeManager.ts | 3 +- packages/core/src/v3/schemas/common.ts | 2 + packages/core/src/v3/utils/timers.ts | 21 ++++++++ references/v3-catalog/src/trigger/other.ts | 16 +++++- 10 files changed, 150 insertions(+), 20 deletions(-) create mode 100644 .changeset/tiny-elephants-scream.md create mode 100644 packages/core/src/v3/utils/timers.ts diff --git a/.changeset/tiny-elephants-scream.md b/.changeset/tiny-elephants-scream.md new file mode 100644 index 000000000..4656e154d --- /dev/null +++ b/.changeset/tiny-elephants-scream.md @@ -0,0 +1,7 @@ +--- +"trigger.dev": patch +"@trigger.dev/core": patch +--- + +- Add graceful exit for prod workers +- Prevent overflow in long waits diff --git a/apps/kubernetes-provider/src/index.ts b/apps/kubernetes-provider/src/index.ts index 569d5035f..7326a0e11 100644 --- a/apps/kubernetes-provider/src/index.ts +++ b/apps/kubernetes-provider/src/index.ts @@ -133,6 +133,7 @@ class KubernetesTaskOperations implements TaskOperations { }, spec: { ...this.#defaultPodSpec, + terminationGracePeriodSeconds: 60 * 60, containers: [ { name: this.#getRunContainerName(opts.runId), diff --git a/apps/webapp/app/v3/services/completeAttempt.server.ts b/apps/webapp/app/v3/services/completeAttempt.server.ts index 56d283c84..1ab4285a2 100644 --- a/apps/webapp/app/v3/services/completeAttempt.server.ts +++ b/apps/webapp/app/v3/services/completeAttempt.server.ts @@ -248,14 +248,52 @@ export class CompleteAttemptService extends BaseService { }, }); - await this._prisma.taskRun.update({ - where: { - id: taskRunAttempt.taskRunId, - }, - data: { - status: "COMPLETED_WITH_ERRORS", - }, - }); + if ( + completion.error.type === "INTERNAL_ERROR" && + completion.error.code === "GRACEFUL_EXIT_TIMEOUT" + ) { + // We need to fail all incomplete spans + const inProgressEvents = await eventRepository.queryIncompleteEvents({ + attemptId: execution.attempt.id, + }); + + logger.debug("Failing in-progress events", { + inProgressEvents: inProgressEvents.map((event) => event.id), + }); + + const exception = { + type: "Graceful exit timeout", + message: completion.error.message, + }; + + await Promise.all( + inProgressEvents.map((event) => { + return eventRepository.crashEvent({ + event: event, + crashedAt: new Date(), + exception, + }); + }) + ); + + await this._prisma.taskRun.update({ + where: { + id: taskRunAttempt.taskRunId, + }, + data: { + status: "SYSTEM_FAILURE", + }, + }); + } else { + await this._prisma.taskRun.update({ + where: { + id: taskRunAttempt.taskRunId, + }, + data: { + status: "COMPLETED_WITH_ERRORS", + }, + }); + } if (!env || env.type !== "DEVELOPMENT") { await ResumeTaskRunDependenciesService.enqueue(taskRunAttempt.id, this._prisma); diff --git a/packages/cli-v3/src/workers/prod/entry-point.ts b/packages/cli-v3/src/workers/prod/entry-point.ts index 745129ee1..1cb2b9e60 100644 --- a/packages/cli-v3/src/workers/prod/entry-point.ts +++ b/packages/cli-v3/src/workers/prod/entry-point.ts @@ -11,7 +11,6 @@ import { import { HttpReply, SimpleLogger, getRandomPortNumber } from "@trigger.dev/core-apps"; import { readFile } from "node:fs/promises"; import { createServer } from "node:http"; -import { z } from "zod"; import { ProdBackgroundWorker } from "./backgroundWorker"; import { TaskMetadataParseError, UncaughtExceptionError } from "../common/errors"; import { setTimeout } from "node:timers/promises"; @@ -58,6 +57,8 @@ class ProdWorker { port: number, private host = "0.0.0.0" ) { + process.on("SIGTERM", this.#handleSignal.bind(this, "SIGTERM")); + this.#coordinatorSocket = this.#createCoordinatorSocket(COORDINATOR_HOST); this.#backgroundWorker = new ProdBackgroundWorker("worker.js", { @@ -150,6 +151,36 @@ class ProdWorker { this.#httpServer = this.#createHttpServer(); } + async #handleSignal(signal: NodeJS.Signals) { + logger.log("Received signal", { signal }); + + if (signal === "SIGTERM") { + if (this.executing) { + const terminationGracePeriodSeconds = 60 * 60; + + logger.log("Waiting for attempt to complete before exiting", { + terminationGracePeriodSeconds, + }); + + // Wait for termination grace period minus 5s to give cleanup a chance to complete + await setTimeout(terminationGracePeriodSeconds * 1000 - 5000); + + logger.log("Termination timeout reached, exiting gracefully."); + } else { + logger.log("Not executing, exiting immediately."); + } + + await this.#exitGracefully(); + } + + logger.log("Unhandled signal", { signal }); + } + + async #exitGracefully() { + await this.#backgroundWorker.close(); + process.exit(0); + } + async #reconnect(isPostStart = false, reconnectImmediately = false) { if (isPostStart) { this.waitForPostStart = false; @@ -206,8 +237,7 @@ class ProdWorker { logger.log("WARNING: Will checkpoint but also requested exit. This won't end well."); } - await this.#backgroundWorker.close(); - process.exit(0); + await this.#exitGracefully(); } this.executing = false; @@ -605,7 +635,6 @@ class ProdWorker { break; } } - logger.log("preStop", { url: req.url }); return reply.text("preStop ok"); } diff --git a/packages/cli-v3/src/workers/prod/worker-facade.ts b/packages/cli-v3/src/workers/prod/worker-facade.ts index 76f5cada0..892fe911e 100644 --- a/packages/cli-v3/src/workers/prod/worker-facade.ts +++ b/packages/cli-v3/src/workers/prod/worker-facade.ts @@ -175,6 +175,23 @@ const zodIpc = new ZodIpcConnection({ CLEANUP: async ({ flush, kill }, sender) => { if (kill) { await tracingSDK.flush(); + + if (_execution) { + // Fail currently executing attempt + await sender.send("TASK_RUN_COMPLETED", { + execution: _execution, + result: { + ok: false, + id: _execution.attempt.id, + error: { + type: "INTERNAL_ERROR", + code: TaskRunErrorCodes.GRACEFUL_EXIT_TIMEOUT, + message: "Worker process killed while attempt in progress.", + }, + }, + }); + } + // Now we need to exit the process await sender.send("READY_TO_DISPOSE", undefined); } else { @@ -186,6 +203,9 @@ const zodIpc = new ZodIpcConnection({ }, }); +// Ignore SIGTERM, handled by entry point +process.on("SIGTERM", async () => {}); + const prodRuntimeManager = new ProdRuntimeManager(zodIpc, { waitThresholdInMs: parseInt(process.env.TRIGGER_RUNTIME_WAIT_THRESHOLD_IN_MS ?? "30000", 10), }); diff --git a/packages/core/src/v3/runtime/devRuntimeManager.ts b/packages/core/src/v3/runtime/devRuntimeManager.ts index f1825130d..ffce394da 100644 --- a/packages/core/src/v3/runtime/devRuntimeManager.ts +++ b/packages/core/src/v3/runtime/devRuntimeManager.ts @@ -5,6 +5,7 @@ import { TaskRunExecutionResult, } from "../schemas"; import { RuntimeManager } from "./manager"; +import { unboundedTimeout } from "../utils/timers"; export class DevRuntimeManager implements RuntimeManager { _taskWaits: Map< @@ -24,15 +25,11 @@ export class DevRuntimeManager implements RuntimeManager { } async waitForDuration(ms: number): Promise { - return new Promise((resolve) => { - setTimeout(resolve, ms); - }); + await unboundedTimeout(ms); } async waitUntil(date: Date): Promise { - return new Promise((resolve) => { - setTimeout(resolve, date.getTime() - Date.now()); - }); + return this.waitForDuration(date.getTime() - Date.now()); } async waitForTask(params: { id: string; ctx: TaskRunContext }): Promise { diff --git a/packages/core/src/v3/runtime/prodRuntimeManager.ts b/packages/core/src/v3/runtime/prodRuntimeManager.ts index 5b0fdab92..6400e3942 100644 --- a/packages/core/src/v3/runtime/prodRuntimeManager.ts +++ b/packages/core/src/v3/runtime/prodRuntimeManager.ts @@ -10,6 +10,7 @@ import { } from "../schemas"; import { ZodIpcConnection } from "../zodIpc"; import { RuntimeManager } from "./manager"; +import { unboundedTimeout } from "../utils/timers"; export type ProdRuntimeManagerOptions = { waitThresholdInMs?: number; @@ -43,7 +44,7 @@ export class ProdRuntimeManager implements RuntimeManager { async waitForDuration(ms: number): Promise { const now = Date.now(); - const resolveAfterDuration = setTimeout(ms, "duration" as const); + const resolveAfterDuration = unboundedTimeout(ms, "duration" as const); if (ms <= this.waitThresholdInMs) { await resolveAfterDuration; diff --git a/packages/core/src/v3/schemas/common.ts b/packages/core/src/v3/schemas/common.ts index 0e393bb7a..e3cd2d5c4 100644 --- a/packages/core/src/v3/schemas/common.ts +++ b/packages/core/src/v3/schemas/common.ts @@ -34,6 +34,7 @@ export const TaskRunErrorCodes = { TASK_RUN_CANCELLED: "TASK_RUN_CANCELLED", TASK_OUTPUT_ERROR: "TASK_OUTPUT_ERROR", HANDLE_ERROR_ERROR: "HANDLE_ERROR_ERROR", + GRACEFUL_EXIT_TIMEOUT: "GRACEFUL_EXIT_TIMEOUT", } as const; export const TaskRunInternalError = z.object({ @@ -49,6 +50,7 @@ export const TaskRunInternalError = z.object({ "TASK_RUN_CANCELLED", "TASK_OUTPUT_ERROR", "HANDLE_ERROR_ERROR", + "GRACEFUL_EXIT_TIMEOUT" ]), message: z.string().optional(), }); diff --git a/packages/core/src/v3/utils/timers.ts b/packages/core/src/v3/utils/timers.ts new file mode 100644 index 000000000..7461a713c --- /dev/null +++ b/packages/core/src/v3/utils/timers.ts @@ -0,0 +1,21 @@ +import { TimerOptions } from "node:timers"; +import { setTimeout } from "node:timers/promises"; + +export async function unboundedTimeout( + delay: number = 0, + value?: T, + options?: TimerOptions +): Promise { + const maxDelay = 2147483647; // Highest value that will fit in a 32-bit signed integer + + const fullTimeouts = Math.floor(delay / maxDelay); + const remainingDelay = delay % maxDelay; + + let lastTimeoutResult = await setTimeout(remainingDelay, value, options); + + for (let i = 0; i < fullTimeouts; i++) { + lastTimeoutResult = await setTimeout(maxDelay, value, options); + } + + return lastTimeoutResult; +} diff --git a/references/v3-catalog/src/trigger/other.ts b/references/v3-catalog/src/trigger/other.ts index aba5dd5a4..813cd3693 100644 --- a/references/v3-catalog/src/trigger/other.ts +++ b/references/v3-catalog/src/trigger/other.ts @@ -1,4 +1,5 @@ -import { task } from "@trigger.dev/sdk/v3"; +import { logger, task, wait } from "@trigger.dev/sdk/v3"; +import { setTimeout } from "node:timers/promises"; export const loggingTask = task({ id: "logging-task-2", @@ -6,3 +7,16 @@ export const loggingTask = task({ console.log("Hello world"); }, }); + +export const waitForever = task({ + id: "wait-forever", + run: async (payload: { freeze?: boolean }) => { + if (payload.freeze) { + await wait.for({ years: 9999 }); + } else { + await logger.trace("Waiting..", async () => { + await setTimeout(2147483647); + }); + } + }, +});