diff --git a/packages/client/src/solid/data.ts b/packages/client/src/solid/data.ts index 6fc56f8517d..2b3acb4ce1c 100644 --- a/packages/client/src/solid/data.ts +++ b/packages/client/src/solid/data.ts @@ -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() // 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>() - // 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>() + const compacting = new Map; request: Promise }>() + 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(sessionID: string, send: () => Promise, gate?: Promise) { + 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() + 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 }) { - const { gate, ...request } = input + prompt(input: SessionPromptInput & { gate?: Promise; prepare?: () => Promise }) { + 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 diff --git a/packages/client/test/solid-compaction.test.ts b/packages/client/test/solid-compaction.test.ts new file mode 100644 index 00000000000..85848a08ae3 --- /dev/null +++ b/packages/client/test/solid-compaction.test.ts @@ -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() + 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() + 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() + const secondResponse = Promise.withResolvers() + 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() + 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() + const requested = Promise.withResolvers() + using fixture = setup(async (request) => { + if (!request.url.endsWith("/api/session")) return undefined + requested.resolve() + return creation.promise + }) + const gate = Promise.withResolvers() + 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) { + const model = Promise.withResolvers() + const response = Promise.withResolvers() + const calls: string[] = [] + const proposals: string[] = [] + const pending: SessionInboxInfo[] = [] + const listeners = new Set[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") +} diff --git a/packages/tui/src/component/prompt/index.tsx b/packages/tui/src/component/prompt/index.tsx index 14fdab17d7c..531bd36c2ff 100644 --- a/packages/tui/src/component/prompt/index.tsx +++ b/packages/tui/src/component/prompt/index.tsx @@ -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) diff --git a/packages/tui/src/routes/session/index.tsx b/packages/tui/src/routes/session/index.tsx index 4bcd0c1820e..0c48c5552fe 100644 --- a/packages/tui/src/routes/session/index.tsx +++ b/packages/tui/src/routes/session/index.tsx @@ -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() }, }, diff --git a/packages/tui/test/compact-admission.test.tsx b/packages/tui/test/compact-admission.test.tsx new file mode 100644 index 00000000000..54c82348b82 --- /dev/null +++ b/packages/tui/test/compact-admission.test.tsx @@ -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() + const model = Promise.withResolvers() + const modelRequested = Promise.withResolvers() + 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() + const first = Promise.withResolvers() + const firstRequested = Promise.withResolvers() + const secondRequested = Promise.withResolvers() + 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() + } + }, +) diff --git a/packages/tui/test/fixture/tui-client.ts b/packages/tui/test/fixture/tui-client.ts index 5fcb6f4bbd1..c29c424a8f1 100644 --- a/packages/tui/test/fixture/tui-client.ts +++ b/packages/tui/test/fixture/tui-client.ts @@ -135,6 +135,8 @@ export function createFetch(override?: FetchHandler, events?: ReturnType