diff --git a/graphql/server-test/__tests__/notify-listener.test.ts b/graphql/server-test/__tests__/notify-listener.test.ts new file mode 100644 index 000000000..5958f9284 --- /dev/null +++ b/graphql/server-test/__tests__/notify-listener.test.ts @@ -0,0 +1,88 @@ +import { Server } from '@constructive-io/graphql-server'; +import { EventEmitter } from 'events'; +import type { PoolClient } from 'pg'; + +class FakeClient extends EventEmitter { + queries: string[] = []; + queryResult: Promise = Promise.resolve({ rows: [] }); + + query(sql: string): Promise { + this.queries.push(sql); + return this.queryResult; + } +} + +const asPoolClient = (client: FakeClient): PoolClient => client as unknown as PoolClient; + +const createServer = () => { + const server = new Server({ pg: { database: 'notify_listener_test' } }); + const reconnects: number[] = []; + jest.spyOn(server, 'addEventListener').mockImplementation(() => { + reconnects.push(Date.now()); + }); + return { server, reconnects }; +}; + +describe('notify listener', () => { + it('attaches the error handler before issuing the LISTEN', () => { + const { server } = createServer(); + const client = new FakeClient(); + const listenerCounts: number[] = []; + jest.spyOn(client, 'query').mockImplementation((sql: string) => { + listenerCounts.push(client.listenerCount('error')); + client.queries.push(sql); + return Promise.resolve({ rows: [] }); + }); + + server.listenForChanges(null, asPoolClient(client), () => {}); + + expect(client.queries).toEqual(['LISTEN "schema:update"']); + expect(listenerCounts).toEqual([1]); + }); + + it('releases and reconnects when the connection drops', () => { + const { server, reconnects } = createServer(); + const client = new FakeClient(); + let released = 0; + + server.listenForChanges(null, asPoolClient(client), () => { + released += 1; + }); + client.emit('error', new Error('Connection terminated unexpectedly')); + + expect(released).toBe(1); + expect(reconnects).toHaveLength(1); + expect(client.listenerCount('error')).toBe(0); + }); + + it('tears down only once when the connection drops repeatedly', () => { + const { server, reconnects } = createServer(); + const client = new FakeClient(); + let released = 0; + + server.listenForChanges(null, asPoolClient(client), () => { + released += 1; + }); + const onError = client.listeners('error')[0] as (e: Error) => void; + onError(new Error('first')); + onError(new Error('second')); + + expect(released).toBe(1); + expect(reconnects).toHaveLength(1); + }); + + it('observes a failing LISTEN and reconnects', async () => { + const { server, reconnects } = createServer(); + const client = new FakeClient(); + client.queryResult = Promise.reject(new Error('LISTEN failed')); + let released = 0; + + server.listenForChanges(null, asPoolClient(client), () => { + released += 1; + }); + await new Promise((resolve) => setImmediate(resolve)); + + expect(released).toBe(1); + expect(reconnects).toHaveLength(1); + }); +}); diff --git a/graphql/server/src/server.ts b/graphql/server/src/server.ts index 8ddd11c48..ea7f1e94c 100644 --- a/graphql/server/src/server.ts +++ b/graphql/server/src/server.ts @@ -267,6 +267,12 @@ class Server { this.listenClient = client; this.listenRelease = release; + // Attached before any query is issued: an 'error' event on a Client with no + // listener is rethrown by Node and takes the whole process down. + client.on('error', (e) => { + this.dropListener(client, 'Error with database notify listener', e); + }); + client.on('notification', ({ channel, payload }) => { if (channel === 'schema:update' && payload) { log.info('schema:update', payload); @@ -274,21 +280,33 @@ class Server { } }); - client.query('LISTEN "schema:update"'); - - client.on('error', (e) => { - if (this.shuttingDown) { - release(); - return; - } - this.error('Error with database notify listener', e); - release(); - this.addEventListener(); + client.query('LISTEN "schema:update"').catch((e) => { + this.dropListener(client, 'Failed to LISTEN for schema:update', e); }); this.log('connected and listening for changes...'); } + private dropListener(client: PoolClient, message: string, err: unknown): void { + // Another path (shutdown, an earlier failure) already tore this client down. + if (this.listenClient !== client) return; + + const release = this.listenRelease; + this.listenClient = null; + this.listenRelease = null; + client.removeAllListeners('notification'); + client.removeAllListeners('error'); + + if (this.shuttingDown) { + release?.(); + return; + } + + this.error(message, err); + release?.(); + this.addEventListener(); + } + async removeEventListener(): Promise { if (!this.listenClient || !this.listenRelease) { return;