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
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
Loading