diff --git a/internal/cloud/code-agent-agentrun-adapter.ts b/internal/cloud/code-agent-agentrun-adapter.ts index 937f37d0..87a40098 100644 --- a/internal/cloud/code-agent-agentrun-adapter.ts +++ b/internal/cloud/code-agent-agentrun-adapter.ts @@ -117,23 +117,50 @@ export async function submitAgentRunChatTurn({ params = {}, options = {}, traceI timeoutMs }); const commandId = requiredString(command?.id, "command.id"); - const mapping = agentRunReusedMapping({ previous: reusable.mapping, run: reusable.run, command, traceId, startedAt, backendProfile, managerUrl, env }); traceStore.append(traceId, { type: "backend", status: "running", label: "agentrun:command:created", - message: `AgentRun command ${commandId} created on reused run ${mapping.runId}; hwlab-cloud-api will not start a new runner Job for this turn.`, - runId: mapping.runId, + message: `AgentRun command ${commandId} created on reused run ${reusable.mapping.runId}; hwlab-cloud-api will ensure a runner Job is available for this turn.`, + runId: reusable.mapping.runId, commandId, backendProfile, - waitingFor: "agentrun-result", + waitingFor: "agentrun-runner-job-ensure", valuesPrinted: false }); + let runnerJob = null; + try { + const runnerJobInput = buildAgentRunRunnerJobInput({ env, traceId, commandId }); + runnerJob = await agentRunJson(fetchImpl, managerUrl, `/api/v1/runs/${encodeURIComponent(reusable.mapping.runId)}/runner-jobs`, { + method: "POST", + body: runnerJobInput, + timeoutMs + }); + } catch (error) { + if (!isAgentRunCommandAlreadyClaimed(error)) throw error; + traceStore.append(traceId, { + type: "backend", + status: "running", + label: "agentrun:runner-job:already-active", + message: `AgentRun command ${commandId} is already claimed by an active runner; no replacement runner Job is needed for this turn.`, + errorCode: error?.code ?? "agentrun_command_already_claimed", + runId: reusable.mapping.runId, + commandId, + runnerId: reusable.mapping.runnerId ?? null, + jobName: reusable.mapping.jobName ?? null, + namespace: reusable.mapping.namespace ?? null, + waitingFor: "agentrun-result", + valuesPrinted: false + }); + } + const mapping = agentRunReusedMapping({ previous: reusable.mapping, run: reusable.run, command, runnerJob, traceId, startedAt, backendProfile, managerUrl, env }); traceStore.append(traceId, { type: "backend", status: "running", - label: "agentrun:runner-job:reused", - message: `AgentRun runner Job ${mapping.jobName ?? "unknown"} is reused for this HWLAB session turn; no new bundle materialization is requested.`, + label: runnerJob ? "agentrun:runner-job:ensured" : "agentrun:runner-job:reused", + message: runnerJob + ? `AgentRun runner Job ${mapping.jobName ?? "unknown"} ensured for reused run ${mapping.runId}; this keeps the persistent session resumable after pod replacement.` + : `AgentRun runner Job ${mapping.jobName ?? "unknown"} is already active for this HWLAB session turn.`, runId: mapping.runId, commandId: mapping.commandId, attemptId: mapping.attemptId, @@ -231,6 +258,12 @@ export async function submitAgentRunChatTurn({ params = {}, options = {}, traceI return decorateAgentRunRunningResult({ base: initialAgentRunChatResult({ params, options, traceId }), mapping, traceStore, traceId }); } +function isAgentRunCommandAlreadyClaimed(error) { + const statusCode = Number(error?.statusCode ?? 0); + const message = String(error?.message ?? error?.agentRunError?.message ?? ""); + return statusCode === 409 && /command\s+[^\s]+\s+is not pending:/u.test(message); +} + export async function syncAgentRunChatResult({ traceId, currentResult = null, options = {}, traceStore = defaultCodeAgentTraceStore }) { const mapped = currentResult ?? await loadPersistedAgentRunResult(traceId, options); if (!mapped?.agentRun?.runId || !mapped?.agentRun?.commandId) return { result: currentResult, runnerTrace: traceStore.snapshot(traceId), found: Boolean(currentResult) }; @@ -1299,7 +1332,7 @@ function agentRunMapping({ env, managerUrl, backendProfile, run, command, runner }; } -function agentRunReusedMapping({ previous = {}, run = {}, command = {}, traceId, startedAt, backendProfile, managerUrl, env }) { +function agentRunReusedMapping({ previous = {}, run = {}, command = {}, runnerJob = null, traceId, startedAt, backendProfile, managerUrl, env }) { return { ...previous, adapter: ADAPTER_ID, @@ -1308,12 +1341,17 @@ function agentRunReusedMapping({ previous = {}, run = {}, command = {}, traceId, providerId: previous.providerId ?? firstNonEmpty(env.HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID, env.AGENTRUN_PROVIDER_ID, DEFAULT_PROVIDER_ID), runId: previous.runId ?? run?.id, commandId: command.id, - status: "runner-job-reused", + attemptId: runnerJob?.attemptId ?? runnerJob?.runner?.attemptId ?? previous.attemptId ?? null, + runnerId: runnerJob?.runnerId ?? runnerJob?.runner?.runnerId ?? previous.runnerId ?? null, + runnerJobId: runnerJob?.id ?? previous.runnerJobId ?? null, + jobName: runnerJob?.jobName ?? runnerJob?.jobIdentity?.name ?? previous.jobName ?? null, + namespace: runnerJob?.namespace ?? runnerJob?.jobIdentity?.namespace ?? previous.namespace ?? DEFAULT_RUNNER_NAMESPACE, + status: runnerJob ? "runner-job-ensured" : "runner-job-reused", runStatus: run?.status ?? previous.runStatus ?? null, commandState: command.state ?? null, terminalStatus: null, lastSeq: previous.lastSeq ?? 0, - runnerJobCount: 0, + runnerJobCount: runnerJob ? Math.max(1, Number(previous.runnerJobCount ?? 0) + 1) : Number(previous.runnerJobCount ?? 0), traceId, sessionId: run?.sessionRef?.sessionId ?? previous.sessionId ?? null, conversationId: run?.sessionRef?.conversationId ?? previous.conversationId ?? null, diff --git a/internal/cloud/server-agent-chat.test.ts b/internal/cloud/server-agent-chat.test.ts index de9a3e18..5eba3836 100644 --- a/internal/cloud/server-agent-chat.test.ts +++ b/internal/cloud/server-agent-chat.test.ts @@ -265,7 +265,8 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte return send({ id: secondTurn ? "cmd_hwlab_adapter_second" : "cmd_hwlab_adapter", runId: "run_hwlab_adapter", state: "pending", type: "turn", seq: secondTurn ? 2 : 1 }); } if (request.method === "POST" && url.pathname === "/api/v1/runs/run_hwlab_adapter/runner-jobs") { - assert.equal(body.commandId, "cmd_hwlab_adapter"); + const secondRunnerJob = body.commandId === "cmd_hwlab_adapter_second"; + assert.ok(body.commandId === "cmd_hwlab_adapter" || secondRunnerJob); assert.deepEqual(body.transientEnv.map((entry) => entry.name), [ "HWLAB_RUNTIME_API_URL", "HWLAB_RUNTIME_WEB_URL", @@ -293,13 +294,13 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte return send({ action: "create-kubernetes-job", runId: "run_hwlab_adapter", - commandId: "cmd_hwlab_adapter", - attemptId: "attempt_hwlab_adapter", - runnerId: "runner_hwlab_adapter", + commandId: body.commandId, + attemptId: secondRunnerJob ? "attempt_hwlab_adapter_second" : "attempt_hwlab_adapter", + runnerId: secondRunnerJob ? "runner_hwlab_adapter_second" : "runner_hwlab_adapter", namespace: "agentrun-v01", - jobName: "agentrun-v01-runner-hwlab-adapter", - jobIdentity: { namespace: "agentrun-v01", name: "agentrun-v01-runner-hwlab-adapter" }, - runner: { attemptId: "attempt_hwlab_adapter", runnerId: "runner_hwlab_adapter" } + jobName: secondRunnerJob ? "agentrun-v01-runner-hwlab-adapter-second" : "agentrun-v01-runner-hwlab-adapter", + jobIdentity: { namespace: "agentrun-v01", name: secondRunnerJob ? "agentrun-v01-runner-hwlab-adapter-second" : "agentrun-v01-runner-hwlab-adapter" }, + runner: { attemptId: secondRunnerJob ? "attempt_hwlab_adapter_second" : "attempt_hwlab_adapter", runnerId: secondRunnerJob ? "runner_hwlab_adapter_second" : "runner_hwlab_adapter" } }); } if (request.method === "GET" && url.pathname === "/api/v1/runs/run_hwlab_adapter/events") { @@ -307,7 +308,7 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte const second = afterSeq >= 5 || url.searchParams.get("afterSeq") === "3" || calls.some((call) => call.path === "/api/v1/runs/run_hwlab_adapter/commands/cmd_hwlab_adapter_second/result"); if (second) return send({ items: [ { id: "evt_old_tail", runId: "run_hwlab_adapter", seq: 6, type: "assistant_message", payload: { commandId: "cmd_hwlab_adapter", text: "旧 command 尾部不应进入第二轮。" }, createdAt: "2026-06-01T00:00:02.500Z" }, - { id: "evt_7", runId: "run_hwlab_adapter", seq: 7, type: "backend_status", payload: { phase: "turn-started", commandId: "cmd_hwlab_adapter_second", attemptId: "attempt_hwlab_adapter", jobName: "agentrun-v01-runner-hwlab-adapter", namespace: "agentrun-v01" }, createdAt: "2026-06-01T00:00:03.000Z" }, + { id: "evt_7", runId: "run_hwlab_adapter", seq: 7, type: "backend_status", payload: { phase: "turn-started", commandId: "cmd_hwlab_adapter_second", attemptId: "attempt_hwlab_adapter_second", runnerId: "runner_hwlab_adapter_second", jobName: "agentrun-v01-runner-hwlab-adapter-second", namespace: "agentrun-v01" }, createdAt: "2026-06-01T00:00:03.000Z" }, { id: "evt_7_prompt", runId: "run_hwlab_adapter", seq: 8, type: "backend_status", payload: { phase: "initial-prompt-assembly", commandId: "cmd_hwlab_adapter_second", initialPromptInjected: false, reason: "thread-resume", initialPrompt: { available: true, bytes: 512, sha256: "prompt-sha", promptRefCount: 1, skillRefCount: 2, valuesPrinted: false } }, createdAt: "2026-06-01T00:00:03.500Z" }, { id: "evt_8", runId: "run_hwlab_adapter", seq: 8, type: "assistant_message", payload: { commandId: "cmd_hwlab_adapter_second", text: "AgentRun adapter 复用已有 runner 完成第二轮。" }, createdAt: "2026-06-01T00:00:04.000Z" }, { id: "evt_9", runId: "run_hwlab_adapter", seq: 9, type: "terminal_status", payload: { commandId: "cmd_hwlab_adapter_second", terminalStatus: "completed" }, createdAt: "2026-06-01T00:00:05.000Z" } @@ -345,9 +346,9 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte return send({ runId: "run_hwlab_adapter", commandId: "cmd_hwlab_adapter_second", - attemptId: "attempt_hwlab_adapter", - runnerId: "runner_hwlab_adapter", - jobName: "agentrun-v01-runner-hwlab-adapter", + attemptId: "attempt_hwlab_adapter_second", + runnerId: "runner_hwlab_adapter_second", + jobName: "agentrun-v01-runner-hwlab-adapter-second", namespace: "agentrun-v01", status: "completed", runStatus: "claimed", @@ -567,15 +568,15 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte assert.equal(secondPayload.providerTrace.traceId, secondTraceId); assert.equal(secondPayload.agentRun.providerTrace.commandId, "cmd_hwlab_adapter_second"); assert.equal(secondPayload.agentRun.providerTrace.traceId, secondTraceId); - assert.equal(secondPayload.agentRun.jobName, "agentrun-v01-runner-hwlab-adapter"); + assert.equal(secondPayload.agentRun.jobName, "agentrun-v01-runner-hwlab-adapter-second"); assert.equal(secondPayload.sessionReuse.reused, true); - assert.equal(secondPayload.agentRun.runnerJobCount, 0); + assert.equal(secondPayload.agentRun.runnerJobCount, 1); assert.match(secondPayload.reply.content, /复用已有 runner/u); assert.ok(secondPayload.runnerTrace.events.some((event) => event.label === "agentrun:run:reused")); - assert.ok(secondPayload.runnerTrace.events.some((event) => event.label === "agentrun:runner-job:reused")); + assert.ok(secondPayload.runnerTrace.events.some((event) => event.label === "agentrun:runner-job:ensured")); assert.equal(secondPayload.runnerTrace.events.some((event) => String(event.text ?? event.message ?? "").includes("旧 command 尾部")), false); assert.equal(calls.filter((call) => call.method === "POST" && call.path === "/api/v1/runs").length, 1); - assert.equal(calls.filter((call) => call.method === "POST" && call.path === "/api/v1/runs/run_hwlab_adapter/runner-jobs").length, 1); + assert.equal(calls.filter((call) => call.method === "POST" && call.path === "/api/v1/runs/run_hwlab_adapter/runner-jobs").length, 2); assert.equal(calls.filter((call) => call.method === "POST" && call.path === "/api/v1/runs/run_hwlab_adapter/commands").length, 3); } finally { await new Promise((resolve, reject) => server.close((error) => (error ? reject(error) : resolve())));