onWait and onResume

This commit is contained in:
Eric Allam
2025-03-22 15:47:59 +00:00
parent e816ba4382
commit aaf2ed8a20
4 changed files with 177 additions and 1 deletions
@@ -7,6 +7,8 @@ import {
AnyOnFailureHookFunction,
AnyOnSuccessHookFunction,
AnyOnCompleteHookFunction,
AnyOnWaitHookFunction,
AnyOnResumeHookFunction,
} from "./types.js";
export class StandardLifecycleHooksManager implements LifecycleHooksManager {
@@ -31,6 +33,13 @@ export class StandardLifecycleHooksManager implements LifecycleHooksManager {
private taskCompleteHooks: Map<string, RegisteredHookFunction<AnyOnCompleteHookFunction>> =
new Map();
private globalWaitHooks: Map<string, RegisteredHookFunction<AnyOnWaitHookFunction>> = new Map();
private taskWaitHooks: Map<string, RegisteredHookFunction<AnyOnWaitHookFunction>> = new Map();
private globalResumeHooks: Map<string, RegisteredHookFunction<AnyOnResumeHookFunction>> =
new Map();
private taskResumeHooks: Map<string, RegisteredHookFunction<AnyOnResumeHookFunction>> = new Map();
registerGlobalStartHook(hook: RegisterHookFunctionParams<AnyOnStartHookFunction>): void {
const id = generateHookId(hook);
@@ -188,6 +197,68 @@ export class StandardLifecycleHooksManager implements LifecycleHooksManager {
getGlobalCompleteHooks(): RegisteredHookFunction<AnyOnCompleteHookFunction>[] {
return Array.from(this.globalCompleteHooks.values());
}
registerGlobalWaitHook(hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>): void {
const id = generateHookId(hook);
this.globalWaitHooks.set(id, {
id,
name: hook.id ?? hook.fn.name ? (hook.fn.name === "" ? undefined : hook.fn.name) : undefined,
fn: hook.fn,
});
}
registerTaskWaitHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>
): void {
const id = generateHookId(hook);
this.taskWaitHooks.set(taskId, {
id,
name: hook.id ?? hook.fn.name ? (hook.fn.name === "" ? undefined : hook.fn.name) : undefined,
fn: hook.fn,
});
}
getTaskWaitHook(taskId: string): AnyOnWaitHookFunction | undefined {
return this.taskWaitHooks.get(taskId)?.fn;
}
getGlobalWaitHooks(): RegisteredHookFunction<AnyOnWaitHookFunction>[] {
return Array.from(this.globalWaitHooks.values());
}
registerGlobalResumeHook(hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>): void {
const id = generateHookId(hook);
this.globalResumeHooks.set(id, {
id,
name: hook.id ?? hook.fn.name ? (hook.fn.name === "" ? undefined : hook.fn.name) : undefined,
fn: hook.fn,
});
}
registerTaskResumeHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>
): void {
const id = generateHookId(hook);
this.taskResumeHooks.set(taskId, {
id,
name: hook.id ?? hook.fn.name ? (hook.fn.name === "" ? undefined : hook.fn.name) : undefined,
fn: hook.fn,
});
}
getTaskResumeHook(taskId: string): AnyOnResumeHookFunction | undefined {
return this.taskResumeHooks.get(taskId)?.fn;
}
getGlobalResumeHooks(): RegisteredHookFunction<AnyOnResumeHookFunction>[] {
return Array.from(this.globalResumeHooks.values());
}
}
export class NoopLifecycleHooksManager implements LifecycleHooksManager {
@@ -285,6 +356,44 @@ export class NoopLifecycleHooksManager implements LifecycleHooksManager {
getGlobalCompleteHooks(): RegisteredHookFunction<AnyOnCompleteHookFunction>[] {
return [];
}
registerGlobalWaitHook(hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>): void {
// Noop
}
registerTaskWaitHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>
): void {
// Noop
}
getTaskWaitHook(taskId: string): AnyOnWaitHookFunction | undefined {
return undefined;
}
getGlobalWaitHooks(): RegisteredHookFunction<AnyOnWaitHookFunction>[] {
return [];
}
registerGlobalResumeHook(hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>): void {
// Noop
}
registerTaskResumeHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>
): void {
// Noop
}
getTaskResumeHook(taskId: string): AnyOnResumeHookFunction | undefined {
return undefined;
}
getGlobalResumeHooks(): RegisteredHookFunction<AnyOnResumeHookFunction>[] {
return [];
}
}
function generateHookId(hook: RegisterHookFunctionParams<any>): string {
@@ -26,6 +26,32 @@ export type OnStartHookFunction<TPayload> = (
export type AnyOnStartHookFunction = OnStartHookFunction<unknown>;
export type TaskWaitHookParams<TPayload = unknown> = {
ctx: TaskRunContext;
payload: TPayload;
task: string;
signal?: AbortSignal;
};
export type OnWaitHookFunction<TPayload> = (
params: TaskWaitHookParams<TPayload>
) => undefined | void | Promise<undefined | void>;
export type AnyOnWaitHookFunction = OnWaitHookFunction<unknown>;
export type TaskResumeHookParams<TPayload = unknown> = {
ctx: TaskRunContext;
payload: TPayload;
task: string;
signal?: AbortSignal;
};
export type OnResumeHookFunction<TPayload> = (
params: TaskResumeHookParams<TPayload>
) => undefined | void | Promise<undefined | void>;
export type AnyOnResumeHookFunction = OnResumeHookFunction<unknown>;
export type TaskFailureHookParams<TPayload = unknown> = {
ctx: TaskRunContext;
payload: TPayload;
@@ -129,4 +155,18 @@ export interface LifecycleHooksManager {
): void;
getTaskCompleteHook(taskId: string): AnyOnCompleteHookFunction | undefined;
getGlobalCompleteHooks(): RegisteredHookFunction<AnyOnCompleteHookFunction>[];
registerGlobalWaitHook(hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>): void;
registerTaskWaitHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnWaitHookFunction>
): void;
getTaskWaitHook(taskId: string): AnyOnWaitHookFunction | undefined;
getGlobalWaitHooks(): RegisteredHookFunction<AnyOnWaitHookFunction>[];
registerGlobalResumeHook(hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>): void;
registerTaskResumeHook(
taskId: string,
hook: RegisterHookFunctionParams<AnyOnResumeHookFunction>
): void;
getTaskResumeHook(taskId: string): AnyOnResumeHookFunction | undefined;
getGlobalResumeHooks(): RegisteredHookFunction<AnyOnResumeHookFunction>[];
}
+25
View File
@@ -10,6 +10,8 @@ import {
type AnyOnSuccessHookFunction,
type AnyOnCompleteHookFunction,
type TaskCompleteResult,
type AnyOnWaitHookFunction,
type AnyOnResumeHookFunction,
} from "@trigger.dev/core/v3";
export type {
@@ -23,6 +25,8 @@ export type {
AnyOnSuccessHookFunction,
AnyOnCompleteHookFunction,
TaskCompleteResult,
AnyOnWaitHookFunction,
AnyOnResumeHookFunction,
};
export function onInit(name: string, fn: AnyOnInitHookFunction): void;
@@ -81,3 +85,24 @@ export function onComplete(
fn: typeof fnOrName === "function" ? fnOrName : fn!,
});
}
export function onWait(name: string, fn: AnyOnWaitHookFunction): void;
export function onWait(fn: AnyOnWaitHookFunction): void;
export function onWait(fnOrName: string | AnyOnWaitHookFunction, fn?: AnyOnWaitHookFunction): void {
lifecycleHooks.registerGlobalWaitHook({
id: typeof fnOrName === "string" ? fnOrName : fnOrName.name ? fnOrName.name : undefined,
fn: typeof fnOrName === "function" ? fnOrName : fn!,
});
}
export function onResume(name: string, fn: AnyOnResumeHookFunction): void;
export function onResume(fn: AnyOnResumeHookFunction): void;
export function onResume(
fnOrName: string | AnyOnResumeHookFunction,
fn?: AnyOnResumeHookFunction
): void {
lifecycleHooks.registerGlobalResumeHook({
id: typeof fnOrName === "string" ? fnOrName : fnOrName.name ? fnOrName.name : undefined,
fn: typeof fnOrName === "function" ? fnOrName : fn!,
});
}
+3 -1
View File
@@ -1,4 +1,4 @@
import { onInit, onStart, onFailure, onSuccess, onComplete } from "./hooks.js";
import { onInit, onStart, onFailure, onSuccess, onComplete, onWait, onResume } from "./hooks.js";
import {
batchTrigger,
batchTriggerAndWait,
@@ -84,4 +84,6 @@ export const tasks = {
onFailure,
onSuccess,
onComplete,
onWait,
onResume,
};