From ee25635c62ee849adb96b155ceb853b71c3d318b Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Thu, 1 Oct 2026 17:34:06 +0200 Subject: [PATCH 1/2] Notifier: 489/400 to incoming SUBSCRIBE, final NOTIFY on UA.stop() - Without a `newSubscribe` listener, an incoming SUBSCRIBE got 405 Method Not Allowed, although SUBSCRIBE is in the Allow header of that very response. What is missing is the event package: 489 Bad Event (RFC 6665 4.2.1.1). - A SUBSCRIBE without Contact or Event, or with a bad Expires, got 405 too; it is a malformed request: 400 Bad Request. - UA.stop() ends every Notifier with a final NOTIFY (Subscription-State: terminated), sent before the UA is closed, as it does for sessions. Subscribers no longer wait for the subscription to expire to learn that we are gone. - IncomingRequest.reply() is declared: every `newSubscribe` listener calls it, and had to cast the request to do so. Co-Authored-By: Claude Opus 5.5 --- src/Notifier.js | 8 +- src/SIPMessage.d.ts | 9 ++ src/UA.js | 31 +++++- src/test/include/ResponderSocket.ts | 160 ++++++++++++++++++++++++++++ src/test/test-Notifier.ts | 133 +++++++++++++++++++++++ 5 files changed, 339 insertions(+), 2 deletions(-) create mode 100644 src/test/include/ResponderSocket.ts create mode 100644 src/test/test-Notifier.ts 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/ResponderSocket.ts b/src/test/include/ResponderSocket.ts new file mode 100644 index 000000000..370547dd4 --- /dev/null +++ b/src/test/include/ResponderSocket.ts @@ -0,0 +1,160 @@ +// ResponderSocket answers every request the UA sends with the response +// built by the given handler. Used to play the server side of a client +// transaction (SUBSCRIBE, PUBLISH...) without a network. + +import type { Socket } from '../../Socket'; + +export type Responder = (request: SentRequest) => string | string[] | null; + +export interface SentRequest { + method: string; + raw: string; + header(name: string): string | undefined; + body: string; +} + +export default class ResponderSocket implements Socket { + url = 'ws://localhost:12345'; + via_transport = 'WS'; + sip_uri = 'sip:localhost:12345;transport=ws'; + + // Every message sent by the UA, requests and responses. + readonly sent: string[] = []; + + // The requests among them. + readonly requests: SentRequest[] = []; + + constructor(private readonly responder: Responder) {} + + connect(): void { + setTimeout(() => { + this.onconnect(); + }, 0); + } + + disconnect(): void { + setTimeout(() => { + this.ondisconnect(); + }, 0); + } + + send(message: string): boolean { + this.sent.push(message); + + // Only requests are answered. + if (message.startsWith('SIP/2.0')) { + return true; + } + + const request = parse(message); + + this.requests.push(request); + + const answer = this.responder(request); + const messages = answer === null ? [] : ([] as string[]).concat(answer); + + for (const data of messages) { + this.inject(data); + } + + return true; + } + + // Deliver a message to the UA, as if the server had sent it. + inject(message: string): void { + setTimeout(() => { + this.ondata(message); + }, 0); + } + + isConnected(): boolean { + return true; + } + + isConnecting(): boolean { + return false; + } + + onconnect(): void {} + + ondisconnect(): void {} + + // eslint-disable-next-line @typescript-eslint/no-unused-vars + ondata(_event: T): void {} +} + +function parse(message: string): SentRequest { + const [head, body = ''] = message.split('\r\n\r\n'); + const lines = head.split('\r\n'); + + return { + method: lines[0].split(' ')[0], + raw: message, + header(name: string): string | undefined { + const prefix = `${name.toLowerCase()}:`; + const line = lines.find(l => l.toLowerCase().startsWith(prefix)); + + return line?.substring(prefix.length).trim(); + }, + body, + }; +} + +/** + * Build the response to `request`: Via, From, To (tagged), Call-ID and CSeq + * copied from it, then `headers`. + */ +export function reply( + request: SentRequest, + status: number, + reason: string, + headers: string[] = [] +): string { + const lines = request.raw.split('\r\n\r\n')[0].split('\r\n'); + const copied = lines.filter(line => /^(via|from|call-id|cseq):/i.test(line)); + let to = lines.find(line => /^to:/i.test(line))!; + + if (!/;tag=/i.test(to)) { + to += ';tag=responder'; + } + + return [ + `SIP/2.0 ${status} ${reason}`, + ...copied, + to, + ...headers, + 'Content-Length: 0', + '', + '', + ].join('\r\n'); +} + +/** + * Build the final NOTIFY of the subscription `subscribe` (any request of + * its dialog), as the notifier of `reply` would send it. + */ +export function finalNotify(subscribe: SentRequest, event: string): string { + const contact = subscribe.header('Contact')!.replace(/^<([^>]*)>.*$/, '$1'); + const from = subscribe.header('From')!; + let to = subscribe.header('To')!; + + if (!/;tag=/i.test(to)) { + to += ';tag=responder'; + } + + return [ + `NOTIFY ${contact} SIP/2.0`, + 'Via: SIP/2.0/WS example.com;branch=z9hG4bKresponder', + 'Max-Forwards: 70', + `From: ${to}`, + `To: ${from}`, + `Call-ID: ${subscribe.header('Call-ID')}`, + 'CSeq: 1 NOTIFY', + `Event: ${event}`, + 'Subscription-State: terminated', + 'Contact: ', + 'Content-Length: 0', + '', + '', + ].join('\r\n'); +} diff --git a/src/test/test-Notifier.ts b/src/test/test-Notifier.ts new file mode 100644 index 000000000..7c9c73168 --- /dev/null +++ b/src/test/test-Notifier.ts @@ -0,0 +1,133 @@ +import './include/common'; +import ResponderSocket, { reply } from './include/ResponderSocket'; + +// eslint-disable-next-line @typescript-eslint/no-require-imports +const JsSIP = require('../JsSIP.js'); +const { UA } = JsSIP; + +interface TestUA { + on(event: string, listener: (...args: never[]) => void): void; + stop(): void; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + notify(...args: any[]): any; +} + +function startUA(): { ua: TestUA; socket: ResponderSocket } { + // NOTIFY requests are accepted. + const socket = new ResponderSocket(request => + request.method === 'NOTIFY' ? reply(request, 200, 'OK') : null + ); + const ua = new UA({ + sockets: socket, + uri: 'sip:alice@example.com', + register: false, + }); + + ua.start(); + + return { ua, socket }; +} + +/** A SUBSCRIBE from bob to alice, without the headers in `omit`. */ +function subscribe(omit: string[] = []): string { + const headers = [ + 'Via: SIP/2.0/WS example.com;branch=z9hG4bKsubscribe', + 'Max-Forwards: 70', + 'From: ;tag=bob', + 'To: ', + 'Call-ID: notifier-test', + 'CSeq: 1 SUBSCRIBE', + 'Contact: ', + 'Event: presence', + 'Expires: 600', + 'Accept: application/pidf+xml', + ].filter(h => !omit.some(name => h.startsWith(`${name}:`))); + + return [ + 'SUBSCRIBE sip:alice@example.com SIP/2.0', + ...headers, + 'Content-Length: 0', + '', + '', + ].join('\r\n'); +} + +/** The status line of the response the UA sent to `method`. */ +function statusOf(socket: ResponderSocket, method: string): string | undefined { + return socket.sent + .find(m => m.startsWith('SIP/2.0') && m.includes(` ${method}\r\n`)) + ?.split('\r\n')[0]; +} + +function stopAndWait(ua: TestUA): Promise { + return new Promise(resolve => { + ua.on('disconnected', () => resolve()); + ua.stop(); + }); +} + +function connected(ua: TestUA): Promise { + return new Promise(resolve => ua.on('connected', () => resolve())); +} + +function later(): Promise { + return new Promise(resolve => setTimeout(resolve, 10)); +} + +describe('incoming SUBSCRIBE', () => { + test('489 when nobody listens to "newSubscribe"', async () => { + const { ua, socket } = startUA(); + + await connected(ua); + socket.inject(subscribe()); + await later(); + + expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 489 Bad Event'); + + await stopAndWait(ua); + }); + + test('400 when the SUBSCRIBE has no Event header', async () => { + const { ua, socket } = startUA(); + + ua.on('newSubscribe', () => { + throw new Error('a malformed SUBSCRIBE must not reach the application'); + }); + + await connected(ua); + socket.inject(subscribe(['Event'])); + await later(); + + expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 400 Bad Request'); + + await stopAndWait(ua); + }); +}); + +describe('Notifier', () => { + test('UA.stop() sends the final NOTIFY', async () => { + const { ua, socket } = startUA(); + const ended: number[] = []; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + let notifier: any; + + ua.on('newSubscribe', (e: { request: unknown }) => { + notifier = ua.notify(e.request, 'application/pidf+xml', {}); + notifier.on('terminated', (code: number) => ended.push(code)); + notifier.start(); + }); + + await connected(ua); + socket.inject(subscribe()); + await later(); + + expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 200 OK'); + + await stopAndWait(ua); + + const notify = socket.requests.find(r => r.method === 'NOTIFY'); + + expect(notify?.header('Subscription-State')).toBe('terminated'); + expect(ended).toEqual([notifier.C.FINAL_NOTIFY_SENT]); + }); +}); From 39e461c633e6252a401cabe99ce7d17a154df055 Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Mon, 5 Oct 2026 14:03:33 +0200 Subject: [PATCH 2/2] Move Notifier tests into test-SubscriberNotifier.ts As asked in the review: the tests now live in the existing file and use LoopSocket with a real Subscriber instead of a dedicated ResponderSocket. LoopSocket.pair() links two UAs, since a stopped UA cannot be its own subscriber. Co-Authored-By: Claude Opus 5.5 --- src/test/include/LoopSocket.ts | 15 ++- src/test/include/ResponderSocket.ts | 160 -------------------------- src/test/test-Notifier.ts | 133 ---------------------- src/test/test-SubscriberNotifier.ts | 167 ++++++++++++++++++++++++++++ 4 files changed, 181 insertions(+), 294 deletions(-) delete mode 100644 src/test/include/ResponderSocket.ts delete mode 100644 src/test/test-Notifier.ts 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/include/ResponderSocket.ts b/src/test/include/ResponderSocket.ts deleted file mode 100644 index 370547dd4..000000000 --- a/src/test/include/ResponderSocket.ts +++ /dev/null @@ -1,160 +0,0 @@ -// ResponderSocket answers every request the UA sends with the response -// built by the given handler. Used to play the server side of a client -// transaction (SUBSCRIBE, PUBLISH...) without a network. - -import type { Socket } from '../../Socket'; - -export type Responder = (request: SentRequest) => string | string[] | null; - -export interface SentRequest { - method: string; - raw: string; - header(name: string): string | undefined; - body: string; -} - -export default class ResponderSocket implements Socket { - url = 'ws://localhost:12345'; - via_transport = 'WS'; - sip_uri = 'sip:localhost:12345;transport=ws'; - - // Every message sent by the UA, requests and responses. - readonly sent: string[] = []; - - // The requests among them. - readonly requests: SentRequest[] = []; - - constructor(private readonly responder: Responder) {} - - connect(): void { - setTimeout(() => { - this.onconnect(); - }, 0); - } - - disconnect(): void { - setTimeout(() => { - this.ondisconnect(); - }, 0); - } - - send(message: string): boolean { - this.sent.push(message); - - // Only requests are answered. - if (message.startsWith('SIP/2.0')) { - return true; - } - - const request = parse(message); - - this.requests.push(request); - - const answer = this.responder(request); - const messages = answer === null ? [] : ([] as string[]).concat(answer); - - for (const data of messages) { - this.inject(data); - } - - return true; - } - - // Deliver a message to the UA, as if the server had sent it. - inject(message: string): void { - setTimeout(() => { - this.ondata(message); - }, 0); - } - - isConnected(): boolean { - return true; - } - - isConnecting(): boolean { - return false; - } - - onconnect(): void {} - - ondisconnect(): void {} - - // eslint-disable-next-line @typescript-eslint/no-unused-vars - ondata(_event: T): void {} -} - -function parse(message: string): SentRequest { - const [head, body = ''] = message.split('\r\n\r\n'); - const lines = head.split('\r\n'); - - return { - method: lines[0].split(' ')[0], - raw: message, - header(name: string): string | undefined { - const prefix = `${name.toLowerCase()}:`; - const line = lines.find(l => l.toLowerCase().startsWith(prefix)); - - return line?.substring(prefix.length).trim(); - }, - body, - }; -} - -/** - * Build the response to `request`: Via, From, To (tagged), Call-ID and CSeq - * copied from it, then `headers`. - */ -export function reply( - request: SentRequest, - status: number, - reason: string, - headers: string[] = [] -): string { - const lines = request.raw.split('\r\n\r\n')[0].split('\r\n'); - const copied = lines.filter(line => /^(via|from|call-id|cseq):/i.test(line)); - let to = lines.find(line => /^to:/i.test(line))!; - - if (!/;tag=/i.test(to)) { - to += ';tag=responder'; - } - - return [ - `SIP/2.0 ${status} ${reason}`, - ...copied, - to, - ...headers, - 'Content-Length: 0', - '', - '', - ].join('\r\n'); -} - -/** - * Build the final NOTIFY of the subscription `subscribe` (any request of - * its dialog), as the notifier of `reply` would send it. - */ -export function finalNotify(subscribe: SentRequest, event: string): string { - const contact = subscribe.header('Contact')!.replace(/^<([^>]*)>.*$/, '$1'); - const from = subscribe.header('From')!; - let to = subscribe.header('To')!; - - if (!/;tag=/i.test(to)) { - to += ';tag=responder'; - } - - return [ - `NOTIFY ${contact} SIP/2.0`, - 'Via: SIP/2.0/WS example.com;branch=z9hG4bKresponder', - 'Max-Forwards: 70', - `From: ${to}`, - `To: ${from}`, - `Call-ID: ${subscribe.header('Call-ID')}`, - 'CSeq: 1 NOTIFY', - `Event: ${event}`, - 'Subscription-State: terminated', - 'Contact: ', - 'Content-Length: 0', - '', - '', - ].join('\r\n'); -} diff --git a/src/test/test-Notifier.ts b/src/test/test-Notifier.ts deleted file mode 100644 index 7c9c73168..000000000 --- a/src/test/test-Notifier.ts +++ /dev/null @@ -1,133 +0,0 @@ -import './include/common'; -import ResponderSocket, { reply } from './include/ResponderSocket'; - -// eslint-disable-next-line @typescript-eslint/no-require-imports -const JsSIP = require('../JsSIP.js'); -const { UA } = JsSIP; - -interface TestUA { - on(event: string, listener: (...args: never[]) => void): void; - stop(): void; - // eslint-disable-next-line @typescript-eslint/no-explicit-any - notify(...args: any[]): any; -} - -function startUA(): { ua: TestUA; socket: ResponderSocket } { - // NOTIFY requests are accepted. - const socket = new ResponderSocket(request => - request.method === 'NOTIFY' ? reply(request, 200, 'OK') : null - ); - const ua = new UA({ - sockets: socket, - uri: 'sip:alice@example.com', - register: false, - }); - - ua.start(); - - return { ua, socket }; -} - -/** A SUBSCRIBE from bob to alice, without the headers in `omit`. */ -function subscribe(omit: string[] = []): string { - const headers = [ - 'Via: SIP/2.0/WS example.com;branch=z9hG4bKsubscribe', - 'Max-Forwards: 70', - 'From: ;tag=bob', - 'To: ', - 'Call-ID: notifier-test', - 'CSeq: 1 SUBSCRIBE', - 'Contact: ', - 'Event: presence', - 'Expires: 600', - 'Accept: application/pidf+xml', - ].filter(h => !omit.some(name => h.startsWith(`${name}:`))); - - return [ - 'SUBSCRIBE sip:alice@example.com SIP/2.0', - ...headers, - 'Content-Length: 0', - '', - '', - ].join('\r\n'); -} - -/** The status line of the response the UA sent to `method`. */ -function statusOf(socket: ResponderSocket, method: string): string | undefined { - return socket.sent - .find(m => m.startsWith('SIP/2.0') && m.includes(` ${method}\r\n`)) - ?.split('\r\n')[0]; -} - -function stopAndWait(ua: TestUA): Promise { - return new Promise(resolve => { - ua.on('disconnected', () => resolve()); - ua.stop(); - }); -} - -function connected(ua: TestUA): Promise { - return new Promise(resolve => ua.on('connected', () => resolve())); -} - -function later(): Promise { - return new Promise(resolve => setTimeout(resolve, 10)); -} - -describe('incoming SUBSCRIBE', () => { - test('489 when nobody listens to "newSubscribe"', async () => { - const { ua, socket } = startUA(); - - await connected(ua); - socket.inject(subscribe()); - await later(); - - expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 489 Bad Event'); - - await stopAndWait(ua); - }); - - test('400 when the SUBSCRIBE has no Event header', async () => { - const { ua, socket } = startUA(); - - ua.on('newSubscribe', () => { - throw new Error('a malformed SUBSCRIBE must not reach the application'); - }); - - await connected(ua); - socket.inject(subscribe(['Event'])); - await later(); - - expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 400 Bad Request'); - - await stopAndWait(ua); - }); -}); - -describe('Notifier', () => { - test('UA.stop() sends the final NOTIFY', async () => { - const { ua, socket } = startUA(); - const ended: number[] = []; - // eslint-disable-next-line @typescript-eslint/no-explicit-any - let notifier: any; - - ua.on('newSubscribe', (e: { request: unknown }) => { - notifier = ua.notify(e.request, 'application/pidf+xml', {}); - notifier.on('terminated', (code: number) => ended.push(code)); - notifier.start(); - }); - - await connected(ua); - socket.inject(subscribe()); - await later(); - - expect(statusOf(socket, 'SUBSCRIBE')).toBe('SIP/2.0 200 OK'); - - await stopAndWait(ua); - - const notify = socket.requests.find(r => r.method === 'NOTIFY'); - - expect(notify?.header('Subscription-State')).toBe('terminated'); - expect(ended).toEqual([notifier.C.FINAL_NOTIFY_SENT]); - }); -}); 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(); + })); });