fix(tui): admit compaction requests optimistically (#45973)

Render queued compaction feedback before model setup, coalesce repeated gestures, and reconcile canonical admissions without restoring consumed rows. Serialize prompt and compaction preparation through the existing per-session admission chain.
This commit is contained in:
Kit Langton 2026-08-28 14:00:09 -04:00 committed by GitHub
parent cd3b12c579
commit 732f949a65
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 758 additions and 61 deletions

View file

@ -29,6 +29,7 @@ import type {
SessionMessageAssistantTool,
SessionInfo,
SessionInboxInfo,
SessionInboxCompaction,
ShellInfo,
SkillInfo,
VcsInfo,
@ -284,12 +285,11 @@ export function createData(config: CreateDataInput) {
setStore("session", "pending", sessionID, index, { ...item, delivery })
}
// Inbox IDs of optimistic prompt admissions still awaiting their durable
// echo. This is the one deliberate piece of in-flight bookkeeping in this
// layer: it exists so a rejection only rolls back rows the server never
// acknowledged, and so a concurrent pending re-fetch cannot wipe a row the
// server does not know about yet. Entries clear on the enqueued echo or on
// rollback — not on POST success, which typically precedes the echo.
// Inbox IDs of optimistic admissions awaiting acknowledgement, so rejection
// only rolls back unacknowledged rows and a pending re-fetch cannot wipe a
// row the server does not know about yet. Prompts clear on their durable
// echo, positive pending read, or rollback; compactions also reconcile the
// POST's canonical ID.
const outbox = new Set<string>()
// Session IDs of optimistic create admissions still awaiting acknowledgement
@ -303,11 +303,12 @@ export function createData(config: CreateDataInput) {
// to exist server-side instead of failing with "not found".
const creating = new Map<string, Promise<unknown>>()
// Per-session send chain: prompts must be admitted in submission order,
// and HTTP gives no ordering across concurrent POSTs. Each prompt waits
// for the previous prompt's POST (settled, so one failure does not block
// the next) before sending its own.
// Per-session send chain: prompts and compactions must be admitted in
// submission order. Each waits for the previous POST to settle, so one
// failure does not block the next.
const sending = new Map<string, Promise<unknown>>()
const compacting = new Map<string, { id: string; observed: Set<string>; request: Promise<SessionInboxCompaction> }>()
onCleanup(() => compacting.clear())
// Register `promise` under `key` until it settles. A later registration
// replaces an earlier one; settlement only clears its own entry.
@ -319,9 +320,24 @@ export function createData(config: CreateDataInput) {
void promise.then(settle, settle)
}
// Capture creation before settlement clears its entry, so dependent RPCs still see a failed create.
function sendAdmission<Value>(sessionID: string, send: () => Promise<Value>, gate?: Promise<unknown>) {
const created = creating.get(sessionID)
const previous = sending.get(sessionID)
const request = Promise.resolve()
.then(() => Promise.all([gate, created, previous]))
.then(send)
track(
sending,
sessionID,
request.catch(() => undefined),
)
return request
}
// Upsert an admitted inbox item into pending, input, and (for user and
// synthetic items) the visible transcript. Used by the inbox.enqueued
// handler and by optimistic prompt admission; the upsert is what reconciles
// handler and by optimistic admission; the upsert is what reconciles
// the durable echo with an optimistic placeholder — the durable payload and
// times replace the client's guess.
function admitLocal(item: SessionInboxInfo) {
@ -334,6 +350,7 @@ export function createData(config: CreateDataInput) {
item.sessionID,
at < 0 ? [...pending, item] : pending.map((entry, index) => (index === at ? item : entry)),
)
if (item.type === "compaction") return
const input = store.session.input[item.sessionID] ?? []
if (!input.includes(item.id)) setStore("session", "input", item.sessionID, [...input, item.id])
materializeInboxMessage(item)
@ -668,6 +685,7 @@ export function createData(config: CreateDataInput) {
draft.push(existing)
message.reindex(draft, index, position)
})
compacting.get(event.data.sessionID)?.observed.add(event.data.inboxID)
return
}
case "session.inbox.delivery.changed":
@ -675,6 +693,7 @@ export function createData(config: CreateDataInput) {
return
case "session.inbox.cancelled": {
retractLocal(event.data.sessionID, event.data.inboxID)
compacting.get(event.data.sessionID)?.observed.add(event.data.inboxID)
return
}
case "session.inbox.enqueued": {
@ -685,6 +704,12 @@ export function createData(config: CreateDataInput) {
timeCreated: event.created,
...event.data.item,
})
if (event.data.item.type === "compaction") {
const active = compacting.get(event.data.sessionID)
active?.observed.add(event.data.inboxID)
if (active && active.id !== event.data.inboxID && outbox.delete(active.id))
removePending(event.data.sessionID, active.id)
}
return
}
case "session.instructions.updated":
@ -983,6 +1008,7 @@ export function createData(config: CreateDataInput) {
time: { created: event.created },
})
})
if (event.data.inputID) compacting.get(event.data.sessionID)?.observed.add(event.data.inputID)
return
case "session.execution.succeeded":
case "session.execution.failed":
@ -1080,6 +1106,7 @@ export function createData(config: CreateDataInput) {
}
message.append(draft, index, failed)
})
if (event.data.inputID) compacting.get(event.data.sessionID)?.observed.add(event.data.inputID)
return
case "permission.asked":
if (store.session.permission[event.data.sessionID]?.some((request) => request.id === event.data.id)) return
@ -1266,12 +1293,17 @@ export function createData(config: CreateDataInput) {
sync(sessionID: string) {
return sync.run(`session.pending:${sessionID}`, async () => {
const pending = await api().session.inbox.list({ sessionID })
// A positive read acknowledges admission even when its SSE echo is delayed.
pending.forEach((item) => outbox.delete(item.id))
// Compactions also coalesce by Session, not just by the proposed ID.
if (pending.some((item) => item.type === "compaction"))
store.session.pending[sessionID]
?.filter((item) => item.type === "compaction")
.forEach((item) => outbox.delete(item.id))
// Keep optimistic rows still awaiting their echo: this fetch may
// have raced ahead of an in-flight admission the server does not
// know about yet.
const inflight = (store.session.pending[sessionID] ?? []).filter(
(item) => outbox.has(item.id) && !pending.some((row) => row.id === item.id),
)
const inflight = (store.session.pending[sessionID] ?? []).filter((item) => outbox.has(item.id))
const merged = inflight.length === 0 ? pending : [...pending, ...inflight]
batch(() => {
setStore("session", "pending", sessionID, reconcile(merged))
@ -1345,13 +1377,56 @@ export function createData(config: CreateDataInput) {
if (fresh) track(creating, id, request)
return { id, request }
},
compact(input: { sessionID: string; model?: ModelRef }) {
const active = compacting.get(input.sessionID)
if (active) return active.request
// A known pending control ID may be consumed while setup waits. Propose
// a fresh ID and let the server coalesce, without duplicating its row.
const id = SessionMessage.ID.create()
if (!store.session.pending[input.sessionID]?.some((item) => item.type === "compaction")) {
outbox.add(id)
admitLocal({
id,
sessionID: input.sessionID,
timeCreated: Date.now(),
type: "compaction",
delivery: "steer",
payload: {},
})
}
// Compaction admission can coalesce onto a different ID. Retire the
// speculative row on an echo, and remember consumed IDs until the POST
// settles so its older response cannot resurrect a queued row.
const observed = new Set<string>()
const request = sendAdmission(input.sessionID, async () => {
if (input.model) await api().session.switchModel({ sessionID: input.sessionID, model: input.model })
return api().session.compact({ sessionID: input.sessionID, id })
})
.then((item) => {
batch(() => {
outbox.delete(id)
if (item.id !== id) removePending(input.sessionID, id)
if (!observed.has(item.id) && !messageIndex.get(input.sessionID)?.has(item.id)) admitLocal(item)
})
return item
})
.catch((error) => {
if (outbox.delete(id)) removePending(input.sessionID, id)
throw error
})
.finally(() => {
if (compacting.get(input.sessionID)?.request === request) compacting.delete(input.sessionID)
})
compacting.set(input.sessionID, { id, observed, request })
return request
},
// Optimistic prompt admission: render the prompt immediately under a
// client-minted ID, send it, and let the durable inbox.enqueued echo
// upsert that same ID with the server's payload. Server admission is
// idempotent per ID, so retrying with the identical payload cannot
// double-admit.
prompt(input: SessionPromptInput & { gate?: Promise<unknown> }) {
const { gate, ...request } = input
prompt(input: SessionPromptInput & { gate?: Promise<unknown>; prepare?: () => Promise<unknown> }) {
const { gate, prepare, ...request } = input
const id = request.id ?? SessionMessage.ID.create()
// A retry may reuse an ID that is already rendered — and possibly
// already durable. Admit optimistically only for new IDs so a failed
@ -1377,25 +1452,15 @@ export function createData(config: CreateDataInput) {
},
})
}
// Wrapped so even a synchronous client failure reaches the rollback.
// The POST additionally waits for the caller's gate, for any
// in-flight optimistic create of this session, and for the previous
// prompt's POST: the row renders now, the send happens once the
// session exists server-side and earlier prompts are admitted.
const previous = sending.get(request.sessionID)
const send = Promise.resolve()
.then(() => Promise.all([gate, creating.get(request.sessionID), previous]))
.then(() => api().session.prompt({ ...request, id }))
track(
sending,
return sendAdmission(
request.sessionID,
send.then(
() => undefined,
() => undefined,
),
)
return send.catch((error) => {
// Roll back only rows this call admitted and the echo has not
async () => {
await prepare?.()
return api().session.prompt({ ...request, id })
},
gate,
).catch((error) => {
// Roll back only rows this call admitted and the server has not
// acknowledged: anything else is server state.
if (fresh && outbox.delete(id)) retractLocal(request.sessionID, id)
throw error

View file

@ -0,0 +1,400 @@
import { expect, test } from "bun:test"
import { createRoot } from "solid-js"
import { createData, type CreateDataInput } from "../src/solid"
import { OpenCode, type OpenCodeEvent, type SessionInboxCompaction, type SessionInboxInfo } from "../src/promise"
test("admits compaction before model setup and serializes the following prompt", async () => {
using fixture = setup()
const compact = fixture.data.session.compact({ sessionID, model: { providerID: "demo", id: "model" } })
const proposed = fixture.data.session.pending.list(sessionID)[0]
expect(proposed).toMatchObject({ type: "compaction", sessionID })
expect(fixture.calls).toEqual([])
expect(fixture.data.session.message.list(sessionID)).toEqual([])
expect(fixture.data.session.status(sessionID)).toBe("idle")
const prompt = fixture.data.session.prompt({ sessionID, text: "Follow up" })
expect(fixture.data.session.message.list(sessionID)).toMatchObject([{ type: "user", text: "Follow up" }])
await wait(() => fixture.calls.length === 1)
expect(fixture.calls).toEqual(["model"])
fixture.model.resolve()
await wait(() => fixture.calls.length === 2)
expect(fixture.calls).toEqual(["model", "compact"])
fixture.response.resolve(Response.json({ data: item(proposed.id) }))
await Promise.all([compact, prompt])
expect(fixture.calls).toEqual(["model", "compact", "prompt"])
expect(fixture.proposals).toEqual([proposed.id])
})
test("coalesces duplicate gestures until the admission request settles", async () => {
using fixture = setup()
const first = fixture.data.session.compact({ sessionID })
expect(fixture.data.session.compact({ sessionID })).toBe(first)
expect(fixture.data.session.pending.list(sessionID)).toHaveLength(1)
await wait(() => fixture.calls.length === 1)
fixture.response.resolve(Response.json({ data: item("msg_canonical") }))
await first
expect(fixture.calls).toEqual(["compact"])
const next = fixture.data.session.compact({ sessionID })
expect(next).not.toBe(first)
await next
expect(fixture.calls).toEqual(["compact", "compact"])
expect(new Set(fixture.proposals).size).toBe(2)
expect(fixture.proposals).not.toContain("msg_canonical")
})
test("substitutes the canonical response ID and reconciles its later echo", async () => {
using fixture = setup()
const request = fixture.data.session.compact({ sessionID })
const proposed = fixture.data.session.pending.list(sessionID)[0].id
await fixture.data.session.pending.sync(sessionID)
expect(fixture.data.session.pending.list(sessionID).map((row) => row.id)).toEqual([proposed])
fixture.response.resolve(Response.json({ data: item("msg_canonical") }))
await request
expect(fixture.data.session.pending.list(sessionID)).toEqual([item("msg_canonical")])
fixture.enqueue("msg_canonical", 20)
expect(fixture.data.session.pending.list(sessionID)).toEqual([item("msg_canonical", 20)])
expect(fixture.data.session.input.list(sessionID)).toEqual([])
})
test.each(["proposed", "canonical"])("adopts the %s echo before the response without duplicating it", async (kind) => {
using fixture = setup()
const request = fixture.data.session.compact({ sessionID })
const id = kind === "proposed" ? fixture.data.session.pending.list(sessionID)[0].id : "msg_canonical"
fixture.enqueue(id, 20)
expect(fixture.data.session.pending.list(sessionID)).toEqual([item(id, 20)])
fixture.response.resolve(Response.json({ data: item(id) }))
await request
expect(fixture.data.session.pending.list(sessionID)).toEqual([item(id, 20)])
})
test.each(["started", "cancelled", "failed"])(
"does not resurrect a canonical item already %s before the response",
async (kind) => {
using fixture = setup()
const request = fixture.data.session.compact({ sessionID })
fixture.enqueue("msg_canonical")
if (kind === "started")
fixture.emit({
...event,
type: "session.compaction.started",
data: { sessionID, inputID: "msg_canonical", reason: "manual" },
})
if (kind === "cancelled")
fixture.emit({ ...event, type: "session.inbox.cancelled", data: { sessionID, inboxID: "msg_canonical" } })
if (kind === "failed")
fixture.emit({
...event,
type: "session.compaction.failed",
data: {
sessionID,
inputID: "msg_canonical",
reason: "manual",
error: { type: "aborted", message: "Cancelled" },
},
})
expect(fixture.data.session.pending.list(sessionID)).toEqual([])
fixture.response.resolve(Response.json({ data: item("msg_canonical") }))
await request
expect(fixture.data.session.pending.list(sessionID)).toEqual([])
if (kind === "started") {
expect(fixture.data.session.message.list(sessionID)).toMatchObject([{ type: "compaction", status: "running" }])
fixture.emit({
...event,
type: "session.compaction.ended",
data: { sessionID, reason: "manual", text: "Summary", recent: "Recent" },
})
expect(fixture.data.session.message.list(sessionID)).toMatchObject([
{ type: "compaction", status: "completed", summary: "Summary" },
])
}
},
)
test.each(["model", "compact"])("rolls back a rejected %s RPC and releases the following prompt", async (rpc) => {
using fixture = setup()
const request = fixture.data.session.compact({ sessionID, model: { providerID: "demo", id: "model" } })
const failed = request.catch((error: unknown) => error)
const prompt = fixture.data.session.prompt({ sessionID, text: "Follow up" })
if (rpc === "model") fixture.model.reject(new Error("Model setup failed"))
if (rpc === "compact") {
fixture.model.resolve()
fixture.response.resolve(new Response("Admission failed", { status: 500 }))
}
expect(await failed).toBeInstanceOf(Error)
await prompt
expect(fixture.data.session.pending.list(sessionID).map((row) => row.type)).toEqual(["user"])
expect(fixture.data.session.message.list(sessionID)).toMatchObject([{ type: "user", text: "Follow up" }])
})
test.each(["proposed", "canonical", "existing"])(
"preserves acknowledged %s compaction after an HTTP error",
async (kind) => {
using fixture = setup()
if (kind === "existing") fixture.enqueue("msg_canonical")
const request = fixture.data.session.compact({ sessionID })
const failed = request.catch((error: unknown) => error)
const id = kind === "proposed" ? fixture.data.session.pending.list(sessionID)[0].id : "msg_canonical"
if (kind !== "existing") fixture.enqueue(id)
expect(fixture.data.session.pending.list(sessionID)).toEqual([item(id)])
fixture.response.resolve(new Response("Lost response", { status: 500 }))
expect(await failed).toBeInstanceOf(Error)
expect(fixture.data.session.pending.list(sessionID)).toEqual([item(id)])
expect(fixture.listeners.size).toBe(1)
},
)
test("uses a fresh control ID when the known pending compaction starts during model setup", async () => {
const proposed = Promise.withResolvers<string>()
using fixture = setup(async (request) => {
if (!request.url.endsWith("/compact")) return undefined
const body = await request.json()
proposed.resolve(body.id)
if (body.id === "msg_existing") return Response.json({ message: "Control ID already consumed" }, { status: 409 })
return Response.json({ data: item(body.id) })
})
fixture.enqueue("msg_existing")
const request = fixture.data.session.compact({ sessionID, model: { providerID: "demo", id: "model" } })
const result = request.catch((error: unknown) => error)
expect(fixture.data.session.pending.list(sessionID)).toEqual([item("msg_existing")])
await wait(() => fixture.calls.includes("model"))
fixture.emit({
...event,
type: "session.compaction.started",
data: { sessionID, inputID: "msg_existing", reason: "manual" },
})
fixture.model.resolve()
expect(await proposed.promise).not.toBe("msg_existing")
expect(await result).toEqual(item(await proposed.promise))
expect(fixture.data.session.pending.list(sessionID)).toEqual([item(await proposed.promise)])
expect(fixture.data.session.message.list(sessionID)).toMatchObject([
{ id: "msg_existing", type: "compaction", status: "running" },
])
})
test.each(["compaction", "canonical compaction", "user"])(
"preserves a fetched durable %s when SSE is delayed and HTTP fails",
async (type) => {
using fixture = setup(async (request) => {
if (request.url.endsWith("/prompt")) return fixture.response.promise
return undefined
})
const request =
type === "user"
? fixture.data.session.prompt({ sessionID, text: "Follow up" })
: fixture.data.session.compact({ sessionID })
const result = request.catch((error: unknown) => error)
const id = type === "canonical compaction" ? "msg_canonical" : fixture.data.session.pending.list(sessionID)[0].id
const durable: SessionInboxInfo =
type === "user" ? { ...item(id, 20), type: "user", payload: { text: "Follow up" } } : item(id, 20)
fixture.pending.push(durable)
await fixture.data.session.pending.sync(sessionID)
expect(fixture.data.session.pending.list(sessionID)).toEqual([durable])
fixture.response.resolve(new Response("Lost response", { status: 500 }))
expect(await result).toBeInstanceOf(Error)
expect(fixture.data.session.pending.list(sessionID)).toEqual([durable])
if (type === "user")
expect(fixture.data.session.message.list(sessionID)).toMatchObject([{ id, type: "user", text: "Follow up" }])
},
)
test("keeps one event listener and removes it when the data owner is disposed during a gate", async () => {
using fixture = setup()
const gate = Promise.withResolvers<void>()
const first = fixture.data.session.prompt({ sessionID, text: "First", gate: gate.promise })
const compact = fixture.data.session.compact({ sessionID })
expect(fixture.listeners.size).toBe(1)
fixture.dispose()
expect(fixture.listeners.size).toBe(0)
gate.resolve()
fixture.response.resolve(Response.json({ data: item("msg_canonical") }))
await Promise.all([first, compact])
expect(fixture.listeners.size).toBe(0)
})
test("routes concurrent compaction observations by session through one listener", async () => {
const firstResponse = Promise.withResolvers<Response>()
const secondResponse = Promise.withResolvers<Response>()
using fixture = setup(async (request) => {
if (!request.url.endsWith("/compact")) return undefined
return request.url.includes(`/session/${sessionID}/`) ? firstResponse.promise : secondResponse.promise
})
const first = fixture.data.session.compact({ sessionID })
const second = fixture.data.session.compact({ sessionID: "ses_other" })
const firstID = fixture.data.session.pending.list(sessionID)[0].id
const secondID = fixture.data.session.pending.list("ses_other")[0].id
expect(fixture.listeners.size).toBe(1)
fixture.emit({ ...event, type: "session.inbox.cancelled", data: { sessionID, inboxID: firstID } })
expect(fixture.data.session.pending.list(sessionID)).toEqual([])
expect(fixture.data.session.pending.list("ses_other").map((row) => row.id)).toEqual([secondID])
firstResponse.resolve(Response.json({ data: item(firstID) }))
secondResponse.resolve(Response.json({ data: { ...item(secondID), sessionID: "ses_other" } }))
await Promise.all([first, second])
expect(fixture.data.session.pending.list(sessionID)).toEqual([])
expect(fixture.data.session.pending.list("ses_other")).toEqual([{ ...item(secondID), sessionID: "ses_other" }])
expect(fixture.listeners.size).toBe(1)
})
test.each(["gate", "prepare"])(
"a preceding prompt's failed %s does not block compaction or following model preparation",
async (kind) => {
using fixture = setup()
const gate = Promise.withResolvers<void>()
const prepared: string[] = []
const first = fixture.data.session
.prompt({
sessionID,
id: "msg_first",
text: "First",
gate: kind === "gate" ? gate.promise : undefined,
prepare: () => {
prepared.push("first")
return gate.promise
},
})
.catch((error: unknown) => error)
const compact = fixture.data.session.compact({ sessionID, model: { providerID: "demo", id: "first" } })
const following = fixture.data.session.prompt({
sessionID,
text: "Follow up",
prepare: () => {
prepared.push("following")
return fixture.api.session.switchModel({ sessionID, model: { providerID: "demo", id: "second" } })
},
})
if (kind === "prepare") await wait(() => prepared.includes("first"))
gate.reject(new Error("Preparation failed"))
expect(await first).toBeInstanceOf(Error)
await wait(() => fixture.calls.includes("model"))
expect(prepared).toEqual(kind === "prepare" ? ["first"] : [])
fixture.model.resolve()
fixture.response.resolve(Response.json({ data: item("msg_canonical") }))
await Promise.all([compact, following])
expect(fixture.calls).toEqual(["model", "compact", "model", "prompt"])
expect(prepared.at(-1)).toBe("following")
expect(fixture.data.session.message.list(sessionID)).toMatchObject([{ type: "user", text: "Follow up" }])
},
)
test("creation failure rejects gated prompt, compaction, and following preparation without sending their RPCs", async () => {
const creation = Promise.withResolvers<Response>()
const requested = Promise.withResolvers<void>()
using fixture = setup(async (request) => {
if (!request.url.endsWith("/api/session")) return undefined
requested.resolve()
return creation.promise
})
const gate = Promise.withResolvers<void>()
const prepared: string[] = []
const created = fixture.data.session.create({ id: sessionID })
const first = fixture.data.session.prompt({ sessionID, text: "First", gate: gate.promise })
const compact = fixture.data.session.compact({ sessionID, model: { providerID: "demo", id: "model" } })
const following = fixture.data.session.prompt({
sessionID,
text: "Follow up",
prepare: async () => {
prepared.push("following")
},
})
const results = Promise.allSettled([created.request, first, compact, following])
await requested.promise
creation.resolve(new Response("Creation failed", { status: 500 }))
expect((await results).map((result) => result.status)).toEqual(["rejected", "rejected", "rejected", "rejected"])
expect(fixture.calls).toEqual([])
expect(prepared).toEqual([])
expect(fixture.data.session.get(sessionID)).toBeUndefined()
expect(fixture.data.session.pending.list(sessionID)).toEqual([])
expect(fixture.listeners.size).toBe(1)
gate.resolve()
})
const sessionID = "ses_compact"
const event = { id: "evt_compact", created: 10, durable: { aggregateID: sessionID, seq: 1, version: 1 } }
const item = (id: string, timeCreated = 10): SessionInboxCompaction => ({
id,
sessionID,
timeCreated,
type: "compaction",
delivery: "steer",
payload: {},
})
function setup(override?: (request: Request) => Promise<Response | undefined>) {
const model = Promise.withResolvers<void>()
const response = Promise.withResolvers<Response>()
const calls: string[] = []
const proposals: string[] = []
const pending: SessionInboxInfo[] = []
const listeners = new Set<Parameters<CreateDataInput["event"]["listen"]>[0]>()
const api = OpenCode.make({
baseUrl: "http://opencode.local",
fetch: async (input, init) => {
const request = input instanceof Request ? input : new Request(input, init)
const overridden = await override?.(request)
if (overridden) return overridden
const rpc = new URL(request.url).pathname.split("/").at(-1)
if (rpc === "inbox") return Response.json({ data: pending })
if (rpc === "model") {
calls.push(rpc)
await model.promise
return new Response(null, { status: 204 })
}
if (rpc === "compact") {
calls.push(rpc)
proposals.push((await request.json()).id)
return (await response.promise).clone()
}
if (rpc === "prompt") {
calls.push(rpc)
return Response.json({
data: { ...item((await request.json()).id), type: "user", payload: { text: "Follow up" } },
})
}
throw new Error(`Unexpected request: ${request.url}`)
},
})
const root = createRoot((dispose) => ({
data: createData({
api: () => api,
directory: "/project",
event: {
on: () => () => {},
listen(handler) {
listeners.add(handler)
return () => listeners.delete(handler)
},
},
}),
dispose,
}))
const emit = (details: OpenCodeEvent) => listeners.forEach((listener) => listener({ name: details.type, details }))
return {
data: root.data,
api,
dispose: root.dispose,
[Symbol.dispose]: root.dispose,
model,
response,
calls,
proposals,
pending,
listeners,
emit,
enqueue(id: string, created = 10) {
emit({
...event,
created,
type: "session.inbox.enqueued",
data: { sessionID, inboxID: id, item: { type: "compaction", delivery: "steer", payload: {} } },
})
},
}
}
async function wait(predicate: () => boolean) {
for (let attempt = 0; attempt < 100; attempt++) {
if (predicate()) return
await Bun.sleep(5)
}
throw new Error("Timed out waiting for request")
}

View file

@ -1318,24 +1318,7 @@ export function Prompt(props: PromptProps) {
restoreEntry()
return true
}
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(target, model)
const switchError = await client.api.session.switchModel({ sessionID: target, model }).then(
() => undefined,
(error) => error,
)
if (switchError) {
cancelCommit()
toast.show({ title: "Failed to switch model", message: errorMessage(switchError), variant: "error" })
restoreEntry()
return true
}
}
const model = { providerID: selection.providerID, id: selection.modelID, variant }
if (session?.revert) {
const error = await client.api.session.revert.commit({ sessionID: target }).then(
() => undefined,
@ -1384,6 +1367,16 @@ export function Prompt(props: PromptProps) {
skills: entry.skills?.length ? entry.skills : undefined,
delivery,
gate: newSession?.gate,
prepare: () => {
// Commit the captured selection after earlier admissions, including
// compaction setup. Cached state may still precede their SSE echoes;
// the server makes an unchanged selection a no-op.
const cancelCommit = local.model.trackSessionCommit(target, model)
return client.api.session.switchModel({ sessionID: target, model }).catch((error) => {
cancelCommit()
throw new Error(`Failed to switch model: ${errorMessage(error)}`, { cause: error })
})
},
})
.catch((error) => {
if (newSession) return newSession.recover(error)

View file

@ -845,18 +845,20 @@ export function Session(props: {
slash: {
name: "compact",
},
run: async () => {
run: () => {
const selection = local.model.current()
if (selection)
await client.api.session.switchModel({
void data.session
.compact({
sessionID: route.sessionID,
model: {
providerID: selection.providerID,
id: selection.modelID,
variant: local.model.variant.current(),
},
model: selection
? {
providerID: selection.providerID,
id: selection.modelID,
variant: local.model.variant.current(),
}
: undefined,
})
await client.api.session.compact({ sessionID: route.sessionID })
.catch((error) => toast.show({ message: errorMessage(error), variant: "error" }))
dialog.clear()
},
},

View file

@ -0,0 +1,235 @@
import { expect, test } from "bun:test"
import { createTestRenderer } from "@opentui/core/testing"
import { InputRenderable } from "@opentui/core"
import { Effect, FileSystem } from "effect"
import { Global } from "@opencode-ai/util/global"
import { createEventStream, createFetch, directory, json } from "./fixture/tui-client"
import { tmpdir } from "./fixture/fixture"
test.each([70, 120])(
"/compact renders before model setup, suppresses repeat gestures, and toasts rollback at %i columns",
async (width) => {
await using state = await tmpdir()
const setup = await createTestRenderer({ width, height: 30, useThread: false, kittyKeyboard: true })
setup.renderer.start()
const ready = Promise.withResolvers<void>()
const model = Promise.withResolvers<Response>()
const modelRequested = Promise.withResolvers<void>()
const events = createEventStream()
const mutations: string[] = []
const sessionID = "ses_compact"
const location = { directory, project: { id: "project", directory, canonical: directory } }
const calls = createFetch(async (url) => {
if (url.pathname === `/api/session/${sessionID}`)
return json({
data: {
id: sessionID,
projectID: "project",
title: "Compact fixture",
model: { providerID: "demo", id: "model" },
location: { directory },
cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: { created: 0, updated: 0 },
},
})
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} })
if (url.pathname === `/api/session/${sessionID}/inbox` || url.pathname === `/api/session/${sessionID}/permission`)
return json({ data: [] })
if (url.pathname === "/api/agent")
return json({ location, data: [{ id: "build", mode: "primary", hidden: false, permissions: [] }] })
if (url.pathname === "/api/provider") return json({ location, data: [{ id: "demo", name: "Demo" }] })
if (url.pathname === "/api/model")
return json({ location, data: [{ id: "model", providerID: "demo", name: "Demo Model", variants: [] }] })
if (url.pathname === `/api/session/${sessionID}/model`) {
mutations.push("model")
modelRequested.resolve()
return model.promise
}
if (url.pathname === `/api/session/${sessionID}/compact`) {
mutations.push("compact")
return json({ data: {} })
}
return undefined
}, events)
const server = Bun.serve({ port: 0, fetch: (request) => calls.fetch(request) })
const { run } = await import("../src/app")
const task = Effect.runPromise(
run({
app: { name: "test", version: "test", channel: "test" },
server: { endpoint: { url: server.url.toString() } },
config: { get: async () => ({ animations: false }), update: async () => ({}) },
packages: { resolve: async () => undefined },
terminalHandoff: async () => ({ renderer: setup.renderer, mode: "dark", complete: ready.resolve }),
args: { sessionID },
log: () => {},
}).pipe(Effect.provide(Global.layerWith({ state: state.path })), Effect.provide(FileSystem.layerNoop({}))),
)
try {
await ready.promise
await setup.waitForFrame((frame) => frame.includes("Demo Model"))
await setup.mockInput.typeText("/compact")
setup.mockInput.pressEnter()
const frame = await setup.waitForFrame((frame) => frame.includes("Compaction queued"))
expect(frame).not.toContain("/compact")
await setup.mockInput.typeText("/compact")
setup.mockInput.pressEnter()
await setup.renderOnce()
expect(setup.captureCharFrame().match(/Compaction queued/g)).toHaveLength(1)
await modelRequested.promise
expect(mutations).toEqual(["model"])
model.resolve(json({ message: "Model setup failed" }, { status: 400 }))
const rejected = await setup.waitForFrame((frame) => frame.includes("Model setup failed"))
expect(rejected).not.toContain("Compaction queued")
expect(mutations).toEqual(["model"])
} finally {
model.resolve(new Response(null, { status: 204 }))
setup.renderer.destroy()
await task
await server.stop()
}
},
)
test.each(["first", "second"])(
"a following prompt commits its selected model after prompt and compaction setup (cached: %s)",
async (initial) => {
await using state = await tmpdir()
const setup = await createTestRenderer({ width: 100, height: 30, useThread: false, kittyKeyboard: true })
setup.renderer.start()
const ready = Promise.withResolvers<void>()
const first = Promise.withResolvers<void>()
const firstRequested = Promise.withResolvers<void>()
const secondRequested = Promise.withResolvers<string>()
const events = createEventStream()
const sessionID = "ses_model_order"
const location = { directory, project: { id: "project", directory, canonical: directory } }
const session = {
id: sessionID,
projectID: "project",
title: "Model ordering fixture",
agent: "build",
model: { providerID: "demo", id: initial },
location: { directory },
cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: { created: 0, updated: 0 },
}
const mutations: string[] = []
const calls = createFetch(async (url, request) => {
if (url.pathname === `/api/session/${sessionID}`) return json({ data: session })
if (url.pathname === `/api/session/${sessionID}/message`) return json({ data: [], cursor: {} })
if (url.pathname === `/api/session/${sessionID}/inbox` || url.pathname === `/api/session/${sessionID}/permission`)
return json({ data: [] })
if (url.pathname === "/api/agent")
return json({ location, data: [{ id: "build", mode: "primary", hidden: false, permissions: [] }] })
if (url.pathname === "/api/provider") return json({ location, data: [{ id: "demo", name: "Demo" }] })
if (url.pathname === "/api/model")
return json({
location,
data: ["first", "second"].map((id) => ({
id,
providerID: "demo",
name: `${id} model`,
variants: [],
cost: [],
time: { released: 0 },
})),
})
if (url.pathname === `/api/session/${sessionID}/model`) {
session.model = (await request.json()).model
mutations.push(`model:${session.model.id}`)
// Delay SSE so local selection must not rely on the cached server model.
return new Response(null, { status: 204 })
}
if (url.pathname === `/api/session/${sessionID}/compact`) {
mutations.push(`compact:${session.model.id}`)
return json({
data: {
id: (await request.json()).id,
sessionID,
type: "compaction",
timeCreated: 10,
payload: {},
delivery: "steer",
},
})
}
if (url.pathname === `/api/session/${sessionID}/prompt`) {
const body = await request.json()
mutations.push(`prompt:${body.text}:${session.model.id}`)
if (body.text === "First prompt") {
firstRequested.resolve()
await first.promise
}
if (body.text === "Second prompt") secondRequested.resolve(session.model.id)
return json({
data: {
id: body.id,
sessionID,
type: "user",
timeCreated: 10,
payload: { text: body.text },
delivery: "steer",
},
})
}
return undefined
}, events)
const server = Bun.serve({ port: 0, fetch: (request) => calls.fetch(request) })
const { run } = await import("../src/app")
const task = Effect.runPromise(
run({
app: { name: "test", version: "test", channel: "test" },
server: { endpoint: { url: server.url.toString() } },
config: { get: async () => ({ animations: false }), update: async () => ({}) },
packages: { resolve: async () => undefined },
terminalHandoff: async () => ({ renderer: setup.renderer, mode: "dark", complete: ready.resolve }),
args: { sessionID },
log: () => {},
}).pipe(Effect.provide(Global.layerWith({ state: state.path })), Effect.provide(FileSystem.layerNoop({}))),
)
const selectModel = async (id: string) => {
await setup.mockInput.typeText("/models")
setup.mockInput.pressEnter()
await setup.waitForFrame(
(frame) => frame.includes("Select model") && setup.renderer.currentFocusedRenderable instanceof InputRenderable,
)
await setup.mockInput.typeText(id)
await setup.renderOnce()
setup.mockInput.pressEnter()
await setup.waitForFrame((frame) => frame.includes(`${id} model`) && !frame.includes("Select model"))
}
try {
await ready.promise
await setup.waitForFrame((frame) => frame.includes(`${initial} model`))
await setup.mockInput.typeText("First prompt")
setup.mockInput.pressEnter()
await firstRequested.promise
const submitted = mutations.length
if (initial !== "first") await selectModel("first")
await setup.mockInput.typeText("/compact")
setup.mockInput.pressEnter()
await setup.waitForFrame((frame) => frame.includes("Compaction queued"))
await selectModel("second")
await setup.mockInput.typeText("Second prompt")
setup.mockInput.pressEnter()
await setup.waitForFrame((frame) => frame.split("\n").slice(0, 20).join("\n").includes("Second prompt"))
expect(mutations).toHaveLength(submitted)
first.resolve()
expect(await secondRequested.promise).toBe("second")
expect(mutations.slice(-4)).toEqual([
"model:first",
"compact:first",
"model:second",
"prompt:Second prompt:second",
])
} finally {
first.resolve()
setup.renderer.destroy()
await task
await server.stop()
}
},
)

View file

@ -135,6 +135,8 @@ export function createFetch(override?: FetchHandler, events?: ReturnType<typeof
})
if (url.pathname === "/api/session") return json({ data: [], cursor: {} })
if (url.pathname === "/api/session/active") return json({ data: {} })
if (request.method === "POST" && /^\/api\/session\/[^/]+\/model$/.test(url.pathname))
return new Response(null, { status: 204 })
if (url.pathname === "/api/permission/request")
return json({
location: { directory, project: { id: "proj_test", directory: worktree, canonical: worktree } },