fix: reensure AgentRun runner for reused sessions
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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())));
|
||||
|
||||
Reference in New Issue
Block a user