diff --git a/packages/tui/src/context/client.tsx b/packages/tui/src/context/client.tsx index 6afee05f8e..a417d70be0 100644 --- a/packages/tui/src/context/client.tsx +++ b/packages/tui/src/context/client.tsx @@ -49,43 +49,57 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( if (history.length > connectionHistoryLimit) history.shift() } - async function runConnection(signal: AbortSignal, attempt: number) { + 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) - throw request.signal.reason instanceof Error - ? request.signal.reason - : new Error("Event stream disconnected") + 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") - throw new Error("Event stream did not start with 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") events.emit(first.value.type, 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) throw new Error("Event stream disconnected") + 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, }) + events.emit(event.value.type, event.value) } + return { error: undefined, connectedAt } } catch (error) { return { error, connectedAt } @@ -103,7 +117,7 @@ export const { use: useClient, provider: ClientProvider } = createSimpleContext( void (async () => { let attempt = 0 while (!abort.signal.aborted && !controller.signal.aborted) { - const result = await runConnection(controller.signal, attempt) + 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