From 9171c1dc89c41a44bf18c7387180ff450da7aa80 Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Thu, 1 Oct 2026 17:30:56 +0200 Subject: [PATCH 1/2] Subscriber: final response in "terminated", unsubscribe on UA.stop() - The `terminated` event gets the final response as a 4th argument. SUBSCRIBE_NON_OK_RESPONSE covers 404 and 489 alike, which call for different reactions (retry later vs. the server has no such event package). - An invalid target throws TypeError in the constructor, as UA.sendMessage() does, instead of sending `SUBSCRIBE undefined`. - terminate() cancels the pending refresh, and the 2xx to the un-SUBSCRIBE no longer schedules a new one (it used to re-SUBSCRIBE after Expires: 0 when the final NOTIFY was lost). If no final NOTIFY comes within 32 s, the subscription ends with UNSUBSCRIBE_TIMEOUT, a code that was declared but never emitted. - UA.stop() unsubscribes every subscription it knows of, as it terminates sessions, so that the server does not keep notifying a gone UA until the subscription expires. A stopped UA drops incoming requests, the final NOTIFY included: such a subscription ends on the 2xx to its un-SUBSCRIBE, with the new code UNSUBSCRIBED_ON_UA_STOP. Co-Authored-By: Claude Opus 5.5 --- src/Subscriber.d.ts | 4 +- src/Subscriber.js | 73 +++++++++++-- src/UA.js | 28 +++++ src/test/include/ResponderSocket.ts | 160 ++++++++++++++++++++++++++++ src/test/test-Subscriber.ts | 160 ++++++++++++++++++++++++++++ 5 files changed, 417 insertions(+), 8 deletions(-) create mode 100644 src/test/include/ResponderSocket.ts create mode 100644 src/test/test-Subscriber.ts diff --git a/src/Subscriber.d.ts b/src/Subscriber.d.ts index 20bb0392a..4837681e4 100644 --- a/src/Subscriber.d.ts +++ b/src/Subscriber.d.ts @@ -1,5 +1,5 @@ import { EventEmitter } from 'events'; -import { IncomingRequest } from './SIPMessage'; +import { IncomingRequest, IncomingResponse } from './SIPMessage'; import { UA } from './UA'; declare enum SubscriberTerminatedCode { @@ -11,6 +11,7 @@ declare enum SubscriberTerminatedCode { UNSUBSCRIBE_TIMEOUT = 5, FINAL_NOTIFY_RECEIVED = 6, WRONG_NOTIFY_RECEIVED = 7, + UNSUBSCRIBED_ON_UA_STOP = 8, } export interface MessageEventMap { @@ -21,6 +22,7 @@ export interface MessageEventMap { terminationCode: SubscriberTerminatedCode, reason: string | undefined, retryAfter: number | undefined, + response: IncomingResponse | undefined, ]; notify: [ isFinal: boolean, diff --git a/src/Subscriber.js b/src/Subscriber.js index 13f8bde89..5cee4e10b 100644 --- a/src/Subscriber.js +++ b/src/Subscriber.js @@ -23,6 +23,7 @@ const C = { UNSUBSCRIBE_TIMEOUT: 5, FINAL_NOTIFY_RECEIVED: 6, WRONG_NOTIFY_RECEIVED: 7, + UNSUBSCRIBED_ON_UA_STOP: 8, // Subscriber states. STATE_PENDING: 0, @@ -33,6 +34,9 @@ const C = { // RFC 6665 3.1.1, default expires value. DEFAULT_EXPIRES_SEC: 900, + + // How long to wait for the final NOTIFY once the un-SUBSCRIBE is accepted. + FINAL_NOTIFY_TIMEOUT_SEC: 32, }; /** @@ -91,7 +95,11 @@ module.exports = class Subscriber extends EventEmitter { } this._ua = ua; - this._target = target; + this._target = this._ua.normalizeTarget(target); + + if (!this._target) { + throw new TypeError(`Invalid target: ${target}`); + } if (!Utils.isDecimal(expires) || expires <= 0) { expires = C.DEFAULT_EXPIRES_SEC; @@ -347,6 +355,10 @@ module.exports = class Subscriber extends EventEmitter { } this._terminated = true; + // No refresh from now on. + clearTimeout(this._expires_timer); + this._expires_timer = null; + // Set header Expires: 0. const headers = this._headers.map(header => { return header.startsWith('Expires') ? 'Expires: 0' : header; @@ -361,7 +373,8 @@ module.exports = class Subscriber extends EventEmitter { _terminateDialog( terminationCode, reason = undefined, - retryAfter = undefined + retryAfter = undefined, + response = undefined ) { // To prevent duplicate emit terminated event. if (this._state === C.STATE_TERMINATED) { @@ -378,9 +391,11 @@ module.exports = class Subscriber extends EventEmitter { this._dialog = null; } + this._ua.destroySubscriber(this); + logger.debug(`emit "terminated" code=${terminationCode}`); - this.emit('terminated', terminationCode, reason, retryAfter); + this.emit('terminated', terminationCode, reason, retryAfter, response); } _sendInitialSubscribe(body, headers) { @@ -395,9 +410,11 @@ module.exports = class Subscriber extends EventEmitter { this._state = C.STATE_WAITING_NOTIFY; + this._ua.newSubscriber(this); + const request = new SIPMessage.OutgoingRequest( JsSIP_C.SUBSCRIBE, - this._ua.normalizeTarget(this._target), + this._target, this._ua, this._params, headers, @@ -477,7 +494,12 @@ module.exports = class Subscriber extends EventEmitter { if (dialog.error) { // OK response without Contact. logger.warn(dialog.error); - this._terminateDialog(C.SUBSCRIBE_WRONG_OK_RESPONSE); + this._terminateDialog( + C.SUBSCRIBE_WRONG_OK_RESPONSE, + undefined, + undefined, + response + ); return; } @@ -496,6 +518,24 @@ module.exports = class Subscriber extends EventEmitter { } } + // Un-SUBSCRIBE accepted: no refresh, only the final NOTIFY to wait for + // (RFC 6665 4.1.2.3). + if (this._terminated) { + // A stopped UA drops every incoming request, the final NOTIFY too. + if (this._ua.status === this._ua.C.STATUS_USER_CLOSED) { + this._terminateDialog( + C.UNSUBSCRIBED_ON_UA_STOP, + undefined, + undefined, + response + ); + } else { + this._scheduleFinalNotifyTimeout(); + } + + return; + } + // Check expires value. const expires_value = response.getHeader('expires'); @@ -515,9 +555,19 @@ module.exports = class Subscriber extends EventEmitter { this._scheduleSubscribe(expires); } } else if (response.status_code === 401 || response.status_code === 407) { - this._terminateDialog(C.SUBSCRIBE_AUTHENTICATION_FAILED); + this._terminateDialog( + C.SUBSCRIBE_AUTHENTICATION_FAILED, + undefined, + undefined, + response + ); } else if (response.status_code >= 300) { - this._terminateDialog(C.SUBSCRIBE_NON_OK_RESPONSE); + this._terminateDialog( + C.SUBSCRIBE_NON_OK_RESPONSE, + undefined, + undefined, + response + ); } } @@ -555,6 +605,15 @@ module.exports = class Subscriber extends EventEmitter { }, timeout); } + _scheduleFinalNotifyTimeout() { + clearTimeout(this._expires_timer); + this._expires_timer = setTimeout(() => { + this._expires_timer = null; + logger.debug('no final NOTIFY received after un-SUBSCRIBE'); + this._terminateDialog(C.UNSUBSCRIBE_TIMEOUT); + }, C.FINAL_NOTIFY_TIMEOUT_SEC * 1000); + } + _parseSubscriptionState(strState) { switch (strState) { case 'pending': { diff --git a/src/UA.js b/src/UA.js index 687ff37a2..e8f9792a3 100644 --- a/src/UA.js +++ b/src/UA.js @@ -76,6 +76,10 @@ module.exports = class UA extends EventEmitter { this._applicants = {}; this._sessions = {}; + + // Subscriber instances with a subscription sent or established. + this._subscribers = new Set(); + this._transport = null; this._contact = null; this._status = C.STATUS_INIT; @@ -326,6 +330,16 @@ module.exports = class UA extends EventEmitter { } } + // Unsubscribe every subscription. The un-SUBSCRIBE must leave before + // the UA is closed, so that the server does not keep sending NOTIFY + // requests to a gone UA until the subscription expires. + for (const subscriber of this._subscribers) { + logger.debug('closing subscriber'); + try { + subscriber.terminate(); + } catch (error) {} + } + // Run _close_ on every applicant. for (const applicant in this._applicants) { if (Object.prototype.hasOwnProperty.call(this._applicants, applicant)) { @@ -505,6 +519,20 @@ module.exports = class UA extends EventEmitter { delete this._applicants[message]; } + /** + * Subscriber sent its initial SUBSCRIBE. + */ + newSubscriber(subscriber) { + this._subscribers.add(subscriber); + } + + /** + * Subscriber terminated. + */ + destroySubscriber(subscriber) { + this._subscribers.delete(subscriber); + } + /** * new RTCSession */ 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-Subscriber.ts b/src/test/test-Subscriber.ts new file mode 100644 index 000000000..2290a9df2 --- /dev/null +++ b/src/test/test-Subscriber.ts @@ -0,0 +1,160 @@ +import './include/common'; +import ResponderSocket, { finalNotify, reply } from './include/ResponderSocket'; +import type { Responder } from './include/ResponderSocket'; + +// eslint-disable-next-line @typescript-eslint/no-require-imports +const JsSIP = require('../JsSIP.js'); +const { UA } = JsSIP; + +const TARGET = 'sip:bob@example.com'; + +// eslint-disable-next-line @typescript-eslint/no-explicit-any +function startUA(responder: Responder): { ua: any; socket: ResponderSocket } { + const socket = new ResponderSocket(responder); + const ua = new UA({ + sockets: socket, + uri: 'sip:alice@example.com', + register: false, + }); + + ua.start(); + + return { ua, socket }; +} + +describe('Subscriber', () => { + test('refuses a target that is not a SIP URI', () => { + const { ua } = startUA(() => null); + + expect(() => + ua.subscribe('sip:', 'presence', 'application/pidf+xml', {}) + ).toThrow(TypeError); + + ua.stop(); + }); + + test('passes the final response with "terminated"', () => + new Promise(resolve => { + const { ua } = startUA(request => + request.method === 'SUBSCRIBE' ? reply(request, 489, 'Bad Event') : null + ); + + ua.on('connected', () => { + const subscriber = ua.subscribe( + TARGET, + 'presence', + 'application/pidf+xml', + {} + ); + + subscriber.on( + 'terminated', + ( + code: number, + reason: string | undefined, + retryAfter: number | undefined, + response: { status_code: number } | undefined + ) => { + expect(code).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE); + expect(reason).toBeUndefined(); + expect(retryAfter).toBeUndefined(); + expect(response?.status_code).toBe(489); + + ua.stop(); + } + ); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => resolve()); + })); + + test('UA.stop() unsubscribes', () => + new Promise(resolve => { + const { ua, socket } = startUA(request => { + if (request.method !== 'SUBSCRIBE') { + return null; + } + + return reply(request, 200, 'OK', [ + 'Contact: ', + `Expires: ${request.header('Expires')}`, + ]); + }); + + ua.on('connected', () => { + const subscriber = ua.subscribe( + TARGET, + 'presence', + 'application/pidf+xml', + { expires: 3600 } + ); + + subscriber.on('terminated', (code: number) => { + expect(code).toBe(subscriber.C.UNSUBSCRIBED_ON_UA_STOP); + }); + + subscriber.on('accepted', () => { + ua.stop(); + + const subscribes = socket.requests.filter( + r => r.method === 'SUBSCRIBE' + ); + + expect(subscribes).toHaveLength(2); + expect(subscribes[0].header('Expires')).toBe('3600'); + expect(subscribes[1].header('Expires')).toBe('0'); + }); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => resolve()); + })); + + test('terminate() ends on the final NOTIFY, without refreshing', () => + new Promise(resolve => { + const { ua, socket } = startUA(request => { + if (request.method !== 'SUBSCRIBE') { + return null; + } + + const ok = reply(request, 200, 'OK', [ + 'Contact: ', + `Expires: ${request.header('Expires')}`, + ]); + + return request.header('Expires') === '0' + ? [ok, finalNotify(request, 'presence')] + : ok; + }); + + ua.on('connected', () => { + const subscriber = ua.subscribe( + TARGET, + 'presence', + 'application/pidf+xml', + { expires: 3600 } + ); + + subscriber.on('accepted', () => subscriber.terminate()); + + subscriber.on('terminated', (code: number) => { + expect(code).toBe(subscriber.C.FINAL_NOTIFY_RECEIVED); + expect( + socket.requests.map(r => [r.method, r.header('Expires')]) + ).toEqual([ + ['SUBSCRIBE', '3600'], + ['SUBSCRIBE', '0'], + ]); + + ua.stop(); + }); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => resolve()); + })); +}); From 3d762a28b452661742a513e75840a495e471b865 Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Mon, 5 Oct 2026 13:49:52 +0200 Subject: [PATCH 2/2] Move Subscriber tests into test-SubscriberNotifier.ts As asked in the review: the tests now live in the existing file and use LoopSocket with a real Notifier instead of a dedicated ResponderSocket. LoopSocket.pair() links two UAs, since a stopped UA cannot be its own notifier. Co-Authored-By: Claude Opus 5.5 --- src/test/include/LoopSocket.ts | 15 ++- src/test/include/ResponderSocket.ts | 160 ---------------------------- src/test/test-Subscriber.ts | 160 ---------------------------- src/test/test-SubscriberNotifier.ts | 152 +++++++++++++++++++++++++- 4 files changed, 165 insertions(+), 322 deletions(-) delete mode 100644 src/test/include/ResponderSocket.ts delete mode 100644 src/test/test-Subscriber.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-Subscriber.ts b/src/test/test-Subscriber.ts deleted file mode 100644 index 2290a9df2..000000000 --- a/src/test/test-Subscriber.ts +++ /dev/null @@ -1,160 +0,0 @@ -import './include/common'; -import ResponderSocket, { finalNotify, reply } from './include/ResponderSocket'; -import type { Responder } from './include/ResponderSocket'; - -// eslint-disable-next-line @typescript-eslint/no-require-imports -const JsSIP = require('../JsSIP.js'); -const { UA } = JsSIP; - -const TARGET = 'sip:bob@example.com'; - -// eslint-disable-next-line @typescript-eslint/no-explicit-any -function startUA(responder: Responder): { ua: any; socket: ResponderSocket } { - const socket = new ResponderSocket(responder); - const ua = new UA({ - sockets: socket, - uri: 'sip:alice@example.com', - register: false, - }); - - ua.start(); - - return { ua, socket }; -} - -describe('Subscriber', () => { - test('refuses a target that is not a SIP URI', () => { - const { ua } = startUA(() => null); - - expect(() => - ua.subscribe('sip:', 'presence', 'application/pidf+xml', {}) - ).toThrow(TypeError); - - ua.stop(); - }); - - test('passes the final response with "terminated"', () => - new Promise(resolve => { - const { ua } = startUA(request => - request.method === 'SUBSCRIBE' ? reply(request, 489, 'Bad Event') : null - ); - - ua.on('connected', () => { - const subscriber = ua.subscribe( - TARGET, - 'presence', - 'application/pidf+xml', - {} - ); - - subscriber.on( - 'terminated', - ( - code: number, - reason: string | undefined, - retryAfter: number | undefined, - response: { status_code: number } | undefined - ) => { - expect(code).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE); - expect(reason).toBeUndefined(); - expect(retryAfter).toBeUndefined(); - expect(response?.status_code).toBe(489); - - ua.stop(); - } - ); - - subscriber.subscribe(); - }); - - ua.on('disconnected', () => resolve()); - })); - - test('UA.stop() unsubscribes', () => - new Promise(resolve => { - const { ua, socket } = startUA(request => { - if (request.method !== 'SUBSCRIBE') { - return null; - } - - return reply(request, 200, 'OK', [ - 'Contact: ', - `Expires: ${request.header('Expires')}`, - ]); - }); - - ua.on('connected', () => { - const subscriber = ua.subscribe( - TARGET, - 'presence', - 'application/pidf+xml', - { expires: 3600 } - ); - - subscriber.on('terminated', (code: number) => { - expect(code).toBe(subscriber.C.UNSUBSCRIBED_ON_UA_STOP); - }); - - subscriber.on('accepted', () => { - ua.stop(); - - const subscribes = socket.requests.filter( - r => r.method === 'SUBSCRIBE' - ); - - expect(subscribes).toHaveLength(2); - expect(subscribes[0].header('Expires')).toBe('3600'); - expect(subscribes[1].header('Expires')).toBe('0'); - }); - - subscriber.subscribe(); - }); - - ua.on('disconnected', () => resolve()); - })); - - test('terminate() ends on the final NOTIFY, without refreshing', () => - new Promise(resolve => { - const { ua, socket } = startUA(request => { - if (request.method !== 'SUBSCRIBE') { - return null; - } - - const ok = reply(request, 200, 'OK', [ - 'Contact: ', - `Expires: ${request.header('Expires')}`, - ]); - - return request.header('Expires') === '0' - ? [ok, finalNotify(request, 'presence')] - : ok; - }); - - ua.on('connected', () => { - const subscriber = ua.subscribe( - TARGET, - 'presence', - 'application/pidf+xml', - { expires: 3600 } - ); - - subscriber.on('accepted', () => subscriber.terminate()); - - subscriber.on('terminated', (code: number) => { - expect(code).toBe(subscriber.C.FINAL_NOTIFY_RECEIVED); - expect( - socket.requests.map(r => [r.method, r.header('Expires')]) - ).toEqual([ - ['SUBSCRIBE', '3600'], - ['SUBSCRIBE', '0'], - ]); - - ua.stop(); - }); - - subscriber.subscribe(); - }); - - ua.on('disconnected', () => resolve()); - })); -}); diff --git a/src/test/test-SubscriberNotifier.ts b/src/test/test-SubscriberNotifier.ts index b96d8fdad..a795aaae3 100644 --- a/src/test/test-SubscriberNotifier.ts +++ b/src/test/test-SubscriberNotifier.ts @@ -20,6 +20,15 @@ const enum STEP { SUBSCRIBER_TERMINATED = 11, } +// Expires header of every SUBSCRIBE sent through the socket. +function sentSubscribeExpires(socket: LoopSocket): string[] { + return jest + .mocked(socket.send) + .mock.calls.map(([message]) => message) + .filter(message => message.startsWith('SUBSCRIBE ')) + .map(message => /\r\nExpires: *(\d+)/i.exec(message)![1]); +} + describe('subscriber/notifier communication', () => { test('should handle subscriber/notifier communication', () => new Promise(resolve => { @@ -104,6 +113,9 @@ describe('subscriber/notifier communication', () => { expect(reason).toBeUndefined(); expect(retryAfter).toBeUndefined(); + // The un-SUBSCRIBE, and no refresh after it. + expect(sentSubscribeExpires(socket)).toEqual(['3600', '0']); + ua.stop(); } ); @@ -164,8 +176,12 @@ describe('subscriber/notifier communication', () => { } // Start JsSIP UA with loop socket. + const socket = new LoopSocket(); // message sending itself, with modified Call-ID + + jest.spyOn(socket, 'send'); + const config = { - sockets: new LoopSocket(), // message sending itself, with modified Call-ID + sockets: socket, uri: REQUEST_URI, contact_uri: CONTACT_URI, register: false, @@ -218,4 +234,138 @@ describe('subscriber/notifier communication', () => { ua.start(); })); + + test('subscriber refuses a target that is not a SIP URI', () => { + const ua = new UA({ + sockets: new LoopSocket(), + uri: 'sip:ikq@example.com', + register: false, + }); + + expect(() => ua.subscribe('sip:', 'weather', 'text/plain', {})).toThrow( + TypeError + ); + }); + + test('subscriber gets the final response with "terminated"', () => + new Promise(resolve => { + const ua = new UA({ + sockets: new LoopSocket(), + uri: 'sip:ikq@example.com', + register: false, + }); + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + ua.on('newSubscribe', (e: any) => { + e.request.reply(489); // "Bad Event" + }); + + ua.on('connected', () => { + const subscriber = ua.subscribe('ikq', 'weather', 'text/plain', {}); + + subscriber.on( + 'terminated', + ( + terminationCode: number, + reason: string | undefined, + retryAfter: number | undefined, + response: { status_code: number } | undefined + ) => { + expect(terminationCode).toBe( + subscriber.C.SUBSCRIBE_NON_OK_RESPONSE + ); + expect(reason).toBeUndefined(); + expect(retryAfter).toBeUndefined(); + expect(response?.status_code).toBe(489); + + ua.stop(); + } + ); + + subscriber.subscribe(); + }); + + ua.on('disconnected', () => { + resolve(); + }); + + ua.start(); + })); + + test('UA.stop() unsubscribes', () => + new Promise(resolve => { + // A stopped UA drops incoming requests, so it cannot be its own + // notifier: the subscriber and the notifier are two UAs. + const [subscriberSocket, notifierSocket] = LoopSocket.pair(); + + jest.spyOn(subscriberSocket, 'send'); + + 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, + }); + + 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.terminate(); + } else { + notifier.notify(); + } + }); + + notifier.start(); + }); + + subscriberUA.on('connected', () => { + const subscriber = subscriberUA.subscribe( + 'sip:ikq@example.com', + 'weather', + 'text/plain', + { expires: 3600 } + ); + + subscriber.on('active', () => { + subscriberUA.stop(); + + expect(sentSubscribeExpires(subscriberSocket)).toEqual(['3600', '0']); + }); + + subscriber.on('terminated', (terminationCode: number) => { + expect(terminationCode).toBe(subscriber.C.UNSUBSCRIBED_ON_UA_STOP); + + notifierUA.stop(); + }); + + subscriber.subscribe(); + }); + + subscriberUA.on('disconnected', onDisconnected); + notifierUA.on('disconnected', onDisconnected); + + notifierUA.start(); + subscriberUA.start(); + })); });