From 16b344d67c9eb2ece1efc1e0cdc8e932bd61b8cc Mon Sep 17 00:00:00 2001 From: nicktrn <55853254+nicktrn@users.noreply.github.com> Date: Tue, 6 May 2025 14:44:46 +0100 Subject: [PATCH] update supervisor and schema --- apps/supervisor/src/workloadServer/index.ts | 26 +++++++++++++++++++ ...nFriendlyId.snapshots.since.$snapshotId.ts | 6 ++--- .../src/entryPoints/managed/execution.ts | 8 +++--- .../src/v3/runEngineWorker/supervisor/http.ts | 15 +++++++++++ .../v3/runEngineWorker/supervisor/schemas.ts | 2 +- 5 files changed, 49 insertions(+), 8 deletions(-) diff --git a/apps/supervisor/src/workloadServer/index.ts b/apps/supervisor/src/workloadServer/index.ts index 2dcf32973..cc41e2bfb 100644 --- a/apps/supervisor/src/workloadServer/index.ts +++ b/apps/supervisor/src/workloadServer/index.ts @@ -17,6 +17,7 @@ import { WorkloadRunAttemptStartRequestBody, type WorkloadRunAttemptStartResponseBody, type WorkloadRunLatestSnapshotResponseBody, + WorkloadRunSnapshotsSinceResponseBody, type WorkloadServerToClientEvents, type WorkloadSuspendRunResponseBody, } from "@trigger.dev/core/v3/workers"; @@ -341,6 +342,31 @@ export class WorkloadServer extends EventEmitter { } satisfies WorkloadRunLatestSnapshotResponseBody); }, }) + .route( + "/api/v1/workload-actions/runs/:runFriendlyId/snapshots/since/:snapshotFriendlyId", + "GET", + { + paramsSchema: WorkloadActionParams, + handler: async ({ req, reply, params }) => { + const sinceSnapshotResponse = await this.workerClient.getSnapshotsSince( + params.runFriendlyId, + params.snapshotFriendlyId, + this.runnerIdFromRequest(req) + ); + + if (!sinceSnapshotResponse.success) { + console.error("Failed to get snapshots since", { + runId: params.runFriendlyId, + error: sinceSnapshotResponse.error, + }); + reply.empty(500); + return; + } + + reply.json(sinceSnapshotResponse.data satisfies WorkloadRunSnapshotsSinceResponseBody); + }, + } + ) .route("/api/v1/workload-actions/runs/:runFriendlyId/logs/debug", "POST", { paramsSchema: WorkloadActionParams.pick({ runFriendlyId: true }), bodySchema: WorkloadDebugLogRequestBody, diff --git a/apps/webapp/app/routes/engine.v1.worker-actions.runs.$runFriendlyId.snapshots.since.$snapshotId.ts b/apps/webapp/app/routes/engine.v1.worker-actions.runs.$runFriendlyId.snapshots.since.$snapshotId.ts index 1531205c2..a79de5869 100644 --- a/apps/webapp/app/routes/engine.v1.worker-actions.runs.$runFriendlyId.snapshots.since.$snapshotId.ts +++ b/apps/webapp/app/routes/engine.v1.worker-actions.runs.$runFriendlyId.snapshots.since.$snapshotId.ts @@ -16,15 +16,15 @@ export const loader = createLoaderWorkerApiRoute( }): Promise> => { const { runFriendlyId, snapshotId } = params; - const executions = await authenticatedWorker.getSnapshotsSince({ + const snapshots = await authenticatedWorker.getSnapshotsSince({ runFriendlyId, snapshotId, }); - if (!executions) { + if (!snapshots) { throw new Error("Failed to retrieve snapshots since given snapshot"); } - return json({ executions }); + return json({ snapshots }); } ); diff --git a/packages/cli-v3/src/entryPoints/managed/execution.ts b/packages/cli-v3/src/entryPoints/managed/execution.ts index 4dda819a0..1c0b636fa 100644 --- a/packages/cli-v3/src/entryPoints/managed/execution.ts +++ b/packages/cli-v3/src/entryPoints/managed/execution.ts @@ -1105,22 +1105,22 @@ export class RunExecution { return; } - const { executions } = response.data; + const { snapshots } = response.data; - if (!executions.length) { + if (!snapshots.length) { this.sendDebugLog(`fetchAndProcessSnapshotChanges: no new snapshots`, { source }); return; } // Only act on the last snapshot - const lastSnapshot = executions[executions.length - 1]; + const lastSnapshot = snapshots[snapshots.length - 1]; if (!lastSnapshot) { this.sendDebugLog(`fetchAndProcessSnapshotChanges: no last snapshot`, { source }); return; } - const previousSnapshots = executions.slice(0, -1); + const previousSnapshots = snapshots.slice(0, -1); // If any previous snapshot is QUEUED or SUSPENDED, deprecate this worker const deprecatedStatus: TaskRunExecutionStatus[] = ["QUEUED", "SUSPENDED"]; diff --git a/packages/core/src/v3/runEngineWorker/supervisor/http.ts b/packages/core/src/v3/runEngineWorker/supervisor/http.ts index 4f899e4f2..43305b456 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/http.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/http.ts @@ -17,6 +17,7 @@ import { WorkerApiDebugLogBody, WorkerApiSuspendRunRequestBody, WorkerApiSuspendRunResponseBody, + WorkerApiRunSnapshotsSinceResponseBody, } from "./schemas.js"; import { SupervisorClientCommonOptions } from "./types.js"; import { getDefaultWorkerHeaders } from "./util.js"; @@ -185,6 +186,20 @@ export class SupervisorHttpClient { ); } + async getSnapshotsSince(runId: string, snapshotId: string, runnerId?: string) { + return wrapZodFetch( + WorkerApiRunSnapshotsSinceResponseBody, + `${this.apiUrl}/engine/v1/worker-actions/runs/${runId}/snapshots/since/${snapshotId}`, + { + method: "GET", + headers: { + ...this.defaultHeaders, + ...this.runnerIdHeader(runnerId), + }, + } + ); + } + async sendDebugLog(runId: string, body: WorkerApiDebugLogBody, runnerId?: string): Promise { try { const res = await wrapZodFetch( diff --git a/packages/core/src/v3/runEngineWorker/supervisor/schemas.ts b/packages/core/src/v3/runEngineWorker/supervisor/schemas.ts index deeaadf47..a49fd53a0 100644 --- a/packages/core/src/v3/runEngineWorker/supervisor/schemas.ts +++ b/packages/core/src/v3/runEngineWorker/supervisor/schemas.ts @@ -164,7 +164,7 @@ export type WorkerApiSuspendCompletionResponseBody = z.infer< >; export const WorkerApiRunSnapshotsSinceResponseBody = z.object({ - executions: z.array(RunExecutionData), + snapshots: z.array(RunExecutionData), }); export type WorkerApiRunSnapshotsSinceResponseBody = z.infer< typeof WorkerApiRunSnapshotsSinceResponseBody