fix: stabilize v2 runtime behavior
This commit is contained in:
parent
ac2a78391f
commit
655adbf46e
35 changed files with 742 additions and 425 deletions
|
|
@ -21,15 +21,10 @@ import type {
|
|||
import { createStore, produce } from "solid-js/store"
|
||||
import { createSimpleContext } from "./helper"
|
||||
import { useSDK } from "./sdk"
|
||||
import { createSignal, onCleanup, onMount } from "solid-js"
|
||||
import { createGlobalEmitter } from "@solid-primitives/event-bus"
|
||||
import { createSignal, onCleanup } from "solid-js"
|
||||
|
||||
export type DataConnectionStatus = "connecting" | "connected" | "reconnecting"
|
||||
export type DataSessionStatus = "idle" | "running"
|
||||
|
||||
export type DataEvent = V2Event
|
||||
type DataEventMap = { [T in DataEvent["type"]]: Extract<DataEvent, { type: T }> }
|
||||
|
||||
type LocationData = {
|
||||
agent?: AgentV2Info[]
|
||||
command?: CommandV2Info[]
|
||||
|
|
@ -41,11 +36,6 @@ type LocationData = {
|
|||
}
|
||||
|
||||
type Data = {
|
||||
connection: {
|
||||
status: DataConnectionStatus
|
||||
attempt: number
|
||||
error?: string
|
||||
}
|
||||
session: {
|
||||
info: Record<string, SessionV2Info>
|
||||
status: Record<string, DataSessionStatus>
|
||||
|
|
@ -71,10 +61,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||
name: "Data",
|
||||
init: () => {
|
||||
const [store, setStore] = createStore<Data>({
|
||||
connection: {
|
||||
status: "connecting",
|
||||
attempt: 0,
|
||||
},
|
||||
session: {
|
||||
info: {},
|
||||
status: {},
|
||||
|
|
@ -89,7 +75,6 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||
})
|
||||
|
||||
const sdk = useSDK()
|
||||
const events = createGlobalEmitter<DataEventMap>()
|
||||
const [defaultLocation, setDefaultLocation] = createSignal<LocationRef>({
|
||||
directory: process.cwd(),
|
||||
})
|
||||
|
|
@ -151,6 +136,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||
|
||||
function handleEvent(event: V2Event) {
|
||||
switch (event.type) {
|
||||
case "session.created":
|
||||
void result.session.refresh(event.data.sessionID)
|
||||
break
|
||||
case "catalog.updated":
|
||||
void Promise.all([
|
||||
result.location.model.refresh(event.location),
|
||||
|
|
@ -477,21 +465,20 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||
])
|
||||
break
|
||||
}
|
||||
events.emit(event.type, event)
|
||||
}
|
||||
|
||||
const result = {
|
||||
on: events.on,
|
||||
listen: events.listen,
|
||||
on: sdk.event.on,
|
||||
listen: sdk.event.listen,
|
||||
connection: {
|
||||
status() {
|
||||
return store.connection.status
|
||||
return sdk.connection.status()
|
||||
},
|
||||
attempt() {
|
||||
return store.connection.attempt
|
||||
return sdk.connection.attempt()
|
||||
},
|
||||
error() {
|
||||
return store.connection.error
|
||||
return sdk.connection.error()
|
||||
},
|
||||
},
|
||||
session: {
|
||||
|
|
@ -687,52 +674,13 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
|
|||
console.error("Failed to refresh default location data", failure.reason)
|
||||
}
|
||||
|
||||
onMount(() => {
|
||||
const controller = new AbortController()
|
||||
onCleanup(() => controller.abort())
|
||||
void (async () => {
|
||||
while (!controller.signal.aborted) {
|
||||
const error = await (async () => {
|
||||
const events = await sdk.client.v2.event.subscribe({
|
||||
signal: controller.signal,
|
||||
sseMaxRetryAttempts: 0,
|
||||
throwOnError: true,
|
||||
})
|
||||
const stream = events.stream[Symbol.asyncIterator]()
|
||||
const first = await stream.next()
|
||||
if (first.done) return new Error("Event stream disconnected")
|
||||
setStore("connection", { status: "connected", attempt: 0, error: undefined })
|
||||
handleEvent(first.value)
|
||||
await bootstrap()
|
||||
while (!controller.signal.aborted) {
|
||||
const event = await stream.next()
|
||||
if (event.done) return new Error("Event stream disconnected")
|
||||
handleEvent(event.value)
|
||||
}
|
||||
})().catch((error) => error)
|
||||
if (controller.signal.aborted) return
|
||||
setStore("connection", {
|
||||
status: "reconnecting",
|
||||
attempt: store.connection.status === "connected" ? 1 : store.connection.attempt + 1,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
})
|
||||
await wait(250, controller.signal)
|
||||
}
|
||||
})()
|
||||
})
|
||||
onCleanup(
|
||||
sdk.event.listen(({ details }) => {
|
||||
handleEvent(details)
|
||||
if (details.type === "server.connected") void bootstrap()
|
||||
}),
|
||||
)
|
||||
|
||||
return result
|
||||
},
|
||||
})
|
||||
|
||||
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()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,31 +1,27 @@
|
|||
import type { Event } from "@opencode-ai/sdk/v2"
|
||||
import type { V2Event } from "@opencode-ai/sdk/v2"
|
||||
import { useSDK } from "./sdk"
|
||||
|
||||
type EventMetadata = {
|
||||
directory: string
|
||||
directory: string | undefined
|
||||
workspace: string | undefined
|
||||
}
|
||||
|
||||
export function useEvent() {
|
||||
const sdk = useSDK()
|
||||
|
||||
function subscribe(handler: (event: Event, metadata: EventMetadata) => void) {
|
||||
return sdk.event.on("event", (event) => {
|
||||
if (event.payload.type === "sync") {
|
||||
return
|
||||
}
|
||||
|
||||
handler(event.payload, { directory: event.directory, workspace: event.workspace })
|
||||
function subscribe(handler: (event: V2Event, metadata: EventMetadata) => void) {
|
||||
return sdk.event.listen(({ details }) => {
|
||||
if (details.type === "server.connected") return
|
||||
handler(details, { directory: details.location?.directory, workspace: details.location?.workspaceID })
|
||||
})
|
||||
}
|
||||
|
||||
function on<T extends Event["type"]>(
|
||||
function on<T extends V2Event["type"]>(
|
||||
type: T,
|
||||
handler: (event: Extract<Event, { type: T }>, metadata: EventMetadata) => void,
|
||||
handler: (event: Extract<V2Event, { type: T }>, metadata: EventMetadata) => void,
|
||||
) {
|
||||
return subscribe((event: Event, metadata: EventMetadata) => {
|
||||
if (event.type !== type) return
|
||||
handler(event as Extract<Event, { type: T }>, metadata)
|
||||
return sdk.event.on(type, (event) => {
|
||||
handler(event, { directory: event.location?.directory, workspace: event.location?.workspaceID })
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -474,7 +474,7 @@ export const { use: useLocal, provider: LocalProvider } = createSimpleContext({
|
|||
}
|
||||
|
||||
event.on("session.deleted", (evt) => {
|
||||
prune(evt.properties.info.id)
|
||||
prune(evt.data.info.id)
|
||||
})
|
||||
|
||||
return {
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
import { batch } from "solid-js"
|
||||
import { batch, onCleanup } from "solid-js"
|
||||
import type { Path, Workspace } from "@opencode-ai/sdk/v2"
|
||||
import { createStore, reconcile } from "solid-js/store"
|
||||
import { createSimpleContext } from "./helper"
|
||||
|
|
@ -67,11 +67,11 @@ export const { use: useProject, provider: ProjectProvider } = createSimpleContex
|
|||
})
|
||||
}
|
||||
|
||||
sdk.event.on("event", (event) => {
|
||||
if (event.payload.type === "workspace.status") {
|
||||
setStore("workspace", "status", event.payload.properties.workspaceID, event.payload.properties.status)
|
||||
}
|
||||
})
|
||||
onCleanup(
|
||||
sdk.event.on("workspace.status", (event) => {
|
||||
setStore("workspace", "status", event.data.workspaceID, event.data.status)
|
||||
}),
|
||||
)
|
||||
|
||||
return {
|
||||
data: store,
|
||||
|
|
|
|||
|
|
@ -1,89 +1,145 @@
|
|||
import type { GlobalEvent, OpencodeClient } from "@opencode-ai/sdk/v2"
|
||||
import { Flag } from "@opencode-ai/core/flag/flag"
|
||||
import type { OpencodeClient, V2Event } from "@opencode-ai/sdk/v2"
|
||||
import { createGlobalEmitter } from "@solid-primitives/event-bus"
|
||||
import { onCleanup, onMount } from "solid-js"
|
||||
import { createStore } from "solid-js/store"
|
||||
import { createSimpleContext } from "./helper"
|
||||
import { batch, onCleanup, onMount } from "solid-js"
|
||||
|
||||
export type SDKConnectionStatus = "connecting" | "connected" | "reconnecting"
|
||||
|
||||
type SDKEventMap = { [Type in V2Event["type"]]: Extract<V2Event, { type: Type }> }
|
||||
const connectTimeout = 2_000
|
||||
|
||||
export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
||||
name: "SDK",
|
||||
init: (props: { client: OpencodeClient }) => {
|
||||
init: (props: { client: OpencodeClient; reload?: () => Promise<OpencodeClient> }) => {
|
||||
const abort = new AbortController()
|
||||
const handlers = new Set<(event: GlobalEvent) => void>()
|
||||
const emitter = {
|
||||
emit(_type: "event", event: GlobalEvent) {
|
||||
for (const handler of handlers) handler(event)
|
||||
},
|
||||
on(_type: "event", handler: (event: GlobalEvent) => void) {
|
||||
handlers.add(handler)
|
||||
return () => {
|
||||
handlers.delete(handler)
|
||||
}
|
||||
},
|
||||
}
|
||||
let client = props.client
|
||||
const events = createGlobalEmitter<SDKEventMap>()
|
||||
const [connection, setConnection] = createStore<{
|
||||
status: SDKConnectionStatus
|
||||
attempt: number
|
||||
error?: string
|
||||
}>({
|
||||
status: "connecting",
|
||||
attempt: 0,
|
||||
})
|
||||
let stream: AbortController | undefined
|
||||
let pending: Promise<void> | undefined
|
||||
|
||||
let queue: GlobalEvent[] = []
|
||||
let timer: Timer | undefined
|
||||
let last = 0
|
||||
const retryDelay = 1000
|
||||
const maxRetryDelay = 30000
|
||||
|
||||
const flush = () => {
|
||||
if (queue.length === 0) return
|
||||
const events = queue
|
||||
queue = []
|
||||
timer = undefined
|
||||
last = Date.now()
|
||||
batch(() => {
|
||||
for (const event of events) emitter.emit("event", event)
|
||||
function start() {
|
||||
stream?.abort()
|
||||
const controller = new AbortController()
|
||||
const current = client
|
||||
let connected!: () => void
|
||||
const ready = new Promise<void>((resolve) => {
|
||||
connected = resolve
|
||||
})
|
||||
}
|
||||
|
||||
const handleEvent = (event: GlobalEvent) => {
|
||||
queue.push(event)
|
||||
const elapsed = Date.now() - last
|
||||
if (timer) return
|
||||
if (elapsed < 16) {
|
||||
timer = setTimeout(flush, 16)
|
||||
return
|
||||
}
|
||||
flush()
|
||||
}
|
||||
|
||||
onMount(() => {
|
||||
stream = controller
|
||||
void (async () => {
|
||||
let attempt = 0
|
||||
while (!abort.signal.aborted) {
|
||||
const events = await props.client.global.event({
|
||||
signal: abort.signal,
|
||||
sseMaxRetryAttempts: 0,
|
||||
})
|
||||
|
||||
if (Flag.OPENCODE_EXPERIMENTAL_WORKSPACES) await props.client.sync.start().catch(() => {})
|
||||
|
||||
for await (const event of events.stream) {
|
||||
if (abort.signal.aborted) break
|
||||
handleEvent(event)
|
||||
}
|
||||
|
||||
if (timer) clearTimeout(timer)
|
||||
if (queue.length > 0) flush()
|
||||
while (!abort.signal.aborted && !controller.signal.aborted) {
|
||||
const connection = new AbortController()
|
||||
const cancel = () => connection.abort(controller.signal.reason)
|
||||
const timeout = setTimeout(() => connection.abort(new Error("Timed out connecting to server")), connectTimeout)
|
||||
controller.signal.addEventListener("abort", cancel, { once: true })
|
||||
const error = await (async () => {
|
||||
const response = await current.v2.event.subscribe({
|
||||
signal: connection.signal,
|
||||
sseMaxRetryAttempts: 0,
|
||||
throwOnError: true,
|
||||
})
|
||||
const iterator = response.stream[Symbol.asyncIterator]()
|
||||
const first = await iterator.next()
|
||||
if (abort.signal.aborted || controller.signal.aborted) return
|
||||
if (first.done)
|
||||
return connection.signal.reason instanceof Error
|
||||
? connection.signal.reason
|
||||
: new Error("Event stream disconnected")
|
||||
clearTimeout(timeout)
|
||||
attempt = 0
|
||||
setConnection({ status: "connected", attempt: 0, error: undefined })
|
||||
events.emit(first.value.type, first.value)
|
||||
connected()
|
||||
while (!abort.signal.aborted && !controller.signal.aborted) {
|
||||
const event = await iterator.next()
|
||||
if (abort.signal.aborted || controller.signal.aborted) return
|
||||
if (event.done) return new Error("Event stream disconnected")
|
||||
events.emit(event.value.type, event.value)
|
||||
}
|
||||
})()
|
||||
.catch((error) => error)
|
||||
.finally(() => {
|
||||
clearTimeout(timeout)
|
||||
controller.signal.removeEventListener("abort", cancel)
|
||||
})
|
||||
if (abort.signal.aborted || controller.signal.aborted) return
|
||||
attempt += 1
|
||||
if (abort.signal.aborted) break
|
||||
await new Promise((resolve) =>
|
||||
setTimeout(resolve, Math.min(retryDelay * 2 ** (attempt - 1), maxRetryDelay)),
|
||||
)
|
||||
setConnection({
|
||||
status: "reconnecting",
|
||||
attempt,
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
})
|
||||
await wait(250, controller.signal)
|
||||
}
|
||||
})().catch(() => {})
|
||||
})
|
||||
})()
|
||||
return ready
|
||||
}
|
||||
|
||||
const reload = props.reload
|
||||
? () => {
|
||||
if (pending) return pending
|
||||
pending = Promise.resolve()
|
||||
.then(props.reload)
|
||||
.then(async (next) => {
|
||||
client = next
|
||||
if (!abort.signal.aborted) await start()
|
||||
})
|
||||
.finally(() => {
|
||||
pending = undefined
|
||||
})
|
||||
return pending
|
||||
}
|
||||
: undefined
|
||||
|
||||
onMount(() => void start())
|
||||
onCleanup(() => {
|
||||
abort.abort()
|
||||
if (timer) clearTimeout(timer)
|
||||
handlers.clear()
|
||||
stream?.abort()
|
||||
events.clear()
|
||||
})
|
||||
|
||||
return {
|
||||
client: props.client,
|
||||
event: emitter,
|
||||
get client() {
|
||||
return client
|
||||
},
|
||||
event: {
|
||||
on: events.on,
|
||||
listen: events.listen,
|
||||
},
|
||||
connection: {
|
||||
status() {
|
||||
return connection.status
|
||||
},
|
||||
attempt() {
|
||||
return connection.attempt
|
||||
},
|
||||
error() {
|
||||
return connection.error
|
||||
},
|
||||
},
|
||||
reload,
|
||||
}
|
||||
},
|
||||
})
|
||||
|
||||
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()
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue