mirror of
https://github.com/ibelick/webclaw.git
synced 2026-08-14 00:57:51 +00:00
fix: improve chat streaming UX and live updates
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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',
|
||||
},
|
||||
})
|
||||
},
|
||||
},
|
||||
},
|
||||
})
|
||||
@@ -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<number | null>(null)
|
||||
const streamIdleTimer = useRef<number | null>(null)
|
||||
const streamRefetchInFlight = useRef(false)
|
||||
const lastStreamStateVersion = useRef<number | null>(null)
|
||||
const lastStreamSeq = useRef<number | null>(null)
|
||||
const streamSourceRef = useRef<EventSource | null>(null)
|
||||
const streamReconnectTimer = useRef<number | null>(null)
|
||||
const streamReconnectAttempt = useRef(0)
|
||||
const streamFinalRefetchTimer = useRef<number | null>(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<React.CSSProperties>(() => ({}), [])
|
||||
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<unknown> } | 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<GatewayMessage>) {
|
||||
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)
|
||||
|
||||
@@ -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<void>
|
||||
sendReq: <TPayload = unknown>(method: string, params?: unknown) => Promise<TPayload>
|
||||
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<string, GatewayClientEntry>()
|
||||
|
||||
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<void> {
|
||||
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<unknown>((resolve, reject) => {
|
||||
waiters.set(connectId, { resolve, reject })
|
||||
})
|
||||
ws.send(JSON.stringify(connectReq))
|
||||
await waitForRes
|
||||
connected = true
|
||||
}
|
||||
|
||||
function sendReq<TPayload = unknown>(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<unknown>((resolve, reject) => {
|
||||
waiters.set(id, { resolve, reject })
|
||||
})
|
||||
ws.send(JSON.stringify(req))
|
||||
return waitForRes as Promise<TPayload>
|
||||
}
|
||||
|
||||
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<GatewayClientHandle> {
|
||||
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<TPayload = unknown>(
|
||||
method: string,
|
||||
params: unknown,
|
||||
key?: string,
|
||||
): Promise<TPayload> {
|
||||
if (key) {
|
||||
const entry = sharedGatewayClients.get(key)
|
||||
if (entry && !entry.client.isClosed()) {
|
||||
await entry.client.connect()
|
||||
return entry.client.sendReq<TPayload>(method, params)
|
||||
}
|
||||
}
|
||||
return gatewayRpc<TPayload>(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,
|
||||
|
||||
+1
-1
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user