From 514a912979ab70d38d857aa825087a9a956d4345 Mon Sep 17 00:00:00 2001 From: TaprootFreak <142087526+TaprootFreak@users.noreply.github.com> Date: Thu, 23 Jul 2026 15:33:20 +0200 Subject: [PATCH] fix(liquidity): keep activation debounce across drain chunks (#4337) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(liquidity): keep activation debounce across drain chunks verifyRule cleared the rule's activation-debounce timer (ruleActivations) unconditionally right before executeRule. After each pipeline completed and the rule returned to Active, the next cron cycle re-armed the timer and had to wait the full lmActivationDelay again before the next chunk. For a redundancy rule whose orders are liquidity-capped (best-price chunk, e.g. an exchange sell), this serialized the drain to one chunk per delay, turning a large offload into a multi-hour trickle. Arm the timer once when the condition first appears and clear it only when the condition resolves (else branch); follow-up chunks then fire at the normal every-minute cadence. Also clear the timer when a rule is Paused so a rule paused by a failed pipeline re-debounces on reactivation — deliberately not for other non-active states, notably Processing (occurs between chunks). * test(liquidity): cover the activation-debounce invariant Regression tests for verifyRule's ruleActivations behaviour: the timer persists across a drain chunk (Active + ongoing condition), is cleared on Paused so a failed-then-reactivated rule re-debounces, and is kept while Processing between chunks. Verified to fail if either fix is reverted. * fix(liquidity): clear activation timer at pause transition, not in verifyRule The earlier verifyRule guard cleared the activation-debounce timer only when the every-minute cron observed the rule Paused. A manual reactivation (reactivateRule) can flip Paused -> Active in a different service before that tick, letting a just-failed rule skip the debounce with a stale timestamp. Move the reset to the actual pause transition: handlePipelineFail now calls a new resetActivation(ruleId) on the liquidity service right after rule.pause(). Revert the verifyRule non-active guard to its original form. Covers manual and cron reactivation alike. Tests updated; adds a handlePipelineFail spec. * style(liquidity): apply prettier formatting to pipeline spec --- ...uidity-management-pipeline.service.spec.ts | 65 +++++++++++ .../liquidity-management-pipeline.service.ts | 5 + .../liquidity-management.service.spec.ts | 101 ++++++++++++++++++ .../services/liquidity-management.service.ts | 8 +- 4 files changed, 177 insertions(+), 2 deletions(-) create mode 100644 src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.spec.ts create mode 100644 src/subdomains/core/liquidity-management/services/liquidity-management.service.spec.ts diff --git a/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.spec.ts b/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.spec.ts new file mode 100644 index 0000000000..9dc1f52c28 --- /dev/null +++ b/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.spec.ts @@ -0,0 +1,65 @@ +import { createMock } from '@golevelup/ts-jest'; +import { NotificationService } from 'src/subdomains/supporting/notification/services/notification.service'; +import { LiquidityManagementOrder } from '../entities/liquidity-management-order.entity'; +import { LiquidityManagementPipeline } from '../entities/liquidity-management-pipeline.entity'; +import { LiquidityManagementRule } from '../entities/liquidity-management-rule.entity'; +import { LiquidityManagementPipelineStatus, LiquidityManagementRuleStatus, LiquidityOptimizationType } from '../enums'; +import { LiquidityActionIntegrationFactory } from '../factories/liquidity-action-integration.factory'; +import { LiquidityManagementOrderRepository } from '../repositories/liquidity-management-order.repository'; +import { LiquidityManagementPipelineRepository } from '../repositories/liquidity-management-pipeline.repository'; +import { LiquidityManagementRuleRepository } from '../repositories/liquidity-management-rule.repository'; +import { LiquidityManagementPipelineService } from './liquidity-management-pipeline.service'; +import { LiquidityManagementService } from './liquidity-management.service'; + +describe('LiquidityManagementPipelineService', () => { + let service: LiquidityManagementPipelineService; + let ruleRepo: LiquidityManagementRuleRepository; + let orderRepo: LiquidityManagementOrderRepository; + let pipelineRepo: LiquidityManagementPipelineRepository; + let actionIntegrationFactory: LiquidityActionIntegrationFactory; + let notificationService: NotificationService; + let liquidityManagementService: LiquidityManagementService; + + beforeEach(() => { + ruleRepo = createMock(); + orderRepo = createMock(); + pipelineRepo = createMock(); + actionIntegrationFactory = createMock(); + notificationService = createMock(); + liquidityManagementService = createMock(); + + service = new LiquidityManagementPipelineService( + ruleRepo, + orderRepo, + pipelineRepo, + actionIntegrationFactory, + notificationService, + liquidityManagementService, + ); + }); + + describe('handlePipelineFail', () => { + it('resets the activation debounce timer when a rule is paused', async () => { + const rule = Object.assign(new LiquidityManagementRule(), { + id: 42, + status: LiquidityManagementRuleStatus.PROCESSING, + sendNotifications: false, + targetFiat: { name: 'EUR' }, + }); + const pipeline = Object.assign(new LiquidityManagementPipeline(), { + id: 1, + type: LiquidityOptimizationType.DEFICIT, + maxAmount: 100, + status: LiquidityManagementPipelineStatus.FAILED, + rule, + }); + const order = Object.assign(new LiquidityManagementOrder(), { errorMessage: 'order failed' }); + + await service['handlePipelineFail'](pipeline, order); + + expect(liquidityManagementService.resetActivation).toHaveBeenCalledWith(42); + expect(rule.status).toBe(LiquidityManagementRuleStatus.PAUSED); + expect(ruleRepo.save).toHaveBeenCalledWith(rule); + }); + }); +}); diff --git a/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.ts b/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.ts index bb347e1049..0f108a73d4 100644 --- a/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.ts +++ b/src/subdomains/core/liquidity-management/services/liquidity-management-pipeline.service.ts @@ -17,6 +17,7 @@ import { LiquidityActionIntegrationFactory } from '../factories/liquidity-action import { LiquidityManagementOrderRepository } from '../repositories/liquidity-management-order.repository'; import { LiquidityManagementPipelineRepository } from '../repositories/liquidity-management-pipeline.repository'; import { LiquidityManagementRuleRepository } from '../repositories/liquidity-management-rule.repository'; +import { LiquidityManagementService } from './liquidity-management.service'; @Injectable() export class LiquidityManagementPipelineService { @@ -28,6 +29,7 @@ export class LiquidityManagementPipelineService { private readonly pipelineRepo: LiquidityManagementPipelineRepository, private readonly actionIntegrationFactory: LiquidityActionIntegrationFactory, private readonly notificationService: NotificationService, + private readonly liquidityManagementService: LiquidityManagementService, ) {} //*** JOBS ***// @@ -264,6 +266,9 @@ export class LiquidityManagementPipelineService { ): Promise { const rule = pipeline.rule.pause(); + // The rule is now paused; clear its activation-debounce timer so a later reactivation re-debounces. + this.liquidityManagementService.resetActivation(rule.id); + await this.ruleRepo.save(rule); const [errorMessage, mailRequest] = this.generateFailMessage(pipeline, order); diff --git a/src/subdomains/core/liquidity-management/services/liquidity-management.service.spec.ts b/src/subdomains/core/liquidity-management/services/liquidity-management.service.spec.ts new file mode 100644 index 0000000000..8a3be2d04d --- /dev/null +++ b/src/subdomains/core/liquidity-management/services/liquidity-management.service.spec.ts @@ -0,0 +1,101 @@ +import { createMock } from '@golevelup/ts-jest'; +import { ConfigService } from 'src/config/config'; +import { SettingService } from 'src/shared/models/setting/setting.service'; +import { Util } from 'src/shared/utils/util'; +import { Price } from 'src/subdomains/supporting/pricing/domain/entities/price'; +import { PricingService } from 'src/subdomains/supporting/pricing/services/pricing.service'; +import { createDefaultLiquidityBalance } from '../__mocks__/liquidity-balance.entity.mock'; +import { LiquidityManagementRule } from '../entities/liquidity-management-rule.entity'; +import { LiquidityManagementRuleStatus, LiquidityOptimizationType } from '../enums'; +import { LiquidityManagementPipelineRepository } from '../repositories/liquidity-management-pipeline.repository'; +import { LiquidityManagementRuleRepository } from '../repositories/liquidity-management-rule.repository'; +import { LiquidityManagementBalanceService } from './liquidity-management-balance.service'; +import { LiquidityManagementService } from './liquidity-management.service'; + +describe('LiquidityManagementService', () => { + let service: LiquidityManagementService; + let ruleRepo: LiquidityManagementRuleRepository; + let pipelineRepo: LiquidityManagementPipelineRepository; + let balanceService: LiquidityManagementBalanceService; + let settingService: SettingService; + let pricingService: PricingService; + let executeRuleSpy: jest.SpyInstance; + + beforeAll(() => { + new ConfigService(); // sets module-level Config (verifyRule reads Config.liquidityManagement) + }); + + beforeEach(() => { + ruleRepo = createMock(); + pipelineRepo = createMock(); + balanceService = createMock(); + settingService = createMock(); + pricingService = createMock(); + + service = new LiquidityManagementService(ruleRepo, pipelineRepo, balanceService, settingService, pricingService); + + executeRuleSpy = jest.spyOn(service as any, 'executeRule').mockResolvedValue(undefined as any); + }); + + function createRule(partial: Partial): LiquidityManagementRule { + return Object.assign(new LiquidityManagementRule(), { + delayActivation: true, + optimal: 0, + ...partial, + }); + } + + describe('verifyRule ruleActivations debounce invariant', () => { + it('keeps the activation timer across a drain chunk', async () => { + const rule = createRule({ + id: 1, + status: LiquidityManagementRuleStatus.ACTIVE, + delayActivation: true, + }); + const balance = createDefaultLiquidityBalance(); + + jest.spyOn(balanceService, 'findRelevantBalance').mockReturnValue(balance); + jest.spyOn(pricingService, 'getPrice').mockResolvedValue(Price.create('EUR', 'ASSET', 1)); + jest.spyOn(rule, 'verify').mockReturnValue({ + action: LiquidityOptimizationType.REDUNDANCY, + minAmount: 0, + maxAmount: 100, + }); + jest.spyOn(balanceService, 'hasPendingOrders').mockResolvedValue(false); + jest.spyOn(settingService, 'get').mockResolvedValue('15'); + + service['ruleActivations'].set(rule.id, Util.minutesBefore(60)); + + await service['verifyRule'](rule, [balance]); + + expect(executeRuleSpy).toHaveBeenCalledTimes(1); + expect(service['ruleActivations'].has(rule.id)).toBe(true); + }); + + it('keeps the activation timer while a rule is processing between chunks', async () => { + const rule = createRule({ + id: 3, + status: LiquidityManagementRuleStatus.PROCESSING, + }); + + service['ruleActivations'].set(rule.id, new Date()); + + await service['verifyRule'](rule, []); + + expect(service['ruleActivations'].has(rule.id)).toBe(true); + expect(executeRuleSpy).not.toHaveBeenCalled(); + }); + }); + + describe('resetActivation', () => { + it('clears the activation timer for the given rule id', () => { + const ruleId = 7; + + service['ruleActivations'].set(ruleId, new Date()); + + service.resetActivation(ruleId); + + expect(service['ruleActivations'].has(ruleId)).toBe(false); + }); + }); +}); diff --git a/src/subdomains/core/liquidity-management/services/liquidity-management.service.ts b/src/subdomains/core/liquidity-management/services/liquidity-management.service.ts index bb7675cb10..3559b5631f 100644 --- a/src/subdomains/core/liquidity-management/services/liquidity-management.service.ts +++ b/src/subdomains/core/liquidity-management/services/liquidity-management.service.ts @@ -97,6 +97,12 @@ export class LiquidityManagementService { return this.executeRule(rule, liquidityState, LiquidityOptimizationType.REDUNDANCY); } + // Clears the activation-debounce timer for a rule. Called when a rule leaves the drain lifecycle + // (e.g. paused after a failed pipeline) so that a subsequent reactivation re-debounces from scratch. + resetActivation(ruleId: number): void { + this.ruleActivations.delete(ruleId); + } + //*** HELPER METHODS ***// private async findRuleByAssetOrThrow(assetId: number): Promise { @@ -144,8 +150,6 @@ export class LiquidityManagementService { const requiredActivationTime = Util.minutesBefore(+delay); if (!rule.delayActivation || this.ruleActivations.get(rule.id) < requiredActivationTime) { - this.ruleActivations.delete(rule.id); - await this.executeRule(rule, result); } } else {