From a0f549d81b117710a7cc024f39f95b8dc63952d8 Mon Sep 17 00:00:00 2001 From: Thomas Kaul <4159106+dtslvr@users.noreply.github.com> Date: Mon, 24 Aug 2026 14:13:47 +0200 Subject: [PATCH] Fix repeated historical market data gathering on weekends --- .../portfolio/current-rate.service.spec.ts | 3 + .../market-data/market-data.service.spec.ts | 116 +++++++++ .../market-data/market-data.service.ts | 28 ++- .../data-gathering.processor.spec.ts | 236 ++++++++++++++++++ .../data-gathering.processor.ts | 12 +- .../data-gathering.service.spec.ts | 189 ++++++++++++++ .../data-gathering/data-gathering.service.ts | 41 ++- libs/common/src/lib/config.ts | 5 +- .../migration.sql | 2 + prisma/schema.prisma | 15 +- 10 files changed, 606 insertions(+), 41 deletions(-) create mode 100644 apps/api/src/services/market-data/market-data.service.spec.ts create mode 100644 apps/api/src/services/queues/data-gathering/data-gathering.processor.spec.ts create mode 100644 apps/api/src/services/queues/data-gathering/data-gathering.service.spec.ts create mode 100644 prisma/migrations/20260824000000_added_is_carried_forward_to_market_data/migration.sql diff --git a/apps/api/src/app/portfolio/current-rate.service.spec.ts b/apps/api/src/app/portfolio/current-rate.service.spec.ts index 5f2358679..ea93e8b4f 100644 --- a/apps/api/src/app/portfolio/current-rate.service.spec.ts +++ b/apps/api/src/app/portfolio/current-rate.service.spec.ts @@ -20,6 +20,7 @@ jest.mock('@ghostfolio/api/services/market-data/market-data.service', () => { createdAt: date, dataSource: DataSource.YAHOO, id: 'aefcbe3a-ee10-4c4f-9f2d-8ffad7b05584', + isCarriedForward: false, marketPrice: 1847.839966, state: 'CLOSE' }); @@ -39,6 +40,7 @@ jest.mock('@ghostfolio/api/services/market-data/market-data.service', () => { dataSource: assetProfileIdentifiers[0].dataSource, date: dateQuery.gte, id: '8fa48fde-f397-4b0d-adbc-fb940e830e6d', + isCarriedForward: false, marketPrice: 1841.823902, state: 'CLOSE', symbol: assetProfileIdentifiers[0].symbol @@ -48,6 +50,7 @@ jest.mock('@ghostfolio/api/services/market-data/market-data.service', () => { dataSource: assetProfileIdentifiers[0].dataSource, date: dateQuery.lt, id: '082d6893-df27-4c91-8a5d-092e84315b56', + isCarriedForward: false, marketPrice: 1847.839966, state: 'CLOSE', symbol: assetProfileIdentifiers[0].symbol diff --git a/apps/api/src/services/market-data/market-data.service.spec.ts b/apps/api/src/services/market-data/market-data.service.spec.ts new file mode 100644 index 000000000..3af85f96e --- /dev/null +++ b/apps/api/src/services/market-data/market-data.service.spec.ts @@ -0,0 +1,116 @@ +import { parseDate } from '@ghostfolio/common/helper'; + +import { MarketDataService } from './market-data.service'; + +describe('MarketDataService', () => { + let marketDataService: MarketDataService; + let prismaService: { + $transaction: jest.Mock; + marketData: { + createMany: jest.Mock; + deleteMany: jest.Mock; + upsert: jest.Mock; + }; + }; + + beforeEach(() => { + prismaService = { + $transaction: jest.fn(), + marketData: { + createMany: jest.fn(), + deleteMany: jest.fn(), + upsert: jest.fn().mockResolvedValue({}) + } + }; + + marketDataService = new MarketDataService(prismaService as any); + }); + + describe('updateMany', () => { + it('does not drop isCarriedForward', async () => { + prismaService.$transaction.mockImplementation((promises) => { + return Promise.all(promises); + }); + + await marketDataService.updateMany({ + data: [ + { + dataSource: 'YAHOO', + date: parseDate('2026-08-22'), + isCarriedForward: true, + marketPrice: 100, + state: 'CLOSE', + symbol: 'AAPL' + } + ] + }); + + expect(prismaService.marketData.upsert).toHaveBeenCalledWith( + expect.objectContaining({ + create: expect.objectContaining({ isCarriedForward: true }), + update: expect.objectContaining({ isCarriedForward: true }) + }) + ); + }); + }); + + describe('replaceForSymbol', () => { + it('does not drop isCarriedForward', async () => { + prismaService.$transaction.mockImplementation((callback) => { + return callback(prismaService); + }); + + await marketDataService.replaceForSymbol({ + data: [ + { + dataSource: 'YAHOO', + date: parseDate('2026-08-21'), + isCarriedForward: false, + marketPrice: 100, + state: 'CLOSE', + symbol: 'AAPL' + }, + { + dataSource: 'YAHOO', + date: parseDate('2026-08-22'), + isCarriedForward: true, + marketPrice: 100, + state: 'CLOSE', + symbol: 'AAPL' + } + ], + dataSource: 'YAHOO', + symbol: 'AAPL' + }); + + const { data } = prismaService.marketData.createMany.mock.calls[0][0]; + + expect(data).toEqual([ + expect.objectContaining({ isCarriedForward: false }), + expect.objectContaining({ isCarriedForward: true }) + ]); + }); + }); + + describe('updateMarketData', () => { + it('resets isCarriedForward for a manually edited market price', async () => { + await marketDataService.updateMarketData({ + data: { marketPrice: 100, state: 'CLOSE' }, + where: { + dataSource_date_symbol: { + dataSource: 'YAHOO', + date: parseDate('2026-08-22'), + symbol: 'AAPL' + } + } + }); + + expect(prismaService.marketData.upsert).toHaveBeenCalledWith( + expect.objectContaining({ + create: expect.objectContaining({ isCarriedForward: false }), + update: expect.objectContaining({ isCarriedForward: false }) + }) + ); + }); + }); +}); diff --git a/apps/api/src/services/market-data/market-data.service.ts b/apps/api/src/services/market-data/market-data.service.ts index ad388ce5c..d5330c801 100644 --- a/apps/api/src/services/market-data/market-data.service.ts +++ b/apps/api/src/services/market-data/market-data.service.ts @@ -189,13 +189,16 @@ export class MarketDataService { }); await prisma.marketData.createMany({ - data: data.map(({ date, marketPrice, state }) => ({ - dataSource, - symbol, - date: date as Date, - marketPrice: marketPrice as number, - state: state as MarketDataState - })), + data: data.map( + ({ date, isCarriedForward, marketPrice, state }) => ({ + dataSource, + symbol, + date: date as Date, + isCarriedForward: isCarriedForward as boolean, + marketPrice: marketPrice as number, + state: state as MarketDataState + }) + ), skipDuplicates: true }); } @@ -233,11 +236,16 @@ export class MarketDataService { create: { dataSource: where.dataSource_date_symbol.dataSource, date: where.dataSource_date_symbol.date, + isCarriedForward: false, marketPrice: data.marketPrice, state: data.state, symbol: where.dataSource_date_symbol.symbol }, - update: { marketPrice: data.marketPrice, state: data.state } + update: { + isCarriedForward: false, + marketPrice: data.marketPrice, + state: data.state + } }); } @@ -251,16 +259,18 @@ export class MarketDataService { data: Prisma.MarketDataUpdateInput[]; }): Promise { const upsertPromises = data.map( - ({ dataSource, date, marketPrice, symbol, state }) => { + ({ dataSource, date, isCarriedForward, marketPrice, symbol, state }) => { return this.prismaService.marketData.upsert({ create: { dataSource: dataSource as DataSource, date: date as Date, + isCarriedForward: isCarriedForward as boolean, marketPrice: marketPrice as number, state: state as MarketDataState, symbol: symbol as string }, update: { + isCarriedForward: isCarriedForward as boolean, marketPrice: marketPrice as number, state: state as MarketDataState }, diff --git a/apps/api/src/services/queues/data-gathering/data-gathering.processor.spec.ts b/apps/api/src/services/queues/data-gathering/data-gathering.processor.spec.ts new file mode 100644 index 000000000..acba8c311 --- /dev/null +++ b/apps/api/src/services/queues/data-gathering/data-gathering.processor.spec.ts @@ -0,0 +1,236 @@ +import { DataGatheringItem } from '@ghostfolio/api/services/interfaces/interfaces'; +import { + getAssetProfileIdentifier, + parseDate +} from '@ghostfolio/common/helper'; + +import { DataSource } from '@prisma/client'; +import { Job } from 'bull'; + +import { DataGatheringProcessor } from './data-gathering.processor'; + +describe('DataGatheringProcessor', () => { + let dataGatheringProcessor: DataGatheringProcessor; + let dataProviderService: { getHistoricalRaw: jest.Mock }; + let marketDataService: { replaceForSymbol: jest.Mock; updateMany: jest.Mock }; + + const createJob = ({ + dataSource, + date, + symbol + }: { + dataSource: DataSource; + date: string; + symbol: string; + }) => { + return { + data: { + dataSource, + symbol, + date: parseDate(date).toISOString() + } + } as unknown as Job; + }; + + const mockHistoricalData = ({ + dataSource, + prices, + symbol + }: { + dataSource: DataSource; + prices: { [date: string]: number }; + symbol: string; + }) => { + const assetProfileIdentifier = getAssetProfileIdentifier({ + dataSource, + symbol + }); + + const historicalData: { + [symbol: string]: { [date: string]: { marketPrice: number } }; + } = { [assetProfileIdentifier]: {} }; + + for (const [date, marketPrice] of Object.entries(prices)) { + historicalData[assetProfileIdentifier][date] = { marketPrice }; + } + + dataProviderService.getHistoricalRaw.mockResolvedValue(historicalData); + }; + + beforeAll(() => { + jest.useFakeTimers().setSystemTime(parseDate('2026-08-24').getTime()); + }); + + beforeEach(() => { + dataProviderService = { getHistoricalRaw: jest.fn() }; + marketDataService = { + replaceForSymbol: jest.fn(), + updateMany: jest.fn() + }; + + dataGatheringProcessor = new DataGatheringProcessor( + null, + dataProviderService as any, + marketDataService as any, + null + ); + }); + + afterAll(() => { + jest.useRealTimers(); + }); + + it('writes an all-real series without carried forward market prices', async () => { + mockHistoricalData({ + dataSource: 'COINGECKO', + symbol: 'bitcoin', + prices: { + '2026-08-17': 1, + '2026-08-18': 2, + '2026-08-19': 3, + '2026-08-20': 4, + '2026-08-21': 5, + '2026-08-22': 6, + '2026-08-23': 7 + } + }); + + await dataGatheringProcessor.gatherHistoricalMarketData( + createJob({ + dataSource: 'COINGECKO', + date: '2026-08-17', + symbol: 'bitcoin' + }) + ); + + const { data } = marketDataService.updateMany.mock.calls[0][0]; + + expect(data).toHaveLength(7); + expect( + data.every(({ isCarriedForward }) => { + return isCarriedForward === false; + }) + ).toBe(true); + }); + + it('fills an interior gap with carried forward market prices', async () => { + mockHistoricalData({ + dataSource: 'YAHOO', + symbol: 'AAPL', + prices: { + '2026-08-17': 1, + '2026-08-18': 2, + '2026-08-21': 5, + '2026-08-22': 6, + '2026-08-23': 7 + } + }); + + await dataGatheringProcessor.gatherHistoricalMarketData( + createJob({ dataSource: 'YAHOO', date: '2026-08-17', symbol: 'AAPL' }) + ); + + const { data } = marketDataService.updateMany.mock.calls[0][0]; + + expect(data).toEqual([ + expect.objectContaining({ + date: parseDate('2026-08-17'), + isCarriedForward: false, + marketPrice: 1 + }), + expect.objectContaining({ + date: parseDate('2026-08-18'), + isCarriedForward: false, + marketPrice: 2 + }), + expect.objectContaining({ + date: parseDate('2026-08-19'), + isCarriedForward: true, + marketPrice: 2 + }), + expect.objectContaining({ + date: parseDate('2026-08-20'), + isCarriedForward: true, + marketPrice: 2 + }), + expect.objectContaining({ + date: parseDate('2026-08-21'), + isCarriedForward: false, + marketPrice: 5 + }), + expect.objectContaining({ + date: parseDate('2026-08-22'), + isCarriedForward: false, + marketPrice: 6 + }), + expect.objectContaining({ + date: parseDate('2026-08-23'), + isCarriedForward: false, + marketPrice: 7 + }) + ]); + }); + + it('fills a trailing gap with carried forward market prices', async () => { + mockHistoricalData({ + dataSource: 'YAHOO', + symbol: 'AAPL', + prices: { + '2026-08-17': 1, + '2026-08-18': 2, + '2026-08-19': 3, + '2026-08-20': 4, + '2026-08-21': 5 + } + }); + + await dataGatheringProcessor.gatherHistoricalMarketData( + createJob({ dataSource: 'YAHOO', date: '2026-08-17', symbol: 'AAPL' }) + ); + + const { data } = marketDataService.updateMany.mock.calls[0][0]; + + expect(data).toHaveLength(7); + expect(data.slice(5)).toEqual([ + expect.objectContaining({ + date: parseDate('2026-08-22'), + isCarriedForward: true, + marketPrice: 5 + }), + expect.objectContaining({ + date: parseDate('2026-08-23'), + isCarriedForward: true, + marketPrice: 5 + }) + ]); + }); + + it('does not fill a leading gap', async () => { + mockHistoricalData({ + dataSource: 'YAHOO', + symbol: 'AAPL', + prices: { + '2026-08-19': 3, + '2026-08-20': 4, + '2026-08-21': 5, + '2026-08-22': 6, + '2026-08-23': 7 + } + }); + + await dataGatheringProcessor.gatherHistoricalMarketData( + createJob({ dataSource: 'YAHOO', date: '2026-08-17', symbol: 'AAPL' }) + ); + + const { data } = marketDataService.updateMany.mock.calls[0][0]; + + expect(data).toHaveLength(5); + expect(data[0]).toEqual( + expect.objectContaining({ + date: parseDate('2026-08-19'), + isCarriedForward: false, + marketPrice: 3 + }) + ); + }); +}); diff --git a/apps/api/src/services/queues/data-gathering/data-gathering.processor.ts b/apps/api/src/services/queues/data-gathering/data-gathering.processor.ts index b49b30cc8..f44a9d56c 100644 --- a/apps/api/src/services/queues/data-gathering/data-gathering.processor.ts +++ b/apps/api/src/services/queues/data-gathering/data-gathering.processor.ts @@ -125,7 +125,6 @@ export class DataGatheringProcessor { const data: Prisma.MarketDataUpdateInput[] = []; let lastMarketPrice: number; - let numberOfMarketDataItemsToKeep = 0; while ( isBefore( @@ -154,24 +153,15 @@ export class DataGatheringProcessor { dataSource, symbol, date: getStartOfUtcDate(currentDate), + isCarriedForward: !marketPriceOfDataProvider, marketPrice: lastMarketPrice, state: 'CLOSE' }); - - if (marketPriceOfDataProvider) { - numberOfMarketDataItemsToKeep = data.length; - } } currentDate = addDays(currentDate, 1); } - // A gap at the end means that the market data is not available yet, in - // contrast to a gap in between, which means that the market was closed. - // Therefore, the market prices which are carried forward after the last - // market price of the data provider are discarded. - data.splice(numberOfMarketDataItemsToKeep); - if (force) { await this.marketDataService.replaceForSymbol({ data, diff --git a/apps/api/src/services/queues/data-gathering/data-gathering.service.spec.ts b/apps/api/src/services/queues/data-gathering/data-gathering.service.spec.ts new file mode 100644 index 000000000..09740726d --- /dev/null +++ b/apps/api/src/services/queues/data-gathering/data-gathering.service.spec.ts @@ -0,0 +1,189 @@ +import { + GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS, + GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_OPTIONS +} from '@ghostfolio/common/config'; +import { parseDate } from '@ghostfolio/common/helper'; + +import { DataGatheringService } from './data-gathering.service'; + +describe('DataGatheringService', () => { + let dataGatheringQueue: { addBulk: jest.Mock; clean: jest.Mock }; + let dataGatheringService: DataGatheringService; + let dataProviderService: { getHistoricalRaw: jest.Mock }; + let prismaService: { marketData: { groupBy: jest.Mock; upsert: jest.Mock } }; + + beforeEach(() => { + dataGatheringQueue = { + addBulk: jest.fn().mockResolvedValue([]), + clean: jest.fn().mockResolvedValue([]) + }; + dataProviderService = { getHistoricalRaw: jest.fn() }; + prismaService = { + marketData: { + groupBy: jest.fn().mockResolvedValue([]), + upsert: jest.fn().mockResolvedValue({}) + } + }; + + dataGatheringService = new DataGatheringService( + null, + dataGatheringQueue as any, + dataProviderService as any, + null, + null, + prismaService as any, + null, + null + ); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + describe('getAssetProfileIdentifiersWithRecentMarketData', () => { + it('excludes carried forward market prices from the query', async () => { + jest.useFakeTimers().setSystemTime(parseDate('2026-08-23').getTime()); + + await dataGatheringService[ + 'getAssetProfileIdentifiersWithRecentMarketData' + ](); + + expect(prismaService.marketData.groupBy).toHaveBeenCalledWith( + expect.objectContaining({ + where: expect.objectContaining({ + isCarriedForward: false, + state: 'CLOSE' + }) + }) + ); + }); + + it('keeps a cryptocurrency with a real market price of yesterday on Sunday', async () => { + jest.useFakeTimers().setSystemTime(parseDate('2026-08-23').getTime()); + + prismaService.marketData.groupBy.mockResolvedValue([ + { + _max: { date: parseDate('2026-08-22') }, + dataSource: 'COINGECKO', + symbol: 'bitcoin' + }, + { + _max: { date: parseDate('2026-08-21') }, + dataSource: 'YAHOO', + symbol: 'AAPL' + } + ]); + + const assetProfileIdentifiers = + await dataGatheringService[ + 'getAssetProfileIdentifiersWithRecentMarketData' + ](); + + expect(assetProfileIdentifiers).toEqual([ + { dataSource: 'COINGECKO', symbol: 'bitcoin' } + ]); + }); + + it('drops a stock with a real market price of Friday on Monday', async () => { + jest.useFakeTimers().setSystemTime(parseDate('2026-08-24').getTime()); + + prismaService.marketData.groupBy.mockResolvedValue([ + { + _max: { date: parseDate('2026-08-23') }, + dataSource: 'COINGECKO', + symbol: 'bitcoin' + }, + { + _max: { date: parseDate('2026-08-21') }, + dataSource: 'YAHOO', + symbol: 'AAPL' + } + ]); + + const assetProfileIdentifiers = + await dataGatheringService[ + 'getAssetProfileIdentifiersWithRecentMarketData' + ](); + + expect(assetProfileIdentifiers).toEqual([ + { dataSource: 'COINGECKO', symbol: 'bitcoin' } + ]); + }); + + it('drops a stock with a late Friday close on Saturday', async () => { + jest.useFakeTimers().setSystemTime(parseDate('2026-08-22').getTime()); + + prismaService.marketData.groupBy.mockResolvedValue([ + { + _max: { date: parseDate('2026-08-20') }, + dataSource: 'YAHOO', + symbol: 'AAPL' + }, + { + _max: { date: parseDate('2026-08-21') }, + dataSource: 'YAHOO', + symbol: 'MSFT' + } + ]); + + const assetProfileIdentifiers = + await dataGatheringService[ + 'getAssetProfileIdentifiersWithRecentMarketData' + ](); + + expect(assetProfileIdentifiers).toEqual([ + { dataSource: 'YAHOO', symbol: 'MSFT' } + ]); + }); + }); + + describe('gatherRecentMarketData', () => { + it('expires completed jobs which are older than the cooldown', async () => { + jest + .spyOn(dataGatheringService as any, 'getCurrencies7D') + .mockResolvedValue([]); + jest + .spyOn(dataGatheringService as any, 'getSymbols7D') + .mockResolvedValue([]); + + await dataGatheringService.gatherRecentMarketData(); + + expect(dataGatheringQueue.clean).toHaveBeenCalledWith( + GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS, + 'completed' + ); + }); + + it('retains completed jobs for the duration of the cooldown', () => { + expect( + GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_OPTIONS.removeOnComplete + ).toEqual({ + age: GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS / 1000 + }); + }); + }); + + describe('gatherSymbolForDate', () => { + it('resets isCarriedForward on a previously carried forward market price', async () => { + dataProviderService.getHistoricalRaw.mockResolvedValue({ + 'YAHOO-AAPL': { + '2026-08-22': { marketPrice: 100 } + } + }); + + await dataGatheringService.gatherSymbolForDate({ + dataSource: 'YAHOO', + date: parseDate('2026-08-22'), + symbol: 'AAPL' + }); + + expect(prismaService.marketData.upsert).toHaveBeenCalledWith( + expect.objectContaining({ + create: expect.objectContaining({ isCarriedForward: false }), + update: { marketPrice: 100, isCarriedForward: false } + }) + ); + }); + }); +}); diff --git a/apps/api/src/services/queues/data-gathering/data-gathering.service.ts b/apps/api/src/services/queues/data-gathering/data-gathering.service.ts index 9596b7614..63e1c022f 100644 --- a/apps/api/src/services/queues/data-gathering/data-gathering.service.ts +++ b/apps/api/src/services/queues/data-gathering/data-gathering.service.ts @@ -11,6 +11,7 @@ import { DATA_GATHERING_QUEUE_PRIORITY_HIGH, DATA_GATHERING_QUEUE_PRIORITY_LOW, DATA_GATHERING_QUEUE_PRIORITY_MEDIUM, + GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS, GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_NAME, GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_OPTIONS, PROPERTY_BENCHMARKS @@ -30,7 +31,14 @@ import { InjectQueue } from '@nestjs/bull'; import { Inject, Injectable, Logger } from '@nestjs/common'; import { Prisma } from '@prisma/client'; import { Job, JobOptions, Queue } from 'bull'; -import { format, min, subDays, subMilliseconds, subYears } from 'date-fns'; +import { + format, + isBefore, + min, + subDays, + subMilliseconds, + subYears +} from 'date-fns'; import { isEmpty } from 'lodash'; import ms, { StringValue } from 'ms'; @@ -252,6 +260,11 @@ export class DataGatheringService { } public async gatherRecentMarketData() { + await this.dataGatheringQueue.clean( + GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS, + 'completed' + ); + await this.gatherSymbols({ dataGatheringItems: await this.getCurrencies7D(), priority: DATA_GATHERING_QUEUE_PRIORITY_HIGH @@ -320,9 +333,10 @@ export class DataGatheringService { dataSource, date, marketPrice, - symbol + symbol, + isCarriedForward: false }, - update: { marketPrice }, + update: { marketPrice, isCarriedForward: false }, where: { dataSource_date_symbol: { dataSource, date, symbol } } }); } @@ -397,22 +411,23 @@ export class DataGatheringService { }); } - private async getAssetProfileIdentifiersWithCompleteMarketData(): Promise< + private async getAssetProfileIdentifiersWithRecentMarketData(): Promise< AssetProfileIdentifier[] > { return ( await this.prismaService.marketData.groupBy({ - _count: true, + _max: { date: true }, by: ['dataSource', 'symbol'], orderBy: [{ symbol: 'asc' }], where: { date: { gt: subDays(resetHours(new Date()), 7) }, + isCarriedForward: false, state: 'CLOSE' } }) ) - .filter(({ _count }) => { - return _count >= 6; + .filter(({ _max }) => { + return !isBefore(_max.date, getStartOfUtcDate(subDays(new Date(), 1))); }) .map(({ dataSource, symbol }) => { return { dataSource, symbol }; @@ -420,13 +435,13 @@ export class DataGatheringService { } private async getCurrencies7D(): Promise { - const assetProfileIdentifiersWithCompleteMarketData = - await this.getAssetProfileIdentifiersWithCompleteMarketData(); + const assetProfileIdentifiersWithRecentMarketData = + await this.getAssetProfileIdentifiersWithRecentMarketData(); return this.exchangeRateDataService .getCurrencyPairs() .filter(({ dataSource, symbol }) => { - return !assetProfileIdentifiersWithCompleteMarketData.some((item) => { + return !assetProfileIdentifiersWithRecentMarketData.some((item) => { return item.dataSource === dataSource && item.symbol === symbol; }); }) @@ -485,8 +500,8 @@ export class DataGatheringService { } ); - const assetProfileIdentifiersWithCompleteMarketData = - await this.getAssetProfileIdentifiersWithCompleteMarketData(); + const assetProfileIdentifiersWithRecentMarketData = + await this.getAssetProfileIdentifiersWithRecentMarketData(); return symbolProfiles .filter(({ dataSource, scraperConfiguration, symbol }) => { @@ -494,7 +509,7 @@ export class DataGatheringService { dataSource === 'MANUAL' && !isEmpty(scraperConfiguration); return ( - !assetProfileIdentifiersWithCompleteMarketData.some((item) => { + !assetProfileIdentifiersWithRecentMarketData.some((item) => { return item.dataSource === dataSource && item.symbol === symbol; }) && (dataSource !== 'MANUAL' || manualDataSourceWithScraperConfiguration) diff --git a/libs/common/src/lib/config.ts b/libs/common/src/lib/config.ts index 2f914eebf..155fef551 100644 --- a/libs/common/src/lib/config.ts +++ b/libs/common/src/lib/config.ts @@ -203,6 +203,7 @@ export const GATHER_ASSET_PROFILE_PROCESS_JOB_OPTIONS: JobOptions = { removeOnComplete: true }; +export const GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS = ms('12 hours'); export const GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_NAME = 'GATHER_HISTORICAL_MARKET_DATA'; export const GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_OPTIONS: JobOptions = { @@ -211,7 +212,9 @@ export const GATHER_HISTORICAL_MARKET_DATA_PROCESS_JOB_OPTIONS: JobOptions = { delay: ms('1 minute'), type: 'exponential' }, - removeOnComplete: true + removeOnComplete: { + age: GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS / 1000 + } }; export const GATHER_STATISTICS_PROCESS_JOB_OPTIONS: JobOptions = { diff --git a/prisma/migrations/20260824000000_added_is_carried_forward_to_market_data/migration.sql b/prisma/migrations/20260824000000_added_is_carried_forward_to_market_data/migration.sql new file mode 100644 index 000000000..1de434fe1 --- /dev/null +++ b/prisma/migrations/20260824000000_added_is_carried_forward_to_market_data/migration.sql @@ -0,0 +1,2 @@ +-- AlterTable +ALTER TABLE "MarketData" ADD COLUMN "isCarriedForward" BOOLEAN NOT NULL DEFAULT false; diff --git a/prisma/schema.prisma b/prisma/schema.prisma index c0dd25a60..600160124 100644 --- a/prisma/schema.prisma +++ b/prisma/schema.prisma @@ -157,13 +157,14 @@ model AuthDevice { } model MarketData { - createdAt DateTime @default(now()) - dataSource DataSource - date DateTime - id String @id @default(uuid()) - marketPrice Float - state MarketDataState @default(CLOSE) - symbol String + createdAt DateTime @default(now()) + dataSource DataSource + date DateTime + id String @id @default(uuid()) + isCarriedForward Boolean @default(false) + marketPrice Float + state MarketDataState @default(CLOSE) + symbol String @@unique([dataSource, date, symbol]) @@index([dataSource])