Merge pull request #1006 from pikasTech/fix/issue1004-steer-short

fix: keep steer short connection running
This commit is contained in:
Lyon
2026-06-06 16:53:20 +08:00
committed by GitHub
6 changed files with 54 additions and 3 deletions
+18 -1
View File
@@ -462,6 +462,9 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte
const agentRunPort = agentRunServer.address().port;
const traceStore = createCodeAgentTraceStore();
let deferOwnerRecord = false;
let releaseOwnerRecord: (() => void) | null = null;
let deferredOwnerRecordStarted = false;
const server = createCloudApiServer({
traceStore,
env: {
@@ -481,6 +484,10 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte
return { ok: true, actor: TEST_AGENT_ACTOR, session: TEST_AUTH_SESSION };
},
async recordAgentSessionOwner(input) {
if (deferOwnerRecord && input.traceId === "trc_server-test-agentrun-adapter" && input.status === "running") {
deferredOwnerRecordStarted = true;
await new Promise<void>((resolve) => { releaseOwnerRecord = resolve; });
}
const record = testAgentSessionRecord(input);
ownerSessions.set(record.id, record);
return record;
@@ -619,7 +626,8 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte
}
);
const steer = await fetch(`http://127.0.0.1:${port}/v1/agent/chat/steer`, {
deferOwnerRecord = true;
const steerRequest = fetch(`http://127.0.0.1:${port}/v1/agent/chat/steer`, {
method: "POST",
headers: { "content-type": "application/json", "x-trace-id": traceId, cookie: "hwlab_session=test-stub-session" },
body: JSON.stringify({
@@ -632,7 +640,16 @@ test("cloud api /v1/agent/chat delegates v0.2 turns to AgentRun v0.1 over adapte
message: "请按 STEER_MARK 调整最终回复"
})
});
const steer = await Promise.race([
steerRequest,
delay(100).then(() => null)
]);
assert.ok(steer, "steer short response must not wait for owner persistence");
assert.equal(steer.status, 202);
for (let index = 0; index < 20 && !deferredOwnerRecordStarted; index += 1) await delay(10);
assert.equal(deferredOwnerRecordStarted, true);
releaseOwnerRecord?.();
deferOwnerRecord = false;
const steerBody = await steer.json();
assert.equal(steerBody.accepted, true);
assert.equal(steerBody.route, "/v1/agent/chat/steer");
+3 -1
View File
@@ -1336,8 +1336,10 @@ export async function handleCodeAgentSteerHttp(request, response, options) {
options,
traceStore
});
await recordCodeAgentSessionOwner({ payload: currentResult, params: { ...params, traceId, ownerUserId: options.actor?.id, ownerRole: options.actor?.role }, options, status: "running" });
sendJson(response, 202, payload);
setImmediate(() => {
void recordCodeAgentSessionOwner({ payload: currentResult, params: { ...params, traceId, ownerUserId: options.actor?.id, ownerRole: options.actor?.role }, options, status: "running" });
});
} catch (error) {
traceStore.append(traceId, {
type: "backend",
@@ -23,9 +23,13 @@ test("issue 853 persists submitted user and running agent messages before long C
test("issue 853 submit failure persists the same trace before refresh", () => {
const source = fs.readFileSync(path.join(srcRoot, "state/workbench.ts"), "utf8");
const failureIndex = source.indexOf('const traceSnapshot = await fetchTraceSnapshot(traceId, traceId, Math.max(8000, state.codeAgentTimeoutMs), undefined, activeProjectId);');
const keepSteerIndex = source.indexOf('const keepSteerActive = steerMode && response.status === 0', failureIndex);
const steerDoneIndex = source.indexOf('dispatch({ type: "chat:steer-done" });', keepSteerIndex);
const failedWithTraceIndex = source.indexOf('const failedWithTrace = runnerTrace ? { ...failed, runnerTrace } : failed;', failureIndex);
const persistIndex = source.indexOf('messages: replaceMessage(visibleMessages, activePending.id, failedWithTrace)', failedWithTraceIndex);
assert.ok(failureIndex > 0);
assert.ok(keepSteerIndex > failureIndex);
assert.ok(steerDoneIndex > keepSteerIndex);
assert.ok(failedWithTraceIndex > failureIndex);
assert.ok(persistIndex > failedWithTraceIndex);
});
@@ -139,6 +139,23 @@ test("conversation selection restores steer mode when switching back to a runnin
assert.equal(composer.targetTraceId, "trc_active");
});
test("steer submit completion keeps the target turn steerable", () => {
const state = baseState({
workspace: workspace("cnv_active", "trc_active"),
currentRequest: { traceId: "trc_active", conversationId: "cnv_active", sessionId: "ses_active", threadId: "thread_active", status: "running" },
messages: [agentMessage({ conversationId: "cnv_active", traceId: "trc_active", status: "running" })],
chatPending: true
});
const next = workbenchReducer(state, { type: "chat:steer-done" });
const composer = composerFromState(next, "cnv_active");
assert.equal(next.chatPending, false);
assert.equal(next.currentRequest?.status, "running");
assert.equal(composer.submitMode, "steer");
assert.equal(composer.targetTraceId, "trc_active");
});
test("composer routes to steer from workspace active status even before running message is restored", () => {
const state = baseState({
workspace: workspace("cnv_active", "trc_active"),
@@ -19,6 +19,7 @@ export type Action =
| { type: "message:trace"; messageId: string; trace: NonNullable<ChatMessage["runnerTrace"]> }
| { type: "message:complete"; messageId: string; message: ChatMessage; availability: CodeAgentAvailability | null; workspace?: WorkbenchState["workspace"] }
| { type: "message:fail"; messageId: string; message: ChatMessage }
| { type: "chat:steer-done" }
| { type: "chat:done" }
| { type: "conversation:select"; conversation: ConversationRecord | null; workspace?: WorkbenchState["workspace"] }
| { type: "conversation:list"; conversations: ConversationRecord[] }
@@ -47,6 +48,7 @@ export function workbenchReducer(state: WorkbenchState, action: Action): Workben
}
case "message:complete": return { ...state, messages: replaceMessage(state.messages, action.messageId, action.message), codeAgentAvailability: action.availability ?? state.codeAgentAvailability, workspace: action.workspace ?? state.workspace };
case "message:fail": return { ...state, messages: replaceMessage(state.messages, action.messageId, action.message) };
case "chat:steer-done": return { ...state, chatPending: false };
case "chat:done": return { ...state, chatPending: false, currentRequest: state.currentRequest ? { ...state.currentRequest, status: "completed" } : null };
case "conversation:select": {
const messages = action.conversation?.messages ?? [];
+10 -1
View File
@@ -307,8 +307,17 @@ export function useWorkbenchStore(enabled: boolean, projectId = WORKBENCH_PROJEC
}
} else {
const traceSnapshot = await fetchTraceSnapshot(traceId, traceId, Math.max(8000, state.codeAgentTimeoutMs), undefined, activeProjectId);
const failed = makeMessage("agent", response.error ?? "Code Agent 请求失败", "failed", { traceId, conversationId, sessionId, threadId, title: "Code Agent 请求失败" });
const runnerTrace = traceSnapshot ? snapshotToRunnerTrace(traceSnapshot) : null;
const keepSteerActive = steerMode && response.status === 0 && !isTerminalStatus(activePending.status) && (!runnerTrace || !isTerminalStatus(runnerTrace.status));
if (keepSteerActive) {
const runningWithTrace = runnerTrace ? { ...activePending, runnerTrace } : activePending;
if (runnerTrace) dispatch({ type: "message:trace", messageId: activePending.id, trace: runnerTrace });
currentWorkspace = await persistConversation({ workspace: currentWorkspace ?? state.workspace, projectId: activeProjectId, conversationId, sessionId, threadId, messages: replaceMessage(visibleMessages, activePending.id, runningWithTrace) });
if (currentWorkspace) dispatch({ type: "workspace:sync", workspace: currentWorkspace });
dispatch({ type: "chat:steer-done" });
return;
}
const failed = makeMessage("agent", response.error ?? "Code Agent 请求失败", "failed", { traceId, conversationId, sessionId, threadId, title: "Code Agent 请求失败" });
const failedWithTrace = runnerTrace ? { ...failed, runnerTrace } : failed;
if (runnerTrace) dispatch({ type: "message:trace", messageId: activePending.id, trace: runnerTrace });
dispatch({ type: "message:fail", messageId: activePending.id, message: failedWithTrace });