Files
pikasTech-HWLAB/internal/harnessrl/workflows.ts
T
2026-07-24 16:27:23 +02:00

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 }) };
}