Загрузка данных
import dayjs from 'dayjs';
import duration from 'dayjs/plugin/duration';
import { parseTimeframe } from 'moex-chart';
import { moexChartTimeConverter, moexChartToIssTimeframe } from '@utils/chartToReqTimeConverter';
import { DEFAULT_SYMBOL } from '@widgets/Chart/const';
import { requestBars, requestRealtimeBars } from '../../requestBars';
import { ChartIndicativeData } from '../../types';
import type { Candle, Timeframes } from 'moex-chart';
dayjs.extend(duration);
interface RealtimeState {
previousTime: number | null;
candles: Candle[];
sessionStart: number | null;
}
function getRequestSymbol(symbolRaw?: string): string | undefined {
const symbol = String(symbolRaw ?? '').trim();
if (!symbol || symbol === DEFAULT_SYMBOL) {
return undefined;
}
return symbol;
}
function getTimeframeSeconds(timeframe: Timeframes): number {
const { candleWidth, dayjsUnit } = parseTimeframe(timeframe);
return dayjs.duration(candleWidth, dayjsUnit).asSeconds();
}
function isValidCandle({ time, open, high, low, close, volume }: Candle): boolean {
const valuesAreValid = [time, open, high, low, close].every(Number.isFinite);
return (
valuesAreValid &&
(volume === undefined || Number.isFinite(volume)) &&
high >= Math.max(open, close) &&
low <= Math.min(open, close)
);
}
function aggregateCandles(candles: Candle[], time: number): Candle | undefined {
const firstCandle = candles[0];
const lastCandle = candles[candles.length - 1];
if (!firstCandle || !lastCandle) {
return undefined;
}
let high = firstCandle.high;
let low = firstCandle.low;
let volume = 0;
for (const candle of candles) {
high = Math.max(high, candle.high);
low = Math.min(low, candle.low);
volume += candle.volume ?? 0;
}
return {
time,
open: firstCandle.open,
high,
low,
close: lastCandle.close,
volume,
};
}
// По хорошему - класс должен быть синглтоном, чтобы кормить MoexChart одинаковой датой,
// и не плодить несколько подключений на одни символа
class DataSourceProvider {
private readonly realtimeStates = new Map<string, RealtimeState>();
private realtimeTimer: ReturnType<typeof setInterval> | null = null;
public getDataSource =
(indicativeData?: ChartIndicativeData, cb?: (timeframe: Timeframes) => void) =>
async (timeframe: Timeframes, symbolId: string, until?: Candle): Promise<Candle[] | null> => {
const symbol = getRequestSymbol(symbolId);
if (!symbol) {
return null;
}
cb?.(timeframe);
const interval = moexChartTimeConverter(timeframe);
const date = until?.time ?? Math.round(Date.now() / 1000);
const data = await requestBars({
currencyPair: symbol.replaceAll(':', '.'),
interval,
periodParams: {
firstDataRequest: true,
to: date,
from: Math.round(Date.now() / 1000),
countBack: 2000,
},
ticker: symbol,
indicativeData,
});
if (data.length === 0) {
return null;
}
const validData = data.filter(isValidCandle);
const invalidCandlesCount = data.length - validData.length;
if (invalidCandlesCount > 0) {
console.error(
`[DataSourceProvider]: получены некорректные свечи для ${symbol}: ${invalidCandlesCount}`,
);
}
if (validData.length === 0) {
return null;
}
const shouldConvolute = moexChartToIssTimeframe(timeframe) !== timeframe;
if (!shouldConvolute) {
if (!until) {
const state = this.getRealtimeState(symbol);
state.previousTime = validData[validData.length - 1]?.time ?? null;
state.candles = [];
state.sessionStart = null;
}
return validData;
}
return this.timeframeConvolution(
validData,
timeframe,
until ? undefined : this.getRealtimeState(symbol),
);
};
public startRealtime({
getSymbols,
getTimeframe,
update,
periodMs = 5000,
indicativeData,
}: {
getSymbols: () => string[];
getTimeframe: () => Timeframes;
update: (symbolId: string, candle: Candle) => void;
periodMs?: number;
indicativeData?: ChartIndicativeData;
}): () => void {
if (this.realtimeTimer) {
clearInterval(this.realtimeTimer);
}
this.realtimeTimer = setInterval(() => {
const timeframe = getTimeframe();
Promise.all(
getSymbols().map((symbolId) =>
this.updateRealtimeSymbol(symbolId, timeframe, update, indicativeData),
),
).catch((error) => {
console.error('[DataSourceProvider]: ошибка получения realtime данных', error);
});
}, periodMs);
return () => {
if (this.realtimeTimer) {
clearInterval(this.realtimeTimer);
}
this.realtimeTimer = null;
};
}
private async updateRealtimeSymbol(
symbolId: string,
timeframe: Timeframes,
update: (symbolId: string, candle: Candle) => void,
indicativeData?: ChartIndicativeData,
): Promise<void> {
const symbol = getRequestSymbol(symbolId);
if (!symbol) {
return;
}
const data = await requestRealtimeBars({
currencyPair: symbol.replaceAll(':', '.'),
interval: moexChartTimeConverter(timeframe),
ticker: symbol,
indicativeData,
});
if (!data) {
return;
}
if (!isValidCandle(data)) {
console.error(`[DataSourceProvider]: получена некорректная realtime свеча для ${symbol}`, data);
return;
}
if (moexChartToIssTimeframe(timeframe) === timeframe) {
this.getRealtimeState(symbol).previousTime = data.time;
update(symbol, data);
return;
}
const candle = this.realtimeConvolution(
timeframe,
data,
this.getRealtimeState(symbol),
);
if (candle) {
update(symbol, candle);
}
}
private getRealtimeState(symbol: string): RealtimeState {
let state = this.realtimeStates.get(symbol);
if (!state) {
state = {
previousTime: null,
candles: [],
sessionStart: null,
};
this.realtimeStates.set(symbol, state);
}
return state;
}
private timeframeConvolution(
data: Candle[],
requestedTimeframe: Timeframes,
realtimeState?: RealtimeState,
): Candle[] {
const timeframeSeconds = getTimeframeSeconds(requestedTimeframe);
const sortedData = [...data].sort((first, second) => first.time - second.time);
const firstCandle = sortedData[0];
if (!firstCandle) {
return [];
}
const result: Candle[] = [];
let sessionStart = firstCandle.time;
let bucketStart = sessionStart;
let previousTime: number | null = null;
let candleGroup: Candle[] = [];
for (const candle of sortedData) {
const isNewSession = previousTime !== null && candle.time - previousTime > timeframeSeconds;
if (isNewSession) {
const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);
if (aggregatedCandle) {
result.push(aggregatedCandle);
}
sessionStart = candle.time;
bucketStart = candle.time;
candleGroup = [candle];
previousTime = candle.time;
continue;
}
const currentBucketStart =
sessionStart + Math.floor((candle.time - sessionStart) / timeframeSeconds) * timeframeSeconds;
if (currentBucketStart !== bucketStart) {
const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);
if (aggregatedCandle) {
result.push(aggregatedCandle);
}
bucketStart = currentBucketStart;
candleGroup = [];
}
candleGroup.push(candle);
previousTime = candle.time;
}
const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);
if (aggregatedCandle) {
result.push(aggregatedCandle);
}
if (realtimeState) {
realtimeState.previousTime = sortedData[sortedData.length - 1]?.time ?? null;
realtimeState.candles = [...candleGroup];
realtimeState.sessionStart = sessionStart;
}
return result;
}
private realtimeConvolution(
timeframe: Timeframes,
data: Candle,
state: RealtimeState,
): Candle | undefined {
const timeframeSeconds = getTimeframeSeconds(timeframe);
const { previousTime } = state;
if (previousTime === null || state.sessionStart === null) {
state.sessionStart = data.time;
state.candles = [data];
} else if (data.time - previousTime > timeframeSeconds) {
state.sessionStart = data.time;
state.candles = [data];
} else {
const previousBucketStart =
state.sessionStart +
Math.floor((previousTime - state.sessionStart) / timeframeSeconds) * timeframeSeconds;
const currentBucketStart =
state.sessionStart +
Math.floor((data.time - state.sessionStart) / timeframeSeconds) * timeframeSeconds;
if (currentBucketStart !== previousBucketStart) {
state.candles = [data];
} else {
const candleIndex = state.candles.findIndex((candle) => candle.time === data.time);
if (candleIndex === -1) {
state.candles.push(data);
} else {
state.candles[candleIndex] = data;
}
}
}
state.previousTime = data.time;
const sessionStart = state.sessionStart;
if (sessionStart === null) {
return undefined;
}
const bucketStart =
sessionStart + Math.floor((data.time - sessionStart) / timeframeSeconds) * timeframeSeconds;
return aggregateCandles(state.candles, bucketStart);
}
}
export { DataSourceProvider };