InspectionObligations GatewayRequester
This commit is contained in:
parent
66ab0cad31
commit
c4c2afe72c
9 changed files with 206 additions and 3 deletions
|
|
@ -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<String, Object> createConsumer(ClearingServiceSettings settings) {
|
||||
return KafkaConsumerFactory.consumer(settings.getKafkaConsumer());
|
||||
}
|
||||
|
||||
@Autowired
|
||||
@Bean("kafkaConsumerGateway")
|
||||
public Supplier<Consumer<String, Object>> 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<String, Object> createProducer(ClearingServiceSettings settings) {
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue,
|
||||
public EventsReceiver(@Qualifier("kafkaConsumer") Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue,
|
||||
IMessageResolver errorResolver,
|
||||
ClearingService clearingService,
|
||||
RegistryService registryService,
|
||||
|
|
|
|||
|
|
@ -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<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue,
|
||||
public LauncherCommandReceiver(@Qualifier("kafkaConsumer") Consumer<String, Object> kafkaQueue, Producer<String, Object> kafkaResponseQueue,
|
||||
ClearingService clearingService) {
|
||||
super(kafkaQueue, kafkaResponseQueue);
|
||||
this.clearingService = clearingService;
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String, Object>> 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<Boolean> gatewayRequestAndWait(Registry om_t) {
|
||||
boolean gatewayReceived = false;
|
||||
gatewayLock.lock();
|
||||
try {
|
||||
//send Request
|
||||
AssetOperationListRequest assetOperationListRequest = new AssetOperationListRequest();
|
||||
Collection<AssetOperationRequest> 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<AssetOperationApprovalRequest> sysReq) {
|
||||
log.trace("gateway answer received, BaseRequest.id={}", sysReq.getId());
|
||||
AssetOperationApprovalRequest requestPayload = sysReq.getRequestPayload();
|
||||
List<SingleAssetResponse> 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<SingleAssetResponse> 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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -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 * * * * *
|
||||
|
|
|
|||
|
|
@ -44,4 +44,9 @@
|
|||
<appender-ref ref="FILE"/>
|
||||
<appender-ref ref="CONSOLE"/>
|
||||
</logger>
|
||||
|
||||
<logger name="ru.spcex.clearing.session.stage.impl.GatewayRequester" level="info" additivity="false">
|
||||
<appender-ref ref="FILE"/>
|
||||
<appender-ref ref="CONSOLE"/>
|
||||
</logger>
|
||||
</configuration>
|
||||
Loading…
Add table
Reference in a new issue