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]); + }); +});