53 lines
3.2 KiB
TypeScript
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;
|
|
}
|
|
}
|