From 8610d90838ba9fb3fff236f4cba6e41ce58bc980 Mon Sep 17 00:00:00 2001 From: Luke Parker <10430890+Hona@users.noreply.github.com> Date: Sat, 15 Aug 2026 15:57:46 +1000 Subject: [PATCH] fix(app): project sent messages through inbox events (#42714) --- .../prompt-input/build-prompt-request.test.ts | 321 ++++++++++ .../prompt-input/build-prompt-request.ts | 115 ++++ .../prompt-input/build-request-parts.test.ts | 396 ------------ .../prompt-input/build-request-parts.ts | 216 ------- .../components/prompt-input/submit.test.ts | 101 +++- .../app/src/components/prompt-input/submit.ts | 86 +-- packages/app/src/context/directory-sync.ts | 38 +- packages/app/src/context/file/path.test.ts | 2 +- .../context/server-session-v2-reducer.test.ts | 26 + .../src/context/server-session-v2-reducer.ts | 17 +- .../app/src/context/server-session.test.ts | 569 ++++++++++-------- packages/app/src/context/server-session.ts | 501 ++++++++------- packages/app/src/context/server-sync.tsx | 5 +- .../app/src/context/sync-optimistic.test.ts | 137 ----- packages/app/src/context/sync.tsx | 108 ---- 15 files changed, 1140 insertions(+), 1498 deletions(-) create mode 100644 packages/app/src/components/prompt-input/build-prompt-request.test.ts create mode 100644 packages/app/src/components/prompt-input/build-prompt-request.ts delete mode 100644 packages/app/src/components/prompt-input/build-request-parts.test.ts delete mode 100644 packages/app/src/components/prompt-input/build-request-parts.ts delete mode 100644 packages/app/src/context/sync-optimistic.test.ts diff --git a/packages/app/src/components/prompt-input/build-prompt-request.test.ts b/packages/app/src/components/prompt-input/build-prompt-request.test.ts new file mode 100644 index 00000000000..58868bc3caa --- /dev/null +++ b/packages/app/src/components/prompt-input/build-prompt-request.test.ts @@ -0,0 +1,321 @@ +import { describe, expect, test } from "bun:test" +import type { Prompt } from "@/context/prompt" +import { buildPromptRequest } from "./build-prompt-request" + +describe("buildPromptRequest", () => { + test("builds text, files, and agents from the prompt", () => { + const prompt: Prompt = [ + { type: "text", content: "hello", start: 0, end: 5 }, + { + type: "file", + path: "src/foo.ts", + content: "@src/foo.ts", + start: 5, + end: 16, + selection: { startLine: 4, startChar: 1, endLine: 6, endChar: 1 }, + }, + { type: "agent", name: "planner", content: "@planner", start: 16, end: 24 }, + ] + + const result = buildPromptRequest({ + prompt, + context: [{ key: "ctx:1", type: "file", path: "src/bar.ts", comment: "check this" }], + images: [ + { type: "image", id: "img_1", filename: "a.png", mime: "image/png", dataUrl: "data:image/png;base64,AAA" }, + ], + text: "hello @src/foo.ts @planner", + sessionDirectory: "/repo", + }) + + expect(result.text).toContain("hello @src/foo.ts @planner") + expect(result.text).toContain("check this") + expect(result.displayText).toBe("hello @src/foo.ts @planner") + expect(result.comments).toMatchObject([{ path: "src/bar.ts", comment: "check this" }]) + expect(result.agents).toEqual([{ name: "planner", mention: { start: 16, end: 24, text: "@planner" } }]) + expect(result.files.some((file) => file.uri.startsWith("file:///repo/src/foo.ts"))).toBe(true) + expect(result.files.find((file) => file.uri.startsWith("file:///repo/src/foo.ts"))?.mention).toEqual({ + start: 5, + end: 16, + text: "@src/foo.ts", + }) + }) + + test("keeps multiple uploaded attachments in order", () => { + const result = buildPromptRequest({ + prompt: [{ type: "text", content: "check these", start: 0, end: 11 }], + context: [], + images: [ + { type: "image", id: "img_1", filename: "a.png", mime: "image/png", dataUrl: "data:image/png;base64,AAA" }, + { + type: "image", + id: "img_2", + filename: "b.pdf", + mime: "application/pdf", + dataUrl: "data:application/pdf;base64,BBB", + }, + ], + text: "check these", + sessionDirectory: "/repo", + }) + + const uploads = result.files.filter((file) => file.uri.startsWith("data:")) + + expect(uploads).toHaveLength(2) + expect(uploads.map((file) => file.name)).toEqual(["a.png", "b.pdf"]) + }) + + test("preserves an external attachment source path for the model", () => { + const result = buildPromptRequest({ + prompt: [], + context: [], + images: [ + { + type: "image", + id: "img_external", + filename: "opencode.global.dat", + sourcePath: "C:\\Users\\Luke\\AppData\\Roaming\\ai.opencode.desktop.beta\\opencode.global.dat", + mime: "text/plain", + dataUrl: "data:text/plain;base64,AAA", + }, + ], + text: "inspect this", + sessionDirectory: "C:\\Repos\\sst\\opencode", + }) + + expect(result.files[0]?.name).toBe( + "C:\\Users\\Luke\\AppData\\Roaming\\ai.opencode.desktop.beta\\opencode.global.dat", + ) + }) + + test("preserves reference aliases as directory files", () => { + const result = buildPromptRequest({ + prompt: [ + { + type: "file", + path: "/repo/../docs", + content: "@docs", + start: 0, + end: 5, + mime: "application/x-directory", + filename: "docs", + }, + ], + context: [], + images: [], + text: "@docs", + sessionDirectory: "/repo/app", + }) + + expect(result.files[0]).toEqual({ + uri: "file:///repo/../docs", + mime: "application/x-directory", + name: "docs", + mention: { start: 0, end: 5, text: "@docs" }, + }) + }) + + test("deduplicates context files when prompt already includes same path", () => { + const prompt: Prompt = [{ type: "file", path: "src/foo.ts", content: "@src/foo.ts", start: 0, end: 11 }] + + const result = buildPromptRequest({ + prompt, + context: [ + { key: "ctx:dup", type: "file", path: "src/foo.ts" }, + { key: "ctx:comment", type: "file", path: "src/foo.ts", comment: "focus here" }, + ], + images: [], + text: "@src/foo.ts", + sessionDirectory: "/repo", + }) + + const fooFiles = result.files.filter((file) => file.uri.startsWith("file:///repo/src/foo.ts")) + + expect(fooFiles).toHaveLength(2) + expect(result.text).toContain("focus here") + }) + + test("adds files for @mentions inside comment text", () => { + const result = buildPromptRequest({ + prompt: [{ type: "text", content: "look", start: 0, end: 4 }], + context: [ + { + key: "ctx:comment-mention", + type: "file", + path: "src/review.ts", + comment: "Compare with @src/shared.ts and @src/review.ts.", + }, + ], + images: [], + text: "look", + sessionDirectory: "/repo", + }) + + expect(result.files).toHaveLength(2) + expect(result.files.some((file) => file.uri === "file:///repo/src/review.ts")).toBe(true) + expect(result.files.some((file) => file.uri === "file:///repo/src/shared.ts")).toBe(true) + }) + + test("handles Windows paths correctly (simulated on macOS)", () => { + const prompt: Prompt = [{ type: "file", path: "src\\foo.ts", content: "@src\\foo.ts", start: 0, end: 11 }] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@src\\foo.ts", + sessionDirectory: "D:\\projects\\myapp", // Windows path + }) + + const file = result.files[0] + expect(file).toBeDefined() + // URL should be parseable + expect(() => new URL(file!.uri)).not.toThrow() + // Should not have encoded backslashes in wrong place + expect(file!.uri).not.toContain("%5C") + // Should have normalized to forward slashes + expect(file!.uri).toContain("/src/foo.ts") + }) + + test("handles Windows absolute path with special characters", () => { + const prompt: Prompt = [{ type: "file", path: "file#name.txt", content: "@file#name.txt", start: 0, end: 14 }] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@file#name.txt", + sessionDirectory: "C:\\Users\\test\\Documents", // Windows path + }) + + const file = result.files[0] + expect(file).toBeDefined() + // URL should be parseable + expect(() => new URL(file!.uri)).not.toThrow() + // Special chars should be encoded + expect(file!.uri).toContain("file%23name.txt") + // Should have Windows drive letter properly encoded + expect(file!.uri).toMatch(/file:\/\/\/[A-Z]:/) + }) + + test("handles Linux absolute paths correctly", () => { + const prompt: Prompt = [{ type: "file", path: "src/app.ts", content: "@src/app.ts", start: 0, end: 10 }] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@src/app.ts", + sessionDirectory: "/home/user/project", + }) + + expect(result.files[0]?.uri).toBe("file:///home/user/project/src/app.ts") + }) + + test("handles macOS paths correctly", () => { + const prompt: Prompt = [{ type: "file", path: "README.md", content: "@README.md", start: 0, end: 9 }] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@README.md", + sessionDirectory: "/Users/kelvin/Projects/opencode", + }) + + expect(result.files[0]?.uri).toBe("file:///Users/kelvin/Projects/opencode/README.md") + }) + + test("handles context files with Windows paths", () => { + const result = buildPromptRequest({ + prompt: [], + context: [ + { key: "ctx:1", type: "file", path: "src\\utils\\helper.ts" }, + { key: "ctx:2", type: "file", path: "test\\unit.test.ts", comment: "check tests" }, + ], + images: [], + text: "test", + sessionDirectory: "D:\\workspace\\app", + }) + + expect(result.files).toHaveLength(2) + + // All file URLs should be valid + result.files.forEach((file) => { + expect(() => new URL(file.uri)).not.toThrow() + expect(file.uri).not.toContain("%5C") // No encoded backslashes + }) + }) + + test("handles absolute Windows paths (user manually specifies full path)", () => { + const prompt: Prompt = [ + { type: "file", path: "D:\\other\\project\\file.ts", content: "@D:\\other\\project\\file.ts", start: 0, end: 25 }, + ] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@D:\\other\\project\\file.ts", + sessionDirectory: "C:\\current\\project", + }) + + const file = result.files[0] + expect(file).toBeDefined() + // Should handle absolute path that differs from sessionDirectory + expect(() => new URL(file!.uri)).not.toThrow() + expect(file!.uri).toContain("/D:/other/project/file.ts") + }) + + test("handles selection with query parameters on Windows", () => { + const prompt: Prompt = [ + { + type: "file", + path: "src\\App.tsx", + content: "@src\\App.tsx", + start: 0, + end: 11, + selection: { startLine: 10, startChar: 0, endLine: 20, endChar: 5 }, + }, + ] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@src\\App.tsx", + sessionDirectory: "C:\\project", + }) + + const file = result.files[0] + expect(file).toBeDefined() + // Should have query parameters + expect(file!.uri).toContain("?start=10&end=20") + // Should be valid URL + expect(() => new URL(file!.uri)).not.toThrow() + // Query params should parse correctly + const url = new URL(file!.uri) + expect(url.searchParams.get("start")).toBe("10") + expect(url.searchParams.get("end")).toBe("20") + }) + + test("handles file paths with dots and special segments on Windows", () => { + const prompt: Prompt = [ + { type: "file", path: "..\\..\\shared\\util.ts", content: "@..\\..\\shared\\util.ts", start: 0, end: 21 }, + ] + + const result = buildPromptRequest({ + prompt, + context: [], + images: [], + text: "@..\\..\\shared\\util.ts", + sessionDirectory: "C:\\projects\\myapp\\src", + }) + + const file = result.files[0] + expect(file).toBeDefined() + // Should be valid URL + expect(() => new URL(file!.uri)).not.toThrow() + // Should preserve .. segments (backend normalizes) + expect(file!.uri).toContain("/..") + }) +}) diff --git a/packages/app/src/components/prompt-input/build-prompt-request.ts b/packages/app/src/components/prompt-input/build-prompt-request.ts new file mode 100644 index 00000000000..255b38a68f6 --- /dev/null +++ b/packages/app/src/components/prompt-input/build-prompt-request.ts @@ -0,0 +1,115 @@ +import { getFilename } from "@opencode-ai/core/util/path" +import type { FileSelection } from "@/context/file" +import { encodeFilePath } from "@/context/file/path" +import type { AgentPart, FileAttachmentPart, ImageAttachmentPart, Prompt } from "@/context/prompt" +import { formatCommentNote, type PromptComment } from "@/utils/comment-note" + +// Network fields feed both boundaries; display fields keep desktop-only rendering details in the local echo. +type PromptRequest = { + text: string + displayText: string + files: { uri: string; mime: string; name?: string; mention?: { start: number; end: number; text: string } }[] + agents: { name: string; mention?: { start: number; end: number; text: string } }[] + comments: PromptComment[] +} + +type ContextFile = { + key: string + type: "file" + path: string + selection?: FileSelection + comment?: string + commentID?: string + commentOrigin?: "review" | "file" + preview?: string +} + +type BuildPromptRequestInput = { + prompt: Prompt + context: ContextFile[] + images: (Omit & { dataUrl: string })[] + text: string + sessionDirectory: string +} + +const absolute = (directory: string, path: string) => { + if (path.startsWith("/")) return path + if (/^[A-Za-z]:[\\/]/.test(path) || /^[A-Za-z]:$/.test(path)) return path + if (path.startsWith("\\\\") || path.startsWith("//")) return path + return `${directory.replace(/[\\/]+$/, "")}/${path}` +} + +const fileQuery = (selection: FileSelection | undefined) => + selection ? `?start=${selection.startLine}&end=${selection.endLine}` : "" + +const mention = /(^|[\s([{"'])@(\S+)/g + +const parseCommentMentions = (comment: string) => { + return Array.from(comment.matchAll(mention)).flatMap((match) => { + const path = (match[2] ?? "").replace(/[.,!?;:)}\]"']+$/, "") + if (!path) return [] + return [path] + }) +} + +const isFileAttachment = (part: Prompt[number]): part is FileAttachmentPart => part.type === "file" +const isAgentAttachment = (part: Prompt[number]): part is AgentPart => part.type === "agent" + +export function buildPromptRequest(input: BuildPromptRequestInput): PromptRequest { + const files = input.prompt.filter(isFileAttachment).map((attachment) => { + const path = absolute(input.sessionDirectory, attachment.path) + return { + uri: attachment.url ?? `file://${encodeFilePath(path)}${fileQuery(attachment.selection)}`, + mime: attachment.mime ?? "text/plain", + name: attachment.filename ?? getFilename(attachment.path), + mention: { start: attachment.start, end: attachment.end, text: attachment.content }, + } + }) + + const agents = input.prompt.filter(isAgentAttachment).map((attachment) => ({ + name: attachment.name, + mention: { start: attachment.start, end: attachment.end, text: attachment.content }, + })) + + const used = new Set(files.map((file) => file.uri)) + const comments: PromptComment[] = [] + const context = input.context.flatMap((item) => { + const path = absolute(input.sessionDirectory, item.path) + const uri = `file://${encodeFilePath(path)}${fileQuery(item.selection)}` + const comment = item.comment?.trim() + if (!comment && used.has(uri)) return [] + used.add(uri) + + const file = { uri, mime: "text/plain", name: getFilename(item.path) } + if (!comment) return [file] + + comments.push({ + path: item.path, + selection: item.selection, + comment, + preview: item.preview, + origin: item.commentOrigin, + }) + const mentions = parseCommentMentions(comment).flatMap((path) => { + const uri = `file://${encodeFilePath(absolute(input.sessionDirectory, path))}` + if (used.has(uri)) return [] + used.add(uri) + return [{ uri, mime: "text/plain", name: getFilename(path) }] + }) + return [file, ...mentions] + }) + + const images = input.images.map((attachment) => ({ + uri: attachment.dataUrl, + mime: attachment.mime, + name: attachment.sourcePath ?? attachment.filename, + })) + + return { + text: [...(input.text.trim() ? [input.text] : []), ...comments.map(formatCommentNote)].join("\n"), + displayText: input.text, + files: [...files, ...context, ...images], + agents, + comments, + } +} diff --git a/packages/app/src/components/prompt-input/build-request-parts.test.ts b/packages/app/src/components/prompt-input/build-request-parts.test.ts deleted file mode 100644 index ab84cb6eae8..00000000000 --- a/packages/app/src/components/prompt-input/build-request-parts.test.ts +++ /dev/null @@ -1,396 +0,0 @@ -import { describe, expect, test } from "bun:test" -import type { Prompt } from "@/context/prompt" -import { buildRequestParts } from "./build-request-parts" - -describe("buildRequestParts", () => { - test("builds typed request and optimistic parts without cast path", () => { - const prompt: Prompt = [ - { type: "text", content: "hello", start: 0, end: 5 }, - { - type: "file", - path: "src/foo.ts", - content: "@src/foo.ts", - start: 5, - end: 16, - selection: { startLine: 4, startChar: 1, endLine: 6, endChar: 1 }, - }, - { type: "agent", name: "planner", content: "@planner", start: 16, end: 24 }, - ] - - const result = buildRequestParts({ - prompt, - context: [{ key: "ctx:1", type: "file", path: "src/bar.ts", comment: "check this" }], - images: [ - { type: "image", id: "img_1", filename: "a.png", mime: "image/png", dataUrl: "data:image/png;base64,AAA" }, - ], - text: "hello @src/foo.ts @planner", - messageID: "msg_1", - sessionID: "ses_1", - sessionDirectory: "/repo", - }) - - expect(result.requestParts[0]?.type).toBe("text") - expect(result.requestParts.some((part) => part.type === "agent")).toBe(true) - expect( - result.requestParts.some((part) => part.type === "file" && part.url.startsWith("file:///repo/src/foo.ts")), - ).toBe(true) - expect(result.requestParts.some((part) => part.type === "text" && part.synthetic)).toBe(true) - expect( - result.requestParts.some( - (part) => - part.type === "text" && - part.synthetic && - part.metadata?.opencodeComment && - (part.metadata.opencodeComment as { comment?: string }).comment === "check this", - ), - ).toBe(true) - - expect(result.optimisticParts).toHaveLength(result.requestParts.length) - expect(result.optimisticParts.every((part) => part.sessionID === "ses_1" && part.messageID === "msg_1")).toBe(true) - }) - - test("keeps multiple uploaded attachments in order", () => { - const result = buildRequestParts({ - prompt: [{ type: "text", content: "check these", start: 0, end: 11 }], - context: [], - images: [ - { type: "image", id: "img_1", filename: "a.png", mime: "image/png", dataUrl: "data:image/png;base64,AAA" }, - { - type: "image", - id: "img_2", - filename: "b.pdf", - mime: "application/pdf", - dataUrl: "data:application/pdf;base64,BBB", - }, - ], - text: "check these", - messageID: "msg_multi", - sessionID: "ses_multi", - sessionDirectory: "/repo", - }) - - const files = result.requestParts.filter((part) => part.type === "file" && part.url.startsWith("data:")) - - expect(files).toHaveLength(2) - expect(files.map((part) => (part.type === "file" ? part.filename : ""))).toEqual(["a.png", "b.pdf"]) - }) - - test("preserves an external attachment source path for the model", () => { - const result = buildRequestParts({ - prompt: [], - context: [], - images: [ - { - type: "image", - id: "img_external", - filename: "opencode.global.dat", - sourcePath: "C:\\Users\\Luke\\AppData\\Roaming\\ai.opencode.desktop.beta\\opencode.global.dat", - mime: "text/plain", - dataUrl: "data:text/plain;base64,AAA", - }, - ], - text: "inspect this", - messageID: "msg_external", - sessionID: "ses_external", - sessionDirectory: "C:\\Repos\\sst\\opencode", - }) - - expect(result.requestParts.find((part) => part.type === "file")?.filename).toBe( - "C:\\Users\\Luke\\AppData\\Roaming\\ai.opencode.desktop.beta\\opencode.global.dat", - ) - }) - - test("preserves reference aliases as directory file parts", () => { - const result = buildRequestParts({ - prompt: [ - { - type: "file", - path: "/repo/../docs", - content: "@docs", - start: 0, - end: 5, - mime: "application/x-directory", - filename: "docs", - }, - ], - context: [], - images: [], - text: "@docs", - messageID: "msg_reference", - sessionID: "ses_reference", - sessionDirectory: "/repo/app", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - expect(filePart.mime).toBe("application/x-directory") - expect(filePart.filename).toBe("docs") - expect(filePart.url).toBe("file:///repo/../docs") - expect(filePart.source?.type).toBe("file") - if (filePart.source?.type === "file") { - expect(filePart.source.path).toBe("/repo/../docs") - expect(filePart.source.text.value).toBe("@docs") - } - } - }) - - test("deduplicates context files when prompt already includes same path", () => { - const prompt: Prompt = [{ type: "file", path: "src/foo.ts", content: "@src/foo.ts", start: 0, end: 11 }] - - const result = buildRequestParts({ - prompt, - context: [ - { key: "ctx:dup", type: "file", path: "src/foo.ts" }, - { key: "ctx:comment", type: "file", path: "src/foo.ts", comment: "focus here" }, - ], - images: [], - text: "@src/foo.ts", - messageID: "msg_2", - sessionID: "ses_2", - sessionDirectory: "/repo", - }) - - const fooFiles = result.requestParts.filter( - (part) => part.type === "file" && part.url.startsWith("file:///repo/src/foo.ts"), - ) - const synthetic = result.requestParts.filter((part) => part.type === "text" && part.synthetic) - - expect(fooFiles).toHaveLength(2) - expect(synthetic).toHaveLength(1) - }) - - test("adds file parts for @mentions inside comment text", () => { - const result = buildRequestParts({ - prompt: [{ type: "text", content: "look", start: 0, end: 4 }], - context: [ - { - key: "ctx:comment-mention", - type: "file", - path: "src/review.ts", - comment: "Compare with @src/shared.ts and @src/review.ts.", - }, - ], - images: [], - text: "look", - messageID: "msg_comment_mentions", - sessionID: "ses_comment_mentions", - sessionDirectory: "/repo", - }) - - const files = result.requestParts.filter((part) => part.type === "file") - expect(files).toHaveLength(2) - expect(files.some((part) => part.type === "file" && part.url === "file:///repo/src/review.ts")).toBe(true) - expect(files.some((part) => part.type === "file" && part.url === "file:///repo/src/shared.ts")).toBe(true) - }) - - test("handles Windows paths correctly (simulated on macOS)", () => { - const prompt: Prompt = [{ type: "file", path: "src\\foo.ts", content: "@src\\foo.ts", start: 0, end: 11 }] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@src\\foo.ts", - messageID: "msg_win_1", - sessionID: "ses_win_1", - sessionDirectory: "D:\\projects\\myapp", // Windows path - }) - - // Should create valid file URLs - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // URL should be parseable - expect(() => new URL(filePart.url)).not.toThrow() - // Should not have encoded backslashes in wrong place - expect(filePart.url).not.toContain("%5C") - // Should have normalized to forward slashes - expect(filePart.url).toContain("/src/foo.ts") - } - }) - - test("handles Windows absolute path with special characters", () => { - const prompt: Prompt = [{ type: "file", path: "file#name.txt", content: "@file#name.txt", start: 0, end: 14 }] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@file#name.txt", - messageID: "msg_win_2", - sessionID: "ses_win_2", - sessionDirectory: "C:\\Users\\test\\Documents", // Windows path - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // URL should be parseable - expect(() => new URL(filePart.url)).not.toThrow() - // Special chars should be encoded - expect(filePart.url).toContain("file%23name.txt") - // Should have Windows drive letter properly encoded - expect(filePart.url).toMatch(/file:\/\/\/[A-Z]:/) - } - }) - - test("handles Linux absolute paths correctly", () => { - const prompt: Prompt = [{ type: "file", path: "src/app.ts", content: "@src/app.ts", start: 0, end: 10 }] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@src/app.ts", - messageID: "msg_linux_1", - sessionID: "ses_linux_1", - sessionDirectory: "/home/user/project", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // URL should be parseable - expect(() => new URL(filePart.url)).not.toThrow() - // Should be a normal Unix path - expect(filePart.url).toBe("file:///home/user/project/src/app.ts") - } - }) - - test("handles macOS paths correctly", () => { - const prompt: Prompt = [{ type: "file", path: "README.md", content: "@README.md", start: 0, end: 9 }] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@README.md", - messageID: "msg_mac_1", - sessionID: "ses_mac_1", - sessionDirectory: "/Users/kelvin/Projects/opencode", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // URL should be parseable - expect(() => new URL(filePart.url)).not.toThrow() - // Should be a normal Unix path - expect(filePart.url).toBe("file:///Users/kelvin/Projects/opencode/README.md") - } - }) - - test("handles context files with Windows paths", () => { - const prompt: Prompt = [] - - const result = buildRequestParts({ - prompt, - context: [ - { key: "ctx:1", type: "file", path: "src\\utils\\helper.ts" }, - { key: "ctx:2", type: "file", path: "test\\unit.test.ts", comment: "check tests" }, - ], - images: [], - text: "test", - messageID: "msg_win_ctx", - sessionID: "ses_win_ctx", - sessionDirectory: "D:\\workspace\\app", - }) - - const fileParts = result.requestParts.filter((part) => part.type === "file") - expect(fileParts).toHaveLength(2) - - // All file URLs should be valid - fileParts.forEach((part) => { - if (part.type === "file") { - expect(() => new URL(part.url)).not.toThrow() - expect(part.url).not.toContain("%5C") // No encoded backslashes - } - }) - }) - - test("handles absolute Windows paths (user manually specifies full path)", () => { - const prompt: Prompt = [ - { type: "file", path: "D:\\other\\project\\file.ts", content: "@D:\\other\\project\\file.ts", start: 0, end: 25 }, - ] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@D:\\other\\project\\file.ts", - messageID: "msg_abs", - sessionID: "ses_abs", - sessionDirectory: "C:\\current\\project", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // Should handle absolute path that differs from sessionDirectory - expect(() => new URL(filePart.url)).not.toThrow() - expect(filePart.url).toContain("/D:/other/project/file.ts") - } - }) - - test("handles selection with query parameters on Windows", () => { - const prompt: Prompt = [ - { - type: "file", - path: "src\\App.tsx", - content: "@src\\App.tsx", - start: 0, - end: 11, - selection: { startLine: 10, startChar: 0, endLine: 20, endChar: 5 }, - }, - ] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@src\\App.tsx", - messageID: "msg_sel", - sessionID: "ses_sel", - sessionDirectory: "C:\\project", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // Should have query parameters - expect(filePart.url).toContain("?start=10&end=20") - // Should be valid URL - expect(() => new URL(filePart.url)).not.toThrow() - // Query params should parse correctly - const url = new URL(filePart.url) - expect(url.searchParams.get("start")).toBe("10") - expect(url.searchParams.get("end")).toBe("20") - } - }) - - test("handles file paths with dots and special segments on Windows", () => { - const prompt: Prompt = [ - { type: "file", path: "..\\..\\shared\\util.ts", content: "@..\\..\\shared\\util.ts", start: 0, end: 21 }, - ] - - const result = buildRequestParts({ - prompt, - context: [], - images: [], - text: "@..\\..\\shared\\util.ts", - messageID: "msg_dots", - sessionID: "ses_dots", - sessionDirectory: "C:\\projects\\myapp\\src", - }) - - const filePart = result.requestParts.find((part) => part.type === "file") - expect(filePart).toBeDefined() - if (filePart?.type === "file") { - // Should be valid URL - expect(() => new URL(filePart.url)).not.toThrow() - // Should preserve .. segments (backend normalizes) - expect(filePart.url).toContain("/..") - } - }) -}) diff --git a/packages/app/src/components/prompt-input/build-request-parts.ts b/packages/app/src/components/prompt-input/build-request-parts.ts deleted file mode 100644 index f23499a643a..00000000000 --- a/packages/app/src/components/prompt-input/build-request-parts.ts +++ /dev/null @@ -1,216 +0,0 @@ -import { getFilename } from "@opencode-ai/core/util/path" -import type { AgentPart as MessageAgentPart, FilePart, Part, TextPart } from "@/types" -import type { FileSelection } from "@/context/file" -import { encodeFilePath } from "@/context/file/path" -import type { AgentPart, FileAttachmentPart, ImageAttachmentPart, Prompt } from "@/context/prompt" -import { Identifier } from "@/utils/id" -import { createCommentMetadata, formatCommentNote } from "@/utils/comment-note" - -type PromptRequestPart = - | (Omit & { id: string }) - | (Omit & { id: string }) - | (Omit & { id: string }) - -type ContextFile = { - key: string - type: "file" - path: string - selection?: FileSelection - comment?: string - commentID?: string - commentOrigin?: "review" | "file" - preview?: string -} - -type BuildRequestPartsInput = { - prompt: Prompt - context: ContextFile[] - images: (Omit & { dataUrl: string })[] - text: string - messageID: string - sessionID: string - sessionDirectory: string -} - -const absolute = (directory: string, path: string) => { - if (path.startsWith("/")) return path - if (/^[A-Za-z]:[\\/]/.test(path) || /^[A-Za-z]:$/.test(path)) return path - if (path.startsWith("\\\\") || path.startsWith("//")) return path - return `${directory.replace(/[\\/]+$/, "")}/${path}` -} - -const fileQuery = (selection: FileSelection | undefined) => - selection ? `?start=${selection.startLine}&end=${selection.endLine}` : "" - -const mention = /(^|[\s([{"'])@(\S+)/g - -const parseCommentMentions = (comment: string) => { - return Array.from(comment.matchAll(mention)).flatMap((match) => { - const path = (match[2] ?? "").replace(/[.,!?;:)}\]"']+$/, "") - if (!path) return [] - return [path] - }) -} - -const isFileAttachment = (part: Prompt[number]): part is FileAttachmentPart => part.type === "file" -const isAgentAttachment = (part: Prompt[number]): part is AgentPart => part.type === "agent" - -const toOptimisticPart = (part: PromptRequestPart, sessionID: string, messageID: string): Part => { - if (part.type === "text") { - return { - id: part.id, - type: "text", - text: part.text, - synthetic: part.synthetic, - ignored: part.ignored, - time: part.time, - metadata: part.metadata, - sessionID, - messageID, - } - } - if (part.type === "file") { - return { - id: part.id, - type: "file", - mime: part.mime, - filename: part.filename, - url: part.url, - source: part.source, - sessionID, - messageID, - } - } - return { - id: part.id, - type: "agent", - name: part.name, - source: part.source, - sessionID, - messageID, - } -} - -export function buildRequestParts(input: BuildRequestPartsInput) { - const requestParts: PromptRequestPart[] = input.text.trim() - ? [ - { - id: Identifier.ascending("part"), - type: "text", - text: input.text, - }, - ] - : [] - - const files = input.prompt.filter(isFileAttachment).map((attachment) => { - const path = absolute(input.sessionDirectory, attachment.path) - const source = attachment.source - ? { - ...attachment.source, - text: { - value: attachment.content, - start: attachment.start, - end: attachment.end, - }, - } - : { - type: "file" as const, - text: { - value: attachment.content, - start: attachment.start, - end: attachment.end, - }, - path, - } - return { - id: Identifier.ascending("part"), - type: "file", - mime: attachment.mime ?? "text/plain", - url: attachment.url ?? `file://${encodeFilePath(path)}${fileQuery(attachment.selection)}`, - filename: attachment.filename ?? getFilename(attachment.path), - source, - } satisfies PromptRequestPart - }) - - const agents = input.prompt.filter(isAgentAttachment).map((attachment) => { - return { - id: Identifier.ascending("part"), - type: "agent", - name: attachment.name, - source: { - value: attachment.content, - start: attachment.start, - end: attachment.end, - }, - } satisfies PromptRequestPart - }) - - const used = new Set(files.map((part) => part.url)) - const context = input.context.flatMap((item) => { - const path = absolute(input.sessionDirectory, item.path) - const url = `file://${encodeFilePath(path)}${fileQuery(item.selection)}` - const comment = item.comment?.trim() - if (!comment && used.has(url)) return [] - used.add(url) - - const filePart = { - id: Identifier.ascending("part"), - type: "file", - mime: "text/plain", - url, - filename: getFilename(item.path), - } satisfies PromptRequestPart - - if (!comment) return [filePart] - - const mentions = parseCommentMentions(comment).flatMap((path) => { - const url = `file://${encodeFilePath(absolute(input.sessionDirectory, path))}` - if (used.has(url)) return [] - used.add(url) - return [ - { - id: Identifier.ascending("part"), - type: "file", - mime: "text/plain", - url, - filename: getFilename(path), - } satisfies PromptRequestPart, - ] - }) - - return [ - { - id: Identifier.ascending("part"), - type: "text", - text: formatCommentNote({ path: item.path, selection: item.selection, comment }), - synthetic: true, - metadata: createCommentMetadata({ - path: item.path, - selection: item.selection, - comment, - preview: item.preview, - origin: item.commentOrigin, - }), - } satisfies PromptRequestPart, - filePart, - ...mentions, - ] - }) - - const images = input.images.map((attachment) => { - return { - id: Identifier.ascending("part"), - type: "file", - mime: attachment.mime, - url: attachment.dataUrl, - filename: attachment.sourcePath ?? attachment.filename, - } satisfies PromptRequestPart - }) - - requestParts.push(...files, ...context, ...agents, ...images) - - return { - requestParts, - optimisticParts: requestParts.map((part) => toOptimisticPart(part, input.sessionID, input.messageID)), - } -} diff --git a/packages/app/src/components/prompt-input/submit.test.ts b/packages/app/src/components/prompt-input/submit.test.ts index ab905c28c67..c54bee7cfae 100644 --- a/packages/app/src/components/prompt-input/submit.test.ts +++ b/packages/app/src/components/prompt-input/submit.test.ts @@ -11,15 +11,17 @@ type SessionCreateInput = { model?: { id: string; providerID: string; variant?: string } location?: { directory: string } } -const optimistic: Array<{ +const admitted: Array<{ directory?: string - sessionID?: string - message: { - agent: string - model: { providerID: string; modelID: string } - variant?: string - } + sessionID: string + messageID: string + text: string + displayText: string + agent: string + model: { providerID: string; modelID: string; variant?: string } + comments: unknown[] }> = [] +const confirmed: unknown[] = [] const storedSessions: Record> = {} const sentShell: Array<{ sessionID: string; id?: string; command: string }> = [] const sentShellDirectories: string[] = [] @@ -35,9 +37,11 @@ const switchedModels: Array<{ const sessionRequestOrder: string[] = [] const updatedDrafts: Array<{ draftID: string; worktree?: string }> = [] const syncedServers: string[] = [] -const optimisticServers: string[] = [] +const admittedServers: string[] = [] const promptCaptures: Array<{ scope?: unknown; target?: unknown }> = [] let serverSessionSyncs = 0 +let restoredPrompts = 0 +let clearEchoCalls = 0 let params: { id?: string } = {} let search: { draftId?: string } = {} @@ -47,6 +51,8 @@ let createSessionGate: Promise | undefined let createWorktreeGate: Promise | undefined let worktreeFailure: Error | undefined let locationFailure: Error | undefined +let promptFailure: Error | undefined +let clearEchoResult = true let worktreeCreates = 0 let activeSDK = "server-a" let activeServerSync = "server-a" @@ -74,7 +80,7 @@ const prompt = { set: () => undefined, }, reset: () => undefined, - set: () => undefined, + set: () => restoredPrompts++, context: { add: () => undefined, remove: () => undefined, @@ -116,7 +122,16 @@ const clientFor = (directory: string) => { sessionRequestOrder.push("prompt") sentPrompts.push(sessionDirectories[(input as { sessionID: string }).sessionID] ?? directory) promptInputs.push(input) - return { data: undefined } + if (promptFailure) throw promptFailure + const prompt = input as { sessionID: string; id: string; text: string } + return { + id: prompt.id, + sessionID: prompt.sessionID, + timeCreated: 1, + type: "user" as const, + delivery: "steer" as const, + payload: { text: prompt.text }, + } }, switchAgent: async (input: { sessionID: string; agent: string }) => { sessionRequestOrder.push("agent") @@ -235,16 +250,27 @@ beforeAll(async () => { return { data: { command: commands, project: "project" }, session: { - optimistic: { - add: (value: { + inbox: { + echo: (value: { directory?: string - sessionID?: string - message: { agent: string; model: { providerID: string; modelID: string; variant?: string } } + sessionID: string + messageID: string + text: string + displayText: string + agent: string + model: { providerID: string; modelID: string; variant?: string } + comments: unknown[] }) => { - optimisticServers.push(server) - optimistic.push(value) + admittedServers.push(server) + admitted.push(value) + }, + confirm: (value: unknown) => { + confirmed.push(value) + }, + clearEcho: () => { + clearEchoCalls++ + return clearEchoResult }, - remove: () => undefined, }, }, set: () => undefined, @@ -304,7 +330,8 @@ beforeAll(async () => { beforeEach(() => { createdSessions.length = 0 - optimistic.length = 0 + admitted.length = 0 + confirmed.length = 0 promotedDrafts.length = 0 updatedDrafts.length = 0 sentCommands.length = 0 @@ -314,8 +341,10 @@ beforeEach(() => { switchedModels.length = 0 sessionRequestOrder.length = 0 syncedServers.length = 0 - optimisticServers.length = 0 + admittedServers.length = 0 promptCaptures.length = 0 + restoredPrompts = 0 + clearEchoCalls = 0 params = {} search = {} sentShell.length = 0 @@ -333,6 +362,8 @@ beforeEach(() => { createWorktreeGate = undefined worktreeFailure = undefined locationFailure = undefined + promptFailure = undefined + clearEchoResult = true worktreeCreates = 0 for (const key of Object.keys(draftServers)) delete draftServers[key] for (const key of Object.keys(sessionDirectories)) delete sessionDirectories[key] @@ -421,7 +452,7 @@ describe("prompt submit worktree selection", () => { expect(updatedDrafts).toEqual([{ draftID: "draft-1", worktree: undefined }]) expect(promotedDrafts).toEqual([{ draftID: "draft-1", server: "project-server-a", sessionId: "session-1" }]) expect(syncedServers.every((server) => server === "server-a")).toBe(true) - expect(optimisticServers).toEqual(["server-a"]) + expect(admittedServers).toEqual(["server-a"]) expect(promptCaptures.at(-1)?.target).toEqual({ server: "project-server-a", scope: ServerScope.local }) expect(submitted).toBe(0) }) @@ -441,13 +472,15 @@ describe("prompt submit worktree selection", () => { await submit.handleSubmit(event) await Bun.sleep(0) - expect(optimistic).toHaveLength(1) - expect(optimistic[0]).toMatchObject({ - message: { - agent: "agent", - model: { providerID: "provider", modelID: "model", variant: "high" }, - }, + expect(admitted).toHaveLength(1) + expect(admitted[0]).toMatchObject({ + sessionID: "session-1", + text: "ls", + agent: "agent", + model: { providerID: "provider", modelID: "model", variant: "high" }, }) + expect(admitted[0]?.messageID).toStartWith("msg_") + expect(confirmed).toMatchObject([{ id: admitted[0]?.messageID, sessionID: "session-1" }]) expect(sentPrompts).toEqual(["/repo/main"]) expect(switchedAgents).toEqual([{ sessionID: "session-1", agent: "agent" }]) expect(switchedModels).toEqual([ @@ -466,6 +499,22 @@ describe("prompt submit worktree selection", () => { expect((promptInputs[0] as { id?: string }).id).toStartWith("msg_") }) + test("keeps a confirmed echo when the prompt response is lost", async () => { + params = { id: "session-1" } + promptFailure = new Error("connection lost") + clearEchoResult = false + const submit = makeSubmit({ + info: () => ({ id: "session-1", agent: "agent", model: { id: "model", providerID: "provider" } }), + }) + + await submit.handleSubmit(event) + await settle() + + expect(admitted).toHaveLength(1) + expect(clearEchoCalls).toBe(1) + expect(restoredPrompts).toBe(0) + }) + test("submits slash commands through the current session API", async () => { params = { id: "session-1" } variant = "high" diff --git a/packages/app/src/components/prompt-input/submit.ts b/packages/app/src/components/prompt-input/submit.ts index da717534f21..e2cdfc0fde3 100644 --- a/packages/app/src/components/prompt-input/submit.ts +++ b/packages/app/src/components/prompt-input/submit.ts @@ -1,10 +1,9 @@ -import type { Message } from "@/types" import type { SessionInfo } from "@opencode-ai/client/promise" import { showToast } from "@/utils/toast" import { base64Encode } from "@opencode-ai/core/util/encode" import { Binary } from "@opencode-ai/core/util/binary" import { useNavigate, useParams, useSearchParams } from "@solidjs/router" -import { batch, startTransition, type Accessor } from "solid-js" +import { startTransition, type Accessor } from "solid-js" import { useTabs } from "@/context/tabs" import { useServerSync, type ServerSync } from "@/context/server-sync" import { useLanguage } from "@/context/language" @@ -15,7 +14,7 @@ import { useSDK, type DirectorySDK } from "@/context/sdk" import { useSync, type DirectorySync } from "@/context/sync" import { Identifier } from "@/utils/id" import { getDirectory } from "@opencode-ai/core/util/path" -import { buildRequestParts } from "./build-request-parts" +import { buildPromptRequest } from "./build-prompt-request" import { setCursorPosition } from "./editor-dom" import { formatServerError } from "@/utils/server-errors" import { ScopedKey } from "@/utils/server-scope" @@ -100,43 +99,22 @@ export async function sendFollowupDraft(input: FollowupSendInput) { dataUrl: await blobDataUrl(attachment.blob, attachment.mime), })), ) - const { requestParts, optimisticParts } = buildRequestParts({ + const request = buildPromptRequest({ prompt: input.draft.prompt, context: input.draft.context, images: encodedImages, text, - sessionID: input.draft.sessionID, - messageID, sessionDirectory: input.draft.sessionDirectory, }) - const message: Message = { - id: messageID, + setBusy() + input.sync.session.inbox.echo({ + directory: input.draft.sessionDirectory, sessionID: input.draft.sessionID, - role: "user", - time: { created: Date.now() }, + messageID, agent: input.draft.agent, model: { ...input.draft.model, variant: input.draft.variant }, - } - - const add = () => - input.sync.session.optimistic.add({ - directory: input.draft.sessionDirectory, - sessionID: input.draft.sessionID, - message, - parts: optimisticParts, - }) - - const remove = () => - input.sync.session.optimistic.remove({ - directory: input.draft.sessionDirectory, - sessionID: input.draft.sessionID, - messageID, - }) - - batch(() => { - setBusy() - add() + ...request, }) try { @@ -159,40 +137,23 @@ export async function sendFollowupDraft(input: FollowupSendInput) { }) } - await input.api.prompt({ + const admitted = await input.api.prompt({ sessionID: input.draft.sessionID, id: messageID, - text: requestParts.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n"), - files: requestParts.flatMap((part) => { - if (part.type !== "file") return [] - const text = part.source?.text - return [ - { - uri: part.url, - name: part.filename, - mention: text ? { start: text.start, end: text.end, text: text.value } : undefined, - }, - ] - }), - agents: requestParts.flatMap((part) => - part.type === "agent" - ? [ - { - name: part.name, - mention: part.source - ? { start: part.source.start, end: part.source.end, text: part.source.value } - : undefined, - }, - ] - : [], - ), + text: request.text, + files: request.files.map((file) => ({ uri: file.uri, name: file.name, mention: file.mention })), + agents: request.agents, }) + input.sync.session.inbox.confirm(admitted) return true } catch (err) { - batch(() => { - setIdle() - remove() + const failed = input.sync.session.inbox.clearEcho({ + directory: input.draft.sessionDirectory, + sessionID: input.draft.sessionID, + messageID, }) + if (!failed) return true + setIdle() throw err } } @@ -538,14 +499,6 @@ export function createPromptSubmit(input: PromptSubmitInput) { const commentItems = context.filter((item) => item.type === "file" && !!item.comment?.trim()) const messageID = Identifier.ascending("message") - const removeOptimisticMessage = () => { - submissionSync.session.optimistic.remove({ - directory: sessionDirectory, - sessionID: session.id, - messageID, - }) - } - for (const item of commentItems) submission.target().context.remove(item.key) clearInput() @@ -565,7 +518,6 @@ export function createPromptSubmit(input: PromptSubmitInput) { title: language.t("prompt.toast.promptSendFailed.title"), description: errorMessage(err), }) - removeOptimisticMessage() if (restoreInput()) restoreCommentItems(submission.target(), commentItems) }) } finally { diff --git a/packages/app/src/context/directory-sync.ts b/packages/app/src/context/directory-sync.ts index f1665e0b390..09e5ec7bba2 100644 --- a/packages/app/src/context/directory-sync.ts +++ b/packages/app/src/context/directory-sync.ts @@ -1,10 +1,10 @@ import { Binary } from "@opencode-ai/core/util/binary" -import type { Message, Part } from "@/types" -import type { SessionInfo } from "@opencode-ai/client/promise" +import type { SessionInboxInfo, SessionInfo } from "@opencode-ai/client/promise" import { createMemo } from "solid-js" import { produce, reconcile, type SetStoreFunction } from "solid-js/store" import type { createServerSdkContext } from "./server-sdk" import type { createServerSyncContextInner } from "./server-sync" +import type { PromptEcho } from "./server-session" import type { State } from "./global-sync/types" const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0) @@ -82,34 +82,16 @@ export const createDirSyncContext = ( const session = serverSync.session.get(sessionID) if (session?.location.directory === directory) return session }, - optimistic: { - add(input: { directory?: string; sessionID: string; message: Message; parts: Part[] }) { - serverSync.session.optimistic.add(input) + inbox: { + echo(input: PromptEcho & { directory?: string }) { + serverSync.session.inbox.echo(input) }, - remove(input: { directory?: string; sessionID: string; messageID: string }) { - serverSync.session.optimistic.remove(input) + confirm(input: SessionInboxInfo) { + return serverSync.session.inbox.confirm(input) + }, + clearEcho(input: { directory?: string; sessionID: string; messageID: string }) { + return serverSync.session.inbox.clearEcho(input) }, - }, - addOptimisticMessage(input: { - sessionID: string - messageID: string - parts: Part[] - agent: string - model: { providerID: string; modelID: string } - variant?: string - }) { - serverSync.session.optimistic.add({ - sessionID: input.sessionID, - message: { - id: input.messageID, - sessionID: input.sessionID, - role: "user", - time: { created: Date.now() }, - agent: input.agent, - model: { ...input.model, variant: input.variant }, - }, - parts: input.parts, - }) }, async sync(sessionID: string, options?: { force?: boolean }) { await serverSync.session.sync(sessionID, options) diff --git a/packages/app/src/context/file/path.test.ts b/packages/app/src/context/file/path.test.ts index 99dd88ae0e9..e89bec7e7f2 100644 --- a/packages/app/src/context/file/path.test.ts +++ b/packages/app/src/context/file/path.test.ts @@ -143,7 +143,7 @@ describe("encodeFilePath", () => { }) test("should handle mixed separator path (Windows + Unix)", () => { - // This is what happens in build-request-parts.ts when concatenating paths + // This is what happens in build-prompt-request.ts when concatenating paths const mixedPath = "D:\\dev\\projects\\opencode/README.bs.md" const result = encodeFilePath(mixedPath) const fileUrl = `file://${result}` diff --git a/packages/app/src/context/server-session-v2-reducer.test.ts b/packages/app/src/context/server-session-v2-reducer.test.ts index 0db44cb95d6..c18e2b668e3 100644 --- a/packages/app/src/context/server-session-v2-reducer.test.ts +++ b/packages/app/src/context/server-session-v2-reducer.test.ts @@ -6,6 +6,32 @@ const event = (input: object) => input as OpenCodeEvent const base = { created: 1, location: { directory: "/repo" }, durable: { aggregateID: "ses_1", seq: 1, version: 1 } } describe("v2 session reducer", () => { + test("moves a repeated inbox payload to the current event position", () => { + const reducer = createV2SessionReducer() + const result = reducer.reduce( + [ + { id: "msg_user", type: "user", text: "local", time: { created: 0 } }, + { id: "msg_agent", type: "agent-switched", agent: "review", time: { created: 1 } }, + ], + event({ + ...base, + id: "evt_admitted", + type: "session.inbox.enqueued", + data: { + sessionID: "ses_1", + inboxID: "msg_user", + item: { type: "user", delivery: "steer", payload: { text: "durable" } }, + }, + }), + ) + + expect(result?.messages).toEqual([ + { id: "msg_agent", type: "agent-switched", agent: "review", time: { created: 1 } }, + { id: "msg_user", type: "user", text: "durable", time: { created: 1 } }, + ]) + expect(result?.touched).toEqual(["msg_user"]) + }) + test("projects promoted input and streaming assistant content", () => { const reducer = createV2SessionReducer() let messages: SessionMessageInfo[] = [] diff --git a/packages/app/src/context/server-session-v2-reducer.ts b/packages/app/src/context/server-session-v2-reducer.ts index b03d1a1e0ee..c94443c8195 100644 --- a/packages/app/src/context/server-session-v2-reducer.ts +++ b/packages/app/src/context/server-session-v2-reducer.ts @@ -1,4 +1,10 @@ -import type { OpenCodeEvent, SessionInboxItem, SessionInfo, SessionMessageInfo } from "@opencode-ai/client/promise" +import type { + OpenCodeEvent, + SessionInboxInfo, + SessionInboxItem, + SessionInfo, + SessionMessageInfo, +} from "@opencode-ai/client/promise" type Assistant = Extract type Compaction = Extract @@ -29,12 +35,14 @@ export function createV2SessionReducer() { }) const append = (message: SessionMessageInfo) => result(source.some((item) => item.id === message.id) ? [...source] : [...source, message], [message.id]) + const replace = (message: SessionMessageInfo) => + result([...source.filter((item) => item.id !== message.id), message], [message.id]) switch (event.type) { case "session.inbox.enqueued": pending.set(key(sessionID, event.data.inboxID), event.data.item) if (event.data.item.type === "user") - return append({ + return replace({ id: event.data.inboxID, type: "user", metadata: event.data.item.payload.metadata, @@ -44,7 +52,7 @@ export function createV2SessionReducer() { time: { created: event.created }, }) if (event.data.item.type !== "synthetic") return result([...source]) - return append({ + return replace({ id: event.data.inboxID, type: "synthetic", metadata: event.data.item.payload.metadata, @@ -480,6 +488,9 @@ export function createV2SessionReducer() { return { reduce, + confirm(item: SessionInboxInfo) { + pending.set(key(item.sessionID, item.id), item) + }, clear(sessionID: string) { for (const id of pending.keys()) { if (id.startsWith(`${sessionID}:`)) pending.delete(id) diff --git a/packages/app/src/context/server-session.test.ts b/packages/app/src/context/server-session.test.ts index 22874d1b5d5..94bc0a9f3d8 100644 --- a/packages/app/src/context/server-session.test.ts +++ b/packages/app/src/context/server-session.test.ts @@ -185,6 +185,16 @@ const textPart = (messageID: string, input: Partial = {}): TextPart => id: `${messageID}:text:${input.id === "pending" ? 1 : 0}`, }) +const promptEcho = (messageID: string, text = "hello") => ({ + sessionID: "child", + messageID, + text, + displayText: text, + agent: "build", + model: { providerID: "provider", modelID: "model" }, + comments: [], +}) + const response = (data: MessageResponse["data"] = [], cursor?: string): MessageResponse => ({ data, response: { headers: new Headers(cursor ? { "x-next-cursor": cursor } : undefined) }, @@ -299,6 +309,26 @@ function setup(sessions: Record) { } describe("server session", () => { + test("hydrates session info after a native session.created event", async () => { + const ctx = setup({ created: session("created") }) + + ctx.store.apply({ + type: "session.created", + properties: { + sessionID: "created", + projectID: "project", + location: { directory: "/repo" }, + slug: "created", + version: "test", + }, + }) + + expect(ctx.store.get("created")).toBeUndefined() + await ctx.store.resolve("created") + expect(ctx.store.get("created")?.location.directory).toBe("/repo") + expect(ctx.get).toEqual([{ sessionID: "created" }]) + }) + test("projects V2 session events into current and legacy message state", () => { const ctx = setup({ child: session("child") }) ctx.store.remember(session("child")) @@ -710,19 +740,17 @@ describe("server session", () => { expect(store.data.part[parent.id]).toBeUndefined() }) - test("does not let an optimistic user suppress initial root backfill", async () => { + test("does not let an admitted user suppress initial root backfill", async () => { const user = userMessage("message-1") - const part = textPart(user.id) const assistants = [assistantMessage("message-2", user.id), assistantMessage("message-3", user.id)] const client = rootMessageClient( [response(assistants.map((info) => ({ info, parts: [] })))], [singleResponse(user)], ) const store = createServerSession(client) - store.optimistic.add({ sessionID: "child", message: user, parts: [part] }) + store.inbox.echo(promptEcho(user.id, "text")) await store.sync("child") - store.optimistic.remove({ sessionID: "child", messageID: user.id }) expect(client.requests).toHaveLength(1) expect(client.rootRequests).toHaveLength(1) @@ -783,28 +811,6 @@ describe("server session", () => { expect(store.data.part[stale.id]).toEqual([freshPart]) }) - test("refreshes a confirmed optimistic parent while preserving pending parts", async () => { - const stale = userMessage("message-1", { summary: { title: "stale", diffs: [] } }) - const fresh = { ...stale, summary: { title: "fresh", diffs: [] } } - const confirmed = textPart(stale.id, { id: "confirmed", text: "stale" }) - const refreshed = { ...confirmed, text: "fresh" } - const pending = textPart(stale.id, { id: "pending", text: "pending" }) - const assistant = assistantMessage("message-2", stale.id) - const client = rootMessageClient( - [response([{ info: stale, parts: [confirmed] }]), response([{ info: assistant, parts: [] }])], - [singleResponse(fresh, [refreshed])], - ) - const store = createServerSession(client) - store.optimistic.add({ sessionID: "child", message: stale, parts: [confirmed, pending] }) - await store.sync("child") - - await store.sync("child", { force: true }) - - expect(client.rootRequests).toEqual([{ sessionID: "child", messageID: stale.id }]) - expect(store.data.message.child).toEqual([fresh, assistant]) - expect(store.data.part[stale.id]).toEqual([refreshed, pending]) - }) - test("uses a parent received by SSE during the replacement load", async () => { const pending = deferredResponse() const user = userMessage("message-1") @@ -1040,30 +1046,6 @@ describe("server session", () => { expect(store.data.part[message.id]).toBeUndefined() }) - test("preserves optimistic parts re-added after removal during a refresh", async () => { - const pending = deferredResponse() - const message = userMessage("message") - const stale = textPart(message.id, { id: "stale", text: "stale" }) - const part = textPart(message.id, { id: "optimistic", text: "optimistic" }) - const store = createServerSession( - messageClient(response([{ info: message, parts: [] }]), pending.promise, response()), - ) - await store.sync("child") - const refreshing = store.sync("child", { force: true }) - - store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } }) - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - pending.resolve(response([{ info: message, parts: [stale] }])) - await refreshing - - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([part]) - - await store.sync("child", { force: true }) - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([part]) - }) - test("drops stale event content omitted by a complete initial page", async () => { const stale = userMessage("stale") const store = createServerSession(messageClient(response())) @@ -1085,170 +1067,309 @@ describe("server session", () => { expect(store.data.message.child).toEqual([live, fetched]) }) - test("does not restore removed optimistic content on refresh", async () => { - const message = userMessage("message") - const part = textPart(message.id, { text: "removed" }) - const kept = { ...message, id: "kept" } - const keptPart = { ...part, id: "kept-part", messageID: kept.id } - const store = createServerSession(messageClient(response([{ info: kept, parts: [] }]))) - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - store.optimistic.add({ sessionID: "child", message: kept, parts: [keptPart] }) - - store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } }) - store.apply({ - type: "message.part.removed", - properties: { sessionID: "child", messageID: kept.id, partID: keptPart.id }, - }) - await store.sync("child", { force: true }) - - expect(store.data.message.child).toEqual([kept]) - expect(store.data.part[message.id]).toBeUndefined() - expect(store.data.part[kept.id]).toBeUndefined() - }) - - test("replaces confirmed optimistic content with the initial page", async () => { - const optimistic = userMessage("message") - const fetched = { ...optimistic, time: { created: 2 } } - const store = createServerSession(messageClient(response([{ info: fetched, parts: [] }]))) - store.optimistic.add({ sessionID: "child", message: optimistic, parts: [] }) - - await store.sync("child") - - expect(store.data.message.child).toEqual([fetched]) - }) - - test("replaces a confirmed optimistic part with fetched content", async () => { - const pending = deferredResponse() - const message = userMessage("message") - const optimistic = textPart(message.id, { text: "optimistic" }) - const fetched = { ...optimistic, text: "fetched" } - const store = createServerSession(messageClient(pending.promise)) - const loading = store.sync("child") - - store.optimistic.add({ sessionID: "child", message, parts: [optimistic] }) - pending.resolve(response([{ info: message, parts: [fetched] }])) - await loading - - expect(store.data.part[message.id]).toEqual([fetched]) - }) - - test("rolls back only unconfirmed optimistic parts", async () => { - const pending = deferredResponse() - const message = userMessage("message") - const confirmed = textPart(message.id, { id: "confirmed", text: "confirmed" }) - const pendingPart = textPart(message.id, { id: "pending", text: "pending" }) - const store = createServerSession(messageClient(pending.promise)) - const loading = store.sync("child") - store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] }) - - pending.resolve(response([{ info: message, parts: [confirmed] }])) - await loading - store.optimistic.remove({ sessionID: "child", messageID: message.id }) - - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([confirmed]) - }) - - test("updates confirmed optimistic parts from later pages", async () => { - const message = userMessage("message") - const confirmed = textPart(message.id, { id: "confirmed", text: "first" }) - const updated = { ...confirmed, text: "updated" } - const pendingPart = textPart(message.id, { id: "pending", text: "pending" }) - const store = createServerSession( - messageClient(response([{ info: message, parts: [confirmed] }]), response([{ info: message, parts: [updated] }])), - ) - store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] }) - await store.sync("child") - - await store.sync("child", { force: true }) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) - - expect(store.data.part[message.id]).toEqual([updated]) - }) - - test("does not restore a confirmed optimistic part after its removal event", async () => { - const message = userMessage("message") - const confirmed = textPart(message.id, { id: "confirmed", text: "confirmed" }) - const pendingPart = textPart(message.id, { id: "pending", text: "pending" }) - const store = createServerSession( - messageClient(response([{ info: message, parts: [confirmed] }]), response([{ info: message, parts: [] }])), - ) - store.optimistic.add({ sessionID: "child", message, parts: [confirmed, pendingPart] }) - await store.sync("child") - store.apply({ - type: "message.part.removed", - properties: { sessionID: "child", messageID: message.id, partID: confirmed.id }, - }) - - await store.sync("child", { force: true }) - - expect(store.data.part[message.id]).toEqual([pendingPart]) - }) - - test("clears delta buffers when removing optimistic content", () => { - const message = userMessage("message") - const part = textPart(message.id, { text: "optimistic" }) + test("echoes a prompt without changing durable message order", () => { const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - store.apply({ - type: "message.part.delta", - properties: { sessionID: "child", messageID: message.id, partID: part.id, field: "text", delta: " delta" }, + + store.inbox.echo({ + ...promptEcho("msg_prompt"), + text: "hello\nThe user made the following comment regarding line 4 of src/foo.ts: check this", + files: [{ uri: "file:///repo/src/foo.ts", mime: "text/plain", name: "foo.ts" }], + agents: [{ name: "explore" }], + comments: [ + { + path: "src/foo.ts", + selection: { startLine: 4, startChar: 1, endLine: 4, endChar: 5 }, + comment: "check this", + preview: "const value = 1", + origin: "review", + }, + ], }) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) + expect(store.data.pending.child).toMatchObject([{ id: "msg_prompt", type: "user", delivery: "steer" }]) + expect(store.data.input.child).toEqual(["msg_prompt"]) + expect(store.data.session_message.child).toBeUndefined() + expect(store.data.message.child?.map((message) => message.id)).toEqual(["msg_prompt"]) + expect(store.data.part.msg_prompt).toMatchObject([ + { id: "msg_prompt:agent:0", type: "agent", name: "explore" }, + { + id: "msg_prompt:comment:0", + type: "text", + synthetic: true, + metadata: { + opencodeComment: { + path: "src/foo.ts", + selection: { startLine: 4, startChar: 1, endLine: 4, endChar: 5 }, + comment: "check this", + preview: "const value = 1", + origin: "review", + }, + }, + }, + { id: "msg_prompt:file:0", type: "file", filename: "foo.ts" }, + { id: "msg_prompt:text:0", type: "text", text: "hello" }, + ]) - expect(store.data.part[message.id]).toBeUndefined() - expect(store.data.part_text_accum_delta[part.id]).toBeUndefined() + store.applyV2({ + id: "evt_prompt", + created: 2, + type: "session.inbox.enqueued", + durable: { aggregateID: "child", seq: 1, version: 1 }, + data: { + sessionID: "child", + inboxID: "msg_prompt", + item: { + type: "user", + delivery: "steer", + payload: { + text: "hello\nThe user made the following comment regarding line 4 of src/foo.ts: check this", + }, + }, + }, + } as OpenCodeEvent) + + expect(store.data.part.msg_prompt).toMatchObject([ + { id: "msg_prompt:comment:0", type: "text", synthetic: true }, + { id: "msg_prompt:text:0", type: "text", text: "hello" }, + ]) }) - test("removes projected messages when rolling back optimistic content", () => { - const message = userMessage("message") - const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [] }) + test("preserves a local echo while message history omits pending input", async () => { + const store = createServerSession(messageClient(response())) + store.inbox.echo(promptEcho("msg_prompt")) + store.inbox.confirm({ + id: "msg_prompt", + sessionID: "child", + timeCreated: 1, + type: "user", + delivery: "steer", + payload: { text: "hello" }, + }) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) + await store.sync("child") + expect(store.data.message.child?.map((message) => message.id)).toEqual(["msg_prompt"]) + expect(store.data.part.msg_prompt).toMatchObject([{ type: "text", text: "hello" }]) + }) + + test("preserves local comment presentation through message refresh", async () => { + const note = "The user made the following comment regarding line 4 of src/foo.ts: check this" + const message = userMessage("msg_prompt") + const store = createServerSession( + messageClient(response([{ info: message, parts: [textPart(message.id, { text: note })] }])), + ) + store.inbox.echo({ + ...promptEcho(message.id), + text: `hello\n${note}`, + comments: [ + { + path: "src/foo.ts", + selection: { startLine: 4, startChar: 1, endLine: 4, endChar: 5 }, + comment: "check this", + origin: "review", + }, + ], + }) + + await store.sync("child") + + expect(store.data.part.msg_prompt).toMatchObject([ + { id: "msg_prompt:comment:0", type: "text", synthetic: true }, + { id: "msg_prompt:text:0", type: "text", text: "hello" }, + ]) + }) + + test("retires an admitted echo absent from authoritative reconnect state", async () => { + const store = createServerSession(messageClient(response())) + store.inbox.echo(promptEcho("msg_prompt")) + store.inbox.confirm({ + id: "msg_prompt", + sessionID: "child", + timeCreated: 1, + type: "user", + delivery: "steer", + payload: { text: "hello" }, + }) + + await Promise.all([store.sync("child"), store.hydrateTransient("child", async () => ({ pending: [], forms: [] }))]) + store.inbox.reconcile("child") + + expect(store.data.pending.child).toEqual([]) + expect(store.data.message.child).toEqual([]) + expect(store.data.part.msg_prompt).toBeUndefined() + }) + + test("retires a stale enqueued message when inbox hydration finishes after history", async () => { + const store = createServerSession(messageClient(response())) + store.applyV2({ + id: "evt_prompt", + created: 1, + type: "session.inbox.enqueued", + durable: { aggregateID: "child", seq: 1, version: 1 }, + data: { + sessionID: "child", + inboxID: "msg_prompt", + item: { type: "user", delivery: "steer", payload: { text: "hello" } }, + }, + } as OpenCodeEvent) + + await store.sync("child") + await store.hydrateTransient("child", async () => ({ pending: [], forms: [] })) + store.inbox.reconcile("child") + + expect(store.data.pending.child).toEqual([]) expect(store.data.session_message.child).toEqual([]) + expect(store.data.message.child).toEqual([]) + expect(store.data.part.msg_prompt).toBeUndefined() }) - test("does not remove content confirmed by a message event", () => { - const message = userMessage("message") - const part = textPart(message.id) + test("deduplicates the durable admission event against its local echo", () => { const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - store.apply({ type: "message.updated", properties: { sessionID: "child", info: message } }) + store.inbox.echo(promptEcho("msg_prompt")) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) + store.applyV2({ + id: "evt_prompt", + created: 2, + type: "session.inbox.enqueued", + durable: { aggregateID: "child", seq: 1, version: 1 }, + data: { + sessionID: "child", + inboxID: "msg_prompt", + item: { type: "user", delivery: "steer", payload: { text: "hello" } }, + }, + } as OpenCodeEvent) - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toBeUndefined() + expect(store.data.pending.child).toHaveLength(1) + expect(store.data.input.child).toEqual(["msg_prompt"]) + expect(store.data.session_message.child?.filter((message) => message.id === "msg_prompt")).toHaveLength(1) + expect(store.data.message.child?.filter((message) => message.id === "msg_prompt")).toHaveLength(1) + expect(store.data.part.msg_prompt).toMatchObject([{ type: "text", text: "hello" }]) }) - test("does not remove parts confirmed by part events", () => { - const message = userMessage("message") - const part = textPart(message.id) + test("uses the prompt response when the admission event was missed", () => { const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - store.apply({ type: "message.updated", properties: { sessionID: "child", info: message } }) - store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } }) + store.inbox.echo(promptEcho("msg_prompt")) + store.inbox.confirm({ + id: "msg_prompt", + sessionID: "child", + timeCreated: 2, + type: "user", + delivery: "steer", + payload: { text: "hello" }, + }) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) + store.applyV2({ + id: "evt_delivered", + created: Date.now() + 1, + type: "session.inbox.delivered", + durable: { aggregateID: "child", seq: 2, version: 1 }, + data: { sessionID: "child", inboxID: "msg_prompt" }, + } as OpenCodeEvent) - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([part]) + expect(store.data.pending.child).toEqual([]) + expect(store.data.input.child).toEqual([]) + expect(store.data.session_message.child).toMatchObject([{ id: "msg_prompt", type: "user", text: "hello" }]) + expect(store.data.message.child?.filter((message) => message.id === "msg_prompt")).toHaveLength(1) + expect(store.data.part.msg_prompt).toMatchObject([{ type: "text", text: "hello" }]) }) - test("treats a part event as confirmation when it precedes the message event", () => { - const message = userMessage("message") - const part = textPart(message.id) + test("keeps a durable admission when the HTTP request later fails", () => { const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - store.apply({ type: "message.part.updated", properties: { sessionID: "child", part, time: 2 } }) + store.inbox.echo(promptEcho("msg_prompt")) + store.applyV2({ + id: "evt_prompt", + created: 2, + type: "session.inbox.enqueued", + durable: { aggregateID: "child", seq: 1, version: 1 }, + data: { + sessionID: "child", + inboxID: "msg_prompt", + item: { type: "user", delivery: "steer", payload: { text: "hello" } }, + }, + } as OpenCodeEvent) - store.optimistic.remove({ sessionID: "child", messageID: message.id }) + expect(store.inbox.clearEcho({ sessionID: "child", messageID: "msg_prompt" })).toBe(false) + expect(store.data.pending.child).toHaveLength(1) + expect(store.data.message.child?.map((message) => message.id)).toEqual(["msg_prompt"]) + }) - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([part]) + test("places durable admission after delayed selection events", () => { + const store = setup({ child: session("child") }).store + store.remember(session("child")) + store.inbox.echo(promptEcho("msg_prompt")) + store.applyV2({ + id: "evt_agent", + created: 1, + type: "session.agent.selected", + durable: { aggregateID: "child", seq: 1, version: 1 }, + data: { sessionID: "child", agent: "review" }, + } as OpenCodeEvent) + store.applyV2({ + id: "evt_model", + created: 2, + type: "session.model.selected", + durable: { aggregateID: "child", seq: 2, version: 1 }, + data: { sessionID: "child", model: { id: "new-model", providerID: "new-provider" } }, + } as OpenCodeEvent) + store.applyV2({ + id: "evt_prompt", + created: 3, + type: "session.inbox.enqueued", + durable: { aggregateID: "child", seq: 3, version: 1 }, + data: { + sessionID: "child", + inboxID: "msg_prompt", + item: { type: "user", delivery: "steer", payload: { text: "hello" } }, + }, + } as OpenCodeEvent) + + expect(store.data.session_message.child?.map((message) => message.type)).toEqual([ + "agent-switched", + "model-switched", + "user", + ]) + expect(store.data.message.child?.find((message) => message.id === "msg_prompt")).toMatchObject({ + agent: "review", + model: { providerID: "new-provider", modelID: "new-model" }, + }) + }) + + test("removes an echoed prompt when submission fails", () => { + const store = setup({ child: session("child") }).store + store.inbox.echo(promptEcho("msg_prompt")) + + expect(store.inbox.clearEcho({ sessionID: "child", messageID: "msg_prompt" })).toBe(true) + + expect(store.data.pending.child).toEqual([]) + expect(store.data.input.child).toEqual([]) + expect(store.data.session_message.child).toBeUndefined() + expect(store.data.message.child).toEqual([]) + expect(store.data.part.msg_prompt).toBeUndefined() + }) + + test("removes a response-confirmed echo when the server cancels it", () => { + const store = setup({ child: session("child") }).store + store.inbox.echo(promptEcho("msg_prompt")) + store.inbox.confirm({ + id: "msg_prompt", + sessionID: "child", + timeCreated: 1, + type: "user", + delivery: "steer", + payload: { text: "hello" }, + }) + + store.applyV2({ + id: "evt_cancelled", + created: 2, + type: "session.inbox.cancelled", + durable: { aggregateID: "child", seq: 2, version: 1 }, + data: { sessionID: "child", inboxID: "msg_prompt" }, + } as OpenCodeEvent) + + expect(store.data.pending.child).toEqual([]) + expect(store.data.message.child).toEqual([]) + expect(store.data.part.msg_prompt).toBeUndefined() }) test("clears stale parts when the initial page has none", async () => { @@ -1469,28 +1590,6 @@ describe("server session", () => { expect(store.data.part[message.id]).toBeUndefined() }) - test("preserves optimistic re-adds across message retries", async () => { - const failed = Promise.withResolvers() - const retried = Promise.withResolvers() - const message = userMessage("message") - const stale = textPart(message.id, { id: "stale", text: "stale" }) - const optimistic = textPart(message.id, { id: "optimistic", text: "optimistic" }) - const client = messageClient(response([{ info: message, parts: [stale] }]), failed.promise, retried.promise) - const store = createServerSession(client, { retry: retryImmediately }) - await store.sync("child") - const loading = store.sync("child", { force: true }) - - store.apply({ type: "message.removed", properties: { sessionID: "child", messageID: message.id } }) - store.optimistic.add({ sessionID: "child", message, parts: [optimistic] }) - failed.reject(new Error("failed to fetch")) - await client.requested(3) - retried.resolve(response([{ info: message, parts: [stale] }])) - await loading - - expect(store.data.message.child).toEqual([message]) - expect(store.data.part[message.id]).toEqual([optimistic]) - }) - test("accepts part omission from a successful retry after an earlier delta", async () => { const failed = Promise.withResolvers() const retried = Promise.withResolvers() @@ -1654,33 +1753,6 @@ describe("server session", () => { expect(store.data.part[message.id]).toBeUndefined() }) - test("does not cache skipped optimistic parts", () => { - const message = userMessage("message") - const part = { id: "part", sessionID: "child", messageID: message.id, type: "step-start" as const } - const store = setup({ child: session("child") }).store - - store.optimistic.add({ sessionID: "child", message, parts: [part] }) - - expect(store.data.part[message.id]).toEqual([]) - }) - - test("clears stale delta buffers when replacing optimistic parts", () => { - const message = userMessage("message") - const stale = textPart(message.id, { id: "stale", text: "stale" }) - const optimistic = textPart(message.id, { id: "optimistic", text: "optimistic" }) - const store = setup({ child: session("child") }).store - store.optimistic.add({ sessionID: "child", message, parts: [stale] }) - store.apply({ - type: "message.part.delta", - properties: { sessionID: "child", messageID: message.id, partID: stale.id, field: "text", delta: " delta" }, - }) - - store.optimistic.add({ sessionID: "child", message, parts: [optimistic] }) - - expect(store.data.part_text_accum_delta[stale.id]).toBeUndefined() - expect(store.data.part_text_accum_delta[optimistic.id]).toBeUndefined() - }) - test("preserves removals during history prepend", async () => { const pending = deferredResponse() const latest = userMessage("message-2", { time: { created: 2 } }) @@ -1906,24 +1978,7 @@ describe("server session", () => { test("preserves pinned session content under server-wide cache pressure", () => { const ctx = setup({}) ctx.store.pin("active") - ctx.store.optimistic.add({ - sessionID: "active", - message: { - id: "message", - sessionID: "active", - role: "assistant", - time: { created: 1 }, - parentID: "parent", - modelID: "model", - providerID: "provider", - mode: "build", - agent: "agent", - path: { cwd: "/repo", root: "/repo" }, - cost: 0, - tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, - }, - parts: [], - }) + ctx.store.inbox.echo({ ...promptEcho("message", "keep"), sessionID: "active" }) for (let index = 0; index < 50; index++) { ctx.store.remember(session(`session-${index}`)) diff --git a/packages/app/src/context/server-session.ts b/packages/app/src/context/server-session.ts index 193dde10e04..19cc5830b07 100644 --- a/packages/app/src/context/server-session.ts +++ b/packages/app/src/context/server-session.ts @@ -18,6 +18,13 @@ import { compareMessages, messageKey, normalizeSessionMessages } from "@/utils/s import { dropSessionCaches, pickSessionCacheEvictions, SESSION_CACHE_LIMIT } from "./global-sync/session-cache" import { createV2SessionReducer, type V2SessionReduction } from "./server-session-v2-reducer" import type { ServerApi } from "@/utils/server" +import { + createCommentMetadata, + formatCommentNote, + parseCommentNote, + readCommentMetadata, + type PromptComment, +} from "@/utils/comment-note" type MessageApi = ServerApi["message"] @@ -28,31 +35,6 @@ const historyMessagePageSize = 200 const sessionInfoLimit = 2_048 const emptyIDs: ReadonlySet = new Set() -function projectMessageSource(message: Message): SessionMessageInfo[] { - if (message.role === "user") { - return [ - { id: `${message.id}:agent`, type: "agent-switched", agent: message.agent, time: message.time }, - { - id: `${message.id}:model`, - type: "model-switched", - model: { id: message.model.modelID, providerID: message.model.providerID, variant: message.model.variant }, - time: message.time, - }, - { id: message.id, type: "user", text: "", time: message.time }, - ] - } - return [ - { - id: message.id, - type: "assistant", - agent: message.agent ?? message.mode, - model: { id: message.modelID, providerID: message.providerID, variant: message.variant }, - content: [], - time: message.time, - }, - ] -} - function needsOlderTurnRoot(source: readonly SessionMessageInfo[]) { const boundary = source.find( (message) => @@ -64,13 +46,6 @@ function needsOlderTurnRoot(source: readonly SessionMessageInfo[]) { return boundary?.type === "assistant" } -type OptimisticItem = { - message: Message - parts: Part[] - confirmedParts?: Part[] - confirmedMessage?: boolean -} - type MessagePage = { session: Message[] part: { id: string; part: Part[] }[] @@ -81,6 +56,18 @@ type MessagePage = { complete: boolean } +export type PromptEcho = { + sessionID: string + messageID: string + text: string + displayText: string + agent: string + model: { providerID: string; modelID: string; variant?: string } + files?: { uri: string; mime: string; name?: string; mention?: { start: number; end: number; text: string } }[] + agents?: { name: string; mention?: { start: number; end: number; text: string } }[] + comments: PromptComment[] +} + // Most markers describe the current HTTP attempt; deltaParts persists non-durable stream state across retries. type MessageLoadState = { touchedMessages: Set @@ -90,7 +77,6 @@ type MessageLoadState = { deltaParts: Map> carriedDeltaParts: Map> removedParts: Map> - optimisticParts: Map> orphanParents: Set clearedMessageParts: Set touchedSource: Set @@ -101,34 +87,6 @@ type MessageLoadBaseline = Pick< "touchedMessages" | "retainedMessages" | "touchedParts" | "clearedMessageParts" > -function mergeOptimisticPage(page: MessagePage, items: OptimisticItem[]) { - if (items.length === 0) return { ...page, observed: [] as { messageID: string; parts: Part[] }[] } - const session = [...page.session] - const part = new Map(page.part.map((item) => [item.id, item.part])) - const observed: { messageID: string; parts: Part[] }[] = [] - for (const item of items) { - const result = Binary.search(session, messageKey(item.message), messageKey) - const found = result.found - if (!found) session.splice(result.index, 0, item.message) - const current = part.get(item.message.id) - const confirmed = found ? item.parts.filter((part) => current?.some((value) => value.id === part.id)) : [] - if (found) observed.push({ messageID: item.message.id, parts: confirmed }) - part.set( - item.message.id, - merge( - found ? (current ?? []) : merge(item.confirmedParts ?? [], current ?? []), - item.parts.filter((part) => !confirmed.includes(part)), - ), - ) - } - return { - ...page, - session, - part: [...part.entries()].sort((a, b) => cmp(a[0], b[0])).map(([id, parts]) => ({ id, part: parts })), - observed, - } -} - function runInflight(map: Map>, key: string, task: () => Promise) { const pending = map.get(key) if (pending) return pending @@ -212,7 +170,6 @@ export function createServerSession( const requests = new Map>() const inflight = new Map>() const inflightTodo = new Map>() - const optimistic = new Map>() const v2 = createV2SessionReducer() const pendingRevision = new Map() const formRevision = new Map() @@ -223,7 +180,45 @@ export function createServerSession( const pendingParts = new Map>>() const orphanParts = new Map>() const removedMessages = new Map>() + const echoes = new Map>() + const messageSnapshots = new Map>() + const settledInputs = new Map>() const deltaBases = new Map() + const markEcho = (sessionID: string, messageID: string) => { + const messages = echoes.get(sessionID) ?? new Map() + messages.set(messageID, "sending") + echoes.set(sessionID, messages) + } + const confirmEcho = (sessionID: string, messageID: string) => { + const messages = echoes.get(sessionID) + if (!messages?.has(messageID)) return false + messages.set(messageID, "admitted") + return true + } + const releaseEcho = (sessionID: string, messageID: string) => { + const messages = echoes.get(sessionID) + const state = messages?.get(messageID) + if (!messages || !state) return + messages.delete(messageID) + if (messages.size === 0) echoes.delete(sessionID) + return state + } + const present = (messageID: string, parts: Part[]) => { + const local = data.part[messageID] ?? [] + const comments = local.filter( + (part) => + part.type === "text" && + part.synthetic && + (readCommentMetadata(part.metadata) !== undefined || parseCommentNote(part.text) !== undefined), + ) + if (!comments.length) return parts + const text = local.find((part) => part.type === "text" && !part.synthetic) + const projected = parts.flatMap((part) => { + if (part.id !== `${messageID}:text:0` || part.type !== "text") return [part] + return text?.type === "text" && text.text ? [{ ...part, text: text.text }] : [] + }) + return merge(projected, comments) + } const deleteMessageParts = ( cache: { part: Record; part_text_accum_delta: Record }, messageID: string, @@ -253,18 +248,6 @@ export function createServerSession( at: {} as Record, }) - const indexProjectedMessage = (message: Message) => { - const current = data.session_message[message.sessionID] ?? [] - if (current.some((item) => item.id === message.id)) return - const projected = projectMessageSource(message) - const projectedIDs = new Set(projected.map((item) => item.id)) - setData( - "session_message", - message.sessionID, - reconcile([...current.filter((item) => !projectedIDs.has(item.id)), ...projected]), - ) - } - const remember = (session: SessionInfo) => { setData("info", session.id, reconcile(session)) infoSeen.delete(session.id) @@ -276,7 +259,7 @@ export function createServerSession( ...inflight.keys(), ...inflightTodo.keys(), ...messageLoads.keys(), - ...optimistic.keys(), + ...echoes.keys(), ...Object.entries(data.permission) .filter(([, items]) => items.length > 0) .map(([sessionID]) => sessionID), @@ -352,65 +335,6 @@ export function createServerSession( return { session, root } } - const clearOptimistic = (sessionID: string, messageID?: string) => { - if (!messageID) { - optimistic.delete(sessionID) - return - } - const items = optimistic.get(sessionID) - if (!items) return - items.delete(messageID) - if (items.size === 0) optimistic.delete(sessionID) - } - - const clearOptimisticPart = (sessionID: string, messageID: string, partID: string) => { - const items = optimistic.get(sessionID) - const item = items?.get(messageID) - if (!items || !item) return - const parts = item.parts.filter((part) => part.id !== partID) - const confirmedParts = item.confirmedParts?.filter((part) => part.id !== partID) - if (parts.length === 0) { - clearOptimistic(sessionID, messageID) - return - } - items.set(messageID, { ...item, parts, confirmedParts, confirmedMessage: true }) - } - - const confirmOptimisticPart = (sessionID: string, messageID: string, part: Part) => { - const items = optimistic.get(sessionID) - const item = items?.get(messageID) - if (!items || !item) return - const parts = item.parts.filter((value) => value.id !== part.id) - if (parts.length === 0) { - clearOptimistic(sessionID, messageID) - return - } - items.set(messageID, { - ...item, - parts, - confirmedParts: merge(item.confirmedParts ?? [], [part]), - confirmedMessage: true, - }) - } - - const confirmOptimistic = (sessionID: string, messageID: string, confirmedParts: Part[]) => { - const items = optimistic.get(sessionID) - const item = items?.get(messageID) - if (!items || !item) return - const confirmed = new Set(confirmedParts.map((part) => part.id)) - const parts = item.parts.filter((part) => !confirmed.has(part.id)) - if (parts.length === 0) { - clearOptimistic(sessionID, messageID) - return - } - items.set(messageID, { - ...item, - parts, - confirmedParts: merge(item.confirmedParts ?? [], confirmedParts), - confirmedMessage: true, - }) - } - const trackPartChange = (sessionID: string, messageID: string, partID: string) => { const load = messageLoads.get(sessionID) if (!load) return @@ -448,14 +372,6 @@ export function createServerSession( const messages = data.message[sessionID] if (messages?.some((message) => message.id === messageID)) load.retainedMessages.add(messageID) } - for (const [messageID, parts] of load.optimisticParts) { - load.removedMessages.delete(messageID) - load.clearedMessageParts.add(messageID) - load.touchedMessages.add(messageID) - const touched = load.touchedParts.get(messageID) ?? new Set() - parts.forEach((partID) => touched.add(partID)) - load.touchedParts.set(messageID, touched) - } baseline?.touchedMessages.forEach((messageID) => load.touchedMessages.add(messageID)) baseline?.retainedMessages.forEach((messageID) => load.retainedMessages.add(messageID)) baseline?.clearedMessageParts.forEach((messageID) => load.clearedMessageParts.add(messageID)) @@ -486,7 +402,9 @@ export function createServerSession( sessionIDs.forEach((sessionID) => { messageHydrationRevision.set(sessionID, (messageHydrationRevision.get(sessionID) ?? 0) + 1) generations.delete(sessionID) - clearOptimistic(sessionID) + echoes.delete(sessionID) + messageSnapshots.delete(sessionID) + settledInputs.delete(sessionID) requests.delete(sessionID) inflight.delete(sessionID) inflightTodo.delete(sessionID) @@ -521,7 +439,7 @@ export function createServerSession( ...inflight.keys(), ...inflightTodo.keys(), ...messageLoads.keys(), - ...optimistic.keys(), + ...echoes.keys(), ...Object.entries(data.permission) .filter(([, items]) => items.length > 0) .map(([sessionID]) => sessionID), @@ -598,9 +516,10 @@ export function createServerSession( ) => { for (const item of items) { if (!messageIDs.has(item.id)) continue - const fetched = load?.clearedMessageParts.has(item.id) - ? [] - : item.part.filter((part) => !SKIP_PARTS.has(part.type)) + const fetched = present( + item.id, + load?.clearedMessageParts.has(item.id) ? [] : item.part.filter((part) => !SKIP_PARTS.has(part.type)), + ) const fetchedIDs = new Set(fetched.map((part) => part.id)) const pending = pendingParts.get(sessionID)?.get(item.id) const touched = new Set([...(load?.touchedParts.get(item.id) ?? []), ...(pending ?? [])]) @@ -651,25 +570,37 @@ export function createServerSession( preserveUnfetched: boolean | ((message: Message) => boolean), cleanupOrphans: boolean, ) => { + if (page.sourceMode === "latest") + messageSnapshots.set(sessionID, new Set((page.source ?? []).map((message) => message.id))) + page.source?.forEach((message) => releaseEcho(sessionID, message.id)) const source = page.source ? (() => { const incoming = new Map(page.source.map((message) => [message.id, message])) const existing = data.session_message[sessionID] ?? [] const boundary = Math.min(...page.source.map((message) => message.time.created)) + const inbox = new Set(data.input[sessionID] ?? []) const current = existing.filter( (message) => !incoming.has(message.id) && + !inbox.has(message.id) && (page.sourceMode === "older" || load?.touchedSource.has(message.id) || (!page.complete && message.time.created < boundary)), ) + // message.list never returns admitted-but-undelivered inbox entries; keep them after the + // fetched history until a delivered or cancelled event resolves them. + const admitted = existing.filter((message) => !incoming.has(message.id) && inbox.has(message.id)) + const combined = + page.sourceMode === "older" + ? [...page.source, ...current, ...admitted] + : [...current, ...page.source, ...admitted] const live = new Map(existing.map((message) => [message.id, message])) - return (page.sourceMode === "older" ? [...page.source, ...current] : [...current, ...page.source]).map( - (message) => (load?.touchedSource.has(message.id) ? (live.get(message.id) ?? message) : message), + return combined.map((message) => + load?.touchedSource.has(message.id) ? (live.get(message.id) ?? message) : message, ) })() : undefined - const projected = + const merged = page.projectSource && source ? (() => { const normalized = normalizeSessionMessages(sessionID, source) @@ -682,16 +613,15 @@ export function createServerSession( } })() : page - const merged = mergeOptimisticPage(projected, [...(optimistic.get(sessionID)?.values() ?? [])]) - merged.observed.forEach((item) => { - if (!load?.clearedMessageParts.has(item.messageID)) confirmOptimistic(sessionID, item.messageID, item.parts) - }) const touchedMessages = new Set([...(load?.touchedMessages ?? []), ...(removedMessages.get(sessionID) ?? [])]) const messages = reconcileFetched(merged.session, data.message[sessionID] ?? [], { touched: touchedMessages, retained: load?.retainedMessages, removed: load?.removedMessages, - preserveUnfetched, + preserveUnfetched: (message) => + echoes.get(sessionID)?.has(message.id) === true || + preserveUnfetched === true || + (typeof preserveUnfetched === "function" && preserveUnfetched(message)), compare: compareMessages, }) batch(() => { @@ -723,7 +653,6 @@ export function createServerSession( deltaParts: new Map(), carriedDeltaParts: new Map(), removedParts: new Map(), - optimisticParts: new Map(), orphanParents: new Set(), clearedMessageParts: new Set(), touchedSource: new Set(), @@ -744,11 +673,7 @@ export function createServerSession( const users = new Set([ ...page.session.filter((message) => message.role === "user").map((message) => message.id), ...(data.message[sessionID] ?? []) - .filter((message) => { - if (message.role !== "user") return false - const item = optimistic.get(sessionID)?.get(message.id) - return load.touchedMessages.has(message.id) && (!item || item.confirmedMessage === true) - }) + .filter((message) => message.role === "user" && load.touchedMessages.has(message.id)) .map((message) => message.id), ]) const parentIDs = [ @@ -893,7 +818,7 @@ export function createServerSession( apply({ type: "message.updated", properties: { sessionID: reduction.sessionID, info: message } }) } for (const messageID of touched) { - const next = normalized.parts.get(messageID) ?? [] + const next = present(messageID, normalized.parts.get(messageID) ?? []) const nextIDs = new Set(next.map((part) => part.id)) for (const part of next) { apply({ type: "message.part.updated", properties: { sessionID: reduction.sessionID, part } }) @@ -926,6 +851,67 @@ export function createServerSession( .catch(() => {}) } + const removeEcho = (sessionID: string, messageID: string) => { + if (!releaseEcho(sessionID, messageID)) return false + pendingRevision.set(sessionID, (pendingRevision.get(sessionID) ?? 0) + 1) + const load = messageLoads.get(sessionID) + load?.touchedMessages.add(messageID) + load?.removedMessages.add(messageID) + load?.clearedMessageParts.add(messageID) + batch(() => { + setData("pending", sessionID, (items) => items?.filter((item) => item.id !== messageID)) + setData("input", sessionID, (items) => items?.filter((id) => id !== messageID)) + setData("message", sessionID, (messages) => messages?.filter((message) => message.id !== messageID)) + setData(produce((draft) => deleteMessageParts(draft, messageID))) + }) + return true + } + + const confirmInbox = (item: SessionInboxInfo) => { + if (!confirmEcho(item.sessionID, item.id)) return false + v2.confirm(item) + pendingRevision.set(item.sessionID, (pendingRevision.get(item.sessionID) ?? 0) + 1) + const current = data.pending[item.sessionID] ?? [] + const index = current.findIndex((entry) => entry.id === item.id) + if (index < 0) setData("pending", item.sessionID, [...current, item]) + if (index >= 0) setData("pending", item.sessionID, index, reconcile(item)) + return true + } + + const reconcileInbox = (sessionID: string) => { + const pending = new Set((data.pending[sessionID] ?? []).map((item) => item.id)) + const fetched = messageSnapshots.get(sessionID) ?? new Set() + const removed = [...(settledInputs.get(sessionID) ?? [])].filter( + (messageID) => !pending.has(messageID) && !fetched.has(messageID), + ) + settledInputs.delete(sessionID) + if (removed.length) { + const ids = new Set(removed) + const source = data.session_message[sessionID] ?? [] + projectV2({ + sessionID, + messages: source.filter((message) => !ids.has(message.id)), + touched: [], + removed: source.filter((message) => ids.has(message.id)).map((message) => message.id), + }) + } + + const messages = echoes.get(sessionID) + if (!messages) return + const projected = new Set((data.session_message[sessionID] ?? []).map((message) => message.id)) + for (const [messageID, state] of messages) { + if (projected.has(messageID)) { + releaseEcho(sessionID, messageID) + continue + } + if (pending.has(messageID)) { + confirmEcho(sessionID, messageID) + continue + } + if (state === "admitted") removeEcho(sessionID, messageID) + } + } + const applyV2 = (event: OpenCodeEvent) => { if (event.type === "form.created") { formRevision.set(event.data.form.sessionID, (formRevision.get(event.data.form.sessionID) ?? 0) + 1) @@ -949,6 +935,9 @@ export function createServerSession( } if (!("data" in event) || !("sessionID" in event.data) || typeof event.data.sessionID !== "string") return const sessionID = event.data.sessionID + if (event.type === "session.inbox.enqueued" || event.type === "session.inbox.delivered") + releaseEcho(sessionID, event.data.inboxID) + if (event.type === "session.inbox.cancelled") removeEcho(sessionID, event.data.inboxID) if ( event.type === "session.inbox.enqueued" || event.type === "session.inbox.delivery.changed" || @@ -960,11 +949,10 @@ export function createServerSession( pendingRevision.set(sessionID, (pendingRevision.get(sessionID) ?? 0) + 1) if (event.type === "session.inbox.enqueued") { const current = data.pending[sessionID] ?? [] - if (!current.some((item) => item.id === event.data.inboxID)) - setData("pending", sessionID, [ - ...current, - { id: event.data.inboxID, sessionID, timeCreated: event.created, ...event.data.item }, - ]) + const item = { id: event.data.inboxID, sessionID, timeCreated: event.created, ...event.data.item } + const index = current.findIndex((entry) => entry.id === event.data.inboxID) + if (index < 0) setData("pending", sessionID, [...current, item]) + if (index >= 0) setData("pending", sessionID, index, reconcile(item)) if (event.data.item.type !== "compaction" && !data.input[sessionID]?.includes(event.data.inboxID)) setData("input", sessionID, [...(data.input[sessionID] ?? []), event.data.inboxID]) } @@ -1101,16 +1089,9 @@ export function createServerSession( } case "message.updated": { const info = (event.properties as { info: Message }).info - indexProjectedMessage(info) const load = messageLoads.get(info.sessionID) load?.touchedMessages.add(info.id) load?.removedMessages.delete(info.id) - const items = optimistic.get(info.sessionID) - const item = items?.get(info.id) - if (items && item) { - if (item.parts.length === 0) clearOptimistic(info.sessionID, info.id) - if (item.parts.length > 0) items.set(info.id, { ...item, confirmedMessage: true }) - } const orphans = orphanParts.get(info.sessionID) orphans?.delete(info.id) if (orphans?.size === 0) orphanParts.delete(info.sessionID) @@ -1123,13 +1104,18 @@ export function createServerSession( return } const result = Binary.search(messages, messageKey(info), messageKey) - if (result.found) setData("message", info.sessionID, result.index, reconcile(info)) - if (!result.found) - setData("message", info.sessionID, (value = []) => { - const next = value.slice() - next.splice(result.index, 0, info) - return next - }) + if (result.found) { + setData("message", info.sessionID, result.index, reconcile(info)) + return + } + // Delivery rewrites time.created, changing the sort key; reposition instead of duplicating. + setData("message", info.sessionID, (value = []) => { + const next = value.slice() + const moved = next.findIndex((message) => message.id === info.id) + if (moved >= 0) next.splice(moved, 1) + next.splice(moved >= 0 && moved < result.index ? result.index - 1 : result.index, 0, info) + return next + }) return } case "message.removed": { @@ -1144,13 +1130,11 @@ export function createServerSession( load?.deltaParts.delete(props.messageID) load?.carriedDeltaParts.delete(props.messageID) load?.removedParts.delete(props.messageID) - load?.optimisticParts.delete(props.messageID) pendingParts.get(props.sessionID)?.delete(props.messageID) if (pendingParts.get(props.sessionID)?.size === 0) pendingParts.delete(props.sessionID) const removedMessagesForSession = removedMessages.get(props.sessionID) ?? new Set() removedMessagesForSession.add(props.messageID) removedMessages.set(props.sessionID, removedMessagesForSession) - clearOptimistic(props.sessionID, props.messageID) setData( produce((draft) => { const messages = draft.message[props.sessionID] @@ -1196,12 +1180,8 @@ export function createServerSession( pending?.delete(part.id) if (pending?.size === 0) pendingParts.get(part.sessionID)?.delete(part.messageID) if (pendingParts.get(part.sessionID)?.size === 0) pendingParts.delete(part.sessionID) - const optimistic = load?.optimisticParts.get(part.messageID) - optimistic?.delete(part.id) - if (optimistic?.size === 0) load?.optimisticParts.delete(part.messageID) deltaBases.delete(part.id) trackPartChange(part.sessionID, part.messageID, part.id) - confirmOptimisticPart(part.sessionID, part.messageID, part) setData( "part_text_accum_delta", produce((draft) => void delete draft[part.id]), @@ -1240,12 +1220,8 @@ export function createServerSession( const parts = load.removedParts.get(props.messageID) ?? new Set() parts.add(props.partID) load.removedParts.set(props.messageID, parts) - const optimistic = load.optimisticParts.get(props.messageID) - optimistic?.delete(props.partID) - if (optimistic?.size === 0) load.optimisticParts.delete(props.messageID) } trackPartChange(props.sessionID, props.messageID, props.partID) - clearOptimisticPart(props.sessionID, props.messageID, props.partID) setData( produce((draft) => { delete draft.part_text_accum_delta[props.partID] @@ -1354,25 +1330,30 @@ export function createServerSession( while (true) { const pendingAt = pendingRevision.get(sessionID) ?? 0 const formAt = formRevision.get(sessionID) ?? 0 + const previous = new Set(data.input[sessionID] ?? []) const result = await load() const pendingStable = (pendingRevision.get(sessionID) ?? 0) === pendingAt const formStable = (formRevision.get(sessionID) ?? 0) === formAt if (pendingStable) { + const current = new Set(result.pending.filter((item) => item.type !== "compaction").map((item) => item.id)) + const settled = settledInputs.get(sessionID) ?? new Set() + previous.forEach((messageID) => { + if (!current.has(messageID)) settled.add(messageID) + }) + if (settled.size) settledInputs.set(sessionID, settled) + result.pending.forEach(v2.confirm) setData("pending", sessionID, reconcile(result.pending)) - setData( - "input", - sessionID, - reconcile(result.pending.filter((item) => item.type !== "compaction").map((item) => item.id)), - ) + setData("input", sessionID, reconcile([...current])) } if (formStable) setData("form", sessionID, reconcile(result.forms)) if (pendingStable && formStable) return } }, refreshPinned(hydrateTransient: (sessionID: string) => Promise) { + const sessions = [...pinned.keys()] return Promise.all( - [...pinned.keys()].flatMap((sessionID) => [sync(sessionID, { force: true }), hydrateTransient(sessionID)]), - ).then(() => undefined) + sessions.flatMap((sessionID) => [sync(sessionID, { force: true }), hydrateTransient(sessionID)]), + ).then(() => sessions.forEach(reconcileInbox)) }, invalidate() { invalidationRevision += 1 @@ -1390,68 +1371,72 @@ export function createServerSession( fresh(sessionID: string, ttl: number) { return Date.now() - (meta.at[sessionID] ?? 0) <= ttl }, - optimistic: { - add(input: { sessionID: string; message: Message; parts: Part[] }) { - const parts = input.parts - .filter((part) => !!part?.id && !SKIP_PARTS.has(part.type)) - .sort((a, b) => cmp(a.id, b.id)) - const load = messageLoads.get(input.sessionID) - if (load?.clearedMessageParts.has(input.message.id)) { - const touched = load.touchedParts.get(input.message.id) ?? new Set() - parts.forEach((part) => touched.add(part.id)) - load.touchedParts.set(input.message.id, touched) + inbox: { + echo(input: PromptEcho) { + const created = Date.now() + const files = input.files?.map((file) => ({ + data: "", + mime: file.mime, + source: { type: "uri" as const, uri: file.uri }, + name: file.name, + mention: file.mention, + })) + const item: SessionInboxInfo = { + id: input.messageID, + sessionID: input.sessionID, + timeCreated: created, + type: "user", + delivery: "steer", + payload: { text: input.text, files, agents: input.agents }, } - if (load) { - load.removedMessages.delete(input.message.id) - load.optimisticParts.set(input.message.id, new Set(parts.map((part) => part.id))) - } - const items = optimistic.get(input.sessionID) - const removedMessagesForSession = removedMessages.get(input.sessionID) - removedMessagesForSession?.delete(input.message.id) - if (removedMessagesForSession?.size === 0) removedMessages.delete(input.sessionID) - if (items) items.set(input.message.id, { ...input, parts, confirmedParts: [] }) - if (!items) - optimistic.set(input.sessionID, new Map([[input.message.id, { ...input, parts, confirmedParts: [] }]])) - indexProjectedMessage(input.message) - setData("message", input.sessionID, (messages = []) => merge(messages, [input.message]).sort(compareMessages)) - setData( - "part_text_accum_delta", - produce((draft) => { - for (const part of [...(data.part[input.message.id] ?? []), ...parts]) { - delete draft[part.id] - deltaBases.delete(part.id) - } - }), - ) - setData("part", input.message.id, parts) + const projected = normalizeSessionMessages(input.sessionID, [ + { id: `${input.messageID}:agent`, type: "agent-switched", agent: input.agent, time: { created } }, + { + id: `${input.messageID}:model`, + type: "model-switched", + model: { + id: input.model.modelID, + providerID: input.model.providerID, + variant: input.model.variant, + }, + time: { created }, + }, + { + id: input.messageID, + type: "user", + text: input.displayText, + files, + agents: input.agents, + time: { created }, + }, + ]) + const message = projected.messages[0]! + const comments: Part[] = input.comments.map((comment, index) => ({ + id: `${input.messageID}:comment:${index}`, + sessionID: input.sessionID, + messageID: input.messageID, + type: "text", + text: formatCommentNote(comment), + synthetic: true, + metadata: createCommentMetadata(comment), + })) + const parts = merge(projected.parts.get(input.messageID) ?? [], comments).sort((a, b) => cmp(a.id, b.id)) + removedMessages.get(input.sessionID)?.delete(input.messageID) + markEcho(input.sessionID, input.messageID) + pendingRevision.set(input.sessionID, (pendingRevision.get(input.sessionID) ?? 0) + 1) + batch(() => { + setData("pending", input.sessionID, (items = []) => [...items.filter((entry) => entry.id !== item.id), item]) + if (!data.input[input.sessionID]?.includes(input.messageID)) + setData("input", input.sessionID, [...(data.input[input.sessionID] ?? []), input.messageID]) + setData("message", input.sessionID, (messages = []) => merge(messages, [message]).sort(compareMessages)) + setData("part", input.messageID, parts) + }) }, - remove(input: { sessionID: string; messageID: string }) { - const item = optimistic.get(input.sessionID)?.get(input.messageID) - if (!item) return - messageLoads.get(input.sessionID)?.optimisticParts.delete(input.messageID) - clearOptimistic(input.sessionID, input.messageID) - if (item.confirmedMessage) { - const partIDs = new Set(item.parts.map((part) => part.id)) - setData( - produce((draft) => { - for (const part of item.parts) { - delete draft.part_text_accum_delta[part.id] - deltaBases.delete(part.id) - } - const parts = draft.part[input.messageID] - if (!parts) return - draft.part[input.messageID] = parts.filter((part) => !partIDs.has(part.id)) - if (draft.part[input.messageID]?.length === 0) delete draft.part[input.messageID] - }), - ) - return - } - const projectedIDs = new Set(projectMessageSource(item.message).map((message) => message.id)) - setData("session_message", input.sessionID, (messages) => - messages?.filter((message) => !projectedIDs.has(message.id)), - ) - setData("message", input.sessionID, (messages) => messages?.filter((message) => message.id !== input.messageID)) - setData(produce((draft) => deleteMessageParts(draft, input.messageID))) + confirm: confirmInbox, + reconcile: reconcileInbox, + clearEcho(input: { sessionID: string; messageID: string }) { + if (echoes.get(input.sessionID)?.get(input.messageID) !== "sending") return false + return removeEcho(input.sessionID, input.messageID) }, }, async todo(sessionID: string, request?: { force?: boolean }) { diff --git a/packages/app/src/context/server-sync.tsx b/packages/app/src/context/server-sync.tsx index 16c085cf5f2..e4befa5a521 100644 --- a/packages/app/src/context/server-sync.tsx +++ b/packages/app/src/context/server-sync.tsx @@ -231,7 +231,10 @@ export function createServerSyncContextInner(serverSDK: ServerSDK) { return { pending, forms } }) } - const hydrateSession = (sessionID: string) => Promise.all([session.sync(sessionID), hydrateSessionState(sessionID)]) + const hydrateSession = async (sessionID: string) => { + await Promise.all([session.sync(sessionID), hydrateSessionState(sessionID)]) + session.inbox.reconcile(sessionID) + } const [configQuery, providerQuery, pathQuery] = useQueries(() => ({ queries: [ diff --git a/packages/app/src/context/sync-optimistic.test.ts b/packages/app/src/context/sync-optimistic.test.ts deleted file mode 100644 index 15adcf441a3..00000000000 --- a/packages/app/src/context/sync-optimistic.test.ts +++ /dev/null @@ -1,137 +0,0 @@ -import { describe, expect, test } from "bun:test" -import type { Message, Part } from "@/types" -import { applyOptimisticAdd, applyOptimisticRemove, mergeOptimisticPage } from "./sync" - -type Text = Extract - -const userMessage = (id: string, sessionID: string, created = 1): Message => ({ - id, - sessionID, - role: "user", - time: { created }, - agent: "assistant", - model: { providerID: "openai", modelID: "gpt" }, -}) - -const textPart = (id: string, sessionID: string, messageID: string): Text => ({ - id, - sessionID, - messageID, - type: "text", - text: id, -}) - -describe("sync optimistic reducers", () => { - test("applyOptimisticAdd inserts by creation time", () => { - const sessionID = "ses_1" - const draft = { - message: { [sessionID]: [userMessage("msg_z", sessionID, 1)] }, - part: {} as Record, - } - - applyOptimisticAdd(draft, { - sessionID, - message: userMessage("msg_a", sessionID, 2), - parts: [textPart("prt_2", sessionID, "msg_a"), textPart("prt_1", sessionID, "msg_a")], - }) - - expect(draft.message[sessionID]?.map((x) => x.id)).toEqual(["msg_z", "msg_a"]) - expect(draft.part.msg_a?.map((x) => x.id)).toEqual(["prt_1", "prt_2"]) - }) - - test("applyOptimisticRemove removes message and part entries", () => { - const sessionID = "ses_1" - const draft = { - message: { [sessionID]: [userMessage("msg_1", sessionID), userMessage("msg_2", sessionID)] }, - part: { - msg_1: [textPart("prt_1", sessionID, "msg_1")], - msg_2: [textPart("prt_2", sessionID, "msg_2")], - } as Record, - } - - applyOptimisticRemove(draft, { sessionID, messageID: "msg_1" }) - - expect(draft.message[sessionID]?.map((x) => x.id)).toEqual(["msg_2"]) - expect(draft.part.msg_1).toBeUndefined() - expect(draft.part.msg_2).toHaveLength(1) - }) - - test("mergeOptimisticPage keeps pending messages in fetched timelines", () => { - const sessionID = "ses_1" - const page = mergeOptimisticPage( - { - session: [userMessage("msg_z", sessionID, 1)], - part: [{ id: "msg_z", part: [textPart("prt_1", sessionID, "msg_z")] }], - complete: true, - }, - [{ message: userMessage("msg_a", sessionID, 2), parts: [textPart("prt_2", sessionID, "msg_a")] }], - ) - - expect(page.session.map((x) => x.id)).toEqual(["msg_z", "msg_a"]) - expect(page.part.find((x) => x.id === "msg_a")?.part.map((x) => x.id)).toEqual(["prt_2"]) - expect(page.confirmed).toEqual([]) - expect(page.complete).toBe(true) - }) - - test("mergeOptimisticPage uses IDs only to break equal-time ties", () => { - const sessionID = "ses_1" - const page = mergeOptimisticPage( - { - session: [userMessage("msg_z", sessionID, 1)], - part: [], - complete: true, - }, - [{ message: userMessage("msg_a", sessionID, 1), parts: [] }], - ) - - expect(page.session.map((message) => message.id)).toEqual(["msg_a", "msg_z"]) - }) - - test("mergeOptimisticPage keeps missing optimistic parts until the server has them", () => { - const sessionID = "ses_1" - const page = mergeOptimisticPage( - { - session: [userMessage("msg_2", sessionID)], - part: [{ id: "msg_2", part: [textPart("prt_2", sessionID, "msg_2")] }], - complete: true, - }, - [ - { - message: userMessage("msg_2", sessionID), - parts: [textPart("prt_1", sessionID, "msg_2"), textPart("prt_2", sessionID, "msg_2")], - }, - ], - ) - - expect(page.part.find((x) => x.id === "msg_2")?.part.map((x) => x.id)).toEqual(["prt_1", "prt_2"]) - expect(page.confirmed).toEqual([]) - }) - - test("mergeOptimisticPage confirms echoed messages once all parts arrive", () => { - const sessionID = "ses_1" - const page = mergeOptimisticPage( - { - session: [userMessage("msg_2", sessionID)], - part: [ - { - id: "msg_2", - part: [{ ...textPart("prt_1", sessionID, "msg_2"), text: "server" }, textPart("prt_2", sessionID, "msg_2")], - }, - ], - complete: true, - }, - [ - { - message: userMessage("msg_2", sessionID), - parts: [textPart("prt_1", sessionID, "msg_2"), textPart("prt_2", sessionID, "msg_2")], - }, - ], - ) - - expect(page.confirmed).toEqual(["msg_2"]) - expect(page.part.find((x) => x.id === "msg_2")?.part).toMatchObject([ - { id: "prt_1", type: "text", text: "server" }, - { id: "prt_2", type: "text", text: "prt_2" }, - ]) - }) -}) diff --git a/packages/app/src/context/sync.tsx b/packages/app/src/context/sync.tsx index d8adf6a14b6..813c5698fd5 100644 --- a/packages/app/src/context/sync.tsx +++ b/packages/app/src/context/sync.tsx @@ -1,114 +1,6 @@ -import { Binary } from "@opencode-ai/core/util/binary" import { createMemo } from "solid-js" import { useServerSync } from "./server-sync" import { useSDK } from "./sdk" -import type { Message, Part } from "@/types" -import { messageKey } from "@/utils/session-message" - -const SKIP_PARTS = new Set(["patch", "step-start", "step-finish"]) - -function sortParts(parts: Part[]) { - return parts.filter((part) => !!part?.id).sort((a, b) => cmp(a.id, b.id)) -} - -const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0) - -type OptimisticStore = { - message: Record - part: Record -} - -type OptimisticAddInput = { - sessionID: string - message: Message - parts: Part[] -} - -type OptimisticRemoveInput = { - sessionID: string - messageID: string -} - -type OptimisticItem = { - message: Message - parts: Part[] -} - -type MessagePage = { - session: Message[] - part: { id: string; part: Part[] }[] - cursor?: string - complete: boolean -} - -const hasParts = (parts: Part[] | undefined, want: Part[]) => { - if (!parts) return want.length === 0 - return want.every((part) => Binary.search(parts, part.id, (item) => item.id).found) -} - -const mergeParts = (parts: Part[] | undefined, want: Part[]) => { - if (!parts) return sortParts(want) - const next = [...parts] - let changed = false - for (const part of want) { - const result = Binary.search(next, part.id, (item) => item.id) - if (result.found) continue - next.splice(result.index, 0, part) - changed = true - } - if (!changed) return parts - return next -} - -export function mergeOptimisticPage(page: MessagePage, items: OptimisticItem[]) { - if (items.length === 0) return { ...page, confirmed: [] as string[] } - - const session = [...page.session] - const part = new Map(page.part.map((item) => [item.id, sortParts(item.part)])) - const confirmed: string[] = [] - - for (const item of items) { - const result = Binary.search(session, messageKey(item.message), messageKey) - const found = result.found - if (!found) session.splice(result.index, 0, item.message) - - const current = part.get(item.message.id) - if (found && hasParts(current, item.parts)) { - confirmed.push(item.message.id) - continue - } - - part.set(item.message.id, mergeParts(current, item.parts)) - } - - return { - cursor: page.cursor, - complete: page.complete, - session, - part: [...part.entries()].sort((a, b) => cmp(a[0], b[0])).map(([id, part]) => ({ id, part })), - confirmed, - } -} - -export function applyOptimisticAdd(draft: OptimisticStore, input: OptimisticAddInput) { - const messages = draft.message[input.sessionID] - if (messages) { - const result = Binary.search(messages, messageKey(input.message), messageKey) - messages.splice(result.index, 0, input.message) - } else { - draft.message[input.sessionID] = [input.message] - } - draft.part[input.message.id] = sortParts(input.parts) -} - -export function applyOptimisticRemove(draft: OptimisticStore, input: OptimisticRemoveInput) { - const messages = draft.message[input.sessionID] - if (messages) { - const index = messages.findIndex((message) => message.id === input.messageID) - if (index >= 0) messages.splice(index, 1) - } - delete draft.part[input.messageID] -} export const useSync = () => { const serverSync = useServerSync()