Compare commits

...

5 Commits

Author SHA1 Message Date
Brendan Allan d33039614f refactor(app): use shared client connection (#43001) 2026-08-17 15:39:20 +08:00
Brendan Allan 01ce44bd13 remove test changes 2026-08-17 15:23:40 +08:00
Brendan Allan 647ae85a04 rename to CreateData from CreateServerData 2026-08-17 15:07:31 +08:00
Brendan Allan a1e018668c refactor(client): remove event coalescing 2026-08-17 14:52:48 +08:00
Brendan Allan d5451cdabe refactor(client): share Solid server data 2026-08-17 14:40:46 +08:00
9 changed files with 1806 additions and 1989 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": {
+1 -58
View File
@@ -1,18 +1,6 @@
import { describe, expect, test } from "bun:test"
import type { OpenCodeEvent } from "@opencode-ai/client/promise"
import { adaptServerEvent, coalesceServerEvents, resumeStreamAfterPageShow } from "./server-sdk"
describe("resumeStreamAfterPageShow", () => {
test("restarts a stream only after a back-forward cache restore", () => {
let starts = 0
const start = () => starts++
resumeStreamAfterPageShow({ persisted: false } as PageTransitionEvent, start)
resumeStreamAfterPageShow({ persisted: true } as PageTransitionEvent, start)
expect(starts).toBe(1)
})
})
import { adaptServerEvent } from "./server-sdk"
describe("adaptServerEvent", () => {
test("preserves current permission requests", () => {
@@ -43,48 +31,3 @@ describe("adaptServerEvent", () => {
})
})
})
describe("current event buffering", () => {
const delta = (id: string, value: string, ordinal = 0) =>
adaptServerEvent({
id,
created: 1,
type: "session.text.delta",
location: { directory: "/repo" },
data: { sessionID: "ses", assistantMessageID: "msg", ordinal, delta: value },
} as OpenCodeEvent)
test("merges adjacent text deltas for the same message and ordinal", () => {
const result = coalesceServerEvents([delta("evt_1", "hello "), delta("evt_2", "world")])
expect(result).toHaveLength(1)
expect(result[0]?.current).toMatchObject({ id: "evt_2", data: { delta: "hello world" } })
expect(result[0]?.properties).toMatchObject({ delta: "hello world" })
})
test("coalesces current tool input deltas by tool ID", () => {
const current = (eventID: string, id: string, delta: string) =>
adaptServerEvent({
id: eventID,
created: 1,
type: "session.tool.input.delta",
location: { directory: "/repo" },
data: { sessionID: "ses", assistantMessageID: "msg", id, delta },
} as OpenCodeEvent)
const result = coalesceServerEvents([
current("evt_1", "call_1", "{"),
current("evt_2", "call_1", "}"),
current("evt_3", "call_2", "[]"),
])
expect(result).toHaveLength(2)
expect(result[0]?.current).toMatchObject({ id: "evt_2", data: { id: "call_1", delta: "{}" } })
expect(result[1]?.current).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(coalesceServerEvents(events).map((event) => event.current?.id)).toEqual(["evt_1", "evt_2", "evt_3"])
})
})
+21 -274
View File
@@ -1,9 +1,8 @@
import type { OpenCodeEvent } from "@opencode-ai/client/promise"
import { createClientConnection, type ClientConnectionStatus } from "@opencode-ai/client/solid"
import type { Event } from "@/types"
import { createGlobalEmitter } from "@solid-primitives/event-bus"
import { makeEventListener } from "@solid-primitives/event-listener"
import { type Accessor, batch, onCleanup, onMount } from "solid-js"
import { createStore } from "solid-js/store"
import { type Accessor, onCleanup } from "solid-js"
import { createApiForServer, type ServerApi } from "@/utils/server"
import { usePlatform } from "./platform"
import { ServerConnection } from "./servers"
@@ -12,78 +11,15 @@ import { ServerScope } from "@/utils/server-scope"
import { useServer } from "./server"
export type ServerEvent = Event & { id?: string; current?: OpenCodeEvent }
type ServerEventMap = { [Type in ServerEvent["type"]]: Extract<ServerEvent, { type: Type }> }
type CurrentDelta = Extract<
OpenCodeEvent,
{ type: "session.text.delta" | "session.reasoning.delta" | "session.tool.input.delta" | "session.compaction.delta" }
>
export function adaptServerEvent(event: OpenCodeEvent): ServerEvent {
return { id: event.id, type: event.type, properties: event.data, current: event } as ServerEvent
}
export function coalesceServerEvents(events: ServerEvent[]) {
const output: ServerEvent[] = []
events.forEach((event) => {
const current = currentDelta(event.current)
if (current) {
const previous = output[output.length - 1]
const prior = currentDelta(previous?.current)
if (
previous &&
prior &&
prior.location?.directory === current.location?.directory &&
currentDeltaKey(prior) === currentDeltaKey(current)
) {
const fragment = currentDeltaFragment(prior) + currentDeltaFragment(current)
const data =
current.type === "session.compaction.delta"
? { ...current.data, text: fragment }
: { ...current.data, delta: fragment }
output[output.length - 1] = {
...event,
properties: data,
current: { ...current, data } as CurrentDelta,
} as ServerEvent
return
}
output.push(event)
return
}
output.push(event)
})
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 resumeStreamAfterPageShow(event: PageTransitionEvent, start: () => unknown) {
if (!event.persisted) return
start()
}
type ServerEventMap = { [Type in ServerEvent["type"]]: Extract<ServerEvent, { type: Type }> }
type ServerEventEmitter = ReturnType<typeof createGlobalEmitter<ServerEventMap>>
type ServerLocationEventEmitter = ReturnType<typeof createGlobalEmitter<{ [directory: string]: ServerEvent }>>
export type ServerConnectionStatus = "connecting" | "connected" | "reconnecting"
export type ServerConnectionStatus = ClientConnectionStatus
type ServerSDKBase = {
server: ServerConnection.Any
scope: ServerScope
@@ -105,227 +41,38 @@ type ServerSDKBase = {
function createServerSdkContextBase(server: ServerConnection.Any, scope: ServerScope): ServerSDKBase {
const platform = usePlatform()
const abort = new AbortController()
const eventFetch = (() => {
if (!platform.fetch || !server) return
try {
const url = new URL(server.http.url)
const loopback = url.hostname === "localhost" || url.hostname === "127.0.0.1" || url.hostname === "::1"
if (url.protocol === "http:" && !loopback) return platform.fetch
} catch {
return
}
})()
const eventApi = createApiForServer({ server: server.http, fetch: eventFetch })
const api = createApiForServer({ server: server.http, fetch: platform.fetch })
const emitter = createGlobalEmitter<ServerEventMap>()
const locations = createGlobalEmitter<{ [directory: string]: ServerEvent }>()
const FLUSH_FRAME_MS = 16
const STREAM_YIELD_MS = 8
const CONNECT_TIMEOUT_MS = 2_000
const RECONNECT_DELAY_MS = 1_000
let queue: ServerEvent[] = []
let buffer: ServerEvent[] = []
let timer: ReturnType<typeof setTimeout> | undefined
let last = 0
function flush() {
if (timer) clearTimeout(timer)
timer = undefined
if (queue.length === 0) return
const events = queue
queue = buffer
buffer = events
queue.length = 0
last = Date.now()
const output = coalesceServerEvents(events)
batch(() => {
output.forEach((event) => {
emitter.emit(event.type, event)
const directory = event.current?.location?.directory
if (directory) locations.emit(directory, event)
})
})
buffer.length = 0
}
function schedule() {
if (timer) return
const elapsed = Date.now() - last
timer = setTimeout(flush, Math.max(0, FLUSH_FRAME_MS - elapsed))
}
function publish(event: OpenCodeEvent) {
queue.push(adaptServerEvent(event))
schedule()
}
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()
}
})
}
let attempt: AbortController | undefined
let run: Promise<void> | undefined
let started = false
let generation = 0
const [connection, setConnection] = createStore<{
status: ServerConnectionStatus
attempt: number
error?: string
}>({ status: "connecting", attempt: 0 })
async function connect(signal: AbortSignal): Promise<{ error: unknown; connectedAt: number | undefined }> {
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")), CONNECT_TIMEOUT_MS)
signal.addEventListener("abort", cancel, { once: true })
try {
// Open the event stream and validate its initial handshake.
const iterator = eventApi.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)
publish(first.value)
connectedAt = Date.now()
setConnection({ status: "connected", attempt: 0, error: undefined })
// Forward events until the stream closes or this connection is cancelled.
let yielded = Date.now()
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 }
publish(event.value)
if (Date.now() - yielded < STREAM_YIELD_MS) continue
yielded = Date.now()
await wait(0, signal)
}
return { error: undefined, connectedAt }
} catch (error) {
return { error, connectedAt }
} finally {
request.abort()
clearTimeout(timeout)
signal.removeEventListener("abort", cancel)
}
}
async function runStream(active: number) {
let retries = 0
// oxlint-disable-next-line no-unmodified-loop-condition -- stop() changes the lifecycle flags and aborts the active request
while (!abort.signal.aborted && started && generation === active) {
setConnection({ status: retries === 0 ? "connecting" : "reconnecting", attempt: retries, error: undefined })
const controller = new AbortController()
attempt = controller
const onAbort = () => controller.abort()
abort.signal.addEventListener("abort", onAbort)
const result = await connect(controller.signal)
abort.signal.removeEventListener("abort", onAbort)
if (abort.signal.aborted || !started || generation !== active) {
if (attempt === controller) attempt = undefined
return
}
if (result.connectedAt !== undefined && Date.now() - result.connectedAt >= 1_000) retries = 0
retries += 1
const message =
result.error === undefined
? undefined
: result.error instanceof Error
? result.error.message
: String(result.error)
console.info("[global-sdk] event stream disconnected", {
url: server.http.url,
fetch: eventFetch ? "platform" : "webview",
attempt: retries,
error: message,
})
setConnection({ status: "reconnecting", attempt: retries, error: message })
await wait(RECONNECT_DELAY_MS, controller.signal)
if (attempt === controller) attempt = undefined
}
}
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
flush()
})
run = current
return run
}
function stop() {
started = false
generation++
attempt?.abort()
}
onMount(() => {
makeEventListener(window, "pagehide", stop)
makeEventListener(window, "pageshow", (event) => resumeStreamAfterPageShow(event, start))
void start()
const connection = createClientConnection(api, {
flushInterval: 16,
pageLifecycle: true,
onEvent(event) {
const adapted = adaptServerEvent(event)
emitter.emit(adapted.type, adapted)
const directory = event.location?.directory
if (directory) locations.emit(directory, adapted)
},
log: {
info(message, data) {
if (message !== "event stream disconnected") return
console.info("[global-sdk] event stream disconnected", { url: server.http.url, ...data })
},
},
})
onCleanup(() => {
stop()
abort.abort()
if (timer) clearTimeout(timer)
timer = undefined
queue = []
buffer = []
emitter.clear()
locations.clear()
})
const api = createApiForServer({ server: server.http, fetch: platform.fetch })
return {
server,
scope,
url: server.http.url,
api,
connection: {
status: () => connection.status,
attempt: () => connection.attempt,
error: () => connection.error,
},
connection,
event: {
on: emitter.on.bind(emitter),
listen: emitter.listen.bind(emitter),
+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:"
}
}
+217
View File
@@ -0,0 +1,217 @@
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
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(() => 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 })
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"
+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