[core] GOWS - fix "Session silently stops sending webhook events after <stream:error> (media ack)" - fix #2151
This commit is contained in:
1 parent
06412c2c1e
commit
d8970326be
3 files changed
+231
-29
No files matched your search
@@ -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<void> {
|
||||
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<any>();
|
||||
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();
|
||||
});
|
||||
});
|
||||
@@ -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<EnginePayload> {
|
||||
_client: grpc.Client;
|
||||
@@ -26,63 +41,79 @@ export class GowsEventStreamObservable extends Observable<EnginePayload> {
|
||||
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');
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
});
|
||||
}
|
||||
Reference in new issue
Block a user