Загрузка данных
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;
}
interface RealtimeState {
prevRealtimeDataArr: Candle[];
prevRealtimeData?: Candle;
realtimeShouldBeConvoluted: boolean;
realtimeSessionStart: 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 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>();
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;
}
const realtimeState = this.getRealtimeState(symbol, timeframe);
if (
realtimeState.prevRealtimeData &&
JSON.stringify(data) === JSON.stringify(realtimeState.prevRealtimeData)
) {
return;
}
if (!realtimeState.realtimeShouldBeConvoluted) {
realtimeState.prevRealtimeData = data;
update(symbol, data);
return;
}
this.realtimeConvolution(timeframe, data, realtimeState, (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) {
realtimeState.realtimeShouldBeConvoluted = false;
if (!until) {
realtimeState.prevRealtimeData = data[data.length - 1];
realtimeState.prevRealtimeDataArr = [];
realtimeState.realtimeSessionStart = null;
}
return data;
}
realtimeState.realtimeShouldBeConvoluted = true;
return this.timeframeConvolution(data, timeframe, !until, realtimeState);
}
private timeframeConvolution(
data: Candle[],
requestedTimeframe: Timeframes,
syncRealtime: boolean,
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 candleGroup: Candle[] = [];
sortedData.forEach((candle, index) => {
const previousCandle = sortedData[index - 1];
const isNewSession = previousCandle && candle.time - previousCandle.time > timeframeSeconds;
if (isNewSession) {
const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);
if (aggregatedCandle) {
result.push(aggregatedCandle);
}
sessionStart = candle.time;
bucketStart = candle.time;
candleGroup = [candle];
return;
}
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);
});
const aggregatedCandle = aggregateCandles(candleGroup, bucketStart);
if (aggregatedCandle) {
result.push(aggregatedCandle);
}
if (syncRealtime) {
realtimeState.realtimeSessionStart = sessionStart;
realtimeState.prevRealtimeDataArr = [...candleGroup];
realtimeState.prevRealtimeData = sortedData[sortedData.length - 1];
}
return result;
}
private realtimeConvolution(
timeframe: Timeframes,
data: Candle,
realtimeState: RealtimeState,
update: (candle: Candle) => void,
): void {
const timeframeSeconds = getTimeframeSeconds(timeframe);
if (!realtimeState.prevRealtimeData || realtimeState.realtimeSessionStart === null) {
realtimeState.realtimeSessionStart = data.time;
realtimeState.prevRealtimeDataArr = [data];
} else {
const isNewSession = data.time - realtimeState.prevRealtimeData.time > timeframeSeconds;
if (isNewSession) {
realtimeState.realtimeSessionStart = data.time;
realtimeState.prevRealtimeDataArr = [data];
} else {
const previousBucketStart =
realtimeState.realtimeSessionStart +
Math.floor(
(realtimeState.prevRealtimeData.time - realtimeState.realtimeSessionStart) / timeframeSeconds,
) *
timeframeSeconds;
const currentBucketStart =
realtimeState.realtimeSessionStart +
Math.floor((data.time - realtimeState.realtimeSessionStart) / timeframeSeconds) * timeframeSeconds;
if (currentBucketStart !== previousBucketStart) {
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 sessionStart = realtimeState.realtimeSessionStart;
if (sessionStart === null) {
return;
}
const bucketStart = sessionStart + Math.floor((data.time - sessionStart) / timeframeSeconds) * timeframeSeconds;
const candle = aggregateCandles(realtimeState.prevRealtimeDataArr, bucketStart);
if (candle) {
update(candle);
}
}
private getRealtimeState(symbol: string, timeframe: Timeframes): RealtimeState {
const key = `${symbol}:${timeframe}`;
let realtimeState = this.realtimeStates.get(key);
if (!realtimeState) {
realtimeState = {
prevRealtimeDataArr: [],
prevRealtimeData: undefined,
realtimeShouldBeConvoluted: moexChartToIssTimeframe(timeframe) !== timeframe,
realtimeSessionStart: null,
};
this.realtimeStates.set(key, realtimeState);
}
return realtimeState;
}
}
export { DataSourceProvider };
import { DEFAULT_SYMBOL } from '@widgets/Chart/const';
import type { Timeframes as TimeframesType } from 'moex-chart';
type DataSourceProvideModule = typeof import('@widgets/Chart/components/MoexChart/dataSourceProvide');
const mockRequestBars = jest.fn();
const mockRequestRealtimeBars = jest.fn();
const mockMoexChartTimeConverter = jest.fn();
const mockParseTimeframe = jest.fn();
const mockMoexChartToIssTimeframe = jest.fn();
jest.mock('moex-chart', () => ({
Timeframes: {
'1m': '1m',
'5m': '5m',
'1h': '1h',
'2h': '2h',
'3h': '3h',
'4h': '4h',
},
parseTimeframe: mockParseTimeframe,
}));
jest.mock('@utils/chartToReqTimeConverter', () => ({
moexChartTimeConverter: mockMoexChartTimeConverter,
moexChartToIssTimeframe: mockMoexChartToIssTimeframe,
}));
jest.mock('../requestBars', () => ({
requestBars: mockRequestBars,
requestRealtimeBars: mockRequestRealtimeBars,
}));
jest.mock('@widgets/Chart/requestBars', () => ({
requestBars: mockRequestBars,
requestRealtimeBars: mockRequestRealtimeBars,
}));
const { DataSourceProvider } = jest.requireActual(
'@widgets/Chart/components/MoexChart/dataSourceProvide',
) as DataSourceProvideModule;
const Timeframes = {
'1m': '1m' as TimeframesType,
'5m': '5m' as TimeframesType,
'1h': '1h' as TimeframesType,
'2h': '2h' as TimeframesType,
'3h': '3h' as TimeframesType,
'4h': '4h' as TimeframesType,
};
const baseTime = Math.floor(Date.parse('2026-05-19T10:00:00Z') / 1000);
const mockBar = {
time: baseTime,
open: 100,
close: 110,
high: 120,
low: 90,
volume: 1000,
};
const createBar = (minute: number, overrides: Partial<typeof mockBar> = {}) => ({
time: baseTime + minute * 60,
open: 100 + minute,
close: 101 + minute,
high: 102 + minute,
low: 99 - minute,
volume: minute + 1,
...overrides,
});
const createHourBar = (startTime: number, hour: number, priceOffset = 0) => ({
time: startTime + hour * 60 * 60,
open: 100 + priceOffset + hour,
close: 101 + priceOffset + hour,
high: 102 + priceOffset + hour,
low: 99 + priceOffset - hour,
volume: hour + 1,
});
const flushPromises = async (): Promise<void> => {
await Promise.resolve();
await Promise.resolve();
await Promise.resolve();
};
const runRealtimeTick = async (periodMs = 1000): Promise<void> => {
jest.advanceTimersByTime(periodMs);
await flushPromises();
};
const configureConvolution = (
requestedTimeframe: TimeframesType,
sourceTimeframe: TimeframesType,
): void => {
mockMoexChartToIssTimeframe.mockImplementation((timeframe: TimeframesType) =>
timeframe === requestedTimeframe ? sourceTimeframe : timeframe,
);
};
describe('DataSourceProvider', () => {
beforeEach(() => {
jest.clearAllMocks();
jest.useFakeTimers();
jest.setSystemTime(new Date('2026-05-19T10:00:00Z'));
mockMoexChartTimeConverter.mockReturnValue('1');
mockMoexChartToIssTimeframe.mockImplementation((timeframe: TimeframesType) => timeframe);
mockParseTimeframe.mockImplementation((timeframe: TimeframesType) => {
switch (timeframe) {
case Timeframes['5m']:
return {
candleWidth: 5,
dayjsUnit: 'minute',
};
case Timeframes['1h']:
return {
candleWidth: 1,
dayjsUnit: 'hour',
};
case Timeframes['2h']:
return {
candleWidth: 2,
dayjsUnit: 'hour',
};
case Timeframes['3h']:
return {
candleWidth: 3,
dayjsUnit: 'hour',
};
case Timeframes['4h']:
return {
candleWidth: 4,
dayjsUnit: 'hour',
};
default:
return {
candleWidth: 1,
dayjsUnit: 'minute',
};
}
});
});
afterEach(() => {
jest.clearAllTimers();
jest.useRealTimers();
});
it('should request chart history data with converted timeframe', async () => {
const mockTimeframeCallback = jest.fn();
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource(undefined, mockTimeframeCallback);
const result = await dataSource(Timeframes['1m'], 'MOEX:SBER');
const now = Math.round(Date.now() / 1000);
expect(mockTimeframeCallback).toHaveBeenCalledWith(Timeframes['1m']);
expect(mockMoexChartTimeConverter).toHaveBeenCalledWith(Timeframes['1m']);
expect(mockRequestBars).toHaveBeenCalledWith({
currencyPair: 'MOEX.SBER',
interval: '1',
periodParams: {
firstDataRequest: true,
to: now,
from: now,
countBack: 2000,
},
ticker: 'MOEX:SBER',
indicativeData: undefined,
});
expect(result).toEqual([mockBar]);
});
it('should normalize symbol before requesting chart history data', async () => {
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
await dataSource(Timeframes['1m'], ' MOEX:SBER ');
expect(mockRequestBars).toHaveBeenCalledWith(
expect.objectContaining({
currencyPair: 'MOEX.SBER',
ticker: 'MOEX:SBER',
}),
);
});
it('should not request chart history data for default symbol', async () => {
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['1m'], DEFAULT_SYMBOL);
expect(result).toBeNull();
expect(mockRequestBars).not.toHaveBeenCalled();
});
it('should not request chart history data for default symbol with surrounding spaces', async () => {
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['1m'], ` ${DEFAULT_SYMBOL} `);
expect(result).toBeNull();
expect(mockRequestBars).not.toHaveBeenCalled();
});
it('should not request chart history data for empty symbol', async () => {
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['1m'], ' ');
expect(result).toBeNull();
expect(mockRequestBars).not.toHaveBeenCalled();
});
it('should request chart history data with indicative data', async () => {
const indicativeData = {
id: 1,
title: 'Test instrument',
secId: 'SBER',
instrumentName: 'SBER',
settlement: 'TQBR',
firmName: 'Test firm',
key: 'SBER_TBQR',
};
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource(indicativeData);
await dataSource(Timeframes['1m'], 'MOEX:SBER');
expect(mockRequestBars).toHaveBeenCalledWith(
expect.objectContaining({
indicativeData,
}),
);
});
it('should use until time when it is provided', async () => {
mockMoexChartTimeConverter.mockReturnValue('5');
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const until = {
time: baseTime - 60,
} as NonNullable<Parameters<typeof dataSource>[2]>;
await dataSource(Timeframes['5m'], 'MOEX:GAZP', until);
expect(mockRequestBars).toHaveBeenCalledWith(
expect.objectContaining({
periodParams: expect.objectContaining({
to: until.time,
}),
}),
);
});
it('should return null when history data is empty', async () => {
mockRequestBars.mockResolvedValue([]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['1m'], 'MOEX:SBER');
expect(result).toBeNull();
});
it('should reuse pending history request for the same symbol, timeframe and until', async () => {
let resolveRequest: ((value: (typeof mockBar)[]) => void) | undefined;
mockRequestBars.mockImplementation(
() =>
new Promise<(typeof mockBar)[]>((resolve) => {
resolveRequest = resolve;
}),
);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const until = {
time: baseTime - 60,
} as NonNullable<Parameters<typeof dataSource>[2]>;
const firstRequest = dataSource(Timeframes['1m'], 'MOEX:SBER', until);
const secondRequest = dataSource(Timeframes['1m'], 'MOEX:SBER', until);
expect(mockRequestBars).toHaveBeenCalledTimes(1);
resolveRequest?.([mockBar]);
await expect(firstRequest).resolves.toEqual([mockBar]);
await expect(secondRequest).resolves.toEqual([mockBar]);
});
it('should not repeat completed history request with the same until', async () => {
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const until = {
time: baseTime - 60,
} as NonNullable<Parameters<typeof dataSource>[2]>;
const firstResult = await dataSource(Timeframes['1m'], 'MOEX:SBER', until);
const secondResult = await dataSource(Timeframes['1m'], 'MOEX:SBER', until);
expect(firstResult).toEqual([mockBar]);
expect(secondResult).toBeNull();
expect(mockRequestBars).toHaveBeenCalledTimes(1);
});
it('should request initial history again after previous initial request is completed', async () => {
mockRequestBars.mockResolvedValue([mockBar]);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
await dataSource(Timeframes['1m'], 'MOEX:SBER');
await dataSource(Timeframes['1m'], 'MOEX:SBER');
expect(mockRequestBars).toHaveBeenCalledTimes(2);
});
it('should request realtime data and update normalized symbol', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
const realtimeBar = {
...mockBar,
volume: 500,
};
mockRequestRealtimeBars.mockResolvedValue(realtimeBar);
const unsubscribe = provider.startRealtime({
getSymbols: () => [' moex:sber '],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).toHaveBeenCalledWith({
currencyPair: 'moex.sber',
interval: '1',
ticker: 'moex:sber',
indicativeData: undefined,
});
expect(mockUpdate).toHaveBeenCalledWith('moex:sber', realtimeBar);
});
it('should not request realtime data for default symbol', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
const unsubscribe = provider.startRealtime({
getSymbols: () => [DEFAULT_SYMBOL],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).not.toHaveBeenCalled();
expect(mockUpdate).not.toHaveBeenCalled();
});
it('should not request realtime data for default symbol with surrounding spaces', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
const unsubscribe = provider.startRealtime({
getSymbols: () => [` ${DEFAULT_SYMBOL} `],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).not.toHaveBeenCalled();
expect(mockUpdate).not.toHaveBeenCalled();
});
it('should not request realtime data for empty normalized symbol', async () => {
const provider = new DataSourceProvider();
const unsubscribe = provider.startRealtime({
getSymbols: () => [' '],
getTimeframe: () => Timeframes['1m'],
update: jest.fn(),
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).not.toHaveBeenCalled();
});
it('should request realtime data with indicative data', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
const indicativeData = {
id: 1,
title: 'Test instrument',
secId: 'SBER',
instrumentName: 'SBER',
settlement: 'TQBR',
firmName: 'Test firm',
key: 'SBER_TBQR',
};
mockRequestRealtimeBars.mockResolvedValue(mockBar);
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
indicativeData,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).toHaveBeenCalledWith(
expect.objectContaining({
indicativeData,
}),
);
});
it('should not call update when realtime data is empty', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
mockRequestRealtimeBars.mockResolvedValue(undefined);
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockUpdate).not.toHaveBeenCalled();
});
it('should not request realtime data when symbols list is empty', async () => {
const provider = new DataSourceProvider();
mockRequestRealtimeBars.mockResolvedValue(mockBar);
const unsubscribe = provider.startRealtime({
getSymbols: () => [],
getTimeframe: () => Timeframes['1m'],
update: jest.fn(),
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).not.toHaveBeenCalled();
});
it('should clear realtime timer on unsubscribe', async () => {
const provider = new DataSourceProvider();
mockRequestRealtimeBars.mockResolvedValue(mockBar);
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: jest.fn(),
periodMs: 1000,
});
unsubscribe();
await runRealtimeTick();
expect(mockRequestRealtimeBars).not.toHaveBeenCalled();
});
it('should replace existing realtime timer when startRealtime is called again', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
mockRequestRealtimeBars.mockResolvedValue(mockBar);
provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).toHaveBeenCalledTimes(1);
expect(mockUpdate).toHaveBeenCalledTimes(1);
});
it('should update distinct realtime candles without convolution', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
const firstCandle = {
...mockBar,
close: 110,
};
const secondCandle = {
...mockBar,
time: mockBar.time + 60,
close: 111,
};
mockRequestRealtimeBars.mockResolvedValueOnce(firstCandle).mockResolvedValueOnce(secondCandle);
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
await runRealtimeTick();
unsubscribe();
expect(mockUpdate).toHaveBeenNthCalledWith(1, 'MOEX:SBER', firstCandle);
expect(mockUpdate).toHaveBeenNthCalledWith(2, 'MOEX:SBER', secondCandle);
});
it('should not update chart for duplicated realtime candle', async () => {
const provider = new DataSourceProvider();
const mockUpdate = jest.fn();
mockRequestRealtimeBars.mockResolvedValue(mockBar);
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['1m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
await runRealtimeTick();
unsubscribe();
expect(mockRequestRealtimeBars).toHaveBeenCalledTimes(2);
expect(mockUpdate).toHaveBeenCalledTimes(1);
expect(mockUpdate).toHaveBeenCalledWith('MOEX:SBER', mockBar);
});
it('should convolve five minute history data from one minute candles', async () => {
configureConvolution(Timeframes['5m'], Timeframes['1m']);
const historyData = Array.from({ length: 11 }, (_, minute) => createBar(minute));
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['5m'], 'MOEX:SBER');
expect(result).toEqual([
{
time: baseTime,
open: 100,
close: 105,
high: 106,
low: 95,
volume: 15,
},
{
time: baseTime + 5 * 60,
open: 105,
close: 110,
high: 111,
low: 90,
volume: 40,
},
{
time: baseTime + 10 * 60,
open: 110,
close: 111,
high: 112,
low: 89,
volume: 11,
},
]);
});
it('should convolve incomplete history candle group', async () => {
configureConvolution(Timeframes['5m'], Timeframes['1m']);
const historyData = Array.from({ length: 8 }, (_, minute) => createBar(minute));
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['5m'], 'MOEX:SBER');
expect(result).toEqual([
{
time: baseTime,
open: 100,
close: 105,
high: 106,
low: 95,
volume: 15,
},
{
time: baseTime + 5 * 60,
open: 105,
close: 108,
high: 109,
low: 92,
volume: 21,
},
]);
});
it('should start a new convolution session after a large gap', async () => {
configureConvolution(Timeframes['5m'], Timeframes['1m']);
const historyData = [
createBar(0),
createBar(1),
createBar(2),
createBar(3),
createBar(4),
createBar(5),
createBar(11),
createBar(12),
];
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['5m'], 'MOEX:SBER');
expect(result).toEqual([
{
time: baseTime,
open: 100,
close: 105,
high: 106,
low: 95,
volume: 15,
},
{
time: baseTime + 5 * 60,
open: 105,
close: 106,
high: 107,
low: 94,
volume: 6,
},
{
time: baseTime + 11 * 60,
open: 111,
close: 113,
high: 114,
low: 87,
volume: 25,
},
]);
});
it.each([
{
timeframe: Timeframes['2h'],
candlesCount: 5,
expectedHourOffsets: [0, 2, 4],
},
{
timeframe: Timeframes['3h'],
candlesCount: 7,
expectedHourOffsets: [0, 3, 6],
},
])(
'should convolve $timeframe history from one hour candles',
async ({ timeframe, candlesCount, expectedHourOffsets }) => {
configureConvolution(timeframe, Timeframes['1h']);
mockMoexChartTimeConverter.mockReturnValue('60');
const sessionStart = Math.floor(Date.parse('2026-07-31T03:00:00Z') / 1000);
const historyData = Array.from({ length: candlesCount }, (_, hour) => createHourBar(sessionStart, hour));
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(timeframe, 'MOEX:SBER');
expect(mockRequestBars).toHaveBeenCalledWith(
expect.objectContaining({
interval: '60',
}),
);
expect(result?.map(({ time }) => time)).toEqual(
expectedHourOffsets.map((hour) => sessionStart + hour * 60 * 60),
);
},
);
it('should align four hour candles to the beginning of the trading session', async () => {
configureConvolution(Timeframes['4h'], Timeframes['1h']);
mockMoexChartTimeConverter.mockReturnValue('60');
const sessionStart = Math.floor(Date.parse('2026-07-31T03:00:00Z') / 1000);
const historyData = Array.from({ length: 12 }, (_, hour) => createHourBar(sessionStart, hour));
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['4h'], 'MOEX:SBER');
expect(result?.map(({ time }) => time)).toEqual([
sessionStart,
sessionStart + 4 * 60 * 60,
sessionStart + 8 * 60 * 60,
]);
});
it('should reset four hour convolution at the next trading session', async () => {
configureConvolution(Timeframes['4h'], Timeframes['1h']);
const firstSessionStart = Math.floor(Date.parse('2026-07-30T03:00:00Z') / 1000);
const secondSessionStart = Math.floor(Date.parse('2026-07-31T03:00:00Z') / 1000);
const historyData = [
...Array.from({ length: 9 }, (_, hour) => createHourBar(firstSessionStart, hour)),
...Array.from({ length: 5 }, (_, hour) => createHourBar(secondSessionStart, hour)),
];
mockRequestBars.mockResolvedValue(historyData);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
const result = await dataSource(Timeframes['4h'], 'MOEX:SBER');
expect(result?.map(({ time }) => time)).toEqual([
firstSessionStart,
firstSessionStart + 4 * 60 * 60,
firstSessionStart + 8 * 60 * 60,
secondSessionStart,
secondSessionStart + 4 * 60 * 60,
]);
});
it('should keep realtime convolution state isolated between symbols', async () => {
configureConvolution(Timeframes['4h'], Timeframes['1h']);
mockMoexChartTimeConverter.mockReturnValue('60');
const sberSessionStart = Math.floor(Date.parse('2026-07-31T03:00:00Z') / 1000);
const compareSessionStart = Math.floor(Date.parse('2026-07-31T04:00:00Z') / 1000);
const sberHistory = Array.from({ length: 8 }, (_, hour) => createHourBar(sberSessionStart, hour));
const compareHistory = Array.from({ length: 8 }, (_, hour) =>
createHourBar(compareSessionStart, hour, 100),
);
mockRequestBars.mockImplementation(({ ticker }: { ticker: string }) =>
Promise.resolve(ticker === 'MOEX:SBER' ? sberHistory : compareHistory),
);
const provider = new DataSourceProvider();
const dataSource = provider.getDataSource();
await dataSource(Timeframes['4h'], 'MOEX:SBER');
await dataSource(Timeframes['4h'], 'MOEX:GAZP');
const realtimeCandle = {
time: sberSessionStart + 9 * 60 * 60,
open: 120,
high: 124,
low: 119,
close: 123,
volume: 50,
};
mockRequestRealtimeBars.mockResolvedValue(realtimeCandle);
const mockUpdate = jest.fn();
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['4h'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
unsubscribe();
expect(mockUpdate).toHaveBeenCalledWith('MOEX:SBER', {
...realtimeCandle,
time: sberSessionStart + 8 * 60 * 60,
});
});
it('should append, replace and reset realtime candles during convolution', async () => {
configureConvolution(Timeframes['5m'], Timeframes['1m']);
mockRequestBars.mockResolvedValue(Array.from({ length: 11 }, (_, minute) => createBar(minute)));
const provider = new DataSourceProvider();
await provider.getDataSource()(Timeframes['5m'], 'MOEX:SBER');
const firstCandle = {
time: baseTime + 10 * 60,
open: 100,
high: 103,
low: 99,
close: 102,
volume: 10,
};
const secondCandle = {
time: baseTime + 11 * 60,
open: 102,
high: 106,
low: 98,
close: 105,
volume: 20,
};
const updatedSecondCandle = {
...secondCandle,
high: 107,
low: 97,
close: 106,
volume: 25,
};
const nextTimeframeCandle = {
time: baseTime + 15 * 60,
open: 106,
high: 108,
low: 105,
close: 107,
volume: 30,
};
mockRequestRealtimeBars
.mockResolvedValueOnce(firstCandle)
.mockResolvedValueOnce(secondCandle)
.mockResolvedValueOnce(updatedSecondCandle)
.mockResolvedValueOnce(nextTimeframeCandle);
const mockUpdate = jest.fn();
const unsubscribe = provider.startRealtime({
getSymbols: () => ['MOEX:SBER'],
getTimeframe: () => Timeframes['5m'],
update: mockUpdate,
periodMs: 1000,
});
await runRealtimeTick();
await runRealtimeTick();
await runRealtimeTick();
await runRealtimeTick();
unsubscribe();
expect(mockUpdate).toHaveBeenNthCalledWith(1, 'MOEX:SBER', {
time: baseTime + 10 * 60,
open: 110,
high: 112,
low: 89,
close: 102,
volume: 21,
});
expect(mockUpdate).toHaveBeenNthCalledWith(2, 'MOEX:SBER', {
time: baseTime + 10 * 60,
open: 110,
high: 106,
low: 89,
close: 105,
volume: 41,
});
expect(mockUpdate).toHaveBeenNthCalledWith(3, 'MOEX:SBER', {
time: baseTime + 10 * 60,
open: 110,
high: 107,
low: 89,
close: 106,
volume: 46,
});
expect(mockUpdate).toHaveBeenNthCalledWith(4, 'MOEX:SBER', nextTimeframeCandle);
});
});