Загрузка данных


import type { Candle as MoexChartCandle } from 'moex-chart';
import type { Candle as ApiCandle } from 'types/Candles';

const getCandleStartTime = (time: string): number => Math.floor(Date.parse(`${time}Z`) / 1000);

export const candleToBar = ({ open, close, high, low, volume, begin }: ApiCandle): MoexChartCandle => ({
  open,
  close,
  high,
  low,
  // TODO временное решение по просьбе PO обнулять volume для прайм инструментов на графике
  // В котировках volume и value значения всегда null
  volume: volume || 0,
  time: getCandleStartTime(begin),
});



import { Timeframes } from 'moex-chart';

// 1 = 1 минута
// 5 = 1 минут
// 10 = 10 минут
// 15 = 15 минут
// 30 = 30 минут
// 45 = 45 минут
// 60 = 1 час
// 240 = 4 часа
// 24 = 1 день
// 7 = 1 неделя
// 31 = 1 месяц
// 4 = 1 квартал

const MOEX_CHART_TIMEFRAMES_INTO_INTERVALS: Record<string, string> = {
  [Timeframes['1m']]: '1',
  [Timeframes['5m']]: '1',
  [Timeframes['10m']]: '1',
  [Timeframes['15m']]: '1',
  [Timeframes['30m']]: '1',
  [Timeframes['45m']]: '1',
  [Timeframes['1h']]: '60',
  [Timeframes['4h']]: '60',
  [Timeframes['1d']]: '24',
  [Timeframes['1w']]: '7',
  [Timeframes['1M']]: '31',
};

const MOEX_CHART_TIMEFRAMES_TO_ISS_POSSIBLE_TIMEFRAMES: Record<string, Timeframes> = {
  [Timeframes['1m']]: Timeframes['1m'],
  [Timeframes['5m']]: Timeframes['1m'],
  [Timeframes['10m']]: Timeframes['1m'],
  [Timeframes['15m']]: Timeframes['1m'],
  [Timeframes['30m']]: Timeframes['1m'],
  [Timeframes['45m']]: Timeframes['1m'],
  [Timeframes['1h']]: Timeframes['1h'],
  [Timeframes['4h']]: Timeframes['1h'],
  [Timeframes['1d']]: Timeframes['1d'],
  [Timeframes['1w']]: Timeframes['1w'],
  [Timeframes['1M']]: Timeframes['1M'],
};

const INTERVALS: Record<string, string> = {
  '1': '1',
  '5': '5',
  '10': '10',
  '15': '15',
  '30': '30',
  '45': '45',
  '60': '60',
  '240': '240',
  '1D': '24',
  '7D': '7',
  '1M': '31',
  '3M': '4',
};

export const chartToReqTimeConverter = (value: string) => INTERVALS[value];

export const moexChartTimeConverter = (timeframe: string) => MOEX_CHART_TIMEFRAMES_INTO_INTERVALS[timeframe];

export const moexChartToIssTimeframe = (timeframe: string) =>
  MOEX_CHART_TIMEFRAMES_TO_ISS_POSSIBLE_TIMEFRAMES[timeframe];



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 HistoryRequestState {
  untilTime?: number;
  request: Promise<Candle[] | null> | 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 getBucketStart(time: number, timeframeSeconds: number, alignmentOffset: number): number {
  return Math.floor((time - alignmentOffset) / timeframeSeconds) * timeframeSeconds + alignmentOffset;
}

// По хорошему - класс должен быть синглтоном, чтобы кормить MoexChart одинаковой датой,
// и не плодить несколько подключений на одни символа
class DataSourceProvider {
  private prevRealtimeDataArr: Candle[] = [];

  private prevRealtimeData: Candle | undefined;

  private realtimeShouldBeConvoluted = false;

  private realtimeTimeframeSeconds = 0;

  private realtimeAlignmentOffset = 0;

  private realtimeTimer: ReturnType<typeof setInterval> | null = null;

  private historyRequests = new Map<string, HistoryRequestState>();

  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;
      }

      const historyRequestKey = `${symbol}:${timeframe}`;
      const historyRequestState = this.historyRequests.get(historyRequestKey);

      if (historyRequestState && historyRequestState.untilTime === until?.time) {
        if (historyRequestState.request) {
          return historyRequestState.request;
        }

        if (until) {
          return null;
        }
      }

      cb?.(timeframe);

      const historyRequest = this.requestHistoryData({
        timeframe,
        symbol,
        until,
        indicativeData,
      });

      this.historyRequests.set(historyRequestKey, {
        untilTime: until?.time,
        request: historyRequest,
      });

      try {
        return await historyRequest;
      } finally {
        if (this.historyRequests.get(historyRequestKey)?.request === historyRequest) {
          if (until) {
            this.historyRequests.set(historyRequestKey, {
              untilTime: until.time,
              request: null,
            });
          } else {
            this.historyRequests.delete(historyRequestKey);
          }
        }
      }
    };

  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();
      const symbolIds = getSymbols();

      Promise.all(
        symbolIds.map(async (symbolId) => {
          const symbol = getRequestSymbol(symbolId);

          if (!symbol) {
            return;
          }

          const data = await requestRealtimeBars({
            currencyPair: symbol.replaceAll(':', '.'),
            interval: moexChartTimeConverter(timeframe),
            ticker: symbol,
            indicativeData,
          });

          if (!data) {
            return;
          }

          if (!this.realtimeShouldBeConvoluted) {
            update(symbol, data);
            return;
          }

          this.realtimeConvolution(data, (candle) => {
            update(symbol, candle);
          });
        }),
      );
    }, periodMs);

    return () => {
      if (this.realtimeTimer) {
        clearInterval(this.realtimeTimer);
      }

      this.realtimeTimer = null;
    };
  }

  private async requestHistoryData({
    timeframe,
    symbol,
    until,
    indicativeData,
  }: {
    timeframe: Timeframes;
    symbol: string;
    until?: Candle;
    indicativeData?: ChartIndicativeData;
  }): Promise<Candle[] | null> {
    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: Date.now(),
        countBack: 2000,
      },
      ticker: symbol,
      indicativeData,
    });

    if (data.length === 0) {
      return null;
    }

    const issTimeframe = moexChartToIssTimeframe(timeframe);

    if (issTimeframe === timeframe) {
      this.realtimeShouldBeConvoluted = false;

      return data;
    }

    this.realtimeShouldBeConvoluted = true;

    return this.timeframeConvolution(data, timeframe);
  }

  private timeframeConvolution(data: Candle[], requestedTimeframe: Timeframes): Candle[] {
    const timeframeSeconds = getTimeframeSeconds(requestedTimeframe);
    const sortedData = [...data].sort((first, second) => first.time - second.time);

    const firstCandle = sortedData[0];

    if (!firstCandle) {
      return [];
    }

    const sessionStartIndex = sortedData.findIndex((candle, index) => {
      const previousCandle = sortedData[index - 1];

      return previousCandle && candle.time - previousCandle.time > timeframeSeconds;
    });

    const sessionStart = sessionStartIndex > 0 ? sortedData[sessionStartIndex] : firstCandle;
    const alignmentOffset = sessionStart.time % timeframeSeconds;

    this.realtimeTimeframeSeconds = timeframeSeconds;
    this.realtimeAlignmentOffset = alignmentOffset;

    const result: Candle[] = [];

    let bucketStart: number | null = null;
    let aggregatedCandle: Candle | null = null;

    sortedData.forEach((candle) => {
      const currentBucketStart = getBucketStart(candle.time, timeframeSeconds, alignmentOffset);

      if (currentBucketStart !== bucketStart) {
        if (aggregatedCandle) {
          result.push(aggregatedCandle);
        }

        bucketStart = currentBucketStart;
        aggregatedCandle = {
          ...candle,
          time: currentBucketStart,
        };

        return;
      }

      if (aggregatedCandle) {
        aggregatedCandle = {
          ...aggregatedCandle,
          high: Math.max(aggregatedCandle.high, candle.high),
          low: Math.min(aggregatedCandle.low, candle.low),
          close: candle.close,
          volume: (aggregatedCandle.volume ?? 0) + (candle.volume ?? 0),
        };
      }
    });

    if (aggregatedCandle) {
      result.push(aggregatedCandle);
    }

    return result;
  }

  private realtimeConvolution(data: Candle, update: (candle: Candle) => void): void {
    const bucketStart = getBucketStart(
      data.time,
      this.realtimeTimeframeSeconds,
      this.realtimeAlignmentOffset,
    );

    if (!this.prevRealtimeData) {
      this.prevRealtimeData = data;
      this.prevRealtimeDataArr = [data];

      update({
        ...data,
        time: bucketStart,
      });

      return;
    }

    const previousBucketStart = getBucketStart(
      this.prevRealtimeData.time,
      this.realtimeTimeframeSeconds,
      this.realtimeAlignmentOffset,
    );

    if (bucketStart !== previousBucketStart) {
      this.prevRealtimeDataArr = [data];
    } else {
      const currentCandleIndex = this.prevRealtimeDataArr.findIndex((candle) => candle.time === data.time);

      if (currentCandleIndex === -1) {
        this.prevRealtimeDataArr.push(data);
      } else {
        this.prevRealtimeDataArr[currentCandleIndex] = data;
      }
    }

    this.prevRealtimeData = data;

    const firstCandle = this.prevRealtimeDataArr[0];
    const lastCandle = this.prevRealtimeDataArr[this.prevRealtimeDataArr.length - 1];

    if (!firstCandle || !lastCandle) {
      return;
    }

    update({
      time: bucketStart,
      open: firstCandle.open,
      high: Math.max(...this.prevRealtimeDataArr.map(({ high }) => high)),
      low: Math.min(...this.prevRealtimeDataArr.map(({ low }) => low)),
      close: lastCandle.close,
      volume: this.prevRealtimeDataArr.reduce(
        (total, candle) => total + (candle.volume ?? 0),
        0,
      ),
    });
  }
}

export { DataSourceProvider };



if (customResolver) {
  try {
    const customBars = await customResolver({
      ticker,
      currencyPair,
      periodParams,
      interval,
    });

    const olderBars = customBars.filter((bar) => bar.time < periodParams.to);

    onHistoryCallback?.(olderBars, {
      noData: olderBars.length === 0,
    });

    return olderBars;
  } catch (e) {
    onHistoryCallback?.([], {
      noData: true,
    });

    return [];
  }
}