Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 21 additions & 3 deletions packages/tui/src/context/client.tsx
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import type { OpenCodeClient, OpenCodeEvent } from "@opencode-ai/client"
import { createGlobalEmitter } from "@solid-primitives/event-bus"
import { onCleanup, onMount } from "solid-js"
import { batch, onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store"
import { errorMessage } from "../util/error"
import { createSimpleContext } from "./helper"
Expand All @@ -25,6 +25,7 @@ type ManagedService = {
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",
Expand All @@ -34,6 +35,8 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
const history: ClientConnectionEvent[] = []
let api = props.api
const events = createGlobalEmitter<ClientEventMap>()
let pending: OpenCodeEvent[] = []
let flushTimer: ReturnType<typeof setTimeout> | undefined
const [connection, setConnection] = createStore<{
status: ClientConnectionStatus
attempt: number
Expand All @@ -49,6 +52,19 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
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

Expand Down Expand Up @@ -80,7 +96,7 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
record("connected", attempt)
connectedAt = Date.now()
log.info("event stream connected")
events.emit(first.value.type, first.value)
emit(first.value)
setConnection({ status: "connected", attempt: 0, error: undefined })

// Forward events until the stream closes or this connection is cancelled.
Expand All @@ -97,7 +113,7 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
seq: event.value.durable.seq,
})

events.emit(event.value.type, event.value)
emit(event.value)
}

return { error: undefined, connectedAt }
Expand Down Expand Up @@ -154,6 +170,8 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext(
onCleanup(() => {
abort.abort()
stream?.abort()
if (flushTimer) clearTimeout(flushTimer)
pending = []
events.clear()
})

Expand Down
Loading
Loading