From e2b10f42a032bd50222e6bd59b7f8ea8ce892e31 Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Thu, 21 Mar 2024 16:10:22 +0000 Subject: [PATCH] improve wait accuracy --- apps/coordinator/src/index.ts | 7 +++++++ .../v3/services/createCheckpoint.server.ts | 2 +- .../src/workers/prod/backgroundWorker.ts | 5 ++--- .../core/src/v3/runtime/prodRuntimeManager.ts | 7 ++++++- packages/core/src/v3/schemas/messages.ts | 1 + packages/core/src/v3/schemas/schemas.ts | 2 ++ packages/core/src/v3/zodSocket.ts | 19 ++++++++++++------- 7 files changed, 31 insertions(+), 12 deletions(-) diff --git a/apps/coordinator/src/index.ts b/apps/coordinator/src/index.ts index 9c7883d3f..91d670889 100644 --- a/apps/coordinator/src/index.ts +++ b/apps/coordinator/src/index.ts @@ -505,6 +505,12 @@ class TaskCoordinator { } confirmCompletion({ didCheckpoint: true, shouldExit: false, checkpoint }); + + if (!checkpoint.docker) { + socket.emit("REQUEST_EXIT", { + version: "v1", + }); + } }); socket.on("WAIT_FOR_DURATION", async (message, callback) => { @@ -550,6 +556,7 @@ class TaskCoordinator { reason: { type: "WAIT_FOR_DURATION", ms: message.ms, + now: message.now, }, }); }); diff --git a/apps/webapp/app/v3/services/createCheckpoint.server.ts b/apps/webapp/app/v3/services/createCheckpoint.server.ts index 9a91e0120..b62bb8e09 100644 --- a/apps/webapp/app/v3/services/createCheckpoint.server.ts +++ b/apps/webapp/app/v3/services/createCheckpoint.server.ts @@ -99,7 +99,7 @@ export class CreateCheckpointService { await marqs?.replaceMessage( attempt.taskRunId, { type: "RESUME_AFTER_DURATION", resumableAttemptId: attempt.id }, - Date.now() + params.reason.ms + params.reason.now + params.reason.ms ); break; } diff --git a/packages/cli-v3/src/workers/prod/backgroundWorker.ts b/packages/cli-v3/src/workers/prod/backgroundWorker.ts index 4fc2bfc25..4593432a6 100644 --- a/packages/cli-v3/src/workers/prod/backgroundWorker.ts +++ b/packages/cli-v3/src/workers/prod/backgroundWorker.ts @@ -17,7 +17,6 @@ import { } from "@trigger.dev/core/v3"; import { Evt } from "evt"; import { ChildProcess, fork } from "node:child_process"; -import { safeDeleteFileSync } from "../../utilities/fileSystem"; import { UncaughtExceptionError } from "../common/errors"; class UnexpectedExitError extends Error { @@ -56,7 +55,7 @@ export class ProdBackgroundWorker { public onTaskHeartbeat: Evt = new Evt(); - public onWaitForDuration: Evt<{ version?: "v1"; ms: number }> = new Evt(); + public onWaitForDuration: Evt<{ version?: "v1"; ms: number; now: number }> = new Evt(); public onWaitForTask: Evt<{ version?: "v1"; id: string }> = new Evt(); public onWaitForBatch: Evt<{ version?: "v1"; id: string; runs: string[] }> = new Evt(); @@ -343,7 +342,7 @@ class TaskRunProcess { public onExit: Evt = new Evt(); public onWaitForBatch: Evt<{ version?: "v1"; id: string; runs: string[] }> = new Evt(); - public onWaitForDuration: Evt<{ version?: "v1"; ms: number }> = new Evt(); + public onWaitForDuration: Evt<{ version?: "v1"; ms: number; now: number }> = new Evt(); public onWaitForTask: Evt<{ version?: "v1"; id: string }> = new Evt(); public preCheckpointNotification = Evt.create<{ willCheckpointAndRestore: boolean }>(); diff --git a/packages/core/src/v3/runtime/prodRuntimeManager.ts b/packages/core/src/v3/runtime/prodRuntimeManager.ts index 0f9889bb6..fe8a1450d 100644 --- a/packages/core/src/v3/runtime/prodRuntimeManager.ts +++ b/packages/core/src/v3/runtime/prodRuntimeManager.ts @@ -49,6 +49,8 @@ export class ProdRuntimeManager implements RuntimeManager { async waitForDuration(ms: number): Promise { let timeout: NodeJS.Timeout | undefined; + const now = Date.now(); + const resolveAfterDuration = new Promise((resolve) => { timeout = setTimeout(resolve, ms); }); @@ -63,7 +65,10 @@ export class ProdRuntimeManager implements RuntimeManager { }); // There is a slight delay before actually checkpointing, so this has a chance to return - const { willCheckpointAndRestore } = await this.ipc.sendWithAck("WAIT_FOR_DURATION", { ms }); + const { willCheckpointAndRestore } = await this.ipc.sendWithAck("WAIT_FOR_DURATION", { + ms, + now, + }); if (!willCheckpointAndRestore) { await resolveAfterDuration; diff --git a/packages/core/src/v3/schemas/messages.ts b/packages/core/src/v3/schemas/messages.ts index fa5de6285..3023ed91b 100644 --- a/packages/core/src/v3/schemas/messages.ts +++ b/packages/core/src/v3/schemas/messages.ts @@ -262,6 +262,7 @@ export const ProdChildToWorkerMessages = { message: z.object({ version: z.literal("v1").default("v1"), ms: z.number(), + now: z.number(), }), callback: z.object({ willCheckpointAndRestore: z.boolean(), diff --git a/packages/core/src/v3/schemas/schemas.ts b/packages/core/src/v3/schemas/schemas.ts index 85c6da359..7bb72ba05 100644 --- a/packages/core/src/v3/schemas/schemas.ts +++ b/packages/core/src/v3/schemas/schemas.ts @@ -204,6 +204,7 @@ export const CoordinatorToPlatformMessages = { z.object({ type: z.literal("WAIT_FOR_DURATION"), ms: z.number(), + now: z.number(), }), z.object({ type: z.literal("WAIT_FOR_BATCH"), @@ -353,6 +354,7 @@ export const ProdWorkerToCoordinatorMessages = { message: z.object({ version: z.literal("v1").default("v1"), ms: z.number(), + now: z.number(), }), callback: z.object({ willCheckpointAndRestore: z.boolean(), diff --git a/packages/core/src/v3/zodSocket.ts b/packages/core/src/v3/zodSocket.ts index 51e267bad..7e9371368 100644 --- a/packages/core/src/v3/zodSocket.ts +++ b/packages/core/src/v3/zodSocket.ts @@ -155,13 +155,18 @@ export class ZodSocketMessageHandler