import { EventDispatcher } from './qik.utils.js';
///////////////////////////////////////////////////
/**
* Creates a new instance of QikSocket a module of the SDK
* that contains all helper functions to do with realtime sockets
* @name socket
* @constructor
* @hideconstructor
* @param {QikAPI} qik A reference to the parent instance of the QikCore module. This module is usually created by a QikCore instance that passes itself in as the first argument.
*/
const QikSocket = function(qik, mode) {
mode = mode || 'production';
if (!qik.auth) {
throw new Error(`Can't Instantiate QikSocket before QikAuth exists`);
}
const windowID = qik.utils.guid();
///////////////////////////////////////////////////
let service = {
debug: false,
url: `wss://iqtm6zjz3l.execute-api.ap-southeast-2.amazonaws.com/${mode}`,
connected: false,
windowID,
autoreconnect:true,
}
let socket;
let timer;
let buckets = {};
const channels = {};
//Create a new dispatcher
const dispatcher = new EventDispatcher();
dispatcher.bootstrap(service);
///////////////////////////////////////////////////
function getSubscriptions() {
return Object.keys(buckets);
}
///////////////////////////////////////////////////
/**
* @description Create an event dispatcher instance and subscribe to a new socket channel.
* This can then be used to receive socket events for this channel
* @alias socket.channel
* @param {String} key The key or name of the channel to subscribe to
* @example
* const socketChannel = await sdk.socket.channel('awesome:notifications')
* socketChannel.addEventListener('message', function(message) {
* console.log('Received a new message from the socket')
* })
*
* // Disconnect from this channel and remove all listeners
* socketChannel.destroy();
*
*/
service.channel = async function(key) {
if (buckets[key]) {
return buckets[key];
}
const bucket = new EventDispatcher();
function dispatchMessage(message) {
bucket.dispatch('message', message);
}
// Destroy this channel - stop receiving its events and tell the
// server to unsubscribe. Only touches this channel's own listeners,
// leaving every other channel (and connection-level listeners) intact.
bucket.destroy = function() {
service.removeEventListener(key, dispatchMessage);
service.unsubscribe(key);
bucket.removeAllListeners();
delete buckets[key];
}
// Cache the fully-formed bucket (with destroy attached) BEFORE the
// connect() await below. A concurrent channel(key) call can hit the
// early return above while we're suspended awaiting connect(); if we
// registered the bucket before assigning destroy, that caller would
// receive a bucket whose destroy() is missing and throw on teardown.
buckets[key] = bucket;
if (!service.connected) {
// Return the socket
await service.connect();
}
// Subscribe to this channel
service.subscribe(key);
// Listen for all relevant events
service.addEventListener(key, dispatchMessage);
return buckets[key];
}
///////////////////////////////////////////////////
function ping() {
broadcast({
action: 'ping',
})
}
function startHeartbeat() {
if (timer) {
return;
}
timer = setInterval(ping, 10000);
}
function stopHeartbeat() {
if (!timer) {
return;
}
clearInterval(timer)
timer = null;
}
///////////////////////////////////////////////////
function socketOpened(event) {
service.debug ? console.log("[socket] Connection open", event) : null;
service.connected = true;
dispatcher.dispatch('connected', event);
// startHeartbeat();
}
async function socketClosed(event) {
if (event.wasClean) {
service.debug ? console.log(`[socket] Connection closed cleanly, code=${event.code} reason=${event.reason}`) : null;
} else {
// e.g. server process killed or network down
// event.code is usually 1006 in this case
service.debug ? console.log('[socket] - connection closed due to error', event) : null
}
service.connected = false;
dispatcher.dispatch('disconnected', event);
socket = null;
if(service.autoreconnect) {
service.reconnect();
}
}
function socketError(error) {
console.log('[event] socketError', event);
service.debug ? console.log("[socket] Error", error) : null;
dispatcher.dispatch('error', error);
}
function socketMessageReceived(event) {
const eventData = JSON.parse(event.data);
service.debug ? console.log("[socket] message received", eventData) : null;
// Dispatch a generic message
dispatcher.dispatch('message', eventData);
if (eventData.channel) {
dispatcher.dispatch(eventData.channel, eventData);
}
}
async function broadcast(data) {
socket.send(JSON.stringify(data));
}
///////////////////////////////////////////////////
function wait() {
return new Promise(function(resolve) {
function check() {
if (service.connected) {
return resolve(socket);
}
setTimeout(check, 1000);
}
check();
})
}
///////////////////////////////////////////////////
/**
* @description Connect to the socket service
* @alias socket.connect
* @example
* const connected = await sdk.socket.connect()
*
*/
service.connect = async function() {
if (socket) {
if (!socket.connected) {
return wait();
}
console.log('[socket] - Socket is already connected');
return socket;
}
const accessToken = qik.auth.getCurrentToken();
if (!accessToken) {
service.debug ? console.log('[socket] - Must be authenticated to connect to socket') : null;
return;
}
socket = new WebSocket(`${service.url}?access_token=${accessToken}&windowid=${windowID}`);
socket.addEventListener('close', socketClosed);
socket.addEventListener('error', socketError);
socket.addEventListener('open', socketOpened);
socket.addEventListener('message', socketMessageReceived);
service.debug ? console.log('[socket] - Connected with event listeners') : null;
return wait();
}
/**
* @description Reconnect to the socket service
* @alias socket.reconnect
* @example
* const connected = await sdk.socket.reconnect()
*
*/
service.reconnect = async function() {
if(service.connected) {
return;
}
service.debug ? console.log(`[socket] reconnecting...`) : null;
// Reconnect all listeners
const result = await service.connect();
// Resubscribe to all previously connected channels
const subscribedChannels = Object.keys(channels);
subscribedChannels.forEach(function(channel) {
service.subscribe(channel);
})
return result;
}
///////////////////////////////////////////////////
/**
* @description Close the connection to the socket service
* @alias socket.close
* @example
* const closed = await sdk.socket.close()
*
*/
service.close = async function() {
if (!socket) {
console.log('[socket] - Socket is not connected');
return;
}
const previousSetting = service.autoreconnect;
service.autoreconnect = false;
socket.close();
socket.removeEventListener('close', socketClosed);
socket.removeEventListener('error', socketError);
socket.removeEventListener('open', socketOpened);
socket.removeEventListener('message', socketMessageReceived);
socket = null;
socketClosed({ wasClean: true });
service.autoreconnect = previousSetting;
}
service.disconnect = service.close;
///////////////////////////////////////////////////
/**
* @description Subscribe to a socket channel
* @alias socket.subscribe
* @param {String} key The name of the channel to subscribe to
* @example
* const subscribed = await sdk.socket.subscribe('some:cool:alert')
*
*/
service.subscribe = async function(channel) {
if (!socket) {
console.log(`[socket] - Can't subscribe to channel as socket is not connected`);
return;
}
broadcast({
action: 'subscribe',
channel,
})
channels[channel] = true;
dispatcher.dispatch('subscribe', channel)
service.debug ? console.log(`[socket] - subscribed to ${channel}`) : null;
}
///////////////////////////////////////////////////
/**
* @description Unsubscribe from a socket channel
* @alias socket.unsubscribe
* @param {String} key The name of the channel to unsubscribe from
* @example
* const unsubscribed = await sdk.socket.unsubscribe('some:cool:alert')
*
*/
service.unsubscribe = async function(channel) {
if (!socket) {
console.log(`[socket] - Can't unsubscribe from channel as socket is not connected`);
return;
}
broadcast({
action: 'unsubscribe',
channel,
})
delete channels[channel];
dispatcher.dispatch('unsubscribe', channel)
service.debug ? console.log(`[socket] - unsubscribed from ${channel}`) : null;
}
///////////////////////////////////////////////////
// Reconnect the socket whenever the authenticated organisation changes.
// The WebSocket authenticates with an access token at connect time, so the
// connection stays bound to whichever organisation was active when it
// opened. Switching organisations issues a new token but leaves the socket
// untouched, so without this it keeps streaming the previous
// organisation's events. Reconnecting re-authenticates with the new token.
let currentOrganisationID = qik.utils.id(qik.auth.getCurrentUser()?.organisation);
qik.auth.addEventListener('change', async function(user) {
const organisationID = qik.utils.id(user?.organisation);
// Ignore auth changes that don't change the organisation (e.g. routine
// access token refreshes) so we don't churn the connection needlessly.
if (organisationID === currentOrganisationID) {
return;
}
currentOrganisationID = organisationID;
// Nothing to reconnect if the socket was never opened.
if (!socket) {
return;
}
service.debug ? console.log('[socket] - organisation changed, reconnecting') : null;
await service.close();
await service.reconnect();
});
///////////////////////////////////////////////////
return service;
}
export default QikSocket;