Files

53 lines
3.2 KiB
TypeScript

import { defineSignal, proxyActivities, setHandler } from "@temporalio/workflow";
type Activities = {
markRunning(input: { runId: string }): Promise<void>;
prepare(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
agent(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
trace(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
diff(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
build(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
collect(input: { runId: string; identity: string }): Promise<Record<string, unknown>>;
markCompleted(input: { runId: string }): Promise<void>;
markCanceled(input: { runId: string }): Promise<void>;
markFailed(input: { runId: string; code: string; message: string }): Promise<void>;
};
export const cancelCaseRun = defineSignal("cancelCaseRun");
export async function caseRunWorkflow(input: { runId: string; activity: { startToCloseTimeoutMs: number; heartbeatTimeoutMs: number; retryInitialIntervalMs: number; retryMaximumIntervalMs: number; retryMaximumAttempts: number } }) {
const activities = proxyActivities<Activities>({
startToCloseTimeout: `${input.activity.startToCloseTimeoutMs} milliseconds`,
heartbeatTimeout: `${input.activity.heartbeatTimeoutMs} milliseconds`,
retry: { initialInterval: `${input.activity.retryInitialIntervalMs} milliseconds`, maximumInterval: `${input.activity.retryMaximumIntervalMs} milliseconds`, maximumAttempts: input.activity.retryMaximumAttempts }
});
let cancelRequested = false;
setHandler(cancelCaseRun, () => { cancelRequested = true; });
const canceled = async () => {
if (!cancelRequested) return false;
await activities.markCanceled({ runId: input.runId });
return true;
};
try {
await activities.markRunning(input);
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.prepare({ runId: input.runId, identity: `${input.runId}:prepare:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.agent({ runId: input.runId, identity: `${input.runId}:agent:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.trace({ runId: input.runId, identity: `${input.runId}:trace:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.diff({ runId: input.runId, identity: `${input.runId}:diff:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.build({ runId: input.runId, identity: `${input.runId}:build:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.collect({ runId: input.runId, identity: `${input.runId}:collect:v1` });
if (await canceled()) return { status: "canceled", runId: input.runId };
await activities.markCompleted({ runId: input.runId });
return { status: "completed", runId: input.runId };
} catch (error: any) {
await activities.markFailed({ runId: input.runId, code: error?.code ?? error?.name ?? "caserun_failed", message: error?.message ?? String(error) });
throw error;
}
}