diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java index 95ad7f18b..1c6662ff3 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/KafkaConfig.java @@ -13,6 +13,7 @@ import ru.spcex.clearing.config.element.ClearingServiceSettings; import ru.spcex.clearing.imdg.IMDGDistributedNames; import ru.spcex.clearing.platform.messaging.config.KafkaConsumerFactory; import ru.spcex.clearing.platform.messaging.config.KafkaProducerFactory; +import ru.spcex.clearing.platform.messaging.config.element.KafkaConsumerSettings; import ru.spcex.clearing.platform.messaging.config.element.KafkaProducerSettings; import ru.spcex.clearing.platform.messaging.service.RequestInfo; import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; @@ -20,15 +21,34 @@ import ru.spcex.platform.imdg.api.Imdg; import ru.spcex.platform.imdg.api.ImdgId; import ru.spcex.platform.imdg.api.ImdgProvider; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Supplier; + @Configuration public class KafkaConfig { @Autowired @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) - @Bean + @Bean("kafkaConsumer") public Consumer createConsumer(ClearingServiceSettings settings) { return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); } + @Autowired + @Bean("kafkaConsumerGateway") + public Supplier> getwaySessionConsumer(ClearingServiceSettings settings) { + AtomicInteger groupId = new AtomicInteger(1); + return () -> { + KafkaConsumerSettings consumer = settings.getKafkaConsumer(); + KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings(); + gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers()); + gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs()); + gatewayConsumer.setAutoOffsetReset(consumer.getAutoOffsetReset()); + gatewayConsumer.setEnableAutoCommit(consumer.getEnableAutoCommit()); + gatewayConsumer.setGroupId(consumer.getGroupId() + "-session-" + groupId.getAndIncrement()); + return KafkaConsumerFactory.consumer(gatewayConsumer); + }; + } + @Autowired @Bean public Producer createProducer(ClearingServiceSettings settings) { diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java index d42e5fcc3..7e9600c2f 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/ClearingServiceSettings.java @@ -14,6 +14,7 @@ public class ClearingServiceSettings { private HazelcastClientParams hazelcast; private KafkaConsumerSettings kafkaConsumer; private KafkaProducerSettings kafkaProducer; + private SessionStageSettings sessionStage; public HazelcastClientParams getHazelcast() { return hazelcast; @@ -38,4 +39,12 @@ public class ClearingServiceSettings { public void setKafkaProducer(KafkaProducerSettings kafkaProducer) { this.kafkaProducer = kafkaProducer; } + + public SessionStageSettings getSessionStage() { + return sessionStage; + } + + public void setSessionStage(SessionStageSettings sessionStage) { + this.sessionStage = sessionStage; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/SessionStageSettings.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/SessionStageSettings.java new file mode 100644 index 000000000..08e7671e0 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/element/SessionStageSettings.java @@ -0,0 +1,14 @@ +package ru.spcex.clearing.config.element; + +public class SessionStageSettings { + //inspection-gateway-timeout + private Integer inspectionGatewayTimeout = 10; + + public Integer getInspectionGatewayTimeout() { + return inspectionGatewayTimeout; + } + + public void setInspectionGatewayTimeout(Integer inspectionGatewayTimeout) { + this.inspectionGatewayTimeout = inspectionGatewayTimeout; + } +} diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java index db3b9959c..4ba7ddd47 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/EventsReceiver.java @@ -6,6 +6,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import ru.spcex.clearing.platform.messaging.domain.BaseRequest; import ru.spcex.clearing.platform.messaging.domain.Consts; @@ -58,7 +59,7 @@ public class EventsReceiver extends QueueConsumer implements InitializingBean { private final PaymentInstructionOutboundService pmtOutboundService; @Autowired - public EventsReceiver(Consumer kafkaQueue, Producer kafkaResponseQueue, + public EventsReceiver(@Qualifier("kafkaConsumer") Consumer kafkaQueue, Producer kafkaResponseQueue, IMessageResolver errorResolver, ClearingService clearingService, RegistryService registryService, diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/LauncherCommandReceiver.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/LauncherCommandReceiver.java index 850f33e89..0d75db589 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/LauncherCommandReceiver.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/LauncherCommandReceiver.java @@ -6,6 +6,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Service; import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; import ru.spcex.clearing.platform.messaging.service.QueueConsumer; @@ -17,7 +18,7 @@ public class LauncherCommandReceiver extends QueueConsumer implements Initializi private final ClearingService clearingService; @Autowired - public LauncherCommandReceiver(Consumer kafkaQueue, Producer kafkaResponseQueue, + public LauncherCommandReceiver(@Qualifier("kafkaConsumer") Consumer kafkaQueue, Producer kafkaResponseQueue, ClearingService clearingService) { super(kafkaQueue, kafkaResponseQueue); this.clearingService = clearingService; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java index 0cf52d5c4..313331997 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/integration/GatewayRequestCreator.java @@ -3,6 +3,9 @@ package ru.spcex.clearing.service.integration; import ru.clearing.classes.statics.data.statement.Statement; import ru.spcex.clearing.platform.messaging.domain.cud.clearing.AssetOperationRequest; import ru.spcex.platform.enumeration.CurrencyCode; +import ru.spcex.platform.enumeration.InOutDirection; + +import java.math.BigDecimal; public class GatewayRequestCreator { public static AssetOperationRequest gatewayRequestPart(Statement stmt, String tradingCode, String tcrCode) { @@ -15,4 +18,21 @@ public class GatewayRequestCreator { req.setDirection(stmt.getInOutDirection()); return req; } + + public static AssetOperationRequest gatewayRequestPart(Long entityId, + Long sessionId, + BigDecimal amount, + InOutDirection direction, + String tradingCode, + String tcrCode) { + AssetOperationRequest req = new AssetOperationRequest(); + req.setSessionId(sessionId); + req.setEntityId(entityId); + req.setAmount(amount); + req.setSecuritySymbol(CurrencyCode.RUB.getKey()); //fixme retrieve security symbol from validator + req.setTradingCode(tradingCode); + req.setCode(tcrCode); + req.setDirection(direction.getKey()); + return req; + } } diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/GatewayRequester.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/GatewayRequester.java new file mode 100644 index 000000000..639d7ff99 --- /dev/null +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/session/stage/impl/GatewayRequester.java @@ -0,0 +1,131 @@ +package ru.spcex.clearing.session.stage.impl; + +import org.apache.kafka.clients.consumer.Consumer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.beans.factory.config.ConfigurableBeanFactory; +import org.springframework.context.annotation.Scope; +import org.springframework.stereotype.Component; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.config.element.ClearingServiceSettings; +import ru.spcex.clearing.config.element.SessionStageSettings; +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.SingleAssetResponse; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.clearing.service.integration.GatewayRequestCreator; +import ru.spcex.platform.enumeration.InOutDirection; +import ru.spcex.platform.utils.log.ExceptionUtils; + +import java.math.BigDecimal; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Condition; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Supplier; + +@Scope(value = ConfigurableBeanFactory.SCOPE_PROTOTYPE) +@Component +public class GatewayRequester extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + private final KafkaSender kafkaSender; + private final Integer GATEWAY_TIMEOUT; + + private final Lock gatewayLock = new ReentrantLock(); + private final Condition gatewayCondition = gatewayLock.newCondition(); + private Long omtId = null; + private SingleAssetResponse assetOperationApprovalRequest = null; + + @Autowired + public GatewayRequester(@Qualifier("kafkaConsumerGateway") Supplier> consumer, + ClearingServiceSettings settings, + KafkaSender kafkaSender) { + super(consumer.get()); + this.kafkaSender = kafkaSender; + SessionStageSettings sessionSettings = settings.getSessionStage(); + this.GATEWAY_TIMEOUT = sessionSettings != null ? sessionSettings.getInspectionGatewayTimeout() : 10; + } + + @Override + public void afterPropertiesSet() throws Exception { + callback(AssetOperationApprovalRequest.class) + .setConsumer(this::receiveGatewayAnswer) + .forDestination(Consts.ASSET_OPERATION_APPROVAL, callbacks::put); + init(); + } + + public Optional gatewayRequestAndWait(Registry om_t) { + boolean gatewayReceived = false; + gatewayLock.lock(); + try { + //send Request + AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest(); + Collection requests = new ArrayList<>(); + AssetOperationRequest req = GatewayRequestCreator.gatewayRequestPart( + om_t.getId(), + null, + BigDecimal.ONE, + InOutDirection.out, + om_t.getTradingCode(), + om_t.getTradingClearingRegistry()); + requests.add(req); + + assetOperationListRequest.setAssetOperationRequests(requests); + log.debug("sending om*t.id={} to gateway", om_t.getId()); + kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, assetOperationListRequest); + this.omtId = om_t.getId(); + gatewayReceived = gatewayCondition.await(GATEWAY_TIMEOUT, TimeUnit.SECONDS); + } catch (InterruptedException e) { + log.error(ExceptionUtils.getStackTrace(e)); + } finally { + gatewayLock.unlock(); + } + if (gatewayReceived) { + boolean approved = assetOperationApprovalRequest.isApproved(); + log.debug("om*t.id={} answer from gateway received. approval: {}", om_t.getId(), approved); + assetOperationApprovalRequest = null; + return Optional.of(approved); + } else { + log.error("om*t.id={} answer from gateway not received.", om_t.getId()); + return Optional.empty(); + } + } + + private void receiveGatewayAnswer(BaseRequest sysReq) { + log.trace("gateway answer received, BaseRequest.id={}", sysReq.getId()); + AssetOperationApprovalRequest requestPayload = sysReq.getRequestPayload(); + List approvals = requestPayload.getApprovals(); + if (approvals == null) { + log.trace("gateway BaseRequest.id={} approvals null", sysReq.getId()); + return; + } + gatewayLock.lock(); + try { + if (omtId == null) { + log.trace("gateway BaseRequest.id={} currently not waiting for OM*T answer", sysReq.getId()); + return; + } + Optional assetResponseFound = + approvals.stream().filter(a -> omtId.equals(a.getEntityId())).findFirst(); + if (assetResponseFound.isPresent()) { + log.debug("gateway BaseRequest.id={} OM*T answer found for OM*T.id={}", sysReq.getId(), omtId); + omtId = null; + assetOperationApprovalRequest = assetResponseFound.get(); + gatewayCondition.signal(); + } + } finally { + gatewayLock.unlock(); + } + } +} diff --git a/clearing-parent/clearing-service/src/main/resources/application.properties b/clearing-parent/clearing-service/src/main/resources/application.properties index f80985c62..79b890a40 100644 --- a/clearing-parent/clearing-service/src/main/resources/application.properties +++ b/clearing-parent/clearing-service/src/main/resources/application.properties @@ -19,4 +19,6 @@ clearing-service.kafka-producer.batch-size=16384 clearing-service.kafka-producer.linger-ms=1 clearing-service.kafka-producer.buffer-memory=33554432 +clearing-service.session-stage.inspection-gateway-timeout=60 + clearing-service.scheduler.check-payment-instruction=*/5 * * * * * diff --git a/clearing-parent/clearing-service/src/main/resources/logback.xml b/clearing-parent/clearing-service/src/main/resources/logback.xml index 0536a83e0..ca7b0bb5b 100644 --- a/clearing-parent/clearing-service/src/main/resources/logback.xml +++ b/clearing-parent/clearing-service/src/main/resources/logback.xml @@ -44,4 +44,9 @@ + + + + + \ No newline at end of file