refactor(client): share Solid server data

This commit is contained in:
Brendan Allan
2026-08-17 14:40:46 +08:00
parent e26473f0bc
commit d5451cdabe
11 changed files with 2001 additions and 1661 deletions
+3
View File
@@ -188,12 +188,15 @@
"@types/bun": "catalog:",
"@typescript/native-preview": "catalog:",
"effect": "catalog:",
"solid-js": "catalog:",
},
"peerDependencies": {
"effect": "4.0.0-beta.101",
"solid-js": ">=1.9.0",
},
"optionalPeers": [
"effect",
"solid-js",
],
},
"packages/codemode": {
+8 -2
View File
@@ -20,6 +20,7 @@
"./promise": "./src/promise/index.ts",
"./promise/api": "./src/promise/api.ts",
"./service": "./src/promise/service.ts",
"./solid": "./src/solid/index.ts",
"./effect": "./src/effect/index.ts",
"./effect/api": "./src/effect/api.ts",
"./effect/service": "./src/effect/service.ts"
@@ -36,11 +37,15 @@
"@opencode-ai/protocol": "workspace:*"
},
"peerDependencies": {
"effect": "4.0.0-beta.101"
"effect": "4.0.0-beta.101",
"solid-js": ">=1.9.0"
},
"peerDependenciesMeta": {
"effect": {
"optional": true
},
"solid-js": {
"optional": true
}
},
"devDependencies": {
@@ -49,6 +54,7 @@
"@tsconfig/bun": "catalog:",
"@types/bun": "catalog:",
"@typescript/native-preview": "catalog:",
"effect": "catalog:"
"effect": "catalog:",
"solid-js": "catalog:"
}
}
+269
View File
@@ -0,0 +1,269 @@
import { batch, onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store"
import type { OpenCodeClient, OpenCodeEvent } from "../promise"
export type ClientConnectionStatus = "connected" | "connecting" | "reconnecting"
export type ClientConnectionEvent = {
readonly type: "client.connection"
readonly created: number
readonly data: {
readonly status: "connecting" | "connected" | "disconnected" | "reconnecting"
readonly attempt: number
readonly error?: string
}
}
export type ClientConnectionOptions = {
readonly reconnect?: (signal: AbortSignal) => Promise<OpenCodeClient>
readonly onEvent: (event: OpenCodeEvent) => void
readonly flushInterval?: number
readonly pageLifecycle?: boolean
readonly log?: {
readonly debug?: (message: string, data?: Readonly<Record<string, unknown>>) => void
readonly info?: (message: string, data?: Readonly<Record<string, unknown>>) => void
}
}
const connectTimeout = 2_000
const reconnectDelay = 1_000
const connectionHistoryLimit = 50
type CurrentDelta = Extract<
OpenCodeEvent,
{ type: "session.text.delta" | "session.reasoning.delta" | "session.tool.input.delta" | "session.compaction.delta" }
>
export function coalesceClientEvents(events: OpenCodeEvent[]) {
return events.reduce<OpenCodeEvent[]>((output, event) => {
const current = currentDelta(event)
const previous = output[output.length - 1]
const prior = currentDelta(previous)
if (
!current ||
!prior ||
previous?.location?.directory !== event.location?.directory ||
currentDeltaKey(prior) !== currentDeltaKey(current)
) {
output.push(event)
return output
}
const fragment = currentDeltaFragment(prior) + currentDeltaFragment(current)
output[output.length - 1] = {
...current,
data:
current.type === "session.compaction.delta"
? { ...current.data, text: fragment }
: { ...current.data, delta: fragment },
} as CurrentDelta
return output
}, [])
}
function currentDelta(event: OpenCodeEvent | undefined): CurrentDelta | undefined {
if (
event?.type === "session.text.delta" ||
event?.type === "session.reasoning.delta" ||
event?.type === "session.tool.input.delta" ||
event?.type === "session.compaction.delta"
)
return event
}
function currentDeltaKey(event: CurrentDelta) {
if (event.type === "session.tool.input.delta")
return `${event.type}:${event.data.sessionID}:${event.data.assistantMessageID}:${event.data.id}`
if (event.type === "session.compaction.delta") return `${event.type}:${event.data.sessionID}`
return `${event.type}:${event.data.sessionID}:${event.data.assistantMessageID}:${event.data.ordinal}`
}
function currentDeltaFragment(event: CurrentDelta) {
return event.type === "session.compaction.delta" ? event.data.text : event.data.delta
}
export function createClientConnection(initialApi: OpenCodeClient, options: ClientConnectionOptions) {
const abort = new AbortController()
const history: ClientConnectionEvent[] = []
const [connection, setConnection] = createStore<{
status: ClientConnectionStatus
attempt: number
error?: string
}>({ status: "connecting", attempt: 0 })
let api = initialApi
let pending: OpenCodeEvent[] = []
let flushTimer: ReturnType<typeof setTimeout> | undefined
let stream: AbortController | undefined
let run: Promise<void> | undefined
let started = false
let generation = 0
function record(status: ClientConnectionEvent["data"]["status"], attempt: number, error?: string) {
history.push({ type: "client.connection", created: Date.now(), data: { status, attempt, error } })
if (history.length > connectionHistoryLimit) history.shift()
}
function publish(event: OpenCodeEvent) {
pending.push(event)
if (flushTimer) return
flushTimer = setTimeout(() => {
flushTimer = undefined
const events = pending
pending = []
batch(() => coalesceClientEvents(events).forEach(options.onEvent))
}, options.flushInterval ?? 10)
}
async function connect(signal: AbortSignal, attempt: number) {
let connectedAt: number | undefined
const request = new AbortController()
const cancel = () => request.abort(signal.reason)
const timeout = setTimeout(() => request.abort(new Error("Timed out connecting to server")), connectTimeout)
signal.addEventListener("abort", cancel, { once: true })
try {
record(attempt === 0 ? "connecting" : "reconnecting", attempt)
options.log?.info?.("event stream connecting", { attempt })
const iterator = api.event.subscribe({ signal: request.signal })[Symbol.asyncIterator]()
const first = await iterator.next()
if (signal.aborted) return { error: undefined, connectedAt }
if (first.done)
return {
error: request.signal.reason instanceof Error ? request.signal.reason : new Error("Event stream disconnected"),
connectedAt,
}
if (first.value.type !== "server.connected")
return { error: new Error("Event stream did not start with server.connected"), connectedAt }
clearTimeout(timeout)
record("connected", attempt)
connectedAt = Date.now()
options.log?.info?.("event stream connected")
publish(first.value)
setConnection({ status: "connected", attempt: 0, error: undefined })
while (!signal.aborted) {
const event = await iterator.next()
if (signal.aborted) return { error: undefined, connectedAt }
if (event.done) return { error: new Error("Event stream disconnected"), connectedAt }
if ("durable" in event.value)
options.log?.debug?.("event", {
type: event.value.type,
aggregateID: event.value.durable.aggregateID,
seq: event.value.durable.seq,
})
publish(event.value)
}
return { error: undefined, connectedAt }
} catch (error) {
return { error, connectedAt }
} finally {
request.abort()
clearTimeout(timeout)
signal.removeEventListener("abort", cancel)
}
}
async function runStream(active: number) {
let attempt = 0
while (!abort.signal.aborted && started && generation === active) {
setConnection({ status: attempt === 0 ? "connecting" : "reconnecting", attempt, error: undefined })
const controller = new AbortController()
stream = controller
const cancel = () => controller.abort(abort.signal.reason)
abort.signal.addEventListener("abort", cancel)
const result = await connect(controller.signal, attempt)
abort.signal.removeEventListener("abort", cancel)
if (abort.signal.aborted || !started || generation !== active) return
if (result.connectedAt !== undefined && Date.now() - result.connectedAt >= reconnectDelay) attempt = 0
attempt += 1
const message = errorMessage(result.error)
record("disconnected", attempt, message)
options.log?.info?.("event stream disconnected", { attempt, error: message })
setConnection({ status: "reconnecting", attempt, error: message })
if (options.reconnect) {
const next = await options.reconnect(controller.signal).catch((error) => {
if (!controller.signal.aborted)
options.log?.info?.("server resolution failed", { attempt, error: errorMessage(error) })
})
if (abort.signal.aborted || controller.signal.aborted || !started || generation !== active) return
if (next) {
api = next
if (attempt === 1) continue
}
}
await wait(reconnectDelay, controller.signal)
}
}
function start() {
if (started) return run
started = true
const active = ++generation
const previous = run
const current = (async () => {
if (previous) await previous
await runStream(active)
})().finally(() => {
if (run !== current) return
run = undefined
})
run = current
return run
}
function stop() {
started = false
generation += 1
stream?.abort()
}
onMount(() => {
if (options.pageLifecycle) {
const pagehide = () => stop()
const pageshow = (event: PageTransitionEvent) => {
if (event.persisted) void start()
}
window.addEventListener("pagehide", pagehide)
window.addEventListener("pageshow", pageshow)
onCleanup(() => {
window.removeEventListener("pagehide", pagehide)
window.removeEventListener("pageshow", pageshow)
})
}
void start()
})
onCleanup(() => {
stop()
abort.abort()
if (flushTimer) clearTimeout(flushTimer)
pending = []
})
return {
status: () => connection.status,
attempt: () => connection.attempt,
error: () => connection.error,
internal: {
history: () => history.slice(),
},
}
}
function errorMessage(error: unknown) {
if (error === undefined) return undefined
if (error instanceof Error) return error.message
return String(error)
}
function wait(delay: number, signal: AbortSignal) {
return new Promise<void>((resolve) => {
const timer = setTimeout(done, delay)
signal.addEventListener("abort", done, { once: true })
function done() {
clearTimeout(timer)
signal.removeEventListener("abort", done)
resolve()
}
})
}
File diff suppressed because it is too large Load Diff
+2
View File
@@ -0,0 +1,2 @@
export * from "./data"
export * from "./connection"
@@ -0,0 +1,44 @@
import { describe, expect, test } from "bun:test"
import type { OpenCodeEvent } from "../src/promise"
import { coalesceClientEvents } from "../src/solid/connection"
describe("coalesceClientEvents", () => {
const delta = (id: string, value: string, ordinal = 0) =>
({
id,
created: 1,
type: "session.text.delta",
location: { directory: "/repo" },
data: { sessionID: "ses", assistantMessageID: "msg", ordinal, delta: value },
}) as OpenCodeEvent
test("merges adjacent deltas for the same stream", () => {
const result = coalesceClientEvents([delta("evt_1", "hello "), delta("evt_2", "world")])
expect(result).toHaveLength(1)
expect(result[0]).toMatchObject({ id: "evt_2", data: { delta: "hello world" } })
})
test("coalesces tool input deltas by tool ID", () => {
const current = (eventID: string, id: string, value: string) =>
({
id: eventID,
created: 1,
type: "session.tool.input.delta",
location: { directory: "/repo" },
data: { sessionID: "ses", assistantMessageID: "msg", id, delta: value },
}) as OpenCodeEvent
const result = coalesceClientEvents([
current("evt_1", "call_1", "{"),
current("evt_2", "call_1", "}"),
current("evt_3", "call_2", "[]"),
])
expect(result).toHaveLength(2)
expect(result[0]).toMatchObject({ id: "evt_2", data: { id: "call_1", delta: "{}" } })
expect(result[1]).toMatchObject({ id: "evt_3", data: { id: "call_2", delta: "[]" } })
})
test("preserves boundaries between distinct delta streams", () => {
const events = [delta("evt_1", "a"), delta("evt_2", "b", 1), delta("evt_3", "c")]
expect(coalesceClientEvents(events).map((event) => event.id)).toEqual(["evt_1", "evt_2", "evt_3"])
})
})
+86
View File
@@ -0,0 +1,86 @@
import { expect, mock, test } from "bun:test"
import type { OpenCodeClient, OpenCodeEvent, Project, SessionInfo } from "../src/promise"
import { createServerData, type CreateServerDataInput } from "../src/solid/data"
import { createRoot } from "solid-js"
const session = {
id: "ses_fork",
parentID: "ses_parent",
projectID: "pro_1",
title: "Fork",
cost: 0,
tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
time: { created: 1, updated: 1 },
location: { directory: "/repo" },
} as SessionInfo
const project = {
id: "pro_1",
name: "Repo",
directory: "/repo",
canonical: "/repo",
vcs: "git",
sandboxes: [],
time: { created: 1, updated: 1 },
} as Project
test("refreshes forked sessions and worktree projects from events", async () => {
const listeners = new Set<(event: { name: OpenCodeEvent["type"]; details: OpenCodeEvent }) => void>()
const sessionGet = mock(async () => session)
const projectList = mock(async () => [project])
const api = {
session: { get: sessionGet },
project: { list: projectList },
} as unknown as OpenCodeClient
const event = {
on: (() => () => {}) as CreateServerDataInput["event"]["on"],
listen(handler: (event: { name: OpenCodeEvent["type"]; details: OpenCodeEvent }) => void) {
listeners.add(handler)
return () => listeners.delete(handler)
},
}
const emit = (details: OpenCodeEvent) => listeners.forEach((listener) => listener({ name: details.type, details }))
await new Promise<void>((resolve) => {
createRoot((dispose) => {
const data = createServerData({ api: () => api, directory: "/repo", event, connection: { status: () => "connected" } })
emit({ type: "session.forked", data: { sessionID: session.id } } as unknown as OpenCodeEvent)
emit({ type: "worktree.updated", data: { projectID: project.id } } as unknown as OpenCodeEvent)
void Promise.all([sessionGet, projectList].map((fn) => fn.mock.results[0]?.value)).then(() => {
expect(data.session.get(session.id)).toEqual(session)
expect(data.project.get(project.id)).toEqual(project)
expect(sessionGet).toHaveBeenCalledTimes(1)
expect(projectList).toHaveBeenCalledTimes(1)
dispose()
resolve()
})
})
})
})
test("resolves location info through the requested ref after canonicalization", async () => {
const requested = { directory: "/repo/../repo" }
const canonical = {
directory: "/repo",
project: { id: "pro_1", directory: "/repo", canonical: "/repo" },
}
const api = {
location: { get: mock(async () => canonical) },
} as unknown as OpenCodeClient
const event = {
on: (() => () => {}) as CreateServerDataInput["event"]["on"],
listen: () => () => {},
} as CreateServerDataInput["event"]
await new Promise<void>((resolve) => {
createRoot((dispose) => {
const data = createServerData({ api: () => api, directory: requested.directory, event })
void data.location.syncInfo(requested).then(() => {
expect(data.location.info(requested)).toEqual(canonical)
expect(data.location.info({ directory: canonical.directory })).toEqual(canonical)
dispose()
resolve()
})
})
})
})
+17 -182
View File
@@ -1,177 +1,39 @@
import type { OpenCodeClient, OpenCodeEvent } from "@opencode-ai/client"
import { createClientConnection } from "@opencode-ai/client/solid"
import { createGlobalEmitter } from "@solid-primitives/event-bus"
import { batch, onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store"
import { errorMessage } from "../util/error"
import { onCleanup } from "solid-js"
import { createSimpleContext } from "./helper"
import { useLog } from "./log"
export type ClientConnectionStatus = "connected" | "connecting" | "reconnecting"
export type ClientConnectionEvent = {
readonly type: "client.connection"
readonly created: number
readonly data: {
readonly status: "connecting" | "connected" | "disconnected" | "reconnecting"
readonly attempt: number
readonly error?: string
}
}
type ManagedService = {
reconnect: (signal: AbortSignal) => Promise<{ api: OpenCodeClient }>
restart: () => Promise<void>
}
type ClientEventMap = { [Type in OpenCodeEvent["type"]]: Extract<OpenCodeEvent, { type: Type }> }
const connectTimeout = 2_000
const connectionHistoryLimit = 50
const eventFlushInterval = 10
export const { use: useClient, provider: ClientProvider } = createSimpleContext({
name: "Client",
init: (props: { api: OpenCodeClient; service?: ManagedService }) => {
const log = useLog({ component: "client" })
const abort = new AbortController()
const history: ClientConnectionEvent[] = []
let api = props.api
const service = props.service
const events = createGlobalEmitter<ClientEventMap>()
let pending: OpenCodeEvent[] = []
let flushTimer: ReturnType<typeof setTimeout> | undefined
const [connection, setConnection] = createStore<{
status: ClientConnectionStatus
attempt: number
error?: string
}>({
status: "connecting",
attempt: 0,
})
let stream: AbortController | undefined
let api = props.api
function record(status: ClientConnectionEvent["data"]["status"], attempt: number, error?: string) {
history.push({ type: "client.connection", created: Date.now(), data: { status, attempt, error } })
if (history.length > connectionHistoryLimit) history.shift()
}
function flushEvents() {
flushTimer = undefined
const queued = pending
pending = []
batch(() => queued.forEach((event) => events.emit(event.type, event)))
}
function emit(event: OpenCodeEvent) {
pending.push(event)
if (flushTimer) return
flushTimer = setTimeout(flushEvents, eventFlushInterval)
}
async function connect(signal: AbortSignal, attempt: number) {
let connectedAt: number | undefined
// Bound the initial handshake and tie this request to the stream lifetime.
const request = new AbortController()
const cancel = () => request.abort(signal.reason)
const timeout = setTimeout(() => request.abort(new Error("Timed out connecting to server")), connectTimeout)
signal.addEventListener("abort", cancel, { once: true })
try {
// Open the event stream and validate its initial handshake.
record(attempt === 0 ? "connecting" : "reconnecting", attempt)
log.info("event stream connecting", { attempt })
const iterator = api.event.subscribe({ signal: request.signal })[Symbol.asyncIterator]()
const first = await iterator.next()
if (signal.aborted) return { error: undefined, connectedAt }
if (first.done) {
const error =
request.signal.reason instanceof Error ? request.signal.reason : new Error("Event stream disconnected")
return { error, connectedAt }
}
if (first.value.type !== "server.connected")
return { error: new Error("Event stream did not start with server.connected"), connectedAt }
// Publish the connected state before forwarding live events.
clearTimeout(timeout)
record("connected", attempt)
connectedAt = Date.now()
log.info("event stream connected")
emit(first.value)
setConnection({ status: "connected", attempt: 0, error: undefined })
// Forward events until the stream closes or this connection is cancelled.
while (!signal.aborted) {
const event = await iterator.next()
if (signal.aborted) return { error: undefined, connectedAt }
if (event.done) return { error: new Error("Event stream disconnected"), connectedAt }
if ("durable" in event.value)
log.debug("event", {
type: event.value.type,
aggregateID: event.value.durable.aggregateID,
seq: event.value.durable.seq,
})
emit(event.value)
}
return { error: undefined, connectedAt }
} catch (error) {
return { error, connectedAt }
} finally {
request.abort()
clearTimeout(timeout)
signal.removeEventListener("abort", cancel)
}
}
function start() {
stream?.abort()
const controller = new AbortController()
stream = controller
void (async () => {
let attempt = 0
while (!abort.signal.aborted && !controller.signal.aborted) {
const result = await connect(controller.signal, attempt)
if (abort.signal.aborted || controller.signal.aborted) return
if (result.connectedAt !== undefined && Date.now() - result.connectedAt >= 1_000) attempt = 0
attempt += 1
const message = errorMessage(result.error)
record("disconnected", attempt, message)
log.info("event stream disconnected", {
attempt,
error: message,
})
setConnection({ status: "reconnecting", attempt, error: message })
// Re-resolve the transport before retrying: the server may have
// moved (service restarted on a new port) or need starting. Static
// transports (--server, standalone) resolve to the same address.
if (props.service) {
const next = await props.service.reconnect(controller.signal).catch((error) => {
if (!controller.signal.aborted)
log.info("server resolution failed", {
attempt,
error: errorMessage(error),
})
})
if (abort.signal.aborted || controller.signal.aborted) return
if (next) {
api = next.api
if (attempt === 1) continue
}
const connection = createClientConnection(api, {
reconnect: service
? async (signal) => {
api = (await service.reconnect(signal)).api
return api
}
await wait(1_000, controller.signal)
}
})()
}
: undefined,
onEvent(event) {
events.emit(event.type, event)
},
log,
})
onMount(start)
onCleanup(() => {
abort.abort()
stream?.abort()
if (flushTimer) clearTimeout(flushTimer)
pending = []
events.clear()
})
@@ -183,35 +45,8 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
on: events.on,
listen: events.listen,
},
connection: {
status() {
return connection.status
},
attempt() {
return connection.attempt
},
error() {
return connection.error
},
internal: {
history() {
return history.slice()
},
},
},
restart: props.service?.restart,
connection,
restart: service?.restart,
}
},
})
function wait(delay: number, signal: AbortSignal) {
return new Promise<void>((resolve) => {
const timer = setTimeout(done, delay)
signal.addEventListener("abort", done, { once: true })
function done() {
clearTimeout(timer)
signal.removeEventListener("abort", done)
resolve()
}
})
}
File diff suppressed because it is too large Load Diff
+8 -4
View File
@@ -1873,15 +1873,14 @@ test("refreshes effective catalog data after catalog updates", async () => {
test("refreshes agents after agent updates", async () => {
const events = createEventStream()
let requests = 0
let agentID = "build"
const calls = createFetch((url) => {
if (url.pathname !== "/api/agent") return
requests++
return json({
location: { directory, project: { id: "proj_test", directory } },
data: [
{
id: requests === 1 ? "build" : "reviewer",
id: agentID,
request: { headers: {}, body: {} },
mode: "primary",
hidden: false,
@@ -1911,6 +1910,11 @@ test("refreshes agents after agent updates", async () => {
try {
await wait(() => data.location.agent.list()?.[0]?.id === "build")
await Bun.sleep(20)
events.emit({ id: "evt_agent_unlocated", created: 0, type: "agent.updated", data: {} })
await Bun.sleep(20)
expect(data.location.agent.list()?.[0]?.id).toBe("build")
agentID = "reviewer"
emitEvent(events, { id: "evt_agent", created: 0, type: "agent.updated", data: {} })
await wait(() => data.location.agent.list()?.[0]?.id === "reviewer")
} finally {
@@ -2801,7 +2805,7 @@ test("renders admitted prompts immediately and tracks them until promoted", asyn
await mounted
const received: string[] = []
const unsubscribe = sync.listen((event) => received.push(event.name))
emitEvent(events, {
events.emit({
id: "evt_admitted_1",
created: 0,
type: "session.inbox.enqueued",
+5
View File
@@ -153,6 +153,11 @@ export function createFetch(override?: FetchHandler, events?: ReturnType<typeof
return json({ location: { directory, project: { id: "proj_test", directory, canonical: directory } }, data: [] })
}
if (url.pathname === "/provider") return json({ all: [], default: {}, connected: [] })
if (url.pathname === "/api/model/default")
return json({
location: { directory, project: { id: "proj_test", directory: worktree, canonical: worktree } },
data: null,
})
if (url.pathname === "/session") return json([])
if (url.pathname === "/vcs") return json({ branch: "main" })
if (url.pathname === "/api/experimental/migration/v1") return json({ status: "completed" })