diff --git a/apps/presentation/dashboard/package.json b/apps/presentation/dashboard/package.json index 7a9d0d4382..50c276b8dd 100644 --- a/apps/presentation/dashboard/package.json +++ b/apps/presentation/dashboard/package.json @@ -61,6 +61,7 @@ "build:chat:vite": "tsc --noEmit && vite build --config vite.chat.config.ts", "smoke:chat-upgrade": "LOOPX_PLAYWRIGHT_PACKAGE=\"$PWD/node_modules/playwright\" node ../../../examples/chat-bundle-upgrade-browser-smoke.mjs", "test:conversation-returns": "node --experimental-strip-types src/data/conversation-returns.test.mjs", + "smoke:conversation-history": "vite build --ssr smoke/conversation-history-smoke.ts --outDir node_modules/.cache/loopx-conversation-history --emptyOutDir && node node_modules/.cache/loopx-conversation-history/conversation-history-smoke.js", "smoke:recent-completions": "tsc --ignoreConfig --target ES2022 --module CommonJS --moduleResolution Node --ignoreDeprecations 6.0 --resolveJsonModule --esModuleInterop --jsx react-jsx --skipLibCheck --strict --types node --outDir /tmp/loopx-recent-completions-smoke smoke/recent-completions-smoke.ts && NODE_PATH=\"$PWD/node_modules\" node /tmp/loopx-recent-completions-smoke/apps/presentation/dashboard/smoke/recent-completions-smoke.js" }, "dependencies": { diff --git a/apps/presentation/dashboard/smoke/conversation-history-http-fixture.py b/apps/presentation/dashboard/smoke/conversation-history-http-fixture.py new file mode 100644 index 0000000000..6289258326 --- /dev/null +++ b/apps/presentation/dashboard/smoke/conversation-history-http-fixture.py @@ -0,0 +1,73 @@ +"""Disposable real Chat HTTP/store fixture; the fault changes reads only.""" +from __future__ import annotations + +import json +from pathlib import Path +import sys +import tempfile +import threading + +sys.path.insert(0, str(Path(__file__).resolve().parents[4])) + +from loopx.chat_runtime import ChatRuntimeController +from loopx.chat_server import ChatHTTPServer, ChatRequestHandler +from loopx.chat_store import ChatSessionStore + + +class HistoryReadFault(ChatRequestHandler): + def _session_snapshot(self, session_id: str) -> None: + if session_id == self.server.unavailable_session_id: + self._send_error("Synthetic history read unavailable", status=503) + else: + super()._session_snapshot(session_id) + + +def main() -> None: + with tempfile.TemporaryDirectory(prefix="loopx-history-http-") as directory: + root = Path(directory) + registry = root / "registry.json" + registry.write_text(json.dumps({"schema_version": "0.1", "goals": [ + {"id": "research", "repo": str(root), "status": "active"}, + ]}), encoding="utf-8") + store = ChatSessionStore(root / "runtime") + runtime = ChatRuntimeController(store=store, codex_bin="missing-codex", registry_path=registry) + for session_id, channel, text in [ + ("old", "goal.research", "Earlier public report"), + ("current", "goal.research", "Current public report"), + ("other-channel", "manager", "Unrelated conversation"), + ]: + store.create_session(goal_id="research", agent_id="codex", adapter_kind="codex_app_server", + upstream_thread_id=session_id, channel_id=channel, session_id=session_id) + store.append_message(session_id, role="agent", text=text, message_id="answer") + before = {str(path.relative_to(store.sessions_root)): path.read_bytes() + for path in store.sessions_root.rglob("*") if path.is_file()} + server = ChatHTTPServer(("127.0.0.1", 0), HistoryReadFault) + server.verbose = False + server.registry_path = registry + server.runtime_root = root / "runtime" + server.chat_store = store + server.runtime_controller = runtime + server.unavailable_session_id = "old" + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + print(json.dumps({"origin": f"http://127.0.0.1:{server.server_port}"}), flush=True) + try: + if sys.stdin.readline().strip() != "recover": + raise ValueError("Expected read recovery") + server.unavailable_session_id = "" + print(json.dumps({"recovered": True}), flush=True) + if sys.stdin.readline().strip() != "inspect": + raise ValueError("Expected final inspection") + after = {str(path.relative_to(store.sessions_root)): path.read_bytes() + for path in store.sessions_root.rglob("*") if path.is_file()} + print(json.dumps({"store_unchanged": before == after, + "turn_count": sum(1 for _ in store.sessions_root.glob("*/turns/*.json"))}), flush=True) + finally: + server.shutdown() + server.server_close() + thread.join(timeout=5) + runtime.close() + + +if __name__ == "__main__": + main() diff --git a/apps/presentation/dashboard/smoke/conversation-history-smoke.ts b/apps/presentation/dashboard/smoke/conversation-history-smoke.ts new file mode 100644 index 0000000000..d8c647f2cd --- /dev/null +++ b/apps/presentation/dashboard/smoke/conversation-history-smoke.ts @@ -0,0 +1,58 @@ +import assert from "node:assert/strict"; +import { spawn } from "node:child_process"; +import { once } from "node:events"; +import { resolve } from "node:path"; +import { createInterface } from "node:readline"; +import { resolveTestPython } from "../../../../scripts/test-python.mjs"; +import { fetchChatHistory } from "../src/data/chat"; + +const repoRoot = resolve(process.cwd(), "../../.."); +const child = spawn(resolveTestPython({ repoRoot }), ["-u", "apps/presentation/dashboard/smoke/conversation-history-http-fixture.py"], + { cwd: repoRoot, stdio: ["pipe", "pipe", "pipe"] }); +const exited = once(child, "exit"); +let stderr = ""; +child.stderr.on("data", chunk => { stderr += String(chunk); }); +const output = createInterface({ input: child.stdout }); +const lines = output[Symbol.asyncIterator](); +const originalFetch = globalThis.fetch; +async function next() { + const line = await lines.next(); + assert.equal(line.done, false, stderr); + return JSON.parse(line.value!); +} +try { + const { origin } = await next(); + const requests: string[] = []; + globalThis.fetch = (input, init) => { + const path = String(input); + assert.ok(!init?.method || init.method === "GET", "History recovery is read-only"); + requests.push(path); + return originalFetch(new URL(path, origin), init); + }; + const options = { agentId: "codex", goalId: "research", channelId: "goal.research" }; + const partial = await fetchChatHistory(options); + assert.deepEqual(partial.unavailableSessionIds, ["old"]); + assert.deepEqual(partial.messages.map(row => row.text), ["Current public report"]); + assert.equal(partial.sessions[0].session_id, "current"); + assert.equal(partial.sessions.some(row => row.session_id === "other-channel"), false); + child.stdin.write("recover\n"); + assert.equal((await next()).recovered, true); + requests.length = 0; + const restored = await fetchChatHistory(options, partial); + assert.deepEqual(requests, ["/api/chat/sessions/old"], "Recovery reads only the missing session"); + assert.deepEqual(restored.unavailableSessionIds, []); + assert.deepEqual(restored.messages.map(row => row.text), ["Earlier public report", "Current public report"]); + assert.deepEqual(restored.messages.map(row => row.session_id), ["old", "current"], "Colliding IDs preserve both sessions"); + requests.length = 0; + await fetchChatHistory(options, restored); + assert.deepEqual(requests, [], "A complete cached history does not poll"); + child.stdin.end("inspect\n"); + assert.deepEqual(await next(), { store_unchanged: true, turn_count: 0 }); + const [exitCode] = await exited; + assert.equal(exitCode, 0, stderr); + console.log("conversation-history: passed (real HTTP/store, partial read, channel isolation, missing-only recovery, scoped identity, zero writes or Turns)"); +} finally { + globalThis.fetch = originalFetch; + output.close(); + if (child.exitCode === null) { child.kill(); await exited; } +} diff --git a/apps/presentation/dashboard/src/data/chat.ts b/apps/presentation/dashboard/src/data/chat.ts index 1a77682671..911c62d25f 100644 --- a/apps/presentation/dashboard/src/data/chat.ts +++ b/apps/presentation/dashboard/src/data/chat.ts @@ -623,7 +623,10 @@ export async function recordProjectionExchange(options: { goalId?: string; question: string; }) { - return requestJson<{ ok: true; schema_version: "loopx_chat_projection_exchange_v1"; session_id: string }>( + return requestJson<{ + ok: true; schema_version: "loopx_chat_projection_exchange_v1"; + session_id: string; user_message_id: string; answer_message_id: string; + }>( "/api/chat/projection-messages", { method: "POST", @@ -742,14 +745,15 @@ export type ChatSessionSnapshot = { active_turn: Record | null; }; -export async function fetchChatSession(sessionId: string) { - return requestJson(`/api/chat/sessions/${sessionId}`); +export async function fetchChatSession(sessionId: string, signal?: AbortSignal) { + return requestJson(`/api/chat/sessions/${sessionId}`, { signal }); } export async function fetchChatSessions(options: { agentId?: string; channelId?: string; goalId?: string; + signal?: AbortSignal; }) { const query = new URLSearchParams(); if (options.agentId) query.set("agent_id", options.agentId); @@ -759,14 +763,14 @@ export async function fetchChatSessions(options: { ok: true; schema_version: "loopx_chat_session_list_v1"; sessions: ChatSessionSummary[]; - }>(`/api/chat/sessions?${query.toString()}`); + }>(`/api/chat/sessions?${query.toString()}`, { signal: options.signal }); } export function mergeChatSessionMessages(snapshots: ChatSessionSnapshot[]) { const messages = new Map(); for (const snapshot of snapshots) { for (const message of snapshot.messages) { - messages.set(message.message_id, { ...message, session_id: snapshot.session.session_id }); + messages.set(`${snapshot.session.session_id}:${message.message_id}`, { ...message, session_id: snapshot.session.session_id }); } } return [...messages.values()].sort((left, right) => @@ -775,6 +779,13 @@ export function mergeChatSessionMessages(snapshots: ChatSessionSnapshot[]) { ); } +export type ChatHistory = { + messages: ChatVisibleMessage[]; + sessions: ChatSessionSummary[]; + snapshots: ChatSessionSnapshot[]; + unavailableSessionIds: string[]; +}; + export async function fetchChatHistory(options: { // An omitted ``agentId`` reads the whole channel transcript. The steward // channel is one conversation across whatever executor it currently @@ -782,15 +793,22 @@ export async function fetchChatHistory(options: { agentId?: string; channelId: string; goalId?: string; -}) { - const listed = await fetchChatSessions(options); - const snapshots = await Promise.all( - listed.sessions.map((session) => fetchChatSession(session.session_id)), - ); +}, previous?: ChatHistory): Promise { + const listed = previous ?? await fetchChatSessions({ ...options, signal: AbortSignal.timeout(5000) }); + const known = new Map(previous?.snapshots.map((snapshot) => [snapshot.session.session_id, snapshot])); + // A failed historical read is not an empty transcript. Retrying only the + // missing snapshots keeps this recovery read-only and bounds repeated work. + const missing = listed.sessions.filter((session) => !known.has(session.session_id)); + const results = await Promise.allSettled(missing.map((session) => fetchChatSession(session.session_id, AbortSignal.timeout(5000)))); + for (const result of results) { + if (result.status === "fulfilled") known.set(result.value.session.session_id, result.value); + } + const snapshots = listed.sessions.flatMap((session) => known.has(session.session_id) ? [known.get(session.session_id)!] : []); return { messages: mergeChatSessionMessages(snapshots), sessions: listed.sessions, snapshots, + unavailableSessionIds: listed.sessions.filter((session) => !known.has(session.session_id)).map((session) => session.session_id), }; } diff --git a/apps/presentation/dashboard/src/data/conversation-returns.test.mjs b/apps/presentation/dashboard/src/data/conversation-returns.test.mjs index 617cc913f4..cd1fcc3425 100644 --- a/apps/presentation/dashboard/src/data/conversation-returns.test.mjs +++ b/apps/presentation/dashboard/src/data/conversation-returns.test.mjs @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { conversationReturnSessions, reconcileConversationReturns } from "./conversation-returns.ts"; +import { conversationReturnSessions, reconcileConversationHistory, reconcileConversationReturns } from "./conversation-returns.ts"; const collaboration = { returns: [] }; const original = [ @@ -39,4 +39,32 @@ assert.equal(hydrated[0], original[0]); assert.deepEqual(conversationReturnSessions(undefined, delivered), ["current"]); assert.deepEqual(conversationReturnSessions(undefined, original), ["current", "old"]); assert.deepEqual(conversationReturnSessions(undefined, [delivered[0], delivered.at(-1)]), []); +const historyRows = [ + { session_id: "old", message_id: "answer", role: "agent", turn_id: "old-turn", text: "Older answer", created_at: "2026-08-01" }, + { session_id: "current", message_id: "answer", role: "agent", turn_id: "current-turn", text: "Stored answer", created_at: "2026-08-02", collaboration }, +]; +const live = [{ sourceSessionId: "current", sourceTurnId: "current-turn", text: "Live answer", pending: true }]; +const createHistory = (row) => ({ sourceSessionId: row.session_id, sourceMessageId: row.message_id, + sourceCreatedAt: row.created_at, text: row.text }); +const recoveredHistory = reconcileConversationHistory(live, historyRows, createHistory); +assert.equal(recoveredHistory.length, 2, "Same Turn keeps its existing live answer"); +assert.equal(recoveredHistory[0].text, "Older answer"); +assert.equal(recoveredHistory[1].text, "Live answer"); +assert.equal(recoveredHistory[1].pending, true); +assert.equal(recoveredHistory[1].collaboration, collaboration); +assert.equal(reconcileConversationHistory(recoveredHistory, historyRows, createHistory), recoveredHistory, "Repeated read is idempotent"); +assert.equal(reconcileConversationHistory(recoveredHistory, [historyRows[0]], createHistory), recoveredHistory, "An incomplete read never retracts known history"); +const bothRoles = [ + { sourceSessionId: "current", sourceTurnId: "same-turn", role: "user", text: "My request" }, + { sourceSessionId: "current", sourceTurnId: "same-turn", role: "assistant", text: "Live answer" }, +]; +const storedRoles = [ + { session_id: "current", message_id: "user-message", turn_id: "same-turn", role: "user", text: "My request" }, + { session_id: "current", message_id: "agent-message", turn_id: "same-turn", role: "agent", text: "Stored answer" }, +]; +const hydratedRoles = reconcileConversationHistory(bothRoles, storedRoles, createHistory); +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"); 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 a03dbcbb81..cea74aec3d 100644 --- a/apps/presentation/dashboard/src/data/conversation-returns.ts +++ b/apps/presentation/dashboard/src/data/conversation-returns.ts @@ -4,6 +4,8 @@ type ConversationMessage = { sourceSessionId?: string; sourceMessageId?: string; sourceTurnId?: string; + sourceCreatedAt?: string; + role?: string; collaboration?: ChatVisibleMessage["collaboration"]; returnDelivery?: ChatVisibleMessage["return_delivery"]; }; @@ -29,23 +31,25 @@ 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.role !== "user" && row.origin !== "manager_followup") - .map((row) => [row.turn_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 seen = new Set(previous.filter((row) => row.sourceSessionId === sessionId).map((row) => row.sourceMessageId)); let changed = false; const updated = previous.map((row) => { if (row.sourceSessionId !== sessionId) return row; - const source = row.sourceMessageId ? byId.get(row.sourceMessageId) : row.sourceTurnId ? byTurn.get(row.sourceTurnId) : undefined; + const source = row.sourceMessageId ? byId.get(row.sourceMessageId) : row.sourceTurnId + ? byTurn.get(`${row.sourceTurnId}:${row.role === "user" ? "user" : "assistant"}`) : undefined; if (!source) return row; // Projection absence is not a retraction: the backend can temporarily be // unable to read collaboration metadata. Keep the last observed receipt and // its outstanding read obligation until a newer observation arrives. const returnDelivery = source.return_delivery ?? row.returnDelivery; const collaboration = source.collaboration ?? row.collaboration; + const sourceCreatedAt = source.created_at ?? row.sourceCreatedAt; if (row.sourceMessageId === source.message_id && JSON.stringify(row.returnDelivery) === JSON.stringify(returnDelivery) - && JSON.stringify(row.collaboration) === JSON.stringify(collaboration)) return row; + && JSON.stringify(row.collaboration) === JSON.stringify(collaboration) && row.sourceCreatedAt === sourceCreatedAt) return row; changed = true; - return { ...row, sourceMessageId: source.message_id, returnDelivery, collaboration }; + return { ...row, sourceMessageId: source.message_id, sourceCreatedAt, returnDelivery, collaboration }; }); for (const message of messages) { if (message.origin !== "manager_followup" || seen.has(message.message_id)) continue; @@ -55,3 +59,19 @@ export function reconcileConversationReturns( } return changed ? updated : previous; } + +/** Recover missing history without retracting messages or overwriting live text. */ +export function reconcileConversationHistory( + previous: T[], messages: ChatVisibleMessage[], createMessage: (message: ChatVisibleMessage) => T, +): T[] { + let updated = previous; + const sessions = new Set(messages.map((message) => message.session_id).filter((id): id is string => Boolean(id))); + for (const sessionId of sessions) { + updated = reconcileConversationReturns(updated, sessionId, messages.filter((message) => message.session_id === sessionId), createMessage); + } + const seen = new Set(updated.map((message) => `${message.sourceSessionId}:${message.sourceMessageId}`)); + const recovered = messages.filter((message) => !seen.has(`${message.session_id}:${message.message_id}`)).map(createMessage); + if (!recovered.length) return updated; + return [...updated, ...recovered].sort((left, right) => !left.sourceCreatedAt ? (right.sourceCreatedAt ? 1 : 0) + : !right.sourceCreatedAt ? -1 : left.sourceCreatedAt.localeCompare(right.sourceCreatedAt)); +} diff --git a/apps/presentation/dashboard/src/data/use-conversation-history.ts b/apps/presentation/dashboard/src/data/use-conversation-history.ts new file mode 100644 index 0000000000..7eb601a95a --- /dev/null +++ b/apps/presentation/dashboard/src/data/use-conversation-history.ts @@ -0,0 +1,80 @@ +import { useCallback, useEffect, useRef, useState } from "react"; +import { fetchChatHistory, type ChatHistory } from "./chat"; + +type HistoryState = { + scope: string; + history: ChatHistory | null; + reading: boolean; + failed: boolean; +}; + +export type ConversationHistoryStatus = Pick, "phase" | "reading" | "sendBlocked" | "retry">; + +/** Read recovery shared by steward and Goal conversations; never executes a Turn. */ +export function useConversationHistory({ agentId, currentAgentId, channelId, goalId, enabled }: { + agentId?: string; + currentAgentId: string; + channelId: string; + goalId?: string; + enabled: boolean; +}) { + const scope = JSON.stringify([agentId ?? null, currentAgentId, channelId, goalId ?? null]); + const [state, setState] = useState(null); + const retryRead = useRef<(() => void) | null>(null); + const retry = useCallback(() => retryRead.current?.(), []); + useEffect(() => { + if (!enabled) return; + let cancelled = false; + let reading = false; + let history: ChatHistory | null = null; + let retryTimer: ReturnType | undefined; + let attempts = 0; + async function read() { + if (cancelled || reading) return; + clearTimeout(retryTimer); + reading = true; + setState({ scope, history, reading: true, failed: false }); + try { + const loaded = await fetchChatHistory({ agentId, channelId, goalId }, history ?? undefined); + if (cancelled) return; + history = loaded; + setState({ scope, history, reading: false, failed: false }); + } catch { + if (cancelled) return; + setState({ scope, history, reading: false, failed: true }); + } finally { + reading = false; + if (!cancelled && (!history || history.unavailableSessionIds.length > 0)) { + // Back off sustained transport failures; a user can retry immediately. + retryTimer = setTimeout(() => void read(), Math.min(3000 * 2 ** attempts++, 30_000)); + } + } + } + retryRead.current = () => void read(); + void read(); + return () => { + cancelled = true; + clearTimeout(retryTimer); + retryRead.current = null; + }; + }, [agentId, channelId, goalId, enabled, scope]); + const current = enabled && state?.scope === scope ? state : null; + const history = current?.history ?? null; + // A channel transcript spans executors; readability of another executor's + // latest session cannot authorize sending into the selected one. + const latest = history?.sessions.find((session) => session.agent_id === currentAgentId); + const currentSessionReadable = history !== null && (!latest || history.snapshots.some((snapshot) => snapshot.session.session_id === latest.session_id)); + return { + history, + currentSession: latest, + phase: !enabled ? "ready" as const : !current ? "loading" as const + : current.failed ? "unavailable" as const + : history?.unavailableSessionIds.length ? (history.snapshots.length ? "partial" as const : "unavailable" as const) + : current.reading ? "loading" as const : "ready" as const, + reading: current?.reading ?? enabled, + sendBlocked: enabled && !currentSessionReadable, + // Missing older history must not restart a recovered current stream. + connectionKey: currentSessionReadable ? `${scope}:${latest?.session_id ?? "empty"}` : null, + retry, + }; +} diff --git a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx index 8ea4878369..d89c924553 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx +++ b/apps/presentation/dashboard/src/features/personal-workspace/i18n.tsx @@ -118,6 +118,12 @@ const en = { "feedback.preparingPreview": "Preparing confirmation preview: {title}", "feedback.previewFailed": "Could not prepare the confirmation preview: {error}", "feedback.sendFailed": "Send failed: {error}", + "history.loading": "Reading conversation history…", + "history.partial": "Some history is temporarily unavailable. Available messages are shown.", + "history.unavailable": "Conversation history is temporarily unavailable. Retrying will only read records.", + "history.currentSessionRecovering": "Recovering the current conversation. You can send once it is restored.", + "history.retry": "Retry history", + "history.retrying": "Reading…", "feedback.sendGenericError": "Could not send the message. Try again later.", "feedback.stale": "State changed, so nothing was applied. Generate a new confirmation preview.", "feedback.taskDraftCreated": "Created a Task draft from the reply. Edit and send it to review the confirmation preview.", @@ -1289,6 +1295,12 @@ const zhCN: Record = { "feedback.preparingPreview": "正在准备确认预览:{title}", "feedback.previewFailed": "无法准备确认预览:{error}", "feedback.sendFailed": "发送失败:{error}", + "history.loading": "正在读取会话记录…", + "history.partial": "部分历史暂时无法读取,已显示可用消息。", + "history.unavailable": "会话历史暂时无法读取,重试只会读取记录。", + "history.currentSessionRecovering": "正在恢复当前会话,恢复后可以继续发送。", + "history.retry": "重试读取", + "history.retrying": "正在读取…", "feedback.sendGenericError": "消息发送失败,请稍后重试。", "feedback.stale": "状态已变化,操作未执行,请重新生成确认预览。", "feedback.taskDraftCreated": "已根据回复生成 Task 草稿。编辑后发送,LoopX 会先展示确认预览。", 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 a6b2deb0ae..9a7befd3ca 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 type { ConversationHistoryStatus } from "../../data/use-conversation-history"; import { GoalDraftCard } from "./goal-draft-card"; import type { GoalDraft } from "../../../../../../loopx/control_plane/collaboration/goal_draft.js"; import { CollaborationCard } from "./collaboration-card"; @@ -728,6 +729,7 @@ function readImageAttachment(file: File, t: WorkspaceTranslate): Promise
+ {conversationHistoryState && conversationHistoryState.phase !== "ready" ? ( +
+
{t(`history.${conversationHistoryState.phase}`)} + {conversationHistoryState.sendBlocked && conversationHistoryState.phase !== "loading" + ? {t("history.currentSessionRecovering")} : null}
+ {conversationHistoryState.phase !== "loading" ? : null} +
+ ) : null} {loopxMode?.session_id === conversationSessionId && loopxMode?.enabled && loopxMode.active_turn_id ? : null} {readOnly ? (
{t("source.readOnlyNoticeTitle")}{t("source.readOnlyNoticeDescription")}
@@ -1939,7 +1953,7 @@ export function PersonalWorkspacePage({ rows={1} value={composer} /> - +
} diff --git a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace.css b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace.css index 3eb57d6b5b..59ce12d24d 100644 --- a/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace.css +++ b/apps/presentation/dashboard/src/features/personal-workspace/personal-workspace.css @@ -60,6 +60,10 @@ .personal-status-source-meta > span.is-error i { background: #dd5b47; } .personal-status-source-meta > span.is-degraded i { background: #e5a33d; } .personal-service-notice { margin: 0; padding: 8px 16px; border-bottom: 1px solid var(--pw-line); background: var(--pw-amber-bg); color: var(--pw-amber); font-size: 13px; line-height: 20px; } +.personal-history-notice { display: flex; align-items: center; justify-content: space-between; gap: 12px; margin-bottom: 10px; color: var(--pw-muted); font-size: 13px; line-height: 1.5; } +.personal-history-notice > div { display: grid; gap: 2px; } +.personal-history-notice button { flex-shrink: 0; padding: 6px 10px; border: 1px solid var(--pw-line); border-radius: 8px; background: var(--pw-card); color: var(--pw-text); cursor: pointer; } +.personal-history-notice button:disabled { opacity: .6; cursor: wait; } .personal-status-source-meta small { min-width: 0; overflow: hidden; white-space: nowrap; text-overflow: ellipsis; } .personal-status-source-meta button { margin-left: auto; width: 22px; height: 22px; color: var(--pw-faint); } .personal-status-source-error { margin: 0 3px; color: var(--pw-red); font-size: 10px; line-height: 1.4; overflow-wrap: anywhere; } diff --git a/apps/presentation/dashboard/src/views/dashboard-page.tsx b/apps/presentation/dashboard/src/views/dashboard-page.tsx index c5812a88be..75e76bfc6c 100644 --- a/apps/presentation/dashboard/src/views/dashboard-page.tsx +++ b/apps/presentation/dashboard/src/views/dashboard-page.tsx @@ -1,5 +1,6 @@ import { normalizeGoalDraft, type GoalDraft } from "../../../../../loopx/control_plane/collaboration/goal_draft.js"; -import { conversationReturnSessions, reconcileConversationReturns } from "../data/conversation-returns"; +import { conversationReturnSessions, reconcileConversationHistory, reconcileConversationReturns } from "../data/conversation-returns"; +import { useConversationHistory } from "../data/use-conversation-history"; import {compactWorkspaceText as compactShareText} from "../features/personal-workspace/personal-workspace-model"; import type { GoalAcceptanceObservation } from "../data/goal-acceptance-observation"; import { attentionDetails, sourceAttention } from "../features/personal-workspace/attention-details"; @@ -43,7 +44,6 @@ import { updateLoopXMode, type LoopXModeSettings, fetchChatCapabilities, - fetchChatHistory, fetchChatSession, fetchChatSessions, interruptChatTurn, @@ -514,6 +514,7 @@ type PersonalManagerMessage = { sourceMessageId?: string; sourceSessionId?: string; sourceTurnId?: string; + sourceCreatedAt?: string; activity?: string[]; agentLabel?: string; attachments?: WorkspaceImageAttachment[]; @@ -1421,6 +1422,13 @@ function PersonalGoalHome({ const managerQuickPrompts = ["我现在该做什么?", "哪些 Goal 在等我?", "Agent 在做什么?"]; const contextMessages = messagesByContext[contextId] ?? []; const contextProposals = proposalsByContext[contextId] ?? []; + const conversationHistory = useConversationHistory({ + agentId: selectedGoal ? selectedAgent.agentId : undefined, + currentAgentId: selectedAgent.agentId, + channelId: selectedGoal ? `goal.${selectedGoal.goalId}` : "manager", + goalId: selectedGoal?.goalId, + enabled: !readOnly && selectedAgent.available, + }); // Who is speaking in the transcript. The manager channel answers as the LoopX // Manager: the executor that served the turn (and the model behind it) belongs @@ -1507,7 +1515,7 @@ function PersonalGoalHome({ const previous = current[contextId] ?? []; const updated = reconcileConversationReturns(previous, sessionId, snapshot.messages, (row) => ({ id: managerMessageId.current++, sourceMessageId: row.message_id, - sourceSessionId: sessionId, + sourceSessionId: sessionId, sourceCreatedAt: row.created_at, role: "assistant" as const, agentLabel: "协作回执", sourceLabel: "协作回执", text: visibleAgentMessage(row.text), lines: [], @@ -1585,56 +1593,46 @@ function PersonalGoalHome({ } }, [selectedAgents]); + useEffect(() => { + if (!conversationHistory.history) return; + const messages = conversationHistory.history.messages; + setMessagesByContext((current) => { + const previous = current[contextId] ?? []; + const updated = reconcileConversationHistory(previous, messages, (message): PersonalManagerMessage => ({ + sourceMessageId: message.message_id, + sourceSessionId: message.session_id, + sourceTurnId: message.turn_id ?? undefined, + sourceCreatedAt: message.created_at, + goalDraft: normalizeGoalDraft(message.goal_draft), + agentLabel: message.role === "user" ? undefined : message.origin === "manager_followup" + ? "协作回执" : answerIdentityLabel(contextId, selectedAgent.label), + attachments: workspaceImageAttachments(message.attachments), + id: managerMessageId.current++, lines: [], + role: message.role === "user" ? "user" : "assistant", + returnDelivery: message.return_delivery, collaboration: message.collaboration, + sourceLabel: message.role === "user" ? undefined : message.role === "error" ? "本地会话记录" + : contextId === "manager" ? `恢复的${t("header.manager")}会话` : `恢复的 ${selectedAgent.label} 会话`, + text: message.role === "user" ? message.text : visibleAgentMessage(message.text), + })); + return updated === previous ? current : { ...current, [contextId]: updated }; + }); + }, [conversationHistory.history, contextId, selectedAgent.label]); + useEffect(() => { if (readOnly) return; if (!selectedAgent.available) return; + if (!conversationHistory.connectionKey || !conversationHistory.history) return; + const history = conversationHistory.history; const targetContextId = contextId; const sessionKey = `${targetContextId}:${selectedAgent.agentId}`; const contextKind = selectedGoal ? "goal" : "manager"; - const channelId = selectedGoal ? `goal.${selectedGoal.goalId}` : "manager"; let cancelled = false; let recoveryController: AbortController | null = null; let latestDiscoveredSessionId: string | null = null; void (async () => { try { - const history = await fetchChatHistory({ - agentId: contextKind === "manager" ? undefined : selectedAgent.agentId, - channelId, - goalId: selectedGoal?.goalId, - }); - if (cancelled) return; - setMessagesByContext((current) => { - if ((current[targetContextId]?.length ?? 0) > 0) return current; - return { - ...current, - [targetContextId]: history.messages.map((message) => ({ - sourceMessageId: message.message_id, - goalDraft: normalizeGoalDraft(message.goal_draft), - sourceSessionId: message.session_id, - agentLabel: message.role === "user" - ? undefined - : message.origin === "manager_followup" - ? "协作回执" - : answerIdentityLabel(targetContextId, selectedAgent.label), - attachments: workspaceImageAttachments(message.attachments), - id: managerMessageId.current++, - lines: [], - role: message.role === "user" ? "user" : "assistant", - returnDelivery: message.return_delivery, - collaboration: message.collaboration, - sourceLabel: message.role === "user" - ? undefined - : message.role === "error" - ? "本地会话记录" - : targetContextId === "manager" - ? `恢复的${t("header.manager")}会话` - : `恢复的 ${selectedAgent.label} 会话`, - text: message.role === "user" ? message.text : visibleAgentMessage(message.text), - })), - }; - }); if (selectedAgent.agentId === "status-only") return; - const latest = history.sessions[0]; + const latest = conversationHistory.currentSession; latestDiscoveredSessionId = latest?.session_id ?? null; if (latest && !latest.resumable) { newSessionRequired.current.add(sessionKey); @@ -1708,7 +1706,7 @@ function PersonalGoalHome({ signal: recoveryController.signal, onDelta: (delta) => { streamedText += delta; - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { text: streamedText, }); }, @@ -1727,7 +1725,7 @@ function PersonalGoalHome({ }, }); if (cancelled) return; - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { lines: streamed.response.gate ? [streamed.response.gate.summary, streamed.response.gate.next_action].filter(Boolean).slice(0, 2) : [], @@ -1766,7 +1764,7 @@ function PersonalGoalHome({ if (cancelled) return; const interrupted = interruptedTurnIds.current.delete(activeTurnId) || (error instanceof ChatApiError && error.payload.error_code === "turn_interrupted"); - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { lines: [], pending: false, reconnect: error instanceof ChatApiError && error.payload.reconnectable === true, @@ -1810,7 +1808,9 @@ function PersonalGoalHome({ cancelled = true; recoveryController?.abort(); }; - }, [contextId, model.goals[0]?.goalId, readOnly, selectedGoal?.goalId, selectedAgent.agentId, selectedAgent.available, selectedAgent.label, selectedAgents]); + // Recovery of an older transcript only hydrates messages. Reconnect the + // executor when the current session becomes readable, not on every retry. + }, [conversationHistory.connectionKey, contextId, model.goals[0]?.goalId, readOnly, selectedGoal?.goalId, selectedAgent.agentId, selectedAgent.available, selectedAgent.label, selectedAgents]); useEffect(() => { if (readOnly) return; @@ -1982,16 +1982,19 @@ function PersonalGoalHome({ return id; } - function updateManagerAssistantMessage( + function updateConversationMessage( targetContextId: string, messageId: number, update: Partial>, ) { setMessagesByContext((messages) => ({ ...messages, - [targetContextId]: (messages[targetContextId] ?? []).map((message) => - message.id === messageId ? { ...message, ...update } : message - ), + [targetContextId]: (messages[targetContextId] ?? []).filter((message) => + // A read may arrive before the original send's storage receipt. Keep + // the live row when its exact persisted identity becomes known. + message.id === messageId || !update.sourceMessageId || !update.sourceSessionId + || message.sourceMessageId !== update.sourceMessageId || message.sourceSessionId !== update.sourceSessionId + ).map((message) => message.id === messageId ? { ...message, ...update } : message), })); } @@ -2059,7 +2062,7 @@ function PersonalGoalHome({ if (selectedRoute.agentId === "status-only" || (!targetGoal && targetContextId !== "manager")) { const answer = personalManagerSnapshot(targetQuestionModel); const usesStatusOnlyRoute = selectedRoute.agentId === "status-only"; - appendManagerAssistantMessage(targetContextId, { + const answerMessageId = appendManagerAssistantMessage(targetContextId, { agentLabel: usesStatusOnlyRoute ? "仅查状态" : "LoopX 管家", lines: answer.lines.slice(0, 3), sourceLabel: usesStatusOnlyRoute ? "LoopX 状态投影 · 仅查状态" : "LoopX 状态投影", @@ -2070,6 +2073,13 @@ function PersonalGoalHome({ contextKind: targetContextId === "manager" ? "manager" : "goal", goalId: targetContextId === "manager" ? undefined : targetContextId, question, + }).then((receipt) => { + updateConversationMessage(targetContextId, userMessageId, { + sourceSessionId: receipt.session_id, sourceMessageId: receipt.user_message_id, + }); + updateConversationMessage(targetContextId, answerMessageId, { + sourceSessionId: receipt.session_id, sourceMessageId: receipt.answer_message_id, + }); }).catch(() => { // The current projection answer remains visible when local history persistence is unavailable. }); @@ -2126,7 +2136,7 @@ function PersonalGoalHome({ onDelta: (delta: string) => { streamedText += delta; if (streamingMessageId !== null) { - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { text: streamedText, }); } @@ -2146,8 +2156,11 @@ function PersonalGoalHome({ })); }, onPhase: (_phase: string, turnId: string) => { + if (submittedTurnId !== turnId) updateConversationMessage(targetContextId, userMessageId, { + sourceTurnId: turnId, sourceSessionId: sessionId, + }); submittedTurnId = turnId; - if (streamingMessageId !== null) updateManagerAssistantMessage(targetContextId, streamingMessageId, { sourceTurnId: turnId, sourceSessionId: sessionId }); + if (streamingMessageId !== null) updateConversationMessage(targetContextId, streamingMessageId, { sourceTurnId: turnId, sourceSessionId: sessionId }); activeTurnIds.current.set(targetContextId, turnId); recordRuntimeBinding(targetContextId, { agentId: selectedRoute.agentId, @@ -2168,7 +2181,7 @@ function PersonalGoalHome({ streamed = await sendChatTurnStreaming(sessionId, question, streamOptions); } const response = streamed.response; - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { goalDraft: response.goal_draft, lines: response.gate ? [response.gate.summary, response.gate.next_action].filter(Boolean).slice(0, 2) : [], pending: false, @@ -2183,7 +2196,7 @@ function PersonalGoalHome({ void fetchChatSession(sessionId).then((stored) => { const answer = stored.messages.find((item) => item.turn_id === streamed.turnId && ["agent", "assistant"].includes(item.role)); - if (answer) updateManagerAssistantMessage(targetContextId, completedMessageId, { + if (answer) updateConversationMessage(targetContextId, completedMessageId, { sourceMessageId: answer.message_id, sourceSessionId: sessionId, collaboration: answer.collaboration, returnDelivery: answer.return_delivery, }); @@ -2194,7 +2207,7 @@ function PersonalGoalHome({ // names the Goal whose workspace holds the card, so a manager-channel // proposal here is only ever a Todo the owner has to be sent to. if (todoProposals.length > 0 && !targetGoal) { - updateManagerAssistantMessage(targetContextId, streamingMessageId, { + updateConversationMessage(targetContextId, streamingMessageId, { lines: ["请进入要修改的 Goal,预览并确认具体变更。"], }); } @@ -2237,7 +2250,7 @@ function PersonalGoalHome({ if (streamingMessageId === null) { appendManagerAssistantMessage(targetContextId, interruptedMessage); } else { - updateManagerAssistantMessage(targetContextId, streamingMessageId, interruptedMessage); + updateConversationMessage(targetContextId, streamingMessageId, interruptedMessage); } return; } @@ -2274,7 +2287,7 @@ function PersonalGoalHome({ if (streamingMessageId === null) { appendManagerAssistantMessage(targetContextId, failureMessage); } else { - updateManagerAssistantMessage(targetContextId, streamingMessageId, failureMessage); + updateConversationMessage(targetContextId, streamingMessageId, failureMessage); } } finally { activeTurnIds.current.delete(targetContextId); @@ -2712,10 +2725,10 @@ function PersonalGoalHome({ signal: controller.signal, onDelta: (delta) => { streamedText += delta; - updateManagerAssistantMessage(run.goalId, messageId, { text: streamedText }); + updateConversationMessage(run.goalId, messageId, { text: streamedText }); }, }); - updateManagerAssistantMessage(run.goalId, messageId, { + updateConversationMessage(run.goalId, messageId, { activity: [], pending: false, text: visibleAgentMessage(streamed.response.message || streamedText.trim()) || `${run.agentLabel} 已完成纠偏。`, @@ -2723,7 +2736,7 @@ function PersonalGoalHome({ } catch (error) { const interrupted = interruptedTurnIds.current.delete(turnId) || (error instanceof ChatApiError && error.payload.error_code === "turn_interrupted"); - updateManagerAssistantMessage(run.goalId, messageId, { + updateConversationMessage(run.goalId, messageId, { activity: [], pending: false, text: interrupted ? [streamedText.trim(), "已中断。你可以在当前会话继续发送消息。"].filter(Boolean).join("\n\n") : error instanceof Error ? error.message : "纠偏回合失败。", @@ -2897,6 +2910,7 @@ function PersonalGoalHome({ managerChannelBinding={managerChannelBinding} managerRuntime={managerRuntime} conversationSessionId={runtimeBindings[contextId]?.sessionId} + conversationHistoryState={conversationHistory} model={workspaceModel} readOnly={readOnly} selectedAgentId={selectedAgent.agentId} 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 c8b1c20fb8..e817bdd465 100644 --- a/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md +++ b/docs/architecture/rfcs/app-conversation-and-async-inbox-v0.md @@ -249,6 +249,20 @@ backs off without resending, and typed terminal blockers stay stopped. This qualifies the persisted recovery boundary, not live provider availability or the complete GQ09 journey. +App history recovery is a shared TS read boundary for steward and Goal channels. +One unavailable older Session must not hide readable messages, lose their +original result locators or restart the current stream. Show incomplete history, +retry missing snapshots with backoff, and preserve live text and drafts. An +unreadable current Session blocks ordinary send until its exact record recovers; +never infer an empty conversation or create a replacement driver from a read +failure. A recovered read does not prove receiver adoption or native GQ02/GQ09. + +App 历史恢复由管家与 Goal 对话共用的 TS 读取边界负责。旧 Session 读取失败不应 +隐藏可用消息、丢失原结果位置或重启当前流;应明确展示历史不完整,退避重读缺失 +快照,并保留实时文本与草稿。当前 Session 无法读取时,普通发送等待原记录恢复, +不能据此推断空会话或新建执行驱动。读取恢复仍不代表接收方采用或原生 GQ02/GQ09 +已经验收。 + Before each extraction report base/head real-call latency, boundary crossings, bytes, owners deleted/retained and compatibility callers. Product delivery must not wait for full Python retirement. Python may retain IO; TS owns migrated diff --git a/examples/loopx-chat-server-smoke.py b/examples/loopx-chat-server-smoke.py index 6a65b16da0..6f0048843c 100644 --- a/examples/loopx-chat-server-smoke.py +++ b/examples/loopx-chat-server-smoke.py @@ -664,6 +664,10 @@ def main() -> None: "现在该做什么?", "当前没有阻塞。", ], projection_snapshot + assert [message["message_id"] for message in projection_snapshot["messages"]] == [ + projection_exchange["user_message_id"], + projection_exchange["answer_message_id"], + ], projection_snapshot code, turn = request_json( f"{base_url}/api/chat/sessions/{session_id}/turns", diff --git a/examples/personal-workspace-browser-smoke.mjs b/examples/personal-workspace-browser-smoke.mjs index d1cc4b1b8a..3cfff873f9 100644 --- a/examples/personal-workspace-browser-smoke.mjs +++ b/examples/personal-workspace-browser-smoke.mjs @@ -13,6 +13,7 @@ import { writeDashboardBrowserCoverage } from "./dashboard-browser-coverage.mjs" import { conversationActivityScenario } from "./personal-workspace-browser/conversation-activity.mjs"; import { chatRecoveryScenario } from "./personal-workspace-browser/chat-recovery.mjs"; import { conversationReturnContinuityScenario } from "./personal-workspace-browser/conversation-return-continuity.mjs"; +import { conversationHistoryRecoveryScenario } from "./personal-workspace-browser/conversation-history-recovery.mjs"; import { executionChipScenario } from "./personal-workspace-browser/execution-chip.mjs"; import { collectCoverage, @@ -47,7 +48,7 @@ import { goalDraftScenario } from "./personal-workspace-browser/goal-draft.mjs"; import { larkCliMissingScenario } from "./personal-workspace-browser/lark-cli-missing.mjs"; import { executionServiceOfflineScenario } from "./personal-workspace-browser/execution-service-offline.mjs"; -const scenarioCatalog = [goalDraftScenario, capabilityScopeScenario, stewardGroupTriggerScenario, conversationInputScenario, goalActivityScenario, conversationActivityScenario, navigationSortingScenario, automationCadenceScenario, chatRecoveryScenario, conversationReturnContinuityScenario, answerPresentationScenario, loopxModeScenario, teamEvidenceScenario, managedGoalResultsScenario, typedActionsScenario, teamPlanScenario, stewardJourneyScenario, executionChipScenario, stewardModelSettingsScenario, progressiveLoadingScenario, workspaceLocaleScenario, newestDraftScenario, larkCliMissingScenario, executionServiceOfflineScenario]; +const scenarioCatalog = [goalDraftScenario, capabilityScopeScenario, stewardGroupTriggerScenario, conversationInputScenario, goalActivityScenario, conversationActivityScenario, navigationSortingScenario, automationCadenceScenario, chatRecoveryScenario, conversationReturnContinuityScenario, conversationHistoryRecoveryScenario, answerPresentationScenario, loopxModeScenario, teamEvidenceScenario, managedGoalResultsScenario, typedActionsScenario, teamPlanScenario, stewardJourneyScenario, executionChipScenario, stewardModelSettingsScenario, progressiveLoadingScenario, workspaceLocaleScenario, newestDraftScenario, larkCliMissingScenario, executionServiceOfflineScenario]; const requestedScenario = process.env.LOOPX_PERSONAL_WORKSPACE_SCENARIO; const scenarios = requestedScenario ? scenarioCatalog.filter((scenario) => scenario.id === requestedScenario) diff --git a/examples/personal-workspace-browser/conversation-history-recovery.mjs b/examples/personal-workspace-browser/conversation-history-recovery.mjs new file mode 100644 index 0000000000..62a17b5524 --- /dev/null +++ b/examples/personal-workspace-browser/conversation-history-recovery.mjs @@ -0,0 +1,102 @@ +import assert from "node:assert/strict"; +import { resolve } from "node:path"; +import { outputDir } from "./fixture.mjs"; +import { openWorkspacePage } from "./scenario-context.mjs"; + +export const conversationHistoryRecoveryScenario = { + id: "conversation-history-recovery", + async run({ browser, collectCoverage, url }) { + const currentId = "session-manager-loopx-manager-codex"; + const oldId = "previous-worker-session"; + let unavailable = oldId; + let failedReads = 0; + const context = await openWorkspacePage(browser, url, { + collectCoverage, + beforeGoto(_api, page) { + for (const [id, agentId, date, text] of [ + ["other-executor-session", "claude-code", "2026-08-14T01:00:00Z", "另一个执行器保留了独立对话。"], + [currentId, "codex", "2026-08-13T01:00:00Z", "目前正在核对微软现金流。"], + [oldId, "codex", "2026-08-12T01:00:00Z", "上一轮已经找到公开财报。"], + ]) { + page.__loopxRuntime.sessions.set(id, { + session_id: id, goal_id: "loopx-manager", agent_id: agentId, adapter_kind: agentId, + channel_id: "manager", status: "ready", active_turn_id: null, resumable: true, + created_at: date, updated_at: date, last_activity_at: date, last_error_code: null, + }); + // Message identity is scoped to a session, even if IDs collide. + page.__loopxRuntime.messages.set(id, [{ + message_id: "answer", turn_id: `turn-${id}`, role: "agent", text, created_at: date, + }]); + } + return page.route("**/api/chat/sessions/*", async (route) => { + if (route.request().method() === "GET" && new URL(route.request().url()).pathname === `/api/chat/sessions/${unavailable}`) { + failedReads += 1; + await route.fulfill({ status: 503, json: { ok: false, error: "Synthetic history read unavailable" } }); + } else await route.fallback(); + }); + }, + }); + const { page, api } = context; + try { + await page.getByRole("navigation", { name: "管家视图" }).getByRole("button", { name: /^(Chat|对话)$/ }).click(); + await page.getByText("目前正在核对微软现金流。", { exact: true }).waitFor({ state: "visible", timeout: 5000 }); + const notice = page.getByRole("status").filter({ hasText: "部分历史暂时无法读取" }); + await notice.waitFor({ state: "visible" }); + await page.screenshot({ path: resolve(outputDir, "conversation-history-partial-desktop.png"), animations: "disabled" }); + await page.setViewportSize({ width: 390, height: 844 }); + assert.ok(await notice.isVisible(), "Recovery feedback stays visible at narrow width"); + assert.ok(await page.locator("body").evaluate(body => body.scrollWidth <= innerWidth + 1), "Recovery must not add horizontal overflow"); + await page.screenshot({ path: resolve(outputDir, "conversation-history-partial-mobile.png"), animations: "disabled" }); + await page.setViewportSize({ width: 1512, height: 982 }); + assert.ok(failedReads > 0, "The unavailable historical read was exercised"); + assert.equal(api.turnRequests.length, 0, "Reading history must not start work"); + unavailable = ""; + await notice.getByRole("button", { name: "重试读取" }).click(); + await page.getByText("上一轮已经找到公开财报。", { exact: true }).waitFor({ state: "visible" }); + await notice.waitFor({ state: "hidden" }); + assert.equal(await page.getByText("目前正在核对微软现金流。", { exact: true }).count(), 1); + assert.equal(api.turnRequests.length, 0, "Retry must not replay a model turn"); + + // The current session's unreadability is a different boundary: keep the + // known transcript, but do not send into a guessed replacement session. + unavailable = currentId; + await page.reload({ waitUntil: "networkidle" }); + await page.getByRole("navigation", { name: "管家视图" }).getByRole("button", { name: /^(Chat|对话)$/ }).click(); + await page.getByText("上一轮已经找到公开财报。", { exact: true }).waitFor({ state: "visible" }); + await page.getByText("正在恢复当前会话,恢复后可以继续发送。", { exact: true }).waitFor({ state: "visible" }); + await page.getByLabel("向 LoopX 发送消息").fill("接着昨天的做。"); + assert.equal(await page.getByRole("button", { name: "发送", exact: true }).isEnabled(), false); + await page.getByLabel("向 LoopX 发送消息").press("Enter"); + assert.equal(api.turnRequests.length, 0, "Enter must not bypass current-session recovery"); + unavailable = ""; + await page.getByRole("button", { name: "重试读取", exact: true }).click(); + await page.getByRole("button", { name: "发送", exact: true }).waitFor({ state: "visible" }); + await page.waitForFunction(() => !document.querySelector('button[aria-label="发送"]')?.disabled); + assert.equal(await page.getByLabel("向 LoopX 发送消息").inputValue(), "接着昨天的做。", "Recovery retains the draft"); + await page.getByRole("button", { name: "发送", exact: true }).click(); + await page.getByText("已沿用当前 Goal 与 Agent Session。接下来会先核对状态,再继续推进。", { exact: true }).waitFor({ state: "visible" }); + assert.equal(api.turnRequests.length, 1); + assert.equal(api.turnRequests[0].sessionId, currentId); + await page.locator(".personal-goal-link").first().click(); + await page.locator(".personal-manager-link").click(); + await page.getByRole("navigation", { name: "管家视图" }).getByRole("button", { name: /^(Chat|对话)$/ }).click(); + await page.getByText("接着昨天的做。", { exact: true }).waitFor({ state: "visible" }); + assert.equal(await page.getByText("接着昨天的做。", { exact: true }).count(), 1, "Rehydration must not duplicate the user's accepted request"); + assert.equal(api.turnRequests.length, 1, "Returning to the conversation must not repeat work"); + await page.getByRole("combobox", { name: "选择聊天 Runtime", exact: true }).click(); + await page.getByRole("option", { name: "仅查状态", exact: true }).click(); + await page.getByLabel("向 LoopX 发送消息").fill("只看当前状态。"); + const projectionSaved = page.waitForResponse(response => new URL(response.url()).pathname === "/api/chat/projection-messages" && response.status() === 201); + await page.getByRole("button", { name: "发送", exact: true }).click(); + await projectionSaved; + await page.locator(".personal-goal-link").first().click(); + await page.locator(".personal-manager-link").click(); + await page.getByRole("navigation", { name: "管家视图" }).getByRole("button", { name: /^(Chat|对话)$/ }).click(); + assert.equal(await page.getByText("只看当前状态。", { exact: true }).count(), 1, "Projection-only messages retain their exact stored identities"); + assert.equal(api.turnRequests.length, 1, "Projection history does not run an Agent"); + return { coverageEntries: context.coverageEntries, note: "One failed history read preserves readable messages; read-only retry recovers session-scoped records, retains the draft and sends exactly once into the recovered current session." }; + } finally { + await context.close(); + } + }, +}; diff --git a/examples/personal-workspace-browser/fixture.mjs b/examples/personal-workspace-browser/fixture.mjs index ccaad4d719..c0aa3e5b3a 100644 --- a/examples/personal-workspace-browser/fixture.mjs +++ b/examples/personal-workspace-browser/fixture.mjs @@ -1468,6 +1468,24 @@ export async function installApi(page, { goalSubagentConfigurationEnabled = true }, status: 200 }); return; } + if (url.pathname === "/api/chat/projection-messages" && request.method() === "POST") { + const body = request.postDataJSON(); + const channel_id = body.context_kind === "manager" ? "manager" : `goal.${body.goal_id}`; + const session_id = `session-projection-${channel_id}`; + if (!sessions.has(session_id)) sessions.set(session_id, { + session_id, goal_id: body.goal_id || "loopx-manager", agent_id: "status-only", adapter_kind: "status_projection", + channel_id, status: "ready", active_turn_id: null, resumable: true, + created_at: "2026-08-13T01:00:00Z", updated_at: "2026-08-13T01:00:00Z", last_activity_at: "2026-08-13T01:00:00Z", last_error_code: null, + }); + const rows = messages.get(session_id) ?? []; + const user_message_id = `${session_id}-${rows.length}`; + const answer_message_id = `${session_id}-${rows.length + 1}`; + rows.push({ message_id: user_message_id, role: "user", text: body.question, created_at: "2026-08-13T01:00:01Z" }, + { message_id: answer_message_id, role: "agent", text: body.answer, created_at: "2026-08-13T01:00:02Z" }); + messages.set(session_id, rows); + await route.fulfill({ status: 201, json: { ok: true, schema_version: "loopx_chat_projection_exchange_v1", session_id, user_message_id, answer_message_id } }); + return; + } if (url.pathname === "/api/chat/sessions" && request.method() === "GET") { const requestedGoal = url.searchParams.get("goal_id"); const requestedAgent = url.searchParams.get("agent_id"); @@ -1484,12 +1502,14 @@ export async function installApi(page, { goalSubagentConfigurationEnabled = true const body = request.postDataJSON(); if (body.context_kind === "manager" && body.goal_id) throw new Error("Global manager request still carries a project anchor"); const resolvedGoalId = body.context_kind === "manager" ? "loopx-manager" : body.goal_id; - const session_id = `session-${body.context_kind}-${resolvedGoalId}-${body.agent_id}`; + // Like the real runtime, an omitted steward endpoint resolves on the host. + 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: body.agent_id, adapter_kind: body.agent_id, 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, 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: body.agent_id, goal_id: body.goal_id, resumed: body.mode === "resume_latest", session_id, session }, status: 201 }); + 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 }); return; } const loopxMode = url.pathname.match(/^\/api\/chat\/sessions\/([^/]+)\/loopx$/); diff --git a/loopx/chat_server.py b/loopx/chat_server.py index d2818075fa..9a8c3b6e1c 100644 --- a/loopx/chat_server.py +++ b/loopx/chat_server.py @@ -646,8 +646,8 @@ def _record_projection_exchange(self) -> None: channel_id=channel_id, ) session_id = str(session["session_id"]) - self.server.chat_store.append_message(session_id, role="user", text=question) - self.server.chat_store.append_message(session_id, role="agent", text=answer) + user_message = self.server.chat_store.append_message(session_id, role="user", text=question) + answer_message = self.server.chat_store.append_message(session_id, role="agent", text=answer) self.server.chat_store.update_session(session_id, last_activity_at=time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())) except Exception as exc: # noqa: BLE001 - local validation response. self._send_error(str(exc)) @@ -657,6 +657,8 @@ def _record_projection_exchange(self) -> None: "ok": True, "schema_version": "loopx_chat_projection_exchange_v1", "session_id": session_id, + "user_message_id": user_message["message_id"], + "answer_message_id": answer_message["message_id"], }, status=201, ) diff --git a/loopx/semantics/project_registry_io_manifest_v1.json b/loopx/semantics/project_registry_io_manifest_v1.json index 78b65b58d2..04a1ce4351 100644 --- a/loopx/semantics/project_registry_io_manifest_v1.json +++ b/loopx/semantics/project_registry_io_manifest_v1.json @@ -439,7 +439,7 @@ }, { "site": "loopx/chat_server.py::.ChatRequestHandler._goal_channel_extension_ready::codec_read:load_registry#1", - "line": 964, + "line": 966, "column": 24, "kind": "codec_read", "api": "load_registry", @@ -455,7 +455,7 @@ }, { "site": "loopx/chat_server.py::.serve_chat::codec_read:load_registry#1", - "line": 1486, + "line": 1488, "column": 16, "kind": "codec_read", "api": "load_registry",