From 09564d5d28313f6143e8f20a161c8df7e1241766 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Tue, 30 Apr 2024 14:05:13 +0100 Subject: [PATCH] wait for post start --- .../cli-v3/src/workers/prod/backgroundWorker.ts | 6 +++--- packages/cli-v3/src/workers/prod/entry-point.ts | 17 ++++++++++++++++- 2 files changed, 19 insertions(+), 4 deletions(-) diff --git a/packages/cli-v3/src/workers/prod/backgroundWorker.ts b/packages/cli-v3/src/workers/prod/backgroundWorker.ts index fe3a2ed7a..8f2bbe897 100644 --- a/packages/cli-v3/src/workers/prod/backgroundWorker.ts +++ b/packages/cli-v3/src/workers/prod/backgroundWorker.ts @@ -442,6 +442,9 @@ class TaskRunProcess { this.onTaskHeartbeat.post(message.id); }, TASKS_READY: async (message) => {}, + WAIT_FOR_TASK: async (message) => { + this.onWaitForTask.post(message); + }, WAIT_FOR_BATCH: async (message) => { this.onWaitForBatch.post(message); }, @@ -467,9 +470,6 @@ class TaskRunProcess { }; } }, - WAIT_FOR_TASK: async (message) => { - this.onWaitForTask.post(message); - }, READY_FOR_CHECKPOINT: async (message) => { this.onReadyForCheckpoint.post(message); }, diff --git a/packages/cli-v3/src/workers/prod/entry-point.ts b/packages/cli-v3/src/workers/prod/entry-point.ts index 0a4d18c91..457ada131 100644 --- a/packages/cli-v3/src/workers/prod/entry-point.ts +++ b/packages/cli-v3/src/workers/prod/entry-point.ts @@ -43,6 +43,7 @@ class ProdWorker { private attemptFriendlyId?: string; private nextResumeAfter?: WaitReason; + private waitForPostStart = false; #httpPort: number; #backgroundWorker: ProdBackgroundWorker; @@ -100,6 +101,7 @@ class ProdWorker { // Worker will resume immediately this.paused = false; this.nextResumeAfter = undefined; + this.waitForPostStart = false; } } @@ -192,6 +194,10 @@ class ProdWorker { } async #reconnect(isPostStart = false, reconnectImmediately = false) { + if (isPostStart) { + this.waitForPostStart = false; + } + this.#coordinatorSocket.close(); if (!reconnectImmediately) { @@ -222,7 +228,7 @@ class ProdWorker { } } - #prepareForWait(reason: WaitReason, willCheckpointAndRestore: boolean) { + async #prepareForWait(reason: WaitReason, willCheckpointAndRestore: boolean) { logger.log(`prepare for ${reason}`, { willCheckpointAndRestore }); this.#backgroundWorker.preCheckpointNotification.post({ willCheckpointAndRestore }); @@ -230,6 +236,7 @@ class ProdWorker { if (willCheckpointAndRestore) { this.paused = true; this.nextResumeAfter = reason; + this.waitForPostStart = true; if (reason === "WAIT_FOR_TASK" || reason === "WAIT_FOR_BATCH") { // Flush before checkpointing so we don't flush the same spans again after restore @@ -255,6 +262,7 @@ class ProdWorker { this.attemptFriendlyId = undefined; if (willCheckpointAndRestore) { + this.waitForPostStart = true; this.#coordinatorSocket.socket.emit("READY_FOR_CHECKPOINT", { version: "v1" }); return; } @@ -263,6 +271,7 @@ class ProdWorker { #resumeAfterDuration() { this.paused = false; this.nextResumeAfter = undefined; + this.waitForPostStart = false; this.#backgroundWorker.waitCompletedNotification(); } @@ -353,6 +362,7 @@ class ProdWorker { this.paused = false; this.nextResumeAfter = undefined; + this.waitForPostStart = false; for (let i = 0; i < message.completions.length; i++) { const completion = message.completions[i]; @@ -434,6 +444,11 @@ class ProdWorker { }, }, onConnection: async (socket, handler, sender, logger) => { + if (this.waitForPostStart) { + logger.log("skip connection handler, waiting for post start hook"); + return; + } + if (this.paused) { if (!this.nextResumeAfter) { return;