From ab76290e9df9b466f9836dec69440a1ee932004a Mon Sep 17 00:00:00 2001 From: root Date: Fri, 17 Jul 2026 03:27:30 +0200 Subject: [PATCH] fix: close HarnessRL authority recovery gaps --- config/hwlab-v03/harnessrl.private.yaml | 20 ----- deploy/deploy.yaml | 14 ++- internal/harnessrl/activities.ts | 110 ++++++++++++++++++++---- internal/harnessrl/contracts.ts | 2 +- internal/harnessrl/harnessrl.test.ts | 91 +++++++++++++++++++- internal/harnessrl/registry.ts | 2 +- internal/harnessrl/runtime.ts | 17 +++- internal/harnessrl/service.ts | 1 + internal/harnessrl/temporal.ts | 21 +++-- internal/harnessrl/workflows.ts | 22 +++-- tools/src/hwlab-caserun-runtime.ts | 15 ++-- tools/src/hwlab-caserun-subject.ts | 9 +- tools/src/hwpod-harness-lib.ts | 2 +- 13 files changed, 258 insertions(+), 68 deletions(-) delete mode 100644 config/hwlab-v03/harnessrl.private.yaml diff --git a/config/hwlab-v03/harnessrl.private.yaml b/config/hwlab-v03/harnessrl.private.yaml deleted file mode 100644 index fcecff52..00000000 --- a/config/hwlab-v03/harnessrl.private.yaml +++ /dev/null @@ -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 diff --git a/deploy/deploy.yaml b/deploy/deploy.yaml index 50b9f462..cf0d5909 100644 --- a/deploy/deploy.yaml +++ b/deploy/deploy.yaml @@ -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: diff --git a/internal/harnessrl/activities.ts b/internal/harnessrl/activities.ts index f60f50bf..62200530 100644 --- a/internal/harnessrl/activities.ts +++ b/internal/harnessrl/activities.ts @@ -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 }) { +export function createHarnessRLActivities(options: { + registry: HarnessRLRegistry; + env?: Record; + stages?: Partial<{ + prepare(context: any, action: string): Promise; + agent(context: any, run: any): Promise; + trace(context: any, run: any): Promise; + diff(context: any, run: any): Promise; + build(context: any, run: any): Promise; + collect(context: any, run: any, evidence: any): Promise; + }>; +}) { 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; 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; 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>) { 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) { @@ -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) { +function contextForRecord(record: any, env: Record, 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 }) { + 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")); } diff --git a/internal/harnessrl/contracts.ts b/internal/harnessrl/contracts.ts index 608e2669..02fb059f 100644 --- a/internal/harnessrl/contracts.ts +++ b/internal/harnessrl/contracts.ts @@ -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; close(): Promise; } diff --git a/internal/harnessrl/harnessrl.test.ts b/internal/harnessrl/harnessrl.test.ts index fad960ba..752ef890 100644 --- a/internal/harnessrl/harnessrl.test.ts +++ b/internal/harnessrl/harnessrl.test.ts @@ -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(); + const acceptedHwpod = new Set(); + 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) {} @@ -95,6 +167,19 @@ class FakeTemporal implements HarnessRLTemporalGateway { async close() {} } +class CrashGapTemporal implements HarnessRLTemporalGateway { + accepted = new Map(); + 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(); events = new Map(); @@ -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() {} } diff --git a/internal/harnessrl/registry.ts b/internal/harnessrl/registry.ts index 7550bb58..e3941775 100644 --- a/internal/harnessrl/registry.ts +++ b/internal/harnessrl/registry.ts @@ -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]); diff --git a/internal/harnessrl/runtime.ts b/internal/harnessrl/runtime.ts index e06c7b3c..24e882d4 100644 --- a/internal/harnessrl/runtime.ts +++ b/internal/harnessrl/runtime.ts @@ -12,9 +12,22 @@ export function harnessRLRuntime(env: Record = 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; +} diff --git a/internal/harnessrl/service.ts b/internal/harnessrl/service.ts index 1cea755e..2bdc3af1 100644 --- a/internal/harnessrl/service.ts +++ b/internal/harnessrl/service.ts @@ -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); }, diff --git a/internal/harnessrl/temporal.ts b/internal/harnessrl/temporal.ts index 3e4b9bae..c24965e7 100644 --- a/internal/harnessrl/temporal.ts +++ b/internal/harnessrl/temporal.ts @@ -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 }): 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"); diff --git a/internal/harnessrl/workflows.ts b/internal/harnessrl/workflows.ts index 3995be71..6bb593b8 100644 --- a/internal/harnessrl/workflows.ts +++ b/internal/harnessrl/workflows.ts @@ -3,6 +3,9 @@ import { defineSignal, proxyActivities, setHandler } from "@temporalio/workflow" type Activities = { markRunning(input: { runId: string }): Promise; prepare(input: { runId: string; identity: string }): Promise>; + agent(input: { runId: string; identity: string }): Promise>; + trace(input: { runId: string; identity: string }): Promise>; + diff(input: { runId: string; identity: string }): Promise>; build(input: { runId: string; identity: string }): Promise>; collect(input: { runId: string; identity: string }): Promise>; markCompleted(input: { runId: string }): Promise; @@ -10,15 +13,14 @@ type Activities = { markFailed(input: { runId: string; code: string; message: string }): Promise; }; -const activities = proxyActivities({ - 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({ + 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` }); diff --git a/tools/src/hwlab-caserun-runtime.ts b/tools/src/hwlab-caserun-runtime.ts index 4744b86b..b04ea888 100644 --- a/tools/src/hwlab-caserun-runtime.ts +++ b/tools/src/hwlab-caserun-runtime.ts @@ -704,7 +704,7 @@ function defaultAgentTask(definition: Record, 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 & Record) { const nextRun = { ...run, ...patch, updatedAt: context.now() } as PreparedCaseRun; await writeRunState(nextRun); diff --git a/tools/src/hwlab-caserun-subject.ts b/tools/src/hwlab-caserun-subject.ts index c9de6088..096c9953 100644 --- a/tools/src/hwlab-caserun-subject.ts +++ b/tools/src/hwlab-caserun-subject.ts @@ -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"; diff --git a/tools/src/hwpod-harness-lib.ts b/tools/src/hwpod-harness-lib.ts index 27f6b008..411fd1ac 100644 --- a/tools/src/hwpod-harness-lib.ts +++ b/tools/src/hwpod-harness-lib.ts @@ -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,