🎩 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.3.0
to
0.4.0
+2
-2
build/channel.d.ts
import { type DXLinkChannel, DXLinkChannelState, type DXLinkChannelMessageListener, type DXLinkChannelStateChangeListener, type DXLinkErrorListener, type DXLinkChannelMessage, type DXLinkError } from '@dxfeed/dxlink-core';
import { type DXLinkWebSocketClientConfig } from './config';
import { type ChannelPayloadMessage, type Message } from './messages';
import { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages';
/**

@@ -18,3 +18,3 @@ * A DXLink channel implementation.

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

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

import type { DXLinkLogLevel } from '@dxfeed/dxlink-core';
import type { DXLinkWebSocketConnector } from './connector';
/**

@@ -34,2 +35,10 @@ * Options for {@link DXLinkWebSocketClient}.

readonly maxReconnectAttempts: number;
/**
* Factory function to create a WebSocket connector.
* This function should return an instance of {@link DXLinkWebSocketConnector} for the given URL.
* This allows for custom WebSocket implementations or configurations.
* @param url The URL to connect to.
* @returns {@link DXLinkWebSocketConnector} instance
*/
readonly connectorFactory: (url: string) => DXLinkWebSocketConnector;
}

@@ -1,9 +0,53 @@

import { type Message } from './messages';
type CloseListener = (reason: string, error: boolean) => void;
import { type DXLinkWebSocketMessage } from './messages';
/**
* Connector for the WebSocket connection.
* Interface for a WebSocket connector that manages the connection to a WebSocket server.
* It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.
*/
export interface DXLinkWebSocketConnector {
/**
* Returns the URL of the WebSocket connection.
*/
getUrl(): string;
/**
* Starts the WebSocket connection.
*/
start(): void;
/**
* Stops the WebSocket connection and cleans up resources.
*/
stop(): void;
/**
* Sends a message over the WebSocket connection.
* @param message The message to send.
*/
sendMessage(message: DXLinkWebSocketMessage): void;
/**
* Sets a listener that is called when the WebSocket connection is opened.
* @param listener The listener function to call when the connection is opened.
*/
setOpenListener(listener: () => void): void;
/**
* Sets a listener that is called when the WebSocket connection is closed.
* @param listener The listener function to call when the connection is closed.
*/
setCloseListener(listener: DXLinkWebSocketCloseListener): void;
/**
* Sets a listener that is called when a message is received from the WebSocket server.
* @param listener The listener function to call when a message is received.
*/
setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void;
}
/**
* Type for a close listener that is called when the WebSocket connection is closed.
* @param reason - The reason for the closure.
* @param error - Indicates if the closure was due to an error.
*/
export type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void;
/**
* Default connector for the WebSocket connection.
* @internal
*/
export declare class WebSocketConnector {
export declare class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {
private readonly url;
private readonly protocols?;
private socket;

@@ -14,9 +58,9 @@ private isAvailable;

private messageListener;
constructor(url: string);
constructor(url: string, protocols?: string | string[] | undefined);
start(): void;
stop(): void;
sendMessage: (message: Message) => void;
sendMessage: (message: DXLinkWebSocketMessage) => void;
setOpenListener: (listener: () => void) => void;
setCloseListener: (listener: CloseListener) => void;
setMessageListener: (listener: (message: Message) => void) => void;
setCloseListener: (listener: DXLinkWebSocketCloseListener) => void;
setMessageListener: (listener: (message: DXLinkWebSocketMessage) => void) => void;
getUrl: () => string;

@@ -28,2 +72,1 @@ private handleOpen;

}
export {};
export { DXLinkWebSocketClient, DXLINK_WS_PROTOCOL_VERSION } from './client';
export { type DXLinkWebSocketClientConfig } from './config';
export { type DXLinkWebSocketConnector, type DXLinkWebSocketCloseListener, DefaultDXLinkWebSocketConnector, } from './connector';
export { type DXLinkWebSocketMessage } from './messages';

@@ -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 WebSocketConnector{constructor(url){this.url=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.socket.addEventListener("close",this.handleClosed),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}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url),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.3.0"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=new dxlinkCore.Scheduler,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=new WebSocketConnector(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,"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.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("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,"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("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,"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,"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,"TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:dxlinkCore.DXLinkLogLevel.WARN,maxReconnectAttempts:-1,...config},this.logger=new dxlinkCore.Logger(this.constructor.name,this.config.logLevel)}};
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.socket.addEventListener("close",this.handleClosed),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.4.0"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=new dxlinkCore.Scheduler,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,"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.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("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,"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("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,"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,"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,"TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"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)}},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 Message } 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: Message) => 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 Message } from './messages'\n\ntype CloseListener = (reason: string, error: boolean) => void\n\n/**\n * Connector for the WebSocket connection.\n * @internal\n */\nexport class WebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: CloseListener | undefined = undefined\n private messageListener: ((message: Message) => void) | undefined = undefined\n\n constructor(private readonly url: string) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url)\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: Message) => {\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: CloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: Message) => 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 this.socket.addEventListener('close', this.handleClosed)\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 Scheduler,\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 { WebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type Message,\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\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 = new Scheduler()\n\n private connector: WebSocketConnector | 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 ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\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 this.connector = new WebSocketConnector(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 '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 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: Message): 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: Message): 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('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(() => this.timeoutCheck(timeoutMills), timeoutMills, 'TIMEOUT')\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('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 '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 '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(() => this.timeoutCheck(timeoutMills), nextTimeout, 'TIMEOUT')\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 '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 Message =\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 = (message: Message): 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: Message): 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","WebSocketConnector","url","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","Scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","processTransportOpen","processMessage","processTransportClose","reconnect","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","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"],"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,QC7BWC,mBASXnD,WAAAA,CAA6BoD,KAAW9C,KAAX8C,SAAA,EAAA9C,KARrB+C,YAAgCC,EAAShD,KAEzCiD,aAAc,EAAKjD,KAEnBkD,kBAAyCF,EAAShD,KAClDmD,mBAA2CH,EAC3CI,KAAAA,qBAA4DJ,EA2BpElD,KAAAA,YAAe4B,eACOsB,IAAhBhD,KAAK+C,QAAyB/C,KAAKiD,aAIvCjD,KAAK+C,OAAOtC,KAAK4C,KAAKC,UAAU5B,SAAQ,EAG1C6B,KAAAA,gBAAmBvC,WACjBhB,KAAKkD,aAAelC,QACtB,EAEAwC,KAAAA,iBAAoBxC,WAClBhB,KAAKmD,cAAgBnC,QAAAA,EAGvByC,KAAAA,mBAAsBzC,WACpBhB,KAAKoD,gBAAkBpC,QACzB,EAAChB,KAED0D,OAAS,IAAM1D,KAAK8C,SAEZa,WAAa,UACCX,IAAhBhD,KAAK+C,SAET/C,KAAKiD,aAAc,EAEnBjD,KAAK+C,OAAOa,oBAAoB,OAAQ5D,KAAK2D,YAE7C3D,KAAK+C,OAAOc,iBAAiB,UAAW7D,KAAK8D,eAC7C9D,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAK+D,cAE3C/D,KAAKkD,iBACP,EAAClD,KAEO+D,aAAgBC,UACFhB,IAAhBhD,KAAK+C,SAET/C,KAAKiE,OAELjE,KAAKmD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBhD,KAAK+C,SAET/C,KAAKiE,OAELjE,KAAKmD,gBAAgB,qBAAqB,GAC5C,EAEQW,KAAAA,cAAiBE,KACvB,IACE,MAAMtC,QAAU2B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ5C,QACT,MAAU,IAAAb,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,QAA6BkC,IAAzBhD,KAAKoD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBxE,KAAKoD,gBAAgB1B,QACtB,CAAC,MAAOD,OACP8C,QAAQ9C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB4D,OAAOhD,QACpF,GAlG0BzB,KAAG8C,IAAHA,GAAc,CAE3C4B,KAAAA,QACsB1B,IAAhBhD,KAAK+C,SAET/C,KAAK+C,OAAS,IAAI4B,UAAU3E,KAAK8C,KAEjC9C,KAAK+C,OAAOc,iBAAiB,OAAQ7D,KAAK2D,YAC1C3D,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAKmE,aAC3CnE,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAK+D,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBhD,KAAK+C,SAET/C,KAAK+C,OAAOa,oBAAoB,OAAQ5D,KAAK2D,YAC7C3D,KAAK+C,OAAOa,oBAAoB,QAAS5D,KAAKmE,aAC9CnE,KAAK+C,OAAOa,oBAAoB,QAAS5D,KAAK+D,cAC9C/D,KAAK+C,OAAOa,oBAAoB,UAAW5D,KAAK8D,eAEhD9D,KAAK+C,OAAOpB,QACZ3B,KAAK+C,YAASC,EACdhD,KAAKiD,aAAc,EACrB,ECNW,MAIP2B,2BAAsD,CAC1DC,gBALwC,MAMxCC,cAJ+B,mDAFS,oCAY7B,MA+CXpF,WAAAA,CAAYK,QA9CKA,KAAAA,YAEAS,EAAAA,KAAAA,YAEAuE,EAAAA,KAAAA,UAAY,IAAIC,WAAWA,UAEpCC,KAAAA,eAEAC,EAAAA,KAAAA,gBAAyCC,WAAAA,sBAAsBC,cAC/DC,KAAAA,kBAA6CT,2BAE7CU,KAAAA,UAA6BC,WAAAA,gBAAgBC,aAGpCC,KAAAA,+BAAiC,IAAIpF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBqF,KAAAA,yBAA2B,IAAIrF,IAMxCsF,KAAAA,kBAAmB,EAAI3F,KAIvB4F,yBAAmB,EAAA5F,KAInB6F,mBAAqB,EAAC7F,KACtB8F,eAAiB,EAAC9F,KAKlB+F,kBAAoB,EAAC/F,KAGrBgG,gBAAkB,EAAChG,KACViG,SAAW,IAAIC,IAAqClG,KAoBrEmG,QAAWrD,MAEL9C,KAAKiF,WAAWvB,WAAaZ,MAGjC9C,KAAKoG,aAELpG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKqG,mBAAmBlB,WAAAA,sBAAsBmB,YAE9CtG,KAAKiF,UAAY,IAAIpC,mBAAmBC,KACxC9C,KAAKiF,UAAU1B,gBAAgBvD,KAAKuG,sBACpCvG,KAAKiF,UAAUxB,mBAAmBzD,KAAKwG,gBACvCxG,KAAKiF,UAAUzB,iBAAiBxD,KAAKyG,uBAGrCzG,KAAKiF,UAAUP,QACjB,EAAC1E,KAED0G,UAAY,KACV,GACE1G,KAAKD,OAAO4G,qBAAuB,GACnC3G,KAAK+F,mBAAqB/F,KAAKD,OAAO4G,qBAKtC,OAHA3G,KAAKQ,OAAOgE,KAAK,uCAEjBxE,KAAKoG,aAIP,GACEpG,KAAKkF,kBAAoBC,WAAAA,sBAAsBC,oBAC5BpC,IAAnBhD,KAAKiF,UAFP,CAMAjF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKiF,UAAUvB,UAExD1D,KAAKiF,UAAUhB,OAGfjE,KAAK+E,UAAUhD,QAGf/B,KAAKqF,kBAAoBT,2BACzB5E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EAGxB3F,KAAK+F,oBAGL/F,KAAKqG,mBAAmBlB,WAAAA,sBAAsBmB,YAC9C,IAAK,MAAMxF,WAAWd,KAAKiG,SAASW,SAC9B9F,QAAQM,aAAelB,WAAkBA,mBAAC0B,QAE9Cd,QAAQmB,yBAMVjC,KAAK+E,UAAU8B,SACb,UACyB7D,IAAnBhD,KAAKiF,WAGTjF,KAAKiF,UAAUP,OAAK,EAEG,IAAzB1E,KAAK+F,kBACL,YArCA,CAqCW,EAIfK,KAAAA,WAAa,KACPpG,KAAKkF,kBAAoBC,WAAAA,sBAAsBC,gBAEnDpF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKiF,WAAWhB,OAChBjE,KAAKiF,eAAYjC,EAGjBhD,KAAK+E,UAAUhD,QAGf/B,KAAKqF,kBAAoBT,2BACzB5E,KAAK6F,mBAAqB,EAC1B7F,KAAK8F,eAAiB,EACtB9F,KAAK2F,kBAAmB,EACxB3F,KAAK+F,kBAAoB,EAEzB/F,KAAKqG,mBAAmBlB,WAAqBA,sBAACC,eAC9CpF,KAAK8G,aAAavB,WAAAA,gBAAgBC,cACpC,EAACxF,KAED2B,MAAQ,KACN3B,KAAKoG,YACP,EAEAW,KAAAA,qBAAuB,IAAM/G,KAAKqF,kBAAiBrF,KACnDgH,mBAAqB,IAAMhH,KAAKkF,gBAChC+B,KAAAA,iCAAoCjG,UAClChB,KAAKyF,+BAA+BxE,IAAID,UAAShB,KACnDkH,oCAAuClG,UACrChB,KAAKyF,+BAA+BtE,OAAOH,UAE7CmG,KAAAA,aAAgBC,QACdpH,KAAK4F,oBAAsBwB,MAEvBpH,KAAKkF,kBAAoBC,WAAAA,sBAAsBkC,WACjDrH,KAAKsH,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuBvH,KAAKsF,UAAStF,KACpDwH,2BAA8BxG,UAC5BhB,KAAK0F,yBAAyBzE,IAAID,UAAShB,KAC7CyH,8BAAiCzG,UAC/BhB,KAAK0F,yBAAyBvE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAEpF0G,KAAAA,YAAc,CAAC9H,QAAiBC,cAC9B,MAAM8H,UAAY3H,KAAKgG,gBACvBhG,KAAKgG,iBAAmB,EAExB,MAAMlF,QAAU,IAAIrB,uBAClBkI,UACA/H,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKiG,SAAS2B,IAAID,UAAW7G,SAI3Bd,KAAKkF,kBAAoBC,WAAqBA,sBAACkC,WAC/CrH,KAAKsF,YAAcC,WAAeA,gBAACsC,YAEnC/G,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,KAAK8H,oBAGL9H,KAAK8F,eAAiBiC,KAAKC,KAAG,EAC/BhI,KAEOsH,gBAAmBF,QACzBpH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAK8G,aAAavB,WAAeA,gBAAC0C,aAElCjI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTsG,aAEJ,EAEQN,KAAAA,aAAgBoB,WACtB,MAAMzF,KAAOzC,KAAKsF,UAElBtF,KAAKsF,UAAY4C,SACjB,IAAK,MAAMlH,YAAgBhB,KAAC0F,yBAC1B,IACE1E,SAASkH,SAAUzF,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EACFvC,KAEOwG,eAAkB9E,UAaxB,GAZA1B,KAAK6F,mBAAqBkC,KAAKC,MAI3BhI,KAAK6F,mBAAqB7F,KAAK8F,gBAAkD,IAAhC9F,KAAKD,OAAOoI,mBAC/DnI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,ICrNmBY,UACd,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDoNJ0H,CAAoB1G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACqI,oBAAoB3G,SAClC,IAAK,aACH,OAAO1B,KAAKsI,wBAAwB5G,SACtC,IAAK,QACH,OAAO1B,KAAKuI,aAAa,CACvB7H,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GCxNsBA,UACX,IAApBA,QAAQZ,QDuNK0H,CAAiB9G,SAAU,CACpC,MAAMZ,QAAUd,KAAKiG,SAASwC,IAAI/G,QAAQZ,SAC1C,QAAgBkC,IAAZlC,QAEF,YADAd,KAAKQ,OAAOgE,KAAK,iDAAkD9C,SAIrE,GC3NJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KDqNAgI,CAA0BhH,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,OAAOgE,KAAK,qBAAsB9C,QAAQhB,KAAI,EACpDV,KAEOqI,oBAAuBM,cAE7B3I,KAAK+E,UAAU6D,OAAO,iBAIpB5I,KAAKkF,kBAAoBC,WAAAA,sBAAsBmB,YAC/CtG,KAAKkF,kBAAoBC,WAAAA,sBAAsBkC,YAE/CrH,KAAKqF,kBAAoB,IACpBrF,KAAKqF,kBACRwD,cAAeF,YAAYG,QAC3BC,uBAAwB/I,KAAKD,OAAOiJ,iBACpCC,uBAAwBN,YAAYK,kBAItChJ,KAAK+F,kBAAoB,OAEQ/C,IAA7BhD,KAAK4F,qBACP5F,KAAKqG,mBAAmBlB,WAAAA,sBAAsBkC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDhJ,KAAK+E,UAAU8B,SAAS,IAAM7G,KAAKmJ,aAAaD,cAAeA,aAAc,UAAS,EAGhFX,KAAAA,aAAgB9G,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,EAGK6G,KAAAA,wBAA0B,EAAGc,gBACnCpJ,KAAKQ,OAAOqB,MAAM,8BAA+BuH,OAGjDpJ,KAAK+E,UAAU6D,OAAO,sBAGlB5I,KAAK2F,iBACP3F,KAAK2F,kBAAmB,EAGV,iBAAVyD,QACFpJ,KAAK4F,yBAAsB5C,GAKjB,eAAVoG,QACFpJ,KAAKqG,mBAAmBlB,WAAqBA,sBAACkC,WAE9CrH,KAAKqJ,yBAGPrJ,KAAK8G,aAAavB,WAAeA,gBAAC6D,OACpC,EAEQC,KAAAA,sBAAwB,KAC9B,IAAK,MAAMvI,WAAed,KAACiG,SAASW,SAE9B9F,QAAQM,aAAelB,WAAAA,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKiG,SAAS9E,OAAOL,QAAQnB,GAKhC,EACFK,KAUOuG,qBAAuB,KAC7BvG,KAAKQ,OAAOqB,MAAM,qBAElB,MAAMyH,aAA6B,CACjC5I,KAAM,QACNI,QAAS,EACTgI,QAAS,GAAG9I,KAAKqF,kBAAkBR,mBAAmB7E,KAAKqF,kBAAkBP,gBAC7EkE,iBAAkBhJ,KAAKD,OAAOiJ,iBAC9BO,uBAAwBvJ,KAAKD,OAAOwJ,wBAItCvJ,KAAK+E,UAAU8B,SACb,KACE,MAAM2C,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,iCAAmC1B,KAAKD,OAAO0J,cAAgB,KAG1EzJ,KAAKF,YAAY0J,cAEjBxJ,KAAKuI,aAAa,CAChB7H,KAAM8I,aAAa/H,MACnBC,QAAS,GAAG8H,aAAa9H,wBAI3B1B,KAAK0G,WACP,EAC4B,IAA5B1G,KAAKD,OAAO0J,cACZ,iBAGFzJ,KAAKF,YAAYwJ,cAEjBtJ,KAAK+E,UAAU8B,SACb,KACE,MAAM2C,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,sCAAwC1B,KAAKD,OAAO0J,cAAgB,KAG/EzJ,KAAKF,YAAY0J,cAEjBxJ,KAAKuI,aAAa,CAChB7H,KAAM8I,aAAa/H,MACnBC,QAAS,GAAG8H,aAAa9H,wBAI3B1B,KAAK0G,WACP,EAC4B,IAA5B1G,KAAKD,OAAO0J,cACZ,2BAG+BzG,IAA7BhD,KAAK4F,qBACP5F,KAAKsH,gBAAgBtH,KAAK4F,oBAC3B,EASKa,KAAAA,sBAAwB,CAACvC,OAAgBzC,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBqC,aAEzBlB,IAAVvB,OACFzB,KAAKuI,aAAa,CAChB7H,KAAM,UACNgB,QAASwC,SAITlE,KAAKsF,YAAcC,WAAAA,gBAAgBC,aAGrC,OAFAxF,KAAK4F,yBAAsB5C,OAC3BhD,KAAKoG,aAIPpG,KAAK0G,WACP,EAEQyC,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBhI,KAAK6F,mBACvC,GAAI6D,qBAAuBR,aAQzB,OAPAlJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BgI,oBAAsB,OAGrD1J,KAAC0G,YAGd,MAAMiD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD1J,KAAK+E,UAAU8B,SAAS,IAAM7G,KAAKmJ,aAAaD,cAAeS,YAAa,UAAS,EACtF3J,KAEO8H,kBAAoB,KAC1B9H,KAAK+E,UAAU8B,SACb,KACE7G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAK8H,mBACP,EACgC,IAAhC9H,KAAKD,OAAOoI,kBACZ,YAEJ,EA/dEnI,KAAKD,OAAS,CACZoI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACf7G,SAAUkH,WAAcA,eAACC,KACzBpD,sBAAuB,KACpB5G,QAGLC,KAAKQ,OAAS,IAAIkC,WAAMA,OAAC1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,SAC9D"}
{"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 this.socket.addEventListener('close', this.handleClosed)\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 Scheduler,\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\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 = new Scheduler()\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 }\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 '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 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('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(() => this.timeoutCheck(timeoutMills), timeoutMills, 'TIMEOUT')\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('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 '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 '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(() => this.timeoutCheck(timeoutMills), nextTimeout, 'TIMEOUT')\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 '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","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","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"],"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,WADAD,KAAAA,SACAC,EAAAA,KAAAA,eAVXC,EAAAA,KAAAA,YAAgCC,EAEhCC,KAAAA,aAAc,EAEdC,KAAAA,kBAAyCF,OACzCG,mBAA0DH,EAASjD,KACnEqD,qBAA2EJ,EAASjD,KA8B5FF,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAEA8B,KAAAA,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QACtB,EAAChB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAAChB,KAED0D,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QACzB,EAEA2C,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,eAC7C/D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAE3ChE,KAAKmD,iBACP,EAEQa,KAAAA,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAAK,EACtCnE,KAEOoE,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAEQW,KAAAA,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,UAAUb,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,GApGgBzB,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,ECrDW,MAIP2B,2BAAsD,CAC1DC,gBALwC,MAMxCC,cAJqB,mDAFmB,0CA2DxCrF,WAAAA,CAAYK,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENgF,UAAY,IAAIC,WAAAA,UAAWjF,KAEpCkF,eAAS,EAAAlF,KAETmF,gBAAyCC,WAAqBA,sBAACC,cAAarF,KAC5EsF,kBAA6CT,2BAA0B7E,KAEvEuF,UAA6BC,WAAeA,gBAACC,aAAYzF,KAGhD0F,+BAAiC,IAAIrF,IAA0CL,KAC/EO,eAAiB,IAAIF,IAA0BL,KAC/C2F,yBAA2B,IAAItF,IAAoCL,KAM5E4F,kBAAmB,EAAI5F,KAIvB6F,yBAAmB,EAAA7F,KAInB8F,mBAAqB,EAAC9F,KACtB+F,eAAiB,EAAC/F,KAKlBgG,kBAAoB,EAAChG,KAGrBiG,gBAAkB,EAACjG,KACVkG,SAAW,IAAIC,IAAqCnG,KAqBrEoG,QAAWtD,MAEL9C,KAAKkF,WAAWvB,WAAab,MAGjC9C,KAAKqG,aAELrG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKsG,mBAAmBlB,WAAAA,sBAAsBmB,YAG9CvG,KAAKkF,UAAYlF,KAAKD,OAAOyG,iBAAiB1D,KAC9C9C,KAAKkF,UAAU1B,gBAAgBxD,KAAKyG,sBACpCzG,KAAKkF,UAAUxB,mBAAmB1D,KAAK0G,gBACvC1G,KAAKkF,UAAUzB,iBAAiBzD,KAAK2G,uBAGrC3G,KAAKkF,UAAUP,QAAK,EACrB3E,KAED4G,UAAY,KACV,GACE5G,KAAKD,OAAO8G,sBAAwB,GACpC7G,KAAKgG,mBAAqBhG,KAAKD,OAAO8G,qBAKtC,OAHA7G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKqG,aAIP,GACErG,KAAKmF,kBAAoBC,WAAAA,sBAAsBC,oBAC5BpC,IAAnBjD,KAAKkF,UAFP,CAMAlF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKkF,UAAUvB,UAExD3D,KAAKkF,UAAUhB,OAGflE,KAAKgF,UAAUjD,QAGf/B,KAAKsF,kBAAoBT,2BACzB7E,KAAK8F,mBAAqB,EAC1B9F,KAAK+F,eAAiB,EACtB/F,KAAK4F,kBAAmB,EAGxB5F,KAAKgG,oBAGLhG,KAAKsG,mBAAmBlB,WAAAA,sBAAsBmB,YAC9C,IAAK,MAAMzF,WAAWd,KAAKkG,SAASY,SAC9BhG,QAAQM,aAAelB,WAAkBA,mBAAC0B,QAE9Cd,QAAQmB,yBAMVjC,KAAKgF,UAAU+B,SACb,UACyB9D,IAAnBjD,KAAKkF,WAGTlF,KAAKkF,UAAUP,OAAK,EAEG,IAAzB3E,KAAKgG,kBACL,YArCA,CAqCW,EAIfK,KAAAA,WAAa,KACPrG,KAAKmF,kBAAoBC,WAAAA,sBAAsBC,gBAEnDrF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKkF,WAAWhB,OAChBlE,KAAKkF,eAAYjC,EAGjBjD,KAAKgF,UAAUjD,QAGf/B,KAAKsF,kBAAoBT,2BACzB7E,KAAK8F,mBAAqB,EAC1B9F,KAAK+F,eAAiB,EACtB/F,KAAK4F,kBAAmB,EACxB5F,KAAKgG,kBAAoB,EAEzBhG,KAAKsG,mBAAmBlB,WAAqBA,sBAACC,eAC9CrF,KAAKgH,aAAaxB,WAAAA,gBAAgBC,cACpC,EAACzF,KAED2B,MAAQ,KACN3B,KAAKqG,YACP,EAEAY,KAAAA,qBAAuB,IAAMjH,KAAKsF,kBAAiBtF,KACnDkH,mBAAqB,IAAMlH,KAAKmF,gBAChCgC,KAAAA,iCAAoCnG,UAClChB,KAAK0F,+BAA+BzE,IAAID,UAAShB,KACnDoH,oCAAuCpG,UACrChB,KAAK0F,+BAA+BvE,OAAOH,UAE7CqG,KAAAA,aAAgBC,QACdtH,KAAK6F,oBAAsByB,MAEvBtH,KAAKmF,kBAAoBC,WAAAA,sBAAsBmC,WACjDvH,KAAKwH,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuBzH,KAAKuF,UAASvF,KACpD0H,2BAA8B1G,UAC5BhB,KAAK2F,yBAAyB1E,IAAID,UACpC2G,KAAAA,8BAAiC3G,UAC/BhB,KAAK2F,yBAAyBxE,OAAOH,UAAShB,KAEhDuB,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAC9EQ,KAAAA,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7F4H,YAAc,CAAChI,QAAiBC,cAC9B,MAAMgI,UAAY7H,KAAKiG,gBACvBjG,KAAKiG,iBAAmB,EAExB,MAAMnF,QAAU,IAAIrB,uBAClBoI,UACAjI,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKkG,SAAS4B,IAAID,UAAW/G,SAI3Bd,KAAKmF,kBAAoBC,WAAqBA,sBAACmC,WAC/CvH,KAAKuF,YAAcC,WAAeA,gBAACuC,YAEnCjH,QAAQkB,UAGHlB,SAGDwF,KAAAA,mBAAsB9D,YAC5B,MAAMC,KAAOzC,KAAKmF,gBAClB,GAAI1C,OAASD,UAAb,CAEAxC,KAAKmF,gBAAkB3C,UACvB,IAAK,MAAMxB,YAAYhB,KAAK0F,+BAC1B1E,SAASwB,UAAWC,KAJE,CAKvB,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKkF,WAAWpF,YAAY4B,SAE5B1B,KAAKgI,oBAGLhI,KAAK+F,eAAiBkC,KAAKC,KAC7B,EAAClI,KAEOwH,gBAAmBF,QACzBtH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAKgH,aAAaxB,WAAAA,gBAAgB2C,aAElCnI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTwG,aACD,EAGKN,KAAAA,aAAgBoB,WACtB,MAAM3F,KAAOzC,KAAKuF,UAElBvF,KAAKuF,UAAY6C,SACjB,IAAK,MAAMpH,YAAYhB,KAAK2F,yBAC1B,IACE3E,SAASoH,SAAU3F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EAGKmE,KAAAA,eAAkBhF,UAaxB,GAZA1B,KAAK8F,mBAAqBmC,KAAKC,MAI3BlI,KAAK8F,mBAAqB9F,KAAK+F,gBAAkD,IAAhC/F,KAAKD,OAAOsI,mBAC/DrI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,ICtNfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDoNJ4H,CAAoB5G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAOV,KAAKuI,oBAAoB7G,SAClC,IAAK,aACH,OAAW1B,KAACwI,wBAAwB9G,SACtC,IAAK,QACH,OAAW1B,KAACyI,aAAa,CACvB/H,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEKgH,GCxNkBhH,UACX,IAApBA,QAAQZ,QDuNK4H,CAAiBhH,SAAU,CACpC,MAAMZ,QAAUd,KAAKkG,SAASyC,IAAIjH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GC3NJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KDqNAkI,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,EAACV,KAEOuI,oBAAuBM,cAE7B7I,KAAKgF,UAAU8D,OAAO,iBAIpB9I,KAAKmF,kBAAoBC,WAAqBA,sBAACmB,YAC/CvG,KAAKmF,kBAAoBC,WAAqBA,sBAACmC,YAE/CvH,KAAKsF,kBAAoB,IACpBtF,KAAKsF,kBACRyD,cAAeF,YAAYG,QAC3BC,uBAAwBjJ,KAAKD,OAAOmJ,iBACpCC,uBAAwBN,YAAYK,kBAItClJ,KAAKgG,kBAAoB,OAEQ/C,IAA7BjD,KAAK6F,qBACP7F,KAAKsG,mBAAmBlB,WAAqBA,sBAACmC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDlJ,KAAKgF,UAAU+B,SAAS,IAAM/G,KAAKqJ,aAAaD,cAAeA,aAAc,UAAS,EACvFpJ,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,EACFzB,KAEOwI,wBAA0B,EAAGc,gBACnCtJ,KAAKQ,OAAOqB,MAAM,8BAA+ByH,OAGjDtJ,KAAKgF,UAAU8D,OAAO,sBAGlB9I,KAAK4F,iBACP5F,KAAK4F,kBAAmB,EAGV,iBAAV0D,QACFtJ,KAAK6F,yBAAsB5C,GAKjB,eAAVqG,QACFtJ,KAAKsG,mBAAmBlB,WAAqBA,sBAACmC,WAE9CvH,KAAKuJ,yBAGPvJ,KAAKgH,aAAaxB,WAAeA,gBAAC8D,OACpC,EAEQC,KAAAA,sBAAwB,KAC9B,IAAK,MAAMzI,WAAed,KAACkG,SAASY,SAE9BhG,QAAQM,aAAelB,WAAAA,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKkG,SAAS/E,OAAOL,QAAQnB,GAKhC,EACFK,KAUOyG,qBAAuB,KAC7BzG,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM2H,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTkI,QAAS,GAAGhJ,KAAKsF,kBAAkBR,mBAAmB9E,KAAKsF,kBAAkBP,gBAC7EmE,iBAAkBlJ,KAAKD,OAAOmJ,iBAC9BO,uBAAwBzJ,KAAKD,OAAO0J,wBAItCzJ,KAAKgF,UAAU+B,SACb,KACE,MAAM2C,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,KAAK4G,WACP,EAC4B,IAA5B5G,KAAKD,OAAO4J,cACZ,iBAGF3J,KAAKF,YAAY0J,cAEjBxJ,KAAKgF,UAAU+B,SACb,KACE,MAAM2C,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,KAAK4G,WACP,EAC4B,IAA5B5G,KAAKD,OAAO4J,cACZ,2BAG+B1G,IAA7BjD,KAAK6F,qBACP7F,KAAKwH,gBAAgBxH,KAAK6F,oBAC3B,EACF7F,KAQO2G,sBAAwB,CAACxC,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAKyI,aAAa,CAChB/H,KAAM,UACNgB,QAASyC,SAITnE,KAAKuF,YAAcC,WAAeA,gBAACC,aAGrC,OAFAzF,KAAK6F,yBAAsB5C,OAC3BjD,KAAKqG,aAIPrG,KAAK4G,WAAS,EAGRyC,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBlI,KAAK8F,mBACvC,GAAI8D,qBAAuBR,aAQzB,OAPApJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BkI,oBAAsB,OAGrD5J,KAAC4G,YAGd,MAAMiD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD5J,KAAKgF,UAAU+B,SAAS,IAAM/G,KAAKqJ,aAAaD,cAAeS,YAAa,UAC9E,EAEQ7B,KAAAA,kBAAoB,KAC1BhI,KAAKgF,UAAU+B,SACb,KACE/G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKgI,mBAAiB,EAEQ,IAAhChI,KAAKD,OAAOsI,kBACZ,YAAW,EA/dbrI,KAAKD,OAAS,CACZsI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACf/G,SAAUoH,WAAAA,eAAeC,KACzBpD,sBAAuB,EACvBL,iBAAmB1D,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,WAAAA,OAAO1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,SAC9D"}

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

import{DXLinkChannelState,Logger,Scheduler,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel}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 WebSocketConnector{constructor(url){this.url=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.socket.addEventListener("close",this.handleClosed),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}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url),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.3.0"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=new Scheduler,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=new WebSocketConnector(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,"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.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("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,"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("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,"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,"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,"TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:DXLinkLogLevel.WARN,maxReconnectAttempts:-1,...config},this.logger=new Logger(this.constructor.name,this.config.logLevel)}}export{DXLINK_WS_PROTOCOL_VERSION,DXLinkWebSocketClient};
import{DXLinkChannelState,Logger,Scheduler,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel}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.socket.addEventListener("close",this.handleClosed),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.4.0"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=new Scheduler,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,"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.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("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,"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("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,"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,"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,"TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"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)}}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 Message } 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: Message) => 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 Message } from './messages'\n\ntype CloseListener = (reason: string, error: boolean) => void\n\n/**\n * Connector for the WebSocket connection.\n * @internal\n */\nexport class WebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: CloseListener | undefined = undefined\n private messageListener: ((message: Message) => void) | undefined = undefined\n\n constructor(private readonly url: string) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url)\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: Message) => {\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: CloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: Message) => 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 this.socket.addEventListener('close', this.handleClosed)\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 Scheduler,\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 { WebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type Message,\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\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 = new Scheduler()\n\n private connector: WebSocketConnector | 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 ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\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 this.connector = new WebSocketConnector(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 '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 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: Message): 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: Message): 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('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(() => this.timeoutCheck(timeoutMills), timeoutMills, 'TIMEOUT')\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('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 '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 '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(() => this.timeoutCheck(timeoutMills), nextTimeout, 'TIMEOUT')\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 '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 Message =\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 = (message: Message): 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: Message): 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","WebSocketConnector","url","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","Scheduler","connector","connectionState","DXLinkConnectionState","NOT_CONNECTED","connectionDetails","authState","DXLinkAuthState","UNAUTHORIZED","connectionStateChangeListeners","authStateChangeListeners","isFirstAuthState","lastSettedAuthToken","lastReceivedMillis","lastSentMillis","reconnectAttempts","globalChannelId","channels","Map","connect","disconnect","setConnectionState","CONNECTING","processTransportOpen","processMessage","processTransportClose","reconnect","maxReconnectAttempts","values","schedule","setAuthState","getConnectionDetails","getConnectionState","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"],"mappings":"gIAmBaA,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,QC7BWC,mBASXnD,WAAAA,CAA6BoD,KAAW9C,KAAX8C,SAAA,EAAA9C,KARrB+C,YAAgCC,EAAShD,KAEzCiD,aAAc,EAAKjD,KAEnBkD,kBAAyCF,EAAShD,KAClDmD,mBAA2CH,EAC3CI,KAAAA,qBAA4DJ,EA2BpElD,KAAAA,YAAe4B,eACOsB,IAAhBhD,KAAK+C,QAAyB/C,KAAKiD,aAIvCjD,KAAK+C,OAAOtC,KAAK4C,KAAKC,UAAU5B,SAAQ,EAG1C6B,KAAAA,gBAAmBvC,WACjBhB,KAAKkD,aAAelC,QACtB,EAEAwC,KAAAA,iBAAoBxC,WAClBhB,KAAKmD,cAAgBnC,QAAAA,EAGvByC,KAAAA,mBAAsBzC,WACpBhB,KAAKoD,gBAAkBpC,QACzB,EAAChB,KAED0D,OAAS,IAAM1D,KAAK8C,SAEZa,WAAa,UACCX,IAAhBhD,KAAK+C,SAET/C,KAAKiD,aAAc,EAEnBjD,KAAK+C,OAAOa,oBAAoB,OAAQ5D,KAAK2D,YAE7C3D,KAAK+C,OAAOc,iBAAiB,UAAW7D,KAAK8D,eAC7C9D,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAK+D,cAE3C/D,KAAKkD,iBACP,EAAClD,KAEO+D,aAAgBC,UACFhB,IAAhBhD,KAAK+C,SAET/C,KAAKiE,OAELjE,KAAKmD,gBAAgBa,GAAGE,QAAQ,GAClC,EAEQC,KAAAA,YAAeC,WACDpB,IAAhBhD,KAAK+C,SAET/C,KAAKiE,OAELjE,KAAKmD,gBAAgB,qBAAqB,GAC5C,EAEQW,KAAAA,cAAiBE,KACvB,IACE,MAAMtC,QAAU2B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ5C,QACT,MAAU,IAAAb,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,QAA6BkC,IAAzBhD,KAAKoD,gBACP,OAAOmB,QAAQC,KAAK,2BAGtBxE,KAAKoD,gBAAgB1B,QACtB,CAAC,MAAOD,OACP8C,QAAQ9C,MAAMA,iBAAiBZ,MAAQY,MAAQ,IAAIZ,MAAM,iBAAmB4D,OAAOhD,QACpF,GAlG0BzB,KAAG8C,IAAHA,GAAc,CAE3C4B,KAAAA,QACsB1B,IAAhBhD,KAAK+C,SAET/C,KAAK+C,OAAS,IAAI4B,UAAU3E,KAAK8C,KAEjC9C,KAAK+C,OAAOc,iBAAiB,OAAQ7D,KAAK2D,YAC1C3D,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAKmE,aAC3CnE,KAAK+C,OAAOc,iBAAiB,QAAS7D,KAAK+D,cAC7C,CAEAE,IAAAA,QACsBjB,IAAhBhD,KAAK+C,SAET/C,KAAK+C,OAAOa,oBAAoB,OAAQ5D,KAAK2D,YAC7C3D,KAAK+C,OAAOa,oBAAoB,QAAS5D,KAAKmE,aAC9CnE,KAAK+C,OAAOa,oBAAoB,QAAS5D,KAAK+D,cAC9C/D,KAAK+C,OAAOa,oBAAoB,UAAW5D,KAAK8D,eAEhD9D,KAAK+C,OAAOpB,QACZ3B,KAAK+C,YAASC,EACdhD,KAAKiD,aAAc,EACrB,ECNW,MAAA2B,2BAA6B,MAIpCC,2BAAsD,CAC1DC,gBALwC,MAMxCC,cAJ+B,gBAUpB,MAAAC,sBA+CXtF,WAAAA,CAAYK,QA9CKA,KAAAA,YAEAS,EAAAA,KAAAA,YAEAyE,EAAAA,KAAAA,UAAY,IAAIC,UAEzBC,KAAAA,eAEAC,EAAAA,KAAAA,gBAAyCC,sBAAsBC,cAC/DC,KAAAA,kBAA6CV,2BAE7CW,KAAAA,UAA6BC,gBAAgBC,aAGpCC,KAAAA,+BAAiC,IAAItF,IACrCE,KAAAA,eAAiB,IAAIF,IACrBuF,KAAAA,yBAA2B,IAAIvF,IAMxCwF,KAAAA,kBAAmB,EAAI7F,KAIvB8F,yBAAmB,EAAA9F,KAInB+F,mBAAqB,EAAC/F,KACtBgG,eAAiB,EAAChG,KAKlBiG,kBAAoB,EAACjG,KAGrBkG,gBAAkB,EAAClG,KACVmG,SAAW,IAAIC,IAAqCpG,KAoBrEqG,QAAWvD,MAEL9C,KAAKmF,WAAWzB,WAAaZ,MAGjC9C,KAAKsG,aAELtG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKuG,mBAAmBlB,sBAAsBmB,YAE9CxG,KAAKmF,UAAY,IAAItC,mBAAmBC,KACxC9C,KAAKmF,UAAU5B,gBAAgBvD,KAAKyG,sBACpCzG,KAAKmF,UAAU1B,mBAAmBzD,KAAK0G,gBACvC1G,KAAKmF,UAAU3B,iBAAiBxD,KAAK2G,uBAGrC3G,KAAKmF,UAAUT,QACjB,EAAC1E,KAED4G,UAAY,KACV,GACE5G,KAAKD,OAAO8G,qBAAuB,GACnC7G,KAAKiG,mBAAqBjG,KAAKD,OAAO8G,qBAKtC,OAHA7G,KAAKQ,OAAOgE,KAAK,uCAEjBxE,KAAKsG,aAIP,GACEtG,KAAKoF,kBAAoBC,sBAAsBC,oBAC5BtC,IAAnBhD,KAAKmF,UAFP,CAMAnF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKmF,UAAUzB,UAExD1D,KAAKmF,UAAUlB,OAGfjE,KAAKiF,UAAUlD,QAGf/B,KAAKuF,kBAAoBV,2BACzB7E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EAGxB7F,KAAKiG,oBAGLjG,KAAKuG,mBAAmBlB,sBAAsBmB,YAC9C,IAAK,MAAM1F,WAAWd,KAAKmG,SAASW,SAC9BhG,QAAQM,aAAelB,mBAAmB0B,QAE9Cd,QAAQmB,yBAMVjC,KAAKiF,UAAU8B,SACb,UACyB/D,IAAnBhD,KAAKmF,WAGTnF,KAAKmF,UAAUT,OAAK,EAEG,IAAzB1E,KAAKiG,kBACL,YArCA,CAqCW,EAIfK,KAAAA,WAAa,KACPtG,KAAKoF,kBAAoBC,sBAAsBC,gBAEnDtF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKmF,WAAWlB,OAChBjE,KAAKmF,eAAYnC,EAGjBhD,KAAKiF,UAAUlD,QAGf/B,KAAKuF,kBAAoBV,2BACzB7E,KAAK+F,mBAAqB,EAC1B/F,KAAKgG,eAAiB,EACtBhG,KAAK6F,kBAAmB,EACxB7F,KAAKiG,kBAAoB,EAEzBjG,KAAKuG,mBAAmBlB,sBAAsBC,eAC9CtF,KAAKgH,aAAavB,gBAAgBC,cACpC,EAAC1F,KAED2B,MAAQ,KACN3B,KAAKsG,YACP,EAEAW,KAAAA,qBAAuB,IAAMjH,KAAKuF,kBAAiBvF,KACnDkH,mBAAqB,IAAMlH,KAAKoF,gBAChC+B,KAAAA,iCAAoCnG,UAClChB,KAAK2F,+BAA+B1E,IAAID,UAAShB,KACnDoH,oCAAuCpG,UACrChB,KAAK2F,+BAA+BxE,OAAOH,UAE7CqG,KAAAA,aAAgBC,QACdtH,KAAK8F,oBAAsBwB,MAEvBtH,KAAKoF,kBAAoBC,sBAAsBkC,WACjDvH,KAAKwH,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuBzH,KAAKwF,UAASxF,KACpD0H,2BAA8B1G,UAC5BhB,KAAK4F,yBAAyB3E,IAAID,UAAShB,KAC7C2H,8BAAiC3G,UAC/BhB,KAAK4F,yBAAyBzE,OAAOH,UAEvCO,KAAAA,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAAShB,KACvFwB,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAEpF4G,KAAAA,YAAc,CAAChI,QAAiBC,cAC9B,MAAMgI,UAAY7H,KAAKkG,gBACvBlG,KAAKkG,iBAAmB,EAExB,MAAMpF,QAAU,IAAIrB,uBAClBoI,UACAjI,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKmG,SAAS2B,IAAID,UAAW/G,SAI3Bd,KAAKoF,kBAAoBC,sBAAsBkC,WAC/CvH,KAAKwF,YAAcC,gBAAgBsC,YAEnCjH,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,KAAKgI,oBAGLhI,KAAKgG,eAAiBiC,KAAKC,KAAG,EAC/BlI,KAEOwH,gBAAmBF,QACzBtH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAKgH,aAAavB,gBAAgB0C,aAElCnI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACTwG,aAEJ,EAEQN,KAAAA,aAAgBoB,WACtB,MAAM3F,KAAOzC,KAAKwF,UAElBxF,KAAKwF,UAAY4C,SACjB,IAAK,MAAMpH,YAAgBhB,KAAC4F,yBAC1B,IACE5E,SAASoH,SAAU3F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EACFvC,KAEO0G,eAAkBhF,UAaxB,GAZA1B,KAAK+F,mBAAqBkC,KAAKC,MAI3BlI,KAAK+F,mBAAqB/F,KAAKgG,gBAAkD,IAAhChG,KAAKD,OAAOsI,mBAC/DrI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,ICrNmBY,UACd,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDoNJ4H,CAAoB5G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAWV,KAACuI,oBAAoB7G,SAClC,IAAK,aACH,OAAO1B,KAAKwI,wBAAwB9G,SACtC,IAAK,QACH,OAAO1B,KAAKyI,aAAa,CACvB/H,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEC,GCxNsBA,UACX,IAApBA,QAAQZ,QDuNK4H,CAAiBhH,SAAU,CACpC,MAAMZ,QAAUd,KAAKmG,SAASwC,IAAIjH,QAAQZ,SAC1C,QAAgBkC,IAAZlC,QAEF,YADAd,KAAKQ,OAAOgE,KAAK,iDAAkD9C,SAIrE,GC3NJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KDqNAkI,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,OAAOgE,KAAK,qBAAsB9C,QAAQhB,KAAI,EACpDV,KAEOuI,oBAAuBM,cAE7B7I,KAAKiF,UAAU6D,OAAO,iBAIpB9I,KAAKoF,kBAAoBC,sBAAsBmB,YAC/CxG,KAAKoF,kBAAoBC,sBAAsBkC,YAE/CvH,KAAKuF,kBAAoB,IACpBvF,KAAKuF,kBACRwD,cAAeF,YAAYG,QAC3BC,uBAAwBjJ,KAAKD,OAAOmJ,iBACpCC,uBAAwBN,YAAYK,kBAItClJ,KAAKiG,kBAAoB,OAEQjD,IAA7BhD,KAAK8F,qBACP9F,KAAKuG,mBAAmBlB,sBAAsBkC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDlJ,KAAKiF,UAAU8B,SAAS,IAAM/G,KAAKqJ,aAAaD,cAAeA,aAAc,UAAS,EAGhFX,KAAAA,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,KAAKiF,UAAU6D,OAAO,sBAGlB9I,KAAK6F,iBACP7F,KAAK6F,kBAAmB,EAGV,iBAAVyD,QACFtJ,KAAK8F,yBAAsB9C,GAKjB,eAAVsG,QACFtJ,KAAKuG,mBAAmBlB,sBAAsBkC,WAE9CvH,KAAKuJ,yBAGPvJ,KAAKgH,aAAavB,gBAAgB6D,OACpC,EAEQC,KAAAA,sBAAwB,KAC9B,IAAK,MAAMzI,WAAed,KAACmG,SAASW,SAE9BhG,QAAQM,aAAelB,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKmG,SAAShF,OAAOL,QAAQnB,GAKhC,EACFK,KAUOyG,qBAAuB,KAC7BzG,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM2H,aAA6B,CACjC9I,KAAM,QACNI,QAAS,EACTkI,QAAS,GAAGhJ,KAAKuF,kBAAkBT,mBAAmB9E,KAAKuF,kBAAkBR,gBAC7EmE,iBAAkBlJ,KAAKD,OAAOmJ,iBAC9BO,uBAAwBzJ,KAAKD,OAAO0J,wBAItCzJ,KAAKiF,UAAU8B,SACb,KACE,MAAM2C,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,KAAK4G,WACP,EAC4B,IAA5B5G,KAAKD,OAAO4J,cACZ,iBAGF3J,KAAKF,YAAY0J,cAEjBxJ,KAAKiF,UAAU8B,SACb,KACE,MAAM2C,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,KAAK4G,WACP,EAC4B,IAA5B5G,KAAKD,OAAO4J,cACZ,2BAG+B3G,IAA7BhD,KAAK8F,qBACP9F,KAAKwH,gBAAgBxH,KAAK8F,oBAC3B,EASKa,KAAAA,sBAAwB,CAACzC,OAAgBzC,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBqC,aAEzBlB,IAAVvB,OACFzB,KAAKyI,aAAa,CAChB/H,KAAM,UACNgB,QAASwC,SAITlE,KAAKwF,YAAcC,gBAAgBC,aAGrC,OAFA1F,KAAK8F,yBAAsB9C,OAC3BhD,KAAKsG,aAIPtG,KAAK4G,WACP,EAEQyC,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBlI,KAAK+F,mBACvC,GAAI6D,qBAAuBR,aAQzB,OAPApJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BkI,oBAAsB,OAGrD5J,KAAC4G,YAGd,MAAMiD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD5J,KAAKiF,UAAU8B,SAAS,IAAM/G,KAAKqJ,aAAaD,cAAeS,YAAa,UAAS,EACtF7J,KAEOgI,kBAAoB,KAC1BhI,KAAKiF,UAAU8B,SACb,KACE/G,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKgI,mBACP,EACgC,IAAhChI,KAAKD,OAAOsI,kBACZ,YAEJ,EA/dErI,KAAKD,OAAS,CACZsI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACf/G,SAAUoH,eAAeC,KACzBpD,sBAAuB,KACpB9G,QAGLC,KAAKQ,OAAS,IAAIkC,OAAO1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,SAC9D"}
{"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 this.socket.addEventListener('close', this.handleClosed)\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 Scheduler,\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\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 = new Scheduler()\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 }\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 '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 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('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(() => this.timeoutCheck(timeoutMills), timeoutMills, 'TIMEOUT')\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('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 '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 '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(() => this.timeoutCheck(timeoutMills), nextTimeout, 'TIMEOUT')\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 '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","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","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"],"mappings":"gIAmBaA,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,WADAD,KAAAA,SACAC,EAAAA,KAAAA,eAVXC,EAAAA,KAAAA,YAAgCC,EAEhCC,KAAAA,aAAc,EAEdC,KAAAA,kBAAyCF,OACzCG,mBAA0DH,EAASjD,KACnEqD,qBAA2EJ,EAASjD,KA8B5FF,YAAe4B,eACOuB,IAAhBjD,KAAKgD,QAAyBhD,KAAKkD,aAIvClD,KAAKgD,OAAOvC,KAAK6C,KAAKC,UAAU7B,SAClC,EAEA8B,KAAAA,gBAAmBxC,WACjBhB,KAAKmD,aAAenC,QACtB,EAAChB,KAEDyD,iBAAoBzC,WAClBhB,KAAKoD,cAAgBpC,QACvB,EAAChB,KAED0D,mBAAsB1C,WACpBhB,KAAKqD,gBAAkBrC,QACzB,EAEA2C,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,eAC7C/D,KAAKgD,OAAOc,iBAAiB,QAAS9D,KAAKgE,cAE3ChE,KAAKmD,iBACP,EAEQa,KAAAA,aAAgBC,UACFhB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgBa,GAAGE,QAAQ,GAAK,EACtCnE,KAEOoE,YAAeC,WACDpB,IAAhBjD,KAAKgD,SAEThD,KAAKkE,OAELlE,KAAKoD,gBAAgB,qBAAqB,GAC5C,EAEQW,KAAAA,cAAiBE,KACvB,IACE,MAAMvC,QAAU4B,KAAKgB,MAAML,GAAGM,MAC9B,GAAuB,iBAAZ7C,QACT,UAAUb,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,GApGgBzB,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,ECrDW,MAAA2B,2BAA6B,MAIpCC,2BAAsD,CAC1DC,gBALwC,MAMxCC,cAJqB,sBAUVC,sBA+CXvF,WAAAA,CAAYK,QAA6CC,KA9CxCD,YAAM,EAAAC,KAENQ,YAAM,EAAAR,KAENkF,UAAY,IAAIC,UAAWnF,KAEpCoF,eAAS,EAAApF,KAETqF,gBAAyCC,sBAAsBC,cAAavF,KAC5EwF,kBAA6CV,2BAA0B9E,KAEvEyF,UAA6BC,gBAAgBC,aAAY3F,KAGhD4F,+BAAiC,IAAIvF,IAA0CL,KAC/EO,eAAiB,IAAIF,IAA0BL,KAC/C6F,yBAA2B,IAAIxF,IAAoCL,KAM5E8F,kBAAmB,EAAI9F,KAIvB+F,yBAAmB,EAAA/F,KAInBgG,mBAAqB,EAAChG,KACtBiG,eAAiB,EAACjG,KAKlBkG,kBAAoB,EAAClG,KAGrBmG,gBAAkB,EAACnG,KACVoG,SAAW,IAAIC,IAAqCrG,KAqBrEsG,QAAWxD,MAEL9C,KAAKoF,WAAWzB,WAAab,MAGjC9C,KAAKuG,aAELvG,KAAKQ,OAAOqB,MAAM,gBAAiBiB,KAGnC9C,KAAKwG,mBAAmBlB,sBAAsBmB,YAG9CzG,KAAKoF,UAAYpF,KAAKD,OAAO2G,iBAAiB5D,KAC9C9C,KAAKoF,UAAU5B,gBAAgBxD,KAAK2G,sBACpC3G,KAAKoF,UAAU1B,mBAAmB1D,KAAK4G,gBACvC5G,KAAKoF,UAAU3B,iBAAiBzD,KAAK6G,uBAGrC7G,KAAKoF,UAAUT,QAAK,EACrB3E,KAED8G,UAAY,KACV,GACE9G,KAAKD,OAAOgH,sBAAwB,GACpC/G,KAAKkG,mBAAqBlG,KAAKD,OAAOgH,qBAKtC,OAHA/G,KAAKQ,OAAOiE,KAAK,uCAEjBzE,KAAKuG,aAIP,GACEvG,KAAKqF,kBAAoBC,sBAAsBC,oBAC5BtC,IAAnBjD,KAAKoF,UAFP,CAMApF,KAAKQ,OAAOqB,MAAM,sBAAuB7B,KAAKoF,UAAUzB,UAExD3D,KAAKoF,UAAUlB,OAGflE,KAAKkF,UAAUnD,QAGf/B,KAAKwF,kBAAoBV,2BACzB9E,KAAKgG,mBAAqB,EAC1BhG,KAAKiG,eAAiB,EACtBjG,KAAK8F,kBAAmB,EAGxB9F,KAAKkG,oBAGLlG,KAAKwG,mBAAmBlB,sBAAsBmB,YAC9C,IAAK,MAAM3F,WAAWd,KAAKoG,SAASY,SAC9BlG,QAAQM,aAAelB,mBAAmB0B,QAE9Cd,QAAQmB,yBAMVjC,KAAKkF,UAAU+B,SACb,UACyBhE,IAAnBjD,KAAKoF,WAGTpF,KAAKoF,UAAUT,OAAK,EAEG,IAAzB3E,KAAKkG,kBACL,YArCA,CAqCW,EAIfK,KAAAA,WAAa,KACPvG,KAAKqF,kBAAoBC,sBAAsBC,gBAEnDvF,KAAKQ,OAAOqB,MAAM,iBAGlB7B,KAAKoF,WAAWlB,OAChBlE,KAAKoF,eAAYnC,EAGjBjD,KAAKkF,UAAUnD,QAGf/B,KAAKwF,kBAAoBV,2BACzB9E,KAAKgG,mBAAqB,EAC1BhG,KAAKiG,eAAiB,EACtBjG,KAAK8F,kBAAmB,EACxB9F,KAAKkG,kBAAoB,EAEzBlG,KAAKwG,mBAAmBlB,sBAAsBC,eAC9CvF,KAAKkH,aAAaxB,gBAAgBC,cACpC,EAAC3F,KAED2B,MAAQ,KACN3B,KAAKuG,YACP,EAEAY,KAAAA,qBAAuB,IAAMnH,KAAKwF,kBAAiBxF,KACnDoH,mBAAqB,IAAMpH,KAAKqF,gBAChCgC,KAAAA,iCAAoCrG,UAClChB,KAAK4F,+BAA+B3E,IAAID,UAAShB,KACnDsH,oCAAuCtG,UACrChB,KAAK4F,+BAA+BzE,OAAOH,UAE7CuG,KAAAA,aAAgBC,QACdxH,KAAK+F,oBAAsByB,MAEvBxH,KAAKqF,kBAAoBC,sBAAsBmC,WACjDzH,KAAK0H,gBAAgBF,MACtB,EAGHG,KAAAA,aAAe,IAAuB3H,KAAKyF,UAASzF,KACpD4H,2BAA8B5G,UAC5BhB,KAAK6F,yBAAyB5E,IAAID,UACpC6G,KAAAA,8BAAiC7G,UAC/BhB,KAAK6F,yBAAyB1E,OAAOH,UAAShB,KAEhDuB,iBAAoBP,UAAkChB,KAAKO,eAAeU,IAAID,UAC9EQ,KAAAA,oBAAuBR,UAAkChB,KAAKO,eAAeY,OAAOH,UAAShB,KAE7F8H,YAAc,CAAClI,QAAiBC,cAC9B,MAAMkI,UAAY/H,KAAKmG,gBACvBnG,KAAKmG,iBAAmB,EAExB,MAAMrF,QAAU,IAAIrB,uBAClBsI,UACAnI,QACAC,WACAG,KAAKF,YACLE,KAAKD,QAaP,OAVAC,KAAKoG,SAAS4B,IAAID,UAAWjH,SAI3Bd,KAAKqF,kBAAoBC,sBAAsBmC,WAC/CzH,KAAKyF,YAAcC,gBAAgBuC,YAEnCnH,QAAQkB,UAGHlB,SAGD0F,KAAAA,mBAAsBhE,YAC5B,MAAMC,KAAOzC,KAAKqF,gBAClB,GAAI5C,OAASD,UAAb,CAEAxC,KAAKqF,gBAAkB7C,UACvB,IAAK,MAAMxB,YAAYhB,KAAK4F,+BAC1B5E,SAASwB,UAAWC,KAJE,CAKvB,EAGK3C,KAAAA,YAAe4B,UACrB1B,KAAKoF,WAAWtF,YAAY4B,SAE5B1B,KAAKkI,oBAGLlI,KAAKiG,eAAiBkC,KAAKC,KAC7B,EAACpI,KAEO0H,gBAAmBF,QACzBxH,KAAKQ,OAAOqB,MAAM,wBAElB7B,KAAKkH,aAAaxB,gBAAgB2C,aAElCrI,KAAKF,YAAY,CACfY,KAAM,OACNI,QAAS,EACT0G,aACD,EAGKN,KAAAA,aAAgBoB,WACtB,MAAM7F,KAAOzC,KAAKyF,UAElBzF,KAAKyF,UAAY6C,SACjB,IAAK,MAAMtH,YAAYhB,KAAK6F,yBAC1B,IACE7E,SAASsH,SAAU7F,KACpB,CAAC,MAAOF,GACPvC,KAAKQ,OAAOiB,MAAM,4BAA6Bc,EAChD,CACF,EAGKqE,KAAAA,eAAkBlF,UAaxB,GAZA1B,KAAKgG,mBAAqBmC,KAAKC,MAI3BpI,KAAKgG,mBAAqBhG,KAAKiG,gBAAkD,IAAhCjG,KAAKD,OAAOwI,mBAC/DvI,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,ICtNfY,UAEoB,IAApBA,QAAQZ,UACU,UAAjBY,QAAQhB,MACU,cAAjBgB,QAAQhB,MACS,SAAjBgB,QAAQhB,MACS,eAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MDoNJ8H,CAAoB9G,SACtB,OAAQA,QAAQhB,MACd,IAAK,QACH,OAAOV,KAAKyI,oBAAoB/G,SAClC,IAAK,aACH,OAAW1B,KAAC0I,wBAAwBhH,SACtC,IAAK,QACH,OAAW1B,KAAC2I,aAAa,CACvBjI,KAAMgB,QAAQD,MACdC,QAASA,QAAQA,UAErB,IAAK,YAEH,YAEKkH,GCxNkBlH,UACX,IAApBA,QAAQZ,QDuNK8H,CAAiBlH,SAAU,CACpC,MAAMZ,QAAUd,KAAKoG,SAASyC,IAAInH,QAAQZ,SAC1C,QAAgBmC,IAAZnC,QAEF,YADAd,KAAKQ,OAAOiE,KAAK,iDAAkD/C,SAIrE,GC3NJA,UAEiB,mBAAjBA,QAAQhB,MACS,mBAAjBgB,QAAQhB,MACS,UAAjBgB,QAAQhB,MACS,oBAAjBgB,QAAQhB,MACS,mBAAjBgB,QAAQhB,KDqNAoI,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,EAACV,KAEOyI,oBAAuBM,cAE7B/I,KAAKkF,UAAU8D,OAAO,iBAIpBhJ,KAAKqF,kBAAoBC,sBAAsBmB,YAC/CzG,KAAKqF,kBAAoBC,sBAAsBmC,YAE/CzH,KAAKwF,kBAAoB,IACpBxF,KAAKwF,kBACRyD,cAAeF,YAAYG,QAC3BC,uBAAwBnJ,KAAKD,OAAOqJ,iBACpCC,uBAAwBN,YAAYK,kBAItCpJ,KAAKkG,kBAAoB,OAEQjD,IAA7BjD,KAAK+F,qBACP/F,KAAKwG,mBAAmBlB,sBAAsBmC,YAKlD,MAAM6B,aAAsD,KAAtCP,YAAYK,kBAAoB,IACtDpJ,KAAKkF,UAAU+B,SAAS,IAAMjH,KAAKuJ,aAAaD,cAAeA,aAAc,UAAS,EACvFtJ,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,EACFzB,KAEO0I,wBAA0B,EAAGc,gBACnCxJ,KAAKQ,OAAOqB,MAAM,8BAA+B2H,OAGjDxJ,KAAKkF,UAAU8D,OAAO,sBAGlBhJ,KAAK8F,iBACP9F,KAAK8F,kBAAmB,EAGV,iBAAV0D,QACFxJ,KAAK+F,yBAAsB9C,GAKjB,eAAVuG,QACFxJ,KAAKwG,mBAAmBlB,sBAAsBmC,WAE9CzH,KAAKyJ,yBAGPzJ,KAAKkH,aAAaxB,gBAAgB8D,OACpC,EAEQC,KAAAA,sBAAwB,KAC9B,IAAK,MAAM3I,WAAed,KAACoG,SAASY,SAE9BlG,QAAQM,aAAelB,mBAAmB0B,OAK9Cd,QAAQkB,UAJNhC,KAAKoG,SAASjF,OAAOL,QAAQnB,GAKhC,EACFK,KAUO2G,qBAAuB,KAC7B3G,KAAKQ,OAAOqB,MAAM,qBAElB,MAAM6H,aAA6B,CACjChJ,KAAM,QACNI,QAAS,EACToI,QAAS,GAAGlJ,KAAKwF,kBAAkBT,mBAAmB/E,KAAKwF,kBAAkBR,gBAC7EoE,iBAAkBpJ,KAAKD,OAAOqJ,iBAC9BO,uBAAwB3J,KAAKD,OAAO4J,wBAItC3J,KAAKkF,UAAU+B,SACb,KACE,MAAM2C,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,KAAK8G,WACP,EAC4B,IAA5B9G,KAAKD,OAAO8J,cACZ,iBAGF7J,KAAKF,YAAY4J,cAEjB1J,KAAKkF,UAAU+B,SACb,KACE,MAAM2C,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,KAAK8G,WACP,EAC4B,IAA5B9G,KAAKD,OAAO8J,cACZ,2BAG+B5G,IAA7BjD,KAAK+F,qBACP/F,KAAK0H,gBAAgB1H,KAAK+F,oBAC3B,EACF/F,KAQO6G,sBAAwB,CAAC1C,OAAgB1C,SAU/C,GATAzB,KAAKQ,OAAOqB,MAAM,oBAAqBsC,aAEzBlB,IAAVxB,OACFzB,KAAK2I,aAAa,CAChBjI,KAAM,UACNgB,QAASyC,SAITnE,KAAKyF,YAAcC,gBAAgBC,aAGrC,OAFA3F,KAAK+F,yBAAsB9C,OAC3BjD,KAAKuG,aAIPvG,KAAK8G,WAAS,EAGRyC,KAAAA,aAAgBD,eACtB,MACMQ,oBADM3B,KAAKC,MACiBpI,KAAKgG,mBACvC,GAAI8D,qBAAuBR,aAQzB,OAPAtJ,KAAKF,YAAY,CACfY,KAAM,QACNI,QAAS,EACTW,MAAO,UACPC,QAAS,6BAA+BoI,oBAAsB,OAGrD9J,KAAC8G,YAGd,MAAMiD,YAAcC,KAAKC,IAAI,IAAKX,aAAeQ,qBACjD9J,KAAKkF,UAAU+B,SAAS,IAAMjH,KAAKuJ,aAAaD,cAAeS,YAAa,UAC9E,EAEQ7B,KAAAA,kBAAoB,KAC1BlI,KAAKkF,UAAU+B,SACb,KACEjH,KAAKF,YAAY,CACfY,KAAM,YACNI,QAAS,IAGXd,KAAKkI,mBAAiB,EAEQ,IAAhClI,KAAKD,OAAOwI,kBACZ,YAAW,EA/dbvI,KAAKD,OAAS,CACZwI,kBAAmB,GACnBa,iBAAkB,GAClBO,uBAAwB,GACxBE,cAAe,GACfjH,SAAUsH,eAAeC,KACzBpD,sBAAuB,EACvBL,iBAAmB5D,KAAQ,IAAID,gCAAgCC,QAC5D/C,QAGLC,KAAKQ,OAAS,IAAIkC,OAAO1C,KAAKN,YAAYiD,KAAM3C,KAAKD,OAAO6C,SAC9D"}

@@ -62,7 +62,7 @@ export interface AuthMessage {

export type ConnectionMessage = SetupMessage | KeepaliveMessage | AuthMessage | AuthStateMessage | ErrorMessage;
export type Message = SetupMessage | KeepaliveMessage | AuthMessage | AuthStateMessage | ErrorMessage | ChannelRequestMessage | ChannelCancelMessage | ChannelOpenedMessage | ChannelClosedMessage | ChannelPayloadMessage | ChannelErrorMessage;
export declare const isConnectionMessage: (message: Message) => message is ConnectionMessage;
export type DXLinkWebSocketMessage = SetupMessage | KeepaliveMessage | AuthMessage | AuthStateMessage | ErrorMessage | ChannelRequestMessage | ChannelCancelMessage | ChannelOpenedMessage | ChannelClosedMessage | ChannelPayloadMessage | ChannelErrorMessage;
export declare const isConnectionMessage: (message: DXLinkWebSocketMessage) => message is ConnectionMessage;
export type ChannelLifecycleMessage = ChannelOpenedMessage | ChannelClosedMessage | ChannelErrorMessage | ChannelRequestMessage | ChannelCancelMessage;
export type ChannelMessage = ChannelLifecycleMessage | ChannelPayloadMessage;
export declare const isChannelMessage: (message: Message) => message is ChannelMessage;
export declare const isChannelMessage: (message: DXLinkWebSocketMessage) => message is ChannelMessage;
export declare const isChannelLifecycleMessage: (message: ChannelMessage) => message is ChannelLifecycleMessage;
{
"name": "@dxfeed/dxlink-websocket-client",
"version": "0.3.0",
"version": "0.4.0",
"private": false,

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

"dependencies": {
"@dxfeed/dxlink-core": "0.3.0"
"@dxfeed/dxlink-core": "0.4.0"
},

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

import { type DXLinkWebSocketClient } from './dxlink';
import { type DXLinkWebSocketClientOptions } from './client';
import { type DXLinkFeed, type DXLinkFeedOptions } from './feed';
import { DXLinkLogLevel } from '@dxfeed/dxlink-core';
import { FeedContract } from './feed-messages';
/**
* The options to use for the {@link DXLink}.
*/
export interface DXLinkOptions {
logLevel?: DXLinkLogLevel;
/**
* The options for the {@link DXLinkWebSocketClient}.
*/
client?: Partial<DXLinkWebSocketClientOptions>;
/**
* The options for the {@link DXLinkFeed}.
*/
feed?: Partial<DXLinkFeedOptions>;
}
/**
* The main entry point for the DXLink API.
* This class provides access to dxFeed services with a single dxLink connection.
* It is recommended to use a single instance of this class per application.
* @example <caption>Creating a new instance of {@link DXLink} with default options.</caption>
* ```typescript
* const dxlink = new DXLink()
* ```
* @example <caption>Connecting to the server.</caption>
* ```typescript
* dxlink.connect('wss://demo.dxfeed.com/dxlink-ws')
* ```
* @example <caption>Creating a new instance of {@link DXLinkFeed}.</caption>
* ```typescript
* const feed = dxlink.createFeed(FeedContract.TICKER)
* ````
* @example <caption>Subscribing to feed events.</caption>
* ```typescript
* feed.addSubscription({ type: 'Quote', symbol: 'AAPL' })
* feed.addEventListener((events) => console.log(events))
* ```
*/
export declare class DXLink {
private readonly options?;
/**
* A singleton instance of {@link DXLink} that is created on the first use of {@link DXLink.getInstance}.
*/
private static instance;
/**
* Returns a default application-wide singleton instance of {@link DXLink}.
* Most applications use only a single data-source and should rely on this method to get one.
* This method creates an endpoint on the first use with a default configuration.
*/
static getInstance(): DXLink;
private readonly client;
/**
* Creates a new instance of {@link DXLink} with the specified options.
* @param options - The options to use for the instance.
*/
constructor(options?: DXLinkOptions | undefined);
/**
* Connects to the server.
* @see {DXLinkWebSocketClient.connect} for more details.
*/
connect(url: string): void;
/**
* Reconnects to the server.
* @see {DXLinkWebSocketClient.reconnect} for more details.
*/
reconnect(): void;
/**
* Disconnects from the server.
* @see {DXLinkWebSocketClient.disconnect} for more details.
*/
disconnect(): void;
/**
* Returns the {@link DXLinkWebSocketClient} instance.
*/
getClient(): DXLinkWebSocketClient;
/**
* Creates a new feed for the specified contract.
* @see {DXLinkFeed} for more details.
*/
createFeed<Contract extends FeedContract>(contract: Contract): DXLinkFeed<Contract>;
/**
* Closes the client and releases all resources.
*/
close(): void;
}
import type { DXLinkLogLevel } from '@dxfeed/dxlink-core';
/**
* Options for {@link DXLinkWebSocketClient}.
*/
export interface DXLinkWebSocketClientConfig {
/**
* Interval in seconds between keepalive messages which are sent to server.
*/
readonly keepaliveInterval: number;
/**
* Timeout in seconds for server to detect that client is disconnected.
* @see {DXLinkConnectionDetails.clientKeepaliveTimeout}
*/
readonly keepaliveTimeout: number;
/**
* Prefered timeout in seconds in for client to detect that server is disconnected.
* @see {DXLinkConnectionDetails.serverKeepaliveTimeout}
*/
readonly acceptKeepaliveTimeout: number;
/**
* Timeout for action which requires update from server.
*/
readonly actionTimeout: number;
/**
* Log level for internal logger.
*/
readonly logLevel: DXLinkLogLevel;
/**
* Maximum number of reconnect attempts.
* If connection is not established after this number of attempts, connection will be closed.
* If 0, then reconnect attempts are not limited.
*/
readonly maxReconnectAttempts: number;
}
import { type ErrorType } from './messages';
/**
* Error type that can be used to handle errors or send them to the remote endpoint.
*/
export type DXLinkErrorType = ErrorType;
/**
* Unified error that can be used to handle errors or send them to the remote endpoint.
* @see {DXLinkChannel.error}
* @see {DXLinkWebSocketClient.addErrorListener}
*/
export interface DXLinkError {
/**
* Type of the error.
* @example 'TIMEOUT'
*/
readonly type: DXLinkErrorType;
/**
* Message of the error with details.
* @example 'Timeout exceeded'
*/
readonly message: string;
}
/**
* Connection state that can be used to check if connection is established and ready to use.
* @see {DXLinkWebSocketClient.getConnectionState}
*/
export declare enum DXLinkConnectionState {
/**
* Client was created and not connected to remote endpoint.
* @see {DXLinkWebSocketClient.connect}
*/
NOT_CONNECTED = "NOT_CONNECTED",
/**
* The {@link DXLinkWebSocketClient.connect} method was called to establish connection or {@link DXLinkWebSocketClient.reconnect} is in progress.
* The connection is not ready to use yet.
*/
CONNECTING = "CONNECTING",
/**
* The connection to remote endpoint is established.
* The connection is ready to use.
*/
CONNECTED = "CONNECTED"
}
/**
* Connection details that can be used for debugging or logging.
*/
export interface DXLinkConnectionDetails {
/**
* Protocol version used for connection to the remote endpoint.
*/
readonly protocolVersion: string;
/**
* Version of the client library.
*/
readonly clientVersion: string;
/**
* Version of the server which the client is connected to.
*/
readonly serverVersion?: string;
/**
* Timeout in seconds for server to detect that client is disconnected.
* If no keepalive message received from client during this time, server will close connection.
*/
readonly clientKeepaliveTimeout?: number;
/**
* Timeout in seconds in for client to detect that server is disconnected.
* If no keepalive message received from server during this time, client will close connection.
*/
readonly serverKeepaliveTimeout?: number;
}
/**
* Listener for connection state changes.
*/
export type DXLinkConnectionStateChangeListener = (state: DXLinkConnectionState, prev: DXLinkConnectionState) => void;
/**
* Authentication state that can be used to check if user is authorized on the remote endpoint.
*/
export declare enum DXLinkAuthState {
/**
* User is unauthorized on the remote endpoint.
*/
UNAUTHORIZED = "UNAUTHORIZED",
/**
* User in the process of authorization, but not yet authorized.
*/
AUTHORIZING = "AUTHORIZING",
/**
* User is authorized on the remote endpoint and can use it.
*/
AUTHORIZED = "AUTHORIZED"
}
/**
* Listener for authentication state changes.
*/
export type DXLinkAuthStateChangeListener = (state: DXLinkAuthState, prev: DXLinkAuthState) => void;
/**
* Message that can be sent to or received from the channel.
* @see {DXLinkChannel.send} and {@link DXLinkChannel.addMessageListener}
*/
export interface DXLinkChannelMessage {
/**
* Type of the message.
*/
readonly type: string;
/**
* Payload of the message.
*/
readonly [key: string]: any;
}
/**
* Listener for messages from the channel.
* @see {DXLinkChannel.addMessageListener}
*/
export type DXLinkChannelMessageListener = (message: DXLinkChannelMessage) => void;
/**
* Listener for channel state changes.
* @see {DXLinkChannel.addStateChangeListener}
*/
export type DXLinkChannelStateChangeListener = (state: DXLinkChannelState, prev: DXLinkChannelState) => void;
/**
* Listener for errors from the server.
* @see {DXLinkWebSocketClient.addErrorListener}
*/
export type DXLinkErrorListener = (error: DXLinkError) => void;
/**
* Channel state that can be used to check if channel is available for sending messages.
* @see {DXLinkChannel.getState}
*/
export declare enum DXLinkChannelState {
/**
* Channel is requested and cannot be used to {@link DXLinkChannel.send} messages yet.
*/
REQUESTED = "REQUESTED",
/**
* Channel is opened and can be used to {@link DXLinkChannel.send} messages.
*/
OPENED = "OPENED",
/**
* Channel was closed by {@link DXLinkChannel.close} or by server.
* Channel cannot be used anymore.
*/
CLOSED = "CLOSED"
}
/**
* Isolated channel to service withing single {@link DXLinkWebSocketClient} connection to remote endpoint.
* @see DXLinkWebSocketClient.openChannel
*/
export interface DXLinkChannel {
/**
* Unique identifier of the channel.
*/
readonly id: number;
/**
* Name of the service that channel is opened to.
*/
readonly service: string;
/**
*
*/
readonly parameters: Record<string, unknown>;
/**
* Send a message to the channel.
*/
send(message: DXLinkChannelMessage): void;
/**
* Add a listener for messages from the channel.
*/
addMessageListener(listener: DXLinkChannelMessageListener): void;
/**
* Remove a listener for messages from the channel.
*/
removeMessageListener(listener: DXLinkChannelMessageListener): void;
/**
* Get channel state that can be used to check if channel is available for sending messages.
*/
getState(): DXLinkChannelState;
/**
* Add a listener for channel state changes.
* If channel is ready to use, listener will be called immediately with {@link DXLinkChannelState.OPENED} state.
* @see {DXLinkChannelState}
* Note: when remote endpoint reconnects, channel will be reopened and listener will be called with {@link DXLinkChannelState.OPENED} state again.
* @see {DXLinkWebSocketClient.getConnectionState}
*/
addStateChangeListener(listener: DXLinkChannelStateChangeListener): void;
/**
* Remove a listener for channel state changes.
*/
removeStateChangeListener(listener: DXLinkChannelStateChangeListener): void;
/**
* Add a listener for errors from the server.
* @see {DXLinkError}
*/
addErrorListener(listener: DXLinkErrorListener): void;
/**
* Remove a listener for errors from the server.
* @see {DXLinkError}
*/
removeErrorListener(listener: DXLinkErrorListener): void;
/**
* Close the channel and free all resources.
* This method does nothing if the channel is already closed.
* The channel {@see DXLinkChannel.getStatus} immediately becomes {@link State#CLOSED CLOSED}.
* @see {DXLinkChannelState}
*/
close(): void;
}
/**
* dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.
*/
export interface DXLinkWebSocketClient {
/**
* Connect to the remote endpoint.
* Connects to the specified remote address. Previously established connections are closed if the new address is different from the old one.
* This method does nothing if address does not change.
* The endpoint {@see DXLinkWebSocketClient.getConnectionState} immediately becomes {@link State#CONNECTING CONNECTING}.
*
* For connection with the authorization token, use {@link DXLinkWebSocketClient.setAuthToken} before calling this method.
* If the token is not set, the connection will be established without authorization.
*
* @param url WebSocket URL to connect to.
*/
connect(url: string): void;
/**
* Reconnect to the remote endpoint.
* This method does nothing if the client is not connected.
* The endpoint {@see DXLinkWebSocketClient.getConnectionState} immediately becomes {@link State#CONNECTING CONNECTING}.
* @see {DXLinkWebSocketClient.connect}
*/
reconnect(): void;
/**
* Disconnect from the remote endpoint.
* This method does nothing if the client is not connected.
* The endpoint {@see DXLinkWebSocketClient.getConnectionState} immediately becomes {@link State#NOT_CONNECTED NOT_CONNECTED}.
* @see {DXLinkWebSocketClient.connect}
*/
disconnect(): void;
/**
* Get connection details that can be used for debugging or logging.
*/
getConnectionDetails(): DXLinkConnectionDetails;
/**
* Get connection state that can be used to check if connection is established and ready to use.
*/
getConnectionState(): DXLinkConnectionState;
/**
* Add a listener for connection state changes.
*/
addConnectionStateChangeListener(listener: DXLinkConnectionStateChangeListener): void;
/**
* Remove a listener for connection state changes.
*/
removeConnectionStateChangeListener(listener: DXLinkConnectionStateChangeListener): void;
/**
* Set authorization token to be used for connection to the remote endpoint.
* This method does nothing if the client is connected.
* @param token Authorization token to be used for connection.
*/
setAuthToken(token: string): void;
/**
* Get authentication state that can be used to check if user is authorized on the remote endpoint.
*/
getAuthState(): DXLinkAuthState;
/**
* Add a listener for authentication state changes.
* When auth state is {@link DXLinkAuthState.UNAUTHORIZED}, you can call {@link DXLinkWebSocketClient.setAuthToken} to authorize the client.
*/
addAuthStateChangeListener(listener: DXLinkAuthStateChangeListener): void;
/**
* Remove a listener for authentication state changes.
*/
removeAuthStateChangeListener(listener: DXLinkAuthStateChangeListener): void;
/**
* Error listener that can be used to handle errors from the server.
*/
addErrorListener(listener: DXLinkErrorListener): void;
/**
* Remove a listener for errors from the server.
*/
removeErrorListener(listener: DXLinkErrorListener): void;
/**
* Open a isolated channel to service withing single {@link DXLinkWebSocketClient} connection to remote endpoint.
* @param service Name of the service to open channel to.
* @param parameters Parameters of the service to open channel to.
* @see {DXLinkChannel}
*/
openChannel(service: string, parameters: Record<string, unknown>): DXLinkChannel;
/**
* Close the client and free all resources.
* This method works the same as {@link DXLinkWebSocketClient.disconnect}.
*/
close(): void;
}
import { type DXLinkChannelMessage } from './dxlink';
export declare enum FeedContract {
'TICKER' = "TICKER",
'HISTORY' = "HISTORY",
'STREAM' = "STREAM",
'AUTO' = "AUTO"
}
export declare enum FeedDataFormat {
'FULL' = "FULL",
'COMPACT' = "COMPACT"
}
export interface FeedParameters {
readonly contract: FeedContract;
}
export interface FeedEventFields {
[eventType: string]: string[];
}
export interface FeedSetupMessage {
readonly type: 'FEED_SETUP';
readonly acceptAggregationPeriod?: number;
readonly acceptDataFormat?: FeedDataFormat;
readonly acceptEventFields?: FeedEventFields;
}
export interface FeedConfigMessage {
readonly type: 'FEED_CONFIG';
readonly aggregationPeriod: number;
readonly dataFormat: FeedDataFormat;
readonly eventFields?: FeedEventFields;
}
export type Subscription = {
readonly type: string;
readonly symbol: string;
};
export type TimeSeriesSubscription = {
readonly type: string;
readonly symbol: string;
readonly fromTime: number;
};
export type IndexedEventSubscription = {
readonly type: string;
readonly symbol: string;
readonly source: string;
};
export interface FeedSubscriptionMessage {
readonly type: 'FEED_SUBSCRIPTION';
readonly add?: (Subscription | TimeSeriesSubscription | IndexedEventSubscription)[];
readonly remove?: (Subscription | TimeSeriesSubscription | IndexedEventSubscription)[];
readonly reset?: boolean;
}
export type FeedEventValue = number | string | boolean;
export interface FeedEventData {
[key: string]: FeedEventValue;
}
export type FeedCompactEventData = [string, FeedEventValue[]];
export interface FeedDataMessage {
readonly type: 'FEED_DATA';
readonly data: FeedEventData[] | FeedCompactEventData;
}
export type FeedMessage = FeedSetupMessage | FeedConfigMessage | FeedSubscriptionMessage | FeedDataMessage;
export declare const isFeedFullData: (data: FeedEventData[] | FeedCompactEventData) => data is FeedEventData[];
export declare const isFeedCompactData: (data: FeedEventData[] | FeedCompactEventData) => data is FeedCompactEventData;
export declare const isFeedMessage: (message: DXLinkChannelMessage) => message is FeedMessage;
import { type DXLinkChannel, DXLinkChannelState, type DXLinkChannelStateChangeListener, type DXLinkWebSocketClient } from './dxlink';
import { type FeedEventFields, FeedContract, FeedDataFormat, type IndexedEventSubscription, type Subscription, type TimeSeriesSubscription, type FeedEventData } from './feed-messages';
import { DXLinkLogLevel } from '@dxfeed/dxlink-core';
/**
* Prefered configuration for the feed channel.
* Server can ignore some of the parameters and use own defaults.
* @see {DXLinkFeed.configure}
*/
export interface FeedAcceptConfig {
/**
* Aggregation period in seconds.
* If not specified, the channel will use the default value.
* If specified as 0, the channel will try not aggregate events.
*/
acceptAggregationPeriod?: number;
/**
* Data format to be used for received events.
* If not specified, the channel will use the default value `FULL`.
*/
acceptDataFormat?: FeedDataFormat;
/**
* Event fields to be included in received events.
* If not specified, the channel will use the default value.
* If specified as an empty array, the channel will try to send events with default fields.
*/
acceptEventFields?: FeedEventFields;
}
/**
* Configuration of the feed channel.
*/
export interface FeedConfig {
/**
* Aggregation period in seconds.
* @example 0.5 - 500 milliseconds.
* @default `NaN`
* @see {FeedAcceptConfig.acceptAggregationPeriod}
*/
readonly aggregationPeriod: number;
/**
* Data format to be used for received events.
* @example `FULL` - object with keys and values.
* @example `COMPACT` - array of values.
* @default `FULL`
* @see {FeedAcceptConfig.acceptDataFormat}
*/
readonly dataFormat: FeedDataFormat;
/**
* Event fields to be included in received events.
* You can specify fields for all event types or for specific event types @see {FeedAcceptConfig.acceptEventFields}.
* @example ```json
* { "Quote": ["eventSymbol", "askPrice", "bidPrice"] }
* ```
* @default `{}`
*/
readonly eventFields: FeedEventFields;
}
/**
* Listner for the feed channel config changes.
*/
export type DXLinkFeedConfigChangeListner = (config: FeedConfig) => void;
/**
* Subscription type by the contract.
*/
export type SubscriptionByContract = {
[FeedContract.AUTO]: Subscription | TimeSeriesSubscription | IndexedEventSubscription;
[FeedContract.TICKER]: Subscription;
[FeedContract.HISTORY]: TimeSeriesSubscription | IndexedEventSubscription;
[FeedContract.STREAM]: Subscription | TimeSeriesSubscription | IndexedEventSubscription;
};
/**
* Listner for the feed channel events received from the channel.
*/
export type DXLinkFeedEventListner = (event: FeedEventData[]) => void;
/**
* dxLink FEED service instance for the specified {@link FeedContract}.
*/
export interface DXLinkFeed<Contract extends FeedContract = FeedContract.AUTO> {
/**
* Unique identifier of the feed channel.
*/
readonly id: number;
/**
* Contract of the feed channel.
* @see {FeedContract}
*/
readonly contract: Contract;
/**
* Get current channel of the feed.
* Note: inaproppriate usage of the channel can lead to unexpected behavior.
* @see {DXLinkChannel}
*/
getChannel(): DXLinkChannel;
/**
* Configure desired configuration of the feed channel.
* @see {FeedAcceptConfig}
*/
configure(acceptConfig: FeedAcceptConfig): void;
/**
* Get current configuration of the feed channel as received from the channel.
*/
getConfig(): FeedConfig;
/**
* Add a listener for the feed channel config changes.
*/
addConfigChangeListener(listener: DXLinkFeedConfigChangeListner): void;
/**
* Remove a listener for the feed channel config changes.
*/
removeConfigChangeListener(listener: DXLinkFeedConfigChangeListner): void;
/**
* Add subscriptions to the feed channel.
* @param subscriptions - Subscriptions to be added.
*/
addSubscriptions(subscriptions: SubscriptionByContract[Contract][]): void;
/**
* Add subscriptions to the feed channel.
* @param subscriptions - Subscriptions to be added.
*/
addSubscriptions(...subscriptions: SubscriptionByContract[Contract][]): void;
/**
* Remove subscriptions from the feed channel.
* @param subscriptions - Subscriptions to be removed.
*/
removeSubscriptions(subscriptions: SubscriptionByContract[Contract][]): void;
/**
* Remove subscriptions from the feed channel.
* @param subscriptions - Subscriptions to be removed.
*/
removeSubscriptions(...subscriptions: SubscriptionByContract[Contract][]): void;
/**
* Remove all active subscriptions from the feed channel.
*/
clearSubscriptions(): void;
/**
* Add a listener for the feed channel events received from the channel.
*/
addEventListener(listener: DXLinkFeedEventListner): void;
/**
* Remove a listener for the feed channel events received from the channel.
*/
removeEventListener(listener: DXLinkFeedEventListner): void;
/**
* Close the feed channel.
*/
close(): void;
}
/**
* Options for the {@link DXLinkFeedImpl} instance.
*/
export interface DXLinkFeedOptions {
/**
* Time in milliseconds to wait for more pending subscriptions before sending them to the channel.
*/
batchSubscriptionsTime: number;
/**
* Maximum size of the subscription chunk to be sent to the channel.
*/
maxSendSubscriptionChunkSize: number;
/**
* Log level for the feed.
*/
logLevel: DXLinkLogLevel;
}
export declare class DXLinkFeedImpl<Contract extends FeedContract> implements DXLinkFeed<Contract> {
readonly contract: Contract;
readonly id: number;
private readonly options;
/**
* dxLink channel instance.
*/
private readonly channel;
/**
* Current accept config of the feed channel.
*/
private acceptConfig;
/**
* Current config of the feed channel.
*/
private config;
private readonly configListeners;
private readonly eventListeners;
/**
* Pending add subscriptions to be sent to the channel.
*/
private readonly pendingAdd;
/**
* Pending remove subscriptions to be sent to the channel.
*/
private readonly pendingRemove;
/**
* Pending reset flag to be sent to the channel.
*/
private pengingReset;
/**
* List of active subscriptions.
* Used to avoid sending the same subscription twice and re-subscribe on the channel re-open.
*/
private readonly subscriptions;
/**
* List of event types which schema was sent to the channel.
*/
private readonly touchedEvents;
/**
* Timeout identifier for the {@link scheduleProcessPendings} method.
*/
private scheduleTimeoutId;
private readonly logger;
/**
* Allows to create {@link DXLinkFeed} instance with the specified {@link FeedContract} for the given {@link DXLinkWebSocketClient}.
*/
constructor(client: DXLinkWebSocketClient, contract: Contract, options?: Partial<DXLinkFeedOptions>);
getChannel: () => DXLinkChannel;
getState: () => DXLinkChannelState;
addStateChangeListener: (listener: DXLinkChannelStateChangeListener) => void;
removeStateChangeListener: (listener: DXLinkChannelStateChangeListener) => void;
getConfig: () => FeedConfig;
addConfigChangeListener: (listener: DXLinkFeedConfigChangeListner) => Set<DXLinkFeedConfigChangeListner>;
removeConfigChangeListener: (listener: DXLinkFeedConfigChangeListner) => boolean;
close: () => void;
configure: (acceptConfig: FeedAcceptConfig) => void;
addSubscriptions(subscriptions: SubscriptionByContract[Contract][]): void;
addSubscriptions(...subscriptions: SubscriptionByContract[Contract][]): void;
removeSubscriptions(subscriptions: SubscriptionByContract[Contract][]): void;
removeSubscriptions(...subscriptions: SubscriptionByContract[Contract][]): void;
clearSubscriptions: () => void;
addEventListener: (listener: DXLinkFeedEventListner) => Set<DXLinkFeedEventListner>;
removeEventListener: (listener: DXLinkFeedEventListner) => boolean;
/**
* Clean the subscription from the fields which are not allowed for the specified contract.
* Note: coze of the TypeScript limitations, we need to clean the subscription from the fields which are not allowed for the specified contract.
*/
private cleanSubscription;
/**
* Process message received in the channel.
*/
private processMessage;
/**
* Process config received from the channel.
*/
private processConfig;
/**
* Parse data received from the channel.
*/
private parseEventData;
/**
* Process data received from the channel.
*/
private processData;
/**
* Process channel status changes from the channel.
*/
private processStatus;
/**
* Send the subscription chunk to the channel.
* @param chunk Subscription chunk to be sent to the channel.
* @param newTouchedEvents List of event types which schema should be sent to the channel before the chunk.
* @returns
*/
private sendSubscriptionChunkAndSchema;
/**
* Process error received from the channel.
*/
private processError;
/**
* Resubscribe to the feed channel subscriptions after the channel re-open.
*/
private resubscribe;
/**
* Schedule sending pending subscriptions to the channel to batch them together to reduce the number of messages.
*/
private scheduleProcessPendings;
/**
* Process pending subscriptions and send them to the channel.
*/
private processPendings;
/**
* Send the `FEED_SETUP` message to the channel with the event fields for the specified event types.
* @param eventTypes List of event type fields to be sent to the channel.
* @param force If `true`, the config will be sent to the channel even if there is no event fields to send.
*/
private sendAcceptConfig;
}
/**
* Scheduler for scheduling callbacks.
* @internal
*/
export declare class Scheduler {
private timeoutIds;
schedule: (callback: () => void, timeout: number, key: string) => string;
cancel: (key: string) => void;
clear: () => void;
}