@dxfeed/dxlink-websocket-client
Advanced tools
| import { DXLinkLogLevel, DXLinkScheduler, DXLinkClient, DXLinkConnectionDetails, DXLinkConnectionState, DXLinkConnectionStateChangeListener, DXLinkAuthState, DXLinkAuthStateChangeListener, DXLinkErrorListener, DXLinkChannelOptions, DXLinkChannel } from '@dxfeed/dxlink-core'; | ||
| interface AuthMessage { | ||
| type: 'AUTH'; | ||
| channel: 0; | ||
| token: string; | ||
| } | ||
| type AuthState = 'AUTHORIZED' | 'UNAUTHORIZED'; | ||
| interface AuthStateMessage { | ||
| type: 'AUTH_STATE'; | ||
| channel: 0; | ||
| state: AuthState; | ||
| } | ||
| interface SetupMessage { | ||
| type: 'SETUP'; | ||
| channel: 0; | ||
| version: string; | ||
| keepaliveTimeout?: number; | ||
| acceptKeepaliveTimeout?: number; | ||
| } | ||
| interface KeepaliveMessage { | ||
| type: 'KEEPALIVE'; | ||
| channel: 0; | ||
| } | ||
| interface ChannelRequestMessage { | ||
| type: 'CHANNEL_REQUEST'; | ||
| channel: number; | ||
| service: string; | ||
| parameters?: Record<string, unknown>; | ||
| } | ||
| interface ChannelCancelMessage { | ||
| type: 'CHANNEL_CANCEL'; | ||
| channel: number; | ||
| } | ||
| interface ChannelOpenedMessage { | ||
| type: 'CHANNEL_OPENED'; | ||
| channel: number; | ||
| service: string; | ||
| parameters?: Record<string, unknown>; | ||
| } | ||
| interface ChannelClosedMessage { | ||
| type: 'CHANNEL_CLOSED'; | ||
| channel: number; | ||
| } | ||
| type ErrorType = 'UNKNOWN' | 'UNSUPPORTED_PROTOCOL' | 'TIMEOUT' | 'UNAUTHORIZED' | 'INVALID_MESSAGE' | 'BAD_ACTION'; | ||
| interface ErrorMessage { | ||
| type: 'ERROR'; | ||
| channel: 0; | ||
| error: ErrorType; | ||
| message: string; | ||
| } | ||
| interface ChannelErrorMessage { | ||
| type: 'ERROR'; | ||
| channel: number; | ||
| error: ErrorType; | ||
| message: string; | ||
| } | ||
| interface ChannelPayloadMessage { | ||
| type: string; | ||
| channel: number; | ||
| [key: string]: unknown; | ||
| } | ||
| type DXLinkWebSocketMessage = SetupMessage | KeepaliveMessage | AuthMessage | AuthStateMessage | ErrorMessage | ChannelRequestMessage | ChannelCancelMessage | ChannelOpenedMessage | ChannelClosedMessage | ChannelPayloadMessage | ChannelErrorMessage; | ||
| /** | ||
| * 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. | ||
| */ | ||
| 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. | ||
| */ | ||
| type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void; | ||
| /** | ||
| * Default connector for the WebSocket connection. | ||
| * @internal | ||
| */ | ||
| declare class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector { | ||
| private readonly url; | ||
| private readonly protocols?; | ||
| private socket; | ||
| private isAvailable; | ||
| private openListener; | ||
| private closeListener; | ||
| private messageListener; | ||
| constructor(url: string, protocols?: string | string[] | undefined); | ||
| start(): void; | ||
| stop(): void; | ||
| sendMessage: (message: DXLinkWebSocketMessage) => void; | ||
| setOpenListener: (listener: () => void) => void; | ||
| setCloseListener: (listener: DXLinkWebSocketCloseListener) => void; | ||
| setMessageListener: (listener: (message: DXLinkWebSocketMessage) => void) => void; | ||
| getUrl: () => string; | ||
| private handleOpen; | ||
| private handleClosed; | ||
| private handleError; | ||
| private handleMessage; | ||
| } | ||
| /** | ||
| * Options for {@link DXLinkWebSocketClient}. | ||
| */ | ||
| 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. | ||
| * Default: -1. By default, reconnect attempts are unlimited. | ||
| */ | ||
| readonly maxReconnectAttempts: number; | ||
| /** | ||
| * Scheduler used by the client for reconnect and timeout handling. | ||
| * If not provided, {@link DefaultDXLinkScheduler} is used. | ||
| */ | ||
| readonly scheduler?: DXLinkScheduler; | ||
| /** | ||
| * 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; | ||
| } | ||
| /** | ||
| * Protocol version that is used by client. | ||
| */ | ||
| declare const DXLINK_WS_PROTOCOL_VERSION = "0.1"; | ||
| /** | ||
| * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services. | ||
| */ | ||
| declare class DXLinkWebSocketClient implements DXLinkClient { | ||
| private readonly config; | ||
| private readonly logger; | ||
| private readonly scheduler; | ||
| private connector; | ||
| private connectionState; | ||
| private connectionDetails; | ||
| private authState; | ||
| private readonly connectionStateChangeListeners; | ||
| private readonly errorListeners; | ||
| private readonly authStateChangeListeners; | ||
| /** | ||
| * Authorization type that was determined by server behavior during setup phase. | ||
| * This value is used to determine if authorization is required or optional or not defined yet. | ||
| */ | ||
| private isFirstAuthState; | ||
| /** | ||
| * Last setted auth token that will be sent to server after connection is established or re-established. | ||
| */ | ||
| private lastSettedAuthToken; | ||
| private lastReceivedMillis; | ||
| private lastSentMillis; | ||
| /** | ||
| * Count of reconnect attempts since last successful connection. | ||
| */ | ||
| private reconnectAttempts; | ||
| private globalChannelId; | ||
| private readonly channels; | ||
| /** | ||
| * Create new instance of {@link DXLinkWebSocketClient}. | ||
| * @param config Configuration of the client. | ||
| */ | ||
| constructor(config?: Partial<DXLinkWebSocketClientConfig>); | ||
| connect: (url: string) => void; | ||
| reconnect: () => void; | ||
| disconnect: () => void; | ||
| close: () => void; | ||
| getConnectionDetails: () => DXLinkConnectionDetails; | ||
| getConnectionState: () => DXLinkConnectionState; | ||
| getScheduler: () => DXLinkScheduler; | ||
| addConnectionStateChangeListener: (listener: DXLinkConnectionStateChangeListener) => Set<DXLinkConnectionStateChangeListener>; | ||
| removeConnectionStateChangeListener: (listener: DXLinkConnectionStateChangeListener) => boolean; | ||
| setAuthToken: (token: string) => void; | ||
| getAuthState: () => DXLinkAuthState; | ||
| addAuthStateChangeListener: (listener: DXLinkAuthStateChangeListener) => Set<DXLinkAuthStateChangeListener>; | ||
| removeAuthStateChangeListener: (listener: DXLinkAuthStateChangeListener) => boolean; | ||
| addErrorListener: (listener: DXLinkErrorListener) => Set<DXLinkErrorListener>; | ||
| removeErrorListener: (listener: DXLinkErrorListener) => boolean; | ||
| openChannel: (service: string, parameters: Record<string, unknown>, options?: DXLinkChannelOptions) => DXLinkChannel; | ||
| private setConnectionState; | ||
| private sendMessage; | ||
| private sendAuthMessage; | ||
| private setAuthState; | ||
| private processMessage; | ||
| private processSetupMessage; | ||
| private publishError; | ||
| private processAuthStateMessage; | ||
| private requestActiveChannels; | ||
| /** | ||
| * Process transport open event from connector. | ||
| * After transport is opened: | ||
| * - setup message is sent to server | ||
| * - auth message is sent to server if auth token is set | ||
| * - wait for setup message from server | ||
| * - wait for auth state message from server | ||
| */ | ||
| private processTransportOpen; | ||
| /** | ||
| * Process transport close event from connector. | ||
| * After transport is closed by server: | ||
| * - reconnect if connection is authorized | ||
| * - disconnect if connection is not authorized | ||
| */ | ||
| private processTransportClose; | ||
| private timeoutCheck; | ||
| private scheduleKeepalive; | ||
| } | ||
| export { DXLINK_WS_PROTOCOL_VERSION, DXLinkWebSocketClient, type DXLinkWebSocketClientConfig, type DXLinkWebSocketCloseListener, type DXLinkWebSocketConnector, type DXLinkWebSocketMessage, DefaultDXLinkWebSocketConnector }; |
+618
| // src/client.ts | ||
| import { | ||
| DXLinkLogLevel, | ||
| Logger as Logger2, | ||
| DefaultDXLinkScheduler, | ||
| DXLinkConnectionState, | ||
| DXLinkAuthState, | ||
| DXLinkChannelState as DXLinkChannelState2 | ||
| } from "@dxfeed/dxlink-core"; | ||
| // src/channel.ts | ||
| import { | ||
| Logger, | ||
| DXLinkChannelState | ||
| } from "@dxfeed/dxlink-core"; | ||
| // src/messages.ts | ||
| var isConnectionMessage = (message) => message.channel === 0 && (message.type === "SETUP" || message.type === "KEEPALIVE" || message.type === "AUTH" || message.type === "AUTH_STATE" || message.type === "ERROR"); | ||
| var isChannelMessage = (message) => message.channel !== 0; | ||
| var isChannelLifecycleMessage = (message) => message.type === "CHANNEL_OPENED" || message.type === "CHANNEL_CLOSED" || message.type === "ERROR" || message.type === "CHANNEL_REQUEST" || message.type === "CHANNEL_CANCEL"; | ||
| // src/channel.ts | ||
| var DXLinkWebSocketChannel = class _DXLinkWebSocketChannel { | ||
| constructor(id, service, parameters, reconnect, sendMessage, config) { | ||
| this.id = id; | ||
| this.service = service; | ||
| this.parameters = parameters; | ||
| this.reconnect = reconnect; | ||
| this.sendMessage = sendMessage; | ||
| this.status = DXLinkChannelState.REQUESTED; | ||
| // Listeners | ||
| this.messageListeners = /* @__PURE__ */ new Set(); | ||
| this.statusListeners = /* @__PURE__ */ new Set(); | ||
| this.errorListeners = /* @__PURE__ */ new Set(); | ||
| this.send = ({ type, ...payload }) => { | ||
| if (this.status !== DXLinkChannelState.OPENED) { | ||
| throw new Error("Channel is not ready"); | ||
| } | ||
| this.sendMessage({ | ||
| 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, message }) => this.send({ | ||
| type: "ERROR", | ||
| error: type, | ||
| message | ||
| }); | ||
| this.close = () => { | ||
| if (this.status === DXLinkChannelState.CLOSED) return; | ||
| 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 (this.errorListeners.size === 0) { | ||
| this.logger.error(`Unhandled error in channel#${this.id}: `, error); | ||
| return; | ||
| } | ||
| for (const listener of this.errorListeners) { | ||
| try { | ||
| listener(error); | ||
| } catch (e) { | ||
| this.logger.error(`Error in channel#${this.id} error listener: `, e); | ||
| } | ||
| } | ||
| }; | ||
| 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.logger = new Logger(`${_DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel); | ||
| } | ||
| }; | ||
| // src/connector.ts | ||
| var DefaultDXLinkWebSocketConnector = class { | ||
| constructor(url, protocols) { | ||
| this.url = url; | ||
| this.protocols = protocols; | ||
| this.socket = void 0; | ||
| this.isAvailable = false; | ||
| this.openListener = void 0; | ||
| this.closeListener = void 0; | ||
| this.messageListener = void 0; | ||
| this.sendMessage = (message) => { | ||
| if (this.socket === void 0 || !this.isAvailable) { | ||
| return; | ||
| } | ||
| 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 = () => { | ||
| if (this.socket === void 0) return; | ||
| this.isAvailable = true; | ||
| this.socket.removeEventListener("open", this.handleOpen); | ||
| this.socket.addEventListener("message", this.handleMessage); | ||
| this.openListener?.(); | ||
| }; | ||
| this.handleClosed = (ev) => { | ||
| if (this.socket === void 0) return; | ||
| this.stop(); | ||
| this.closeListener?.(ev.reason, false); | ||
| }; | ||
| this.handleError = (_ev) => { | ||
| if (this.socket === void 0) return; | ||
| this.stop(); | ||
| this.closeListener?.("Unable to connect", true); | ||
| }; | ||
| this.handleMessage = (ev) => { | ||
| try { | ||
| const message = JSON.parse(ev.data); | ||
| if (typeof message !== "object") { | ||
| throw new Error("Unexpected message: " + typeof message); | ||
| } | ||
| if (typeof message.type !== "string") { | ||
| throw new Error("Unexpected message type: " + typeof message.type); | ||
| } | ||
| if (typeof message.channel !== "number") { | ||
| throw new Error("Unexpected message channel: " + typeof message.channel); | ||
| } | ||
| if (this.messageListener === void 0) { | ||
| return console.warn("No message listener set"); | ||
| } | ||
| this.messageListener(message); | ||
| } catch (error) { | ||
| console.error(error instanceof Error ? error : new Error("Parsing error:" + String(error))); | ||
| } | ||
| }; | ||
| } | ||
| start() { | ||
| if (this.socket !== void 0) return; | ||
| 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() { | ||
| if (this.socket === void 0) return; | ||
| 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 = false; | ||
| } | ||
| }; | ||
| // src/version.ts | ||
| var VERSION = false ? "local-unknown" : "0.8.1"; | ||
| // src/client.ts | ||
| var DXLINK_WS_PROTOCOL_VERSION = "0.1"; | ||
| var CLIENT_VERSION = `DXF-JS/${VERSION}`; | ||
| var DXLWS_SCHEDULER_KEY_RECONNECT = "DXLWS_RECONNECT"; | ||
| var DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = "DXLWS_SETUP_TIMEOUT"; | ||
| var DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = "DXLWS_AUTH_STATE_TIMEOUT"; | ||
| var DXLWS_SCHEDULER_KEY_TIMEOUT = "DXLWS_TIMEOUT"; | ||
| var DXLWS_SCHEDULER_KEY_KEEPALIVE = "DXLWS_KEEPALIVE"; | ||
| var DEFAULT_CONNECTION_DETAILS = { | ||
| protocolVersion: DXLINK_WS_PROTOCOL_VERSION, | ||
| clientVersion: CLIENT_VERSION | ||
| }; | ||
| var DXLinkWebSocketClient = class { | ||
| /** | ||
| * Create new instance of {@link DXLinkWebSocketClient}. | ||
| * @param config Configuration of the client. | ||
| */ | ||
| constructor(config) { | ||
| this.connectionState = DXLinkConnectionState.NOT_CONNECTED; | ||
| this.connectionDetails = DEFAULT_CONNECTION_DETAILS; | ||
| this.authState = DXLinkAuthState.UNAUTHORIZED; | ||
| // Listeners | ||
| this.connectionStateChangeListeners = /* @__PURE__ */ new Set(); | ||
| this.errorListeners = /* @__PURE__ */ new Set(); | ||
| this.authStateChangeListeners = /* @__PURE__ */ new Set(); | ||
| /** | ||
| * Authorization type that was determined by server behavior during setup phase. | ||
| * This value is used to determine if authorization is required or optional or not defined yet. | ||
| */ | ||
| this.isFirstAuthState = true; | ||
| // Stats for keepalive | ||
| // TODO: mb move to connector | ||
| this.lastReceivedMillis = 0; | ||
| this.lastSentMillis = 0; | ||
| /** | ||
| * Count of reconnect attempts since last successful connection. | ||
| */ | ||
| this.reconnectAttempts = 0; | ||
| // Channels | ||
| this.globalChannelId = 1; | ||
| this.channels = /* @__PURE__ */ new Map(); | ||
| this.connect = (url) => { | ||
| if (this.connector?.getUrl() === url) return; | ||
| 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) { | ||
| this.logger.warn("Max reconnect attempts reached"); | ||
| this.disconnect(); | ||
| return; | ||
| } | ||
| if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED || this.connector === void 0) | ||
| return; | ||
| 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 = true; | ||
| this.reconnectAttempts++; | ||
| this.setConnectionState(DXLinkConnectionState.CONNECTING); | ||
| for (const channel of this.channels.values()) { | ||
| if (channel.getState() === DXLinkChannelState2.CLOSED) continue; | ||
| if (!channel.reconnect) { | ||
| channel.processStatusClosed(); | ||
| continue; | ||
| } | ||
| channel.processStatusRequested(); | ||
| } | ||
| this.scheduler.schedule( | ||
| () => { | ||
| if (this.connector === void 0) return; | ||
| this.connector.start(); | ||
| }, | ||
| this.reconnectAttempts * 1e3, | ||
| DXLWS_SCHEDULER_KEY_RECONNECT | ||
| ); | ||
| }; | ||
| this.disconnect = () => { | ||
| if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return; | ||
| 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 = true; | ||
| this.reconnectAttempts = 0; | ||
| this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED); | ||
| this.setAuthState(DXLinkAuthState.UNAUTHORIZED); | ||
| }; | ||
| this.close = () => { | ||
| this.disconnect(); | ||
| }; | ||
| this.getConnectionDetails = () => this.connectionDetails; | ||
| this.getConnectionState = () => this.connectionState; | ||
| this.getScheduler = () => this.scheduler; | ||
| this.addConnectionStateChangeListener = (listener) => this.connectionStateChangeListeners.add(listener); | ||
| this.removeConnectionStateChangeListener = (listener) => this.connectionStateChangeListeners.delete(listener); | ||
| this.setAuthToken = (token) => { | ||
| this.lastSettedAuthToken = token; | ||
| if (this.connectionState === DXLinkConnectionState.CONNECTED) { | ||
| this.sendAuthMessage(token); | ||
| } | ||
| }; | ||
| this.getAuthState = () => this.authState; | ||
| this.addAuthStateChangeListener = (listener) => this.authStateChangeListeners.add(listener); | ||
| this.removeAuthStateChangeListener = (listener) => this.authStateChangeListeners.delete(listener); | ||
| this.addErrorListener = (listener) => this.errorListeners.add(listener); | ||
| this.removeErrorListener = (listener) => this.errorListeners.delete(listener); | ||
| this.openChannel = (service, parameters, options) => { | ||
| const channelId = this.globalChannelId; | ||
| this.globalChannelId += 2; | ||
| const channel = new DXLinkWebSocketChannel( | ||
| channelId, | ||
| service, | ||
| parameters, | ||
| options?.reconnect ?? true, | ||
| this.sendMessage, | ||
| this.config | ||
| ); | ||
| this.channels.set(channelId, channel); | ||
| if (this.connectionState === DXLinkConnectionState.CONNECTED && this.authState === DXLinkAuthState.AUTHORIZED) { | ||
| channel.request(); | ||
| } | ||
| return channel; | ||
| }; | ||
| this.setConnectionState = (newStatus) => { | ||
| const prev = this.connectionState; | ||
| if (prev === newStatus) return; | ||
| 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 | ||
| }); | ||
| }; | ||
| 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) => { | ||
| this.lastReceivedMillis = Date.now(); | ||
| if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1e3) { | ||
| this.sendMessage({ | ||
| type: "KEEPALIVE", | ||
| channel: 0 | ||
| }); | ||
| } | ||
| if (isConnectionMessage(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 (isChannelMessage(message)) { | ||
| const channel = this.channels.get(message.channel); | ||
| if (channel === void 0) { | ||
| this.logger.warn("Received lifecycle message for unknown channel", message); | ||
| return; | ||
| } | ||
| if (isChannelLifecycleMessage(message)) { | ||
| switch (message.type) { | ||
| case "CHANNEL_OPENED": | ||
| return channel.processStatusOpened(); | ||
| case "CHANNEL_CLOSED": | ||
| return channel.processStatusClosed(); | ||
| case "ERROR": | ||
| return channel.processError({ | ||
| type: message.error, | ||
| message: message.message | ||
| }); | ||
| } | ||
| return; | ||
| } | ||
| return channel.processPayloadMessage(message); | ||
| } | ||
| this.logger.warn("Unhandeled message", message.type); | ||
| }; | ||
| this.processSetupMessage = (serverSetup) => { | ||
| this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT); | ||
| if (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; | ||
| if (this.lastSettedAuthToken === void 0) { | ||
| this.setConnectionState(DXLinkConnectionState.CONNECTED); | ||
| } | ||
| } | ||
| const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1e3; | ||
| this.scheduler.schedule( | ||
| () => this.timeoutCheck(timeoutMills), | ||
| timeoutMills, | ||
| DXLWS_SCHEDULER_KEY_TIMEOUT | ||
| ); | ||
| }; | ||
| this.publishError = (error) => { | ||
| this.logger.debug("Publishing error", error); | ||
| if (this.errorListeners.size === 0) { | ||
| this.logger.error("Unhandled dxLink error", error); | ||
| return; | ||
| } | ||
| for (const listener of this.errorListeners) { | ||
| try { | ||
| listener(error); | ||
| } catch (e) { | ||
| this.logger.error("Error listener error", e); | ||
| } | ||
| } | ||
| }; | ||
| this.processAuthStateMessage = ({ state }) => { | ||
| this.logger.debug("Received auth state message", state); | ||
| this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT); | ||
| if (this.isFirstAuthState) { | ||
| this.isFirstAuthState = false; | ||
| } else { | ||
| if (state === "UNAUTHORIZED") { | ||
| this.lastSettedAuthToken = void 0; | ||
| } | ||
| } | ||
| if (state === "AUTHORIZED") { | ||
| this.setConnectionState(DXLinkConnectionState.CONNECTED); | ||
| this.requestActiveChannels(); | ||
| } | ||
| this.setAuthState(DXLinkAuthState[state]); | ||
| }; | ||
| this.requestActiveChannels = () => { | ||
| for (const channel of this.channels.values()) { | ||
| if (channel.getState() === DXLinkChannelState2.CLOSED) { | ||
| this.channels.delete(channel.id); | ||
| continue; | ||
| } | ||
| channel.request(); | ||
| } | ||
| }; | ||
| /** | ||
| * Process transport open event from connector. | ||
| * After transport is opened: | ||
| * - setup message is sent to server | ||
| * - auth message is sent to server if auth token is set | ||
| * - wait for setup message from server | ||
| * - wait for auth state message from server | ||
| */ | ||
| 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(); | ||
| }, | ||
| this.config.actionTimeout * 1e3, | ||
| DXLWS_SCHEDULER_KEY_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(); | ||
| }, | ||
| this.config.actionTimeout * 1e3, | ||
| DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT | ||
| ); | ||
| if (this.lastSettedAuthToken !== void 0) { | ||
| this.sendAuthMessage(this.lastSettedAuthToken); | ||
| } | ||
| }; | ||
| /** | ||
| * Process transport close event from connector. | ||
| * After transport is closed by server: | ||
| * - reconnect if connection is authorized | ||
| * - disconnect if connection is not authorized | ||
| */ | ||
| this.processTransportClose = (reason, error) => { | ||
| this.logger.debug("Connection closed", reason); | ||
| if (error !== void 0) { | ||
| this.publishError({ | ||
| type: "UNKNOWN", | ||
| message: reason | ||
| }); | ||
| } | ||
| if (this.authState === DXLinkAuthState.UNAUTHORIZED) { | ||
| this.lastSettedAuthToken = void 0; | ||
| this.disconnect(); | ||
| return; | ||
| } | ||
| this.reconnect(); | ||
| }; | ||
| this.timeoutCheck = (timeoutMills) => { | ||
| const now = Date.now(); | ||
| const noKeepaliveDuration = now - this.lastReceivedMillis; | ||
| if (noKeepaliveDuration >= timeoutMills) { | ||
| this.sendMessage({ | ||
| type: "ERROR", | ||
| channel: 0, | ||
| error: "TIMEOUT", | ||
| message: "No keepalive received for " + noKeepaliveDuration + "ms" | ||
| }); | ||
| return this.reconnect(); | ||
| } | ||
| const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration); | ||
| this.scheduler.schedule( | ||
| () => this.timeoutCheck(timeoutMills), | ||
| nextTimeout, | ||
| DXLWS_SCHEDULER_KEY_TIMEOUT | ||
| ); | ||
| }; | ||
| this.scheduleKeepalive = () => { | ||
| this.scheduler.schedule( | ||
| () => { | ||
| this.sendMessage({ | ||
| type: "KEEPALIVE", | ||
| channel: 0 | ||
| }); | ||
| this.scheduleKeepalive(); | ||
| }, | ||
| this.config.keepaliveInterval * 1e3, | ||
| DXLWS_SCHEDULER_KEY_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 Logger2(this.constructor.name, this.config.logLevel); | ||
| this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler(); | ||
| } | ||
| }; | ||
| export { | ||
| DXLINK_WS_PROTOCOL_VERSION, | ||
| DXLinkWebSocketClient, | ||
| DefaultDXLinkWebSocketConnector | ||
| }; | ||
| //# sourceMappingURL=index.mjs.map |
| {"version":3,"sources":["../src/client.ts","../src/channel.ts","../src/messages.ts","../src/connector.ts","../src/version.ts"],"sourcesContent":["import {\n DXLinkLogLevel,\n type DXLinkLogger,\n Logger,\n type DXLinkScheduler,\n DefaultDXLinkScheduler,\n type DXLinkConnectionDetails,\n DXLinkConnectionState,\n type DXLinkConnectionStateChangeListener,\n type DXLinkErrorListener,\n DXLinkAuthState,\n DXLinkChannelState,\n type DXLinkAuthStateChangeListener,\n type DXLinkChannel,\n type DXLinkChannelOptions,\n type DXLinkError,\n type DXLinkClient,\n} from '@dxfeed/dxlink-core'\n\nimport { DXLinkWebSocketChannel } from './channel'\nimport type { DXLinkWebSocketClientConfig } from './config'\nimport { type DXLinkWebSocketConnector, DefaultDXLinkWebSocketConnector } from './connector'\nimport {\n type AuthStateMessage,\n type ErrorMessage,\n type DXLinkWebSocketMessage,\n type SetupMessage,\n isChannelLifecycleMessage,\n isChannelMessage,\n isConnectionMessage,\n} from './messages'\nimport { VERSION } from './version'\n\n/**\n * Protocol version that is used by client.\n */\nexport const DXLINK_WS_PROTOCOL_VERSION = '0.1'\n\nconst CLIENT_VERSION = `DXF-JS/${VERSION}`\n\n// Scheduler keys\nconst DXLWS_SCHEDULER_KEY_RECONNECT = 'DXLWS_RECONNECT'\nconst DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT = 'DXLWS_SETUP_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT = 'DXLWS_AUTH_STATE_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_TIMEOUT = 'DXLWS_TIMEOUT'\nconst DXLWS_SCHEDULER_KEY_KEEPALIVE = 'DXLWS_KEEPALIVE'\n\nconst DEFAULT_CONNECTION_DETAILS: DXLinkConnectionDetails = {\n protocolVersion: DXLINK_WS_PROTOCOL_VERSION,\n clientVersion: CLIENT_VERSION,\n}\n\n/**\n * dxLink WebSocket client that can be used to connect to the remote dxLink WebSocket endpoint and open channels to services.\n */\nexport class DXLinkWebSocketClient implements DXLinkClient {\n private readonly config: DXLinkWebSocketClientConfig\n\n private readonly logger: DXLinkLogger\n\n private readonly scheduler: DXLinkScheduler\n\n private connector: DXLinkWebSocketConnector | undefined\n\n private connectionState: DXLinkConnectionState = DXLinkConnectionState.NOT_CONNECTED\n private connectionDetails: DXLinkConnectionDetails = DEFAULT_CONNECTION_DETAILS\n\n private authState: DXLinkAuthState = DXLinkAuthState.UNAUTHORIZED\n\n // Listeners\n private readonly connectionStateChangeListeners = new Set<DXLinkConnectionStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n private readonly authStateChangeListeners = new Set<DXLinkAuthStateChangeListener>()\n\n /**\n * Authorization type that was determined by server behavior during setup phase.\n * This value is used to determine if authorization is required or optional or not defined yet.\n */\n private isFirstAuthState = true\n /**\n * Last setted auth token that will be sent to server after connection is established or re-established.\n */\n private lastSettedAuthToken: string | undefined\n\n // Stats for keepalive\n // TODO: mb move to connector\n private lastReceivedMillis = 0\n private lastSentMillis = 0\n\n /**\n * Count of reconnect attempts since last successful connection.\n */\n private reconnectAttempts = 0\n\n // Channels\n private globalChannelId = 1\n private readonly channels = new Map<number, DXLinkWebSocketChannel>()\n\n /**\n * Create new instance of {@link DXLinkWebSocketClient}.\n * @param config Configuration of the client.\n */\n constructor(config?: Partial<DXLinkWebSocketClientConfig>) {\n this.config = {\n keepaliveInterval: 30,\n keepaliveTimeout: 60,\n acceptKeepaliveTimeout: 60,\n actionTimeout: 10,\n logLevel: DXLinkLogLevel.WARN,\n maxReconnectAttempts: -1,\n connectorFactory: (url) => new DefaultDXLinkWebSocketConnector(url),\n ...config,\n }\n\n this.logger = new Logger(this.constructor.name, this.config.logLevel)\n this.scheduler = this.config.scheduler ?? new DefaultDXLinkScheduler()\n }\n\n connect = (url: string) => {\n // Do nothing if already connected to the same url\n if (this.connector?.getUrl() === url) return\n\n // Disconnect from previous connection if any exists\n this.disconnect()\n\n this.logger.debug('Connecting to', url)\n\n // Immediately set connection state to CONNECTING\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n\n // Create new connector\n this.connector = this.config.connectorFactory(url)\n this.connector.setOpenListener(this.processTransportOpen)\n this.connector.setMessageListener(this.processMessage)\n this.connector.setCloseListener(this.processTransportClose)\n\n // Initiate websocket connection\n this.connector.start()\n }\n\n reconnect = () => {\n if (\n this.config.maxReconnectAttempts >= 0 &&\n this.reconnectAttempts >= this.config.maxReconnectAttempts\n ) {\n this.logger.warn('Max reconnect attempts reached')\n\n this.disconnect()\n return\n }\n\n if (\n this.connectionState === DXLinkConnectionState.NOT_CONNECTED ||\n this.connector === undefined\n )\n return\n\n this.logger.debug('Trying to reconnect', this.connector.getUrl())\n\n this.connector.stop()\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n\n // Increase reconnect attempts counter\n this.reconnectAttempts++\n\n // Update state for connection and channels\n this.setConnectionState(DXLinkConnectionState.CONNECTING)\n for (const channel of this.channels.values()) {\n if (channel.getState() === DXLinkChannelState.CLOSED) continue\n\n if (!channel.reconnect) {\n channel.processStatusClosed()\n continue\n }\n\n channel.processStatusRequested()\n }\n\n // Schedule reconnect attempt after some time\n // Additionally, task will be executed in case when tab is active again\n // coz browser sometimes doesn't run scheduled tasks when tab is inactive\n this.scheduler.schedule(\n () => {\n if (this.connector === undefined) return\n\n // Start new connection attempt\n this.connector.start()\n },\n this.reconnectAttempts * 1000,\n DXLWS_SCHEDULER_KEY_RECONNECT\n )\n }\n\n disconnect = () => {\n if (this.connectionState === DXLinkConnectionState.NOT_CONNECTED) return\n\n this.logger.debug('Disconnecting')\n\n // Destroy connector\n this.connector?.stop()\n this.connector = undefined\n\n // Clear all timeouts\n this.scheduler.clear()\n\n // Set initial state\n this.connectionDetails = DEFAULT_CONNECTION_DETAILS\n this.lastReceivedMillis = 0\n this.lastSentMillis = 0\n this.isFirstAuthState = true\n this.reconnectAttempts = 0\n\n this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED)\n this.setAuthState(DXLinkAuthState.UNAUTHORIZED)\n }\n\n close = () => {\n this.disconnect()\n }\n\n getConnectionDetails = () => this.connectionDetails\n getConnectionState = () => this.connectionState\n getScheduler = () => this.scheduler\n addConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.add(listener)\n removeConnectionStateChangeListener = (listener: DXLinkConnectionStateChangeListener) =>\n this.connectionStateChangeListeners.delete(listener)\n\n setAuthToken = (token: string): void => {\n this.lastSettedAuthToken = token\n\n if (this.connectionState === DXLinkConnectionState.CONNECTED) {\n this.sendAuthMessage(token)\n }\n }\n\n getAuthState = (): DXLinkAuthState => this.authState\n addAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.add(listener)\n removeAuthStateChangeListener = (listener: DXLinkAuthStateChangeListener) =>\n this.authStateChangeListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n openChannel = (\n service: string,\n parameters: Record<string, unknown>,\n options?: DXLinkChannelOptions\n ): DXLinkChannel => {\n const channelId = this.globalChannelId\n this.globalChannelId += 2\n\n const channel = new DXLinkWebSocketChannel(\n channelId,\n service,\n parameters,\n options?.reconnect ?? true,\n this.sendMessage,\n this.config\n )\n\n this.channels.set(channelId, channel)\n\n // Send channel request if connection is already established\n if (\n this.connectionState === DXLinkConnectionState.CONNECTED &&\n this.authState === DXLinkAuthState.AUTHORIZED\n ) {\n channel.request()\n }\n\n return channel\n }\n\n private setConnectionState = (newStatus: DXLinkConnectionState) => {\n const prev = this.connectionState\n if (prev === newStatus) return\n\n this.connectionState = newStatus\n for (const listener of this.connectionStateChangeListeners) {\n listener(newStatus, prev)\n }\n }\n\n private sendMessage = (message: DXLinkWebSocketMessage): void => {\n this.connector?.sendMessage(message)\n\n this.scheduleKeepalive()\n\n // TODO: mb move to connector\n this.lastSentMillis = Date.now()\n }\n\n private sendAuthMessage = (token: string): void => {\n this.logger.debug('Sending auth message')\n\n this.setAuthState(DXLinkAuthState.AUTHORIZING)\n\n this.sendMessage({\n type: 'AUTH',\n channel: 0,\n token,\n })\n }\n\n private setAuthState = (newState: DXLinkAuthState): void => {\n const prev = this.authState\n\n this.authState = newState\n for (const listener of this.authStateChangeListeners) {\n try {\n listener(newState, prev)\n } catch (e) {\n this.logger.error('Auth state listener error', e)\n }\n }\n }\n\n private processMessage = (message: DXLinkWebSocketMessage): void => {\n this.lastReceivedMillis = Date.now()\n\n // Send keepalive message if no messages sent for a while (keepaliveInterval)\n // Because browser sometimes doesn't run scheduled tasks when tab is inactive\n if (this.lastReceivedMillis - this.lastSentMillis >= this.config.keepaliveInterval * 1000) {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n }\n\n // Connection messages are messages that are sent to the channel 0\n if (isConnectionMessage(message)) {\n switch (message.type) {\n case 'SETUP':\n return this.processSetupMessage(message)\n case 'AUTH_STATE':\n return this.processAuthStateMessage(message)\n case 'ERROR':\n return this.publishError({\n type: message.error,\n message: message.message,\n })\n case 'KEEPALIVE':\n // Ignore keepalive messages coz they are used only to maintain connection\n return\n }\n } else if (isChannelMessage(message)) {\n const channel = this.channels.get(message.channel)\n if (channel === undefined) {\n this.logger.warn('Received lifecycle message for unknown channel', message)\n return\n }\n\n if (isChannelLifecycleMessage(message)) {\n switch (message.type) {\n case 'CHANNEL_OPENED':\n return channel.processStatusOpened()\n case 'CHANNEL_CLOSED':\n return channel.processStatusClosed()\n case 'ERROR':\n return channel.processError({\n type: message.error,\n message: message.message,\n })\n }\n return\n }\n\n return channel.processPayloadMessage(message)\n }\n\n this.logger.warn('Unhandeled message', message.type)\n }\n\n private processSetupMessage = (serverSetup: SetupMessage): void => {\n // Clear setup timeout check from connect method\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT)\n\n // Mark connection as connected after first setup message and subsequent ones\n if (\n this.connectionState === DXLinkConnectionState.CONNECTING ||\n this.connectionState === DXLinkConnectionState.CONNECTED\n ) {\n this.connectionDetails = {\n ...this.connectionDetails,\n serverVersion: serverSetup.version,\n clientKeepaliveTimeout: this.config.keepaliveTimeout,\n serverKeepaliveTimeout: serverSetup.keepaliveTimeout,\n }\n\n // Reset reconnect attempts counter after successful connection\n this.reconnectAttempts = 0\n\n if (this.lastSettedAuthToken === undefined) {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n }\n }\n\n // Connection maintance: Setup keepalive timeout check\n const timeoutMills = (serverSetup.keepaliveTimeout ?? 60) * 1000\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n timeoutMills,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private publishError = (error: DXLinkError): void => {\n this.logger.debug('Publishing error', error)\n\n if (this.errorListeners.size === 0) {\n this.logger.error('Unhandled dxLink error', error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error('Error listener error', e)\n }\n }\n }\n\n private processAuthStateMessage = ({ state }: AuthStateMessage): void => {\n this.logger.debug('Received auth state message', state)\n\n // Clear auth state timeout check\n this.scheduler.cancel(DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT)\n\n // Ignore first auth state message because it is sent during connection setup\n if (this.isFirstAuthState) {\n this.isFirstAuthState = false\n } else {\n // Reset auth token if server rejected it\n if (state === 'UNAUTHORIZED') {\n this.lastSettedAuthToken = undefined\n }\n }\n\n // Request active channels if connection is authorized\n if (state === 'AUTHORIZED') {\n this.setConnectionState(DXLinkConnectionState.CONNECTED)\n\n this.requestActiveChannels()\n }\n\n this.setAuthState(DXLinkAuthState[state])\n }\n\n private requestActiveChannels = (): void => {\n for (const channel of this.channels.values()) {\n // clear closed channels\n if (channel.getState() === DXLinkChannelState.CLOSED) {\n this.channels.delete(channel.id)\n continue\n }\n\n channel.request()\n }\n }\n\n /**\n * Process transport open event from connector.\n * After transport is opened:\n * - setup message is sent to server\n * - auth message is sent to server if auth token is set\n * - wait for setup message from server\n * - wait for auth state message from server\n */\n private processTransportOpen = (): void => {\n this.logger.debug('Connection opened')\n\n const setupMessage: SetupMessage = {\n type: 'SETUP',\n channel: 0,\n version: `${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,\n keepaliveTimeout: this.config.keepaliveTimeout,\n acceptKeepaliveTimeout: this.config.acceptKeepaliveTimeout,\n }\n\n // Setup timeout check\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No setup message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no setup message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_SETUP_TIMEOUT\n )\n\n this.sendMessage(setupMessage)\n\n this.scheduler.schedule(\n () => {\n const errorMessage: ErrorMessage = {\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No auth state message received for ' + this.config.actionTimeout + 's',\n }\n\n this.sendMessage(errorMessage)\n\n this.publishError({\n type: errorMessage.error,\n message: `${errorMessage.message} from server`,\n })\n\n // Disconnect if no auth state message received\n this.reconnect()\n },\n this.config.actionTimeout * 1000,\n DXLWS_SCHEDULER_KEY_AUTH_STATE_TIMEOUT\n )\n\n if (this.lastSettedAuthToken !== undefined) {\n this.sendAuthMessage(this.lastSettedAuthToken)\n }\n }\n\n /**\n * Process transport close event from connector.\n * After transport is closed by server:\n * - reconnect if connection is authorized\n * - disconnect if connection is not authorized\n */\n private processTransportClose = (reason: string, error: boolean): void => {\n this.logger.debug('Connection closed', reason)\n\n if (error !== undefined) {\n this.publishError({\n type: 'UNKNOWN',\n message: reason,\n })\n }\n\n if (this.authState === DXLinkAuthState.UNAUTHORIZED) {\n this.lastSettedAuthToken = undefined\n this.disconnect()\n return\n }\n\n this.reconnect()\n }\n\n private timeoutCheck = (timeoutMills: number) => {\n const now = Date.now()\n const noKeepaliveDuration = now - this.lastReceivedMillis\n if (noKeepaliveDuration >= timeoutMills) {\n this.sendMessage({\n type: 'ERROR',\n channel: 0,\n error: 'TIMEOUT',\n message: 'No keepalive received for ' + noKeepaliveDuration + 'ms',\n })\n\n return this.reconnect()\n }\n\n const nextTimeout = Math.max(200, timeoutMills - noKeepaliveDuration)\n this.scheduler.schedule(\n () => this.timeoutCheck(timeoutMills),\n nextTimeout,\n DXLWS_SCHEDULER_KEY_TIMEOUT\n )\n }\n\n private scheduleKeepalive = () => {\n this.scheduler.schedule(\n () => {\n this.sendMessage({\n type: 'KEEPALIVE',\n channel: 0,\n })\n\n this.scheduleKeepalive()\n },\n this.config.keepaliveInterval * 1000,\n DXLWS_SCHEDULER_KEY_KEEPALIVE\n )\n }\n}\n","import {\n type DXLinkLogger,\n Logger,\n type DXLinkChannel,\n DXLinkChannelState,\n type DXLinkChannelMessageListener,\n type DXLinkChannelStateChangeListener,\n type DXLinkErrorListener,\n type DXLinkChannelMessage,\n type DXLinkError,\n} from '@dxfeed/dxlink-core'\n\nimport { type DXLinkWebSocketClientConfig } from './config'\nimport { type ChannelPayloadMessage, type DXLinkWebSocketMessage } from './messages'\n\n/**\n * A DXLink channel implementation.\n * @internal\n */\nexport class DXLinkWebSocketChannel implements DXLinkChannel {\n private status = DXLinkChannelState.REQUESTED\n\n // Listeners\n private readonly messageListeners = new Set<DXLinkChannelMessageListener>()\n private readonly statusListeners = new Set<DXLinkChannelStateChangeListener>()\n private readonly errorListeners = new Set<DXLinkErrorListener>()\n\n private logger: DXLinkLogger\n\n constructor(\n public readonly id: number,\n public readonly service: string,\n public readonly parameters: Record<string, unknown>,\n public readonly reconnect: boolean,\n private readonly sendMessage: (message: DXLinkWebSocketMessage) => void,\n config: DXLinkWebSocketClientConfig\n ) {\n this.logger = new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`, config.logLevel)\n }\n\n send = ({ type, ...payload }: DXLinkChannelMessage) => {\n if (this.status !== DXLinkChannelState.OPENED) {\n throw new Error('Channel is not ready')\n }\n\n this.sendMessage({\n type,\n channel: this.id,\n ...payload,\n })\n }\n\n addMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.add(listener)\n removeMessageListener = (listener: DXLinkChannelMessageListener) =>\n this.messageListeners.delete(listener)\n\n getState = () => this.status\n addStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.add(listener)\n removeStateChangeListener = (listener: DXLinkChannelStateChangeListener) =>\n this.statusListeners.delete(listener)\n\n addErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.add(listener)\n removeErrorListener = (listener: DXLinkErrorListener) => this.errorListeners.delete(listener)\n\n error = ({ type, message }: DXLinkError) =>\n this.send({\n type: 'ERROR',\n error: type,\n message,\n })\n\n close = () => {\n if (this.status === DXLinkChannelState.CLOSED) return\n\n this.logger.debug(`Closing by user`)\n\n // We can think that channel is closed already\n this.setStatus(DXLinkChannelState.CLOSED)\n\n this.clear()\n\n this.sendMessage({\n type: 'CHANNEL_CANCEL',\n channel: this.id,\n })\n }\n\n request = () => {\n this.logger.debug('Requesting')\n\n this.sendMessage({\n type: 'CHANNEL_REQUEST',\n channel: this.id,\n service: this.service,\n parameters: this.parameters,\n })\n\n this.processStatusRequested()\n }\n\n processPayloadMessage = (message: ChannelPayloadMessage) => {\n for (const listener of this.messageListeners) {\n listener(message)\n }\n }\n\n processStatusOpened = () => {\n this.logger.debug('Opened')\n\n this.setStatus(DXLinkChannelState.OPENED)\n }\n\n processStatusRequested = () => {\n this.setStatus(DXLinkChannelState.REQUESTED)\n }\n\n processStatusClosed = () => {\n this.logger.debug('Closed by remote endpoint')\n\n this.setStatus(DXLinkChannelState.CLOSED)\n this.clear()\n }\n\n processError = (error: DXLinkError) => {\n if (this.errorListeners.size === 0) {\n this.logger.error(`Unhandled error in channel#${this.id}: `, error)\n return\n }\n\n for (const listener of this.errorListeners) {\n try {\n listener(error)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} error listener: `, e)\n }\n }\n }\n\n private setStatus = (newStatus: DXLinkChannelState) => {\n if (this.status === newStatus) return\n\n const prev = this.status\n this.status = newStatus\n for (const listener of this.statusListeners) {\n try {\n listener(newStatus, prev)\n } catch (e) {\n this.logger.error(`Error in channel#${this.id} status listener: `, e)\n }\n }\n }\n\n private clear = () => {\n this.messageListeners.clear()\n this.statusListeners.clear()\n // TODO: rethink approach when error came after channel is closed\n // this.errorListeners.clear()\n }\n}\n","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","import { type DXLinkWebSocketMessage } from './messages'\n\n/**\n * Interface for a WebSocket connector that manages the connection to a WebSocket server.\n * It provides methods to start and stop the connection, send messages, and set listeners for open, close, and message events.\n */\nexport interface DXLinkWebSocketConnector {\n /**\n * Returns the URL of the WebSocket connection.\n */\n getUrl(): string\n /**\n * Starts the WebSocket connection.\n */\n start(): void\n /**\n * Stops the WebSocket connection and cleans up resources.\n */\n stop(): void\n /**\n * Sends a message over the WebSocket connection.\n * @param message The message to send.\n */\n sendMessage(message: DXLinkWebSocketMessage): void\n /**\n * Sets a listener that is called when the WebSocket connection is opened.\n * @param listener The listener function to call when the connection is opened.\n */\n setOpenListener(listener: () => void): void\n /**\n * Sets a listener that is called when the WebSocket connection is closed.\n * @param listener The listener function to call when the connection is closed.\n */\n setCloseListener(listener: DXLinkWebSocketCloseListener): void\n /**\n * Sets a listener that is called when a message is received from the WebSocket server.\n * @param listener The listener function to call when a message is received.\n */\n setMessageListener(listener: (message: DXLinkWebSocketMessage) => void): void\n}\n\n/**\n * Type for a close listener that is called when the WebSocket connection is closed.\n * @param reason - The reason for the closure.\n * @param error - Indicates if the closure was due to an error.\n */\nexport type DXLinkWebSocketCloseListener = (reason: string, error: boolean) => void\n\n/**\n * Default connector for the WebSocket connection.\n * @internal\n */\nexport class DefaultDXLinkWebSocketConnector implements DXLinkWebSocketConnector {\n private socket: WebSocket | undefined = undefined\n\n private isAvailable = false\n\n private openListener: (() => void) | undefined = undefined\n private closeListener: DXLinkWebSocketCloseListener | undefined = undefined\n private messageListener: ((message: DXLinkWebSocketMessage) => void) | undefined = undefined\n\n constructor(\n private readonly url: string,\n private readonly protocols?: string | string[]\n ) {}\n\n start() {\n if (this.socket !== undefined) return\n\n this.socket = new WebSocket(this.url, this.protocols)\n\n this.socket.addEventListener('open', this.handleOpen)\n this.socket.addEventListener('error', this.handleError)\n this.socket.addEventListener('close', this.handleClosed)\n }\n\n stop() {\n if (this.socket === undefined) return\n\n this.socket.removeEventListener('open', this.handleOpen)\n this.socket.removeEventListener('error', this.handleError)\n this.socket.removeEventListener('close', this.handleClosed)\n this.socket.removeEventListener('message', this.handleMessage)\n\n this.socket.close()\n this.socket = undefined\n this.isAvailable = false\n }\n\n sendMessage = (message: DXLinkWebSocketMessage) => {\n if (this.socket === undefined || !this.isAvailable) {\n return\n }\n\n this.socket.send(JSON.stringify(message))\n }\n\n setOpenListener = (listener: () => void) => {\n this.openListener = listener\n }\n\n setCloseListener = (listener: DXLinkWebSocketCloseListener) => {\n this.closeListener = listener\n }\n\n setMessageListener = (listener: (message: DXLinkWebSocketMessage) => void) => {\n this.messageListener = listener\n }\n\n getUrl = () => this.url\n\n private handleOpen = () => {\n if (this.socket === undefined) return\n\n this.isAvailable = true\n\n this.socket.removeEventListener('open', this.handleOpen)\n\n this.socket.addEventListener('message', this.handleMessage)\n\n this.openListener?.()\n }\n\n private handleClosed = (ev: CloseEvent) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.(ev.reason, false)\n }\n\n private handleError = (_ev: Event) => {\n if (this.socket === undefined) return\n\n this.stop()\n\n this.closeListener?.('Unable to connect', true)\n }\n\n private handleMessage = (ev: MessageEvent) => {\n try {\n const message = JSON.parse(ev.data)\n if (typeof message !== 'object') {\n throw new Error('Unexpected message: ' + typeof message)\n }\n if (typeof message.type !== 'string') {\n throw new Error('Unexpected message type: ' + typeof message.type)\n }\n if (typeof message.channel !== 'number') {\n throw new Error('Unexpected message channel: ' + typeof message.channel)\n }\n\n // TODO: validate message properly (e.g. check for type and other fields)\n\n if (this.messageListener === undefined) {\n return console.warn('No message listener set')\n }\n\n this.messageListener(message)\n } catch (error) {\n console.error(error instanceof Error ? error : new Error('Parsing error:' + String(error)))\n }\n }\n}\n","/**\n * Version of the dxlink-websocket-client package.\n * `__DXLINK_VERSION__` is replaced with the package version at build time via\n * tsup's `define` (see tsup.config.ts), so it is inlined into the bundle as a\n * string literal. In dev/test runs the constant is undefined and we fall back\n * to 'local-unknown'; `typeof` keeps that safe even though the global is only\n * declared, not defined.\n * @internal\n */\ndeclare const __DXLINK_VERSION__: string | undefined\n\nexport const VERSION =\n typeof __DXLINK_VERSION__ === 'undefined' ? 'local-unknown' : __DXLINK_VERSION__\n"],"mappings":";AAAA;AAAA,EACE;AAAA,EAEA,UAAAA;AAAA,EAEA;AAAA,EAEA;AAAA,EAGA;AAAA,EACA,sBAAAC;AAAA,OAMK;;;ACjBP;AAAA,EAEE;AAAA,EAEA;AAAA,OAMK;;;ACyFA,IAAM,sBAAsB,CACjC,YAEA,QAAQ,YAAY,MACnB,QAAQ,SAAS,WAChB,QAAQ,SAAS,eACjB,QAAQ,SAAS,UACjB,QAAQ,SAAS,gBACjB,QAAQ,SAAS;AAWd,IAAM,mBAAmB,CAAC,YAC/B,QAAQ,YAAY;AAEf,IAAM,4BAA4B,CACvC,YAEA,QAAQ,SAAS,oBACjB,QAAQ,SAAS,oBACjB,QAAQ,SAAS,WACjB,QAAQ,SAAS,qBACjB,QAAQ,SAAS;;;AD7GZ,IAAM,yBAAN,MAAM,wBAAgD;AAAA,EAU3D,YACkB,IACA,SACA,YACA,WACC,aACjB,QACA;AANgB;AACA;AACA;AACA;AACC;AAdnB,SAAQ,SAAS,mBAAmB;AAGpC;AAAA,SAAiB,mBAAmB,oBAAI,IAAkC;AAC1E,SAAiB,kBAAkB,oBAAI,IAAsC;AAC7E,SAAiB,iBAAiB,oBAAI,IAAyB;AAe/D,gBAAO,CAAC,EAAE,MAAM,GAAG,QAAQ,MAA4B;AACrD,UAAI,KAAK,WAAW,mBAAmB,QAAQ;AAC7C,cAAM,IAAI,MAAM,sBAAsB;AAAA,MACxC;AAEA,WAAK,YAAY;AAAA,QACf;AAAA,QACA,SAAS,KAAK;AAAA,QACd,GAAG;AAAA,MACL,CAAC;AAAA,IACH;AAEA,8BAAqB,CAAC,aACpB,KAAK,iBAAiB,IAAI,QAAQ;AACpC,iCAAwB,CAAC,aACvB,KAAK,iBAAiB,OAAO,QAAQ;AAEvC,oBAAW,MAAM,KAAK;AACtB,kCAAyB,CAAC,aACxB,KAAK,gBAAgB,IAAI,QAAQ;AACnC,qCAA4B,CAAC,aAC3B,KAAK,gBAAgB,OAAO,QAAQ;AAEtC,4BAAmB,CAAC,aAAkC,KAAK,eAAe,IAAI,QAAQ;AACtF,+BAAsB,CAAC,aAAkC,KAAK,eAAe,OAAO,QAAQ;AAE5F,iBAAQ,CAAC,EAAE,MAAM,QAAQ,MACvB,KAAK,KAAK;AAAA,MACR,MAAM;AAAA,MACN,OAAO;AAAA,MACP;AAAA,IACF,CAAC;AAEH,iBAAQ,MAAM;AACZ,UAAI,KAAK,WAAW,mBAAmB,OAAQ;AAE/C,WAAK,OAAO,MAAM,iBAAiB;AAGnC,WAAK,UAAU,mBAAmB,MAAM;AAExC,WAAK,MAAM;AAEX,WAAK,YAAY;AAAA,QACf,MAAM;AAAA,QACN,SAAS,KAAK;AAAA,MAChB,CAAC;AAAA,IACH;AAEA,mBAAU,MAAM;AACd,WAAK,OAAO,MAAM,YAAY;AAE9B,WAAK,YAAY;AAAA,QACf,MAAM;AAAA,QACN,SAAS,KAAK;AAAA,QACd,SAAS,KAAK;AAAA,QACd,YAAY,KAAK;AAAA,MACnB,CAAC;AAED,WAAK,uBAAuB;AAAA,IAC9B;AAEA,iCAAwB,CAAC,YAAmC;AAC1D,iBAAW,YAAY,KAAK,kBAAkB;AAC5C,iBAAS,OAAO;AAAA,MAClB;AAAA,IACF;AAEA,+BAAsB,MAAM;AAC1B,WAAK,OAAO,MAAM,QAAQ;AAE1B,WAAK,UAAU,mBAAmB,MAAM;AAAA,IAC1C;AAEA,kCAAyB,MAAM;AAC7B,WAAK,UAAU,mBAAmB,SAAS;AAAA,IAC7C;AAEA,+BAAsB,MAAM;AAC1B,WAAK,OAAO,MAAM,2BAA2B;AAE7C,WAAK,UAAU,mBAAmB,MAAM;AACxC,WAAK,MAAM;AAAA,IACb;AAEA,wBAAe,CAAC,UAAuB;AACrC,UAAI,KAAK,eAAe,SAAS,GAAG;AAClC,aAAK,OAAO,MAAM,8BAA8B,KAAK,EAAE,MAAM,KAAK;AAClE;AAAA,MACF;AAEA,iBAAW,YAAY,KAAK,gBAAgB;AAC1C,YAAI;AACF,mBAAS,KAAK;AAAA,QAChB,SAAS,GAAG;AACV,eAAK,OAAO,MAAM,oBAAoB,KAAK,EAAE,qBAAqB,CAAC;AAAA,QACrE;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,YAAY,CAAC,cAAkC;AACrD,UAAI,KAAK,WAAW,UAAW;AAE/B,YAAM,OAAO,KAAK;AAClB,WAAK,SAAS;AACd,iBAAW,YAAY,KAAK,iBAAiB;AAC3C,YAAI;AACF,mBAAS,WAAW,IAAI;AAAA,QAC1B,SAAS,GAAG;AACV,eAAK,OAAO,MAAM,oBAAoB,KAAK,EAAE,sBAAsB,CAAC;AAAA,QACtE;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,QAAQ,MAAM;AACpB,WAAK,iBAAiB,MAAM;AAC5B,WAAK,gBAAgB,MAAM;AAAA,IAG7B;AA1HE,SAAK,SAAS,IAAI,OAAO,GAAG,wBAAuB,IAAI,IAAI,EAAE,IAAI,OAAO,IAAI,OAAO,QAAQ;AAAA,EAC7F;AA0HF;;;AE5GO,IAAM,kCAAN,MAA0E;AAAA,EAS/E,YACmB,KACA,WACjB;AAFiB;AACA;AAVnB,SAAQ,SAAgC;AAExC,SAAQ,cAAc;AAEtB,SAAQ,eAAyC;AACjD,SAAQ,gBAA0D;AAClE,SAAQ,kBAA2E;AA8BnF,uBAAc,CAAC,YAAoC;AACjD,UAAI,KAAK,WAAW,UAAa,CAAC,KAAK,aAAa;AAClD;AAAA,MACF;AAEA,WAAK,OAAO,KAAK,KAAK,UAAU,OAAO,CAAC;AAAA,IAC1C;AAEA,2BAAkB,CAAC,aAAyB;AAC1C,WAAK,eAAe;AAAA,IACtB;AAEA,4BAAmB,CAAC,aAA2C;AAC7D,WAAK,gBAAgB;AAAA,IACvB;AAEA,8BAAqB,CAAC,aAAwD;AAC5E,WAAK,kBAAkB;AAAA,IACzB;AAEA,kBAAS,MAAM,KAAK;AAEpB,SAAQ,aAAa,MAAM;AACzB,UAAI,KAAK,WAAW,OAAW;AAE/B,WAAK,cAAc;AAEnB,WAAK,OAAO,oBAAoB,QAAQ,KAAK,UAAU;AAEvD,WAAK,OAAO,iBAAiB,WAAW,KAAK,aAAa;AAE1D,WAAK,eAAe;AAAA,IACtB;AAEA,SAAQ,eAAe,CAAC,OAAmB;AACzC,UAAI,KAAK,WAAW,OAAW;AAE/B,WAAK,KAAK;AAEV,WAAK,gBAAgB,GAAG,QAAQ,KAAK;AAAA,IACvC;AAEA,SAAQ,cAAc,CAAC,QAAe;AACpC,UAAI,KAAK,WAAW,OAAW;AAE/B,WAAK,KAAK;AAEV,WAAK,gBAAgB,qBAAqB,IAAI;AAAA,IAChD;AAEA,SAAQ,gBAAgB,CAAC,OAAqB;AAC5C,UAAI;AACF,cAAM,UAAU,KAAK,MAAM,GAAG,IAAI;AAClC,YAAI,OAAO,YAAY,UAAU;AAC/B,gBAAM,IAAI,MAAM,yBAAyB,OAAO,OAAO;AAAA,QACzD;AACA,YAAI,OAAO,QAAQ,SAAS,UAAU;AACpC,gBAAM,IAAI,MAAM,8BAA8B,OAAO,QAAQ,IAAI;AAAA,QACnE;AACA,YAAI,OAAO,QAAQ,YAAY,UAAU;AACvC,gBAAM,IAAI,MAAM,iCAAiC,OAAO,QAAQ,OAAO;AAAA,QACzE;AAIA,YAAI,KAAK,oBAAoB,QAAW;AACtC,iBAAO,QAAQ,KAAK,yBAAyB;AAAA,QAC/C;AAEA,aAAK,gBAAgB,OAAO;AAAA,MAC9B,SAAS,OAAO;AACd,gBAAQ,MAAM,iBAAiB,QAAQ,QAAQ,IAAI,MAAM,mBAAmB,OAAO,KAAK,CAAC,CAAC;AAAA,MAC5F;AAAA,IACF;AAAA,EAlGG;AAAA,EAEH,QAAQ;AACN,QAAI,KAAK,WAAW,OAAW;AAE/B,SAAK,SAAS,IAAI,UAAU,KAAK,KAAK,KAAK,SAAS;AAEpD,SAAK,OAAO,iBAAiB,QAAQ,KAAK,UAAU;AACpD,SAAK,OAAO,iBAAiB,SAAS,KAAK,WAAW;AACtD,SAAK,OAAO,iBAAiB,SAAS,KAAK,YAAY;AAAA,EACzD;AAAA,EAEA,OAAO;AACL,QAAI,KAAK,WAAW,OAAW;AAE/B,SAAK,OAAO,oBAAoB,QAAQ,KAAK,UAAU;AACvD,SAAK,OAAO,oBAAoB,SAAS,KAAK,WAAW;AACzD,SAAK,OAAO,oBAAoB,SAAS,KAAK,YAAY;AAC1D,SAAK,OAAO,oBAAoB,WAAW,KAAK,aAAa;AAE7D,SAAK,OAAO,MAAM;AAClB,SAAK,SAAS;AACd,SAAK,cAAc;AAAA,EACrB;AA4EF;;;ACxJO,IAAM,UACX,QAA4C,kBAAkB;;;AJwBzD,IAAM,6BAA6B;AAE1C,IAAM,iBAAiB,UAAU,OAAO;AAGxC,IAAM,gCAAgC;AACtC,IAAM,oCAAoC;AAC1C,IAAM,yCAAyC;AAC/C,IAAM,8BAA8B;AACpC,IAAM,gCAAgC;AAEtC,IAAM,6BAAsD;AAAA,EAC1D,iBAAiB;AAAA,EACjB,eAAe;AACjB;AAKO,IAAM,wBAAN,MAAoD;AAAA;AAAA;AAAA;AAAA;AAAA,EA+CzD,YAAY,QAA+C;AAtC3D,SAAQ,kBAAyC,sBAAsB;AACvE,SAAQ,oBAA6C;AAErD,SAAQ,YAA6B,gBAAgB;AAGrD;AAAA,SAAiB,iCAAiC,oBAAI,IAAyC;AAC/F,SAAiB,iBAAiB,oBAAI,IAAyB;AAC/D,SAAiB,2BAA2B,oBAAI,IAAmC;AAMnF;AAAA;AAAA;AAAA;AAAA,SAAQ,mBAAmB;AAQ3B;AAAA;AAAA,SAAQ,qBAAqB;AAC7B,SAAQ,iBAAiB;AAKzB;AAAA;AAAA;AAAA,SAAQ,oBAAoB;AAG5B;AAAA,SAAQ,kBAAkB;AAC1B,SAAiB,WAAW,oBAAI,IAAoC;AAsBpE,mBAAU,CAAC,QAAgB;AAEzB,UAAI,KAAK,WAAW,OAAO,MAAM,IAAK;AAGtC,WAAK,WAAW;AAEhB,WAAK,OAAO,MAAM,iBAAiB,GAAG;AAGtC,WAAK,mBAAmB,sBAAsB,UAAU;AAGxD,WAAK,YAAY,KAAK,OAAO,iBAAiB,GAAG;AACjD,WAAK,UAAU,gBAAgB,KAAK,oBAAoB;AACxD,WAAK,UAAU,mBAAmB,KAAK,cAAc;AACrD,WAAK,UAAU,iBAAiB,KAAK,qBAAqB;AAG1D,WAAK,UAAU,MAAM;AAAA,IACvB;AAEA,qBAAY,MAAM;AAChB,UACE,KAAK,OAAO,wBAAwB,KACpC,KAAK,qBAAqB,KAAK,OAAO,sBACtC;AACA,aAAK,OAAO,KAAK,gCAAgC;AAEjD,aAAK,WAAW;AAChB;AAAA,MACF;AAEA,UACE,KAAK,oBAAoB,sBAAsB,iBAC/C,KAAK,cAAc;AAEnB;AAEF,WAAK,OAAO,MAAM,uBAAuB,KAAK,UAAU,OAAO,CAAC;AAEhE,WAAK,UAAU,KAAK;AAGpB,WAAK,UAAU,MAAM;AAGrB,WAAK,oBAAoB;AACzB,WAAK,qBAAqB;AAC1B,WAAK,iBAAiB;AACtB,WAAK,mBAAmB;AAGxB,WAAK;AAGL,WAAK,mBAAmB,sBAAsB,UAAU;AACxD,iBAAW,WAAW,KAAK,SAAS,OAAO,GAAG;AAC5C,YAAI,QAAQ,SAAS,MAAMC,oBAAmB,OAAQ;AAEtD,YAAI,CAAC,QAAQ,WAAW;AACtB,kBAAQ,oBAAoB;AAC5B;AAAA,QACF;AAEA,gBAAQ,uBAAuB;AAAA,MACjC;AAKA,WAAK,UAAU;AAAA,QACb,MAAM;AACJ,cAAI,KAAK,cAAc,OAAW;AAGlC,eAAK,UAAU,MAAM;AAAA,QACvB;AAAA,QACA,KAAK,oBAAoB;AAAA,QACzB;AAAA,MACF;AAAA,IACF;AAEA,sBAAa,MAAM;AACjB,UAAI,KAAK,oBAAoB,sBAAsB,cAAe;AAElE,WAAK,OAAO,MAAM,eAAe;AAGjC,WAAK,WAAW,KAAK;AACrB,WAAK,YAAY;AAGjB,WAAK,UAAU,MAAM;AAGrB,WAAK,oBAAoB;AACzB,WAAK,qBAAqB;AAC1B,WAAK,iBAAiB;AACtB,WAAK,mBAAmB;AACxB,WAAK,oBAAoB;AAEzB,WAAK,mBAAmB,sBAAsB,aAAa;AAC3D,WAAK,aAAa,gBAAgB,YAAY;AAAA,IAChD;AAEA,iBAAQ,MAAM;AACZ,WAAK,WAAW;AAAA,IAClB;AAEA,gCAAuB,MAAM,KAAK;AAClC,8BAAqB,MAAM,KAAK;AAChC,wBAAe,MAAM,KAAK;AAC1B,4CAAmC,CAAC,aAClC,KAAK,+BAA+B,IAAI,QAAQ;AAClD,+CAAsC,CAAC,aACrC,KAAK,+BAA+B,OAAO,QAAQ;AAErD,wBAAe,CAAC,UAAwB;AACtC,WAAK,sBAAsB;AAE3B,UAAI,KAAK,oBAAoB,sBAAsB,WAAW;AAC5D,aAAK,gBAAgB,KAAK;AAAA,MAC5B;AAAA,IACF;AAEA,wBAAe,MAAuB,KAAK;AAC3C,sCAA6B,CAAC,aAC5B,KAAK,yBAAyB,IAAI,QAAQ;AAC5C,yCAAgC,CAAC,aAC/B,KAAK,yBAAyB,OAAO,QAAQ;AAE/C,4BAAmB,CAAC,aAAkC,KAAK,eAAe,IAAI,QAAQ;AACtF,+BAAsB,CAAC,aAAkC,KAAK,eAAe,OAAO,QAAQ;AAE5F,uBAAc,CACZ,SACA,YACA,YACkB;AAClB,YAAM,YAAY,KAAK;AACvB,WAAK,mBAAmB;AAExB,YAAM,UAAU,IAAI;AAAA,QAClB;AAAA,QACA;AAAA,QACA;AAAA,QACA,SAAS,aAAa;AAAA,QACtB,KAAK;AAAA,QACL,KAAK;AAAA,MACP;AAEA,WAAK,SAAS,IAAI,WAAW,OAAO;AAGpC,UACE,KAAK,oBAAoB,sBAAsB,aAC/C,KAAK,cAAc,gBAAgB,YACnC;AACA,gBAAQ,QAAQ;AAAA,MAClB;AAEA,aAAO;AAAA,IACT;AAEA,SAAQ,qBAAqB,CAAC,cAAqC;AACjE,YAAM,OAAO,KAAK;AAClB,UAAI,SAAS,UAAW;AAExB,WAAK,kBAAkB;AACvB,iBAAW,YAAY,KAAK,gCAAgC;AAC1D,iBAAS,WAAW,IAAI;AAAA,MAC1B;AAAA,IACF;AAEA,SAAQ,cAAc,CAAC,YAA0C;AAC/D,WAAK,WAAW,YAAY,OAAO;AAEnC,WAAK,kBAAkB;AAGvB,WAAK,iBAAiB,KAAK,IAAI;AAAA,IACjC;AAEA,SAAQ,kBAAkB,CAAC,UAAwB;AACjD,WAAK,OAAO,MAAM,sBAAsB;AAExC,WAAK,aAAa,gBAAgB,WAAW;AAE7C,WAAK,YAAY;AAAA,QACf,MAAM;AAAA,QACN,SAAS;AAAA,QACT;AAAA,MACF,CAAC;AAAA,IACH;AAEA,SAAQ,eAAe,CAAC,aAAoC;AAC1D,YAAM,OAAO,KAAK;AAElB,WAAK,YAAY;AACjB,iBAAW,YAAY,KAAK,0BAA0B;AACpD,YAAI;AACF,mBAAS,UAAU,IAAI;AAAA,QACzB,SAAS,GAAG;AACV,eAAK,OAAO,MAAM,6BAA6B,CAAC;AAAA,QAClD;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,iBAAiB,CAAC,YAA0C;AAClE,WAAK,qBAAqB,KAAK,IAAI;AAInC,UAAI,KAAK,qBAAqB,KAAK,kBAAkB,KAAK,OAAO,oBAAoB,KAAM;AACzF,aAAK,YAAY;AAAA,UACf,MAAM;AAAA,UACN,SAAS;AAAA,QACX,CAAC;AAAA,MACH;AAGA,UAAI,oBAAoB,OAAO,GAAG;AAChC,gBAAQ,QAAQ,MAAM;AAAA,UACpB,KAAK;AACH,mBAAO,KAAK,oBAAoB,OAAO;AAAA,UACzC,KAAK;AACH,mBAAO,KAAK,wBAAwB,OAAO;AAAA,UAC7C,KAAK;AACH,mBAAO,KAAK,aAAa;AAAA,cACvB,MAAM,QAAQ;AAAA,cACd,SAAS,QAAQ;AAAA,YACnB,CAAC;AAAA,UACH,KAAK;AAEH;AAAA,QACJ;AAAA,MACF,WAAW,iBAAiB,OAAO,GAAG;AACpC,cAAM,UAAU,KAAK,SAAS,IAAI,QAAQ,OAAO;AACjD,YAAI,YAAY,QAAW;AACzB,eAAK,OAAO,KAAK,kDAAkD,OAAO;AAC1E;AAAA,QACF;AAEA,YAAI,0BAA0B,OAAO,GAAG;AACtC,kBAAQ,QAAQ,MAAM;AAAA,YACpB,KAAK;AACH,qBAAO,QAAQ,oBAAoB;AAAA,YACrC,KAAK;AACH,qBAAO,QAAQ,oBAAoB;AAAA,YACrC,KAAK;AACH,qBAAO,QAAQ,aAAa;AAAA,gBAC1B,MAAM,QAAQ;AAAA,gBACd,SAAS,QAAQ;AAAA,cACnB,CAAC;AAAA,UACL;AACA;AAAA,QACF;AAEA,eAAO,QAAQ,sBAAsB,OAAO;AAAA,MAC9C;AAEA,WAAK,OAAO,KAAK,sBAAsB,QAAQ,IAAI;AAAA,IACrD;AAEA,SAAQ,sBAAsB,CAAC,gBAAoC;AAEjE,WAAK,UAAU,OAAO,iCAAiC;AAGvD,UACE,KAAK,oBAAoB,sBAAsB,cAC/C,KAAK,oBAAoB,sBAAsB,WAC/C;AACA,aAAK,oBAAoB;AAAA,UACvB,GAAG,KAAK;AAAA,UACR,eAAe,YAAY;AAAA,UAC3B,wBAAwB,KAAK,OAAO;AAAA,UACpC,wBAAwB,YAAY;AAAA,QACtC;AAGA,aAAK,oBAAoB;AAEzB,YAAI,KAAK,wBAAwB,QAAW;AAC1C,eAAK,mBAAmB,sBAAsB,SAAS;AAAA,QACzD;AAAA,MACF;AAGA,YAAM,gBAAgB,YAAY,oBAAoB,MAAM;AAC5D,WAAK,UAAU;AAAA,QACb,MAAM,KAAK,aAAa,YAAY;AAAA,QACpC;AAAA,QACA;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,eAAe,CAAC,UAA6B;AACnD,WAAK,OAAO,MAAM,oBAAoB,KAAK;AAE3C,UAAI,KAAK,eAAe,SAAS,GAAG;AAClC,aAAK,OAAO,MAAM,0BAA0B,KAAK;AACjD;AAAA,MACF;AAEA,iBAAW,YAAY,KAAK,gBAAgB;AAC1C,YAAI;AACF,mBAAS,KAAK;AAAA,QAChB,SAAS,GAAG;AACV,eAAK,OAAO,MAAM,wBAAwB,CAAC;AAAA,QAC7C;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,0BAA0B,CAAC,EAAE,MAAM,MAA8B;AACvE,WAAK,OAAO,MAAM,+BAA+B,KAAK;AAGtD,WAAK,UAAU,OAAO,sCAAsC;AAG5D,UAAI,KAAK,kBAAkB;AACzB,aAAK,mBAAmB;AAAA,MAC1B,OAAO;AAEL,YAAI,UAAU,gBAAgB;AAC5B,eAAK,sBAAsB;AAAA,QAC7B;AAAA,MACF;AAGA,UAAI,UAAU,cAAc;AAC1B,aAAK,mBAAmB,sBAAsB,SAAS;AAEvD,aAAK,sBAAsB;AAAA,MAC7B;AAEA,WAAK,aAAa,gBAAgB,KAAK,CAAC;AAAA,IAC1C;AAEA,SAAQ,wBAAwB,MAAY;AAC1C,iBAAW,WAAW,KAAK,SAAS,OAAO,GAAG;AAE5C,YAAI,QAAQ,SAAS,MAAMA,oBAAmB,QAAQ;AACpD,eAAK,SAAS,OAAO,QAAQ,EAAE;AAC/B;AAAA,QACF;AAEA,gBAAQ,QAAQ;AAAA,MAClB;AAAA,IACF;AAUA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,SAAQ,uBAAuB,MAAY;AACzC,WAAK,OAAO,MAAM,mBAAmB;AAErC,YAAM,eAA6B;AAAA,QACjC,MAAM;AAAA,QACN,SAAS;AAAA,QACT,SAAS,GAAG,KAAK,kBAAkB,eAAe,IAAI,KAAK,kBAAkB,aAAa;AAAA,QAC1F,kBAAkB,KAAK,OAAO;AAAA,QAC9B,wBAAwB,KAAK,OAAO;AAAA,MACtC;AAGA,WAAK,UAAU;AAAA,QACb,MAAM;AACJ,gBAAM,eAA6B;AAAA,YACjC,MAAM;AAAA,YACN,SAAS;AAAA,YACT,OAAO;AAAA,YACP,SAAS,mCAAmC,KAAK,OAAO,gBAAgB;AAAA,UAC1E;AAEA,eAAK,YAAY,YAAY;AAE7B,eAAK,aAAa;AAAA,YAChB,MAAM,aAAa;AAAA,YACnB,SAAS,GAAG,aAAa,OAAO;AAAA,UAClC,CAAC;AAGD,eAAK,UAAU;AAAA,QACjB;AAAA,QACA,KAAK,OAAO,gBAAgB;AAAA,QAC5B;AAAA,MACF;AAEA,WAAK,YAAY,YAAY;AAE7B,WAAK,UAAU;AAAA,QACb,MAAM;AACJ,gBAAM,eAA6B;AAAA,YACjC,MAAM;AAAA,YACN,SAAS;AAAA,YACT,OAAO;AAAA,YACP,SAAS,wCAAwC,KAAK,OAAO,gBAAgB;AAAA,UAC/E;AAEA,eAAK,YAAY,YAAY;AAE7B,eAAK,aAAa;AAAA,YAChB,MAAM,aAAa;AAAA,YACnB,SAAS,GAAG,aAAa,OAAO;AAAA,UAClC,CAAC;AAGD,eAAK,UAAU;AAAA,QACjB;AAAA,QACA,KAAK,OAAO,gBAAgB;AAAA,QAC5B;AAAA,MACF;AAEA,UAAI,KAAK,wBAAwB,QAAW;AAC1C,aAAK,gBAAgB,KAAK,mBAAmB;AAAA,MAC/C;AAAA,IACF;AAQA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA,SAAQ,wBAAwB,CAAC,QAAgB,UAAyB;AACxE,WAAK,OAAO,MAAM,qBAAqB,MAAM;AAE7C,UAAI,UAAU,QAAW;AACvB,aAAK,aAAa;AAAA,UAChB,MAAM;AAAA,UACN,SAAS;AAAA,QACX,CAAC;AAAA,MACH;AAEA,UAAI,KAAK,cAAc,gBAAgB,cAAc;AACnD,aAAK,sBAAsB;AAC3B,aAAK,WAAW;AAChB;AAAA,MACF;AAEA,WAAK,UAAU;AAAA,IACjB;AAEA,SAAQ,eAAe,CAAC,iBAAyB;AAC/C,YAAM,MAAM,KAAK,IAAI;AACrB,YAAM,sBAAsB,MAAM,KAAK;AACvC,UAAI,uBAAuB,cAAc;AACvC,aAAK,YAAY;AAAA,UACf,MAAM;AAAA,UACN,SAAS;AAAA,UACT,OAAO;AAAA,UACP,SAAS,+BAA+B,sBAAsB;AAAA,QAChE,CAAC;AAED,eAAO,KAAK,UAAU;AAAA,MACxB;AAEA,YAAM,cAAc,KAAK,IAAI,KAAK,eAAe,mBAAmB;AACpE,WAAK,UAAU;AAAA,QACb,MAAM,KAAK,aAAa,YAAY;AAAA,QACpC;AAAA,QACA;AAAA,MACF;AAAA,IACF;AAEA,SAAQ,oBAAoB,MAAM;AAChC,WAAK,UAAU;AAAA,QACb,MAAM;AACJ,eAAK,YAAY;AAAA,YACf,MAAM;AAAA,YACN,SAAS;AAAA,UACX,CAAC;AAED,eAAK,kBAAkB;AAAA,QACzB;AAAA,QACA,KAAK,OAAO,oBAAoB;AAAA,QAChC;AAAA,MACF;AAAA,IACF;AArfE,SAAK,SAAS;AAAA,MACZ,mBAAmB;AAAA,MACnB,kBAAkB;AAAA,MAClB,wBAAwB;AAAA,MACxB,eAAe;AAAA,MACf,UAAU,eAAe;AAAA,MACzB,sBAAsB;AAAA,MACtB,kBAAkB,CAAC,QAAQ,IAAI,gCAAgC,GAAG;AAAA,MAClE,GAAG;AAAA,IACL;AAEA,SAAK,SAAS,IAAIC,QAAO,KAAK,YAAY,MAAM,KAAK,OAAO,QAAQ;AACpE,SAAK,YAAY,KAAK,OAAO,aAAa,IAAI,uBAAuB;AAAA,EACvE;AAyeF;","names":["Logger","DXLinkChannelState","DXLinkChannelState","Logger"]} |
+1
-1
@@ -1,2 +0,2 @@ | ||
| var dxlinkCore=require("@dxfeed/dxlink-core");class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=dxlinkCore.DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==dxlinkCore.DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==dxlinkCore.DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(dxlinkCore.DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(dxlinkCore.DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new dxlinkCore.Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.1"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=dxlinkCore.DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=dxlinkCore.DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(dxlinkCore.DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(dxlinkCore.DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.authState===dxlinkCore.DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(dxlinkCore.DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTING&&this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(dxlinkCore.DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===dxlinkCore.DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:dxlinkCore.DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new dxlinkCore.Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new dxlinkCore.DefaultDXLinkScheduler}},exports.DefaultDXLinkWebSocketConnector=DefaultDXLinkWebSocketConnector; | ||
| var dxlinkCore=require("@dxfeed/dxlink-core");class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=dxlinkCore.DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==dxlinkCore.DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==dxlinkCore.DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(dxlinkCore.DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(dxlinkCore.DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(dxlinkCore.DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new dxlinkCore.Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.9.0"};exports.DXLINK_WS_PROTOCOL_VERSION="0.1",exports.DXLinkWebSocketClient=class{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=dxlinkCore.DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=dxlinkCore.DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==dxlinkCore.DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(dxlinkCore.DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(dxlinkCore.DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===dxlinkCore.DXLinkConnectionState.CONNECTED&&this.authState===dxlinkCore.DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(dxlinkCore.DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTING&&this.connectionState!==dxlinkCore.DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(dxlinkCore.DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(dxlinkCore.DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==dxlinkCore.DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===dxlinkCore.DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:dxlinkCore.DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new dxlinkCore.Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new dxlinkCore.DefaultDXLinkScheduler}},exports.DefaultDXLinkWebSocketConnector=DefaultDXLinkWebSocketConnector; | ||
| //# sourceMappingURL=index.js.map |
@@ -1,2 +0,2 @@ | ||
| import{DXLinkChannelState,Logger,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel,DefaultDXLinkScheduler}from"@dxfeed/dxlink-core";class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DXLINK_WS_PROTOCOL_VERSION="0.1",DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.8.1"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===DXLinkConnectionState.CONNECTED&&this.authState===DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==DXLinkConnectionState.CONNECTING&&this.connectionState!==DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new DefaultDXLinkScheduler}}export{DXLINK_WS_PROTOCOL_VERSION,DXLinkWebSocketClient,DefaultDXLinkWebSocketConnector}; | ||
| import{DXLinkChannelState,Logger,DXLinkConnectionState,DXLinkAuthState,DXLinkLogLevel,DefaultDXLinkScheduler}from"@dxfeed/dxlink-core";class DXLinkWebSocketChannel{constructor(id,service,parameters,reconnect,sendMessage,config){this.id=void 0,this.service=void 0,this.parameters=void 0,this.reconnect=void 0,this.sendMessage=void 0,this.status=DXLinkChannelState.REQUESTED,this.messageListeners=new Set,this.statusListeners=new Set,this.errorListeners=new Set,this.logger=void 0,this.send=({type:type,...payload})=>{if(this.status!==DXLinkChannelState.OPENED)throw new Error("Channel is not ready");this.sendMessage({type:type,channel:this.id,...payload})},this.addMessageListener=listener=>this.messageListeners.add(listener),this.removeMessageListener=listener=>this.messageListeners.delete(listener),this.getState=()=>this.status,this.addStateChangeListener=listener=>this.statusListeners.add(listener),this.removeStateChangeListener=listener=>this.statusListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.error=({type:type,message:message})=>this.send({type:"ERROR",error:type,message:message}),this.close=()=>{this.status!==DXLinkChannelState.CLOSED&&(this.logger.debug("Closing by user"),this.setStatus(DXLinkChannelState.CLOSED),this.clear(),this.sendMessage({type:"CHANNEL_CANCEL",channel:this.id}))},this.request=()=>{this.logger.debug("Requesting"),this.sendMessage({type:"CHANNEL_REQUEST",channel:this.id,service:this.service,parameters:this.parameters}),this.processStatusRequested()},this.processPayloadMessage=message=>{for(const listener of this.messageListeners)listener(message)},this.processStatusOpened=()=>{this.logger.debug("Opened"),this.setStatus(DXLinkChannelState.OPENED)},this.processStatusRequested=()=>{this.setStatus(DXLinkChannelState.REQUESTED)},this.processStatusClosed=()=>{this.logger.debug("Closed by remote endpoint"),this.setStatus(DXLinkChannelState.CLOSED),this.clear()},this.processError=error=>{if(0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error(`Error in channel#${this.id} error listener: `,e)}else this.logger.error(`Unhandled error in channel#${this.id}: `,error)},this.setStatus=newStatus=>{if(this.status===newStatus)return;const prev=this.status;this.status=newStatus;for(const listener of this.statusListeners)try{listener(newStatus,prev)}catch(e){this.logger.error(`Error in channel#${this.id} status listener: `,e)}},this.clear=()=>{this.messageListeners.clear(),this.statusListeners.clear()},this.id=id,this.service=service,this.parameters=parameters,this.reconnect=reconnect,this.sendMessage=sendMessage,this.logger=new Logger(`${DXLinkWebSocketChannel.name}#${id} ${service}`,config.logLevel)}}class DefaultDXLinkWebSocketConnector{constructor(url,protocols){this.url=void 0,this.protocols=void 0,this.socket=void 0,this.isAvailable=!1,this.openListener=void 0,this.closeListener=void 0,this.messageListener=void 0,this.sendMessage=message=>{void 0!==this.socket&&this.isAvailable&&this.socket.send(JSON.stringify(message))},this.setOpenListener=listener=>{this.openListener=listener},this.setCloseListener=listener=>{this.closeListener=listener},this.setMessageListener=listener=>{this.messageListener=listener},this.getUrl=()=>this.url,this.handleOpen=()=>{void 0!==this.socket&&(this.isAvailable=!0,this.socket.removeEventListener("open",this.handleOpen),this.socket.addEventListener("message",this.handleMessage),this.openListener?.())},this.handleClosed=ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.(ev.reason,!1))},this.handleError=_ev=>{void 0!==this.socket&&(this.stop(),this.closeListener?.("Unable to connect",!0))},this.handleMessage=ev=>{try{const message=JSON.parse(ev.data);if("object"!=typeof message)throw new Error("Unexpected message: "+typeof message);if("string"!=typeof message.type)throw new Error("Unexpected message type: "+typeof message.type);if("number"!=typeof message.channel)throw new Error("Unexpected message channel: "+typeof message.channel);if(void 0===this.messageListener)return console.warn("No message listener set");this.messageListener(message)}catch(error){console.error(error instanceof Error?error:new Error("Parsing error:"+String(error)))}},this.url=url,this.protocols=protocols}start(){void 0===this.socket&&(this.socket=new WebSocket(this.url,this.protocols),this.socket.addEventListener("open",this.handleOpen),this.socket.addEventListener("error",this.handleError),this.socket.addEventListener("close",this.handleClosed))}stop(){void 0!==this.socket&&(this.socket.removeEventListener("open",this.handleOpen),this.socket.removeEventListener("error",this.handleError),this.socket.removeEventListener("close",this.handleClosed),this.socket.removeEventListener("message",this.handleMessage),this.socket.close(),this.socket=void 0,this.isAvailable=!1)}}const DXLINK_WS_PROTOCOL_VERSION="0.1",DEFAULT_CONNECTION_DETAILS={protocolVersion:"0.1",clientVersion:"DXF-JS/0.9.0"};class DXLinkWebSocketClient{constructor(config){this.config=void 0,this.logger=void 0,this.scheduler=void 0,this.connector=void 0,this.connectionState=DXLinkConnectionState.NOT_CONNECTED,this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.authState=DXLinkAuthState.UNAUTHORIZED,this.connectionStateChangeListeners=new Set,this.errorListeners=new Set,this.authStateChangeListeners=new Set,this.isFirstAuthState=!0,this.lastSettedAuthToken=void 0,this.lastReceivedMillis=0,this.lastSentMillis=0,this.reconnectAttempts=0,this.globalChannelId=1,this.channels=new Map,this.connect=url=>{this.connector?.getUrl()!==url&&(this.disconnect(),this.logger.debug("Connecting to",url),this.setConnectionState(DXLinkConnectionState.CONNECTING),this.connector=this.config.connectorFactory(url),this.connector.setOpenListener(this.processTransportOpen),this.connector.setMessageListener(this.processMessage),this.connector.setCloseListener(this.processTransportClose),this.connector.start())},this.reconnect=()=>{if(this.config.maxReconnectAttempts>=0&&this.reconnectAttempts>=this.config.maxReconnectAttempts)return this.logger.warn("Max reconnect attempts reached"),void this.disconnect();if(this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&void 0!==this.connector){this.logger.debug("Trying to reconnect",this.connector.getUrl()),this.connector.stop(),this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts++,this.setConnectionState(DXLinkConnectionState.CONNECTING);for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED&&(channel.reconnect?channel.processStatusRequested():channel.processStatusClosed());this.scheduler.schedule(()=>{void 0!==this.connector&&this.connector.start()},1e3*this.reconnectAttempts,"DXLWS_RECONNECT")}},this.disconnect=()=>{this.connectionState!==DXLinkConnectionState.NOT_CONNECTED&&(this.logger.debug("Disconnecting"),this.connector?.stop(),this.connector=void 0,this.scheduler.clear(),this.connectionDetails=DEFAULT_CONNECTION_DETAILS,this.lastReceivedMillis=0,this.lastSentMillis=0,this.isFirstAuthState=!0,this.reconnectAttempts=0,this.setConnectionState(DXLinkConnectionState.NOT_CONNECTED),this.setAuthState(DXLinkAuthState.UNAUTHORIZED))},this.close=()=>{this.disconnect()},this.getConnectionDetails=()=>this.connectionDetails,this.getConnectionState=()=>this.connectionState,this.getScheduler=()=>this.scheduler,this.addConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.add(listener),this.removeConnectionStateChangeListener=listener=>this.connectionStateChangeListeners.delete(listener),this.setAuthToken=token=>{this.lastSettedAuthToken=token,this.connectionState===DXLinkConnectionState.CONNECTED&&this.sendAuthMessage(token)},this.getAuthState=()=>this.authState,this.addAuthStateChangeListener=listener=>this.authStateChangeListeners.add(listener),this.removeAuthStateChangeListener=listener=>this.authStateChangeListeners.delete(listener),this.addErrorListener=listener=>this.errorListeners.add(listener),this.removeErrorListener=listener=>this.errorListeners.delete(listener),this.openChannel=(service,parameters,options)=>{const channelId=this.globalChannelId;this.globalChannelId+=2;const channel=new DXLinkWebSocketChannel(channelId,service,parameters,options?.reconnect??!0,this.sendMessage,this.config);return this.channels.set(channelId,channel),this.connectionState===DXLinkConnectionState.CONNECTED&&this.authState===DXLinkAuthState.AUTHORIZED&&channel.request(),channel},this.setConnectionState=newStatus=>{const prev=this.connectionState;if(prev!==newStatus){this.connectionState=newStatus;for(const listener of this.connectionStateChangeListeners)listener(newStatus,prev)}},this.sendMessage=message=>{this.connector?.sendMessage(message),this.scheduleKeepalive(),this.lastSentMillis=Date.now()},this.sendAuthMessage=token=>{this.logger.debug("Sending auth message"),this.setAuthState(DXLinkAuthState.AUTHORIZING),this.sendMessage({type:"AUTH",channel:0,token:token})},this.setAuthState=newState=>{const prev=this.authState;this.authState=newState;for(const listener of this.authStateChangeListeners)try{listener(newState,prev)}catch(e){this.logger.error("Auth state listener error",e)}},this.processMessage=message=>{if(this.lastReceivedMillis=Date.now(),this.lastReceivedMillis-this.lastSentMillis>=1e3*this.config.keepaliveInterval&&this.sendMessage({type:"KEEPALIVE",channel:0}),(message=>0===message.channel&&("SETUP"===message.type||"KEEPALIVE"===message.type||"AUTH"===message.type||"AUTH_STATE"===message.type||"ERROR"===message.type))(message))switch(message.type){case"SETUP":return this.processSetupMessage(message);case"AUTH_STATE":return this.processAuthStateMessage(message);case"ERROR":return this.publishError({type:message.error,message:message.message});case"KEEPALIVE":return}else if((message=>0!==message.channel)(message)){const channel=this.channels.get(message.channel);if(void 0===channel)return void this.logger.warn("Received lifecycle message for unknown channel",message);if((message=>"CHANNEL_OPENED"===message.type||"CHANNEL_CLOSED"===message.type||"ERROR"===message.type||"CHANNEL_REQUEST"===message.type||"CHANNEL_CANCEL"===message.type)(message)){switch(message.type){case"CHANNEL_OPENED":return channel.processStatusOpened();case"CHANNEL_CLOSED":return channel.processStatusClosed();case"ERROR":return channel.processError({type:message.error,message:message.message})}return}return channel.processPayloadMessage(message)}this.logger.warn("Unhandeled message",message.type)},this.processSetupMessage=serverSetup=>{this.scheduler.cancel("DXLWS_SETUP_TIMEOUT"),this.connectionState!==DXLinkConnectionState.CONNECTING&&this.connectionState!==DXLinkConnectionState.CONNECTED||(this.connectionDetails={...this.connectionDetails,serverVersion:serverSetup.version,clientKeepaliveTimeout:this.config.keepaliveTimeout,serverKeepaliveTimeout:serverSetup.keepaliveTimeout},this.reconnectAttempts=0,void 0===this.lastSettedAuthToken&&this.setConnectionState(DXLinkConnectionState.CONNECTED));const timeoutMills=1e3*(serverSetup.keepaliveTimeout??60);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),timeoutMills,"DXLWS_TIMEOUT")},this.publishError=error=>{if(this.logger.debug("Publishing error",error),0!==this.errorListeners.size)for(const listener of this.errorListeners)try{listener(error)}catch(e){this.logger.error("Error listener error",e)}else this.logger.error("Unhandled dxLink error",error)},this.processAuthStateMessage=({state:state})=>{this.logger.debug("Received auth state message",state),this.scheduler.cancel("DXLWS_AUTH_STATE_TIMEOUT"),this.isFirstAuthState?this.isFirstAuthState=!1:"UNAUTHORIZED"===state&&(this.lastSettedAuthToken=void 0),"AUTHORIZED"===state&&(this.setConnectionState(DXLinkConnectionState.CONNECTED),this.requestActiveChannels()),this.setAuthState(DXLinkAuthState[state])},this.requestActiveChannels=()=>{for(const channel of this.channels.values())channel.getState()!==DXLinkChannelState.CLOSED?channel.request():this.channels.delete(channel.id)},this.processTransportOpen=()=>{this.logger.debug("Connection opened");const setupMessage={type:"SETUP",channel:0,version:`${this.connectionDetails.protocolVersion}-${this.connectionDetails.clientVersion}`,keepaliveTimeout:this.config.keepaliveTimeout,acceptKeepaliveTimeout:this.config.acceptKeepaliveTimeout};this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No setup message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_SETUP_TIMEOUT"),this.sendMessage(setupMessage),this.scheduler.schedule(()=>{const errorMessage={type:"ERROR",channel:0,error:"TIMEOUT",message:"No auth state message received for "+this.config.actionTimeout+"s"};this.sendMessage(errorMessage),this.publishError({type:errorMessage.error,message:`${errorMessage.message} from server`}),this.reconnect()},1e3*this.config.actionTimeout,"DXLWS_AUTH_STATE_TIMEOUT"),void 0!==this.lastSettedAuthToken&&this.sendAuthMessage(this.lastSettedAuthToken)},this.processTransportClose=(reason,error)=>{if(this.logger.debug("Connection closed",reason),void 0!==error&&this.publishError({type:"UNKNOWN",message:reason}),this.authState===DXLinkAuthState.UNAUTHORIZED)return this.lastSettedAuthToken=void 0,void this.disconnect();this.reconnect()},this.timeoutCheck=timeoutMills=>{const noKeepaliveDuration=Date.now()-this.lastReceivedMillis;if(noKeepaliveDuration>=timeoutMills)return this.sendMessage({type:"ERROR",channel:0,error:"TIMEOUT",message:"No keepalive received for "+noKeepaliveDuration+"ms"}),this.reconnect();const nextTimeout=Math.max(200,timeoutMills-noKeepaliveDuration);this.scheduler.schedule(()=>this.timeoutCheck(timeoutMills),nextTimeout,"DXLWS_TIMEOUT")},this.scheduleKeepalive=()=>{this.scheduler.schedule(()=>{this.sendMessage({type:"KEEPALIVE",channel:0}),this.scheduleKeepalive()},1e3*this.config.keepaliveInterval,"DXLWS_KEEPALIVE")},this.config={keepaliveInterval:30,keepaliveTimeout:60,acceptKeepaliveTimeout:60,actionTimeout:10,logLevel:DXLinkLogLevel.WARN,maxReconnectAttempts:-1,connectorFactory:url=>new DefaultDXLinkWebSocketConnector(url),...config},this.logger=new Logger(this.constructor.name,this.config.logLevel),this.scheduler=this.config.scheduler??new DefaultDXLinkScheduler}}export{DXLINK_WS_PROTOCOL_VERSION,DXLinkWebSocketClient,DefaultDXLinkWebSocketConnector}; | ||
| //# sourceMappingURL=index.module.js.map |
+2
-2
| { | ||
| "name": "@dxfeed/dxlink-websocket-client", | ||
| "version": "0.8.1", | ||
| "version": "0.9.0", | ||
| "private": false, | ||
@@ -33,3 +33,3 @@ "publishConfig": { | ||
| "dependencies": { | ||
| "@dxfeed/dxlink-core": "0.8.1" | ||
| "@dxfeed/dxlink-core": "0.9.0" | ||
| }, | ||
@@ -36,0 +36,0 @@ "author": "Dmitry Petrov <dmitry.petrov@devexperts.com>", |
| import type { DXLinkClient } from '@dxfeed/dxlink-core'; | ||
| /** | ||
| * Adapter to connect dxLink client to gRPC DevTools. | ||
| */ | ||
| export declare function enableGrpcDevTools(client: DXLinkClient): void; |
| export {}; |
Major refactor
Supply chain riskPackage has recently undergone a major refactor. It may be unstable or indicate significant internal changes. Use caution when updating to versions that include significant changes.
210746
55.71%17
6.25%1049
136.26%1
Infinity%+ Added
- Removed
Updated