отправить ответный запрос с gateway в клиринг сервис

This commit is contained in:
etreschenkov 2023-06-21 15:35:28 +03:00
parent 11fbb281a6
commit 43dde65830
3 changed files with 61 additions and 25 deletions

View file

@ -21,6 +21,7 @@ import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.ExportToFileRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.balance.StatementRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationListRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
@ -80,7 +81,8 @@ public class Sdf06Executor {
@Autowired
public Sdf06Executor(ImdgProvider imdgProvider,
IMessageResolver messageResolver,
@Qualifier("sdf06ValidatorNew") Function<SDf06, IValidator> sDf06Validator, KafkaSender kafkaSender) {
@Qualifier("sdf06ValidatorNew") Function<SDf06, IValidator> sDf06Validator,
KafkaSender kafkaSender) {
this.imdgProvider = imdgProvider;
this.statementImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Statement, Statement.class);
this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class);
@ -158,7 +160,9 @@ public class Sdf06Executor {
if (requests.size() > 0) {
this.sdf06GroupId = groupId;
this.sdf07GroupId = sdf07GroupId;
kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, requests);
AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest();
assetOperationListRequest.setAssetOperationRequests(requests);
kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest);
} else if (sdf07WasCreated) {
sendToExporter(groupId);
}

View file

@ -16,17 +16,20 @@ import ru.spcex.clearing.gatewayapi.enums.OutboundRequestType;
import ru.spcex.clearing.gatewayapi.service.builders.OutboundRequestBuilder;
import ru.spcex.clearing.platform.messaging.domain.BaseRequest;
import ru.spcex.clearing.platform.messaging.domain.Consts;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationListRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.AssetOperationApprovalRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.GatewayTaskRequest;
import ru.spcex.clearing.platform.messaging.domain.cud.gateway.SingleAssetResponse;
import ru.spcex.clearing.platform.messaging.domain.cud.utilities.LimExportedRequest;
import ru.spcex.clearing.platform.messaging.service.QueueConsumer;
import ru.spcex.clearing.platform.messaging.service.RequestInfoUpdate;
import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender;
import ru.spcex.clearing.util.security.UserRoleVerification;
import ru.spcex.platform.enumeration.Section;
import ru.spcex.platform.enumeration.Task;
import java.util.Collections;
import java.util.Map;
import java.util.*;
@Service
public class GatewayService extends QueueConsumer implements InitializingBean {
@ -35,16 +38,19 @@ public class GatewayService extends QueueConsumer implements InitializingBean {
private final UserRoleVerification userRoleVerification;
private final RestTemplate restTemplate;
private final InboundServerSettings inboundServerSettings;
private final KafkaSender kafkaSender;
public GatewayService(Consumer<String, Object> kafkaQueue,
Producer<String, Object> kafkaProducer,
RestTemplate restTemplate,
UserRoleVerification userRoleVerification,
GatewayApiSettings gatewayApiSettings) {
GatewayApiSettings gatewayApiSettings,
KafkaSender kafkaSender) {
super(kafkaQueue, kafkaProducer);
this.restTemplate = restTemplate;
this.userRoleVerification = userRoleVerification;
this.inboundServerSettings = gatewayApiSettings.getInboundServer();
this.kafkaSender = kafkaSender;
}
@Override
@ -58,7 +64,7 @@ public class GatewayService extends QueueConsumer implements InitializingBean {
callback(LimExportedRequest.class)
.setConsumer(this::requestOnLimit)
.forDestination(Consts.LIM_EXPORTED, callbacks::put);
callback(AssetOperationRequest.class)
callback(AssetOperationListRequest.class)
.setConsumer(this::requestOnAsset)
.forDestination(Consts.ASSET_OPERATION, callbacks::put);
init();
@ -138,33 +144,41 @@ public class GatewayService extends QueueConsumer implements InitializingBean {
log.debug("LOSC task 2/2 complete (MKR section), response: {}", bodyResponse);
}
public void requestOnAsset(BaseRequest<AssetOperationRequest> request) {
RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(request);
if (requestInfoUpdate != null) return;
AssetOperationRequest assetOperationRequest = request.getRequestPayload();
Map<String, Object> content = Map.of("direction", assetOperationRequest.getDirection(),
"amount", assetOperationRequest.getAmount(),
"code", assetOperationRequest.getCode(),
"asset", assetOperationRequest.getSecuritySymbol(),
"firm_id", assetOperationRequest.getTradingCode());
public void requestOnAsset(BaseRequest<AssetOperationListRequest> request) {
Collection<AssetOperationRequest> assetOperations = request.getRequestPayload().getAssetOperationRequests();
List<SingleAssetResponse> responses = new ArrayList<>();
AssetOperationApprovalRequest assetOperationApprovalRequest = new AssetOperationApprovalRequest();
for (AssetOperationRequest assetOperation : assetOperations) {
Map<String, Object> content = Map.of("direction", assetOperation.getDirection(),
"amount", assetOperation.getAmount(),
"code", assetOperation.getCode(),
"asset", assetOperation.getSecuritySymbol(),
"firm_id", assetOperation.getTradingCode());
// if (assetOperationRequest.getAmount() != null) {
// content.put("amount", assetOperationRequest.getAmount());
// } else {
// content.put("quantity", assetOperationRequest.getQuantity());
// }
OutboundRequest outboundRequest = OutboundRequestBuilder.builder()
.section(Section.MKR.getKey())
.type(OutboundRequestType.ASSET_OPERATION)
.content(content).build();
OutboundRequest outboundRequest = OutboundRequestBuilder.builder()
.section(Section.MKR.getKey())
.type(OutboundRequestType.ASSET_OPERATION)
.content(content).build();
String url = formingInboundUrl(inboundServerSettings.getPathLOCM());
String url = formingInboundUrl(inboundServerSettings.getPathLOCM());
HttpEntity<OutboundRequest> r = makeDefaultRequest(outboundRequest);
ResponseEntity<Map> response = restTemplate.exchange(url, HttpMethod.POST, r, Map.class);
Map bodyResponse = response.getBody();
log.debug("Response from service: {}", bodyResponse);
HttpEntity<OutboundRequest> r = makeDefaultRequest(outboundRequest);
ResponseEntity<Map> response = restTemplate.exchange(url, HttpMethod.POST, r, Map.class);
Map bodyResponse = response.getBody();
SingleAssetResponse singleAssetResponse = new SingleAssetResponse();
singleAssetResponse.setApproved(Boolean.parseBoolean((String) bodyResponse.get("result")));
singleAssetResponse.setStatementId(assetOperation.getStatementId());
responses.add(singleAssetResponse);
log.debug("Response from service: {}", bodyResponse);
}
assetOperationApprovalRequest.setApprovals(responses);
kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION_APPROVAL, responses);
}
private String formingInboundUrl(String path) {

View file

@ -0,0 +1,18 @@
package ru.spcex.clearing.platform.messaging.domain.cud.clearing;
import com.fasterxml.jackson.annotation.JsonProperty;
import java.util.Collection;
public class AssetOperationListRequest {
@JsonProperty
private Collection<AssetOperationRequest> assetOperationRequests;
public Collection<AssetOperationRequest> getAssetOperationRequests() {
return assetOperationRequests;
}
public void setAssetOperationRequests(Collection<AssetOperationRequest> assetOperationRequests) {
this.assetOperationRequests = assetOperationRequests;
}
}