From 9171c1dc89c41a44bf18c7387180ff450da7aa80 Mon Sep 17 00:00:00 2001 From: Emmanuel BUU Date: Thu, 1 Oct 2026 17:30:56 +0200 Subject: [PATCH] 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()); + })); +});