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


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