From 2faa6c496e27c37d21a27c388a8fe8ec831a4f58 Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Thu, 1 Oct 2026 03:04:14 +0800 Subject: [PATCH] fix(chat): steer active Codex turns from the shared composer Signed-off-by: huangruiteng --- apps/presentation/dashboard/src/data/chat.ts | 5 + .../personal-workspace-page.tsx | 78 ++++++++++---- .../dashboard/src/views/dashboard-page.tsx | 17 ++- .../app-conversation-and-async-inbox-v0.md | 10 ++ .../attached-host-follow-up.mjs | 18 ++-- .../chat-recovery.mjs | 12 ++- .../composer-session-admission.mjs | 101 +++++++++++++++--- .../personal-workspace-browser/fixture.mjs | 4 +- .../steward-journey.mjs | 3 + 9 files changed, 195 insertions(+), 53 deletions(-) diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 691de82ec7..3fb3f14bc5 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -703,6 +703,11 @@ export function chatSessionQueuesFollowUps(session: Pick) { + return session.session_mode !== "attached_host" && session.adapter_kind === "codex_app_server"; +} + export type ManagerRuntimeSessionReadback = { schema_version: "manager_runtime_session_readback_v0"; runtime_profile: "restricted" | "trusted_owner"; diff --git a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx index 1b65c66819..116d07c326 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace-page.tsx @@ -750,6 +750,7 @@ function readImageAttachment(file: File, t: WorkspaceTranslate): Promise()); const [actionDraft, setActionDraft] = useState(null); const [loopxMode, setLoopxMode] = useState(null); const [loopxDelivery, setLoopxDelivery] = useState<"queue" | "inbox" | "steer">("queue"); @@ -848,8 +853,9 @@ export function PersonalWorkspacePage({ input.style.height = "auto"; input.style.height = `${Math.min(input.scrollHeight, 120)}px`; }, [composer, selectedGoalId, managerChatOpen]); - function setComposerDraft(key: string, value: string) { + function setComposerDraft(key: string, value: string, expectedValue?: string) { setDrafts((current) => { + if (expectedValue !== undefined && current[key] !== expectedValue) return current; const next = { ...current }; if (value) { next[key] = value; @@ -1027,18 +1033,17 @@ export function PersonalWorkspacePage({ setGoalConversationReceiptVisible(true); } }, [goalMessages, selectedGoal, selectedGoalTab]); - // A managed runtime Session admits one Turn at a time. While the current - // Session shows a Turn in flight, a new message would only be rejected, so - // the composer waits and points to the reply's own adjust/interrupt - // controls. Two deliveries stay open because the service queues them behind - // the running Turn: LoopX mode through its own queue, and any message to an - // attached host Session. + // One composer for both conversations. Running managed Codex work receives + // exact-turn instructions; attached hosts and LoopX mode keep their queues. const loopxDeliveryOpen = Boolean(conversationSessionId && loopxMode?.session_id === conversationSessionId && loopxMode?.enabled && loopxMode.active_turn_id); - const conversationTurnRunning = !loopxDeliveryOpen && !conversationQueuesFollowUps && Boolean(conversationSessionId) - && managerMessages.some((message) => message.pending && Boolean(message.sourceTurnId) - && message.sourceSessionId === conversationSessionId); - const composerBlocked = sending || conversationTurnRunning; + const runningMessage = managerMessages.find((message) => message.pending && Boolean(message.sourceTurnId) + && message.sourceSessionId === conversationSessionId); + const conversationTurnRunning = !loopxDeliveryOpen && !conversationQueuesFollowUps && Boolean(runningMessage); + const steeringTurnId = conversationTurnRunning && conversationSupportsSteering && !readOnly + && callbacks.onSteerConversationTurn ? runningMessage?.sourceTurnId : undefined; + const composerBlocked = steering || (!steeringTurnId && (sending || conversationTurnRunning)); + const quickPromptBlocked = steering || sending || conversationTurnRunning; const managerChatItems = useMemo( () => items.filter((item) => item.kind === "message" || (item.kind === "proposal" && (sessionProposalIds.includes(item.proposal.previewId) @@ -1626,7 +1631,36 @@ export function PersonalWorkspacePage({ async function sendMessage(messageOverride?: string) { const pendingImages = messageOverride ? [] : imageAttachments; const message = (messageOverride ?? composer).trim() || (pendingImages.length ? t("composer.imageAnalysisPrompt") : ""); - if (!message || sending || conversationHistoryState?.sendBlocked) return; + if (!message || composerBlocked || conversationHistoryState?.sendBlocked) return; + const previousSteering = steeringRequests.current.get(composerDraftKey); + const retry = previousSteering && previousSteering.sessionId === conversationSessionId && previousSteering.text === message + ? previousSteering : undefined; + if ((retry || steeringTurnId) && conversationSessionId && callbacks.onSteerConversationTurn) { + if (pendingImages.length) { + setImageAttachmentError(locale === "zh-CN" ? "本轮追加指令暂不支持图片,图片和草稿已保留。" : "This turn accepts text instructions only. Images and draft retained."); + return; + } + const request = retry ?? { sessionId: conversationSessionId, turnId: steeringTurnId!, text: message, id: crypto.randomUUID() }; + steeringRequests.current.set(composerDraftKey, request); + setSteering(true); + setActionFeedback(null); + setImageAttachmentError(null); + try { + await callbacks.onSteerConversationTurn(selectedGoalId ?? "manager", request.turnId, message, request.id); + steeringRequests.current.delete(composerDraftKey); + if (!messageOverride) setComposerDraft(composerDraftKey, "", composer); + setActionFeedback(locale === "zh-CN" ? "执行器已接收本轮追加指令。" : "The executor accepted instructions for this turn."); + } catch (error) { + // Unknown delivery retries the original Turn even after it completes. + // A confirmed non-delivery may use a new ingress after recovery. + if (error instanceof ChatApiError && error.payload.delivery_state === "not_delivered") { + steeringRequests.current.delete(composerDraftKey); + } + setActionFeedback(error instanceof Error ? error.message : t("feedback.sendGenericError")); + } finally { setSteering(false); } + return; + } + if (sending) return; followConversationRef.current = true; setShowLatestMessage(false); if (loopxMode?.session_id === conversationSessionId && loopxMode?.enabled && loopxMode.active_turn_id && conversationSessionId) { @@ -1969,16 +2003,16 @@ export function PersonalWorkspacePage({ {locale === "zh-CN" ? "快捷提问" : "Suggestions"} {selectedGoal ? (
- - + + - - + +
) : (
- - + +
)} @@ -1991,7 +2025,9 @@ export function PersonalWorkspacePage({ ))} : null} {imageAttachmentError ?

{imageAttachmentError}

: null} - {conversationTurnRunning ?

{t("composer.turnRunning")}

: null} + {conversationTurnRunning ?

{steeringTurnId + ? (locale === "zh-CN" ? "本轮进行中 · 发消息可调整当前工作" : "Turn in progress · send instructions to adjust this work") + : t("composer.turnRunning")}

: null}
{ @@ -2034,7 +2070,9 @@ export function PersonalWorkspacePage({ />
- {conversationOpen ?
{sending + {conversationOpen ?
{steering + ? (locale === "zh-CN" ? "正在发送本轮追加指令…" : "Sending instructions for this turn…") + : sending && !steeringTurnId ? (locale === "zh-CN" ? "正在回复 · 修改当前任务请使用“调整本轮”" : "Reply in progress · use Adjust turn to change the current task") : (locale === "zh-CN" ? "Enter 发送 · Shift+Enter 换行" : "Enter to send · Shift+Enter for a new line")}
: null} } diff --git a/apps/presentation/dashboard/src/views/dashboard-page.tsx b/apps/presentation/dashboard/src/views/dashboard-page.tsx index 258c208df4..8481280032 100644 --- a/apps/presentation/dashboard/src/views/dashboard-page.tsx +++ b/apps/presentation/dashboard/src/views/dashboard-page.tsx @@ -56,6 +56,7 @@ import { resumeChatTurnStreaming, sendChatTurnStreaming, chatSessionQueuesFollowUps, + chatSessionSupportsSteering, selectAvailableChatAgent, sessionInvalidatedByPayload, todoNoWriteReceiptFromPayload, @@ -1410,6 +1411,7 @@ function PersonalGoalHome({ // Bound Sessions whose mode queues a message sent while a Turn runs, read // from the Session owner each time this page binds a Session. const [followUpQueueSessionIds, setFollowUpQueueSessionIds] = useState>(() => new Set()); + const [steeringSessionIds, setSteeringSessionIds] = useState>(() => new Set()); const [executionSessions, setExecutionSessions] = useState([]); // Bumped when the service reports a running Turn this page did not know // about, so the Turn recovery effect re-reads the Session and adopts it. @@ -1559,6 +1561,14 @@ function PersonalGoalHome({ function recordSessionAdmission(session: ChatSessionSummary) { const queues = chatSessionQueuesFollowUps(session); + const supportsSteering = chatSessionSupportsSteering(session); + setSteeringSessionIds((current) => { + if (current.has(session.session_id) === supportsSteering) return current; + const next = new Set(current); + if (supportsSteering) next.add(session.session_id); + else next.delete(session.session_id); + return next; + }); setFollowUpQueueSessionIds((current) => { if (current.has(session.session_id) === queues) return current; const next = new Set(current); @@ -2939,9 +2949,9 @@ function PersonalGoalHome({ }, onSteerConversationTurn: async (targetContextId, turnId, message, ingressId) => { const binding = runtimeBindings[targetContextId]; - if (!binding?.sessionId || binding.turnId !== turnId || activeTurnIds.current.get(targetContextId) !== turnId) { - throw new Error("本轮已结束或已被新的回合取代,追加指令未发送,草稿已保留。"); - } + if (!binding?.sessionId) throw new Error("当前会话不可用,追加指令未发送,草稿已保留。"); + // The service owns exact-turn admission and durable retry. A delivered + // ingress may be read back after completion; never retarget it locally. await steerChatTurn(binding.sessionId, turnId, message, ingressId); const id = managerMessageId.current++; setMessagesByContext(current => { @@ -3081,6 +3091,7 @@ function PersonalGoalHome({ managerRuntime={managerRuntime} conversationSessionId={runtimeBindings[contextId]?.sessionId} conversationQueuesFollowUps={followUpQueueSessionIds.has(runtimeBindings[contextId]?.sessionId ?? "")} + conversationSupportsSteering={steeringSessionIds.has(runtimeBindings[contextId]?.sessionId ?? "")} conversationHistoryState={conversationHistory} model={workspaceModel} readOnly={readOnly} diff --git a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md index 32acdb3ea5..07bcc83a5a 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -126,6 +126,16 @@ receipt → observed work or actionable failure → readable answer in the same - Before dispatch, cancel only session preparation and state that the request was not submitted. After acceptance, existing exact-turn steering/interrupt controls own effects; stopping observation is not stopping the worker. +- During a managed Codex Turn, the ordinary composer sends text instructions to + that exact Turn, without a second adjustment form or another Turn submission. + Attached-host messages keep their next-Turn queue semantics; unsupported + managed adapters keep the draft without advertising native steering. The + idle composer and initial presentation are unchanged. This changes the former + managed-running Send lockout; LoopX mode retains its explicit delivery choice. + A lost or mismatched receipt preserves draft and ingress identity, including + retry after completion. Only a confirmed non-delivery permits a new ingress. + Acceptance means the executor received the instructions, not that it adopted + them or that delegated/team work stopped. Live adoption stays a release gate. - The compact receipt and full conversation offer the same controls. Failure ends the live indicator, preserves the request/partial answer and names the next supported action. A completed delegation still shows receiver adoption diff --git a/examples/personal-workspace-browser/attached-host-follow-up.mjs b/examples/personal-workspace-browser/attached-host-follow-up.mjs index 2da1f536fb..039d93d743 100644 --- a/examples/personal-workspace-browser/attached-host-follow-up.mjs +++ b/examples/personal-workspace-browser/attached-host-follow-up.mjs @@ -2,7 +2,7 @@ import { openWorkspacePage } from "./scenario-context.mjs"; // An attached_host Session queues a follow-up while its host runs a Turn // (ChatRuntimeController.submit_turn -> enqueue_attached_agent_turn), whereas a -// managed_runtime Session accepts one Turn at a time. The composer must follow +// managed_runtime Session accepts one Turn at a time with native steering when supported. The composer must follow // the Session's typed mode, not read every running Turn as a rejected send. const attachedGoal = { id: "product-release", label: "Product Release" }; const managedGoal = { id: "research-monitor", label: "Research Monitor" }; @@ -11,7 +11,7 @@ const sessionIdFor = (goalId) => `session-goal-${goalId}-codex`; function runningSession(goalId, mode) { return { - session_id: sessionIdFor(goalId), goal_id: goalId, agent_id: "codex", adapter_kind: "codex", + session_id: sessionIdFor(goalId), goal_id: goalId, agent_id: "codex", adapter_kind: "codex_app_server", channel_id: `goal.${goalId}`, status: "busy", active_turn_id: runningTurnId(goalId), last_error_code: null, created_at: "2026-08-13T01:00:00Z", updated_at: "2026-08-13T01:00:01Z", last_activity_at: "2026-08-13T01:00:01Z", resumable: true, session_mode: mode, host_surface: mode === "attached_host" ? "codex_app" : null, @@ -72,12 +72,12 @@ export const attachedHostFollowUpScenario = { }; try { await openGoalChat(managedGoal); - await turnRunningHint.waitFor({ state: "visible", timeout: 5_000 }); - await composer.fill("托管会话回合进行中不应发送"); - if (!(await sendButton.isDisabled())) throw new Error("A managed_runtime composer stayed sendable while its Turn was running"); - await sendButton.click({ force: true }); - await page.waitForTimeout(300); - if (posts.length) throw new Error(`A managed_runtime Session was sent ${posts.length} message(s) into its running Turn`); + await page.locator(".personal-composer-status", { hasText: "发消息可调整当前工作" }).waitFor(); + await composer.fill("当前工作先检查依赖"); + if (await sendButton.isDisabled()) throw new Error("Managed Codex instructions were blocked"); + // This scenario isolates attached queues; exact native steering and retry + // are checked by composer-session-admission. No new managed Turn is posted. + if (posts.length) throw new Error("Managed readback created another Turn"); await composer.fill(""); await openGoalChat(attachedGoal); @@ -102,7 +102,7 @@ export const attachedHostFollowUpScenario = { } return { coverageEntries: context.coverageEntries, - note: "managed-runtime-running-turn: Send stays closed with the wait hint and no message reaches the running Turn. " + note: "managed-runtime-running-turn: native instructions stay available without creating another Turn. " + "attached-host-follow-up: while the host runs a Turn, Send stays open, one POST reaches the Session queue and the page follows the queued Turn.", }; }, diff --git a/examples/personal-workspace-browser/chat-recovery.mjs b/examples/personal-workspace-browser/chat-recovery.mjs index 791c23b3e1..5fd873d273 100644 --- a/examples/personal-workspace-browser/chat-recovery.mjs +++ b/examples/personal-workspace-browser/chat-recovery.mjs @@ -42,6 +42,8 @@ export const chatRecoveryScenario = { if (!(await page.locator(".personal-home-lanes").isVisible())) throw new Error("Manager send replaced the home lane overview"); const managerUrlBefore = page.url(); await page.getByRole("button", { name: "询问全局待办", exact: true }).click(); + // This is a fresh question, not an instruction to the shortcut's active Turn. + await page.locator(".personal-message-pending").waitFor({ state: "hidden" }); await page.getByLabel("向 LoopX 发送消息").fill("我现在该做什么?只读回答,不要创建或修改任何状态。"); await page.getByRole("button", { name: "发送", exact: true }).click(); await page.getByText("管家已读取当前授权范围的 Goal 证据。", { exact: true }).waitFor({ state: "visible" }); @@ -292,7 +294,7 @@ export const chatRecoveryScenario = { for (const suffix of ["a", "b"]) { const sessionId = `session-run-action-${suffix}`; page.__loopxRuntime.sessions.set(sessionId, { - session_id: sessionId, goal_id: actionGoalId, agent_id: "codex", adapter_kind: "codex", + session_id: sessionId, goal_id: actionGoalId, agent_id: "codex", adapter_kind: "codex_app_server", channel_id: `task.run-action-${suffix}`, status: "ready", active_turn_id: null, last_error_code: null, created_at: "2026-08-13T01:00:00Z", updated_at: "2026-08-13T01:00:00Z", last_activity_at: "2026-08-13T01:00:00Z", resumable: true, }); @@ -352,13 +354,15 @@ export const chatRecoveryScenario = { if (await closeButton.isDisabled()) throw new Error("Session B's own close success left its button disabled"); pass("run-action-ownership", "Late Session close results report on, and release the guard of, only the Run that issued them"); - // The Chat service accepts one Turn per Session. After a reload the page - // only learns about a running Turn from the Session snapshot, so the - // composer must wait for it instead of sending into a 409. + // An executor without native steering accepts one Turn at a time. Its + // recovered and 409-handoff paths must keep Send closed; supported Codex + // instructions are covered by composer-session-admission. const turnsBeforeRunningCheck = api.turnRequests.length; await page.getByLabel("向 LoopX 发送消息").fill("刷新后验证中断控制:输入框应等待本轮。"); await page.getByRole("button", { name: "发送", exact: true }).click(); while (api.turnRequests.length === turnsBeforeRunningCheck) await page.waitForTimeout(50); + const unsupportedSessionId = api.turnRequests.at(-1).sessionId; + page.__loopxRuntime.sessions.set(unsupportedSessionId, { ...page.__loopxRuntime.sessions.get(unsupportedSessionId), adapter_kind: "external" }); await page.reload({ waitUntil: "domcontentloaded" }); await page.getByTestId("personal-goal-home").waitFor({ state: "visible" }); await page.locator(".personal-goal-link").first().click(); diff --git a/examples/personal-workspace-browser/composer-session-admission.mjs b/examples/personal-workspace-browser/composer-session-admission.mjs index f40d93d5af..315e24c4ce 100644 --- a/examples/personal-workspace-browser/composer-session-admission.mjs +++ b/examples/personal-workspace-browser/composer-session-admission.mjs @@ -1,10 +1,9 @@ import { openWorkspacePage } from "./scenario-context.mjs"; -// The Chat service admits a message sent while a Turn runs according to the -// Session's mode (`ChatRuntimeController.submit_turn`): a managed runtime -// Session answers 409, so the composer waits; an attached host Session queues -// the message behind the host's running Turn, so the composer stays open. -// Both cases use the same Goal Session and running Turn; only the mode differs. +// The ordinary composer steers a managed Codex Turn through exact-turn ingress; +// attached hosts retain queued follow-ups and unsupported adapters remain blocked. +// Browser routes are synthetic; the production HTTP/store/subprocess contract is +// independently covered by tests/test_chat_turn_steering.py. export const composerSessionAdmissionScenario = { id: "composer-session-admission", async run({ browser, collectCoverage, url }) { @@ -32,13 +31,14 @@ export const composerSessionAdmissionScenario = { // Start the Turn outside this page and reload, so the page learns of it // only from the Session read. Its stream is held so the Turn stays running. - const reloadWithRunningTurn = async (sessionMode, turnId) => { + const reloadWithRunningTurn = async (sessionMode, turnId, adapterKind = "codex_app_server") => { page.__loopxRuntime.turnMessages.set(turnId, "另一入口发起的运行中回合"); page.__loopxRuntime.sessions.set(sessionId, { ...page.__loopxRuntime.sessions.get(sessionId), active_turn_id: turnId, status: "busy", session_mode: sessionMode, + adapter_kind: adapterKind, host_surface: sessionMode === "attached_host" ? "codex_cli" : null, }); const heldEvents = []; @@ -55,7 +55,7 @@ export const composerSessionAdmissionScenario = { await page.route(routes.at(-1), (held) => { heldEvents.push(held); }); await route.fulfill({ contentType: "application/json", status: 202, json: { ok: true, schema_version: "loopx_chat_turn_accepted_v1", session_id: sessionId, turn_id: queuedTurnId, - created: true, status: "queued", events_url: `/api/chat/sessions/${sessionId}/turns/${queuedTurnId}/events`, + created: true, status: sessionMode === "attached_host" ? "queued" : "running", events_url: `/api/chat/sessions/${sessionId}/turns/${queuedTurnId}/events`, } }); }); await page.reload({ waitUntil: "domcontentloaded" }); @@ -63,7 +63,7 @@ export const composerSessionAdmissionScenario = { // The recovered reply is the page's own evidence that it saw the Turn run. await page.getByRole("button", { name: "中断本轮", exact: true }).waitFor({ state: "visible", timeout: 10_000 }); return { - posts, + posts, heldEvents, turnId, async release() { for (const pattern of routes) await page.unroute(pattern); await Promise.all(heldEvents.splice(0).map((held) => held.abort().catch(() => {}))); @@ -72,14 +72,85 @@ export const composerSessionAdmissionScenario = { }; const managed = await reloadWithRunningTurn("managed_runtime", `turn-managed-${Date.now()}`); - await turnRunningHint.waitFor({ state: "visible", timeout: 5_000 }); - await composerInput.fill("托管回合运行中不应发送"); - if (!(await sendButton.isDisabled())) throw new Error("A managed runtime Session's composer stayed sendable while its Turn ran"); - await sendButton.click({ force: true }); - await page.waitForTimeout(300); - if (managed.posts.length !== 0) throw new Error(`The composer posted ${managed.posts.length} times into a running managed Turn`); + const adjustments = []; + let delayedReceipt; + await page.route(`**/api/chat/sessions/${sessionId}/turns/*/steer`, async (route) => { + const body = route.request().postDataJSON(); + const target = new URL(route.request().url()).pathname.split("/")[6]; + adjustments.push({ ...body, turnId: target }); + if ([1, 6].includes(adjustments.length)) return route.fulfill({ status: 409, json: { ok: false, error: "接收状态未确认" } }); + if (adjustments.length === 4) return route.fulfill({ status: 409, json: { ok: false, error: "本次未送达", delivery_state: "not_delivered" } }); + const receipt = { ok: true, session_id: sessionId, turn_id: adjustments.length === 2 ? "wrong-turn" : target, + client_ingress_id: body.client_ingress_id, status: "delivered" }; + if (adjustments.length === 3) { delayedReceipt = () => route.fulfill({ json: receipt }); return; } + await route.fulfill({ json: receipt }); + }); + const instruction = "先做中文,别发布。"; + await page.locator(".personal-composer-status", { hasText: "发消息可调整当前工作" }).waitFor(); + await composerInput.fill(instruction); + if (await sendButton.isDisabled()) throw new Error("The native Turn cannot receive instructions from its ordinary composer"); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "接收状态未确认" }).waitFor(); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "回执不匹配" }).waitFor(); + if (await composerInput.inputValue() !== instruction) throw new Error("Unconfirmed steering lost its draft"); + await sendButton.click(); + for (let attempt = 0; attempt < 100 && !delayedReceipt; attempt += 1) await page.waitForTimeout(50); + if (!delayedReceipt) throw new Error("Retry did not reach the original Turn"); + await composerInput.fill("先写下另一条指令"); + await delayedReceipt(); + await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + if (await composerInput.inputValue() !== "先写下另一条指令") throw new Error("Delivery erased the draft typed while it was sending"); + if (new Set(adjustments.slice(0, 3).map(row => row.client_ingress_id)).size !== 1) throw new Error("Uncertain retries minted new ingress identities"); + await composerInput.fill(instruction); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "本次未送达" }).waitFor(); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + if (adjustments[3].client_ingress_id === adjustments[4].client_ingress_id) throw new Error("Confirmed non-delivery could not retry after recovery"); + + // Lose a response, then let the original Turn finish. Retry reads its own + // receipt; it must neither create another Turn nor steer a newer one. + await composerInput.fill(instruction); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "接收状态未确认" }).waitFor(); + page.__loopxRuntime.sessions.set(sessionId, { ...page.__loopxRuntime.sessions.get(sessionId), active_turn_id: null, status: "ready" }); + const completed = { event_id: "done", sequence: 1, kind: "turn.completed", payload: { response: { schema_version: "loopx_chat_agent_response_v0", message: "本轮已完成", proposals: [], gate: null } } }; + await Promise.all(managed.heldEvents.splice(0).map(route => route.fulfill({ contentType: "text/event-stream", + body: `id: done\nevent: turn.completed\ndata: ${JSON.stringify(completed)}\n\n` }))); + await page.getByText("本轮已完成", { exact: true }).waitFor(); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + if (adjustments.length !== 7 || adjustments[5].client_ingress_id !== adjustments[6].client_ingress_id + || adjustments.some(row => row.turnId !== managed.turnId) || managed.posts.length) { + throw new Error("Steering retry changed its identity/target or started another Turn"); + } + // A Turn sent from this page keeps its original send promise pending. + // That promise must not block the same composer's native instructions. + await composerInput.fill("给 LoopX 做份社区问卷,先给我草稿。"); + await sendButton.click(); + await page.locator(".personal-composer-status", { hasText: "发消息可调整当前工作" }).waitFor(); + await composerInput.fill(instruction); + if (await sendButton.isDisabled()) throw new Error("Original send promise blocks native instructions"); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + if (adjustments.length !== 8 || managed.posts.length !== 1 + || adjustments[7].turnId === managed.turnId || await composerInput.inputValue()) { + throw new Error("Instructions during the original send started new work or failed to clear the submitted draft"); + } await managed.release(); - notes.push("managed-runtime: a running Turn keeps Send closed and posts nothing"); + notes.push("managed Codex: ordinary composer steers its exact Turn while the original send waits, retaining drafts and retry identity, including after completion"); + + const unsupported = await reloadWithRunningTurn("managed_runtime", `turn-unsupported-${Date.now()}`, "external"); + await turnRunningHint.waitFor(); + await composerInput.fill("不应向不支持的执行器追加指令"); + if (!await sendButton.isDisabled()) throw new Error("Unsupported executor was offered native steering"); + await sendButton.click({ force: true }); + await page.waitForTimeout(100); + if (unsupported.posts.length || adjustments.length !== 8) throw new Error("Unsupported executor received an effect"); + await composerInput.fill(""); + await unsupported.release(); + notes.push("unsupported managed adapter: draft retained and no Turn or steering effect"); const attached = await reloadWithRunningTurn("attached_host", `turn-attached-${Date.now()}`); const followUp = "宿主回合运行中排队的下一条消息"; diff --git a/examples/personal-workspace-browser/fixture.mjs b/examples/personal-workspace-browser/fixture.mjs index c0aa3e5b3a..ca1cf8c2c7 100644 --- a/examples/personal-workspace-browser/fixture.mjs +++ b/examples/personal-workspace-browser/fixture.mjs @@ -1506,7 +1506,7 @@ export async function installApi(page, { goalSubagentConfigurationEnabled = true const resolvedAgentId = body.agent_id ?? managerChannelBinding?.executor_endpoint ?? state.machineNamespaces?.steward_executor?.executor_endpoint ?? "codex"; const session_id = `session-${body.context_kind}-${resolvedGoalId}-${resolvedAgentId}`; const existing = body.mode === "resume_latest" ? sessions.get(session_id) : null; - const session = existing ?? { session_id, goal_id: resolvedGoalId, agent_id: resolvedAgentId, adapter_kind: resolvedAgentId, channel_id: body.context_kind === "manager" ? "manager" : `goal.${body.goal_id}`, status: "ready", active_turn_id: null, last_error_code: null, created_at: "2026-08-13T01:00:00Z", updated_at: "2026-08-13T01:00:00Z", last_activity_at: "2026-08-13T01:00:00Z", resumable: true, ...(body.context_kind === "manager" ? { manager_runtime: { schema_version: "manager_runtime_session_readback_v0", runtime_profile: "restricted", configuration_revision: "absent", status: "ready", sandbox: "read-only", standing_grant: "none", tool_classes: ["loopx_core"] } } : {}) }; + const session = existing ?? { session_id, goal_id: resolvedGoalId, agent_id: resolvedAgentId, adapter_kind: resolvedAgentId === "codex" ? "codex_app_server" : resolvedAgentId === "claude-code" ? "claude_code_cli" : "acp", channel_id: body.context_kind === "manager" ? "manager" : `goal.${body.goal_id}`, status: "ready", active_turn_id: null, last_error_code: null, created_at: "2026-08-13T01:00:00Z", updated_at: "2026-08-13T01:00:00Z", last_activity_at: "2026-08-13T01:00:00Z", resumable: true, ...(body.context_kind === "manager" ? { manager_runtime: { schema_version: "manager_runtime_session_readback_v0", runtime_profile: "restricted", configuration_revision: "absent", status: "ready", sandbox: "read-only", standing_grant: "none", tool_classes: ["loopx_core"] } } : {}) }; sessions.set(session_id, session); messages.set(session_id, messages.get(session_id) ?? []); await route.fulfill({ contentType: "application/json", json: { ok: true, agent_id: resolvedAgentId, goal_id: body.goal_id, resumed: body.mode === "resume_latest", session_id, session }, status: 201 }); @@ -1828,7 +1828,7 @@ export async function installApi(page, { goalSubagentConfigurationEnabled = true session_id: sessionId, goal_id: preview.normalized_parameters.goal_id, agent_id: "codex", - adapter_kind: "codex", + adapter_kind: "codex_app_server", channel_id: `goal.${preview.normalized_parameters.goal_id}`, active_turn_id: null, status: "ready", diff --git a/examples/personal-workspace-browser/steward-journey.mjs b/examples/personal-workspace-browser/steward-journey.mjs index 6a7f174bf4..a95079dbfd 100644 --- a/examples/personal-workspace-browser/steward-journey.mjs +++ b/examples/personal-workspace-browser/steward-journey.mjs @@ -209,6 +209,9 @@ export const stewardJourneyScenario = { gate_prompt_sent: Boolean(chipTurn), composer_after_click: composerAfterChip, }); + // A new question follows the chip response; an in-flight message would + // instead adjust the native Codex Turn through the ordinary composer. + await page.locator(".personal-message-pending").waitFor({ state: "hidden" }); await composer.fill(`${STEWARD_PROMPT}:请给我一份当前 Goal 的下一步。`); await page.getByRole("button", { name: "发送", exact: true }).click(); let turn;