diff --git a/internal/cloud/access-control.test.ts b/internal/cloud/access-control.test.ts index f9e81c37..24b89703 100644 --- a/internal/cloud/access-control.test.ts +++ b/internal/cloud/access-control.test.ts @@ -677,6 +677,186 @@ test("workbench workspace status clears completed AgentRun active trace on read" } }); +test("workbench workspace clears stale continuation after AgentRun thread resume failure", async () => { + const staleThreadId = "019e0000-0000-7000-8000-000000000195"; + const agentRunCalls = []; + const createRunInputs = []; + const agentSessions = new Map(); + const agentRunServer = createServer(async (request, response) => { + const url = new URL(request.url || "/", "http://127.0.0.1"); + agentRunCalls.push({ method: request.method, path: url.pathname, search: url.search }); + const send = (data) => { + response.writeHead(200, { "content-type": "application/json" }); + response.end(`${JSON.stringify({ ok: true, data })}\n`); + }; + if (request.method === "GET" && url.pathname === "/api/v1/runs/run_workspace_thread_resume_failed/events") { + return send({ items: [ + { id: "evt_thread_resume_failed", runId: "run_workspace_thread_resume_failed", seq: 1, type: "terminal_status", payload: { commandId: "cmd_workspace_thread_resume_failed", terminalStatus: "failed" }, createdAt: "2026-06-01T00:00:01.000Z" } + ] }); + } + if (request.method === "GET" && url.pathname === "/api/v1/runs/run_workspace_thread_resume_failed/commands/cmd_workspace_thread_resume_failed/result") { + return send({ + runId: "run_workspace_thread_resume_failed", + commandId: "cmd_workspace_thread_resume_failed", + status: "failed", + runStatus: "claimed", + commandState: "failed", + terminalStatus: "failed", + completed: false, + failureKind: "thread-resume-failed", + failureMessage: `codex app-server thread/resume failed for existing thread: no rollout found for thread id ${staleThreadId}`, + lastSeq: 1, + sessionRef: { sessionId: "ses_issue195_stale", conversationId: "cnv_issue195_stale", threadId: staleThreadId } + }); + } + if (request.method === "POST" && url.pathname === "/api/v1/runs") { + const body = await requestJson(request); + createRunInputs.push(body); + return send({ id: "run_workspace_after_stale", status: "pending", backendProfile: "deepseek", sessionRef: { sessionId: body.sessionRef?.sessionId } }); + } + if (request.method === "POST" && url.pathname === "/api/v1/runs/run_workspace_after_stale/commands") { + return send({ id: "cmd_workspace_after_stale", runId: "run_workspace_after_stale", state: "pending", type: "turn", seq: 1 }); + } + if (request.method === "POST" && url.pathname === "/api/v1/runs/run_workspace_after_stale/runner-jobs") { + return send({ + action: "create-kubernetes-job", + runId: "run_workspace_after_stale", + commandId: "cmd_workspace_after_stale", + attemptId: "attempt_workspace_after_stale", + runnerId: "runner_workspace_after_stale", + namespace: "agentrun-v01", + jobName: "agentrun-v01-runner-workspace-after-stale" + }); + } + response.writeHead(404, { "content-type": "application/json" }); + response.end(`${JSON.stringify({ ok: false, message: `unexpected ${request.method} ${url.pathname}` })}\n`); + }); + await new Promise((resolve) => agentRunServer.listen(0, "127.0.0.1", resolve)); + + const agentRunPort = agentRunServer.address().port; + const accessController = createAccessController({ + env: { + HWLAB_ACCESS_CONTROL_REQUIRED: "1", + HWLAB_BOOTSTRAP_ADMIN_USERNAME: "admin", + HWLAB_BOOTSTRAP_ADMIN_PASSWORD: "admin-pass" + }, + now: () => "2026-06-01T00:00:00.000Z" + }); + accessController.getAgentSessionByTraceId = async (traceId) => agentSessions.get(traceId) ?? null; + const server = createCloudApiServer({ + env: { + HWLAB_ACCESS_CONTROL_REQUIRED: "1", + HWLAB_BOOTSTRAP_ADMIN_USERNAME: "admin", + HWLAB_BOOTSTRAP_ADMIN_PASSWORD: "admin-pass", + HWLAB_CODE_AGENT_ADAPTER: "agentrun-v01", + AGENTRUN_MGR_URL: `http://127.0.0.1:${agentRunPort}`, + HWLAB_CODE_AGENT_AGENTRUN_ALLOW_NON_K3S_URL: "1", + HWLAB_CODE_AGENT_AGENTRUN_SOURCE_COMMIT: "0123456789abcdef0123456789abcdef01234567", + HWLAB_CODE_AGENT_DEFAULT_PROVIDER_PROFILE: "deepseek" + }, + accessController, + now: () => "2026-06-01T00:00:00.000Z" + }); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + + try { + const { port } = server.address(); + const adminLogin = await postJson(port, "/auth/login", { username: "admin", password: "admin-pass" }); + const alice = await postJson(port, "/v1/admin/users", { username: "alice-ws-thread-resume", password: "alice-pass" }, adminLogin.cookie); + assert.equal(alice.status, 201); + const aliceLogin = await postJson(port, "/auth/login", { username: "alice-ws-thread-resume", password: "alice-pass" }); + const workspace = await getJson(port, "/v1/workbench/workspace?projectId=prj_device_pod_workbench", aliceLogin.cookie); + assert.equal(workspace.status, 200); + const conversation = await putJson(port, "/v1/agent/conversations/cnv_issue195_stale", { + projectId: "prj_device_pod_workbench", + sessionId: "ses_issue195_stale", + threadId: staleThreadId, + sessionStatus: "active", + lastTraceId: "trc_issue195_thread_resume_failed", + messages: [{ role: "agent", text: "stale continuation still running", status: "running", traceId: "trc_issue195_thread_resume_failed" }] + }, aliceLogin.cookie); + assert.equal(conversation.status, 200); + + agentSessions.set("trc_issue195_thread_resume_failed", { + id: "ses_issue195_stale", + ownerUserId: alice.body.user.id, + conversationId: "cnv_issue195_stale", + threadId: staleThreadId, + status: "active", + session: { + messageId: "msg_issue195_thread_resume_failed", + agentRun: { + runId: "run_workspace_thread_resume_failed", + commandId: "cmd_workspace_thread_resume_failed", + backendProfile: "deepseek", + managerUrl: `http://127.0.0.1:${agentRunPort}`, + lastSeq: 0, + sessionId: "ses_issue195_stale", + conversationId: "cnv_issue195_stale", + threadId: staleThreadId + } + }, + updatedAt: "2026-06-01T00:00:00.000Z" + }); + + const update = await patchJson(port, `/v1/workbench/workspace/${workspace.body.workspace.workspaceId}`, { + expectedRevision: 1, + selectedConversationId: "cnv_issue195_stale", + selectedAgentSessionId: "ses_issue195_stale", + activeTraceId: "trc_issue195_thread_resume_failed", + providerProfile: "deepseek", + sessionStatus: "running", + updatedByClient: "test-suite" + }, aliceLogin.cookie); + assert.equal(update.status, 200); + assert.equal(update.body.workspace.selectedConversationId, "cnv_issue195_stale"); + assert.equal(update.body.workspace.activeTraceId, "trc_issue195_thread_resume_failed"); + + const restored = await getJson(port, "/v1/workbench/workspace?projectId=prj_device_pod_workbench", aliceLogin.cookie); + assert.equal(restored.status, 200); + assert.equal(restored.body.workspace.activeTraceId, null); + assert.equal(restored.body.workspace.selectedConversationId, null); + assert.equal(restored.body.workspace.selectedAgentSessionId, null); + assert.equal(restored.body.workspace.selectedConversation, null); + assert.equal(restored.body.workspace.workspace.sessionStatus, "failed"); + assert.equal(restored.body.workspace.workspace.lastTraceId, "trc_issue195_thread_resume_failed"); + assert.equal(restored.body.workspace.workspace.staleContinuationCleared, true); + assert.equal(restored.body.workspace.workspace.staleContinuationReason, "thread-resume-failed"); + assert.equal(restored.body.workspace.workspace.staleConversationId, "cnv_issue195_stale"); + assert.equal(restored.body.workspace.workspace.staleAgentSessionId, "ses_issue195_stale"); + assert.equal(restored.body.workspace.workspace.staleThreadId, staleThreadId); + assert.equal(restored.body.workspace.workspace.recoveryAction, "new-session-on-next-turn"); + + const failedConversation = await getJson(port, "/v1/agent/conversations/cnv_issue195_stale", aliceLogin.cookie); + assert.equal(failedConversation.status, 200); + assert.equal(failedConversation.body.conversation.status, "failed"); + assert.equal(failedConversation.body.conversation.threadId, staleThreadId); + + const next = await postJson(port, "/v1/agent/chat", { + message: "new turn after stale continuation was cleared", + traceId: "trc_issue195_after_stale", + workspaceId: workspace.body.workspace.workspaceId, + expectedWorkspaceRevision: restored.body.workspace.revision, + shortConnection: true + }, aliceLogin.cookie, { prefer: "respond-async", "x-trace-id": "trc_issue195_after_stale" }); + assert.equal(next.status, 202); + assert.equal(next.body.status, "running"); + assert.equal(next.body.traceId, "trc_issue195_after_stale"); + await waitForCondition(() => createRunInputs.length === 1); + assert.equal(createRunInputs.length, 1); + assert.notEqual(createRunInputs[0].sessionRef?.threadId, staleThreadId); + assert.equal(createRunInputs[0].sessionRef?.threadId, undefined); + assert.equal(createRunInputs[0].sessionRef?.conversationId, undefined); + + const afterNext = await getJson(port, "/v1/workbench/workspace?projectId=prj_device_pod_workbench", aliceLogin.cookie); + assert.equal(afterNext.body.workspace.activeTraceId, "trc_issue195_after_stale"); + assert.ok(agentRunCalls.some((call) => call.path === "/api/v1/runs/run_workspace_thread_resume_failed/commands/cmd_workspace_thread_resume_failed/result")); + } finally { + await new Promise((resolve, reject) => server.close((error) => (error ? reject(error) : resolve()))); + await new Promise((resolve, reject) => agentRunServer.close((error) => (error ? reject(error) : resolve()))); + } +}); + test("workbench workspace status repairs terminal selected conversation after active trace was cleared", async () => { const agentRunCalls = []; const agentSessions = new Map(); @@ -1928,3 +2108,11 @@ async function requestJson(request) { const text = Buffer.concat(chunks).toString("utf8").trim(); return text ? JSON.parse(text) : {}; } + +async function waitForCondition(predicate, timeoutMs = 500) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } +} diff --git a/internal/cloud/access-control.ts b/internal/cloud/access-control.ts index 99e6cf91..94e362ed 100644 --- a/internal/cloud/access-control.ts +++ b/internal/cloud/access-control.ts @@ -862,17 +862,21 @@ class AccessController { if (!isTerminalCodeAgentResult(result)) return workspace; const now = this.now(); - const selectedConversationId = safeConversationIdLocal(result.conversationId) ? result.conversationId : workspace.selectedConversationId; - const selectedAgentSessionId = safeAgentSessionId(result.sessionId ?? result.session?.sessionId ?? result.sessionReuse?.sessionId) || workspace.selectedAgentSessionId; + const staleContinuation = isThreadResumeFailedResult(result); + const resultConversationId = safeConversationIdLocal(result.conversationId) ? result.conversationId : workspace.selectedConversationId; + const resultAgentSessionId = safeAgentSessionId(result.sessionId ?? result.session?.sessionId ?? result.sessionReuse?.sessionId) || workspace.selectedAgentSessionId; + const selectedConversationId = staleContinuation ? null : resultConversationId; + const selectedAgentSessionId = staleContinuation ? null : resultAgentSessionId; await this.syncTerminalWorkbenchConversation({ workspace, actor, result, activeTraceId: traceId, - selectedConversationId, - selectedAgentSessionId, + selectedConversationId: resultConversationId, + selectedAgentSessionId: resultAgentSessionId, now }); + const staleThreadId = threadResumeFailureThreadId(result); const updated = await this.store.updateWorkspace?.({ workspaceId: workspace.id, ownerUserId: workspace.ownerUserId, @@ -890,6 +894,14 @@ class AccessController { sessionStatus: terminalWorkbenchSessionStatus(result), lastTraceId: traceId, syncedTraceStatus: result.status ?? result.agentRun?.terminalStatus ?? null, + ...(staleContinuation ? { + staleContinuationCleared: true, + staleContinuationReason: "thread-resume-failed", + staleConversationId: resultConversationId, + staleAgentSessionId: resultAgentSessionId, + staleThreadId, + recoveryAction: "new-session-on-next-turn" + } : {}), updatedAt: now, source: "workbench-status-sync", valuesRedacted: true, @@ -1775,6 +1787,12 @@ function redactedWorkspaceJson(value = {}) { providerProfile: textOr(workspace.providerProfile, ""), sessionStatus: textOr(workspace.sessionStatus, ""), lastTraceId: textOr(workspace.lastTraceId, ""), + staleContinuationCleared: workspace.staleContinuationCleared === true, + staleContinuationReason: textOr(workspace.staleContinuationReason, ""), + staleConversationId: textOr(workspace.staleConversationId, ""), + staleAgentSessionId: textOr(workspace.staleAgentSessionId, ""), + staleThreadId: textOr(workspace.staleThreadId, ""), + recoveryAction: textOr(workspace.recoveryAction, ""), messages: Array.isArray(workspace.messages) ? workspace.messages.slice(-50).map(redactConversationMessage).filter(Boolean) : undefined, actor: workspace.actor && typeof workspace.actor === "object" ? publicActor(workspace.actor) : undefined, updatedAt: textOr(workspace.updatedAt, ""), @@ -1804,6 +1822,29 @@ function isTerminalCodeAgentResult(result = {}) { const terminalStatus = textOr(result.agentRun?.terminalStatus, "").toLowerCase(); return [status, terminalStatus].some((value) => ["completed", "failed", "blocked", "canceled", "cancelled", "timeout", "error"].includes(value)); } +function isThreadResumeFailedResult(result = {}) { + const values = [ + result.error?.code, + result.error?.category, + result.blocker?.code, + result.blocker?.category, + result.providerTrace?.failureKind, + result.agentRun?.providerTrace?.failureKind, + result.agentRun?.failureKind + ].map((value) => textOr(value, "").toLowerCase().replace(/_/gu, "-")); + return values.some((value) => value === "thread-resume-failed") + || /no rollout found for thread id|thread\/resume failed/iu.test(String(result.error?.message ?? result.blocker?.message ?? result.providerTrace?.failureMessage ?? result.agentRun?.providerTrace?.failureMessage ?? "")); +} +function threadResumeFailureThreadId(result = {}) { + return boundedText(textOr( + result.providerTrace?.threadId + ?? result.agentRun?.providerTrace?.threadId + ?? result.session?.threadId + ?? result.sessionReuse?.threadId + ?? result.threadId, + "" + ), 240) || null; +} function terminalWorkbenchSessionStatus(result = {}) { const sessionStatus = textOr(result.session?.status ?? result.sessionSummary?.status ?? result.sessionLifecycleStatus ?? result.runnerTrace?.sessionStatus, "").toLowerCase(); if (sessionStatus && sessionStatus !== "running" && sessionStatus !== "busy" && sessionStatus !== "pending") return sessionStatus === "cancelled" ? "canceled" : sessionStatus; diff --git a/internal/cloud/code-agent-agentrun-adapter.ts b/internal/cloud/code-agent-agentrun-adapter.ts index d6c30d67..78ab76cc 100644 --- a/internal/cloud/code-agent-agentrun-adapter.ts +++ b/internal/cloud/code-agent-agentrun-adapter.ts @@ -795,7 +795,7 @@ function agentRunFailureAttribution({ code, message, canceled = false } = {}) { return { category: "thread_resume_failed", retryable: true, - userMessage: "AgentRun 复用的 Codex thread 已失效;当前标准是终止本轮并保留原 session 指针,不启动新的 thread/start,也不拼接历史上下文。请重新发起新会话。", + userMessage: "AgentRun 复用的 Codex thread 已失效;当前 turn 会按失败终止,Workbench 会清理 stale continuation 指针,下一轮自动从新 session/thread 启动。", summary: message || "AgentRun thread/resume failed for an existing thread." }; }