account-service-v2
This commit is contained in:
parent
cc95a67759
commit
ccef108e13
8 changed files with 97 additions and 64 deletions
|
|
@ -10,9 +10,11 @@ import org.slf4j.Logger;
|
|||
import org.slf4j.LoggerFactory;
|
||||
import org.springframework.beans.factory.InitializingBean;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.util.Assert;
|
||||
import ru.clearing.classes.statics.data.account.Account;
|
||||
import ru.clearing.classes.statics.data.company.Company;
|
||||
import ru.clearing.classes.statics.data.registry.TradingClearingRegistry;
|
||||
import ru.spcex.clearing.account.errors.AccountError;
|
||||
import ru.spcex.clearing.account.model.ValidationResult;
|
||||
import ru.spcex.clearing.account.service.v2.GatewayRequestCreator;
|
||||
import ru.spcex.clearing.account.service.v2.facade.ClientCodeFacade;
|
||||
|
|
@ -26,6 +28,8 @@ import ru.spcex.clearing.platform.messaging.domain.cud.account.TkrAccount;
|
|||
import ru.spcex.clearing.platform.messaging.domain.cud.account.TkrAccountsGatewayRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SendTkrRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.Tkr;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.DefaultResponses;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||
|
|
@ -74,15 +78,15 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali
|
|||
.forDestination(Consts.DESTINATION_CLIENT_CODE_NEW, callbacks::put);
|
||||
//from gateway requests
|
||||
callback(TkrAccountsGatewayRequest.class)
|
||||
.setFunction(this::clientCodeNewFromGateway)
|
||||
.setConsumer(this::clientCodeNewFromGateway)
|
||||
.forDestination(Consts.DESTINATION_CLIENT_CODE_NEW_FROM_GATEWAY, callbacks::put);
|
||||
callback(TkrAccountsGatewayRequest.class)
|
||||
.setFunction(this::clientCodeNewFromGateway)
|
||||
.setConsumer(this::clientCodeNewFromGateway)
|
||||
.forDestination(Consts.DESTINATION_CLIENT_CODE_UPDATE_FROM_GATEWAY, callbacks::put);
|
||||
init();
|
||||
}
|
||||
|
||||
private RequestInfoUpdate clientCodeNewFromGateway(BaseRequest<TkrAccountsGatewayRequest> tkrRequest) {
|
||||
private void clientCodeNewFromGateway(BaseRequest<TkrAccountsGatewayRequest> tkrRequest) {
|
||||
clientCodeFacade.lock();
|
||||
//валидация запроса
|
||||
TkrAccountsGatewayRequest gatewayRequest = tkrRequest.getRequestPayload();
|
||||
|
|
@ -113,7 +117,6 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali
|
|||
.toList();
|
||||
sendTkrRequest.setTkrs(tkrs);
|
||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||
return null;
|
||||
}
|
||||
|
||||
//unwrap запроса
|
||||
|
|
@ -132,7 +135,7 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali
|
|||
clientCodeNewRequest.setCompanyId(company.getId());
|
||||
clientCodeNewRequest.setCode(tkrAccount.getClientCode());
|
||||
clientCodeNewRequest.setDepoAccountId(depoAccountId);
|
||||
Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB);
|
||||
Account moneyAccount = accountByCurrency.get(CurrencyCode.RUB.getKey());
|
||||
clientCodeNewRequest.setMoneyAccountId(moneyAccount.getId());
|
||||
|
||||
List<Long> foreignCurrencyList = accountByCurrency.entrySet()
|
||||
|
|
@ -149,20 +152,20 @@ public class ClientCodeMessageListener extends QueueConsumer implements Initiali
|
|||
log.debug("Success created new clientCode by tkr account: {}", tkrAccount.getClientCode());
|
||||
}
|
||||
kafkaSender.sendRequestToQueue(Consts.ACCOUNTS_TO_GATEWAY, sendTkrRequest);
|
||||
return null;
|
||||
}
|
||||
|
||||
private RequestInfoUpdate clientCodeNewFromBackend(BaseRequest<ClientCodeNewRequest> clientCodeNewRequest) {
|
||||
private PayloadInfo clientCodeNewFromBackend(BaseRequest<ClientCodeNewRequest> clientCodeNewRequest) {
|
||||
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(clientCodeNewRequest);
|
||||
if (requestInfoUpdate != null) return requestInfoUpdate;
|
||||
ClientCodeNewRequest request = clientCodeNewRequest.getRequestPayload();
|
||||
ValidationResult validationResult = clientCodeValidator.checkBackendNewRequest(request);
|
||||
if (!validationResult.isValid()) {
|
||||
return makeErrorResponse(clientCodeNewRequest, validationResult);
|
||||
return DefaultResponses.build(clientCodeNewRequest.getId(),
|
||||
validationResult.errorMsg().orElse(""), DefaultResponses.ERROR);
|
||||
}
|
||||
|
||||
clientCodeFacade.createClientCode(request, validationResult.validator());
|
||||
return null;
|
||||
return DefaultResponses.build(clientCodeNewRequest.getId(), DefaultResponses.SUCCESS);
|
||||
}
|
||||
|
||||
private RequestInfoUpdate makeErrorResponse(BaseRequest<ClientCodeNewRequest> clientCodeNewRequest,
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ import org.springframework.stereotype.Service;
|
|||
import ru.spcex.clearing.imdg.IMDGDistributedNames;
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.domain.Consts;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
|
||||
import ru.spcex.clearing.platform.messaging.service.RequestInfo;
|
||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||
|
|
@ -32,15 +33,15 @@ public class RequestInfoAccepter extends QueueConsumer implements InitializingBe
|
|||
|
||||
@Override
|
||||
public void afterPropertiesSet() throws Exception {
|
||||
callback(RequestInfoUpdate.class)
|
||||
callback(PayloadInfo.class)
|
||||
.setConsumer(this::updateRequestInfo)
|
||||
.forDestination(Consts.REQUEST_INFO_UPDATE, callbacks::put);
|
||||
init();
|
||||
}
|
||||
|
||||
private void updateRequestInfo(BaseRequest<RequestInfoUpdate> requestInfoUpdateBaseRequest) {
|
||||
private void updateRequestInfo(BaseRequest<PayloadInfo> requestInfoUpdateBaseRequest) {
|
||||
//will throw exception for any class other than RequestInfoUpdate
|
||||
RequestInfoUpdate statusInfo = requestInfoUpdateBaseRequest.getRequestPayload();
|
||||
PayloadInfo statusInfo = requestInfoUpdateBaseRequest.getRequestPayload();
|
||||
if (statusInfo == null || statusInfo.getId() == null) {
|
||||
log.warn("BaseRequest id={}, RequestPayload={}: requestInfo id was null",
|
||||
requestInfoUpdateBaseRequest.getId(), statusInfo);
|
||||
|
|
|
|||
|
|
@ -1,6 +1,7 @@
|
|||
package ru.spcex.clearing.platform.messaging.logic.functional;
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
|
@ -19,5 +20,5 @@ public interface BuilderConsumerStep<T1> {
|
|||
/**
|
||||
* с обновлением RequestInfo в мапе
|
||||
*/
|
||||
BuilderDestinationStep setFunction(Function<BaseRequest<T1>, Object> function);
|
||||
BuilderDestinationStep setFunction(Function<BaseRequest<T1>, PayloadInfo> function);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,7 +1,10 @@
|
|||
package ru.spcex.clearing.platform.messaging.logic.functional;
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.DefaultResponses;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
|
||||
import java.util.Optional;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
|
|
@ -19,7 +22,7 @@ public class ConsumerSpecificClass<T> implements BuilderConsumerStep<T>, Builder
|
|||
* действие-обработчик реквеста (метод из конкретного менеджера)
|
||||
*/
|
||||
private java.util.function.Consumer<BaseRequest<T>> consumer;
|
||||
private java.util.function.Function<BaseRequest<T>, Object> function;
|
||||
private java.util.function.Function<BaseRequest<T>, PayloadInfo> function;
|
||||
|
||||
private ConsumerSpecificClass(Class<T> clazz) {
|
||||
this.clazz = clazz;
|
||||
|
|
@ -40,7 +43,7 @@ public class ConsumerSpecificClass<T> implements BuilderConsumerStep<T>, Builder
|
|||
}
|
||||
|
||||
@Override
|
||||
public BuilderDestinationStep setFunction(Function<BaseRequest<T>, Object> function) {
|
||||
public BuilderDestinationStep setFunction(Function<BaseRequest<T>, PayloadInfo> function) {
|
||||
this.function = function;
|
||||
return this;
|
||||
}
|
||||
|
|
@ -50,18 +53,22 @@ public class ConsumerSpecificClass<T> implements BuilderConsumerStep<T>, Builder
|
|||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public Object acceptRaw(Object obj) {
|
||||
public PayloadInfo acceptRaw(Object obj) {
|
||||
if (this.consumer != null) {
|
||||
this.consumer.accept((BaseRequest<T>) obj);
|
||||
return null;
|
||||
return DefaultResponses.build(((BaseRequest<?>) obj).getId(), DefaultResponses.SUCCESS);
|
||||
} else {
|
||||
return this.function.apply((BaseRequest<T>) obj);
|
||||
PayloadInfo res = this.function.apply((BaseRequest<T>) obj);
|
||||
return res == null ?
|
||||
DefaultResponses.build(((BaseRequest<?>) obj).getId(), DefaultResponses.SUCCESS) :
|
||||
res;
|
||||
}
|
||||
}
|
||||
|
||||
public Consumer<BaseRequest<T>> getConsumer() {
|
||||
return consumer;
|
||||
}
|
||||
|
||||
public Class<T> getClazz() {
|
||||
return clazz;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,30 @@
|
|||
package ru.spcex.clearing.platform.messaging.logic.response;
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||
|
||||
public enum DefaultResponses {
|
||||
SUCCESS,
|
||||
ERROR,
|
||||
;
|
||||
public static PayloadInfo build(Long id, String message, DefaultResponses defaultResponse) {
|
||||
PayloadInfo info = build(id, defaultResponse);
|
||||
info.setMessage(message);
|
||||
return info;
|
||||
}
|
||||
public static PayloadInfo build(Long id, DefaultResponses defaultResponse) {
|
||||
PayloadInfo info = new RequestInfoUpdate();
|
||||
info.setId(id);
|
||||
switch (defaultResponse) {
|
||||
case SUCCESS: {
|
||||
info.setStatus(Status.Success);
|
||||
break;
|
||||
}
|
||||
case ERROR: {
|
||||
info.setStatus(Status.Error);
|
||||
break;
|
||||
}
|
||||
}
|
||||
return info;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,15 @@
|
|||
package ru.spcex.clearing.platform.messaging.logic.response;
|
||||
|
||||
|
||||
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
|
||||
import ru.spcex.clearing.platform.messaging.service.Status;
|
||||
|
||||
public interface PayloadInfo {
|
||||
Long getId();
|
||||
Status getStatus();
|
||||
String getMessage();
|
||||
|
||||
PayloadInfo setId(Long id);
|
||||
PayloadInfo setStatus(Status status);
|
||||
RequestInfoUpdate setMessage(String message);
|
||||
}
|
||||
|
|
@ -18,6 +18,8 @@ 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.clearing.platform.messaging.logic.response.DefaultResponses;
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
import ru.spcex.clearing.platform.messaging.service.sender.CorrelationHeader;
|
||||
import ru.spcex.platform.utils.enumeration.EnumMessage;
|
||||
import ru.spcex.platform.utils.enumeration.IMessageResolver;
|
||||
|
|
@ -108,7 +110,7 @@ public class QueueConsumer implements AutoCloseable {
|
|||
((BaseRequest<?>) o).setCorrelationId(correlationId.value());
|
||||
}
|
||||
if (needsProcessing(next.topic(), (BaseRequest<?>) o)) {
|
||||
Object topicResponse = callback.acceptRaw(o);
|
||||
PayloadInfo topicResponse = callback.acceptRaw(o);
|
||||
if (producer != null) {
|
||||
sendResponse((BaseRequest<?>) o, topicResponse, correlationId);
|
||||
}
|
||||
|
|
@ -117,12 +119,13 @@ public class QueueConsumer implements AutoCloseable {
|
|||
lastErrors = 0;
|
||||
} catch (Throwable e) {
|
||||
log.error("Listener {} last topic \"{}\", error: {}",
|
||||
QueueConsumer.this.getClass().getName(), lastTopic, ExceptionUtils.getStackTrace(e));
|
||||
QueueConsumer.this.getClass().getName(), lastTopic, ExceptionUtils.getStackTrace(e));
|
||||
if (lastTopic == null) {
|
||||
log.debug("main thread stacktrace: {}", ExceptionUtils.getStackTrace(debugCallerStacktrace));
|
||||
}
|
||||
if (producer != null && o != null) {
|
||||
sendErrorResponse((BaseRequest<?>) o, correlationId);
|
||||
BaseRequest<?> br = (BaseRequest<?>) o;
|
||||
sendResponse(br, DefaultResponses.build(br.getId(), DefaultResponses.ERROR), correlationId);
|
||||
}
|
||||
if (lastErrors++ > 20) {
|
||||
log.warn("Too many error at row, {}. Sleep.", lastErrors);
|
||||
|
|
@ -139,57 +142,23 @@ public class QueueConsumer implements AutoCloseable {
|
|||
if (!closed.get()) throw e;
|
||||
} catch (Throwable e) {
|
||||
log.error("QueueConsumer error: {}; main thread stacktrace: {}",
|
||||
ExceptionUtils.getStackTrace(e), ExceptionUtils.getStackTrace(debugCallerStacktrace));
|
||||
ExceptionUtils.getStackTrace(e), ExceptionUtils.getStackTrace(debugCallerStacktrace));
|
||||
} finally {
|
||||
consumer.close();
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
//мб перенести в другой класс
|
||||
private void sendErrorResponse(BaseRequest<?> o, Header correlationId) {
|
||||
private void sendResponse(BaseRequest<?> o, PayloadInfo response, Header correlationId) {
|
||||
outputExecutor.submit(() -> {
|
||||
try {
|
||||
Future<RecordMetadata> send;
|
||||
//default response
|
||||
BaseRequest<RequestInfoUpdate> req = new BaseRequest<>();
|
||||
RequestInfoUpdate success = new RequestInfoUpdate();
|
||||
success.setId(o.getId());
|
||||
success.setStatus(Status.Error);
|
||||
req.setRequestPayload(success);
|
||||
req.setId(o.getId());
|
||||
req.setActionType(ActionType.SYSTEM);
|
||||
ProducerRecord<String, Object> respRec = new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req);
|
||||
if (correlationId != null) {
|
||||
respRec.headers().add(KafkaHeaders.CORRELATION_ID, correlationId.value());
|
||||
}
|
||||
send = producer.send(respRec);
|
||||
send.get();
|
||||
} catch (Exception e) {
|
||||
if (e instanceof InterruptedException) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
log.error(ExceptionUtils.getStackTrace(e));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
protected void sendResponse(BaseRequest<?> o, Object response, Header correlationId) {
|
||||
outputExecutor.submit(() -> {
|
||||
try {
|
||||
Future<RecordMetadata> send;
|
||||
BaseRequest<Object> req = new BaseRequest<>();
|
||||
BaseRequest<PayloadInfo> 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);
|
||||
}
|
||||
req.setRequestPayload(response);
|
||||
|
||||
ProducerRecord<String, Object> respRec = new ProducerRecord<>(Consts.REQUEST_INFO_UPDATE, req);
|
||||
if (correlationId != null) {
|
||||
respRec.headers().add(new CorrelationHeader(correlationId.value().clone()));
|
||||
|
|
@ -225,9 +194,9 @@ public class QueueConsumer implements AutoCloseable {
|
|||
String errorMsg = messageResolver.resolve(validationError.get());
|
||||
log.warn("validation error for {} error={}, id={}: {}", req.getClass().getSimpleName(), validationError.get().getSubject(), userRequest.getId(), errorMsg);
|
||||
return new RequestInfoUpdate()
|
||||
.setId(userRequest.getId())
|
||||
.setStatus(Status.Error)
|
||||
.setMessage(errorMsg);
|
||||
.setId(userRequest.getId())
|
||||
.setStatus(Status.Error)
|
||||
.setMessage(errorMsg);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
|
|
|
|||
|
|
@ -1,6 +1,9 @@
|
|||
package ru.spcex.clearing.platform.messaging.service;
|
||||
|
||||
public class RequestInfoUpdate {
|
||||
import ru.spcex.clearing.platform.messaging.logic.response.PayloadInfo;
|
||||
import ru.spcex.platform.classes.base.interfaces.WithId;
|
||||
|
||||
public class RequestInfoUpdate implements PayloadInfo {
|
||||
private Long id;
|
||||
private Status status;
|
||||
private String message;
|
||||
|
|
@ -13,10 +16,12 @@ public class RequestInfoUpdate {
|
|||
this.message = message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Long getId() {
|
||||
return id;
|
||||
}
|
||||
|
||||
@Override
|
||||
public RequestInfoUpdate setId(Long id) {
|
||||
this.id = id;
|
||||
return this;
|
||||
|
|
@ -26,6 +31,7 @@ public class RequestInfoUpdate {
|
|||
return status;
|
||||
}
|
||||
|
||||
@Override
|
||||
public RequestInfoUpdate setStatus(Status status) {
|
||||
this.status = status;
|
||||
return this;
|
||||
|
|
@ -35,6 +41,7 @@ public class RequestInfoUpdate {
|
|||
return message;
|
||||
}
|
||||
|
||||
@Override
|
||||
public RequestInfoUpdate setMessage(String message) {
|
||||
this.message = message;
|
||||
return this;
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue