94 lines
5.7 KiB
TypeScript
94 lines
5.7 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; details?: unknown }): 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 };
|
|
const agent = await activities.agent({ runId: input.runId, identity: `${input.runId}:agent:v1` });
|
|
const agentFailure = caseRunAgentFailure(agent);
|
|
if (agentFailure) {
|
|
await activities.markFailed({ runId: input.runId, ...agentFailure });
|
|
return { status: "failed", runId: input.runId };
|
|
}
|
|
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 };
|
|
const built = 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 };
|
|
const buildFailure = caseRunBuildFailure(built);
|
|
if (buildFailure) {
|
|
await activities.markFailed({ runId: input.runId, ...buildFailure });
|
|
return { status: "failed", runId: input.runId };
|
|
}
|
|
await activities.markCompleted({ runId: input.runId });
|
|
return { status: "completed", runId: input.runId };
|
|
} catch (error: any) {
|
|
await activities.markFailed({ runId: input.runId, ...caseRunWorkflowFailure(error) });
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
export function caseRunAgentFailure(result: any) {
|
|
const agent = result?.agent ?? {};
|
|
const status = String(agent.stageStatus ?? "").toLowerCase();
|
|
if (agent.timedOut === true || status === "timeout" || status === "timed_out") return { code: "agentrun_task_timeout", message: "AgentRun task did not reach a terminal result within the CaseRun timeout", details: { traceId: agent.traceId ?? null, status } };
|
|
if (agent.error || ["failed", "cancelled", "canceled"].includes(status)) return { code: String(agent.error?.code ?? "agentrun_task_failed"), message: String(agent.error?.message ?? `AgentRun task reached ${status || "failed"}`), details: { traceId: agent.traceId ?? null, status, error: agent.error ?? null } };
|
|
return null;
|
|
}
|
|
|
|
export function caseRunBuildFailure(result: any) {
|
|
const validation = result?.validation ?? {};
|
|
if (validation.status !== "blocked" && !validation.blocker) return null;
|
|
const blocker = validation.blocker ?? {};
|
|
return { code: String(blocker.code ?? "hwpod_validation_blocked"), message: String(blocker.summary ?? "HWPOD post-agent validation was blocked"), details: blocker.details ?? blocker };
|
|
}
|
|
|
|
export function caseRunWorkflowFailure(error: any) {
|
|
let current = error;
|
|
let fallback = { code: "caserun_failed", message: String(error), details: undefined as unknown };
|
|
for (let depth = 0; current && depth < 8; depth += 1) {
|
|
const code = String(current.code ?? current.type ?? current.name ?? "").trim();
|
|
const message = String(current.message ?? "").trim();
|
|
const details = Array.isArray(current.details) && current.details.length === 1 ? current.details[0] : current.details;
|
|
if (code || message) fallback = { code: code || fallback.code, message: message || fallback.message, details };
|
|
if (code && !["ActivityFailure", "ApplicationFailure", "Error"].includes(code) && message !== "Activity task failed") {
|
|
return { code, message: message || code, ...(details === undefined ? {} : { details }) };
|
|
}
|
|
current = current.cause;
|
|
}
|
|
return { code: fallback.code, message: fallback.message, ...(fallback.details === undefined ? {} : { details: fallback.details }) };
|
|
}
|