fix: close HarnessRL authority recovery gaps

This commit is contained in:
root
2026-07-17 03:27:30 +02:00
parent 2563036317
commit ab76290e9d
13 changed files with 258 additions and 68 deletions
-20
View File
@@ -1,20 +0,0 @@
apiVersion: hwlab.pikastech.local/v1alpha1
kind: HarnessRLConsumerConfig
metadata:
name: harnessrl-v03
spec:
api:
serviceId: hwlab-harnessrl-api
port: 6675
worker:
serviceId: hwlab-harnessrl-worker
healthPort: 6676
registry:
databaseUrl: secretRef:hwlab-harnessrl-db/database-url
temporal:
address: temporal-frontend.temporal.svc.cluster.local:7233
namespace: unidesk
taskQueue: hwlab-v03-harnessrl
consumer:
cloudApiEnv:
HARNESSRL_API_URL: http://hwlab-harnessrl-api.hwlab-v03.svc.cluster.local:6675
+12 -2
View File
@@ -1058,19 +1058,29 @@ lanes:
- serviceId: hwlab-harnessrl-api
replicas: 1
env:
HARNESSRL_DATABASE_URL: secretRef:hwlab-harnessrl-db/database-url
HARNESSRL_DATABASE_URL: secretRef:hwlab-cloud-api-v03-db/database-url
HARNESSRL_TEMPORAL_ADDRESS: temporal-frontend.temporal.svc.cluster.local:7233
HARNESSRL_TEMPORAL_NAMESPACE: unidesk
HARNESSRL_TEMPORAL_TASK_QUEUE: hwlab-v03-harnessrl
HARNESSRL_ACTIVITY_START_TO_CLOSE_TIMEOUT_MS: "1800000"
HARNESSRL_ACTIVITY_HEARTBEAT_TIMEOUT_MS: "30000"
HARNESSRL_ACTIVITY_RETRY_INITIAL_INTERVAL_MS: "1000"
HARNESSRL_ACTIVITY_RETRY_MAXIMUM_INTERVAL_MS: "30000"
HARNESSRL_ACTIVITY_RETRY_MAXIMUM_ATTEMPTS: "5"
HARNESSRL_API_PORT: "6675"
OTEL_SERVICE_NAME: hwlab-harnessrl-api
- serviceId: hwlab-harnessrl-worker
replicas: 1
env:
HARNESSRL_DATABASE_URL: secretRef:hwlab-harnessrl-db/database-url
HARNESSRL_DATABASE_URL: secretRef:hwlab-cloud-api-v03-db/database-url
HARNESSRL_TEMPORAL_ADDRESS: temporal-frontend.temporal.svc.cluster.local:7233
HARNESSRL_TEMPORAL_NAMESPACE: unidesk
HARNESSRL_TEMPORAL_TASK_QUEUE: hwlab-v03-harnessrl
HARNESSRL_ACTIVITY_START_TO_CLOSE_TIMEOUT_MS: "1800000"
HARNESSRL_ACTIVITY_HEARTBEAT_TIMEOUT_MS: "30000"
HARNESSRL_ACTIVITY_RETRY_INITIAL_INTERVAL_MS: "1000"
HARNESSRL_ACTIVITY_RETRY_MAXIMUM_INTERVAL_MS: "30000"
HARNESSRL_ACTIVITY_RETRY_MAXIMUM_ATTEMPTS: "5"
HARNESSRL_WORKER_HEALTH_PORT: "6676"
OTEL_SERVICE_NAME: hwlab-harnessrl-worker
services:
+93 -17
View File
@@ -3,13 +3,32 @@ import { createHash } from "node:crypto";
import { mkdir, readFile, writeFile } from "node:fs/promises";
import path from "node:path";
import { buildCaseRun, collectCaseRun, prepareCaseRun } from "../../tools/src/hwlab-caserun-lib.ts";
import { buildCaseRun, collectAgentDiff, collectAgentTraceEvidence, collectCaseRun, prepareCaseRun, runAgentTaskStage } from "../../tools/src/hwlab-caserun-lib.ts";
import type { HarnessRLRegistry } from "./contracts.ts";
import { caseRunInternalFetch } from "../cloud/server-caserun-http.ts";
export function createHarnessRLActivities(options: { registry: HarnessRLRegistry; env?: Record<string, string | undefined> }) {
export function createHarnessRLActivities(options: {
registry: HarnessRLRegistry;
env?: Record<string, string | undefined>;
stages?: Partial<{
prepare(context: any, action: string): Promise<any>;
agent(context: any, run: any): Promise<any>;
trace(context: any, run: any): Promise<any>;
diff(context: any, run: any): Promise<any>;
build(context: any, run: any): Promise<any>;
collect(context: any, run: any, evidence: any): Promise<any>;
}>;
}) {
const env = options.env ?? process.env;
const registry = options.registry;
const stages = {
prepare: options.stages?.prepare ?? prepareCaseRun,
agent: options.stages?.agent ?? runAgentTaskStage,
trace: options.stages?.trace ?? collectAgentTraceEvidence,
diff: options.stages?.diff ?? collectAgentDiff,
build: options.stages?.build ?? buildCaseRun,
collect: options.stages?.collect ?? collectCaseRun
};
return {
async markRunning({ runId }: { runId: string }) {
await transition(registry, runId, "running", "prepare", { worker: "hwlab-harnessrl-worker" });
@@ -25,17 +44,47 @@ export function createHarnessRLActivities(options: { registry: HarnessRLRegistry
await transition(registry, input.runId, "running", "prepared", { mode: "software-smoke" });
return output;
}
const context = contextForRecord(record, env);
const prepared = await prepareCaseRun(context, "harnessrl-worker");
const context = contextForRecord(record, env, input.identity);
const prepared = await stages.prepare(context, "harnessrl-worker");
const output = { mode: definition.mode, prepared } as Record<string, unknown>;
await transition(registry, input.runId, "running", "prepared", { runDir: prepared.runDir, specPath: prepared.specPath });
return output;
});
},
async agent(input: { runId: string; identity: string }) {
return idempotent(registry, input, async () => {
const record = required(await registry.getRun(input.runId));
const prepared = activityOutput(required(await registry.getActivityResult(input.runId, "prepare", `${input.runId}:prepare:v1`))) as any;
if (prepared.mode === "software-smoke") return { mode: prepared.mode, run: prepared.run, authority: { agentRun: "not-applicable" } };
const stage = await stages.agent(contextForRecord(record, env, input.identity), prepared.prepared.run);
await transition(registry, input.runId, "running", "agent-completed", { traceId: stage.agent?.traceId ?? null, sessionId: stage.agent?.sessionId ?? null });
return { mode: prepared.mode, run: stage.run, agent: stage.agent };
});
},
async trace(input: { runId: string; identity: string }) {
return idempotent(registry, input, async () => {
const record = required(await registry.getRun(input.runId));
const agent = activityOutput(required(await registry.getActivityResult(input.runId, "agent", `${input.runId}:agent:v1`))) as any;
if (agent.mode === "software-smoke") return agent;
const stage = await stages.trace(contextForRecord(record, env, input.identity), agent.run);
await transition(registry, input.runId, "running", "trace-completed", { traceId: stage.trace?.traceId ?? null });
return { ...agent, run: stage.run, trace: stage.trace };
});
},
async diff(input: { runId: string; identity: string }) {
return idempotent(registry, input, async () => {
const record = required(await registry.getRun(input.runId));
const traced = activityOutput(required(await registry.getActivityResult(input.runId, "trace", `${input.runId}:trace:v1`))) as any;
if (traced.mode === "software-smoke") return traced;
const stage = await stages.diff(contextForRecord(record, env, input.identity), traced.run);
await transition(registry, input.runId, "running", "diff-completed", { diffSha256: stage.diff?.diffPatchSha256 ?? null });
return { ...traced, run: stage.run, diff: stage.diff };
});
},
async build(input: { runId: string; identity: string }) {
return idempotent(registry, input, async () => {
const record = required(await registry.getRun(input.runId));
const prepared = required(await registry.getActivityResult(input.runId, "prepare", `${input.runId}:prepare:v1`)).output as any;
const prepared = activityOutput(required(await registry.getActivityResult(input.runId, "diff", `${input.runId}:diff:v1`))) as any;
activityHeartbeat({ runId: input.runId, stage: "build" });
if (prepared.mode === "software-smoke") {
const artifact = {
@@ -50,23 +99,30 @@ export function createHarnessRLActivities(options: { registry: HarnessRLRegistry
await transition(registry, input.runId, "running", "build-completed", { artifactSha256: artifactRef.sha256 });
return output;
}
const context = contextForRecord(record, env);
const build = await buildCaseRun(context, prepared.prepared.run);
const context = contextForRecord(record, env, input.identity);
const build = await stages.build(context, prepared.run);
await transition(registry, input.runId, "running", "build-completed", { jobId: build.summary?.jobId ?? build.evidence?.keilJob?.jobId ?? null });
return { mode: prepared.mode, run: build.run ?? prepared.prepared.run, build };
return { mode: prepared.mode, run: build.run ?? prepared.run, build, agent: prepared.agent, trace: prepared.trace, diff: prepared.diff };
});
},
async collect(input: { runId: string; identity: string }) {
return idempotent(registry, input, async () => {
const record = required(await registry.getRun(input.runId));
const built = required(await registry.getActivityResult(input.runId, "build", `${input.runId}:build:v1`)).output as any;
const built = activityOutput(required(await registry.getActivityResult(input.runId, "build", `${input.runId}:build:v1`))) as any;
activityHeartbeat({ runId: input.runId, stage: "collect" });
let result: Record<string, unknown>;
if (built.mode === "software-smoke") result = { summary: built.summary, evidence: built.evidence, run: built.run };
else {
const context = contextForRecord(record, env);
const collected = await collectCaseRun(context, built.run, built.build.evidence);
result = { summary: built.build.summary, evidence: built.build.evidence, run: collected.run ?? built.run, collected: { status: collected.status, summary: collected.summary ?? null } };
const context = contextForRecord(record, env, input.identity);
const collected = await stages.collect(context, built.run, built.build.evidence);
const manifest = await manifestRefs(collected.summary);
result = {
summary: built.build.summary, evidence: built.build.evidence, run: collected.run ?? built.run,
authority: { agentRun: authorityRefs(built), hwpod: built.build.evidence?.operation ?? built.build.evidence?.validation ?? null },
artifactManifestPath: manifest.path, artifactManifestSha256: manifest.sha256,
aggregate: manifest.aggregate,
collected: { status: collected.status, summary: collected.summary ?? null }
};
}
await registry.updateRun(input.runId, { result });
await transition(registry, input.runId, "running", "collect-completed", {});
@@ -85,9 +141,11 @@ export function createHarnessRLActivities(options: { registry: HarnessRLRegistry
async function idempotent(registry: HarnessRLRegistry, input: { runId: string; identity: string }, execute: () => Promise<Record<string, unknown>>) {
const activityName = input.identity.split(":").at(-2) ?? "activity";
const existing = await registry.getActivityResult(input.runId, activityName, input.identity);
if (existing) return existing.output;
if ((existing?.output as any)?.state === "completed") return (existing!.output as any).result;
if (!existing) await registry.saveActivityResult({ runId: input.runId, activityName, identity: input.identity, output: { state: "accepted", authorityIdentity: input.identity } });
const output = await execute();
return (await registry.saveActivityResult({ runId: input.runId, activityName, identity: input.identity, output })).output;
const saved = await registry.saveActivityResult({ runId: input.runId, activityName, identity: input.identity, output: { state: "completed", authorityIdentity: input.identity, result: output } });
return activityOutput(saved);
}
async function transition(registry: HarnessRLRegistry, runId: string, status: any, stage: string, payload: Record<string, unknown>) {
@@ -95,16 +153,34 @@ async function transition(registry: HarnessRLRegistry, runId: string, status: an
await registry.appendEvent(runId, status, stage, payload);
}
function contextForRecord(record: any, env: Record<string, string | undefined>) {
function contextForRecord(record: any, env: Record<string, string | undefined>, identity: string) {
const apiUrl = String(env.HWLAB_CASERUN_INTERNAL_API_URL ?? env.HWLAB_RUNTIME_INTERNAL_API_URL ?? record.runtimeApiUrl).replace(/\/+$/u, "");
const traceId = `trc_harnessrl_${createHash("sha256").update(identity).digest("hex").slice(0, 24)}`;
return {
parsed: { _: [], caseId: record.caseId, runId: record.runId, runDir: record.runDir, caseRepo: record.caseRepo, apiUrl, noCaseRepoRecord: true, stateDir: path.join(path.dirname(record.runDir), "cli") },
env: { ...process.env, ...env, HWLAB_CASE_REPO: record.caseRepo, HWLAB_RUNTIME_API_URL: apiUrl, HWLAB_CASERUN_PUBLIC_RUNTIME_API_URL: record.runtimeApiUrl, HWLAB_RUNTIME_ENDPOINT_LOCKED: "1" },
parsed: { _: [], caseId: record.caseId, runId: record.runId, runDir: record.runDir, caseRepo: record.caseRepo, apiUrl, traceId, operationIdentity: identity, noCaseRepoRecord: false, stateDir: path.join(path.dirname(record.runDir), "cli") },
env: { ...process.env, ...env, HWLAB_CASE_REPO: record.caseRepo, HWLAB_RUNTIME_API_URL: apiUrl, HWLAB_CASERUN_PUBLIC_RUNTIME_API_URL: record.runtimeApiUrl, HWLAB_RUNTIME_ENDPOINT_LOCKED: "1", HWLAB_HWPOD_OPERATION_IDENTITY: identity },
fetchImpl: caseRunInternalFetch(globalThis.fetch, apiUrl, env), cwd: record.caseRepo,
now: () => new Date().toISOString(), sleep: (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)), argv: [], rest: ["worker", record.caseId]
};
}
function activityOutput(result: { output: Record<string, unknown> }) {
return (result.output as any)?.state === "completed" ? (result.output as any).result : result.output;
}
async function manifestRefs(summary: any) {
const manifestPath = String(summary?.artifactManifestPath ?? "");
if (!manifestPath) return { path: null, sha256: null, aggregate: null };
const body = await readFile(manifestPath);
const manifest = JSON.parse(body.toString("utf8"));
return { path: manifestPath, sha256: createHash("sha256").update(body).digest("hex"), aggregate: manifest.aggregate ?? null };
}
function authorityRefs(built: any) {
const agent = built.agent ?? built.run?.agent ?? {};
return { conversationId: agent.conversationId ?? null, sessionId: agent.sessionId ?? null, threadId: agent.threadId ?? null, traceId: agent.traceId ?? null, resultUrl: agent.resultUrl ?? null, commandStatus: agent.commandStatus ?? null, final: agent.finalResponse ?? agent.result ?? null, diffSha256: built.diff?.diffPatchSha256 ?? built.run?.agentDiff?.diffPatchSha256 ?? null };
}
async function loadDefinition(caseRepo: string, caseId: string) {
return JSON.parse(await readFile(path.join(caseRepo, "cases", caseId, "case.json"), "utf8"));
}
+1 -1
View File
@@ -57,7 +57,7 @@ export interface HarnessRLRegistry {
}
export interface HarnessRLTemporalGateway {
start(input: { workflowId: string; runId: string }): Promise<{ workflowRunId: string }>;
start(input: { workflowId: string; runId: string }): Promise<{ workflowRunId: string; reused?: boolean }>;
requestCancel(workflowId: string): Promise<void>;
close(): Promise<void>;
}
+88 -3
View File
@@ -35,11 +35,18 @@ test("worker restart reuses idempotent prepare/build/collect activity results",
const firstWorker = createHarnessRLActivities({ registry, env: {} });
await firstWorker.markRunning({ runId: "worker-restart" });
const prepared = await firstWorker.prepare({ runId: "worker-restart", identity: "worker-restart:prepare:v1" });
await firstWorker.agent({ runId: "worker-restart", identity: "worker-restart:agent:v1" });
await firstWorker.trace({ runId: "worker-restart", identity: "worker-restart:trace:v1" });
await firstWorker.diff({ runId: "worker-restart", identity: "worker-restart:diff:v1" });
const built = await firstWorker.build({ runId: "worker-restart", identity: "worker-restart:build:v1" });
const restartedWorker = createHarnessRLActivities({ registry, env: {} });
assert.deepEqual(await restartedWorker.prepare({ runId: "worker-restart", identity: "worker-restart:prepare:v1" }), prepared);
assert.deepEqual(await restartedWorker.build({ runId: "worker-restart", identity: "worker-restart:build:v1" }), built);
await restartedWorker.agent({ runId: "worker-restart", identity: "worker-restart:agent:v1" });
await restartedWorker.trace({ runId: "worker-restart", identity: "worker-restart:trace:v1" });
await restartedWorker.diff({ runId: "worker-restart", identity: "worker-restart:diff:v1" });
const rebuilt = await restartedWorker.build({ runId: "worker-restart", identity: "worker-restart:build:v1" });
assert.equal(rebuilt.mode, built.mode);
await restartedWorker.collect({ runId: "worker-restart", identity: "worker-restart:collect:v1" });
await restartedWorker.markCompleted({ runId: "worker-restart" });
@@ -48,7 +55,64 @@ test("worker restart reuses idempotent prepare/build/collect activity results",
assert.equal((record?.result as any)?.summary?.mode, "software-smoke");
const artifact = JSON.parse(await readFile(path.join(record!.runDir, "software-smoke.json"), "utf8"));
assert.equal(artifact.checks.temporalWorkerRunning, true);
assert.equal(registry.activityResults.size, 3);
assert.equal(registry.activityResults.size, 6);
});
test("AgentRun and HWPOD crash gaps reuse stable authoritative operation identities", async () => {
const fixture = await hardwareFixture();
const registry = new MemoryRegistry();
await registry.createRun({ runId: "authority-gap", caseId: fixture.caseId, workflowId: "workflow-authority-gap", caseRepo: fixture.root, runDir: path.join(fixture.stateRoot, "runs", "authority-gap"), runtimeApiUrl: "http://127.0.0.1:6667" });
const acceptedAgent = new Set<string>();
const acceptedHwpod = new Set<string>();
let agentCrash = true;
let hwpodCrash = true;
const stages = {
async prepare(context: any) { return { run: { runId: "authority-gap", caseId: fixture.caseId, runDir: context.parsed.runDir } }; },
async agent(context: any, run: any) {
acceptedAgent.add(context.parsed.traceId);
if (agentCrash) { agentCrash = false; throw new Error("worker exited after AgentRun accepted"); }
return { run: { ...run, agent: { traceId: context.parsed.traceId, sessionId: "ses-authoritative", resultUrl: "/v1/agentrun/result" } }, agent: { traceId: context.parsed.traceId, sessionId: "ses-authoritative", resultUrl: "/v1/agentrun/result" } };
},
async trace(_context: any, run: any) { return { run, trace: { traceId: run.agent.traceId, status: "completed" } }; },
async diff(_context: any, run: any) { return { run, diff: { diffPatchSha256: "d".repeat(64) } }; },
async build(context: any, run: any) {
acceptedHwpod.add(context.parsed.operationIdentity);
if (hwpodCrash) { hwpodCrash = false; throw new Error("worker exited after HWPOD accepted"); }
return { run, summary: { jobId: "job-authoritative" }, evidence: { operation: { operationId: context.parsed.operationIdentity, jobId: "job-authoritative" } } };
},
async collect(_context: any, run: any, evidence: any) {
const manifestPath = path.join(run.runDir, "artifact-manifest.json");
await mkdir(run.runDir, { recursive: true });
await writeFile(manifestPath, `${JSON.stringify({ aggregate: { path: "aggregate.md", sha256: "a".repeat(64) } })}\n`);
return { run, status: "completed", summary: { artifactManifestPath: manifestPath }, evidence };
}
};
const worker = createHarnessRLActivities({ registry, env: {}, stages });
await worker.prepare({ runId: "authority-gap", identity: "authority-gap:prepare:v1" });
await assert.rejects(worker.agent({ runId: "authority-gap", identity: "authority-gap:agent:v1" }), /AgentRun accepted/);
await worker.agent({ runId: "authority-gap", identity: "authority-gap:agent:v1" });
await worker.trace({ runId: "authority-gap", identity: "authority-gap:trace:v1" });
await worker.diff({ runId: "authority-gap", identity: "authority-gap:diff:v1" });
await assert.rejects(worker.build({ runId: "authority-gap", identity: "authority-gap:build:v1" }), /HWPOD accepted/);
await worker.build({ runId: "authority-gap", identity: "authority-gap:build:v1" });
const result = await worker.collect({ runId: "authority-gap", identity: "authority-gap:collect:v1" });
assert.equal(acceptedAgent.size, 1);
assert.equal(acceptedHwpod.size, 1);
assert.equal((result as any).authority.agentRun.sessionId, "ses-authoritative");
assert.equal((result as any).authority.hwpod.jobId, "job-authoritative");
assert.equal(typeof (result as any).artifactManifestSha256, "string");
assert.equal((result as any).aggregate.sha256, "a".repeat(64));
});
test("Temporal start crash gap reuses the existing workflow identity", async () => {
const fixture = await softwareSmokeFixture();
const registry = new MemoryRegistry();
const temporal = new CrashGapTemporal();
const service = createHarnessRLService({ registry, temporal, cwd: fixture.root, caseRepo: fixture.root, stateRoot: fixture.stateRoot });
await assert.rejects(service.submit({ caseId: fixture.caseId, runId: "start-gap", runtimeApiUrl: "http://harnessrl.test" }), /API exited/);
const recovered = await service.submit({ caseId: fixture.caseId, runId: "start-gap", runtimeApiUrl: "http://harnessrl.test" });
assert.equal(recovered?.workflowRunId, "temporal-start-gap");
assert.equal(temporal.accepted.size, 1);
});
test("cancel converges to durable canceled through workflow signal", async () => {
@@ -87,6 +151,14 @@ async function softwareSmokeFixture() {
return { root, caseId, stateRoot: path.join(root, ".state", "harnessrl") };
}
async function hardwareFixture() {
const root = await mkdtemp(path.join(os.tmpdir(), "harnessrl-authority-"));
const caseId = "hardware-authority";
await mkdir(path.join(root, "cases", caseId), { recursive: true });
await writeFile(path.join(root, "cases", caseId, "case.json"), `${JSON.stringify({ caseId, title: "Authoritative fixture", mode: "compile-only" }, null, 2)}\n`);
return { root, caseId, stateRoot: path.join(root, ".state", "harnessrl") };
}
class FakeTemporal implements HarnessRLTemporalGateway {
starts = 0;
constructor(private readonly onCancel?: (workflowId: string) => Promise<void>) {}
@@ -95,6 +167,19 @@ class FakeTemporal implements HarnessRLTemporalGateway {
async close() {}
}
class CrashGapTemporal implements HarnessRLTemporalGateway {
accepted = new Map<string, string>();
crashed = false;
async start(input: { workflowId: string }) {
const workflowRunId = this.accepted.get(input.workflowId) ?? "temporal-start-gap";
this.accepted.set(input.workflowId, workflowRunId);
if (!this.crashed) { this.crashed = true; throw new Error("API exited after Temporal accepted"); }
return { workflowRunId, reused: true };
}
async requestCancel() {}
async close() {}
}
class MemoryRegistry implements HarnessRLRegistry {
runs = new Map<string, CaseRunRecord>();
events = new Map<string, CaseRunEvent[]>();
@@ -123,7 +208,7 @@ class MemoryRegistry implements HarnessRLRegistry {
}
async listEvents(runId: string) { return structuredClone(this.events.get(runId) ?? []); }
async getActivityResult(runId: string, activityName: string, identity: string) { return structuredClone(this.activityResults.get(`${runId}:${activityName}:${identity}`) ?? null); }
async saveActivityResult(result: ActivityResult) { const key = `${result.runId}:${result.activityName}:${result.identity}`; if (!this.activityResults.has(key)) this.activityResults.set(key, structuredClone(result)); return structuredClone(this.activityResults.get(key)!); }
async saveActivityResult(result: ActivityResult) { const key = `${result.runId}:${result.activityName}:${result.identity}`; this.activityResults.set(key, structuredClone(result)); return structuredClone(this.activityResults.get(key)!); }
async health() { return { ok: true as const }; }
async close() {}
}
+1 -1
View File
@@ -98,7 +98,7 @@ export class PostgresHarnessRLRegistry implements HarnessRLRegistry {
await this.ensureSchema();
const saved = await this.pool.query(
`INSERT INTO harnessrl_activity_results (run_id,activity_name,identity,output) VALUES ($1,$2,$3,$4)
ON CONFLICT (run_id,activity_name,identity) DO UPDATE SET output=harnessrl_activity_results.output RETURNING *`,
ON CONFLICT (run_id,activity_name,identity) DO UPDATE SET output=EXCLUDED.output RETURNING *`,
[result.runId, result.activityName, result.identity, result.output]
);
return activityRow(saved.rows[0]);
+15 -2
View File
@@ -12,9 +12,22 @@ export function harnessRLRuntime(env: Record<string, string | undefined> = proce
const temporalAddress = String(env.HARNESSRL_TEMPORAL_ADDRESS ?? env.TEMPORAL_ADDRESS ?? "");
const temporalNamespace = String(env.HARNESSRL_TEMPORAL_NAMESPACE ?? "unidesk");
const taskQueue = String(env.HARNESSRL_TEMPORAL_TASK_QUEUE ?? "hwlab-v03-harnessrl");
const temporal = createHarnessRLTemporalGateway({ address: temporalAddress, namespace: temporalNamespace, taskQueue });
const activity = {
startToCloseTimeoutMs: positiveInteger(env.HARNESSRL_ACTIVITY_START_TO_CLOSE_TIMEOUT_MS, 1_800_000, "HARNESSRL_ACTIVITY_START_TO_CLOSE_TIMEOUT_MS"),
heartbeatTimeoutMs: positiveInteger(env.HARNESSRL_ACTIVITY_HEARTBEAT_TIMEOUT_MS, 30_000, "HARNESSRL_ACTIVITY_HEARTBEAT_TIMEOUT_MS"),
retryInitialIntervalMs: positiveInteger(env.HARNESSRL_ACTIVITY_RETRY_INITIAL_INTERVAL_MS, 1_000, "HARNESSRL_ACTIVITY_RETRY_INITIAL_INTERVAL_MS"),
retryMaximumIntervalMs: positiveInteger(env.HARNESSRL_ACTIVITY_RETRY_MAXIMUM_INTERVAL_MS, 30_000, "HARNESSRL_ACTIVITY_RETRY_MAXIMUM_INTERVAL_MS"),
retryMaximumAttempts: positiveInteger(env.HARNESSRL_ACTIVITY_RETRY_MAXIMUM_ATTEMPTS, 5, "HARNESSRL_ACTIVITY_RETRY_MAXIMUM_ATTEMPTS")
};
const temporal = createHarnessRLTemporalGateway({ address: temporalAddress, namespace: temporalNamespace, taskQueue, activity });
return {
cwd, caseRepo, stateRoot, registry, temporalAddress, temporalNamespace, taskQueue,
cwd, caseRepo, stateRoot, registry, temporalAddress, temporalNamespace, taskQueue, activity,
service: createHarnessRLService({ registry, temporal, cwd, caseRepo, stateRoot })
};
}
function positiveInteger(value: string | undefined, fallback: number, name: string) {
const parsed = Number.parseInt(String(value ?? fallback), 10);
if (!Number.isSafeInteger(parsed) || parsed <= 0) throw Object.assign(new Error(`${name} must be a positive integer`), { code: "harnessrl_temporal_config_invalid", details: { name } });
return parsed;
}
+1
View File
@@ -33,6 +33,7 @@ export function createHarnessRLService(options: { registry: HarnessRLRegistry; t
await registry.appendEvent(runId, "queued", "queued", { caseId });
const started = await temporal.start({ workflowId, runId });
await registry.updateRun(runId, { workflowRunId: started.workflowRunId });
if (started.reused) await registry.appendEvent(runId, "queued", "workflow-reused", { workflowId, workflowRunId: started.workflowRunId });
}
return await registry.getRun(runId);
},
+14 -7
View File
@@ -2,7 +2,7 @@ import { Client, Connection } from "@temporalio/client";
import type { HarnessRLTemporalGateway } from "./contracts.ts";
export function createHarnessRLTemporalGateway(options: { address: string; namespace: string; taskQueue: string }): HarnessRLTemporalGateway {
export function createHarnessRLTemporalGateway(options: { address: string; namespace: string; taskQueue: string; activity: Record<string, number> }): HarnessRLTemporalGateway {
if (!options.address) throw codedError("harnessrl_temporal_address_required", "HARNESSRL_TEMPORAL_ADDRESS is required");
let connection: Connection | undefined;
let client: Client | undefined;
@@ -13,12 +13,19 @@ export function createHarnessRLTemporalGateway(options: { address: string; names
};
return {
async start(input) {
const handle = await (await getClient()).workflow.start("caseRunWorkflow", {
taskQueue: options.taskQueue,
workflowId: input.workflowId,
args: [{ runId: input.runId }]
});
return { workflowRunId: handle.firstExecutionRunId };
const temporalClient = await getClient();
try {
const handle = await temporalClient.workflow.start("caseRunWorkflow", {
taskQueue: options.taskQueue,
workflowId: input.workflowId,
args: [{ runId: input.runId, activity: options.activity }]
});
return { workflowRunId: handle.firstExecutionRunId, reused: false };
} catch (error: any) {
if (error?.name !== "WorkflowExecutionAlreadyStartedError") throw error;
const description = await temporalClient.workflow.getHandle(input.workflowId).describe();
return { workflowRunId: description.runId, reused: true };
}
},
async requestCancel(workflowId) {
await (await getClient()).workflow.getHandle(workflowId).signal("cancelCaseRun");
+15 -7
View File
@@ -3,6 +3,9 @@ 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>;
@@ -10,15 +13,14 @@ type Activities = {
markFailed(input: { runId: string; code: string; message: string }): Promise<void>;
};
const activities = proxyActivities<Activities>({
startToCloseTimeout: "30 minutes",
heartbeatTimeout: "30 seconds",
retry: { initialInterval: "1 second", maximumInterval: "30 seconds", maximumAttempts: 5 }
});
export const cancelCaseRun = defineSignal("cancelCaseRun");
export async function caseRunWorkflow(input: { runId: string }) {
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 () => {
@@ -31,6 +33,12 @@ export async function caseRunWorkflow(input: { runId: string }) {
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` });
+10 -5
View File
@@ -704,7 +704,7 @@ function defaultAgentTask(definition: Record<string, unknown>, caseId: string):
};
}
async function runAgentTaskStage(context: CaseContext, run: PreparedCaseRun) {
export async function runAgentTaskStage(context: CaseContext, run: PreparedCaseRun) {
if (!run.agentTask) {
run = await updateRun(context, run, { agentTask: defaultAgentTask(run.definition, run.caseId) });
}
@@ -794,7 +794,7 @@ async function runAgentTaskStage(context: CaseContext, run: PreparedCaseRun) {
return { run: nextRun, agent };
}
async function collectAgentTraceEvidence(context: CaseContext, run: PreparedCaseRun) {
export async function collectAgentTraceEvidence(context: CaseContext, run: PreparedCaseRun) {
const agent = run.agent;
const traceId = text(agent?.traceId);
if (!agent || !traceId) {
@@ -915,7 +915,7 @@ function effectiveAgentPollIntervalMs(context: CaseContext, agentTask?: Prepared
return Math.max(numberOption(context.parsed.agentPollIntervalMs ?? context.parsed.pollIntervalMs) ?? agentTask?.pollIntervalMs ?? DEFAULT_POLL_INTERVAL_MS, 250);
}
async function collectAgentDiff(context: CaseContext, run: PreparedCaseRun) {
export async function collectAgentDiff(context: CaseContext, run: PreparedCaseRun) {
const status = await requestSubjectCommand(context, run, ["-C", run.subject.worktreePath, "status", "--short"]);
const diffStat = await requestSubjectCommand(context, run, ["-C", run.subject.worktreePath, "diff", "--stat"]);
const diffPatch = await requestSubjectCommand(context, run, ["-C", run.subject.worktreePath, "diff", "--binary"]);
@@ -965,7 +965,7 @@ export async function requestSubjectProcess(context: CaseContext, run: PreparedC
const hwpodId = text(document.metadata.name) || text(document.metadata.uid) || run.caseId;
const plan = {
contractVersion: "hwpod-node-ops-v1",
planId: `case_subject_cmd_${randomUUID()}`,
planId: stableOperationIdentity(context, `subject-${command}-${argv.join("-")}`),
hwpodId,
nodeId: document.spec.nodeBinding.nodeId,
intent: "cmd.run",
@@ -1009,7 +1009,7 @@ async function requestSubjectCommand(context: CaseContext, run: PreparedCaseRun,
const hwpodId = text(document.metadata.name) || text(document.metadata.uid) || run.caseId;
const plan = {
contractVersion: "hwpod-node-ops-v1",
planId: `case_subject_cmd_${randomUUID()}`,
planId: stableOperationIdentity(context, `subject-git-${argv.join("-")}`),
hwpodId,
nodeId: document.spec.nodeBinding.nodeId,
intent: "cmd.run",
@@ -1032,6 +1032,11 @@ async function requestSubjectCommand(context: CaseContext, run: PreparedCaseRun,
return { ok: response.status >= 200 && response.status < 300 && exitCode === 0, httpStatus: response.status, exitCode, stdout: text(result?.stdout), stderr: text(result?.stderr), body: compactObject(body) };
}
function stableOperationIdentity(context: CaseContext, suffix: string) {
const base = text(context.parsed.operationIdentity);
return base ? `${base}:${slug(suffix)}` : `case_subject_cmd_${randomUUID()}`;
}
async function updateRun(context: CaseContext, run: PreparedCaseRun, patch: Partial<PreparedCaseRun> & Record<string, unknown>) {
const nextRun = { ...run, ...patch, updatedAt: context.now() } as PreparedCaseRun;
await writeRunState(nextRun);
+7 -2
View File
@@ -103,7 +103,7 @@ export async function requestSubjectWorktree(context: CaseContext, input: { case
}, ...(sidecarOp ? [sidecarOp] : [])];
const plan = {
contractVersion: "hwpod-node-ops-v1",
planId: `case_subject_${randomUUID()}`,
planId: stablePlanId(context, "subject-prepare"),
hwpodId,
nodeId: document.spec.nodeBinding.nodeId,
intent: "cmd.run",
@@ -408,7 +408,7 @@ export async function requestKeilJobStatus(context: CaseContext, run: PreparedCa
const hwpodId = text(document.metadata.name) || text(document.metadata.uid) || run.caseId;
const plan = {
contractVersion: "hwpod-node-ops-v1",
planId: `case_job_${randomUUID()}`,
planId: stablePlanId(context, `job-${jobId}`),
hwpodId,
nodeId: document.spec.nodeBinding.nodeId,
intent: "cmd.run",
@@ -433,6 +433,11 @@ export async function requestKeilJobStatus(context: CaseContext, run: PreparedCa
return { status: response.status, body: parseJsonMaybe(textBody) ?? { raw: textBody }, plan };
}
function stablePlanId(context: CaseContext, suffix: string) {
const identity = String((context.parsed as any).operationIdentity ?? "").trim();
return identity ? `${identity}:${suffix.replace(/[^a-z0-9_.:-]+/giu, "-")}` : `case_plan_${randomUUID()}`;
}
export function keilJobStatus(body: any) {
const parsed = parseJsonMaybe(hwpodResultTexts(body)[0]) ?? body;
const status = text(parsed?.status ?? parsed?.data?.status ?? parsed?.job?.status) || "unknown";
+1 -1
View File
@@ -151,7 +151,7 @@ export function compileHwpodNodeOpsPlan({ document, specPath = DEFAULT_HWPOD_SPE
}));
return {
contractVersion: HWPOD_NODE_OPS_CONTRACT_VERSION,
planId: `hwpod_plan_${randomUUID()}`,
planId: String(process.env.HWLAB_HWPOD_OPERATION_IDENTITY ?? "").trim() || `hwpod_plan_${randomUUID()}`,
hwpodId,
nodeId,
intent: normalizedIntent,