diff --git a/apps/webclaw/src/routeTree.gen.ts b/apps/webclaw/src/routeTree.gen.ts index 97661c6..9fc7637 100644 --- a/apps/webclaw/src/routeTree.gen.ts +++ b/apps/webclaw/src/routeTree.gen.ts @@ -13,6 +13,7 @@ import { Route as NewRouteImport } from './routes/new' import { Route as ConnectRouteImport } from './routes/connect' import { Route as IndexRouteImport } from './routes/index' import { Route as ChatSessionKeyRouteImport } from './routes/chat/$sessionKey' +import { Route as ApiStreamRouteImport } from './routes/api/stream' import { Route as ApiSessionsRouteImport } from './routes/api/sessions' import { Route as ApiSendRouteImport } from './routes/api/send' import { Route as ApiPingRouteImport } from './routes/api/ping' @@ -39,6 +40,11 @@ const ChatSessionKeyRoute = ChatSessionKeyRouteImport.update({ path: '/chat/$sessionKey', getParentRoute: () => rootRouteImport, } as any) +const ApiStreamRoute = ApiStreamRouteImport.update({ + id: '/api/stream', + path: '/api/stream', + getParentRoute: () => rootRouteImport, +} as any) const ApiSessionsRoute = ApiSessionsRouteImport.update({ id: '/api/sessions', path: '/api/sessions', @@ -74,6 +80,7 @@ export interface FileRoutesByFullPath { '/api/ping': typeof ApiPingRoute '/api/send': typeof ApiSendRoute '/api/sessions': typeof ApiSessionsRoute + '/api/stream': typeof ApiStreamRoute '/chat/$sessionKey': typeof ChatSessionKeyRoute } export interface FileRoutesByTo { @@ -85,6 +92,7 @@ export interface FileRoutesByTo { '/api/ping': typeof ApiPingRoute '/api/send': typeof ApiSendRoute '/api/sessions': typeof ApiSessionsRoute + '/api/stream': typeof ApiStreamRoute '/chat/$sessionKey': typeof ChatSessionKeyRoute } export interface FileRoutesById { @@ -97,6 +105,7 @@ export interface FileRoutesById { '/api/ping': typeof ApiPingRoute '/api/send': typeof ApiSendRoute '/api/sessions': typeof ApiSessionsRoute + '/api/stream': typeof ApiStreamRoute '/chat/$sessionKey': typeof ChatSessionKeyRoute } export interface FileRouteTypes { @@ -110,6 +119,7 @@ export interface FileRouteTypes { | '/api/ping' | '/api/send' | '/api/sessions' + | '/api/stream' | '/chat/$sessionKey' fileRoutesByTo: FileRoutesByTo to: @@ -121,6 +131,7 @@ export interface FileRouteTypes { | '/api/ping' | '/api/send' | '/api/sessions' + | '/api/stream' | '/chat/$sessionKey' id: | '__root__' @@ -132,6 +143,7 @@ export interface FileRouteTypes { | '/api/ping' | '/api/send' | '/api/sessions' + | '/api/stream' | '/chat/$sessionKey' fileRoutesById: FileRoutesById } @@ -144,6 +156,7 @@ export interface RootRouteChildren { ApiPingRoute: typeof ApiPingRoute ApiSendRoute: typeof ApiSendRoute ApiSessionsRoute: typeof ApiSessionsRoute + ApiStreamRoute: typeof ApiStreamRoute ChatSessionKeyRoute: typeof ChatSessionKeyRoute } @@ -177,6 +190,13 @@ declare module '@tanstack/react-router' { preLoaderRoute: typeof ChatSessionKeyRouteImport parentRoute: typeof rootRouteImport } + '/api/stream': { + id: '/api/stream' + path: '/api/stream' + fullPath: '/api/stream' + preLoaderRoute: typeof ApiStreamRouteImport + parentRoute: typeof rootRouteImport + } '/api/sessions': { id: '/api/sessions' path: '/api/sessions' @@ -224,6 +244,7 @@ const rootRouteChildren: RootRouteChildren = { ApiPingRoute: ApiPingRoute, ApiSendRoute: ApiSendRoute, ApiSessionsRoute: ApiSessionsRoute, + ApiStreamRoute: ApiStreamRoute, ChatSessionKeyRoute: ChatSessionKeyRoute, } export const routeTree = rootRouteImport diff --git a/apps/webclaw/src/routes/api/send.ts b/apps/webclaw/src/routes/api/send.ts index b05f09a..5243ba1 100644 --- a/apps/webclaw/src/routes/api/send.ts +++ b/apps/webclaw/src/routes/api/send.ts @@ -1,7 +1,7 @@ import { randomUUID } from 'node:crypto' import { createFileRoute } from '@tanstack/react-router' import { json } from '@tanstack/react-start' -import { gatewayRpc } from '../../server/gateway' +import { gatewayRpc, gatewayRpcShared } from '../../server/gateway' type SessionsResolveResponse = { ok?: boolean @@ -70,18 +70,22 @@ export const Route = createFileRoute('/api/send')({ sessionKey = 'main' } - const res = await gatewayRpc<{ runId: string }>('chat.send', { + const res = await gatewayRpcShared<{ runId: string }>( + 'chat.send', + { sessionKey, message, thinking, attachments, - deliver: false, + deliver: true, timeoutMs: 120_000, idempotencyKey: typeof body.idempotencyKey === 'string' ? body.idempotencyKey : randomUUID(), - }) + }, + sessionKey, + ) return json({ ok: true, ...res, sessionKey }) } catch (err) { diff --git a/apps/webclaw/src/routes/api/stream.ts b/apps/webclaw/src/routes/api/stream.ts new file mode 100644 index 0000000..ffba067 --- /dev/null +++ b/apps/webclaw/src/routes/api/stream.ts @@ -0,0 +1,109 @@ +import { createFileRoute } from '@tanstack/react-router' +import { acquireGatewayClient, gatewayRpcShared } from '../../server/gateway' + +type StreamEventPayload = { + event: string + payload?: unknown + seq?: number + stateVersion?: number +} + +export const Route = createFileRoute('/api/stream')({ + server: { + handlers: { + GET: ({ request }) => { + const url = new URL(request.url) + const sessionKey = url.searchParams.get('sessionKey')?.trim() || '' + const friendlyId = url.searchParams.get('friendlyId')?.trim() || '' + const encoder = new TextEncoder() + + let releaseClient: (() => void) | null = null + let closed = false + + const stream = new ReadableStream({ + start(controller) { + function send(data: StreamEventPayload) { + if (closed) return + try { + controller.enqueue( + encoder.encode(`data: ${JSON.stringify(data)}\n\n`), + ) + } catch { + closed = true + } + } + + const heartbeat = setInterval(() => { + controller.enqueue(encoder.encode('event: ping\ndata: {}\n\n')) + }, 15000) + + const key = sessionKey || friendlyId + if (key) { + void acquireGatewayClient(key, { + onEvent(event) { + send({ + event: event.event, + payload: event.payload, + seq: event.seq, + stateVersion: event.stateVersion, + }) + }, + onError(error) { + send({ event: 'error', payload: error.message }) + }, + }) + .then((handle) => { + if (closed) { + handle.release() + return + } + releaseClient = handle.release + if (sessionKey) { + void gatewayRpcShared( + 'chat.history', + { sessionKey, limit: 1 }, + sessionKey, + ) + } + }) + .catch((error: unknown) => { + if (closed) return + const message = error instanceof Error ? error.message : String(error) + send({ event: 'error', payload: message }) + }) + } + + request.signal.addEventListener( + 'abort', + () => { + if (closed) return + closed = true + clearInterval(heartbeat) + releaseClient?.() + try { + controller.close() + } catch { + return + } + }, + { once: true }, + ) + }, + cancel() { + if (closed) return + closed = true + releaseClient?.() + }, + }) + + return new Response(stream, { + headers: { + 'Content-Type': 'text/event-stream', + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + }, + }) + }, + }, + }, +}) diff --git a/apps/webclaw/src/screens/chat/chat-screen.tsx b/apps/webclaw/src/screens/chat/chat-screen.tsx index 6137d7c..c6bcc06 100644 --- a/apps/webclaw/src/screens/chat/chat-screen.tsx +++ b/apps/webclaw/src/screens/chat/chat-screen.tsx @@ -11,6 +11,7 @@ import { useQuery, useQueryClient } from '@tanstack/react-query' import { deriveFriendlyIdFromKey, + getMessageTimestamp, isMissingGatewayAuth, readError, textFromMessage, @@ -22,6 +23,7 @@ import { clearHistoryMessages, fetchGatewayStatus, removeHistoryMessageByClientId, + updateHistoryMessages, updateHistoryMessageByClientId, updateSessionLastMessage, } from './chat-queries' @@ -47,7 +49,7 @@ import { useChatMobile } from './hooks/use-chat-mobile' import { useChatSessions } from './hooks/use-chat-sessions' import type { AttachmentFile } from '@/components/attachment-button' import type { ChatComposerHelpers } from './components/chat-composer' -import type { HistoryResponse } from './types' +import type { GatewayMessage, HistoryResponse } from './types' import { useExport } from '@/hooks/use-export' import { cn } from '@/lib/utils' @@ -81,8 +83,15 @@ export function ChatScreen({ const [pinToTop, setPinToTop] = useState( () => hasPendingSend() || hasPendingGeneration(), ) - const streamTimer = useRef(null) const streamIdleTimer = useRef(null) + const streamRefetchInFlight = useRef(false) + const lastStreamStateVersion = useRef(null) + const lastStreamSeq = useRef(null) + const streamSourceRef = useRef(null) + const streamReconnectTimer = useRef(null) + const streamReconnectAttempt = useRef(0) + const streamFinalRefetchTimer = useRef(null) + const lastStreamFinalRunId = useRef('') const lastAssistantSignature = useRef('') const refreshHistoryRef = useRef<() => void>(() => {}) const pendingStartRef = useRef(false) @@ -157,32 +166,242 @@ export function ChatScreen({ navigate({ to: '/new', replace: true }) }, [navigate]) const streamStop = useCallback(() => { - if (streamTimer.current) { - window.clearInterval(streamTimer.current) - streamTimer.current = null - } if (streamIdleTimer.current) { window.clearTimeout(streamIdleTimer.current) streamIdleTimer.current = null } + if (streamReconnectTimer.current) { + window.clearTimeout(streamReconnectTimer.current) + streamReconnectTimer.current = null + } + if (streamFinalRefetchTimer.current) { + window.clearTimeout(streamFinalRefetchTimer.current) + streamFinalRefetchTimer.current = null + } + if (streamSourceRef.current) { + streamSourceRef.current.close() + streamSourceRef.current = null + } + streamRefetchInFlight.current = false }, []) const streamFinish = useCallback(() => { streamStop() setPendingGeneration(false) setWaitingForResponse(false) }, [streamStop]) - const streamStart = useCallback(() => { - if (!activeFriendlyId || isNewChat) return - if (streamTimer.current) window.clearInterval(streamTimer.current) - streamTimer.current = window.setInterval(() => { - refreshHistoryRef.current() - }, 350) - }, [activeFriendlyId, isNewChat]) const stableContentStyle = useMemo(() => ({}), []) refreshHistoryRef.current = function refreshHistory() { void historyQuery.refetch() } + useEffect(function setupStream() { + if (!activeFriendlyId || isNewChat || isRedirecting) return + let cancelled = false + + function startStream() { + if (cancelled) return + if (streamSourceRef.current) { + streamSourceRef.current.close() + streamSourceRef.current = null + } + const params = new URLSearchParams() + const streamSessionKey = resolvedSessionKey || sessionKeyForHistory + if (streamSessionKey) params.set('sessionKey', streamSessionKey) + if (activeFriendlyId) params.set('friendlyId', activeFriendlyId) + const source = new EventSource(`/api/stream?${params.toString()}`) + streamSourceRef.current = source + + function handleStreamEvent(event: MessageEvent) { + try { + const parsed = JSON.parse(String(event.data || '{}')) as { + event?: string + payload?: unknown + seq?: number + stateVersion?: number + } + if ( + typeof parsed.stateVersion === 'number' && + parsed.stateVersion === lastStreamStateVersion.current + ) { + return + } + if (typeof parsed.stateVersion === 'number') { + lastStreamStateVersion.current = parsed.stateVersion + } + if (typeof parsed.seq === 'number') { + if (parsed.seq === lastStreamSeq.current) return + lastStreamSeq.current = parsed.seq + } + + if (parsed.event === 'chat.history') { + const payload = parsed.payload as { messages?: Array } | null + if (payload && Array.isArray(payload.messages)) { + console.info('[stream] apply history payload', { + count: payload.messages.length, + }) + queryClient.setQueryData( + chatQueryKeys.history(activeFriendlyId, sessionKeyForHistory), + { + sessionKey: sessionKeyForHistory, + messages: payload.messages, + }, + ) + return + } + return + } + if (!parsed.event) return + if (parsed.event === 'chat') { + const payload = parsed.payload as + | { + runId?: string + sessionKey?: string + state?: string + message?: GatewayMessage + } + | null + if (payload?.message && typeof payload.message === 'object') { + const payloadSessionKey = payload.sessionKey + if ( + payloadSessionKey && + resolvedSessionKey && + payloadSessionKey !== resolvedSessionKey && + payloadSessionKey !== sessionKeyForHistory + ) { + return + } + const streamRunId = + typeof payload.runId === 'string' ? payload.runId : '' + const nextMessage = { + ...payload.message, + __streamRunId: streamRunId || undefined, + } + function upsert(messages: Array) { + if (streamRunId) { + const index = messages.findIndex( + (message) => + (message as { __streamRunId?: string }).__streamRunId === + streamRunId, + ) + if (index >= 0) { + const next = [...messages] + next[index] = nextMessage + return next + } + } + if (nextMessage.role === 'assistant') { + const nextTime = getMessageTimestamp(nextMessage) + const index = [...messages] + .reverse() + .findIndex((message) => message.role === 'assistant') + if (index >= 0) { + const target = messages.length - 1 - index + const targetTime = getMessageTimestamp(messages[target]) + if (Math.abs(nextTime - targetTime) <= 15000) { + const next = [...messages] + next[target] = nextMessage + return next + } + } + } + return [...messages, nextMessage] + } + + updateHistoryMessages( + queryClient, + activeFriendlyId, + sessionKeyForHistory, + upsert, + ) + if (payloadSessionKey && payloadSessionKey !== sessionKeyForHistory) { + updateHistoryMessages( + queryClient, + activeFriendlyId, + payloadSessionKey, + upsert, + ) + } + if (payloadSessionKey) { + updateSessionLastMessage( + queryClient, + payloadSessionKey, + activeFriendlyId, + nextMessage, + ) + } + if (payload.state === 'final') { + const nextRunId = streamRunId || 'final' + if (lastStreamFinalRunId.current !== nextRunId) { + lastStreamFinalRunId.current = nextRunId + if (streamFinalRefetchTimer.current) { + window.clearTimeout(streamFinalRefetchTimer.current) + } + streamFinalRefetchTimer.current = window.setTimeout(() => { + streamFinalRefetchTimer.current = null + refreshHistoryRef.current() + }, 350) + } + } + } + return + } + if (!parsed.event.startsWith('chat.')) { + return + } + } catch { + // ignore parse errors + } + if (streamRefetchInFlight.current) return + streamRefetchInFlight.current = true + const refetchStart = performance.now() + Promise.resolve(refreshHistoryRef.current()).finally(() => { + streamRefetchInFlight.current = false + void refetchStart + }) + } + + function handleStreamOpen() { + streamReconnectAttempt.current = 0 + console.info('[stream] open') + refreshHistoryRef.current() + } + + function handleStreamError() { + console.info('[stream] error') + if (cancelled) return + if (streamReconnectTimer.current) return + if (streamSourceRef.current) { + streamSourceRef.current.close() + streamSourceRef.current = null + } + streamReconnectAttempt.current += 1 + const backoff = Math.min(8000, 1000 * streamReconnectAttempt.current) + streamReconnectTimer.current = window.setTimeout(() => { + streamReconnectTimer.current = null + startStream() + }, backoff) + } + + source.addEventListener('message', handleStreamEvent) + source.addEventListener('open', handleStreamOpen) + source.addEventListener('error', handleStreamError) + } + + startStream() + + return () => { + cancelled = true + streamStop() + } + }, [ + activeFriendlyId, + isNewChat, + isRedirecting, + resolvedSessionKey, + sessionKeyForHistory, + streamStop, + ]) + useEffect(() => { if (isRedirecting) { if (error) setError(null) @@ -279,7 +498,7 @@ export function ChatScreen({ } streamIdleTimer.current = window.setTimeout(() => { streamFinish() - }, 4000) + }, 12000) } }, [historyMessages, streamFinish]) @@ -406,7 +625,7 @@ export function ChatScreen({ }) .then(async (res) => { if (!res.ok) throw new Error(await readError(res)) - streamStart() + refreshHistoryRef.current() }) .catch((err) => { const messageText = err instanceof Error ? err.message : String(err) diff --git a/apps/webclaw/src/server/gateway.ts b/apps/webclaw/src/server/gateway.ts index e6d6e5e..ccbd5f6 100644 --- a/apps/webclaw/src/server/gateway.ts +++ b/apps/webclaw/src/server/gateway.ts @@ -10,7 +10,13 @@ type GatewayFrame = payload?: unknown error?: { code: string; message: string; details?: unknown } } - | { type: 'event'; event: string; payload?: unknown; seq?: number } + | { + type: 'event' + event: string + payload?: unknown + seq?: number + stateVersion?: number + } type ConnectParams = { minProtocol: number @@ -33,6 +39,44 @@ type GatewayWaiter = { handleMessage: (evt: MessageEvent) => void } +type GatewayEventFrame = { + type: 'event' + event: string + payload?: unknown + seq?: number + stateVersion?: number +} + +type GatewayEventStreamOptions = { + sessionKey?: string + friendlyId?: string + signal?: AbortSignal + onEvent: (event: GatewayEventFrame) => void + onError?: (error: Error) => void +} + +type GatewayClient = { + connect: () => Promise + sendReq: (method: string, params?: unknown) => Promise + close: () => void + setOnEvent: (handler?: (event: GatewayEventFrame) => void) => void + setOnError: (handler?: (error: Error) => void) => void + isClosed: () => boolean +} + +type GatewayClientEntry = { + key: string + refs: number + client: GatewayClient +} + +type GatewayClientHandle = { + client: GatewayClient + release: () => void +} + +const sharedGatewayClients = new Map() + function getGatewayConfig() { const url = process.env.CLAWDBOT_GATEWAY_URL?.trim() || 'ws://127.0.0.1:18789' const token = process.env.CLAWDBOT_GATEWAY_TOKEN?.trim() || '' @@ -69,6 +113,293 @@ function buildConnectParams(token: string, password: string): ConnectParams { } } +async function connectGateway(ws: WebSocket): Promise { + const { token, password } = getGatewayConfig() + await wsOpen(ws) + const connectId = randomUUID() + const connectParams = buildConnectParams(token, password) + const connectReq: GatewayFrame = { + type: 'req', + id: connectId, + method: 'connect', + params: connectParams, + } + const waiter = createGatewayWaiter() + ws.addEventListener('message', waiter.handleMessage) + ws.send(JSON.stringify(connectReq)) + await waiter.waitForRes(connectId) + ws.removeEventListener('message', waiter.handleMessage) +} + +function createGatewayClient(): GatewayClient { + const { url, token, password } = getGatewayConfig() + const ws = new WebSocket(url) + let closed = false + let connected = false + let onEvent: ((event: GatewayEventFrame) => void) | undefined + let onError: ((error: Error) => void) | undefined + const waiters = new Map< + string, + { + resolve: (v: unknown) => void + reject: (e: Error) => void + } + >() + + function rejectAll(error: Error) { + for (const [, waiter] of waiters) { + waiter.reject(error) + } + waiters.clear() + } + + function handleMessage(evt: MessageEvent) { + try { + const data = typeof evt.data === 'string' ? evt.data : '' + const parsed = JSON.parse(data) as GatewayFrame + if (parsed.type === 'event') { + if (onEvent) onEvent(parsed) + return + } + if (parsed.type !== 'res') return + const waiter = waiters.get(parsed.id) + if (!waiter) return + waiters.delete(parsed.id) + if (parsed.ok) waiter.resolve(parsed.payload) + else waiter.reject(new Error(parsed.error?.message ?? 'gateway error')) + } catch { + // ignore parse errors + } + } + + function handleError(err: Event) { + if (onError) { + onError( + new Error(`Gateway client error: ${String((err as any)?.message ?? err)}`), + ) + } + } + + function handleClose() { + if (closed) return + closed = true + rejectAll(new Error('Gateway client closed')) + } + + ws.addEventListener('message', handleMessage) + ws.addEventListener('error', handleError) + ws.addEventListener('close', handleClose) + + async function connect() { + if (connected || closed) return + await wsOpen(ws) + const connectId = randomUUID() + const connectParams = buildConnectParams(token, password) + const connectReq: GatewayFrame = { + type: 'req', + id: connectId, + method: 'connect', + params: connectParams, + } + const waitForRes = new Promise((resolve, reject) => { + waiters.set(connectId, { resolve, reject }) + }) + ws.send(JSON.stringify(connectReq)) + await waitForRes + connected = true + } + + function sendReq(method: string, params?: unknown) { + if (closed) { + return Promise.reject(new Error('Gateway client closed')) + } + const id = randomUUID() + const req: GatewayFrame = { + type: 'req', + id, + method, + params, + } + const waitForRes = new Promise((resolve, reject) => { + waiters.set(id, { resolve, reject }) + }) + ws.send(JSON.stringify(req)) + return waitForRes as Promise + } + + function close() { + if (closed) return + closed = true + ws.removeEventListener('message', handleMessage) + ws.removeEventListener('error', handleError) + ws.removeEventListener('close', handleClose) + rejectAll(new Error('Gateway client closed')) + void wsClose(ws) + } + + function setOnEvent(handler?: (event: GatewayEventFrame) => void) { + onEvent = handler + } + + function setOnError(handler?: (error: Error) => void) { + onError = handler + } + + function isClosed() { + return closed + } + + return { connect, sendReq, close, setOnEvent, setOnError, isClosed } +} + +export async function acquireGatewayClient( + key: string, + options?: { + onEvent?: (event: GatewayEventFrame) => void + onError?: (error: Error) => void + }, +): Promise { + const existing = sharedGatewayClients.get(key) + if (existing && !existing.client.isClosed()) { + existing.refs += 1 + if (options?.onEvent) existing.client.setOnEvent(options.onEvent) + if (options?.onError) existing.client.setOnError(options.onError) + return { + client: existing.client, + release: function release() { + releaseGatewayClient(key) + }, + } + } + + const client = createGatewayClient() + if (options?.onEvent) client.setOnEvent(options.onEvent) + if (options?.onError) client.setOnError(options.onError) + await client.connect() + sharedGatewayClients.set(key, { key, refs: 1, client }) + return { + client, + release: function release() { + releaseGatewayClient(key) + }, + } +} + +function releaseGatewayClient(key: string) { + const entry = sharedGatewayClients.get(key) + if (!entry) return + entry.refs -= 1 + if (entry.refs > 0) return + entry.client.close() + sharedGatewayClients.delete(key) +} + +export async function gatewayRpcShared( + method: string, + params: unknown, + key?: string, +): Promise { + if (key) { + const entry = sharedGatewayClients.get(key) + if (entry && !entry.client.isClosed()) { + await entry.client.connect() + return entry.client.sendReq(method, params) + } + } + return gatewayRpc(method, params) +} + +export function gatewayEventStream({ + sessionKey, + friendlyId, + signal, + onEvent, + onError, +}: GatewayEventStreamOptions) { + const { url } = getGatewayConfig() + const ws = new WebSocket(url) + let closed = false + + function handleMessage(evt: MessageEvent) { + try { + const data = typeof evt.data === 'string' ? evt.data : '' + const parsed = JSON.parse(data) as GatewayFrame + if (parsed.type !== 'event') return + onEvent(parsed) + } catch { + // ignore parse errors + } + } + + function handleError(err: Event) { + if (onError) { + onError( + new Error(`Gateway event stream error: ${String((err as any)?.message ?? err)}`), + ) + } + } + + function handleClose() { + if (closed) return + closed = true + } + + ws.addEventListener('message', handleMessage) + ws.addEventListener('error', handleError) + ws.addEventListener('close', handleClose) + + void connectGateway(ws) + .then(async () => { + if (!sessionKey && !friendlyId) return + const subscribeReq: GatewayFrame = { + type: 'req', + id: randomUUID(), + method: 'chat.subscribe', + params: { + sessionKey: sessionKey || undefined, + friendlyId: friendlyId || undefined, + }, + } + const waiter = createGatewayWaiter() + ws.addEventListener('message', waiter.handleMessage) + try { + ws.send(JSON.stringify(subscribeReq)) + await waiter.waitForRes(subscribeReq.id) + } catch (err) { + if (onError) { + onError(err instanceof Error ? err : new Error(String(err))) + } + close() + } finally { + ws.removeEventListener('message', waiter.handleMessage) + } + }) + .catch((err) => { + if (onError) onError(err instanceof Error ? err : new Error(String(err))) + }) + + if (signal) { + signal.addEventListener( + 'abort', + () => { + close() + }, + { once: true }, + ) + } + + function close() { + if (closed) return + closed = true + ws.removeEventListener('message', handleMessage) + ws.removeEventListener('error', handleError) + ws.removeEventListener('close', handleClose) + void wsClose(ws) + } + + return close +} + function createGatewayWaiter(): GatewayWaiter { const waiters = new Map< string, diff --git a/lefthook.yml b/lefthook.yml index d8139d3..61039fe 100644 --- a/lefthook.yml +++ b/lefthook.yml @@ -3,5 +3,5 @@ pre-commit: commands: eslint-webclaw: glob: "apps/webclaw/**/*.{js,jsx,ts,tsx}" - run: "pnpm -C apps/webclaw eslint --fix {staged_files}" + run: "pnpm -C apps/webclaw exec eslint --fix ." stage_fixed: true