Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion src/Notifier.js
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ module.exports = class Notifier extends EventEmitter {
error.message
);

request.reply(405);
request.reply(400);

return;
}
Expand Down Expand Up @@ -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);
}
}

/**
Expand Down Expand Up @@ -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);
Expand Down
9 changes: 9 additions & 0 deletions src/SIPMessage.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
31 changes: 30 additions & 1 deletion src/UA.js
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {};

Expand Down Expand Up @@ -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 =
Expand Down Expand Up @@ -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
*/
Expand Down Expand Up @@ -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;
}
Expand Down
160 changes: 160 additions & 0 deletions src/test/include/ResponderSocket.ts
Original file line number Diff line number Diff line change
@@ -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<T>(_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: <sip:responder@example.com;transport=ws>',
'Content-Length: 0',
'',
'',
].join('\r\n');
}
133 changes: 133 additions & 0 deletions src/test/test-Notifier.ts
Original file line number Diff line number Diff line change
@@ -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: <sip:bob@example.com>;tag=bob',
'To: <sip:alice@example.com>',
'Call-ID: notifier-test',
'CSeq: 1 SUBSCRIBE',
'Contact: <sip:bob@example.com;transport=ws>',
'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<void> {
return new Promise(resolve => {
ua.on('disconnected', () => resolve());
ua.stop();
});
}

function connected(ua: TestUA): Promise<void> {
return new Promise(resolve => ua.on('connected', () => resolve()));
}

function later(): Promise<void> {
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]);
});
});