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
4 changes: 3 additions & 1 deletion src/Subscriber.d.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { EventEmitter } from 'events';
import { IncomingRequest } from './SIPMessage';
import { IncomingRequest, IncomingResponse } from './SIPMessage';
import { UA } from './UA';

declare enum SubscriberTerminatedCode {
Expand All @@ -11,6 +11,7 @@ declare enum SubscriberTerminatedCode {
UNSUBSCRIBE_TIMEOUT = 5,
FINAL_NOTIFY_RECEIVED = 6,
WRONG_NOTIFY_RECEIVED = 7,
UNSUBSCRIBED_ON_UA_STOP = 8,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

UA_STOPPED

}

export interface MessageEventMap {
Expand All @@ -21,6 +22,7 @@ export interface MessageEventMap {
terminationCode: SubscriberTerminatedCode,
reason: string | undefined,
retryAfter: number | undefined,
response: IncomingResponse | undefined,
];
notify: [
isFinal: boolean,
Expand Down
73 changes: 66 additions & 7 deletions src/Subscriber.js
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
};

/**
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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) {
Expand All @@ -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) {
Expand All @@ -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,
Expand Down Expand Up @@ -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;
}
Expand All @@ -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');

Expand All @@ -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
);
}
}

Expand Down Expand Up @@ -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': {
Expand Down
28 changes: 28 additions & 0 deletions src/UA.js
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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)) {
Expand Down Expand Up @@ -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
*/
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');
}
Loading