Merge pull request #711 from pikasTech/fix-issue195-stale-continuation
fix: recover from stale thread continuation
This commit is contained in:
@@ -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));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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."
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user