From 6ad28b6818bb83e7ef686d2609becf4d72ccf808 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Wed, 8 May 2024 08:38:01 +0100 Subject: [PATCH] fresh attempt processes in dev and prod --- packages/cli-v3/src/workers/common/errors.ts | 32 +++ .../src/workers/dev/backgroundWorker.ts | 214 ++++++++++----- .../src/workers/prod/backgroundWorker.ts | 255 ++++++++++-------- 3 files changed, 330 insertions(+), 171 deletions(-) diff --git a/packages/cli-v3/src/workers/common/errors.ts b/packages/cli-v3/src/workers/common/errors.ts index 4017d3cc6..4ba1a8ef3 100644 --- a/packages/cli-v3/src/workers/common/errors.ts +++ b/packages/cli-v3/src/workers/common/errors.ts @@ -21,3 +21,35 @@ export class TaskMetadataParseError extends Error { this.name = "TaskMetadataParseError"; } } + +export class UnexpectedExitError extends Error { + constructor(public code: number) { + super(`Unexpected exit with code ${code}`); + + this.name = "UnexpectedExitError"; + } +} + +export class CleanupProcessError extends Error { + constructor() { + super("Cancelled"); + + this.name = "CleanupProcessError"; + } +} + +export class CancelledProcessError extends Error { + constructor() { + super("Cancelled"); + + this.name = "CancelledProcessError"; + } +} + +export class SigKillTimeoutProcessError extends Error { + constructor() { + super("Process kill timeout"); + + this.name = "SigKillTimeoutProcessError"; + } +} diff --git a/packages/cli-v3/src/workers/dev/backgroundWorker.ts b/packages/cli-v3/src/workers/dev/backgroundWorker.ts index 2ded486b1..b1867355b 100644 --- a/packages/cli-v3/src/workers/dev/backgroundWorker.ts +++ b/packages/cli-v3/src/workers/dev/backgroundWorker.ts @@ -1,5 +1,4 @@ import { - APIError, BackgroundWorkerProperties, BackgroundWorkerServerMessages, CreateBackgroundWorkerResponse, @@ -39,7 +38,14 @@ import { import { safeDeleteFileSync } from "../../utilities/fileSystem.js"; import { installPackages } from "../../utilities/installPackages.js"; import { logger } from "../../utilities/logger.js"; -import { TaskMetadataParseError, UncaughtExceptionError } from "../common/errors.js"; +import { + CancelledProcessError, + CleanupProcessError, + SigKillTimeoutProcessError, + TaskMetadataParseError, + UncaughtExceptionError, + UnexpectedExitError, +} from "../common/errors.js"; import { CliApiClient } from "../../apiClient.js"; export type CurrentWorkers = BackgroundWorkerCoordinator["currentWorkers"]; @@ -246,30 +252,6 @@ export class BackgroundWorkerCoordinator { } } -class UnexpectedExitError extends Error { - constructor(public code: number) { - super(`Unexpected exit with code ${code}`); - - this.name = "UnexpectedExitError"; - } -} - -class CleanupProcessError extends Error { - constructor() { - super("Cancelled"); - - this.name = "CleanupProcessError"; - } -} - -class CancelledProcessError extends Error { - constructor() { - super("Cancelled"); - - this.name = "CancelledProcessError"; - } -} - export type BackgroundWorkerParams = { env: Record; dependencies?: Record; @@ -295,6 +277,7 @@ export class BackgroundWorker { public metadata: BackgroundWorkerProperties | undefined; _taskRunProcesses: Map = new Map(); + private _taskRunProcessesBeingKilled: Set = new Set(); private _closed: boolean = false; @@ -422,39 +405,106 @@ export class BackgroundWorker { throw new Error("Worker not registered"); } - if (!this._taskRunProcesses.has(payload.execution.run.id)) { - const taskRunProcess = new TaskRunProcess( - payload.execution.run.id, - payload.execution.run.isTest, - this.path, - { - ...this.params.env, - ...(payload.environment ?? {}), - ...this.#readEnvVars(), - }, - this.metadata, - this.params, - messageId - ); + this._closed = false; - taskRunProcess.onExit.attach(() => { - this._taskRunProcesses.delete(payload.execution.run.id); - }); - - taskRunProcess.onTaskHeartbeat.attach((id) => { - this.onTaskHeartbeat.post(id); - }); - - taskRunProcess.onTaskRunHeartbeat.attach((id) => { - this.onTaskRunHeartbeat.post(id); - }); - - await taskRunProcess.initialize(); - - this._taskRunProcesses.set(payload.execution.run.id, taskRunProcess); + if (this._taskRunProcesses.has(payload.execution.run.id)) { + return this._taskRunProcesses.get(payload.execution.run.id) as TaskRunProcess; } - return this._taskRunProcesses.get(payload.execution.run.id) as TaskRunProcess; + await this.#killCurrentTaskRunProcessBeforeAttempt(payload.execution.run.id); + + const taskRunProcess = new TaskRunProcess( + payload.execution.run.id, + payload.execution.run.isTest, + this.path, + { + ...this.params.env, + ...(payload.environment ?? {}), + ...this.#readEnvVars(), + }, + this.metadata, + this.params, + messageId + ); + + taskRunProcess.onExit.attach(({ pid }) => { + this._taskRunProcesses.delete(payload.execution.run.id); + if (pid) { + this._taskRunProcessesBeingKilled.delete(pid); + } + }); + + taskRunProcess.onIsBeingKilled.attach((pid) => { + if (pid) { + this._taskRunProcessesBeingKilled.add(pid); + } + }); + + taskRunProcess.onTaskHeartbeat.attach((id) => { + this.onTaskHeartbeat.post(id); + }); + + taskRunProcess.onTaskRunHeartbeat.attach((id) => { + this.onTaskRunHeartbeat.post(id); + }); + + await taskRunProcess.initialize(); + + this._taskRunProcesses.set(payload.execution.run.id, taskRunProcess); + + return taskRunProcess; + } + + async #killCurrentTaskRunProcessBeforeAttempt(runId: string) { + const taskRunProcess = this._taskRunProcesses.get(runId); + + if (!taskRunProcess) { + return; + } + + if (taskRunProcess.isBeingKilled) { + if (this._taskRunProcessesBeingKilled.size > 1) { + // If there's more than one being killed, wait for graceful exit + try { + await taskRunProcess.onExit.waitFor(5_000); + } catch (error) { + console.error("TaskRunProcess graceful kill timeout exceeded", error); + + try { + const forcedKill = taskRunProcess.onExit.waitFor(5_000); + taskRunProcess.kill("SIGKILL"); + await forcedKill; + } catch (error) { + console.error("TaskRunProcess forced kill timeout exceeded", error); + throw new SigKillTimeoutProcessError(); + } + } + } else { + // If there's only one or none being killed, don't do anything so we can create a fresh one in parallel + } + } else { + // It's not being killed, so kill it + if (this._taskRunProcessesBeingKilled.size > 0) { + // If there's one being killed already, wait for graceful exit + try { + await taskRunProcess.onExit.waitFor(5_000); + } catch (error) { + console.error("TaskRunProcess graceful kill timeout exceeded", error); + + try { + const forcedKill = taskRunProcess.onExit.waitFor(5_000); + taskRunProcess.kill("SIGKILL"); + await forcedKill; + } catch (error) { + console.error("TaskRunProcess forced kill timeout exceeded", error); + throw new SigKillTimeoutProcessError(); + } + } + } else { + // There's none being killed yet, so we can kill it without waiting. We still set a timeout to kill it forcefully just in case it sticks around. + taskRunProcess.kill("SIGTERM", 5_000).catch(() => {}); + } + } } async cancelRun(taskRunId: string) { @@ -567,8 +617,8 @@ export class BackgroundWorker { const taskRunProcess = await this.#initializeTaskRunProcess(payload, messageId); const result = await taskRunProcess.executeTaskRun(payload); - // Kill the worker if the task was successful or if it's not going to be retried); - await taskRunProcess.cleanup(result.ok || result.retry === undefined); + // Always kill the worker + await taskRunProcess.cleanup(true); if (result.ok) { return result; @@ -669,6 +719,7 @@ class TaskRunProcess { }); private _sender: ZodMessageSender; private _child: ChildProcess | undefined; + private _childPid?: number; private _attemptPromises: Map< string, { resolver: (value: TaskRunExecutionResult) => void; rejecter: (err?: any) => void } @@ -682,7 +733,9 @@ class TaskRunProcess { */ public onTaskHeartbeat: Evt = new Evt(); public onTaskRunHeartbeat: Evt = new Evt(); - public onExit: Evt = new Evt(); + public onExit: Evt<{ code: number | null; signal: NodeJS.Signals | null; pid?: number }> = + new Evt(); + public onIsBeingKilled: Evt = new Evt(); constructor( private runId: string, @@ -736,6 +789,7 @@ class TaskRunProcess { ? ["--inspect-brk", "--trace-uncaught", "--no-warnings=ExperimentalWarning"] : ["--trace-uncaught", "--no-warnings=ExperimentalWarning"], }); + this._childPid = this._child?.pid; this._child.on("message", this.#handleMessage.bind(this)); this._child.on("exit", this.#handleExit.bind(this)); @@ -748,6 +802,11 @@ class TaskRunProcess { return; } + if (kill) { + this._isBeingKilled = true; + this.onIsBeingKilled.post(this._child?.pid); + } + logger.debug(`[${this.runId}] cleaning up task run process`, { kill }); await this._sender.send("CLEANUP", { @@ -755,7 +814,7 @@ class TaskRunProcess { kill, }); - this._isBeingKilled = kill; + // FIXME: Something broke READY_TO_DISPOSE. We never receive it, so we always have to kill the process after the timeout below. // Set a timeout to kill the child process if it hasn't been killed within 5 seconds setTimeout(() => { @@ -867,8 +926,8 @@ class TaskRunProcess { } } - async #handleExit(code: number) { - logger.debug(`[${this.runId}] task run process exiting`, { code }); + async #handleExit(code: number | null, signal: NodeJS.Signals | null) { + logger.debug(`[${this.runId}] task run process exiting`, { code, signal }); // Go through all the attempts currently pending and reject them for (const [id, status] of this._attemptStatuses.entries()) { @@ -888,12 +947,12 @@ class TaskRunProcess { } else if (this._isBeingKilled) { rejecter(new CleanupProcessError()); } else { - rejecter(new UnexpectedExitError(code)); + rejecter(new UnexpectedExitError(code ?? -1)); } } } - this.onExit.post(code); + this.onExit.post({ code, signal, pid: this.pid }); } #handleLog(data: Buffer) { @@ -939,6 +998,33 @@ class TaskRunProcess { this._child?.kill(); } } + + async kill(signal?: number | NodeJS.Signals, timeoutInMs?: number) { + logger.debug(`[${this.runId}] killing task run process`, { + signal, + timeoutInMs, + pid: this.pid, + }); + + this._isBeingKilled = true; + + const killTimeout = this.onExit.waitFor(timeoutInMs); + + this.onIsBeingKilled.post(this._child?.pid); + this._child?.kill(signal); + + if (timeoutInMs) { + await killTimeout; + } + } + + get isBeingKilled() { + return this._isBeingKilled || this._child?.killed; + } + + get pid() { + return this._childPid; + } } function formatErrorLog(error: TaskRunError) { diff --git a/packages/cli-v3/src/workers/prod/backgroundWorker.ts b/packages/cli-v3/src/workers/prod/backgroundWorker.ts index 91c898e98..7536777a2 100644 --- a/packages/cli-v3/src/workers/prod/backgroundWorker.ts +++ b/packages/cli-v3/src/workers/prod/backgroundWorker.ts @@ -21,39 +21,14 @@ import { ZodIpcConnection } from "@trigger.dev/core/v3/zodIpc"; import type { InferSocketMessageSchema } from "@trigger.dev/core/v3/zodSocket"; import { Evt } from "evt"; import { ChildProcess, fork } from "node:child_process"; -import { TaskMetadataParseError, UncaughtExceptionError } from "../common/errors"; - -class UnexpectedExitError extends Error { - constructor(public code: number) { - super(`Unexpected exit with code ${code}`); - - this.name = "UnexpectedExitError"; - } -} - -class CleanupProcessError extends Error { - constructor() { - super("Cancelled"); - - this.name = "CleanupProcessError"; - } -} - -class CancelledProcessError extends Error { - constructor() { - super("Cancelled"); - - this.name = "CancelledProcessError"; - } -} - -class SigKillTimeoutProcessError extends Error { - constructor() { - super("Process kill timeout"); - - this.name = "SigKillTimeoutProcessError"; - } -} +import { + CancelledProcessError, + CleanupProcessError, + SigKillTimeoutProcessError, + TaskMetadataParseError, + UncaughtExceptionError, + UnexpectedExitError, +} from "../common/errors"; type BackgroundWorkerParams = { env: Record; @@ -104,6 +79,7 @@ export class ProdBackgroundWorker { public tasks: Array = []; _taskRunProcess: TaskRunProcess | undefined; + private _taskRunProcessesBeingKilled: Set = new Set(); private _closed: boolean = false; @@ -135,14 +111,16 @@ export class ProdBackgroundWorker { await this.flushTelemetry(); } + const currentTaskRunProcess = this._taskRunProcess; + try { - const initialExit = this._taskRunProcess.onExit.waitFor(5_000); - this._taskRunProcess.kill(initialSignal); + const initialExit = currentTaskRunProcess.onExit.waitFor(5_000); + currentTaskRunProcess.kill(initialSignal); await initialExit; } catch (error) { // Try again with SIGKILL - const forcedExit = this._taskRunProcess.onExit.waitFor(5_000); - this._taskRunProcess.kill("SIGKILL"); + const forcedExit = currentTaskRunProcess.onExit.waitFor(5_000); + currentTaskRunProcess.kill("SIGKILL"); await forcedExit; } @@ -250,7 +228,7 @@ export class ProdBackgroundWorker { this._taskRunProcess?.waitCompletedNotification(); } - async #initializeTaskRunProcess( + async #getFreshTaskRunProcess( payload: ProdTaskRunExecutionPayload, messageId?: string ): Promise { @@ -261,93 +239,136 @@ export class ProdBackgroundWorker { this._closed = false; - // If the child process is currently being killed, we should wait for it to be dead before creating a fresh one (with a sensible timeout) - if (this._taskRunProcess?.isBeingKilled) { - try { - await this._taskRunProcess.onExit.waitFor(5_000); - } catch (error) { - console.error("TaskRunProcess graceful kill timeout exceeded", error); + await this.#killCurrentTaskRunProcessBeforeAttempt(); - try { - const forcedKill = this._taskRunProcess.onExit.waitFor(5_000); - this._taskRunProcess.kill("SIGKILL"); - await forcedKill; - } catch (error) { - console.error("TaskRunProcess forced kill timeout exceeded", error); - throw new SigKillTimeoutProcessError(); - } + const taskRunProcess = new TaskRunProcess( + payload.execution.run.id, + payload.execution.run.isTest, + this.path, + { + ...this.params.env, + ...(payload.environment ?? {}), + }, + metadata, + this.params, + messageId + ); + + taskRunProcess.onExit.attach(({ pid }) => { + this._taskRunProcess = undefined; + if (pid) { + this._taskRunProcessesBeingKilled.delete(pid); } - } + }); - if (!this._taskRunProcess) { - const taskRunProcess = new TaskRunProcess( - payload.execution.run.id, - payload.execution.run.isTest, - this.path, - { - ...this.params.env, - ...(payload.environment ?? {}), - }, - metadata, - this.params, - messageId - ); + taskRunProcess.onIsBeingKilled.attach((pid) => { + if (pid) { + this._taskRunProcessesBeingKilled.add(pid); + } + }); - taskRunProcess.onExit.attach(() => { - this._taskRunProcess = undefined; - }); + taskRunProcess.onTaskHeartbeat.attach((id) => { + this.onTaskHeartbeat.post(id); + }); - taskRunProcess.onTaskHeartbeat.attach((id) => { - this.onTaskHeartbeat.post(id); - }); + taskRunProcess.onTaskRunHeartbeat.attach((id) => { + this.onTaskRunHeartbeat.post(id); + }); - taskRunProcess.onTaskRunHeartbeat.attach((id) => { - this.onTaskRunHeartbeat.post(id); - }); + taskRunProcess.onWaitForBatch.attach((message) => { + this.onWaitForBatch.post(message); + }); - taskRunProcess.onWaitForBatch.attach((message) => { - this.onWaitForBatch.post(message); - }); + taskRunProcess.onWaitForDuration.attach((message) => { + this.onWaitForDuration.post(message); + }); - taskRunProcess.onWaitForDuration.attach((message) => { - this.onWaitForDuration.post(message); - }); + taskRunProcess.onWaitForTask.attach((message) => { + this.onWaitForTask.post(message); + }); - taskRunProcess.onWaitForTask.attach((message) => { - this.onWaitForTask.post(message); - }); + taskRunProcess.onReadyForCheckpoint.attach((message) => { + this.onReadyForCheckpoint.post(message); + }); - taskRunProcess.onReadyForCheckpoint.attach((message) => { - this.onReadyForCheckpoint.post(message); - }); + taskRunProcess.onCancelCheckpoint.attach((message) => { + this.onCancelCheckpoint.post(message); + }); - taskRunProcess.onCancelCheckpoint.attach((message) => { - this.onCancelCheckpoint.post(message); - }); + // Notify down the chain + this.preCheckpointNotification.attach((message) => { + taskRunProcess.preCheckpointNotification.post(message); + }); + this.checkpointCanceledNotification.attach((message) => { + taskRunProcess.checkpointCanceledNotification.post(message); + }); - // Notify down the chain - this.preCheckpointNotification.attach((message) => { - taskRunProcess.preCheckpointNotification.post(message); - }); - this.checkpointCanceledNotification.attach((message) => { - taskRunProcess.checkpointCanceledNotification.post(message); - }); + await taskRunProcess.initialize(); - await taskRunProcess.initialize(); - - this._taskRunProcess = taskRunProcess; - } + this._taskRunProcess = taskRunProcess; return this._taskRunProcess; } - // We need to fork the process before we can execute any tasks + async #killCurrentTaskRunProcessBeforeAttempt() { + if (!this._taskRunProcess) { + return; + } + + const currentTaskRunProcess = this._taskRunProcess; + + if (currentTaskRunProcess.isBeingKilled) { + if (this._taskRunProcessesBeingKilled.size > 1) { + // If there's more than one being killed, wait for graceful exit + try { + await currentTaskRunProcess.onExit.waitFor(5_000); + } catch (error) { + console.error("TaskRunProcess graceful kill timeout exceeded", error); + + try { + const forcedKill = currentTaskRunProcess.onExit.waitFor(5_000); + currentTaskRunProcess.kill("SIGKILL"); + await forcedKill; + } catch (error) { + console.error("TaskRunProcess forced kill timeout exceeded", error); + throw new SigKillTimeoutProcessError(); + } + } + } else { + // If there's only one or none being killed, don't do anything so we can create a fresh one in parallel + } + } else { + // It's not being killed, so kill it + if (this._taskRunProcessesBeingKilled.size > 0) { + // If there's one being killed already, wait for graceful exit + try { + await currentTaskRunProcess.onExit.waitFor(5_000); + } catch (error) { + console.error("TaskRunProcess graceful kill timeout exceeded", error); + + try { + const forcedKill = currentTaskRunProcess.onExit.waitFor(5_000); + currentTaskRunProcess.kill("SIGKILL"); + await forcedKill; + } catch (error) { + console.error("TaskRunProcess forced kill timeout exceeded", error); + throw new SigKillTimeoutProcessError(); + } + } + } else { + // There's none being killed yet, so we can kill it without waiting. We still set a timeout to kill it forcefully just in case it sticks around. + currentTaskRunProcess.kill("SIGTERM", 5_000).catch(() => {}); + } + } + } + + // We need to fork the process before we can execute any tasks, use a fresh process for each execution async executeTaskRun( payload: ProdTaskRunExecutionPayload, messageId?: string ): Promise { try { - const taskRunProcess = await this.#initializeTaskRunProcess(payload, messageId); + const taskRunProcess = await this.#getFreshTaskRunProcess(payload, messageId); const result = await taskRunProcess.executeTaskRun(payload); @@ -483,6 +504,7 @@ class TaskRunProcess { typeof ProdWorkerToChildMessages >; private _child?: ChildProcess; + private _childPid?: number; private _attemptPromises: Map< string, @@ -498,7 +520,9 @@ class TaskRunProcess { */ public onTaskHeartbeat: Evt = new Evt(); public onTaskRunHeartbeat: Evt = new Evt(); - public onExit: Evt<{ code: number | null; signal: NodeJS.Signals | null }> = new Evt(); + public onExit: Evt<{ code: number | null; signal: NodeJS.Signals | null; pid?: number }> = + new Evt(); + public onIsBeingKilled: Evt = new Evt(); public onWaitForBatch: Evt< InferSocketMessageSchema @@ -538,6 +562,7 @@ class TaskRunProcess { ...(this.worker.debugOtel ? { OTEL_LOG_LEVEL: "debug" } : {}), }, }); + this._childPid = this._child?.pid; this._ipc = new ZodIpcConnection({ listenSchema: ProdChildToWorkerMessages, @@ -652,7 +677,10 @@ class TaskRunProcess { return; } - this._isBeingKilled = kill; + if (kill) { + this._isBeingKilled = true; + this.onIsBeingKilled.post(this._child?.pid); + } await this._ipc?.sendWithAck("CLEANUP", { flush: true, @@ -736,7 +764,7 @@ class TaskRunProcess { } } - this.onExit.post({ code, signal }); + this.onExit.post({ code, signal, pid: this.pid }); } #handleLog(data: Buffer) { @@ -769,11 +797,24 @@ class TaskRunProcess { ); } - kill(signal?: number | NodeJS.Signals) { + async kill(signal?: number | NodeJS.Signals, timeoutInMs?: number) { + this._isBeingKilled = true; + + const killTimeout = this.onExit.waitFor(timeoutInMs); + + this.onIsBeingKilled.post(this._child?.pid); this._child?.kill(signal); + + if (timeoutInMs) { + await killTimeout; + } } get isBeingKilled() { return this._isBeingKilled || this._child?.killed; } + + get pid() { + return this._childPid; + } }