From 0ec2901eb2c9479d9e375eaa322db75ab021cf7e Mon Sep 17 00:00:00 2001 From: pawana_backbase Date: Mon, 20 Jul 2026 15:27:00 +0530 Subject: [PATCH 1/2] TAR-944 : Upgrade BJS to update Investment Caboose Data --- .../configuration/InvestmentClientConfig.java | 71 ++++++++++--------- ...tmentIngestionConfigurationProperties.java | 2 - ...InvestmentRestServiceApiConfiguration.java | 35 +++++---- .../InvestmentServiceConfiguration.java | 24 ++++++- .../InvestmentWebClientConfiguration.java | 14 ++-- .../InvestmentWebClientProperties.java | 38 ++++++---- .../investment/service/AsyncTaskService.java | 17 +++-- .../service/InvestmentClientService.java | 11 +-- .../service/InvestmentPortfolioService.java | 8 ++- .../InvestmentClientConfigTest.java | 2 +- .../InvestmentServiceConfigurationTest.java | 6 +- 11 files changed, 143 insertions(+), 85 deletions(-) diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentClientConfig.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentClientConfig.java index 019fa5bdd..d267bee6f 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentClientConfig.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentClientConfig.java @@ -25,7 +25,9 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; import org.springframework.http.MediaType; +import org.springframework.http.client.reactive.ReactorClientHttpConnector; import org.springframework.http.codec.json.Jackson2JsonDecoder; import org.springframework.http.codec.json.Jackson2JsonEncoder; import org.springframework.web.reactive.function.client.WebClient; @@ -41,9 +43,11 @@ * {@link InvestmentWebClientConfiguration} which should be imported alongside this class. * The WebClient connection pool prevents resource exhaustion and 503 errors by limiting: * */ @Configuration @@ -60,8 +64,9 @@ public InvestmentClientConfig() { /** * Configuration for Investment service REST client (ClientApi). */ - @Bean - @ConditionalOnMissingBean + @Bean("investmentApiClient") + @ConditionalOnMissingBean(name = "investmentApiClient") + @Primary public ApiClient investmentApiClient(WebClient interServiceWebClient, @Qualifier("investmentHttpClient") HttpClient investmentHttpClient, ObjectMapper objectMapper, @@ -69,7 +74,8 @@ public ApiClient investmentApiClient(WebClient interServiceWebClient, ObjectMapper mapper = objectMapper.copy(); mapper.setSerializationInclusion(Include.NON_EMPTY); mapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS); - Builder mutate = interServiceWebClient.mutate(); + Builder mutate = interServiceWebClient.mutate() + .clientConnector(new ReactorClientHttpConnector(investmentHttpClient)); mutate.codecs(clientCodecConfigurer -> { Jackson2JsonEncoder encoder = new Jackson2JsonEncoder(mapper, MediaType.APPLICATION_JSON); Jackson2JsonDecoder decoder = new Jackson2JsonDecoder(mapper, MediaType.APPLICATION_JSON); @@ -81,80 +87,81 @@ public ApiClient investmentApiClient(WebClient interServiceWebClient, } @Bean - @ConditionalOnMissingBean - public ClientApi clientApi(ApiClient investmentApiClient) { + @Primary + public ClientApi clientApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new ClientApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public InvestmentProductsApi investmentProductsApi(ApiClient investmentApiClient) { + @Primary + public InvestmentProductsApi investmentProductsApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new InvestmentProductsApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public PortfolioApi portfolioApi(ApiClient investmentApiClient) { + @Primary + public PortfolioApi portfolioApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new PortfolioApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public FinancialAdviceApi financialAdviceApi(ApiClient investmentApiClient) { + @Primary + public FinancialAdviceApi financialAdviceApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new FinancialAdviceApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public AssetUniverseApi assetUniverseApi(ApiClient investmentApiClient) { + @Primary + public AssetUniverseApi assetUniverseApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new AssetUniverseApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public AllocationsApi allocationsApi(ApiClient investmentApiClient) { + @Primary + public AllocationsApi allocationsApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new AllocationsApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public InvestmentApi investmentApi(ApiClient investmentApiClient) { + @Primary + public InvestmentApi investmentApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new InvestmentApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public ContentApi contentApi(ApiClient investmentApiClient) { + @Primary + public ContentApi contentApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new ContentApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public PaymentsApi paymentsApi(ApiClient investmentApiClient) { + @Primary + public PaymentsApi paymentsApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new PaymentsApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public PortfolioTradingAccountsApi portfolioTradingAccountsApi(ApiClient investmentApiClient) { + @Primary + public PortfolioTradingAccountsApi portfolioTradingAccountsApi( + @Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new PortfolioTradingAccountsApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public CurrencyApi currencyApi(ApiClient investmentApiClient) { + @Primary + public CurrencyApi currencyApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new CurrencyApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public RiskAssessmentApi riskAssessmentApi(ApiClient investmentApiClient) { + @Primary + public RiskAssessmentApi riskAssessmentApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new RiskAssessmentApi(investmentApiClient); } @Bean - @ConditionalOnMissingBean - public AsyncBulkGroupsApi asyncBulkGroupsApi(ApiClient investmentApiClient) { + @Primary + public AsyncBulkGroupsApi asyncBulkGroupsApi(@Qualifier("investmentApiClient") ApiClient investmentApiClient) { return new AsyncBulkGroupsApi(investmentApiClient); } diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentIngestionConfigurationProperties.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentIngestionConfigurationProperties.java index 49e75c73a..6281b9d37 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentIngestionConfigurationProperties.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentIngestionConfigurationProperties.java @@ -2,7 +2,6 @@ package com.backbase.stream.configuration; import lombok.Data; -import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.context.properties.ConfigurationProperties; /** @@ -13,7 +12,6 @@ * {@link IngestConfigProperties}. */ @Data -@ConditionalOnBean(InvestmentServiceConfiguration.class) @ConfigurationProperties(prefix = "backbase.bootstrap.ingestions.investment") public class InvestmentIngestionConfigurationProperties { diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentRestServiceApiConfiguration.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentRestServiceApiConfiguration.java index cac0326d1..65d85466b 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentRestServiceApiConfiguration.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentRestServiceApiConfiguration.java @@ -18,6 +18,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Primary; import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter; import org.springframework.web.client.RestTemplate; @@ -37,8 +38,9 @@ public class InvestmentRestServiceApiConfiguration { /** * Configuration for Investment service REST client (ClientApi). */ - @Bean - @ConditionalOnMissingBean + @Bean("restInvestmentApiClient") + @ConditionalOnMissingBean(name = "restInvestmentApiClient") + @Primary public com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient( @Qualifier("interServiceRestTemplate") RestTemplate restTemplate, @Qualifier("restInvestmentObjectMapper") ObjectMapper restInvestmentObjectMapper) { @@ -65,47 +67,54 @@ public ObjectMapper restInvestmentObjectMapper(ObjectMapper legacyObjectMapper) } @Bean - @ConditionalOnMissingBean + @Primary public com.backbase.investment.api.service.sync.v1.ContentApi restContentApi( - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { return new com.backbase.investment.api.service.sync.v1.ContentApi(restInvestmentApiClient); } @Bean - @ConditionalOnMissingBean + @Primary public com.backbase.investment.api.service.sync.v1.AssetUniverseApi restAssetUniverseApi( - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { return new com.backbase.investment.api.service.sync.v1.AssetUniverseApi(restInvestmentApiClient); } @Bean - public InvestmentRestNewsContentService investmentNewsContentService(ContentApi restContentApi, - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { + @Primary + public InvestmentRestNewsContentService investmentNewsContentService( + @Qualifier("restContentApi") ContentApi restContentApi, + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { return new InvestmentRestNewsContentService(restContentApi, restInvestmentApiClient); } @Bean - public InvestmentRestDocumentContentService investmentRestContentDocumentService(ContentApi restContentApi, - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { + @Primary + public InvestmentRestDocumentContentService investmentRestContentDocumentService( + @Qualifier("restContentApi") ContentApi restContentApi, + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { return new InvestmentRestDocumentContentService(restContentApi, restInvestmentApiClient); } @Bean + @Primary public InvestmentRestAssetUniverseService investmentRestAssetUniverseService( - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient, + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient, IngestConfigProperties portfolioProperties) { return new InvestmentRestAssetUniverseService(restInvestmentApiClient, portfolioProperties); } @Bean + @Primary public InvestmentRestModelPortfolioService investmentRestModelPortfolioService( - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient) { return new InvestmentRestModelPortfolioService(restInvestmentApiClient); } @Bean + @Primary public InvestmentRestProductPortfolioService investmentRestProductPortfolioService( - com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient, + @Qualifier("restInvestmentApiClient") com.backbase.investment.api.service.sync.ApiClient restInvestmentApiClient, IngestConfigProperties portfolioProperties) { return new InvestmentRestProductPortfolioService(restInvestmentApiClient, portfolioProperties); } diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentServiceConfiguration.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentServiceConfiguration.java index 467392d24..954c55a85 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentServiceConfiguration.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentServiceConfiguration.java @@ -39,10 +39,13 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; +import org.springframework.context.annotation.Primary; @Import({ DbsApiClientsAutoConfiguration.class, - InvestmentClientConfig.class + InvestmentClientConfig.class, + InvestmentWebClientConfiguration.class, + InvestmentRestServiceApiConfiguration.class }) @EnableConfigurationProperties({ InvestmentIngestionConfigurationProperties.class, @@ -50,15 +53,17 @@ }) @RequiredArgsConstructor @Configuration -@ConditionalOnProperty(name = "backbase.bootstrap.ingestions.investment.enabled") +@ConditionalOnProperty(name = "backbase.bootstrap.ingestions.investment.enabled", havingValue = "true") public class InvestmentServiceConfiguration { @Bean + @Primary public InvestmentClientService investmentClientService(ClientApi clientApi) { return new InvestmentClientService(clientApi); } @Bean + @Primary public InvestmentModelPortfolioService investmentModelPortfolioService(FinancialAdviceApi financialAdviceApi, InvestmentRestModelPortfolioService investmentRestModelPortfolioService, IngestConfigProperties portfolioProperties) { @@ -67,6 +72,7 @@ public InvestmentModelPortfolioService investmentModelPortfolioService(Financial } @Bean + @Primary public InvestmentPortfolioProductService investmentPortfolioProductService( InvestmentProductsApi investmentProductsApi, IngestConfigProperties portfolioProperties, InvestmentModelPortfolioService modelPortfolioService, @@ -76,6 +82,7 @@ public InvestmentPortfolioProductService investmentPortfolioProductService( } @Bean + @Primary public InvestmentPortfolioService investmentPortfolioService(PortfolioApi portfolioApi, PaymentsApi paymentsApi, PortfolioTradingAccountsApi portfolioTradingAccountsApi, IngestConfigProperties portfolioProperties) { @@ -84,27 +91,32 @@ public InvestmentPortfolioService investmentPortfolioService(PortfolioApi portfo } @Bean + @Primary public InvestmentAssetUniverseService investmentAssetUniverseService(AssetUniverseApi assetUniverseApi, InvestmentRestAssetUniverseService investmentRestAssetUniverseService) { return new InvestmentAssetUniverseService(assetUniverseApi, investmentRestAssetUniverseService); } @Bean + @Primary public AsyncTaskService asyncTaskService(AsyncBulkGroupsApi asyncBulkGroupsApi) { return new AsyncTaskService(asyncBulkGroupsApi); } @Bean + @Primary public InvestmentAssetPriceService investmentAssetPriceService(AssetUniverseApi assetUniverseApi) { return new InvestmentAssetPriceService(assetUniverseApi); } @Bean + @Primary public InvestmentIntradayAssetPriceService investmentIntradayAssetPriceService(AssetUniverseApi assetUniverseApi) { return new InvestmentIntradayAssetPriceService(assetUniverseApi); } @Bean + @Primary public InvestmentPortfolioAllocationService investmentPortfolioAllocationService(AllocationsApi allocationsApi, AssetUniverseApi assetUniverseApi, InvestmentApi investmentApi, IngestConfigProperties portfolioProperties) { @@ -113,22 +125,26 @@ public InvestmentPortfolioAllocationService investmentPortfolioAllocationService } @Bean + @Primary public InvestmentCurrencyService investmentCurrencyService(CurrencyApi currencyApi) { return new InvestmentCurrencyService(currencyApi); } @Bean + @Primary public InvestmentRiskAssessmentService investmentRiskAssessmentService(RiskAssessmentApi riskAssessmentApi) { return new InvestmentRiskAssessmentService(riskAssessmentApi); } @Bean + @Primary public InvestmentRiskQuestionaryService investmentRiskQuestionaryService(RiskAssessmentApi riskAssessmentApi, IngestConfigProperties portfolioProperties) { return new InvestmentRiskQuestionaryService(riskAssessmentApi, portfolioProperties); } @Bean + @Primary public InvestmentSaga investmentSaga(InvestmentClientService investmentClientService, InvestmentRiskAssessmentService investmentRiskAssessmentService, InvestmentRiskQuestionaryService investmentRiskQuestionaryService, @@ -145,7 +161,8 @@ public InvestmentSaga investmentSaga(InvestmentClientService investmentClientSer } @Bean - public InvestmentAssetUniverseSaga investmentStaticDataSaga( + @Primary + public InvestmentAssetUniverseSaga investmentAssetUniverseSaga( InvestmentAssetUniverseService investmentAssetUniverseService, InvestmentAssetPriceService investmentAssetPriceService, InvestmentIntradayAssetPriceService investmentIntradayAssetPriceService, @@ -159,6 +176,7 @@ public InvestmentAssetUniverseSaga investmentStaticDataSaga( } @Bean + @Primary public InvestmentContentSaga investmentContentSaga( InvestmentRestNewsContentService investmentRestNewsContentService, InvestmentRestDocumentContentService investmentRestDocumentContentService, diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientConfiguration.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientConfiguration.java index 4ade2ad4b..34252faea 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientConfiguration.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientConfiguration.java @@ -7,7 +7,7 @@ import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Qualifier; -import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import reactor.netty.http.client.HttpClient; @@ -31,9 +31,14 @@ */ @Slf4j @Configuration -@EnableConfigurationProperties(InvestmentWebClientProperties.class) public class InvestmentWebClientConfiguration { + @Bean("investmentWebClientProperties") + @ConfigurationProperties("backbase.communication.services.investment.http-client") + public InvestmentWebClientProperties investmentWebClientProperties() { + return new InvestmentWebClientProperties(); + } + /** * Dedicated {@link ConnectionProvider} for the Investment service client pool. * @@ -45,7 +50,8 @@ public class InvestmentWebClientConfiguration { * @return investment-specific ConnectionProvider */ @Bean("investmentConnectionProvider") - public ConnectionProvider investmentConnectionProvider(InvestmentWebClientProperties props) { + public ConnectionProvider investmentConnectionProvider( + @Qualifier("investmentWebClientProperties") InvestmentWebClientProperties props) { ConnectionProvider provider = ConnectionProvider.builder("investment-client-pool") .maxConnections(props.getMaxConnections()) .maxIdleTime(Duration.ofMinutes(props.getMaxIdleTimeMinutes())) @@ -81,7 +87,7 @@ public ConnectionProvider investmentConnectionProvider(InvestmentWebClientProper @Bean("investmentHttpClient") public HttpClient investmentHttpClient( @Qualifier("investmentConnectionProvider") ConnectionProvider connectionProvider, - InvestmentWebClientProperties props) { + @Qualifier("investmentWebClientProperties") InvestmentWebClientProperties props) { return HttpClient.create(connectionProvider) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, props.getConnectTimeoutSeconds() * 1000) .option(ChannelOption.SO_KEEPALIVE, true) diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientProperties.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientProperties.java index af6c23775..5dfa89492 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientProperties.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/configuration/InvestmentWebClientProperties.java @@ -1,39 +1,42 @@ package com.backbase.stream.configuration; import lombok.Data; -import org.springframework.boot.context.properties.ConfigurationProperties; /** * Configuration properties for the Investment service HTTP client connection pool and timeouts. * - *

All values can be overridden via {@code application.yml} / {@code application.properties} - * using the prefix {@code backbase.communication.services.investment.http-client}. + *

Bound via {@code backbase.communication.services.investment.http-client} and + * {@code backbase.communication.services.investment-caboose.http-client}. * - *

Example: + *

Example (values shown are illustrative; defaults are in field JavaDoc below): *

  * backbase:
  *   communication:
  *     services:
  *       investment:
  *         http-client:
- *           max-connections: 20
+ *           max-connections: 50
  *           max-idle-time-minutes: 5
- *           max-pending-acquires: 100
- *           pending-acquire-timeout-millis: 45000
+ *           max-life-time-minutes: 30
+ *           max-pending-acquires: -1
+ *           pending-acquire-timeout-millis: 90000
+ *           evict-in-background-seconds: 120
  *           connect-timeout-seconds: 10
  *           read-timeout-seconds: 30
  *           write-timeout-seconds: 30
+ *       investment-caboose:
+ *         http-client:
+ *           max-connections: 50
  * 
*/ @Data -@ConfigurationProperties(prefix = "backbase.communication.services.investment.http-client") public class InvestmentWebClientProperties { /** * Maximum number of open TCP connections to the Investment service. * Limiting this prevents the service from being overwhelmed (which causes 503 responses). */ - private int maxConnections = 20; + private int maxConnections = 50; /** * Maximum time (in minutes) that a connection can remain idle in the pool before being evicted. @@ -46,15 +49,24 @@ public class InvestmentWebClientProperties { private long maxLifeTimeMinutes = 30; /** - * Maximum number of requests that can be queued waiting for a connection. - * Bounds the in-memory queue so callers receive a fast failure rather than an unbounded backlog. + * Maximum number of requests that can wait for a free connection. + * + *

{@code -1} means no queue limit — callers block up to + * {@link #pendingAcquireTimeoutMillis} instead of failing fast with + * "Pending acquire queue has reached its maximum size". + * + *

Pair with application-level {@code flatMap} concurrency limits in the ingestion + * services so the queue does not grow without bound. */ - private int maxPendingAcquires = 100; + private int maxPendingAcquires = -1; /** * Maximum time (in milliseconds) a request will wait to acquire a connection from the pool. + * + *

When all {@link #maxConnections} are in use, new requests wait here rather than + * opening additional connections or failing immediately. */ - private long pendingAcquireTimeoutMillis = 45_000; + private long pendingAcquireTimeoutMillis = 90_000; /** * Background eviction interval (in seconds) for idle/expired connections in the pool. diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/AsyncTaskService.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/AsyncTaskService.java index 32c3c7cba..6f2b9c5b6 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/AsyncTaskService.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/AsyncTaskService.java @@ -22,20 +22,23 @@ public Mono> checkPriceAsyncTasksFinished(List as return Mono.just(List.of()); } return Flux.interval(java.time.Duration.ofSeconds(5)) - // Poll the status of all tasks .flatMap( tick -> Flux.fromIterable(asyncTasks) .flatMap(gr -> this.groupResultStatus(gr.getUuid())) .collectList() + .doOnNext(results -> { + long pending = results.stream() + .filter(gr -> "PENDING".equalsIgnoreCase(gr.getStatus())) + .count(); + if (pending > 0) { + log.info("Waiting for price async tasks: pending={}/{}", pending, asyncTasks.size()); + } + }) ) .filter(results -> results.stream().noneMatch(gr -> "PENDING".equalsIgnoreCase(gr.getStatus()))) .next() - .timeout(java.time.Duration.ofMinutes(5)) - .doOnSuccess(tasks -> { - log.info("Prices tasks finished added"); - log.debug("Price async tasks failure: {}", - tasks.stream().filter(gr -> "FAILURE".equalsIgnoreCase(gr.getStatus())).toList()); - }) + .timeout(java.time.Duration.ofMinutes(10)) + .doOnSuccess(tasks -> log.info("Prices tasks finished, added")) .onErrorResume(throwable -> { log.error("Timeout or error waiting for GroupResult tasks: taskIds={}", asyncTasks.stream().map(GroupResult::getUuid).toList(), throwable); diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentClientService.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentClientService.java index 0e8779a68..7d02a08af 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentClientService.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentClientService.java @@ -54,6 +54,7 @@ public class InvestmentClientService { *

  • Limits concurrent requests to avoid overwhelming the service and triggering 503 errors
  • *
  • Implements exponential backoff retry with 503 and 409 conflict handling
  • *
  • Processes clients sequentially through the same ExecutorService to maintain order and prevent race conditions
  • + *
  • Failed clients are logged and skipped to prevent batch failures
  • * * * @param clientUsers the list of clients to upsert @@ -99,10 +100,12 @@ public Mono> upsertClients(List clientUsers) { .doOnSuccess(upsertedClient -> log.debug( "Successfully upserted client: investmentClientId={}, internalUserId={}", upsertedClient.getInvestmentClientId(), upsertedClient.getInternalUserId())) - .doOnError(throwable -> log.error( - "Failed to upsert client: internalUserId={}, externalUserId={}, legalEntityExternalId={}", - clientUser.getInternalUserId(), clientUser.getExternalUserId(), - clientUser.getLegalEntityId(), throwable)); + .onErrorResume(throwable -> { + log.warn("Skipping client due to error: internalUserId={}, externalUserId={}, legalEntityExternalId={}", + clientUser.getInternalUserId(), clientUser.getExternalUserId(), + clientUser.getLegalEntityId(), throwable); + return Mono.empty(); + }); }, 5) // Max 5 concurrent requests to avoid overwhelming the service diff --git a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentPortfolioService.java b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentPortfolioService.java index 0006816e7..358bf7ad6 100644 --- a/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentPortfolioService.java +++ b/stream-investment/investment-core/src/main/java/com/backbase/stream/investment/service/InvestmentPortfolioService.java @@ -83,9 +83,11 @@ public Mono> upsertPortfolios(List log.debug( "Successfully upserted investment portfolio: portfolioUuid={}, externalId={}, name={}", ip.getPortfolio().getUuid(), ip.getPortfolio().getExternalId(), ip.getPortfolio().getName())) - .doOnError(throwable -> log.error( - "Failed to upsert investment portfolio: arrangementExternalId={}, arrangementName={}", - arrangement.getExternalId(), arrangement.getName(), throwable)); + .onErrorResume(throwable -> { + log.warn("Skipping investment portfolio due to error: arrangementExternalId={}, arrangementName={}", + arrangement.getExternalId(), arrangement.getName(), throwable); + return Mono.empty(); + }); }) .collectList(); } diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentClientConfigTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentClientConfigTest.java index 98242d1b2..eee081c42 100644 --- a/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentClientConfigTest.java +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentClientConfigTest.java @@ -71,7 +71,7 @@ void investmentApiClient_returnsNonNullApiClientWithCodecsConfigured() { WebClient webClient = WebClient.builder().build(); ObjectMapper objectMapper = new ObjectMapper(); DateFormat dateFormat = mock(DateFormat.class); - HttpClient httpClient = mock(HttpClient.class); // declared parameter; not used by the method body + HttpClient httpClient = HttpClient.create(); ApiClient result = config.investmentApiClient(webClient, httpClient, objectMapper, dateFormat); diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentServiceConfigurationTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentServiceConfigurationTest.java index 0af267eb0..91cffa2b3 100644 --- a/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentServiceConfigurationTest.java +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentServiceConfigurationTest.java @@ -180,9 +180,9 @@ void investmentSaga_returnsInvestmentSaga() { } @Test - @DisplayName("investmentStaticDataSaga — returns non-null InvestmentAssetUniverseSaga") - void investmentStaticDataSaga_returnsAssetUniverseSaga() { - assertThat(config.investmentStaticDataSaga( + @DisplayName("investmentAssetUniverseSaga — returns non-null InvestmentAssetUniverseSaga") + void investmentAssetUniverseSaga_returnsAssetUniverseSaga() { + assertThat(config.investmentAssetUniverseSaga( mock(InvestmentAssetUniverseService.class), mock(InvestmentAssetPriceService.class), mock(InvestmentIntradayAssetPriceService.class), From 045de3d5986c56816f6bbd920b24eafa6c147d0e Mon Sep 17 00:00:00 2001 From: pawana_backbase Date: Mon, 20 Jul 2026 18:18:25 +0530 Subject: [PATCH 2/2] TAR-944 : Upgrade BJS to update Investment Caboose Data - Tests --- CHANGELOG.md | 4 + .../InvestmentWebClientConfigurationTest.java | 57 ++++ .../service/AsyncTaskServiceTest.java | 149 +++++++++++ .../service/InvestmentClientServiceTest.java | 248 +++++++++++++++--- .../InvestmentPortfolioServiceTest.java | 74 ++++++ 5 files changed, 492 insertions(+), 40 deletions(-) create mode 100644 stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentWebClientConfigurationTest.java create mode 100644 stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/AsyncTaskServiceTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index 85a78875f..d27c3d0e6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,10 @@ # Changelog All notable changes to this project will be documented in this file. +## [10.6.2] +### Changed +- Updated stream-investment to be able to seed to investment caboose + ## [10.6.1] ### Changed - Updated plan manager ingestion logic to attempt to retrieve plans if they don't exist in the pre-filled map. diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentWebClientConfigurationTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentWebClientConfigurationTest.java new file mode 100644 index 000000000..ef5d1f8ce --- /dev/null +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/configuration/InvestmentWebClientConfigurationTest.java @@ -0,0 +1,57 @@ +package com.backbase.stream.configuration; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import reactor.netty.http.client.HttpClient; +import reactor.netty.resources.ConnectionProvider; + +/** + * Unit tests for {@link InvestmentWebClientConfiguration}. + * + *

    Each {@code @Bean} factory method is invoked directly (no Spring context) to verify + * that it returns a non-null, correctly typed bean. + */ +@DisplayName("InvestmentWebClientConfiguration") +class InvestmentWebClientConfigurationTest { + + private InvestmentWebClientConfiguration config; + + @BeforeEach + void setUp() { + config = new InvestmentWebClientConfiguration(); + } + + @Test + @DisplayName("investmentWebClientProperties — returns non-null properties with defaults") + void investmentWebClientProperties_returnsNonNullProperties() { + InvestmentWebClientProperties properties = config.investmentWebClientProperties(); + + assertThat(properties).isNotNull(); + assertThat(properties.getMaxConnections()).isPositive(); + assertThat(properties.getConnectTimeoutSeconds()).isPositive(); + } + + @Test + @DisplayName("investmentConnectionProvider — returns configured connection provider") + void investmentConnectionProvider_returnsConfiguredProvider() { + InvestmentWebClientProperties properties = config.investmentWebClientProperties(); + + ConnectionProvider provider = config.investmentConnectionProvider(properties); + + assertThat(provider).isNotNull(); + } + + @Test + @DisplayName("investmentHttpClient — returns configured HttpClient") + void investmentHttpClient_returnsConfiguredHttpClient() { + InvestmentWebClientProperties properties = config.investmentWebClientProperties(); + ConnectionProvider provider = config.investmentConnectionProvider(properties); + + HttpClient httpClient = config.investmentHttpClient(provider, properties); + + assertThat(httpClient).isNotNull(); + } +} diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/AsyncTaskServiceTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/AsyncTaskServiceTest.java new file mode 100644 index 000000000..841194e47 --- /dev/null +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/AsyncTaskServiceTest.java @@ -0,0 +1,149 @@ +package com.backbase.stream.investment.service; + +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.backbase.investment.api.service.v1.AsyncBulkGroupsApi; +import com.backbase.investment.api.service.v1.model.GroupResult; +import java.time.Duration; +import java.util.List; +import java.util.UUID; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; +import reactor.core.publisher.Mono; +import reactor.test.StepVerifier; + +@DisplayName("AsyncTaskService") +class AsyncTaskServiceTest { + + private AsyncBulkGroupsApi asyncBulkGroupsApi; + private AsyncTaskService service; + + @BeforeEach + void setUp() { + asyncBulkGroupsApi = mock(AsyncBulkGroupsApi.class); + service = new AsyncTaskService(asyncBulkGroupsApi); + } + + @Nested + @DisplayName("checkPriceAsyncTasksFinished") + class CheckPriceAsyncTasksFinishedTests { + + @Test + @DisplayName("empty task list — returns empty list immediately") + void emptyTaskList_returnsEmptyList() { + StepVerifier.create(service.checkPriceAsyncTasksFinished(List.of())) + .expectNext(List.of()) + .verifyComplete(); + } + + @Test + @DisplayName("all tasks completed on first poll — returns polled results") + void allTasksCompletedOnFirstPoll_returnsPolledResults() { + UUID uuid = UUID.randomUUID(); + GroupResult inputTask = new GroupResult(uuid, "PENDING", List.of()); + GroupResult completed = new GroupResult(uuid, "COMPLETED", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())).thenReturn(Mono.just(completed)); + + StepVerifier.withVirtualTime(() -> service.checkPriceAsyncTasksFinished(List.of(inputTask))) + .thenAwait(Duration.ofSeconds(5)) + .expectNextMatches(results -> results.size() == 1 + && "COMPLETED".equalsIgnoreCase(results.getFirst().getStatus())) + .verifyComplete(); + } + + @Test + @DisplayName("tasks transition from pending to completed — polls until finished") + void tasksTransitionFromPendingToCompleted_pollsUntilFinished() { + UUID uuid = UUID.randomUUID(); + GroupResult inputTask = new GroupResult(uuid, "PENDING", List.of()); + GroupResult completed = new GroupResult(uuid, "COMPLETED", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())) + .thenReturn(Mono.just(new GroupResult(uuid, "PENDING", List.of()))) + .thenReturn(Mono.just(completed)); + + StepVerifier.withVirtualTime(() -> service.checkPriceAsyncTasksFinished(List.of(inputTask))) + .thenAwait(Duration.ofSeconds(10)) + .expectNextMatches(results -> results.size() == 1 + && "COMPLETED".equalsIgnoreCase(results.getFirst().getStatus())) + .verifyComplete(); + } + + @Test + @DisplayName("pending status is case-insensitive — completes when status is lowercase") + void pendingStatusCaseInsensitive_completesWhenNoLongerPending() { + UUID uuid = UUID.randomUUID(); + GroupResult inputTask = new GroupResult(uuid, "PENDING", List.of()); + GroupResult completed = new GroupResult(uuid, "completed", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())) + .thenReturn(Mono.just(new GroupResult(uuid, "pending", List.of()))) + .thenReturn(Mono.just(completed)); + + StepVerifier.withVirtualTime(() -> service.checkPriceAsyncTasksFinished(List.of(inputTask))) + .thenAwait(Duration.ofSeconds(10)) + .expectNextMatches(results -> results.size() == 1 + && "completed".equals(results.getFirst().getStatus())) + .verifyComplete(); + } + + @Test + @DisplayName("timeout waiting for tasks — returns original task list") + void timeoutWaitingForTasks_returnsOriginalTaskList() { + UUID uuid = UUID.randomUUID(); + GroupResult inputTask = new GroupResult(uuid, "PENDING", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())) + .thenReturn(Mono.just(new GroupResult(uuid, "PENDING", List.of()))); + + StepVerifier.withVirtualTime(() -> service.checkPriceAsyncTasksFinished(List.of(inputTask))) + .thenAwait(Duration.ofMinutes(11)) + .expectNextMatches(results -> results.size() == 1 + && uuid.equals(results.getFirst().getUuid())) + .verifyComplete(); + } + + @Test + @DisplayName("API error while polling — returns original task list") + void apiErrorWhilePolling_returnsOriginalTaskList() { + UUID uuid = UUID.randomUUID(); + GroupResult inputTask = new GroupResult(uuid, "PENDING", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())) + .thenReturn(Mono.error(new RuntimeException("bulk group lookup failed"))); + + StepVerifier.withVirtualTime(() -> service.checkPriceAsyncTasksFinished(List.of(inputTask))) + .thenAwait(Duration.ofSeconds(5)) + .expectNextMatches(results -> results.size() == 1 + && uuid.equals(results.getFirst().getUuid())) + .verifyComplete(); + } + } + + @Nested + @DisplayName("groupResultStatus") + class GroupResultStatusTests { + + @Test + @DisplayName("delegates to AsyncBulkGroupsApi") + void delegatesToAsyncBulkGroupsApi() { + UUID uuid = UUID.randomUUID(); + GroupResult groupResult = new GroupResult(uuid, "COMPLETED", List.of()); + + when(asyncBulkGroupsApi.getBulkGroup(uuid.toString())).thenReturn(Mono.just(groupResult)); + + StepVerifier.create(service.groupResultStatus(uuid)) + .expectNext(groupResult) + .verifyComplete(); + + verify(asyncBulkGroupsApi, times(1)).getBulkGroup(eq(uuid.toString())); + } + } +} diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentClientServiceTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentClientServiceTest.java index a391ee1ab..874cc0c7e 100644 --- a/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentClientServiceTest.java +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentClientServiceTest.java @@ -1,77 +1,245 @@ package com.backbase.stream.investment.service; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyList; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.isNull; +import static org.mockito.Mockito.lenient; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + import com.backbase.investment.api.service.v1.ClientApi; import com.backbase.investment.api.service.v1.model.ClientCreate; -import com.backbase.investment.api.service.v1.model.ClientCreateRequest; import com.backbase.investment.api.service.v1.model.OASClient; import com.backbase.investment.api.service.v1.model.OASClientUpdateRequest; +import com.backbase.investment.api.service.v1.model.PaginatedOASClientList; import com.backbase.investment.api.service.v1.model.PatchedOASClientUpdateRequest; +import com.backbase.stream.investment.ClientUser; import java.nio.charset.StandardCharsets; +import java.util.List; import java.util.UUID; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; -import org.mockito.Mockito; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; import org.springframework.web.reactive.function.client.WebClientResponseException; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; -import static org.mockito.ArgumentMatchers.*; -import static org.mockito.Mockito.when; - +@DisplayName("InvestmentClientService") class InvestmentClientServiceTest { - ClientApi clientApi; - InvestmentClientService service; + private ClientApi clientApi; + private InvestmentClientService service; @BeforeEach void setUp() { - clientApi = Mockito.mock(ClientApi.class); + clientApi = mock(ClientApi.class); service = new InvestmentClientService(clientApi); } - @Test - void createClient_success() { - ClientCreateRequest request = new ClientCreateRequest(); - ClientCreate created = new ClientCreate(UUID.randomUUID()); - when(clientApi.createClient(any())).thenReturn(Mono.just(created)); + @Nested + @DisplayName("upsertClients") + class UpsertClientsTests { + + @Test + @DisplayName("no existing client — creates new client") + void noExistingClient_createsNewClient() { + String internalUserId = "internal-1"; + String externalUserId = "external-1"; + String legalEntityId = "le-1"; + UUID clientUuid = UUID.randomUUID(); + + ClientUser clientUser = ClientUser.builder() + .internalUserId(internalUserId) + .externalUserId(externalUserId) + .legalEntityId(legalEntityId) + .build(); + + stubEmptyClientList(internalUserId); + when(clientApi.createClient(any())).thenReturn(Mono.just(createdClient(clientUuid, internalUserId))); + + StepVerifier.create(service.upsertClients(List.of(clientUser))) + .expectNextMatches(results -> results.size() == 1 + && clientUuid.equals(results.getFirst().getInvestmentClientId()) + && internalUserId.equals(results.getFirst().getInternalUserId()) + && externalUserId.equals(results.getFirst().getExternalUserId()) + && legalEntityId.equals(results.getFirst().getLegalEntityId())) + .verifyComplete(); + + verify(clientApi).createClient(any()); + verify(clientApi, never()).patchClient(any(), any()); + } + + @Test + @DisplayName("duplicate internalUserId — deduplicates and processes once") + void duplicateInternalUserId_deduplicatesClients() { + String internalUserId = "duplicate-id"; + UUID clientUuid = UUID.randomUUID(); + + ClientUser first = ClientUser.builder() + .internalUserId(internalUserId) + .externalUserId("external-a") + .legalEntityId("le-a") + .build(); + ClientUser duplicate = ClientUser.builder() + .internalUserId(internalUserId) + .externalUserId("external-b") + .legalEntityId("le-b") + .build(); + + stubEmptyClientList(internalUserId); + when(clientApi.createClient(any())).thenReturn(Mono.just(createdClient(clientUuid, internalUserId))); + + StepVerifier.create(service.upsertClients(List.of(first, duplicate))) + .expectNextMatches(results -> results.size() == 1) + .verifyComplete(); + + verify(clientApi, times(1)).listClients(any(), any(), any(), any(), any(), + eq(internalUserId), any(), any(), any(), any(), any(), any()); + } + + @Test + @DisplayName("client upsert fails — skips client and returns empty list") + void clientUpsertFails_skipsClient() { + ClientUser clientUser = ClientUser.builder() + .internalUserId("failing-id") + .externalUserId("external-fail") + .legalEntityId("le-fail") + .build(); + + stubEmptyClientList("failing-id"); + when(clientApi.createClient(any())) + .thenReturn(Mono.error(new RuntimeException("create failed"))); + + StepVerifier.create(service.upsertClients(List.of(clientUser))) + .expectNextMatches(List::isEmpty) + .verifyComplete(); + } + } + + @Nested + @DisplayName("getClient") + class GetClientTests { + + @Test + @DisplayName("client found — returns client") + void clientFound_returnsClient() { + UUID uuid = UUID.randomUUID(); + OASClient client = new OASClient(); + when(clientApi.getClient(eq(uuid), anyList(), isNull(), isNull())).thenReturn(Mono.just(client)); + + StepVerifier.create(service.getClient(uuid)) + .expectNext(client) + .verifyComplete(); + } + + @Test + @DisplayName("client not found — returns empty Mono") + void clientNotFound_returnsEmpty() { + UUID uuid = UUID.randomUUID(); + WebClientResponseException notFound = new WebClientResponseException( + 404, "Not Found", new HttpHeaders(), new byte[0], StandardCharsets.UTF_8); + when(clientApi.getClient(eq(uuid), anyList(), isNull(), isNull())).thenReturn(Mono.error(notFound)); + + StepVerifier.create(service.getClient(uuid)) + .verifyComplete(); + } + + @Test + @DisplayName("non-404 error — propagates error") + void nonNotFoundError_propagatesError() { + UUID uuid = UUID.randomUUID(); + when(clientApi.getClient(eq(uuid), anyList(), isNull(), isNull())) + .thenReturn(Mono.error(new RuntimeException("service unavailable"))); + StepVerifier.create(service.getClient(uuid)) + .expectErrorMatches(e -> e instanceof RuntimeException + && "service unavailable".equals(e.getMessage())) + .verify(); + } } - @Test - void getClient_notFoundReturnsEmpty() { - UUID uuid = UUID.randomUUID(); - WebClientResponseException notFound = new WebClientResponseException( - 404, "Not Found", new HttpHeaders(), new byte[0], StandardCharsets.UTF_8); - when(clientApi.getClient(eq(uuid), anyList(), any(), any())).thenReturn(Mono.error(notFound)); + @Nested + @DisplayName("patchClient") + class PatchClientTests { - StepVerifier.create(service.getClient(uuid)) - .verifyComplete(); + @Test + @DisplayName("patch succeeds — returns updated client") + void patchClient_success() { + UUID uuid = UUID.randomUUID(); + PatchedOASClientUpdateRequest patch = new PatchedOASClientUpdateRequest(); + OASClient updated = new OASClient(); + when(clientApi.patchClient(eq(uuid), any())).thenReturn(Mono.just(updated)); + + StepVerifier.create(service.patchClient(uuid, patch)) + .expectNext(updated) + .verifyComplete(); + } + + @Test + @DisplayName("patchClient propagates WebClientResponseException errors") + void patchClient_webClientError_propagatesError() { + UUID uuid = UUID.randomUUID(); + PatchedOASClientUpdateRequest patch = new PatchedOASClientUpdateRequest(); + when(clientApi.patchClient(eq(uuid), any())) + .thenReturn(Mono.error(WebClientResponseException.create( + HttpStatus.BAD_REQUEST.value(), "Bad Request", + HttpHeaders.EMPTY, "invalid patch".getBytes(StandardCharsets.UTF_8), + StandardCharsets.UTF_8))); + + StepVerifier.create(service.patchClient(uuid, patch)) + .expectError(WebClientResponseException.class) + .verify(); + } } - @Test - void patchClient_success() { - UUID uuid = UUID.randomUUID(); - PatchedOASClientUpdateRequest patch = new PatchedOASClientUpdateRequest(); - OASClient updated = new OASClient(); - when(clientApi.patchClient(eq(uuid), any())).thenReturn(Mono.just(updated)); + @Nested + @DisplayName("updateClient") + class UpdateClientTests { + + @Test + @DisplayName("update succeeds — returns updated client") + void updateClient_success() { + UUID uuid = UUID.randomUUID(); + OASClientUpdateRequest update = new OASClientUpdateRequest(); + OASClient updated = new OASClient(); + when(clientApi.updateClient(eq(uuid), any())).thenReturn(Mono.just(updated)); - StepVerifier.create(service.patchClient(uuid, patch)) - .expectNext(updated) - .verifyComplete(); + StepVerifier.create(service.updateClient(uuid, update)) + .expectNext(updated) + .verifyComplete(); + } + + @Test + @DisplayName("update propagates WebClientResponseException errors") + void updateClient_webClientError_propagatesError() { + UUID uuid = UUID.randomUUID(); + OASClientUpdateRequest update = new OASClientUpdateRequest(); + when(clientApi.updateClient(eq(uuid), any())) + .thenReturn(Mono.error(WebClientResponseException.create( + HttpStatus.BAD_REQUEST.value(), "Bad Request", + HttpHeaders.EMPTY, "invalid update".getBytes(StandardCharsets.UTF_8), + StandardCharsets.UTF_8))); + + StepVerifier.create(service.updateClient(uuid, update)) + .expectError(WebClientResponseException.class) + .verify(); + } } - @Test - void updateClient_success() { - UUID uuid = UUID.randomUUID(); - OASClientUpdateRequest update = new OASClientUpdateRequest(); - OASClient updated = new OASClient(); - when(clientApi.updateClient(eq(uuid), any())).thenReturn(Mono.just(updated)); + private void stubEmptyClientList(String internalUserId) { + lenient().when(clientApi.listClients(any(), any(), any(), any(), any(), + eq(internalUserId), any(), any(), any(), any(), any(), any())) + .thenAnswer(invocation -> Mono.just(new PaginatedOASClientList().results(List.of()))); + } - StepVerifier.create(service.updateClient(uuid, update)) - .expectNext(updated) - .verifyComplete(); + private ClientCreate createdClient(UUID clientUuid, String internalUserId) { + return new ClientCreate(clientUuid).internalUserId(internalUserId); } } - diff --git a/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentPortfolioServiceTest.java b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentPortfolioServiceTest.java index 2276346a3..06a925af4 100644 --- a/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentPortfolioServiceTest.java +++ b/stream-investment/investment-core/src/test/java/com/backbase/stream/investment/service/InvestmentPortfolioServiceTest.java @@ -740,6 +740,80 @@ void upsertPortfolios_emptyArrangements_returnsEmptyList() { verifyNoInteractions(portfolioApi); } + + @Test + @DisplayName("single arrangement fails — skips failed arrangement and returns successful ones") + void upsertPortfolios_singleFailure_skipsFailedArrangement() { + UUID portfolioUuid = UUID.randomUUID(); + UUID productId = UUID.randomUUID(); + String successExternalId = "EXT-BATCH-OK"; + String failExternalId = "EXT-BATCH-FAIL"; + String leExternalId = "LE-BATCH"; + UUID clientUuid = UUID.randomUUID(); + + InvestmentArrangement successArrangement = buildArrangement( + successExternalId, "Success Portfolio", productId, leExternalId); + InvestmentArrangement failArrangement = buildArrangement( + failExternalId, "Fail Portfolio", productId, leExternalId); + + PaginatedPortfolioListList emptyList = Mockito.mock(PaginatedPortfolioListList.class); + when(emptyList.getResults()).thenReturn(List.of()); + when(portfolioApi.listPortfolios(isNull(), isNull(), isNull(), + isNull(), eq(successExternalId), isNull(), isNull(), eq(1), + isNull(), isNull(), isNull(), isNull())) + .thenReturn(Mono.just(emptyList)); + when(portfolioApi.listPortfolios(isNull(), isNull(), isNull(), + isNull(), eq(failExternalId), isNull(), isNull(), eq(1), + isNull(), isNull(), isNull(), isNull())) + .thenReturn(Mono.error(new RuntimeException("list portfolios failed"))); + + PortfolioList created = buildPortfolioList(portfolioUuid, successExternalId, + OffsetDateTime.now().minusMonths(6)); + when(portfolioApi.createPortfolio(any(), isNull(), isNull(), isNull())) + .thenReturn(Mono.just(created)); + + Map> clientsByLeExternalId = Map.of(leExternalId, List.of(clientUuid)); + + StepVerifier.create(service.upsertPortfolios( + List.of(successArrangement, failArrangement), clientsByLeExternalId)) + .expectNextMatches(list -> list.size() == 1 + && portfolioUuid.equals(list.getFirst().getPortfolio().getUuid())) + .verifyComplete(); + } + + @Test + @DisplayName("successful upsert — maps initial cash and withdrawal amount from arrangement") + void upsertPortfolios_success_mapsCashAndWithdrawalAmount() { + UUID portfolioUuid = UUID.randomUUID(); + UUID productId = UUID.randomUUID(); + String externalId = "EXT-BATCH-CASH"; + String leExternalId = "LE-CASH"; + UUID clientUuid = UUID.randomUUID(); + + InvestmentArrangement arrangement = buildArrangement( + externalId, "Cash Portfolio", productId, leExternalId); + when(arrangement.getInitialCash()).thenReturn(BigDecimal.valueOf(25_000)); + when(arrangement.getWithdrawalAmount()).thenReturn(BigDecimal.valueOf(1_500)); + + PaginatedPortfolioListList emptyList = Mockito.mock(PaginatedPortfolioListList.class); + when(emptyList.getResults()).thenReturn(List.of()); + when(portfolioApi.listPortfolios(isNull(), isNull(), isNull(), + isNull(), eq(externalId), isNull(), isNull(), eq(1), + isNull(), isNull(), isNull(), isNull())) + .thenReturn(Mono.just(emptyList)); + + PortfolioList created = buildPortfolioList(portfolioUuid, externalId, + OffsetDateTime.now().minusMonths(6)); + when(portfolioApi.createPortfolio(any(), isNull(), isNull(), isNull())) + .thenReturn(Mono.just(created)); + + StepVerifier.create(service.upsertPortfolios( + List.of(arrangement), Map.of(leExternalId, List.of(clientUuid)))) + .expectNextMatches(list -> list.size() == 1 + && BigDecimal.valueOf(25_000).equals(list.getFirst().getInitialCash()) + && BigDecimal.valueOf(1_500).equals(list.getFirst().getWithdrawalAmount())) + .verifyComplete(); + } } // =========================================================================