From ef29c0b340ccbd373c0c1a04006c37659b7fc614 Mon Sep 17 00:00:00 2001 From: kostyastruga Date: Thu, 6 Aug 2026 16:01:02 +0300 Subject: [PATCH 1/4] Fix local count --- README.md | 4 +- .../config/payment/PaymentResourceConfig.java | 7 +- .../LocalResultStorageRepository.java | 37 +++-- .../LocalCountAggregatorDecorator.java | 15 +- .../LocalSumAggregatorDecorator.java | 18 ++- .../LocalUniqueValueAggregatorDecorator.java | 11 +- .../DatabasePaymentFieldResolver.java | 12 ++ .../handler/FraudInspectorHandler.java | 12 ++ .../LocalResultStorageRepositoryTest.java | 118 ++++++++++++++ tests/k6/README.md | 147 ++++++++++++++++++ tests/k6/config.js | 24 +++ tests/k6/inspection-load.js | 37 +++++ tests/k6/lib/client.js | 88 +++++++++++ tests/k6/lib/payloads.js | 89 +++++++++++ tests/k6/lib/provision.js | 39 +++++ tests/k6/mixed-load.js | 55 +++++++ tests/k6/run.sh | 26 ++++ tests/k6/smoke.js | 51 ++++++ 18 files changed, 758 insertions(+), 32 deletions(-) create mode 100644 src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java create mode 100644 tests/k6/README.md create mode 100644 tests/k6/config.js create mode 100644 tests/k6/inspection-load.js create mode 100644 tests/k6/lib/client.js create mode 100644 tests/k6/lib/payloads.js create mode 100644 tests/k6/lib/provision.js create mode 100644 tests/k6/mixed-load.js create mode 100755 tests/k6/run.sh create mode 100644 tests/k6/smoke.js diff --git a/README.md b/README.md index f139e48d..8f58148a 100644 --- a/README.md +++ b/README.md @@ -21,6 +21,8 @@ protocol for managing a set of patterns and bindings Instruction for testing in https://github.com/valitydev/fraudbusters-compose +System integration and load tests using `fraudbusters-api` and k6: +[tests/k6/README.md](tests/k6/README.md). + ### License [Apache 2.0 License.](/LICENSE) - diff --git a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java index 256f2ac4..032b9595 100644 --- a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java +++ b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java @@ -7,6 +7,7 @@ import dev.vality.fraudbusters.domain.FraudResult; import dev.vality.fraudbusters.resource.payment.handler.FraudInspectorHandler; import dev.vality.fraudbusters.stream.impl.TemplateVisitorImpl; +import io.micrometer.core.instrument.MeterRegistry; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -24,14 +25,16 @@ public InspectorProxySrv.Iface fraudInspectorHandler( CheckedResultToRiskScoreConverter checkedResultToRiskScoreConverter, ContextToFraudRequestConverter requestConverter, TemplateVisitorImpl templateVisitor, - WbListServiceSrv.Iface wbListServiceSrv) { + WbListServiceSrv.Iface wbListServiceSrv, + MeterRegistry meterRegistry) { return new FraudInspectorHandler( resultTopic, checkedResultToRiskScoreConverter, requestConverter, templateVisitor, kafkaFraudResultTemplate, - wbListServiceSrv + wbListServiceSrv, + meterRegistry ); } diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepository.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepository.java index f0a96d8d..8e5de83d 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepository.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepository.java @@ -42,7 +42,8 @@ public Integer countOperationByFieldWithGroupBy( List fieldModels) { List checkedPayments = localStorage.get(); int count = (int) checkedPayments.stream() - .filter(checkedPayment -> filterPaymentByValue(from, to, fieldModels, checkedPayment)) + .filter(checkedPayment -> filterPaymentByValue( + fieldName, value, from, to, fieldModels, checkedPayment)) .count(); log.debug("LocalResultStorageRepository countOperationByFieldWithGroupBy: {}", count); return count; @@ -57,7 +58,8 @@ public Long sumOperationByFieldWithGroupBy( List fieldModels) { List checkedPayments = localStorage.get(); long sum = checkedPayments.stream() - .filter(checkedPayment -> filterPaymentByValue(from, to, fieldModels, checkedPayment)) + .filter(checkedPayment -> filterPaymentByValue( + fieldName, value, from, to, fieldModels, checkedPayment)) .mapToLong(CheckedPayment::getAmount) .sum(); log.debug("LocalResultStorageRepository sumOperationByFieldWithGroupBy: {}", sum); @@ -88,7 +90,8 @@ public Integer uniqCountOperationWithGroupBy( List fieldModels) { List checkedPayments = localStorage.get(); long count = checkedPayments.stream() - .filter(checkedPayment -> filterPaymentByValue(from, to, fieldModels, checkedPayment)) + .filter(checkedPayment -> filterPaymentByValue( + fieldNameBy, value, from, to, fieldModels, checkedPayment)) .map(checkedPayment -> paymentFieldValueResolver.resolve(fieldNameCount, checkedPayment)) .distinct() .count(); @@ -97,12 +100,15 @@ public Integer uniqCountOperationWithGroupBy( } private boolean filterPaymentByValue( + String fieldName, + Object value, Long from, Long to, List fieldModels, CheckedPayment checkedPayment) { return checkedPayment.getEventTime() >= from && checkedPayment.getEventTime() <= to + && paymentFieldValueFilter.filter(fieldName, value, checkedPayment) && fieldModels.stream() .allMatch(fieldModel -> paymentFieldValueFilter.filter( fieldModel.getName(), @@ -121,6 +127,8 @@ public Integer countOperationSuccessWithGroupBy( List checkedPayments = localStorage.get(); return (int) checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -136,6 +144,8 @@ public Integer countOperationPendingWithGroupBy(String fieldName, Object value, List checkedPayments = localStorage.get(); return (int) checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -156,6 +166,8 @@ public Integer countOperationErrorWithGroupBy( List checkedPayments = localStorage.get(); return (int) checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -177,6 +189,8 @@ public Integer countOperationErrorWithGroupBy( List checkedPayments = localStorage.get(); return (int) checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -195,6 +209,8 @@ public Long sumOperationSuccessWithGroupBy( List checkedPayments = localStorage.get(); return checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -216,6 +232,8 @@ public Long sumOperationErrorWithGroupBy( List checkedPayments = localStorage.get(); return checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -236,6 +254,8 @@ public Long sumOperationErrorWithGroupBy(String fieldName, List checkedPayments = localStorage.get(); return checkedPayments.stream() .filter(checkedPayment -> filterByStatusAndFields( + fieldName, + value, from, to, fieldModels, @@ -247,19 +267,14 @@ public Long sumOperationErrorWithGroupBy(String fieldName, } private boolean filterByStatusAndFields( + String fieldName, + Object value, Long from, Long to, List fieldModels, CheckedPayment checkedPayment, PaymentStatus paymentStatus) { - return checkedPayment.getEventTime() >= from - && checkedPayment.getEventTime() <= to - && fieldModels.stream() - .allMatch(fieldModel -> paymentFieldValueFilter.filter( - fieldModel.getName(), - fieldModel.getValue(), - checkedPayment - )) + return filterPaymentByValue(fieldName, value, from, to, fieldModels, checkedPayment) && paymentStatus.name().equals(checkedPayment.getPaymentStatus()); } diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java index 1f15279e..43399d8d 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java @@ -75,7 +75,8 @@ public Integer countError( Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Integer localCount = localStorageRepository.countOperationErrorWithGroupBy( checkedField.name(), resolve.getValue(), @@ -106,7 +107,8 @@ public Integer countError(PaymentCheckedField paymentCheckedField, PaymentModel Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(paymentCheckedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Integer localCount = localStorageRepository.countOperationErrorWithGroupBy( paymentCheckedField.name(), resolve.getValue(), @@ -172,12 +174,13 @@ private Integer getCount( Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Integer count = aggregateFunction.accept( - resolve.getName(), + checkedField.name(), resolve.getValue(), - timeBound.getLeft().toEpochMilli(), - timeBound.getRight().toEpochMilli(), + timeBound.getLeft().getEpochSecond(), + timeBound.getRight().getEpochSecond(), eventFields ); log.debug( diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java index 3f13e342..778f9249 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java @@ -39,7 +39,8 @@ public Double sum( FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Long localSum = localStorageRepository.sumOperationByFieldWithGroupBy( checkedField.name(), resolve.getValue(), @@ -81,7 +82,8 @@ public Double sumError(PaymentCheckedField checkedField, Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Long localSum = localStorageRepository.sumOperationErrorWithGroupBy( checkedField.name(), resolve.getValue(), @@ -114,7 +116,8 @@ public Double sumError(PaymentCheckedField checkedField, Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Long localSum = localStorageRepository.sumOperationErrorWithGroupBy( checkedField.name(), resolve.getValue(), @@ -176,12 +179,13 @@ private Double getSum( Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List eventFields = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Long sum = aggregateFunction.accept( - resolve.getName(), + checkedField.name(), resolve.getValue(), - timeBound.getLeft().toEpochMilli(), - timeBound.getRight().toEpochMilli(), + timeBound.getLeft().getEpochSecond(), + timeBound.getRight().getEpochSecond(), eventFields ); double resultSum = withCurrent diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java index 70ff3df7..2543e6f5 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java @@ -39,13 +39,14 @@ public Integer countUniqueValue( Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); FieldModel resolve = databasePaymentFieldResolver.resolve(countField, paymentModel); - List fieldModels = databasePaymentFieldResolver.resolveListFields(paymentModel, list); + List fieldModels = + databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); Integer localUniqCountOperation = localStorageRepository.uniqCountOperationWithGroupBy( - resolve.getName(), + countField.name(), resolve.getValue(), - databasePaymentFieldResolver.resolve(onField), - timeBound.getLeft().toEpochMilli(), - timeBound.getRight().toEpochMilli(), + onField.name(), + timeBound.getLeft().getEpochSecond(), + timeBound.getRight().getEpochSecond(), fieldModels ); int result = localUniqCountOperation + uniq; diff --git a/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java b/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java index c115d330..f98f15bb 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java @@ -25,6 +25,18 @@ public List resolveListFields(PaymentModel model, List(); } + @NotNull + public List resolveListFieldsForLocalStorage( + PaymentModel model, + List list) { + if (list != null) { + return list.stream() + .map(field -> new FieldModel(field.name(), resolve(field, model).getValue())) + .collect(Collectors.toList()); + } + return new ArrayList<>(); + } + public FieldModel resolve(PaymentCheckedField field, PaymentModel model) { if (field == null) { throw new UnknownFieldException(); diff --git a/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java b/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java index e66412b6..e274f8e6 100644 --- a/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java +++ b/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java @@ -13,6 +13,8 @@ import dev.vality.fraudbusters.domain.FraudResult; import dev.vality.fraudbusters.fraud.model.PaymentModel; import dev.vality.fraudbusters.stream.TemplateVisitor; +import io.micrometer.core.instrument.MeterRegistry; +import io.micrometer.core.instrument.Timer; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.thrift.TException; @@ -28,9 +30,12 @@ public class FraudInspectorHandler implements InspectorProxySrv.Iface { private final TemplateVisitor templateVisitor; private final KafkaTemplate kafkaFraudResultTemplate; private final WbListServiceSrv.Iface wbListServiceSrv; + private final MeterRegistry meterRegistry; @Override public RiskScore inspectPayment(Context context) throws TException { + Timer.Sample sample = Timer.start(meterRegistry); + String outcome = "success"; try { FraudRequest model = requestConverter.convert(context); if (model != null) { @@ -42,8 +47,15 @@ public RiskScore inspectPayment(Context context) throws TException { } return RiskScore.high; } catch (Exception e) { + outcome = "error"; log.error("Error when inspectPayment() e: ", e); throw new TException("Error when inspectPayment() e: ", e); + } finally { + sample.stop(Timer.builder("fraudbusters.online.inspect") + .description("Online Fraudbusters inspectPayment latency") + .tag("method", "inspectPayment") + .tag("outcome", outcome) + .register(meterRegistry)); } } diff --git a/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java b/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java new file mode 100644 index 00000000..2b3f85ff --- /dev/null +++ b/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java @@ -0,0 +1,118 @@ +package dev.vality.fraudbusters.fraud.localstorage; + +import dev.vality.damsel.fraudbusters.PaymentStatus; +import dev.vality.fraudbusters.domain.CheckedPayment; +import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.filter.PaymentFieldValueFilter; +import dev.vality.fraudbusters.fraud.filter.PaymentFieldValueResolver; +import dev.vality.fraudbusters.fraud.model.FieldModel; +import dev.vality.fraudbusters.fraud.model.PaymentModel; +import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +class LocalResultStorageRepositoryTest { + + private static final long FROM = 50L; + private static final long TO = 150L; + private static final String TOKEN = "token-a"; + + private LocalResultStorageRepository repository; + + @BeforeEach + void setUp() { + LocalResultStorage storage = new LocalResultStorage(); + PaymentFieldValueResolver resolver = new PaymentFieldValueResolver(); + repository = new LocalResultStorageRepository( + storage, + new PaymentFieldValueFilter(resolver), + resolver + ); + + storage.get().addAll(List.of( + payment(TOKEN, "same@example.com", PaymentStatus.captured, null, 10L), + payment(TOKEN, "same@example.com", PaymentStatus.pending, null, 20L), + payment(TOKEN, "other@example.com", PaymentStatus.failed, "E1", 30L), + payment(TOKEN, "other@example.com", PaymentStatus.failed, "E2", 40L), + payment("token-b", "one@example.com", PaymentStatus.captured, null, 100L), + payment("token-c", "two@example.com", PaymentStatus.failed, "E1", 200L) + )); + } + + @Test + void shouldApplyMainFieldFilterToGeneralAggregations() { + assertEquals(4, repository.countOperationByFieldWithGroupBy( + PaymentCheckedField.CARD_TOKEN.name(), TOKEN, FROM, TO, List.of())); + assertEquals(100L, repository.sumOperationByFieldWithGroupBy( + PaymentCheckedField.CARD_TOKEN.name(), TOKEN, FROM, TO, List.of())); + assertEquals(2, repository.uniqCountOperationWithGroupBy( + PaymentCheckedField.CARD_TOKEN.name(), + TOKEN, + PaymentCheckedField.EMAIL.name(), + FROM, + TO, + List.of() + )); + } + + @Test + void shouldApplyMainFieldFilterToStatusCounts() { + String fieldName = PaymentCheckedField.CARD_TOKEN.name(); + + assertEquals(1, repository.countOperationSuccessWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of())); + assertEquals(1, repository.countOperationPendingWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of())); + assertEquals(2, repository.countOperationErrorWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of())); + assertEquals(1, repository.countOperationErrorWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of(), "E1")); + } + + @Test + void shouldApplyMainFieldFilterToStatusSums() { + String fieldName = PaymentCheckedField.CARD_TOKEN.name(); + + assertEquals(10L, repository.sumOperationSuccessWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of())); + assertEquals(70L, repository.sumOperationErrorWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of())); + assertEquals(30L, repository.sumOperationErrorWithGroupBy( + fieldName, TOKEN, FROM, TO, List.of(), "E1")); + } + + @Test + void shouldResolveDifferentFieldNamesForClickHouseAndLocalStorage() { + DatabasePaymentFieldResolver resolver = new DatabasePaymentFieldResolver(); + PaymentModel payment = new PaymentModel(); + payment.setCardToken(TOKEN); + + List clickHouseFields = + resolver.resolveListFields(payment, List.of(PaymentCheckedField.CARD_TOKEN)); + List localFields = + resolver.resolveListFieldsForLocalStorage(payment, List.of(PaymentCheckedField.CARD_TOKEN)); + + assertEquals("cardToken", clickHouseFields.get(0).getName()); + assertEquals("CARD_TOKEN", localFields.get(0).getName()); + } + + private CheckedPayment payment( + String cardToken, + String email, + PaymentStatus status, + String errorCode, + Long amount) { + CheckedPayment payment = new CheckedPayment(); + payment.setEventTime(100L); + payment.setCardToken(cardToken); + payment.setEmail(email); + payment.setPaymentStatus(status.name()); + payment.setErrorCode(errorCode); + payment.setAmount(amount); + return payment; + } +} diff --git a/tests/k6/README.md b/tests/k6/README.md new file mode 100644 index 00000000..509d9193 --- /dev/null +++ b/tests/k6/README.md @@ -0,0 +1,147 @@ +# Fraudbusters k6 tests + +Набор использует реальные публичные HTTP-фасады: + +- `fraudbusters-api`: `/inspect-payment` и `/payments`; +- `fraudbusters-management`: создание правил/references и чтение control plane. + +Тем самым запрос k6 проходит через Swagger API, Thrift, Fraudbusters runtime, +Kafka configuration flow, внешние rule dependencies и хранилища, поднятые в +тестовом окружении. + +Архитектурно это позволяет измерять два уровня отдельно: + +- latency внешнего контракта `fraudbusters-api → Thrift → fraudbusters`; +- control/data plane всей системы вместе с management, Kafka и хранилищами. + +## Сценарии + +| Файл | Назначение | +|---|---| +| `smoke.js` | Сквозной интеграционный тест: создать правило и reference, дождаться активации, проверить `fatal/high`, загрузить payment event | +| `inspection-load.js` | Нагрузка только на критический онлайн-путь `/inspect-payment` | +| `mixed-load.js` | Смешанная нагрузка: inspection, event ingestion и management reads | + +## Предварительные условия + +Поднять полный контур, например из соседнего `fraudbusters-compose`, включая: + +- `fraudbusters-api` на host-порту `9999`; +- `fraudbusters-management` на host-порту `8085`; +- Fraudbusters, Kafka, ClickHouse, PostgreSQL, `wb-list` и остальные зависимости. + +В текущем `fraudbusters-compose/docker-compose.yml` определения основных +Fraudbusters-сервисов закомментированы. Перед запуском тестов нужен рабочий +compose/profile или развернутое тестовое окружение. + +Если включена авторизация, передать JWT: + +```bash +export FB_API_TOKEN=... +export FB_MANAGEMENT_TOKEN=... +``` + +## Запуск + +k6 не установлен локально, поэтому `run.sh` запускает официальный Docker image: + +```bash +./tests/k6/run.sh smoke.js +./tests/k6/run.sh inspection-load.js +./tests/k6/run.sh mixed-load.js +``` + +При локально установленном k6: + +```bash +k6 run tests/k6/smoke.js +``` + +## Настройка нагрузки + +```bash +FB_RATE=100 \ +FB_DURATION=10m \ +FB_PREALLOCATED_VUS=100 \ +FB_MAX_VUS=500 \ +./tests/k6/run.sh inspection-load.js +``` + +Для mixed workload: + +```bash +FB_INSPECT_RATE=100 \ +FB_INGEST_RATE=20 \ +FB_MANAGEMENT_VUS=5 \ +FB_DURATION=15m \ +./tests/k6/run.sh mixed-load.js +``` + +Основные переменные: + +| Переменная | Default | Назначение | +|---|---:|---| +| `FB_API_URL` | `http://host.docker.internal:9999` | fraudbusters-api | +| `FB_MANAGEMENT_URL` | `http://host.docker.internal:8085/fb-management/v1` | fraudbusters-management | +| `FB_ACTIVATION_WAIT_SECONDS` | `5` | Ожидание eventual consistency после создания reference | +| `FB_PARTY_ID`, `FB_SHOP_ID` | `k6-load-*` | Предварительно настроенный merchant для load-сценариев | +| `FB_CLEANUP` | `false` | Удалять созданные smoke-test template/reference | + +## Подготовка данных для load-тестов + +`inspection-load.js` и `mixed-load.js` намеренно не изменяют конфигурацию правил +во время нагрузки. До теста нужно один раз создать template/reference для +`FB_PARTY_ID` и `FB_SHOP_ID`, либо использовать merchant с уже активными +правилами. + +Для сравнимых прогонов фиксировать: + +- набор и сложность правил; +- объём истории ClickHouse; +- ответы/latency `wb-list`, Columbus и Trusted Tokens; +- число partitions и consumer concurrency Kafka; +- CPU/memory limits всех сервисов. + +Не смешивать тест максимальной пропускной способности inspection с массовым +созданием правил: это разные профили нагрузки и разные bottlenecks. + +## План покрытия всей системы + +Текущий каркас покрывает базовый сквозной путь и готов для расширения. Полный +system-test набор рекомендуется разделить на независимые профили: + +| Профиль | Подготовка | Проверяемая цепочка | Главная метрика | +|---|---|---|---| +| Simple rule | Template + reference | management → Kafka → rule pools → API → inspector | inspection latency | +| Aggregates | История платежей + aggregate rule | API → Thrift → Fraudo → ClickHouse | query latency, ClickHouse load | +| Black/white list | Записи через management | management → Kafka → wb-list → inspector | wb-list latency/errors | +| Grey list | RowInfo count/TTL + история | wb-list + ClickHouse → inspector | combined dependency latency | +| Geo IP | Rule `countryBy("ip")` | inspector → Columbus | cache hit/miss latency | +| Trusted token | Trusted-token conditions | inspector → trusted-tokens | dependency latency/errors | +| Configuration churn | Создание/удаление rules и references | management → Kafka → runtime/read-model | activation lag | +| Historical/emulation | Prepared datasets | management → Fraudbusters → ClickHouse | long-query latency | +| Failure modes | Controlled dependency delay/errors | circuit breakers, timeouts, recovery | error rate, recovery time | + +Для каждого dependency-профиля стоит иметь два запуска: + +1. **Integration smoke** с детерминированными данными и строгими checks. +2. **Isolated load profile**, где setup выполнен заранее и во время измерения + меняется только исследуемая нагрузка. + +Существующие сценарии в соседнем +`fraudbusters-compose/e2e-test/test/` уже содержат полезные fixtures для +aggregates, Columbus, Trusted Tokens и списков. Их следует переносить в k6 +по одному, сохраняя те же ожидаемые risk scores. + +## Ограничения интерпретации результатов + +Нагрузка через `fraudbusters-api` измеряет пользовательский REST-путь целиком, +включая JSON conversion и Thrift proxy. Для локализации bottleneck нужны +одновременные метрики: + +- `fraudbusters-api`: HTTP latency, errors, connection pools; +- `fraudbusters`: inspection/function timers, JVM, rule pool sizes; +- внешние сервисы: latency/error rate по `wb-list`, Columbus, Trusted Tokens; +- ClickHouse: query latency, running queries, CPU/IO; +- Kafka: producer latency, consumer lag и activation lag; +- PostgreSQL management: query latency и pool saturation. diff --git a/tests/k6/config.js b/tests/k6/config.js new file mode 100644 index 00000000..8ca544af --- /dev/null +++ b/tests/k6/config.js @@ -0,0 +1,24 @@ +function env(name, fallback) { + return __ENV[name] || fallback; +} + +export const config = { + apiBaseUrl: env("FB_API_URL", "http://host.docker.internal:9999"), + managementBaseUrl: env("FB_MANAGEMENT_URL", "http://host.docker.internal:8085/fb-management/v1"), + apiToken: env("FB_API_TOKEN", env("FB_TOKEN", "")), + managementToken: env("FB_MANAGEMENT_TOKEN", env("FB_TOKEN", "")), + activationWaitSeconds: Number(env("FB_ACTIVATION_WAIT_SECONDS", "5")), + cleanup: env("FB_CLEANUP", "false") === "true", + partyId: env("FB_PARTY_ID", "k6-load-party"), + shopId: env("FB_SHOP_ID", "k6-load-shop"), + templateId: env("FB_TEMPLATE_ID", "k6-load-template"), +}; + +export function authHeaders(token) { + const headers = { "Content-Type": "application/json" }; + if (token) { + headers.Authorization = `Bearer ${token}`; + } + return headers; +} + diff --git a/tests/k6/inspection-load.js b/tests/k6/inspection-load.js new file mode 100644 index 00000000..33a646e6 --- /dev/null +++ b/tests/k6/inspection-load.js @@ -0,0 +1,37 @@ +import { check } from "k6"; +import { config } from "./config.js"; +import { inspectPayment } from "./lib/client.js"; +import { inspectionRequest } from "./lib/payloads.js"; + +const rate = Number(__ENV.FB_RATE || "20"); +const duration = __ENV.FB_DURATION || "2m"; +const preAllocatedVUs = Number(__ENV.FB_PREALLOCATED_VUS || "20"); +const maxVUs = Number(__ENV.FB_MAX_VUS || "100"); + +export const options = { + scenarios: { + inspections: { + executor: "constant-arrival-rate", + rate, + timeUnit: "1s", + duration, + preAllocatedVUs, + maxVUs, + gracefulStop: "10s", + }, + }, + thresholds: { + checks: ["rate>0.99"], + http_req_failed: ["rate<0.01"], + "http_req_duration{endpoint:inspect-payment}": ["p(95)<500", "p(99)<1000"], + }, +}; + +export default function () { + const response = inspectPayment(inspectionRequest(config.partyId, config.shopId, 100)); + check(response, { + "inspect status is 200": (r) => r.status === 200, + "risk score is returned": (r) => ["low", "high", "fatal"].includes(r.json("result")), + }); +} + diff --git a/tests/k6/lib/client.js b/tests/k6/lib/client.js new file mode 100644 index 00000000..80103cd0 --- /dev/null +++ b/tests/k6/lib/client.js @@ -0,0 +1,88 @@ +import http from "k6/http"; +import { check, fail } from "k6"; +import { authHeaders, config } from "../config.js"; + +function jsonRequest(method, url, body, token, tags) { + const params = { + headers: authHeaders(token), + tags, + timeout: "30s", + }; + return http.request(method, url, body === null ? null : JSON.stringify(body), params); +} + +export function inspectPayment(body) { + return jsonRequest("POST", `${config.apiBaseUrl}/inspect-payment`, body, config.apiToken, { + service: "fraudbusters-api", + endpoint: "inspect-payment", + }); +} + +export function ingestPayments(body) { + return jsonRequest("POST", `${config.apiBaseUrl}/payments`, body, config.apiToken, { + service: "fraudbusters-api", + endpoint: "payments", + }); +} + +export function createTemplate(body) { + return jsonRequest( + "POST", + `${config.managementBaseUrl}/payments-templates`, + body, + config.managementToken, + { service: "fraudbusters-management", endpoint: "create-template" }, + ); +} + +export function createReferences(body) { + return jsonRequest( + "POST", + `${config.managementBaseUrl}/payments-references`, + body, + config.managementToken, + { service: "fraudbusters-management", endpoint: "create-reference" }, + ); +} + +export function filterTemplates(size = 20) { + return jsonRequest( + "GET", + `${config.managementBaseUrl}/payments-templates?size=${size}`, + null, + config.managementToken, + { service: "fraudbusters-management", endpoint: "filter-templates" }, + ); +} + +export function removeTemplate(id) { + return jsonRequest( + "DELETE", + `${config.managementBaseUrl}/payments-templates/${encodeURIComponent(id)}`, + null, + config.managementToken, + { service: "fraudbusters-management", endpoint: "remove-template" }, + ); +} + +export function removeReference(id) { + return jsonRequest( + "DELETE", + `${config.managementBaseUrl}/payments-references/${encodeURIComponent(id)}`, + null, + config.managementToken, + { service: "fraudbusters-management", endpoint: "remove-reference" }, + ); +} + +export function requireStatus(response, expected, label) { + const expectedStatuses = Array.isArray(expected) ? expected : [expected]; + const ok = check(response, { + [`${label}: expected status`]: (r) => expectedStatuses.includes(r.status), + }); + if (!ok) { + fail(`${label} failed: status=${response.status}, body=${response.body}`); + } + return response; +} + diff --git a/tests/k6/lib/payloads.js b/tests/k6/lib/payloads.js new file mode 100644 index 00000000..74189aa3 --- /dev/null +++ b/tests/k6/lib/payloads.js @@ -0,0 +1,89 @@ +function randomSuffix() { + return `${__VU}-${__ITER}-${Date.now()}-${Math.floor(Math.random() * 100000)}`; +} + +export function payment(partyId, shopId, amount = 100, overrides = {}) { + const suffix = randomSuffix(); + return { + id: overrides.id || `k6-payment-${suffix}`, + payerType: overrides.payerType || "customer", + merchant: { + id: partyId, + shop: { + id: shopId, + name: "k6 shop", + category: "test", + location: "test", + }, + }, + provider: { + id: "k6-provider", + terminalId: "k6-terminal", + country: "RUS", + }, + paymentResource: { + type: "bank_card", + cardToken: overrides.cardToken || `k6-card-${suffix}`, + lastDigits: "4242", + bin: "424242", + countryCode: "RUS", + bankName: "k6 bank", + paymentSystem: "visa", + cardType: "credit", + }, + cash: { + amount, + currency: overrides.currency || "RUB", + }, + customer: { + name: "k6 customer", + device: { + ip: overrides.ip || "192.0.2.10", + fingerprint: overrides.fingerprint || `k6-fingerprint-${suffix}`, + }, + contact: { + email: overrides.email || `k6-${suffix}@example.test`, + phone: "+70000000000", + }, + }, + createdAt: new Date().toISOString(), + description: "k6 fraudbusters test", + }; +} + +export function inspectionRequest(partyId, shopId, amount, overrides = {}) { + return { payment: payment(partyId, shopId, amount, overrides) }; +} + +export function paymentChangeRequest(partyId, shopId, amount, status = "processed", overrides = {}) { + return { + paymentsChanges: [ + { + payment: payment(partyId, shopId, amount, overrides), + paymentStatus: status, + eventTime: new Date().toISOString(), + }, + ], + }; +} + +export function template(id, expression) { + return { + id, + template: expression, + lastUpdateDate: new Date().toISOString(), + modifiedByUser: "k6", + }; +} + +export function reference(id, partyId, shopId, templateId) { + return { + id, + partyId, + shopId, + templateId, + lastUpdateDate: new Date().toISOString(), + modifiedByUser: "k6", + }; +} + diff --git a/tests/k6/lib/provision.js b/tests/k6/lib/provision.js new file mode 100644 index 00000000..5a34b6de --- /dev/null +++ b/tests/k6/lib/provision.js @@ -0,0 +1,39 @@ +import { sleep } from "k6"; +import { config } from "../config.js"; +import { + createReferences, + createTemplate, + removeReference, + removeTemplate, + requireStatus, +} from "./client.js"; +import { reference, template } from "./payloads.js"; + +export function provisionRule(prefix, expression = "rule:k6_amount:amount() >= 20 -> decline;") { + const runId = `${prefix}-${Date.now()}`; + const state = { + templateId: `${runId}-template`, + referenceId: `${runId}-reference`, + partyId: `${runId}-party`, + shopId: `${runId}-shop`, + }; + + requireStatus(createTemplate(template(state.templateId, expression)), [200, 201], "create template"); + requireStatus( + createReferences([reference(state.referenceId, state.partyId, state.shopId, state.templateId)]), + [200, 201], + "create reference", + ); + + sleep(config.activationWaitSeconds); + return state; +} + +export function cleanupRule(state) { + if (!config.cleanup || !state) { + return; + } + removeReference(state.referenceId); + removeTemplate(state.templateId); +} + diff --git a/tests/k6/mixed-load.js b/tests/k6/mixed-load.js new file mode 100644 index 00000000..477d31f8 --- /dev/null +++ b/tests/k6/mixed-load.js @@ -0,0 +1,55 @@ +import { check } from "k6"; +import { config } from "./config.js"; +import { filterTemplates, ingestPayments, inspectPayment } from "./lib/client.js"; +import { inspectionRequest, paymentChangeRequest } from "./lib/payloads.js"; + +export const options = { + scenarios: { + inspections: { + executor: "constant-arrival-rate", + exec: "inspection", + rate: Number(__ENV.FB_INSPECT_RATE || "20"), + timeUnit: "1s", + duration: __ENV.FB_DURATION || "2m", + preAllocatedVUs: 20, + maxVUs: 100, + }, + event_ingestion: { + executor: "constant-arrival-rate", + exec: "eventIngestion", + rate: Number(__ENV.FB_INGEST_RATE || "5"), + timeUnit: "1s", + duration: __ENV.FB_DURATION || "2m", + preAllocatedVUs: 5, + maxVUs: 30, + }, + management_reads: { + executor: "constant-vus", + exec: "managementRead", + vus: Number(__ENV.FB_MANAGEMENT_VUS || "2"), + duration: __ENV.FB_DURATION || "2m", + }, + }, + thresholds: { + http_req_failed: ["rate<0.01"], + "http_req_duration{endpoint:inspect-payment}": ["p(95)<500", "p(99)<1000"], + "http_req_duration{endpoint:payments}": ["p(95)<1000"], + "http_req_duration{endpoint:filter-templates}": ["p(95)<1000"], + }, +}; + +export function inspection() { + const response = inspectPayment(inspectionRequest(config.partyId, config.shopId, 100)); + check(response, { "inspection succeeds": (r) => r.status === 200 }); +} + +export function eventIngestion() { + const response = ingestPayments(paymentChangeRequest(config.partyId, config.shopId, 100, "processed")); + check(response, { "event ingestion succeeds": (r) => [200, 201].includes(r.status) }); +} + +export function managementRead() { + const response = filterTemplates(20); + check(response, { "management read succeeds": (r) => [200, 201].includes(r.status) }); +} + diff --git a/tests/k6/run.sh b/tests/k6/run.sh new file mode 100755 index 00000000..3b72b50b --- /dev/null +++ b/tests/k6/run.sh @@ -0,0 +1,26 @@ +#!/usr/bin/env sh +set -eu + +SCRIPT="${1:-smoke.js}" +SCRIPT_DIR=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) + +docker run --rm -i \ + -v "${SCRIPT_DIR}:/scripts:ro" \ + -e FB_API_URL="${FB_API_URL:-http://host.docker.internal:9999}" \ + -e FB_MANAGEMENT_URL="${FB_MANAGEMENT_URL:-http://host.docker.internal:8085/fb-management/v1}" \ + -e FB_API_TOKEN="${FB_API_TOKEN:-${FB_TOKEN:-}}" \ + -e FB_MANAGEMENT_TOKEN="${FB_MANAGEMENT_TOKEN:-${FB_TOKEN:-}}" \ + -e FB_ACTIVATION_WAIT_SECONDS="${FB_ACTIVATION_WAIT_SECONDS:-5}" \ + -e FB_CLEANUP="${FB_CLEANUP:-false}" \ + -e FB_PARTY_ID="${FB_PARTY_ID:-k6-load-party}" \ + -e FB_SHOP_ID="${FB_SHOP_ID:-k6-load-shop}" \ + -e FB_TEMPLATE_ID="${FB_TEMPLATE_ID:-k6-load-template}" \ + -e FB_RATE="${FB_RATE:-20}" \ + -e FB_DURATION="${FB_DURATION:-2m}" \ + -e FB_PREALLOCATED_VUS="${FB_PREALLOCATED_VUS:-20}" \ + -e FB_MAX_VUS="${FB_MAX_VUS:-100}" \ + -e FB_INSPECT_RATE="${FB_INSPECT_RATE:-20}" \ + -e FB_INGEST_RATE="${FB_INGEST_RATE:-5}" \ + -e FB_MANAGEMENT_VUS="${FB_MANAGEMENT_VUS:-2}" \ + grafana/k6:latest run "/scripts/${SCRIPT}" + diff --git a/tests/k6/smoke.js b/tests/k6/smoke.js new file mode 100644 index 00000000..4d6a4af7 --- /dev/null +++ b/tests/k6/smoke.js @@ -0,0 +1,51 @@ +import { check, group } from "k6"; +import { cleanupRule, provisionRule } from "./lib/provision.js"; +import { ingestPayments, inspectPayment, requireStatus } from "./lib/client.js"; +import { inspectionRequest, paymentChangeRequest } from "./lib/payloads.js"; + +export const options = { + vus: 1, + iterations: 1, + thresholds: { + checks: ["rate==1"], + http_req_failed: ["rate==0"], + "http_req_duration{endpoint:inspect-payment}": ["p(95)<1000"], + }, +}; + +export function setup() { + return provisionRule("k6-smoke"); +} + +export default function (state) { + group("configured rule returns fatal", () => { + const response = requireStatus( + inspectPayment(inspectionRequest(state.partyId, state.shopId, 100)), + 200, + "inspect fatal payment", + ); + check(response, { "risk score is fatal": (r) => r.json("result") === "fatal" }); + }); + + group("configured rule leaves low amount at default", () => { + const response = requireStatus( + inspectPayment(inspectionRequest(state.partyId, state.shopId, 10)), + 200, + "inspect default payment", + ); + check(response, { "risk score is high": (r) => r.json("result") === "high" }); + }); + + group("payment event reaches runtime ingestion API", () => { + requireStatus( + ingestPayments(paymentChangeRequest(state.partyId, state.shopId, 100, "processed")), + [200, 201], + "ingest payment", + ); + }); +} + +export function teardown(state) { + cleanupRule(state); +} + From bd418fd0588301d1266d227bd3e20c9ff5a2e414 Mon Sep 17 00:00:00 2001 From: kostyastruga Date: Mon, 10 Aug 2026 10:28:04 +0300 Subject: [PATCH 2/4] Remove unrelated load tests and metrics --- README.md | 4 +- .../config/payment/PaymentResourceConfig.java | 7 +- .../handler/FraudInspectorHandler.java | 12 -- tests/k6/README.md | 147 ------------------ tests/k6/config.js | 24 --- tests/k6/inspection-load.js | 37 ----- tests/k6/lib/client.js | 88 ----------- tests/k6/lib/payloads.js | 89 ----------- tests/k6/lib/provision.js | 39 ----- tests/k6/mixed-load.js | 55 ------- tests/k6/run.sh | 26 ---- tests/k6/smoke.js | 51 ------ 12 files changed, 3 insertions(+), 576 deletions(-) delete mode 100644 tests/k6/README.md delete mode 100644 tests/k6/config.js delete mode 100644 tests/k6/inspection-load.js delete mode 100644 tests/k6/lib/client.js delete mode 100644 tests/k6/lib/payloads.js delete mode 100644 tests/k6/lib/provision.js delete mode 100644 tests/k6/mixed-load.js delete mode 100755 tests/k6/run.sh delete mode 100644 tests/k6/smoke.js diff --git a/README.md b/README.md index 8f58148a..f139e48d 100644 --- a/README.md +++ b/README.md @@ -21,8 +21,6 @@ protocol for managing a set of patterns and bindings Instruction for testing in https://github.com/valitydev/fraudbusters-compose -System integration and load tests using `fraudbusters-api` and k6: -[tests/k6/README.md](tests/k6/README.md). - ### License [Apache 2.0 License.](/LICENSE) + diff --git a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java index 032b9595..256f2ac4 100644 --- a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java +++ b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentResourceConfig.java @@ -7,7 +7,6 @@ import dev.vality.fraudbusters.domain.FraudResult; import dev.vality.fraudbusters.resource.payment.handler.FraudInspectorHandler; import dev.vality.fraudbusters.stream.impl.TemplateVisitorImpl; -import io.micrometer.core.instrument.MeterRegistry; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -25,16 +24,14 @@ public InspectorProxySrv.Iface fraudInspectorHandler( CheckedResultToRiskScoreConverter checkedResultToRiskScoreConverter, ContextToFraudRequestConverter requestConverter, TemplateVisitorImpl templateVisitor, - WbListServiceSrv.Iface wbListServiceSrv, - MeterRegistry meterRegistry) { + WbListServiceSrv.Iface wbListServiceSrv) { return new FraudInspectorHandler( resultTopic, checkedResultToRiskScoreConverter, requestConverter, templateVisitor, kafkaFraudResultTemplate, - wbListServiceSrv, - meterRegistry + wbListServiceSrv ); } diff --git a/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java b/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java index e274f8e6..e66412b6 100644 --- a/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java +++ b/src/main/java/dev/vality/fraudbusters/resource/payment/handler/FraudInspectorHandler.java @@ -13,8 +13,6 @@ import dev.vality.fraudbusters.domain.FraudResult; import dev.vality.fraudbusters.fraud.model.PaymentModel; import dev.vality.fraudbusters.stream.TemplateVisitor; -import io.micrometer.core.instrument.MeterRegistry; -import io.micrometer.core.instrument.Timer; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.thrift.TException; @@ -30,12 +28,9 @@ public class FraudInspectorHandler implements InspectorProxySrv.Iface { private final TemplateVisitor templateVisitor; private final KafkaTemplate kafkaFraudResultTemplate; private final WbListServiceSrv.Iface wbListServiceSrv; - private final MeterRegistry meterRegistry; @Override public RiskScore inspectPayment(Context context) throws TException { - Timer.Sample sample = Timer.start(meterRegistry); - String outcome = "success"; try { FraudRequest model = requestConverter.convert(context); if (model != null) { @@ -47,15 +42,8 @@ public RiskScore inspectPayment(Context context) throws TException { } return RiskScore.high; } catch (Exception e) { - outcome = "error"; log.error("Error when inspectPayment() e: ", e); throw new TException("Error when inspectPayment() e: ", e); - } finally { - sample.stop(Timer.builder("fraudbusters.online.inspect") - .description("Online Fraudbusters inspectPayment latency") - .tag("method", "inspectPayment") - .tag("outcome", outcome) - .register(meterRegistry)); } } diff --git a/tests/k6/README.md b/tests/k6/README.md deleted file mode 100644 index 509d9193..00000000 --- a/tests/k6/README.md +++ /dev/null @@ -1,147 +0,0 @@ -# Fraudbusters k6 tests - -Набор использует реальные публичные HTTP-фасады: - -- `fraudbusters-api`: `/inspect-payment` и `/payments`; -- `fraudbusters-management`: создание правил/references и чтение control plane. - -Тем самым запрос k6 проходит через Swagger API, Thrift, Fraudbusters runtime, -Kafka configuration flow, внешние rule dependencies и хранилища, поднятые в -тестовом окружении. - -Архитектурно это позволяет измерять два уровня отдельно: - -- latency внешнего контракта `fraudbusters-api → Thrift → fraudbusters`; -- control/data plane всей системы вместе с management, Kafka и хранилищами. - -## Сценарии - -| Файл | Назначение | -|---|---| -| `smoke.js` | Сквозной интеграционный тест: создать правило и reference, дождаться активации, проверить `fatal/high`, загрузить payment event | -| `inspection-load.js` | Нагрузка только на критический онлайн-путь `/inspect-payment` | -| `mixed-load.js` | Смешанная нагрузка: inspection, event ingestion и management reads | - -## Предварительные условия - -Поднять полный контур, например из соседнего `fraudbusters-compose`, включая: - -- `fraudbusters-api` на host-порту `9999`; -- `fraudbusters-management` на host-порту `8085`; -- Fraudbusters, Kafka, ClickHouse, PostgreSQL, `wb-list` и остальные зависимости. - -В текущем `fraudbusters-compose/docker-compose.yml` определения основных -Fraudbusters-сервисов закомментированы. Перед запуском тестов нужен рабочий -compose/profile или развернутое тестовое окружение. - -Если включена авторизация, передать JWT: - -```bash -export FB_API_TOKEN=... -export FB_MANAGEMENT_TOKEN=... -``` - -## Запуск - -k6 не установлен локально, поэтому `run.sh` запускает официальный Docker image: - -```bash -./tests/k6/run.sh smoke.js -./tests/k6/run.sh inspection-load.js -./tests/k6/run.sh mixed-load.js -``` - -При локально установленном k6: - -```bash -k6 run tests/k6/smoke.js -``` - -## Настройка нагрузки - -```bash -FB_RATE=100 \ -FB_DURATION=10m \ -FB_PREALLOCATED_VUS=100 \ -FB_MAX_VUS=500 \ -./tests/k6/run.sh inspection-load.js -``` - -Для mixed workload: - -```bash -FB_INSPECT_RATE=100 \ -FB_INGEST_RATE=20 \ -FB_MANAGEMENT_VUS=5 \ -FB_DURATION=15m \ -./tests/k6/run.sh mixed-load.js -``` - -Основные переменные: - -| Переменная | Default | Назначение | -|---|---:|---| -| `FB_API_URL` | `http://host.docker.internal:9999` | fraudbusters-api | -| `FB_MANAGEMENT_URL` | `http://host.docker.internal:8085/fb-management/v1` | fraudbusters-management | -| `FB_ACTIVATION_WAIT_SECONDS` | `5` | Ожидание eventual consistency после создания reference | -| `FB_PARTY_ID`, `FB_SHOP_ID` | `k6-load-*` | Предварительно настроенный merchant для load-сценариев | -| `FB_CLEANUP` | `false` | Удалять созданные smoke-test template/reference | - -## Подготовка данных для load-тестов - -`inspection-load.js` и `mixed-load.js` намеренно не изменяют конфигурацию правил -во время нагрузки. До теста нужно один раз создать template/reference для -`FB_PARTY_ID` и `FB_SHOP_ID`, либо использовать merchant с уже активными -правилами. - -Для сравнимых прогонов фиксировать: - -- набор и сложность правил; -- объём истории ClickHouse; -- ответы/latency `wb-list`, Columbus и Trusted Tokens; -- число partitions и consumer concurrency Kafka; -- CPU/memory limits всех сервисов. - -Не смешивать тест максимальной пропускной способности inspection с массовым -созданием правил: это разные профили нагрузки и разные bottlenecks. - -## План покрытия всей системы - -Текущий каркас покрывает базовый сквозной путь и готов для расширения. Полный -system-test набор рекомендуется разделить на независимые профили: - -| Профиль | Подготовка | Проверяемая цепочка | Главная метрика | -|---|---|---|---| -| Simple rule | Template + reference | management → Kafka → rule pools → API → inspector | inspection latency | -| Aggregates | История платежей + aggregate rule | API → Thrift → Fraudo → ClickHouse | query latency, ClickHouse load | -| Black/white list | Записи через management | management → Kafka → wb-list → inspector | wb-list latency/errors | -| Grey list | RowInfo count/TTL + история | wb-list + ClickHouse → inspector | combined dependency latency | -| Geo IP | Rule `countryBy("ip")` | inspector → Columbus | cache hit/miss latency | -| Trusted token | Trusted-token conditions | inspector → trusted-tokens | dependency latency/errors | -| Configuration churn | Создание/удаление rules и references | management → Kafka → runtime/read-model | activation lag | -| Historical/emulation | Prepared datasets | management → Fraudbusters → ClickHouse | long-query latency | -| Failure modes | Controlled dependency delay/errors | circuit breakers, timeouts, recovery | error rate, recovery time | - -Для каждого dependency-профиля стоит иметь два запуска: - -1. **Integration smoke** с детерминированными данными и строгими checks. -2. **Isolated load profile**, где setup выполнен заранее и во время измерения - меняется только исследуемая нагрузка. - -Существующие сценарии в соседнем -`fraudbusters-compose/e2e-test/test/` уже содержат полезные fixtures для -aggregates, Columbus, Trusted Tokens и списков. Их следует переносить в k6 -по одному, сохраняя те же ожидаемые risk scores. - -## Ограничения интерпретации результатов - -Нагрузка через `fraudbusters-api` измеряет пользовательский REST-путь целиком, -включая JSON conversion и Thrift proxy. Для локализации bottleneck нужны -одновременные метрики: - -- `fraudbusters-api`: HTTP latency, errors, connection pools; -- `fraudbusters`: inspection/function timers, JVM, rule pool sizes; -- внешние сервисы: latency/error rate по `wb-list`, Columbus, Trusted Tokens; -- ClickHouse: query latency, running queries, CPU/IO; -- Kafka: producer latency, consumer lag и activation lag; -- PostgreSQL management: query latency и pool saturation. diff --git a/tests/k6/config.js b/tests/k6/config.js deleted file mode 100644 index 8ca544af..00000000 --- a/tests/k6/config.js +++ /dev/null @@ -1,24 +0,0 @@ -function env(name, fallback) { - return __ENV[name] || fallback; -} - -export const config = { - apiBaseUrl: env("FB_API_URL", "http://host.docker.internal:9999"), - managementBaseUrl: env("FB_MANAGEMENT_URL", "http://host.docker.internal:8085/fb-management/v1"), - apiToken: env("FB_API_TOKEN", env("FB_TOKEN", "")), - managementToken: env("FB_MANAGEMENT_TOKEN", env("FB_TOKEN", "")), - activationWaitSeconds: Number(env("FB_ACTIVATION_WAIT_SECONDS", "5")), - cleanup: env("FB_CLEANUP", "false") === "true", - partyId: env("FB_PARTY_ID", "k6-load-party"), - shopId: env("FB_SHOP_ID", "k6-load-shop"), - templateId: env("FB_TEMPLATE_ID", "k6-load-template"), -}; - -export function authHeaders(token) { - const headers = { "Content-Type": "application/json" }; - if (token) { - headers.Authorization = `Bearer ${token}`; - } - return headers; -} - diff --git a/tests/k6/inspection-load.js b/tests/k6/inspection-load.js deleted file mode 100644 index 33a646e6..00000000 --- a/tests/k6/inspection-load.js +++ /dev/null @@ -1,37 +0,0 @@ -import { check } from "k6"; -import { config } from "./config.js"; -import { inspectPayment } from "./lib/client.js"; -import { inspectionRequest } from "./lib/payloads.js"; - -const rate = Number(__ENV.FB_RATE || "20"); -const duration = __ENV.FB_DURATION || "2m"; -const preAllocatedVUs = Number(__ENV.FB_PREALLOCATED_VUS || "20"); -const maxVUs = Number(__ENV.FB_MAX_VUS || "100"); - -export const options = { - scenarios: { - inspections: { - executor: "constant-arrival-rate", - rate, - timeUnit: "1s", - duration, - preAllocatedVUs, - maxVUs, - gracefulStop: "10s", - }, - }, - thresholds: { - checks: ["rate>0.99"], - http_req_failed: ["rate<0.01"], - "http_req_duration{endpoint:inspect-payment}": ["p(95)<500", "p(99)<1000"], - }, -}; - -export default function () { - const response = inspectPayment(inspectionRequest(config.partyId, config.shopId, 100)); - check(response, { - "inspect status is 200": (r) => r.status === 200, - "risk score is returned": (r) => ["low", "high", "fatal"].includes(r.json("result")), - }); -} - diff --git a/tests/k6/lib/client.js b/tests/k6/lib/client.js deleted file mode 100644 index 80103cd0..00000000 --- a/tests/k6/lib/client.js +++ /dev/null @@ -1,88 +0,0 @@ -import http from "k6/http"; -import { check, fail } from "k6"; -import { authHeaders, config } from "../config.js"; - -function jsonRequest(method, url, body, token, tags) { - const params = { - headers: authHeaders(token), - tags, - timeout: "30s", - }; - return http.request(method, url, body === null ? null : JSON.stringify(body), params); -} - -export function inspectPayment(body) { - return jsonRequest("POST", `${config.apiBaseUrl}/inspect-payment`, body, config.apiToken, { - service: "fraudbusters-api", - endpoint: "inspect-payment", - }); -} - -export function ingestPayments(body) { - return jsonRequest("POST", `${config.apiBaseUrl}/payments`, body, config.apiToken, { - service: "fraudbusters-api", - endpoint: "payments", - }); -} - -export function createTemplate(body) { - return jsonRequest( - "POST", - `${config.managementBaseUrl}/payments-templates`, - body, - config.managementToken, - { service: "fraudbusters-management", endpoint: "create-template" }, - ); -} - -export function createReferences(body) { - return jsonRequest( - "POST", - `${config.managementBaseUrl}/payments-references`, - body, - config.managementToken, - { service: "fraudbusters-management", endpoint: "create-reference" }, - ); -} - -export function filterTemplates(size = 20) { - return jsonRequest( - "GET", - `${config.managementBaseUrl}/payments-templates?size=${size}`, - null, - config.managementToken, - { service: "fraudbusters-management", endpoint: "filter-templates" }, - ); -} - -export function removeTemplate(id) { - return jsonRequest( - "DELETE", - `${config.managementBaseUrl}/payments-templates/${encodeURIComponent(id)}`, - null, - config.managementToken, - { service: "fraudbusters-management", endpoint: "remove-template" }, - ); -} - -export function removeReference(id) { - return jsonRequest( - "DELETE", - `${config.managementBaseUrl}/payments-references/${encodeURIComponent(id)}`, - null, - config.managementToken, - { service: "fraudbusters-management", endpoint: "remove-reference" }, - ); -} - -export function requireStatus(response, expected, label) { - const expectedStatuses = Array.isArray(expected) ? expected : [expected]; - const ok = check(response, { - [`${label}: expected status`]: (r) => expectedStatuses.includes(r.status), - }); - if (!ok) { - fail(`${label} failed: status=${response.status}, body=${response.body}`); - } - return response; -} - diff --git a/tests/k6/lib/payloads.js b/tests/k6/lib/payloads.js deleted file mode 100644 index 74189aa3..00000000 --- a/tests/k6/lib/payloads.js +++ /dev/null @@ -1,89 +0,0 @@ -function randomSuffix() { - return `${__VU}-${__ITER}-${Date.now()}-${Math.floor(Math.random() * 100000)}`; -} - -export function payment(partyId, shopId, amount = 100, overrides = {}) { - const suffix = randomSuffix(); - return { - id: overrides.id || `k6-payment-${suffix}`, - payerType: overrides.payerType || "customer", - merchant: { - id: partyId, - shop: { - id: shopId, - name: "k6 shop", - category: "test", - location: "test", - }, - }, - provider: { - id: "k6-provider", - terminalId: "k6-terminal", - country: "RUS", - }, - paymentResource: { - type: "bank_card", - cardToken: overrides.cardToken || `k6-card-${suffix}`, - lastDigits: "4242", - bin: "424242", - countryCode: "RUS", - bankName: "k6 bank", - paymentSystem: "visa", - cardType: "credit", - }, - cash: { - amount, - currency: overrides.currency || "RUB", - }, - customer: { - name: "k6 customer", - device: { - ip: overrides.ip || "192.0.2.10", - fingerprint: overrides.fingerprint || `k6-fingerprint-${suffix}`, - }, - contact: { - email: overrides.email || `k6-${suffix}@example.test`, - phone: "+70000000000", - }, - }, - createdAt: new Date().toISOString(), - description: "k6 fraudbusters test", - }; -} - -export function inspectionRequest(partyId, shopId, amount, overrides = {}) { - return { payment: payment(partyId, shopId, amount, overrides) }; -} - -export function paymentChangeRequest(partyId, shopId, amount, status = "processed", overrides = {}) { - return { - paymentsChanges: [ - { - payment: payment(partyId, shopId, amount, overrides), - paymentStatus: status, - eventTime: new Date().toISOString(), - }, - ], - }; -} - -export function template(id, expression) { - return { - id, - template: expression, - lastUpdateDate: new Date().toISOString(), - modifiedByUser: "k6", - }; -} - -export function reference(id, partyId, shopId, templateId) { - return { - id, - partyId, - shopId, - templateId, - lastUpdateDate: new Date().toISOString(), - modifiedByUser: "k6", - }; -} - diff --git a/tests/k6/lib/provision.js b/tests/k6/lib/provision.js deleted file mode 100644 index 5a34b6de..00000000 --- a/tests/k6/lib/provision.js +++ /dev/null @@ -1,39 +0,0 @@ -import { sleep } from "k6"; -import { config } from "../config.js"; -import { - createReferences, - createTemplate, - removeReference, - removeTemplate, - requireStatus, -} from "./client.js"; -import { reference, template } from "./payloads.js"; - -export function provisionRule(prefix, expression = "rule:k6_amount:amount() >= 20 -> decline;") { - const runId = `${prefix}-${Date.now()}`; - const state = { - templateId: `${runId}-template`, - referenceId: `${runId}-reference`, - partyId: `${runId}-party`, - shopId: `${runId}-shop`, - }; - - requireStatus(createTemplate(template(state.templateId, expression)), [200, 201], "create template"); - requireStatus( - createReferences([reference(state.referenceId, state.partyId, state.shopId, state.templateId)]), - [200, 201], - "create reference", - ); - - sleep(config.activationWaitSeconds); - return state; -} - -export function cleanupRule(state) { - if (!config.cleanup || !state) { - return; - } - removeReference(state.referenceId); - removeTemplate(state.templateId); -} - diff --git a/tests/k6/mixed-load.js b/tests/k6/mixed-load.js deleted file mode 100644 index 477d31f8..00000000 --- a/tests/k6/mixed-load.js +++ /dev/null @@ -1,55 +0,0 @@ -import { check } from "k6"; -import { config } from "./config.js"; -import { filterTemplates, ingestPayments, inspectPayment } from "./lib/client.js"; -import { inspectionRequest, paymentChangeRequest } from "./lib/payloads.js"; - -export const options = { - scenarios: { - inspections: { - executor: "constant-arrival-rate", - exec: "inspection", - rate: Number(__ENV.FB_INSPECT_RATE || "20"), - timeUnit: "1s", - duration: __ENV.FB_DURATION || "2m", - preAllocatedVUs: 20, - maxVUs: 100, - }, - event_ingestion: { - executor: "constant-arrival-rate", - exec: "eventIngestion", - rate: Number(__ENV.FB_INGEST_RATE || "5"), - timeUnit: "1s", - duration: __ENV.FB_DURATION || "2m", - preAllocatedVUs: 5, - maxVUs: 30, - }, - management_reads: { - executor: "constant-vus", - exec: "managementRead", - vus: Number(__ENV.FB_MANAGEMENT_VUS || "2"), - duration: __ENV.FB_DURATION || "2m", - }, - }, - thresholds: { - http_req_failed: ["rate<0.01"], - "http_req_duration{endpoint:inspect-payment}": ["p(95)<500", "p(99)<1000"], - "http_req_duration{endpoint:payments}": ["p(95)<1000"], - "http_req_duration{endpoint:filter-templates}": ["p(95)<1000"], - }, -}; - -export function inspection() { - const response = inspectPayment(inspectionRequest(config.partyId, config.shopId, 100)); - check(response, { "inspection succeeds": (r) => r.status === 200 }); -} - -export function eventIngestion() { - const response = ingestPayments(paymentChangeRequest(config.partyId, config.shopId, 100, "processed")); - check(response, { "event ingestion succeeds": (r) => [200, 201].includes(r.status) }); -} - -export function managementRead() { - const response = filterTemplates(20); - check(response, { "management read succeeds": (r) => [200, 201].includes(r.status) }); -} - diff --git a/tests/k6/run.sh b/tests/k6/run.sh deleted file mode 100755 index 3b72b50b..00000000 --- a/tests/k6/run.sh +++ /dev/null @@ -1,26 +0,0 @@ -#!/usr/bin/env sh -set -eu - -SCRIPT="${1:-smoke.js}" -SCRIPT_DIR=$(CDPATH= cd -- "$(dirname -- "$0")" && pwd) - -docker run --rm -i \ - -v "${SCRIPT_DIR}:/scripts:ro" \ - -e FB_API_URL="${FB_API_URL:-http://host.docker.internal:9999}" \ - -e FB_MANAGEMENT_URL="${FB_MANAGEMENT_URL:-http://host.docker.internal:8085/fb-management/v1}" \ - -e FB_API_TOKEN="${FB_API_TOKEN:-${FB_TOKEN:-}}" \ - -e FB_MANAGEMENT_TOKEN="${FB_MANAGEMENT_TOKEN:-${FB_TOKEN:-}}" \ - -e FB_ACTIVATION_WAIT_SECONDS="${FB_ACTIVATION_WAIT_SECONDS:-5}" \ - -e FB_CLEANUP="${FB_CLEANUP:-false}" \ - -e FB_PARTY_ID="${FB_PARTY_ID:-k6-load-party}" \ - -e FB_SHOP_ID="${FB_SHOP_ID:-k6-load-shop}" \ - -e FB_TEMPLATE_ID="${FB_TEMPLATE_ID:-k6-load-template}" \ - -e FB_RATE="${FB_RATE:-20}" \ - -e FB_DURATION="${FB_DURATION:-2m}" \ - -e FB_PREALLOCATED_VUS="${FB_PREALLOCATED_VUS:-20}" \ - -e FB_MAX_VUS="${FB_MAX_VUS:-100}" \ - -e FB_INSPECT_RATE="${FB_INSPECT_RATE:-20}" \ - -e FB_INGEST_RATE="${FB_INGEST_RATE:-5}" \ - -e FB_MANAGEMENT_VUS="${FB_MANAGEMENT_VUS:-2}" \ - grafana/k6:latest run "/scripts/${SCRIPT}" - diff --git a/tests/k6/smoke.js b/tests/k6/smoke.js deleted file mode 100644 index 4d6a4af7..00000000 --- a/tests/k6/smoke.js +++ /dev/null @@ -1,51 +0,0 @@ -import { check, group } from "k6"; -import { cleanupRule, provisionRule } from "./lib/provision.js"; -import { ingestPayments, inspectPayment, requireStatus } from "./lib/client.js"; -import { inspectionRequest, paymentChangeRequest } from "./lib/payloads.js"; - -export const options = { - vus: 1, - iterations: 1, - thresholds: { - checks: ["rate==1"], - http_req_failed: ["rate==0"], - "http_req_duration{endpoint:inspect-payment}": ["p(95)<1000"], - }, -}; - -export function setup() { - return provisionRule("k6-smoke"); -} - -export default function (state) { - group("configured rule returns fatal", () => { - const response = requireStatus( - inspectPayment(inspectionRequest(state.partyId, state.shopId, 100)), - 200, - "inspect fatal payment", - ); - check(response, { "risk score is fatal": (r) => r.json("result") === "fatal" }); - }); - - group("configured rule leaves low amount at default", () => { - const response = requireStatus( - inspectPayment(inspectionRequest(state.partyId, state.shopId, 10)), - 200, - "inspect default payment", - ); - check(response, { "risk score is high": (r) => r.json("result") === "high" }); - }); - - group("payment event reaches runtime ingestion API", () => { - requireStatus( - ingestPayments(paymentChangeRequest(state.partyId, state.shopId, 100, "processed")), - [200, 201], - "ingest payment", - ); - }); -} - -export function teardown(state) { - cleanupRule(state); -} - From 588ec104c49c7b666c842fc0af176935ab99e8fd Mon Sep 17 00:00:00 2001 From: kostyastruga Date: Mon, 10 Aug 2026 10:52:17 +0300 Subject: [PATCH 3/4] Extract local payment field resolver --- .../config/payment/PaymentFraudoConfig.java | 10 +++-- .../LocalPaymentFieldResolver.java | 39 +++++++++++++++++++ .../LocalCountAggregatorDecorator.java | 29 +++++++------- .../LocalSumAggregatorDecorator.java | 32 +++++++-------- .../LocalUniqueValueAggregatorDecorator.java | 13 +++---- .../DatabasePaymentFieldResolver.java | 12 ------ .../LocalResultStorageRepositoryTest.java | 3 +- 7 files changed, 81 insertions(+), 57 deletions(-) create mode 100644 src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java diff --git a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java index a69e0f48..699099fe 100644 --- a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java +++ b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java @@ -2,6 +2,7 @@ import dev.vality.damsel.wb_list.WbListServiceSrv; import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.localstorage.LocalPaymentFieldResolver; import dev.vality.fraudbusters.fraud.localstorage.LocalResultStorageRepository; import dev.vality.fraudbusters.fraud.localstorage.aggregator.LocalCountAggregatorDecorator; import dev.vality.fraudbusters.fraud.localstorage.aggregator.LocalSumAggregatorDecorator; @@ -139,6 +140,7 @@ public CountPaymentAggregator countResultAggr RefundRepository refundRepository, ChargebackRepository chargebackRepository, DatabasePaymentFieldResolver databasePaymentFieldResolver, + LocalPaymentFieldResolver localPaymentFieldResolver, TimeBoundaryService timeBoundaryService) { CountAggregatorImpl countAggregatorDecorator = new CountAggregatorImpl( @@ -150,7 +152,7 @@ public CountPaymentAggregator countResultAggr ); return new LocalCountAggregatorDecorator( countAggregatorDecorator, - databasePaymentFieldResolver, + localPaymentFieldResolver, localResultStorageRepository, timeBoundaryService ); @@ -163,6 +165,7 @@ public SumPaymentAggregator sumResultAggregat RefundRepository refundRepository, ChargebackRepository chargebackRepository, DatabasePaymentFieldResolver databasePaymentFieldResolver, + LocalPaymentFieldResolver localPaymentFieldResolver, TimeBoundaryService timeBoundaryService) { SumAggregatorImpl sumAggregator = new SumAggregatorImpl( @@ -174,7 +177,7 @@ public SumPaymentAggregator sumResultAggregat ); return new LocalSumAggregatorDecorator( sumAggregator, - databasePaymentFieldResolver, + localPaymentFieldResolver, localResultStorageRepository, timeBoundaryService ); @@ -185,13 +188,14 @@ public UniqueValueAggregator uniqueValueResul LocalResultStorageRepository localResultStorageRepository, PaymentRepository fraudResultRepository, DatabasePaymentFieldResolver databasePaymentFieldResolver, + LocalPaymentFieldResolver localPaymentFieldResolver, TimeBoundaryService timeBoundaryService) { UniqueValueAggregatorImpl uniqueValueAggregator = new UniqueValueAggregatorImpl(databasePaymentFieldResolver, fraudResultRepository, timeBoundaryService); return new LocalUniqueValueAggregatorDecorator( uniqueValueAggregator, - databasePaymentFieldResolver, + localPaymentFieldResolver, localResultStorageRepository, timeBoundaryService ); diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java new file mode 100644 index 00000000..6584b9bb --- /dev/null +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java @@ -0,0 +1,39 @@ +package dev.vality.fraudbusters.fraud.localstorage; + +import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.model.FieldModel; +import dev.vality.fraudbusters.fraud.model.PaymentModel; +import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; +import lombok.RequiredArgsConstructor; +import org.jetbrains.annotations.NotNull; +import org.springframework.stereotype.Service; + +import java.util.ArrayList; +import java.util.List; +import java.util.stream.Collectors; + +@Service +@RequiredArgsConstructor +public class LocalPaymentFieldResolver { + + private final DatabasePaymentFieldResolver databasePaymentFieldResolver; + + public FieldModel resolve(PaymentCheckedField field, PaymentModel model) { + FieldModel fieldModel = databasePaymentFieldResolver.resolve(field, model); + return new FieldModel(resolveName(field), fieldModel.getValue()); + } + + public String resolveName(PaymentCheckedField field) { + return field.name(); + } + + @NotNull + public List resolveListFields(PaymentModel model, List fields) { + if (fields != null) { + return fields.stream() + .map(field -> resolve(field, model)) + .collect(Collectors.toList()); + } + return new ArrayList<>(); + } +} diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java index 43399d8d..5be2ddff 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalCountAggregatorDecorator.java @@ -4,11 +4,11 @@ import dev.vality.fraudbusters.exception.RuleFunctionException; import dev.vality.fraudbusters.fraud.AggregateGroupingFunction; import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.localstorage.LocalPaymentFieldResolver; import dev.vality.fraudbusters.fraud.localstorage.LocalResultStorageRepository; import dev.vality.fraudbusters.fraud.model.FieldModel; import dev.vality.fraudbusters.fraud.model.PaymentModel; import dev.vality.fraudbusters.fraud.payment.aggregator.clickhouse.CountAggregatorImpl; -import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; import dev.vality.fraudbusters.service.TimeBoundaryService; import dev.vality.fraudbusters.util.TimestampUtil; import dev.vality.fraudo.model.TimeWindow; @@ -25,7 +25,7 @@ public class LocalCountAggregatorDecorator implements CountPaymentAggregator { private final CountAggregatorImpl countAggregator; - private final DatabasePaymentFieldResolver databasePaymentFieldResolver; + private final LocalPaymentFieldResolver localPaymentFieldResolver; private final LocalResultStorageRepository localStorageRepository; private final TimeBoundaryService timeBoundaryService; @@ -36,11 +36,11 @@ public Integer count( TimeWindow timeWindow, List list) { Integer count = countAggregator.count(checkedField, paymentModel, timeWindow, list); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); Integer localCount = localStorageRepository.countOperationByField( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond() @@ -74,11 +74,10 @@ public Integer countError( Integer countError = countAggregator.countError(checkedField, paymentModel, timeWindow, errorCode, list); Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Integer localCount = localStorageRepository.countOperationErrorWithGroupBy( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), @@ -106,11 +105,10 @@ public Integer countError(PaymentCheckedField paymentCheckedField, PaymentModel Integer countError = countAggregator.countError(paymentCheckedField, paymentModel, timeWindow, list); Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(paymentCheckedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(paymentCheckedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Integer localCount = localStorageRepository.countOperationErrorWithGroupBy( - paymentCheckedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), @@ -173,11 +171,10 @@ private Integer getCount( try { Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Integer count = aggregateFunction.accept( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java index 778f9249..503b3273 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalSumAggregatorDecorator.java @@ -4,11 +4,11 @@ import dev.vality.fraudbusters.exception.RuleFunctionException; import dev.vality.fraudbusters.fraud.AggregateGroupingFunction; import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.localstorage.LocalPaymentFieldResolver; import dev.vality.fraudbusters.fraud.localstorage.LocalResultStorageRepository; import dev.vality.fraudbusters.fraud.model.FieldModel; import dev.vality.fraudbusters.fraud.model.PaymentModel; import dev.vality.fraudbusters.fraud.payment.aggregator.clickhouse.SumAggregatorImpl; -import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; import dev.vality.fraudbusters.service.TimeBoundaryService; import dev.vality.fraudbusters.util.TimestampUtil; import dev.vality.fraudo.model.TimeWindow; @@ -25,7 +25,7 @@ public class LocalSumAggregatorDecorator implements SumPaymentAggregator { private final SumAggregatorImpl sumAggregatorImpl; - private final DatabasePaymentFieldResolver databasePaymentFieldResolver; + private final LocalPaymentFieldResolver localPaymentFieldResolver; private final LocalResultStorageRepository localStorageRepository; private final TimeBoundaryService timeBoundaryService; @@ -36,13 +36,12 @@ public Double sum( TimeWindow timeWindow, List list) { Double sum = sumAggregatorImpl.sum(checkedField, paymentModel, timeWindow, list); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Long localSum = localStorageRepository.sumOperationByFieldWithGroupBy( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), @@ -81,11 +80,10 @@ public Double sumError(PaymentCheckedField checkedField, Double sumError = sumAggregatorImpl.sumError(checkedField, paymentModel, timeWindow, errorCode, list); Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Long localSum = localStorageRepository.sumOperationErrorWithGroupBy( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), @@ -115,11 +113,10 @@ public Double sumError(PaymentCheckedField checkedField, Double sumError = sumAggregatorImpl.sumError(checkedField, paymentModel, timeWindow, list); Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Long localSum = localStorageRepository.sumOperationErrorWithGroupBy( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), @@ -178,11 +175,10 @@ private Double getSum( try { Instant now = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(now, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(checkedField, paymentModel); - List eventFields = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(checkedField, paymentModel); + List eventFields = localPaymentFieldResolver.resolveListFields(paymentModel, list); Long sum = aggregateFunction.accept( - checkedField.name(), + resolve.getName(), resolve.getValue(), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java index 2543e6f5..3bc30bc4 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/aggregator/LocalUniqueValueAggregatorDecorator.java @@ -3,11 +3,11 @@ import dev.vality.fraudbusters.domain.TimeBound; import dev.vality.fraudbusters.exception.RuleFunctionException; import dev.vality.fraudbusters.fraud.constant.PaymentCheckedField; +import dev.vality.fraudbusters.fraud.localstorage.LocalPaymentFieldResolver; import dev.vality.fraudbusters.fraud.localstorage.LocalResultStorageRepository; import dev.vality.fraudbusters.fraud.model.FieldModel; import dev.vality.fraudbusters.fraud.model.PaymentModel; import dev.vality.fraudbusters.fraud.payment.aggregator.clickhouse.UniqueValueAggregatorImpl; -import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; import dev.vality.fraudbusters.service.TimeBoundaryService; import dev.vality.fraudbusters.util.TimestampUtil; import dev.vality.fraudo.aggregator.UniqueValueAggregator; @@ -23,7 +23,7 @@ public class LocalUniqueValueAggregatorDecorator implements UniqueValueAggregator { private final UniqueValueAggregatorImpl uniqueValueAggregator; - private final DatabasePaymentFieldResolver databasePaymentFieldResolver; + private final LocalPaymentFieldResolver localPaymentFieldResolver; private final LocalResultStorageRepository localStorageRepository; private final TimeBoundaryService timeBoundaryService; @@ -38,13 +38,12 @@ public Integer countUniqueValue( Integer uniq = uniqueValueAggregator.countUniqueValue(countField, paymentModel, onField, timeWindow, list); Instant timestamp = TimestampUtil.instantFromPaymentModel(paymentModel); TimeBound timeBound = timeBoundaryService.getBoundary(timestamp, timeWindow); - FieldModel resolve = databasePaymentFieldResolver.resolve(countField, paymentModel); - List fieldModels = - databasePaymentFieldResolver.resolveListFieldsForLocalStorage(paymentModel, list); + FieldModel resolve = localPaymentFieldResolver.resolve(countField, paymentModel); + List fieldModels = localPaymentFieldResolver.resolveListFields(paymentModel, list); Integer localUniqCountOperation = localStorageRepository.uniqCountOperationWithGroupBy( - countField.name(), + resolve.getName(), resolve.getValue(), - onField.name(), + localPaymentFieldResolver.resolveName(onField), timeBound.getLeft().getEpochSecond(), timeBound.getRight().getEpochSecond(), fieldModels diff --git a/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java b/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java index f98f15bb..c115d330 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/payment/resolver/DatabasePaymentFieldResolver.java @@ -25,18 +25,6 @@ public List resolveListFields(PaymentModel model, List(); } - @NotNull - public List resolveListFieldsForLocalStorage( - PaymentModel model, - List list) { - if (list != null) { - return list.stream() - .map(field -> new FieldModel(field.name(), resolve(field, model).getValue())) - .collect(Collectors.toList()); - } - return new ArrayList<>(); - } - public FieldModel resolve(PaymentCheckedField field, PaymentModel model) { if (field == null) { throw new UnknownFieldException(); diff --git a/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java b/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java index 2b3f85ff..d4e894b0 100644 --- a/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java +++ b/src/test/java/dev/vality/fraudbusters/fraud/localstorage/LocalResultStorageRepositoryTest.java @@ -88,13 +88,14 @@ void shouldApplyMainFieldFilterToStatusSums() { @Test void shouldResolveDifferentFieldNamesForClickHouseAndLocalStorage() { DatabasePaymentFieldResolver resolver = new DatabasePaymentFieldResolver(); + LocalPaymentFieldResolver localResolver = new LocalPaymentFieldResolver(resolver); PaymentModel payment = new PaymentModel(); payment.setCardToken(TOKEN); List clickHouseFields = resolver.resolveListFields(payment, List.of(PaymentCheckedField.CARD_TOKEN)); List localFields = - resolver.resolveListFieldsForLocalStorage(payment, List.of(PaymentCheckedField.CARD_TOKEN)); + localResolver.resolveListFields(payment, List.of(PaymentCheckedField.CARD_TOKEN)); assertEquals("cardToken", clickHouseFields.get(0).getName()); assertEquals("CARD_TOKEN", localFields.get(0).getName()); From 66e32634403546554103d861a3eef841b7023958 Mon Sep 17 00:00:00 2001 From: kostyastruga Date: Mon, 10 Aug 2026 11:10:38 +0300 Subject: [PATCH 4/4] Register local payment field resolver --- .../fraudbusters/config/payment/PaymentFraudoConfig.java | 6 ++++++ .../fraud/localstorage/LocalPaymentFieldResolver.java | 2 -- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java index 699099fe..5683fce6 100644 --- a/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java +++ b/src/main/java/dev/vality/fraudbusters/config/payment/PaymentFraudoConfig.java @@ -100,6 +100,12 @@ public FieldResolver paymentModelFieldResolve return new PaymentModelFieldResolver(); } + @Bean + public LocalPaymentFieldResolver localPaymentFieldResolver( + DatabasePaymentFieldResolver databasePaymentFieldResolver) { + return new LocalPaymentFieldResolver(databasePaymentFieldResolver); + } + @Bean public InListFinder paymentInListFinder( WbListServiceSrv.Iface wbListServiceSrv, diff --git a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java index 6584b9bb..96bf1851 100644 --- a/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java +++ b/src/main/java/dev/vality/fraudbusters/fraud/localstorage/LocalPaymentFieldResolver.java @@ -6,13 +6,11 @@ import dev.vality.fraudbusters.fraud.payment.resolver.DatabasePaymentFieldResolver; import lombok.RequiredArgsConstructor; import org.jetbrains.annotations.NotNull; -import org.springframework.stereotype.Service; import java.util.ArrayList; import java.util.List; import java.util.stream.Collectors; -@Service @RequiredArgsConstructor public class LocalPaymentFieldResolver {