diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 969814a740..c7b57a1519 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -1032,12 +1032,12 @@ export async function interruptChatTurn(sessionId: string, turnId: string) { } export async function steerChatTurn(sessionId: string, turnId: string, message: string, ingressId: string) { - const receipt = await requestJson<{ ok: boolean; session_id: string; turn_id: string; client_ingress_id: string; status: string }>( + const receipt = await requestJson<{ ok: boolean; session_id: string; turn_id: string; client_ingress_id: string; status: string; created: boolean }>( `/api/chat/sessions/${sessionId}/turns/${turnId}/steer`, { method: "POST", body: JSON.stringify({ message, client_ingress_id: ingressId }) }, ); if (receipt.ok !== true || receipt.session_id !== sessionId || receipt.turn_id !== turnId - || receipt.client_ingress_id !== ingressId || receipt.status !== "delivered") { + || receipt.client_ingress_id !== ingressId || receipt.status !== "delivered" || typeof receipt.created !== "boolean") { throw new ChatApiError("追加指令的回执不匹配,请保留草稿并检查当前状态。", { error_code: "steer_receipt_mismatch" }); } return receipt; diff --git a/apps/presentation/dashboard/src/data/conversation-returns.test.mjs b/apps/presentation/dashboard/src/data/conversation-returns.test.mjs index cd1fcc3425..c5b097b132 100644 --- a/apps/presentation/dashboard/src/data/conversation-returns.test.mjs +++ b/apps/presentation/dashboard/src/data/conversation-returns.test.mjs @@ -67,4 +67,16 @@ assert.equal(hydratedRoles.length, 2); assert.equal(hydratedRoles[0].sourceMessageId, "user-message"); assert.equal(hydratedRoles[1].sourceMessageId, "agent-message"); assert.equal(hydratedRoles[1].text, "Live answer"); +// The initial optimistic request is the first user message of its Turn. +// Later instructions must retain independent stored identities and remain visible. +const withInstructions = [storedRoles[0], + { session_id: "current", message_id: "instruction-1", turn_id: "same-turn", role: "user", text: "Chinese first" }, + { session_id: "current", message_id: "instruction-2", turn_id: "same-turn", role: "user", text: "Do not publish" }, + storedRoles[1]]; +const recoveredInstructions = reconcileConversationHistory(bothRoles, withInstructions, createHistory); +assert.equal(recoveredInstructions.length, 4); +assert.equal(recoveredInstructions.find(row => row.text === "My request").sourceMessageId, "user-message"); +assert.equal(recoveredInstructions.find(row => row.text === "Chinese first").sourceMessageId, "instruction-1"); +assert.equal(recoveredInstructions.find(row => row.text === "Do not publish").sourceMessageId, "instruction-2"); +assert.equal(reconcileConversationHistory(recoveredInstructions, withInstructions, createHistory), recoveredInstructions); console.log("conversation-returns: passed (session isolation, late return, deduplication, transport uncertainty, stream preservation and watch retirement)"); diff --git a/apps/presentation/dashboard/src/data/conversation-returns.ts b/apps/presentation/dashboard/src/data/conversation-returns.ts index cea74aec3d..5fd82d3900 100644 --- a/apps/presentation/dashboard/src/data/conversation-returns.ts +++ b/apps/presentation/dashboard/src/data/conversation-returns.ts @@ -31,8 +31,14 @@ export function reconcileConversationReturns( createReply: (message: ChatVisibleMessage) => T, ): T[] { const byId = new Map(messages.map((row) => [row.message_id, row])); - const byTurn = new Map(messages.filter((row) => row.origin !== "manager_followup") - .map((row) => [`${row.turn_id}:${row.role === "user" ? "user" : "assistant"}`, row])); + const byTurn = new Map(); + for (const row of messages) { + if (row.origin === "manager_followup") continue; + const key = `${row.turn_id}:${row.role === "user" ? "user" : "assistant"}`; + // The Turn starts with its user request. Later user instructions have their + // own message IDs; they must not hydrate that original optimistic request. + if (row.role !== "user" || !byTurn.has(key)) byTurn.set(key, row); + } const seen = new Set(previous.filter((row) => row.sourceSessionId === sessionId).map((row) => row.sourceMessageId)); let changed = false; const updated = previous.map((row) => { diff --git a/apps/presentation/dashboard/src/features/personal-workspace/composer-steering-recovery.ts b/apps/presentation/dashboard/src/features/personal-workspace/composer-steering-recovery.ts new file mode 100644 index 0000000000..bfd75598da --- /dev/null +++ b/apps/presentation/dashboard/src/features/personal-workspace/composer-steering-recovery.ts @@ -0,0 +1,36 @@ +/** Client retry identity only; the Chat ingress/store still owns delivery. */ +export type ComposerSteeringRequest = { + sessionId: string; + turnId: string; + text: string; + id: string; +}; + +const storageKey = "loopx-pw-composer-steering"; + +export function readComposerSteeringRequests(): Map { + try { + const raw = window.sessionStorage.getItem(storageKey); + const entries: unknown = raw ? JSON.parse(raw) : []; + if (!Array.isArray(entries)) return new Map(); + return new Map(entries.flatMap((entry): [string, ComposerSteeringRequest][] => { + if (!Array.isArray(entry) || entry.length !== 2 || typeof entry[0] !== "string") return []; + const request = entry[1]; + if (!request || typeof request !== "object" || Array.isArray(request) + || ![request.sessionId, request.turnId, request.text, request.id] + .every(value => typeof value === "string" && value.length > 0)) return []; + return [[entry[0], { sessionId: request.sessionId, turnId: request.turnId, text: request.text, id: request.id }]]; + })); + } catch { + return new Map(); + } +} + +export function persistComposerSteeringRequests(requests: ReadonlyMap) { + try { + window.sessionStorage.setItem(storageKey, JSON.stringify([...requests])); + } catch { + // As with composer drafts, unavailable browser storage retains only memory. + // No provider action or delivery conclusion follows from this cache. + } +} 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 116d07c326..ae53146737 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 @@ -1,4 +1,5 @@ import { goalCreateRequest } from "./goal-create-request"; +import { persistComposerSteeringRequests, readComposerSteeringRequests, type ComposerSteeringRequest } from "./composer-steering-recovery"; import type { ConversationHistoryStatus } from "../../data/use-conversation-history"; import { GoalDraftCard } from "./goal-draft-card"; import type { GoalDraft } from "../../../../../../loopx/control_plane/collaboration/goal_draft.js"; @@ -810,7 +811,8 @@ export function PersonalWorkspacePage({ }); const [sending, setSending] = useState(false); const [steering, setSteering] = useState(false); - const steeringRequests = useRef(new Map()); + const [restoredSteeringRequests] = useState(readComposerSteeringRequests); + const steeringRequests = useRef(restoredSteeringRequests); const [actionDraft, setActionDraft] = useState(null); const [loopxMode, setLoopxMode] = useState(null); const [loopxDelivery, setLoopxDelivery] = useState<"queue" | "inbox" | "steer">("queue"); @@ -873,6 +875,17 @@ export function PersonalWorkspacePage({ function setComposer(value: string) { setComposerDraft(composerDraftKey, value); } + function retainSteeringRequest(key: string, request: ComposerSteeringRequest) { + steeringRequests.current.set(key, request); + // Persist before sending: reload after provider acceptance must replay the + // original ingress, never silently submit another Turn. + persistComposerSteeringRequests(steeringRequests.current); + } + function retireSteeringRequest(key: string, id: string) { + if (steeringRequests.current.get(key)?.id !== id) return; + steeringRequests.current.delete(key); + persistComposerSteeringRequests(steeringRequests.current); + } async function reviewGoalDraft(draft: GoalDraft, edit = false, draftId = "") { // Source message + reviewed contents survive retry without merging distinct requests. if (!edit && !draft.question && draft.completion_criteria.trim()) { @@ -1641,20 +1654,20 @@ export function PersonalWorkspacePage({ return; } const request = retry ?? { sessionId: conversationSessionId, turnId: steeringTurnId!, text: message, id: crypto.randomUUID() }; - steeringRequests.current.set(composerDraftKey, request); + retainSteeringRequest(composerDraftKey, request); setSteering(true); setActionFeedback(null); setImageAttachmentError(null); try { await callbacks.onSteerConversationTurn(selectedGoalId ?? "manager", request.turnId, message, request.id); - steeringRequests.current.delete(composerDraftKey); + retireSteeringRequest(composerDraftKey, request.id); 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); + retireSteeringRequest(composerDraftKey, request.id); } setActionFeedback(error instanceof Error ? error.message : t("feedback.sendGenericError")); } finally { setSteering(false); } diff --git a/apps/presentation/dashboard/src/views/dashboard-page.tsx b/apps/presentation/dashboard/src/views/dashboard-page.tsx index 8481280032..fb41de5ba3 100644 --- a/apps/presentation/dashboard/src/views/dashboard-page.tsx +++ b/apps/presentation/dashboard/src/views/dashboard-page.tsx @@ -2952,7 +2952,13 @@ function PersonalGoalHome({ 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 receipt = await steerChatTurn(binding.sessionId, turnId, message, ingressId); + if (receipt.created === false) { + // A replay reads an existing delivery; its message belongs to the + // stored transcript, not a second optimistic user bubble. + if (targetContextId === contextId) await conversationHistory.refresh(); + return; + } const id = managerMessageId.current++; setMessagesByContext(current => { const messages = current[targetContextId] ?? []; 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 07bcc83a5a..39f6e4d19a 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -133,7 +133,16 @@ receipt → observed work or actionable failure → readable answer in the same 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. + retry after completion or page reload in the same browser tab. The shared + composer persists the original Session/Turn/text/ingress before dispatch and + reuses it only for the matching Session and text; restoring the cache never + sends automatically. Only a confirmed non-delivery permits a new ingress. + Replaying a delivered receipt reads the current stored transcript rather than + reusing an earlier snapshot or inserting another local copy of the instruction. + History keeps the initial user request and later instructions as distinct messages + even while the initial request is still being reconciled with its stored identity. + Unavailable browser storage keeps current-page retry behavior but cannot + promise reload recovery. This tab-local cache is not delivery authority. 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 diff --git a/examples/personal-workspace-browser/chat-recovery.mjs b/examples/personal-workspace-browser/chat-recovery.mjs index 5fd873d273..3c9c557037 100644 --- a/examples/personal-workspace-browser/chat-recovery.mjs +++ b/examples/personal-workspace-browser/chat-recovery.mjs @@ -529,7 +529,7 @@ export const chatRecoveryScenario = { const body = route.request().postDataJSON(); const [sessionId, turnId] = new URL(route.request().url()).pathname.match(/sessions\/([^/]+)\/turns\/([^/]+)\/steer/).slice(1); steers.push({ sessionId, turnId }); - await route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: turnId, client_ingress_id: body.client_ingress_id, status: "delivered" } }); + await route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: turnId, client_ingress_id: body.client_ingress_id, status: "delivered", created: true } }); }); const interruptsBefore = api.interrupts.length; await composerInput.fill(draft); diff --git a/examples/personal-workspace-browser/composer-session-admission.mjs b/examples/personal-workspace-browser/composer-session-admission.mjs index 315e24c4ce..446fb81436 100644 --- a/examples/personal-workspace-browser/composer-session-admission.mjs +++ b/examples/personal-workspace-browser/composer-session-admission.mjs @@ -7,7 +7,13 @@ import { openWorkspacePage } from "./scenario-context.mjs"; export const composerSessionAdmissionScenario = { id: "composer-session-admission", async run({ browser, collectCoverage, url }) { - const context = await openWorkspacePage(browser, url, { collectCoverage }); + const context = await openWorkspacePage(browser, url, { collectCoverage, + beforeGoto: async (_api, page) => page.addInitScript(() => { + if (!sessionStorage.getItem("loopx-pw-composer-steering")) { + sessionStorage.setItem("loopx-pw-composer-steering", JSON.stringify([["invalid", { id: "incomplete" }]])); + } + }), + }); const { page } = context; const notes = []; const composerInput = page.getByLabel("向 LoopX 发送消息"); @@ -50,6 +56,8 @@ export const composerSessionAdmissionScenario = { const body = route.request().postDataJSON(); posts.push(body.message); const queuedTurnId = `${turnId}-queued-${posts.length}`; + page.__loopxRuntime.messages.get(sessionId).push({ message_id: `${queuedTurnId}-user`, turn_id: queuedTurnId, + role: "user", text: body.message, created_at: "2026-08-13T01:00:01Z" }); // A queued follow-up waits behind the running Turn; its stream is held too. routes.push(`**/api/chat/sessions/${sessionId}/turns/${queuedTurnId}/events`); await page.route(routes.at(-1), (held) => { heldEvents.push(held); }); @@ -78,10 +86,14 @@ export const composerSessionAdmissionScenario = { 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 === 9) { + page.__loopxRuntime.messages.get(sessionId).push({ message_id: "stored-current-page", turn_id: target, + role: "user", text: body.message, created_at: "2026-08-13T01:00:02Z" }); + } + if ([1, 6, 9].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" }; + client_ingress_id: body.client_ingress_id, status: "delivered", created: ![7, 10].includes(adjustments.length) }; if (adjustments.length === 3) { delayedReceipt = () => route.fulfill({ json: receipt }); return; } await route.fulfill({ json: receipt }); }); @@ -91,6 +103,13 @@ export const composerSessionAdmissionScenario = { 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 Promise.all(managed.heldEvents.splice(0).map(route => route.abort().catch(() => {}))); + await page.reload({ waitUntil: "domcontentloaded" }); + await openGoalChat(); + await page.getByRole("button", { name: "中断本轮", exact: true }).waitFor(); + if (adjustments.length !== 1 || managed.posts.length || await composerInput.inputValue() !== instruction) { + throw new Error("Restoring an unconfirmed instruction sent automatically or lost its draft"); + } await sendButton.click(); await page.getByRole("status").filter({ hasText: "回执不匹配" }).waitFor(); if (await composerInput.inputValue() !== instruction) throw new Error("Unconfirmed steering lost its draft"); @@ -115,16 +134,30 @@ export const composerSessionAdmissionScenario = { 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" }); + page.__loopxRuntime.messages.get(sessionId).push({ message_id: "stored-correction", turn_id: managed.turnId, + role: "user", text: instruction, created_at: "2026-08-13T01:00:01Z" }); 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 page.reload({ waitUntil: "domcontentloaded" }); + await openGoalChat(); + if (adjustments.length !== 6 || managed.posts.length) throw new Error("Reload automatically replayed uncertain work"); + if (await composerInput.inputValue() !== instruction) throw new Error("Reload lost the unconfirmed correction draft"); + const storedInstruction = page.locator(".personal-message").getByText(instruction, { exact: true }); + await storedInstruction.waitFor(); await sendButton.click(); await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + if (await storedInstruction.count() !== 1) { + throw new Error("Reading a delivered receipt duplicated its stored instruction in the conversation"); + } 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"); } + if (await page.evaluate(() => JSON.parse(sessionStorage.getItem("loopx-pw-composer-steering")).length)) { + throw new Error("Accepted instructions left an uncertain retry cached"); + } // 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 做份社区问卷,先给我草稿。"); @@ -138,8 +171,22 @@ export const composerSessionAdmissionScenario = { || 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"); } + // Without reloading, a delivered replay must fetch the new stored message, + // not reuse the history snapshot taken before the instruction was sent. + const currentPageInstruction = "问卷里补上使用频率。"; + await composerInput.fill(currentPageInstruction); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "接收状态未确认" }).waitFor(); + await sendButton.click(); + await page.getByRole("status").filter({ hasText: "执行器已接收本轮追加指令" }).waitFor(); + const currentPageStored = page.locator(".personal-message").getByText(currentPageInstruction, { exact: true }); + await currentPageStored.waitFor({ timeout: 5000 }); + if (await currentPageStored.count() !== 1 || adjustments.length !== 10 + || adjustments[8].client_ingress_id !== adjustments[9].client_ingress_id || managed.posts.length !== 1) { + throw new Error("Delivered replay did not read back one stored instruction on the current page"); + } await managed.release(); - notes.push("managed Codex: ordinary composer steers its exact Turn while the original send waits, retaining drafts and retry identity, including after completion"); + notes.push("managed Codex: ordinary composer steers its exact Turn while the original send waits, retaining drafts and retry identity, including after completion and reload; restored requests never dispatch automatically"); const unsupported = await reloadWithRunningTurn("managed_runtime", `turn-unsupported-${Date.now()}`, "external"); await turnRunningHint.waitFor(); @@ -147,7 +194,7 @@ export const composerSessionAdmissionScenario = { 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"); + if (unsupported.posts.length || adjustments.length !== 10) 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"); diff --git a/examples/personal-workspace-browser/conversation-activity.mjs b/examples/personal-workspace-browser/conversation-activity.mjs index 7939f1e690..954ffae369 100644 --- a/examples/personal-workspace-browser/conversation-activity.mjs +++ b/examples/personal-workspace-browser/conversation-activity.mjs @@ -54,7 +54,7 @@ export const conversationActivityScenario = { const [sessionId, turnId] = new URL(route.request().url()).pathname.match(/sessions\/([^/]+)\/turns\/([^/]+)\/steer/).slice(1); if (adjustments.length === 1) return route.fulfill({ status: 409, json: { ok: false, error: "执行器暂时未确认接收,草稿已保留。" } }); if (adjustments.length === 3) return route.fulfill({ status: 409, json: { ok: false, error: "执行器暂时不可用。", error_code: "live_steering_session_not_attached", delivery_state: "not_delivered" } }); - return route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: adjustments.length === 2 ? "wrong-turn" : turnId, client_ingress_id: body.client_ingress_id, status: "delivered" } }); + return route.fulfill({ json: { ok: true, session_id: sessionId, turn_id: adjustments.length === 2 ? "wrong-turn" : turnId, client_ingress_id: body.client_ingress_id, status: "delivered", created: true } }); }); await pending.getByRole("button", { name: "调整本轮", exact: true }).click(); await pending.getByLabel("追加给本轮的指令").fill("先核对依赖,再继续当前任务。");