From 40f248ad7c5f90954d3f6dd369d47ab51855ea3b Mon Sep 17 00:00:00 2001 From: AKurakin Date: Mon, 5 May 2025 18:29:51 +0300 Subject: [PATCH] =?UTF-8?q?utility-service=20http://jira.mfd.msk:8088/brow?= =?UTF-8?q?se/CLS-839=20=D1=81=D0=B2=D0=B5=D1=80=D0=BA=D0=B0=20AM=5FT,AM?= =?UTF-8?q?=5FF=20=D0=BF=D0=BE=20GVTF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../utility/service/AmatAmafChecker.java | 223 ++++++++++++++++++ .../utility/service/AmatAmafCheckerTest.java | 147 ++++++++++++ .../cud/schedule/LauncherCommandRequest.java | 13 +- .../spcex/platform/utils/collection/Pair.java | 19 ++ 4 files changed, 401 insertions(+), 1 deletion(-) create mode 100644 clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/AmatAmafChecker.java create mode 100644 clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/AmatAmafCheckerTest.java diff --git a/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/AmatAmafChecker.java b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/AmatAmafChecker.java new file mode 100644 index 000000000..2c59775ab --- /dev/null +++ b/clearing-parent/utility-service/src/main/java/ru/spcex/clearing/utility/service/AmatAmafChecker.java @@ -0,0 +1,223 @@ +package ru.spcex.clearing.utility.service; + +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.InitializingBean; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.stereotype.Service; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.clearing.classes.statics.data.registry.TradingClearingRegistry; +import ru.spcex.clearing.imdg.IMDGDistributedNames; +import ru.spcex.clearing.platform.messaging.domain.BaseRequest; +import ru.spcex.clearing.platform.messaging.domain.Consts; +import ru.spcex.clearing.platform.messaging.domain.cud.schedule.LauncherCommandRequest; +import ru.spcex.clearing.platform.messaging.domain.cud.utilities.NotificationNewRequest; +import ru.spcex.clearing.platform.messaging.serialization.LogFormatter; +import ru.spcex.clearing.platform.messaging.service.QueueConsumer; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.*; +import ru.spcex.platform.imdg.api.Imdg; +import ru.spcex.platform.imdg.api.ImdgProvider; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicate; +import ru.spcex.platform.imdg.api.predicate.ImdgPredicateBuilder; +import ru.spcex.platform.imdg.api.predicate.specific.RegistryCodeSqlBuilder; +import ru.spcex.platform.imdg.iml.hazelcast.adapter.predicate.ImdgPredicateBuilderHazelcast; +import ru.spcex.platform.utils.collection.Pair; +import ru.spcex.platform.utils.number.BigDecimalUtil; + +import java.math.BigDecimal; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import java.util.*; +import java.util.stream.Collectors; + +@Service +public class AmatAmafChecker extends QueueConsumer implements InitializingBean { + private final Logger log = LoggerFactory.getLogger(getClass()); + + static final DateTimeFormatter MESSAGE_TIME_FORMATTER = DateTimeFormatter.ofPattern("yyyy.MM.dd HH:mm:ss"); + + private final Imdg registryImdg; + private final Imdg tradingClearingRegistryImdg; + private final ImdgProvider imdgProvider; + private final KafkaSender kafkaSender; + + @Autowired + public AmatAmafChecker(ImdgProvider imdgProvider, + KafkaSender kafkaSender, + Producer kafkaProducer, + Consumer kafkaQueue) { + super(kafkaQueue, kafkaProducer); + this.imdgProvider = imdgProvider; + this.kafkaSender = kafkaSender; + this.registryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_Registry, Registry.class); + this.tradingClearingRegistryImdg = imdgProvider.getImdg(IMDGDistributedNames.Map_TradingClearingRegistry, TradingClearingRegistry.class); + } + + @Override + public void afterPropertiesSet() throws Exception { + callback(LauncherCommandRequest.class) + .setConsumer(this::checkOnDate) + .forDestination(Task.GVTF.topic(), callbacks::put); + init(); + } + + public void checkOnDate(BaseRequest gvtfEvent) { + List specialId = null; + if (gvtfEvent.getRequestPayload() != null && gvtfEvent.getRequestPayload().getIds() != null) { + specialId = gvtfEvent.getRequestPayload().getIds(); + } + log.debug("Check AM*T and AM*F (parameter id's = {})", specialId); + + ImdgPredicateBuilder builder = ImdgPredicateBuilderHazelcast.instance(); + Collection checkRegistries; + final LocalDateTime checkTime = LocalDateTime.now(); + + ImdgPredicate AM_T_AM_F_predicate = RegistryCodeSqlBuilder.getInstance( + RegistryTradingParams.AM_T, + RegistryTradingParams.AM_F) + .buildPredicate(builder); + if (specialId == null) { + long timer = System.currentTimeMillis(); + Collection activeTcr = tradingClearingRegistryImdg.getCollectionObjectsByPredicate( + builder.equals("status", ServiceStatus.Active.getKey()) + ); + List tcrActiveCode = activeTcr.stream() + .map(tcr -> tcr.getCode()) + .filter(Objects::nonNull) + .collect(Collectors.toList()); + log.debug("Found {} code's from active TradingClearingRegistry", tcrActiveCode.size()); + ImdgPredicate predicate = builder.and( + AM_T_AM_F_predicate, + builder.equals("accountType", AccountType.Info.getKey()), + builder.in("tradingClearingRegistry", tcrActiveCode.toArray(new String[0])) + ); + checkRegistries = registryImdg.getCollectionObjectsByPredicate(predicate); + timer = System.currentTimeMillis() - timer; + log.debug("Select {} Registry by query: {}. Time {} ms.", checkRegistries.size(), predicate, timer); + } else { + log.debug("Selecting {} AM_T and by tradingClearingRegistry and securitySymbol AM_F for his.", specialId); + long timer = System.currentTimeMillis(); + checkRegistries = new ArrayList<>(specialId.size() * 2); + for (Long am_tId : specialId) { + if (am_tId == null) { // never + log.warn("Received argument ID was null"); + continue; + } + Registry am_tRegistry = registryImdg.getSingleObjectByID(am_tId); + if (am_tRegistry == null) { + log.warn("Registry not found by ID={}", am_tId); + } else { + checkRegistries.add(am_tRegistry); + ImdgPredicate predicate = builder.and( + AM_T_AM_F_predicate, + builder.equals("accountType", AccountType.Info.getKey()), + builder.equals("tradingClearingRegistry", am_tRegistry.getTradingClearingRegistry()), + builder.equals("securitySymbol", am_tRegistry.getSecuritySymbol()), + builder.not(builder.equals("id", am_tRegistry.getId())) // counter-registry + ); + Collection otherRegistry = registryImdg.getCollectionObjectsByPredicate(predicate); + log.trace("For Registry.id={} found {} other Registry's by query {}", am_tId, otherRegistry.size(), predicate); + for (Registry otherReg : otherRegistry) { + if (!am_tId.equals(otherReg.getId())) + checkRegistries.add(am_tRegistry); + } + } + } + timer = System.currentTimeMillis() - timer; + log.debug("Select {} Registry and hi's pairs by ID collection. Time {} ms.", checkRegistries.size(), timer); + } + + // Сверка Registry + List mismatch = checkRegisters(checkRegistries); + + if (mismatch.isEmpty()) { + log.debug("Check registry result: OK"); + sendMessage(makeSuccessMessage(checkTime), Priority.LOW); + } else { + log.debug("Check registry result: not match"); + sendMessage(makeErrorMessage(checkTime, mismatch), Priority.HIGH); + } + } + + String makeSuccessMessage(LocalDateTime checkTime) { + + return String.format(""" + Проведена сверка по регистрам AMAT/AMAF. + %s. + Результат сверки: расхождения отсутствуют.""", + MESSAGE_TIME_FORMATTER.format(checkTime)); + } + + String makeErrorMessage(LocalDateTime checkTime, List mismatch) { + assert !mismatch.isEmpty(); + String msgItems = mismatch.stream().collect(Collectors.joining("; ")); + return String.format(""" + Проведена сверка по регистрам AMAT/AMAF. + %s. + Результат сверки: %s.""", + MESSAGE_TIME_FORMATTER.format(checkTime), msgItems); + } + + protected List checkRegisters(Collection checkRegistries) { + List resultMismatch = new ArrayList<>(); + Map, List> regGroup = checkRegistries.stream().collect(Collectors.groupingBy( + r -> new Pair<>(r.getTradingClearingRegistry(), r.getSecuritySymbol()) + )); + log.trace("Make {} group of TradingClearingRegistry+SecuritySymbol from {} Registry", regGroup.size(), checkRegistries.size()); + int equalsCount = 0; + int mismatchCount = 0; + for (Map.Entry, List> group : regGroup.entrySet()) { + BigDecimal am_tSumm = BigDecimal.ZERO; + BigDecimal am_fSumm = BigDecimal.ZERO; + for (Registry r : group.getValue()) { + if (match(RegistryTradingParams.AM_T, r)) { + am_tSumm = am_tSumm.add(BigDecimalUtil.safeBD(r.getBalance())); + } + if (match(RegistryTradingParams.AM_F, r)) { + am_fSumm = am_fSumm.add(BigDecimalUtil.safeBD(r.getBalance())); + } + } + log.trace("Group {} size of {} item: AM_T summ={}, AM_F summ={}", + group.getKey(), group.getValue().size(), am_tSumm, am_fSumm); + boolean balanceOk = am_tSumm.compareTo(am_fSumm) == 0; + if (balanceOk) { + equalsCount++; + } else { + mismatchCount++; + String describe = String.format("по ТКР %s выявлено расхождение на регистрах AMAT/AMAF по %s на сумму %s", + group.getKey().getFirst(), // TradingClearingRegistry + group.getKey().getSecond(), // SecuritySymbol + am_tSumm.subtract(am_fSumm) + ); + resultMismatch.add(describe); + } + } + log.debug("Check result: equalsCount = {}, mismatchCount = {}", equalsCount, mismatchCount); + return resultMismatch; + } + + protected void sendMessage(String message, Priority priority) { + final String destination = Consts.NOTIFICATION_NEW; + NotificationNewRequest request = new NotificationNewRequest(); + request.setObjectType(ObjectType.rgst.getKey()); + request.setPriority(priority.getKey()); + request.setComment(message); + log.debug("Send message to kafka \"{}\": {}", destination, LogFormatter.toStringWrapper(request)); + kafkaSender.sendRequestToQueue(destination, request); + } + + protected boolean match(RegistryTradingParams params, Registry reg) { + if (params.registryDesignation() != null && !params.registryDesignation().equalsByKey(reg.getRegistryDesignation())) + return false; + if (params.registryCapacity() != null && !params.registryCapacity().equalsByKey(reg.getRegistryCapacity())) + return false; + if (params.registryInstrumentType() != null && !params.registryInstrumentType().equalsByKey(reg.getRegistryInstrumentType())) + return false; + if (params.registryUnit() != null && !params.registryUnit().equalsByKey(reg.getRegistryUnit())) + return false; + return true; + } +} diff --git a/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/AmatAmafCheckerTest.java b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/AmatAmafCheckerTest.java new file mode 100644 index 000000000..c41feb424 --- /dev/null +++ b/clearing-parent/utility-service/src/test/java/ru/spcex/clearing/utility/service/AmatAmafCheckerTest.java @@ -0,0 +1,147 @@ +package ru.spcex.clearing.utility.service; + +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.producer.Producer; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; +import ru.clearing.classes.statics.data.registry.Registry; +import ru.spcex.clearing.platform.messaging.service.sender.KafkaSender; +import ru.spcex.platform.enumeration.RegistryDesignation; +import ru.spcex.platform.enumeration.RegistryInstrumentType; +import ru.spcex.platform.enumeration.RegistryTradingParams; +import ru.spcex.platform.enumeration.RegistryUnit; +import ru.spcex.platform.imdg.api.ImdgProvider; + +import java.math.BigDecimal; +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.*; + +class AmatAmafCheckerTest { + + AmatAmafChecker instance() { + ImdgProvider imdgProvider = Mockito.mock(ImdgProvider.class); + KafkaSender kafkaSender = Mockito.mock(KafkaSender.class); + Producer kafkaProducer = Mockito.mock(Producer.class); + Consumer kafkaQueue = Mockito.mock(Consumer.class); + return new AmatAmafChecker(imdgProvider, kafkaSender, kafkaProducer, kafkaQueue); + } + + @Test + void checkOnDate() { + AmatAmafChecker checker = instance(); + + List regSrc = new ArrayList<>(); + List tkr = Arrays.asList("0004CAV00003", "0482CAT00001", "0095CAV00002"); + List securitySymbol = Arrays.asList("RUB", "CNY"); + int i = 1; + for (String tcrCode : tkr) + for (String secSymbol : securitySymbol) { + Registry reg1 = new Registry(); + reg1.setId(1L); + reg1.setRegistryDesignation("A"); + reg1.setRegistryInstrumentType("M"); + reg1.setRegistryCapacity("A"); + reg1.setRegistryUnit("T"); + reg1.setRegistryCode("AMAT"); // RegistryUtil.clearingCode(reg1) + reg1.setTradingClearingRegistry(tcrCode); + reg1.setSecuritySymbol(secSymbol); + reg1.setBalance(BigDecimal.valueOf(100 * i)); + if ("RUB".equals(secSymbol)) + reg1.setBalanceRub(reg1.getBalance()); + else + reg1.setBalanceRub(reg1.getBalance().multiply(BigDecimal.valueOf(0.1))); + regSrc.add(reg1); + Registry reg2 = new Registry(); + reg2.setId(1L); + reg2.setRegistryDesignation("A"); + reg2.setRegistryInstrumentType("M"); + reg2.setRegistryCapacity("A"); + reg2.setRegistryUnit("F"); + reg2.setRegistryCode("AMAF"); + reg2.setTradingClearingRegistry(tcrCode); + reg2.setSecuritySymbol(secSymbol); + reg2.setBalance(reg1.getBalance()); + reg2.setBalanceRub(reg1.getBalanceRub()); + regSrc.add(reg2); + } + + // All OK + List errors = checker.checkRegisters(regSrc); + assertEquals(0, errors.size()); + + + // Error in data + regSrc.get(0).setBalance(regSrc.get(0).getBalance().add(BigDecimal.valueOf(100))); + regSrc.get(4).setBalance(BigDecimal.valueOf(10.25)); + errors = checker.checkRegisters(regSrc); + assertEquals(2, errors.size()); + assertEquals("[по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100, по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75]", + errors.toString()); + + // next test makeErrorMessage: + assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" + + "2025.05.05 12:37:56.\n" + + "Результат сверки: по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100; по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75.", + checker.makeErrorMessage( + LocalDateTime.of(2025, 5, 5, 12, 37, 56), errors)); + } + + @Test + void makeSuccessMessage() { + AmatAmafChecker checker = instance(); + assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" + + "2025.05.05 12:37:56.\n" + + "Результат сверки: расхождения отсутствуют.", + checker.makeSuccessMessage(LocalDateTime.of(2025, 5, 5, 12, 37, 56))); + } + + @Test + void makeErrorMessage() { + AmatAmafChecker checker = instance(); + + assertEquals("Проведена сверка по регистрам AMAT/AMAF.\n" + + "2025.05.05 12:37:56.\n" + + "Результат сверки: по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100; по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75.", + checker.makeErrorMessage( + LocalDateTime.of(2025, 5, 5, 12, 37, 56), + Arrays.asList( + "по ТКР 0004CAV00003 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму 100", + "по ТКР 0482CAT00001 выявлено расхождение на регистрах AMAT/AMAF по RUB на сумму -89.75" + ) + )); + } + + + @Test + void match() { + AmatAmafChecker checker = instance(); + + Registry reg1 = new Registry(); + reg1.setId(1L); + reg1.setRegistryDesignation("A"); + reg1.setRegistryInstrumentType("M"); + reg1.setRegistryCapacity("A"); + reg1.setRegistryUnit("T"); + reg1.setRegistryCode("AMAT"); // RegistryUtil.clearingCode(reg1) + Registry reg2 = new Registry(); + reg2.setId(1L); + reg2.setRegistryDesignation("A"); + reg2.setRegistryInstrumentType("M"); + reg2.setRegistryCapacity("A"); + reg2.setRegistryUnit("F"); + reg2.setRegistryCode("AMAF"); + + assertTrue(checker.match(RegistryTradingParams.AM_T, reg1)); + assertTrue(checker.match(RegistryTradingParams.AM_F, reg2)); + assertTrue(checker.match(RegistryTradingParams.AM__, reg1)); + assertTrue(checker.match(RegistryTradingParams.AM__, reg2)); + + assertFalse(checker.match(RegistryTradingParams.AS__, reg1)); + assertFalse(checker.match(RegistryTradingParams.AM_F, reg1)); + assertFalse(checker.match(RegistryTradingParams.AM_T, reg2)); + } +} \ No newline at end of file diff --git a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/LauncherCommandRequest.java b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/LauncherCommandRequest.java index be464313c..b68f3fd13 100644 --- a/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/LauncherCommandRequest.java +++ b/platform-parent/platform-messaging/src/main/java/ru/spcex/clearing/platform/messaging/domain/cud/schedule/LauncherCommandRequest.java @@ -1,6 +1,5 @@ package ru.spcex.clearing.platform.messaging.domain.cud.schedule; -import com.fasterxml.jackson.annotation.JsonFormat; import com.fasterxml.jackson.annotation.JsonProperty; import com.fasterxml.jackson.databind.annotation.JsonDeserialize; import com.fasterxml.jackson.databind.annotation.JsonSerialize; @@ -12,6 +11,7 @@ import ru.spcex.clearing.platform.messaging.domain.json.serialize.LocalTimeSeria import java.math.BigDecimal; import java.time.LocalDate; import java.time.LocalTime; +import java.util.List; public class LauncherCommandRequest { @@ -66,6 +66,9 @@ public class LauncherCommandRequest { @JsonProperty public LocalDate toDate; + @JsonProperty + private List ids; + public Long getUserId() { return userId; } @@ -233,4 +236,12 @@ public class LauncherCommandRequest { public void setToDate(LocalDate toDate) { this.toDate = toDate; } + + public List getIds() { + return ids; + } + + public void setIds(List ids) { + this.ids = ids; + } } diff --git a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/Pair.java b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/Pair.java index 83b3ca0d6..a558c51b3 100644 --- a/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/Pair.java +++ b/platform-parent/platform-utils/src/main/java/ru/spcex/platform/utils/collection/Pair.java @@ -1,5 +1,6 @@ package ru.spcex.platform.utils.collection; +import java.util.Objects; import java.util.function.Consumer; public class Pair { @@ -40,4 +41,22 @@ public class Pair { secondConsumer.accept(second); } } + + @Override + public boolean equals(Object o) { + if (this == o) return true; + if (o == null || getClass() != o.getClass()) return false; + Pair pair = (Pair) o; + return Objects.equals(first, pair.first) && Objects.equals(second, pair.second); + } + + @Override + public int hashCode() { + return (first == null ? 0 : first.hashCode()) * 31 ^ (second == null ? 0 : second.hashCode()); + } + + @Override + public String toString() { + return "Pair{" + first + "; " + second + "}"; + } }