chore(webapp,core): remove the end-of-life v3 (engine V1) execution stack (#4236)
## Summary v3 (the engine that ran the SDK v3 era, internally `RunEngineVersion.V1`) is end-of-life. Following the removal of the v3 execution apps ([#4194](https://github.com/triggerdotdev/trigger.dev/pull/4194)) and the legacy dev websocket ([#4198](https://github.com/triggerdotdev/trigger.dev/pull/4198)), this removes the remaining v3 execution stack from the server. Clients still on v3 (an old SDK or CLI that has not upgraded) keep getting a clear "upgrade to v4" response. Triggers, batch triggers, reschedules, and deploys that resolve to v3 are rejected with a graceful 4xx pointing at the migration guide, never a 5xx, so a stale client cannot affect server health. Self-hosted instances still running v3 should stay on the 4.5.x release line until they migrate. ## What is removed - The MarQS queue and its shared/dev queue consumers. - The v3 socket.io namespaces (coordinator, provider, shared-queue) and the v3 run lifecycle services (attempt, checkpoint, and batch-resume). - The graphile-worker background job system; all live jobs already run on `@trigger.dev/redis-worker`. - The `DEPRECATE_V3_ENABLED` flag: v3 is now rejected unconditionally, so the flag is gone. - Unused v3 exports from `@trigger.dev/core` (the `v3/zodNamespace` subpath and the legacy socket message catalogs) and the now-dead MarQS environment variables. ## What stays The v4 engine is untouched. The graceful v3 rejection boundary stays, `determineEngineVersion` still detects a v3 project so it can reject it, and the batch service plus batch-completion worker stay for current clients. Live queue concurrency limits and metrics now read from the v4 run engine instead of MarQS, and a brand-new dev environment now defaults to v4. ## Dependency cleanup Removes webapp dependencies left unused by this change: `seedrandom` and `semver` (only the removed v3 code used them) plus a set that was already dead, their orphaned `@types` packages, and two dead files. Adds a `knip:deps` script and a `knip.json` config so unused dependencies can be found the same way going forward.
This commit is contained in:
@@ -47,7 +47,6 @@
|
||||
"./v3/test": "./src/v3/test/index.ts",
|
||||
"./v3/zodfetch": "./src/v3/zodfetch.ts",
|
||||
"./v3/zodMessageHandler": "./src/v3/zodMessageHandler.ts",
|
||||
"./v3/zodNamespace": "./src/v3/zodNamespace.ts",
|
||||
"./v3/zodSocket": "./src/v3/zodSocket.ts",
|
||||
"./v3/zodIpc": "./src/v3/zodIpc.ts",
|
||||
"./v3/utils/timers": "./src/v3/utils/timers.ts",
|
||||
@@ -134,9 +133,6 @@
|
||||
"v3/zodMessageHandler": [
|
||||
"dist/commonjs/v3/zodMessageHandler.d.ts"
|
||||
],
|
||||
"v3/zodNamespace": [
|
||||
"dist/commonjs/v3/zodNamespace.d.ts"
|
||||
],
|
||||
"v3/zodSocket": [
|
||||
"dist/commonjs/v3/zodSocket.d.ts"
|
||||
],
|
||||
@@ -538,17 +534,6 @@
|
||||
"default": "./dist/commonjs/v3/zodMessageHandler.js"
|
||||
}
|
||||
},
|
||||
"./v3/zodNamespace": {
|
||||
"import": {
|
||||
"@triggerdotdev/source": "./src/v3/zodNamespace.ts",
|
||||
"types": "./dist/esm/v3/zodNamespace.d.ts",
|
||||
"default": "./dist/esm/v3/zodNamespace.js"
|
||||
},
|
||||
"require": {
|
||||
"types": "./dist/commonjs/v3/zodNamespace.d.ts",
|
||||
"default": "./dist/commonjs/v3/zodNamespace.js"
|
||||
}
|
||||
},
|
||||
"./v3/zodSocket": {
|
||||
"import": {
|
||||
"@triggerdotdev/source": "./src/v3/zodSocket.ts",
|
||||
|
||||
@@ -1,110 +1,10 @@
|
||||
import { z } from "zod";
|
||||
import { ImportTaskFileErrors, WorkerManifest } from "./build.js";
|
||||
import {
|
||||
MachinePreset,
|
||||
TaskRunExecution,
|
||||
TaskRunExecutionResult,
|
||||
TaskRunFailedExecutionResult,
|
||||
TaskRunInternalError,
|
||||
V3TaskRunExecution,
|
||||
} from "./common.js";
|
||||
import { TaskResource } from "./resources.js";
|
||||
import {
|
||||
EnvironmentType,
|
||||
V3ProdTaskRunExecution,
|
||||
V3ProdTaskRunExecutionPayload,
|
||||
RunEngineVersionSchema,
|
||||
TaskRunExecutionLazyAttemptPayload,
|
||||
TaskRunExecutionMetrics,
|
||||
WaitReason,
|
||||
} from "./schemas.js";
|
||||
import { TaskRunExecution, TaskRunExecutionResult } from "./common.js";
|
||||
import { RunEngineVersionSchema, TaskRunExecutionMetrics } from "./schemas.js";
|
||||
import { CompletedWaitpoint } from "./runEngine.js";
|
||||
import { DebugLogPropertiesInput } from "../runEngineWorker/supervisor/schemas.js";
|
||||
|
||||
export const AckCallbackResult = z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
stderr: z.string().optional(),
|
||||
}),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
}),
|
||||
]);
|
||||
|
||||
export type AckCallbackResult = z.infer<typeof AckCallbackResult>;
|
||||
|
||||
export const BackgroundWorkerServerMessages = z.discriminatedUnion("type", [
|
||||
z.object({
|
||||
type: z.literal("CANCEL_ATTEMPT"),
|
||||
taskAttemptId: z.string(),
|
||||
taskRunId: z.string(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("SCHEDULE_ATTEMPT"),
|
||||
image: z.string(),
|
||||
version: z.string(),
|
||||
machine: MachinePreset,
|
||||
nextAttemptNumber: z.number().optional(),
|
||||
// identifiers
|
||||
id: z.string().optional(), // TODO: Remove this completely in a future release
|
||||
envId: z.string(),
|
||||
envType: EnvironmentType,
|
||||
orgId: z.string(),
|
||||
projectId: z.string(),
|
||||
runId: z.string(),
|
||||
dequeuedAt: z.number().optional(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("EXECUTE_RUN_LAZY_ATTEMPT"),
|
||||
payload: TaskRunExecutionLazyAttemptPayload,
|
||||
}),
|
||||
]);
|
||||
|
||||
export type BackgroundWorkerServerMessages = z.infer<typeof BackgroundWorkerServerMessages>;
|
||||
|
||||
export const serverWebsocketMessages = {
|
||||
SERVER_READY: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
id: z.string(),
|
||||
}),
|
||||
BACKGROUND_WORKER_MESSAGE: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
data: BackgroundWorkerServerMessages,
|
||||
}),
|
||||
};
|
||||
|
||||
export const BackgroundWorkerClientMessages = z.discriminatedUnion("type", [
|
||||
z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
type: z.literal("TASK_RUN_COMPLETED"),
|
||||
completion: TaskRunExecutionResult,
|
||||
execution: V3TaskRunExecution,
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
type: z.literal("TASK_RUN_FAILED_TO_RUN"),
|
||||
completion: TaskRunFailedExecutionResult,
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
type: z.literal("TASK_HEARTBEAT"),
|
||||
id: z.string(),
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
type: z.literal("TASK_RUN_HEARTBEAT"),
|
||||
id: z.string(),
|
||||
}),
|
||||
]);
|
||||
|
||||
export type BackgroundWorkerClientMessages = z.infer<typeof BackgroundWorkerClientMessages>;
|
||||
|
||||
export const ServerBackgroundWorker = z.object({
|
||||
id: z.string(),
|
||||
version: z.string(),
|
||||
@@ -114,23 +14,6 @@ export const ServerBackgroundWorker = z.object({
|
||||
|
||||
export type ServerBackgroundWorker = z.infer<typeof ServerBackgroundWorker>;
|
||||
|
||||
export const clientWebsocketMessages = {
|
||||
READY_FOR_TASKS: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
inProgressRuns: z.string().array().optional(),
|
||||
}),
|
||||
BACKGROUND_WORKER_DEPRECATED: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
}),
|
||||
BACKGROUND_WORKER_MESSAGE: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
data: BackgroundWorkerClientMessages,
|
||||
}),
|
||||
};
|
||||
|
||||
export const UncaughtExceptionMessage = z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
error: z.object({
|
||||
@@ -234,667 +117,3 @@ export const WorkerToExecutorMessageCatalog = {
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const ProviderToPlatformMessages = {
|
||||
LOG: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
data: z.string(),
|
||||
}),
|
||||
},
|
||||
LOG_WITH_ACK: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
data: z.string(),
|
||||
}),
|
||||
callback: z.object({
|
||||
status: z.literal("ok"),
|
||||
}),
|
||||
},
|
||||
WORKER_CRASHED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
reason: z.string().optional(),
|
||||
exitCode: z.number().optional(),
|
||||
message: z.string().optional(),
|
||||
logs: z.string().optional(),
|
||||
/** This means we should only update the error if one exists */
|
||||
overrideCompletion: z.boolean().optional(),
|
||||
errorCode: TaskRunInternalError.shape.code.optional(),
|
||||
}),
|
||||
},
|
||||
INDEXING_FAILED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
deploymentId: z.string(),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
stderr: z.string().optional(),
|
||||
}),
|
||||
overrideCompletion: z.boolean().optional(),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const PlatformToProviderMessages = {
|
||||
INDEX: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
imageTag: z.string(),
|
||||
shortCode: z.string(),
|
||||
apiKey: z.string(),
|
||||
apiUrl: z.string(),
|
||||
// identifiers
|
||||
envId: z.string(),
|
||||
envType: EnvironmentType,
|
||||
orgId: z.string(),
|
||||
projectId: z.string(),
|
||||
deploymentId: z.string(),
|
||||
}),
|
||||
callback: AckCallbackResult,
|
||||
},
|
||||
RESTORE: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
type: z.enum(["DOCKER", "KUBERNETES"]),
|
||||
location: z.string(),
|
||||
reason: z.string().optional(),
|
||||
imageRef: z.string(),
|
||||
attemptNumber: z.number().optional(),
|
||||
machine: MachinePreset,
|
||||
// identifiers
|
||||
checkpointId: z.string(),
|
||||
envId: z.string(),
|
||||
envType: EnvironmentType,
|
||||
orgId: z.string(),
|
||||
projectId: z.string(),
|
||||
runId: z.string(),
|
||||
}),
|
||||
},
|
||||
PRE_PULL_DEPLOYMENT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
imageRef: z.string(),
|
||||
shortCode: z.string(),
|
||||
// identifiers
|
||||
envId: z.string(),
|
||||
envType: EnvironmentType,
|
||||
orgId: z.string(),
|
||||
projectId: z.string(),
|
||||
deploymentId: z.string(),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
const CreateWorkerMessage = z.object({
|
||||
projectRef: z.string(),
|
||||
envId: z.string(),
|
||||
deploymentId: z.string(),
|
||||
metadata: z.object({
|
||||
cliPackageVersion: z.string().optional(),
|
||||
contentHash: z.string(),
|
||||
packageVersion: z.string(),
|
||||
tasks: TaskResource.array(),
|
||||
}),
|
||||
});
|
||||
|
||||
export const CoordinatorToPlatformMessages = {
|
||||
LOG: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
metadata: z.any(),
|
||||
text: z.string(),
|
||||
}),
|
||||
},
|
||||
CREATE_WORKER: {
|
||||
message: z.discriminatedUnion("version", [
|
||||
CreateWorkerMessage.extend({
|
||||
version: z.literal("v1"),
|
||||
}),
|
||||
CreateWorkerMessage.extend({
|
||||
version: z.literal("v2"),
|
||||
supportsLazyAttempts: z.boolean(),
|
||||
}),
|
||||
]),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
}),
|
||||
]),
|
||||
},
|
||||
CREATE_TASK_RUN_ATTEMPT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
envId: z.string(),
|
||||
}),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
reason: z.string().optional(),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
executionPayload: V3ProdTaskRunExecutionPayload,
|
||||
}),
|
||||
]),
|
||||
},
|
||||
// Deprecated: Only workers without lazy attempt support will use this
|
||||
READY_FOR_EXECUTION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
totalCompletions: z.number(),
|
||||
}),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
payload: V3ProdTaskRunExecutionPayload,
|
||||
}),
|
||||
]),
|
||||
},
|
||||
READY_FOR_LAZY_ATTEMPT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
envId: z.string(),
|
||||
totalCompletions: z.number(),
|
||||
}),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
reason: z.string().optional(),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
lazyPayload: TaskRunExecutionLazyAttemptPayload,
|
||||
}),
|
||||
]),
|
||||
},
|
||||
READY_FOR_RESUME: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptFriendlyId: z.string(),
|
||||
type: WaitReason,
|
||||
}),
|
||||
},
|
||||
TASK_RUN_COMPLETED: {
|
||||
message: z.object({
|
||||
version: z.enum(["v1", "v2"]).default("v1"),
|
||||
execution: V3ProdTaskRunExecution,
|
||||
completion: TaskRunExecutionResult,
|
||||
checkpoint: z
|
||||
.object({
|
||||
docker: z.boolean(),
|
||||
location: z.string(),
|
||||
})
|
||||
.optional(),
|
||||
}),
|
||||
},
|
||||
TASK_RUN_COMPLETED_WITH_ACK: {
|
||||
message: z.object({
|
||||
version: z.enum(["v1", "v2"]).default("v2"),
|
||||
execution: V3ProdTaskRunExecution,
|
||||
completion: TaskRunExecutionResult,
|
||||
checkpoint: z
|
||||
.object({
|
||||
docker: z.boolean(),
|
||||
location: z.string(),
|
||||
})
|
||||
.optional(),
|
||||
}),
|
||||
callback: AckCallbackResult,
|
||||
},
|
||||
TASK_RUN_FAILED_TO_RUN: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
completion: TaskRunFailedExecutionResult,
|
||||
}),
|
||||
},
|
||||
TASK_HEARTBEAT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
},
|
||||
TASK_RUN_HEARTBEAT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
}),
|
||||
},
|
||||
CHECKPOINT_CREATED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string().optional(),
|
||||
attemptFriendlyId: z.string(),
|
||||
docker: z.boolean(),
|
||||
location: z.string(),
|
||||
reason: z.discriminatedUnion("type", [
|
||||
z.object({
|
||||
type: z.literal("WAIT_FOR_DURATION"),
|
||||
ms: z.number(),
|
||||
now: z.number(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("WAIT_FOR_BATCH"),
|
||||
batchFriendlyId: z.string(),
|
||||
runFriendlyIds: z.string().array(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("WAIT_FOR_TASK"),
|
||||
friendlyId: z.string(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("RETRYING_AFTER_FAILURE"),
|
||||
attemptNumber: z.number(),
|
||||
}),
|
||||
z.object({
|
||||
type: z.literal("MANUAL"),
|
||||
/** If unspecified it will be restored immediately, e.g. for live migration */
|
||||
restoreAtUnixTimeMs: z.number().optional(),
|
||||
}),
|
||||
]),
|
||||
}),
|
||||
callback: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
keepRunAlive: z.boolean(),
|
||||
}),
|
||||
},
|
||||
INDEXING_FAILED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
deploymentId: z.string(),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
stderr: z.string().optional(),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
RUN_CRASHED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const PlatformToCoordinatorMessages = {
|
||||
/** @deprecated use RESUME_AFTER_DEPENDENCY_WITH_ACK instead */
|
||||
RESUME_AFTER_DEPENDENCY: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
attemptId: z.string(),
|
||||
attemptFriendlyId: z.string(),
|
||||
completions: TaskRunExecutionResult.array(),
|
||||
executions: TaskRunExecution.array(),
|
||||
}),
|
||||
},
|
||||
RESUME_AFTER_DEPENDENCY_WITH_ACK: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
attemptId: z.string(),
|
||||
attemptFriendlyId: z.string(),
|
||||
completions: TaskRunExecutionResult.array(),
|
||||
executions: TaskRunExecution.array(),
|
||||
}),
|
||||
callback: AckCallbackResult,
|
||||
},
|
||||
RESUME_AFTER_DURATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptId: z.string(),
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
},
|
||||
REQUEST_ATTEMPT_CANCELLATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptId: z.string(),
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
},
|
||||
REQUEST_RUN_CANCELLATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
delayInMs: z.number().optional(),
|
||||
}),
|
||||
},
|
||||
READY_FOR_RETRY: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
}),
|
||||
},
|
||||
DYNAMIC_CONFIG: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
checkpointThresholdInMs: z.number(),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const ClientToSharedQueueMessages = {
|
||||
READY_FOR_TASKS: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
}),
|
||||
},
|
||||
BACKGROUND_WORKER_DEPRECATED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
}),
|
||||
},
|
||||
BACKGROUND_WORKER_MESSAGE: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
data: BackgroundWorkerClientMessages,
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const SharedQueueToClientMessages = {
|
||||
SERVER_READY: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
id: z.string(),
|
||||
}),
|
||||
},
|
||||
BACKGROUND_WORKER_MESSAGE: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
backgroundWorkerId: z.string(),
|
||||
data: BackgroundWorkerServerMessages,
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
const IndexTasksMessage = z.object({
|
||||
version: z.literal("v1"),
|
||||
deploymentId: z.string(),
|
||||
tasks: TaskResource.array(),
|
||||
packageVersion: z.string(),
|
||||
});
|
||||
|
||||
export const ProdWorkerToCoordinatorMessages = {
|
||||
TEST: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
}),
|
||||
callback: z.void(),
|
||||
},
|
||||
INDEX_TASKS: {
|
||||
message: z.discriminatedUnion("version", [
|
||||
IndexTasksMessage.extend({
|
||||
version: z.literal("v1"),
|
||||
}),
|
||||
IndexTasksMessage.extend({
|
||||
version: z.literal("v2"),
|
||||
supportsLazyAttempts: z.boolean(),
|
||||
}),
|
||||
]),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
}),
|
||||
]),
|
||||
},
|
||||
// Deprecated: Only workers without lazy attempt support will use this
|
||||
READY_FOR_EXECUTION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
totalCompletions: z.number(),
|
||||
}),
|
||||
},
|
||||
READY_FOR_LAZY_ATTEMPT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
totalCompletions: z.number(),
|
||||
startTime: z.number().optional(),
|
||||
}),
|
||||
},
|
||||
READY_FOR_RESUME: {
|
||||
message: z.discriminatedUnion("version", [
|
||||
z.object({
|
||||
version: z.literal("v1"),
|
||||
attemptFriendlyId: z.string(),
|
||||
type: WaitReason,
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v2"),
|
||||
attemptFriendlyId: z.string(),
|
||||
attemptNumber: z.number(),
|
||||
type: WaitReason,
|
||||
}),
|
||||
]),
|
||||
},
|
||||
READY_FOR_CHECKPOINT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
}),
|
||||
},
|
||||
CANCEL_CHECKPOINT: {
|
||||
message: z
|
||||
.discriminatedUnion("version", [
|
||||
z.object({
|
||||
version: z.literal("v1"),
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v2"),
|
||||
reason: WaitReason.optional(),
|
||||
}),
|
||||
])
|
||||
.default({ version: "v1" }),
|
||||
callback: z.object({
|
||||
version: z.literal("v2").default("v2"),
|
||||
checkpointCanceled: z.boolean(),
|
||||
reason: WaitReason.optional(),
|
||||
}),
|
||||
},
|
||||
TASK_HEARTBEAT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
},
|
||||
TASK_RUN_HEARTBEAT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
}),
|
||||
},
|
||||
TASK_RUN_COMPLETED: {
|
||||
message: z.object({
|
||||
version: z.enum(["v1", "v2"]).default("v1"),
|
||||
execution: V3ProdTaskRunExecution,
|
||||
completion: TaskRunExecutionResult,
|
||||
}),
|
||||
callback: z.object({
|
||||
willCheckpointAndRestore: z.boolean(),
|
||||
shouldExit: z.boolean(),
|
||||
}),
|
||||
},
|
||||
TASK_RUN_FAILED_TO_RUN: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
completion: TaskRunFailedExecutionResult,
|
||||
}),
|
||||
},
|
||||
WAIT_FOR_DURATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
ms: z.number(),
|
||||
now: z.number(),
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
callback: z.object({
|
||||
willCheckpointAndRestore: z.boolean(),
|
||||
}),
|
||||
},
|
||||
WAIT_FOR_TASK: {
|
||||
message: z.object({
|
||||
version: z.enum(["v1", "v2"]).default("v1"),
|
||||
friendlyId: z.string(),
|
||||
// This is the attempt that is waiting
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
callback: z.object({
|
||||
willCheckpointAndRestore: z.boolean(),
|
||||
}),
|
||||
},
|
||||
WAIT_FOR_BATCH: {
|
||||
message: z.object({
|
||||
version: z.enum(["v1", "v2"]).default("v1"),
|
||||
batchFriendlyId: z.string(),
|
||||
runFriendlyIds: z.string().array(),
|
||||
// This is the attempt that is waiting
|
||||
attemptFriendlyId: z.string(),
|
||||
}),
|
||||
callback: z.object({
|
||||
willCheckpointAndRestore: z.boolean(),
|
||||
}),
|
||||
},
|
||||
INDEXING_FAILED: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
deploymentId: z.string(),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
stderr: z.string().optional(),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
CREATE_TASK_RUN_ATTEMPT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
}),
|
||||
callback: z.discriminatedUnion("success", [
|
||||
z.object({
|
||||
success: z.literal(false),
|
||||
reason: z.string().optional(),
|
||||
}),
|
||||
z.object({
|
||||
success: z.literal(true),
|
||||
executionPayload: V3ProdTaskRunExecutionPayload,
|
||||
}),
|
||||
]),
|
||||
},
|
||||
UNRECOVERABLE_ERROR: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
error: z.object({
|
||||
name: z.string(),
|
||||
message: z.string(),
|
||||
stack: z.string().optional(),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
SET_STATE: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptFriendlyId: z.string().optional(),
|
||||
attemptNumber: z.string().optional(),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
// TODO: The coordinator can only safely use v1 worker messages, higher versions will need a new flag, e.g. SUPPORTS_VERSIONED_MESSAGES
|
||||
export const CoordinatorToProdWorkerMessages = {
|
||||
RESUME_AFTER_DEPENDENCY: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptId: z.string(),
|
||||
completions: TaskRunExecutionResult.array(),
|
||||
executions: TaskRunExecution.array(),
|
||||
}),
|
||||
},
|
||||
RESUME_AFTER_DURATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptId: z.string(),
|
||||
}),
|
||||
},
|
||||
// Deprecated: Only workers without lazy attempt support will use this
|
||||
EXECUTE_TASK_RUN: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
executionPayload: V3ProdTaskRunExecutionPayload,
|
||||
}),
|
||||
},
|
||||
EXECUTE_TASK_RUN_LAZY_ATTEMPT: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
lazyPayload: TaskRunExecutionLazyAttemptPayload,
|
||||
}),
|
||||
},
|
||||
REQUEST_ATTEMPT_CANCELLATION: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
attemptId: z.string(),
|
||||
}),
|
||||
},
|
||||
REQUEST_EXIT: {
|
||||
message: z.discriminatedUnion("version", [
|
||||
z.object({
|
||||
version: z.literal("v1"),
|
||||
}),
|
||||
z.object({
|
||||
version: z.literal("v2"),
|
||||
delayInMs: z.number().optional(),
|
||||
}),
|
||||
]),
|
||||
},
|
||||
READY_FOR_RETRY: {
|
||||
message: z.object({
|
||||
version: z.literal("v1").default("v1"),
|
||||
runId: z.string(),
|
||||
}),
|
||||
},
|
||||
};
|
||||
|
||||
export const ProdWorkerSocketData = z.object({
|
||||
contentHash: z.string(),
|
||||
projectRef: z.string(),
|
||||
envId: z.string(),
|
||||
runId: z.string(),
|
||||
attemptFriendlyId: z.string().optional(),
|
||||
attemptNumber: z.string().optional(),
|
||||
podName: z.string(),
|
||||
deploymentId: z.string(),
|
||||
deploymentVersion: z.string(),
|
||||
requiresCheckpointResumeWithMessage: z.string().optional(),
|
||||
});
|
||||
|
||||
export const CoordinatorSocketData = z.object({
|
||||
supportsDynamicConfig: z.string().optional(),
|
||||
});
|
||||
|
||||
@@ -1,203 +0,0 @@
|
||||
import type { DisconnectReason, Namespace, Server, Socket } from "socket.io";
|
||||
import { ZodMessageSender } from "./zodMessageHandler.js";
|
||||
import type {
|
||||
ZodMessageCatalogToSocketIoEvents,
|
||||
ZodSocketMessageCatalogSchema,
|
||||
ZodSocketMessageHandlers,
|
||||
} from "./zodSocket.js";
|
||||
import { ZodSocketMessageHandler } from "./zodSocket.js";
|
||||
// @ts-ignore
|
||||
import type { DefaultEventsMap, EventsMap } from "socket.io/dist/typed-events";
|
||||
import type { z } from "zod";
|
||||
import type { StructuredLogger } from "./utils/structuredLogger.js";
|
||||
import { SimpleStructuredLogger } from "./utils/structuredLogger.js";
|
||||
|
||||
interface ExtendedError extends Error {
|
||||
data?: any;
|
||||
}
|
||||
|
||||
export type ZodNamespaceSocket<
|
||||
TClientMessages extends ZodSocketMessageCatalogSchema,
|
||||
TServerMessages extends ZodSocketMessageCatalogSchema,
|
||||
TServerSideEvents extends EventsMap = DefaultEventsMap,
|
||||
TSocketData extends z.ZodObject<any, any, any> = any,
|
||||
> = Socket<
|
||||
ZodMessageCatalogToSocketIoEvents<TClientMessages>,
|
||||
ZodMessageCatalogToSocketIoEvents<TServerMessages>,
|
||||
TServerSideEvents,
|
||||
z.infer<TSocketData>
|
||||
>;
|
||||
|
||||
interface ZodNamespaceOptions<
|
||||
TClientMessages extends ZodSocketMessageCatalogSchema,
|
||||
TServerMessages extends ZodSocketMessageCatalogSchema,
|
||||
TServerSideEvents extends EventsMap = DefaultEventsMap,
|
||||
TSocketData extends z.ZodObject<any, any, any> = any,
|
||||
> {
|
||||
io: Server;
|
||||
name: string;
|
||||
clientMessages: TClientMessages;
|
||||
serverMessages: TServerMessages;
|
||||
socketData?: TSocketData;
|
||||
handlers?: ZodSocketMessageHandlers<TClientMessages>;
|
||||
authToken?: string;
|
||||
logger?: StructuredLogger;
|
||||
logHandlerPayloads?: boolean;
|
||||
preAuth?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
next: (err?: ExtendedError) => void,
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
postAuth?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
next: (err?: ExtendedError) => void,
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onConnection?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
handler: ZodSocketMessageHandler<TClientMessages>,
|
||||
sender: ZodMessageSender<TServerMessages>,
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onDisconnect?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
reason: DisconnectReason,
|
||||
description: any,
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
onError?: (
|
||||
socket: ZodNamespaceSocket<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>,
|
||||
err: Error,
|
||||
logger: StructuredLogger
|
||||
) => Promise<void>;
|
||||
}
|
||||
|
||||
export class ZodNamespace<
|
||||
TClientMessages extends ZodSocketMessageCatalogSchema,
|
||||
TServerMessages extends ZodSocketMessageCatalogSchema,
|
||||
TSocketData extends z.ZodObject<any, any, any> = any,
|
||||
TServerSideEvents extends EventsMap = DefaultEventsMap,
|
||||
> {
|
||||
#logger: StructuredLogger;
|
||||
#handler: ZodSocketMessageHandler<TClientMessages>;
|
||||
sender: ZodMessageSender<TServerMessages>;
|
||||
|
||||
io: Server;
|
||||
namespace: Namespace<
|
||||
ZodMessageCatalogToSocketIoEvents<TClientMessages>,
|
||||
ZodMessageCatalogToSocketIoEvents<TServerMessages>,
|
||||
TServerSideEvents,
|
||||
z.infer<TSocketData>
|
||||
>;
|
||||
|
||||
constructor(
|
||||
opts: ZodNamespaceOptions<TClientMessages, TServerMessages, TServerSideEvents, TSocketData>
|
||||
) {
|
||||
this.#logger =
|
||||
opts.logger ??
|
||||
new SimpleStructuredLogger(`ns-${opts.name}`, undefined, {
|
||||
namespace: opts.name,
|
||||
});
|
||||
|
||||
this.#handler = new ZodSocketMessageHandler({
|
||||
schema: opts.clientMessages,
|
||||
handlers: opts.handlers,
|
||||
logPayloads: opts.logHandlerPayloads,
|
||||
});
|
||||
|
||||
this.io = opts.io;
|
||||
|
||||
this.namespace = this.io.of(opts.name);
|
||||
|
||||
// FIXME: There's a bug here, this sender should not accept Socket schemas with callbacks
|
||||
this.sender = new ZodMessageSender({
|
||||
schema: opts.serverMessages,
|
||||
sender: async (message) => {
|
||||
return new Promise((resolve, reject) => {
|
||||
try {
|
||||
// @ts-expect-error
|
||||
this.namespace.emit(message.type, message.payload);
|
||||
resolve();
|
||||
} catch (err) {
|
||||
reject(err);
|
||||
}
|
||||
});
|
||||
},
|
||||
});
|
||||
|
||||
if (opts.preAuth) {
|
||||
this.namespace.use(async (socket, next) => {
|
||||
const logger = this.#logger.child({ socketId: socket.id, socketStage: "preAuth" });
|
||||
|
||||
if (typeof opts.preAuth === "function") {
|
||||
await opts.preAuth(socket, next, logger);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
if (opts.authToken) {
|
||||
this.namespace.use((socket, next) => {
|
||||
const logger = this.#logger.child({ socketId: socket.id, socketStage: "auth" });
|
||||
|
||||
const { auth } = socket.handshake;
|
||||
|
||||
if (!("token" in auth)) {
|
||||
logger.error("no token");
|
||||
return socket.disconnect(true);
|
||||
}
|
||||
|
||||
if (auth.token !== opts.authToken) {
|
||||
logger.error("invalid token");
|
||||
return socket.disconnect(true);
|
||||
}
|
||||
|
||||
logger.info("success");
|
||||
|
||||
next();
|
||||
|
||||
return;
|
||||
});
|
||||
}
|
||||
|
||||
if (opts.postAuth) {
|
||||
this.namespace.use(async (socket, next) => {
|
||||
const logger = this.#logger.child({ socketId: socket.id, socketStage: "auth" });
|
||||
|
||||
if (typeof opts.postAuth === "function") {
|
||||
await opts.postAuth(socket, next, logger);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
this.namespace.on("connection", async (socket) => {
|
||||
const logger = this.#logger.child({ socketId: socket.id, socketStage: "connection" });
|
||||
logger.info("connected");
|
||||
|
||||
this.#handler.registerHandlers(socket, logger);
|
||||
|
||||
socket.on("disconnect", async (reason, description) => {
|
||||
logger.info("disconnect", { reason, description });
|
||||
|
||||
if (opts.onDisconnect) {
|
||||
await opts.onDisconnect(socket, reason, description, logger);
|
||||
}
|
||||
});
|
||||
|
||||
socket.on("error", async (error) => {
|
||||
logger.error("error", { error });
|
||||
|
||||
if (opts.onError) {
|
||||
await opts.onError(socket, error, logger);
|
||||
}
|
||||
});
|
||||
|
||||
if (opts.onConnection) {
|
||||
await opts.onConnection(socket, this.#handler, this.sender, logger);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
fetchSockets() {
|
||||
return this.namespace.fetchSockets();
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,6 @@
|
||||
# Redis Worker
|
||||
|
||||
`@trigger.dev/redis-worker` - custom Redis-based background job system. **This replaces graphile-worker/zodworker** for all new background job needs.
|
||||
`@trigger.dev/redis-worker` - custom Redis-based background job system. This is the background job system for the webapp and run engine.
|
||||
|
||||
## Key Files
|
||||
|
||||
@@ -12,7 +12,7 @@
|
||||
|
||||
Used by the webapp for background jobs (alerting, batch processing, common tasks) and by the run engine for TTL expiration and batch operations.
|
||||
|
||||
All new background jobs in the webapp should use redis-worker. Do NOT add new jobs to zodworker (`@internal/zodworker`) or graphile-worker.
|
||||
All background jobs in the webapp use redis-worker.
|
||||
|
||||
## Testing
|
||||
|
||||
|
||||
Reference in New Issue
Block a user