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 2726478a3..e5bdfcc49 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 @@ -2,6 +2,8 @@ package ru.spcex.clearing.config; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.producer.Producer; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.ConfigurableBeanFactory; import org.springframework.context.annotation.Bean; @@ -26,19 +28,31 @@ import java.util.function.Supplier; @Configuration public class KafkaConfig { - @Autowired - @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) - @Bean("kafkaConsumer") - public Consumer createConsumer(ClearingServiceSettings settings) { - return KafkaConsumerFactory.consumer(settings.getKafkaConsumer()); + private final Logger log = LoggerFactory.getLogger(getClass()); + private final KafkaConsumerSettings kafkaSettings; + + public KafkaConfig(ClearingServiceSettings settings) { + this.kafkaSettings = settings.getKafkaConsumer(); + Integer gtwTimeout = settings.getSessionStage().getInspectionGatewayTimeout(); + if (gtwTimeout != null) { + int maxPollIntervalMs = gtwTimeout * 1000 + 1000; + log.debug("setting max.poll.interval.ms {}", maxPollIntervalMs); + this.kafkaSettings.setManuallyMaxPollIntervalMs(maxPollIntervalMs); + } + } + + + @Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE) + @Bean("kafkaConsumer") + public Consumer createConsumer() { + return KafkaConsumerFactory.consumer(kafkaSettings); } - @Autowired @Bean("kafkaConsumerGateway") - public Supplier> getwaySessionConsumer(ClearingServiceSettings settings) { + public Supplier> getwaySessionConsumer() { AtomicInteger groupId = new AtomicInteger(1); return () -> { - KafkaConsumerSettings consumer = settings.getKafkaConsumer(); + KafkaConsumerSettings consumer = kafkaSettings; KafkaConsumerSettings gatewayConsumer = new KafkaConsumerSettings(); gatewayConsumer.setBootstrapServers(consumer.getBootstrapServers()); gatewayConsumer.setSessionTimeoutMs(consumer.getSessionTimeoutMs()); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java index c1f8999b2..ad1b60132 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/config/ValidationConfig.java @@ -368,6 +368,7 @@ public class ValidationConfig { new PresentById(IMDGDistributedNames.Map_Registry, ClearingError.RecordNotFound, true), StatusExtractValidationRule.RegistryCodeCheck, StatusExtractValidationRule.ContractCheck, + StatusExtractValidationRule.RequestStatusValid, StatusExtractValidationRule.RegistryStatusValid ); }; diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java index 32a6038c5..3c58d46ef 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/RegistryService.java @@ -48,16 +48,11 @@ import ru.spcex.platform.utils.validation.IValidator; import java.math.BigDecimal; import java.time.Instant; import java.time.LocalDate; -import java.util.Collection; -import java.util.List; -import java.util.Map; -import java.util.Optional; +import java.util.*; import java.util.function.Function; -import java.util.function.Supplier; import java.util.stream.Collectors; import java.util.stream.Stream; -import static ru.spcex.clearing.session.stage.impl.GatewayRequester.mapError; import static ru.spcex.platform.utils.number.BigDecimalUtil.safeBD; @Service @@ -290,23 +285,11 @@ public class RegistryService { return reqHelp.error(req.getId(), err.get()); } Registry rgs = validator.getStored(Stored.PresentById); - Supplier gtwBuilder = () -> GatewayRequestCreator.from(rgs, InOutDirection.in); - Optional gatewayOk; if (trdTime.isTradingTime()) { - gatewayOk = gateway.gatewayRequestAndWait(gtwBuilder); + Long gtwReqId = kafkaSender.sendRequestToQueue(Consts.ASSET_OPERATION, gtwReq(rgs)); + log.debug("send request to {} id={}", Consts.ASSET_OPERATION, gtwReqId); } else { log.debug("RegistryChangeStatusExtractRequest rgs.id={} not sending gateway request", rgs.getId()); - gatewayOk = Optional.of(true); - } - if (gatewayOk.isEmpty() || !gatewayOk.get()) { - String gtwErr = msgResolver.resolve(mapError(gatewayOk)); - log.error("{}.id={} {}", - rgs.getRegistryCode(), - rgs.getId(), - gtwErr); - notification.sendNotification(ObjectType.rgst, - "Отметка о получении выписки: %s".formatted(gtwErr), - Priority.HIGH); } rgs.setRegistryStatus(payload.getRegistryStatus()); log.debug("changing registry.id={} status to {}", rgs.getId(), payload.getRegistryStatus()); @@ -321,6 +304,13 @@ public class RegistryService { return reqHelp.success(req.getId()); } + private AssetOperationListRequest gtwReq(Registry rgs) { + AssetOperationRequest item = GatewayRequestCreator.from(rgs, InOutDirection.in); + AssetOperationListRequest gtwReq = new AssetOperationListRequest(); + gtwReq.setAssetOperationRequests(Collections.singletonList(item)); + return gtwReq; + } + public RequestInfoUpdate identificationFunds(BaseRequest req) { if (!rights.userHasRole(req.getUserId(), UserRole.Admin)) { return reqHelp.error(req.getId(), ClearingError.UserVerifyDenial); diff --git a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/StatusExtractValidationRule.java b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/StatusExtractValidationRule.java index 0bd51ec41..d31159f21 100644 --- a/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/StatusExtractValidationRule.java +++ b/clearing-parent/clearing-service/src/main/java/ru/spcex/clearing/service/validation/StatusExtractValidationRule.java @@ -63,7 +63,7 @@ public enum StatusExtractValidationRule implements IValidationRule validate(ImdgValidationContext context) { RegistryChangeStatusExtractRequest validatedObject = context.getValidatedObject(); @@ -76,6 +76,17 @@ public enum StatusExtractValidationRule implements IValidationRule validate(ImdgValidationContext context) { + Registry rgs = context.getStoredObject(Stored.PresentById); + RegistryStatus status = IEnumKey.getEnumByKey(RegistryStatus.class, rgs.getRegistryStatus()); + if (RegistryStatus.OK.equals(status)) { + return of(ClearingError.ObligationsAlreadyCalculated); + } + return empty(); + } } ; private final static Logger log = LoggerFactory.getLogger(StatusExtractValidationRule.class); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java index 26e7209f8..1d765f2ad 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/KafkaConsumerFactory.java @@ -22,6 +22,9 @@ public class KafkaConsumerFactory { props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString()); props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString()); props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset()); + if (kafkaSettings.getMaxPollIntervalMs() != null) { + props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs())); + } props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); return new KafkaConsumer<>(props); @@ -36,6 +39,9 @@ public class KafkaConsumerFactory { } props.put("enable.auto.commit", kafkaSettings.getEnableAutoCommit().toString()); props.put("session.timeout.ms", kafkaSettings.getSessionTimeoutMs().toString()); + if (kafkaSettings.getMaxPollIntervalMs() != null) { + props.put("max.poll.interval.ms", String.valueOf(kafkaSettings.getMaxPollIntervalMs())); + } props.put("auto.offset.reset", kafkaSettings.getAutoOffsetReset()); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // props.put("value.deserializer", JsonDeserializer.class.getName()); diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java index 8e6a6fc73..a6eab1d81 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/config/element/KafkaConsumerSettings.java @@ -6,6 +6,7 @@ public class KafkaConsumerSettings { private Boolean enableAutoCommit; private Integer sessionTimeoutMs; private String autoOffsetReset; + private Integer maxPollIntervalMs = null; public String getBootstrapServers() { @@ -47,4 +48,13 @@ public class KafkaConsumerSettings { public void setAutoOffsetReset(String autoOffsetReset) { this.autoOffsetReset = autoOffsetReset; } + + + public Integer getMaxPollIntervalMs() { + return maxPollIntervalMs; + } + + public void setManuallyMaxPollIntervalMs(Integer maxPollIntervalMs) { + this.maxPollIntervalMs = maxPollIntervalMs; + } }