From ccef108e13adfcc1a25374e7e2afb9df4e8733af Mon Sep 17 00:00:00 2001 From: etreschenkov Date: Thu, 27 Jun 2024 18:05:35 +0300 Subject: [PATCH] account-service-v2 --- .../listeners/ClientCodeMessageListener.java | 21 ++++--- .../service/RequestInfoAccepter.java | 7 ++- .../logic/functional/BuilderConsumerStep.java | 3 +- .../functional/ConsumerSpecificClass.java | 17 ++++-- .../logic/response/DefaultResponses.java | 30 ++++++++++ .../messaging/logic/response/PayloadInfo.java | 15 +++++ .../messaging/service/QueueConsumer.java | 59 +++++-------------- .../messaging/service/RequestInfoUpdate.java | 9 ++- 8 files changed, 97 insertions(+), 64 deletions(-) create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/DefaultResponses.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/PayloadInfo.java diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java index ee8a81f81..271a0b0a8 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/v2/listeners/ClientCodeMessageListener.java @@ -10,9 +10,11 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Service; +import org.springframework.util.Assert; import ru.clearing.classes.statics.data.account.Account; import ru.clearing.classes.statics.data.company.Company; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.account.model.ValidationResult; import ru.spcex.clearing.account.service.v2.GatewayRequestCreator; import ru.spcex.clearing.account.service.v2.facade.ClientCodeFacade; @@ -26,6 +28,8 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.TkrAccount; import ru.spcex.clearing.platform.messaging.domain.cud.account.TkrAccountsGatewayRequest; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SendTkrRequest; import ru.spcex.clearing.platform.messaging.domain.cud.gateway.Tkr; +import ru.spcex.clearing.platform.messaging.logic.response.DefaultResponses; +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; import ru.spcex.clearing.platform.messaging.service.Status; @@ -74,15 +78,15 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali .forDestination(Consts.DESTINATION_CLIENT_CODE_NEW, callbacks::put); //from gateway requests callback(TkrAccountsGatewayRequest.class) - .setFunction(this::clientCodeNewFromGateway) + .setConsumer(this::clientCodeNewFromGateway) .forDestination(Consts.DESTINATION_CLIENT_CODE_NEW_FROM_GATEWAY, callbacks::put); callback(TkrAccountsGatewayRequest.class) - .setFunction(this::clientCodeNewFromGateway) + .setConsumer(this::clientCodeNewFromGateway) .forDestination(Consts.DESTINATION_CLIENT_CODE_UPDATE_FROM_GATEWAY, callbacks::put); init(); } - private RequestInfoUpdate clientCodeNewFromGateway(BaseRequest tkrRequest) { + private void clientCodeNewFromGateway(BaseRequest tkrRequest) { clientCodeFacade.lock(); //валидация запроса TkrAccountsGatewayRequest gatewayRequest = tkrRequest.getRequestPayload(); @@ -113,7 +117,6 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali .toList(); sendTkrRequest.setTkrs(tkrs); kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); - return null; } //unwrap запроса @@ -132,7 +135,7 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali clientCodeNewRequest.setCompanyId(company.getId()); clientCodeNewRequest.setCode(tkrAccount.getClientCode()); clientCodeNewRequest.setDepoAccountId(depoAccountId); - Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB); + Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB.getKey()); clientCodeNewRequest.setMoneyAccountId(moneyAccount.getId()); List foreignCurrencyList = accountByCurrency.entrySet() @@ -149,20 +152,20 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali log.debug("Success created new clientCode by tkr account: {}", tkrAccount.getClientCode()); } kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest); - return null; } - private RequestInfoUpdate clientCodeNewFromBackend(BaseRequest clientCodeNewRequest) { + private PayloadInfo clientCodeNewFromBackend(BaseRequest clientCodeNewRequest) { RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(clientCodeNewRequest); if (requestInfoUpdate != null) return requestInfoUpdate; ClientCodeNewRequest request = clientCodeNewRequest.getRequestPayload(); ValidationResult validationResult = clientCodeValidator.checkBackendNewRequest(request); if (!validationResult.isValid()) { - return makeErrorResponse(clientCodeNewRequest, validationResult); + return DefaultResponses.build(clientCodeNewRequest.getId(), + validationResult.errorMsg().orElse(""), DefaultResponses.ERROR); } clientCodeFacade.createClientCode(request, validationResult.validator()); - return null; + return DefaultResponses.build(clientCodeNewRequest.getId(), DefaultResponses.SUCCESS); } private RequestInfoUpdate makeErrorResponse(BaseRequest clientCodeNewRequest, diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java index b5a795465..4cb222aeb 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/RequestInfoAccepter.java @@ -8,6 +8,7 @@ import org.springframework.stereotype.Service; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; @@ -32,15 +33,15 @@ public class RequestInfoAccepter extends QueueConsumer implements InitializingBe @Override public void afterPropertiesSet() throws Exception { - callback(RequestInfoUpdate.class) + callback(PayloadInfo.class) .setConsumer(this::updateRequestInfo) .forDestination(Consts.REQUEST_INFO_UPDATE, callbacks::put); init(); } - private void updateRequestInfo(BaseRequest requestInfoUpdateBaseRequest) { + private void updateRequestInfo(BaseRequest requestInfoUpdateBaseRequest) { //will throw exception for any class other than RequestInfoUpdate - RequestInfoUpdate statusInfo = requestInfoUpdateBaseRequest.getRequestPayload(); + PayloadInfo statusInfo = requestInfoUpdateBaseRequest.getRequestPayload(); if (statusInfo == null || statusInfo.getId() == null) { log.warn("BaseRequest id={}, RequestPayload={}: requestInfo id was null", requestInfoUpdateBaseRequest.getId(), statusInfo); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java index 0a7bb7f01..ef6ab5abf 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/BuilderConsumerStep.java @@ -1,6 +1,7 @@ package ru.spcex.clearing.platform.messaging.logic.functional; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; import java.util.function.Consumer; import java.util.function.Function; @@ -19,5 +20,5 @@ public interface BuilderConsumerStep { /** * с обновлением RequestInfo в мапе */ - BuilderDestinationStep setFunction(Function, Object> function); + BuilderDestinationStep setFunction(Function, PayloadInfo> function); } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java index 0b6e96404..cb31c53f1 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/functional/ConsumerSpecificClass.java @@ -1,7 +1,10 @@ package ru.spcex.clearing.platform.messaging.logic.functional; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.logic.response.DefaultResponses; +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; +import java.util.Optional; import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.function.Function; @@ -19,7 +22,7 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder * действие-обработчик реквеста (метод из конкретного менеджера) */ private java.util.function.Consumer> consumer; - private java.util.function.Function, Object> function; + private java.util.function.Function, PayloadInfo> function; private ConsumerSpecificClass(Class clazz) { this.clazz = clazz; @@ -40,7 +43,7 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder } @Override - public BuilderDestinationStep setFunction(Function, Object> function) { + public BuilderDestinationStep setFunction(Function, PayloadInfo> function) { this.function = function; return this; } @@ -50,18 +53,22 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder } @SuppressWarnings("unchecked") - public Object acceptRaw(Object obj) { + public PayloadInfo acceptRaw(Object obj) { if (this.consumer != null) { this.consumer.accept((BaseRequest) obj); - return null; + return DefaultResponses.build(((BaseRequest) obj).getId(), DefaultResponses.SUCCESS); } else { - return this.function.apply((BaseRequest) obj); + PayloadInfo res = this.function.apply((BaseRequest) obj); + return res == null ? + DefaultResponses.build(((BaseRequest) obj).getId(), DefaultResponses.SUCCESS) : + res; } } public Consumer> getConsumer() { return consumer; } + public Class getClazz() { return clazz; } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/DefaultResponses.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/DefaultResponses.java new file mode 100644 index 000000000..6f149573b --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/DefaultResponses.java @@ -0,0 +1,30 @@ +package ru.spcex.clearing.platform.messaging.logic.response; + +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; + +public enum DefaultResponses { + SUCCESS, + ERROR, + ; + public static PayloadInfo build(Long id, String message, DefaultResponses defaultResponse) { + PayloadInfo info = build(id, defaultResponse); + info.setMessage(message); + return info; + } + public static PayloadInfo build(Long id, DefaultResponses defaultResponse) { + PayloadInfo info = new RequestInfoUpdate(); + info.setId(id); + switch (defaultResponse) { + case SUCCESS: { + info.setStatus(Status.Success); + break; + } + case ERROR: { + info.setStatus(Status.Error); + break; + } + } + return info; + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/PayloadInfo.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/PayloadInfo.java new file mode 100644 index 000000000..47f0db736 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/logic/response/PayloadInfo.java @@ -0,0 +1,15 @@ +package ru.spcex.clearing.platform.messaging.logic.response; + + +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; + +public interface PayloadInfo { + Long getId(); + Status getStatus(); + String getMessage(); + + PayloadInfo setId(Long id); + PayloadInfo setStatus(Status status); + RequestInfoUpdate setMessage(String message); +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index bfe4c6cb4..78ff7b536 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java @@ -18,6 +18,8 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.logic.functional.BuilderConsumerStep; import ru.spcex.clearing.platform.messaging.logic.functional.ConsumerSpecificClass; +import ru.spcex.clearing.platform.messaging.logic.response.DefaultResponses; +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; import ru.spcex.clearing.platform.messaging.service.sender.CorrelationHeader; import ru.spcex.platform.utils.enumeration.EnumMessage; import ru.spcex.platform.utils.enumeration.IMessageResolver; @@ -108,7 +110,7 @@ public class QueueConsumer implements AutoCloseable { ((BaseRequest) o).setCorrelationId(correlationId.value()); } if (needsProcessing(next.topic(), (BaseRequest) o)) { - Object topicResponse = callback.acceptRaw(o); + PayloadInfo topicResponse = callback.acceptRaw(o); if (producer != null) { sendResponse((BaseRequest) o, topicResponse, correlationId); } @@ -117,12 +119,13 @@ public class QueueConsumer implements AutoCloseable { lastErrors = 0; } catch (Throwable e) { log.error("Listener {} last topic \"{}\", error: {}", - QueueConsumer.this.getClass().getName(), lastTopic, ExceptionUtils.getStackTrace(e)); + QueueConsumer.this.getClass().getName(), lastTopic, ExceptionUtils.getStackTrace(e)); if (lastTopic == null) { log.debug("main thread stacktrace: {}", ExceptionUtils.getStackTrace(debugCallerStacktrace)); } if (producer != null && o != null) { - sendErrorResponse((BaseRequest) o, correlationId); + BaseRequest br = (BaseRequest) o; + sendResponse(br, DefaultResponses.build(br.getId(), DefaultResponses.ERROR), correlationId); } if (lastErrors++ > 20) { log.warn("Too many error at row, {}. Sleep.", lastErrors); @@ -139,57 +142,23 @@ public class QueueConsumer implements AutoCloseable { if (!closed.get()) throw e; } catch (Throwable e) { log.error("QueueConsumer error: {}; main thread stacktrace: {}", - ExceptionUtils.getStackTrace(e), ExceptionUtils.getStackTrace(debugCallerStacktrace)); + ExceptionUtils.getStackTrace(e), ExceptionUtils.getStackTrace(debugCallerStacktrace)); } finally { consumer.close(); } }); } - //мб перенести в другой класс - private void sendErrorResponse(BaseRequest o, Header correlationId) { + private void sendResponse(BaseRequest o, PayloadInfo response, Header correlationId) { outputExecutor.submit(() -> { try { Future send; - //default response - BaseRequest req = new BaseRequest<>(); - RequestInfoUpdate success = new RequestInfoUpdate(); - success.setId(o.getId()); - success.setStatus(Status.Error); - req.setRequestPayload(success); - req.setId(o.getId()); - req.setActionType(ActionType.SYSTEM); - ProducerRecord respRec = new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req); - if (correlationId != null) { - respRec.headers().add(KafkaHeaders.CORRELATION_ID, correlationId.value()); - } - send = producer.send(respRec); - send.get(); - } catch (Exception e) { - if (e instanceof InterruptedException) { - Thread.currentThread().interrupt(); - } - log.error(ExceptionUtils.getStackTrace(e)); - } - }); - } - protected void sendResponse(BaseRequest o, Object response, Header correlationId) { - outputExecutor.submit(() -> { - try { - Future send; - BaseRequest req = new BaseRequest<>(); + BaseRequest req = new BaseRequest<>(); req.setId(o.getId()); req.setActionType(ActionType.SYSTEM); - if (response != null) { - req.setRequestPayload(response); - } else { - //default response - RequestInfoUpdate success = new RequestInfoUpdate(); - success.setId(o.getId()); - success.setStatus(Status.Success); - req.setRequestPayload(success); - } + req.setRequestPayload(response); + ProducerRecord respRec = new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req); if (correlationId != null) { respRec.headers().add(new CorrelationHeader(correlationId.value().clone())); @@ -225,9 +194,9 @@ public class QueueConsumer implements AutoCloseable { String errorMsg = messageResolver.resolve(validationError.get()); log.warn("validation error for {} error={}, id={}: {}", req.getClass().getSimpleName(), validationError.get().getSubject(), userRequest.getId(), errorMsg); return new RequestInfoUpdate() - .setId(userRequest.getId()) - .setStatus(Status.Error) - .setMessage(errorMsg); + .setId(userRequest.getId()) + .setStatus(Status.Error) + .setMessage(errorMsg); } } return null; diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java index 77444032c..f198eb2e8 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java @@ -1,6 +1,9 @@ package ru.spcex.clearing.platform.messaging.service; -public class RequestInfoUpdate { +import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo; +import ru.spcex.platform.classes.base.interfaces.WithId; + +public class RequestInfoUpdate implements PayloadInfo { private Long id; private Status status; private String message; @@ -13,10 +16,12 @@ public class RequestInfoUpdate { this.message = message; } + @Override public Long getId() { return id; } + @Override public RequestInfoUpdate setId(Long id) { this.id = id; return this; @@ -26,6 +31,7 @@ public class RequestInfoUpdate { return status; } + @Override public RequestInfoUpdate setStatus(Status status) { this.status = status; return this; @@ -35,6 +41,7 @@ public class RequestInfoUpdate { return message; } + @Override public RequestInfoUpdate setMessage(String message) { this.message = message; return this;