v3: checkpoint failover and misc fixes (#1157)

* configurable checkpoint registry namespace

* add missing task create await

* remove unused messages

* changeset

* update self-hosting docs

* capture and display stderr for failed deploys

* add missing lockfile changes

* stderr changeset

* fix cli stderr message

* update error logs label
This commit is contained in:
nicktrn
2024-06-11 14:06:50 +01:00
committed by GitHub
parent 36ac79ac66
commit 68d32429b6
19 changed files with 217 additions and 105 deletions
+7
View File
@@ -0,0 +1,7 @@
---
"@trigger.dev/core-apps": patch
"trigger.dev": patch
"@trigger.dev/core": patch
---
Capture and display stderr on index failures
+7
View File
@@ -0,0 +1,7 @@
---
"@trigger.dev/core-apps": patch
"@trigger.dev/core": patch
---
- Fix uncaught provider exception
- Remove unused provider messages
+2 -1
View File
@@ -43,6 +43,7 @@ const SIMULATE_CHECKPOINT_FAILURE_SECONDS = parseInt(
);
const REGISTRY_HOST = process.env.REGISTRY_HOST || "localhost:5000";
const REGISTRY_NAMESPACE = process.env.REGISTRY_NAMESPACE || "trigger";
const CHECKPOINT_PATH = process.env.CHECKPOINT_PATH || "/checkpoints";
const REGISTRY_TLS_VERIFY = process.env.REGISTRY_TLS_VERIFY === "false" ? "false" : "true";
@@ -179,7 +180,7 @@ class Checkpointer {
}
#getImageRef(projectRef: string, deploymentVersion: string, shortCode: string) {
return `${REGISTRY_HOST}/trigger/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
return `${REGISTRY_HOST}/${REGISTRY_NAMESPACE}/${projectRef}:${deploymentVersion}.prod-${shortCode}`;
}
#getExportLocation(projectRef: string, deploymentVersion: string, shortCode: string) {
+17 -31
View File
@@ -85,37 +85,23 @@ class DockerTaskOperations implements TaskOperations {
port: COORDINATOR_PORT,
});
try {
logger.debug(
await execa("docker", [
"run",
"--network=host",
"--rm",
`--env=INDEX_TASKS=true`,
`--env=TRIGGER_SECRET_KEY=${opts.apiKey}`,
`--env=TRIGGER_API_URL=${opts.apiUrl}`,
`--env=TRIGGER_ENV_ID=${opts.envId}`,
`--env=OTEL_EXPORTER_OTLP_ENDPOINT=${OTEL_EXPORTER_OTLP_ENDPOINT}`,
`--env=POD_NAME=${containerName}`,
`--env=COORDINATOR_HOST=${COORDINATOR_HOST}`,
`--env=COORDINATOR_PORT=${COORDINATOR_PORT}`,
`--name=${containerName}`,
`${opts.imageRef}`,
])
);
} catch (error: any) {
if (!isExecaChildProcess(error)) {
throw error;
}
logger.error("Index failed:", {
opts,
exitCode: error.exitCode,
escapedCommand: error.escapedCommand,
stdout: error.stdout,
stderr: error.stderr,
});
}
logger.debug(
await execa("docker", [
"run",
"--network=host",
"--rm",
`--env=INDEX_TASKS=true`,
`--env=TRIGGER_SECRET_KEY=${opts.apiKey}`,
`--env=TRIGGER_API_URL=${opts.apiUrl}`,
`--env=TRIGGER_ENV_ID=${opts.envId}`,
`--env=OTEL_EXPORTER_OTLP_ENDPOINT=${OTEL_EXPORTER_OTLP_ENDPOINT}`,
`--env=POD_NAME=${containerName}`,
`--env=COORDINATOR_HOST=${COORDINATOR_HOST}`,
`--env=COORDINATOR_PORT=${COORDINATOR_PORT}`,
`--name=${containerName}`,
`${opts.imageRef}`,
])
);
}
async create(opts: TaskOperationsCreateOptions) {
@@ -20,6 +20,17 @@ export function DeploymentError({ errorData }: DeploymentErrorProps) {
maxLines={20}
/>
)}
{errorData.stderr && (
<>
<DeploymentErrorHeader title="Error logs:" />
<CodeBlock
showCopyButton={false}
showLineNumbers={false}
code={errorData.stderr}
maxLines={20}
/>
</>
)}
</div>
);
}
@@ -17,6 +17,7 @@ export type ErrorData = {
name: string;
message: string;
stack?: string;
stderr?: string;
};
export class DeploymentPresenter {
@@ -177,17 +178,20 @@ export class DeploymentPresenter {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stack: createTaskMetadataFailedErrorStack(parsedError.data),
stderr: parsedErrorData.data.stderr,
};
} else {
return {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stderr: parsedErrorData.data.stderr,
};
}
} else {
return {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stderr: parsedErrorData.data.stderr,
};
}
}
@@ -196,6 +200,7 @@ export class DeploymentPresenter {
name: parsedErrorData.data.name,
message: parsedErrorData.data.message,
stack: parsedErrorData.data.stack,
stderr: parsedErrorData.data.stderr,
};
}
}
@@ -1,14 +1,28 @@
import { PerformDeploymentAlertsService } from "./alerts/performDeploymentAlerts.server";
import { BaseService } from "./baseService.server";
import { logger } from "~/services/logger.server";
import { WorkerDeploymentStatus } from "@trigger.dev/database";
const FINAL_DEPLOYMENT_STATUSES: WorkerDeploymentStatus[] = [
"CANCELED",
"DEPLOYED",
"FAILED",
"TIMED_OUT",
];
export class DeploymentIndexFailed extends BaseService {
public async call(
maybeFriendlyId: string,
error: { name: string; message: string; stack?: string }
error: {
name: string;
message: string;
stack?: string;
stderr?: string;
}
) {
const isFriendlyId = maybeFriendlyId.startsWith("deployment_");
const deployment = await this._prisma.workerDeployment.update({
const deployment = await this._prisma.workerDeployment.findUnique({
where: isFriendlyId
? {
friendlyId: maybeFriendlyId,
@@ -16,6 +30,25 @@ export class DeploymentIndexFailed extends BaseService {
: {
id: maybeFriendlyId,
},
});
if (!deployment) {
logger.error("Worker deployment not found", { maybeFriendlyId });
return;
}
if (FINAL_DEPLOYMENT_STATUSES.includes(deployment.status)) {
logger.error("Worker deployment already in final state", {
id: deployment.id,
status: deployment.status,
});
return;
}
const failedDeployment = await this._prisma.workerDeployment.update({
where: {
id: deployment.id,
},
data: {
status: "FAILED",
failedAt: new Date(),
@@ -23,8 +56,8 @@ export class DeploymentIndexFailed extends BaseService {
},
});
await PerformDeploymentAlertsService.enqueue(deployment.id, this._prisma);
await PerformDeploymentAlertsService.enqueue(failedDeployment.id, this._prisma);
return deployment;
return failedDeployment;
}
}
+3 -1
View File
@@ -206,6 +206,8 @@ scp -3 root@<webapp_machine>:docker/.env root@<worker_machine>:docker/.env
Checkpointing allows you to save the state of a running container to disk and restore it later. This can be useful for
long-running tasks that need to be paused and resumed without losing state. Think fan-out and fan-in, or long waits in email campaigns.
The checkpoints will be pushed to the same registry as the deployed images. Please see the [Registry setup](#registry-setup) section for more information.
### Requirements
- Debian, **NOT** a derivative like Ubuntu
@@ -225,7 +227,7 @@ sudo apt-get install criu
2. Tweak the config so we can successfully checkpoint our workloads
```bash
mkdir /etc/criu
mkdir -p /etc/criu
cat << EOF >/etc/criu/runc.conf
tcp-close
+4
View File
@@ -507,6 +507,10 @@ async function _deployCommand(dir: string, options: DeployCommandOptions) {
await preExitTasks();
if (finishedDeployment.errorData.stderr) {
log.error(`stderr:\n${finishedDeployment.errorData.stderr}`);
}
throw new SkipLoggingError(
`Deployment encountered an error: ${finishedDeployment.errorData.name}`
);
@@ -78,6 +78,7 @@ export class ProdBackgroundWorker {
private _onClose: Evt<void> = new Evt();
public tasks: Array<TaskMetadataWithFilePath> = [];
public stderr: Array<string> = [];
_taskRunProcess: TaskRunProcess | undefined;
private _taskRunProcessesBeingKilled: Map<number, TaskRunProcess> = new Map();
@@ -161,6 +162,23 @@ export class ProdBackgroundWorker {
reject(new Error("Worker timed out"));
}, 10_000);
child.stdout?.on("data", (data) => {
console.log(data.toString());
});
child.stderr?.on("data", (data) => {
console.error(data.toString());
this.stderr.push(data.toString());
});
child.on("exit", (code) => {
if (!resolved) {
clearTimeout(timeout);
resolved = true;
reject(new Error(`Worker exited with code ${code}`));
}
});
new ZodIpcConnection({
listenSchema: ProdChildToWorkerMessages,
emitSchema: ProdWorkerToChildMessages,
@@ -192,22 +210,6 @@ export class ProdBackgroundWorker {
},
},
});
child.stdout?.on("data", (data) => {
console.log(data.toString());
});
child.stderr?.on("data", (data) => {
console.error(data.toString());
});
child.on("exit", (code) => {
if (!resolved) {
clearTimeout(timeout);
resolved = true;
reject(new Error(`Worker exited with code ${code}`));
}
});
});
this._initialized = true;
@@ -634,6 +634,8 @@ class ProdWorker {
process.exit(1);
}
} catch (e) {
const stderr = this.#backgroundWorker.stderr.join("\n");
if (e instanceof TaskMetadataParseError) {
logger.error("tasks metadata parse error", {
zodIssues: e.zodIssues,
@@ -647,6 +649,7 @@ class ProdWorker {
name: "TaskMetadataParseError",
message: "There was an error parsing the task metadata",
stack: JSON.stringify({ zodIssues: e.zodIssues, tasks: e.tasks }),
stderr,
},
});
} else if (e instanceof UncaughtExceptionError) {
@@ -654,6 +657,7 @@ class ProdWorker {
name: e.originalError.name,
message: e.originalError.message,
stack: e.originalError.stack,
stderr,
};
logger.error("uncaught exception", { originalError: error });
@@ -668,6 +672,7 @@ class ProdWorker {
name: e.name,
message: e.message,
stack: e.stack,
stderr,
};
logger.error("error", { error });
@@ -686,6 +691,7 @@ class ProdWorker {
error: {
name: "Error",
message: e,
stderr,
},
});
} else {
@@ -697,6 +703,7 @@ class ProdWorker {
error: {
name: "Error",
message: "Unknown error",
stderr,
},
});
}
+40 -17
View File
@@ -12,6 +12,8 @@ import { ZodMessageSender } from "@trigger.dev/core/v3/zodMessageHandler";
import { ZodSocketConnection } from "@trigger.dev/core/v3/zodSocket";
import { getRandomPortNumber, HttpReply, getTextBody } from "./http";
import { SimpleLogger } from "./logger";
import { isExecaChildProcess } from "./checkpoints";
import { setTimeout } from "node:timers/promises";
const HTTP_SERVER_PORT = Number(process.env.HTTP_SERVER_PORT || getRandomPortNumber());
const MACHINE_NAME = process.env.MACHINE_NAME || "local";
@@ -122,7 +124,7 @@ export class ProviderShell implements Provider {
BACKGROUND_WORKER_MESSAGE: async (message) => {
if (message.data.type === "SCHEDULE_ATTEMPT") {
try {
this.tasks.create({
await this.tasks.create({
image: message.data.image,
machine: message.data.machine,
version: message.data.version,
@@ -172,21 +174,6 @@ export class ProviderShell implements Provider {
"x-trigger-provider-type": this.options.type,
},
handlers: {
DELETE: async (message) => {
this.tasks.delete({ runId: message.name });
return {
message: "delete request received",
};
},
GET: async (message) => {
this.tasks.get({ runId: message.name });
},
HEALTH: async (message) => {
return {
status: "ok",
};
},
INDEX: async (message) => {
try {
await this.tasks.index({
@@ -202,7 +189,43 @@ export class ProviderShell implements Provider {
deploymentId: message.deploymentId,
});
} catch (error) {
logger.error("index failed", error);
if (isExecaChildProcess(error)) {
logger.error("Index failed", {
socketMessage: message,
exitCode: error.exitCode,
escapedCommand: error.escapedCommand,
stdout: error.stdout,
stderr: error.stderr,
});
if (error.exitCode === 111) {
logger.error("Index failure already reported by the worker", {
socketMessage: message,
});
// Add a brief delay to avoid messaging race conditions
await setTimeout(2000);
}
function normalizeStderr(stderr: string) {
return stderr
.split("\n")
.map((line) => line.trim())
.filter((line) => line.length > 0)
.join("\n");
}
return {
success: false,
error: {
name: "Index error",
message: `Crashed with exit code ${error.exitCode}`,
stderr: normalizeStderr(error.stderr),
},
};
} else {
logger.error("Index failed", error);
}
if (error instanceof Error) {
return {
+1
View File
@@ -162,6 +162,7 @@ export const DeploymentErrorData = z.object({
name: z.string(),
message: z.string(),
stack: z.string().optional(),
stderr: z.string().optional(),
});
export const GetDeploymentResponseBody = z.object({
+4 -24
View File
@@ -341,20 +341,13 @@ export const ProviderToPlatformMessages = {
name: z.string(),
message: z.string(),
stack: z.string().optional(),
stderr: z.string().optional(),
}),
}),
},
};
export const PlatformToProviderMessages = {
HEALTH: {
message: z.object({
version: z.literal("v1").default("v1"),
}),
callback: z.object({
status: z.literal("ok"),
}),
},
INDEX: {
message: z.object({
version: z.literal("v1").default("v1"),
@@ -376,6 +369,7 @@ export const PlatformToProviderMessages = {
name: z.string(),
message: z.string(),
stack: z.string().optional(),
stderr: z.string().optional(),
}),
}),
z.object({
@@ -383,7 +377,6 @@ export const PlatformToProviderMessages = {
}),
]),
},
// TODO: this should be a shared queue message instead
RESTORE: {
message: z.object({
version: z.literal("v1").default("v1"),
@@ -401,21 +394,6 @@ export const PlatformToProviderMessages = {
runId: z.string(),
}),
},
DELETE: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
callback: z.object({
message: z.string(),
}),
},
GET: {
message: z.object({
version: z.literal("v1").default("v1"),
name: z.string(),
}),
},
};
const CreateWorkerMessage = z.object({
@@ -582,6 +560,7 @@ export const CoordinatorToPlatformMessages = {
name: z.string(),
message: z.string(),
stack: z.string().optional(),
stderr: z.string().optional(),
}),
}),
},
@@ -825,6 +804,7 @@ export const ProdWorkerToCoordinatorMessages = {
name: z.string(),
message: z.string(),
stack: z.string().optional(),
stderr: z.string().optional(),
}),
}),
},
+24 -8
View File
@@ -3090,6 +3090,9 @@ importers:
'@sindresorhus/slugify':
specifier: ^2.2.1
version: 2.2.1
'@t3-oss/env-core':
specifier: ^0.10.1
version: 0.10.1(typescript@5.3.3)(zod@3.22.3)
'@traceloop/instrumentation-openai':
specifier: ^0.3.9
version: 0.3.9(@opentelemetry/api@1.4.1)
@@ -3135,6 +3138,9 @@ importers:
yt-dlp-wrap:
specifier: ^2.3.12
version: 2.3.12
zod:
specifier: 3.22.3
version: 3.22.3
devDependencies:
'@opentelemetry/core':
specifier: ^1.22.0
@@ -15291,6 +15297,19 @@ packages:
defer-to-connect: 2.0.1
dev: false
/@t3-oss/env-core@0.10.1(typescript@5.3.3)(zod@3.22.3):
resolution: {integrity: sha512-GcKZiCfWks5CTxhezn9k5zWX3sMDIYf6Kaxy2Gx9YEQftFcz8hDRN56hcbylyAO3t4jQnQ5ifLawINsNgCDpOg==}
peerDependencies:
typescript: '>=5.0.0'
zod: ^3.0.0
peerDependenciesMeta:
typescript:
optional: true
dependencies:
typescript: 5.3.3
zod: 3.22.3
dev: false
/@tabler/icons-react@2.40.0(react@18.2.0):
resolution: {integrity: sha512-C+dDOZowFbwI3LGQP0fdua+hOPkGkW7XeMcRXTSdEKc5fD75W6zRO5nXnWivIMRKsi/Y26EDmnQo15N8JX378w==}
peerDependencies:
@@ -15376,7 +15395,7 @@ packages:
'@graphql-typed-document-node/core': 3.2.0(graphql@16.6.0)
axios: 1.4.0
graphql: 16.6.0
zod: 3.22.4
zod: 3.22.3
transitivePeerDependencies:
- debug
dev: false
@@ -15386,7 +15405,7 @@ packages:
dependencies:
'@graphql-typed-document-node/core': 3.2.0(graphql@16.6.0)
graphql: 16.6.0
zod: 3.22.4
zod: 3.22.3
dev: false
/@testing-library/dom@8.19.1:
@@ -27074,7 +27093,7 @@ packages:
workerd: 1.20231030.0
ws: 8.16.0
youch: 3.3.3
zod: 3.22.4
zod: 3.22.3
transitivePeerDependencies:
- bufferutil
- supports-color
@@ -27097,7 +27116,7 @@ packages:
workerd: 1.20231030.0
ws: 8.16.0
youch: 3.3.3
zod: 3.22.4
zod: 3.22.3
transitivePeerDependencies:
- bufferutil
- supports-color
@@ -36355,7 +36374,7 @@ packages:
/zod-error@1.5.0:
resolution: {integrity: sha512-zzopKZ/skI9iXpqCEPj+iLCKl9b88E43ehcU+sbRoHuwGd9F1IDVGQ70TyO6kmfiRL1g4IXkjsXK+g1gLYl4WQ==}
dependencies:
zod: 3.22.4
zod: 3.22.3
dev: false
/zod-validation-error@1.5.0(zod@3.22.3):
@@ -36377,8 +36396,5 @@ packages:
/zod@3.22.3:
resolution: {integrity: sha512-EjIevzuJRiRPbVH4mGc8nApb/lVLKVpmUhAaR5R5doKGfAnGJ6Gr3CViAVjP+4FWSxCsybeWQdcgCtbX+7oZug==}
/zod@3.22.4:
resolution: {integrity: sha512-iC+8Io04lddc+mVqQ9AZ7OQ2MrUKGN+oIQyq1vemgt46jwCwLfhq7/pwnBnNXXXZb8VTVLKwp9EDkx+ryxIWmg==}
/zwitch@2.0.4:
resolution: {integrity: sha512-bXE4cR/kVZhKZX/RjPEflHaKVhUVl85noU3v6b8apfQEc1x4A+zBxjZ4lN8LqGd6WZ3dl98pY4o717VFmoPp+A==}
+3 -1
View File
@@ -15,6 +15,7 @@
"@react-email/components": "^0.0.17",
"@react-email/render": "^0.0.7",
"@sindresorhus/slugify": "^2.2.1",
"@t3-oss/env-core": "^0.10.1",
"@traceloop/instrumentation-openai": "^0.3.9",
"@trigger.dev/core": "workspace:^3.0.0-beta.0",
"@trigger.dev/sdk": "workspace:^3.0.0-beta.0",
@@ -29,7 +30,8 @@
"server-only": "^0.0.1",
"stripe": "^12.14.0",
"typeorm": "^0.3.20",
"yt-dlp-wrap": "^2.3.12"
"yt-dlp-wrap": "^2.3.12",
"zod": "3.22.3"
},
"devDependencies": {
"@opentelemetry/api": "^1.8.0",
+24
View File
@@ -0,0 +1,24 @@
import { createEnv } from "@t3-oss/env-core";
import { logger, task } from "@trigger.dev/sdk/v3";
import { z } from "zod";
// uncomment to trigger deploy failures
export const env = createEnv({
clientPrefix: "NEXT_PUBLIC_",
server: {
SECRET_DATABASE_URL: z.string().url(),
},
client: {
// NEXT_PUBLIC_SOME_PUBKEY: z.string().min(1),
},
runtimeEnv: {
// NEXT_PUBLIC_SOME_PUBKEY: process.env.NEXT_PUBLIC_SOME_PUBKEY,
},
});
export const simplestTask = task({
id: "t3-env-test",
run: async (payload: any) => {
logger.info("Environment variables", env);
},
});
+1 -1
View File
@@ -46,7 +46,7 @@ export const config: TriggerConfig = {
},
additionalPackages: ["wrangler@3.35.0", "pg@8.11.5"],
additionalFiles: ["./wrangler/wrangler.toml"],
dependenciesToBundle: [/@sindresorhus/, "escape-string-regexp"],
dependenciesToBundle: [/@sindresorhus/, "escape-string-regexp", "@t3-oss/env-core"],
instrumentations: [new OpenAIInstrumentation()],
logLevel: "info",
onStart: async (payload, { ctx }) => {
+2 -1
View File
@@ -13,6 +13,7 @@
"@trigger.dev/sdk/v3/*": ["../../packages/trigger-sdk/src/v3/*"]
},
"emitDecoratorMetadata": true,
"experimentalDecorators": true
"experimentalDecorators": true,
"moduleResolution": "node16"
}
}