Compare commits

...
Sign in to create a new pull request.

1 commit

Author SHA1 Message Date
ialbert
56c7573fd3 IKafkaRequestStatusProcessor implementation 2023-12-04 19:17:41 +03:00
2 changed files with 65 additions and 0 deletions

View file

@ -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<RequestInfo> 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);
}
}

View file

@ -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) {}
}