diff --git a/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/messaging/KafkaRequestStatusProcessorImpl.java b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/messaging/KafkaRequestStatusProcessorImpl.java new file mode 100644 index 000000000..4f772a455 --- /dev/null +++ b/clearing-parent/clearing-validation/src/main/java/ru/spcex/clearing/messaging/KafkaRequestStatusProcessorImpl.java @@ -0,0 +1,57 @@ +package ru.spcex.clearing.messaging; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.service.IKafkaRequestStatusProcessor; +import ru.spcex.clearing.platform.messaging.service.RequestInfo; +import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate; +import ru.spcex.clearing.platform.messaging.service.Status; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.utils.log.ExceptionUtils; + +public class KafkaRequestStatusProcessorImpl implements IKafkaRequestStatusProcessor { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final Imdg rqstInfImdg; + + public KafkaRequestStatusProcessorImpl(ImdgProvider imdgProvider) { + this.rqstInfImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_RequestInfo, RequestInfo.class); + } + + @Override + public void updateRequestStatus(BaseRequest o, Object response) { + try { + RequestInfoUpdate updateEvent; + if (response == null) { + //default response + updateEvent = new RequestInfoUpdate(); + updateEvent.setId(o.getId()); + updateEvent.setStatus(Status.Success); + } else { + updateEvent = (RequestInfoUpdate) response; + } + RequestInfo rqstInfo = rqstInfImdg.getSingleObjectByID(o.getId()); + if (rqstInfo == null) { + log.trace("cannot find requestInfo.id={} for update '{}'", o.getId(), updateEvent.getStatus()); + return; + } + rqstInfo.setStatus(updateEvent.getStatus()); + rqstInfo.setMessage(updateEvent.getMessage()); + rqstInfImdg.update(rqstInfo); + } catch (Exception e) { + log.error(ExceptionUtils.getStackTrace(e)); + } + } + + @Override + public void updateRequestStatusOnError(BaseRequest o, Object response) { + RequestInfoUpdate updateEvent; + //default response + updateEvent = new RequestInfoUpdate(); + updateEvent.setId(o.getId()); + updateEvent.setStatus(Status.Error); + updateRequestStatus(o, updateEvent); + } +} diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/IKafkaRequestStatusProcessor.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/IKafkaRequestStatusProcessor.java new file mode 100644 index 000000000..e3547a95b --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/service/IKafkaRequestStatusProcessor.java @@ -0,0 +1,8 @@ +package ru.spcex.clearing.platform.messaging.service; + +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; + +public interface IKafkaRequestStatusProcessor { + default void updateRequestStatus(BaseRequest o, Object response) {} + default void updateRequestStatusOnError(BaseRequest o, Object response) {} +}