mirror of https://github.com/ghostfolio/ghostfolio
committed by
GitHub
13 changed files with 377 additions and 75 deletions
@ -0,0 +1,6 @@ |
|||
import { WatchlistResponse } from '@ghostfolio/common/interfaces'; |
|||
|
|||
export interface WatchlistValue { |
|||
expiration: number; |
|||
watchlist: WatchlistResponse['watchlist']; |
|||
} |
|||
@ -0,0 +1,3 @@ |
|||
export interface WatchlistQueueJob { |
|||
userId: string; |
|||
} |
|||
@ -0,0 +1,31 @@ |
|||
import { WATCHLIST_COMPUTATION_QUEUE } from '@ghostfolio/common/config'; |
|||
|
|||
import { InjectQueue } from '@nestjs/bull'; |
|||
import { Injectable } from '@nestjs/common'; |
|||
import { JobOptions, Queue } from 'bull'; |
|||
|
|||
import { WatchlistQueueJob } from './interfaces/watchlist-queue-job.interface'; |
|||
|
|||
@Injectable() |
|||
export class WatchlistComputationService { |
|||
public constructor( |
|||
@InjectQueue(WATCHLIST_COMPUTATION_QUEUE) |
|||
private readonly watchlistComputationQueue: Queue |
|||
) {} |
|||
|
|||
public async addJobToQueue({ |
|||
data, |
|||
name, |
|||
opts |
|||
}: { |
|||
data: WatchlistQueueJob; |
|||
name: string; |
|||
opts?: JobOptions; |
|||
}) { |
|||
return this.watchlistComputationQueue.add(name, data, opts); |
|||
} |
|||
|
|||
public async getJob(jobId: string) { |
|||
return this.watchlistComputationQueue.getJob(jobId); |
|||
} |
|||
} |
|||
@ -0,0 +1,52 @@ |
|||
import { RedisCacheModule } from '@ghostfolio/api/app/redis-cache/redis-cache.module'; |
|||
import { BenchmarkModule } from '@ghostfolio/api/services/benchmark/benchmark.module'; |
|||
import { ConfigurationModule } from '@ghostfolio/api/services/configuration/configuration.module'; |
|||
import { DataProviderModule } from '@ghostfolio/api/services/data-provider/data-provider.module'; |
|||
import { MarketDataModule } from '@ghostfolio/api/services/market-data/market-data.module'; |
|||
import { PrismaModule } from '@ghostfolio/api/services/prisma/prisma.module'; |
|||
import { WatchlistComputationService } from '@ghostfolio/api/services/queues/watchlist/watchlist-computation.service'; |
|||
import { SymbolProfileModule } from '@ghostfolio/api/services/symbol-profile/symbol-profile.module'; |
|||
import { |
|||
DEFAULT_PROCESSOR_WATCHLIST_COMPUTATION_TIMEOUT, |
|||
WATCHLIST_COMPUTATION_QUEUE |
|||
} from '@ghostfolio/common/config'; |
|||
|
|||
import { BullAdapter } from '@bull-board/api/bullAdapter'; |
|||
import { BullBoardModule } from '@bull-board/nestjs'; |
|||
import { BullModule } from '@nestjs/bull'; |
|||
import { Module } from '@nestjs/common'; |
|||
|
|||
import { WatchlistProcessor } from './watchlist.processor'; |
|||
|
|||
@Module({ |
|||
exports: [BullModule, WatchlistComputationService], |
|||
imports: [ |
|||
BenchmarkModule, |
|||
BullBoardModule.forFeature({ |
|||
adapter: BullAdapter, |
|||
name: WATCHLIST_COMPUTATION_QUEUE, |
|||
options: { |
|||
displayName: 'Watchlist Computation', |
|||
readOnlyMode: process.env.BULL_BOARD_IS_READ_ONLY !== 'false' |
|||
} |
|||
}), |
|||
BullModule.registerQueue({ |
|||
name: WATCHLIST_COMPUTATION_QUEUE, |
|||
settings: { |
|||
lockDuration: parseInt( |
|||
process.env.PROCESSOR_WATCHLIST_COMPUTATION_TIMEOUT ?? |
|||
DEFAULT_PROCESSOR_WATCHLIST_COMPUTATION_TIMEOUT.toString(), |
|||
10 |
|||
) |
|||
} |
|||
}), |
|||
ConfigurationModule, |
|||
DataProviderModule, |
|||
MarketDataModule, |
|||
PrismaModule, |
|||
RedisCacheModule, |
|||
SymbolProfileModule |
|||
], |
|||
providers: [WatchlistComputationService, WatchlistProcessor] |
|||
}) |
|||
export class WatchlistComputationQueueModule {} |
|||
@ -0,0 +1,152 @@ |
|||
import { WatchlistValue } from '@ghostfolio/api/app/endpoints/watchlist/interfaces/watchlist-value.interface'; |
|||
import { RedisCacheService } from '@ghostfolio/api/app/redis-cache/redis-cache.service'; |
|||
import { BenchmarkService } from '@ghostfolio/api/services/benchmark/benchmark.service'; |
|||
import { ConfigurationService } from '@ghostfolio/api/services/configuration/configuration.service'; |
|||
import { DataProviderService } from '@ghostfolio/api/services/data-provider/data-provider.service'; |
|||
import { MarketDataService } from '@ghostfolio/api/services/market-data/market-data.service'; |
|||
import { PrismaService } from '@ghostfolio/api/services/prisma/prisma.service'; |
|||
import { SymbolProfileService } from '@ghostfolio/api/services/symbol-profile/symbol-profile.service'; |
|||
import { |
|||
CACHE_TTL_INFINITE, |
|||
DEFAULT_PROCESSOR_WATCHLIST_COMPUTATION_CONCURRENCY, |
|||
WATCHLIST_COMPUTATION_QUEUE, |
|||
WATCHLIST_PROCESS_JOB_NAME |
|||
} from '@ghostfolio/common/config'; |
|||
import { WatchlistResponse } from '@ghostfolio/common/interfaces'; |
|||
|
|||
import { Process, Processor } from '@nestjs/bull'; |
|||
import { Injectable, Logger } from '@nestjs/common'; |
|||
import { Job } from 'bull'; |
|||
import { addMilliseconds } from 'date-fns'; |
|||
|
|||
import { WatchlistQueueJob } from './interfaces/watchlist-queue-job.interface'; |
|||
|
|||
@Injectable() |
|||
@Processor(WATCHLIST_COMPUTATION_QUEUE) |
|||
export class WatchlistProcessor { |
|||
private readonly logger = new Logger(WatchlistProcessor.name); |
|||
|
|||
public constructor( |
|||
private readonly benchmarkService: BenchmarkService, |
|||
private readonly configurationService: ConfigurationService, |
|||
private readonly dataProviderService: DataProviderService, |
|||
private readonly marketDataService: MarketDataService, |
|||
private readonly prismaService: PrismaService, |
|||
private readonly redisCacheService: RedisCacheService, |
|||
private readonly symbolProfileService: SymbolProfileService |
|||
) {} |
|||
|
|||
@Process({ |
|||
concurrency: parseInt( |
|||
process.env.PROCESSOR_WATCHLIST_COMPUTATION_CONCURRENCY ?? |
|||
DEFAULT_PROCESSOR_WATCHLIST_COMPUTATION_CONCURRENCY.toString(), |
|||
10 |
|||
), |
|||
name: WATCHLIST_PROCESS_JOB_NAME |
|||
}) |
|||
public async calculateWatchlist(job: Job<WatchlistQueueJob>) { |
|||
try { |
|||
const startTime = performance.now(); |
|||
const { userId } = job.data; |
|||
|
|||
this.logger.log( |
|||
`Watchlist calculation of user '${userId}' has been started` |
|||
); |
|||
|
|||
const user = await this.prismaService.user.findUnique({ |
|||
select: { |
|||
watchlist: { |
|||
select: { dataSource: true, symbol: true } |
|||
} |
|||
}, |
|||
where: { id: userId } |
|||
}); |
|||
|
|||
const [assetProfiles, quotes] = await Promise.all([ |
|||
this.symbolProfileService.getSymbolProfiles(user.watchlist), |
|||
this.dataProviderService.getQuotes({ |
|||
items: user.watchlist.map(({ dataSource, symbol }) => { |
|||
return { dataSource, symbol }; |
|||
}) |
|||
}) |
|||
]); |
|||
|
|||
let isComplete = user.watchlist.length > 0; |
|||
|
|||
const watchlist: WatchlistResponse['watchlist'] = await Promise.all( |
|||
user.watchlist.map(async ({ dataSource, symbol }) => { |
|||
const assetProfile = assetProfiles.find((profile) => { |
|||
return ( |
|||
profile.dataSource === dataSource && profile.symbol === symbol |
|||
); |
|||
}); |
|||
|
|||
const [allTimeHigh, trends] = await Promise.all([ |
|||
this.marketDataService.getMax({ |
|||
dataSource, |
|||
symbol |
|||
}), |
|||
this.benchmarkService.getBenchmarkTrends({ dataSource, symbol }) |
|||
]); |
|||
|
|||
if (!allTimeHigh?.marketPrice || !quotes[symbol]?.marketPrice) { |
|||
isComplete = false; |
|||
} |
|||
|
|||
const performancePercent = |
|||
this.benchmarkService.calculateChangeInPercentage( |
|||
allTimeHigh?.marketPrice, |
|||
quotes[symbol]?.marketPrice |
|||
); |
|||
|
|||
return { |
|||
dataSource, |
|||
symbol, |
|||
marketCondition: |
|||
this.benchmarkService.getMarketCondition(performancePercent), |
|||
name: assetProfile?.name, |
|||
performances: { |
|||
allTimeHigh: { |
|||
performancePercent, |
|||
date: allTimeHigh?.date |
|||
} |
|||
}, |
|||
trend50d: trends.trend50d, |
|||
trend200d: trends.trend200d |
|||
}; |
|||
}) |
|||
); |
|||
|
|||
const sortedWatchlist = watchlist.sort((a, b) => { |
|||
return a.name.localeCompare(b.name); |
|||
}); |
|||
|
|||
this.logger.log( |
|||
`Watchlist calculation of user '${userId}' has been completed in ${( |
|||
(performance.now() - startTime) / |
|||
1000 |
|||
).toFixed(3)} seconds` |
|||
); |
|||
|
|||
const expiration = addMilliseconds( |
|||
new Date(), |
|||
isComplete ? this.configurationService.get('CACHE_QUOTES_TTL') : 0 |
|||
); |
|||
|
|||
await this.redisCacheService.set( |
|||
this.redisCacheService.getWatchlistKey({ userId }), |
|||
JSON.stringify({ |
|||
expiration: expiration.getTime(), |
|||
watchlist: sortedWatchlist |
|||
} as WatchlistValue), |
|||
CACHE_TTL_INFINITE |
|||
); |
|||
|
|||
return sortedWatchlist; |
|||
} catch (error) { |
|||
this.logger.error(error); |
|||
|
|||
throw new Error(error); |
|||
} |
|||
} |
|||
} |
|||
Loading…
Reference in new issue