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
15 changes: 14 additions & 1 deletion src/test/include/LoopSocket.ts
Original file line number Diff line number Diff line change
@@ -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';

Expand All @@ -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();
Expand All @@ -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;
Expand Down
167 changes: 167 additions & 0 deletions src/test/test-SubscriberNotifier.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,15 @@ const enum STEP {
SUBSCRIBER_TERMINATED = 11,
}

// Status line of the response to `method` sent through the socket.
function sentStatus(socket: LoopSocket, method: string): string | undefined {
return jest
.mocked(socket.send)
.mock.calls.map(([message]) => message)
.find(m => m.startsWith('SIP/2.0') && m.includes(` ${method}\r\n`))
?.split('\r\n')[0];
}

describe('subscriber/notifier communication', () => {
test('should handle subscriber/notifier communication', () =>
new Promise<void>(resolve => {
Expand Down Expand Up @@ -218,4 +227,162 @@ describe('subscriber/notifier communication', () => {

ua.start();
}));

test('notifier answers 489 when nobody listens to "newSubscribe"', () =>
new Promise<void>(resolve => {
const socket = new LoopSocket();

jest.spyOn(socket, 'send');

const ua = new UA({
sockets: socket,
uri: 'sip:ikq@example.com',
register: false,
});

ua.on('connected', () => {
const subscriber = ua.subscribe('ikq', 'weather', 'text/plain', {});

subscriber.on('terminated', (terminationCode: number) => {
expect(terminationCode).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE);
expect(sentStatus(socket, 'SUBSCRIBE')).toBe('SIP/2.0 489 Bad Event');

ua.stop();
});

subscriber.subscribe();
});

ua.on('disconnected', () => {
resolve();
});

ua.start();
}));

test('notifier answers 400 to a SUBSCRIBE without Event header', () =>
new Promise<void>(resolve => {
const socket = new LoopSocket();
const send = socket.send.bind(socket);

// The Event header is removed from the SUBSCRIBE on its way.
jest
.spyOn(socket, 'send')
.mockImplementation(message =>
send(
message.startsWith('SUBSCRIBE ')
? message.replace(/\r\nEvent:[^\r]*/i, '')
: message
)
);

const ua = new UA({
sockets: socket,
uri: 'sip:ikq@example.com',
register: false,
});

ua.on('newSubscribe', () => {
throw new Error('a malformed SUBSCRIBE must not reach the application');
});

ua.on('connected', () => {
const subscriber = ua.subscribe('ikq', 'weather', 'text/plain', {});

subscriber.on('terminated', (terminationCode: number) => {
expect(terminationCode).toBe(subscriber.C.SUBSCRIBE_NON_OK_RESPONSE);
expect(sentStatus(socket, 'SUBSCRIBE')).toBe(
'SIP/2.0 400 Bad Request'
);

ua.stop();
});

subscriber.subscribe();
});

ua.on('disconnected', () => {
resolve();
});

ua.start();
}));

test('UA.stop() sends the final NOTIFY', () =>
new Promise<void>(resolve => {
// A stopped UA drops incoming requests, so it cannot be its own
// subscriber: the subscriber and the notifier are two UAs.
const [subscriberSocket, notifierSocket] = LoopSocket.pair();

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,
});

const ended: string[] = [];
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.notify();
}
});

notifier.on('terminated', (terminationCode: number) => {
expect(terminationCode).toBe(notifier.C.FINAL_NOTIFY_SENT);
ended.push('notifier');
});

notifier.start();
});

subscriberUA.on('connected', () => {
const subscriber = subscriberUA.subscribe(
'sip:ikq@example.com',
'weather',
'text/plain',
{ expires: 3600 }
);

subscriber.on('active', () => {
notifierUA.stop();
});

subscriber.on('terminated', (terminationCode: number) => {
expect(terminationCode).toBe(subscriber.C.FINAL_NOTIFY_RECEIVED);
expect(ended).toEqual(['notifier']);

subscriberUA.stop();
});

subscriber.subscribe();
});

subscriberUA.on('disconnected', onDisconnected);
notifierUA.on('disconnected', onDisconnected);

notifierUA.start();
subscriberUA.start();
}));
});