qik.socket.js

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;