From f253bdaeef10b528e1a062e5d387349e58520079 Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Mon, 17 Aug 2026 17:20:29 -0400 Subject: [PATCH] feat(tui): wire engine data layer --- packages/client/src/solid/engine-data.ts | 35 +++++++--- packages/client/src/solid/engine/engine.ts | 1 + packages/tui/src/component/prompt/index.tsx | 7 +- packages/tui/src/context/data.tsx | 4 +- packages/tui/test/fixture/tui-client.ts | 74 ++++++++++++++++++++- 5 files changed, 106 insertions(+), 15 deletions(-) diff --git a/packages/client/src/solid/engine-data.ts b/packages/client/src/solid/engine-data.ts index 73fc6ffa67c..98c63c1ee04 100644 --- a/packages/client/src/solid/engine-data.ts +++ b/packages/client/src/solid/engine-data.ts @@ -72,9 +72,12 @@ export function createEngineData(config: CreateDataInput) { const update = (sessionID: string, view: Engine.SessionView) => { batch(() => { - setViews(sessionID, reconcile(view)) - legacy.session.remember(view.session) - view.children.forEach(legacy.session.remember) + setViews(sessionID, reconcile(structuredClone(view))) + const current = legacy.session.get(sessionID) + if (!current || current.time.updated <= view.session.time.updated) { + legacy.session.remember(structuredClone(view.session)) + } + view.children.forEach((child) => legacy.session.remember(structuredClone(child))) }) } @@ -101,40 +104,50 @@ export function createEngineData(config: CreateDataInput) { ...legacy, session: { ...legacy.session, - sync(sessionID: string) { + sync(sessionID: string, _options?: { readonly children?: boolean }) { return ensure(sessionID).then(() => undefined) }, status(sessionID: string) { - return views[sessionID]?.active ?? legacy.session.status(sessionID) + if (views[sessionID]?.active === "running") return "running" + return legacy.session.status(sessionID) }, input: { list(sessionID: string) { - return views[sessionID]?.pending.filter((item) => item.type !== "compaction").map((item) => item.id) ?? [] + return ( + views[sessionID]?.pending.filter((item) => item.type !== "compaction").map((item) => item.id) ?? + legacy.session.input.list(sessionID) + ) }, has(sessionID: string, inboxID: string) { - return views[sessionID]?.pending.some((item) => item.type !== "compaction" && item.id === inboxID) ?? false + return ( + views[sessionID]?.pending.some((item) => item.type !== "compaction" && item.id === inboxID) ?? + legacy.session.input.has(sessionID, inboxID) + ) }, }, pending: { list(sessionID: string) { - return views[sessionID]?.pending ?? [] + void ensure(sessionID) + return [...(views[sessionID]?.pending ?? [])] }, sync(sessionID: string) { return ensure(sessionID).then(() => undefined) }, - invalidate() {}, + invalidate(_sessionID: string) {}, }, message: { list(sessionID: string) { - return views[sessionID]?.messages ?? [] + void ensure(sessionID) + return [...(views[sessionID]?.messages ?? [])] }, get(sessionID: string, messageID: string) { + void ensure(sessionID) return views[sessionID]?.messages.find((message) => message.id === messageID) }, sync(sessionID: string) { return ensure(sessionID).then(() => undefined) }, - invalidate() {}, + invalidate(_sessionID: string) {}, }, async prompt(input: SessionPromptInput) { return (await ensure(input.sessionID)).submit({ diff --git a/packages/client/src/solid/engine/engine.ts b/packages/client/src/solid/engine/engine.ts index c8e5382041e..b1e62a058ed 100644 --- a/packages/client/src/solid/engine/engine.ts +++ b/packages/client/src/solid/engine/engine.ts @@ -276,6 +276,7 @@ export function render(state: Pick [ ...state.folded.messages, ...pending.flatMap((item): ReadonlyArray => { + if (item.type !== "compaction" && item.delivery === "queue") return [] if (messageIDs.has(item.id)) return [] const message = SessionFold.messageFromInbox(item) return message ? [message] : [] diff --git a/packages/tui/src/component/prompt/index.tsx b/packages/tui/src/component/prompt/index.tsx index e3e37358a8d..9061db81201 100644 --- a/packages/tui/src/component/prompt/index.tsx +++ b/packages/tui/src/component/prompt/index.tsx @@ -208,6 +208,11 @@ export function Prompt(props: PromptProps) { const config = useConfig().data const dialog = useDialog() const toast = useToast() + onCleanup( + data.session.failures.listen((failure) => { + toast.show({ title: "Prompt rejected", message: failure.reason, variant: "error" }) + }), + ) const status = createMemo(() => data.session.status(props.sessionID ?? "")) const history = usePromptHistory() const stash = usePromptStash() @@ -1304,7 +1309,7 @@ export function Prompt(props: PromptProps) { return false } } - const error = await client.api.session + const error = await data.session .prompt({ sessionID, text: inputText, diff --git a/packages/tui/src/context/data.tsx b/packages/tui/src/context/data.tsx index fa164735e4e..0285d31e413 100644 --- a/packages/tui/src/context/data.tsx +++ b/packages/tui/src/context/data.tsx @@ -1,4 +1,4 @@ -import { createData } from "@opencode-ai/client/solid" +import { createEngineData } from "@opencode-ai/client/solid" import type { Plugin } from "@opencode-ai/plugin/tui" import { createSimpleContext } from "./helper" import { useClient } from "./client" @@ -10,7 +10,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({ name: "Data", init: () => { const client = useClient() - const data = createData({ + const data = createEngineData({ api: () => client.api, event: client.event, connection: client.connection, diff --git a/packages/tui/test/fixture/tui-client.ts b/packages/tui/test/fixture/tui-client.ts index 51f79a6a840..fe86f8d3b48 100644 --- a/packages/tui/test/fixture/tui-client.ts +++ b/packages/tui/test/fixture/tui-client.ts @@ -16,6 +16,9 @@ export function createEventStream() { const encoder = new TextEncoder() const v2 = new Set>() const pending: Uint8Array[] = [] + const logs = new Map>>() + const logSeq = new Map() + const logHistory = new Map>() const response = ( controllers: Set>, queued: Uint8Array[], @@ -27,7 +30,8 @@ export function createEventStream() { start(controller) { current = controller controllers.add(controller) - if (initial) controller.enqueue(encoder.encode(`data: ${JSON.stringify(initial)}\n\n`)) + const values = Array.isArray(initial) ? initial : initial ? [initial] : [] + for (const value of values) controller.enqueue(encoder.encode(`data: ${JSON.stringify(value)}\n\n`)) for (const chunk of queued.splice(0)) controller.enqueue(chunk) }, cancel() { @@ -53,13 +57,44 @@ export function createEventStream() { return { emit(event: OpenCodeEvent) { send(v2, pending, event) + const sessionID = + "durable" in event + ? event.durable.aggregateID + : "sessionID" in event.data && typeof event.data.sessionID === "string" + ? event.data.sessionID + : undefined + if (!sessionID) return + const seq = (logSeq.get(sessionID) ?? 0) + 1 + const item = "durable" in event ? { ...event, durable: { ...event.durable, seq } } : event + if ("durable" in event) { + logSeq.set(sessionID, seq) + logHistory.set(sessionID, [...(logHistory.get(sessionID) ?? []), { seq, event: item }]) + } + const controllers = logs.get(sessionID) + if (controllers) send(controllers, [], item) }, v2() { return response(v2, pending, { id: "evt_connected", type: "server.connected", data: {} }) }, + log(sessionID: string, after: number) { + const controllers = logs.get(sessionID) ?? new Set>() + logs.set(sessionID, controllers) + return response( + controllers, + [], + [ + ...(logHistory.get(sessionID) ?? []).filter((entry) => entry.seq > after).map((entry) => entry.event), + { type: "log.synced", aggregateID: sessionID, seq: logSeq.get(sessionID) ?? 0 }, + ], + ) + }, disconnect() { for (const controller of v2) controller.close() v2.clear() + for (const controllers of logs.values()) { + for (const controller of controllers) controller.close() + controllers.clear() + } }, } } @@ -68,6 +103,7 @@ export type FetchHandler = (url: URL, request: Request) => Response | undefined export function createFetch(override?: FetchHandler, events?: ReturnType) { const session = [] as URL[] + const sessionEvents = events ?? createEventStream() async function fetch(input: RequestInfo | URL, init?: RequestInit) { const request = input instanceof Request ? input : new Request(input, init) const url = new URL(request.url) @@ -75,6 +111,42 @@ export function createFetch(override?: FetchHandler, events?: ReturnType { + const response = await override?.(new URL(path, url), new Request(new URL(path, url))) + if (!response) return fallback + const body = await response.json() + if (typeof body !== "object" || body === null || !("data" in body)) return fallback + return body.data + } + const children = await read(`/api/session?parentID=${encodeURIComponent(sessionID)}`, []) + const messages = await read(`/api/session/${encodeURIComponent(sessionID)}/message`, []) + return json({ + data: { + session: await read(`/api/session/${encodeURIComponent(sessionID)}`, { + id: sessionID, + projectID: "proj_test", + location: { directory }, + cost: 0, + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + time: { created: 0, updated: 0 }, + }), + children: Array.isArray(children) + ? children.filter( + (child) => + typeof child === "object" && child !== null && "parentID" in child && child.parentID === sessionID, + ) + : [], + inbox: await read(`/api/session/${encodeURIComponent(sessionID)}/inbox`, []), + messages: Array.isArray(messages) ? messages.toReversed() : [], + seq: 0, + }, + }) + } + const log = url.pathname.match(/^\/api\/experimental\/session\/([^/]+)\/log$/) + if (log) return sessionEvents.log(decodeURIComponent(log[1]), Number(url.searchParams.get("after") ?? 0)) if ( [