| import WebSocketFactory from './lib/websocket-factory'; |
| import { CHANNEL_EVENTS, CONNECTION_STATE, DEFAULT_VERSION, DEFAULT_TIMEOUT, DEFAULT_VSN, VSN_1_0_0, VSN_2_0_0, } from './lib/constants'; |
| import Serializer from './lib/serializer'; |
| import { httpEndpointURL } from './lib/transformers'; |
| import RealtimeChannel from './RealtimeChannel'; |
| import SocketAdapter from './phoenix/socketAdapter'; |
| |
| const CONNECTION_TIMEOUTS = { |
| HEARTBEAT_INTERVAL: 25000, |
| RECONNECT_DELAY: 10, |
| HEARTBEAT_TIMEOUT_FALLBACK: 100, |
| }; |
| const RECONNECT_INTERVALS = [1000, 2000, 5000, 10000]; |
| const DEFAULT_RECONNECT_FALLBACK = 10000; |
| function createMemorySessionStorage() { |
| const store = new Map(); |
| return { |
| get length() { |
| return store.size; |
| }, |
| clear() { |
| store.clear(); |
| }, |
| getItem(key) { |
| return store.has(key) ? store.get(key) : null; |
| }, |
| key(index) { |
| var _a; |
| return (_a = Array.from(store.keys())[index]) !== null && _a !== void 0 ? _a : null; |
| }, |
| removeItem(key) { |
| store.delete(key); |
| }, |
| setItem(key, value) { |
| store.set(key, String(value)); |
| }, |
| }; |
| } |
| function resolveSessionStorage() { |
| try { |
| if (typeof globalThis !== 'undefined' && globalThis.sessionStorage) { |
| return globalThis.sessionStorage; |
| } |
| } |
| catch (_a) { |
| |
| } |
| return createMemorySessionStorage(); |
| } |
| const WORKER_SCRIPT = ` |
| addEventListener("message", (e) => { |
| if (e.data.event === "start") { |
| setInterval(() => postMessage({ event: "keepAlive" }), e.data.interval); |
| } |
| });`; |
| export default class RealtimeClient { |
| get endPoint() { |
| return this.socketAdapter.endPoint; |
| } |
| get timeout() { |
| return this.socketAdapter.timeout; |
| } |
| get transport() { |
| return this.socketAdapter.transport; |
| } |
| get heartbeatCallback() { |
| return this.socketAdapter.heartbeatCallback; |
| } |
| get heartbeatIntervalMs() { |
| return this.socketAdapter.heartbeatIntervalMs; |
| } |
| get heartbeatTimer() { |
| if (this.worker) { |
| return this._workerHeartbeatTimer; |
| } |
| return this.socketAdapter.heartbeatTimer; |
| } |
| get pendingHeartbeatRef() { |
| if (this.worker) { |
| return this._pendingWorkerHeartbeatRef; |
| } |
| return this.socketAdapter.pendingHeartbeatRef; |
| } |
| get reconnectTimer() { |
| return this.socketAdapter.reconnectTimer; |
| } |
| get vsn() { |
| return this.socketAdapter.vsn; |
| } |
| get encode() { |
| return this.socketAdapter.encode; |
| } |
| get decode() { |
| return this.socketAdapter.decode; |
| } |
| get reconnectAfterMs() { |
| return this.socketAdapter.reconnectAfterMs; |
| } |
| get sendBuffer() { |
| return this.socketAdapter.sendBuffer; |
| } |
| get stateChangeCallbacks() { |
| return this.socketAdapter.stateChangeCallbacks; |
| } |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| constructor(endPoint, options) { |
| var _a; |
| this.channels = new Array(); |
| this.accessTokenValue = null; |
| this.accessToken = null; |
| this.apiKey = null; |
| this.httpEndpoint = ''; |
| |
| this.headers = {}; |
| this.params = {}; |
| this.ref = 0; |
| this.serializer = new Serializer(); |
| this._manuallySetToken = false; |
| this._authPromise = null; |
| this._workerHeartbeatTimer = undefined; |
| this._pendingWorkerHeartbeatRef = null; |
| this._pendingDisconnectTimer = null; |
| this._disconnectOnEmptyChannelsAfterMs = 0; |
| |
| |
| |
| |
| |
| this._resolveFetch = (customFetch) => { |
| if (customFetch) { |
| return (...args) => customFetch(...args); |
| } |
| return (...args) => fetch(...args); |
| }; |
| |
| if (!((_a = options === null || options === void 0 ? void 0 : options.params) === null || _a === void 0 ? void 0 : _a.apikey)) { |
| throw new Error('API key is required to connect to Realtime'); |
| } |
| this.apiKey = options.params.apikey; |
| const socketAdapterOptions = this._initializeOptions(options); |
| this.socketAdapter = new SocketAdapter(endPoint, socketAdapterOptions); |
| this.httpEndpoint = httpEndpointURL(endPoint); |
| this.fetch = this._resolveFetch(options === null || options === void 0 ? void 0 : options.fetch); |
| } |
| |
| |
| |
| |
| |
| connect() { |
| |
| if (this.isConnecting() || this.isDisconnecting() || this.isConnected()) { |
| return; |
| } |
| |
| |
| |
| if (this.accessToken && !this._authPromise) { |
| this._setAuthSafely('connect'); |
| } |
| this._setupConnectionHandlers(); |
| try { |
| this.socketAdapter.connect(); |
| } |
| catch (error) { |
| const errorMessage = error.message; |
| throw new Error(`WebSocket not available: ${errorMessage}`); |
| } |
| this._handleNodeJsRaceCondition(); |
| } |
| |
| |
| |
| |
| |
| |
| endpointURL() { |
| return this.socketAdapter.endPointURL(); |
| } |
| |
| |
| |
| |
| |
| |
| |
| |
| async disconnect(code, reason) { |
| this._cancelPendingDisconnect(); |
| if (this.isDisconnecting()) { |
| return 'ok'; |
| } |
| return await this.socketAdapter.disconnect(() => { |
| clearInterval(this._workerHeartbeatTimer); |
| this._terminateWorker(); |
| }, code, reason); |
| } |
| |
| |
| |
| |
| |
| getChannels() { |
| return this.channels; |
| } |
| |
| |
| |
| |
| |
| |
| async removeChannel(channel) { |
| const status = await channel.unsubscribe(); |
| if (status === 'ok') { |
| channel.teardown(); |
| } |
| return status; |
| } |
| |
| |
| |
| |
| |
| async removeAllChannels() { |
| const promises = this.channels.map(async (channel) => { |
| const result = await channel.unsubscribe(); |
| channel.teardown(); |
| return result; |
| }); |
| const result = await Promise.all(promises); |
| await this.disconnect(); |
| return result; |
| } |
| |
| |
| |
| |
| |
| |
| |
| log(kind, msg, data) { |
| this.socketAdapter.log(kind, msg, data); |
| } |
| |
| |
| |
| |
| |
| connectionState() { |
| return this.socketAdapter.connectionState() || CONNECTION_STATE.closed; |
| } |
| |
| |
| |
| |
| |
| isConnected() { |
| return this.socketAdapter.isConnected(); |
| } |
| |
| |
| |
| |
| |
| isConnecting() { |
| return this.socketAdapter.isConnecting(); |
| } |
| |
| |
| |
| |
| |
| isDisconnecting() { |
| return this.socketAdapter.isDisconnecting(); |
| } |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| channel(topic, params = { config: {} }) { |
| const realtimeTopic = `realtime:${topic}`; |
| const exists = this.getChannels().find((c) => c.topic === realtimeTopic); |
| if (!exists) { |
| const chan = new RealtimeChannel(`realtime:${topic}`, params, this); |
| this._cancelPendingDisconnect(); |
| this.channels.push(chan); |
| return chan; |
| } |
| else { |
| return exists; |
| } |
| } |
| |
| |
| |
| |
| |
| |
| |
| push(data) { |
| this.socketAdapter.push(data); |
| } |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| async setAuth(token = null) { |
| this._authPromise = this._performAuth(token); |
| try { |
| await this._authPromise; |
| } |
| finally { |
| this._authPromise = null; |
| } |
| } |
| |
| |
| |
| |
| |
| _isManualToken() { |
| return this._manuallySetToken; |
| } |
| |
| |
| |
| |
| |
| async sendHeartbeat() { |
| this.socketAdapter.sendHeartbeat(); |
| } |
| |
| |
| |
| |
| |
| |
| onHeartbeat(callback) { |
| this.socketAdapter.heartbeatCallback = this._wrapHeartbeatCallback(callback); |
| } |
| |
| |
| |
| |
| |
| _makeRef() { |
| return this.socketAdapter.makeRef(); |
| } |
| |
| |
| |
| |
| |
| |
| |
| _remove(channel) { |
| this.channels = this.channels.filter((c) => c.topic !== channel.topic); |
| if (this.channels.length === 0) { |
| this.log('transport', 'no channels remaining, scheduling disconnect'); |
| this._schedulePendingDisconnect(); |
| } |
| } |
| |
| _schedulePendingDisconnect() { |
| this._cancelPendingDisconnect(); |
| if (this._disconnectOnEmptyChannelsAfterMs === 0) { |
| this.log('transport', 'disconnecting immediately - no channels'); |
| this.disconnect(); |
| return; |
| } |
| this._pendingDisconnectTimer = setTimeout(() => { |
| this._pendingDisconnectTimer = null; |
| if (this.channels.length === 0) { |
| this.log('transport', 'deferred disconnect fired - no channels, disconnecting'); |
| this.disconnect(); |
| } |
| }, this._disconnectOnEmptyChannelsAfterMs); |
| this.log('transport', `deferred disconnect scheduled in ${this._disconnectOnEmptyChannelsAfterMs}ms`); |
| } |
| |
| _cancelPendingDisconnect() { |
| if (this._pendingDisconnectTimer !== null) { |
| this.log('transport', 'pending disconnect cancelled - channel activity detected'); |
| clearTimeout(this._pendingDisconnectTimer); |
| this._pendingDisconnectTimer = null; |
| } |
| } |
| |
| |
| |
| |
| async _performAuth(token = null) { |
| let tokenToSend; |
| let isManualToken = false; |
| if (token) { |
| tokenToSend = token; |
| |
| isManualToken = true; |
| } |
| else if (this.accessToken) { |
| |
| try { |
| tokenToSend = await this.accessToken(); |
| } |
| catch (e) { |
| this.log('error', 'Error fetching access token from callback', e); |
| |
| tokenToSend = this.accessTokenValue; |
| } |
| } |
| else { |
| tokenToSend = this.accessTokenValue; |
| } |
| |
| if (isManualToken) { |
| this._manuallySetToken = true; |
| } |
| else if (this.accessToken) { |
| |
| this._manuallySetToken = false; |
| } |
| if (this.accessTokenValue != tokenToSend) { |
| this.accessTokenValue = tokenToSend; |
| this.channels.forEach((channel) => { |
| const payload = { |
| access_token: tokenToSend, |
| version: DEFAULT_VERSION, |
| }; |
| tokenToSend && channel.updateJoinPayload(payload); |
| if (channel.joinedOnce && channel.channelAdapter.isJoined()) { |
| channel.channelAdapter.push(CHANNEL_EVENTS.access_token, { |
| access_token: tokenToSend, |
| }); |
| } |
| }); |
| } |
| } |
| |
| |
| |
| |
| async _waitForAuthIfNeeded() { |
| if (this._authPromise) { |
| await this._authPromise; |
| } |
| } |
| |
| |
| |
| |
| _setAuthSafely(context = 'general') { |
| |
| if (!this._isManualToken()) { |
| this.setAuth().catch((e) => { |
| this.log('error', `Error setting auth in ${context}`, e); |
| }); |
| } |
| } |
| |
| _setupConnectionHandlers() { |
| this.socketAdapter.onOpen(() => { |
| const authPromise = this._authPromise || |
| (this.accessToken && !this.accessTokenValue ? this.setAuth() : Promise.resolve()); |
| authPromise.catch((e) => { |
| this.log('error', 'error waiting for auth on connect', e); |
| }); |
| if (this.worker && !this.workerRef) { |
| this._startWorkerHeartbeat(); |
| } |
| }); |
| this.socketAdapter.onClose(() => { |
| if (this.worker && this.workerRef) { |
| this._terminateWorker(); |
| } |
| }); |
| this.socketAdapter.onMessage((message) => { |
| if (message.ref && message.ref === this._pendingWorkerHeartbeatRef) { |
| this._pendingWorkerHeartbeatRef = null; |
| } |
| }); |
| } |
| |
| _handleNodeJsRaceCondition() { |
| if (this.socketAdapter.isConnected()) { |
| |
| this.socketAdapter.getSocket().onConnOpen(); |
| } |
| } |
| |
| _wrapHeartbeatCallback(heartbeatCallback) { |
| return (status, latency) => { |
| if (status === 'disconnected') |
| return; |
| if (status == 'sent') |
| this._setAuthSafely(); |
| if (heartbeatCallback) |
| heartbeatCallback(status, latency); |
| }; |
| } |
| |
| _startWorkerHeartbeat() { |
| if (this.workerUrl) { |
| this.log('worker', `starting worker for from ${this.workerUrl}`); |
| } |
| else { |
| this.log('worker', `starting default worker`); |
| } |
| const objectUrl = this._workerObjectUrl(this.workerUrl); |
| this.workerRef = new Worker(objectUrl); |
| this.workerRef.onerror = (error) => { |
| this.log('worker', 'worker error', error.message); |
| this._terminateWorker(); |
| this.disconnect(); |
| }; |
| this.workerRef.onmessage = (event) => { |
| if (event.data.event === 'keepAlive') { |
| this.sendHeartbeat(); |
| } |
| }; |
| this.workerRef.postMessage({ |
| event: 'start', |
| interval: this.heartbeatIntervalMs, |
| }); |
| } |
| |
| |
| |
| |
| _terminateWorker() { |
| if (this.workerRef) { |
| this.log('worker', 'terminating worker'); |
| this.workerRef.terminate(); |
| this.workerRef = undefined; |
| } |
| } |
| |
| _workerObjectUrl(url) { |
| let result_url; |
| if (url) { |
| result_url = url; |
| } |
| else { |
| const blob = new Blob([WORKER_SCRIPT], { type: 'application/javascript' }); |
| result_url = URL.createObjectURL(blob); |
| } |
| return result_url; |
| } |
| |
| |
| |
| |
| _initializeOptions(options) { |
| var _a, _b, _c, _d, _e, _f, _g, _h, _j, _k, _l, _m; |
| this.worker = (_a = options === null || options === void 0 ? void 0 : options.worker) !== null && _a !== void 0 ? _a : false; |
| this.accessToken = (_b = options === null || options === void 0 ? void 0 : options.accessToken) !== null && _b !== void 0 ? _b : null; |
| const result = {}; |
| result.timeout = (_c = options === null || options === void 0 ? void 0 : options.timeout) !== null && _c !== void 0 ? _c : DEFAULT_TIMEOUT; |
| result.heartbeatIntervalMs = |
| (_d = options === null || options === void 0 ? void 0 : options.heartbeatIntervalMs) !== null && _d !== void 0 ? _d : CONNECTION_TIMEOUTS.HEARTBEAT_INTERVAL; |
| this._disconnectOnEmptyChannelsAfterMs = |
| (_e = options === null || options === void 0 ? void 0 : options.disconnectOnEmptyChannelsAfterMs) !== null && _e !== void 0 ? _e : 2 * ((_f = options === null || options === void 0 ? void 0 : options.heartbeatIntervalMs) !== null && _f !== void 0 ? _f : CONNECTION_TIMEOUTS.HEARTBEAT_INTERVAL); |
| |
| result.transport = (_g = options === null || options === void 0 ? void 0 : options.transport) !== null && _g !== void 0 ? _g : WebSocketFactory.getWebSocketConstructor(); |
| result.params = options === null || options === void 0 ? void 0 : options.params; |
| result.logger = options === null || options === void 0 ? void 0 : options.logger; |
| result.heartbeatCallback = this._wrapHeartbeatCallback(options === null || options === void 0 ? void 0 : options.heartbeatCallback); |
| result.sessionStorage = (_h = options === null || options === void 0 ? void 0 : options.sessionStorage) !== null && _h !== void 0 ? _h : resolveSessionStorage(); |
| result.reconnectAfterMs = |
| (_j = options === null || options === void 0 ? void 0 : options.reconnectAfterMs) !== null && _j !== void 0 ? _j : ((tries) => { |
| return RECONNECT_INTERVALS[tries - 1] || DEFAULT_RECONNECT_FALLBACK; |
| }); |
| let defaultEncode; |
| let defaultDecode; |
| const vsn = (_k = options === null || options === void 0 ? void 0 : options.vsn) !== null && _k !== void 0 ? _k : DEFAULT_VSN; |
| switch (vsn) { |
| case VSN_1_0_0: |
| defaultEncode = (payload, callback) => { |
| return callback(JSON.stringify(payload)); |
| }; |
| defaultDecode = (payload, callback) => { |
| return callback(JSON.parse(payload)); |
| }; |
| break; |
| case VSN_2_0_0: |
| defaultEncode = this.serializer.encode.bind(this.serializer); |
| defaultDecode = this.serializer.decode.bind(this.serializer); |
| break; |
| default: |
| throw new Error(`Unsupported serializer version: ${result.vsn}`); |
| } |
| result.vsn = vsn; |
| result.encode = (_l = options === null || options === void 0 ? void 0 : options.encode) !== null && _l !== void 0 ? _l : defaultEncode; |
| result.decode = (_m = options === null || options === void 0 ? void 0 : options.decode) !== null && _m !== void 0 ? _m : defaultDecode; |
| result.beforeReconnect = this._reconnectAuth.bind(this); |
| if ((options === null || options === void 0 ? void 0 : options.logLevel) || (options === null || options === void 0 ? void 0 : options.log_level)) { |
| this.logLevel = options.logLevel || options.log_level; |
| result.params = Object.assign(Object.assign({}, result.params), { log_level: this.logLevel }); |
| } |
| |
| if (this.worker) { |
| if (typeof window !== 'undefined' && !window.Worker) { |
| throw new Error('Web Worker is not supported'); |
| } |
| this.workerUrl = options === null || options === void 0 ? void 0 : options.workerUrl; |
| result.autoSendHeartbeat = !this.worker; |
| } |
| return result; |
| } |
| |
| async _reconnectAuth() { |
| await this._waitForAuthIfNeeded(); |
| if (!this.isConnected()) { |
| this.connect(); |
| } |
| } |
| } |
| |