use tryCatch in TaskExecutor
This commit is contained in:
@@ -37,6 +37,7 @@ import {
|
||||
stringifyIO,
|
||||
} from "../utils/ioSerialization.js";
|
||||
import { calculateNextRetryDelay } from "../utils/retries.js";
|
||||
import { tryCatch } from "../tryCatch.js";
|
||||
|
||||
export type TaskExecutorOptions = {
|
||||
tracingSDK: TracingSDK;
|
||||
@@ -108,14 +109,15 @@ export class TaskExecutor {
|
||||
let parsedPayload: any;
|
||||
let initOutput: any;
|
||||
|
||||
try {
|
||||
await runTimelineMetrics.measureMetric("trigger.dev/execution", "payload", async () => {
|
||||
const [inputError, payloadResult] = await tryCatch(
|
||||
runTimelineMetrics.measureMetric("trigger.dev/execution", "payload", async () => {
|
||||
const payloadPacket = await conditionallyImportPacket(originalPacket, this._tracer);
|
||||
parsedPayload = await parsePacket(payloadPacket);
|
||||
});
|
||||
} catch (inputError) {
|
||||
recordSpanException(span, inputError);
|
||||
return await parsePacket(payloadPacket);
|
||||
})
|
||||
);
|
||||
|
||||
if (inputError) {
|
||||
recordSpanException(span, inputError);
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.TASK_INPUT_ERROR,
|
||||
@@ -123,121 +125,123 @@ export class TaskExecutor {
|
||||
);
|
||||
}
|
||||
|
||||
parsedPayload = await this.#parsePayload(parsedPayload);
|
||||
parsedPayload = await this.#parsePayload(payloadResult);
|
||||
|
||||
const executeTask = async (payload: any) => {
|
||||
try {
|
||||
initOutput = await this.#callInitFunctions(payload, ctx, signal);
|
||||
const [runError, output] = await tryCatch(
|
||||
(async () => {
|
||||
initOutput = await this.#callInitFunctions(payload, ctx, signal);
|
||||
|
||||
if (execution.attempt.number === 1) {
|
||||
await this.#callOnStartFunctions(payload, ctx, initOutput, signal);
|
||||
}
|
||||
|
||||
const output = await this.#callRun(payload, ctx, initOutput, signal);
|
||||
|
||||
try {
|
||||
const stringifiedOutput = await stringifyIO(output);
|
||||
|
||||
const finalOutput = await conditionallyExportPacket(
|
||||
stringifiedOutput,
|
||||
`${execution.attempt.id}/output`,
|
||||
this._tracer
|
||||
);
|
||||
|
||||
const attributes = await createPacketAttributes(
|
||||
finalOutput,
|
||||
SemanticInternalAttributes.OUTPUT,
|
||||
SemanticInternalAttributes.OUTPUT_TYPE
|
||||
);
|
||||
|
||||
if (attributes) {
|
||||
span.setAttributes(attributes);
|
||||
if (execution.attempt.number === 1) {
|
||||
await this.#callOnStartFunctions(payload, ctx, initOutput, signal);
|
||||
}
|
||||
|
||||
await this.#callOnSuccessFunctions(payload, output, ctx, initOutput, signal);
|
||||
return await this.#callRun(payload, ctx, initOutput, signal);
|
||||
})()
|
||||
);
|
||||
|
||||
// Call onComplete with success result
|
||||
await this.#callOnCompleteFunctions(
|
||||
payload,
|
||||
{ ok: true, data: output },
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
if (runError) {
|
||||
const [handleErrorError, handleErrorResult] = await tryCatch(
|
||||
this.#handleError(execution, runError, payload, ctx, initOutput, signal)
|
||||
);
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
id: execution.run.id,
|
||||
output: finalOutput.data,
|
||||
outputType: finalOutput.dataType,
|
||||
} satisfies TaskRunExecutionResult;
|
||||
} catch (outputError) {
|
||||
recordSpanException(span, outputError);
|
||||
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.TASK_OUTPUT_ERROR,
|
||||
outputError
|
||||
);
|
||||
}
|
||||
} catch (runError) {
|
||||
try {
|
||||
const handleErrorResult = await this.#handleError(
|
||||
execution,
|
||||
runError,
|
||||
payload,
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
|
||||
recordSpanException(span, handleErrorResult.error ?? runError);
|
||||
|
||||
if (handleErrorResult.status !== "retry") {
|
||||
await this.#callOnFailureFunctions(
|
||||
payload,
|
||||
handleErrorResult.error ?? runError,
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
|
||||
// Call onComplete with error result
|
||||
await this.#callOnCompleteFunctions(
|
||||
payload,
|
||||
{ ok: false, error: handleErrorResult.error ?? runError },
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
id: execution.run.id,
|
||||
ok: false,
|
||||
error: sanitizeError(
|
||||
handleErrorResult.error
|
||||
? parseError(handleErrorResult.error)
|
||||
: parseError(runError)
|
||||
),
|
||||
retry: handleErrorResult.status === "retry" ? handleErrorResult.retry : undefined,
|
||||
skippedRetrying: handleErrorResult.status === "skipped",
|
||||
} satisfies TaskRunExecutionResult;
|
||||
} catch (handleErrorError) {
|
||||
if (handleErrorError) {
|
||||
recordSpanException(span, handleErrorError);
|
||||
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.HANDLE_ERROR_ERROR,
|
||||
handleErrorError
|
||||
);
|
||||
}
|
||||
} finally {
|
||||
await this.#callTaskCleanup(payload, ctx, initOutput, signal);
|
||||
await this.#blockForWaitUntil();
|
||||
|
||||
span.setAttributes(runTimelineMetrics.convertMetricsToSpanAttributes());
|
||||
recordSpanException(span, handleErrorResult.error ?? runError);
|
||||
|
||||
if (handleErrorResult.status !== "retry") {
|
||||
await this.#callOnFailureFunctions(
|
||||
payload,
|
||||
handleErrorResult.error ?? runError,
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
|
||||
await this.#callOnCompleteFunctions(
|
||||
payload,
|
||||
{ ok: false, error: handleErrorResult.error ?? runError },
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
id: execution.run.id,
|
||||
ok: false,
|
||||
error: sanitizeError(
|
||||
handleErrorResult.error
|
||||
? parseError(handleErrorResult.error)
|
||||
: parseError(runError)
|
||||
),
|
||||
retry: handleErrorResult.status === "retry" ? handleErrorResult.retry : undefined,
|
||||
skippedRetrying: handleErrorResult.status === "skipped",
|
||||
} satisfies TaskRunExecutionResult;
|
||||
}
|
||||
|
||||
const [outputError, stringifiedOutput] = await tryCatch(stringifyIO(output));
|
||||
|
||||
if (outputError) {
|
||||
recordSpanException(span, outputError);
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.TASK_OUTPUT_ERROR,
|
||||
outputError
|
||||
);
|
||||
}
|
||||
|
||||
const [exportError, finalOutput] = await tryCatch(
|
||||
conditionallyExportPacket(
|
||||
stringifiedOutput,
|
||||
`${execution.attempt.id}/output`,
|
||||
this._tracer
|
||||
)
|
||||
);
|
||||
|
||||
if (exportError) {
|
||||
recordSpanException(span, exportError);
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.TASK_OUTPUT_ERROR,
|
||||
exportError
|
||||
);
|
||||
}
|
||||
|
||||
const [attrError, attributes] = await tryCatch(
|
||||
createPacketAttributes(
|
||||
finalOutput,
|
||||
SemanticInternalAttributes.OUTPUT,
|
||||
SemanticInternalAttributes.OUTPUT_TYPE
|
||||
)
|
||||
);
|
||||
|
||||
if (!attrError && attributes) {
|
||||
span.setAttributes(attributes);
|
||||
}
|
||||
|
||||
await this.#callOnSuccessFunctions(payload, output, ctx, initOutput, signal);
|
||||
await this.#callOnCompleteFunctions(
|
||||
payload,
|
||||
{ ok: true, data: output },
|
||||
ctx,
|
||||
initOutput,
|
||||
signal
|
||||
);
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
id: execution.run.id,
|
||||
output: finalOutput.data,
|
||||
outputType: finalOutput.dataType,
|
||||
} satisfies TaskRunExecutionResult;
|
||||
};
|
||||
|
||||
const globalMiddlewareHooks = lifecycleHooks.getGlobalMiddlewareHooks();
|
||||
@@ -297,18 +301,22 @@ export class TaskExecutor {
|
||||
};
|
||||
},
|
||||
async () => {
|
||||
try {
|
||||
output = await executeTask(payload);
|
||||
} catch (error) {
|
||||
const [error, result] = await tryCatch(executeTask(payload));
|
||||
if (error) {
|
||||
executeError = error;
|
||||
} else {
|
||||
output = result;
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
try {
|
||||
await runner();
|
||||
} catch (error) {
|
||||
return this.#internalErrorResult(execution, TaskRunErrorCodes.TASK_MIDDLEWARE_ERROR, error);
|
||||
const [runnerError] = await tryCatch(runner());
|
||||
if (runnerError) {
|
||||
return this.#internalErrorResult(
|
||||
execution,
|
||||
TaskRunErrorCodes.TASK_MIDDLEWARE_ERROR,
|
||||
runnerError
|
||||
);
|
||||
}
|
||||
|
||||
if (executeError) {
|
||||
@@ -348,27 +356,28 @@ export class TaskExecutor {
|
||||
// Store global hook results in an array
|
||||
const globalResults = [];
|
||||
for (const hook of globalInitHooks) {
|
||||
const result = await this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
const result = await hook.fn({ payload, ctx, signal, task: this.task.id });
|
||||
const [hookError, result] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
const result = await hook.fn({ payload, ctx, signal, task: this.task.id });
|
||||
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
span.setAttributes(flattenAttributes(result));
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
span.setAttributes(flattenAttributes(result));
|
||||
return result;
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
return {};
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
return {};
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
// Only include object results
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
|
||||
if (!hookError && result && typeof result === "object" && !Array.isArray(result)) {
|
||||
globalResults.push(result);
|
||||
}
|
||||
}
|
||||
@@ -377,28 +386,34 @@ export class TaskExecutor {
|
||||
const mergedGlobalResults = Object.assign({}, ...globalResults);
|
||||
|
||||
if (taskInitHook) {
|
||||
const taskResult = await this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
const result = await taskInitHook({ payload, ctx, signal, task: this.task.id });
|
||||
const [hookError, taskResult] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
const result = await taskInitHook({ payload, ctx, signal, task: this.task.id });
|
||||
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
span.setAttributes(flattenAttributes(result));
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
span.setAttributes(flattenAttributes(result));
|
||||
return result;
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
return {};
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
return {};
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
|
||||
// Only merge if taskResult is an object
|
||||
if (taskResult && typeof taskResult === "object" && !Array.isArray(taskResult)) {
|
||||
if (
|
||||
!hookError &&
|
||||
taskResult &&
|
||||
typeof taskResult === "object" &&
|
||||
!Array.isArray(taskResult)
|
||||
) {
|
||||
return { ...mergedGlobalResults, ...taskResult };
|
||||
}
|
||||
|
||||
@@ -412,7 +427,6 @@ export class TaskExecutor {
|
||||
|
||||
if (result && typeof result === "object" && !Array.isArray(result)) {
|
||||
span.setAttributes(flattenAttributes(result));
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -449,45 +463,51 @@ export class TaskExecutor {
|
||||
"success",
|
||||
async () => {
|
||||
for (const hook of globalSuccessHooks) {
|
||||
await this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({
|
||||
payload,
|
||||
output,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({
|
||||
payload,
|
||||
output,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onSuccess functions
|
||||
}
|
||||
|
||||
if (taskSuccessHook) {
|
||||
await this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskSuccessHook({
|
||||
payload,
|
||||
output,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskSuccessHook({
|
||||
payload,
|
||||
output,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onSuccess functions
|
||||
}
|
||||
}
|
||||
);
|
||||
@@ -522,8 +542,8 @@ export class TaskExecutor {
|
||||
"failure",
|
||||
async () => {
|
||||
for (const hook of globalFailureHooks) {
|
||||
try {
|
||||
await this._tracer.startActiveSpan(
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({
|
||||
@@ -540,15 +560,14 @@ export class TaskExecutor {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
);
|
||||
} catch {
|
||||
// Ignore errors from onFailure functions
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onFailure functions
|
||||
}
|
||||
|
||||
if (taskFailureHook) {
|
||||
try {
|
||||
await this._tracer.startActiveSpan(
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskFailureHook({
|
||||
@@ -565,10 +584,9 @@ export class TaskExecutor {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
);
|
||||
} catch {
|
||||
// Ignore errors from onFailure functions
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onFailure functions
|
||||
}
|
||||
}
|
||||
);
|
||||
@@ -586,11 +604,11 @@ export class TaskExecutor {
|
||||
return payload;
|
||||
}
|
||||
|
||||
try {
|
||||
return await this.task.fns.parsePayload(payload);
|
||||
} catch (e) {
|
||||
throw new TaskPayloadParsedError(e);
|
||||
const [parseError, result] = await tryCatch(this.task.fns.parsePayload(payload));
|
||||
if (parseError) {
|
||||
throw new TaskPayloadParsedError(parseError);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
async #callOnStartFunctions(
|
||||
@@ -614,36 +632,40 @@ export class TaskExecutor {
|
||||
"start",
|
||||
async () => {
|
||||
for (const hook of globalStartHooks) {
|
||||
await this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({ payload, ctx, signal, task: this.task.id, init: initOutput });
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({ payload, ctx, signal, task: this.task.id, init: initOutput });
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
}
|
||||
|
||||
if (taskStartHook) {
|
||||
await this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskStartHook({
|
||||
payload,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskStartHook({
|
||||
payload,
|
||||
ctx,
|
||||
signal,
|
||||
task: this.task.id,
|
||||
init: initOutput,
|
||||
});
|
||||
},
|
||||
}
|
||||
{
|
||||
attributes: {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
)
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -880,8 +902,8 @@ export class TaskExecutor {
|
||||
"complete",
|
||||
async () => {
|
||||
for (const hook of globalCompleteHooks) {
|
||||
try {
|
||||
await this._tracer.startActiveSpan(
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
hook.name ?? "global",
|
||||
async (span) => {
|
||||
await hook.fn({
|
||||
@@ -898,15 +920,14 @@ export class TaskExecutor {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
);
|
||||
} catch {
|
||||
// Ignore errors from onComplete functions
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onComplete functions
|
||||
}
|
||||
|
||||
if (taskCompleteHook) {
|
||||
try {
|
||||
await this._tracer.startActiveSpan(
|
||||
const [hookError] = await tryCatch(
|
||||
this._tracer.startActiveSpan(
|
||||
"task",
|
||||
async (span) => {
|
||||
await taskCompleteHook({
|
||||
@@ -923,10 +944,9 @@ export class TaskExecutor {
|
||||
[SemanticInternalAttributes.STYLE_ICON]: "tabler-function",
|
||||
},
|
||||
}
|
||||
);
|
||||
} catch {
|
||||
// Ignore errors from onComplete functions
|
||||
}
|
||||
)
|
||||
);
|
||||
// Ignore errors from onComplete functions
|
||||
}
|
||||
}
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user