From 70a66292b288ab34dd6ad7173d4449ee275adb0d Mon Sep 17 00:00:00 2001 From: Phoenix-Forrest-Lin <113887049+PhoenixForrestLin@users.noreply.github.com> Date: Sun, 5 Apr 2026 01:14:24 +0800 Subject: [PATCH] fix(server-*,stage-ui): WebSocket cannot reconnect correctly, improved log (#1563) --- packages/server-runtime/src/index.test.ts | 15 +- packages/server-runtime/src/index.ts | 70 ++++++- packages/server-sdk/src/client.ts | 59 ++++++ packages/server-sdk/test/client.test.ts | 83 ++++++++ packages/server-sdk/tsconfig.json | 8 + .../stores/mods/api/channel-server.test.ts | 192 +++++++++++++++++- .../src/stores/mods/api/channel-server.ts | 81 ++++++-- .../src/stores/mods/api/context-bridge.ts | 25 ++- 8 files changed, 500 insertions(+), 33 deletions(-) diff --git a/packages/server-runtime/src/index.test.ts b/packages/server-runtime/src/index.test.ts index 5834db207..9a9a73afa 100644 --- a/packages/server-runtime/src/index.test.ts +++ b/packages/server-runtime/src/index.test.ts @@ -2,7 +2,7 @@ import type { WebSocketBaseEvent, WebSocketEvents } from '@proj-airi/server-shar import { describe, expect, it } from 'vitest' -import { resolveDeliveryConfig, selectConsumerPeerId } from './index' +import { detectHeartbeatControlFrame, resolveDeliveryConfig, selectConsumerPeerId } from './index' function createInputTextEvent( overrides: Partial> = {}, @@ -195,3 +195,16 @@ describe('selectConsumerPeerId', () => { expect(secondSelectedPeerId).toBe('stage-window-a') }) }) + +describe('detectHeartbeatControlFrame', () => { + it('recognizes raw websocket control frame text without treating it as protocol JSON', () => { + expect(detectHeartbeatControlFrame('ping')).toBe('ping') + expect(detectHeartbeatControlFrame('pong')).toBe('pong') + }) + + it('ignores non-control payloads', () => { + expect(detectHeartbeatControlFrame('')).toBeUndefined() + expect(detectHeartbeatControlFrame('🩵')).toBeUndefined() + expect(detectHeartbeatControlFrame('{"type":"transport:connection:heartbeat"}')).toBeUndefined() + }) +}) diff --git a/packages/server-runtime/src/index.ts b/packages/server-runtime/src/index.ts index 810ddc4bd..b303534db 100644 --- a/packages/server-runtime/src/index.ts +++ b/packages/server-runtime/src/index.ts @@ -127,6 +127,12 @@ function send(peer: Peer, event: WebSocketEvent> | strin peer.send(typeof event === 'string' ? event : stringify(event)) } +export function detectHeartbeatControlFrame(text: string): MessageHeartbeatKind | undefined { + if (text === MessageHeartbeatKind.Ping || text === MessageHeartbeatKind.Pong) { + return text + } +} + export function resolveDeliveryConfig(event: WebSocketEvent): DeliveryConfig | undefined { const eventMetadata = getProtocolEventMetadata(event.type) const defaultDelivery = eventMetadata?.delivery @@ -539,6 +545,32 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () => let event: WebSocketEvent try { + const text = message.text() + const controlFrame = detectHeartbeatControlFrame(text) + + // Some websocket runtimes surface control frames as plain text messages instead of + // exposing them through dedicated ping/pong hooks. Treat those payloads as transport + // liveness only so they do not leak into the application event protocol. + if (controlFrame) { + if (authenticatedPeer) { + authenticatedPeer.lastHeartbeatAt = Date.now() + authenticatedPeer.missedHeartbeats = 0 + + if (authenticatedPeer.healthy === false && authenticatedPeer.name && authenticatedPeer.identity) { + authenticatedPeer.healthy = true + logger.withFields({ peer: peer.id, peerName: authenticatedPeer.name }) + .debug('ping/pong recovered, marking healthy') + broadcastToAuthenticated({ + type: 'registry:modules:health:healthy', + data: { name: authenticatedPeer.name, index: authenticatedPeer.index, identity: authenticatedPeer.identity }, + metadata: createServerEventMetadata(instanceId), + }) + } + } + + return + } + // NOTICE: SDK clients send events using superjson.stringify, so we must use // superjson.parse here instead of message.json() (which uses JSON.parse). // Using JSON.parse on a superjson-encoded string returns the wrapper object @@ -547,7 +579,6 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () => // However, external clients may send plain JSON (not superjson-encoded). // superjson.parse on plain JSON returns undefined since there is no `json` wrapper key. // In that case, fall back to JSON.parse so external clients can interoperate. - const text = message.text() const parsed = parse(text) const potentialEvent = (parsed && typeof parsed === 'object' && 'type' in parsed) ? parsed @@ -881,10 +912,45 @@ export function setupApp(options?: AppOptions): { app: H3, closeAllPeers: () => }, close: (peer, details) => { const p = peers.get(peer.id) + const now = Date.now() + const safeDetails = details ?? {} + const closeCode = typeof safeDetails.code === 'number' ? safeDetails.code : undefined + const closeReason = typeof safeDetails.reason === 'string' ? safeDetails.reason : undefined + const closeWasClean = typeof (safeDetails as { wasClean?: unknown }).wasClean === 'boolean' + ? (safeDetails as { wasClean?: unknown }).wasClean + : undefined + const heartbeatLastSeenAt = p?.lastHeartbeatAt + const heartbeatSilentForMs = heartbeatLastSeenAt ? now - heartbeatLastSeenAt : undefined + const likelyHeartbeatExpiry = Boolean( + p + && typeof heartbeatSilentForMs === 'number' + && heartbeatSilentForMs > heartbeatTtlMs, + ) + const likelySilentNetworkClose = closeCode === 1005 + if (p) unregisterModulePeer(p, 'connection closed') - logger.withFields({ peer: peer.id, peerRemote: peer.remoteAddress, details, activePeers: peers.size }).log('closed') + logger.withFields({ + peer: peer.id, + peerRemote: peer.remoteAddress, + details, + closeCode, + closeReason, + closeWasClean, + activePeers: peers.size, + peerAuthenticated: p?.authenticated, + peerName: p?.name, + peerIndex: p?.index, + peerHealthy: p?.healthy, + peerMissedHeartbeats: p?.missedHeartbeats, + heartbeatLastSeenAt, + heartbeatSilentForMs, + heartbeatTtlMs, + healthCheckIntervalMs, + likelyHeartbeatExpiry, + likelySilentNetworkClose, + }).log('closed') peers.delete(peer.id) }, })) diff --git a/packages/server-sdk/src/client.ts b/packages/server-sdk/src/client.ts index ec5a96244..511fda904 100644 --- a/packages/server-sdk/src/client.ts +++ b/packages/server-sdk/src/client.ts @@ -49,6 +49,7 @@ export interface ClientOptions { name: string token?: string websocketConstructor?: WebSocketLikeConstructor + connectTimeoutMs?: number possibleEvents?: Array> identity?: MetadataEventSource @@ -148,6 +149,7 @@ export class Client { this.opts = { url: 'ws://localhost:6121/ws', + connectTimeoutMs: 15_000, onAnyMessage: () => {}, onAnySend: () => {}, possibleEvents: [], @@ -372,6 +374,19 @@ export class Client { this.connectionAttempt = attempt const isCurrentSocket = () => this.websocket === ws + const connectTimeoutMs = this.opts.connectTimeoutMs + const connectTimer = setTimeout(() => { + if (ws.readyState === WebSocket.OPEN) { + return + } + + ws.close() + deferred.reject(new Error(`Connection timeout after ${connectTimeoutMs}ms`)) + }, connectTimeoutMs) + + const clearConnectTimer = () => { + clearTimeout(connectTimer) + } ws.onmessage = (event: WebSocketMessageEventLike) => { if (!isCurrentSocket()) { @@ -382,6 +397,8 @@ export class Client { } ws.onerror = (event: any) => { + clearConnectTimer() + if (!isCurrentSocket()) { return } @@ -397,6 +414,8 @@ export class Client { } ws.onclose = () => { + clearConnectTimer() + if (!isCurrentSocket()) { return } @@ -411,6 +430,7 @@ export class Client { if (wasReady && this.opts.autoReconnect) { this.pendingReconnect = true + this.transitionTo('idle') void this.connect() return } @@ -419,6 +439,8 @@ export class Client { } ws.onopen = () => { + clearConnectTimer() + if (!isCurrentSocket()) { return } @@ -671,6 +693,42 @@ export class Client { return } + if (this.status === 'ready') { + return + } + + if (this.connectionAttempt) { + this.connectionAttempt.announced = true + } + + this.reconnectAttempts = 0 + this.transitionTo('ready') + this.resolveAttempt() + this.opts.onReady?.() + return + } + + case 'registry:modules:sync': { + // Fallback: If the status is stuck at 'announcing' but the sync already contains this module, + // it means the announce succeeded; the server simply didn't send back 'module:announced' + if (this.status !== 'announcing' || !this.connectionAttempt) { + return + } + + const modules = (data.data as any)?.modules as Array<{ + name: string + identity?: { id?: string } + }> ?? [] + + const selfRegistered = modules.some( + m => m.name === this.opts.name + && m.identity?.id === this.identity.id, + ) + + if (!selfRegistered) { + return + } + if (this.connectionAttempt) { this.connectionAttempt.announced = true } @@ -829,6 +887,7 @@ export class Client { return } + this.transitionTo('idle') void this.connect() } } diff --git a/packages/server-sdk/test/client.test.ts b/packages/server-sdk/test/client.test.ts index 10a3b2b03..06e3b84ba 100644 --- a/packages/server-sdk/test/client.test.ts +++ b/packages/server-sdk/test/client.test.ts @@ -299,4 +299,87 @@ describe('client', () => { dispose() }) + + it('retries after connect timeout and eventually connects on a later socket', async () => { + vi.useFakeTimers() + + const client = new Client({ + autoConnect: false, + autoReconnect: true, + connectTimeoutMs: 50, + name: 'test-plugin', + }) + + const connecting = client.connect() + const firstSocket = lastSocket() + const firstCloseSpy = vi.spyOn(firstSocket, 'close') + + await vi.advanceTimersByTimeAsync(50) + expect(firstCloseSpy).toHaveBeenCalledTimes(1) + + await vi.advanceTimersByTimeAsync(1_000) + expect(MockWebSocket.instances).toHaveLength(2) + + const secondSocket = lastSocket() + emitOpen(secondSocket) + const announceEvent = parseSent(secondSocket) + + emitMessage(secondSocket, { + type: 'module:announced', + data: { + name: 'test-plugin', + identity: announceEvent.data.identity, + }, + metadata: { + source: { kind: 'plugin', plugin: { id: 'server' }, id: 'server-1' }, + event: { id: 'announce-retry-1' }, + }, + }) + + await expect(connecting).resolves.toBeUndefined() + expect(client.connectionStatus).toBe('ready') + }) + + it('does not emit onReady twice when sync fallback already moved status to ready', async () => { + const onReady = vi.fn() + const client = new Client({ + autoConnect: false, + autoReconnect: false, + name: 'test-plugin', + onReady, + }) + + const connecting = client.connect() + const socket = lastSocket() + emitOpen(socket) + + const announceEvent = parseSent(socket) + const selfIdentity = announceEvent.data.identity + + emitMessage(socket, { + type: 'registry:modules:sync', + data: { + modules: [{ name: 'test-plugin', identity: selfIdentity }], + }, + metadata: { + source: { kind: 'plugin', plugin: { id: 'server' }, id: 'server-1' }, + event: { id: 'sync-1' }, + }, + }) + + emitMessage(socket, { + type: 'module:announced', + data: { + name: 'test-plugin', + identity: selfIdentity, + }, + metadata: { + source: { kind: 'plugin', plugin: { id: 'server' }, id: 'server-1' }, + event: { id: 'announce-1' }, + }, + }) + + await expect(connecting).resolves.toBeUndefined() + expect(onReady).toHaveBeenCalledTimes(1) + }) }) diff --git a/packages/server-sdk/tsconfig.json b/packages/server-sdk/tsconfig.json index 00dcfd807..07486031b 100644 --- a/packages/server-sdk/tsconfig.json +++ b/packages/server-sdk/tsconfig.json @@ -6,6 +6,14 @@ ], "module": "ESNext", "moduleResolution": "bundler", + "paths": { + "@proj-airi/server-shared": [ + "../server-shared/src/index.ts" + ], + "@proj-airi/server-shared/*": [ + "../server-shared/src/*" + ] + }, "esModuleInterop": true, "forceConsistentCasingInFileNames": true, "isolatedModules": true, diff --git a/packages/stage-ui/src/stores/mods/api/channel-server.test.ts b/packages/stage-ui/src/stores/mods/api/channel-server.test.ts index d74de0427..52948e533 100644 --- a/packages/stage-ui/src/stores/mods/api/channel-server.test.ts +++ b/packages/stage-ui/src/stores/mods/api/channel-server.test.ts @@ -48,8 +48,8 @@ const serverSdkMocks = vi.hoisted(() => { return true } - close() { - this.options.onClose?.() + close(code?: number, reason?: string) { + this.options.onClose?.(code, reason) } emit(type: string, data: any) { @@ -64,12 +64,24 @@ const serverSdkMocks = vi.hoisted(() => { } simulateTransientDisconnect() { - this.options.onClose?.() + this.options.onClose?.(1005, '') + } + + simulateClose(code?: number, reason?: string) { + this.options.onClose?.(code, reason) } simulateReconnectReady() { this.options.onReady?.() } + + simulateError(error: unknown) { + this.options.onError?.(error) + } + + simulateStateChange(previousStatus: string, status: string) { + this.options.onStateChange?.({ previousStatus, status }) + } } return { @@ -162,4 +174,178 @@ describe('channel-server store reconnect', () => { }), ])) }) + + it('uses explicit heartbeat settings to avoid client/server timeout mismatch', async () => { + const store = useModsServerChannelStore() + + const initializePromise = store.initialize({ token: 'secret' }) + const client = serverSdkMocks.MockClient.instances[0] + + client.simulateAuthenticated() + await initializePromise + + expect(client.options.heartbeat).toEqual({ + readTimeout: 60_000, + pingInterval: 20_000, + }) + }) + + it('notifies onReconnected callbacks when the websocket becomes ready again', async () => { + const store = useModsServerChannelStore() + const onReconnected = vi.fn() + store.onReconnected(onReconnected) + + const initializePromise = store.initialize({ token: 'secret' }) + const client = serverSdkMocks.MockClient.instances[0] + + client.simulateAuthenticated() + await initializePromise + + client.simulateTransientDisconnect() + client.simulateReconnectReady() + client.simulateReconnectReady() + + expect(onReconnected).toHaveBeenCalledTimes(1) + }) + + it('does not notify onReconnected on first authenticated->ready flow and only on subsequent ready events', async () => { + const store = useModsServerChannelStore() + const onReconnected = vi.fn() + store.onReconnected(onReconnected) + + const initializePromise = store.initialize({ token: 'secret' }) + const client = serverSdkMocks.MockClient.instances[0] + + client.simulateAuthenticated() + client.simulateReconnectReady() + expect(onReconnected).toHaveBeenCalledTimes(0) + + client.simulateReconnectReady() + expect(onReconnected).toHaveBeenCalledTimes(1) + + await initializePromise + }) + + it('continues invoking remaining onReconnected callbacks when one throws', async () => { + const store = useModsServerChannelStore() + const successfulCallback = vi.fn() + const consoleErrorSpy = vi.spyOn(console, 'error').mockImplementation(() => {}) + + store.onReconnected(() => { + throw new Error('boom') + }) + store.onReconnected(successfulCallback) + + const initializePromise = store.initialize({ token: 'secret' }) + const client = serverSdkMocks.MockClient.instances[0] + + client.simulateAuthenticated() + await initializePromise + + client.simulateTransientDisconnect() + client.simulateReconnectReady() + client.simulateReconnectReady() + + expect(successfulCallback).toHaveBeenCalledTimes(1) + expect(consoleErrorSpy).toHaveBeenCalledTimes(1) + + consoleErrorSpy.mockRestore() + }) + + it('allows initialize retry after first handshake close before any successful connection', async () => { + const store = useModsServerChannelStore() + + const firstInitializePromise = store.initialize({ token: 'invalid-token' }) + const firstClient = serverSdkMocks.MockClient.instances[0] + + firstClient.simulateClose(1008, 'invalid token') + + const secondInitializePromise = store.initialize({ token: 'valid-token' }) + const secondClient = serverSdkMocks.MockClient.instances[1] + + expect(secondInitializePromise).not.toBe(firstInitializePromise) + expect(secondClient).toBeDefined() + + secondClient.simulateAuthenticated() + await secondInitializePromise + + expect(store.connected).toBe(true) + }) + + it('allows initialize retry when sdk enters failed after a previous successful connection', async () => { + const store = useModsServerChannelStore() + + const firstInitializePromise = store.initialize({ token: 'secret' }) + const firstClient = serverSdkMocks.MockClient.instances[0] + + firstClient.simulateAuthenticated() + await firstInitializePromise + + firstClient.simulateStateChange('reconnecting', 'failed') + + const secondInitializePromise = store.initialize({ token: 'secret-rotated' }) + const secondClient = serverSdkMocks.MockClient.instances[1] + + expect(secondInitializePromise).not.toBe(firstInitializePromise) + expect(secondClient).toBeDefined() + + secondClient.simulateAuthenticated() + await secondInitializePromise + + expect(store.connected).toBe(true) + }) + + it('keeps the initialize lock on recoverable onError so auto-reconnect does not spawn a second client', () => { + const store = useModsServerChannelStore() + const consoleDebugSpy = vi.spyOn(console, 'debug').mockImplementation(() => {}) + + const firstInitializePromise = store.initialize({ token: 'secret' }) + const firstClient = serverSdkMocks.MockClient.instances[0] + + firstClient.simulateError(new Error('temporary websocket glitch')) + + const secondInitializePromise = store.initialize({ token: 'secret' }) + + expect(secondInitializePromise).toBeInstanceOf(Promise) + expect(firstInitializePromise).toBeInstanceOf(Promise) + expect(serverSdkMocks.MockClient.instances).toHaveLength(1) + + consoleDebugSpy.mockRestore() + }) + + it('does not flush queued events on reconnect authenticated before ready', async () => { + const store = useModsServerChannelStore() + + const initializePromise = store.initialize({ token: 'secret' }) + const client = serverSdkMocks.MockClient.instances[0] + + client.simulateAuthenticated() + await initializePromise + client.simulateReconnectReady() + + client.simulateTransientDisconnect() + + store.send({ + type: 'spark:notify', + data: { message: 'reconnect-authenticated-queued' }, + } as any) + + expect(store.pendingSendCount).toBe(1) + + client.simulateAuthenticated() + + expect(store.connected).toBe(false) + expect(store.pendingSendCount).toBe(1) + + client.simulateReconnectReady() + + expect(store.connected).toBe(true) + expect(store.pendingSendCount).toBe(0) + expect(client.sent).toEqual(expect.arrayContaining([ + expect.objectContaining({ + type: 'spark:notify', + data: { message: 'reconnect-authenticated-queued' }, + }), + ])) + }) }) diff --git a/packages/stage-ui/src/stores/mods/api/channel-server.ts b/packages/stage-ui/src/stores/mods/api/channel-server.ts index f0d22f1de..a4ed25be7 100644 --- a/packages/stage-ui/src/stores/mods/api/channel-server.ts +++ b/packages/stage-ui/src/stores/mods/api/channel-server.ts @@ -37,8 +37,10 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se const client = ref() const initializing = ref | null>(null) const websocketConstructor = ref() + const hasEverConnected = ref(false) const pendingSend = ref>([]) const pendingSendCount = computed(() => pendingSend.value.length) + const reconnectedCallbacks = new Set<() => void>() const defaultWebSocketUrl = import.meta.env.VITE_AIRI_WS_URL || 'ws://localhost:6121/ws' const websocketUrl = useLocalStorage('settings/connection/websocket-url', defaultWebSocketUrl) @@ -94,6 +96,11 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se url: websocketUrl.value || defaultWebSocketUrl, token: options?.token, websocketConstructor: websocketConstructor.value, + heartbeat: { + // Keep client and server heartbeat windows aligned to reduce false-positive disconnects. + readTimeout: 60_000, + pingInterval: 20_000, + }, possibleEvents, onAnyMessage: (event) => { if (REPLAYABLE_EVENT_TYPES.has(event.type as keyof WebSocketEvents)) @@ -106,37 +113,67 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se }, onError: (error) => { connected.value = false - initializing.value = null - clearListeners() - replayableEvents.clear() - - console.warn('WebSocket server connection error:', error) + // Do not clear listeners or replay cache here. + // onError may be recoverable while the SDK is reconnecting. + if (import.meta.env.DEV) { + // eslint-disable-next-line no-console + console.debug('WebSocket server connection error:', error) + } }, onClose: () => { connected.value = false - initializing.value = null - clearListeners() - replayableEvents.clear() - console.warn('WebSocket server connection closed') + if (!hasEverConnected.value) { + // First handshake failed: clear lock so initialize() can be retried externally. + initializing.value = null + } + // Runtime disconnect: keep initialize/listeners for SDK auto-reconnect. + // Terminal failure: handled by onStateChange status === 'failed'. + }, + onStateChange: ({ status }) => { + if (status === 'failed') { + // SDK entered terminal state (auth terminal / retries exhausted / autoReconnect disabled). + connected.value = false + initializing.value = null + console.warn('WebSocket server connection failed') + } }, onReady: () => { + const isReconnect = hasEverConnected.value + + hasEverConnected.value = true connected.value = true flush() initializeListeners() + + if (isReconnect) { + for (const callback of reconnectedCallbacks) { + try { + callback() + } + catch (error) { + console.error('Error in reconnected callback:', error) + } + } + } + if (isReconnect && import.meta.env.DEV) { + // eslint-disable-next-line no-console + console.debug('WebSocket server connection re-established') + } }, }) client.value.onEvent('module:authenticated', (event) => { if (event.data.authenticated) { - connected.value = true - flush() - initializeListeners() + if (!hasEverConnected.value) { + // First connection can flush immediately after authentication. + connected.value = true + flush() + initializeListeners() + } + // On reconnect, wait for onReady (after announce) before flushing business events. resolve() - // eslint-disable-next-line no-console - console.log('WebSocket server connection established and authenticated') - return } @@ -236,6 +273,14 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se return registerListener(type, callback) } + function onReconnected(callback: () => void) { + reconnectedCallbacks.add(callback) + + return () => { + reconnectedCallbacks.delete(callback) + } + } + function sendContextUpdate(message: InputContextUpdate) { const id = nanoid() send({ @@ -246,6 +291,9 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se function dispose() { flush() + hasEverConnected.value = false + connected.value = false + initializing.value = null clearListeners() replayableEvents.clear() @@ -253,8 +301,6 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se client.value.close() client.value = undefined } - connected.value = false - initializing.value = null } watch(websocketUrl, (newUrl, oldUrl) => { @@ -277,6 +323,7 @@ export const useModsServerChannelStore = defineStore('mods:channels:proj-airi:se sendContextUpdate, onContextUpdate, onEvent, + onReconnected, getPendingSendSnapshot: () => [...pendingSend.value], dispose, } diff --git a/packages/stage-ui/src/stores/mods/api/context-bridge.ts b/packages/stage-ui/src/stores/mods/api/context-bridge.ts index 08b4753df..081f222c6 100644 --- a/packages/stage-ui/src/stores/mods/api/context-bridge.ts +++ b/packages/stage-ui/src/stores/mods/api/context-bridge.ts @@ -63,18 +63,23 @@ export const useContextBridgeStore = defineStore('mods:api:context-bridge', () = await mutex.acquire() try { + const registerConsumers = () => { + for (const consumerEvent of consumerRegistrationEvents) { + serverChannelStore.send({ + type: 'module:consumer:register', + data: { + event: consumerEvent, + mode: 'consumer-group', + group: 'chat-ingestion', + }, + }) + } + } + await serverChannelStore.ensureConnected() - for (const consumerEvent of consumerRegistrationEvents) { - serverChannelStore.send({ - type: 'module:consumer:register', - data: { - event: consumerEvent, - mode: 'consumer-group', - group: 'chat-ingestion', - }, - }) - } + registerConsumers() + disposeHookFns.value.push(serverChannelStore.onReconnected(() => registerConsumers())) let isProcessingRemoteStream = false