🎩 You're Invited:Meet the Socket team at Black Hat in Las Vegas, August 3-6.RSVP
Sign In

@dxfeed/dxlink-websocket-client

Package Overview
Dependencies
Maintainers
3
Versions
19
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

@dxfeed/dxlink-websocket-client - npm Package Compare versions

Comparing version
0.8.0
to
0.8.1
+2
-1
build/channel.d.ts

@@ -12,2 +12,3 @@ import { type DXLinkChannel, DXLinkChannelState, type DXLinkChannelMessageListener, type DXLinkChannelStateChangeListener, type DXLinkErrorListener, type DXLinkChannelMessage, type DXLinkError } from '@dxfeed/dxlink-core';

readonly parameters: Record<string, unknown>;
readonly reconnect: boolean;
private readonly sendMessage;

@@ -19,3 +20,3 @@ private status;

private logger;
constructor(id: number, service: string, parameters: Record<string, unknown>, sendMessage: (message: DXLinkWebSocketMessage) => void, config: DXLinkWebSocketClientConfig);
constructor(id: number, service: string, parameters: Record<string, unknown>, reconnect: boolean, sendMessage: (message: DXLinkWebSocketMessage) => void, config: DXLinkWebSocketClientConfig);
send: ({ type, ...payload }: DXLinkChannelMessage) => void;

@@ -22,0 +23,0 @@ addMessageListener: (listener: DXLinkChannelMessageListener) => Set<DXLinkChannelMessageListener>;

@@ -1,2 +0,2 @@

import { type DXLinkScheduler, type DXLinkConnectionDetails, DXLinkConnectionState, type DXLinkConnectionStateChangeListener, type DXLinkErrorListener, DXLinkAuthState, type DXLinkAuthStateChangeListener, type DXLinkChannel, type DXLinkClient } from '@dxfeed/dxlink-core';
import { type DXLinkScheduler, type DXLinkConnectionDetails, DXLinkConnectionState, type DXLinkConnectionStateChangeListener, type DXLinkErrorListener, DXLinkAuthState, type DXLinkAuthStateChangeListener, type DXLinkChannel, type DXLinkChannelOptions, type DXLinkClient } from '@dxfeed/dxlink-core';
import type { DXLinkWebSocketClientConfig } from './config';

@@ -58,3 +58,3 @@ /**

removeErrorListener: (listener: DXLinkErrorListener) => boolean;
openChannel: (service: string, parameters: Record<string, unknown>) => DXLinkChannel;
openChannel: (service: string, parameters: Record<string, unknown>, options?: DXLinkChannelOptions) => DXLinkChannel;
private setConnectionState;

@@ -61,0 +61,0 @@ private sendMessage;

@@ -1,2 +0,2 @@

var dxlinkCore=require("@dxfeed/dxlink-core");class DXLinkWebSocketChannel{constructor(id,service,parameters,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.sendMessage=void 0,this.status=dxlinkCore.DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==dxlinkCore.DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==dxlinkCore.DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(dxlinkCore.DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(dxlinkCore.DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.sendMessage=sendMessage,this.logger=new dxlinkCore.Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.0"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=dxlinkCore.DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=dxlinkCore.DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED&&channel.processStatusRequested();this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(dxlinkCore.DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(dxlinkCore.DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.authState===dxlinkCore.DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(dxlinkCore.DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTING&&this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(dxlinkCore.DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===dxlinkCore.DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:dxlinkCore.DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new dxlinkCore.Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new dxlinkCore.DefaultDXLinkScheduler}},exports.DefaultDXLinkWebSocketConnector=DefaultDXLinkWebSocketConnector;
var dxlinkCore=require("@dxfeed/dxlink-core");class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=dxlinkCore.DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==dxlinkCore.DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==dxlinkCore.DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(dxlinkCore.DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(dxlinkCore.DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new dxlinkCore.Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.1"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=dxlinkCore.DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=dxlinkCore.DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(dxlinkCore.DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(dxlinkCore.DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.authState===dxlinkCore.DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(dxlinkCore.DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTING&&this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(dxlinkCore.DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===dxlinkCore.DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:dxlinkCore.DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new dxlinkCore.Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new dxlinkCore.DefaultDXLinkScheduler}},exports.DefaultDXLinkWebSocketConnector=DefaultDXLinkWebSocketConnector;
//# sourceMappingURL=index.js.map

@@ -1,1 +0,1 @@

{"version":3,"file":"index.js","sources":["../src/channel.ts","../src/connector.ts","../src/client.ts","../src/messages.ts"],"sourcesContent":["import {\n type DXLinkLogger,\n Logger,\n type DXLinkChannel,\n DXLinkChannelState,\n type DXLinkChannelMessageListener,\n type DXLinkChannelStateChangeListener,\n type DXLinkErrorListener,\n type DXLinkChannelMessage,\n type DXLinkError,\n} from '@dxfeed/dxlink-core'\n\nimport { type DXLinkWebSocketClientConfig } from './config'\nimport { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages'\n\n/**\n * A DXLink channel implementation.\n * @internal\n */\nexport class DXLinkWebSocketChannel implements DXLinkChannel {\n private status = DXLinkChannelState.REQUESTED\n\n // Listeners\n private readonly messageListeners = new Set<DXLinkChannelMessageListener>()\n private readonly statusListeners = new Set<DXLinkChannelStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n\n private logger: DXLinkLogger\n\n constructor(\n public readonly id: number,\n public readonly service: string,\n public readonly parameters: Record<string, unknown>,\n private readonly sendMessage: (message: DXLinkWebSocketMessage) => void,\n config: DXLinkWebSocketClientConfig\n ) {\n this.logger = new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel)\n }\n\n send = ({ type, ...payload }: DXLinkChannelMessage) => {\n if (this.status !== DXLinkChannelState.OPENED) {\n throw new Error('Channel is not ready')\n }\n\n this.sendMessage({\n type,\n channel: this.id,\n ...payload,\n })\n }\n\n addMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.add(listener)\n removeMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.delete(listener)\n\n getState = () => this.status\n addStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.add(listener)\n removeStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n error = ({ type, message }: DXLinkError) =>\n this.send({\n type: 'ERROR',\n error: type,\n message,\n })\n\n close = () => {\n if (this.status === DXLinkChannelState.CLOSED) return\n\n this.logger.debug(`Closing by user`)\n\n // We can think that channel is closed already\n this.setStatus(DXLinkChannelState.CLOSED)\n\n this.clear()\n\n this.sendMessage({\n type: 'CHANNEL_CANCEL',\n channel: this.id,\n })\n }\n\n request = () => {\n this.logger.debug('Requesting')\n\n this.sendMessage({\n type: 'CHANNEL_REQUEST',\n channel: this.id,\n service: this.service,\n parameters: this.parameters,\n })\n\n this.processStatusRequested()\n }\n\n processPayloadMessage = (message: ChannelPayloadMessage) => {\n for (const listener of this.messageListeners) {\n listener(message)\n }\n }\n\n processStatusOpened = () => {\n this.logger.debug('Opened')\n\n this.setStatus(DXLinkChannelState.OPENED)\n }\n\n processStatusRequested = () => {\n this.setStatus(DXLinkChannelState.REQUESTED)\n }\n\n processStatusClosed = () => {\n this.logger.debug('Closed by remote endpoint')\n\n this.setStatus(DXLinkChannelState.CLOSED)\n this.clear()\n }\n\n processError = (error: DXLinkError) => {\n if (this.errorListeners.size === 0) {\n this.logger.error(`Unhandled error in channel#${this.id}: `, error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} error listener: `, e)\n }\n }\n }\n\n private setStatus = (newStatus: DXLinkChannelState) => {\n if (this.status === newStatus) return\n\n const prev = this.status\n this.status = newStatus\n for (const listener of this.statusListeners) {\n try {\n listener(newStatus, prev)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} status listener: `, e)\n }\n }\n }\n\n private clear = () => {\n this.messageListeners.clear()\n this.statusListeners.clear()\n // TODO: rethink approach when error came after channel is closed\n // this.errorListeners.clear()\n }\n}\n","import { type DXLinkWebSocketMessage } from './messages'\n\n/**\n * Interface for a WebSocket connector that manages the connection to a WebSocket server.\n * It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.\n */\nexport interface DXLinkWebSocketConnector {\n /**\n * Returns the URL of the WebSocket connection.\n */\n getUrl(): string\n /**\n * Starts the WebSocket connection.\n */\n start(): void\n /**\n * Stops the WebSocket connection and cleans up resources.\n */\n stop(): void\n /**\n * Sends a message over the WebSocket connection.\n * @param message The message to send.\n */\n sendMessage(message: DXLinkWebSocketMessage): void\n /**\n * Sets a listener that is called when the WebSocket connection is opened.\n * @param listener The listener function to call when the connection is opened.\n */\n setOpenListener(listener: () => void): void\n /**\n * Sets a listener that is called when the WebSocket connection is closed.\n * @param listener The listener function to call when the connection is closed.\n */\n setCloseListener(listener: DXLinkWebSocketCloseListener): void\n /**\n * Sets a listener that is called when a message is received from the WebSocket server.\n * @param listener The listener function to call when a message is received.\n */\n setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void\n}\n\n/**\n * Type for a close listener that is called when the WebSocket connection is closed.\n * @param reason - The reason for the closure.\n * @param error - Indicates if the closure was due to an error.\n */\nexport type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void\n\n/**\n * Default connector for the WebSocket connection.\n * @internal\n */\nexport class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: DXLinkWebSocketCloseListener | undefined = undefined\n private messageListener: ((message: DXLinkWebSocketMessage) => void) | undefined = undefined\n\n constructor(\n private readonly url: string,\n private readonly protocols?: string | string[]\n ) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url, this.protocols)\n\n this.socket.addEventListener('open', this.handleOpen)\n this.socket.addEventListener('error', this.handleError)\n this.socket.addEventListener('close', this.handleClosed)\n }\n\n stop() {\n if (this.socket === undefined) return\n\n this.socket.removeEventListener('open', this.handleOpen)\n this.socket.removeEventListener('error', this.handleError)\n this.socket.removeEventListener('close', this.handleClosed)\n this.socket.removeEventListener('message', this.handleMessage)\n\n this.socket.close()\n this.socket = undefined\n this.isAvailable = false\n }\n\n sendMessage = (message: DXLinkWebSocketMessage) => {\n if (this.socket === undefined || !this.isAvailable) {\n return\n }\n\n this.socket.send(JSON.stringify(message))\n }\n\n setOpenListener = (listener: () => void) => {\n this.openListener = listener\n }\n\n setCloseListener = (listener: DXLinkWebSocketCloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: DXLinkWebSocketMessage) => void) => {\n this.messageListener = listener\n }\n\n getUrl = () => this.url\n\n private handleOpen = () => {\n if (this.socket === undefined) return\n\n this.isAvailable = true\n\n this.socket.removeEventListener('open', this.handleOpen)\n\n this.socket.addEventListener('message', this.handleMessage)\n\n this.openListener?.()\n }\n\n private handleClosed = (ev: CloseEvent) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.(ev.reason, false)\n }\n\n private handleError = (_ev: Event) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.('Unable to connect', true)\n }\n\n private handleMessage = (ev: MessageEvent) => {\n try {\n const message = JSON.parse(ev.data)\n if (typeof message !== 'object') {\n throw new Error('Unexpected message: ' + typeof message)\n }\n if (typeof message.type !== 'string') {\n throw new Error('Unexpected message type: ' + typeof message.type)\n }\n if (typeof message.channel !== 'number') {\n throw new Error('Unexpected message channel: ' + typeof message.channel)\n }\n\n // TODO: validate message properly (e.g. check for type and other fields)\n\n if (this.messageListener === undefined) {\n return console.warn('No message listener set')\n }\n\n this.messageListener(message)\n } catch (error) {\n console.error(error instanceof Error ? error : new Error('Parsing error:' + String(error)))\n }\n }\n}\n","import {\n DXLinkLogLevel,\n type DXLinkLogger,\n Logger,\n type DXLinkScheduler,\n DefaultDXLinkScheduler,\n type DXLinkConnectionDetails,\n DXLinkConnectionState,\n type DXLinkConnectionStateChangeListener,\n type DXLinkErrorListener,\n DXLinkAuthState,\n DXLinkChannelState,\n type DXLinkAuthStateChangeListener,\n type DXLinkChannel,\n type DXLinkError,\n type DXLinkClient,\n} from '@dxfeed/dxlink-core'\n\nimport { DXLinkWebSocketChannel } from './channel'\nimport type { DXLinkWebSocketClientConfig } from './config'\nimport { type DXLinkWebSocketConnector, DefaultDXLinkWebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type DXLinkWebSocketMessage,\n type SetupMessage,\n isChannelLifecycleMessage,\n isChannelMessage,\n isConnectionMessage,\n} from './messages'\nimport { VERSION } from './version'\n\n/**\n * Protocol version that is used by client.\n */\nexport const DXLINK_WS_PROTOCOL_VERSION = '0.1'\n\nconst CLIENT_VERSION = `DXF-JS/${VERSION}`\n\n// Scheduler keys\nconst DXLWS_SCHEDULER_KEY_RECONNECT = 'DXLWS_RECONNECT'\nconst DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = 'DXLWS_SETUP_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = 'DXLWS_AUTH_STATE_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_TIMEOUT = 'DXLWS_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_KEEPALIVE = 'DXLWS_KEEPALIVE'\n\nconst DEFAULT_CONNECTION_DETAILS: DXLinkConnectionDetails = {\n protocolVersion: DXLINK_WS_PROTOCOL_VERSION,\n clientVersion: CLIENT_VERSION,\n}\n\n/**\n * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.\n */\nexport class DXLinkWebSocketClient implements DXLinkClient {\n private readonly config: DXLinkWebSocketClientConfig\n\n private readonly logger: DXLinkLogger\n\n private readonly scheduler: DXLinkScheduler\n\n private connector: DXLinkWebSocketConnector | undefined\n\n private connectionState: DXLinkConnectionState = DXLinkConnectionState.NOT_CONNECTED\n private connectionDetails: DXLinkConnectionDetails = DEFAULT_CONNECTION_DETAILS\n\n private authState: DXLinkAuthState = DXLinkAuthState.UNAUTHORIZED\n\n // Listeners\n private readonly connectionStateChangeListeners = new Set<DXLinkConnectionStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n private readonly authStateChangeListeners = new Set<DXLinkAuthStateChangeListener>()\n\n /**\n * Authorization type that was determined by server behavior during setup phase.\n * This value is used to determine if authorization is required or optional or not defined yet.\n */\n private isFirstAuthState = true\n /**\n * Last setted auth token that will be sent to server after connection is established or re-established.\n */\n private lastSettedAuthToken: string | undefined\n\n // Stats for keepalive\n // TODO: mb move to connector\n private lastReceivedMillis = 0\n private lastSentMillis = 0\n\n /**\n * Count of reconnect attempts since last successful connection.\n */\n private reconnectAttempts = 0\n\n // Channels\n private globalChannelId = 1\n private readonly channels = new Map<number, DXLinkWebSocketChannel>()\n\n /**\n * Create new instance of {@link DXLinkWebSocketClient}.\n * @param config Configuration of the client.\n */\n constructor(config?: Partial<DXLinkWebSocketClientConfig>) {\n this.config = {\n keepaliveInterval: 30,\n keepaliveTimeout: 60,\n acceptKeepaliveTimeout: 60,\n actionTimeout: 10,\n logLevel: DXLinkLogLevel.WARN,\n maxReconnectAttempts: -1,\n connectorFactory: (url) => new DefaultDXLinkWebSocketConnector(url),\n ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\n this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler()\n }\n\n connect = (url: string) => {\n // Do nothing if already connected to the same url\n if (this.connector?.getUrl() === url) return\n\n // Disconnect from previous connection if any exists\n this.disconnect()\n\n this.logger.debug('Connecting to', url)\n\n // Immediately set connection state to CONNECTING\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n\n // Create new connector\n this.connector = this.config.connectorFactory(url)\n this.connector.setOpenListener(this.processTransportOpen)\n this.connector.setMessageListener(this.processMessage)\n this.connector.setCloseListener(this.processTransportClose)\n\n // Initiate websocket connection\n this.connector.start()\n }\n\n reconnect = () => {\n if (\n this.config.maxReconnectAttempts >= 0 &&\n this.reconnectAttempts >= this.config.maxReconnectAttempts\n ) {\n this.logger.warn('Max reconnect attempts reached')\n\n this.disconnect()\n return\n }\n\n if (\n this.connectionState === DXLinkConnectionState.NOT_CONNECTED ||\n this.connector === undefined\n )\n return\n\n this.logger.debug('Trying to reconnect', this.connector.getUrl())\n\n this.connector.stop()\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n\n // Increase reconnect attempts counter\n this.reconnectAttempts++\n\n // Update state for connection and channels\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n for (const channel of this.channels.values()) {\n if (channel.getState() === DXLinkChannelState.CLOSED) continue\n\n channel.processStatusRequested()\n }\n\n // Schedule reconnect attempt after some time\n // Additionally, task will be executed in case when tab is active again\n // coz browser sometimes doesn't run scheduled tasks when tab is inactive\n this.scheduler.schedule(\n () => {\n if (this.connector === undefined) return\n\n // Start new connection attempt\n this.connector.start()\n },\n this.reconnectAttempts * 1000,\n DXLWS_SCHEDULER_KEY_RECONNECT\n )\n }\n\n disconnect = () => {\n if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return\n\n this.logger.debug('Disconnecting')\n\n // Destroy connector\n this.connector?.stop()\n this.connector = undefined\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n this.reconnectAttempts = 0\n\n this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED)\n this.setAuthState(DXLinkAuthState.UNAUTHORIZED)\n }\n\n close = () => {\n this.disconnect()\n }\n\n getConnectionDetails = () => this.connectionDetails\n getConnectionState = () => this.connectionState\n getScheduler = () => this.scheduler\n addConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.add(listener)\n removeConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.delete(listener)\n\n setAuthToken = (token: string): void => {\n this.lastSettedAuthToken = token\n\n if (this.connectionState === DXLinkConnectionState.CONNECTED) {\n this.sendAuthMessage(token)\n }\n }\n\n getAuthState = (): DXLinkAuthState => this.authState\n addAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.add(listener)\n removeAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n openChannel = (service: string, parameters: Record<string, unknown>): DXLinkChannel => {\n const channelId = this.globalChannelId\n this.globalChannelId += 2\n\n const channel = new DXLinkWebSocketChannel(\n channelId,\n service,\n parameters,\n this.sendMessage,\n this.config\n )\n\n this.channels.set(channelId, channel)\n\n // Send channel request if connection is already established\n if (\n this.connectionState === DXLinkConnectionState.CONNECTED &&\n this.authState === DXLinkAuthState.AUTHORIZED\n ) {\n channel.request()\n }\n\n return channel\n }\n\n private setConnectionState = (newStatus: DXLinkConnectionState) => {\n const prev = this.connectionState\n if (prev === newStatus) return\n\n this.connectionState = newStatus\n for (const listener of this.connectionStateChangeListeners) {\n listener(newStatus, prev)\n }\n }\n\n private sendMessage = (message: DXLinkWebSocketMessage): void => {\n this.connector?.sendMessage(message)\n\n this.scheduleKeepalive()\n\n // TODO: mb move to connector\n this.lastSentMillis = Date.now()\n }\n\n private sendAuthMessage = (token: string): void => {\n this.logger.debug('Sending auth message')\n\n this.setAuthState(DXLinkAuthState.AUTHORIZING)\n\n this.sendMessage({\n type: 'AUTH',\n channel: 0,\n token,\n })\n }\n\n private setAuthState = (newState: DXLinkAuthState): void => {\n const prev = this.authState\n\n this.authState = newState\n for (const listener of this.authStateChangeListeners) {\n try {\n listener(newState, prev)\n } catch (e) {\n this.logger.error('Auth state listener error', e)\n }\n }\n }\n\n private processMessage = (message: DXLinkWebSocketMessage): void => {\n this.lastReceivedMillis = Date.now()\n\n // Send keepalive message if no messages sent for a while (keepaliveInterval)\n // Because browser sometimes doesn't run scheduled tasks when tab is inactive\n if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1000) {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n }\n\n // Connection messages are messages that are sent to the channel 0\n if (isConnectionMessage(message)) {\n switch (message.type) {\n case 'SETUP':\n return this.processSetupMessage(message)\n case 'AUTH_STATE':\n return this.processAuthStateMessage(message)\n case 'ERROR':\n return this.publishError({\n type: message.error,\n message: message.message,\n })\n case 'KEEPALIVE':\n // Ignore keepalive messages coz they are used only to maintain connection\n return\n }\n } else if (isChannelMessage(message)) {\n const channel = this.channels.get(message.channel)\n if (channel === undefined) {\n this.logger.warn('Received lifecycle message for unknown channel', message)\n return\n }\n\n if (isChannelLifecycleMessage(message)) {\n switch (message.type) {\n case 'CHANNEL_OPENED':\n return channel.processStatusOpened()\n case 'CHANNEL_CLOSED':\n return channel.processStatusClosed()\n case 'ERROR':\n return channel.processError({\n type: message.error,\n message: message.message,\n })\n }\n return\n }\n\n return channel.processPayloadMessage(message)\n }\n\n this.logger.warn('Unhandeled message', message.type)\n }\n\n private processSetupMessage = (serverSetup: SetupMessage): void => {\n // Clear setup timeout check from connect method\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT)\n\n // Mark connection as connected after first setup message and subsequent ones\n if (\n this.connectionState === DXLinkConnectionState.CONNECTING ||\n this.connectionState === DXLinkConnectionState.CONNECTED\n ) {\n this.connectionDetails = {\n ...this.connectionDetails,\n serverVersion: serverSetup.version,\n clientKeepaliveTimeout: this.config.keepaliveTimeout,\n serverKeepaliveTimeout: serverSetup.keepaliveTimeout,\n }\n\n // Reset reconnect attempts counter after successful connection\n this.reconnectAttempts = 0\n\n if (this.lastSettedAuthToken === undefined) {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n }\n }\n\n // Connection maintance: Setup keepalive timeout check\n const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1000\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n timeoutMills,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private publishError = (error: DXLinkError): void => {\n this.logger.debug('Publishing error', error)\n\n if (this.errorListeners.size === 0) {\n this.logger.error('Unhandled dxLink error', error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error('Error listener error', e)\n }\n }\n }\n\n private processAuthStateMessage = ({ state }: AuthStateMessage): void => {\n this.logger.debug('Received auth state message', state)\n\n // Clear auth state timeout check\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT)\n\n // Ignore first auth state message because it is sent during connection setup\n if (this.isFirstAuthState) {\n this.isFirstAuthState = false\n } else {\n // Reset auth token if server rejected it\n if (state === 'UNAUTHORIZED') {\n this.lastSettedAuthToken = undefined\n }\n }\n\n // Request active channels if connection is authorized\n if (state === 'AUTHORIZED') {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n\n this.requestActiveChannels()\n }\n\n this.setAuthState(DXLinkAuthState[state])\n }\n\n private requestActiveChannels = (): void => {\n for (const channel of this.channels.values()) {\n // clear closed channels\n if (channel.getState() === DXLinkChannelState.CLOSED) {\n this.channels.delete(channel.id)\n continue\n }\n\n channel.request()\n }\n }\n\n /**\n * Process transport open event from connector.\n * After transport is opened:\n * - setup message is sent to server\n * - auth message is sent to server if auth token is set\n * - wait for setup message from server\n * - wait for auth state message from server\n */\n private processTransportOpen = (): void => {\n this.logger.debug('Connection opened')\n\n const setupMessage: SetupMessage = {\n type: 'SETUP',\n channel: 0,\n version: `${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,\n keepaliveTimeout: this.config.keepaliveTimeout,\n acceptKeepaliveTimeout: this.config.acceptKeepaliveTimeout,\n }\n\n // Setup timeout check\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No setup message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no setup message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT\n )\n\n this.sendMessage(setupMessage)\n\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No auth state message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no auth state message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT\n )\n\n if (this.lastSettedAuthToken !== undefined) {\n this.sendAuthMessage(this.lastSettedAuthToken)\n }\n }\n\n /**\n * Process transport close event from connector.\n * After transport is closed by server:\n * - reconnect if connection is authorized\n * - disconnect if connection is not authorized\n */\n private processTransportClose = (reason: string, error: boolean): void => {\n this.logger.debug('Connection closed', reason)\n\n if (error !== undefined) {\n this.publishError({\n type: 'UNKNOWN',\n message: reason,\n })\n }\n\n if (this.authState === DXLinkAuthState.UNAUTHORIZED) {\n this.lastSettedAuthToken = undefined\n this.disconnect()\n return\n }\n\n this.reconnect()\n }\n\n private timeoutCheck = (timeoutMills: number) => {\n const now = Date.now()\n const noKeepaliveDuration = now - this.lastReceivedMillis\n if (noKeepaliveDuration >= timeoutMills) {\n this.sendMessage({\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No keepalive received for ' + noKeepaliveDuration + 'ms',\n })\n\n return this.reconnect()\n }\n\n const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration)\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n nextTimeout,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private scheduleKeepalive = () => {\n this.scheduler.schedule(\n () => {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n\n this.scheduleKeepalive()\n },\n this.config.keepaliveInterval * 1000,\n DXLWS_SCHEDULER_KEY_KEEPALIVE\n )\n }\n}\n","export interface AuthMessage {\n type: 'AUTH'\n channel: 0\n token: string\n}\n\nexport type AuthState = 'AUTHORIZED' | 'UNAUTHORIZED'\n\nexport interface AuthStateMessage {\n type: 'AUTH_STATE'\n channel: 0\n state: AuthState\n}\n\nexport interface SetupMessage {\n type: 'SETUP'\n channel: 0\n version: string\n keepaliveTimeout?: number\n acceptKeepaliveTimeout?: number\n}\n\nexport interface KeepaliveMessage {\n type: 'KEEPALIVE'\n channel: 0\n}\n\nexport interface ChannelRequestMessage {\n type: 'CHANNEL_REQUEST'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelCancelMessage {\n type: 'CHANNEL_CANCEL'\n channel: number\n}\n\nexport interface ChannelOpenedMessage {\n type: 'CHANNEL_OPENED'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelClosedMessage {\n type: 'CHANNEL_CLOSED'\n channel: number\n}\n\nexport type ErrorType =\n | 'UNKNOWN'\n | 'UNSUPPORTED_PROTOCOL'\n | 'TIMEOUT'\n | 'UNAUTHORIZED'\n | 'INVALID_MESSAGE'\n | 'BAD_ACTION'\n\nexport interface ErrorMessage {\n type: 'ERROR'\n channel: 0\n error: ErrorType\n message: string\n}\n\nexport interface ChannelErrorMessage {\n type: 'ERROR'\n channel: number\n error: ErrorType\n message: string\n}\n\nexport interface ChannelPayloadMessage {\n type: string\n channel: number\n [key: string]: unknown\n}\n\nexport type ConnectionMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n\nexport type DXLinkWebSocketMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelPayloadMessage\n | ChannelErrorMessage\n\nexport const isConnectionMessage = (\n message: DXLinkWebSocketMessage\n): message is ConnectionMessage =>\n message.channel === 0 &&\n (message.type === 'SETUP' ||\n message.type === 'KEEPALIVE' ||\n message.type === 'AUTH' ||\n message.type === 'AUTH_STATE' ||\n message.type === 'ERROR')\n\nexport type ChannelLifecycleMessage =\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n\nexport type ChannelMessage = ChannelLifecycleMessage | ChannelPayloadMessage\n\nexport const isChannelMessage = (message: DXLinkWebSocketMessage): message is ChannelMessage =>\n message.channel !== 0\n\nexport const isChannelLifecycleMessage = (\n message: ChannelMessage\n): message is ChannelLifecycleMessage =>\n message.type === 'CHANNEL_OPENED' ||\n message.type === 'CHANNEL_CLOSED' ||\n message.type === 'ERROR' ||\n message.type === 'CHANNEL_REQUEST' ||\n message.type === 'CHANNEL_CANCEL'\n"],"names":["DXLinkWebSocketChannel","constructor","id","service","parameters","sendMessage","config","this","status","DXLinkChannelState","REQUESTED","messageListeners","Set","statusListeners","errorListeners","logger","send","type","payload","OPENED","Error","channel","addMessageListener","listener","add","removeMessageListener","delete","getState","addStateChangeListener","removeStateChangeListener","addErrorListener","removeErrorListener","error","message","close","CLOSED","debug","setStatus","clear","request","processStatusRequested","processPayloadMessage","processStatusOpened","processStatusClosed","processError","size","e","newStatus","prev","Logger","name","logLevel","DefaultDXLinkWebSocketConnector","url","protocols","socket","undefined","isAvailable","openListener","closeListener","messageListener","JSON","stringify","setOpenListener","setCloseListener","setMessageListener","getUrl","handleOpen","removeEventListener","addEventListener","handleMessage","handleClosed","ev","stop","reason","handleError","_ev","parse","data","console","warn","String","start","WebSocket","DEFAULT_CONNECTION_DETAILS","protocolVersion","clientVersion","scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","connectorFactory","processTransportOpen","processMessage","processTransportClose","reconnect","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","getScheduler","addConnectionStateChangeListener","removeConnectionStateChangeListener","setAuthToken","token","CONNECTED","sendAuthMessage","getAuthState","addAuthStateChangeListener","removeAuthStateChangeListener","openChannel","channelId","set","AUTHORIZED","scheduleKeepalive","Date","now","AUTHORIZING","newState","keepaliveInterval","isConnectionMessage","processSetupMessage","processAuthStateMessage","publishError","isChannelMessage","get","isChannelLifecycleMessage","serverSetup","cancel","serverVersion","version","clientKeepaliveTimeout","keepaliveTimeout","serverKeepaliveTimeout","timeoutMills","timeoutCheck","state","requestActiveChannels","setupMessage","acceptKeepaliveTimeout","errorMessage","actionTimeout","noKeepaliveDuration","nextTimeout","Math","max","DXLinkLogLevel","WARN","DefaultDXLinkScheduler"],"mappings":"oDAmBaA,uBAUXC,WAAAA,CACkBC,GACAC,QACAC,WACCC,YACjBC,aAJgBJ,QAAA,EAAAK,KACAJ,aAAA,EAAAI,KACAH,gBACCC,EAAAA,KAAAA,iBAbXG,EAAAA,KAAAA,OAASC,8BAAmBC,eAGnBC,iBAAmB,IAAIC,IACvBC,KAAAA,gBAAkB,IAAID,IAAuCL,KAC7DO,eAAiB,IAAIF,IAE9BG,KAAAA,mBAYRC,KAAO,EAAGC,aAASC,YACjB,GAAIX,KAAKC,SAAWC,WAAAA,mBAAmBU,OACrC,MAAM,IAAIC,MAAM,wBAGlBb,KAAKF,YAAY,CACfY,UACAI,QAASd,KAAKL,MACXgB,SACJ,EACFX,KAEDe,mBAAsBC,UACpBhB,KAAKI,iBAAiBa,IAAID,UAC5BE,KAAAA,sBAAyBF,UACvBhB,KAAKI,iBAAiBe,OAAOH,UAE/BI,KAAAA,SAAW,IAAMpB,KAAKC,OACtBoB,KAAAA,uBAA0BL,UACxBhB,KAAKM,gBAAgBW,IAAID,UAAShB,KACpCsB,0BAA6BN,UAC3BhB,KAAKM,gBAAgBa,OAAOH,eAE9BO,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAC9EQ,KAAAA,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7FyB,MAAQ,EAAGf,UAAMgB,mBACf1B,KAAKS,KAAK,CACRC,KAAM,QACNe,MAAOf,KACPgB,kBAGJC,KAAAA,MAAQ,KACF3B,KAAKC,SAAWC,WAAAA,mBAAmB0B,SAEvC5B,KAAKQ,OAAOqB,MAAM,mBAGlB7B,KAAK8B,UAAU5B,WAAAA,mBAAmB0B,QAElC5B,KAAK+B,QAEL/B,KAAKF,YAAY,CACfY,KAAM,iBACNI,QAASd,KAAKL,KAElB,OAEAqC,QAAU,KACRhC,KAAKQ,OAAOqB,MAAM,cAElB7B,KAAKF,YAAY,CACfY,KAAM,kBACNI,QAASd,KAAKL,GACdC,QAASI,KAAKJ,QACdC,WAAYG,KAAKH,aAGnBG,KAAKiC,wBAAsB,OAG7BC,sBAAyBR,UACvB,IAAK,MAAMV,YAAgBhB,KAACI,iBAC1BY,SAASU,QACV,EACF1B,KAEDmC,oBAAsB,KACpBnC,KAAKQ,OAAOqB,MAAM,UAElB7B,KAAK8B,UAAU5B,8BAAmBU,SAGpCqB,KAAAA,uBAAyB,KACvBjC,KAAK8B,UAAU5B,WAAAA,mBAAmBC,UAAS,OAG7CiC,oBAAsB,KACpBpC,KAAKQ,OAAOqB,MAAM,6BAElB7B,KAAK8B,UAAU5B,WAAkBA,mBAAC0B,QAClC5B,KAAK+B,cAGPM,aAAgBZ,QACd,GAAiC,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAYhB,KAAKO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKL,sBAAuB4C,EACnE,MATDvC,KAAKQ,OAAOiB,MAAM,8BAA8BzB,KAAKL,OAAQ8B,MAU9D,EAGKK,KAAAA,UAAaU,YACnB,GAAIxC,KAAKC,SAAWuC,UAAW,OAE/B,MAAMC,KAAOzC,KAAKC,OAClBD,KAAKC,OAASuC,UACd,IAAK,MAAMxB,YAAYhB,KAAKM,gBAC1B,IACEU,SAASwB,UAAWC,KACrB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKL,uBAAwB4C,EACpE,CACF,EAGKR,KAAAA,MAAQ,KACd/B,KAAKI,iBAAiB2B,QACtB/B,KAAKM,gBAAgByB,OAGvB,EAhIkB/B,KAAEL,GAAFA,GACAK,KAAOJ,QAAPA,QACAI,KAAUH,WAAVA,WACCG,KAAWF,YAAXA,YAGjBE,KAAKQ,OAAS,IAAIkC,WAAMA,OAAC,GAAGjD,uBAAuBkD,QAAQhD,MAAMC,UAAWG,OAAO6C,SACrF,QCeWC,gCASXnD,WAAAA,CACmBoD,IACAC,WAA6B/C,KAD7B8C,SAAA,EAAA9C,KACA+C,eAAA,EAAA/C,KAVXgD,YAAgCC,EAASjD,KAEzCkD,aAAc,EAAKlD,KAEnBmD,kBAAyCF,EACzCG,KAAAA,mBAA0DH,EAC1DI,KAAAA,qBAA2EJ,EA8BnFnD,KAAAA,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAAC1B,KAEDwD,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QAAAA,EACrBhB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAEA0C,KAAAA,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QAAAA,EAGzB2C,KAAAA,OAAS,IAAM3D,KAAK8C,IAEZc,KAAAA,WAAa,UACCX,IAAhBjD,KAAKgD,SAEThD,KAAKkD,aAAc,EAEnBlD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAE7C5D,KAAKgD,OAAOc,iBAAiB,UAAW9D,KAAK+D,eAE7C/D,KAAKmD,iBACP,EAACnD,KAEOgE,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAACpD,KAEO+D,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,MAAM,IAAIb,MAAM,8BAAgCa,SAElD,GAA4B,iBAAjBA,QAAQhB,KACjB,MAAU,IAAAG,MAAM,mCAAqCa,QAAQhB,MAE/D,GAA+B,iBAApBgB,QAAQZ,QACjB,MAAM,IAAID,MAAM,sCAAwCa,QAAQZ,SAKlE,QAA6BmC,IAAzBjD,KAAKqD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBzE,KAAKqD,gBAAgB3B,QACtB,CAAC,MAAOD,OACP+C,QAAQ/C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB6D,OAAOjD,QACpF,GAnGgBzB,KAAG8C,IAAHA,IACA9C,KAAS+C,UAATA,SAChB,CAEH4B,KAAAA,QACsB1B,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAS,IAAI4B,UAAU5E,KAAK8C,IAAK9C,KAAK+C,WAE3C/C,KAAKgD,OAAOc,iBAAiB,OAAQ9D,KAAK4D,YAC1C5D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKoE,aAC3CpE,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAC7C5D,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKoE,aAC9CpE,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKgE,cAC9ChE,KAAKgD,OAAOa,oBAAoB,UAAW7D,KAAK+D,eAEhD/D,KAAKgD,OAAOrB,QACZ3B,KAAKgD,YAASC,EACdjD,KAAKkD,aAAc,EACrB,ECpDW,MAWP2B,2BAAsD,CAC1DC,gBAZwC,MAaxCC,cAXqB,mDAFmB,0CAkExCrF,WAAAA,CAAYK,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENgF,eAAS,EAAAhF,KAElBiF,eAAS,EAAAjF,KAETkF,gBAAyCC,WAAqBA,sBAACC,cAAapF,KAC5EqF,kBAA6CR,2BAA0B7E,KAEvEsF,UAA6BC,WAAAA,gBAAgBC,aAGpCC,KAAAA,+BAAiC,IAAIpF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBqF,KAAAA,yBAA2B,IAAIrF,IAMxCsF,KAAAA,kBAAmB,EAInBC,KAAAA,yBAIAC,EAAAA,KAAAA,mBAAqB,EACrBC,KAAAA,eAAiB,EAKjBC,KAAAA,kBAAoB,EAGpBC,KAAAA,gBAAkB,EACTC,KAAAA,SAAW,IAAIC,IAsBhCC,KAAAA,QAAWrD,MAEL9C,KAAKiF,WAAWtB,WAAab,MAGjC9C,KAAKoG,aAELpG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKqG,mBAAmBlB,WAAqBA,sBAACmB,YAG9CtG,KAAKiF,UAAYjF,KAAKD,OAAOwG,iBAAiBzD,KAC9C9C,KAAKiF,UAAUzB,gBAAgBxD,KAAKwG,sBACpCxG,KAAKiF,UAAUvB,mBAAmB1D,KAAKyG,gBACvCzG,KAAKiF,UAAUxB,iBAAiBzD,KAAK0G,uBAGrC1G,KAAKiF,UAAUN,QACjB,EAAC3E,KAED2G,UAAY,KACV,GACE3G,KAAKD,OAAO6G,sBAAwB,GACpC5G,KAAK+F,mBAAqB/F,KAAKD,OAAO6G,qBAKtC,OAHA5G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKoG,aAIP,GACEpG,KAAKkF,kBAAoBC,WAAqBA,sBAACC,oBAC5BnC,IAAnBjD,KAAKiF,UAFP,CAMAjF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKiF,UAAUtB,UAExD3D,KAAKiF,UAAUf,OAGflE,KAAKgF,UAAUjD,QAGf/B,KAAKqF,kBAAoBR,2BACzB7E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EAGxB3F,KAAK+F,oBAGL/F,KAAKqG,mBAAmBlB,WAAAA,sBAAsBmB,YAC9C,IAAK,MAAMxF,WAAWd,KAAKiG,SAASY,SAC9B/F,QAAQM,aAAelB,WAAkBA,mBAAC0B,QAE9Cd,QAAQmB,yBAMVjC,KAAKgF,UAAU8B,SACb,UACyB7D,IAAnBjD,KAAKiF,WAGTjF,KAAKiF,UAAUN,OAAK,EAEG,IAAzB3E,KAAK+F,kBAtJ2B,kBAkHhC,CAuCJ,EAAC/F,KAEDoG,WAAa,KACPpG,KAAKkF,kBAAoBC,WAAqBA,sBAACC,gBAEnDpF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKiF,WAAWf,OAChBlE,KAAKiF,eAAYhC,EAGjBjD,KAAKgF,UAAUjD,QAGf/B,KAAKqF,kBAAoBR,2BACzB7E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EACxB3F,KAAK+F,kBAAoB,EAEzB/F,KAAKqG,mBAAmBlB,WAAqBA,sBAACC,eAC9CpF,KAAK+G,aAAaxB,WAAAA,gBAAgBC,cACpC,EAACxF,KAED2B,MAAQ,KACN3B,KAAKoG,YACP,EAEAY,KAAAA,qBAAuB,IAAMhH,KAAKqF,kBAAiBrF,KACnDiH,mBAAqB,IAAMjH,KAAKkF,gBAAelF,KAC/CkH,aAAe,IAAMlH,KAAKgF,UAC1BmC,KAAAA,iCAAoCnG,UAClChB,KAAKyF,+BAA+BxE,IAAID,UAAShB,KACnDoH,oCAAuCpG,UACrChB,KAAKyF,+BAA+BtE,OAAOH,UAE7CqG,KAAAA,aAAgBC,QACdtH,KAAK4F,oBAAsB0B,MAEvBtH,KAAKkF,kBAAoBC,WAAqBA,sBAACoC,WACjDvH,KAAKwH,gBAAgBF,MACtB,EACFtH,KAEDyH,aAAe,IAAuBzH,KAAKsF,UAC3CoC,KAAAA,2BAA8B1G,UAC5BhB,KAAK0F,yBAAyBzE,IAAID,UAAShB,KAC7C2H,8BAAiC3G,UAC/BhB,KAAK0F,yBAAyBvE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAEpF4G,KAAAA,YAAc,CAAChI,QAAiBC,cAC9B,MAAMgI,UAAY7H,KAAKgG,gBACvBhG,KAAKgG,iBAAmB,EAExB,MAAMlF,QAAU,IAAIrB,uBAClBoI,UACAjI,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKiG,SAAS6B,IAAID,UAAW/G,SAI3Bd,KAAKkF,kBAAoBC,WAAAA,sBAAsBoC,WAC/CvH,KAAKsF,YAAcC,WAAAA,gBAAgBwC,YAEnCjH,QAAQkB,UAGHlB,SAGDuF,KAAAA,mBAAsB7D,YAC5B,MAAMC,KAAOzC,KAAKkF,gBAClB,GAAIzC,OAASD,UAAb,CAEAxC,KAAKkF,gBAAkB1C,UACvB,IAAK,MAAMxB,YAAYhB,KAAKyF,+BAC1BzE,SAASwB,UAAWC,KAFtB,CAGC,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKiF,WAAWnF,YAAY4B,SAE5B1B,KAAKgI,oBAGLhI,KAAK8F,eAAiBmC,KAAKC,KAC7B,EAEQV,KAAAA,gBAAmBF,QACzBtH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAK+G,aAAaxB,WAAeA,gBAAC4C,aAElCnI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTwG,aACD,EACFtH,KAEO+G,aAAgBqB,WACtB,MAAM3F,KAAOzC,KAAKsF,UAElBtF,KAAKsF,UAAY8C,SACjB,IAAK,MAAMpH,YAAgBhB,KAAC0F,yBAC1B,IACE1E,SAASoH,SAAU3F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EACFvC,KAEOyG,eAAkB/E,UAaxB,GAZA1B,KAAK6F,mBAAqBoC,KAAKC,MAI3BlI,KAAK6F,mBAAqB7F,KAAK8F,gBAAkD,IAAhC9F,KAAKD,OAAOsI,mBAC/DrI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IChOfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MD8NJ4H,CAAoB5G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACuI,oBAAoB7G,SAClC,IAAK,aACH,OAAW1B,KAACwI,wBAAwB9G,SACtC,IAAK,QACH,OAAO1B,KAAKyI,aAAa,CACvB/H,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GClOsBA,UACX,IAApBA,QAAQZ,QDiOK4H,CAAiBhH,SAAU,CACpC,MAAMZ,QAAUd,KAAKiG,SAAS0C,IAAIjH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GCrOJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KD+NAkI,CAA0BlH,SAAU,CACtC,OAAQA,QAAQhB,MACd,IAAK,iBACH,OAAOI,QAAQqB,sBACjB,IAAK,iBACH,OAAOrB,QAAQsB,sBACjB,IAAK,QACH,OAAOtB,QAAQuB,aAAa,CAC1B3B,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAGvB,MACD,CAED,OAAOZ,QAAQoB,sBAAsBR,QACtC,CAED1B,KAAKQ,OAAOiE,KAAK,qBAAsB/C,QAAQhB,KACjD,EAEQ6H,KAAAA,oBAAuBM,cAE7B7I,KAAKgF,UAAU8D,OA7UuB,uBAiVpC9I,KAAKkF,kBAAoBC,WAAqBA,sBAACmB,YAC/CtG,KAAKkF,kBAAoBC,WAAqBA,sBAACoC,YAE/CvH,KAAKqF,kBAAoB,IACpBrF,KAAKqF,kBACR0D,cAAeF,YAAYG,QAC3BC,uBAAwBjJ,KAAKD,OAAOmJ,iBACpCC,uBAAwBN,YAAYK,kBAItClJ,KAAK+F,kBAAoB,OAEQ9C,IAA7BjD,KAAK4F,qBACP5F,KAAKqG,mBAAmBlB,WAAqBA,sBAACoC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDlJ,KAAKgF,UAAU8B,SACb,IAAM9G,KAAKqJ,aAAaD,cACxBA,aArW8B,gBAwWlC,EAACpJ,KAEOyI,aAAgBhH,QAGtB,GAFAzB,KAAKQ,OAAOqB,MAAM,mBAAoBJ,OAEL,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAYhB,KAAKO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,uBAAwBc,EAC3C,MATDvC,KAAKQ,OAAOiB,MAAM,yBAA0BA,MAU7C,EAGK+G,KAAAA,wBAA0B,EAAGc,gBACnCtJ,KAAKQ,OAAOqB,MAAM,8BAA+ByH,OAGjDtJ,KAAKgF,UAAU8D,OAhY4B,4BAmYvC9I,KAAK2F,iBACP3F,KAAK2F,kBAAmB,EAGV,iBAAV2D,QACFtJ,KAAK4F,yBAAsB3C,GAKjB,eAAVqG,QACFtJ,KAAKqG,mBAAmBlB,WAAAA,sBAAsBoC,WAE9CvH,KAAKuJ,yBAGPvJ,KAAK+G,aAAaxB,WAAAA,gBAAgB+D,OAAM,EACzCtJ,KAEOuJ,sBAAwB,KAC9B,IAAK,MAAMzI,WAAWd,KAAKiG,SAASY,SAE9B/F,QAAQM,aAAelB,WAAAA,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKiG,SAAS9E,OAAOL,QAAQnB,GAKhC,EAWK6G,KAAAA,qBAAuB,KAC7BxG,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM2H,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTkI,QAAS,GAAGhJ,KAAKqF,kBAAkBP,mBAAmB9E,KAAKqF,kBAAkBN,gBAC7EmE,iBAAkBlJ,KAAKD,OAAOmJ,iBAC9BO,uBAAwBzJ,KAAKD,OAAO0J,wBAItCzJ,KAAKgF,UAAU8B,SACb,KACE,MAAM4C,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,iCAAmC1B,KAAKD,OAAO4J,cAAgB,KAG1E3J,KAAKF,YAAY4J,cAEjB1J,KAAKyI,aAAa,CAChB/H,KAAMgJ,aAAajI,MACnBC,QAAS,GAAGgI,aAAahI,wBAI3B1B,KAAK2G,WAAS,EAEY,IAA5B3G,KAAKD,OAAO4J,cA1cwB,uBA8ctC3J,KAAKF,YAAY0J,cAEjBxJ,KAAKgF,UAAU8B,SACb,KACE,MAAM4C,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,sCAAwC1B,KAAKD,OAAO4J,cAAgB,KAG/E3J,KAAKF,YAAY4J,cAEjB1J,KAAKyI,aAAa,CAChB/H,KAAMgJ,aAAajI,MACnBC,QAAS,GAAGgI,aAAahI,wBAI3B1B,KAAK2G,WAAS,EAEY,IAA5B3G,KAAKD,OAAO4J,cAle6B,iCAseV1G,IAA7BjD,KAAK4F,qBACP5F,KAAKwH,gBAAgBxH,KAAK4F,oBAC3B,EACF5F,KAQO0G,sBAAwB,CAACvC,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAKyI,aAAa,CAChB/H,KAAM,UACNgB,QAASyC,SAITnE,KAAKsF,YAAcC,WAAAA,gBAAgBC,aAGrC,OAFAxF,KAAK4F,yBAAsB3C,OAC3BjD,KAAKoG,aAIPpG,KAAK2G,WACP,EAEQ0C,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBlI,KAAK6F,mBACvC,GAAI+D,qBAAuBR,aAQzB,OAPApJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BkI,oBAAsB,OAGzD5J,KAAK2G,YAGd,MAAMkD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD5J,KAAKgF,UAAU8B,SACb,IAAM9G,KAAKqJ,aAAaD,cACxBS,YAphB8B,gBAuhBlC,EAAC7J,KAEOgI,kBAAoB,KAC1BhI,KAAKgF,UAAU8B,SACb,KACE9G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKgI,mBACP,EACgC,IAAhChI,KAAKD,OAAOsI,kBAliBoB,kBAqiBpC,EA3eErI,KAAKD,OAAS,CACZsI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACf/G,SAAUoH,WAAAA,eAAeC,KACzBrD,sBAAuB,EACvBL,iBAAmBzD,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,WAAMA,OAAC1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,UAC5D5C,KAAKgF,UAAYhF,KAAKD,OAAOiF,WAAa,IAAIkF,iCAChD"}
{"version":3,"file":"index.js","sources":["../src/channel.ts","../src/connector.ts","../src/client.ts","../src/messages.ts"],"sourcesContent":["import {\n type DXLinkLogger,\n Logger,\n type DXLinkChannel,\n DXLinkChannelState,\n type DXLinkChannelMessageListener,\n type DXLinkChannelStateChangeListener,\n type DXLinkErrorListener,\n type DXLinkChannelMessage,\n type DXLinkError,\n} from '@dxfeed/dxlink-core'\n\nimport { type DXLinkWebSocketClientConfig } from './config'\nimport { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages'\n\n/**\n * A DXLink channel implementation.\n * @internal\n */\nexport class DXLinkWebSocketChannel implements DXLinkChannel {\n private status = DXLinkChannelState.REQUESTED\n\n // Listeners\n private readonly messageListeners = new Set<DXLinkChannelMessageListener>()\n private readonly statusListeners = new Set<DXLinkChannelStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n\n private logger: DXLinkLogger\n\n constructor(\n public readonly id: number,\n public readonly service: string,\n public readonly parameters: Record<string, unknown>,\n public readonly reconnect: boolean,\n private readonly sendMessage: (message: DXLinkWebSocketMessage) => void,\n config: DXLinkWebSocketClientConfig\n ) {\n this.logger = new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel)\n }\n\n send = ({ type, ...payload }: DXLinkChannelMessage) => {\n if (this.status !== DXLinkChannelState.OPENED) {\n throw new Error('Channel is not ready')\n }\n\n this.sendMessage({\n type,\n channel: this.id,\n ...payload,\n })\n }\n\n addMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.add(listener)\n removeMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.delete(listener)\n\n getState = () => this.status\n addStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.add(listener)\n removeStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n error = ({ type, message }: DXLinkError) =>\n this.send({\n type: 'ERROR',\n error: type,\n message,\n })\n\n close = () => {\n if (this.status === DXLinkChannelState.CLOSED) return\n\n this.logger.debug(`Closing by user`)\n\n // We can think that channel is closed already\n this.setStatus(DXLinkChannelState.CLOSED)\n\n this.clear()\n\n this.sendMessage({\n type: 'CHANNEL_CANCEL',\n channel: this.id,\n })\n }\n\n request = () => {\n this.logger.debug('Requesting')\n\n this.sendMessage({\n type: 'CHANNEL_REQUEST',\n channel: this.id,\n service: this.service,\n parameters: this.parameters,\n })\n\n this.processStatusRequested()\n }\n\n processPayloadMessage = (message: ChannelPayloadMessage) => {\n for (const listener of this.messageListeners) {\n listener(message)\n }\n }\n\n processStatusOpened = () => {\n this.logger.debug('Opened')\n\n this.setStatus(DXLinkChannelState.OPENED)\n }\n\n processStatusRequested = () => {\n this.setStatus(DXLinkChannelState.REQUESTED)\n }\n\n processStatusClosed = () => {\n this.logger.debug('Closed by remote endpoint')\n\n this.setStatus(DXLinkChannelState.CLOSED)\n this.clear()\n }\n\n processError = (error: DXLinkError) => {\n if (this.errorListeners.size === 0) {\n this.logger.error(`Unhandled error in channel#${this.id}: `, error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} error listener: `, e)\n }\n }\n }\n\n private setStatus = (newStatus: DXLinkChannelState) => {\n if (this.status === newStatus) return\n\n const prev = this.status\n this.status = newStatus\n for (const listener of this.statusListeners) {\n try {\n listener(newStatus, prev)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} status listener: `, e)\n }\n }\n }\n\n private clear = () => {\n this.messageListeners.clear()\n this.statusListeners.clear()\n // TODO: rethink approach when error came after channel is closed\n // this.errorListeners.clear()\n }\n}\n","import { type DXLinkWebSocketMessage } from './messages'\n\n/**\n * Interface for a WebSocket connector that manages the connection to a WebSocket server.\n * It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.\n */\nexport interface DXLinkWebSocketConnector {\n /**\n * Returns the URL of the WebSocket connection.\n */\n getUrl(): string\n /**\n * Starts the WebSocket connection.\n */\n start(): void\n /**\n * Stops the WebSocket connection and cleans up resources.\n */\n stop(): void\n /**\n * Sends a message over the WebSocket connection.\n * @param message The message to send.\n */\n sendMessage(message: DXLinkWebSocketMessage): void\n /**\n * Sets a listener that is called when the WebSocket connection is opened.\n * @param listener The listener function to call when the connection is opened.\n */\n setOpenListener(listener: () => void): void\n /**\n * Sets a listener that is called when the WebSocket connection is closed.\n * @param listener The listener function to call when the connection is closed.\n */\n setCloseListener(listener: DXLinkWebSocketCloseListener): void\n /**\n * Sets a listener that is called when a message is received from the WebSocket server.\n * @param listener The listener function to call when a message is received.\n */\n setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void\n}\n\n/**\n * Type for a close listener that is called when the WebSocket connection is closed.\n * @param reason - The reason for the closure.\n * @param error - Indicates if the closure was due to an error.\n */\nexport type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void\n\n/**\n * Default connector for the WebSocket connection.\n * @internal\n */\nexport class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: DXLinkWebSocketCloseListener | undefined = undefined\n private messageListener: ((message: DXLinkWebSocketMessage) => void) | undefined = undefined\n\n constructor(\n private readonly url: string,\n private readonly protocols?: string | string[]\n ) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url, this.protocols)\n\n this.socket.addEventListener('open', this.handleOpen)\n this.socket.addEventListener('error', this.handleError)\n this.socket.addEventListener('close', this.handleClosed)\n }\n\n stop() {\n if (this.socket === undefined) return\n\n this.socket.removeEventListener('open', this.handleOpen)\n this.socket.removeEventListener('error', this.handleError)\n this.socket.removeEventListener('close', this.handleClosed)\n this.socket.removeEventListener('message', this.handleMessage)\n\n this.socket.close()\n this.socket = undefined\n this.isAvailable = false\n }\n\n sendMessage = (message: DXLinkWebSocketMessage) => {\n if (this.socket === undefined || !this.isAvailable) {\n return\n }\n\n this.socket.send(JSON.stringify(message))\n }\n\n setOpenListener = (listener: () => void) => {\n this.openListener = listener\n }\n\n setCloseListener = (listener: DXLinkWebSocketCloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: DXLinkWebSocketMessage) => void) => {\n this.messageListener = listener\n }\n\n getUrl = () => this.url\n\n private handleOpen = () => {\n if (this.socket === undefined) return\n\n this.isAvailable = true\n\n this.socket.removeEventListener('open', this.handleOpen)\n\n this.socket.addEventListener('message', this.handleMessage)\n\n this.openListener?.()\n }\n\n private handleClosed = (ev: CloseEvent) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.(ev.reason, false)\n }\n\n private handleError = (_ev: Event) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.('Unable to connect', true)\n }\n\n private handleMessage = (ev: MessageEvent) => {\n try {\n const message = JSON.parse(ev.data)\n if (typeof message !== 'object') {\n throw new Error('Unexpected message: ' + typeof message)\n }\n if (typeof message.type !== 'string') {\n throw new Error('Unexpected message type: ' + typeof message.type)\n }\n if (typeof message.channel !== 'number') {\n throw new Error('Unexpected message channel: ' + typeof message.channel)\n }\n\n // TODO: validate message properly (e.g. check for type and other fields)\n\n if (this.messageListener === undefined) {\n return console.warn('No message listener set')\n }\n\n this.messageListener(message)\n } catch (error) {\n console.error(error instanceof Error ? error : new Error('Parsing error:' + String(error)))\n }\n }\n}\n","import {\n DXLinkLogLevel,\n type DXLinkLogger,\n Logger,\n type DXLinkScheduler,\n DefaultDXLinkScheduler,\n type DXLinkConnectionDetails,\n DXLinkConnectionState,\n type DXLinkConnectionStateChangeListener,\n type DXLinkErrorListener,\n DXLinkAuthState,\n DXLinkChannelState,\n type DXLinkAuthStateChangeListener,\n type DXLinkChannel,\n type DXLinkChannelOptions,\n type DXLinkError,\n type DXLinkClient,\n} from '@dxfeed/dxlink-core'\n\nimport { DXLinkWebSocketChannel } from './channel'\nimport type { DXLinkWebSocketClientConfig } from './config'\nimport { type DXLinkWebSocketConnector, DefaultDXLinkWebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type DXLinkWebSocketMessage,\n type SetupMessage,\n isChannelLifecycleMessage,\n isChannelMessage,\n isConnectionMessage,\n} from './messages'\nimport { VERSION } from './version'\n\n/**\n * Protocol version that is used by client.\n */\nexport const DXLINK_WS_PROTOCOL_VERSION = '0.1'\n\nconst CLIENT_VERSION = `DXF-JS/${VERSION}`\n\n// Scheduler keys\nconst DXLWS_SCHEDULER_KEY_RECONNECT = 'DXLWS_RECONNECT'\nconst DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = 'DXLWS_SETUP_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = 'DXLWS_AUTH_STATE_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_TIMEOUT = 'DXLWS_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_KEEPALIVE = 'DXLWS_KEEPALIVE'\n\nconst DEFAULT_CONNECTION_DETAILS: DXLinkConnectionDetails = {\n protocolVersion: DXLINK_WS_PROTOCOL_VERSION,\n clientVersion: CLIENT_VERSION,\n}\n\n/**\n * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.\n */\nexport class DXLinkWebSocketClient implements DXLinkClient {\n private readonly config: DXLinkWebSocketClientConfig\n\n private readonly logger: DXLinkLogger\n\n private readonly scheduler: DXLinkScheduler\n\n private connector: DXLinkWebSocketConnector | undefined\n\n private connectionState: DXLinkConnectionState = DXLinkConnectionState.NOT_CONNECTED\n private connectionDetails: DXLinkConnectionDetails = DEFAULT_CONNECTION_DETAILS\n\n private authState: DXLinkAuthState = DXLinkAuthState.UNAUTHORIZED\n\n // Listeners\n private readonly connectionStateChangeListeners = new Set<DXLinkConnectionStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n private readonly authStateChangeListeners = new Set<DXLinkAuthStateChangeListener>()\n\n /**\n * Authorization type that was determined by server behavior during setup phase.\n * This value is used to determine if authorization is required or optional or not defined yet.\n */\n private isFirstAuthState = true\n /**\n * Last setted auth token that will be sent to server after connection is established or re-established.\n */\n private lastSettedAuthToken: string | undefined\n\n // Stats for keepalive\n // TODO: mb move to connector\n private lastReceivedMillis = 0\n private lastSentMillis = 0\n\n /**\n * Count of reconnect attempts since last successful connection.\n */\n private reconnectAttempts = 0\n\n // Channels\n private globalChannelId = 1\n private readonly channels = new Map<number, DXLinkWebSocketChannel>()\n\n /**\n * Create new instance of {@link DXLinkWebSocketClient}.\n * @param config Configuration of the client.\n */\n constructor(config?: Partial<DXLinkWebSocketClientConfig>) {\n this.config = {\n keepaliveInterval: 30,\n keepaliveTimeout: 60,\n acceptKeepaliveTimeout: 60,\n actionTimeout: 10,\n logLevel: DXLinkLogLevel.WARN,\n maxReconnectAttempts: -1,\n connectorFactory: (url) => new DefaultDXLinkWebSocketConnector(url),\n ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\n this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler()\n }\n\n connect = (url: string) => {\n // Do nothing if already connected to the same url\n if (this.connector?.getUrl() === url) return\n\n // Disconnect from previous connection if any exists\n this.disconnect()\n\n this.logger.debug('Connecting to', url)\n\n // Immediately set connection state to CONNECTING\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n\n // Create new connector\n this.connector = this.config.connectorFactory(url)\n this.connector.setOpenListener(this.processTransportOpen)\n this.connector.setMessageListener(this.processMessage)\n this.connector.setCloseListener(this.processTransportClose)\n\n // Initiate websocket connection\n this.connector.start()\n }\n\n reconnect = () => {\n if (\n this.config.maxReconnectAttempts >= 0 &&\n this.reconnectAttempts >= this.config.maxReconnectAttempts\n ) {\n this.logger.warn('Max reconnect attempts reached')\n\n this.disconnect()\n return\n }\n\n if (\n this.connectionState === DXLinkConnectionState.NOT_CONNECTED ||\n this.connector === undefined\n )\n return\n\n this.logger.debug('Trying to reconnect', this.connector.getUrl())\n\n this.connector.stop()\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n\n // Increase reconnect attempts counter\n this.reconnectAttempts++\n\n // Update state for connection and channels\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n for (const channel of this.channels.values()) {\n if (channel.getState() === DXLinkChannelState.CLOSED) continue\n\n if (!channel.reconnect) {\n channel.processStatusClosed()\n continue\n }\n\n channel.processStatusRequested()\n }\n\n // Schedule reconnect attempt after some time\n // Additionally, task will be executed in case when tab is active again\n // coz browser sometimes doesn't run scheduled tasks when tab is inactive\n this.scheduler.schedule(\n () => {\n if (this.connector === undefined) return\n\n // Start new connection attempt\n this.connector.start()\n },\n this.reconnectAttempts * 1000,\n DXLWS_SCHEDULER_KEY_RECONNECT\n )\n }\n\n disconnect = () => {\n if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return\n\n this.logger.debug('Disconnecting')\n\n // Destroy connector\n this.connector?.stop()\n this.connector = undefined\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n this.reconnectAttempts = 0\n\n this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED)\n this.setAuthState(DXLinkAuthState.UNAUTHORIZED)\n }\n\n close = () => {\n this.disconnect()\n }\n\n getConnectionDetails = () => this.connectionDetails\n getConnectionState = () => this.connectionState\n getScheduler = () => this.scheduler\n addConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.add(listener)\n removeConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.delete(listener)\n\n setAuthToken = (token: string): void => {\n this.lastSettedAuthToken = token\n\n if (this.connectionState === DXLinkConnectionState.CONNECTED) {\n this.sendAuthMessage(token)\n }\n }\n\n getAuthState = (): DXLinkAuthState => this.authState\n addAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.add(listener)\n removeAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n openChannel = (\n service: string,\n parameters: Record<string, unknown>,\n options?: DXLinkChannelOptions\n ): DXLinkChannel => {\n const channelId = this.globalChannelId\n this.globalChannelId += 2\n\n const channel = new DXLinkWebSocketChannel(\n channelId,\n service,\n parameters,\n options?.reconnect ?? true,\n this.sendMessage,\n this.config\n )\n\n this.channels.set(channelId, channel)\n\n // Send channel request if connection is already established\n if (\n this.connectionState === DXLinkConnectionState.CONNECTED &&\n this.authState === DXLinkAuthState.AUTHORIZED\n ) {\n channel.request()\n }\n\n return channel\n }\n\n private setConnectionState = (newStatus: DXLinkConnectionState) => {\n const prev = this.connectionState\n if (prev === newStatus) return\n\n this.connectionState = newStatus\n for (const listener of this.connectionStateChangeListeners) {\n listener(newStatus, prev)\n }\n }\n\n private sendMessage = (message: DXLinkWebSocketMessage): void => {\n this.connector?.sendMessage(message)\n\n this.scheduleKeepalive()\n\n // TODO: mb move to connector\n this.lastSentMillis = Date.now()\n }\n\n private sendAuthMessage = (token: string): void => {\n this.logger.debug('Sending auth message')\n\n this.setAuthState(DXLinkAuthState.AUTHORIZING)\n\n this.sendMessage({\n type: 'AUTH',\n channel: 0,\n token,\n })\n }\n\n private setAuthState = (newState: DXLinkAuthState): void => {\n const prev = this.authState\n\n this.authState = newState\n for (const listener of this.authStateChangeListeners) {\n try {\n listener(newState, prev)\n } catch (e) {\n this.logger.error('Auth state listener error', e)\n }\n }\n }\n\n private processMessage = (message: DXLinkWebSocketMessage): void => {\n this.lastReceivedMillis = Date.now()\n\n // Send keepalive message if no messages sent for a while (keepaliveInterval)\n // Because browser sometimes doesn't run scheduled tasks when tab is inactive\n if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1000) {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n }\n\n // Connection messages are messages that are sent to the channel 0\n if (isConnectionMessage(message)) {\n switch (message.type) {\n case 'SETUP':\n return this.processSetupMessage(message)\n case 'AUTH_STATE':\n return this.processAuthStateMessage(message)\n case 'ERROR':\n return this.publishError({\n type: message.error,\n message: message.message,\n })\n case 'KEEPALIVE':\n // Ignore keepalive messages coz they are used only to maintain connection\n return\n }\n } else if (isChannelMessage(message)) {\n const channel = this.channels.get(message.channel)\n if (channel === undefined) {\n this.logger.warn('Received lifecycle message for unknown channel', message)\n return\n }\n\n if (isChannelLifecycleMessage(message)) {\n switch (message.type) {\n case 'CHANNEL_OPENED':\n return channel.processStatusOpened()\n case 'CHANNEL_CLOSED':\n return channel.processStatusClosed()\n case 'ERROR':\n return channel.processError({\n type: message.error,\n message: message.message,\n })\n }\n return\n }\n\n return channel.processPayloadMessage(message)\n }\n\n this.logger.warn('Unhandeled message', message.type)\n }\n\n private processSetupMessage = (serverSetup: SetupMessage): void => {\n // Clear setup timeout check from connect method\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT)\n\n // Mark connection as connected after first setup message and subsequent ones\n if (\n this.connectionState === DXLinkConnectionState.CONNECTING ||\n this.connectionState === DXLinkConnectionState.CONNECTED\n ) {\n this.connectionDetails = {\n ...this.connectionDetails,\n serverVersion: serverSetup.version,\n clientKeepaliveTimeout: this.config.keepaliveTimeout,\n serverKeepaliveTimeout: serverSetup.keepaliveTimeout,\n }\n\n // Reset reconnect attempts counter after successful connection\n this.reconnectAttempts = 0\n\n if (this.lastSettedAuthToken === undefined) {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n }\n }\n\n // Connection maintance: Setup keepalive timeout check\n const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1000\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n timeoutMills,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private publishError = (error: DXLinkError): void => {\n this.logger.debug('Publishing error', error)\n\n if (this.errorListeners.size === 0) {\n this.logger.error('Unhandled dxLink error', error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error('Error listener error', e)\n }\n }\n }\n\n private processAuthStateMessage = ({ state }: AuthStateMessage): void => {\n this.logger.debug('Received auth state message', state)\n\n // Clear auth state timeout check\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT)\n\n // Ignore first auth state message because it is sent during connection setup\n if (this.isFirstAuthState) {\n this.isFirstAuthState = false\n } else {\n // Reset auth token if server rejected it\n if (state === 'UNAUTHORIZED') {\n this.lastSettedAuthToken = undefined\n }\n }\n\n // Request active channels if connection is authorized\n if (state === 'AUTHORIZED') {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n\n this.requestActiveChannels()\n }\n\n this.setAuthState(DXLinkAuthState[state])\n }\n\n private requestActiveChannels = (): void => {\n for (const channel of this.channels.values()) {\n // clear closed channels\n if (channel.getState() === DXLinkChannelState.CLOSED) {\n this.channels.delete(channel.id)\n continue\n }\n\n channel.request()\n }\n }\n\n /**\n * Process transport open event from connector.\n * After transport is opened:\n * - setup message is sent to server\n * - auth message is sent to server if auth token is set\n * - wait for setup message from server\n * - wait for auth state message from server\n */\n private processTransportOpen = (): void => {\n this.logger.debug('Connection opened')\n\n const setupMessage: SetupMessage = {\n type: 'SETUP',\n channel: 0,\n version: `${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,\n keepaliveTimeout: this.config.keepaliveTimeout,\n acceptKeepaliveTimeout: this.config.acceptKeepaliveTimeout,\n }\n\n // Setup timeout check\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No setup message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no setup message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT\n )\n\n this.sendMessage(setupMessage)\n\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No auth state message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no auth state message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT\n )\n\n if (this.lastSettedAuthToken !== undefined) {\n this.sendAuthMessage(this.lastSettedAuthToken)\n }\n }\n\n /**\n * Process transport close event from connector.\n * After transport is closed by server:\n * - reconnect if connection is authorized\n * - disconnect if connection is not authorized\n */\n private processTransportClose = (reason: string, error: boolean): void => {\n this.logger.debug('Connection closed', reason)\n\n if (error !== undefined) {\n this.publishError({\n type: 'UNKNOWN',\n message: reason,\n })\n }\n\n if (this.authState === DXLinkAuthState.UNAUTHORIZED) {\n this.lastSettedAuthToken = undefined\n this.disconnect()\n return\n }\n\n this.reconnect()\n }\n\n private timeoutCheck = (timeoutMills: number) => {\n const now = Date.now()\n const noKeepaliveDuration = now - this.lastReceivedMillis\n if (noKeepaliveDuration >= timeoutMills) {\n this.sendMessage({\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No keepalive received for ' + noKeepaliveDuration + 'ms',\n })\n\n return this.reconnect()\n }\n\n const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration)\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n nextTimeout,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private scheduleKeepalive = () => {\n this.scheduler.schedule(\n () => {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n\n this.scheduleKeepalive()\n },\n this.config.keepaliveInterval * 1000,\n DXLWS_SCHEDULER_KEY_KEEPALIVE\n )\n }\n}\n","export interface AuthMessage {\n type: 'AUTH'\n channel: 0\n token: string\n}\n\nexport type AuthState = 'AUTHORIZED' | 'UNAUTHORIZED'\n\nexport interface AuthStateMessage {\n type: 'AUTH_STATE'\n channel: 0\n state: AuthState\n}\n\nexport interface SetupMessage {\n type: 'SETUP'\n channel: 0\n version: string\n keepaliveTimeout?: number\n acceptKeepaliveTimeout?: number\n}\n\nexport interface KeepaliveMessage {\n type: 'KEEPALIVE'\n channel: 0\n}\n\nexport interface ChannelRequestMessage {\n type: 'CHANNEL_REQUEST'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelCancelMessage {\n type: 'CHANNEL_CANCEL'\n channel: number\n}\n\nexport interface ChannelOpenedMessage {\n type: 'CHANNEL_OPENED'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelClosedMessage {\n type: 'CHANNEL_CLOSED'\n channel: number\n}\n\nexport type ErrorType =\n | 'UNKNOWN'\n | 'UNSUPPORTED_PROTOCOL'\n | 'TIMEOUT'\n | 'UNAUTHORIZED'\n | 'INVALID_MESSAGE'\n | 'BAD_ACTION'\n\nexport interface ErrorMessage {\n type: 'ERROR'\n channel: 0\n error: ErrorType\n message: string\n}\n\nexport interface ChannelErrorMessage {\n type: 'ERROR'\n channel: number\n error: ErrorType\n message: string\n}\n\nexport interface ChannelPayloadMessage {\n type: string\n channel: number\n [key: string]: unknown\n}\n\nexport type ConnectionMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n\nexport type DXLinkWebSocketMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelPayloadMessage\n | ChannelErrorMessage\n\nexport const isConnectionMessage = (\n message: DXLinkWebSocketMessage\n): message is ConnectionMessage =>\n message.channel === 0 &&\n (message.type === 'SETUP' ||\n message.type === 'KEEPALIVE' ||\n message.type === 'AUTH' ||\n message.type === 'AUTH_STATE' ||\n message.type === 'ERROR')\n\nexport type ChannelLifecycleMessage =\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n\nexport type ChannelMessage = ChannelLifecycleMessage | ChannelPayloadMessage\n\nexport const isChannelMessage = (message: DXLinkWebSocketMessage): message is ChannelMessage =>\n message.channel !== 0\n\nexport const isChannelLifecycleMessage = (\n message: ChannelMessage\n): message is ChannelLifecycleMessage =>\n message.type === 'CHANNEL_OPENED' ||\n message.type === 'CHANNEL_CLOSED' ||\n message.type === 'ERROR' ||\n message.type === 'CHANNEL_REQUEST' ||\n message.type === 'CHANNEL_CANCEL'\n"],"names":["DXLinkWebSocketChannel","constructor","id","service","parameters","reconnect","sendMessage","config","this","status","DXLinkChannelState","REQUESTED","messageListeners","Set","statusListeners","errorListeners","logger","send","type","payload","OPENED","Error","channel","addMessageListener","listener","add","removeMessageListener","delete","getState","addStateChangeListener","removeStateChangeListener","addErrorListener","removeErrorListener","error","message","close","CLOSED","debug","setStatus","clear","request","processStatusRequested","processPayloadMessage","processStatusOpened","processStatusClosed","processError","size","e","newStatus","prev","Logger","name","logLevel","DefaultDXLinkWebSocketConnector","url","protocols","socket","undefined","isAvailable","openListener","closeListener","messageListener","JSON","stringify","setOpenListener","setCloseListener","setMessageListener","getUrl","handleOpen","removeEventListener","addEventListener","handleMessage","handleClosed","ev","stop","reason","handleError","_ev","parse","data","console","warn","String","start","WebSocket","DEFAULT_CONNECTION_DETAILS","protocolVersion","clientVersion","scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","connectorFactory","processTransportOpen","processMessage","processTransportClose","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","getScheduler","addConnectionStateChangeListener","removeConnectionStateChangeListener","setAuthToken","token","CONNECTED","sendAuthMessage","getAuthState","addAuthStateChangeListener","removeAuthStateChangeListener","openChannel","options","channelId","set","AUTHORIZED","scheduleKeepalive","Date","now","AUTHORIZING","newState","keepaliveInterval","isConnectionMessage","processSetupMessage","processAuthStateMessage","publishError","isChannelMessage","get","isChannelLifecycleMessage","serverSetup","cancel","serverVersion","version","clientKeepaliveTimeout","keepaliveTimeout","serverKeepaliveTimeout","timeoutMills","timeoutCheck","state","requestActiveChannels","setupMessage","acceptKeepaliveTimeout","errorMessage","actionTimeout","noKeepaliveDuration","nextTimeout","Math","max","DXLinkLogLevel","WARN","DefaultDXLinkScheduler"],"mappings":"oDAmBaA,uBAUXC,WAAAA,CACkBC,GACAC,QACAC,WACAC,UACCC,YACjBC,QAAmCC,KALnBN,QAAA,EAAAM,KACAL,aACAC,EAAAA,KAAAA,gBACAC,EAAAA,KAAAA,sBACCC,iBAAA,EAAAE,KAdXC,OAASC,WAAAA,mBAAmBC,UAGnBC,KAAAA,iBAAmB,IAAIC,IAAmCL,KAC1DM,gBAAkB,IAAID,SACtBE,eAAiB,IAAIF,IAE9BG,KAAAA,YAaRC,EAAAA,KAAAA,KAAO,EAAGC,aAASC,YACjB,GAAIX,KAAKC,SAAWC,WAAAA,mBAAmBU,OACrC,MAAU,IAAAC,MAAM,wBAGlBb,KAAKF,YAAY,CACfY,UACAI,QAASd,KAAKN,MACXiB,WAIPI,KAAAA,mBAAsBC,UACpBhB,KAAKI,iBAAiBa,IAAID,UAAShB,KACrCkB,sBAAyBF,UACvBhB,KAAKI,iBAAiBe,OAAOH,UAAShB,KAExCoB,SAAW,IAAMpB,KAAKC,OAAMD,KAC5BqB,uBAA0BL,UACxBhB,KAAKM,gBAAgBW,IAAID,UAC3BM,KAAAA,0BAA6BN,UAC3BhB,KAAKM,gBAAgBa,OAAOH,UAE9BO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7FyB,MAAQ,EAAGf,UAAMgB,mBACf1B,KAAKS,KAAK,CACRC,KAAM,QACNe,MAAOf,KACPgB,kBAGJC,KAAAA,MAAQ,KACF3B,KAAKC,SAAWC,WAAAA,mBAAmB0B,SAEvC5B,KAAKQ,OAAOqB,MAAM,mBAGlB7B,KAAK8B,UAAU5B,8BAAmB0B,QAElC5B,KAAK+B,QAEL/B,KAAKF,YAAY,CACfY,KAAM,iBACNI,QAASd,KAAKN,OAEjBM,KAEDgC,QAAU,KACRhC,KAAKQ,OAAOqB,MAAM,cAElB7B,KAAKF,YAAY,CACfY,KAAM,kBACNI,QAASd,KAAKN,GACdC,QAASK,KAAKL,QACdC,WAAYI,KAAKJ,aAGnBI,KAAKiC,wBAAsB,EAC5BjC,KAEDkC,sBAAyBR,UACvB,IAAK,MAAMV,iBAAiBZ,iBAC1BY,SAASU,QACV,EAGHS,KAAAA,oBAAsB,KACpBnC,KAAKQ,OAAOqB,MAAM,UAElB7B,KAAK8B,UAAU5B,8BAAmBU,SACnCZ,KAEDiC,uBAAyB,KACvBjC,KAAK8B,UAAU5B,WAAAA,mBAAmBC,UAAS,EAG7CiC,KAAAA,oBAAsB,KACpBpC,KAAKQ,OAAOqB,MAAM,6BAElB7B,KAAK8B,UAAU5B,WAAAA,mBAAmB0B,QAClC5B,KAAK+B,OAAK,EACX/B,KAEDqC,aAAgBZ,QACd,GAAiC,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,iBAAiBT,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKN,sBAAuB6C,EACnE,MATDvC,KAAKQ,OAAOiB,MAAM,8BAA8BzB,KAAKN,OAAQ+B,MAU9D,EACFzB,KAEO8B,UAAaU,YACnB,GAAIxC,KAAKC,SAAWuC,UAAW,OAE/B,MAAMC,KAAOzC,KAAKC,OAClBD,KAAKC,OAASuC,UACd,IAAK,MAAMxB,YAAYhB,KAAKM,gBAC1B,IACEU,SAASwB,UAAWC,KACrB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKN,uBAAwB6C,EACpE,CACF,EACFvC,KAEO+B,MAAQ,KACd/B,KAAKI,iBAAiB2B,QACtB/B,KAAKM,gBAAgByB,OAAK,EA9HV/B,KAAEN,GAAFA,GACAM,KAAOL,QAAPA,QACAK,KAAUJ,WAAVA,WACAI,KAASH,UAATA,UACCG,KAAWF,YAAXA,YAGjBE,KAAKQ,OAAS,IAAIkC,WAAMA,OAAC,GAAGlD,uBAAuBmD,QAAQjD,MAAMC,UAAWI,OAAO6C,SACrF,QCcWC,gCASXpD,WAAAA,CACmBqD,IACAC,WAA6B/C,KAD7B8C,SAAA,EAAA9C,KACA+C,eAAA,EAAA/C,KAVXgD,YAAgCC,EAASjD,KAEzCkD,aAAc,EAAKlD,KAEnBmD,kBAAyCF,EACzCG,KAAAA,mBAA0DH,EAC1DI,KAAAA,qBAA2EJ,EA8BnFnD,KAAAA,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAAC1B,KAEDwD,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QAAAA,EACrBhB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAEA0C,KAAAA,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QAAAA,EAGzB2C,KAAAA,OAAS,IAAM3D,KAAK8C,IAEZc,KAAAA,WAAa,UACCX,IAAhBjD,KAAKgD,SAEThD,KAAKkD,aAAc,EAEnBlD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAE7C5D,KAAKgD,OAAOc,iBAAiB,UAAW9D,KAAK+D,eAE7C/D,KAAKmD,iBACP,EAACnD,KAEOgE,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAACpD,KAEO+D,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,MAAM,IAAIb,MAAM,8BAAgCa,SAElD,GAA4B,iBAAjBA,QAAQhB,KACjB,MAAU,IAAAG,MAAM,mCAAqCa,QAAQhB,MAE/D,GAA+B,iBAApBgB,QAAQZ,QACjB,MAAM,IAAID,MAAM,sCAAwCa,QAAQZ,SAKlE,QAA6BmC,IAAzBjD,KAAKqD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBzE,KAAKqD,gBAAgB3B,QACtB,CAAC,MAAOD,OACP+C,QAAQ/C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB6D,OAAOjD,QACpF,GAnGgBzB,KAAG8C,IAAHA,IACA9C,KAAS+C,UAATA,SAChB,CAEH4B,KAAAA,QACsB1B,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAS,IAAI4B,UAAU5E,KAAK8C,IAAK9C,KAAK+C,WAE3C/C,KAAKgD,OAAOc,iBAAiB,OAAQ9D,KAAK4D,YAC1C5D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKoE,aAC3CpE,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAC7C5D,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKoE,aAC9CpE,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKgE,cAC9ChE,KAAKgD,OAAOa,oBAAoB,UAAW7D,KAAK+D,eAEhD/D,KAAKgD,OAAOrB,QACZ3B,KAAKgD,YAASC,EACdjD,KAAKkD,aAAc,EACrB,ECnDW,MAWP2B,2BAAsD,CAC1DC,gBAZwC,MAaxCC,cAXqB,mDAFmB,0CAkExCtF,WAAAA,CAAYM,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENgF,eAAS,EAAAhF,KAElBiF,eAAS,EAAAjF,KAETkF,gBAAyCC,WAAAA,sBAAsBC,cAAapF,KAC5EqF,kBAA6CR,2BAA0B7E,KAEvEsF,UAA6BC,WAAeA,gBAACC,aAGpCC,KAAAA,+BAAiC,IAAIpF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBqF,KAAAA,yBAA2B,IAAIrF,IAMxCsF,KAAAA,kBAAmB,EAInBC,KAAAA,yBAIAC,EAAAA,KAAAA,mBAAqB,EACrBC,KAAAA,eAAiB,EAAC9F,KAKlB+F,kBAAoB,EAAC/F,KAGrBgG,gBAAkB,EAAChG,KACViG,SAAW,IAAIC,IAAqClG,KAsBrEmG,QAAWrD,MAEL9C,KAAKiF,WAAWtB,WAAab,MAGjC9C,KAAKoG,aAELpG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKqG,mBAAmBlB,WAAqBA,sBAACmB,YAG9CtG,KAAKiF,UAAYjF,KAAKD,OAAOwG,iBAAiBzD,KAC9C9C,KAAKiF,UAAUzB,gBAAgBxD,KAAKwG,sBACpCxG,KAAKiF,UAAUvB,mBAAmB1D,KAAKyG,gBACvCzG,KAAKiF,UAAUxB,iBAAiBzD,KAAK0G,uBAGrC1G,KAAKiF,UAAUN,QAAK,EAGtB9E,KAAAA,UAAY,KACV,GACEG,KAAKD,OAAO4G,sBAAwB,GACpC3G,KAAK+F,mBAAqB/F,KAAKD,OAAO4G,qBAKtC,OAHA3G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKoG,aAIP,GACEpG,KAAKkF,kBAAoBC,WAAqBA,sBAACC,oBAC5BnC,IAAnBjD,KAAKiF,UAFP,CAMAjF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKiF,UAAUtB,UAExD3D,KAAKiF,UAAUf,OAGflE,KAAKgF,UAAUjD,QAGf/B,KAAKqF,kBAAoBR,2BACzB7E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EAGxB3F,KAAK+F,oBAGL/F,KAAKqG,mBAAmBlB,iCAAsBmB,YAC9C,IAAK,MAAMxF,WAAed,KAACiG,SAASW,SAC9B9F,QAAQM,aAAelB,WAAkBA,mBAAC0B,SAEzCd,QAAQjB,UAKbiB,QAAQmB,yBAJNnB,QAAQsB,uBAUZpC,KAAKgF,UAAU6B,SACb,UACyB5D,IAAnBjD,KAAKiF,WAGTjF,KAAKiF,UAAUN,OAAK,EAEG,IAAzB3E,KAAK+F,kBA3J2B,kBAkHhC,CA0C6B,EAIjCK,KAAAA,WAAa,KACPpG,KAAKkF,kBAAoBC,WAAAA,sBAAsBC,gBAEnDpF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKiF,WAAWf,OAChBlE,KAAKiF,eAAYhC,EAGjBjD,KAAKgF,UAAUjD,QAGf/B,KAAKqF,kBAAoBR,2BACzB7E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EACxB3F,KAAK+F,kBAAoB,EAEzB/F,KAAKqG,mBAAmBlB,WAAqBA,sBAACC,eAC9CpF,KAAK8G,aAAavB,WAAeA,gBAACC,cAAY,EAC/CxF,KAED2B,MAAQ,KACN3B,KAAKoG,YAAU,OAGjBW,qBAAuB,IAAM/G,KAAKqF,kBAClC2B,KAAAA,mBAAqB,IAAMhH,KAAKkF,gBAAelF,KAC/CiH,aAAe,IAAMjH,KAAKgF,UAC1BkC,KAAAA,iCAAoClG,UAClChB,KAAKyF,+BAA+BxE,IAAID,UAC1CmG,KAAAA,oCAAuCnG,UACrChB,KAAKyF,+BAA+BtE,OAAOH,UAAShB,KAEtDoH,aAAgBC,QACdrH,KAAK4F,oBAAsByB,MAEvBrH,KAAKkF,kBAAoBC,WAAAA,sBAAsBmC,WACjDtH,KAAKuH,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuBxH,KAAKsF,UAAStF,KACpDyH,2BAA8BzG,UAC5BhB,KAAK0F,yBAAyBzE,IAAID,UACpC0G,KAAAA,8BAAiC1G,UAC/BhB,KAAK0F,yBAAyBvE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7F2H,YAAc,CACZhI,QACAC,WACAgI,WAEA,MAAMC,UAAY7H,KAAKgG,gBACvBhG,KAAKgG,iBAAmB,EAExB,MAAMlF,QAAU,IAAItB,uBAClBqI,UACAlI,QACAC,WACAgI,SAAS/H,YAAa,EACtBG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKiG,SAAS6B,IAAID,UAAW/G,SAI3Bd,KAAKkF,kBAAoBC,WAAqBA,sBAACmC,WAC/CtH,KAAKsF,YAAcC,WAAeA,gBAACwC,YAEnCjH,QAAQkB,UAGHlB,SAGDuF,KAAAA,mBAAsB7D,YAC5B,MAAMC,KAAOzC,KAAKkF,gBAClB,GAAIzC,OAASD,UAAb,CAEAxC,KAAKkF,gBAAkB1C,UACvB,IAAK,MAAMxB,YAAgBhB,KAACyF,+BAC1BzE,SAASwB,UAAWC,KAFtB,CAGC,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKiF,WAAWnF,YAAY4B,SAE5B1B,KAAKgI,oBAGLhI,KAAK8F,eAAiBmC,KAAKC,KAAG,EAGxBX,KAAAA,gBAAmBF,QACzBrH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAK8G,aAAavB,2BAAgB4C,aAElCnI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTuG,aACD,EAGKP,KAAAA,aAAgBsB,WACtB,MAAM3F,KAAOzC,KAAKsF,UAElBtF,KAAKsF,UAAY8C,SACjB,IAAK,MAAMpH,YAAgBhB,KAAC0F,yBAC1B,IACE1E,SAASoH,SAAU3F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EAGKkE,KAAAA,eAAkB/E,UAaxB,GAZA1B,KAAK6F,mBAAqBoC,KAAKC,MAI3BlI,KAAK6F,mBAAqB7F,KAAK8F,gBAAkD,IAAhC9F,KAAKD,OAAOsI,mBAC/DrI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IC3OfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDyOJ4H,CAAoB5G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACuI,oBAAoB7G,SAClC,IAAK,aACH,OAAO1B,KAAKwI,wBAAwB9G,SACtC,IAAK,QACH,OAAW1B,KAACyI,aAAa,CACvB/H,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GC7OsBA,UACX,IAApBA,QAAQZ,QD4OK4H,CAAiBhH,SAAU,CACpC,MAAMZ,QAAUd,KAAKiG,SAAS0C,IAAIjH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GChPJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KD0OAkI,CAA0BlH,SAAU,CACtC,OAAQA,QAAQhB,MACd,IAAK,iBACH,OAAOI,QAAQqB,sBACjB,IAAK,iBACH,OAAOrB,QAAQsB,sBACjB,IAAK,QACH,OAAOtB,QAAQuB,aAAa,CAC1B3B,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAGvB,MACD,CAED,OAAOZ,QAAQoB,sBAAsBR,QACtC,CAED1B,KAAKQ,OAAOiE,KAAK,qBAAsB/C,QAAQhB,KAAI,EAG7C6H,KAAAA,oBAAuBM,cAE7B7I,KAAKgF,UAAU8D,OAvVuB,uBA2VpC9I,KAAKkF,kBAAoBC,WAAAA,sBAAsBmB,YAC/CtG,KAAKkF,kBAAoBC,WAAqBA,sBAACmC,YAE/CtH,KAAKqF,kBAAoB,IACpBrF,KAAKqF,kBACR0D,cAAeF,YAAYG,QAC3BC,uBAAwBjJ,KAAKD,OAAOmJ,iBACpCC,uBAAwBN,YAAYK,kBAItClJ,KAAK+F,kBAAoB,OAEQ9C,IAA7BjD,KAAK4F,qBACP5F,KAAKqG,mBAAmBlB,WAAqBA,sBAACmC,YAKlD,MAAM8B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDlJ,KAAKgF,UAAU6B,SACb,IAAM7G,KAAKqJ,aAAaD,cACxBA,aA/W8B,gBAkXlC,EAACpJ,KAEOyI,aAAgBhH,QAGtB,GAFAzB,KAAKQ,OAAOqB,MAAM,mBAAoBJ,OAEL,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAgBhB,KAACO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,uBAAwBc,EAC3C,MATDvC,KAAKQ,OAAOiB,MAAM,yBAA0BA,MAU7C,EAGK+G,KAAAA,wBAA0B,EAAGc,gBACnCtJ,KAAKQ,OAAOqB,MAAM,8BAA+ByH,OAGjDtJ,KAAKgF,UAAU8D,OA1Y4B,4BA6YvC9I,KAAK2F,iBACP3F,KAAK2F,kBAAmB,EAGV,iBAAV2D,QACFtJ,KAAK4F,yBAAsB3C,GAKjB,eAAVqG,QACFtJ,KAAKqG,mBAAmBlB,WAAqBA,sBAACmC,WAE9CtH,KAAKuJ,yBAGPvJ,KAAK8G,aAAavB,WAAeA,gBAAC+D,OACpC,EAACtJ,KAEOuJ,sBAAwB,KAC9B,IAAK,MAAMzI,WAAed,KAACiG,SAASW,SAE9B9F,QAAQM,aAAelB,WAAkBA,mBAAC0B,OAK9Cd,QAAQkB,UAJNhC,KAAKiG,SAAS9E,OAAOL,QAAQpB,GAKhC,EAWK8G,KAAAA,qBAAuB,KAC7BxG,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM2H,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTkI,QAAS,GAAGhJ,KAAKqF,kBAAkBP,mBAAmB9E,KAAKqF,kBAAkBN,gBAC7EmE,iBAAkBlJ,KAAKD,OAAOmJ,iBAC9BO,uBAAwBzJ,KAAKD,OAAO0J,wBAItCzJ,KAAKgF,UAAU6B,SACb,KACE,MAAM6C,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,iCAAmC1B,KAAKD,OAAO4J,cAAgB,KAG1E3J,KAAKF,YAAY4J,cAEjB1J,KAAKyI,aAAa,CAChB/H,KAAMgJ,aAAajI,MACnBC,QAAS,GAAGgI,aAAahI,wBAI3B1B,KAAKH,WACP,EAC4B,IAA5BG,KAAKD,OAAO4J,cApdwB,uBAwdtC3J,KAAKF,YAAY0J,cAEjBxJ,KAAKgF,UAAU6B,SACb,KACE,MAAM6C,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,sCAAwC1B,KAAKD,OAAO4J,cAAgB,KAG/E3J,KAAKF,YAAY4J,cAEjB1J,KAAKyI,aAAa,CAChB/H,KAAMgJ,aAAajI,MACnBC,QAAS,GAAGgI,aAAahI,wBAI3B1B,KAAKH,WAAS,EAEY,IAA5BG,KAAKD,OAAO4J,cA5e6B,iCAgfV1G,IAA7BjD,KAAK4F,qBACP5F,KAAKuH,gBAAgBvH,KAAK4F,oBAC3B,OASKc,sBAAwB,CAACvC,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAKyI,aAAa,CAChB/H,KAAM,UACNgB,QAASyC,SAITnE,KAAKsF,YAAcC,WAAeA,gBAACC,aAGrC,OAFAxF,KAAK4F,yBAAsB3C,OAC3BjD,KAAKoG,aAIPpG,KAAKH,WACP,EAEQwJ,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBlI,KAAK6F,mBACvC,GAAI+D,qBAAuBR,aAQzB,OAPApJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BkI,oBAAsB,OAGzD5J,KAAKH,YAGd,MAAMgK,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD5J,KAAKgF,UAAU6B,SACb,IAAM7G,KAAKqJ,aAAaD,cACxBS,YA9hB8B,gBA+hBH,EAE9B7J,KAEOgI,kBAAoB,KAC1BhI,KAAKgF,UAAU6B,SACb,KACE7G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKgI,mBAAiB,EAEQ,IAAhChI,KAAKD,OAAOsI,kBA5iBoB,kBA6iBH,EAnf/BrI,KAAKD,OAAS,CACZsI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACf/G,SAAUoH,WAAcA,eAACC,KACzBtD,sBAAuB,EACvBJ,iBAAmBzD,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,WAAAA,OAAO1C,KAAKP,YAAYkD,KAAM3C,KAAKD,OAAO6C,UAC5D5C,KAAKgF,UAAYhF,KAAKD,OAAOiF,WAAa,IAAIkF,iCAChD"}

@@ -1,2 +0,2 @@

import{DXLinkChannelState,Logger,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel,DefaultDXLinkScheduler}from"@dxfeed/dxlink-core";class DXLinkWebSocketChannel{constructor(id,service,parameters,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.sendMessage=void 0,this.status=DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.sendMessage=sendMessage,this.logger=new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DXLINK_WS_PROTOCOL_VERSION="0.1",DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.0"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED&&channel.processStatusRequested();this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===DXLinkConnectionState.CONNECTED&&this.authState===DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==DXLinkConnectionState.CONNECTING&&this.connectionState!==DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new DefaultDXLinkScheduler}}export{DXLINK_WS_PROTOCOL_VERSION,DXLinkWebSocketClient,DefaultDXLinkWebSocketConnector};
import{DXLinkChannelState,Logger,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel,DefaultDXLinkScheduler}from"@dxfeed/dxlink-core";class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DXLINK_WS_PROTOCOL_VERSION="0.1",DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.1"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===DXLinkConnectionState.CONNECTED&&this.authState===DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==DXLinkConnectionState.CONNECTING&&this.connectionState!==DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new DefaultDXLinkScheduler}}export{DXLINK_WS_PROTOCOL_VERSION,DXLinkWebSocketClient,DefaultDXLinkWebSocketConnector};
//# sourceMappingURL=index.module.js.map

@@ -1,1 +0,1 @@

{"version":3,"file":"index.module.js","sources":["../src/channel.ts","../src/connector.ts","../src/client.ts","../src/messages.ts"],"sourcesContent":["import {\n type DXLinkLogger,\n Logger,\n type DXLinkChannel,\n DXLinkChannelState,\n type DXLinkChannelMessageListener,\n type DXLinkChannelStateChangeListener,\n type DXLinkErrorListener,\n type DXLinkChannelMessage,\n type DXLinkError,\n} from '@dxfeed/dxlink-core'\n\nimport { type DXLinkWebSocketClientConfig } from './config'\nimport { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages'\n\n/**\n * A DXLink channel implementation.\n * @internal\n */\nexport class DXLinkWebSocketChannel implements DXLinkChannel {\n private status = DXLinkChannelState.REQUESTED\n\n // Listeners\n private readonly messageListeners = new Set<DXLinkChannelMessageListener>()\n private readonly statusListeners = new Set<DXLinkChannelStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n\n private logger: DXLinkLogger\n\n constructor(\n public readonly id: number,\n public readonly service: string,\n public readonly parameters: Record<string, unknown>,\n private readonly sendMessage: (message: DXLinkWebSocketMessage) => void,\n config: DXLinkWebSocketClientConfig\n ) {\n this.logger = new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel)\n }\n\n send = ({ type, ...payload }: DXLinkChannelMessage) => {\n if (this.status !== DXLinkChannelState.OPENED) {\n throw new Error('Channel is not ready')\n }\n\n this.sendMessage({\n type,\n channel: this.id,\n ...payload,\n })\n }\n\n addMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.add(listener)\n removeMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.delete(listener)\n\n getState = () => this.status\n addStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.add(listener)\n removeStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n error = ({ type, message }: DXLinkError) =>\n this.send({\n type: 'ERROR',\n error: type,\n message,\n })\n\n close = () => {\n if (this.status === DXLinkChannelState.CLOSED) return\n\n this.logger.debug(`Closing by user`)\n\n // We can think that channel is closed already\n this.setStatus(DXLinkChannelState.CLOSED)\n\n this.clear()\n\n this.sendMessage({\n type: 'CHANNEL_CANCEL',\n channel: this.id,\n })\n }\n\n request = () => {\n this.logger.debug('Requesting')\n\n this.sendMessage({\n type: 'CHANNEL_REQUEST',\n channel: this.id,\n service: this.service,\n parameters: this.parameters,\n })\n\n this.processStatusRequested()\n }\n\n processPayloadMessage = (message: ChannelPayloadMessage) => {\n for (const listener of this.messageListeners) {\n listener(message)\n }\n }\n\n processStatusOpened = () => {\n this.logger.debug('Opened')\n\n this.setStatus(DXLinkChannelState.OPENED)\n }\n\n processStatusRequested = () => {\n this.setStatus(DXLinkChannelState.REQUESTED)\n }\n\n processStatusClosed = () => {\n this.logger.debug('Closed by remote endpoint')\n\n this.setStatus(DXLinkChannelState.CLOSED)\n this.clear()\n }\n\n processError = (error: DXLinkError) => {\n if (this.errorListeners.size === 0) {\n this.logger.error(`Unhandled error in channel#${this.id}: `, error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} error listener: `, e)\n }\n }\n }\n\n private setStatus = (newStatus: DXLinkChannelState) => {\n if (this.status === newStatus) return\n\n const prev = this.status\n this.status = newStatus\n for (const listener of this.statusListeners) {\n try {\n listener(newStatus, prev)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} status listener: `, e)\n }\n }\n }\n\n private clear = () => {\n this.messageListeners.clear()\n this.statusListeners.clear()\n // TODO: rethink approach when error came after channel is closed\n // this.errorListeners.clear()\n }\n}\n","import { type DXLinkWebSocketMessage } from './messages'\n\n/**\n * Interface for a WebSocket connector that manages the connection to a WebSocket server.\n * It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.\n */\nexport interface DXLinkWebSocketConnector {\n /**\n * Returns the URL of the WebSocket connection.\n */\n getUrl(): string\n /**\n * Starts the WebSocket connection.\n */\n start(): void\n /**\n * Stops the WebSocket connection and cleans up resources.\n */\n stop(): void\n /**\n * Sends a message over the WebSocket connection.\n * @param message The message to send.\n */\n sendMessage(message: DXLinkWebSocketMessage): void\n /**\n * Sets a listener that is called when the WebSocket connection is opened.\n * @param listener The listener function to call when the connection is opened.\n */\n setOpenListener(listener: () => void): void\n /**\n * Sets a listener that is called when the WebSocket connection is closed.\n * @param listener The listener function to call when the connection is closed.\n */\n setCloseListener(listener: DXLinkWebSocketCloseListener): void\n /**\n * Sets a listener that is called when a message is received from the WebSocket server.\n * @param listener The listener function to call when a message is received.\n */\n setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void\n}\n\n/**\n * Type for a close listener that is called when the WebSocket connection is closed.\n * @param reason - The reason for the closure.\n * @param error - Indicates if the closure was due to an error.\n */\nexport type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void\n\n/**\n * Default connector for the WebSocket connection.\n * @internal\n */\nexport class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: DXLinkWebSocketCloseListener | undefined = undefined\n private messageListener: ((message: DXLinkWebSocketMessage) => void) | undefined = undefined\n\n constructor(\n private readonly url: string,\n private readonly protocols?: string | string[]\n ) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url, this.protocols)\n\n this.socket.addEventListener('open', this.handleOpen)\n this.socket.addEventListener('error', this.handleError)\n this.socket.addEventListener('close', this.handleClosed)\n }\n\n stop() {\n if (this.socket === undefined) return\n\n this.socket.removeEventListener('open', this.handleOpen)\n this.socket.removeEventListener('error', this.handleError)\n this.socket.removeEventListener('close', this.handleClosed)\n this.socket.removeEventListener('message', this.handleMessage)\n\n this.socket.close()\n this.socket = undefined\n this.isAvailable = false\n }\n\n sendMessage = (message: DXLinkWebSocketMessage) => {\n if (this.socket === undefined || !this.isAvailable) {\n return\n }\n\n this.socket.send(JSON.stringify(message))\n }\n\n setOpenListener = (listener: () => void) => {\n this.openListener = listener\n }\n\n setCloseListener = (listener: DXLinkWebSocketCloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: DXLinkWebSocketMessage) => void) => {\n this.messageListener = listener\n }\n\n getUrl = () => this.url\n\n private handleOpen = () => {\n if (this.socket === undefined) return\n\n this.isAvailable = true\n\n this.socket.removeEventListener('open', this.handleOpen)\n\n this.socket.addEventListener('message', this.handleMessage)\n\n this.openListener?.()\n }\n\n private handleClosed = (ev: CloseEvent) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.(ev.reason, false)\n }\n\n private handleError = (_ev: Event) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.('Unable to connect', true)\n }\n\n private handleMessage = (ev: MessageEvent) => {\n try {\n const message = JSON.parse(ev.data)\n if (typeof message !== 'object') {\n throw new Error('Unexpected message: ' + typeof message)\n }\n if (typeof message.type !== 'string') {\n throw new Error('Unexpected message type: ' + typeof message.type)\n }\n if (typeof message.channel !== 'number') {\n throw new Error('Unexpected message channel: ' + typeof message.channel)\n }\n\n // TODO: validate message properly (e.g. check for type and other fields)\n\n if (this.messageListener === undefined) {\n return console.warn('No message listener set')\n }\n\n this.messageListener(message)\n } catch (error) {\n console.error(error instanceof Error ? error : new Error('Parsing error:' + String(error)))\n }\n }\n}\n","import {\n DXLinkLogLevel,\n type DXLinkLogger,\n Logger,\n type DXLinkScheduler,\n DefaultDXLinkScheduler,\n type DXLinkConnectionDetails,\n DXLinkConnectionState,\n type DXLinkConnectionStateChangeListener,\n type DXLinkErrorListener,\n DXLinkAuthState,\n DXLinkChannelState,\n type DXLinkAuthStateChangeListener,\n type DXLinkChannel,\n type DXLinkError,\n type DXLinkClient,\n} from '@dxfeed/dxlink-core'\n\nimport { DXLinkWebSocketChannel } from './channel'\nimport type { DXLinkWebSocketClientConfig } from './config'\nimport { type DXLinkWebSocketConnector, DefaultDXLinkWebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type DXLinkWebSocketMessage,\n type SetupMessage,\n isChannelLifecycleMessage,\n isChannelMessage,\n isConnectionMessage,\n} from './messages'\nimport { VERSION } from './version'\n\n/**\n * Protocol version that is used by client.\n */\nexport const DXLINK_WS_PROTOCOL_VERSION = '0.1'\n\nconst CLIENT_VERSION = `DXF-JS/${VERSION}`\n\n// Scheduler keys\nconst DXLWS_SCHEDULER_KEY_RECONNECT = 'DXLWS_RECONNECT'\nconst DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = 'DXLWS_SETUP_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = 'DXLWS_AUTH_STATE_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_TIMEOUT = 'DXLWS_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_KEEPALIVE = 'DXLWS_KEEPALIVE'\n\nconst DEFAULT_CONNECTION_DETAILS: DXLinkConnectionDetails = {\n protocolVersion: DXLINK_WS_PROTOCOL_VERSION,\n clientVersion: CLIENT_VERSION,\n}\n\n/**\n * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.\n */\nexport class DXLinkWebSocketClient implements DXLinkClient {\n private readonly config: DXLinkWebSocketClientConfig\n\n private readonly logger: DXLinkLogger\n\n private readonly scheduler: DXLinkScheduler\n\n private connector: DXLinkWebSocketConnector | undefined\n\n private connectionState: DXLinkConnectionState = DXLinkConnectionState.NOT_CONNECTED\n private connectionDetails: DXLinkConnectionDetails = DEFAULT_CONNECTION_DETAILS\n\n private authState: DXLinkAuthState = DXLinkAuthState.UNAUTHORIZED\n\n // Listeners\n private readonly connectionStateChangeListeners = new Set<DXLinkConnectionStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n private readonly authStateChangeListeners = new Set<DXLinkAuthStateChangeListener>()\n\n /**\n * Authorization type that was determined by server behavior during setup phase.\n * This value is used to determine if authorization is required or optional or not defined yet.\n */\n private isFirstAuthState = true\n /**\n * Last setted auth token that will be sent to server after connection is established or re-established.\n */\n private lastSettedAuthToken: string | undefined\n\n // Stats for keepalive\n // TODO: mb move to connector\n private lastReceivedMillis = 0\n private lastSentMillis = 0\n\n /**\n * Count of reconnect attempts since last successful connection.\n */\n private reconnectAttempts = 0\n\n // Channels\n private globalChannelId = 1\n private readonly channels = new Map<number, DXLinkWebSocketChannel>()\n\n /**\n * Create new instance of {@link DXLinkWebSocketClient}.\n * @param config Configuration of the client.\n */\n constructor(config?: Partial<DXLinkWebSocketClientConfig>) {\n this.config = {\n keepaliveInterval: 30,\n keepaliveTimeout: 60,\n acceptKeepaliveTimeout: 60,\n actionTimeout: 10,\n logLevel: DXLinkLogLevel.WARN,\n maxReconnectAttempts: -1,\n connectorFactory: (url) => new DefaultDXLinkWebSocketConnector(url),\n ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\n this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler()\n }\n\n connect = (url: string) => {\n // Do nothing if already connected to the same url\n if (this.connector?.getUrl() === url) return\n\n // Disconnect from previous connection if any exists\n this.disconnect()\n\n this.logger.debug('Connecting to', url)\n\n // Immediately set connection state to CONNECTING\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n\n // Create new connector\n this.connector = this.config.connectorFactory(url)\n this.connector.setOpenListener(this.processTransportOpen)\n this.connector.setMessageListener(this.processMessage)\n this.connector.setCloseListener(this.processTransportClose)\n\n // Initiate websocket connection\n this.connector.start()\n }\n\n reconnect = () => {\n if (\n this.config.maxReconnectAttempts >= 0 &&\n this.reconnectAttempts >= this.config.maxReconnectAttempts\n ) {\n this.logger.warn('Max reconnect attempts reached')\n\n this.disconnect()\n return\n }\n\n if (\n this.connectionState === DXLinkConnectionState.NOT_CONNECTED ||\n this.connector === undefined\n )\n return\n\n this.logger.debug('Trying to reconnect', this.connector.getUrl())\n\n this.connector.stop()\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n\n // Increase reconnect attempts counter\n this.reconnectAttempts++\n\n // Update state for connection and channels\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n for (const channel of this.channels.values()) {\n if (channel.getState() === DXLinkChannelState.CLOSED) continue\n\n channel.processStatusRequested()\n }\n\n // Schedule reconnect attempt after some time\n // Additionally, task will be executed in case when tab is active again\n // coz browser sometimes doesn't run scheduled tasks when tab is inactive\n this.scheduler.schedule(\n () => {\n if (this.connector === undefined) return\n\n // Start new connection attempt\n this.connector.start()\n },\n this.reconnectAttempts * 1000,\n DXLWS_SCHEDULER_KEY_RECONNECT\n )\n }\n\n disconnect = () => {\n if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return\n\n this.logger.debug('Disconnecting')\n\n // Destroy connector\n this.connector?.stop()\n this.connector = undefined\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n this.reconnectAttempts = 0\n\n this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED)\n this.setAuthState(DXLinkAuthState.UNAUTHORIZED)\n }\n\n close = () => {\n this.disconnect()\n }\n\n getConnectionDetails = () => this.connectionDetails\n getConnectionState = () => this.connectionState\n getScheduler = () => this.scheduler\n addConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.add(listener)\n removeConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.delete(listener)\n\n setAuthToken = (token: string): void => {\n this.lastSettedAuthToken = token\n\n if (this.connectionState === DXLinkConnectionState.CONNECTED) {\n this.sendAuthMessage(token)\n }\n }\n\n getAuthState = (): DXLinkAuthState => this.authState\n addAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.add(listener)\n removeAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n openChannel = (service: string, parameters: Record<string, unknown>): DXLinkChannel => {\n const channelId = this.globalChannelId\n this.globalChannelId += 2\n\n const channel = new DXLinkWebSocketChannel(\n channelId,\n service,\n parameters,\n this.sendMessage,\n this.config\n )\n\n this.channels.set(channelId, channel)\n\n // Send channel request if connection is already established\n if (\n this.connectionState === DXLinkConnectionState.CONNECTED &&\n this.authState === DXLinkAuthState.AUTHORIZED\n ) {\n channel.request()\n }\n\n return channel\n }\n\n private setConnectionState = (newStatus: DXLinkConnectionState) => {\n const prev = this.connectionState\n if (prev === newStatus) return\n\n this.connectionState = newStatus\n for (const listener of this.connectionStateChangeListeners) {\n listener(newStatus, prev)\n }\n }\n\n private sendMessage = (message: DXLinkWebSocketMessage): void => {\n this.connector?.sendMessage(message)\n\n this.scheduleKeepalive()\n\n // TODO: mb move to connector\n this.lastSentMillis = Date.now()\n }\n\n private sendAuthMessage = (token: string): void => {\n this.logger.debug('Sending auth message')\n\n this.setAuthState(DXLinkAuthState.AUTHORIZING)\n\n this.sendMessage({\n type: 'AUTH',\n channel: 0,\n token,\n })\n }\n\n private setAuthState = (newState: DXLinkAuthState): void => {\n const prev = this.authState\n\n this.authState = newState\n for (const listener of this.authStateChangeListeners) {\n try {\n listener(newState, prev)\n } catch (e) {\n this.logger.error('Auth state listener error', e)\n }\n }\n }\n\n private processMessage = (message: DXLinkWebSocketMessage): void => {\n this.lastReceivedMillis = Date.now()\n\n // Send keepalive message if no messages sent for a while (keepaliveInterval)\n // Because browser sometimes doesn't run scheduled tasks when tab is inactive\n if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1000) {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n }\n\n // Connection messages are messages that are sent to the channel 0\n if (isConnectionMessage(message)) {\n switch (message.type) {\n case 'SETUP':\n return this.processSetupMessage(message)\n case 'AUTH_STATE':\n return this.processAuthStateMessage(message)\n case 'ERROR':\n return this.publishError({\n type: message.error,\n message: message.message,\n })\n case 'KEEPALIVE':\n // Ignore keepalive messages coz they are used only to maintain connection\n return\n }\n } else if (isChannelMessage(message)) {\n const channel = this.channels.get(message.channel)\n if (channel === undefined) {\n this.logger.warn('Received lifecycle message for unknown channel', message)\n return\n }\n\n if (isChannelLifecycleMessage(message)) {\n switch (message.type) {\n case 'CHANNEL_OPENED':\n return channel.processStatusOpened()\n case 'CHANNEL_CLOSED':\n return channel.processStatusClosed()\n case 'ERROR':\n return channel.processError({\n type: message.error,\n message: message.message,\n })\n }\n return\n }\n\n return channel.processPayloadMessage(message)\n }\n\n this.logger.warn('Unhandeled message', message.type)\n }\n\n private processSetupMessage = (serverSetup: SetupMessage): void => {\n // Clear setup timeout check from connect method\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT)\n\n // Mark connection as connected after first setup message and subsequent ones\n if (\n this.connectionState === DXLinkConnectionState.CONNECTING ||\n this.connectionState === DXLinkConnectionState.CONNECTED\n ) {\n this.connectionDetails = {\n ...this.connectionDetails,\n serverVersion: serverSetup.version,\n clientKeepaliveTimeout: this.config.keepaliveTimeout,\n serverKeepaliveTimeout: serverSetup.keepaliveTimeout,\n }\n\n // Reset reconnect attempts counter after successful connection\n this.reconnectAttempts = 0\n\n if (this.lastSettedAuthToken === undefined) {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n }\n }\n\n // Connection maintance: Setup keepalive timeout check\n const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1000\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n timeoutMills,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private publishError = (error: DXLinkError): void => {\n this.logger.debug('Publishing error', error)\n\n if (this.errorListeners.size === 0) {\n this.logger.error('Unhandled dxLink error', error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error('Error listener error', e)\n }\n }\n }\n\n private processAuthStateMessage = ({ state }: AuthStateMessage): void => {\n this.logger.debug('Received auth state message', state)\n\n // Clear auth state timeout check\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT)\n\n // Ignore first auth state message because it is sent during connection setup\n if (this.isFirstAuthState) {\n this.isFirstAuthState = false\n } else {\n // Reset auth token if server rejected it\n if (state === 'UNAUTHORIZED') {\n this.lastSettedAuthToken = undefined\n }\n }\n\n // Request active channels if connection is authorized\n if (state === 'AUTHORIZED') {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n\n this.requestActiveChannels()\n }\n\n this.setAuthState(DXLinkAuthState[state])\n }\n\n private requestActiveChannels = (): void => {\n for (const channel of this.channels.values()) {\n // clear closed channels\n if (channel.getState() === DXLinkChannelState.CLOSED) {\n this.channels.delete(channel.id)\n continue\n }\n\n channel.request()\n }\n }\n\n /**\n * Process transport open event from connector.\n * After transport is opened:\n * - setup message is sent to server\n * - auth message is sent to server if auth token is set\n * - wait for setup message from server\n * - wait for auth state message from server\n */\n private processTransportOpen = (): void => {\n this.logger.debug('Connection opened')\n\n const setupMessage: SetupMessage = {\n type: 'SETUP',\n channel: 0,\n version: `${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,\n keepaliveTimeout: this.config.keepaliveTimeout,\n acceptKeepaliveTimeout: this.config.acceptKeepaliveTimeout,\n }\n\n // Setup timeout check\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No setup message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no setup message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT\n )\n\n this.sendMessage(setupMessage)\n\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No auth state message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no auth state message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT\n )\n\n if (this.lastSettedAuthToken !== undefined) {\n this.sendAuthMessage(this.lastSettedAuthToken)\n }\n }\n\n /**\n * Process transport close event from connector.\n * After transport is closed by server:\n * - reconnect if connection is authorized\n * - disconnect if connection is not authorized\n */\n private processTransportClose = (reason: string, error: boolean): void => {\n this.logger.debug('Connection closed', reason)\n\n if (error !== undefined) {\n this.publishError({\n type: 'UNKNOWN',\n message: reason,\n })\n }\n\n if (this.authState === DXLinkAuthState.UNAUTHORIZED) {\n this.lastSettedAuthToken = undefined\n this.disconnect()\n return\n }\n\n this.reconnect()\n }\n\n private timeoutCheck = (timeoutMills: number) => {\n const now = Date.now()\n const noKeepaliveDuration = now - this.lastReceivedMillis\n if (noKeepaliveDuration >= timeoutMills) {\n this.sendMessage({\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No keepalive received for ' + noKeepaliveDuration + 'ms',\n })\n\n return this.reconnect()\n }\n\n const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration)\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n nextTimeout,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private scheduleKeepalive = () => {\n this.scheduler.schedule(\n () => {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n\n this.scheduleKeepalive()\n },\n this.config.keepaliveInterval * 1000,\n DXLWS_SCHEDULER_KEY_KEEPALIVE\n )\n }\n}\n","export interface AuthMessage {\n type: 'AUTH'\n channel: 0\n token: string\n}\n\nexport type AuthState = 'AUTHORIZED' | 'UNAUTHORIZED'\n\nexport interface AuthStateMessage {\n type: 'AUTH_STATE'\n channel: 0\n state: AuthState\n}\n\nexport interface SetupMessage {\n type: 'SETUP'\n channel: 0\n version: string\n keepaliveTimeout?: number\n acceptKeepaliveTimeout?: number\n}\n\nexport interface KeepaliveMessage {\n type: 'KEEPALIVE'\n channel: 0\n}\n\nexport interface ChannelRequestMessage {\n type: 'CHANNEL_REQUEST'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelCancelMessage {\n type: 'CHANNEL_CANCEL'\n channel: number\n}\n\nexport interface ChannelOpenedMessage {\n type: 'CHANNEL_OPENED'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelClosedMessage {\n type: 'CHANNEL_CLOSED'\n channel: number\n}\n\nexport type ErrorType =\n | 'UNKNOWN'\n | 'UNSUPPORTED_PROTOCOL'\n | 'TIMEOUT'\n | 'UNAUTHORIZED'\n | 'INVALID_MESSAGE'\n | 'BAD_ACTION'\n\nexport interface ErrorMessage {\n type: 'ERROR'\n channel: 0\n error: ErrorType\n message: string\n}\n\nexport interface ChannelErrorMessage {\n type: 'ERROR'\n channel: number\n error: ErrorType\n message: string\n}\n\nexport interface ChannelPayloadMessage {\n type: string\n channel: number\n [key: string]: unknown\n}\n\nexport type ConnectionMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n\nexport type DXLinkWebSocketMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelPayloadMessage\n | ChannelErrorMessage\n\nexport const isConnectionMessage = (\n message: DXLinkWebSocketMessage\n): message is ConnectionMessage =>\n message.channel === 0 &&\n (message.type === 'SETUP' ||\n message.type === 'KEEPALIVE' ||\n message.type === 'AUTH' ||\n message.type === 'AUTH_STATE' ||\n message.type === 'ERROR')\n\nexport type ChannelLifecycleMessage =\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n\nexport type ChannelMessage = ChannelLifecycleMessage | ChannelPayloadMessage\n\nexport const isChannelMessage = (message: DXLinkWebSocketMessage): message is ChannelMessage =>\n message.channel !== 0\n\nexport const isChannelLifecycleMessage = (\n message: ChannelMessage\n): message is ChannelLifecycleMessage =>\n message.type === 'CHANNEL_OPENED' ||\n message.type === 'CHANNEL_CLOSED' ||\n message.type === 'ERROR' ||\n message.type === 'CHANNEL_REQUEST' ||\n message.type === 'CHANNEL_CANCEL'\n"],"names":["DXLinkWebSocketChannel","constructor","id","service","parameters","sendMessage","config","this","status","DXLinkChannelState","REQUESTED","messageListeners","Set","statusListeners","errorListeners","logger","send","type","payload","OPENED","Error","channel","addMessageListener","listener","add","removeMessageListener","delete","getState","addStateChangeListener","removeStateChangeListener","addErrorListener","removeErrorListener","error","message","close","CLOSED","debug","setStatus","clear","request","processStatusRequested","processPayloadMessage","processStatusOpened","processStatusClosed","processError","size","e","newStatus","prev","Logger","name","logLevel","DefaultDXLinkWebSocketConnector","url","protocols","socket","undefined","isAvailable","openListener","closeListener","messageListener","JSON","stringify","setOpenListener","setCloseListener","setMessageListener","getUrl","handleOpen","removeEventListener","addEventListener","handleMessage","handleClosed","ev","stop","reason","handleError","_ev","parse","data","console","warn","String","start","WebSocket","DXLINK_WS_PROTOCOL_VERSION","DEFAULT_CONNECTION_DETAILS","protocolVersion","clientVersion","DXLinkWebSocketClient","scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","connectorFactory","processTransportOpen","processMessage","processTransportClose","reconnect","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","getScheduler","addConnectionStateChangeListener","removeConnectionStateChangeListener","setAuthToken","token","CONNECTED","sendAuthMessage","getAuthState","addAuthStateChangeListener","removeAuthStateChangeListener","openChannel","channelId","set","AUTHORIZED","scheduleKeepalive","Date","now","AUTHORIZING","newState","keepaliveInterval","isConnectionMessage","processSetupMessage","processAuthStateMessage","publishError","isChannelMessage","get","isChannelLifecycleMessage","serverSetup","cancel","serverVersion","version","clientKeepaliveTimeout","keepaliveTimeout","serverKeepaliveTimeout","timeoutMills","timeoutCheck","state","requestActiveChannels","setupMessage","acceptKeepaliveTimeout","errorMessage","actionTimeout","noKeepaliveDuration","nextTimeout","Math","max","DXLinkLogLevel","WARN","DefaultDXLinkScheduler"],"mappings":"6IAmBaA,uBAUXC,WAAAA,CACkBC,GACAC,QACAC,WACCC,YACjBC,aAJgBJ,QAAA,EAAAK,KACAJ,aAAA,EAAAI,KACAH,gBACCC,EAAAA,KAAAA,iBAbXG,EAAAA,KAAAA,OAASC,mBAAmBC,eAGnBC,iBAAmB,IAAIC,IACvBC,KAAAA,gBAAkB,IAAID,IAAuCL,KAC7DO,eAAiB,IAAIF,IAE9BG,KAAAA,mBAYRC,KAAO,EAAGC,aAASC,YACjB,GAAIX,KAAKC,SAAWC,mBAAmBU,OACrC,MAAM,IAAIC,MAAM,wBAGlBb,KAAKF,YAAY,CACfY,UACAI,QAASd,KAAKL,MACXgB,SACJ,EACFX,KAEDe,mBAAsBC,UACpBhB,KAAKI,iBAAiBa,IAAID,UAC5BE,KAAAA,sBAAyBF,UACvBhB,KAAKI,iBAAiBe,OAAOH,UAE/BI,KAAAA,SAAW,IAAMpB,KAAKC,OACtBoB,KAAAA,uBAA0BL,UACxBhB,KAAKM,gBAAgBW,IAAID,UAAShB,KACpCsB,0BAA6BN,UAC3BhB,KAAKM,gBAAgBa,OAAOH,eAE9BO,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAC9EQ,KAAAA,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7FyB,MAAQ,EAAGf,UAAMgB,mBACf1B,KAAKS,KAAK,CACRC,KAAM,QACNe,MAAOf,KACPgB,kBAGJC,KAAAA,MAAQ,KACF3B,KAAKC,SAAWC,mBAAmB0B,SAEvC5B,KAAKQ,OAAOqB,MAAM,mBAGlB7B,KAAK8B,UAAU5B,mBAAmB0B,QAElC5B,KAAK+B,QAEL/B,KAAKF,YAAY,CACfY,KAAM,iBACNI,QAASd,KAAKL,KAElB,OAEAqC,QAAU,KACRhC,KAAKQ,OAAOqB,MAAM,cAElB7B,KAAKF,YAAY,CACfY,KAAM,kBACNI,QAASd,KAAKL,GACdC,QAASI,KAAKJ,QACdC,WAAYG,KAAKH,aAGnBG,KAAKiC,wBAAsB,OAG7BC,sBAAyBR,UACvB,IAAK,MAAMV,YAAgBhB,KAACI,iBAC1BY,SAASU,QACV,EACF1B,KAEDmC,oBAAsB,KACpBnC,KAAKQ,OAAOqB,MAAM,UAElB7B,KAAK8B,UAAU5B,mBAAmBU,SAGpCqB,KAAAA,uBAAyB,KACvBjC,KAAK8B,UAAU5B,mBAAmBC,UAAS,OAG7CiC,oBAAsB,KACpBpC,KAAKQ,OAAOqB,MAAM,6BAElB7B,KAAK8B,UAAU5B,mBAAmB0B,QAClC5B,KAAK+B,cAGPM,aAAgBZ,QACd,GAAiC,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAYhB,KAAKO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKL,sBAAuB4C,EACnE,MATDvC,KAAKQ,OAAOiB,MAAM,8BAA8BzB,KAAKL,OAAQ8B,MAU9D,EAGKK,KAAAA,UAAaU,YACnB,GAAIxC,KAAKC,SAAWuC,UAAW,OAE/B,MAAMC,KAAOzC,KAAKC,OAClBD,KAAKC,OAASuC,UACd,IAAK,MAAMxB,YAAYhB,KAAKM,gBAC1B,IACEU,SAASwB,UAAWC,KACrB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKL,uBAAwB4C,EACpE,CACF,EAGKR,KAAAA,MAAQ,KACd/B,KAAKI,iBAAiB2B,QACtB/B,KAAKM,gBAAgByB,OAGvB,EAhIkB/B,KAAEL,GAAFA,GACAK,KAAOJ,QAAPA,QACAI,KAAUH,WAAVA,WACCG,KAAWF,YAAXA,YAGjBE,KAAKQ,OAAS,IAAIkC,OAAO,GAAGjD,uBAAuBkD,QAAQhD,MAAMC,UAAWG,OAAO6C,SACrF,QCeWC,gCASXnD,WAAAA,CACmBoD,IACAC,WAA6B/C,KAD7B8C,SAAA,EAAA9C,KACA+C,eAAA,EAAA/C,KAVXgD,YAAgCC,EAASjD,KAEzCkD,aAAc,EAAKlD,KAEnBmD,kBAAyCF,EACzCG,KAAAA,mBAA0DH,EAC1DI,KAAAA,qBAA2EJ,EA8BnFnD,KAAAA,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAAC1B,KAEDwD,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QAAAA,EACrBhB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAEA0C,KAAAA,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QAAAA,EAGzB2C,KAAAA,OAAS,IAAM3D,KAAK8C,IAEZc,KAAAA,WAAa,UACCX,IAAhBjD,KAAKgD,SAEThD,KAAKkD,aAAc,EAEnBlD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAE7C5D,KAAKgD,OAAOc,iBAAiB,UAAW9D,KAAK+D,eAE7C/D,KAAKmD,iBACP,EAACnD,KAEOgE,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAACpD,KAEO+D,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,MAAM,IAAIb,MAAM,8BAAgCa,SAElD,GAA4B,iBAAjBA,QAAQhB,KACjB,MAAU,IAAAG,MAAM,mCAAqCa,QAAQhB,MAE/D,GAA+B,iBAApBgB,QAAQZ,QACjB,MAAM,IAAID,MAAM,sCAAwCa,QAAQZ,SAKlE,QAA6BmC,IAAzBjD,KAAKqD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBzE,KAAKqD,gBAAgB3B,QACtB,CAAC,MAAOD,OACP+C,QAAQ/C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB6D,OAAOjD,QACpF,GAnGgBzB,KAAG8C,IAAHA,IACA9C,KAAS+C,UAATA,SAChB,CAEH4B,KAAAA,QACsB1B,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAS,IAAI4B,UAAU5E,KAAK8C,IAAK9C,KAAK+C,WAE3C/C,KAAKgD,OAAOc,iBAAiB,OAAQ9D,KAAK4D,YAC1C5D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKoE,aAC3CpE,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAC7C5D,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKoE,aAC9CpE,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKgE,cAC9ChE,KAAKgD,OAAOa,oBAAoB,UAAW7D,KAAK+D,eAEhD/D,KAAKgD,OAAOrB,QACZ3B,KAAKgD,YAASC,EACdjD,KAAKkD,aAAc,EACrB,ECpDW,MAAA2B,2BAA6B,MAWpCC,2BAAsD,CAC1DC,gBAZwC,MAaxCC,cAXqB,sBAiBVC,sBA+CXvF,WAAAA,CAAYK,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENkF,eAAS,EAAAlF,KAElBmF,eAAS,EAAAnF,KAEToF,gBAAyCC,sBAAsBC,cAAatF,KAC5EuF,kBAA6CT,2BAA0B9E,KAEvEwF,UAA6BC,gBAAgBC,aAGpCC,KAAAA,+BAAiC,IAAItF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBuF,KAAAA,yBAA2B,IAAIvF,IAMxCwF,KAAAA,kBAAmB,EAInBC,KAAAA,yBAIAC,EAAAA,KAAAA,mBAAqB,EACrBC,KAAAA,eAAiB,EAKjBC,KAAAA,kBAAoB,EAGpBC,KAAAA,gBAAkB,EACTC,KAAAA,SAAW,IAAIC,IAsBhCC,KAAAA,QAAWvD,MAEL9C,KAAKmF,WAAWxB,WAAab,MAGjC9C,KAAKsG,aAELtG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKuG,mBAAmBlB,sBAAsBmB,YAG9CxG,KAAKmF,UAAYnF,KAAKD,OAAO0G,iBAAiB3D,KAC9C9C,KAAKmF,UAAU3B,gBAAgBxD,KAAK0G,sBACpC1G,KAAKmF,UAAUzB,mBAAmB1D,KAAK2G,gBACvC3G,KAAKmF,UAAU1B,iBAAiBzD,KAAK4G,uBAGrC5G,KAAKmF,UAAUR,QACjB,EAAC3E,KAED6G,UAAY,KACV,GACE7G,KAAKD,OAAO+G,sBAAwB,GACpC9G,KAAKiG,mBAAqBjG,KAAKD,OAAO+G,qBAKtC,OAHA9G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKsG,aAIP,GACEtG,KAAKoF,kBAAoBC,sBAAsBC,oBAC5BrC,IAAnBjD,KAAKmF,UAFP,CAMAnF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKmF,UAAUxB,UAExD3D,KAAKmF,UAAUjB,OAGflE,KAAKkF,UAAUnD,QAGf/B,KAAKuF,kBAAoBT,2BACzB9E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EAGxB7F,KAAKiG,oBAGLjG,KAAKuG,mBAAmBlB,sBAAsBmB,YAC9C,IAAK,MAAM1F,WAAWd,KAAKmG,SAASY,SAC9BjG,QAAQM,aAAelB,mBAAmB0B,QAE9Cd,QAAQmB,yBAMVjC,KAAKkF,UAAU8B,SACb,UACyB/D,IAAnBjD,KAAKmF,WAGTnF,KAAKmF,UAAUR,OAAK,EAEG,IAAzB3E,KAAKiG,kBAtJ2B,kBAkHhC,CAuCJ,EAACjG,KAEDsG,WAAa,KACPtG,KAAKoF,kBAAoBC,sBAAsBC,gBAEnDtF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKmF,WAAWjB,OAChBlE,KAAKmF,eAAYlC,EAGjBjD,KAAKkF,UAAUnD,QAGf/B,KAAKuF,kBAAoBT,2BACzB9E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EACxB7F,KAAKiG,kBAAoB,EAEzBjG,KAAKuG,mBAAmBlB,sBAAsBC,eAC9CtF,KAAKiH,aAAaxB,gBAAgBC,cACpC,EAAC1F,KAED2B,MAAQ,KACN3B,KAAKsG,YACP,EAEAY,KAAAA,qBAAuB,IAAMlH,KAAKuF,kBAAiBvF,KACnDmH,mBAAqB,IAAMnH,KAAKoF,gBAAepF,KAC/CoH,aAAe,IAAMpH,KAAKkF,UAC1BmC,KAAAA,iCAAoCrG,UAClChB,KAAK2F,+BAA+B1E,IAAID,UAAShB,KACnDsH,oCAAuCtG,UACrChB,KAAK2F,+BAA+BxE,OAAOH,UAE7CuG,KAAAA,aAAgBC,QACdxH,KAAK8F,oBAAsB0B,MAEvBxH,KAAKoF,kBAAoBC,sBAAsBoC,WACjDzH,KAAK0H,gBAAgBF,MACtB,EACFxH,KAED2H,aAAe,IAAuB3H,KAAKwF,UAC3CoC,KAAAA,2BAA8B5G,UAC5BhB,KAAK4F,yBAAyB3E,IAAID,UAAShB,KAC7C6H,8BAAiC7G,UAC/BhB,KAAK4F,yBAAyBzE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAEpF8G,KAAAA,YAAc,CAAClI,QAAiBC,cAC9B,MAAMkI,UAAY/H,KAAKkG,gBACvBlG,KAAKkG,iBAAmB,EAExB,MAAMpF,QAAU,IAAIrB,uBAClBsI,UACAnI,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKmG,SAAS6B,IAAID,UAAWjH,SAI3Bd,KAAKoF,kBAAoBC,sBAAsBoC,WAC/CzH,KAAKwF,YAAcC,gBAAgBwC,YAEnCnH,QAAQkB,UAGHlB,SAGDyF,KAAAA,mBAAsB/D,YAC5B,MAAMC,KAAOzC,KAAKoF,gBAClB,GAAI3C,OAASD,UAAb,CAEAxC,KAAKoF,gBAAkB5C,UACvB,IAAK,MAAMxB,YAAYhB,KAAK2F,+BAC1B3E,SAASwB,UAAWC,KAFtB,CAGC,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKmF,WAAWrF,YAAY4B,SAE5B1B,KAAKkI,oBAGLlI,KAAKgG,eAAiBmC,KAAKC,KAC7B,EAEQV,KAAAA,gBAAmBF,QACzBxH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAKiH,aAAaxB,gBAAgB4C,aAElCrI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACT0G,aACD,EACFxH,KAEOiH,aAAgBqB,WACtB,MAAM7F,KAAOzC,KAAKwF,UAElBxF,KAAKwF,UAAY8C,SACjB,IAAK,MAAMtH,YAAgBhB,KAAC4F,yBAC1B,IACE5E,SAASsH,SAAU7F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EACFvC,KAEO2G,eAAkBjF,UAaxB,GAZA1B,KAAK+F,mBAAqBoC,KAAKC,MAI3BpI,KAAK+F,mBAAqB/F,KAAKgG,gBAAkD,IAAhChG,KAAKD,OAAOwI,mBAC/DvI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IChOfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MD8NJ8H,CAAoB9G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACyI,oBAAoB/G,SAClC,IAAK,aACH,OAAW1B,KAAC0I,wBAAwBhH,SACtC,IAAK,QACH,OAAO1B,KAAK2I,aAAa,CACvBjI,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GClOsBA,UACX,IAApBA,QAAQZ,QDiOK8H,CAAiBlH,SAAU,CACpC,MAAMZ,QAAUd,KAAKmG,SAAS0C,IAAInH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GCrOJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KD+NAoI,CAA0BpH,SAAU,CACtC,OAAQA,QAAQhB,MACd,IAAK,iBACH,OAAOI,QAAQqB,sBACjB,IAAK,iBACH,OAAOrB,QAAQsB,sBACjB,IAAK,QACH,OAAOtB,QAAQuB,aAAa,CAC1B3B,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAGvB,MACD,CAED,OAAOZ,QAAQoB,sBAAsBR,QACtC,CAED1B,KAAKQ,OAAOiE,KAAK,qBAAsB/C,QAAQhB,KACjD,EAEQ+H,KAAAA,oBAAuBM,cAE7B/I,KAAKkF,UAAU8D,OA7UuB,uBAiVpChJ,KAAKoF,kBAAoBC,sBAAsBmB,YAC/CxG,KAAKoF,kBAAoBC,sBAAsBoC,YAE/CzH,KAAKuF,kBAAoB,IACpBvF,KAAKuF,kBACR0D,cAAeF,YAAYG,QAC3BC,uBAAwBnJ,KAAKD,OAAOqJ,iBACpCC,uBAAwBN,YAAYK,kBAItCpJ,KAAKiG,kBAAoB,OAEQhD,IAA7BjD,KAAK8F,qBACP9F,KAAKuG,mBAAmBlB,sBAAsBoC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDpJ,KAAKkF,UAAU8B,SACb,IAAMhH,KAAKuJ,aAAaD,cACxBA,aArW8B,gBAwWlC,EAACtJ,KAEO2I,aAAgBlH,QAGtB,GAFAzB,KAAKQ,OAAOqB,MAAM,mBAAoBJ,OAEL,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAYhB,KAAKO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,uBAAwBc,EAC3C,MATDvC,KAAKQ,OAAOiB,MAAM,yBAA0BA,MAU7C,EAGKiH,KAAAA,wBAA0B,EAAGc,gBACnCxJ,KAAKQ,OAAOqB,MAAM,8BAA+B2H,OAGjDxJ,KAAKkF,UAAU8D,OAhY4B,4BAmYvChJ,KAAK6F,iBACP7F,KAAK6F,kBAAmB,EAGV,iBAAV2D,QACFxJ,KAAK8F,yBAAsB7C,GAKjB,eAAVuG,QACFxJ,KAAKuG,mBAAmBlB,sBAAsBoC,WAE9CzH,KAAKyJ,yBAGPzJ,KAAKiH,aAAaxB,gBAAgB+D,OAAM,EACzCxJ,KAEOyJ,sBAAwB,KAC9B,IAAK,MAAM3I,WAAWd,KAAKmG,SAASY,SAE9BjG,QAAQM,aAAelB,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKmG,SAAShF,OAAOL,QAAQnB,GAKhC,EAWK+G,KAAAA,qBAAuB,KAC7B1G,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM6H,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACToI,QAAS,GAAGlJ,KAAKuF,kBAAkBR,mBAAmB/E,KAAKuF,kBAAkBP,gBAC7EoE,iBAAkBpJ,KAAKD,OAAOqJ,iBAC9BO,uBAAwB3J,KAAKD,OAAO4J,wBAItC3J,KAAKkF,UAAU8B,SACb,KACE,MAAM4C,aAA6B,CACjClJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,iCAAmC1B,KAAKD,OAAO8J,cAAgB,KAG1E7J,KAAKF,YAAY8J,cAEjB5J,KAAK2I,aAAa,CAChBjI,KAAMkJ,aAAanI,MACnBC,QAAS,GAAGkI,aAAalI,wBAI3B1B,KAAK6G,WAAS,EAEY,IAA5B7G,KAAKD,OAAO8J,cA1cwB,uBA8ctC7J,KAAKF,YAAY4J,cAEjB1J,KAAKkF,UAAU8B,SACb,KACE,MAAM4C,aAA6B,CACjClJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,sCAAwC1B,KAAKD,OAAO8J,cAAgB,KAG/E7J,KAAKF,YAAY8J,cAEjB5J,KAAK2I,aAAa,CAChBjI,KAAMkJ,aAAanI,MACnBC,QAAS,GAAGkI,aAAalI,wBAI3B1B,KAAK6G,WAAS,EAEY,IAA5B7G,KAAKD,OAAO8J,cAle6B,iCAseV5G,IAA7BjD,KAAK8F,qBACP9F,KAAK0H,gBAAgB1H,KAAK8F,oBAC3B,EACF9F,KAQO4G,sBAAwB,CAACzC,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAK2I,aAAa,CAChBjI,KAAM,UACNgB,QAASyC,SAITnE,KAAKwF,YAAcC,gBAAgBC,aAGrC,OAFA1F,KAAK8F,yBAAsB7C,OAC3BjD,KAAKsG,aAIPtG,KAAK6G,WACP,EAEQ0C,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBpI,KAAK+F,mBACvC,GAAI+D,qBAAuBR,aAQzB,OAPAtJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BoI,oBAAsB,OAGzD9J,KAAK6G,YAGd,MAAMkD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD9J,KAAKkF,UAAU8B,SACb,IAAMhH,KAAKuJ,aAAaD,cACxBS,YAphB8B,gBAuhBlC,EAAC/J,KAEOkI,kBAAoB,KAC1BlI,KAAKkF,UAAU8B,SACb,KACEhH,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKkI,mBACP,EACgC,IAAhClI,KAAKD,OAAOwI,kBAliBoB,kBAqiBpC,EA3eEvI,KAAKD,OAAS,CACZwI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACfjH,SAAUsH,eAAeC,KACzBrD,sBAAuB,EACvBL,iBAAmB3D,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,OAAO1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,UAC5D5C,KAAKkF,UAAYlF,KAAKD,OAAOmF,WAAa,IAAIkF,sBAChD"}
{"version":3,"file":"index.module.js","sources":["../src/channel.ts","../src/connector.ts","../src/client.ts","../src/messages.ts"],"sourcesContent":["import {\n type DXLinkLogger,\n Logger,\n type DXLinkChannel,\n DXLinkChannelState,\n type DXLinkChannelMessageListener,\n type DXLinkChannelStateChangeListener,\n type DXLinkErrorListener,\n type DXLinkChannelMessage,\n type DXLinkError,\n} from '@dxfeed/dxlink-core'\n\nimport { type DXLinkWebSocketClientConfig } from './config'\nimport { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages'\n\n/**\n * A DXLink channel implementation.\n * @internal\n */\nexport class DXLinkWebSocketChannel implements DXLinkChannel {\n private status = DXLinkChannelState.REQUESTED\n\n // Listeners\n private readonly messageListeners = new Set<DXLinkChannelMessageListener>()\n private readonly statusListeners = new Set<DXLinkChannelStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n\n private logger: DXLinkLogger\n\n constructor(\n public readonly id: number,\n public readonly service: string,\n public readonly parameters: Record<string, unknown>,\n public readonly reconnect: boolean,\n private readonly sendMessage: (message: DXLinkWebSocketMessage) => void,\n config: DXLinkWebSocketClientConfig\n ) {\n this.logger = new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel)\n }\n\n send = ({ type, ...payload }: DXLinkChannelMessage) => {\n if (this.status !== DXLinkChannelState.OPENED) {\n throw new Error('Channel is not ready')\n }\n\n this.sendMessage({\n type,\n channel: this.id,\n ...payload,\n })\n }\n\n addMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.add(listener)\n removeMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.delete(listener)\n\n getState = () => this.status\n addStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.add(listener)\n removeStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n error = ({ type, message }: DXLinkError) =>\n this.send({\n type: 'ERROR',\n error: type,\n message,\n })\n\n close = () => {\n if (this.status === DXLinkChannelState.CLOSED) return\n\n this.logger.debug(`Closing by user`)\n\n // We can think that channel is closed already\n this.setStatus(DXLinkChannelState.CLOSED)\n\n this.clear()\n\n this.sendMessage({\n type: 'CHANNEL_CANCEL',\n channel: this.id,\n })\n }\n\n request = () => {\n this.logger.debug('Requesting')\n\n this.sendMessage({\n type: 'CHANNEL_REQUEST',\n channel: this.id,\n service: this.service,\n parameters: this.parameters,\n })\n\n this.processStatusRequested()\n }\n\n processPayloadMessage = (message: ChannelPayloadMessage) => {\n for (const listener of this.messageListeners) {\n listener(message)\n }\n }\n\n processStatusOpened = () => {\n this.logger.debug('Opened')\n\n this.setStatus(DXLinkChannelState.OPENED)\n }\n\n processStatusRequested = () => {\n this.setStatus(DXLinkChannelState.REQUESTED)\n }\n\n processStatusClosed = () => {\n this.logger.debug('Closed by remote endpoint')\n\n this.setStatus(DXLinkChannelState.CLOSED)\n this.clear()\n }\n\n processError = (error: DXLinkError) => {\n if (this.errorListeners.size === 0) {\n this.logger.error(`Unhandled error in channel#${this.id}: `, error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} error listener: `, e)\n }\n }\n }\n\n private setStatus = (newStatus: DXLinkChannelState) => {\n if (this.status === newStatus) return\n\n const prev = this.status\n this.status = newStatus\n for (const listener of this.statusListeners) {\n try {\n listener(newStatus, prev)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} status listener: `, e)\n }\n }\n }\n\n private clear = () => {\n this.messageListeners.clear()\n this.statusListeners.clear()\n // TODO: rethink approach when error came after channel is closed\n // this.errorListeners.clear()\n }\n}\n","import { type DXLinkWebSocketMessage } from './messages'\n\n/**\n * Interface for a WebSocket connector that manages the connection to a WebSocket server.\n * It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.\n */\nexport interface DXLinkWebSocketConnector {\n /**\n * Returns the URL of the WebSocket connection.\n */\n getUrl(): string\n /**\n * Starts the WebSocket connection.\n */\n start(): void\n /**\n * Stops the WebSocket connection and cleans up resources.\n */\n stop(): void\n /**\n * Sends a message over the WebSocket connection.\n * @param message The message to send.\n */\n sendMessage(message: DXLinkWebSocketMessage): void\n /**\n * Sets a listener that is called when the WebSocket connection is opened.\n * @param listener The listener function to call when the connection is opened.\n */\n setOpenListener(listener: () => void): void\n /**\n * Sets a listener that is called when the WebSocket connection is closed.\n * @param listener The listener function to call when the connection is closed.\n */\n setCloseListener(listener: DXLinkWebSocketCloseListener): void\n /**\n * Sets a listener that is called when a message is received from the WebSocket server.\n * @param listener The listener function to call when a message is received.\n */\n setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void\n}\n\n/**\n * Type for a close listener that is called when the WebSocket connection is closed.\n * @param reason - The reason for the closure.\n * @param error - Indicates if the closure was due to an error.\n */\nexport type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void\n\n/**\n * Default connector for the WebSocket connection.\n * @internal\n */\nexport class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: DXLinkWebSocketCloseListener | undefined = undefined\n private messageListener: ((message: DXLinkWebSocketMessage) => void) | undefined = undefined\n\n constructor(\n private readonly url: string,\n private readonly protocols?: string | string[]\n ) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url, this.protocols)\n\n this.socket.addEventListener('open', this.handleOpen)\n this.socket.addEventListener('error', this.handleError)\n this.socket.addEventListener('close', this.handleClosed)\n }\n\n stop() {\n if (this.socket === undefined) return\n\n this.socket.removeEventListener('open', this.handleOpen)\n this.socket.removeEventListener('error', this.handleError)\n this.socket.removeEventListener('close', this.handleClosed)\n this.socket.removeEventListener('message', this.handleMessage)\n\n this.socket.close()\n this.socket = undefined\n this.isAvailable = false\n }\n\n sendMessage = (message: DXLinkWebSocketMessage) => {\n if (this.socket === undefined || !this.isAvailable) {\n return\n }\n\n this.socket.send(JSON.stringify(message))\n }\n\n setOpenListener = (listener: () => void) => {\n this.openListener = listener\n }\n\n setCloseListener = (listener: DXLinkWebSocketCloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: DXLinkWebSocketMessage) => void) => {\n this.messageListener = listener\n }\n\n getUrl = () => this.url\n\n private handleOpen = () => {\n if (this.socket === undefined) return\n\n this.isAvailable = true\n\n this.socket.removeEventListener('open', this.handleOpen)\n\n this.socket.addEventListener('message', this.handleMessage)\n\n this.openListener?.()\n }\n\n private handleClosed = (ev: CloseEvent) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.(ev.reason, false)\n }\n\n private handleError = (_ev: Event) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.('Unable to connect', true)\n }\n\n private handleMessage = (ev: MessageEvent) => {\n try {\n const message = JSON.parse(ev.data)\n if (typeof message !== 'object') {\n throw new Error('Unexpected message: ' + typeof message)\n }\n if (typeof message.type !== 'string') {\n throw new Error('Unexpected message type: ' + typeof message.type)\n }\n if (typeof message.channel !== 'number') {\n throw new Error('Unexpected message channel: ' + typeof message.channel)\n }\n\n // TODO: validate message properly (e.g. check for type and other fields)\n\n if (this.messageListener === undefined) {\n return console.warn('No message listener set')\n }\n\n this.messageListener(message)\n } catch (error) {\n console.error(error instanceof Error ? error : new Error('Parsing error:' + String(error)))\n }\n }\n}\n","import {\n DXLinkLogLevel,\n type DXLinkLogger,\n Logger,\n type DXLinkScheduler,\n DefaultDXLinkScheduler,\n type DXLinkConnectionDetails,\n DXLinkConnectionState,\n type DXLinkConnectionStateChangeListener,\n type DXLinkErrorListener,\n DXLinkAuthState,\n DXLinkChannelState,\n type DXLinkAuthStateChangeListener,\n type DXLinkChannel,\n type DXLinkChannelOptions,\n type DXLinkError,\n type DXLinkClient,\n} from '@dxfeed/dxlink-core'\n\nimport { DXLinkWebSocketChannel } from './channel'\nimport type { DXLinkWebSocketClientConfig } from './config'\nimport { type DXLinkWebSocketConnector, DefaultDXLinkWebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type DXLinkWebSocketMessage,\n type SetupMessage,\n isChannelLifecycleMessage,\n isChannelMessage,\n isConnectionMessage,\n} from './messages'\nimport { VERSION } from './version'\n\n/**\n * Protocol version that is used by client.\n */\nexport const DXLINK_WS_PROTOCOL_VERSION = '0.1'\n\nconst CLIENT_VERSION = `DXF-JS/${VERSION}`\n\n// Scheduler keys\nconst DXLWS_SCHEDULER_KEY_RECONNECT = 'DXLWS_RECONNECT'\nconst DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = 'DXLWS_SETUP_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = 'DXLWS_AUTH_STATE_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_TIMEOUT = 'DXLWS_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_KEEPALIVE = 'DXLWS_KEEPALIVE'\n\nconst DEFAULT_CONNECTION_DETAILS: DXLinkConnectionDetails = {\n protocolVersion: DXLINK_WS_PROTOCOL_VERSION,\n clientVersion: CLIENT_VERSION,\n}\n\n/**\n * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.\n */\nexport class DXLinkWebSocketClient implements DXLinkClient {\n private readonly config: DXLinkWebSocketClientConfig\n\n private readonly logger: DXLinkLogger\n\n private readonly scheduler: DXLinkScheduler\n\n private connector: DXLinkWebSocketConnector | undefined\n\n private connectionState: DXLinkConnectionState = DXLinkConnectionState.NOT_CONNECTED\n private connectionDetails: DXLinkConnectionDetails = DEFAULT_CONNECTION_DETAILS\n\n private authState: DXLinkAuthState = DXLinkAuthState.UNAUTHORIZED\n\n // Listeners\n private readonly connectionStateChangeListeners = new Set<DXLinkConnectionStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n private readonly authStateChangeListeners = new Set<DXLinkAuthStateChangeListener>()\n\n /**\n * Authorization type that was determined by server behavior during setup phase.\n * This value is used to determine if authorization is required or optional or not defined yet.\n */\n private isFirstAuthState = true\n /**\n * Last setted auth token that will be sent to server after connection is established or re-established.\n */\n private lastSettedAuthToken: string | undefined\n\n // Stats for keepalive\n // TODO: mb move to connector\n private lastReceivedMillis = 0\n private lastSentMillis = 0\n\n /**\n * Count of reconnect attempts since last successful connection.\n */\n private reconnectAttempts = 0\n\n // Channels\n private globalChannelId = 1\n private readonly channels = new Map<number, DXLinkWebSocketChannel>()\n\n /**\n * Create new instance of {@link DXLinkWebSocketClient}.\n * @param config Configuration of the client.\n */\n constructor(config?: Partial<DXLinkWebSocketClientConfig>) {\n this.config = {\n keepaliveInterval: 30,\n keepaliveTimeout: 60,\n acceptKeepaliveTimeout: 60,\n actionTimeout: 10,\n logLevel: DXLinkLogLevel.WARN,\n maxReconnectAttempts: -1,\n connectorFactory: (url) => new DefaultDXLinkWebSocketConnector(url),\n ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\n this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler()\n }\n\n connect = (url: string) => {\n // Do nothing if already connected to the same url\n if (this.connector?.getUrl() === url) return\n\n // Disconnect from previous connection if any exists\n this.disconnect()\n\n this.logger.debug('Connecting to', url)\n\n // Immediately set connection state to CONNECTING\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n\n // Create new connector\n this.connector = this.config.connectorFactory(url)\n this.connector.setOpenListener(this.processTransportOpen)\n this.connector.setMessageListener(this.processMessage)\n this.connector.setCloseListener(this.processTransportClose)\n\n // Initiate websocket connection\n this.connector.start()\n }\n\n reconnect = () => {\n if (\n this.config.maxReconnectAttempts >= 0 &&\n this.reconnectAttempts >= this.config.maxReconnectAttempts\n ) {\n this.logger.warn('Max reconnect attempts reached')\n\n this.disconnect()\n return\n }\n\n if (\n this.connectionState === DXLinkConnectionState.NOT_CONNECTED ||\n this.connector === undefined\n )\n return\n\n this.logger.debug('Trying to reconnect', this.connector.getUrl())\n\n this.connector.stop()\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n\n // Increase reconnect attempts counter\n this.reconnectAttempts++\n\n // Update state for connection and channels\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n for (const channel of this.channels.values()) {\n if (channel.getState() === DXLinkChannelState.CLOSED) continue\n\n if (!channel.reconnect) {\n channel.processStatusClosed()\n continue\n }\n\n channel.processStatusRequested()\n }\n\n // Schedule reconnect attempt after some time\n // Additionally, task will be executed in case when tab is active again\n // coz browser sometimes doesn't run scheduled tasks when tab is inactive\n this.scheduler.schedule(\n () => {\n if (this.connector === undefined) return\n\n // Start new connection attempt\n this.connector.start()\n },\n this.reconnectAttempts * 1000,\n DXLWS_SCHEDULER_KEY_RECONNECT\n )\n }\n\n disconnect = () => {\n if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return\n\n this.logger.debug('Disconnecting')\n\n // Destroy connector\n this.connector?.stop()\n this.connector = undefined\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n this.reconnectAttempts = 0\n\n this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED)\n this.setAuthState(DXLinkAuthState.UNAUTHORIZED)\n }\n\n close = () => {\n this.disconnect()\n }\n\n getConnectionDetails = () => this.connectionDetails\n getConnectionState = () => this.connectionState\n getScheduler = () => this.scheduler\n addConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.add(listener)\n removeConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.delete(listener)\n\n setAuthToken = (token: string): void => {\n this.lastSettedAuthToken = token\n\n if (this.connectionState === DXLinkConnectionState.CONNECTED) {\n this.sendAuthMessage(token)\n }\n }\n\n getAuthState = (): DXLinkAuthState => this.authState\n addAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.add(listener)\n removeAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n openChannel = (\n service: string,\n parameters: Record<string, unknown>,\n options?: DXLinkChannelOptions\n ): DXLinkChannel => {\n const channelId = this.globalChannelId\n this.globalChannelId += 2\n\n const channel = new DXLinkWebSocketChannel(\n channelId,\n service,\n parameters,\n options?.reconnect ?? true,\n this.sendMessage,\n this.config\n )\n\n this.channels.set(channelId, channel)\n\n // Send channel request if connection is already established\n if (\n this.connectionState === DXLinkConnectionState.CONNECTED &&\n this.authState === DXLinkAuthState.AUTHORIZED\n ) {\n channel.request()\n }\n\n return channel\n }\n\n private setConnectionState = (newStatus: DXLinkConnectionState) => {\n const prev = this.connectionState\n if (prev === newStatus) return\n\n this.connectionState = newStatus\n for (const listener of this.connectionStateChangeListeners) {\n listener(newStatus, prev)\n }\n }\n\n private sendMessage = (message: DXLinkWebSocketMessage): void => {\n this.connector?.sendMessage(message)\n\n this.scheduleKeepalive()\n\n // TODO: mb move to connector\n this.lastSentMillis = Date.now()\n }\n\n private sendAuthMessage = (token: string): void => {\n this.logger.debug('Sending auth message')\n\n this.setAuthState(DXLinkAuthState.AUTHORIZING)\n\n this.sendMessage({\n type: 'AUTH',\n channel: 0,\n token,\n })\n }\n\n private setAuthState = (newState: DXLinkAuthState): void => {\n const prev = this.authState\n\n this.authState = newState\n for (const listener of this.authStateChangeListeners) {\n try {\n listener(newState, prev)\n } catch (e) {\n this.logger.error('Auth state listener error', e)\n }\n }\n }\n\n private processMessage = (message: DXLinkWebSocketMessage): void => {\n this.lastReceivedMillis = Date.now()\n\n // Send keepalive message if no messages sent for a while (keepaliveInterval)\n // Because browser sometimes doesn't run scheduled tasks when tab is inactive\n if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1000) {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n }\n\n // Connection messages are messages that are sent to the channel 0\n if (isConnectionMessage(message)) {\n switch (message.type) {\n case 'SETUP':\n return this.processSetupMessage(message)\n case 'AUTH_STATE':\n return this.processAuthStateMessage(message)\n case 'ERROR':\n return this.publishError({\n type: message.error,\n message: message.message,\n })\n case 'KEEPALIVE':\n // Ignore keepalive messages coz they are used only to maintain connection\n return\n }\n } else if (isChannelMessage(message)) {\n const channel = this.channels.get(message.channel)\n if (channel === undefined) {\n this.logger.warn('Received lifecycle message for unknown channel', message)\n return\n }\n\n if (isChannelLifecycleMessage(message)) {\n switch (message.type) {\n case 'CHANNEL_OPENED':\n return channel.processStatusOpened()\n case 'CHANNEL_CLOSED':\n return channel.processStatusClosed()\n case 'ERROR':\n return channel.processError({\n type: message.error,\n message: message.message,\n })\n }\n return\n }\n\n return channel.processPayloadMessage(message)\n }\n\n this.logger.warn('Unhandeled message', message.type)\n }\n\n private processSetupMessage = (serverSetup: SetupMessage): void => {\n // Clear setup timeout check from connect method\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT)\n\n // Mark connection as connected after first setup message and subsequent ones\n if (\n this.connectionState === DXLinkConnectionState.CONNECTING ||\n this.connectionState === DXLinkConnectionState.CONNECTED\n ) {\n this.connectionDetails = {\n ...this.connectionDetails,\n serverVersion: serverSetup.version,\n clientKeepaliveTimeout: this.config.keepaliveTimeout,\n serverKeepaliveTimeout: serverSetup.keepaliveTimeout,\n }\n\n // Reset reconnect attempts counter after successful connection\n this.reconnectAttempts = 0\n\n if (this.lastSettedAuthToken === undefined) {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n }\n }\n\n // Connection maintance: Setup keepalive timeout check\n const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1000\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n timeoutMills,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private publishError = (error: DXLinkError): void => {\n this.logger.debug('Publishing error', error)\n\n if (this.errorListeners.size === 0) {\n this.logger.error('Unhandled dxLink error', error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error('Error listener error', e)\n }\n }\n }\n\n private processAuthStateMessage = ({ state }: AuthStateMessage): void => {\n this.logger.debug('Received auth state message', state)\n\n // Clear auth state timeout check\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT)\n\n // Ignore first auth state message because it is sent during connection setup\n if (this.isFirstAuthState) {\n this.isFirstAuthState = false\n } else {\n // Reset auth token if server rejected it\n if (state === 'UNAUTHORIZED') {\n this.lastSettedAuthToken = undefined\n }\n }\n\n // Request active channels if connection is authorized\n if (state === 'AUTHORIZED') {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n\n this.requestActiveChannels()\n }\n\n this.setAuthState(DXLinkAuthState[state])\n }\n\n private requestActiveChannels = (): void => {\n for (const channel of this.channels.values()) {\n // clear closed channels\n if (channel.getState() === DXLinkChannelState.CLOSED) {\n this.channels.delete(channel.id)\n continue\n }\n\n channel.request()\n }\n }\n\n /**\n * Process transport open event from connector.\n * After transport is opened:\n * - setup message is sent to server\n * - auth message is sent to server if auth token is set\n * - wait for setup message from server\n * - wait for auth state message from server\n */\n private processTransportOpen = (): void => {\n this.logger.debug('Connection opened')\n\n const setupMessage: SetupMessage = {\n type: 'SETUP',\n channel: 0,\n version: `${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,\n keepaliveTimeout: this.config.keepaliveTimeout,\n acceptKeepaliveTimeout: this.config.acceptKeepaliveTimeout,\n }\n\n // Setup timeout check\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No setup message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no setup message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT\n )\n\n this.sendMessage(setupMessage)\n\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No auth state message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no auth state message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT\n )\n\n if (this.lastSettedAuthToken !== undefined) {\n this.sendAuthMessage(this.lastSettedAuthToken)\n }\n }\n\n /**\n * Process transport close event from connector.\n * After transport is closed by server:\n * - reconnect if connection is authorized\n * - disconnect if connection is not authorized\n */\n private processTransportClose = (reason: string, error: boolean): void => {\n this.logger.debug('Connection closed', reason)\n\n if (error !== undefined) {\n this.publishError({\n type: 'UNKNOWN',\n message: reason,\n })\n }\n\n if (this.authState === DXLinkAuthState.UNAUTHORIZED) {\n this.lastSettedAuthToken = undefined\n this.disconnect()\n return\n }\n\n this.reconnect()\n }\n\n private timeoutCheck = (timeoutMills: number) => {\n const now = Date.now()\n const noKeepaliveDuration = now - this.lastReceivedMillis\n if (noKeepaliveDuration >= timeoutMills) {\n this.sendMessage({\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No keepalive received for ' + noKeepaliveDuration + 'ms',\n })\n\n return this.reconnect()\n }\n\n const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration)\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n nextTimeout,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private scheduleKeepalive = () => {\n this.scheduler.schedule(\n () => {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n\n this.scheduleKeepalive()\n },\n this.config.keepaliveInterval * 1000,\n DXLWS_SCHEDULER_KEY_KEEPALIVE\n )\n }\n}\n","export interface AuthMessage {\n type: 'AUTH'\n channel: 0\n token: string\n}\n\nexport type AuthState = 'AUTHORIZED' | 'UNAUTHORIZED'\n\nexport interface AuthStateMessage {\n type: 'AUTH_STATE'\n channel: 0\n state: AuthState\n}\n\nexport interface SetupMessage {\n type: 'SETUP'\n channel: 0\n version: string\n keepaliveTimeout?: number\n acceptKeepaliveTimeout?: number\n}\n\nexport interface KeepaliveMessage {\n type: 'KEEPALIVE'\n channel: 0\n}\n\nexport interface ChannelRequestMessage {\n type: 'CHANNEL_REQUEST'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelCancelMessage {\n type: 'CHANNEL_CANCEL'\n channel: number\n}\n\nexport interface ChannelOpenedMessage {\n type: 'CHANNEL_OPENED'\n channel: number\n service: string\n parameters?: Record<string, unknown>\n}\n\nexport interface ChannelClosedMessage {\n type: 'CHANNEL_CLOSED'\n channel: number\n}\n\nexport type ErrorType =\n | 'UNKNOWN'\n | 'UNSUPPORTED_PROTOCOL'\n | 'TIMEOUT'\n | 'UNAUTHORIZED'\n | 'INVALID_MESSAGE'\n | 'BAD_ACTION'\n\nexport interface ErrorMessage {\n type: 'ERROR'\n channel: 0\n error: ErrorType\n message: string\n}\n\nexport interface ChannelErrorMessage {\n type: 'ERROR'\n channel: number\n error: ErrorType\n message: string\n}\n\nexport interface ChannelPayloadMessage {\n type: string\n channel: number\n [key: string]: unknown\n}\n\nexport type ConnectionMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n\nexport type DXLinkWebSocketMessage =\n | SetupMessage\n | KeepaliveMessage\n | AuthMessage\n | AuthStateMessage\n | ErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelPayloadMessage\n | ChannelErrorMessage\n\nexport const isConnectionMessage = (\n message: DXLinkWebSocketMessage\n): message is ConnectionMessage =>\n message.channel === 0 &&\n (message.type === 'SETUP' ||\n message.type === 'KEEPALIVE' ||\n message.type === 'AUTH' ||\n message.type === 'AUTH_STATE' ||\n message.type === 'ERROR')\n\nexport type ChannelLifecycleMessage =\n | ChannelOpenedMessage\n | ChannelClosedMessage\n | ChannelErrorMessage\n | ChannelRequestMessage\n | ChannelCancelMessage\n\nexport type ChannelMessage = ChannelLifecycleMessage | ChannelPayloadMessage\n\nexport const isChannelMessage = (message: DXLinkWebSocketMessage): message is ChannelMessage =>\n message.channel !== 0\n\nexport const isChannelLifecycleMessage = (\n message: ChannelMessage\n): message is ChannelLifecycleMessage =>\n message.type === 'CHANNEL_OPENED' ||\n message.type === 'CHANNEL_CLOSED' ||\n message.type === 'ERROR' ||\n message.type === 'CHANNEL_REQUEST' ||\n message.type === 'CHANNEL_CANCEL'\n"],"names":["DXLinkWebSocketChannel","constructor","id","service","parameters","reconnect","sendMessage","config","this","status","DXLinkChannelState","REQUESTED","messageListeners","Set","statusListeners","errorListeners","logger","send","type","payload","OPENED","Error","channel","addMessageListener","listener","add","removeMessageListener","delete","getState","addStateChangeListener","removeStateChangeListener","addErrorListener","removeErrorListener","error","message","close","CLOSED","debug","setStatus","clear","request","processStatusRequested","processPayloadMessage","processStatusOpened","processStatusClosed","processError","size","e","newStatus","prev","Logger","name","logLevel","DefaultDXLinkWebSocketConnector","url","protocols","socket","undefined","isAvailable","openListener","closeListener","messageListener","JSON","stringify","setOpenListener","setCloseListener","setMessageListener","getUrl","handleOpen","removeEventListener","addEventListener","handleMessage","handleClosed","ev","stop","reason","handleError","_ev","parse","data","console","warn","String","start","WebSocket","DXLINK_WS_PROTOCOL_VERSION","DEFAULT_CONNECTION_DETAILS","protocolVersion","clientVersion","DXLinkWebSocketClient","scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","connectorFactory","processTransportOpen","processMessage","processTransportClose","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","getScheduler","addConnectionStateChangeListener","removeConnectionStateChangeListener","setAuthToken","token","CONNECTED","sendAuthMessage","getAuthState","addAuthStateChangeListener","removeAuthStateChangeListener","openChannel","options","channelId","set","AUTHORIZED","scheduleKeepalive","Date","now","AUTHORIZING","newState","keepaliveInterval","isConnectionMessage","processSetupMessage","processAuthStateMessage","publishError","isChannelMessage","get","isChannelLifecycleMessage","serverSetup","cancel","serverVersion","version","clientKeepaliveTimeout","keepaliveTimeout","serverKeepaliveTimeout","timeoutMills","timeoutCheck","state","requestActiveChannels","setupMessage","acceptKeepaliveTimeout","errorMessage","actionTimeout","noKeepaliveDuration","nextTimeout","Math","max","DXLinkLogLevel","WARN","DefaultDXLinkScheduler"],"mappings":"6IAmBaA,uBAUXC,WAAAA,CACkBC,GACAC,QACAC,WACAC,UACCC,YACjBC,QAAmCC,KALnBN,QAAA,EAAAM,KACAL,aACAC,EAAAA,KAAAA,gBACAC,EAAAA,KAAAA,sBACCC,iBAAA,EAAAE,KAdXC,OAASC,mBAAmBC,UAGnBC,KAAAA,iBAAmB,IAAIC,IAAmCL,KAC1DM,gBAAkB,IAAID,SACtBE,eAAiB,IAAIF,IAE9BG,KAAAA,YAaRC,EAAAA,KAAAA,KAAO,EAAGC,aAASC,YACjB,GAAIX,KAAKC,SAAWC,mBAAmBU,OACrC,MAAU,IAAAC,MAAM,wBAGlBb,KAAKF,YAAY,CACfY,UACAI,QAASd,KAAKN,MACXiB,WAIPI,KAAAA,mBAAsBC,UACpBhB,KAAKI,iBAAiBa,IAAID,UAAShB,KACrCkB,sBAAyBF,UACvBhB,KAAKI,iBAAiBe,OAAOH,UAAShB,KAExCoB,SAAW,IAAMpB,KAAKC,OAAMD,KAC5BqB,uBAA0BL,UACxBhB,KAAKM,gBAAgBW,IAAID,UAC3BM,KAAAA,0BAA6BN,UAC3BhB,KAAKM,gBAAgBa,OAAOH,UAE9BO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7FyB,MAAQ,EAAGf,UAAMgB,mBACf1B,KAAKS,KAAK,CACRC,KAAM,QACNe,MAAOf,KACPgB,kBAGJC,KAAAA,MAAQ,KACF3B,KAAKC,SAAWC,mBAAmB0B,SAEvC5B,KAAKQ,OAAOqB,MAAM,mBAGlB7B,KAAK8B,UAAU5B,mBAAmB0B,QAElC5B,KAAK+B,QAEL/B,KAAKF,YAAY,CACfY,KAAM,iBACNI,QAASd,KAAKN,OAEjBM,KAEDgC,QAAU,KACRhC,KAAKQ,OAAOqB,MAAM,cAElB7B,KAAKF,YAAY,CACfY,KAAM,kBACNI,QAASd,KAAKN,GACdC,QAASK,KAAKL,QACdC,WAAYI,KAAKJ,aAGnBI,KAAKiC,wBAAsB,EAC5BjC,KAEDkC,sBAAyBR,UACvB,IAAK,MAAMV,iBAAiBZ,iBAC1BY,SAASU,QACV,EAGHS,KAAAA,oBAAsB,KACpBnC,KAAKQ,OAAOqB,MAAM,UAElB7B,KAAK8B,UAAU5B,mBAAmBU,SACnCZ,KAEDiC,uBAAyB,KACvBjC,KAAK8B,UAAU5B,mBAAmBC,UAAS,EAG7CiC,KAAAA,oBAAsB,KACpBpC,KAAKQ,OAAOqB,MAAM,6BAElB7B,KAAK8B,UAAU5B,mBAAmB0B,QAClC5B,KAAK+B,OAAK,EACX/B,KAEDqC,aAAgBZ,QACd,GAAiC,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,iBAAiBT,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKN,sBAAuB6C,EACnE,MATDvC,KAAKQ,OAAOiB,MAAM,8BAA8BzB,KAAKN,OAAQ+B,MAU9D,EACFzB,KAEO8B,UAAaU,YACnB,GAAIxC,KAAKC,SAAWuC,UAAW,OAE/B,MAAMC,KAAOzC,KAAKC,OAClBD,KAAKC,OAASuC,UACd,IAAK,MAAMxB,YAAYhB,KAAKM,gBAC1B,IACEU,SAASwB,UAAWC,KACrB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,oBAAoBzB,KAAKN,uBAAwB6C,EACpE,CACF,EACFvC,KAEO+B,MAAQ,KACd/B,KAAKI,iBAAiB2B,QACtB/B,KAAKM,gBAAgByB,OAAK,EA9HV/B,KAAEN,GAAFA,GACAM,KAAOL,QAAPA,QACAK,KAAUJ,WAAVA,WACAI,KAASH,UAATA,UACCG,KAAWF,YAAXA,YAGjBE,KAAKQ,OAAS,IAAIkC,OAAO,GAAGlD,uBAAuBmD,QAAQjD,MAAMC,UAAWI,OAAO6C,SACrF,QCcWC,gCASXpD,WAAAA,CACmBqD,IACAC,WAA6B/C,KAD7B8C,SAAA,EAAA9C,KACA+C,eAAA,EAAA/C,KAVXgD,YAAgCC,EAASjD,KAEzCkD,aAAc,EAAKlD,KAEnBmD,kBAAyCF,EACzCG,KAAAA,mBAA0DH,EAC1DI,KAAAA,qBAA2EJ,EA8BnFnD,KAAAA,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAAC1B,KAEDwD,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QAAAA,EACrBhB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAEA0C,KAAAA,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QAAAA,EAGzB2C,KAAAA,OAAS,IAAM3D,KAAK8C,IAEZc,KAAAA,WAAa,UACCX,IAAhBjD,KAAKgD,SAEThD,KAAKkD,aAAc,EAEnBlD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAE7C5D,KAAKgD,OAAOc,iBAAiB,UAAW9D,KAAK+D,eAE7C/D,KAAKmD,iBACP,EAACnD,KAEOgE,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAACpD,KAEO+D,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,MAAM,IAAIb,MAAM,8BAAgCa,SAElD,GAA4B,iBAAjBA,QAAQhB,KACjB,MAAU,IAAAG,MAAM,mCAAqCa,QAAQhB,MAE/D,GAA+B,iBAApBgB,QAAQZ,QACjB,MAAM,IAAID,MAAM,sCAAwCa,QAAQZ,SAKlE,QAA6BmC,IAAzBjD,KAAKqD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBzE,KAAKqD,gBAAgB3B,QACtB,CAAC,MAAOD,OACP+C,QAAQ/C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB6D,OAAOjD,QACpF,GAnGgBzB,KAAG8C,IAAHA,IACA9C,KAAS+C,UAATA,SAChB,CAEH4B,KAAAA,QACsB1B,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAS,IAAI4B,UAAU5E,KAAK8C,IAAK9C,KAAK+C,WAE3C/C,KAAKgD,OAAOc,iBAAiB,OAAQ9D,KAAK4D,YAC1C5D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKoE,aAC3CpE,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBjD,KAAKgD,SAEThD,KAAKgD,OAAOa,oBAAoB,OAAQ7D,KAAK4D,YAC7C5D,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKoE,aAC9CpE,KAAKgD,OAAOa,oBAAoB,QAAS7D,KAAKgE,cAC9ChE,KAAKgD,OAAOa,oBAAoB,UAAW7D,KAAK+D,eAEhD/D,KAAKgD,OAAOrB,QACZ3B,KAAKgD,YAASC,EACdjD,KAAKkD,aAAc,EACrB,ECnDW,MAAA2B,2BAA6B,MAWpCC,2BAAsD,CAC1DC,gBAZwC,MAaxCC,cAXqB,sBAiBVC,sBA+CXxF,WAAAA,CAAYM,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENkF,eAAS,EAAAlF,KAElBmF,eAAS,EAAAnF,KAEToF,gBAAyCC,sBAAsBC,cAAatF,KAC5EuF,kBAA6CT,2BAA0B9E,KAEvEwF,UAA6BC,gBAAgBC,aAGpCC,KAAAA,+BAAiC,IAAItF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBuF,KAAAA,yBAA2B,IAAIvF,IAMxCwF,KAAAA,kBAAmB,EAInBC,KAAAA,yBAIAC,EAAAA,KAAAA,mBAAqB,EACrBC,KAAAA,eAAiB,EAAChG,KAKlBiG,kBAAoB,EAACjG,KAGrBkG,gBAAkB,EAAClG,KACVmG,SAAW,IAAIC,IAAqCpG,KAsBrEqG,QAAWvD,MAEL9C,KAAKmF,WAAWxB,WAAab,MAGjC9C,KAAKsG,aAELtG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKuG,mBAAmBlB,sBAAsBmB,YAG9CxG,KAAKmF,UAAYnF,KAAKD,OAAO0G,iBAAiB3D,KAC9C9C,KAAKmF,UAAU3B,gBAAgBxD,KAAK0G,sBACpC1G,KAAKmF,UAAUzB,mBAAmB1D,KAAK2G,gBACvC3G,KAAKmF,UAAU1B,iBAAiBzD,KAAK4G,uBAGrC5G,KAAKmF,UAAUR,QAAK,EAGtB9E,KAAAA,UAAY,KACV,GACEG,KAAKD,OAAO8G,sBAAwB,GACpC7G,KAAKiG,mBAAqBjG,KAAKD,OAAO8G,qBAKtC,OAHA7G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKsG,aAIP,GACEtG,KAAKoF,kBAAoBC,sBAAsBC,oBAC5BrC,IAAnBjD,KAAKmF,UAFP,CAMAnF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKmF,UAAUxB,UAExD3D,KAAKmF,UAAUjB,OAGflE,KAAKkF,UAAUnD,QAGf/B,KAAKuF,kBAAoBT,2BACzB9E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EAGxB7F,KAAKiG,oBAGLjG,KAAKuG,mBAAmBlB,sBAAsBmB,YAC9C,IAAK,MAAM1F,WAAed,KAACmG,SAASW,SAC9BhG,QAAQM,aAAelB,mBAAmB0B,SAEzCd,QAAQjB,UAKbiB,QAAQmB,yBAJNnB,QAAQsB,uBAUZpC,KAAKkF,UAAU6B,SACb,UACyB9D,IAAnBjD,KAAKmF,WAGTnF,KAAKmF,UAAUR,OAAK,EAEG,IAAzB3E,KAAKiG,kBA3J2B,kBAkHhC,CA0C6B,EAIjCK,KAAAA,WAAa,KACPtG,KAAKoF,kBAAoBC,sBAAsBC,gBAEnDtF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKmF,WAAWjB,OAChBlE,KAAKmF,eAAYlC,EAGjBjD,KAAKkF,UAAUnD,QAGf/B,KAAKuF,kBAAoBT,2BACzB9E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EACxB7F,KAAKiG,kBAAoB,EAEzBjG,KAAKuG,mBAAmBlB,sBAAsBC,eAC9CtF,KAAKgH,aAAavB,gBAAgBC,cAAY,EAC/C1F,KAED2B,MAAQ,KACN3B,KAAKsG,YAAU,OAGjBW,qBAAuB,IAAMjH,KAAKuF,kBAClC2B,KAAAA,mBAAqB,IAAMlH,KAAKoF,gBAAepF,KAC/CmH,aAAe,IAAMnH,KAAKkF,UAC1BkC,KAAAA,iCAAoCpG,UAClChB,KAAK2F,+BAA+B1E,IAAID,UAC1CqG,KAAAA,oCAAuCrG,UACrChB,KAAK2F,+BAA+BxE,OAAOH,UAAShB,KAEtDsH,aAAgBC,QACdvH,KAAK8F,oBAAsByB,MAEvBvH,KAAKoF,kBAAoBC,sBAAsBmC,WACjDxH,KAAKyH,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuB1H,KAAKwF,UAASxF,KACpD2H,2BAA8B3G,UAC5BhB,KAAK4F,yBAAyB3E,IAAID,UACpC4G,KAAAA,8BAAiC5G,UAC/BhB,KAAK4F,yBAAyBzE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7F6H,YAAc,CACZlI,QACAC,WACAkI,WAEA,MAAMC,UAAY/H,KAAKkG,gBACvBlG,KAAKkG,iBAAmB,EAExB,MAAMpF,QAAU,IAAItB,uBAClBuI,UACApI,QACAC,WACAkI,SAASjI,YAAa,EACtBG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKmG,SAAS6B,IAAID,UAAWjH,SAI3Bd,KAAKoF,kBAAoBC,sBAAsBmC,WAC/CxH,KAAKwF,YAAcC,gBAAgBwC,YAEnCnH,QAAQkB,UAGHlB,SAGDyF,KAAAA,mBAAsB/D,YAC5B,MAAMC,KAAOzC,KAAKoF,gBAClB,GAAI3C,OAASD,UAAb,CAEAxC,KAAKoF,gBAAkB5C,UACvB,IAAK,MAAMxB,YAAgBhB,KAAC2F,+BAC1B3E,SAASwB,UAAWC,KAFtB,CAGC,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKmF,WAAWrF,YAAY4B,SAE5B1B,KAAKkI,oBAGLlI,KAAKgG,eAAiBmC,KAAKC,KAAG,EAGxBX,KAAAA,gBAAmBF,QACzBvH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAKgH,aAAavB,gBAAgB4C,aAElCrI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTyG,aACD,EAGKP,KAAAA,aAAgBsB,WACtB,MAAM7F,KAAOzC,KAAKwF,UAElBxF,KAAKwF,UAAY8C,SACjB,IAAK,MAAMtH,YAAgBhB,KAAC4F,yBAC1B,IACE5E,SAASsH,SAAU7F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EAGKoE,KAAAA,eAAkBjF,UAaxB,GAZA1B,KAAK+F,mBAAqBoC,KAAKC,MAI3BpI,KAAK+F,mBAAqB/F,KAAKgG,gBAAkD,IAAhChG,KAAKD,OAAOwI,mBAC/DvI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IC3OfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDyOJ8H,CAAoB9G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACyI,oBAAoB/G,SAClC,IAAK,aACH,OAAO1B,KAAK0I,wBAAwBhH,SACtC,IAAK,QACH,OAAW1B,KAAC2I,aAAa,CACvBjI,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GC7OsBA,UACX,IAApBA,QAAQZ,QD4OK8H,CAAiBlH,SAAU,CACpC,MAAMZ,QAAUd,KAAKmG,SAAS0C,IAAInH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GChPJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KD0OAoI,CAA0BpH,SAAU,CACtC,OAAQA,QAAQhB,MACd,IAAK,iBACH,OAAOI,QAAQqB,sBACjB,IAAK,iBACH,OAAOrB,QAAQsB,sBACjB,IAAK,QACH,OAAOtB,QAAQuB,aAAa,CAC1B3B,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAGvB,MACD,CAED,OAAOZ,QAAQoB,sBAAsBR,QACtC,CAED1B,KAAKQ,OAAOiE,KAAK,qBAAsB/C,QAAQhB,KAAI,EAG7C+H,KAAAA,oBAAuBM,cAE7B/I,KAAKkF,UAAU8D,OAvVuB,uBA2VpChJ,KAAKoF,kBAAoBC,sBAAsBmB,YAC/CxG,KAAKoF,kBAAoBC,sBAAsBmC,YAE/CxH,KAAKuF,kBAAoB,IACpBvF,KAAKuF,kBACR0D,cAAeF,YAAYG,QAC3BC,uBAAwBnJ,KAAKD,OAAOqJ,iBACpCC,uBAAwBN,YAAYK,kBAItCpJ,KAAKiG,kBAAoB,OAEQhD,IAA7BjD,KAAK8F,qBACP9F,KAAKuG,mBAAmBlB,sBAAsBmC,YAKlD,MAAM8B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDpJ,KAAKkF,UAAU6B,SACb,IAAM/G,KAAKuJ,aAAaD,cACxBA,aA/W8B,gBAkXlC,EAACtJ,KAEO2I,aAAgBlH,QAGtB,GAFAzB,KAAKQ,OAAOqB,MAAM,mBAAoBJ,OAEL,IAA7BzB,KAAKO,eAAe+B,KAKxB,IAAK,MAAMtB,YAAgBhB,KAACO,eAC1B,IACES,SAASS,MACV,CAAC,MAAOc,GACPvC,KAAKQ,OAAOiB,MAAM,uBAAwBc,EAC3C,MATDvC,KAAKQ,OAAOiB,MAAM,yBAA0BA,MAU7C,EAGKiH,KAAAA,wBAA0B,EAAGc,gBACnCxJ,KAAKQ,OAAOqB,MAAM,8BAA+B2H,OAGjDxJ,KAAKkF,UAAU8D,OA1Y4B,4BA6YvChJ,KAAK6F,iBACP7F,KAAK6F,kBAAmB,EAGV,iBAAV2D,QACFxJ,KAAK8F,yBAAsB7C,GAKjB,eAAVuG,QACFxJ,KAAKuG,mBAAmBlB,sBAAsBmC,WAE9CxH,KAAKyJ,yBAGPzJ,KAAKgH,aAAavB,gBAAgB+D,OACpC,EAACxJ,KAEOyJ,sBAAwB,KAC9B,IAAK,MAAM3I,WAAed,KAACmG,SAASW,SAE9BhG,QAAQM,aAAelB,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKmG,SAAShF,OAAOL,QAAQpB,GAKhC,EAWKgH,KAAAA,qBAAuB,KAC7B1G,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM6H,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACToI,QAAS,GAAGlJ,KAAKuF,kBAAkBR,mBAAmB/E,KAAKuF,kBAAkBP,gBAC7EoE,iBAAkBpJ,KAAKD,OAAOqJ,iBAC9BO,uBAAwB3J,KAAKD,OAAO4J,wBAItC3J,KAAKkF,UAAU6B,SACb,KACE,MAAM6C,aAA6B,CACjClJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,iCAAmC1B,KAAKD,OAAO8J,cAAgB,KAG1E7J,KAAKF,YAAY8J,cAEjB5J,KAAK2I,aAAa,CAChBjI,KAAMkJ,aAAanI,MACnBC,QAAS,GAAGkI,aAAalI,wBAI3B1B,KAAKH,WACP,EAC4B,IAA5BG,KAAKD,OAAO8J,cApdwB,uBAwdtC7J,KAAKF,YAAY4J,cAEjB1J,KAAKkF,UAAU6B,SACb,KACE,MAAM6C,aAA6B,CACjClJ,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,sCAAwC1B,KAAKD,OAAO8J,cAAgB,KAG/E7J,KAAKF,YAAY8J,cAEjB5J,KAAK2I,aAAa,CAChBjI,KAAMkJ,aAAanI,MACnBC,QAAS,GAAGkI,aAAalI,wBAI3B1B,KAAKH,WAAS,EAEY,IAA5BG,KAAKD,OAAO8J,cA5e6B,iCAgfV5G,IAA7BjD,KAAK8F,qBACP9F,KAAKyH,gBAAgBzH,KAAK8F,oBAC3B,OASKc,sBAAwB,CAACzC,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAK2I,aAAa,CAChBjI,KAAM,UACNgB,QAASyC,SAITnE,KAAKwF,YAAcC,gBAAgBC,aAGrC,OAFA1F,KAAK8F,yBAAsB7C,OAC3BjD,KAAKsG,aAIPtG,KAAKH,WACP,EAEQ0J,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBpI,KAAK+F,mBACvC,GAAI+D,qBAAuBR,aAQzB,OAPAtJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BoI,oBAAsB,OAGzD9J,KAAKH,YAGd,MAAMkK,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD9J,KAAKkF,UAAU6B,SACb,IAAM/G,KAAKuJ,aAAaD,cACxBS,YA9hB8B,gBA+hBH,EAE9B/J,KAEOkI,kBAAoB,KAC1BlI,KAAKkF,UAAU6B,SACb,KACE/G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKkI,mBAAiB,EAEQ,IAAhClI,KAAKD,OAAOwI,kBA5iBoB,kBA6iBH,EAnf/BvI,KAAKD,OAAS,CACZwI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACfjH,SAAUsH,eAAeC,KACzBtD,sBAAuB,EACvBJ,iBAAmB3D,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,OAAO1C,KAAKP,YAAYkD,KAAM3C,KAAKD,OAAO6C,UAC5D5C,KAAKkF,UAAYlF,KAAKD,OAAOmF,WAAa,IAAIkF,sBAChD"}
{
"name": "@dxfeed/dxlink-websocket-client",
"version": "0.8.0",
"version": "0.8.1",
"private": false,

@@ -33,3 +33,3 @@ "publishConfig": {

"dependencies": {
"@dxfeed/dxlink-core": "0.8.0"
"@dxfeed/dxlink-core": "0.8.1"
},

@@ -36,0 +36,0 @@ "author": "Dmitry Petrov <dmitry.petrov@devexperts.com>",