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 f1e8d01cd..5829a3e3d 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 @@ -21,15 +21,7 @@ import { Process, Processor } from '@nestjs/bull'; import { Injectable, Logger } from '@nestjs/common'; import { Prisma } from '@prisma/client'; import { Job } from 'bull'; -import { - addDays, - format, - getDate, - getMonth, - getYear, - isBefore, - parseISO -} from 'date-fns'; +import { addDays, format, isBefore, parseISO } from 'date-fns'; import { DataGatheringService } from './data-gathering.service'; @@ -126,19 +118,7 @@ export class DataGatheringProcessor { const data: Prisma.MarketDataUpdateInput[] = []; let lastMarketPrice: number; - while ( - isBefore( - currentDate, - new Date( - Date.UTC( - getYear(new Date()), - getMonth(new Date()), - getDate(new Date()), - 0 - ) - ) - ) - ) { + while (isBefore(currentDate, getStartOfUtcDate(new Date()))) { const marketPriceOfDataProvider = historicalData[assetProfileIdentifier]?.[ format(currentDate, DATE_FORMAT) 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 index 7830d428e..b0e8faa81 100644 --- 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 @@ -42,8 +42,8 @@ describe('DataGatheringService', () => { }); describe('getAssetProfileIdentifiersWithRecentMarketData', () => { - it('excludes carried forward market prices from the query', async () => { - jest.useFakeTimers().setSystemTime(parseDate('2026-08-23').getTime()); + it('queries real market prices since the start of yesterday (UTC)', async () => { + jest.useFakeTimers().setSystemTime(new Date('2026-08-24T14:00:00.000Z')); await dataGatheringService[ 'getAssetProfileIdentifiersWithRecentMarketData' @@ -51,80 +51,19 @@ describe('DataGatheringService', () => { expect(prismaService.marketData.groupBy).toHaveBeenCalledWith( expect.objectContaining({ - where: expect.objectContaining({ + where: { + date: { gte: new Date('2026-08-23T00:00:00.000Z') }, 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()); - + it('maps the query result to asset profile identifiers', async () => { 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' - } + { dataSource: 'COINGECKO', symbol: 'bitcoin' }, + { dataSource: 'YAHOO', symbol: 'AAPL' } ]); const assetProfileIdentifiers = @@ -133,7 +72,8 @@ describe('DataGatheringService', () => { ](); expect(assetProfileIdentifiers).toEqual([ - { dataSource: 'YAHOO', symbol: 'MSFT' } + { dataSource: 'COINGECKO', symbol: 'bitcoin' }, + { dataSource: 'YAHOO', symbol: 'AAPL' } ]); }); }); 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 010fe5f90..d41061734 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 @@ -31,14 +31,7 @@ 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, - isBefore, - min, - subDays, - subMilliseconds, - subYears -} from 'date-fns'; +import { format, min, subDays, subMilliseconds, subYears } from 'date-fns'; import { isEmpty } from 'lodash'; import ms, { StringValue } from 'ms'; @@ -269,15 +262,21 @@ export class DataGatheringService { age: GATHER_HISTORICAL_MARKET_DATA_COOLDOWN_IN_MS / 1000 }; + const assetProfileIdentifiersWithRecentMarketData = + await this.getAssetProfileIdentifiersWithRecentMarketData(); + await this.gatherSymbols({ removeOnComplete, - dataGatheringItems: await this.getCurrencies7D(), + dataGatheringItems: await this.getCurrencies7D({ + assetProfileIdentifiersWithRecentMarketData + }), priority: DATA_GATHERING_QUEUE_PRIORITY_HIGH }); await this.gatherSymbols({ removeOnComplete, dataGatheringItems: await this.getSymbols7D({ + assetProfileIdentifiersWithRecentMarketData, withUserSubscription: true }), priority: DATA_GATHERING_QUEUE_PRIORITY_MEDIUM @@ -286,6 +285,7 @@ export class DataGatheringService { await this.gatherSymbols({ removeOnComplete, dataGatheringItems: await this.getSymbols7D({ + assetProfileIdentifiersWithRecentMarketData, withUserSubscription: false }), priority: DATA_GATHERING_QUEUE_PRIORITY_LOW @@ -426,28 +426,24 @@ export class DataGatheringService { > { return ( await this.prismaService.marketData.groupBy({ - _max: { date: true }, by: ['dataSource', 'symbol'], orderBy: [{ symbol: 'asc' }], where: { - date: { gt: subDays(resetHours(new Date()), 7) }, + date: { gte: getStartOfUtcDate(subDays(new Date(), 1)) }, isCarriedForward: false, state: 'CLOSE' } }) - ) - .filter(({ _max }) => { - return !isBefore(_max.date, getStartOfUtcDate(subDays(new Date(), 1))); - }) - .map(({ dataSource, symbol }) => { - return { dataSource, symbol }; - }); + ).map(({ dataSource, symbol }) => { + return { dataSource, symbol }; + }); } - private async getCurrencies7D(): Promise { - const assetProfileIdentifiersWithRecentMarketData = - await this.getAssetProfileIdentifiersWithRecentMarketData(); - + private async getCurrencies7D({ + assetProfileIdentifiersWithRecentMarketData + }: { + assetProfileIdentifiersWithRecentMarketData: AssetProfileIdentifier[]; + }): Promise { return this.exchangeRateDataService .getCurrencyPairs() .filter(({ dataSource, symbol }) => { @@ -499,8 +495,10 @@ export class DataGatheringService { } private async getSymbols7D({ + assetProfileIdentifiersWithRecentMarketData, withUserSubscription = false }: { + assetProfileIdentifiersWithRecentMarketData: AssetProfileIdentifier[]; withUserSubscription?: boolean; }): Promise { const symbolProfiles = @@ -510,9 +508,6 @@ export class DataGatheringService { } ); - const assetProfileIdentifiersWithRecentMarketData = - await this.getAssetProfileIdentifiersWithRecentMarketData(); - return symbolProfiles .filter(({ dataSource, scraperConfiguration, symbol }) => { const manualDataSourceWithScraperConfiguration =