From 3ee570da81176c9463a72db0c384e6cebf85990a Mon Sep 17 00:00:00 2001 From: xBillCipher <0xdream@proton.me> Date: Mon, 14 Aug 2023 15:12:14 +0700 Subject: [PATCH 1/2] fix bug generage duplicated timeframe --- .../processor/timeFrame.build.processor.ts | 19 +++++++++++-------- .../src/lib/util/util.service.ts | 2 +- 2 files changed, 12 insertions(+), 9 deletions(-) diff --git a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts index 193b9ee..a41cb4d 100644 --- a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts +++ b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts @@ -90,18 +90,21 @@ export class TimeFrameBuildProcessor { const id = this.utilService.generateAggregatedId(current) removeItems[wallet].push(id) operations.push({ - create: { + update: { _id: id, }, }) - operations.push(current) + operations.push({ + doc: current, + doc_as_upsert: true, + }) }) }), ) }), ) - const createResponse = await this.esService.bulk({ + const updateResponse = await this.esService.bulk({ index: this.utilService.aggregatedDataIndex, operations: operations, refresh: true, @@ -142,17 +145,17 @@ export class TimeFrameBuildProcessor { conflicts: 'proceed', refresh: true, }) - const [skipped, inserted] = [ - createResponse.items.filter((c) => c?.create?.status === 409).length, - createResponse.items.filter((c) => c?.create?.status === 201).length, + const [updated, inserted] = [ + updateResponse.items.filter((c) => c?.update?.status === 200).length, + updateResponse.items.filter((c) => c?.update?.status === 201).length, ] this.logger.debug(` [update] info of lp performance tranche: ${tranche} - skipped: ${skipped} + skipped: ${updated} inserted: ${inserted} deleted: ${deleteResponse.deleted} - failed: ${createResponse.items.length - skipped - inserted} + failed: ${updateResponse.items.length - updated - inserted} `) } } diff --git a/packages/llp-aggregator-services/src/lib/util/util.service.ts b/packages/llp-aggregator-services/src/lib/util/util.service.ts index a7b758d..b477635 100644 --- a/packages/llp-aggregator-services/src/lib/util/util.service.ts +++ b/packages/llp-aggregator-services/src/lib/util/util.service.ts @@ -184,7 +184,7 @@ export class UtilService { } generateAggregatedId(item: AggreatedData, chainId = this.chainId) { - const based = `${item.wallet}_${chainId}_${item.tranche}_${item.from}_${item.to}_${item.valueMovement.fee}_${item.valueMovement.pnl}_${item.valueMovement.price}` + const based = `${item.wallet}_${chainId}_${item.tranche}_${item.from}_${item.to}` return utils.id(based) } From 13cdebd688ae30bdb48de2fd80217939e63d30c4 Mon Sep 17 00:00:00 2001 From: xBillCipher <0xdream@proton.me> Date: Tue, 5 Sep 2023 10:49:39 +0700 Subject: [PATCH 2/2] array empty and object null check --- .../src/guard/apikey.guard.ts | 2 +- .../src/crawler/prices.crawler.processor.ts | 2 +- .../processor/timeFrame.build.processor.ts | 90 +++++++++++-------- .../src/processor/timeFrame.cron.processor.ts | 2 +- 4 files changed, 54 insertions(+), 42 deletions(-) diff --git a/apps/llp-performance-aggregate-api/src/guard/apikey.guard.ts b/apps/llp-performance-aggregate-api/src/guard/apikey.guard.ts index 07cedd1..69cd056 100644 --- a/apps/llp-performance-aggregate-api/src/guard/apikey.guard.ts +++ b/apps/llp-performance-aggregate-api/src/guard/apikey.guard.ts @@ -7,7 +7,7 @@ export class ApiKeyGuard implements CanActivate { canActivate(context: ExecutionContext): boolean { const req = context.switchToHttp().getRequest() - const key = req.headers['x-api-key'] + const key = req.headers['x-api-key'] || req['query']?.['apiKey'] if (!key) { return false } diff --git a/apps/llp-performance-aggregate-worker/src/crawler/prices.crawler.processor.ts b/apps/llp-performance-aggregate-worker/src/crawler/prices.crawler.processor.ts index 3ba3460..3bc3ae4 100644 --- a/apps/llp-performance-aggregate-worker/src/crawler/prices.crawler.processor.ts +++ b/apps/llp-performance-aggregate-worker/src/crawler/prices.crawler.processor.ts @@ -101,7 +101,7 @@ export class PricesCrawlerProcessor { } } ` - const response = await this.graphqlClient.request(query, { + const response: any = await this.graphqlClient.request(query, { tranche: job.tranche, take: take, timestamp: lastSynced, diff --git a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts index a41cb4d..db867ba 100644 --- a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts +++ b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.build.processor.ts @@ -9,6 +9,10 @@ import { UtilService } from 'llp-aggregator-services/dist/util' import { Injectable, Logger } from '@nestjs/common' import { ElasticsearchService } from '@nestjs/elasticsearch' import { WorkerService } from '../worker.service' +import { + BulkResponse, + DeleteByQueryResponse, +} from '@elastic/elasticsearch/lib/api/types' @Injectable() export class TimeFrameBuildProcessor { @@ -104,58 +108,66 @@ export class TimeFrameBuildProcessor { }), ) - const updateResponse = await this.esService.bulk({ - index: this.utilService.aggregatedDataIndex, - operations: operations, - refresh: true, - }) - const deleteResponse = await this.esService.deleteByQuery({ - index: this.utilService.aggregatedDataIndex, - query: { - bool: { - must: [ - ...Object.entries(removeItems).map(([wallet, ids]) => ({ - bool: { - must: [ - { - bool: { - must_not: { - terms: { - _id: ids, + let updateResponse: BulkResponse, deleteResponse: DeleteByQueryResponse + if (operations.length) { + updateResponse = await this.esService.bulk({ + index: this.utilService.aggregatedDataIndex, + operations: operations, + refresh: true, + }) + } + const removeEntities = Object.entries(removeItems) + if (removeEntities.length) { + deleteResponse = await this.esService.deleteByQuery({ + index: this.utilService.aggregatedDataIndex, + query: { + bool: { + must: [ + ...Object.entries(removeItems).map(([wallet, ids]) => ({ + bool: { + must: [ + { + bool: { + must_not: { + terms: { + _id: ids, + }, }, }, }, - }, - { - term: { - wallet: wallet, + { + term: { + wallet: wallet, + }, }, - }, - ], - }, - })), - { - term: { - tranche: tranche, + ], + }, + })), + { + term: { + tranche: tranche, + }, }, - }, - ], + ], + }, }, - }, - conflicts: 'proceed', - refresh: true, - }) + conflicts: 'proceed', + refresh: true, + }) + } const [updated, inserted] = [ - updateResponse.items.filter((c) => c?.update?.status === 200).length, - updateResponse.items.filter((c) => c?.update?.status === 201).length, + updateResponse?.items.filter((c) => c?.update?.status === 200).length || + 0, + updateResponse?.items.filter((c) => c?.update?.status === 201).length || + 0, ] this.logger.debug(` [update] info of lp performance tranche: ${tranche} skipped: ${updated} inserted: ${inserted} - deleted: ${deleteResponse.deleted} - failed: ${updateResponse.items.length - updated - inserted} + deleted: ${deleteResponse?.deleted} + failed: ${(updateResponse?.items.length || 0) - updated - inserted} `) } } diff --git a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.cron.processor.ts b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.cron.processor.ts index fde331c..f151e4f 100644 --- a/apps/llp-performance-aggregate-worker/src/processor/timeFrame.cron.processor.ts +++ b/apps/llp-performance-aggregate-worker/src/processor/timeFrame.cron.processor.ts @@ -34,7 +34,7 @@ export class TimeFrameCronProcessor { wallets.push(wallet) }), ) - if (!wallets) { + if (!wallets || !wallets.length) { return } await this.redisService.client.sadd(