diff --git a/src/Notifier.js b/src/Notifier.js index deea8bbf8..66997ea04 100644 --- a/src/Notifier.js +++ b/src/Notifier.js @@ -49,7 +49,7 @@ module.exports = class Notifier extends EventEmitter { error.message ); - request.reply(405); + request.reply(400); return; } @@ -223,6 +223,11 @@ module.exports = class Notifier extends EventEmitter { } this.receiveRequest(this._initial_subscribe); + + // A fetch (Expires: 0) is already over. + if (this._state !== C.STATE_TERMINATED) { + this._ua.newNotifier(this); + } } /** @@ -315,6 +320,7 @@ module.exports = class Notifier extends EventEmitter { this._dialog.terminate(); this._dialog = null; } + this._ua.destroyNotifier(this); logger.debug(`emit "terminated" code=${termination_code}`); this.emit('terminated', termination_code); diff --git a/src/SIPMessage.d.ts b/src/SIPMessage.d.ts index 31a20a98e..ad4de068a 100644 --- a/src/SIPMessage.d.ts +++ b/src/SIPMessage.d.ts @@ -24,6 +24,15 @@ declare class IncomingMessage { export class IncomingRequest extends IncomingMessage { ruri: URI; + + reply( + code: number, + reason?: string, + extraHeaders?: string[], + body?: string, + onSuccess?: () => void, + onFailure?: () => void + ): void; } export class IncomingResponse extends IncomingMessage { diff --git a/src/UA.js b/src/UA.js index 687ff37a2..1d022ae22 100644 --- a/src/UA.js +++ b/src/UA.js @@ -72,6 +72,9 @@ module.exports = class UA extends EventEmitter { this._dynConfiguration = {}; this._dialogs = {}; + // Notifier instances with an established subscription. + this._notifiers = new Set(); + // User actions outside any session/dialog (MESSAGE/OPTIONS). this._applicants = {}; @@ -335,6 +338,16 @@ module.exports = class UA extends EventEmitter { } } + // End every subscription we serve with a final NOTIFY, sent before the + // UA is closed: the subscribers would otherwise wait for the + // subscription to expire. + for (const notifier of this._notifiers) { + logger.debug('closing notifier'); + try { + notifier.terminate(); + } catch (error) {} + } + this._status = C.STATUS_USER_CLOSED; const num_transactions = @@ -520,6 +533,20 @@ module.exports = class UA extends EventEmitter { delete this._sessions[session.id]; } + /** + * Notifier accepted its initial SUBSCRIBE. + */ + newNotifier(notifier) { + this._notifiers.add(notifier); + } + + /** + * Notifier terminated. + */ + destroyNotifier(notifier) { + this._notifiers.delete(notifier); + } + /** * Registered */ @@ -616,8 +643,10 @@ module.exports = class UA extends EventEmitter { message.init_incoming(request); } else if (method === JsSIP_C.SUBSCRIBE) { + // SUBSCRIBE is an allowed method: without a listener it is the event + // package that is not supported (RFC 6665 4.2.1.1). if (this.listeners('newSubscribe').length === 0) { - request.reply(405); + request.reply(489); return; } diff --git a/src/test/include/LoopSocket.ts b/src/test/include/LoopSocket.ts index 0bc912104..daa0e9c7c 100644 --- a/src/test/include/LoopSocket.ts +++ b/src/test/include/LoopSocket.ts @@ -1,5 +1,6 @@ // LoopSocket send message itself. // Used P2P logic: message call-id is modified in each leg. +// LoopSocket.pair() links two sockets, so that two UAs talk to each other. import type { Socket } from '../../Socket'; @@ -8,6 +9,18 @@ export default class LoopSocket implements Socket { via_transport = 'WS'; sip_uri = 'sip:localhost:12345;transport=ws'; + private peer: LoopSocket = this; + + static pair(): [LoopSocket, LoopSocket] { + const a = new LoopSocket(); + const b = new LoopSocket(); + + a.peer = b; + b.peer = a; + + return [a, b]; + } + connect(): void { setTimeout(() => { this.onconnect(); @@ -20,7 +33,7 @@ export default class LoopSocket implements Socket { const new_message = this.modifyCallId(message); setTimeout(() => { - this.ondata(new_message); + this.peer.ondata(new_message); }, 0); return true; diff --git a/src/test/test-SubscriberNotifier.ts b/src/test/test-SubscriberNotifier.ts index b96d8fdad..e873a0467 100644 --- a/src/test/test-SubscriberNotifier.ts +++ b/src/test/test-SubscriberNotifier.ts @@ -20,6 +20,15 @@ const enum STEP { SUBSCRIBER_TERMINATED = 11, } +// Status line of the response to `method` sent through the socket. +function sentStatus(socket: LoopSocket, method: string): string | undefined { + return jest + .mocked(socket.send) + .mock.calls.map(([message]) => message) + .find(m => m.startsWith('SIP/2.0') && m.includes(` ${method}\r\n`)) + ?.split('\r\n')[0]; +} + describe('subscriber/notifier communication', () => { test('should handle subscriber/notifier communication', () => new Promise(resolve => { @@ -218,4 +227,162 @@ describe('subscriber/notifier communication', () => { ua.start(); })); + + test('notifier answers 489 when nobody listens to "newSubscribe"', () => + new Promise(resolve => { + const socket = new LoopSocket(); + + jest.spyOn(socket, 'send'); + + const ua = new UA({ + sockets: socket, + uri: 'sip:ikq@example.com', + register: false, + }); + + ua.on('connected', () => { + const subscriber = ua.subscribe('ikq', 'weather', 'text/plain', {}); + + subscriber.on('terminated', (terminationCode: number) => { + expect(terminationCode).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE); + expect(sentStatus(socket, 'SUBSCRIBE')).toBe('SIP/2.0 489 Bad Event'); + + ua.stop(); + }); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => { + resolve(); + }); + + ua.start(); + })); + + test('notifier answers 400 to a SUBSCRIBE without Event header', () => + new Promise(resolve => { + const socket = new LoopSocket(); + const send = socket.send.bind(socket); + + // The Event header is removed from the SUBSCRIBE on its way. + jest + .spyOn(socket, 'send') + .mockImplementation(message => + send( + message.startsWith('SUBSCRIBE ') + ? message.replace(/\r\nEvent:[^\r]*/i, '') + : message + ) + ); + + const ua = new UA({ + sockets: socket, + uri: 'sip:ikq@example.com', + register: false, + }); + + ua.on('newSubscribe', () => { + throw new Error('a malformed SUBSCRIBE must not reach the application'); + }); + + ua.on('connected', () => { + const subscriber = ua.subscribe('ikq', 'weather', 'text/plain', {}); + + subscriber.on('terminated', (terminationCode: number) => { + expect(terminationCode).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE); + expect(sentStatus(socket, 'SUBSCRIBE')).toBe( + 'SIP/2.0 400 Bad Request' + ); + + ua.stop(); + }); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => { + resolve(); + }); + + ua.start(); + })); + + test('UA.stop() sends the final NOTIFY', () => + new Promise(resolve => { + // A stopped UA drops incoming requests, so it cannot be its own + // subscriber: the subscriber and the notifier are two UAs. + const [subscriberSocket, notifierSocket] = LoopSocket.pair(); + + const subscriberUA = new UA({ + sockets: subscriberSocket, + uri: 'sip:alice@example.com', + contact_uri: 'sip:alice@abcdefabcdef.invalid;transport=ws', + register: false, + }); + + const notifierUA = new UA({ + sockets: notifierSocket, + uri: 'sip:ikq@example.com', + contact_uri: 'sip:ikq@abcdefabcdef.invalid;transport=ws', + register: false, + }); + + const ended: string[] = []; + let disconnected = 0; + + const onDisconnected = (): void => { + if (++disconnected === 2) { + resolve(); + } + }; + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + notifierUA.on('newSubscribe', (e: any) => { + const notifier = notifierUA.notify(e.request, 'text/plain', { + pending: false, + }); + + notifier.on('subscribe', (isUnsubscribe: boolean) => { + if (!isUnsubscribe) { + notifier.notify(); + } + }); + + notifier.on('terminated', (terminationCode: number) => { + expect(terminationCode).toBe(notifier.C.FINAL_NOTIFY_SENT); + ended.push('notifier'); + }); + + notifier.start(); + }); + + subscriberUA.on('connected', () => { + const subscriber = subscriberUA.subscribe( + 'sip:ikq@example.com', + 'weather', + 'text/plain', + { expires: 3600 } + ); + + subscriber.on('active', () => { + notifierUA.stop(); + }); + + subscriber.on('terminated', (terminationCode: number) => { + expect(terminationCode).toBe(subscriber.C.FINAL_NOTIFY_RECEIVED); + expect(ended).toEqual(['notifier']); + + subscriberUA.stop(); + }); + + subscriber.subscribe(); + }); + + subscriberUA.on('disconnected', onDisconnected); + notifierUA.on('disconnected', onDisconnected); + + notifierUA.start(); + subscriberUA.start(); + })); });