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/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..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(); + })); });