From 3ee096d50c942a5741fd880cb22e61664e28950b Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Mon, 25 Mar 2024 14:30:50 +0000 Subject: [PATCH] skip connection handler when waiting for post start hook --- .../cli-v3/src/workers/prod/entry-point.ts | 30 +++++++++++++++---- 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/packages/cli-v3/src/workers/prod/entry-point.ts b/packages/cli-v3/src/workers/prod/entry-point.ts index 3e780ac8b..39a853798 100644 --- a/packages/cli-v3/src/workers/prod/entry-point.ts +++ b/packages/cli-v3/src/workers/prod/entry-point.ts @@ -42,6 +42,7 @@ class ProdWorker { private attemptFriendlyId?: string; private nextResumeAfter?: WaitReason; + private waitForPostStart = false; #httpPort: number; #backgroundWorker: ProdBackgroundWorker; @@ -79,6 +80,16 @@ class ProdWorker { }); this.#backgroundWorker.onCancelCheckpoint.attach(async (message) => { + logger.log("onCancelCheckpoint() clearing paused state, don't wait for post start hook", { + paused: this.paused, + nextResumeAfter: this.nextResumeAfter, + waitForPostStart: this.waitForPostStart, + }); + + this.paused = false; + this.nextResumeAfter = undefined; + this.waitForPostStart = false; + this.#coordinatorSocket.socket.emit("CANCEL_CHECKPOINT", { version: "v1" }); }); @@ -88,7 +99,6 @@ class ProdWorker { return; } - // TODO: Switch to .send() once coordinator uses zod handler for all messages const { willCheckpointAndRestore } = await this.#coordinatorSocket.socket.emitWithAck( "WAIT_FOR_DURATION", { @@ -106,7 +116,6 @@ class ProdWorker { return; } - // TODO: Switch to .send() once coordinator uses zod handler for all messages const { willCheckpointAndRestore } = await this.#coordinatorSocket.socket.emitWithAck( "WAIT_FOR_TASK", { @@ -124,7 +133,6 @@ class ProdWorker { return; } - // TODO: Switch to .send() once coordinator uses zod handler for all messages const { willCheckpointAndRestore } = await this.#coordinatorSocket.socket.emitWithAck( "WAIT_FOR_BATCH", { @@ -140,7 +148,11 @@ class ProdWorker { this.#httpServer = this.#createHttpServer(); } - async #reconnect() { + async #reconnect(isPostStart = false) { + if (isPostStart) { + this.waitForPostStart = false; + } + this.#coordinatorSocket.close(); if (!this.runningInKubernetes) { @@ -177,6 +189,7 @@ class ProdWorker { if (willCheckpointAndRestore) { this.paused = true; this.nextResumeAfter = reason; + this.waitForPostStart = true; } } @@ -197,8 +210,8 @@ class ProdWorker { this.attemptFriendlyId = undefined; if (willCheckpointAndRestore) { + this.waitForPostStart = true; this.#coordinatorSocket.socket.emit("READY_FOR_CHECKPOINT", { version: "v1" }); - this.#coordinatorSocket.close(); return; } } @@ -378,6 +391,11 @@ class ProdWorker { }, }, onConnection: async (socket, handler, sender, logger) => { + if (this.waitForPostStart) { + logger.log("skip connection handler, waiting for post start hook"); + return; + } + if (process.env.INDEX_TASKS === "true") { try { const taskResources = await this.#initializeWorker(); @@ -596,7 +614,7 @@ class ProdWorker { break; } case "restore": { - await this.#reconnect(); + await this.#reconnect(true); break; } default: {