From d8970326be6f0352fa3099684da5f7922bf4fb56 Mon Sep 17 00:00:00 2001 From: devlikepro Date: Fri, 10 Jul 2026 12:46:26 +0700 Subject: [PATCH] [core] GOWS - fix "Session silently stops sending webhook events after (media ack)" - fix #2151 --- .../gows/GowsEventStreamObservable.test.ts | 167 ++++++++++++++++++ .../engines/gows/GowsEventStreamObservable.ts | 89 +++++++--- src/core/engines/gows/clients.ts | 4 + 3 files changed, 231 insertions(+), 29 deletions(-) create mode 100644 src/core/engines/gows/GowsEventStreamObservable.test.ts diff --git a/src/core/engines/gows/GowsEventStreamObservable.test.ts b/src/core/engines/gows/GowsEventStreamObservable.test.ts new file mode 100644 index 00000000..9961222d --- /dev/null +++ b/src/core/engines/gows/GowsEventStreamObservable.test.ts @@ -0,0 +1,167 @@ +import * as grpc from '@grpc/grpc-js'; +import { + GowsEventStreamObservable, + GowsStreamEndedError, +} from '@waha/core/engines/gows/GowsEventStreamObservable'; +import { EventEmitter } from 'events'; +import { merge, Subject } from 'rxjs'; +import { retry } from 'rxjs/operators'; + +/** + * Mimics grpc.ClientReadableStream. + * The important part is failWithStatus - it reproduces what grpc-js does in + * client.js makeServerStreamRequest.onReceiveStatus on a non OK status: + * push(null) schedules 'end' on the next tick, then 'error' is emitted + * synchronously in the current tick. + */ +class FakeStream extends EventEmitter { + public cancelled = false; + + cancel() { + this.cancelled = true; + } + + failWithStatus(err: any) { + process.nextTick(() => this.emit('end')); + this.emit('error', err); + } + + endCleanly() { + process.nextTick(() => this.emit('end')); + } +} + +class FakeClient { + public closed = false; + + close() { + this.closed = true; + } +} + +function buildLogger(): any { + function noop() { + return undefined; + } + return { + debug: noop, + info: noop, + warn: noop, + error: noop, + setBindings: noop, + }; +} + +function drainTicks(): Promise { + return new Promise((resolve) => setImmediate(resolve)); +} + +function observe(stream: FakeStream, client: FakeClient) { + const observable = new GowsEventStreamObservable(buildLogger(), () => ({ + client: client as any, + stream: stream as any, + })); + // Do not wait a real second for the client to close + observable.CLIENT_CLOSE_TIMEOUT = 0; + + const next = jest.fn(); + const error = jest.fn(); + const complete = jest.fn(); + const subscription = observable.subscribe({ + next: next, + error: error, + complete: complete, + }); + return { + next: next, + error: error, + complete: complete, + subscription: subscription, + }; +} + +describe('GowsEventStreamObservable', () => { + it('errors (not completes) when the stream fails with a non OK status', async () => { + const stream = new FakeStream(); + const client = new FakeClient(); + const { error, complete } = observe(stream, client); + + const err = { code: grpc.status.UNAVAILABLE, message: 'unavailable' }; + stream.failWithStatus(err); + await drainTicks(); + + expect(error).toHaveBeenCalledTimes(1); + expect(error).toHaveBeenCalledWith(err); + expect(complete).not.toHaveBeenCalled(); + expect(stream.cancelled).toBe(true); + expect(client.closed).toBe(true); + }); + + it('errors (not completes) when the stream ends cleanly', async () => { + const stream = new FakeStream(); + const client = new FakeClient(); + const { error, complete } = observe(stream, client); + + stream.endCleanly(); + await drainTicks(); + + expect(error).toHaveBeenCalledTimes(1); + expect(error.mock.calls[0][0]).toBeInstanceOf(GowsStreamEndedError); + expect(complete).not.toHaveBeenCalled(); + }); + + it('does not error when we cancel the stream ourselves', async () => { + const stream = new FakeStream(); + const client = new FakeClient(); + const { error, complete, subscription } = observe(stream, client); + + subscription.unsubscribe(); + stream.emit('error', { code: grpc.status.CANCELLED }); + stream.emit('end'); + await drainTicks(); + + expect(error).not.toHaveBeenCalled(); + expect(complete).not.toHaveBeenCalled(); + expect(stream.cancelled).toBe(true); + expect(client.closed).toBe(true); + }); + + it('reconnects via retry() when merged with a never ending subject', async () => { + const streams = [new FakeStream(), new FakeStream()]; + const clients = [new FakeClient(), new FakeClient()]; + let attempt = 0; + + const observable = new GowsEventStreamObservable(buildLogger(), () => { + const index = attempt; + attempt += 1; + return { client: clients[index] as any, stream: streams[index] as any }; + }); + observable.CLIENT_CLOSE_TIMEOUT = 0; + + const local$ = new Subject(); + const next = jest.fn(); + const error = jest.fn(); + const subscription = merge(observable, local$) + .pipe(retry({ delay: 1 })) + .subscribe({ next: next, error: error }); + + streams[0].failWithStatus({ code: grpc.status.INTERNAL }); + await drainTicks(); + await new Promise((resolve) => setTimeout(resolve, 20)); + + expect(attempt).toBe(2); + expect(error).not.toHaveBeenCalled(); + + streams[1].emit('data', { + toObject: () => ({ event: 'Message', data: '{"id":"1"}' }), + }); + await drainTicks(); + + expect(next).toHaveBeenCalledTimes(1); + expect(next.mock.calls[0][0]).toEqual({ + event: 'Message', + data: { id: '1' }, + }); + subscription.unsubscribe(); + }); +}); diff --git a/src/core/engines/gows/GowsEventStreamObservable.ts b/src/core/engines/gows/GowsEventStreamObservable.ts index 2c44d766..b4fa70de 100644 --- a/src/core/engines/gows/GowsEventStreamObservable.ts +++ b/src/core/engines/gows/GowsEventStreamObservable.ts @@ -1,14 +1,29 @@ import * as grpc from '@grpc/grpc-js'; +import { rand } from '@waha/core/auth/config'; import { messages } from '@waha/core/engines/gows/grpc/gows'; import { EnginePayload } from '@waha/structures/webhooks.dto'; import { sleep } from '@waha/utils/promiseTimeout'; import { Logger } from 'pino'; import { Observable } from 'rxjs'; -import { rand } from '@waha/core/auth/config'; + +/** + * Raised when the gRPC stream ends without an error. + * The engine event stream is expected to live as long as the session, + * so a clean end still means we lost the events and have to reconnect. + */ +export class GowsStreamEndedError extends Error { + constructor() { + super('gRPC event stream ended'); + this.name = 'GowsStreamEndedError'; + } +} /** * Observable that listens to a gRPC stream and emits EnginePayload objects. * Pass a factory function that returns a client and a stream. + * + * The observable always terminates with an error, never with a completion, + * so that an upstream retry() reconnects the stream. */ export class GowsEventStreamObservable extends Observable { _client: grpc.Client; @@ -26,63 +41,79 @@ export class GowsEventStreamObservable extends Observable { logger.setBindings({ id: rand() }); const { client, stream } = factory(); this._client = client; + const closeTimeout = this.CLIENT_CLOSE_TIMEOUT; let closed = false; - const cleanup = async (reason: string) => { + let terminated = false; + let tearingDown = false; + + async function cleanup(reason: string) { if (closed) { return; } closed = true; - logger.debug({ reason }, 'Cancelling gRPC stream...'); + logger.debug({ reason: reason }, 'Cancelling gRPC stream...'); try { stream.cancel(); } catch (err) { - logger.warn({ err }, 'Failed to cancel gRPC stream'); + logger.warn({ err: err }, 'Failed to cancel gRPC stream'); } - logger.debug({ reason }, 'Closing gRPC client...'); + logger.debug({ reason: reason }, 'Closing gRPC client...'); try { client.close(); } catch (err) { - logger.warn({ err }, 'Failed to close gRPC client'); + logger.warn({ err: err }, 'Failed to close gRPC client'); } - await sleep(this.CLIENT_CLOSE_TIMEOUT); - }; + await sleep(closeTimeout); + } + + // Must run synchronously from the stream handlers. + // grpc-js calls stream.push(null) - which schedules 'end' on the next tick - + // and only then emits 'error' in the same tick. Erroring the subscriber + // right away wins that race, otherwise 'end' completes the observable + // and the upstream retry() never reconnects. + function terminate(err: Error) { + if (terminated) { + return; + } + terminated = true; + // Erroring the subscriber runs the teardown below, which cleans up. + subscriber.error(err); + } stream.on('data', (raw) => { setImmediate(() => { const obj = raw.toObject(); obj.data = JSON.parse(obj.data); - subscriber?.next(obj); + subscriber.next(obj); }); }); - stream.on('end', (...args) => { - logger.debug('Stream ended', args); - subscriber?.complete(); - subscriber = null; - void cleanup('end'); - }); - - stream.on('error', async (err: any) => { - const CLIENT_CANCELLED_CODE = grpc.status.CANCELLED; - if (err.code === CLIENT_CANCELLED_CODE) { - logger.debug('Stream cancelled by client'); - await cleanup('cancelled'); + stream.on('end', () => { + if (tearingDown || terminated) { + logger.debug('Stream ended'); return; } - logger.error(err, 'Stream error'); - await cleanup('error'); - // Give some time to node event loop to process the error - await sleep(100); - subscriber?.error(err); - subscriber = null; + logger.error('Stream ended unexpectedly, reconnecting...'); + terminate(new GowsStreamEndedError()); }); - return async () => { - await cleanup('teardown'); + stream.on('error', (err: any) => { + if (tearingDown || terminated) { + // We cancelled the stream ourselves, no need to reconnect + logger.debug({ err: err }, 'Stream cancelled by client'); + return; + } + logger.error(err, 'Stream error, reconnecting...'); + terminate(err); + }); + + return () => { + tearingDown = true; + void cleanup('teardown'); }; }); } diff --git a/src/core/engines/gows/clients.ts b/src/core/engines/gows/clients.ts index 3cc5bbd4..2dff0bee 100644 --- a/src/core/engines/gows/clients.ts +++ b/src/core/engines/gows/clients.ts @@ -18,5 +18,9 @@ export function BuildEventStreamClient( return new messages.EventStreamClient(address, credentials, { 'grpc.max_send_message_length': 128 * 1024 * 1024, 'grpc.max_receive_message_length': 128 * 1024 * 1024, + // Detect a hung server, so the stream errors out and we reconnect + 'grpc.keepalive_time_ms': 30_000, + 'grpc.keepalive_timeout_ms': 10_000, + 'grpc.keepalive_permit_without_calls': 1, }); }