import type {
Controller,
Manager,
Middleware,
} from '@data-client/react';
import { actionTypes, getDefaultManagers } from '@data-client/react';
import { ReconnectingSocket } from './socket';
import { Ticker } from './resources';
const { SUBSCRIBE, UNSUBSCRIBE } = actionTypes;
interface Product {
channel: string;
count: number;
}
/** Pushes Coinbase prices into the store over one socket */
export class StreamManager implements Manager {
protected socket = new ReconnectingSocket(
'wss://ws-feed.exchange.coinbase.com',
);
declare protected controller: Controller;
protected products = new Map<string, Product>();
middleware: Middleware = controller => {
this.controller = controller;
return next => async action => {
if (
(action.type === SUBSCRIBE || action.type === UNSUBSCRIBE) &&
'channel' in action.endpoint
) {
const { productId } = action.args[0];
if (action.type === SUBSCRIBE) {
const { channel } = action.endpoint as { channel: string };
this.subscribe(productId, channel);
} else this.unsubscribe(productId);
return;
}
return next(action);
};
};
/** Shares one socket subscription among a product's components */
protected subscribe(productId: string, channel: string) {
const product = this.products.get(productId);
if (product) {
product.count++;
} else {
this.products.set(productId, { channel, count: 1 });
this.send('subscribe', channel, productId);
}
}
protected unsubscribe(productId: string) {
const product = this.products.get(productId);
if (product && --product.count === 0) {
this.products.delete(productId);
this.send('unsubscribe', product.channel, productId);
}
}
init() {
this.socket.onopen = () => {
// a new socket has no subscriptions yet
for (const [productId, { channel }] of this.products)
this.send('subscribe', channel, productId);
};
this.socket.onmessage = ({ type, product_id, price, time }) => {
if (type !== 'ticker') return;
this.controller.set(
Ticker,
{ product_id },
{ product_id, price, time },
);
};
// without subscriptions, quiet is expected
this.socket.expectsMessages = () => this.products.size > 0;
this.socket.open();
}
cleanup() {
this.socket.close();
}
protected send(
type: 'subscribe' | 'unsubscribe',
channel: string,
productId: string,
) {
this.socket.send({
type,
product_ids: [productId],
channels: [channel],
});
}
}
// passed to <DataProvider managers={getManagers()}>
export default function getManagers() {
return [new StreamManager(), ...getDefaultManagers()];
}
/** A WebSocket that reconnects after drops, going offline and silence,
* and pauses while the page is hidden */
export class ReconnectingSocket {
onopen = () => {};
onmessage = (message: any) => {};
/** Whether silence means a dead connection */
expectsMessages = () => true;
declare protected socket: WebSocket;
protected attempts = 0;
declare protected openedAt: number | undefined;
/** Pending reconnect, or the watchdog while connected */
declare protected timer: ReturnType<typeof setTimeout>;
constructor(protected url: string) {}
open() {
if (!document.hidden) this.connect();
addEventListener('online', this.reconnect);
// a socket can take minutes to notice the network is gone
addEventListener('offline', this.stop);
document.addEventListener('visibilitychange', this.onVisibility);
}
close() {
removeEventListener('online', this.reconnect);
removeEventListener('offline', this.stop);
document.removeEventListener(
'visibilitychange',
this.onVisibility,
);
this.stop();
}
/** Dropped until open; onopen is the place to (re)send state */
send(message: object) {
if (this.socket?.readyState === WebSocket.OPEN)
this.socket.send(JSON.stringify(message));
}
protected connect() {
this.stop();
this.socket = new WebSocket(this.url);
this.openedAt = undefined;
this.socket.onopen = () => {
this.openedAt = Date.now();
this.onopen();
};
this.socket.onmessage = event => {
this.watch();
this.onmessage(JSON.parse(event.data));
};
// after a failed connect, an error or a server close
this.socket.onclose = () => this.retry();
this.watch();
}
/** A socket can stay open on a dead network (like after sleep),
* so reconnect when messages stop arriving */
protected watch() {
clearTimeout(this.timer);
this.timer = setTimeout(() => {
if (this.expectsMessages()) this.retry();
else this.watch();
}, 30_000);
}
/** Reconnects after a delay that doubles with each attempt, until a
* socket stays up long enough to count as working */
protected retry() {
this.stop();
if (
this.openedAt !== undefined &&
Date.now() - this.openedAt > 10_000
)
this.attempts = 0;
const delay = Math.min(30_000, 1000 * 2 ** this.attempts);
this.attempts++;
this.timer = setTimeout(this.reconnect, delay);
}
protected reconnect = () => {
if (
!document.hidden &&
this.socket?.readyState !== WebSocket.OPEN
)
this.connect();
};
/** Background tabs don't need prices */
protected onVisibility = () => {
if (document.hidden) this.stop();
else this.reconnect();
};
/** Stops the socket without waiting for it to finish closing,
* which it can't do offline */
protected stop = () => {
clearTimeout(this.timer);
if (!this.socket) return;
this.socket.onmessage = this.socket.onclose = null;
this.socket.close();
};
}