117 lines
2.9 KiB
TypeScript
117 lines
2.9 KiB
TypeScript
import { createOpencodeClient } from "@opencode-ai/sdk/v2"
|
|
import type { GlobalEvent, Event } from "@opencode-ai/sdk/v2"
|
|
import { createSimpleContext } from "./helper"
|
|
import { createGlobalEmitter } from "@solid-primitives/event-bus"
|
|
import { batch, onCleanup, onMount } from "solid-js"
|
|
|
|
export type EventSource = {
|
|
subscribe: (handler: (event: GlobalEvent) => void) => Promise<() => void>
|
|
}
|
|
|
|
export const { use: useSDK, provider: SDKProvider } = createSimpleContext({
|
|
name: "SDK",
|
|
init: (props: {
|
|
url: string
|
|
directory?: string
|
|
fetch?: typeof fetch
|
|
headers?: RequestInit["headers"]
|
|
events?: EventSource
|
|
}) => {
|
|
const abort = new AbortController()
|
|
let sse: AbortController | undefined
|
|
|
|
function createSDK() {
|
|
return createOpencodeClient({
|
|
baseUrl: props.url,
|
|
signal: abort.signal,
|
|
directory: props.directory,
|
|
fetch: props.fetch,
|
|
headers: props.headers,
|
|
})
|
|
}
|
|
|
|
let sdk = createSDK()
|
|
|
|
const emitter = createGlobalEmitter<{
|
|
event: GlobalEvent
|
|
}>()
|
|
|
|
let queue: GlobalEvent[] = []
|
|
let timer: Timer | undefined
|
|
let last = 0
|
|
|
|
const flush = () => {
|
|
if (queue.length === 0) return
|
|
const events = queue
|
|
queue = []
|
|
timer = undefined
|
|
last = Date.now()
|
|
// Batch all event emissions so all store updates result in a single render
|
|
batch(() => {
|
|
for (const event of events) {
|
|
emitter.emit("event", event)
|
|
}
|
|
})
|
|
}
|
|
|
|
const handleEvent = (event: GlobalEvent) => {
|
|
queue.push(event)
|
|
const elapsed = Date.now() - last
|
|
|
|
if (timer) return
|
|
// If we just flushed recently (within 16ms), batch this with future events
|
|
// Otherwise, process immediately to avoid latency
|
|
if (elapsed < 16) {
|
|
timer = setTimeout(flush, 16)
|
|
return
|
|
}
|
|
flush()
|
|
}
|
|
|
|
function startSSE() {
|
|
sse?.abort()
|
|
const ctrl = new AbortController()
|
|
sse = ctrl
|
|
;(async () => {
|
|
while (true) {
|
|
if (abort.signal.aborted || ctrl.signal.aborted) break
|
|
const events = await sdk.global.event({ signal: ctrl.signal })
|
|
|
|
for await (const event of events.stream) {
|
|
if (ctrl.signal.aborted) break
|
|
handleEvent(event)
|
|
}
|
|
|
|
if (timer) clearTimeout(timer)
|
|
if (queue.length > 0) flush()
|
|
}
|
|
})().catch(() => {})
|
|
}
|
|
|
|
onMount(async () => {
|
|
if (props.events) {
|
|
const unsub = await props.events.subscribe(handleEvent)
|
|
onCleanup(unsub)
|
|
} else {
|
|
startSSE()
|
|
}
|
|
})
|
|
|
|
onCleanup(() => {
|
|
abort.abort()
|
|
sse?.abort()
|
|
if (timer) clearTimeout(timer)
|
|
})
|
|
|
|
return {
|
|
get client() {
|
|
return sdk
|
|
},
|
|
directory: props.directory,
|
|
event: emitter,
|
|
fetch: props.fetch ?? fetch,
|
|
url: props.url,
|
|
}
|
|
},
|
|
})
|