From 5f2abbd2d9b30d2fa56434d2a8723d1930406f6b Mon Sep 17 00:00:00 2001 From: ialbert Date: Tue, 11 Oct 2022 20:18:06 +0300 Subject: [PATCH] http://jira.mfd.msk:8088/browse/CLS-151 --- .../clearing/account/config/KafkaConfig.java | 3 +- .../queue/QueueExceptionHandler.java | 4 +- .../response/system/RequestInfoResponse.java | 51 +++++++++++++ .../system/RequestStatusController.java | 48 ++++++++++++ .../backendapi/service/impl/OperatorImpl.java | 3 +- .../balance/config/KafkaSenderConfig.java | 3 +- .../dbf/importer/config/KafkaConfig.java | 3 +- .../clearing/imdg/IMDGDistributedNames.java | 1 + .../platform/messaging/domain/Consts.java | 2 + .../domain/json/serialize/EnumSerializer.java | 19 +++++ .../logic/functional/BuilderConsumerStep.java | 2 + .../functional/ConsumerSpecificClass.java | 26 ++++++- .../messaging/service/QueueConsumer.java | 75 ++++++++++++++++--- .../messaging/service/RequestInfoUpdate.java | 22 ++++++ 14 files changed, 243 insertions(+), 19 deletions(-) create mode 100644 clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/response/system/RequestInfoResponse.java create mode 100644 clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/system/RequestStatusController.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/serialize/EnumSerializer.java create mode 100644 platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java diff --git a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java index f1cac72b8..58a707200 100644 --- a/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java +++ b/clearing-parent/account-service/src/main/java/ru/spcex/clearing/account/config/KafkaConfig.java @@ -8,6 +8,7 @@ import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Scope; import ru.spcex.clearing.account.config.settings.AccountServiceSettings; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; import ru.spcex.clearing.platform.messaging.service.RequestInfo; @@ -41,7 +42,7 @@ public class KafkaConfig { .producer(kafkaProducer) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); return imdg::insert; }) .build(); diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/QueueExceptionHandler.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/QueueExceptionHandler.java index 5b6503559..70dafb130 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/QueueExceptionHandler.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/queue/QueueExceptionHandler.java @@ -39,7 +39,7 @@ public class QueueExceptionHandler extends ResponseEntityExceptionHandler { @ResponseStatus(HttpStatus.BAD_REQUEST) @ResponseBody public BasicSpcexResponse handleValidationException(ActionValidationException e) { - log.error("Default scoring exception handler: {}", ExceptionUtils.getStackTrace(e)); + log.error("Default exception handler: {}", ExceptionUtils.getStackTrace(e)); EnumMessage firstError = e.getErrors().iterator().next(); BasicSpcexResponse errorResponse = new BasicSpcexResponse(); errorResponse.setCode(firstError.getSubject().getId()); @@ -68,7 +68,7 @@ public class QueueExceptionHandler extends ResponseEntityExceptionHandler { @ResponseStatus(HttpStatus.INTERNAL_SERVER_ERROR) @ResponseBody public BasicSpcexResponse handleDefault(Throwable e) { - log.error("Default scoring exception handler: {}", ExceptionUtils.getStackTrace(e)); + log.error("Default exception handler: {}", ExceptionUtils.getStackTrace(e)); BasicSpcexResponse errorResponse = new BasicSpcexResponse(); errorResponse.setCode(-1); errorResponse.setMessage("internal server error"); diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/response/system/RequestInfoResponse.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/response/system/RequestInfoResponse.java new file mode 100644 index 000000000..ed21ea44d --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/response/system/RequestInfoResponse.java @@ -0,0 +1,51 @@ +package ru.spcex.clearing.backendapi.controller.response.system; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fasterxml.jackson.databind.annotation.JsonSerialize; +import io.swagger.annotations.ApiModel; +import io.swagger.annotations.ApiModelProperty; +import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse; +import ru.spcex.clearing.platform.messaging.domain.json.serialize.EnumSerializer; +import ru.spcex.clearing.platform.messaging.service.Status; + + +/** + * Банковские реквизиты для перечисления денежных средств + * @ApiModel(description = "Ответ в результате отправки операции в топик Kafka.") + * @ApiModelProperty(value = "Информация о принятом запросе") + **/ +@ApiModel(description = "Ответ при получении статуса запроса.") +public class RequestInfoResponse extends BasicSpcexResponse { + + @JsonProperty + @ApiModelProperty(value = "Поля объекта") + private RequestInfoPayload payload; + + public RequestInfoPayload getPayload() { + return payload; + } + + public void setStatus(Status status) { + this.payload = new RequestInfoPayload(); + payload.setStatus(status); + } + + public void setPayload(RequestInfoPayload payload) { + this.payload = payload; + } + + private static class RequestInfoPayload { + @JsonProperty + @JsonSerialize(using = EnumSerializer.class) + private Status status; + + public Status getStatus() { + return status; + } + + public void setStatus(Status status) { + this.status = status; + } + } + +} \ No newline at end of file diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/system/RequestStatusController.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/system/RequestStatusController.java new file mode 100644 index 000000000..cebfc2d57 --- /dev/null +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/controller/system/RequestStatusController.java @@ -0,0 +1,48 @@ +package ru.spcex.clearing.backendapi.controller.system; + +import io.swagger.annotations.ApiOperation; +import io.swagger.annotations.ApiParam; +import io.swagger.annotations.ApiResponse; +import io.swagger.annotations.ApiResponses; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Controller; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RequestMethod; +import org.springframework.web.bind.annotation.ResponseBody; +import ru.spcex.clearing.backendapi.controller.response.BasicSpcexResponse; +import ru.spcex.clearing.backendapi.controller.response.system.RequestInfoResponse; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; + +@Controller +@RequestMapping("/requests") +public class RequestStatusController { + + private final Imdg imdg; + + @Autowired + public RequestStatusController(ImdgProvider imdgProvider) { + this.imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + } + + @ApiOperation(value = "get request status.") + @ApiResponses(value = {@ApiResponse(code = 200, message = "OK", response = RequestInfoResponse.class), @ApiResponse(code = 400, message = "Ошибка валидации", response = BasicSpcexResponse.class)}) + @RequestMapping(value = "/{id}", method = RequestMethod.GET) + @ResponseBody + public RequestInfoResponse getById(@ApiParam(value = "Идентификатор объекта", required = true, example = "1234") + @PathVariable("id") Long id) { + RequestInfo singleObjectByID; + if (id == null || (singleObjectByID = imdg.getSingleObjectByID(id)) == null) { + RequestInfoResponse response = new RequestInfoResponse(); + response.setCode(-1L); + response.setMessage("cannot find request with id='" + id + "'."); + return response; + } + RequestInfoResponse response = new RequestInfoResponse(); + response.setStatus(singleObjectByID.getStatus()); + return response; + } +} diff --git a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/impl/OperatorImpl.java b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/impl/OperatorImpl.java index dabfee7e6..216b25f58 100644 --- a/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/impl/OperatorImpl.java +++ b/clearing-parent/backend-api/src/main/java/ru/spcex/clearing/backendapi/service/impl/OperatorImpl.java @@ -8,6 +8,7 @@ import ru.spcex.clearing.backendapi.controller.response.cud.QueueSuccessResponse import ru.spcex.clearing.backendapi.domain.actions.IAction; import ru.spcex.clearing.backendapi.errors.ActionValidationException; import ru.spcex.clearing.backendapi.service.IOperator; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.platform.imdg.api.Imdg; @@ -46,7 +47,7 @@ public class OperatorImpl implements IOperator { } private void saveRequestToStorage(String destination, BaseRequest request) { - Imdg requestStorage = imdgProvider.getImdg(destination, RequestInfo.class); + Imdg requestStorage = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); RequestInfo requestInfo = RequestInfo.create(request.getId()); requestStorage.insert(requestInfo); } diff --git a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java index fa621992c..a32e3bc5e 100644 --- a/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java +++ b/clearing-parent/balance-service/src/main/java/ru/spcex/clearing/balance/config/KafkaSenderConfig.java @@ -4,6 +4,7 @@ import org.apache.kafka.clients.producer.Producer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; import ru.spcex.platform.imdg.api.Imdg; @@ -22,7 +23,7 @@ public class KafkaSenderConfig { .producer(kafkaProducer) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); return imdg::insert; }) .build(); diff --git a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java index 86358d7c1..bf533336b 100644 --- a/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java +++ b/clearing-parent/dbf-importer/src/main/java/ru/spcex/clearing/dbf/importer/config/KafkaConfig.java @@ -5,6 +5,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import ru.spcex.clearing.dbf.importer.config.settings.ImportDBFServiceSettings; +import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; @@ -28,7 +29,7 @@ public class KafkaConfig { .producer(kafkaProducer) .idGenerator(imdgIdGenerator::nextId) .imdgProvider(s -> { - Imdg imdg = imdgProvider.getImdg(s, RequestInfo.class); + Imdg imdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); return imdg::insert; }) .build(); diff --git a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java index b3c3b9d8d..d0c03210f 100644 --- a/platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java +++ b/platform-parent/platform-imdg-api/src/main/java/ru/spcex/clearing/imdg/IMDGDistributedNames.java @@ -101,6 +101,7 @@ public final class IMDGDistributedNames { public static final String Map_AccountStatusDictionary = "Map_AccountStatusDictionary"; public static final String Map_KeyRate = "Map_KeyRate"; public static final String Map_MoneyMarketSecurity = "Map_MoneyMarketSecurity"; + public static final String Map_RequestInfo = "Map_RequestInfo"; public static final String MAP_SEQUENCE_NAME = "MAP_SEQUENCE_NAME"; // todo вынести idGenerator отдельно и завернуть в метод, чтобы не напрямую обращаться. private IMDGDistributedNames() { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java index 212900ad1..3eca7e9a6 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/Consts.java @@ -33,4 +33,6 @@ public interface Consts { String ACCOUNT_NEW = "account-new"; String BALANCE_ACCOUNT_NEW = "balance-account-new"; + String REQUEST_INFO_UPDATE = "request-info-update"; + } diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/serialize/EnumSerializer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/serialize/EnumSerializer.java new file mode 100644 index 000000000..6cf93cc26 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/json/serialize/EnumSerializer.java @@ -0,0 +1,19 @@ +package ru.spcex.clearing.platform.messaging.domain.json.serialize; + +import com.fasterxml.jackson.core.JsonGenerator; +import com.fasterxml.jackson.databind.JsonSerializer; +import com.fasterxml.jackson.databind.SerializerProvider; + +import java.io.IOException; +import java.time.format.DateTimeFormatter; + +public class EnumSerializer extends JsonSerializer> { + private static final DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd"); + + + @Override + public void serialize(Enum value, JsonGenerator gen, SerializerProvider serializers) throws IOException { + if (value == null) return; + gen.writeString(value.name()); + } +} 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 88d687f2c..224b7b1c7 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 @@ -3,7 +3,9 @@ package ru.spcex.clearing.platform.messaging.logic.functional; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import java.util.function.Consumer; +import java.util.function.Function; public interface BuilderConsumerStep { BuilderDestinationStep setConsumer(Consumer> consumer); + BuilderDestinationStep setFunction(Function, Object> 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 597cf8b04..0b6e96404 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 @@ -4,6 +4,7 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import java.util.function.BiConsumer; import java.util.function.Consumer; +import java.util.function.Function; /** * вспомогательный класс @@ -18,6 +19,7 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder * действие-обработчик реквеста (метод из конкретного менеджера) */ private java.util.function.Consumer> consumer; + private java.util.function.Function, Object> function; private ConsumerSpecificClass(Class clazz) { this.clazz = clazz; @@ -29,7 +31,17 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder @Override public BuilderDestinationStep setConsumer(Consumer> consumer) { - this.consumer = consumer; + if (this.consumer == null) { + this.consumer = consumer; + } else { + this.consumer = this.consumer.andThen(consumer); + } + return this; + } + + @Override + public BuilderDestinationStep setFunction(Function, Object> function) { + this.function = function; return this; } @@ -38,10 +50,18 @@ public class ConsumerSpecificClass implements BuilderConsumerStep, Builder } @SuppressWarnings("unchecked") - public void acceptRaw(Object obj) { - this.consumer.accept((BaseRequest) obj); + public Object acceptRaw(Object obj) { + if (this.consumer != null) { + this.consumer.accept((BaseRequest) obj); + return null; + } else { + return this.function.apply((BaseRequest) obj); + } } + 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/service/QueueConsumer.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/QueueConsumer.java index 7a1f96915..da27d7857 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 @@ -5,10 +5,14 @@ import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; +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.errors.WakeupException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; 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.platform.utils.log.ExceptionUtils; @@ -19,6 +23,7 @@ import java.util.HashMap; import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicBoolean; /** @@ -28,33 +33,40 @@ public class QueueConsumer implements AutoCloseable { private final Logger log = LoggerFactory.getLogger(getClass()); private final AtomicBoolean closed = new AtomicBoolean(false); private final Consumer consumer; - private final ExecutorService executor; + private Producer producer; + private final ExecutorService inputExecutor; + private final ExecutorService outputExecutor; protected final Map> callbacks; private final ObjectMapper json; + protected boolean supportStartOffsetTimeWindow; public QueueConsumer(Consumer kafkaQueue) { this.consumer = kafkaQueue; this.callbacks = new HashMap<>(); - this.executor = Executors.newSingleThreadExecutor(); + this.inputExecutor = Executors.newSingleThreadExecutor(); + this.outputExecutor = Executors.newSingleThreadExecutor(); this.json = new ObjectMapper(); + this.supportStartOffsetTimeWindow = false; + } + + public QueueConsumer(Consumer kafkaQueue, Producer kafkaResponseQueue) { + this(kafkaQueue); + this.producer = kafkaResponseQueue; } protected boolean needsProcessing(String topicName, BaseRequest request) { return true; } - protected boolean supportStartOffsetTimeWindow() { - return false; - } - public void init() { - executor.submit(() -> { + inputExecutor.submit(() -> { try { - if (supportStartOffsetTimeWindow()) { + if (supportStartOffsetTimeWindow) { consumer.subscribe(callbacks.keySet(), new OffsetChanger(consumer, callbacks.keySet())); } else { consumer.subscribe(callbacks.keySet()); } + Object o = null; while (!closed.get()) { try { ConsumerRecords records = consumer.poll(Duration.of(10, ChronoUnit.SECONDS)); @@ -62,13 +74,19 @@ public class QueueConsumer implements AutoCloseable { ConsumerSpecificClass callback = callbacks.get(next.topic()); Class clazz = callback.getClazz(); JavaType payloadType = json.getTypeFactory().constructParametricType(BaseRequest.class, clazz); - Object o = json.readValue((String) next.value(), payloadType); + o = json.readValue((String) next.value(), payloadType); if (needsProcessing(next.topic(), (BaseRequest) o)) { - callback.acceptRaw(o); + Object topicResponse = callback.acceptRaw(o); + if (producer != null) { + sendResponse((BaseRequest) o, topicResponse); + } } } } catch (Throwable e) { log.error(ExceptionUtils.getStackTrace(e)); + if (producer != null) { + sendErrorResponse((BaseRequest) o); + } } } } catch (WakeupException e) { @@ -81,6 +99,43 @@ public class QueueConsumer implements AutoCloseable { }); } + //мб перенести в другой класс + private void sendErrorResponse(BaseRequest o) { + outputExecutor.submit(() -> { + try { + Future send; + //default response + RequestInfoUpdate success = new RequestInfoUpdate(); + success.setId(o.getId()); + success.setStatus(Status.Error); + send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, success)); + send.get(); + } catch (Exception e) { + log.error(ExceptionUtils.getStackTrace(e)); + } + }); + } + + private void sendResponse(BaseRequest o, Object response) { + outputExecutor.submit(() -> { + try { + Future send; + if (response != null) { + send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, response)); + } else { + //default response + RequestInfoUpdate success = new RequestInfoUpdate(); + success.setId(o.getId()); + success.setStatus(Status.Success); + send = producer.send(new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, success)); + } + send.get(); + } catch (Exception e) { + log.error(ExceptionUtils.getStackTrace(e)); + } + }); + } + protected BuilderConsumerStep callback(Class clazz) { return ConsumerSpecificClass.build(clazz); } 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 new file mode 100644 index 000000000..d83b8e9c1 --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/RequestInfoUpdate.java @@ -0,0 +1,22 @@ +package ru.spcex.clearing.platform.messaging.service; + +public class RequestInfoUpdate { + private Long id; + private Status status; + + public Long getId() { + return id; + } + + public void setId(Long id) { + this.id = id; + } + + public Status getStatus() { + return status; + } + + public void setStatus(Status status) { + this.status = status; + } +}