diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java index d9c164f9b..177779d87 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/executors/Sdf06Executor.java @@ -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 sDf06Validator, KafkaSender kafkaSender) { + @Qualifier("sdf06ValidatorNew") Function 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); } diff --git a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java index 37db4afbd..3366425fa 100644 --- a/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java +++ b/clearing-parent/gateway-api/src/main/java/ru/spcex/clearing/gatewayapi/service/GatewayService.java @@ -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 kafkaQueue, Producer 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 request) { - RequestInfoUpdate requestInfoUpdate = userRoleVerification.validateRoleAndGetResult(request); - if (requestInfoUpdate != null) return; - - AssetOperationRequest assetOperationRequest = request.getRequestPayload(); - Map content = Map.of("direction", assetOperationRequest.getDirection(), - "amount", assetOperationRequest.getAmount(), - "code", assetOperationRequest.getCode(), - "asset", assetOperationRequest.getSecuritySymbol(), - "firm_id", assetOperationRequest.getTradingCode()); + public void requestOnAsset(BaseRequest request) { + Collection assetOperations = request.getRequestPayload().getAssetOperationRequests(); + List responses = new ArrayList<>(); + AssetOperationApprovalRequest assetOperationApprovalRequest = new AssetOperationApprovalRequest(); + for (AssetOperationRequest assetOperation : assetOperations) { + Map 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 r = makeDefaultRequest(outboundRequest); - ResponseEntity response = restTemplate.exchange(url, HttpMethod.POST, r, Map.class); - Map bodyResponse = response.getBody(); - log.debug("Response from service: {}", bodyResponse); + HttpEntity r = makeDefaultRequest(outboundRequest); + ResponseEntity 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) { diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/AssetOperationListRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/AssetOperationListRequest.java new file mode 100644 index 000000000..dcf7158fa --- /dev/null +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/clearing/AssetOperationListRequest.java @@ -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 assetOperationRequests; + + public Collection getAssetOperationRequests() { + return assetOperationRequests; + } + + public void setAssetOperationRequests(Collection assetOperationRequests) { + this.assetOperationRequests = assetOperationRequests; + } +}