Compare commits
1 commit
10_01_2024
...
kafka-requ
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
56c7573fd3 |
2 changed files with 65 additions and 0 deletions
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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) {}
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue