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


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

const DAY_SECONDS = 24 * 60 * 60;

interface HistoryRequestState {
  untilTime?: number;
  request: Promise<Candle[] | null> | null;
}

interface RealtimeState {
  prevRealtimeDataArr: Candle[];
  prevRealtimeData?: Candle;
}

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 getSessionOffset(data: Candle[]): number {
  const firstCandle = data[0];

  if (!firstCandle) {
    return 0;
  }

  return data.reduce((minOffset, candle) => {
    const offset = ((candle.time % DAY_SECONDS) + DAY_SECONDS) % DAY_SECONDS;

    return Math.min(minOffset, offset);
  }, ((firstCandle.time % DAY_SECONDS) + DAY_SECONDS) % DAY_SECONDS);
}

function getBucketStart(time: number, timeframeSeconds: number, sessionOffset: number): number {
  const dayStart = Math.floor(time / DAY_SECONDS) * DAY_SECONDS;
  let sessionStart = dayStart + sessionOffset;

  if (time < sessionStart) {
    sessionStart -= DAY_SECONDS;
  }

  return sessionStart + Math.floor((time - sessionStart) / timeframeSeconds) * timeframeSeconds;
}

function aggregateCandles(candles: Candle[], time: number): Candle | undefined {
  const firstCandle = candles[0];
  const lastCandle = candles[candles.length - 1];

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

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

// По хорошему - класс должен быть синглтоном, чтобы кормить MoexChart одинаковой датой,
// и не плодить несколько подключений на одни символа
class DataSourceProvider {
  private realtimeTimer: ReturnType<typeof setInterval> | null = null;

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

  private realtimeStates = new Map<string, RealtimeState>();

  private sessionOffset: number | 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;
      }

      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();
      const shouldConvolute = moexChartToIssTimeframe(timeframe) !== timeframe;

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

          const realtimeState = this.getRealtimeState(symbol, timeframe);

          if (
            realtimeState.prevRealtimeData &&
            JSON.stringify(data) === JSON.stringify(realtimeState.prevRealtimeData)
          ) {
            return;
          }

          if (!shouldConvolute) {
            realtimeState.prevRealtimeData = data;
            update(symbol, data);
            return;
          }

          this.realtimeConvolution(symbol, timeframe, 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: Math.round(Date.now() / 1000),
        countBack: 2000,
      },
      ticker: symbol,
      indicativeData,
    });

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

    const issTimeframe = moexChartToIssTimeframe(timeframe);
    const realtimeState = this.getRealtimeState(symbol, timeframe);

    if (issTimeframe === timeframe) {
      if (!until) {
        realtimeState.prevRealtimeData = data[data.length - 1];
        realtimeState.prevRealtimeDataArr = [];
      }

      return data;
    }

    if (this.sessionOffset === null) {
      this.sessionOffset = getSessionOffset(data);
    }

    console.log(
      '[DataSourceProvider][history-before-convolution]',
      JSON.stringify(
        {
          symbol,
          timeframe,
          issTimeframe,
          sessionOffset: this.sessionOffset,
          sessionTime: new Date(this.sessionOffset * 1000).toISOString().slice(11, 19),
          first: data.slice(0, 12).map((candle) => ({
            ...candle,
            utc: new Date(candle.time * 1000).toISOString(),
          })),
          last: data.slice(-12).map((candle) => ({
            ...candle,
            utc: new Date(candle.time * 1000).toISOString(),
          })),
        },
        null,
        2,
      ),
    );

    const result = this.timeframeConvolution(data, timeframe, symbol, !until);

    console.log(
      '[DataSourceProvider][history-after-convolution]',
      JSON.stringify(
        {
          symbol,
          timeframe,
          sessionOffset: this.sessionOffset,
          sessionTime: new Date(this.sessionOffset * 1000).toISOString().slice(11, 19),
          first: result.slice(0, 12).map((candle) => ({
            ...candle,
            utc: new Date(candle.time * 1000).toISOString(),
          })),
          last: result.slice(-12).map((candle) => ({
            ...candle,
            utc: new Date(candle.time * 1000).toISOString(),
          })),
        },
        null,
        2,
      ),
    );

    return result;
  }

  private timeframeConvolution(
    data: Candle[],
    requestedTimeframe: Timeframes,
    symbol: string,
    syncRealtime: boolean,
  ): Candle[] {
    const timeframeSeconds = getTimeframeSeconds(requestedTimeframe);
    const sessionOffset = this.sessionOffset ?? getSessionOffset(data);
    const sortedData = [...data].sort((first, second) => first.time - second.time);
    const result: Candle[] = [];

    let bucketStart: number | null = null;
    let candleGroup: Candle[] = [];

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

      if (bucketStart !== null && currentBucketStart !== bucketStart) {
        const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);

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

        candleGroup = [];
      }

      bucketStart = currentBucketStart;
      candleGroup.push(candle);
    });

    if (bucketStart !== null) {
      const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);

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

    if (syncRealtime) {
      const realtimeState = this.getRealtimeState(symbol, requestedTimeframe);

      realtimeState.prevRealtimeDataArr = candleGroup.slice();
      realtimeState.prevRealtimeData = sortedData[sortedData.length - 1];
    }

    return result;
  }

  private realtimeConvolution(
    symbol: string,
    timeframe: Timeframes,
    data: Candle,
    update: (candle: Candle) => void,
  ): void {
    const timeframeSeconds = getTimeframeSeconds(timeframe);
    const sessionOffset = this.sessionOffset ?? getSessionOffset([data]);
    const realtimeState = this.getRealtimeState(symbol, timeframe);
    const bucketStart = getBucketStart(data.time, timeframeSeconds, sessionOffset);
    const previousBucketStart = realtimeState.prevRealtimeData
      ? getBucketStart(realtimeState.prevRealtimeData.time, timeframeSeconds, sessionOffset)
      : null;

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

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

    realtimeState.prevRealtimeData = data;

    const candle = aggregateCandles(realtimeState.prevRealtimeDataArr, bucketStart);

    if (candle) {
      update(candle);
    }
  }

  private getRealtimeState(symbol: string, timeframe: Timeframes): RealtimeState {
    const key = `${symbol}:${timeframe}`;
    let state = this.realtimeStates.get(key);

    if (!state) {
      state = {
        prevRealtimeDataArr: [],
      };

      this.realtimeStates.set(key, state);
    }

    return state;
  }
}

export { DataSourceProvider };