diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java index 86c89312f..383db9f8b 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/service/TradingClearingRegistryService.java @@ -2,6 +2,9 @@ package ru.spcex.clearing.account.service; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.header.Header; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; @@ -16,6 +19,7 @@ import ru.clearing.classes.statics.data.company.relation.Relation; import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; import ru.spcex.clearing.account.errors.AccountError; import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.ActionType; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; import ru.spcex.clearing.platform.messaging.domain.cud.common.CommonIdRequest; @@ -23,6 +27,8 @@ import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingR import ru.spcex.clearing.platform.messaging.domain.cud.registry.TradingClearingRegistryUpdateRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; +import ru.spcex.clearing.platform.messaging.service.sender.CorrelationHeader; import ru.spcex.clearing.util.security.UserRoleVerification; import ru.spcex.clearing.util.services.RequestHelper; import ru.spcex.clearing.validation.common.ValidationHelper; @@ -35,6 +41,7 @@ import ru.spcex.platform.imdg.api.ImdgProvider; import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; import ru.spcex.platform.utils.enumeration.IMessageResolver; +import ru.spcex.platform.utils.log.ExceptionUtils; import ru.spcex.platform.utils.validation.IValidator; import java.time.Instant; @@ -42,6 +49,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.List; import java.util.Map; +import java.util.concurrent.Future; import java.util.function.Function; @Service @@ -65,6 +73,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini private final Function tradingClearingRegistryBlockRequestValidator; private final IMessageResolver messageResolver; + private final Producer kafkaProducer; public TradingClearingRegistryService(Consumer kafkaQueue, Producer kafkaProducer, @@ -80,6 +89,7 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini @Qualifier("tradingClearingRegistryBlockRequest") Function tradingClearingRegistryBlockRequestValidator) { super(kafkaQueue, kafkaProducer); + this.kafkaProducer = kafkaProducer; this.validationHelper = validationHelper; this.userRoleVerification = userRoleVerification; this.requestHelper = requestHelper; @@ -265,10 +275,16 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini public RequestInfoUpdate tradingClearingRegistryNew(BaseRequest userRequest) { log.debug("TradingClearingRegistryNewRequest received"); RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest); - if (requestInfoUpdate != null) return requestInfoUpdate; + if (requestInfoUpdate != null) { + sendResponse(userRequest, requestInfoUpdate, null); + return requestInfoUpdate; + } requestInfoUpdate = validationHelper.validateTillFirstError(userRequest, tradingClearingRegistryNewRequestValidator); - if (requestInfoUpdate != null) return requestInfoUpdate; + if (requestInfoUpdate != null) { + sendResponse(userRequest, requestInfoUpdate, null); + return requestInfoUpdate; + } TradingClearingRegistryNewRequest req = userRequest.getRequestPayload(); @@ -313,11 +329,38 @@ public class TradingClearingRegistryService extends QueueConsumer implements Ini log.debug("successfully processed, id {}", id); - //todo this.DESTINATION_TRADING_CLEARING_REGISTRY_REPLY + sendResponse(userRequest, null, null); // to DESTINATION_TRADING_CLEARING_REGISTRY_REPLY return null; } + + private void sendResponse(BaseRequest o, Object response, Header correlationId) { + try { + Future send; + 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); + } + ProducerRecord respRec = new ProducerRecord<>(Consts.DESTINATION_TRADING_CLEARING_REGISTRY_REPLY, req); + if (correlationId != null) { + respRec.headers().add(new CorrelationHeader(correlationId.value().clone())); + } + send = kafkaProducer.send(respRec); + send.get(); + } catch (Exception e) { + log.error(ExceptionUtils.getStackTrace(e)); + } + } + public RequestInfoUpdate tradingClearingRegistryUpdate(BaseRequest userRequest) { log.debug("TradingClearingRegistryUpdateRequest received"); RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(userRequest);