Compare commits

...

1 Commits

Author SHA1 Message Date
Kit Langton d4bd41c4ba feat(tui): optimistic prompt submission
Echo user prompts into the local data store with a client-generated
message ID before server admission. The editor clears and the message
renders immediately; admission runs in a per-session background queue
that preserves steer order. The session.inbox.enqueued echo replaces
the local copy with server truth, wholesale sync replaces re-append
unconfirmed echoes, and failures roll back the echo, restore the
captured prompt, and toast.
2026-08-15 14:32:38 -04:00
3 changed files with 552 additions and 89 deletions
+136 -68
View File
@@ -57,7 +57,10 @@ import {
type LocalAttachment,
} from "./local-attachment"
import { useData } from "../../context/data"
import { usePromptRef } from "../../context/prompt"
import { useLocation } from "../../context/location"
import type { PromptFileAttachment, PromptSkillAttachment, SkillInfo } from "@opencode-ai/client"
import { SessionMessage } from "@opencode-ai/schema/session-message"
import { Keymap, type KeymapCommand } from "../../context/keymap"
import { abbreviateHome } from "../../runtime"
import { Slot } from "../../plugin/render"
@@ -100,6 +103,59 @@ export type PromptRef = {
const DRAFT_RETENTION_MIN_CHARS = 20
// Serialize background prompt submissions per session so admission order matches
// the on-screen optimistic order, even across Prompt remounts and route changes.
const submitTails = new Map<string, Promise<void>>()
function enqueueSubmit(sessionID: string, task: () => Promise<void>) {
const tail = (submitTails.get(sessionID) ?? Promise.resolve()).then(task, task)
submitTails.set(sessionID, tail)
void tail.finally(() => {
if (submitTails.get(sessionID) === tail) submitTails.delete(sessionID)
})
return tail
}
// Approximate the server's materialized attachment shape for the local echo. Pasted
// data: URIs carry real content so images preview immediately; file references render
// as labels until the admission echo replaces them with server truth.
function optimisticFiles(files: PromptInfo["files"]): PromptFileAttachment[] | undefined {
if (!files?.length) return undefined
return files.map((file) => {
const match = /^data:([^;,]*);base64,(.*)$/.exec(file.uri)
if (match)
return {
data: match[2] ?? "",
mime: match[1] || "application/octet-stream",
source: { type: "inline" as const },
name: file.name,
description: file.description,
mention: file.mention,
}
return {
data: "",
mime: file.uri.endsWith("/") ? "application/x-directory" : "text/plain",
source: { type: "uri" as const, uri: file.uri },
name: file.name,
description: file.description,
mention: file.mention,
}
})
}
function optimisticSkills(
skills: PromptInfo["skills"],
available: SkillInfo[],
): PromptSkillAttachment[] | undefined {
if (!skills?.length) return undefined
return skills.map((attachment) => ({
id: attachment.id,
name: available.find((skill) => skill.id === attachment.id)?.name ?? attachment.id,
text: "",
mention: attachment.mention,
}))
}
function randomIndex(count: number) {
if (count <= 0) return 0
return Math.floor(Math.random() * count)
@@ -201,6 +257,7 @@ export function Prompt(props: PromptProps) {
const editor = useEditorContext()
const route = useRoute()
const data = useData()
const activePrompt = usePromptRef()
const directoryRecents = useDirectoryRecents()
const keymapCommands = Keymap.useCommands()
const currentLocation = useLocation()
@@ -1070,6 +1127,20 @@ export function Prompt(props: PromptProps) {
}
})
// Return a failed background submission to its author. When the target tab's
// editor is live and empty, restore directly; otherwise stash a draft for the
// next mount. A non-empty editor is never clobbered—prompt history retains the
// failed prompt either way.
function restorePrompt(sessionID: string, snapshot: PromptInfo) {
const active = route.data
if (active.type === "session" && active.sessionID === sessionID) {
const live = activePrompt.current
if (live && !live.current.text.trim()) live.set(snapshot)
return
}
saveDraft(sessionID, { prompt: snapshot, cursor: snapshot.text.length })
}
let submitting = false
async function submit(delivery: SessionInbox.Delivery = "steer") {
// Prevent overlapping invocations (e.g. a double-pressed Enter, or the
@@ -1170,7 +1241,6 @@ export function Prompt(props: PromptProps) {
const variant = selection.variant
let sessionID = props.sessionID
let session = sessionID ? data.session.get(sessionID) : undefined
let finishMoveProgress = false
if (sessionID == null) {
const directory = await move.getDirectory()
@@ -1204,7 +1274,6 @@ export function Prompt(props: PromptProps) {
}
sessionID = created.id
session = created
}
// Capture mode before it gets reset
@@ -1245,70 +1314,70 @@ export function Prompt(props: PromptProps) {
})
} else {
move.startSubmit()
if (!session) {
await data.session.sync(sessionID)
session = data.session.get(sessionID)
}
if (session?.agent !== agent.id) {
await client.api.session.switchAgent({ sessionID, agent: agent.id })
}
if (
session?.model?.providerID !== selection.providerID ||
session.model.id !== selection.modelID ||
(session.model.variant ?? "default") !== (variant ?? "default")
) {
const model = { providerID: selection.providerID, id: selection.modelID, variant }
const cancelCommit = local.model.trackSessionCommit(sessionID, model)
await client.api.session.switchModel({ sessionID, model }).catch((error) => {
cancelCommit()
throw error
})
}
if (session?.revert) {
const error = await client.api.session.revert.commit({ sessionID }).then(
() => undefined,
(error) => error,
)
if (error) {
toast.show({ title: "Failed to commit revert", message: errorMessage(error), variant: "error" })
return false
}
}
if (pendingEditorSelection) {
// Keep editor context hidden while admitting it before the corresponding user prompt.
const error = await client.api.session
.synthetic({
sessionID,
text: formatEditorContext(pendingEditorSelection),
resume: false,
// Echo the prompt locally and admit it in the background: the editor clears and
// the message renders immediately, while a per-session queue preserves admission
// order. Failure rolls the echo back and restores the captured prompt.
const submitSessionID = sessionID
const messageID = SessionMessage.ID.create()
const snapshot = structuredClone(unwrap(store.prompt))
const editorContextText = pendingEditorSelection ? formatEditorContext(pendingEditorSelection) : undefined
const targetAgent = agent.id
const model = { providerID: selection.providerID, id: selection.modelID, variant }
data.session.optimistic.prompt({
sessionID: submitSessionID,
messageID,
delivery,
text: inputText,
files: optimisticFiles(snapshot.files),
agents: snapshot.agents?.length ? snapshot.agents : undefined,
skills: optimisticSkills(snapshot.skills, data.location.skill.list(currentLocation.ref) ?? []),
})
// Mark the editor context sent with the echo so a rapid follow-up submit
// does not re-attach the same selection while admission is in flight.
if (editorContextText) editor.markSelectionSent()
void enqueueSubmit(submitSessionID, async () => {
const error = await (async () => {
let session = data.session.get(submitSessionID)
if (!session) {
await data.session.sync(submitSessionID)
session = data.session.get(submitSessionID)
}
if (session?.agent !== targetAgent) {
await client.api.session.switchAgent({ sessionID: submitSessionID, agent: targetAgent })
}
if (
session?.model?.providerID !== model.providerID ||
session.model.id !== model.id ||
(session.model.variant ?? "default") !== (model.variant ?? "default")
) {
const cancelCommit = local.model.trackSessionCommit(submitSessionID, model)
await client.api.session.switchModel({ sessionID: submitSessionID, model }).catch((error) => {
cancelCommit()
throw error
})
}
if (session?.revert) await client.api.session.revert.commit({ sessionID: submitSessionID })
// Keep editor context hidden while admitting it before the corresponding user prompt.
if (editorContextText)
await client.api.session.synthetic({ sessionID: submitSessionID, text: editorContextText, resume: false })
await client.api.session.prompt({
sessionID: submitSessionID,
id: messageID,
text: inputText,
files: snapshot.files,
agents: snapshot.agents,
skills: snapshot.skills?.length ? snapshot.skills : undefined,
delivery,
})
.then(
() => undefined,
(error) => error,
)
if (error) {
toast.show({ title: "Failed to send editor context", message: errorMessage(error), variant: "error" })
return false
}
}
const error = await client.api.session
.prompt({
sessionID,
text: inputText,
files: store.prompt.files,
agents: store.prompt.agents,
skills: store.prompt.skills?.length ? store.prompt.skills : undefined,
delivery,
})
.then(
})().then(
() => undefined,
(error) => error,
)
if (error) {
if (error === undefined) return
data.session.optimistic.rollback(submitSessionID, messageID)
restorePrompt(submitSessionID, snapshot)
toast.show({ title: "Failed to send prompt", message: errorMessage(error), variant: "error" })
return false
}
if (pendingEditorSelection) editor.markSelectionSent()
})
}
history.append({
...store.prompt,
@@ -1319,15 +1388,14 @@ export function Prompt(props: PromptProps) {
setStore("extmarkToPart", new Map())
props.onSubmit?.()
// temporary hack to make sure the message is sent
if (!props.sessionID) {
if (pendingEditorSelection) editor.preserveSelectionFromNewSession()
setTimeout(() => {
route.navigate({
type: "session",
sessionID,
})
}, 50)
// The optimistic echo is already in the data store, so the session route
// renders the prompt immediately; admission continues in the background.
route.navigate({
type: "session",
sessionID,
})
}
input.clear()
if (finishMoveProgress) move.finishSubmit()
+140 -21
View File
@@ -20,6 +20,7 @@ import type {
ProviderInfo,
ReferenceInfo,
SessionMessageInfo,
SessionMessageUser,
SessionMessageAssistant,
SessionMessageAssistantReasoning,
SessionMessageAssistantText,
@@ -40,7 +41,7 @@ import { useClient } from "./client"
import { nonEmptyToolContent } from "../util/tool-display"
import type { SessionInbox } from "@opencode-ai/schema/session-inbox"
import { Worktree } from "@opencode-ai/schema/worktree"
import { createEffect, createSignal, onCleanup } from "solid-js"
import { batch, createEffect, createSignal, onCleanup } from "solid-js"
export type DataSessionStatus = "idle" | "running"
@@ -161,6 +162,25 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
const messageIndex = new Map<string, Map<string, number>>()
const sync = createSync()
// Optimistic prompt echoes: user prompts applied locally before server admission,
// keyed sessionID -> messageID. The session.inbox.enqueued echo carrying the same
// ID replaces the local copy with server truth (admission materializes file
// attachments and expands skills). Entries survive wholesale sync replaces until
// the server confirms them or the submitter rolls back.
type OptimisticPrompt = { message: SessionMessageInfo; pending: SessionInboxInfo }
const optimisticPrompts = new Map<string, Map<string, OptimisticPrompt>>()
function confirmOptimistic(sessionID: string, messageID: string) {
const entries = optimisticPrompts.get(sessionID)
if (!entries?.delete(messageID)) return false
if (entries.size === 0) optimisticPrompts.delete(sessionID)
return true
}
function optimisticMessages(sessionID: string) {
return [...(optimisticPrompts.get(sessionID)?.values() ?? [])]
}
function setSessionActive(sessionID: string, status: DataSessionStatus) {
setStore("session", "active", sessionID, status)
}
@@ -188,6 +208,17 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
)
}
function removeMessage(sessionID: string, messageID: string) {
if (!messageIndex.get(sessionID)?.has(messageID)) return
message.update(sessionID, (draft, index) => {
const position = index.get(messageID)
if (position === undefined) return
draft.splice(position, 1)
index.delete(messageID)
message.reindex(draft, index, position)
})
}
function removePermission(sessionID: string, requestID: string) {
const requests = store.session.permission[sessionID]
if (!requests?.some((request) => request.id === requestID)) return
@@ -313,6 +344,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
function removeSession(sessionID: string) {
messageIndex.delete(sessionID)
optimisticPrompts.delete(sessionID)
sync.invalidate(`session:${sessionID}`)
sync.invalidate(`session.pending:${sessionID}`)
sync.invalidate(`session.message:${sessionID}`)
@@ -471,6 +503,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
break
}
case "session.inbox.delivered": {
// Delivery implies the message is projected server-side, so future message
// fetches include it and the optimistic entry no longer needs re-appending.
confirmOptimistic(event.data.sessionID, event.data.inboxID)
const admitted = store.session.input[event.data.sessionID]?.includes(event.data.inboxID) ?? false
removePending(event.data.sessionID, event.data.inboxID)
message.update(event.data.sessionID, (draft, index) => {
@@ -489,25 +524,31 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
updatePending(event.data.sessionID, event.data.inboxID, event.data.delivery)
break
case "session.inbox.cancelled": {
confirmOptimistic(event.data.sessionID, event.data.inboxID)
removePending(event.data.sessionID, event.data.inboxID)
if (messageIndex.get(event.data.sessionID)?.has(event.data.inboxID))
message.update(event.data.sessionID, (draft, index) => {
const position = index.get(event.data.inboxID)
if (position === undefined) return
draft.splice(position, 1)
index.delete(event.data.inboxID)
message.reindex(draft, index, position)
})
removeMessage(event.data.sessionID, event.data.inboxID)
break
}
case "session.inbox.enqueued": {
const item = event.data.item
addPending({
// The admission echo is authoritative for an optimistic local copy: it
// materializes file attachments and expands skills, so replace in place.
const confirmed = confirmOptimistic(event.data.sessionID, event.data.inboxID)
const pendingItem: SessionInboxInfo = {
id: event.data.inboxID,
sessionID: event.data.sessionID,
timeCreated: event.created,
...item,
})
}
const pendingList = store.session.pending[event.data.sessionID]
if (confirmed && pendingList?.some((pending) => pending.id === event.data.inboxID))
setStore(
"session",
"pending",
event.data.sessionID,
pendingList.map((pending) => (pending.id === event.data.inboxID ? pendingItem : pending)),
)
else addPending(pendingItem)
if (!store.session.input[event.data.sessionID]?.includes(event.data.inboxID))
setStore("session", "input", event.data.sessionID, [
...(store.session.input[event.data.sessionID] ?? []),
@@ -515,9 +556,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
])
if (item.type !== "user" && item.type !== "synthetic") break
message.update(event.data.sessionID, (draft, index) => {
message.append(
draft,
index,
const next: SessionMessageInfo =
item.type === "user"
? {
id: event.data.inboxID,
@@ -530,8 +569,13 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
type: "synthetic",
...item.payload,
time: { created: event.created },
},
)
}
const position = index.get(event.data.inboxID)
if (confirmed && position !== undefined) {
draft[position] = next
return
}
message.append(draft, index, next)
})
break
}
@@ -825,22 +869,30 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
if (store.session.info[event.data.sessionID])
setStore("session", "info", event.data.sessionID, "revert", undefined)
break
case "session.revert.committed":
case "session.revert.committed": {
if (store.session.info[event.data.sessionID]) {
setStore("session", "info", event.data.sessionID, "revert", undefined)
}
// Unconfirmed optimistic prompts postdate the revert boundary but were never
// part of the reverted history: the server admits them after the commit.
const local = optimisticPrompts.get(event.data.sessionID)
setStore(
"session",
"input",
event.data.sessionID,
(store.session.input[event.data.sessionID] ?? []).filter((id) => id < event.data.to),
(store.session.input[event.data.sessionID] ?? []).filter(
(id) => id < event.data.to || local?.has(id) === true,
),
)
message.update(event.data.sessionID, (draft, index) => {
const position = draft.findIndex((item) => item.id >= event.data.to)
if (position === -1) return
for (const item of draft.splice(position)) index.delete(item.id)
const dropped = draft.splice(position)
for (const item of dropped) index.delete(item.id)
for (const item of dropped) if (local?.has(item.id)) message.append(draft, index, item)
})
break
}
case "session.compaction.delta":
message.update(event.data.sessionID, (draft) => {
const current = message.compaction(draft)
@@ -1013,13 +1065,74 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
return store.session.input[sessionID]?.includes(inboxID) ?? false
},
},
optimistic: {
// Locally echo a user prompt before server admission. The session.inbox.enqueued
// echo carrying the same message ID replaces the copy with server truth; rollback
// removes the echo when submission fails.
prompt(input: {
sessionID: string
messageID: string
delivery: SessionInbox.Delivery
text: string
files?: SessionMessageUser["files"]
agents?: SessionMessageUser["agents"]
skills?: SessionMessageUser["skills"]
}) {
const created = Date.now()
const pendingItem: SessionInboxInfo = {
id: input.messageID,
sessionID: input.sessionID,
timeCreated: created,
type: "user",
payload: { text: input.text, files: input.files, agents: input.agents, skills: input.skills },
delivery: input.delivery,
}
const messageItem: SessionMessageInfo = {
id: input.messageID,
type: "user",
text: input.text,
files: input.files,
agents: input.agents,
skills: input.skills,
time: { created },
}
const entries = optimisticPrompts.get(input.sessionID) ?? new Map<string, OptimisticPrompt>()
optimisticPrompts.set(input.sessionID, entries)
entries.set(input.messageID, { message: messageItem, pending: pendingItem })
batch(() => {
addPending(pendingItem)
if (!store.session.input[input.sessionID]?.includes(input.messageID))
setStore("session", "input", input.sessionID, [
...(store.session.input[input.sessionID] ?? []),
input.messageID,
])
message.update(input.sessionID, (draft, index) => message.append(draft, index, messageItem))
})
},
rollback(sessionID: string, messageID: string) {
if (!confirmOptimistic(sessionID, messageID)) return
batch(() => {
removePending(sessionID, messageID)
removeMessage(sessionID, messageID)
})
},
},
pending: {
list(sessionID: string) {
return store.session.pending[sessionID] ?? []
},
sync(sessionID: string) {
return sync.run(`session.pending:${sessionID}`, async () => {
const pending = await client.api.session.inbox.list({ sessionID })
const fetched = await client.api.session.inbox.list({ sessionID })
// Keep unconfirmed optimistic prompts pending across the wholesale replace.
// Server presence here is not treated as confirmation: until the enqueued
// echo or a projected message arrives, the ledger still guards message.sync.
const pending = [
...fetched,
...optimisticMessages(sessionID)
.map((entry) => entry.pending)
.filter((entry) => !fetched.some((item) => item.id === entry.id)),
]
setStore("session", "pending", sessionID, reconcile(pending))
setStore(
"session",
@@ -1069,9 +1182,15 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
},
sync(sessionID: string) {
return sync.run(`session.message:${sessionID}`, async () => {
const messages = (
const fetched = (
await client.api.message.list({ sessionID, limit: 200, order: "desc" })
).data.toReversed()
// A wholesale replace would drop optimistic prompts the server has not
// admitted yet. A fetched ID is server confirmation; the rest re-append.
for (const entry of optimisticMessages(sessionID))
if (fetched.some((item) => item.id === entry.message.id))
confirmOptimistic(sessionID, entry.message.id)
const messages = [...fetched, ...optimisticMessages(sessionID).map((entry) => entry.message)]
messageIndex.set(sessionID, new Map(messages.map((message, index) => [message.id, index])))
setStore("session", "message", sessionID, reconcile(messages))
})
@@ -0,0 +1,276 @@
/** @jsxImportSource @opentui/solid */
import { expect, test } from "bun:test"
import { testRender } from "@opentui/solid"
import type { OpenCodeEvent, SessionMessageInfo } from "@opencode-ai/client"
import { createEffect, type ParentProps } from "solid-js"
import { ConfigProvider } from "../../../src/config"
import { ClientProvider, useClient } from "../../../src/context/client"
import { DataProvider as DataProviderBase, useData } from "../../../src/context/data"
import { LocationProvider, useLocation } from "../../../src/context/location"
import { createApi, createEventStream, createFetch, directory, json } from "../../fixture/tui-client"
import { TestTuiContexts } from "../../fixture/tui-environment"
import { createTuiResolvedConfig } from "../../fixture/tui-runtime"
async function wait(fn: () => boolean, timeout = 2000) {
const start = Date.now()
while (!fn()) {
if (Date.now() - start > timeout) throw new Error("timed out waiting for condition")
await Bun.sleep(10)
}
}
function emitEvent(events: ReturnType<typeof createEventStream>, event: OpenCodeEvent) {
events.emit({ ...event, location: { directory } })
}
const config = createTuiResolvedConfig()
function DataProvider(props: ParentProps) {
return (
<ConfigProvider config={config}>
<DataProviderBase>
<LocationProvider>
<SyncLocation />
{props.children}
</LocationProvider>
</DataProviderBase>
</ConfigProvider>
)
}
function SyncLocation() {
const data = useData()
const location = useLocation()
createEffect(() => location.set(data.location.default()))
return null
}
function durable(sessionID: string, seq = 0): { aggregateID: string; seq: number; version: 1 } {
return { aggregateID: sessionID, seq, version: 1 }
}
type Harness = {
data: ReturnType<typeof useData>
client: ReturnType<typeof useClient>
}
async function renderData(fetch: ReturnType<typeof createFetch>["fetch"]) {
const harness = {} as Harness
function Probe() {
harness.client = useClient()
harness.data = useData()
return <box />
}
const app = await testRender(() => (
<TestTuiContexts>
<ClientProvider api={createApi(fetch)}>
<DataProvider>
<Probe />
</DataProvider>
</ClientProvider>
</TestTuiContexts>
))
await wait(() => harness.client.connection.status() === "connected")
return { app, ...harness }
}
test("echoes an optimistic prompt and replaces it with the admission echo", async () => {
const events = createEventStream()
const sessionID = "session-optimistic-echo"
const calls = createFetch(undefined, events)
const { app, data } = await renderData(calls.fetch)
try {
data.session.optimistic.prompt({
sessionID,
messageID: "msg_optimistic",
delivery: "steer",
text: "Hello",
files: [{ data: "", mime: "text/plain", source: { type: "uri", uri: "file:///tmp/a.ts" }, name: "a.ts" }],
})
const echoed = data.session.message.get(sessionID, "msg_optimistic")
expect(echoed?.type === "user" && echoed.text).toBe("Hello")
expect(echoed?.type === "user" && echoed.files?.[0]?.data).toBe("")
expect(data.session.pending.list(sessionID).map((item) => item.id)).toEqual(["msg_optimistic"])
expect(data.session.input.has(sessionID, "msg_optimistic")).toBe(true)
// Admission materializes attachments, so the echo must be replaced in place.
emitEvent(events, {
id: "evt_admitted",
created: 9,
type: "session.inbox.enqueued",
durable: durable(sessionID),
data: {
sessionID,
inboxID: "msg_optimistic",
item: {
type: "user",
payload: {
text: "Hello",
files: [{ data: "QUJD", mime: "text/plain", source: { type: "uri", uri: "file:///tmp/a.ts" }, name: "a.ts" }],
},
delivery: "steer",
},
},
})
await wait(() => {
const message = data.session.message.get(sessionID, "msg_optimistic")
return message?.type === "user" && message.files?.[0]?.data === "QUJD"
})
expect(data.session.message.list(sessionID)).toHaveLength(1)
expect(data.session.pending.list(sessionID)).toHaveLength(1)
expect(data.session.pending.list(sessionID)[0]?.timeCreated).toBe(9)
} finally {
app.renderer.destroy()
}
})
test("rolls back a failed optimistic prompt", async () => {
const events = createEventStream()
const sessionID = "session-optimistic-rollback"
const calls = createFetch(undefined, events)
const { app, data } = await renderData(calls.fetch)
try {
data.session.optimistic.prompt({
sessionID,
messageID: "msg_failed",
delivery: "queue",
text: "Will fail",
})
expect(data.session.message.get(sessionID, "msg_failed")).toBeDefined()
data.session.optimistic.rollback(sessionID, "msg_failed")
expect(data.session.message.get(sessionID, "msg_failed")).toBeUndefined()
expect(data.session.pending.list(sessionID)).toHaveLength(0)
expect(data.session.input.has(sessionID, "msg_failed")).toBe(false)
// Rollback of an unknown or already-settled echo is a no-op.
data.session.optimistic.rollback(sessionID, "msg_failed")
expect(data.session.message.list(sessionID)).toHaveLength(0)
} finally {
app.renderer.destroy()
}
})
test("optimistic prompts survive sync replaces until the server confirms them", async () => {
const events = createEventStream()
const sessionID = "session-optimistic-sync"
let serverMessages: SessionMessageInfo[] = []
const calls = createFetch((url) => {
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: serverMessages, cursor: {} })
if (url.pathname === `/api/session/${sessionID}/inbox`) return json({ data: [] })
}, events)
const { app, data } = await renderData(calls.fetch)
try {
data.session.optimistic.prompt({
sessionID,
messageID: "msg_pending",
delivery: "steer",
text: "Survive the sync",
})
// A wholesale replace from an empty server page keeps the unconfirmed echo.
await data.session.message.sync(sessionID)
expect(data.session.message.get(sessionID, "msg_pending")).toBeDefined()
await data.session.pending.sync(sessionID)
expect(data.session.pending.list(sessionID).map((item) => item.id)).toEqual(["msg_pending"])
expect(data.session.input.has(sessionID, "msg_pending")).toBe(true)
// A fetched page containing the ID is server confirmation: the projected copy
// wins and later rollback attempts become no-ops.
serverMessages = [{ id: "msg_pending", type: "user", text: "Survive the sync", time: { created: 5 } }]
data.session.message.invalidate(sessionID)
await data.session.message.sync(sessionID)
const confirmed = data.session.message.get(sessionID, "msg_pending")
expect(confirmed?.type === "user" && confirmed.time.created).toBe(5)
expect(data.session.message.list(sessionID)).toHaveLength(1)
data.session.optimistic.rollback(sessionID, "msg_pending")
expect(data.session.message.get(sessionID, "msg_pending")).toBeDefined()
} finally {
app.renderer.destroy()
}
})
test("cancellation clears the optimistic echo for good", async () => {
const events = createEventStream()
const sessionID = "session-optimistic-cancel"
const calls = createFetch((url) => {
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} })
}, events)
const { app, data } = await renderData(calls.fetch)
try {
data.session.optimistic.prompt({
sessionID,
messageID: "msg_cancelled",
delivery: "queue",
text: "Cancel me",
})
emitEvent(events, {
id: "evt_cancelled",
created: 2,
type: "session.inbox.cancelled",
durable: durable(sessionID),
data: { sessionID, inboxID: "msg_cancelled" },
})
await wait(() => data.session.message.get(sessionID, "msg_cancelled") === undefined)
expect(data.session.pending.list(sessionID)).toHaveLength(0)
// The ledger entry is gone too: a sync replace must not resurrect the echo.
await data.session.message.sync(sessionID)
expect(data.session.message.get(sessionID, "msg_cancelled")).toBeUndefined()
} finally {
app.renderer.destroy()
}
})
test("revert commit preserves unconfirmed optimistic prompts", async () => {
const events = createEventStream()
const sessionID = "session-optimistic-revert"
const calls = createFetch(undefined, events)
const { app, data } = await renderData(calls.fetch)
try {
for (const [seq, id] of [
[0, "msg_1"],
[1, "msg_2"],
] as const) {
emitEvent(events, {
id: `evt_seed_${id}`,
created: seq + 1,
type: "session.inbox.enqueued",
durable: durable(sessionID, seq),
data: { sessionID, inboxID: id, item: { type: "user", payload: { text: id }, delivery: "steer" } },
})
}
await wait(() => data.session.message.list(sessionID).length === 2)
data.session.optimistic.prompt({
sessionID,
messageID: "msg_9",
delivery: "steer",
text: "After the revert boundary",
})
emitEvent(events, {
id: "evt_revert",
created: 4,
type: "session.revert.committed",
durable: durable(sessionID, 2),
data: { sessionID, to: "msg_2" },
})
await wait(() => data.session.message.get(sessionID, "msg_2") === undefined)
expect(data.session.message.get(sessionID, "msg_1")).toBeDefined()
expect(data.session.message.get(sessionID, "msg_9")).toBeDefined()
expect(data.session.input.has(sessionID, "msg_9")).toBe(true)
} finally {
app.renderer.destroy()
}
})